From 88ddaed09bc996812d97c4ee89725acaf7191920 Mon Sep 17 00:00:00 2001 From: Julien Neuhart Date: Thu, 3 Sep 2026 16:04:20 +0200 Subject: [PATCH] fix(supervisor): report healthy during planned process restarts The eager restart fired after --chromium-restart-after or --libreoffice-restart-after conversions made Healthy() report false, so a client probing /health between two conversions got a 503 from an otherwise serving node. Tasks arriving during that window are requeued by acquireSlot, not rejected. Track whether the in-flight restart is planned and keep reporting healthy for those. Unplanned restarts still report unhealthy so load balancers get honest information. Closes #1648 --- pkg/gotenberg/supervisor.go | 55 +++++++--- pkg/gotenberg/supervisor_test.go | 134 ++++++++++++++++++++++- test/integration/features/health.feature | 16 ++- test/integration/scenario/scenario.go | 111 +++++++++++++++++++ 4 files changed, 295 insertions(+), 21 deletions(-) diff --git a/pkg/gotenberg/supervisor.go b/pkg/gotenberg/supervisor.go index 57006619..1d568033 100644 --- a/pkg/gotenberg/supervisor.go +++ b/pkg/gotenberg/supervisor.go @@ -59,8 +59,9 @@ type ProcessSupervisor interface { // Healthy checks and returns the health status of the managed [Process]. // // A non-started process is considered healthy (startup is deferred until - // the first request). Returns false if the process is currently restarting - // or is reported unhealthy by the underlying [Process]. + // the first request), as is one going through a planned restart, since it + // keeps serving traffic. Returns false during an unplanned restart or when + // the underlying [Process] reports unhealthy. Healthy() bool // Run executes a provided task while managing the state of the [Process]. @@ -103,6 +104,17 @@ const healthCheckCacheTTL = 2 * time.Second // this. See https://github.com/gotenberg/gotenberg/issues/1561. const healthFailureThreshold = 2 +// Restart reasons, also reported as the gotenberg.process.start.reason span +// attribute by [processSupervisor.tracedLaunch]. Only +// [restartReasonMaxRequests] is a planned restart: it fires on a healthy +// process that reached its conversion limit, so the node keeps serving +// traffic throughout. The others signal a process that cannot serve. +const ( + restartReasonFirstStart = "first_start" + restartReasonUnhealthy = "unhealthy" + restartReasonMaxRequests = "max_requests" +) + type processSupervisor struct { logger *slog.Logger engine string @@ -118,11 +130,16 @@ type processSupervisor struct { // transient failure (such as a cold-start timeout) must not poison the // supervisor for the rest of the container's lifetime. See // https://github.com/gotenberg/gotenberg/issues/1538. - firstStartMu sync.Mutex - reqCounter atomic.Int64 - reqQueueSize atomic.Int64 - restartsCounter atomic.Int64 - isRestarting atomic.Bool + firstStartMu sync.Mutex + reqCounter atomic.Int64 + reqQueueSize atomic.Int64 + restartsCounter atomic.Int64 + isRestarting atomic.Bool + // restartPlanned records whether the in-flight restart is a planned one + // (see [restartReasonMaxRequests]). Written before isRestarting and never + // cleared, so a reader that observed isRestarting always sees the matching + // kind. See [processSupervisor.Healthy]. + restartPlanned atomic.Bool activeTasks atomic.Int64 restartMutex sync.Mutex idleShutdownTimeout time.Duration @@ -234,9 +251,17 @@ func (s *processSupervisor) Healthy() bool { } if s.isRestarting.Load() { - // A restarting process is not yet healthy. This gives load balancers - // honest information so they can avoid routing traffic to this node. - return false + // A planned restart is routine maintenance: the process reached the + // limit set by --chromium-restart-after (env CHROMIUM_RESTART_AFTER) or + // --libreoffice-restart-after (env LIBREOFFICE_RESTART_AFTER) while + // healthy. Tasks arriving during it are requeued by acquireSlot, not + // rejected, so the node still serves traffic and must report healthy. A + // probe sent between two conversions used to fail here. + // See https://github.com/gotenberg/gotenberg/issues/1648. + // + // An unplanned restart keeps reporting unhealthy, which gives load + // balancers honest information so they can avoid routing traffic here. + return s.restartPlanned.Load() } // Cache hit: a recent probe succeeded. Skip the CDP roundtrip so probe @@ -484,7 +509,7 @@ func (s *processSupervisor) ensureStarted(ctx context.Context) error { return nil } - err := s.tracedLaunch(ctx, "first_start", func() error { + err := s.tracedLaunch(ctx, restartReasonFirstStart, func() error { return s.runWithDeadline(ctx, s.Launch) }) if err != nil { @@ -525,7 +550,7 @@ func (s *processSupervisor) ensureHealthy(ctx context.Context) error { s.logger.DebugContext(context.Background(), "process is unhealthy, cannot handle task, restarting...") - if err := s.doRestart(ctx, "unhealthy"); err != nil { + if err := s.doRestart(ctx, restartReasonUnhealthy); err != nil { return fmt.Errorf("process restart before task: %w", err) } @@ -548,7 +573,7 @@ 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(), "max_requests") + restartErr := s.doRestartLocked(context.Background(), restartReasonMaxRequests) s.restartMutex.Unlock() if restartErr != nil { s.logger.ErrorContext(context.Background(), fmt.Sprintf("process restart after task: %v", restartErr)) @@ -571,6 +596,10 @@ func (s *processSupervisor) doRestart(ctx context.Context, reason string) error // doRestartLocked performs the restart drain logic. The caller must hold restartMutex. func (s *processSupervisor) doRestartLocked(ctx context.Context, reason string) error { + // Publish the kind before raising the flag. [processSupervisor.Healthy] + // reads restartPlanned only after it observes isRestarting, so this + // ordering keeps it from pairing a new restart with a stale kind. + s.restartPlanned.Store(reason == restartReasonMaxRequests) s.isRestarting.Store(true) defer s.isRestarting.Store(false) diff --git a/pkg/gotenberg/supervisor_test.go b/pkg/gotenberg/supervisor_test.go index f9bdfa0c..cc553ad3 100644 --- a/pkg/gotenberg/supervisor_test.go +++ b/pkg/gotenberg/supervisor_test.go @@ -162,11 +162,12 @@ func TestProcessSupervisor_restart(t *testing.T) { func TestProcessSupervisor_Healthy(t *testing.T) { for _, tc := range []struct { - scenario string - initiallyStarted bool - initiallyRestarting bool - processHealthy bool - expectHealthy bool + scenario string + initiallyStarted bool + initiallyRestarting bool + initiallyRestartPlanned bool + processHealthy bool + expectHealthy bool }{ { scenario: "non-started process is healthy", @@ -179,6 +180,13 @@ func TestProcessSupervisor_Healthy(t *testing.T) { initiallyRestarting: true, expectHealthy: false, }, + { + scenario: "process going through a planned restart is healthy", + initiallyStarted: true, + initiallyRestarting: true, + initiallyRestartPlanned: true, + expectHealthy: true, + }, { scenario: "process reports as healthy", initiallyStarted: true, @@ -208,6 +216,9 @@ func TestProcessSupervisor_Healthy(t *testing.T) { if tc.initiallyRestarting { ps.isRestarting.Store(true) } + if tc.initiallyRestartPlanned { + ps.restartPlanned.Store(true) + } healthy := ps.Healthy() @@ -257,6 +268,119 @@ func TestProcessSupervisor_Healthy_ConsecutiveFailures(t *testing.T) { } } +// TestProcessSupervisor_Healthy_PlannedRestart reproduces +// https://github.com/gotenberg/gotenberg/issues/1648. It drives the real +// Run() path until the maximum request limit triggers the eager restart, then +// asserts the supervisor reports healthy while that restart is in flight. +// Tasks arriving during it are requeued by acquireSlot, not rejected, so the +// node still serves traffic. +func TestProcessSupervisor_Healthy_PlannedRestart(t *testing.T) { + logger := slog.New(slog.DiscardHandler) + + const maxReqLimit = 10 + + restarting := make(chan struct{}) + release := make(chan struct{}) + + var ( + starts atomic.Int64 + signalOne sync.Once + ) + process := &ProcessMock{ + StartMock: func(_ *slog.Logger) error { + // Hold the restart open so the assertions below run inside the + // window that used to report unhealthy. + if starts.Add(1) > 1 { + signalOne.Do(func() { close(restarting) }) + <-release + } + + return nil + }, + StopMock: func(_ *slog.Logger) error { return nil }, + HealthyMock: func(_ *slog.Logger) bool { return true }, + } + + ps := NewProcessSupervisor(logger, "test", process, maxReqLimit, 0, 1, 0).(*processSupervisor) + + for i := range maxReqLimit { + err := ps.Run(context.Background(), logger, func() error { return nil }) + if err != nil { + t.Fatalf("task %d: expected no error but got: %v", i+1, err) + } + } + + select { + case <-restarting: + case <-time.After(10 * time.Second): + t.Fatalf("expected an eager restart after %d tasks", maxReqLimit) + } + + if !ps.isRestarting.Load() { + t.Fatal("expected the supervisor to be restarting") + } + + if !ps.restartPlanned.Load() { + t.Fatal("expected the restart to be flagged as planned") + } + + if !ps.Healthy() { + t.Fatal("expected a planned restart to report healthy") + } + + close(release) +} + +// TestProcessSupervisor_Healthy_UnplannedRestart verifies the counterpart of +// [TestProcessSupervisor_Healthy_PlannedRestart]: a restart triggered by an +// unhealthy process keeps reporting unhealthy, so load balancers get honest +// information. +func TestProcessSupervisor_Healthy_UnplannedRestart(t *testing.T) { + logger := slog.New(slog.DiscardHandler) + + restarting := make(chan struct{}) + release := make(chan struct{}) + + var signalOne sync.Once + process := &ProcessMock{ + StartMock: func(_ *slog.Logger) error { + signalOne.Do(func() { close(restarting) }) + <-release + + return nil + }, + StopMock: func(_ *slog.Logger) error { return nil }, + HealthyMock: func(_ *slog.Logger) bool { return false }, + } + + ps := NewProcessSupervisor(logger, "test", process, 0, 0, 1, 0).(*processSupervisor) + ps.firstStart.Store(true) + + go func() { + _ = ps.ensureHealthy(context.Background()) + }() + + select { + case <-restarting: + case <-time.After(10 * time.Second): + t.Fatal("expected an unhealthy restart to be triggered") + } + + if !ps.isRestarting.Load() { + t.Fatal("expected the supervisor to be restarting") + } + + if ps.restartPlanned.Load() { + t.Fatal("expected the restart not to be flagged as planned") + } + + if ps.Healthy() { + t.Fatal("expected an unplanned restart to report unhealthy") + } + + close(release) +} + // TestProcessSupervisor_Healthy_CachesPositiveResult verifies that a // successful probe is cached for [healthCheckCacheTTL] so subsequent // supervisor.Healthy() calls do not re-issue the underlying process diff --git a/test/integration/features/health.feature b/test/integration/features/health.feature index 6fbbac29..88d583ca 100644 --- a/test/integration/features/health.feature +++ b/test/integration/features/health.feature @@ -1,6 +1,5 @@ # TODO: # 1. Check if down for each module. -# 2. Restarting modules do not make health check fail. @health Feature: /health @@ -106,7 +105,18 @@ Feature: /health When I make a "HEAD" request to Gotenberg at the "/foo/health" endpoint Then the response status code should be 200 + # A planned restart, the eager cycle after LIBREOFFICE_RESTART_AFTER + # conversions, must not fail the health check: requests arriving during it + # are requeued, not rejected. Setting the limit to 1 restarts LibreOffice + # after every conversion, so each probe lands right on a restart. + # See https://github.com/gotenberg/gotenberg/issues/1648. + Scenario: GET /health (Planned LibreOffice Restart) + Given I have a Gotenberg container with the following environment variable(s): + | LIBREOFFICE_RESTART_AFTER | 1 | + When I make 5 sequential "POST" requests to Gotenberg at the "/forms/libreoffice/convert" endpoint, probing "/health" after each, with the following form data and header(s): + | files | testdata/page_1.docx | file | + Then all probe response status codes should be 200 + # TODO: -# 1. Check if down for each module. -# 2. Restarting modules do not make health check fail. \ No newline at end of file +# 1. Check if down for each module. \ No newline at end of file diff --git a/test/integration/scenario/scenario.go b/test/integration/scenario/scenario.go index 9b8cd428..ec360bd4 100644 --- a/test/integration/scenario/scenario.go +++ b/test/integration/scenario/scenario.go @@ -88,6 +88,7 @@ func findScenarioLine(filePath, name string) int { type scenario struct { resp *httptest.ResponseRecorder concurrentResps []*httptest.ResponseRecorder + probeResps []*httptest.ResponseRecorder workdir string teststoreDir string gotenbergContainer testcontainers.Container @@ -99,6 +100,7 @@ type scenario struct { func (s *scenario) reset(ctx context.Context) error { s.resp = httptest.NewRecorder() s.concurrentResps = nil + s.probeResps = nil err := os.RemoveAll(s.workdir) if err != nil { @@ -470,6 +472,113 @@ func (s *scenario) iMakeConcurrentRequestsToGotenberg(ctx context.Context, count return nil } +// iMakeSequentialRequestsToGotenbergProbing mirrors the client loop from +// https://github.com/gotenberg/gotenberg/issues/1648: a conversion, then a +// probe, repeated. It records every probe response so a scenario can assert +// that a planned process restart never fails the probe. Requests are +// sequential on purpose, since the bug only surfaces between two conversions. +func (s *scenario) iMakeSequentialRequestsToGotenbergProbing(ctx context.Context, count int, method, endpoint, probeEndpoint string, dataTable *godog.Table) error { + if s.gotenbergContainer == nil { + return errors.New("no Gotenberg container") + } + + fields := make(map[string][]string) + files := make(map[string][]string) + headers := make(map[string]string) + + for _, row := range dataTable.Rows { + name := row.Cells[0].Value + value := row.Cells[1].Value + kind := row.Cells[2].Value + + switch kind { + case "field": + fields[name] = append(fields[name], value) + case "file": + wd, err := os.Getwd() + if err != nil { + return fmt.Errorf("get current directory: %w", err) + } + value = fmt.Sprintf("%s/%s", wd, value) + files[name] = append(files[name], value) + case "header": + headers[name] = value + default: + return fmt.Errorf("unexpected %q %q", kind, value) + } + } + + base, err := containerHttpEndpoint(ctx, s.gotenbergContainer, "3000") + if err != nil { + return fmt.Errorf("get container HTTP endpoint: %w", err) + } + + record := func(resp *http.Response) (*httptest.ResponseRecorder, error) { + defer resp.Body.Close() + + body, readErr := io.ReadAll(resp.Body) + if readErr != nil { + return nil, fmt.Errorf("read response body: %w", readErr) + } + + rec := httptest.NewRecorder() + rec.Code = resp.StatusCode + for key, values := range resp.Header { + for _, v := range values { + rec.Header().Add(key, v) + } + } + _, writeErr := rec.Body.Write(body) + if writeErr != nil { + return nil, fmt.Errorf("write response body: %w", writeErr) + } + + return rec, nil + } + + s.probeResps = make([]*httptest.ResponseRecorder, 0, count) + + for i := range count { + resp, reqErr := doFormDataRequest(method, fmt.Sprintf("%s%s", base, endpoint), fields, files, headers) + if reqErr != nil { + return fmt.Errorf("request %d: do request: %w", i+1, reqErr) + } + + rec, recErr := record(resp) + if recErr != nil { + return fmt.Errorf("request %d: %w", i+1, recErr) + } + s.resp = rec + + probeResp, probeErr := doRequest(http.MethodGet, fmt.Sprintf("%s%s", base, probeEndpoint), nil, nil) + if probeErr != nil { + return fmt.Errorf("probe %d: do request: %w", i+1, probeErr) + } + + probeRec, probeRecErr := record(probeResp) + if probeRecErr != nil { + return fmt.Errorf("probe %d: %w", i+1, probeRecErr) + } + s.probeResps = append(s.probeResps, probeRec) + } + + return nil +} + +func (s *scenario) allProbeResponseStatusCodesShouldBe(expected int) error { + if len(s.probeResps) == 0 { + return errors.New("no probe responses recorded") + } + + for i, resp := range s.probeResps { + if resp.Code != expected { + return fmt.Errorf("probe %d: expected status %d, got %d %q", i+1, expected, resp.Code, resp.Body.String()) + } + } + + return nil +} + func (s *scenario) allConcurrentResponseStatusCodesShouldBe(expected int) error { if len(s.concurrentResps) == 0 { return errors.New("no concurrent responses recorded") @@ -1644,10 +1753,12 @@ func InitializeScenario(ctx *godog.ScenarioContext) { ctx.When(`^I make a "(GET|HEAD)" request to Gotenberg at the "([^"]*)" endpoint with the following header\(s\):$`, s.iMakeARequestToGotenbergWithTheFollowingHeaders) ctx.When(`^I make a "(POST)" request to Gotenberg at the "([^"]*)" endpoint with the following form data and header\(s\):$`, s.iMakeARequestToGotenbergWithTheFollowingFormDataAndHeaders) ctx.When(`^I make (\d+) concurrent "(POST)" requests to Gotenberg at the "([^"]*)" endpoint with the following form data and header\(s\):$`, s.iMakeConcurrentRequestsToGotenberg) + ctx.When(`^I make (\d+) sequential "(POST)" requests to Gotenberg at the "([^"]*)" endpoint, probing "([^"]*)" after each, with the following form data and header\(s\):$`, s.iMakeSequentialRequestsToGotenbergProbing) ctx.When(`^I wait for the asynchronous request to the webhook$`, s.iWaitForTheAsynchronousRequestToWebhook) ctx.Then(`^the Gotenberg container (should|should NOT) log the following entries:$`, s.theGotenbergContainerShouldLogTheFollowingEntries) ctx.Then(`^the response status code should be (\d+)$`, s.theResponseStatusCodeShouldBe) ctx.Then(`^all concurrent response status codes should be (\d+)$`, s.allConcurrentResponseStatusCodesShouldBe) + ctx.Then(`^all probe response status codes should be (\d+)$`, s.allProbeResponseStatusCodesShouldBe) ctx.Then(`^all concurrent responses should have (\d+) PDF\(s\)$`, s.allConcurrentResponsesShouldHavePdfs) ctx.Then(`^the (response|webhook request|file request|server request) header "([^"]*)" should be "([^"]*)"$`, s.theHeaderValueShouldBe) ctx.Then(`^the webhook request header "([^"]*)" should carry trace id "([^"]*)"$`, s.theWebhookRequestHeaderShouldCarryTraceID)