mirror of
https://github.com/gotenberg/gotenberg.git
synced 2026-10-08 05:23:18 +01:00
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.
This commit is contained in:
@@ -115,6 +115,17 @@ const (
|
|||||||
restartReasonMaxRequests = "max_requests"
|
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 {
|
type processSupervisor struct {
|
||||||
logger *slog.Logger
|
logger *slog.Logger
|
||||||
engine string
|
engine string
|
||||||
@@ -152,6 +163,9 @@ type processSupervisor struct {
|
|||||||
consecutiveHealthFailures atomic.Int64 // reset to 0 on every successful probe
|
consecutiveHealthFailures atomic.Int64 // reset to 0 on every successful probe
|
||||||
idleMu sync.Mutex // protects idleStopChan
|
idleMu sync.Mutex // protects idleStopChan
|
||||||
idleStopChan chan struct{} // signal to stop the idle ticker goroutine
|
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
|
// NewProcessSupervisor initializes a new [ProcessSupervisor]. engine names the
|
||||||
@@ -175,6 +189,7 @@ func NewProcessSupervisor(logger *slog.Logger, engine string, process Process, m
|
|||||||
maxQueueSize: maxQueueSize,
|
maxQueueSize: maxQueueSize,
|
||||||
maxConcurrency: maxConcurrency,
|
maxConcurrency: maxConcurrency,
|
||||||
idleShutdownTimeout: idleShutdownTimeout,
|
idleShutdownTimeout: idleShutdownTimeout,
|
||||||
|
eagerRestartTimeout: defaultEagerRestartTimeout,
|
||||||
}
|
}
|
||||||
b.reqCounter.Store(0)
|
b.reqCounter.Store(0)
|
||||||
b.reqQueueSize.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
|
// maybeRestartAfterTask checks if the maximum request limit has been reached
|
||||||
// and, if so, triggers an asynchronous restart. If a restart is initiated, it
|
// and, if so, triggers an asynchronous restart bounded by
|
||||||
// takes ownership of the caller's semaphore slot (the caller must not release
|
// [eagerRestartTimeout]. If a restart is initiated, it takes ownership of the
|
||||||
// it). Returns true if ownership was taken.
|
// caller's semaphore slot (the caller must not release it). Returns true if
|
||||||
|
// ownership was taken.
|
||||||
func (s *processSupervisor) maybeRestartAfterTask(logger *slog.Logger) bool {
|
func (s *processSupervisor) maybeRestartAfterTask(logger *slog.Logger) bool {
|
||||||
if s.maxReqLimit <= 0 || s.reqCounter.Load() < s.maxReqLimit {
|
if s.maxReqLimit <= 0 || s.reqCounter.Load() < s.maxReqLimit {
|
||||||
return false
|
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...")
|
s.logger.DebugContext(context.Background(), "max request limit reached, restarting eagerly...")
|
||||||
|
|
||||||
go func() {
|
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()
|
s.restartMutex.Unlock()
|
||||||
if restartErr != nil {
|
if restartErr != nil {
|
||||||
s.logger.ErrorContext(context.Background(), fmt.Sprintf("process restart after task: %v", restartErr))
|
s.logger.ErrorContext(context.Background(), fmt.Sprintf("process restart after task: %v", restartErr))
|
||||||
|
|||||||
@@ -381,6 +381,91 @@ func TestProcessSupervisor_Healthy_UnplannedRestart(t *testing.T) {
|
|||||||
close(release)
|
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
|
// TestProcessSupervisor_Healthy_CachesPositiveResult verifies that a
|
||||||
// successful probe is cached for [healthCheckCacheTTL] so subsequent
|
// successful probe is cached for [healthCheckCacheTTL] so subsequent
|
||||||
// supervisor.Healthy() calls do not re-issue the underlying process
|
// supervisor.Healthy() calls do not re-issue the underlying process
|
||||||
|
|||||||
Reference in New Issue
Block a user