diff --git a/FOLLOWUPS.md b/FOLLOWUPS.md index 859e395d..d51ca623 100644 --- a/FOLLOWUPS.md +++ b/FOLLOWUPS.md @@ -542,3 +542,36 @@ initial panel plus two fix rounds; do not start a fourth panel cycle for these n No residual acceptance or merge authority is recorded here. The merge decision still belongs to the operator's governed path. Windows visible-window acceptance and effective hook migration remain separate from the green portability tests. + +## Fleet task coordination residuals (PR #288, 2026-09-08) + +Final panel reviewed code head `9a243af21bd7f72f94e5c5af0e7bad0e4d60b1fd`. +Three cycles completed; AGENTS.md Review-cycle discipline requires residual P2s +and nits to be recorded for judgment rather than another panel loop. These are +proposed deferrals, not accepted risk or merge permission. No installation. + +- **P2, Windows branch spelling:** CmdRequest retains supplied spelling rather + than the canonical branch spelling resolved by Git. On case-insensitive Windows + ref lookup, Task can resolve task but compare ownership/activity under a different + key. Fix canonical identity while retaining deleted-branch replay behavior and + add Windows coverage before portability claims. Normal hook effects remain + subject to their own branch guard; request itself grants no lease or execution. + Source: Codex review comment 3954511775. Windows live parity remains unproved. +- **P2, incomplete request evidence:** strictDispatchRows validates base dispatch + fields, not all request-specific types. Missing/invalid numeric at can default to + zero and admit old matching activity while complete remains true. Validate + request_id, worker and finite timestamps before deriving status/replaying. + Source: Codex review comment 3954511783. This is a status-integrity defect, not + evidence of acceptance or authority. Must be resolved or explicitly accepted + before treating the new board as operationally reliable. +- **P2 assessment, busy-store diagnostics:** ErrKeyBusy bubbles up as generic exit + 4 in request/dispatch/reassign/undispatch rather than an actionable refusal. + The lock callback does not execute, so no conflicting write is authorized. + Normalize error classification and test contention response in a follow-on. + Source: Copilot review comment 3954505362 and its suppressed sibling comments. +- **Cosmetic:** scoped row loops retain redundant repo/branch/relationship tests. + Claude final review finds no blocking defects; leave this harmless redundancy. + +Code verification: Fleet race tests, lint/vet, both harness regression suites and +full CI/fuzz/hygiene passed. Those checks do not cover or dismiss the residuals +above. No additional panel request should be sent for this PR under the current cap. diff --git a/WORK.md b/WORK.md new file mode 100644 index 00000000..87907686 --- /dev/null +++ b/WORK.md @@ -0,0 +1,63 @@ + +# Work: Retry-safe Fleet assignments and truthful observation + +Work-ID: fleet-task-coordination +Status: active +Subject: git:f1cea9a96ef77680ab3a6e05ed1614d6b8048e40 +Stop-at: reviewed-change + +## Outcome + +A supervisor can retry an assignment without overwriting or duplicating work and +inspect hook-observed activity without claiming delivery, acceptance or termination. + +## Preserve + +- Existing branch leases, hooks, Gate authority and legacy dispatch behavior except + refusing uncorrelated replacement of a request-bound or unreadable assignment. +- Root cmd/triage/labels/mismatches.jsonl and earlier worktree friction logs are unrelated. +- One existing dispatch store and hook-owned facts; no second editable ledger. + +## Change + +- `FOLLOWUPS.md`: final-panel residuals and written deferral rationale. +- `cmd/fleet/internal/verbs/request.go`: immutable request IDs and serialized replay. +- `cmd/fleet/internal/verbs/role.go`: subscribe generated hooks to write events. +- `cmd/fleet/internal/verbs/role_hooks_test.go`: generated subscription and rebind coverage. +- `cmd/fleet/internal/verbs/status.go`: read-only plain-language request observations. +- `cmd/fleet/internal/verbs/work.go`: protect request records from legacy mutations. +- `cmd/fleet/internal/verbs/verbs.go`: command entrypoints before lazy migration. +- `cmd/fleet/internal/fleet/session.go`: merge per-branch observations under the session lock. +- `cmd/fleet/internal/codex/task_activity_test.go`: patch adapter lease and evidence proof. +- `cmd/fleet/internal/fleet/hook.go`: completed write-tool observation in session record. +- `cmd/fleet/internal/mcp/mcp.go`: equivalent request/status tool entrypoints. +- `cmd/fleet/internal/verbs/request_test.go`: replay, conflicts, races and read-only proof. +- `cmd/fleet/internal/fleet/task_activity_test.go`: post-tool activity provenance. +- `cmd/fleet/internal/mcp/task_test.go`: non-migrating observation through JSON-RPC. +- `cmd/fleet/README.md`: interface, current limits and next adapter proof. + +## Prove + +- Green: go test -race ./cmd/fleet/...; go vet ./cmd/fleet/...; golangci-lint run ./cmd/fleet/... +- Green: both cmd/fleet/testdata/run-suite.sh harness suites; focused real Git CLI exercise. +- Red: competing requests, changed replay, damaged evidence and late unrelated activity never yield false acceptance. + +## Stop + +- No live worker launches, stops, lease transfers or installed hook changes from this PR. +- Delivery, semantic acceptance, correlated answers and effect-safe replacement remain + follow-on implementation under tsk_01M1ZJVZ1ZDHJC1PR1AZGE47TC, not claimed complete. + +## Evidence + +- Verified: focused Fleet race tests, root Go vet/lint, and Claude regression scenarios pass. +- Verified: separate-process replay/conflict tests and Codex regression scenarios pass. +- Verified: incomplete MCP status remains parseable JSON; compiled-binary fixture smoke passes for both adapter shapes. +- Verified: initial-head full-module CI and all three configured reviewers completed. +- Verified: final code head 9a243af has green CI and completed three-member panel. +- Residual: FOLLOWUPS.md records Windows spelling, incomplete request evidence and lock diagnostics for judgment. + +## Handoff + +- Last: second panel consolidated; per-branch evidence, MultiEdit subscription and scoped legacy maintenance fixed. +- Next: judge recorded residuals; no more panel cycles, no merge or live installation. Broader adapter work remains open. diff --git a/cmd/fleet/README.md b/cmd/fleet/README.md new file mode 100644 index 00000000..20a3a545 --- /dev/null +++ b/cmd/fleet/README.md @@ -0,0 +1,82 @@ +# Fleet task coordination: first implementation increment + +This adds retry-safe local assignments and a read-only observation view. It does +not yet implement the four-interaction product: launch/delivery, semantic worker +acceptance, correlated questions, safe stop and replacement remain adapter work. +Do not activate a live trial or present this as cross-harness lifecycle parity. + +The approved direction is [cc-skills PR #60](https://github.com/itsHabib/cc-skills/pull/60): +one lead, one active worker, task-owned workspace and natural interaction through +the supervisor skill. The interfaces below are for the supervisor/adapter, not a +set of commands the operator should have to learn. + +## Record once, retry safely + +From the task's repository checkout, with a known worker: + +```sh +fleet request my-branch --id navigation-fix-1 --worker SESSION \ + --for supervisor:ivy --brief 'Reproduce and fix the navigation failure; return focused checks.' +fleet status +fleet status --json +``` + +Use the discovered executable path if Fleet is not on PATH. The installed Mac +entrypoint during this build was `/Users/mh/.fleet/bin/fleet`; the new commands +are not available there until a reviewed release is installed. + +Equivalent MCP tools are `fleet_request` (requires caller cwd) and `fleet_status`. +The supervisor chooses a stable request ID before calling. Same repo + same ID + +same branch/worker/lead/brief returns the existing assignment without renewing its +timestamp, resetting its initial head, posting a message or acquiring a lease. +Full-ID retries remain valid after branch deletion or session-record cleanup; a +short session prefix must still resolve uniquely. Supply a branch name, not a +numbered change. Changing the payload under that ID refuses. A second assignment for the same +branch refuses, as do unknown ownership, an unavailable worker and an applicable +stop flag. A recorded assignment is not an execution reservation; the ordinary +hook/lease guard still controls actual effects. + +Records extend the existing `dispatch` row with `request_id` and `worker`. +One dispatch-store lock serializes decisions across processes, followed by the +existing branch lock for ownership inspection. Legacy dispatch/reassign/undispatch +cannot overwrite or delete these records, including with `--take`. They remain +retained until a correlated lifecycle operation is implemented. Do not remove +records manually to reuse IDs. Ordinary legacy records remain supported. + +`request` is effectful and performs the existing lazy key migration before lease +inspection. Retained collisions refuse; failed requests can leave a migration +marker/lock but no new assignment. No GitHub write or worker launch occurs. + +## Observe without claiming more than the evidence + +`status` bypasses migration in both CLI and MCP and restores read-only mode after +rendering. JSON is `fleet-task-status-v1`, scoped to local request-bound assignments; +it is not the full portfolio inventory or permission to dispatch. `complete` +means the assignment sources were readable, not that every task is healthy. + +- **Queued:** the assignment exists; delivery and acceptance are unconfirmed. +- **Activity observed:** the selected worker has a matching post-dispatch write-tool + event on the task branch. This is not a claim of successful edits or acceptance. +- **Status needs checking:** conflicting/unreadable ownership, stop flag or missing + worker liveness. A stopped/dead session does not establish command quiescence. + +The existing hook-owned session record carries `last_writes`, keyed by branch, +with time and tool-use ID; `last_write` remains for compatibility. Writes on a +second branch do not erase the first branch observation. Read-only commands, old activity and another branch do not count. +JSON keeps IDs and evidence timestamps for debugging; terminal output does not +require the operator to interpret internal session IDs. No new agent-written +progress ledger, acceptance claim, done state or automatic takeover is introduced. + +Generated role bindings subscribe Codex write events and supplement Claude +file-write post-tool events alongside its global Bash hook. Existing bindings +need regeneration and harness reload when this release is installed. This PR +does not edit installed hooks. Terminal observations include the activity age. + +## Verification + +Run `go test -race ./cmd/fleet/...`, `go vet ./cmd/fleet/...` and +`golangci-lint run ./cmd/fleet/...`, then both `testdata/run-suite.sh` variants +(default and `codex`). New tests cover real Git state, separate-process replay and +conflicts, immutable payloads, legacy-writer protection, damaged evidence, +post-tool provenance and non-migrating JSON-RPC observation. Harness event +fixtures are not proof of actual live Claude/Codex delivery or stop behavior. diff --git a/cmd/fleet/internal/codex/task_activity_test.go b/cmd/fleet/internal/codex/task_activity_test.go new file mode 100644 index 00000000..eb6e8623 --- /dev/null +++ b/cmd/fleet/internal/codex/task_activity_test.go @@ -0,0 +1,48 @@ +package codex + +import ( + "github.com/itsHabib/workbench/cmd/fleet/internal/fleet" + "os" + "path/filepath" + "testing" +) + +func TestPatchAdapterEnforcesLeaseAndRecordsActivity(t *testing.T) { + oldState, oldOrg := fleet.State, fleet.OrgState + root := t.TempDir() + fleet.State, fleet.OrgState = filepath.Join(root, "state"), filepath.Join(root, "org") + t.Cleanup(func() { fleet.State, fleet.OrgState = oldState, oldOrg }) + repo := filepath.Join(root, "repo") + if err := os.MkdirAll(filepath.Join(repo, ".git", "refs", "heads"), 0700); err != nil { + t.Fatal(err) + } + if err := os.WriteFile(filepath.Join(repo, ".git", "HEAD"), []byte("ref: refs/heads/task\n"), 0600); err != nil { + t.Fatal(err) + } + key := fleet.Scope(repo, "task") + if err := fleet.WriteJSON(fleet.Path("sessions", "holder.json"), fleet.Rec{"session": "holder", "pid": os.Getpid(), "pid_kind": "harness", "last_event_at": fleet.Now()}); err != nil { + t.Fatal(err) + } + if err := fleet.WriteLease(key, fleet.LeaseRecord(key, "holder", "", repo, nil)); err != nil { + t.Fatal(err) + } + ev := fleet.Event{"session_id": "worker", "cwd": repo, "hook_event_name": "PreToolUse", "tool_name": "apply_patch", "tool_input": fleet.Rec{"command": "*** Begin Patch\n*** Add File: file.txt\n+hello\n*** End Patch"}} + if v := Run(ev); v.Code != 2 { + t.Fatalf("patch escaped foreign lease: %v", v) + } + if fleet.S(fleet.Lease(key), "session") != "holder" { + t.Fatal("foreign lease changed") + } + ev["session_id"] = "holder" + if v := Run(ev); v.Code != 0 { + t.Fatal(v) + } + ev["hook_event_name"] = "PostToolUse" + if v := Run(ev); v.Code != 0 { + t.Fatal(v) + } + write := fleet.M(fleet.M(fleet.SessionRecord("holder"), "last_writes"), key) + if fleet.S(write, "key") != key || fleet.F(write, "at") == 0 { + t.Fatalf("patch observation absent: %v", write) + } +} diff --git a/cmd/fleet/internal/fleet/hook.go b/cmd/fleet/internal/fleet/hook.go index 53fd5365..a1d14bc7 100644 --- a/cmd/fleet/internal/fleet/hook.go +++ b/cmd/fleet/internal/fleet/hook.go @@ -379,7 +379,7 @@ func recordInflight(ev Event, sid, cmd string) { } func onPostTool(ev Event, sid string) *Verdict { - rec := TouchSession(sid, ev, Rec{}) + rec := TouchSession(sid, ev, postWriteEvidence(ev, sid)) cmd := S(M(ev, "tool_input"), "command") if cmd != "" { start := S(ev, "cwd") @@ -452,3 +452,24 @@ func Exit(v *Verdict) { } os.Exit(v.Code) } + +// postWriteEvidence records observed tool activity, not acceptance or success. +// It is telemetry in the existing session record; missing evidence never permits +// a tool and does not imply that a queued assignment started. +func postWriteEvidence(ev Event, sid string) Rec { + tool := S(ev, "tool_name") + cmd := S(M(ev, "tool_input"), "command") + if !IsWrite(tool, cmd) { + return Rec{} + } + target := preToolTarget(ev, SessionRecord(sid), tool, cmd) + branch := BranchOf(target) + if branch == "" { + return Rec{} + } + key := Scope(target, branch) + if key == "" { + return Rec{} + } + return Rec{"last_write": Rec{"key": key, "at": Now(), "tool_use_id": ev["tool_use_id"]}} +} diff --git a/cmd/fleet/internal/fleet/session.go b/cmd/fleet/internal/fleet/session.go index 73a00d33..7eb5d550 100644 --- a/cmd/fleet/internal/fleet/session.go +++ b/cmd/fleet/internal/fleet/session.go @@ -62,6 +62,16 @@ func touchSessionLocked(sid string, ev Event, fields Rec) (Rec, error) { pid, kind := HarnessPid(true) rec["pid"], rec["pid_kind"] = float64(pid), kind } + // Merge branch observations under the session lock so concurrent tools on + // different task branches cannot erase each other's evidence. + if write := M(fields, "last_write"); S(write, "key") != "" { + writes := M(rec, "last_writes") + if writes == nil { + writes = Rec{} + } + writes[S(write, "key")] = write + rec["last_writes"] = writes + } for k, v := range fields { rec[k] = v } diff --git a/cmd/fleet/internal/fleet/task_activity_test.go b/cmd/fleet/internal/fleet/task_activity_test.go new file mode 100644 index 00000000..089601f2 --- /dev/null +++ b/cmd/fleet/internal/fleet/task_activity_test.go @@ -0,0 +1,39 @@ +package fleet + +import ( + "os" + "path/filepath" + "testing" +) + +func TestPostWriteEvidenceOnlyObservedTarget(t *testing.T) { + oldState, oldOrg := State, OrgState + root := t.TempDir() + State, OrgState = filepath.Join(root, "fleet"), filepath.Join(root, "org") + t.Cleanup(func() { State, OrgState = oldState, oldOrg }) + repo := filepath.Join(root, "repo") + if err := os.MkdirAll(filepath.Join(repo, ".git", "refs", "heads"), 0700); err != nil { + t.Fatal(err) + } + _ = os.WriteFile(filepath.Join(repo, ".git", "HEAD"), []byte("ref: refs/heads/task\n"), 0600) + ev := Event{"session_id": "worker", "cwd": repo, "hook_event_name": "PostToolUse", "tool_name": "Read", "tool_input": Rec{"file_path": filepath.Join(repo, "file")}} + if len(postWriteEvidence(ev, "worker")) != 0 { + t.Fatal("read counted as write activity") + } + ev["tool_name"] = "Edit" + got := postWriteEvidence(ev, "worker") + if S(M(got, "last_write"), "key") != Scope(repo, "task") || F(M(got, "last_write"), "at") == 0 { + t.Fatal(got) + } + // Actual hook persists this fact for both adapter faces; it is not an agent-written acceptance. + if v := Run(ev); v.Code != 0 { + t.Fatal(v) + } + if S(M(SessionRecord("worker"), "last_write"), "key") != Scope(repo, "task") { + t.Fatal(SessionRecord("worker")) + } + ev["tool_input"] = Rec{"file_path": filepath.Join(root, "outside")} + if len(postWriteEvidence(ev, "worker")) != 0 { + t.Fatal("outside target attributed to task") + } +} diff --git a/cmd/fleet/internal/mcp/mcp.go b/cmd/fleet/internal/mcp/mcp.go index 7cf15901..a2e1b766 100644 --- a/cmd/fleet/internal/mcp/mcp.go +++ b/cmd/fleet/internal/mcp/mcp.go @@ -67,6 +67,11 @@ var tools = []schema{ "for": str("accountable role; default: the dispatcher"), "due": str("duration like 45m or 2h"), "slot": str("free slot to place the work in (fleet_slots)"), "brief": str("one line the slot's session reads at start"), "reply_to": str("your session id, handed to the seat as its address for questions"), "take": schema{"type": "boolean", "description": "rewrite a row that has live hands"}, "cwd": cwdArg}, "required": []any{"change", "as", "cwd"}}}, + {"name": "fleet_request", + "description": "Record one retry-safe local assignment for a known worker; does not deliver, accept, launch or transfer a lease.", + "inputSchema": schema{"type": "object", "properties": schema{"change": str("branch"), "id": str("stable repo-scoped request ID"), "worker": str("known session ID or unique prefix"), "for": str("accountable lead"), "brief": str("bounded assignment"), "cwd": cwdArg}, "required": []any{"change", "id", "worker", "for", "brief", "cwd"}}}, + {"name": "fleet_status", "description": "Read-only local request board; activity is not acceptance or completion.", + "inputSchema": schema{"type": "object", "properties": schema{}}}, {"name": "fleet_work", "description": "Every ownership row on this machine with its observed state: dead, late, undeclared (need a decision); working, idle, dispatched, done.", "inputSchema": schema{"type": "object", "properties": schema{"for": str("only rows this role is accountable for"), "cwd": cwdArg}}}, @@ -259,6 +264,10 @@ func dispatch(name string, a map[string]any) (string, bool) { return runVerb(func() error { return verbs.CmdDispatch(s("change"), s("as"), s("for"), s("due"), s("slot"), s("brief"), "mcp", s("reply_to"), take) }) + case "fleet_request": + return runVerb(func() error { return verbs.CmdRequest(s("change"), s("id"), s("worker"), s("for"), s("brief")) }) + case "fleet_status": + return taskStatusJSON() case "fleet_work": return js(verbs.WorkRows(s("for"))), false case "fleet_reassign": @@ -299,7 +308,9 @@ func handle(msg map[string]any) map[string]any { case "tools/call": name, _ := params["name"].(string) args, _ := params["arguments"].(map[string]any) - fleet.MigrateLegacyKeys() + if name != "fleet_status" && name != "fleet_request" { + fleet.MigrateLegacyKeys() + } text, isErr, err := safeCall(name, args) if errors.Is(err, errUnknownTool) { return rpcError(id, -32602, fmt.Sprintf("unknown tool %s", fleet.PyRepr(name))) @@ -364,3 +375,14 @@ func handleSafe(msg map[string]any) (resp map[string]any) { }() return handle(msg) } + +// Preserve a parseable packet even when its sources are incomplete. The tool's +// isError flag carries the failure; appending prose would corrupt the JSON. +func taskStatusJSON() (string, bool) { + var buf strings.Builder + previous := verbs.Out + verbs.Out = &buf + defer func() { verbs.Out = previous }() + err := verbs.CmdStatus([]string{"--json"}) + return buf.String(), err != nil +} diff --git a/cmd/fleet/internal/mcp/task_test.go b/cmd/fleet/internal/mcp/task_test.go new file mode 100644 index 00000000..5c530e80 --- /dev/null +++ b/cmd/fleet/internal/mcp/task_test.go @@ -0,0 +1,46 @@ +package mcp + +import ( + "encoding/json" + "os" + "path/filepath" + "strings" + "testing" + + "github.com/itsHabib/workbench/cmd/fleet/internal/fleet" +) + +func TestTaskStatusMCPDoesNotInitializeOrMigrate(t *testing.T) { + old := fleet.State + fleet.State = filepath.Join(t.TempDir(), "absent") + t.Cleanup(func() { fleet.State = old }) + var out strings.Builder + Serve(strings.NewReader("{\"jsonrpc\":\"2.0\",\"id\":1,\"method\":\"tools/call\",\"params\":{\"name\":\"fleet_status\",\"arguments\":{}}}\n"), &out) + if !strings.Contains(out.String(), "fleet-task-status-v1") || strings.Contains(out.String(), `"isError":true`) { + t.Fatal(out.String()) + } + if _, err := os.Stat(fleet.State); !os.IsNotExist(err) { + t.Fatalf("read initialized state: %v", err) + } + _ = fleet.WriteJSON(fleet.Path("stop", "legacy.json"), fleet.Rec{"branch": "legacy", "reason": "keep"}) + out.Reset() + Serve(strings.NewReader("{\"jsonrpc\":\"2.0\",\"id\":2,\"method\":\"tools/call\",\"params\":{\"name\":\"fleet_status\",\"arguments\":{}}}\n"), &out) + var reply struct { + Result struct { + Content []struct{ Text string } + IsError bool + } + } + if err := json.Unmarshal([]byte(out.String()), &reply); err != nil { + t.Fatal(err) + } + if len(reply.Result.Content) != 1 || !json.Valid([]byte(reply.Result.Content[0].Text)) || !reply.Result.IsError { + t.Fatal("incomplete status is not a parseable error packet", out.String()) + } + if !strings.Contains(out.String(), "legacy state needs reconciliation") { + t.Fatal(out.String()) + } + if _, err := os.Stat(fleet.Path("stop", "legacy.json")); err != nil { + t.Fatal("legacy moved", err) + } +} diff --git a/cmd/fleet/internal/verbs/request.go b/cmd/fleet/internal/verbs/request.go new file mode 100644 index 00000000..6b3be122 --- /dev/null +++ b/cmd/fleet/internal/verbs/request.go @@ -0,0 +1,232 @@ +package verbs + +import ( + "fmt" + "os" + "path/filepath" + "regexp" + "strings" + + "github.com/itsHabib/workbench/cmd/fleet/internal/fleet" +) + +var requestID = regexp.MustCompile(`^[A-Za-z0-9][A-Za-z0-9._-]{0,95}$`) + +// CmdRequest declares an immutable, retry-safe assignment in the existing dispatch +// store. It does not launch, message, accept, acquire a lease, or place a worktree. +// The ID is repo-scoped. Replaying it preserves the original head and timestamp. +func CmdRequest(change, id, worker, lead, brief string) error { + if fleet.ReadOnly { + return refuse("fleet request: cannot dispatch in read-only mode") + } + + if !requestID.MatchString(id) || worker == "" || lead == "" || strings.TrimSpace(brief) == "" { + return refuse("fleet request: requires --id (1-96 letters/digits/._-), --worker, --for and --brief") + } + branch := strings.TrimSpace(change) + if branch == "" || strings.HasPrefix(branch, "#") { + return refuse("fleet request: provide a branch name rather than a numbered change") + } + rid := fleet.RepoID(cwd()) + if rid == "" { + return refuse("fleet request: run inside the target repository") + } + // Ownership and activity are keyed by the branch as git spells it, never by the + // caller's spelling: on a case-insensitive filesystem two spellings would + // otherwise key two assignments for one branch. + requested := branch + branch, found := canonicalBranch(cwd(), branch) + wanted := fleet.Rec{"request_id": id, "repo": rid, "change": branch, "requested": requested, + "worker": worker, "for": lead, "brief": strings.TrimSpace(brief), "relationship": "implementation"} + replay := false + err := fleet.KeyLock("dispatch", func() error { + // Reconcile existing keys like legacy dispatch; a retained collision refuses. + fleet.MigrateLegacyKeys() + if fleet.MigrationPending() { + return refuse("fleet request: legacy ownership needs reconciliation before assigning work") + } + rows, err := strictDispatchRows() + if err != nil { + return err + } + for _, row := range rows { + if fleet.S(row, "repo") != rid { + continue + } + if fleet.S(row, "request_id") == id { + if err := validateReplay(row, wanted, worker, found); err != nil { + return err + } + replay = true + return nil + } + } + _, _, head, err := resolveDispatchTarget("request", branch) + if err != nil { + return err + } + sid, err := findSession(worker) + if err != nil { + return err + } + wanted["worker"] = sid + return createRequest(rows, wanted, head) + }) + if err != nil { + return err + } + if replay { + say("Assignment already recorded; no new dispatch or delivery. Inspect `fleet status` for observed activity.") + return nil + } + say("Queued: %s. Assignment recorded; worker delivery and acceptance are not yet confirmed.", branch) + return nil +} + +// A full retained ID survives session cleanup; prefixes must still resolve +// uniquely. No live branch or worker is required for an identical replay. +func validateReplay(row, wanted fleet.Rec, worker string, resolved bool) error { + if worker != fleet.S(row, "worker") { + sid, err := findSession(worker) + if err != nil { + return err + } + wanted["worker"] = sid + } + if !sameRequest(row, wanted, resolved) { + return refuse("fleet request: request ID already has different work or recipient; inspect `fleet status`") + } + return nil +} + +// sameRequest is identity for a replay. The change is compared by git's spelling +// when the caller's name resolved to a ref now — two refs that differ only by +// case are different work on a case-sensitive filesystem — and by the caller's +// original spelling when it did not, which is the deleted-branch retry the +// request contract promises. Never by a case-folded match, which would let an +// ID recorded for Nav-Fix be replayed against a distinct nav-fix. +func sameRequest(a, b fleet.Rec, resolved bool) bool { + for _, key := range []string{"request_id", "repo", "worker", "for", "brief", "relationship"} { + if fleet.S(a, key) != fleet.S(b, key) { + return false + } + } + if resolved { + return fleet.S(a, "change") == fleet.S(b, "change") + } + return fleet.S(a, "requested") == fleet.S(b, "requested") +} + +func createRequest(rows []fleet.Rec, wanted fleet.Rec, head string) error { + rid, branch, sid := fleet.S(wanted, "repo"), fleet.S(wanted, "change"), fleet.S(wanted, "worker") + for _, row := range rows { + if fleet.S(row, "repo") == rid && fleet.S(row, "change") == branch { + return refuse("fleet request: this branch already has an assignment; inspect `fleet status` and existing work before assigning again") + } + } + if !fleet.SessionAlive(fleet.SessionRecord(sid)) { + return refuse("fleet request: selected worker is not observably live; choose a reachable worker before queuing") + } + key := "repo:" + rid + ":" + branch + // Serialize with the hook's first-write lease acquisition. This check does not + // reserve the branch: assignment is accountability, not execution authority. + return fleet.KeyLock(key, func() error { + if fleet.MigrationPending() { + return refuse("fleet request: legacy ownership changed during dispatch; inspect state") + } + if _, err := os.Lstat(dispatchFile(rid, branch, "implementation")); !os.IsNotExist(err) { + return refuse("fleet request: assignment path already exists or is unreadable; inspect state") + } + if reason := fleet.CheckStop(key, branch, sid); reason != "" { + return refuse("fleet request: this branch is stopped; no new assignment was recorded") + } + lease := fleet.Lease(key) + if fleet.IsMalformed(lease) { + return refuse("fleet request: branch ownership is unreadable; inspect it before dispatch") + } + if lease != nil && fleet.S(lease, "session") != sid { + return refuse("fleet request: another session holds this branch; no assignment or takeover performed") + } + wanted["at"], wanted["head_at_dispatch"], wanted["by"] = fleet.Now(), head, dispatcher("") + if err := fleet.WriteJSON(dispatchFile(rid, branch, "implementation"), wanted); err != nil { + return err + } + fleet.ObserveAction("dispatch", wanted) + return nil + }) +} + +// strictDispatchRows refuses gaps instead of letting a damaged row disappear and +// allowing its request ID or branch to be assigned a second time. +func strictDispatchRows() ([]fleet.Rec, error) { + entries, err := os.ReadDir(dispatchDir()) + if os.IsNotExist(err) { + return nil, nil + } + if err != nil { + return nil, err + } + var rows []fleet.Rec + for _, entry := range entries { + if !strings.HasSuffix(entry.Name(), ".json") { + continue + } + row := fleet.ReadJSON(filepath.Join(dispatchDir(), entry.Name())) + if row == nil || fleet.S(row, "repo") == "" || fleet.S(row, "change") == "" || fleet.S(row, "relationship") == "" { + return nil, fmt.Errorf("assignment evidence unreadable: %s", entry.Name()) + } + // A request-bound row carries who it went to and when; without either it is + // damaged evidence, and damaged evidence must not read as a complete status. + if id, present := row["request_id"]; present { + ids, ok := id.(string) + if !ok || !requestID.MatchString(ids) || fleet.S(row, "worker") == "" || fleet.S(row, "for") == "" || fleet.F(row, "at") <= 0 { + return nil, fmt.Errorf("assignment evidence damaged: %s", entry.Name()) + } + } + rows = append(rows, row) + } + return rows, nil +} + +// canonicalBranch is the branch as git lists it. An exact match wins; else a +// unique case-insensitive match is the same ref spelled differently; else the +// caller's spelling stands (a branch that does not exist yet keeps its name). +func canonicalBranch(dir, cand string) (string, bool) { + rc, out := gitTry(dir, gitTimeout, "for-each-ref", "--format=%(refname:short)", "refs/heads") + if rc != 0 { + return cand, false + } + match := "" + for _, name := range strings.Split(out, "\n") { + name = strings.TrimSpace(name) + if name == cand { + return cand, true + } + if name != "" && strings.EqualFold(name, cand) { + if match != "" { + return cand, true // two refs differ only by case: the caller's exact name stands + } + match = name + } + } + if match != "" { + return match, true + } + return cand, false +} + +func dispatchRequest(args []string) error { + vals := map[string]string{} + for _, flag := range []string{"--id", "--worker", "--for", "--brief"} { + value, err := optValue(args, flag, "request") + if err != nil { + return err + } + vals[flag] = value + } + pos := positional(args, "--id", "--worker", "--for", "--brief") + if len(pos) != 1 { + return refuse("usage: fleet request --id --worker --for --brief ") + } + return CmdRequest(pos[0], vals["--id"], vals["--worker"], vals["--for"], vals["--brief"]) +} diff --git a/cmd/fleet/internal/verbs/request_test.go b/cmd/fleet/internal/verbs/request_test.go new file mode 100644 index 00000000..dcd473f2 --- /dev/null +++ b/cmd/fleet/internal/verbs/request_test.go @@ -0,0 +1,558 @@ +package verbs + +import ( + "bytes" + "encoding/json" + "io" + "os" + "os/exec" + "path/filepath" + "sync" + "testing" + + "github.com/itsHabib/workbench/cmd/fleet/internal/fleet" +) + +func requestFixture(t *testing.T) (string, string) { + t.Helper() + oldState, oldOrg, oldOut := fleet.State, fleet.OrgState, Out + oldCwd, err := os.Getwd() + if err != nil { + t.Fatal(err) + } + root := t.TempDir() + fleet.State, fleet.OrgState, Out = filepath.Join(root, "state"), filepath.Join(root, "org"), io.Discard + t.Setenv("FLEET_GITHUB", "off") + t.Setenv("FLEET_WATCH", "off") + t.Cleanup(func() { fleet.State, fleet.OrgState, Out = oldState, oldOrg, oldOut; _ = os.Chdir(oldCwd) }) + repo := filepath.Join(root, "repo") + if err := os.MkdirAll(repo, 0700); err != nil { + t.Fatal(err) + } + runGit(t, repo, "init", "-b", "task") + runGit(t, repo, "-c", "user.name=Test", "-c", "user.email=test@example.invalid", "commit", "--allow-empty", "-m", "fixture") + if err := os.Chdir(repo); err != nil { + t.Fatal(err) + } + sid := "11111111-1111-1111-1111-111111111111" + if err := fleet.WriteJSON(fleet.Path("sessions", sid+".json"), fleet.Rec{ + "session": sid, "cwd": repo, "repo": fleet.RepoID(repo), "branch": "task", "last_event_at": fleet.Now(), "pid_kind": "harness", "pid": os.Getpid(), + }); err != nil { + t.Fatal(err) + } + return repo, sid +} + +func runGit(t *testing.T, repo string, args ...string) { + t.Helper() + cmd := exec.Command("git", append([]string{"-C", repo}, args...)...) + if out, err := cmd.CombinedOutput(); err != nil { + t.Fatalf("git %v: %s: %v", args, out, err) + } +} + +func requestFile(repo string) string { + return dispatchFile(fleet.RepoID(repo), "task", "implementation") +} +func readBytes(t *testing.T, path string) []byte { + t.Helper() + b, e := os.ReadFile(path) + if e != nil { + t.Fatal(e) + } + return b +} + +func TestRequestReplayPreservesOriginalAfterHeadChanges(t *testing.T) { + repo, sid := requestFixture(t) + if err := CmdRequest("task", "fix-1", sid, "lead", "fix it"); err != nil { + t.Fatal(err) + } + before := readBytes(t, requestFile(repo)) + runGit(t, repo, "-c", "user.name=Test", "-c", "user.email=test@example.invalid", "commit", "--allow-empty", "-m", "new head") + if err := CmdRequest("task", "fix-1", sid[:8], "lead", "fix it"); err != nil { + t.Fatal(err) + } + if !bytes.Equal(before, readBytes(t, requestFile(repo))) { + t.Fatal("replay rewrote head, timestamp or assignment") + } + if err := CmdRequest("task", "fix-1", sid, "lead", "different task"); err == nil { + t.Fatal("changed replay accepted") + } + if !bytes.Equal(before, readBytes(t, requestFile(repo))) { + t.Fatal("conflict changed record") + } + if err := CmdRequest("task", "fix-2", sid, "lead", "fix it"); err == nil { + t.Fatal("second request replaced occupied assignment") + } +} + +func TestConcurrentRequestOnlyOneWins(t *testing.T) { + repo, sid := requestFixture(t) + start := make(chan struct{}) + results := make(chan error, 2) + var wg sync.WaitGroup + for _, id := range []string{"one", "two"} { + wg.Add(1) + go func(id string) { defer wg.Done(); <-start; results <- CmdRequest("task", id, sid, "lead", "fix it") }(id) + } + close(start) + wg.Wait() + close(results) + successes := 0 + for err := range results { + if err == nil { + successes++ + } + } + if successes != 1 { + t.Fatalf("%d assignments won", successes) + } + if fleet.ReadJSON(requestFile(repo)) == nil { + t.Fatal("no surviving assignment") + } +} + +func TestRequestProtectsLegacyMutations(t *testing.T) { + repo, sid := requestFixture(t) + if err := CmdRequest("task", "one", sid, "lead", "fix it"); err != nil { + t.Fatal(err) + } + before := readBytes(t, requestFile(repo)) + for _, fn := range []func() error{ + func() error { return CmdDispatch("task", "implementation", "lead", "", "", "changed", "", "", true) }, + func() error { return CmdReassign("task", "other") }, + func() error { return cmdUndispatch("task", "") }, + } { + if err := fn(); err == nil { + t.Fatal("uncorrelated mutation accepted") + } + } + if !bytes.Equal(before, readBytes(t, requestFile(repo))) { + t.Fatal("immutable request changed") + } +} + +func TestRequestRejectsForeignLeaseAndMalformedStore(t *testing.T) { + for _, mode := range []string{"foreign", "malformed", "collision", "legacy"} { + t.Run(mode, func(t *testing.T) { + repo, sid := requestFixture(t) + rid := fleet.RepoID(repo) + switch mode { + case "foreign": + _ = fleet.WriteJSON(fleet.KeyFile("leases", "repo:"+rid+":task"), fleet.Rec{"session": "other", "key": "repo:" + rid + ":task"}) + case "malformed": + _ = os.MkdirAll(dispatchDir(), 0700) + _ = os.WriteFile(filepath.Join(dispatchDir(), "broken.json"), []byte("{"), 0600) + case "collision": + _ = fleet.WriteJSON(requestFile(repo), fleet.Rec{"repo": "other", "change": "other", "relationship": "implementation"}) + case "legacy": + _ = fleet.WriteJSON(fleet.Path("stop", "task.json"), fleet.Rec{"repo": rid, "branch": "task", "reason": "legacy"}) + } + if err := CmdRequest("task", "one", sid, "lead", "fix it"); err == nil { + t.Fatalf("%s accepted", mode) + } + if mode != "collision" { + if _, err := os.Stat(requestFile(repo)); !os.IsNotExist(err) { + t.Fatalf("%s wrote assignment", mode) + } + } + }) + } +} + +func TestStatusDoesNotInventAcceptanceOrCompletion(t *testing.T) { + repo, sid := requestFixture(t) + if err := CmdRequest("task", "one", sid, "lead", "fix it"); err != nil { + t.Fatal(err) + } + d := fleet.ReadJSON(requestFile(repo)) + key := "repo:" + fleet.RepoID(repo) + ":task" + rec := fleet.SessionRecord(sid) + rec["turn_open"] = true + for _, c := range []struct { + name, key string + at float64 + want string + }{ + {"no activity", "", 0, "Queued"}, + {"old activity", key, fleet.F(d, "at") - 1, "Queued"}, + {"other branch", key + "-other", fleet.Now(), "Queued"}, + {"this branch", key, fleet.Now(), "Activity observed"}, + } { + t.Run(c.name, func(t *testing.T) { + rec["last_write"] = fleet.Rec{"key": c.key, "at": c.at} + _ = fleet.WriteJSON(fleet.Path("sessions", sid+".json"), rec) + row := requestStatus(d, fleet.Now()) + if row["status"] != c.want { + t.Fatalf("%s: %v", c.name, row) + } + }) + } + _ = fleet.WriteJSON(fleet.KeyFile("stop", key), fleet.Rec{"key": key, "reason": "stop"}) + if row := requestStatus(d, fleet.Now()); row["status"] == "Stopped" { + t.Fatal("stop flag reported termination") + } + // A dead session does not prove clean stop either. + fleet.Unlink(fleet.KeyFile("stop", key)) + rec["ended"] = true + _ = fleet.WriteJSON(fleet.Path("sessions", sid+".json"), rec) + if row := requestStatus(d, fleet.Now()); row["status"] != "Status needs checking" { + t.Fatal(row) + } +} + +func TestStatusReadOnlyAndGapReporting(t *testing.T) { + repo, sid := requestFixture(t) + if err := CmdRequest("task", "one", sid, "lead", "fix it"); err != nil { + t.Fatal(err) + } + before := snapshotFiles(t, fleet.State) + var b bytes.Buffer + Out = &b + if err := Dispatch([]string{"status", "--json"}); err != nil { + t.Fatal(err) + } + var packet map[string]any + if err := json.Unmarshal(b.Bytes(), &packet); err != nil { + t.Fatal(err) + } + if packet["schema"] != "fleet-task-status-v1" || packet["complete"] != true { + t.Fatal(packet) + } + after := snapshotFiles(t, fleet.State) + if !sameFiles(before, after) { + t.Fatal("status mutated state") + } + _ = os.WriteFile(requestFile(repo), []byte("{"), 0600) + b.Reset() + if err := Dispatch([]string{"status", "--json"}); err == nil { + t.Fatal("damaged assignment hidden") + } + if err := json.Unmarshal(b.Bytes(), &packet); err != nil { + t.Fatal(err) + } + if packet["complete"] != false { + t.Fatal(packet) + } +} + +func snapshotFiles(t *testing.T, root string) map[string]string { + t.Helper() + m := map[string]string{} + err := filepath.WalkDir(root, func(path string, d os.DirEntry, err error) error { + if err != nil { + return err + } + if d.IsDir() { + return nil + } + b, err := os.ReadFile(path) + m[path] = string(b) + return err + }) + if err != nil { + t.Fatal(err) + } + return m +} +func sameFiles(a, b map[string]string) bool { + if len(a) != len(b) { + return false + } + for k, v := range a { + if b[k] != v { + return false + } + } + return true +} + +func TestRequestCLIRejectsIncompleteAndUnsafeIDs(t *testing.T) { + _, sid := requestFixture(t) + for _, args := range [][]string{ + {"request", "task", "--id", "../escape", "--worker", sid, "--for", "lead", "--brief", "fix"}, + {"request", "task", "--id", "one", "--worker", sid, "--brief", "fix"}, + {"request", "task", "--id", "one", "--worker", sid, "--for", "lead", "--brief", "fix", "--take"}, + } { + if err := Dispatch(args); err == nil { + t.Fatalf("accepted %v", args) + } + } +} + +// Replays also cross the process boundary: each contender has independent globals +// and file descriptors, as two supervisor harnesses would. +func TestRequestConcurrentProcesses(t *testing.T) { + for _, same := range []bool{false, true} { + t.Run(map[bool]string{false: "different IDs", true: "same ID"}[same], func(t *testing.T) { + repo, sid := requestFixture(t) + ids := []string{"one", "two"} + if same { + ids[1] = "one" + } + var cmds []*exec.Cmd + for _, id := range ids { + cmd := exec.Command(os.Args[0], "-test.run=^TestRequestHelperProcess$") + cmd.Dir = repo + cmd.Env = append(os.Environ(), "FLEET_REQUEST_HELPER=1", "FLEET_STATE="+fleet.State, "ORG_STATE="+fleet.OrgState, "FLEET_REQUEST_ID="+id, "FLEET_REQUEST_WORKER="+sid) + cmds = append(cmds, cmd) + if err := cmd.Start(); err != nil { + t.Fatal(err) + } + } + successes := 0 + for _, cmd := range cmds { + if cmd.Wait() == nil { + successes++ + } + } + want := 1 + if same { + want = 2 + } + if successes != want { + t.Fatalf("%d successful processes, want %d", successes, want) + } + rows, err := strictDispatchRows() + if err != nil || len(rows) != 1 { + t.Fatalf("rows %v, err %v", rows, err) + } + }) + } +} + +func TestRequestHelperProcess(_ *testing.T) { + if os.Getenv("FLEET_REQUEST_HELPER") != "1" { + return + } + Out = io.Discard + if err := CmdRequest("task", os.Getenv("FLEET_REQUEST_ID"), os.Getenv("FLEET_REQUEST_WORKER"), "lead", "fix it"); err != nil { + os.Exit(1) + } + os.Exit(0) +} + +func TestRequestReplayAfterBranchAndSessionRemoval(t *testing.T) { + repo, sid := requestFixture(t) + if err := CmdRequest("task", "replay", sid, "lead", "fix it"); err != nil { + t.Fatal(err) + } + before := readBytes(t, requestFile(repo)) + runGit(t, repo, "checkout", "--detach") + runGit(t, repo, "branch", "-D", "task") + fleet.Unlink(fleet.Path("sessions", sid+".json")) + if err := CmdRequest("task", "replay", sid, "lead", "fix it"); err != nil { + t.Fatal(err) + } + if !bytes.Equal(before, readBytes(t, requestFile(repo))) { + t.Fatal("replay changed record") + } + if err := CmdRequest("task", "replay", sid, "lead", "changed"); err == nil { + t.Fatal("conflicting replay accepted") + } +} + +func TestLegacyDispatchIgnoresUnrelatedDamage(t *testing.T) { + repo, _ := requestFixture(t) + if err := os.MkdirAll(dispatchDir(), 0700); err != nil { + t.Fatal(err) + } + if err := os.WriteFile(filepath.Join(dispatchDir(), "unrelated.json"), []byte("broken"), 0600); err != nil { + t.Fatal(err) + } + rid := fleet.RepoID(repo) + if _, err := replaceableDispatch(rid, "task", "verify"); err != nil { + t.Fatal(err) + } + path := dispatchFile(rid, "task", "verify") + if err := os.WriteFile(path, []byte("broken"), 0600); err != nil { + t.Fatal(err) + } + if _, err := replaceableDispatch(rid, "task", "verify"); err == nil { + t.Fatal("target damage ignored") + } +} + +func TestUndispatchLegacySiblingPreservesRequest(t *testing.T) { + repo, sid := requestFixture(t) + if err := CmdRequest("task", "one", sid, "lead", "fix it"); err != nil { + t.Fatal(err) + } + before := readBytes(t, requestFile(repo)) + path := dispatchFile(fleet.RepoID(repo), "task", "verify") + if err := fleet.WriteJSON(path, fleet.Rec{"repo": fleet.RepoID(repo), "change": "task", "relationship": "verify"}); err != nil { + t.Fatal(err) + } + if err := cmdUndispatch("task", "verify"); err != nil { + t.Fatal(err) + } + if _, err := os.Stat(path); !os.IsNotExist(err) { + t.Fatal("sibling not removed") + } + if !bytes.Equal(before, readBytes(t, requestFile(repo))) { + t.Fatal("request changed") + } +} + +func TestStatusExemptsRevokeRecipient(t *testing.T) { + repo, sid := requestFixture(t) + if err := CmdRequest("task", "one", sid, "lead", "fix it"); err != nil { + t.Fatal(err) + } + key := "repo:" + fleet.RepoID(repo) + ":task" + if err := fleet.WriteJSON(fleet.KeyFile("stop", key), fleet.Rec{"key": key, "except": sid}); err != nil { + t.Fatal(err) + } + if row := requestStatus(fleet.ReadJSON(requestFile(repo)), fleet.Now()); row["status"] != "Queued" { + t.Fatalf("recipient falsely stopped: %v", row) + } +} + +func TestLegacyMaintenanceIgnoresUnrelatedDamage(t *testing.T) { + repo, _ := requestFixture(t) + rid := fleet.RepoID(repo) + path := dispatchFile(rid, "task", "verify") + if err := fleet.WriteJSON(path, fleet.Rec{"repo": rid, "change": "task", "relationship": "verify"}); err != nil { + t.Fatal(err) + } + if err := os.WriteFile(filepath.Join(dispatchDir(), "unrelated.json"), []byte("broken"), 0600); err != nil { + t.Fatal(err) + } + if err := CmdReassign("task", "new-lead"); err != nil { + t.Fatal(err) + } + if fleet.S(fleet.ReadJSON(path), "for") != "new-lead" { + t.Fatal("not reassigned") + } + if err := cmdUndispatch("task", "verify"); err != nil { + t.Fatal(err) + } + if err := os.WriteFile(path, []byte("broken"), 0600); err != nil { + t.Fatal(err) + } + if err := CmdReassign("task", "new-lead"); err == nil { + t.Fatal("damaged target accepted") + } + if err := cmdUndispatch("task", "verify"); err == nil { + t.Fatal("damaged target removed") + } +} + +func TestActivitySurvivesWritesOnAnotherBranch(t *testing.T) { + repo, sid := requestFixture(t) + if err := CmdRequest("task", "one", sid, "lead", "fix it"); err != nil { + t.Fatal(err) + } + ev := fleet.Event{"session_id": sid, "cwd": repo, "hook_event_name": "PostToolUse", "tool_name": "Edit", "tool_input": fleet.Rec{"file_path": filepath.Join(repo, "file")}} + if v := fleet.Run(ev); v.Code != 0 { + t.Fatal(v) + } + runGit(t, repo, "checkout", "-b", "other") + if v := fleet.Run(ev); v.Code != 0 { + t.Fatal(v) + } + if row := requestStatus(fleet.ReadJSON(requestFile(repo)), fleet.Now()); row["status"] != "Activity observed" { + t.Fatalf("prior observation lost: %v", row) + } +} + +func TestRequestKeysByGitSpellingNotCallerSpelling(t *testing.T) { + repo, sid := requestFixture(t) + runGit(t, repo, "branch", "Nav-Fix") + if err := CmdRequest("nav-fix", "case-1", sid, "lead", "fix it"); err != nil { + t.Fatal(err) + } + if _, err := os.Stat(dispatchFile(fleet.RepoID(repo), "Nav-Fix", "implementation")); err != nil { + t.Fatalf("assignment keyed by the caller's spelling, not git's: %v", err) + } + row := fleet.ReadJSON(dispatchFile(fleet.RepoID(repo), "Nav-Fix", "implementation")) + if fleet.S(row, "change") != "Nav-Fix" { + t.Fatalf("change recorded as %q, want git's spelling", fleet.S(row, "change")) + } + // A second spelling of the same branch is the same assignment, not a second one. + if err := CmdRequest("NAV-FIX", "case-2", sid, "lead", "fix it"); err == nil { + t.Fatal("second assignment for one branch under another spelling was accepted") + } +} + +func TestStatusRefusesDamagedRequestEvidence(t *testing.T) { + repo, sid := requestFixture(t) + if err := CmdRequest("task", "one", sid, "lead", "fix it"); err != nil { + t.Fatal(err) + } + row := fleet.ReadJSON(requestFile(repo)) + for _, damage := range []string{"worker", "at", "for"} { + broken := fleet.Rec{} + for k, v := range row { + broken[k] = v + } + delete(broken, damage) + if err := fleet.WriteJSON(requestFile(repo), broken); err != nil { + t.Fatal(err) + } + var b bytes.Buffer + Out = &b + err := Dispatch([]string{"status", "--json"}) + var packet map[string]any + _ = json.Unmarshal(b.Bytes(), &packet) + if err == nil || packet["complete"] != false { + t.Fatalf("row missing %q reported complete: err=%v packet=%v", damage, err, packet) + } + if err := CmdRequest("task", "two", sid, "lead", "fix it"); err == nil { + t.Fatalf("request accepted beside damaged evidence (missing %q)", damage) + } + } +} + +func TestRequestReplaySurvivesBranchDeletionUnderCallerSpelling(t *testing.T) { + repo, sid := requestFixture(t) + runGit(t, repo, "branch", "Nav-Fix") + if err := CmdRequest("nav-fix", "case-3", sid, "lead", "fix it"); err != nil { + t.Fatal(err) + } + runGit(t, repo, "branch", "-D", "Nav-Fix") + if err := CmdRequest("nav-fix", "case-3", sid, "lead", "fix it"); err != nil { + t.Fatalf("identical retry after branch deletion refused: %v", err) + } +} + +func TestStatusRefusesNonStringRequestID(t *testing.T) { + repo, sid := requestFixture(t) + if err := CmdRequest("task", "one", sid, "lead", "fix it"); err != nil { + t.Fatal(err) + } + row := fleet.ReadJSON(requestFile(repo)) + row["request_id"] = 7.0 + if err := fleet.WriteJSON(requestFile(repo), row); err != nil { + t.Fatal(err) + } + var b bytes.Buffer + Out = &b + err := Dispatch([]string{"status", "--json"}) + var packet map[string]any + _ = json.Unmarshal(b.Bytes(), &packet) + if err == nil || packet["complete"] != false { + t.Fatalf("non-string request_id read as complete: err=%v packet=%v", err, packet) + } +} + +func TestReplayNeverMatchesADistinctRefByCase(t *testing.T) { + // The judge's counterexample: on a case-sensitive filesystem Nav-Fix and + // nav-fix are two branches. An ID recorded for one must not replay against + // the other. Exercised on sameRequest directly so it holds on every platform. + stored := fleet.Rec{"request_id": "x", "repo": "r", "change": "Nav-Fix", "requested": "Nav-Fix", "worker": "w", "for": "l", "brief": "b", "relationship": "implementation"} + other := fleet.Rec{"request_id": "x", "repo": "r", "change": "nav-fix", "requested": "nav-fix", "worker": "w", "for": "l", "brief": "b", "relationship": "implementation"} + if sameRequest(stored, other, true) { + t.Fatal("an ID for Nav-Fix replayed against the distinct branch nav-fix") + } + // The deleted-branch retry: the caller's spelling is compared exactly. + retry := fleet.Rec{"request_id": "x", "repo": "r", "change": "nav-fix", "requested": "nav-fix", "worker": "w", "for": "l", "brief": "b", "relationship": "implementation"} + recorded := fleet.Rec{"request_id": "x", "repo": "r", "change": "Nav-Fix", "requested": "nav-fix", "worker": "w", "for": "l", "brief": "b", "relationship": "implementation"} + if !sameRequest(recorded, retry, false) { + t.Fatal("identical retry after branch deletion refused") + } + if sameRequest(recorded, fleet.Rec{"request_id": "x", "repo": "r", "change": "Nav-Fix", "requested": "Nav-Fix", "worker": "w", "for": "l", "brief": "b", "relationship": "implementation"}, false) { + t.Fatal("a differently spelled retry after deletion was accepted as the same work") + } +} diff --git a/cmd/fleet/internal/verbs/role.go b/cmd/fleet/internal/verbs/role.go index f701e076..8577d90c 100644 --- a/cmd/fleet/internal/verbs/role.go +++ b/cmd/fleet/internal/verbs/role.go @@ -137,7 +137,7 @@ func codexHooks(command string) (string, map[string]any, error) { hooks = map[string]any{} data["hooks"] = hooks } - specs := [][2]string{{"SessionStart", ""}, {"UserPromptSubmit", ""}, {"PreToolUse", "^(Bash|Edit|Write)$"}, {"PostToolUse", "^Bash$"}, {"Stop", ""}, {"SessionEnd", ""}} + specs := [][2]string{{"SessionStart", ""}, {"UserPromptSubmit", ""}, {"PreToolUse", "^(Bash|Edit|Write|MultiEdit|NotebookEdit|apply_patch)$"}, {"PostToolUse", "^(Bash|Edit|Write|MultiEdit|NotebookEdit|apply_patch)$"}, {"Stop", ""}, {"SessionEnd", ""}} for _, spec := range specs { event, matcher := spec[0], spec[1] groups, err := withoutFleetHandlers(hooks[event], target, event) @@ -153,6 +153,28 @@ func codexHooks(command string) (string, map[string]any, error) { return target, data, nil } +// claudeWriteHooks supplements the global Bash hook with completed file writes. +// Keep unrelated local hooks intact and replace our marked group on rebind. +func claudeWriteHooks(data map[string]any, target, command string) error { + hooks, ok := data["hooks"].(map[string]any) + if data["hooks"] != nil && !ok { + return refuse("fleet role: %s hooks is not an object", target) + } + if hooks == nil { + hooks = map[string]any{} + data["hooks"] = hooks + } + groups, err := withoutFleetHandlers(hooks["PostToolUse"], target, "PostToolUse") + if err != nil { + return err + } + hooks["PostToolUse"] = append(groups, map[string]any{ + "matcher": "^(Edit|Write|MultiEdit|NotebookEdit)$", + "hooks": []any{map[string]any{"type": "command", "command": command, "statusMessage": fleetHookMark}}, + }) + return nil +} + // codexRules is the Codex projection of a lane's denies: every `Bash(:*)` // deny becomes an execpolicy prefix rule. Denies with a wildcard mid-pattern have no // prefix form and stay Claude-only. @@ -382,6 +404,9 @@ func cmdRole(checkout, role string, force bool, tenant, slot string) error { if err != nil { return err } + if err := claudeWriteHooks(existing, settingsTarget, strings.TrimSuffix(hookCommand(), " codex")); err != nil { + return err + } if err := os.MkdirAll(fleet.OrgState, 0o755); err != nil { return err } diff --git a/cmd/fleet/internal/verbs/role_hooks_test.go b/cmd/fleet/internal/verbs/role_hooks_test.go new file mode 100644 index 00000000..b1d97a0f --- /dev/null +++ b/cmd/fleet/internal/verbs/role_hooks_test.go @@ -0,0 +1,54 @@ +package verbs + +import ( + "regexp" + "testing" +) + +func TestGeneratedCodexHooksSubscribeToWrites(t *testing.T) { + t.Setenv("CODEX_HOME", t.TempDir()) + _, data, err := codexHooks("fleet hook codex") + if err != nil { + t.Fatal(err) + } + hooks := data["hooks"].(map[string]any) + for _, event := range []string{"PreToolUse", "PostToolUse"} { + groups := hooks[event].([]any) + group := groups[len(groups)-1].(map[string]any) + matcher := regexp.MustCompile(group["matcher"].(string)) + for _, tool := range []string{"Bash", "Edit", "Write", "MultiEdit", "NotebookEdit", "apply_patch"} { + if !matcher.MatchString(tool) { + t.Errorf("%s drops %s", event, tool) + } + } + if matcher.MatchString("Read") { + t.Error("read subscribed as write") + } + } +} + +func TestClaudeWriteHooksPreserveOtherHandlersOnRebind(t *testing.T) { + other := map[string]any{"matcher": "Read", "hooks": []any{map[string]any{"type": "command", "command": "other"}}} + data := map[string]any{"hooks": map[string]any{"PostToolUse": []any{other}}} + for range 2 { + if err := claudeWriteHooks(data, "fixture", "fleet hook"); err != nil { + t.Fatal(err) + } + } + groups := data["hooks"].(map[string]any)["PostToolUse"].([]any) + if len(groups) != 2 { + t.Fatalf("lost or duplicated handlers: %v", groups) + } + matcher := regexp.MustCompile(groups[1].(map[string]any)["matcher"].(string)) + for _, tool := range []string{"Edit", "Write", "MultiEdit", "NotebookEdit"} { + if !matcher.MatchString(tool) { + t.Errorf("drops %s", tool) + } + } + if matcher.MatchString("Bash") { + t.Fatal("duplicates global Bash hook") + } + if groups[0].(map[string]any)["matcher"] != "Read" { + t.Fatal("unrelated hook changed") + } +} diff --git a/cmd/fleet/internal/verbs/status.go b/cmd/fleet/internal/verbs/status.go new file mode 100644 index 00000000..aafd458d --- /dev/null +++ b/cmd/fleet/internal/verbs/status.go @@ -0,0 +1,113 @@ +package verbs + +import ( + "os" + "strings" + + "github.com/itsHabib/workbench/cmd/fleet/internal/fleet" +) + +// CmdStatus is a read-only, local request view. It neither infers semantic worker +// acceptance nor claims process termination from a stop flag. +func CmdStatus(args []string) error { + if len(args) > 1 || (len(args) == 1 && args[0] != "--json") { + return refuse("usage: fleet status [--json]") + } + // CLI and MCP execute verbs serially; this scoped flag is not goroutine-local. + before := fleet.ReadOnly + fleet.ReadOnly = true + defer func() { fleet.ReadOnly = before }() + pending := statusMigrationPending() + rows, err := strictDispatchRows() + out := []fleet.Rec{} + gaps := []string{} + if err != nil { + gaps = append(gaps, err.Error()) + } + now := fleet.Now() + for _, row := range rows { + if fleet.S(row, "request_id") != "" { + out = append(out, requestStatus(row, now)) + } + } + if pending || statusMigrationPending() { + gaps = append(gaps, "legacy state needs reconciliation") + } + packet := fleet.Rec{"schema": "fleet-task-status-v1", "at": now, "complete": len(gaps) == 0, + "scope": "local request-bound assignments only", "tasks": out, "gaps": gaps} + if len(args) == 1 { + say("%s", jsonIndent(packet)) + } + if len(args) == 0 { + renderStatus(out, gaps) + } + if len(gaps) > 0 { + return refuse("status needs checking: assignment evidence incomplete") + } + return nil +} + +func requestStatus(d fleet.Rec, now float64) fleet.Rec { + sid, rid, branch := fleet.S(d, "worker"), fleet.S(d, "repo"), fleet.S(d, "change") + key := "repo:" + rid + ":" + branch + row := fleet.Rec{"request_id": d["request_id"], "work": branch, "repo": rid, + "worker": sid, "lead": d["for"], "brief": d["brief"], "status": "Queued", + "needs": "Worker delivery and acceptance are unconfirmed", "next": "Deliver the brief through the worker's supported harness", "verified_at": now} + rec, lease := fleet.SessionRecord(sid), fleet.Lease(key) + write := fleet.M(fleet.M(rec, "last_writes"), key) + if write == nil { + write = fleet.M(rec, "last_write") + } + switch { + case rec == nil || fleet.IsMalformed(lease): + taskState(row, "Status needs checking", "Worker or ownership evidence is unavailable", "Inspect the evidence; do not redispatch") + case lease != nil && fleet.S(lease, "session") != sid: + taskState(row, "Status needs checking", "Another session holds the branch", "Resolve the ownership conflict") + case fleet.CheckStop(key, branch, sid) != "": + taskState(row, "Status needs checking", "A stop flag exists; process termination is unconfirmed", "Inspect the affected session before continuing") + case !fleet.SessionAlive(rec): + taskState(row, "Status needs checking", "Worker is no longer observably live", "Recover its work before choosing a replacement") + case fleet.S(write, "key") == key && fleet.F(write, "at") >= fleet.F(d, "at"): + row["activity_at"] = write["at"] + taskState(row, "Activity observed", "Acceptance and completion remain unconfirmed", "Read the worker's result and current checks") + } + return row +} + +func taskState(row fleet.Rec, state, needs, next string) { + row["status"], row["needs"], row["next"] = state, needs, next +} + +func renderStatus(rows []fleet.Rec, gaps []string) { + if len(gaps) > 0 { + say("Status needs checking: %s", strings.Join(gaps, "; ")) + } + say("Local assignments — recording work is not delivery or acceptance.") + if len(rows) == 0 { + say("No request-bound work recorded.") + return + } + for _, row := range rows { + if at := fleet.F(row, "activity_at"); at > 0 { + row["status"] = fleet.S(row, "status") + " (" + fleet.FmtAge(fleet.F(row, "verified_at")-at) + " ago)" + } + // IDs remain in JSON details. Do not guess a model/person name from a session. + say("%s — %s\n Needs: %s\n Next: %s", terminalText(fleet.S(row, "work")), fleet.S(row, "status"), terminalText(fleet.S(row, "needs")), terminalText(fleet.S(row, "next"))) + } +} + +func terminalText(s string) string { + return strings.Map(func(r rune) rune { + if r < 32 || r == 127 { + return ' ' + } + return r + }, s) +} + +func statusMigrationPending() bool { + if _, err := os.Stat(fleet.State); os.IsNotExist(err) { + return false + } + return fleet.MigrationPending() +} diff --git a/cmd/fleet/internal/verbs/verbs.go b/cmd/fleet/internal/verbs/verbs.go index 420656ca..1defa2c6 100644 --- a/cmd/fleet/internal/verbs/verbs.go +++ b/cmd/fleet/internal/verbs/verbs.go @@ -79,6 +79,9 @@ revoke / handoff act on the repo you are standing in. ` + "`main`" + ` in two re fleet reassign --for move a change's rows to another accountable role (splitting a hub is this plus one roles.map line) fleet undispatch [--as ] retire a change's rows fleet sync [--repo ] refresh the cache of open changes and the rows other machines declared on them + fleet request --id --worker --for --brief + record one retry-safe local assignment; does not launch a worker + fleet status [--json] read-only request board; queued is not accepted or running fleet inspect-hooks --config read-only static hook inventory; no execution or migration fleet report [--since 24h | --snapshot] derived telemetry or JSON observations, without writing state fleet shadow-report [--since 24h] [--json] the day's numbers from 'fleet hook --shadow' running beside the installed hook @@ -122,6 +125,12 @@ func Dispatch(args []string) error { if len(args) == 0 { return exitCode(2, usage) } + if args[0] == "status" { + return CmdStatus(args[1:]) + } + if args[0] == "request" { + return dispatchRequest(args[1:]) + } if args[0] == "report" { return cmdReport(args[1:]) } diff --git a/cmd/fleet/internal/verbs/work.go b/cmd/fleet/internal/verbs/work.go index a51fbe6d..3b049533 100644 --- a/cmd/fleet/internal/verbs/work.go +++ b/cmd/fleet/internal/verbs/work.go @@ -97,6 +97,12 @@ func liveHands(key string) string { // slot is named. Placement is `assign`, under the slot's lock, before the row is // written, so a refused placement leaves no row behind. func CmdDispatch(change, rel, forRole, due, slot, brief, by, replyTo string, take bool) error { + return fleet.KeyLock("dispatch", func() error { + return cmdDispatch(change, rel, forRole, due, slot, brief, by, replyTo, take) + }) +} + +func cmdDispatch(change, rel, forRole, due, slot, brief, by, replyTo string, take bool) error { usage := `usage: fleet dispatch --as [--for ] [--due 45m] [--slot ] [--brief ""] [--reply-to ] [--take]` if change == "" || rel == "" { return refuse("%s", usage) @@ -123,7 +129,10 @@ func CmdDispatch(change, rel, forRole, due, slot, brief, by, replyTo string, tak forRole = by } key := "repo:" + rid + ":" + branch - existing := fleet.ReadJSON(dispatchFile(rid, branch, rel)) + existing, err := replaceableDispatch(rid, branch, rel) + if err != nil { + return err + } if hands := liveHands(key); hands != "" && existing != nil && !take { return refuse("fleet dispatch: %s/%s already has live hands (%s, for %s); `--take` rewrites the row, or `fleet work` to see it", branch, rel, fleet.Short(hands), fleet.S(existing, "for")) @@ -170,6 +179,10 @@ func dispatcher(by string) string { // CmdReassign moves every row of a change to another accountable role. Two // commands split a hub: bind the second role's directory, then reassign its rows. func CmdReassign(change, forRole string) error { + return fleet.KeyLock("dispatch", func() error { return cmdReassign(change, forRole) }) +} + +func cmdReassign(change, forRole string) error { if change == "" || forRole == "" { return refuse("usage: fleet reassign --for ") } @@ -177,8 +190,17 @@ func CmdReassign(change, forRole string) error { if err != nil { return err } + rows, err := scopedDispatchRows(rid, branch, "") + if err != nil { + return err + } + for _, r := range rows { + if fleet.S(r, "repo") == rid && fleet.S(r, "change") == branch && fleet.S(r, "request_id") != "" { + return refuse("fleet: request-bound assignments require a correlated lifecycle action; inspect `fleet status`") + } + } n := 0 - for _, r := range dispatchRows() { + for _, r := range rows { if fleet.S(r, "repo") != rid || fleet.S(r, "change") != branch { continue } @@ -199,6 +221,10 @@ func CmdReassign(change, forRole string) error { // cmdUndispatch retires a change's rows (one relationship, or all of them). func cmdUndispatch(change, rel string) error { + return fleet.KeyLock("dispatch", func() error { return undispatch(change, rel) }) +} + +func undispatch(change, rel string) error { if change == "" { return refuse("usage: fleet undispatch [--as ]") } @@ -206,8 +232,17 @@ func cmdUndispatch(change, rel string) error { if err != nil { return err } + rows, err := scopedDispatchRows(rid, branch, rel) + if err != nil { + return err + } + for _, r := range rows { + if fleet.S(r, "repo") == rid && fleet.S(r, "change") == branch && (rel == "" || fleet.S(r, "relationship") == rel) && fleet.S(r, "request_id") != "" { + return refuse("fleet: request-bound assignments require a correlated lifecycle action; inspect `fleet status`") + } + } n := 0 - for _, r := range dispatchRows() { + for _, r := range rows { if fleet.S(r, "repo") != rid || fleet.S(r, "change") != branch || (rel != "" && fleet.S(r, "relationship") != rel) { continue } @@ -477,3 +512,53 @@ func cmdWork(forRole string, asJSON bool) error { say("%s", scope) return nil } + +func replaceableDispatch(rid, branch, rel string) (fleet.Rec, error) { + path := dispatchFile(rid, branch, rel) + if _, err := os.Lstat(path); os.IsNotExist(err) { + return nil, nil + } + existing := fleet.ReadJSON(path) + if existing == nil || fleet.S(existing, "repo") != rid || fleet.S(existing, "change") != branch || fleet.S(existing, "relationship") != rel { + return nil, refuse("fleet dispatch: selected assignment is unreadable or has conflicting identity") + } + if fleet.S(existing, "request_id") != "" { + return nil, refuse("fleet dispatch: this assignment is request-bound; inspect `fleet status` rather than replacing it") + } + return existing, nil +} + +// scopedDispatchRows preserves unrelated legacy maintenance while refusing every +// damaged target before any mutation. Filename candidates catch unreadable rows; +// decoded identity also catches selected rows stored under an unexpected name. +func scopedDispatchRows(rid, branch, rel string) ([]fleet.Rec, error) { + entries, err := os.ReadDir(dispatchDir()) + if os.IsNotExist(err) { + return nil, nil + } + if err != nil { + return nil, err + } + var rows []fleet.Rec + prefix := fleet.Safe(rid + "__" + branch + "__") + for _, entry := range entries { + if !strings.HasSuffix(entry.Name(), ".json") { + continue + } + path := filepath.Join(dispatchDir(), entry.Name()) + candidate := strings.HasPrefix(entry.Name(), prefix) + if rel != "" { + candidate = path == dispatchFile(rid, branch, rel) + } + row := fleet.ReadJSON(path) + selected := fleet.S(row, "repo") == rid && fleet.S(row, "change") == branch && (rel == "" || fleet.S(row, "relationship") == rel) + if !candidate && !selected { + continue + } + if !selected || fleet.S(row, "relationship") == "" || path != dispatchFile(rid, branch, fleet.S(row, "relationship")) { + return nil, refuse("fleet: selected assignment is unreadable or has conflicting identity") + } + rows = append(rows, row) + } + return rows, nil +} diff --git a/cmd/fleet/testdata/test.sh b/cmd/fleet/testdata/test.sh index c9f19426..ac7dd2f8 100755 --- a/cmd/fleet/testdata/test.sh +++ b/cmd/fleet/testdata/test.sh @@ -1124,7 +1124,7 @@ printf '%s\n' '{"jsonrpc":"2.0","id":1,"method":"initialize","params":{"protocol '{"jsonrpc":"2.0","id":7,"method":"tools/call","params":{"name":"fleet_take","arguments":{"resource":"slot:hyper"}}}' \ "{\"jsonrpc\":\"2.0\",\"id\":8,\"method\":\"tools/call\",\"params\":{\"name\":\"fleet_who\",\"arguments\":{\"name\":\"feat/q\",\"cwd\":\"$SL1\"}}}" \ | (cd "$work" && "$PY" "$here/fleet-mcp.py") > "$work/mcp.out" 2>"$work/mcp.err" -"$PY" - "$work/mcp.out" "$S20" <<'PY' && echo " ok fleet-mcp: initialize, 11 one-line tools, who resolves in the caller's cwd, refusal is isError with the CLI's text, not-done is an answer, missing arg is -32602, acting tools need cwd" || { echo " FAIL fleet-mcp: $(cat "$work/mcp.out" "$work/mcp.err")"; fails=$((fails+1)); } +"$PY" - "$work/mcp.out" "$S20" <<'PY' && echo " ok fleet-mcp: initialize, 13 one-line tools, who resolves in the caller's cwd, refusal is isError with the CLI's text, not-done is an answer, missing arg is -32602, acting tools need cwd" || { echo " FAIL fleet-mcp: $(cat "$work/mcp.out" "$work/mcp.err")"; fails=$((fails+1)); } import json, sys by = {} for line in open(sys.argv[1], encoding="utf-8"): @@ -1132,7 +1132,7 @@ for line in open(sys.argv[1], encoding="utf-8"): bad = [] if by[1]["result"]["serverInfo"]["name"] != "fleet": bad.append("initialize") tools = by[2]["result"]["tools"] -if len(tools) != 11 or any("\n" in t["description"] or len(t["description"]) > 160 for t in tools): bad.append("tools: %d, long or multi-line description" % len(tools)) +if len(tools) != 13 or any("\n" in t["description"] or len(t["description"]) > 160 for t in tools): bad.append("tools: %d, long or multi-line description" % len(tools)) if sys.argv[2] not in by[3]["result"]["content"][0]["text"] or by[3]["result"]["isError"]: bad.append("who") if not by[4]["result"]["isError"] or "busy" not in by[4]["result"]["content"][0]["text"]: bad.append("assign refusal") if by[5]["result"]["isError"] or '"ok": false' not in by[5]["result"]["content"][0]["text"]: bad.append("done") @@ -1889,7 +1889,7 @@ try: except subprocess.TimeoutExpired as e: b_rc, b_out = None, "second watcher kept running for 15s: " + str(e) pA.kill(); pA.wait() -report(b_rc not in (None, 0) and "already" in b_out, "a second watcher is refused for the first's whole lifetime, not only after its first heartbeat", f"B rc={b_rc} {b_out[:160]!r}") +report(b_rc == 1 and "watcher lock unavailable" in b_out and "already running" not in b_out, "a second watcher is refused for the first's whole lifetime, not only after its first heartbeat", f"B rc={b_rc} {b_out[:160]!r}") # 8. the last word in Codex's transcript shape, and the harness's own statement of it tp = os.path.join(work, "codex-transcript.jsonl") open(tp, "w").write(json.dumps({"type": "response_item", "payload": {"type": "message", "role": "assistant", "content": [{"type": "output_text", "text": "done: pushed the fix"}]}}) + "\n") diff --git a/friction-log.md b/friction-log.md index 69be2b35..ef153954 100644 --- a/friction-log.md +++ b/friction-log.md @@ -360,6 +360,38 @@ Second occurrence of the `#214` entry above, one failure mode further in. own comment so the attestation fires, and put the focus areas in a second comment. Costs nothing and keeps the panel complete. +## 2026-09-08 — Fleet assignment is not worker acceptance + +- **Observed:** dispatch stores accountability but exposes no retry identity; a + retry can rewrite queued work. Existing CLI helpers do not establish end-to-end + cross-harness acceptance or effect-safe stop. An operator should not infer those + from a successful command or copied session ID. +- **Change:** request-bound rows, immutable payload checks, cross-process dispatch + serialization and non-migrating status output. Hook-owned post-tool evidence is + explicitly activity, not success or semantic acceptance. Legacy mutations cannot + replace a request record. No second editable ledger added. +- **Validation boundary:** fixture and real-process tests prove these local + contracts, not live model delivery/replacement. This is the first build increment; + actual adapters and correlated lifecycle remain under the natural-coordination + Dossier task. No hooks installed or live agents controlled by this change. +- **Tooling:** full-module tests were cost-guarded; followed the requested focused + Fleet test path and left the full suite to CI. Root vet/lint and both harness + regression suites were run. +- **Review integration gaps:** generated Codex hooks only subscribed to Bash + post-tool events; direct hook tests hid missing file-edit observations. Generated + subscriptions now cover writes, with Claude local file-write supplementation + and rebind tests. Existing installations still require regeneration/reload. + Replay now survives branch/session cleanup, selective legacy retirement preserves + request siblings, and status honors the revoke recipient exemption. These are + regression-tested without taking over any live session. +- **Second review:** one latest-write field let activity on another branch erase + observed task progress; observations now merge per branch under the existing + session lock. Added MultiEdit to generated subscriptions and scoped legacy + maintenance validation. Two reviewers reported raw apply_patch bypassing the + classifier; the actual Codex adapter already expands patches into per-file Edit + events. A direct adapter regression now proves foreign-holder refusal and + post-tool evidence, without duplicating parsing in the policy layer. + ## 2026-09-09 — watcher diagnostics turn unknown evidence into liveness claims Work log a3f579e reported a zero heartbeat as 56 years old and a stale-board