diff --git a/BURNIN.md b/BURNIN.md index b4af557..9f445e7 100644 --- a/BURNIN.md +++ b/BURNIN.md @@ -2587,3 +2587,23 @@ file. The ring earned its keep again: one distinct message with a count of 34, rather than 34 overwrites of one slot. + +**Fixed, not yet proven live.** The request is now stamped +(`herdr.Session.HandoffRequestedAt`) and the wait is bounded by +`watchHandoff` in the worker: + +- the lease renews while Orchestra is explicitly waiting, because a quiet pane + is the answer the agent was asked for โ€” the renewal gate's new case, bounded + by `handoffAnswerTimeout`; +- the request is re-sent once at `handoffRetryAfter` (4 minutes), with the same + reason it was first asked with; +- at 10 minutes the worker nacks with failure class `handoff_unanswered`, and + the coordinator emits `TaskReleased reason=handoff_unanswered` rather than + letting the lease die as generic idleness. + +`DebtClassForFailureClass` knows the class, so a harness that repeatedly +ignores handoff requests now accumulates in the debt ledger instead of hiding +inside `lease_expired`. + +Runtime proof still owed: force a request the agent will not answer, and read +the release event rather than the worker state file. diff --git a/cmd/orchestra-worker/handoff_test.go b/cmd/orchestra-worker/handoff_test.go new file mode 100644 index 0000000..48efc8f --- /dev/null +++ b/cmd/orchestra-worker/handoff_test.go @@ -0,0 +1,110 @@ +package main + +import ( + "context" + "encoding/json" + "net/http" + "net/http/httptest" + "path/filepath" + "strings" + "testing" + "time" + + "orchestra/internal/domain" + "orchestra/internal/federation" + "orchestra/internal/herdr" +) + +// F62. Run 16: the agent was asked to hand off, never wrote HANDOFF.md, +// renewals stopped, and the lease died as ordinary idleness. Waiting is now +// bounded: re-ask once, then give the task up with a class that says why. +func TestUnansweredHandoffIsRetriedThenGivenUp(t *testing.T) { + var nack map[string]any + w, backend, _, done := phaseWorker(t, func(rw http.ResponseWriter, r *http.Request) { + if strings.HasSuffix(r.URL.Path, "/nack") { + _ = json.NewDecoder(r.Body).Decode(&nack) + } + rw.Write([]byte(`{}`)) + }) + defer done() + ctx := context.Background() + requested := func(ago time.Duration) herdr.Session { + s := w.sessions["task"] + s.HandoffRequested, s.HandoffReason = true, "phase_changed" + s.HandoffRequestedAt = time.Now().UTC().Add(-ago) + w.sessions["task"] = s + return s + } + + // Still inside the answering window: nothing said, nothing given up. + if s, gaveUp := w.watchHandoff(ctx, "task", requested(time.Minute)); gaveUp || s.HandoffRetried { + t.Fatalf("gave up while still waiting: gaveUp=%v session=%+v", gaveUp, s) + } + if len(backend.prompts) != 0 { + t.Fatalf("re-asked too early: %q", backend.prompts) + } + + // Past the retry point: asked again, exactly once. + s, gaveUp := w.watchHandoff(ctx, "task", requested(handoffRetryAfter+time.Minute)) + if gaveUp || !s.HandoffRetried || len(backend.prompts) != 1 { + t.Fatalf("retry: gaveUp=%v retried=%v prompts=%q", gaveUp, s.HandoffRetried, backend.prompts) + } + if _, gaveUp = w.watchHandoff(ctx, "task", s); gaveUp || len(backend.prompts) != 1 { + t.Fatalf("re-asked every tick: %q", backend.prompts) + } + + // Past the bound: a causal reclaim, and no lease left to renew. + s = requested(handoffAnswerTimeout + time.Second) + s.HandoffRetried = true + w.sessions["task"] = s + if _, gaveUp = w.watchHandoff(ctx, "task", s); !gaveUp { + t.Fatal("an unanswered handoff waited forever") + } + if nack["failure_class"] != "handoff_unanswered" { + t.Fatalf("nack = %+v", nack) + } + if detail, _ := nack["last_error"].(string); !strings.Contains(detail, "phase_changed") { + t.Fatalf("the reclaim does not name the request: %q", detail) + } + if _, held := w.leases["task"]; held { + t.Fatal("the given-up task kept its lease") + } +} + +// The lease must survive the wait it was asked to make: an idle pane is the +// answer Orchestra requested, not evidence of an agent that stopped working. +func TestWaitingForAHandoffKeepsTheLease(t *testing.T) { + renewals := 0 + api := httptest.NewServer(http.HandlerFunc(func(rw http.ResponseWriter, r *http.Request) { + renewals++ + rw.Write([]byte(`{}`)) + })) + defer api.Close() + backend := &recordingBackend{status: "idle", progress: "same screen"} + w := &worker{ + api: federation.Client{BaseURL: api.URL, WorkerID: "h", Token: "t"}, + backend: backend, + harness: "claude", + sessions: map[string]herdr.Session{"task": {PaneID: "pane", HandoffRequested: true, HandoffRequestedAt: time.Now().UTC()}}, + leases: map[string]lease{"task": {Epoch: "e", Version: 1, Until: time.Now(), ProgressSHA: domain.Hash([]byte("same screen"))}}, + quarantined: map[string]bool{}, + statePath: filepath.Join(t.TempDir(), "state.json"), + } + w.renewLeases(context.Background()) + if renewals != 1 { + t.Fatalf("a lease waiting on a requested handoff renewed %d times, want 1", renewals) + } + + // Past the bound the exemption stops: watchHandoff has given the task up + // by then, and nothing keeps an unanswered request alive. + s := w.sessions["task"] + s.HandoffRequestedAt = time.Now().UTC().Add(-handoffAnswerTimeout - time.Second) + w.sessions["task"] = s + l := w.leases["task"] + l.Until = time.Now() + w.leases["task"] = l + w.renewLeases(context.Background()) + if renewals != 1 { + t.Fatalf("the exemption outlived its bound: renewals=%d", renewals) + } +} diff --git a/cmd/orchestra-worker/main.go b/cmd/orchestra-worker/main.go index adddae1..bd9389a 100644 --- a/cmd/orchestra-worker/main.go +++ b/cmd/orchestra-worker/main.go @@ -174,6 +174,7 @@ func releaseBackoff(attempts int) time.Duration { } return d } + type projectConfig struct { Repo string `json:"repo"` Root string `json:"worktree_root"` @@ -673,6 +674,13 @@ func (w *worker) releaseReady(ctx context.Context) { w.advanceRelease(ctx, id, s) continue } + if s.HandoffRequested { + next, gaveUp := w.watchHandoff(ctx, id, s) + if gaveUp { + continue + } + s = next + } w.rotationTick(ctx, id, s) } } @@ -772,7 +780,7 @@ func (w *worker) rotationTick(ctx context.Context, id string, s herdr.Session) { w.recordError(fmt.Errorf("rotation %s threshold prompt: %w", id, err)) return } - s.HandoffRequested, s.HandoffReason = true, d.Reason + s.HandoffRequested, s.HandoffReason, s.HandoffRequestedAt = true, d.Reason, time.Now().UTC() w.sessions[id] = s _ = w.save() } @@ -1166,6 +1174,70 @@ func (w *worker) submit(ctx context.Context, id string, s herdr.Session, e compl // the directory also holds for a worktree that has no inner .gitignore. const stageExclude = ":!.orchestra" +// F62: a requested handoff nobody answers was invisible. Renewals stopped, +// the lease expired, and the task lost an attempt with nothing on record +// saying a handoff had ever been asked for โ€” worker health showed only "agent +// status idle and pane unchanged", 34 times in run 16. +const ( + handoffRetryAfter = 4 * time.Minute + handoffAnswerTimeout = 10 * time.Minute +) + +// watchHandoff bounds the wait for an agent's handoff answer: re-send the +// request once, then give the task up with a class that names the cause. It +// returns the session to keep using and whether the task was given up. +func (w *worker) watchHandoff(ctx context.Context, id string, s herdr.Session) (herdr.Session, bool) { + if s.HandoffRequestedAt.IsZero() { + // A session persisted before the stamp existed, or requested by a path + // that does not set it. Start the clock now rather than time out a + // request retroactively. + s.HandoffRequestedAt = time.Now().UTC() + w.sessions[id] = s + _ = w.save() + return s, false + } + waited := time.Since(s.HandoffRequestedAt) + if waited < handoffRetryAfter { + return s, false + } + if waited < handoffAnswerTimeout { + if s.HandoffRetried || w.executionBackend() == nil { + return s, false + } + a := herdr.CLIAdapter{Backend: w.executionBackend(), Harness: w.harness} + var err error + if s.HandoffReason != "" { + err = a.RequestHandoffReason(ctx, s, s.HandoffReason, nil) + } else { + err = a.RequestHandoff(ctx, s) + } + if err != nil { + w.recordError(fmt.Errorf("handoff %s re-request: %w", id, err)) + return s, false + } + s.HandoffRetried = true + w.sessions[id] = s + _ = w.save() + w.recordError(fmt.Errorf("handoff %s (%s) unanswered for %s: request re-sent", id, s.HandoffReason, waited.Round(time.Second))) + return s, false + } + l, ok := w.leases[id] + if !ok { + return s, false + } + detail := fmt.Sprintf("handoff requested (%s) and unanswered for %s", s.HandoffReason, waited.Round(time.Second)) + if err := w.api.NackStart(ctx, id, l.Epoch, l.Version, "handoff_unanswered", detail, w.sessionEvidence(ctx, id, s)); err != nil { + w.recordError(fmt.Errorf("handoff timeout %s: %w", id, err)) + return s, false + } + // The coordinator answers with TaskReleased; its replay quarantines the + // pane. Drop the lease here so nothing renews it in the meantime. + delete(w.leases, id) + _ = w.save() + w.recordError(errors.New(detail)) + return s, true +} + func (w *worker) renewLeases(ctx context.Context) { if w.executionBackend() == nil { return @@ -1199,6 +1271,11 @@ func (w *worker) renewLeases(ctx context.Context) { case l.ProgressSHA == "": // First renewal has no baseline to compare against. Record one and // allow this renewal; the next one must show real movement. + case s.HandoffRequested && !s.HandoffRequestedAt.IsZero() && time.Since(s.HandoffRequestedAt) < handoffAnswerTimeout: + // Orchestra told this agent to stop and write its handoff. A quiet + // pane is the answer it was asked for, so the lease is held while + // the wait is explicitly bounded (F62). Past the bound the case + // stops matching and watchHandoff has already given the task up. default: w.recordError(fmt.Errorf("lease %s not renewed: agent status %s and pane unchanged since the last renewal", taskID, status)) continue @@ -1932,7 +2009,7 @@ func (w *worker) federatedTurn(ctx context.Context, id string, a herdr.Adapter, w.recordError(fmt.Errorf("reconcile failure handoff %s: %w", id, err)) return } - s.HandoffRequested, s.HandoffReason = true, "reconcile_failure" + s.HandoffRequested, s.HandoffReason, s.HandoffRequestedAt = true, "reconcile_failure", time.Now().UTC() w.sessions[id] = s _ = w.save() return @@ -2009,7 +2086,7 @@ func (w *worker) rotateForPhase(ctx context.Context, id string, a herdr.Adapter, w.recordError(fmt.Errorf("phase rotation %s: %w", id, err)) return } - s.HandoffRequested, s.HandoffReason = true, "phase_changed" + s.HandoffRequested, s.HandoffReason, s.HandoffRequestedAt = true, "phase_changed", time.Now().UTC() w.sessions[id] = s _ = w.save() log.Printf("phase changed for %s: session rotating", id) diff --git a/cmd/orchestra/main.go b/cmd/orchestra/main.go index 118ad07..e1b1c41 100644 --- a/cmd/orchestra/main.go +++ b/cmd/orchestra/main.go @@ -1589,6 +1589,12 @@ func main() { case "invalid_handoff": typ = "TaskBlocked" p, _ = json.Marshal(map[string]any{"blocker": b.LastError, "block_reason": string(domain.BlockReasonHandoffValidation), "harness_id": parts[3], "lease_epoch": b.LeaseEpoch, "expected_version": t.Version, "lifecycle_phase": "launch_nacked", "last_error": b.LastError, "session_evidence": b.SessionEvidence}) + case "handoff_unanswered": + // F62: not a launch failure. The agent was asked to hand off + // and never did, so the reclaim says exactly that instead of + // arriving as an ordinary idle expiry. + typ = "TaskReleased" + p, _ = json.Marshal(map[string]any{"reason": "handoff_unanswered", "failure_class": b.FailureClass, "harness_id": parts[3], "lease_epoch": b.LeaseEpoch, "expected_version": t.Version, "lifecycle_phase": "handoff_unanswered", "last_error": b.LastError, "session_evidence": b.SessionEvidence}) case "launch_uncertain": typ = "TaskNeedsAttention" p, _ = json.Marshal(map[string]any{"blocker": b.LastError, "block_reason": string(domain.BlockReasonLeaseFailure), "harness_id": parts[3], "lease_epoch": b.LeaseEpoch, "expected_version": t.Version, "lifecycle_phase": "launch_uncertain", "last_error": b.LastError, "session_evidence": b.SessionEvidence}) diff --git a/internal/domain/debt.go b/internal/domain/debt.go index 49be075..8a6ffca 100644 --- a/internal/domain/debt.go +++ b/internal/domain/debt.go @@ -251,7 +251,7 @@ func DebtClassForBlockReason(r BlockReason) (DebtClass, bool) { // classes a worker actually emits are listed; an unknown one is not guessed at. func DebtClassForFailureClass(f string) (DebtClass, bool) { switch f { - case "retry_limit", "launch_failed", "launch_transient", "launch_uncertain", "prompt_not_submitted", "lease_expired": + case "retry_limit", "launch_failed", "launch_transient", "launch_uncertain", "prompt_not_submitted", "lease_expired", "handoff_unanswered": return DebtOperational, true case "invalid_handoff": return DebtCorrectness, true diff --git a/internal/herdr/herdr.go b/internal/herdr/herdr.go index 6765049..d8dc56f 100644 --- a/internal/herdr/herdr.go +++ b/internal/herdr/herdr.go @@ -174,6 +174,14 @@ type Session struct { // its ยง6.1 handoff (HandoffFile) โ€” avoids re-sending the same prompt // every tick while Release keeps waiting for the file to appear. HandoffRequested bool `json:"handoff_requested,omitempty"` + // HandoffRequestedAt stamps that prompt. A request nobody answers used to + // end as an ordinary idle expiry, indistinguishable from an agent that + // never started (F62); the stamp is what makes the wait bounded and the + // giving-up causal. + HandoffRequestedAt time.Time `json:"handoff_requested_at,omitempty"` + // HandoffRetried records that the request was re-sent once, so a session + // waiting on an answer is not re-prompted every tick. + HandoffRetried bool `json:"handoff_retried,omitempty"` // HandoffReason is selected by the coordinator when it asks for the // semantic report. The checkout worker, rather than the harness, copies // it into the canonical handoff it seals at release time.