From d79e174c6f8d5b4ba411c9ccb684cefedee7c2f9 Mon Sep 17 00:00:00 2001 From: Julien Neuhart Date: Thu, 3 Sep 2026 16:10:16 +0200 Subject: [PATCH] fix(supervisor): bound the eager restart with a deadline maybeRestartAfterTask ran its restart on a bare background context while the drain loop in doRestartLocked selects only on ctx.Done(). A task that never completed blocked the drain forever, pinning isRestarting. Since a planned restart now reports healthy, that left the node claiming health for good. LibreOffice escaped it because its concurrency of 1 drains no slots, but Chromium defaults to 6. Give that restart its own deadline. On expiry it aborts and the next task retries it. --- pkg/gotenberg/supervisor.go | 27 ++++++++-- pkg/gotenberg/supervisor_test.go | 85 ++++++++++++++++++++++++++++++++ 2 files changed, 108 insertions(+), 4 deletions(-) diff --git a/pkg/gotenberg/supervisor.go b/pkg/gotenberg/supervisor.go index 1d568033..0c875bcd 100644 --- a/pkg/gotenberg/supervisor.go +++ b/pkg/gotenberg/supervisor.go @@ -115,6 +115,17 @@ const ( restartReasonMaxRequests = "max_requests" ) +// defaultEagerRestartTimeout bounds the restart triggered after the maximum +// request limit. That restart runs on a background context, unlike the one from +// ensureHealthy which inherits the request deadline, so without a deadline of +// its own the drain loop in [processSupervisor.doRestartLocked] would wait +// forever on a task that never completes. That would pin isRestarting and, +// with it, the health reported by [processSupervisor.Healthy]. Sized well above +// --api-timeout (30s by default) plus the engine start timeouts (20s by +// default) so it never fires while tasks are merely slow. The eager restart is +// opportunistic: on expiry it aborts, and the next task retries it. +const defaultEagerRestartTimeout = 2 * time.Minute + type processSupervisor struct { logger *slog.Logger engine string @@ -152,6 +163,9 @@ type processSupervisor struct { consecutiveHealthFailures atomic.Int64 // reset to 0 on every successful probe idleMu sync.Mutex // protects idleStopChan idleStopChan chan struct{} // signal to stop the idle ticker goroutine + // eagerRestartTimeout bounds the restart from maybeRestartAfterTask. + // Defaults to [defaultEagerRestartTimeout]; only tests shorten it. + eagerRestartTimeout time.Duration } // NewProcessSupervisor initializes a new [ProcessSupervisor]. engine names the @@ -175,6 +189,7 @@ func NewProcessSupervisor(logger *slog.Logger, engine string, process Process, m maxQueueSize: maxQueueSize, maxConcurrency: maxConcurrency, idleShutdownTimeout: idleShutdownTimeout, + eagerRestartTimeout: defaultEagerRestartTimeout, } b.reqCounter.Store(0) b.reqQueueSize.Store(0) @@ -558,9 +573,10 @@ func (s *processSupervisor) ensureHealthy(ctx context.Context) error { } // maybeRestartAfterTask checks if the maximum request limit has been reached -// and, if so, triggers an asynchronous restart. If a restart is initiated, it -// takes ownership of the caller's semaphore slot (the caller must not release -// it). Returns true if ownership was taken. +// and, if so, triggers an asynchronous restart bounded by +// [eagerRestartTimeout]. If a restart is initiated, it takes ownership of the +// caller's semaphore slot (the caller must not release it). Returns true if +// ownership was taken. func (s *processSupervisor) maybeRestartAfterTask(logger *slog.Logger) bool { if s.maxReqLimit <= 0 || s.reqCounter.Load() < s.maxReqLimit { return false @@ -573,7 +589,10 @@ func (s *processSupervisor) maybeRestartAfterTask(logger *slog.Logger) bool { s.logger.DebugContext(context.Background(), "max request limit reached, restarting eagerly...") go func() { - restartErr := s.doRestartLocked(context.Background(), restartReasonMaxRequests) + ctx, cancel := context.WithTimeout(context.Background(), s.eagerRestartTimeout) + defer cancel() + + restartErr := s.doRestartLocked(ctx, restartReasonMaxRequests) s.restartMutex.Unlock() if restartErr != nil { s.logger.ErrorContext(context.Background(), fmt.Sprintf("process restart after task: %v", restartErr)) diff --git a/pkg/gotenberg/supervisor_test.go b/pkg/gotenberg/supervisor_test.go index cc553ad3..793f6f75 100644 --- a/pkg/gotenberg/supervisor_test.go +++ b/pkg/gotenberg/supervisor_test.go @@ -381,6 +381,91 @@ func TestProcessSupervisor_Healthy_UnplannedRestart(t *testing.T) { close(release) } +// TestProcessSupervisor_doRestartLocked_DrainDeadline verifies that a drain +// unable to acquire every slot gives up on the context deadline and clears +// isRestarting. Without a deadline on the eager restart, a task that never +// completes would pin the flag and, since a planned restart reports healthy, +// leave the supervisor claiming health forever. +func TestProcessSupervisor_doRestartLocked_DrainDeadline(t *testing.T) { + logger := slog.New(slog.DiscardHandler) + + var starts atomic.Int64 + process := &ProcessMock{ + StartMock: func(_ *slog.Logger) error { + starts.Add(1) + + return nil + }, + StopMock: func(_ *slog.Logger) error { return nil }, + HealthyMock: func(_ *slog.Logger) bool { return true }, + } + + // A concurrency of 2 makes the drain acquire one slot on top of the one the + // triggering task hands over. Fill the semaphore so it never can, mimicking + // a concurrent task that never completes. + ps := NewProcessSupervisor(logger, "test", process, 1, 0, 2, 0).(*processSupervisor) + ps.semaphore <- struct{}{} + ps.semaphore <- struct{}{} + + ctx, cancel := context.WithTimeout(context.Background(), 50*time.Millisecond) + defer cancel() + + err := ps.doRestartLocked(ctx, restartReasonMaxRequests) + if err == nil { + t.Fatal("expected the drain to fail on the context deadline") + } + + if ps.isRestarting.Load() { + t.Fatal("expected isRestarting to be cleared after a failed drain") + } + + if starts.Load() != 0 { + t.Fatalf("expected no restart attempt after a failed drain but got %d", starts.Load()) + } +} + +// TestProcessSupervisor_maybeRestartAfterTask_Bounded verifies that the eager +// restart runs under a deadline. A concurrent task that never completes blocks +// the drain, and without a bound the restart goroutine would wait forever with +// isRestarting pinned, leaving Healthy() reporting a planned restart for good. +func TestProcessSupervisor_maybeRestartAfterTask_Bounded(t *testing.T) { + logger := slog.New(slog.DiscardHandler) + + process := &ProcessMock{ + StartMock: func(_ *slog.Logger) error { return nil }, + StopMock: func(_ *slog.Logger) error { return nil }, + HealthyMock: func(_ *slog.Logger) bool { return true }, + } + + ps := NewProcessSupervisor(logger, "test", process, 1, 0, 2, 0).(*processSupervisor) + ps.eagerRestartTimeout = 100 * time.Millisecond + ps.firstStart.Store(true) + + // Wedge one slot so the drain, which needs one on top of the slot the + // triggering task hands over, can never complete. + ps.semaphore <- struct{}{} + + err := ps.Run(context.Background(), logger, func() error { return nil }) + if err != nil { + t.Fatalf("expected no error but got: %v", err) + } + + waitFor := func(what string, want bool) { + t.Helper() + + deadline := time.Now().Add(5 * time.Second) + for ps.isRestarting.Load() != want { + if time.Now().After(deadline) { + t.Fatalf("timed out waiting for the eager restart to %s", what) + } + time.Sleep(5 * time.Millisecond) + } + } + + waitFor("start", true) + waitFor("give up on its deadline", false) +} + // TestProcessSupervisor_Healthy_CachesPositiveResult verifies that a // successful probe is cached for [healthCheckCacheTTL] so subsequent // supervisor.Healthy() calls do not re-issue the underlying process