Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
28 changes: 27 additions & 1 deletion internal/daemon/pool.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 (
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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)
Expand Down
30 changes: 30 additions & 0 deletions internal/daemon/pool_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ package daemon
import (
"context"
"errors"
"fmt"
"sync"
"sync/atomic"
"testing"
Expand Down Expand Up @@ -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},
Expand Down
Loading