mirror of
https://github.com/gotenberg/gotenberg.git
synced 2026-10-07 21:13:18 +01:00
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
This commit is contained in:
@@ -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)
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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.
|
||||
# 1. Check if down for each module.
|
||||
@@ -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)
|
||||
|
||||
Reference in New Issue
Block a user