diff --git a/.commitrail/INDEX.md b/.commitrail/INDEX.md index cc542838e..dd48bd337 100644 --- a/.commitrail/INDEX.md +++ b/.commitrail/INDEX.md @@ -8,7 +8,7 @@ for current product capability. |---|---|---| | [WS-DB-002](initiatives/WS-DB-002/OVERVIEW.md) | Complete | Shared UUIDv7 record generation, native-UUID relationships and fresh v0.1 baseline; natural-owner retry custody and aligned CI/local setup | | [WS-MCP-002](initiatives/WS-MCP-002/OVERVIEW.md) | Planned | Nine tools through WS-MCP-002-03: self-service and administrative reads; 18 tools remain and WS-MCP-002-04 administrative grant mutations are next | -| [WS-CLI-001](initiatives/WS-CLI-001/OVERVIEW.md) | Planned | Public Go CLI self-service, project inspection, manager browsing and contributor discovery delivered through WS-CLI-001-05; further public journeys and binary distribution remain | +| [WS-CLI-001](initiatives/WS-CLI-001/OVERVIEW.md) | Planned | Public Go CLI self-service, project inspection, manager browsing, contributor discovery and claim/start delivered through WS-CLI-001-06; further public journeys and binary distribution remain | | [WS-ARCH-001](initiatives/WS-ARCH-001/OVERVIEW.md) | Planned | Source storage, inert AUTH contracts, TASK request reservation, hidden exact AUTH preparation and the REV-04C hidden FinalAcceptance/TASK/CON participant are delivered; TASK-before-CHECKERS reservation/current-read custody and ordered admission INSERTs are delivered by 04E1B-B1; 04E1B-B2 adds exact source preparation without publication; remaining hidden routing handlers are next. True admission remains independent of shared acceptance; false activation still requires exact AUTH receipt custody, database complete-set enforcement, currentness race proof and atomic routing activation. | | [WS-ART-001](initiatives/WS-ART-001/OVERVIEW.md) | Planned | Exact checker input/output custody and packet foundations are delivered; routing integration, remediation and public intake remain. | | [WS-AUTH-001](initiatives/WS-AUTH-001/OVERVIEW.md) | Planned | AUTH-19A commitments, TASK request reservation, ARCH-04E2-A hidden strict PREP matching and the REV-04C hidden FinalAcceptance/TASK/CON participant are delivered; the action remains unavailable, with mandatory exact AUTH receipt input, database closure, audit/outbox and scoped activation still required for the automated path. | diff --git a/.commitrail/initiatives/WS-CLI-001/OVERVIEW.md b/.commitrail/initiatives/WS-CLI-001/OVERVIEW.md index 7dfbdade5..722379af3 100644 --- a/.commitrail/initiatives/WS-CLI-001/OVERVIEW.md +++ b/.commitrail/initiatives/WS-CLI-001/OVERVIEW.md @@ -9,7 +9,8 @@ self-profile editing through the public REST API; [WS-CLI-001-03](WS-CLI-001-03.md), exact-project inspection with server-owned full/minimal disclosure; [WS-CLI-001-04](WS-CLI-001-04.md), public manager task pagination and detail; - [WS-CLI-001-05](WS-CLI-001-05.md), contributor ready-task discovery and instructions. + [WS-CLI-001-05](WS-CLI-001-05.md), contributor ready-task discovery and instructions; + [WS-CLI-001-06](WS-CLI-001-06.md), public contributor claim/start with explicit retry keys. ## Current boundary @@ -27,6 +28,9 @@ reads only the existing public project's server-selected disclosure shape. manager queue/detail journey with server-owned authority on each page. `task ready PROJECT_ID` and `task show TASK_ID` add the contributor projection with exact Submitter authority and server-owned assignment visibility. +`task claim TASK_ID --idempotency-key UUID` and `task start` add the public +contributor writes, with server-owned fresh authority, assignment and lineage. +Exact caller keys support manual replay; uncertainty never triggers automatic retries. All have text/JSON output and built-binary integration proof. Mutations preserve omitted/null semantics and explicitly report uncertain outcomes without automatic retries. @@ -71,11 +75,14 @@ CLIs. Keep the package independent of backend and MCP runtime dependencies. detail, passing opaque continuation unchanged without automatic pagination. 5. **WS-CLI-001-05:** Discover one ready-task page and inspect contributor instructions through public reads, without management metadata or task writes. -6. **Later governed-work commands:** Add project setup, contributor task mutations, submission, +6. **WS-CLI-001-06:** Claim ready work and start the caller's assignment through + public POSTs with caller-supplied keys, exact response validation and explicit + uncertain outcomes. No local authority decisions or operator override. +7. **Later governed-work commands:** Add project setup, submission, review, revision, and contribution reads/writes only as their actual public contracts and authority boundaries become available. Split by user journey, not one PR per endpoint or one giant catalogue PR. -7. **Optional TUI:** Add a focused public queue/evidence view after its API +8. **Optional TUI:** Add a focused public queue/evidence view after its API workflow is complete. Never require a TUI for agents or scripts. ## Risks and proof diff --git a/.commitrail/initiatives/WS-CLI-001/WS-CLI-001-06.md b/.commitrail/initiatives/WS-CLI-001/WS-CLI-001-06.md new file mode 100644 index 000000000..bad8a8ae5 --- /dev/null +++ b/.commitrail/initiatives/WS-CLI-001/WS-CLI-001-06.md @@ -0,0 +1,139 @@ +# WS-CLI-001-06 — Claim and start contributor work + +- Initiative: `WS-CLI-001` +- Durable disposition: `Complete` +- Intended merge outcome: Claim a ready task and start the caller's assignment + through the existing public REST operations. + +## Intent + +Contributor discovery is delivered. It does not reserve work. Add the next +bounded journey without implementing authority or lifecycle in the CLI. +`tasks/router.py` exposes claim and start with a required UUID `Idempotency-Key` +and optional reason of at most 1000 characters. `AuthorizedTaskCommands` owns +fresh AUTH, assignment, locked policy lineage, transaction and transition audit. +`TaskCommandReplay` scopes receipts by actor/action/key and verifies task/reason, +current authority and state. A receipt is not permission and replay may fail +after state changes; neither CLI retries nor local success caching are justified. + +## Bounded change + +Allowed: cohesive `cli/internal/` client/command code, process/API integration +tests, CLI/root README, affected roadmap claims, this record, overview and index. +Prohibited: backend/product, MCP, workflows, dependencies, operator overrides, +hidden APIs, credential persistence, local role/policy decisions, preflight, +automatic retries, automatic key generation, TUI or distribution work. + +`workstream task claim TASK_ID --idempotency-key UUID [--reason TEXT]` sends one +POST to `/api/v1/tasks/{task_id}/claim`; `task start` sends one POST to the public +start route. Require a caller-supplied key so a lost response does not lose the +retry identity. Validate selectors/key and bounded UTF-8 reason locally, preserve +supported UUID spelling on the wire and compare returned identity by value. +Send `{}` when reason is absent and preserve explicit empty reason. No actor or +project override and no arbitrary endpoint. Reuse the existing safe transport. +Mutation bodies must not be rewindable by the HTTP transport: a caller's +idempotency key does not permit silent HTTP/2 GOAWAY/REFUSED_STREAM replay. +After a sent body loses its response, return an unknown outcome for manual retry. + +Validate the exact contributor mutation TaskResponse and claim assignment +envelope, including UUID/timestamp/array/null shapes, requested task identity, +assignment task/project/contribution-policy identity, contributor equal to +assigned-by, active assignment, non-null assigned/accepted times, no released +time, required claimed/in-progress result, and absent management-only fields. +Preserve raw successful JSON; +text prints all returned public fields safely. Do not infer the caller's actor +ID from a JWT or add a whoami request. Workstream owns caller binding. + +Transport/send/read failures, redirects, unexpected success statuses, +malformed/oversized success, 5xx and noncanonical/intermediary 4xx produce +`outcome_unknown: true`. Only complete strictly decoded canonical API 4xx +envelopes establish known denial. Extend the shared transport deliberately for +these POST writes; retain the profile payload and error contract while honoring +its existing one-request/no-retry promise. Task-specific text directs +inspection with `task show TASK_ID` and +manual replay of the unchanged action/task/reason/key, not `whoami` or a new +key. Inspection is observational, not proof of rollback or global ordering. +No false rollback, automatic retry, +or guarantee that a later state-dependent replay will succeed. + +## Acceptance criteria + +- Built-process tests in the new `test_contributor_task_mutations.py` prove both + fixed routes, unchanged bearer/key/body, compact UUID parity, no preflight or + retry, local invalid input without network, complete safe text/raw JSON, + malformed/substituted/management replies, and unknown-outcome handling. +- Extend `contributor_task_journey.py` inside the existing one-bootstrap real + FastAPI/PostgreSQL test. Public create/screen/release and canonical approved + guide prerequisites remain; inference/storage fixtures do not certify S3. + CLI claim/start must be OpenAPI-visible. Preserve every existing read control. +- Prove persisted own claim/start and same-key response parity; changed reason + with the same key conflicts; changed task with that key cannot succeed. Active + same-project peer authority cannot start another assignment. Absent, Reviewer- + only, foreign-project, revoked and suspended authority deny writes/replays. + Positive controls reach each relevant guard; restore authority before testing + suspension. Inspect the task/assignment through public reads to prove denied + attempts do not mutate them. For revoke/suspend first observe canonical + post-administration state through public management reads, then baseline + denied attempts against that state, not the previous claimed/in-progress + state. Invalidation is asynchronous: ready/authority_revoked is its completed + effect, not an immediate API guarantee. This fixture has no background + dispatcher; do not add private processing to simulate that completion. + Claim replay after start need not succeed. +- Run Go verify/tidy/vet/build, Ruff, complete process/API suite via isolated + PostgreSQL, existing workflow guards, links, stale scans and Commitrail checks. + Full hosted suite remains blocking; no test/coverage quota or deadline change. + +## Evidence + +Verification commands; exact execution and reviewer/hosted results belong in the +PR. Process tests prove wire, output, strict response and uncertainty boundaries. +The real public API journey reuses the existing bootstrap and approved-guide +prerequisites; it does not certify deployed Flow or real guide inference/storage. + +```sh +cd cli +go mod verify +go mod tidy -diff +go vet ./... +go build -trimpath -o /tmp/workstream-cli ./cmd/workstream +cd ../backend +WORKSTREAM_CLI_EXECUTABLE=/tmp/workstream-cli python scripts/run_isolated_tests.py --metadata-json /tmp/workstream-cli-isolation.json --timeout-seconds 240 -- python -m pytest -q ../cli/tests/integration +``` + +## Risk and review routing + +Risk `L1`: credential transport, caller-isolated mutation and misleading retry +success. Plan review precedes implementation. Focused tracks: security; +architecture/documentation; QA/test delta in three bounded assignments. The +lead runs unchanged CI guards once. No broad product rewrite or nine-reviewer +fanout. Human focus: explicit retry key, fresh server authority, assignment and +lineage identity, uncertain outcome, source availability versus distribution. +The typed mutation projection and process/real-API proof exceed the default +500-line guideline, but remain one assignment journey and one unchanged public +owner boundary. Splitting claim from start would duplicate setup and obscure +their replay/state relationship. + +Reject automatic keys/retries because they obscure lost-response recovery; +reject manager/operator routes or local authorization because they change the +contract. No new human design decision is required. Human approval/merge remain +the final GitHub authority. Submission and further public journeys remain later. + +Plan corrections `PLAN-SEC-CLI06-001/002` specify deterministic uncertainty and +exact composite claim identity. `PLAN-FIXTURE-CLI06-001` distinguishes legitimate +administrative assignment invalidation from denied-command effects. The +implementation proof must independently exercise these boundaries. +`TEST-CLI-003` adds exact-1000-rune success with the unchanged POST body and +invalid-UTF-8 POSIX argument rejection with no request, distinguishing both +off-by-one and omitted-validation defects in the existing process cases. + +External findings `EXT-CLI06-HTTP2-REPLAY` and `EXT-CLI06-NULL-DETAILS` require +real TLS/HTTP2 process proof that graceful GOAWAY after body receipt sends only +one request, plus an otherwise complete canonical 4xx with null `details` that +must remain an unknown outcome. The former removes transport rewind capability; +the latter strengthens proof of the existing strict envelope guard. Neither +changes public API scope, backend retry authority or workflow settings. +Tracing that same shared transport also reproduced `EXT-CLI06-PATCH-REPLAY`: +the existing profile PATCH body was rewindable too. Apply the one-shot body rule +to both existing mutation methods, with a dedicated profile process regression +using the same real TLS/HTTP2 peer. Preserve profile omission/null and error +classification, and leave read transport behavior unchanged. diff --git a/README.md b/README.md index 03e6e6220..f63702a40 100644 --- a/README.md +++ b/README.md @@ -317,6 +317,9 @@ the caller's current authority, preserving contributor-minimal disclosure. provide a paginated manager list-to-detail journey through public reads. `workstream task ready PROJECT_ID` and `task show TASK_ID` provide contributor discovery and instructions, with live Submitter authority and no task claim. +`workstream task claim TASK_ID --idempotency-key UUID` and `task start` add +contributor writes through existing public APIs, with server-owned authority, +explicit caller retry keys and uncertain-outcome reporting without automatic retries. All support human-readable and JSON output, using the caller's Flow token. The first source package is buildable; further workflow commands and published binaries remain planned. diff --git a/cli/README.md b/cli/README.md index 2b095b04f..5069b3f4f 100644 --- a/cli/README.md +++ b/cli/README.md @@ -3,7 +3,7 @@ An independent Go client for Workstream's public REST API, for humans and agents using the terminal. It provides human self-profile reads and editing, plus exact-project inspection, authority reads, manager task browsing and -contributor work discovery: +contributor work discovery, claim and start: | Command | Public API | |---|---| @@ -15,6 +15,8 @@ contributor work discovery: | `workstream project task PROJECT_ID TASK_ID` | `GET /api/v1/projects/PROJECT_ID/tasks/TASK_ID` | | `workstream task ready PROJECT_ID` | `GET /api/v1/projects/PROJECT_ID/tasks/ready` | | `workstream task show TASK_ID` | `GET /api/v1/tasks/TASK_ID` | +| `workstream task claim TASK_ID --idempotency-key UUID` | `POST /api/v1/tasks/TASK_ID/claim` | +| `workstream task start TASK_ID --idempotency-key UUID` | `POST /api/v1/tasks/TASK_ID/start` | Workstream verifies the caller's Flow bearer and owns identity resolution, authorization and lifecycle decisions. Reading a profile can admit a first-time @@ -151,9 +153,47 @@ than silently displayed or ignored. Discovery is live, not a reservation or a claimability guarantee. A later claim must revalidate current authority and state. Cursors are action/project/limit bound and cannot be reused as manager cursors; the CLI never decodes them or -automatically fetches another page. These commands do not claim/start tasks, +automatically fetches another page. These reads do not claim/start tasks, upload submissions or complete unfinished acceptance integration. +## Claim and start contributor work + +```sh +workstream task claim TASK_ID --idempotency-key CLAIM_UUID --reason 'Begin this work' --output json +workstream task start TASK_ID --idempotency-key START_UUID --output json +``` + +Supply your own UUID key and retain it with the action, task and optional reason. +These commands send exactly one public POST, with no preflight, automatic key, +retry or operator override. Reason is optional (at most 1000 UTF-8 characters); +an omitted flag sends `{}`, while an explicit empty flag sends an empty string. +The caller's Flow bearer and key are forwarded unchanged. Only Workstream +decides whether current identity, lifecycle, exact Submitter grant, task state, +assignment ownership and locked policy permit the write. + +Claim returns the contributor-safe task and its assignment; start returns the +contributor-safe task in progress. The CLI validates the requested task identity, +assignment/task/project/policy consistency, claim contributor/assigner identity, +active/unreleased assignment and required timestamps before output. It rejects +management-only fields. JSON preserves the API object; text prints every public +field with escaped terminal controls. This does not expose submission or review +commands, or activate unfinished product lifecycle work. + +The API scopes keys by actor and action, and checks current authority before +recovering a committed result. An exact retry can recover the same result only +while the required state remains current. Changed reason/task conflicts; +claim replay after start can be denied. A key never grants permission. Revocation +and suspension deny further writes/replays; assignment invalidation is a separate +asynchronous consequence, not a CLI effect or immediate API guarantee. + +If the response is lost, malformed, oversized, redirected, an unexpected success +status, a server error or a noncanonical/intermediary denial, the CLI exits +nonzero with `error.outcome_unknown: true`. Only a complete strictly decoded +canonical Workstream 4xx error envelope establishes a known denial. Inspect with +`workstream task show TASK_ID`; this observes current state, not rollback or +global ordering. If manually retrying, preserve the unchanged action, task, +reason and key. Do not invent a new key or assume recovery will still succeed. + ## Edit your profile ```sh @@ -213,6 +253,13 @@ same-project non-owner concealment, independent authorized cursor substitution, exact Submitter grants, Reviewer-only/foreign/revoked denial and suspension after restored positive authority. The process fixture verifies the separate contributor shapes, complete field output and refusal of management-only data. +Contributor mutation proof extends the same API/bootstrap journey with CLI +claim/start, persisted assignment and locked lineage, exact replay parity, +reason/task mismatch, same-project non-owner denial, foreign/Reviewer-only +authority, revocation and suspension. Denied writes are compared with the +publicly observed post-administration task baseline, not an assumed pre-revoke +state. Process tests prove the exact POST/key/body, strict claim/start response +identity, safe text, canonical errors and uncertain/no-retry behavior. Local Flow-compatible tokens are test fixtures, not deployed-provider proof. No coverage percentage or test-count target is used. diff --git a/cli/internal/api/client.go b/cli/internal/api/client.go index a670a0dcc..92596588a 100644 --- a/cli/internal/api/client.go +++ b/cli/internal/api/client.go @@ -35,13 +35,18 @@ type Failure struct { Status int `json:"status,omitempty"` CorrelationID string `json:"correlation_id,omitempty"` OutcomeUnknown bool `json:"outcome_unknown,omitempty"` + RecoveryHint string `json:"-"` } func (f *Failure) Error() string { if f.OutcomeUnknown { known := *f known.OutcomeUnknown = false - return known.Error() + "; update outcome unknown; read whoami before retrying" + hint := f.RecoveryHint + if hint == "" { + hint = "update outcome unknown; read whoami before retrying" + } + return known.Error() + "; " + hint } if f.CorrelationID != "" { return fmt.Sprintf("%s (HTTP %d; correlation %s)", f.Code, f.Status, f.CorrelationID) @@ -300,6 +305,11 @@ func (c *Client) safeMetadata(value string) bool { } func (c *Client) request(ctx context.Context, method, path, query string, body []byte) (json.RawMessage, error) { + return c.requestWithKey(ctx, method, path, query, body, "") +} + +func (c *Client) requestWithKey(ctx context.Context, method, path, query string, body []byte, key string) (json.RawMessage, error) { + mutation := method == http.MethodPatch || method == http.MethodPost target := c.origin + path if query != "" { target += "?" + query @@ -312,20 +322,29 @@ func (c *Client) request(ctx context.Context, method, path, query string, body [ if err != nil { return nil, &Failure{Code: "invalid_request"} } + if mutation { + // NewRequest makes bytes.Reader bodies rewindable. Disable transport + // replay after HTTP/2 GOAWAY/REFUSED_STREAM (or HTTP/1 connection + // failure); write retries belong to the caller, even with a retry key. + req.GetBody = nil + } req.Header.Set("Authorization", "Bearer "+c.token) req.Header.Set("Accept", "application/json") req.Header.Set("Accept-Encoding", "identity") + if key != "" { + req.Header.Set("Idempotency-Key", key) + } if body != nil { req.Header.Set("Content-Type", "application/json") } response, err := c.http.Do(req) if err != nil { - return nil, &Failure{Code: "service_unavailable", OutcomeUnknown: method == http.MethodPatch} + return nil, &Failure{Code: "service_unavailable", OutcomeUnknown: mutation} } defer response.Body.Close() responseBody, err := io.ReadAll(io.LimitReader(response.Body, maxResponseBytes+1)) if err != nil || len(responseBody) > maxResponseBytes { - return nil, &Failure{Code: "invalid_api_response", Status: response.StatusCode, OutcomeUnknown: method == http.MethodPatch} + return nil, &Failure{Code: "invalid_api_response", Status: response.StatusCode, OutcomeUnknown: mutation} } correlation := response.Header.Get("X-Correlation-ID") if !safeCorrelation.MatchString(correlation) || !c.safeMetadata(correlation) { @@ -350,14 +369,17 @@ func (c *Client) request(ctx context.Context, method, path, query string, body [ code = envelope.Error.Code } } + if method == http.MethodPost { + knownEnvelope = canonicalTaskError(responseBody, response.Header.Get("Content-Type")) + } // A complete API 4xx envelope is a known denial, independent of code // spelling/redaction. A gateway reply without it cannot prove rollback. return nil, &Failure{Code: code, Status: response.StatusCode, CorrelationID: correlation, - OutcomeUnknown: method == http.MethodPatch && (response.StatusCode < 400 || response.StatusCode >= 500 || !knownEnvelope)} + OutcomeUnknown: mutation && (response.StatusCode < 400 || response.StatusCode >= 500 || !knownEnvelope)} } mediaType, _, err := mime.ParseMediaType(response.Header.Get("Content-Type")) if err != nil || mediaType != "application/json" || !jsontext.Value(responseBody).IsValid() { - return nil, &Failure{Code: "invalid_api_response", Status: response.StatusCode, CorrelationID: correlation, OutcomeUnknown: method == http.MethodPatch} + return nil, &Failure{Code: "invalid_api_response", Status: response.StatusCode, CorrelationID: correlation, OutcomeUnknown: mutation} } return json.RawMessage(responseBody), nil } diff --git a/cli/internal/api/task_mutations.go b/cli/internal/api/task_mutations.go new file mode 100644 index 000000000..59acaf23b --- /dev/null +++ b/cli/internal/api/task_mutations.go @@ -0,0 +1,194 @@ +package api + +import ( + "context" + "encoding/json" + "errors" + "mime" + "net/http" + "net/url" + "regexp" + "unicode/utf8" +) + +// MutationTask is the contributor-safe TaskResponse of the two public writes, +// not the management detail or contributor read projection. +type MutationTask struct { + ID string `json:"id"` + ProjectID string `json:"project_id"` + LockedContributionPolicyVersionID string `json:"locked_contribution_policy_version_id"` + LockedGuideVersion string `json:"locked_guide_version"` + LockedReviewPolicyID string `json:"locked_review_policy_id"` + LockedReviewPolicyGeneration int `json:"locked_review_policy_generation"` + LockedReviewPolicyHash string `json:"locked_review_policy_hash"` + LockedRevisionPolicyID string `json:"locked_revision_policy_id"` + LockedRevisionPolicyGeneration int `json:"locked_revision_policy_generation"` + LockedRevisionPolicyHash string `json:"locked_revision_policy_hash"` + LockedPaymentPolicyVersion *string `json:"locked_payment_policy_version"` + SourceType string `json:"source_type"` + Title string `json:"title"` + Description string `json:"description"` + TaskType *string `json:"task_type"` + Difficulty *string `json:"difficulty"` + SkillTags []string `json:"skill_tags"` + EstimatedTimeMinutes *int `json:"estimated_time_minutes"` + BaseAmount *string `json:"base_amount"` + Currency *string `json:"currency"` + PayoutType *string `json:"payout_type"` + Status string `json:"status"` + AcceptanceCriteria *string `json:"acceptance_criteria"` + RejectionCriteria *string `json:"rejection_criteria"` + DeadlineAt *string `json:"deadline_at"` + CreatedAt string `json:"created_at"` + UpdatedAt string `json:"updated_at"` +} + +type TaskAssignment struct { + ID string `json:"id"` + TaskID string `json:"task_id"` + ProjectID string `json:"project_id"` + SubmitterContributionPolicyVersionID string `json:"submitter_contribution_policy_version_id"` + ContributorID string `json:"contributor_id"` + AssignedBy string `json:"assigned_by"` + AssignedAt string `json:"assigned_at"` + AcceptedAt string `json:"accepted_at"` + ReleasedAt *string `json:"released_at"` + Status string `json:"status"` +} + +type ClaimedTask struct { + Task MutationTask + Assignment TaskAssignment +} + +var policyHash = regexp.MustCompile(`^sha256:[0-9a-f]{64}$`) +var decimalAmount = regexp.MustCompile(`^[+-]?[0-9]+(?:\.[0-9]+)?(?:[eE][+-]?[0-9]+)?$`) + +func (c *Client) taskWrite(ctx context.Context, selector, action, key string, reason *string) (json.RawMessage, error) { + if _, ok := uuidIdentity(selector); !ok || len(selector) > 100 { + return nil, errors.New("TASK_ID must be a UUID") + } + if _, ok := uuidIdentity(key); !ok || len(key) > 100 { + return nil, errors.New("--idempotency-key must be a UUID") + } + fields := map[string]string{} + if reason != nil { + if !utf8.ValidString(*reason) || utf8.RuneCountInString(*reason) > 1000 { + return nil, errors.New("reason must contain at most 1000 valid UTF-8 characters") + } + fields["reason"] = *reason + } + body, err := json.Marshal(fields) + if err != nil || len(body) > maxUpdateBytes { + return nil, errors.New("task request exceeds the request size limit") + } + raw, err := c.requestWithKey(ctx, http.MethodPost, "/api/v1/tasks/"+url.PathEscape(selector)+"/"+action, "", body, key) + return raw, taskWriteFailure(err) +} + +func taskWriteFailure(err error) error { + var failure *Failure + if errors.As(err, &failure) && failure.OutcomeUnknown { + failure.RecoveryHint = "task outcome unknown; inspect task show (observation, not rollback proof); replay only the unchanged action, task, reason and idempotency key" + } + return err +} + +func (c *Client) ClaimTask(ctx context.Context, selector, key string, reason *string) (Result[ClaimedTask], error) { + var result Result[ClaimedTask] + raw, err := c.taskWrite(ctx, selector, "claim", key, reason) + if err != nil { + return result, err + } + var envelope struct { + Task json.RawMessage `json:"task"` + Assignment json.RawMessage `json:"assignment"` + } + var value ClaimedTask + err = decode(raw, &envelope, []string{"task", "assignment"}, nil) + if err == nil { + value.Task, err = decodeMutationTask(envelope.Task, selector, "claimed") + } + if err == nil { + err = decodeTaskFields(envelope.Assignment, &value.Assignment, + []string{"id", "task_id", "project_id", "submitter_contribution_policy_version_id", "contributor_id", "assigned_by", "assigned_at", "accepted_at", "status"}, nil) + } + a := value.Assignment + if err != nil || !sameUUID(a.TaskID, value.Task.ID) || !sameUUID(a.ProjectID, value.Task.ProjectID) || + !sameUUID(a.SubmitterContributionPolicyVersionID, value.Task.LockedContributionPolicyVersionID) || + !sameUUID(a.ContributorID, a.AssignedBy) || !validUUID(a.ID) || a.Status != "active" || + !validTime(a.AssignedAt) || !validTime(a.AcceptedAt) || a.ReleasedAt != nil { + return result, taskWriteFailure(&Failure{Code: "invalid_api_response", OutcomeUnknown: true}) + } + return Result[ClaimedTask]{Raw: raw, Value: value}, nil +} + +func (c *Client) StartTask(ctx context.Context, selector, key string, reason *string) (Result[MutationTask], error) { + var result Result[MutationTask] + raw, err := c.taskWrite(ctx, selector, "start", key, reason) + if err != nil { + return result, err + } + value, err := decodeMutationTask(raw, selector, "in_progress") + if err != nil { + return result, taskWriteFailure(&Failure{Code: "invalid_api_response", OutcomeUnknown: true}) + } + return Result[MutationTask]{Raw: raw, Value: value}, nil +} + +func decodeMutationTask(raw json.RawMessage, selector, status string) (MutationTask, error) { + var task MutationTask + err := decodeTaskFields(raw, &task, []string{ + "id", "project_id", "locked_contribution_policy_version_id", "locked_guide_version", + "locked_review_policy_id", "locked_review_policy_hash", "locked_revision_policy_id", + "locked_revision_policy_hash", "source_type", "title", "description", "status", "created_at", "updated_at", + }, []string{"skill_tags", "locked_review_policy_generation", "locked_revision_policy_generation"}) + if err != nil || !sameUUID(task.ID, selector) || !validUUID(task.ProjectID) || + !validUUID(task.LockedContributionPolicyVersionID) || !validUUID(task.LockedReviewPolicyID) || + !validUUID(task.LockedRevisionPolicyID) || task.LockedGuideVersion == "" || + task.LockedReviewPolicyGeneration < 1 || task.LockedRevisionPolicyGeneration < 1 || + !policyHash.MatchString(task.LockedReviewPolicyHash) || !policyHash.MatchString(task.LockedRevisionPolicyHash) || + task.Status != status || !validTime(task.CreatedAt) || !validTime(task.UpdatedAt) || + (task.DeadlineAt != nil && !validTime(*task.DeadlineAt)) || + (task.BaseAmount != nil && !decimalAmount.MatchString(*task.BaseAmount)) { + return task, &Failure{Code: "invalid_api_response"} + } + return task, nil +} + +func validUUID(value string) bool { + _, valid := uuidIdentity(value) + return valid +} + +func sameUUID(left, right string) bool { + a, validA := uuidIdentity(left) + b, validB := uuidIdentity(right) + return validA && validB && a == b +} + +// A task write needs a complete public ApiError, not a gateway's code-shaped +// JSON. Profile PATCH deliberately retains its existing response contract. +func canonicalTaskError(raw []byte, contentType string) bool { + mediaType, _, err := mime.ParseMediaType(contentType) + if err != nil || mediaType != "application/json" { + return false + } + var envelope struct { + Error json.RawMessage `json:"error"` + Detail json.RawMessage `json:"detail"` + } + if decode(raw, &envelope, []string{"error"}, nil) != nil { + return false + } + var value struct { + Code *string `json:"code"` + Message *string `json:"message"` + Details map[string]json.RawMessage `json:"details"` + CorrelationID *string `json:"correlation_id"` + Retryable *bool `json:"retryable"` + } + return decode(envelope.Error, &value, []string{"code", "message", "details", "correlation_id", "retryable"}, nil) == nil && + value.Code != nil && *value.Code != "" && value.Message != nil && value.Details != nil && + value.CorrelationID != nil && validUUID(*value.CorrelationID) && value.Retryable != nil +} diff --git a/cli/internal/command/command.go b/cli/internal/command/command.go index 5b98be02e..bd0b21775 100644 --- a/cli/internal/command/command.go +++ b/cli/internal/command/command.go @@ -170,7 +170,7 @@ func Run(args []string, stdout, stderr io.Writer, getenv environment, version st }, }) root.AddCommand(project) - addContributorTaskReads(root, client, &output, stdout) + addContributorTasks(root, client, &output, stdout) err := root.ExecuteContext(context.Background()) if err == nil { diff --git a/cli/internal/command/contributor_tasks.go b/cli/internal/command/contributor_tasks.go index d3ee026ae..10c58dc91 100644 --- a/cli/internal/command/contributor_tasks.go +++ b/cli/internal/command/contributor_tasks.go @@ -8,7 +8,7 @@ import ( "github.com/spf13/cobra" ) -func addContributorTaskReads(root *cobra.Command, client func() (*api.Client, error), output *string, stdout io.Writer) { +func addContributorTasks(root *cobra.Command, client func() (*api.Client, error), output *string, stdout io.Writer) { task := &cobra.Command{Use: "task", Short: "Discover and inspect contributor work"} var limit int var cursor string @@ -75,5 +75,6 @@ func addContributorTaskReads(root *cobra.Command, client func() (*api.Client, er return err }, }) + addContributorTaskWrites(task, client, output, stdout) root.AddCommand(task) } diff --git a/cli/internal/command/task_mutations.go b/cli/internal/command/task_mutations.go new file mode 100644 index 000000000..72df02eb9 --- /dev/null +++ b/cli/internal/command/task_mutations.go @@ -0,0 +1,68 @@ +package command + +import ( + "fmt" + "io" + + "github.com/Flow-Research/workstream/cli/internal/api" + "github.com/spf13/cobra" +) + +func addContributorTaskWrites(task *cobra.Command, client func() (*api.Client, error), output *string, stdout io.Writer) { + for _, action := range []string{"claim", "start"} { + var key, reason string + cmd := &cobra.Command{ + Use: action + " TASK_ID --idempotency-key UUID", Short: action + " a task under current Submitter authority", Args: cobra.ExactArgs(1), + RunE: func(cmd *cobra.Command, args []string) error { + var note *string + if cmd.Flags().Changed("reason") { + note = &reason + } + apiClient, err := client() + if err != nil { + return err + } + if action == "start" { + result, err := apiClient.StartTask(cmd.Context(), args[0], key, note) + if err != nil { + return err + } + if *output == "json" { + return writeJSON(stdout, result.Raw) + } + return writeMutationTask(stdout, result.Value) + } + result, err := apiClient.ClaimTask(cmd.Context(), args[0], key, note) + if err != nil { + return err + } + if *output == "json" { + return writeJSON(stdout, result.Raw) + } + if err := writeMutationTask(stdout, result.Value.Task); err != nil { + return err + } + a := result.Value.Assignment + _, err = fmt.Fprintf(stdout, "Assignment: %s\nAssignment task: %s\nAssignment project: %s\nSubmitter policy: %s\nContributor: %s\nAssigned by: %s\nAssigned at: %s\nAccepted at: %s\nReleased at: %s\nAssignment status: %s\n", + safeText(a.ID), safeText(a.TaskID), safeText(a.ProjectID), safeText(a.SubmitterContributionPolicyVersionID), safeText(a.ContributorID), safeText(a.AssignedBy), safeText(a.AssignedAt), safeText(a.AcceptedAt), optional(a.ReleasedAt), safeText(a.Status)) + return err + }, + } + cmd.Flags().StringVar(&key, "idempotency-key", "", "Caller-supplied UUID; retain it for an exact manual retry") + cmd.Flags().StringVar(&reason, "reason", "", "Optional transition reason (at most 1000 characters)") + task.AddCommand(cmd) + } +} + +func writeMutationTask(w io.Writer, t api.MutationTask) error { + if err := writeTaskSummary(w, api.TaskSummary{ + TaskID: t.ID, ProjectID: t.ProjectID, Title: t.Title, TaskType: t.TaskType, + Difficulty: t.Difficulty, SkillTags: t.SkillTags, EstimatedTimeMinutes: t.EstimatedTimeMinutes, + Status: t.Status, DeadlineAt: t.DeadlineAt, CreatedAt: t.CreatedAt, UpdatedAt: t.UpdatedAt, + }); err != nil { + return err + } + _, err := fmt.Fprintf(w, "Description: %s\nAcceptance criteria: %s\nRejection criteria: %s\nSource type: %s\nContribution policy: %s\nGuide version: %s\nReview policy: %s\nReview generation: %d\nReview hash: %s\nRevision policy: %s\nRevision generation: %d\nRevision hash: %s\nPayment policy: %s\nBase amount: %s\nCurrency: %s\nPayout type: %s\n", + safeText(t.Description), optional(t.AcceptanceCriteria), optional(t.RejectionCriteria), safeText(t.SourceType), safeText(t.LockedContributionPolicyVersionID), safeText(t.LockedGuideVersion), safeText(t.LockedReviewPolicyID), t.LockedReviewPolicyGeneration, safeText(t.LockedReviewPolicyHash), safeText(t.LockedRevisionPolicyID), t.LockedRevisionPolicyGeneration, safeText(t.LockedRevisionPolicyHash), optional(t.LockedPaymentPolicyVersion), optional(t.BaseAmount), optional(t.Currency), optional(t.PayoutType)) + return err +} diff --git a/cli/tests/integration/conftest.py b/cli/tests/integration/conftest.py index 22aedd8d4..8b144bd11 100644 --- a/cli/tests/integration/conftest.py +++ b/cli/tests/integration/conftest.py @@ -18,7 +18,10 @@ def cli(tmp_path: Path) -> Callable[..., subprocess.CompletedProcess[str]]: ) def invoke( - origin: str, token: str, *args: str, extra_env: dict[str, str] | None = None + origin: str, + token: str, + *args: str | bytes, + extra_env: dict[str, str] | None = None, ) -> subprocess.CompletedProcess[str]: env = { key: value diff --git a/cli/tests/integration/contributor_task_journey.py b/cli/tests/integration/contributor_task_journey.py index 3a230b681..69b0c56df 100644 --- a/cli/tests/integration/contributor_task_journey.py +++ b/cli/tests/integration/contributor_task_journey.py @@ -1,4 +1,4 @@ -"""Real public contributor reads over canonical approved-guide prerequisites. +"""Real public contributor discovery and writes over approved prerequisites. Only upstream guide inference/storage are scripted fixtures. Task, assignment, grant and lifecycle operations exercise the real API with no route overrides. @@ -37,6 +37,35 @@ def read(command, presented=token, status=None): assert result.returncode == 0 and result.stderr == "", result.stderr return json.loads(result.stdout) + async def mutate(action, task_id, key, reason=None, presented=token, status=None): + # Read the current public management projection AFTER administrator + # changes; async invalidation is not simulated in this fixture. + path = f"/api/v1/projects/{project}/tasks/{task_id}" + before = await direct.get(path, headers=manager) + assert before.status_code == 200, before.text + flags = () if reason is None else ("--reason", reason) + result = cli( + origin, + presented, + "task", + action, + task_id.replace("-", ""), + "--idempotency-key", + str(key), + *flags, + "-o", + "json", + ) + if status is not None: + assert result.returncode == 1 and result.stdout == "", result.stderr + error = json.loads(result.stderr)["error"] + assert error["status"] == status and "outcome_unknown" not in error + after = await direct.get(path, headers=manager) + assert after.status_code == 200 and after.json() == before.json() + return error + assert result.returncode == 0 and result.stderr == "", result.stderr + return json.loads(result.stdout) + async def grant(project, actor, role="submitter"): return await post( f"/api/v1/projects/{project}/role-grants", @@ -104,6 +133,10 @@ async def grant(project, actor, role="submitter"): ) actor_id = profiles["cli-outsider"]["actor_profile_id"] + # Ready persisted work does not permit writes without Submitter authority. + for action in ("claim", "start"): + await mutate(action, ready_ids[0], uuid4(), status=403) + # Stored ready resources exist, but no matching grant permits either read. read(("ready", project), status=404) read(("show", ready_ids[0]), status=404) @@ -112,6 +145,14 @@ async def grant(project, actor, role="submitter"): await grant(project, peer_profiles["cli-task-reviewer"], role="reviewer") for command in (("ready", project), ("show", ready_ids[0])): read(command, peer_tokens["cli-task-reviewer"], status=404) + for action in ("claim", "start"): + await mutate( + action, + ready_ids[0], + uuid4(), + presented=peer_tokens["cli-task-reviewer"], + status=403, + ) cursor = None observed = [] @@ -164,6 +205,33 @@ async def detail(task_id, presented=token): read(("show", draft["id"]), status=404) read(("ready", foreign_project), status=404) read(("show", foreign_task_id), status=404) + # The same caller has Submitter A, but not B. Use B's management baseline. + for action in ("claim", "start"): + before = await direct.get( + f"/api/v1/projects/{foreign_project}/tasks/{foreign_task_id}", + headers=manager, + ) + denied = cli( + origin, + token, + "task", + action, + foreign_task_id, + "--idempotency-key", + str(uuid4()), + "-o", + "json", + ) + assert denied.returncode == 1 and denied.stdout == "" + assert json.loads(denied.stderr)["error"]["status"] == 403 + after = await direct.get( + f"/api/v1/projects/{foreign_project}/tasks/{foreign_task_id}", + headers=manager, + ) + assert ( + before.status_code == after.status_code == 200 + and before.json() == after.json() + ) # Each cursor substitution retains authority before testing the codec. await grant(foreign_project, actor_id) read(("ready", foreign_project)) @@ -207,12 +275,25 @@ async def detail(task_id, presented=token): assert management.returncode == 1 and management.stdout == "" assert json.loads(management.stderr)["error"]["status"] == 422 - claimed = await post( - f"/api/v1/tasks/{ready_ids[0]}/claim", - {"reason": "Assignment visibility control"}, - {"Authorization": f"Bearer {token}"}, - ) + claim_key, start_key = uuid4(), uuid4() + reason = "Assignment visibility control" + claimed = await mutate("claim", ready_ids[0], claim_key, reason) assert claimed["assignment"]["contributor_id"] == actor_id + assert claimed["assignment"]["task_id"] == claimed["task"]["id"] == ready_ids[0] + assert ( + claimed["assignment"]["project_id"] + == claimed["task"]["project_id"] + == project + ) + assert ( + claimed["assignment"]["submitter_contribution_policy_version_id"] + == claimed["task"]["locked_contribution_policy_version_id"] + ) + assert await mutate("claim", ready_ids[0], claim_key, reason) == claimed + mismatch = await mutate("claim", ready_ids[0], claim_key, "Changed", status=409) + assert mismatch["code"] == "idempotency_mismatch" + mismatch = await mutate("claim", ready_ids[1], claim_key, reason, status=409) + assert mismatch["code"] == "idempotency_mismatch" assert (await detail(ready_ids[0]))["status"] == "claimed" # Same-project authority was demonstrated before and after this ownership change. read(("show", ready_ids[0]), peer_tokens["cli-task-peer"], status=404) @@ -221,16 +302,42 @@ async def detail(task_id, presented=token): page = read(("ready", project), presented) assert {item["task_id"] for item in page["items"]} == set(ready_ids[1:]) + await mutate( + "start", + ready_ids[0], + uuid4(), + presented=peer_tokens["cli-task-peer"], + status=403, + ) + started = await mutate("start", ready_ids[0], start_key, reason) + assert started["id"] == ready_ids[0] and started["status"] == "in_progress" + assert ( + started["locked_contribution_policy_version_id"] + == claimed["task"]["locked_contribution_policy_version_id"] + ) + assert (await detail(ready_ids[0]))["status"] == "in_progress" + assert await mutate("start", ready_ids[0], start_key, reason) == started + assert (await mutate("start", ready_ids[0], start_key, "Changed", status=409))[ + "code" + ] == "idempotency_mismatch" + await mutate("claim", ready_ids[0], claim_key, reason, status=403) + await post( f"/api/v1/projects/{project}/role-grants/{first_grant['id']}/revoke", {"reason": "Discovery must reauthorize"}, ) read(("ready", project, "--limit", "1", "--cursor", first_cursor), status=404) read(("show", ready_ids[1]), status=404) + await mutate("start", ready_ids[0], start_key, reason, status=403) + await mutate("claim", ready_ids[1], uuid4(), status=403) # Manager authority cannot substitute for a revoked Submitter grant. await grant(project, actor_id) read(("ready", project, "--limit", "1", "--cursor", first_cursor)) await detail(ready_ids[1]) + # Establish fresh successful write authority before lifecycle denial. + second_key = uuid4() + second_claim = await mutate("claim", ready_ids[1], second_key) + assert second_claim["assignment"]["contributor_id"] == actor_id # The same freshly authorized actor is suspended: stale revocation cannot # masquerade as lifecycle denial. Restore it for the outer journey. await post( @@ -240,6 +347,8 @@ async def detail(task_id, presented=token): ) read(("ready", project), status=404) read(("show", ready_ids[1]), status=404) + await mutate("claim", ready_ids[1], second_key, status=403) + await mutate("start", ready_ids[1], uuid4(), status=403) await post( f"/api/v1/actors/{actor_id}/reactivate", {"reason": "Continue independent CLI proof"}, diff --git a/cli/tests/integration/http2_fixture.py b/cli/tests/integration/http2_fixture.py new file mode 100644 index 000000000..b75c7418d --- /dev/null +++ b/cli/tests/integration/http2_fixture.py @@ -0,0 +1,152 @@ +"""A loopback TLS peer that receives a body, then sends graceful HTTP/2 GOAWAY. + +Only the frames needed for this transport regression are implemented. There is +no application handler or HTTP client replacement: the built CLI negotiates +TLS/HTTP2, writes its request, and decides whether to replay on a new connection. +""" + +from contextlib import contextmanager +from datetime import datetime, timedelta, timezone +from ipaddress import ip_address +import socket +import ssl +from threading import Event, Thread + +from cryptography import x509 +from cryptography.hazmat.primitives import hashes, serialization +from cryptography.hazmat.primitives.asymmetric import ec +from cryptography.x509.oid import NameOID + +PREFACE = b"PRI * HTTP/2.0\r\n\r\nSM\r\n\r\n" + + +def receive(connection, size): + value = bytearray() + while len(value) < size: + chunk = connection.recv(size - len(value)) + if not chunk: + raise EOFError + value.extend(chunk) + return bytes(value) + + +def send_frame(connection, kind, flags=0, stream=0, payload=b""): + connection.sendall( + len(payload).to_bytes(3, "big") + + bytes((kind, flags)) + + stream.to_bytes(4, "big") + + payload + ) + + +def receive_request(connection): + assert connection.selected_alpn_protocol() == "h2" + assert receive(connection, len(PREFACE)) == PREFACE + send_frame(connection, 4) # SETTINGS + stream_id = None + body = bytearray() + while True: + header = receive(connection, 9) + length = int.from_bytes(header[:3], "big") + kind, flags = header[3:5] + stream = int.from_bytes(header[5:], "big") & 0x7FFFFFFF + assert length <= 16384 + payload = receive(connection, length) + if kind == 4 and not flags & 1: + send_frame(connection, 4, 1) # SETTINGS ACK + elif kind == 1: # HEADERS: no HPACK decoding is needed to count requests. + assert stream_id is None and stream > 0 and flags & 4 + stream_id = stream + assert not flags & 1 # These writes must carry their nonempty JSON body. + elif kind == 0: + assert stream == stream_id and not flags & 8 # No padded DATA. + body.extend(payload) + assert len(body) <= 8192 + if flags & 1: # END_STREAM: the complete body reached the peer. + return bytes(body) + + +@contextmanager +def goaway_fixture(tmp_path): + """Observe every fresh TLS connection until the built CLI exits.""" + key = ec.generate_private_key(ec.SECP256R1()) + name = x509.Name([x509.NameAttribute(NameOID.COMMON_NAME, "CLI test loopback")]) + now = datetime.now(timezone.utc) + certificate = ( + x509.CertificateBuilder() + .subject_name(name) + .issuer_name(name) + .public_key(key.public_key()) + .serial_number(x509.random_serial_number()) + .not_valid_before(now - timedelta(minutes=1)) + .not_valid_after(now + timedelta(hours=1)) + .add_extension(x509.BasicConstraints(ca=True, path_length=0), critical=True) + .add_extension( + x509.SubjectAlternativeName([x509.IPAddress(ip_address("127.0.0.1"))]), + critical=False, + ) + .sign(key, hashes.SHA256()) + ) + cert_path, key_path = tmp_path / "peer.pem", tmp_path / "peer-key.pem" + cert_path.write_bytes(certificate.public_bytes(serialization.Encoding.PEM)) + key_path.write_bytes( + key.private_bytes( + serialization.Encoding.PEM, + serialization.PrivateFormat.PKCS8, + serialization.NoEncryption(), + ) + ) + key_path.chmod(0o600) + context = ssl.SSLContext(ssl.PROTOCOL_TLS_SERVER) + context.load_cert_chain(cert_path, key_path) + context.set_alpn_protocols(["h2"]) + stopped = Event() + bodies, connections, errors = [], [], [] + listener = socket.socket() + listener.bind(("127.0.0.1", 0)) + listener.listen() + listener.settimeout(0.1) + + def serve(): + while not stopped.is_set(): + try: + raw, _address = listener.accept() + except TimeoutError: + continue + try: + raw.settimeout(3) + with raw, context.wrap_socket(raw, server_side=True) as connection: + connections.append(connection.selected_alpn_protocol()) + bodies.append(receive_request(connection)) + # last_stream_id=0, NO_ERROR: none of the streams is reported + # processed. A rewindable request is retried by Go here. + if len(bodies) > 1: + # A replay is already observable. Close without another + # GOAWAY so the broken client cannot loop until timeout. + continue + send_frame(connection, 7, payload=b"\x00" * 8) + # Let the client process GOAWAY before closing the TLS peer. + try: + while connection.recv(4096): + pass + except (TimeoutError, ConnectionResetError): + pass + except Exception as error: # Surface fixture failures to the main test. + errors.append(error) + stopped.set() + + thread = Thread(target=serve, daemon=True) + thread.start() + try: + yield ( + f"https://127.0.0.1:{listener.getsockname()[1]}", + {"SSL_CERT_FILE": str(cert_path)}, + bodies, + connections, + ) + finally: + stopped.set() + thread.join(timeout=4) + listener.close() + assert not thread.is_alive(), "HTTP/2 fixture did not terminate" + assert not errors, errors diff --git a/cli/tests/integration/test_contributor_task_mutations.py b/cli/tests/integration/test_contributor_task_mutations.py new file mode 100644 index 000000000..ce6c2e093 --- /dev/null +++ b/cli/tests/integration/test_contributor_task_mutations.py @@ -0,0 +1,386 @@ +"""Built-process proof of public contributor writes, not local authority rules.""" + +from copy import deepcopy +import json +from urllib.parse import quote + +from http2_fixture import goaway_fixture +from test_http_boundary import ( + ACTOR, + PROFILE, + PROJECT, + TOKEN, + assert_failure, + http_fixture, +) +from test_task_http_boundary import TASK, SUMMARY_TEXT + +KEY = "018f0ebc-7966-7e8d-bc4d-1cae1e000004" +POLICY = "018f0ebc-7966-7e8d-bc4d-1cae1e000005" +ASSIGNMENT = "018f0ebc-7966-7e8d-bc4d-1cae1e000006" +TASK_RESPONSE = { + "id": TASK, + "project_id": PROJECT, + "locked_contribution_policy_version_id": POLICY, + "locked_guide_version": "guide é\n", + "locked_review_policy_id": POLICY, + "locked_review_policy_generation": 1, + "locked_review_policy_hash": "sha256:" + "a" * 64, + "locked_revision_policy_id": POLICY, + "locked_revision_policy_generation": 2, + "locked_revision_policy_hash": "sha256:" + "b" * 64, + "source_type": "manual", + "title": "Work é\n\x1b[31m", + "description": "Instructions\n\x1b[32m", + "skill_tags": ["analysis", "tag\n"], + "status": "claimed", + "created_at": PROFILE["created_at"], + "updated_at": "2026-10-02T00:00:00Z", +} +ASSIGNMENT_RESPONSE = { + "id": ASSIGNMENT, + "task_id": TASK, + "project_id": PROJECT, + "submitter_contribution_policy_version_id": POLICY, + "contributor_id": ACTOR, + "assigned_by": ACTOR, + "assigned_at": PROFILE["created_at"], + "accepted_at": PROFILE["created_at"], + "status": "active", +} + + +def body(action): + task = deepcopy(TASK_RESPONSE) + if action == "claim": + return {"task": task, "assignment": deepcopy(ASSIGNMENT_RESPONSE)} + task["status"] = "in_progress" + return task + + +def invoke(cli, origin, action, *flags, selector=TASK, key=KEY, output="json"): + return cli( + origin, + TOKEN, + "-o", + output, + "task", + action, + selector, + "--idempotency-key", + key, + *flags, + ) + + +def canonical_error(code="authorization_denied"): + return { + "error": { + "code": code, + "message": "Denied", + "details": {}, + "correlation_id": KEY, + "retryable": False, + } + } + + +def test_task_writes_preserve_body_key_selector_and_json(cli): + with http_fixture() as (origin, response, requests): + for action in ("claim", "start"): + response["body"] = json.dumps(body(action)).encode() + for selector in ( + TASK, + TASK.replace("-", ""), + "{" + TASK + "}", + "urn:uuid:" + TASK, + ): + for flags, expected in ( + ((), {}), + (("--reason", ""), {"reason": ""}), + (("--reason", "Note é\n"), {"reason": "Note é\n"}), + ): + result = invoke( + cli, + origin, + action, + *flags, + selector=selector, + key=KEY.replace("-", ""), + ) + assert result.returncode == 0 and result.stderr == "", result.stderr + assert result.stdout.strip().encode() == response["body"] + assert requests[-1] == ( + "POST", + f"/api/v1/tasks/{quote(selector, safe=':')}/{action}", + "Bearer " + TOKEN, + ) + assert response["commands"][-1] == ( + "application/json", + [KEY.replace("-", "")], + expected, + ) + boundary = "é" * 1000 + result = invoke(cli, origin, action, "--reason", boundary) + assert result.returncode == 0 and result.stderr == "", result.stderr + assert result.stdout.strip().encode() == response["body"] + assert response["commands"][-1] == ( + "application/json", + [KEY], + {"reason": boundary}, + ) + assert len(requests) == len(response["commands"]) == 26 + + +def test_task_write_text_is_complete_and_escaped(cli): + optional = { + "task_type": "evaluation", + "difficulty": "medium", + "estimated_time_minutes": 17, + "deadline_at": "2026-10-03T00:00:00Z", + "acceptance_criteria": "Accurate", + "rejection_criteria": "Missing", + "locked_payment_policy_version": "payment\n", + "base_amount": "12.50", + "currency": "USD\n", + "payout_type": "fixed\n", + } + with http_fixture() as (origin, response, requests): + for action in ("claim", "start"): + value = body(action) + task = value["task"] if action == "claim" else value + task.update(optional) + response["body"] = json.dumps(value).encode() + result = invoke(cli, origin, action, output="text") + assert result.returncode == 0 and result.stderr == "" + status = "claimed" if action == "claim" else "in_progress" + expected = SUMMARY_TEXT.replace("Status: draft", f"Status: {status}") + expected += ( + "Description: Instructions\\u000A\\u001B[32m\nAcceptance criteria: Accurate\n" + "Rejection criteria: Missing\nSource type: manual\n" + f"Contribution policy: {POLICY}\nGuide version: guide é\\u000A\n" + f"Review policy: {POLICY}\nReview generation: 1\nReview hash: sha256:{'a' * 64}\n" + f"Revision policy: {POLICY}\nRevision generation: 2\nRevision hash: sha256:{'b' * 64}\n" + "Payment policy: payment\\u000A\nBase amount: 12.50\nCurrency: USD\\u000A\nPayout type: fixed\\u000A\n" + ) + if action == "claim": + expected += ( + f"Assignment: {ASSIGNMENT}\nAssignment task: {TASK}\nAssignment project: {PROJECT}\n" + f"Submitter policy: {POLICY}\nContributor: {ACTOR}\nAssigned by: {ACTOR}\n" + "Assigned at: 2026-10-01T00:00:00Z\nAccepted at: 2026-10-01T00:00:00Z\n" + "Released at: —\nAssignment status: active\n" + ) + assert result.stdout == expected + # Explicit nullable fields remain valid and raw JSON is not rewritten. + for name in optional: + task[name] = None + response["body"] = json.dumps(value).encode() + result = invoke(cli, origin, action) + assert ( + result.returncode == 0 + and result.stdout.strip().encode() == response["body"] + ) + assert len(requests) == 4 + + +def test_task_write_arguments_fail_before_network(cli): + with http_fixture() as (origin, _response, requests): + for action in ("claim", "start"): + for selector, key, flags in ( + ("not-uuid", KEY, ()), + (TASK + "/start", KEY, ()), + (TASK, "", ()), + (TASK, "not-uuid", ()), + (TASK, KEY, ("--reason", "é" * 1001)), + (TASK, KEY, ("--reason", b"invalid-\xff")), + (TASK, KEY, ("--actor-id", ACTOR)), + (TASK, KEY, ("--project-id", PROJECT)), + (TASK, KEY, ("--endpoint", "/private")), + ): + assert_failure( + invoke(cli, origin, action, *flags, selector=selector, key=key), + "invalid_arguments", + 2, + ) + assert_failure( + cli(origin, TOKEN, "-o", "json", "task", action, TASK), + "invalid_arguments", + 2, + ) + assert requests == [] + + +def test_task_mutation_malformed_and_substituted_success_is_unknown(cli): + with http_fixture() as (origin, response, requests): + variants = [ + {"id": PROJECT}, + {"project_id": "bad"}, + {"locked_contribution_policy_version_id": "bad"}, + {"locked_guide_version": None}, + {"locked_review_policy_id": "bad"}, + {"locked_review_policy_generation": None}, + {"locked_review_policy_generation": 0}, + {"locked_review_policy_hash": "bad"}, + {"locked_revision_policy_id": None}, + {"locked_revision_policy_generation": "1"}, + {"locked_revision_policy_hash": "bad"}, + {"status": "ready"}, + {"title": None}, + {"description": None}, + {"skill_tags": [None]}, + {"skill_tags": None}, + {"created_at": "bad"}, + {"updated_at": None}, + {"deadline_at": "bad"}, + {"base_amount": "NaN"}, + {"estimated_time_minutes": "10"}, + ] + for action in ("claim", "start"): + for delta in variants + [ + {name: None} + for name in ( + "created_by", + "assigned_to", + "source_ref", + "source_payload_hash", + "import_batch_id", + "external_task_id", + "locked_guide_source_snapshot_id", + "locked_guide_source_snapshot_hash", + "locked_effective_project_submission_artifact_policy_id", + "locked_effective_project_submission_artifact_policy_hash", + "locked_pre_submit_checker_policy_id", + "locked_pre_submit_checker_bundle_hash", + ) + ]: + value = body(action) + (value["task"] if action == "claim" else value).update(delta) + response["body"] = json.dumps(value).encode() + result = invoke(cli, origin, action) + assert_failure(result, "invalid_api_response") + assert json.loads(result.stderr)["error"]["outcome_unknown"] is True + value = body(action) + for malformed in ( + b"null", + b"[]", + b"{}", + json.dumps(value) + .encode() + .replace(b'"id":', b'"id": "' + TASK.encode() + b'", "id":', 1), + ): + response["body"] = malformed + assert_failure(invoke(cli, origin, action), "invalid_api_response") + for delta in ( + {"id": "bad"}, + {"task_id": PROJECT}, + {"project_id": TASK}, + {"submitter_contribution_policy_version_id": TASK}, + {"assigned_by": PROJECT}, + {"contributor_id": "bad"}, + {"status": "released"}, + {"assigned_at": None}, + {"accepted_at": None}, + {"accepted_at": "bad"}, + {"released_at": PROFILE["created_at"]}, + {"private": "secret"}, + ): + value = body("claim") + value["assignment"].update(delta) + response["body"] = json.dumps(value).encode() + result = invoke(cli, origin, "claim") + assert_failure(result, "invalid_api_response") + assert json.loads(result.stderr)["error"]["outcome_unknown"] is True + assert len(requests) == 2 * (len(variants) + 12 + 4) + 12 + + +def test_task_write_uncertainty_and_no_retry(cli): + with http_fixture() as (origin, response, requests): + cases = ( + (403, canonical_error(), "application/json", False), + (409, canonical_error("idempotency_mismatch"), "application/json", False), + (422, canonical_error("validation_error"), "application/json", False), + ( + 403, + {"error": canonical_error()["error"] | {"details": None}}, + "application/json", + True, + ), + (408, {"error": {"code": "gateway_timeout"}}, "application/json", True), + (403, canonical_error(), "text/html", True), + (403, canonical_error() | {"unexpected": True}, "application/json", True), + (503, canonical_error(), "application/json", True), + (302, {}, "application/json", True), + (201, body("claim"), "application/json", True), + (200, body("claim"), "text/html", True), + (200, {}, "application/json", True), + ) + for action in ("claim", "start"): + for status, value, content_type, unknown in cases: + response.update( + status=status, + body=json.dumps(value).encode(), + headers={ + "Content-Type": content_type, + "Location": origin + "/private", + }, + ) + result = invoke(cli, origin, action) + assert result.returncode == 1 and result.stdout == "" + assert ( + json.loads(result.stderr)["error"].get("outcome_unknown", False) + is unknown + ) + response.update( + status=200, + body=b" " * 65537, + headers={"Content-Type": "application/json"}, + ) + result = invoke(cli, origin, action) + assert_failure(result, "invalid_api_response") + assert json.loads(result.stderr)["error"]["outcome_unknown"] is True + response["drop"] = True + result = invoke( + cli, + origin, + action, + "--reason", + "Received before disconnect", + output="text", + ) + assert result.returncode == 1 and result.stdout == "" + assert ( + "task outcome unknown" in result.stderr + and "unchanged action, task, reason and idempotency key" + in result.stderr + ) + assert "whoami" not in result.stderr + assert response["commands"][-1][2] == { + "reason": "Received before disconnect" + } + response["drop"] = False + assert len(requests) == len(response["commands"]) == 2 * (len(cases) + 2) + + +def test_task_posts_are_not_replayed_after_http2_goaway(cli, tmp_path): + for action in ("claim", "start"): + with goaway_fixture(tmp_path) as (origin, env, bodies, connections): + result = cli( + origin, + TOKEN, + "-o", + "json", + "task", + action, + TASK, + "--idempotency-key", + KEY, + "--reason", + "Received before GOAWAY", + extra_env=env, + ) + assert_failure(result, "service_unavailable") + assert json.loads(result.stderr)["error"]["outcome_unknown"] is True + assert connections == ["h2"] + assert [json.loads(value) for value in bodies] == [ + {"reason": "Received before GOAWAY"} + ] diff --git a/cli/tests/integration/test_http_boundary.py b/cli/tests/integration/test_http_boundary.py index b18bf0fce..5db183157 100644 --- a/cli/tests/integration/test_http_boundary.py +++ b/cli/tests/integration/test_http_boundary.py @@ -9,6 +9,8 @@ import time from urllib.parse import parse_qs, urlsplit +from http2_fixture import goaway_fixture + TOKEN = "caller.flow.token" ACTOR = "018f0ebc-7966-7e8d-bc4d-1cae1e000001" PROJECT = "018f0ebc-7966-7e8d-bc4d-1cae1e000002" @@ -45,6 +47,7 @@ def http_fixture(): "delay": 0, "drop": False, "updates": [], + "commands": [], } class Handler(BaseHTTPRequestHandler): @@ -72,6 +75,17 @@ def do_PATCH(self): # noqa: N802 - standard HTTP handler interface ) self.do_GET() + def do_POST(self): # noqa: N802 - standard HTTP handler interface + body = self.rfile.read(int(self.headers.get("Content-Length", "0"))) + response["commands"].append( + ( + self.headers.get("Content-Type"), + self.headers.get_all("Idempotency-Key"), + json.loads(body), + ) + ) + self.do_GET() + def do_CONNECT(self): # noqa: N802 - captures attempted HTTPS proxy use requests.append( (self.command, self.path, self.headers.get("Authorization")) @@ -208,6 +222,27 @@ def test_project_show_invalid_arguments_and_redirect_are_bounded(cli): assert len(requests) == 2 +def test_profile_patch_is_not_replayed_after_http2_goaway(cli, tmp_path): + with goaway_fixture(tmp_path) as (origin, env, bodies, connections): + result = cli( + origin, + TOKEN, + "-o", + "json", + "profile", + "update", + "--display-name", + "Received before GOAWAY", + extra_env=env, + ) + assert_failure(result, "service_unavailable") + assert json.loads(result.stderr)["error"]["outcome_unknown"] is True + assert [json.loads(value) for value in bodies] == [ + {"display_name": "Received before GOAWAY"} + ] + assert connections == ["h2"] + + def test_profile_update_sends_only_selected_fields(cli): with http_fixture() as (origin, response, requests): for flags, expected in ( diff --git a/cli/tests/integration/test_public_self_service.py b/cli/tests/integration/test_public_self_service.py index 92f2c2f40..18cd2618e 100644 --- a/cli/tests/integration/test_public_self_service.py +++ b/cli/tests/integration/test_public_self_service.py @@ -128,6 +128,13 @@ async def test_installed_cli_uses_only_public_profile_and_project_context( "/api/v1/tasks/{task_id}", ): assert "get" in specification.json()["paths"][path] + for operation in ("claim", "start"): + assert ( + "post" + in specification.json()["paths"][ + f"/api/v1/tasks/{{task_id}}/{operation}" + ] + ) profiles: dict[str, dict] = {} for name, token in tokens.items(): result = cli(origin, token, "whoami", "--output", "json") diff --git a/docs/roadmap_status.md b/docs/roadmap_status.md index f282b110e..e9eee24f9 100644 --- a/docs/roadmap_status.md +++ b/docs/roadmap_status.md @@ -137,7 +137,8 @@ The [independent MCP package](../mcp_server/README.md) implements nine self-serv and administrative read tools through WS-MCP-002-03. The [Go CLI](../cli/README.md) provides caller-profile and exact-project authorization reads plus human self-profile editing, exact-project inspection and manager -task queue/detail reads and contributor ready-work/instruction reads, with +task queue/detail reads, contributor ready-work/instruction reads and public +claim/start commands with caller-supplied retry keys, with text/JSON output. Project inspection preserves server-selected full/minimal fields; no public project-list route is invented. CLI write uncertainty is explicit and never automatically @@ -309,6 +310,9 @@ cannot be reused as post-submission review-gate evidence. See the [CLI contributor discovery](../.commitrail/initiatives/WS-CLI-001/WS-CLI-001-05.md) adds ready-task pages and contributor instructions through public reads, preserving Submitter scope and assignment visibility without claiming work. + [CLI contributor claim/start](../.commitrail/initiatives/WS-CLI-001/WS-CLI-001-06.md) + adds the existing public writes with explicit retry keys, server-owned + assignment/lineage, fresh-authority replay and uncertain-outcome handling. Built-binary HTTP integration and isolated real-API proof accompany the package. Further public commands, optional TUI and binary distribution