Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
19 commits
Select commit Hold shift + click to select a range
9aec673
fix(daemon): synchronize TestPoolDrainKillsStraggler on active workers
jatmn Aug 18, 2026
e6c720c
fix(daemon): close launch and drain race
jatmn Aug 18, 2026
7620006
fix(daemon): wait for late launch cleanup during drain
jatmn Aug 18, 2026
9c61e27
fix(imageinput): remove vulnerable PDF parser
jatmn Aug 18, 2026
e8c13d1
fix: restore in-process PDF fallback and make drain errors terminal
jatmn Aug 18, 2026
6e79c83
fix(daemon): always decrement in-flight launch count
jatmn Aug 18, 2026
ca71906
fix(daemon,imageinput): bound late-launch drain wait and restore page…
jatmn Aug 18, 2026
afa5acb
fix(imageinput): fall back to pdfinfo when in-process page count is zero
jatmn Aug 18, 2026
4a2b4d3
chore: split PDF changes into separate PR
jatmn Aug 19, 2026
aa49bba
fix(daemon): interrupt retry delays while draining
jatmn Aug 19, 2026
25920e3
docs(daemon): clarify bounded late-launch drain cleanup
jatmn Aug 19, 2026
270d971
refactor(daemon): remove unused pool tracker
jatmn Aug 19, 2026
3e5cf78
docs(daemon): describe bounded late cleanup
jatmn Aug 19, 2026
f4c2a5f
test(daemon): prove drain waits for late worker reap
jatmn Aug 19, 2026
9c13f8b
test(daemon): bound pool synchronization waits
jatmn Aug 20, 2026
b85d335
test(daemon): assert late drain result
jatmn Aug 20, 2026
e1db5c2
fix(daemon): publish drain state before cancellation
jatmn Aug 20, 2026
5c191da
test(daemon): anchor blocked-launch drain timing
jatmn Aug 28, 2026
d5e2d6f
fix(daemon): preserve terminal outcomes across drain boundaries
jatmn Sep 12, 2026
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
129 changes: 103 additions & 26 deletions internal/daemon/pool.go
Original file line number Diff line number Diff line change
Expand Up @@ -92,13 +92,15 @@ type Pool struct {
opts PoolOptions
slots chan struct{}

mu sync.Mutex
draining bool
active map[int]WorkerHandle // worker id -> handle, for drain/kill + status
nextID int

drainOnce sync.Once
drained chan struct{}
mu sync.Mutex
draining bool
active map[int]WorkerHandle // worker id -> handle, for drain/kill + status
launching int // launchers in progress; Drain must not mistake these for idle
nextID int

drainStartOnce sync.Once
drainOnce sync.Once
drained chan struct{}
}

// workerStat tracks one in-flight request's restart count (local to Run).
Expand Down Expand Up @@ -197,22 +199,38 @@ func (p *Pool) Run(ctx context.Context, spec WorkerSpec, sink Sink) (int, error)
return 0, ErrPoolDraining
}
code, err := p.runOnce(ctx, stat.id, spec, sink)
// Fully completed work remains successful even during the drain grace
// window. Failed or killed workers must still take the terminal path.
if err == nil && code == 0 {
return 0, nil
}
// A failed run can finish just as Drain starts. Recheck before retrying
// or classifying the failure as permanent.
if p.isDraining() {
return 0, ErrPoolDraining
}
Comment thread
jatmn marked this conversation as resolved.
switch {
case err != nil:
lastErr = err
if ctx.Err() != nil {
return 0, ctx.Err()
}
// Drain is terminal: do not backoff/retry, and do not wrap the
// shutdown error as ErrPermanent when attempts are exhausted.
if errors.Is(err, ErrPoolDraining) {
return 0, ErrPoolDraining
}
p.logf("worker %d launch/run error: %v", stat.id, err)
case code == 0:
return 0, nil // clean success
case code == ExitPermanent:
p.logf("worker %d exited permanently (code=%d) — not retrying", stat.id, code)
return code, ErrPermanent
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)
if !p.sleep(ctx, p.opts.TempfailDelay) {
if p.isDraining() {
return 0, ErrPoolDraining
}
return 0, ctx.Err()
}
continue // tempfail retries do not count against the crash backoff
Expand All @@ -227,9 +245,18 @@ func (p *Pool) Run(ctx context.Context, spec WorkerSpec, sink Sink) (int, error)
delay := p.opts.Backoff(stat.restarts)
p.logf("worker %d restart %d after backoff %s", stat.id, stat.restarts, delay)
if !p.sleep(ctx, delay) {
if p.isDraining() {
return 0, ErrPoolDraining
}
return 0, ctx.Err()
}
}
// On the final tempfail attempt, sleep can select an expired timer while
// the drain channel is also ready. There is no next iteration to recheck
// shutdown, so do it before reporting attempt exhaustion.
if p.isDraining() {
return 0, ErrPoolDraining
}
if lastErr == nil {
lastErr = ErrPermanent
}
Expand All @@ -239,11 +266,43 @@ func (p *Pool) Run(ctx context.Context, spec WorkerSpec, sink Sink) (int, error)
// runOnce launches a single worker, pumps its output to sink, and returns its
// exit code. The worker handle is tracked so Drain can kill it.
func (p *Pool) runOnce(ctx context.Context, id int, spec WorkerSpec, sink Sink) (int, error) {
p.mu.Lock()
if p.draining {
p.mu.Unlock()
return 0, ErrPoolDraining
}
p.launching++
p.mu.Unlock()
inLaunch := true
defer func() {
if inLaunch {
p.mu.Lock()
p.launching--
p.mu.Unlock()
}
}()
handle, err := p.opts.Launcher(ctx, spec)
p.mu.Lock()
draining := p.draining
if err == nil && !draining {
// Key by the monotonic worker id, not the reusable OS pid, so an older
// worker's untrack cannot remove a newly launched handle (D10).
p.active[id] = handle
p.launching--
inLaunch = false
}
p.mu.Unlock()
if err != nil {
if draining {
return 0, ErrPoolDraining
}
return 0, err
}
p.track(id, handle)
if draining {
_ = handle.Kill()
_, _ = handle.Wait()
return 0, ErrPoolDraining
}
defer p.untrack(id)

// Pump stdout lines until the stream ends.
Expand Down Expand Up @@ -281,15 +340,6 @@ func (p *Pool) newStat() *workerStat {
return &workerStat{id: p.nextID}
}

// track/untrack key the active set by the pool's monotonic worker id, not the OS
// pid: the OS can reuse a pid the instant a worker exits, so a pid key could collide
// a finished worker with a freshly-launched one and drop the wrong handle (D10).
func (p *Pool) track(id int, h WorkerHandle) {
p.mu.Lock()
p.active[id] = h
p.mu.Unlock()
}

func (p *Pool) untrack(id int) {
p.mu.Lock()
delete(p.active, id)
Expand All @@ -315,25 +365,37 @@ func (p *Pool) sleep(ctx context.Context, d time.Duration) bool {
return true
case <-ctx.Done():
return false
case <-p.drained:
return false
}
}

// Drain stops accepting new work, gives in-flight workers a grace window
// (KillTimeout) to finish on their own, then force-kills any straggler. It
// returns as soon as the pool is idle (graceful) or the window elapses. Safe to
// call once; subsequent calls are no-ops.
func (p *Pool) Drain() {
p.drainOnce.Do(func() {
// beginDrain publishes the terminal pool state before callers cancel work that
// may be waiting in Run. Publishing this separately from Drain's bounded
// cleanup makes the shutdown result deterministic for every Run wakeup.
func (p *Pool) beginDrain() {
p.drainStartOnce.Do(func() {
p.mu.Lock()
p.draining = true
p.mu.Unlock()
close(p.drained)
})
}

// Drain stops accepting new work, gives in-flight workers a grace window
// (KillTimeout) to finish on their own, then force-kills any straggler. It
// returns as soon as the pool is idle, the grace window elapses and existing
// workers are force-killed, or one separately bounded late-launch cleanup wait
// elapses. A launcher that ignores cancellation may finish its cleanup after
// Drain returns. Safe to call once; subsequent calls are no-ops.
func (p *Pool) Drain() {
p.beginDrain()
p.drainOnce.Do(func() {
// Grace window: poll until idle or the deadline.
deadline := time.Now().Add(p.opts.KillTimeout)
for time.Now().Before(deadline) {
p.mu.Lock()
n := len(p.active)
n := len(p.active) + p.launching
p.mu.Unlock()
if n == 0 {
return // all workers drained gracefully
Expand All @@ -352,6 +414,21 @@ func (p *Pool) Drain() {
p.logf("drain: killing straggler worker pid=%d", h.Pid())
_ = h.Kill()
}

// Force-kill only covers handles already in active. A launcher still
// inside Launcher has no handle yet; wait one more KillTimeout for that
// late-launch path to finish kill+wait. Do not wait forever: a Launcher
// that ignores ctx would otherwise wedge shutdown.
deadline = time.Now().Add(p.opts.KillTimeout)
for time.Now().Before(deadline) {
p.mu.Lock()
n := p.launching
p.mu.Unlock()
if n == 0 {
return
}
time.Sleep(5 * time.Millisecond)
}
})
}

Expand Down
Loading
Loading