diff --git a/internal/daemon/pool.go b/internal/daemon/pool.go index 6de38315b..9d7839656 100644 --- a/internal/daemon/pool.go +++ b/internal/daemon/pool.go @@ -17,6 +17,28 @@ const ( ExitPermanent = 76 ) +// Exit codes `zero exec` (internal/cli/exec.go) uses to report a run that +// finished, or was stopped on purpose, rather than crashed. They are duplicated +// here because the daemon cannot import the CLI package; keep them in sync. +const ( + exitUsage = 2 // bad flags or arguments: a retry fails the same way + exitProvider = 3 // provider failure the run already surfaced + exitIncomplete = 4 // the agent ran (edits, commands) and stopped unfinished + exitInterrupted = 130 // SIGINT: someone stopped it on purpose +) + +// isFinalExit reports whether a worker exit code means the run is over and must +// not be retried. Retrying would repeat the agent's side effects and re-bill the +// provider, and append a second run to the same session stream. Only crash-type +// exits (1, signals) and launch/read errors are worth another attempt. +func isFinalExit(code int) bool { + switch code { + case exitUsage, exitProvider, exitIncomplete, exitInterrupted: + return true + } + return false +} + // defaultPoolSize, defaultMaxAttempts and defaultKillTimeout are used when a // PoolOptions field is left zero. const ( @@ -176,7 +198,8 @@ type Sink interface { // Run leases a worker slot and dispatches spec to a worker, streaming its // stream-json lines to sink. It is at-least-once with bounded retries: a worker // that crashes (non-zero, non-permanent) is retried on a fresh worker after a -// backoff; ExitPermanent or exhausting MaxAttempts returns ErrPermanent. Run +// backoff; ExitPermanent, a final `zero exec` exit (usage, provider, incomplete, +// interrupted) or exhausting MaxAttempts returns ErrPermanent. Run // queues when all slots are busy. The returned int is the final worker exit code. func (p *Pool) Run(ctx context.Context, spec WorkerSpec, sink Sink) (int, error) { if err := p.acquire(ctx); err != nil { @@ -209,6 +232,9 @@ func (p *Pool) Run(ctx context.Context, spec WorkerSpec, sink Sink) (int, error) case code == ExitPermanent: p.logf("worker %d exited permanently (code=%d) — not retrying", stat.id, code) return code, ErrPermanent + case isFinalExit(code): + p.logf("worker %d finished with exit code %d — not retrying", stat.id, code) + return code, fmt.Errorf("%w: worker %d exited code=%d", ErrPermanent, stat.id, code) case code == ExitTempfail: lastErr = fmt.Errorf("worker %d tempfail (code=%d)", stat.id, code) p.logf("worker %d tempfail — retry after %s", stat.id, p.opts.TempfailDelay) diff --git a/internal/daemon/pool_test.go b/internal/daemon/pool_test.go index 2270bbf49..288f59858 100644 --- a/internal/daemon/pool_test.go +++ b/internal/daemon/pool_test.go @@ -3,6 +3,7 @@ package daemon import ( "context" "errors" + "fmt" "sync" "sync/atomic" "testing" @@ -148,6 +149,35 @@ func TestPoolRunPermanentStopsImmediately(t *testing.T) { } } +// A `zero exec` run that ended with usage (2), provider (3), incomplete (4) or +// interrupted (130) is over: retrying would repeat its side effects and merge a +// second run into the same session stream. +func TestPoolRunFinalExitCodesDoNotRetry(t *testing.T) { + for _, code := range []int{2, 3, 4, 130} { + t.Run(fmt.Sprintf("exit_%d", code), func(t *testing.T) { + launcher, calls := seqLauncher( + &fakeWorker{pid: 1, exitCode: code, out: []string{"first run"}}, + &fakeWorker{pid: 2, exitCode: 0, out: []string{"second run"}}, + ) + pool, _ := NewPool(PoolOptions{Size: 1, Launcher: launcher, MaxAttempts: 5, Backoff: func(int) time.Duration { return 0 }}) + sink := &collectSink{} + got, err := pool.Run(context.Background(), WorkerSpec{Session: "s"}, sink) + if !errors.Is(err, ErrPermanent) { + t.Fatalf("Run err = %v, want ErrPermanent", err) + } + if got != code { + t.Fatalf("Run code = %d, want the worker's exit code %d", got, code) + } + if *calls != 1 { + t.Fatalf("launcher calls = %d, want 1 (no retry)", *calls) + } + if len(sink.lines) != 1 || sink.lines[0] != "first run" { + t.Fatalf("sink lines = %v, want only the first run's output", sink.lines) + } + }) + } +} + func TestPoolRunCapExhausted(t *testing.T) { launcher, calls := seqLauncher( &fakeWorker{pid: 1, exitCode: 1},