diff --git a/caddy/app.go b/caddy/app.go index 8dd40490d6..cf7368a64e 100644 --- a/caddy/app.go +++ b/caddy/app.go @@ -23,6 +23,7 @@ import ( var ( options []frankenphp.Option optionsMU sync.RWMutex + activeApp atomic.Pointer[FrankenPHPApp] ) // EXPERIMENTAL: RegisterWorkers provides a way for extensions to register frankenphp.Workers @@ -140,6 +141,9 @@ func (f *FrankenPHPApp) collectOptions(repl *caddy.Replacer, keep bool) ([]frank frankenphp.WithMaxIdleTime(f.MaxIdleTime), frankenphp.WithMaxRequests(f.MaxRequests), ) + if f.httpApp != nil { + opts = append(opts, frankenphp.WithWorkerRequestDrainTimeout(time.Duration(f.httpApp.GracePeriod))) + } usedWorkerNames := make(map[string]bool, len(f.Workers)) @@ -182,8 +186,10 @@ func (f *FrankenPHPApp) Start() error { // if FrankenPHP is currently running, shut it down first // this will happen in admin API reloads and caddy tests + activeApp.Store(f) frankenphp.Shutdown() if err := frankenphp.Init(f.opts...); err != nil { + activeApp.CompareAndSwap(f, nil) return err } @@ -193,6 +199,9 @@ func (f *FrankenPHPApp) Start() error { } func (f *FrankenPHPApp) Stop() error { + f.hasStarted.Store(false) + activeApp.CompareAndSwap(f, nil) + if f.logger.Enabled(f.ctx, slog.LevelInfo) { f.logger.LogAttrs(f.ctx, slog.LevelInfo, "FrankenPHP stopped 🐘") } @@ -215,32 +224,39 @@ func (f *FrankenPHPApp) Stop() error { // register workers and servers for "php" and "php_server" modules func (f *FrankenPHPApp) collectModuleOptions(repl *caddy.Replacer, usedWorkerNames map[string]bool, keep bool) ([]frankenphp.Option, error) { opts := make([]frankenphp.Option, 0, len(f.modules)) - serversByIndex := make(map[int]*frankenphp.Server, len(f.modules)) + type registeredServer struct { + server *frankenphp.Server + name string + } + serversByIndex := make(map[int]registeredServer, len(f.modules)) for _, module := range f.modules { // modules with the same server_idx should share the same server instance // example: the worker { match * } rule adds 2 "php" subroutes to the caddy handler // the 2 handlers belong to the same "php_server" and must therefore share workers if module.ServerIndex != 0 { - if server, ok := serversByIndex[module.ServerIndex]; ok { + if registered, ok := serversByIndex[module.ServerIndex]; ok { if keep { - module.server = server + module.server = registered.server + module.reloadName = registered.name } continue } } - server, moduleOpts, err := f.collectModule(repl, module, usedWorkerNames) + serverName := f.resolveServerName(module) + server, moduleOpts, err := f.collectModule(repl, module, usedWorkerNames, serverName) if err != nil { return nil, err } if keep { module.server = server + module.reloadName = serverName } if module.ServerIndex != 0 { - serversByIndex[module.ServerIndex] = server + serversByIndex[module.ServerIndex] = registeredServer{server: server, name: serverName} } opts = append(opts, moduleOpts...) } @@ -248,8 +264,7 @@ func (f *FrankenPHPApp) collectModuleOptions(repl *caddy.Replacer, usedWorkerNam return opts, nil } -func (f *FrankenPHPApp) collectModule(repl *caddy.Replacer, module *FrankenPHPModule, usedWorkerNames map[string]bool) (*frankenphp.Server, []frankenphp.Option, error) { - serverName := f.resolveServerName(module) +func (f *FrankenPHPApp) collectModule(repl *caddy.Replacer, module *FrankenPHPModule, usedWorkerNames map[string]bool, serverName string) (*frankenphp.Server, []frankenphp.Option, error) { server, err := frankenphp.NewServer( module.resolvedDocumentRoot, frankenphp.WithServerName(serverName), diff --git a/caddy/module.go b/caddy/module.go index 20dcec9ee2..199b206cb6 100644 --- a/caddy/module.go +++ b/caddy/module.go @@ -59,6 +59,7 @@ type FrankenPHPModule struct { requestEnv frankenphp.PreparedEnv requestOptions []frankenphp.RequestOption server *frankenphp.Server + reloadName string logger *slog.Logger app *FrankenPHPApp } @@ -187,11 +188,8 @@ func needReplacement(s string) bool { // ServeHTTP implements caddyhttp.MiddlewareHandler. func (f *FrankenPHPModule) ServeHTTP(w http.ResponseWriter, r *http.Request, _ caddyhttp.Handler) error { if !f.app.hasStarted.Load() { - // stall any incoming request if FrankenPHP has not started yet, blocking for up to 10 seconds - select { - case <-f.app.started: - case <-time.After(10 * time.Second): - return caddyhttp.Error(http.StatusServiceUnavailable, frankenphp.ErrNotRunning) + if err := f.app.waitForStartup(r.Context()); err != nil { + return err } } @@ -226,9 +224,17 @@ func (f *FrankenPHPModule) ServeHTTP(w http.ResponseWriter, r *http.Request, _ c } err := f.server.ServeHTTP(w, r, opts...) + if err != nil { + if errors.Is(err, frankenphp.ErrNotRunning) { + if app := activeApp.Load(); app != nil && app != f.app { + return f.requestReload(w, r) + } + return caddyhttp.Error(http.StatusServiceUnavailable, err) + } - if _, rejected := errors.AsType[frankenphp.ErrRejected](err); err != nil && !rejected { - return caddyhttp.Error(http.StatusInternalServerError, err) + if _, rejected := errors.AsType[frankenphp.ErrRejected](err); !rejected { + return caddyhttp.Error(http.StatusInternalServerError, err) + } } return nil diff --git a/caddy/reload.go b/caddy/reload.go new file mode 100644 index 0000000000..fef44c1eae --- /dev/null +++ b/caddy/reload.go @@ -0,0 +1,83 @@ +package caddy + +import ( + "context" + "net/http" + "time" + + "github.com/caddyserver/caddy/v2/modules/caddyhttp" + "github.com/dunglas/frankenphp" +) + +// requestReload dispatches r on the new runtime. Safe because the request +// never started executing: its body is unread and no response was written. +func (f *FrankenPHPModule) requestReload(w http.ResponseWriter, r *http.Request) error { + app := activeApp.Load() + if app == nil { + return caddyhttp.Error(http.StatusServiceUnavailable, frankenphp.ErrNotRunning) + } + if err := app.waitForStartup(r.Context()); err != nil { + return err + } + if !app.hasStarted.Load() || f.server == nil { + return caddyhttp.Error(http.StatusServiceUnavailable, frankenphp.ErrNotRunning) + } + + // Runtime-generated server_N names depend on registration order across reloads. + name := f.reloadName + if old := f.app.reloadModule(name); old == nil || old.server != f.server { + return caddyhttp.Error(http.StatusServiceUnavailable, frankenphp.ErrNotRunning) + } + if m := app.reloadModule(name); m != nil { + return m.ServeHTTP(w, r, nil) + } + return caddyhttp.Error(http.StatusServiceUnavailable, frankenphp.ErrNotRunning) +} + +func (f *FrankenPHPApp) waitForStartup(ctx context.Context) error { + select { + case <-f.started: + return nil + default: + } + + wait := f.MaxWaitTime + if f.httpApp != nil { + if grace := time.Duration(f.httpApp.GracePeriod); grace > 0 && (wait == 0 || grace < wait) { + wait = grace + } + } + var expired <-chan time.Time + if wait > 0 { + timer := time.NewTimer(wait) + defer timer.Stop() + expired = timer.C + } + select { + case <-f.started: + return nil + case <-ctx.Done(): + return ctx.Err() + case <-f.ctx.Done(): + return caddyhttp.Error(http.StatusServiceUnavailable, frankenphp.ErrNotRunning) + case <-expired: + return caddyhttp.Error(http.StatusServiceUnavailable, frankenphp.ErrMaxWaitTimeExceeded) + } +} + +func (f *FrankenPHPApp) reloadModule(name string) *FrankenPHPModule { + var target *FrankenPHPModule + for _, m := range f.modules { + if m.server == nil || m.reloadName != name { + continue + } + // A php_server may embed several handlers sharing one server. + if target != nil && target.server != m.server { + return nil + } + if target == nil { + target = m + } + } + return target +} diff --git a/caddy/reload_test.go b/caddy/reload_test.go new file mode 100644 index 0000000000..c12d77f866 --- /dev/null +++ b/caddy/reload_test.go @@ -0,0 +1,307 @@ +package caddy + +import ( + "context" + "encoding/json" + "fmt" + "log/slog" + "net/http" + "net/http/httptest" + "os" + "path/filepath" + "strings" + "sync" + "testing" + "time" + + "github.com/caddyserver/caddy/v2" + "github.com/caddyserver/caddy/v2/modules/caddyhttp" + _ "github.com/caddyserver/caddy/v2/modules/caddyhttp/headers" + "github.com/dunglas/frankenphp" + "github.com/stretchr/testify/require" +) + +var reloadTestEntered, reloadTestRelease chan struct{} + +type reloadTestPause struct { + Pause bool `json:"pause"` +} + +func (*reloadTestPause) CaddyModule() caddy.ModuleInfo { + return caddy.ModuleInfo{ + ID: "http.handlers.frankenphp_test_reload_pause", + New: func() caddy.Module { return new(reloadTestPause) }, + } +} + +func (p *reloadTestPause) ServeHTTP(w http.ResponseWriter, r *http.Request, next caddyhttp.Handler) error { + if p.Pause { + close(reloadTestEntered) + <-reloadTestRelease + } + return next.ServeHTTP(w, r) +} + +func init() { + caddy.RegisterModule(new(reloadTestPause)) +} + +func TestReloadPreservesRequestAndResponse(t *testing.T) { + for _, duringRequest := range []bool{false, true} { + t.Run(fmt.Sprintf("during_request=%t", duringRequest), func(t *testing.T) { + t.Cleanup(func() { + require.NoError(t, caddy.Stop()) + frankenphp.Shutdown() + activeApp.Store(nil) + }) + reloadTestEntered, reloadTestRelease = make(chan struct{}), make(chan struct{}) + defer func() { + select { + case <-reloadTestRelease: + default: + close(reloadTestRelease) + } + }() + load := func(policy string, pause bool) *caddyhttp.Server { + t.Helper() + config := fmt.Sprintf(`{ + "admin":{"disabled":true,"config":{"persist":false}}, + "logging":{"logs":{"default":{"level":"ERROR"}}}, + "apps":{ + "frankenphp":{"num_threads":2}, + "http":{"servers":{"test":{ + "listen":["127.0.0.1:0"], + "automatic_https":{"disable":true}, + "routes":[ + {"match":[{"path":["/static"]}],"handle":[{"handler":"static_response","body":%q}]}, + {"match":[{"host":["app.example"]}],"handle":[ + {"handler":"headers","request":{"set":{"Host":["internal.example"]},"add":{"X-Tag":["tag"]}},"response":{"add":{"X-Trace":[%q]}}}, + {"handler":"headers","response":{"set":{"X-Policy":[%q]},"deferred":true}}, + {"handler":"rewrite","uri":"/reload.php"}, + {"handler":"frankenphp_test_reload_pause","pause":%t}, + {"handler":"php","root":"../testdata"} + ]}] + }}} + } + }`, policy, policy, policy, pause) + require.NoError(t, caddy.Load([]byte(config), true)) + app := activeApp.Load() + require.Equal(t, "app.example", app.modules[0].server.Name()) + return app.httpApp.Servers["test"] + } + request := func(server *caddyhttp.Server, response *httptest.ResponseRecorder) { + r := httptest.NewRequest(http.MethodPost, "http://app.example/api?query=value", strings.NewReader("payload")) + server.ServeHTTP(response, r) + } + + old := load("old", duringRequest) + forwarded := httptest.NewRecorder() + done := make(chan struct{}) + if duringRequest { + go func() { + defer close(done) + request(old, forwarded) + }() + select { + case <-reloadTestEntered: + case <-time.After(5 * time.Second): + t.Fatal("request did not reach the old PHP handler") + } + } + + current := load("new", false) + static := httptest.NewRecorder() + old.ServeHTTP(static, httptest.NewRequest(http.MethodGet, "http://app.example/static", nil)) + require.Equal(t, "old", static.Body.String()) + if duringRequest { + close(reloadTestRelease) + select { + case <-done: + case <-time.After(5 * time.Second): + t.Fatal("request handoff did not finish") + } + } else { + request(old, forwarded) + } + direct := httptest.NewRecorder() + request(current, direct) + require.Equal(t, http.StatusOK, forwarded.Code) + require.Equal(t, direct.Body.String(), forwarded.Body.String()) + // keeps the old tree's header state + require.Equal(t, []string{"old"}, forwarded.Header().Values("X-Trace")) + require.Equal(t, "old", forwarded.Header().Get("X-Policy")) + var result map[string]string + require.NoError(t, json.Unmarshal(forwarded.Body.Bytes(), &result)) + require.Equal(t, map[string]string{ + "host": "internal.example", + "tag": "tag", + "uri": "/api?query=value", + "body": "payload", + }, result) + }) + } +} + +func newReloadTestApp(root string, worker bool) *FrankenPHPApp { + app := &FrankenPHPApp{ + NumThreads: 1, + MaxThreads: 2, + ctx: context.Background(), + logger: slog.New(slog.DiscardHandler), + started: make(chan any), + } + module := &FrankenPHPModule{Name: "site", resolvedDocumentRoot: root, app: app} + if worker { + module.Workers = []workerConfig{{Name: "site-worker", FileName: filepath.Join(root, "index.php"), Num: 1}} + } + app.modules = []*FrankenPHPModule{module} + return app +} + +func newReloadTestRequest() *http.Request { + r := httptest.NewRequest(http.MethodGet, "http://app.example/index.php", nil) + ctx := context.WithValue(r.Context(), caddy.ReplacerCtxKey, caddy.NewReplacer()) + ctx = context.WithValue(ctx, caddyhttp.OriginalRequestCtxKey, *r) + return r.WithContext(ctx) +} + +func waitForReloadTest(t *testing.T, done <-chan struct{}) { + t.Helper() + select { + case <-done: + case <-time.After(5 * time.Second): + t.Fatal("reload test did not finish") + } +} + +type reloadTestWaitingContext struct { + context.Context + waiting chan struct{} + once sync.Once +} + +func (c *reloadTestWaitingContext) Done() <-chan struct{} { + c.once.Do(func() { close(c.waiting) }) + return c.Context.Done() +} + +func TestReloadUsesNewRootWhileShutdownDrains(t *testing.T) { + for _, worker := range []bool{false, true} { + t.Run(fmt.Sprintf("worker=%t", worker), func(t *testing.T) { + root, err := filepath.EvalSymlinks(t.TempDir()) + require.NoError(t, err) + oldRoot, newRoot := filepath.Join(root, "old"), filepath.Join(root, "new") + for path, body := range map[string]string{oldRoot: "OLD", newRoot: "NEW"} { + require.NoError(t, os.Mkdir(path, 0700)) + script := fmt.Sprintf(" 0 { @@ -489,6 +490,10 @@ func Shutdown() { // shutdown without any locking (for internal use) func shutdown() { + // Reject new requests before hooks or PHP threads begin draining. + // Requests already executing can finish on their original server. + unregisterServers() + // call the shutdown hooks (mainly useful for extensions) for _, fn := range onServerShutdown { fn() @@ -496,7 +501,9 @@ func shutdown() { drainWatchers() drainPHPThreads() - unregisterServers() + for _, worker := range workers { + go worker.drainRequests(globalCtx, workerRequestDrainTimeout) + } metrics.Shutdown() diff --git a/options.go b/options.go index 5566ed41cf..c6ff80192e 100644 --- a/options.go +++ b/options.go @@ -37,6 +37,7 @@ type opt struct { maxIdleTime time.Duration maxRequests int servers []*Server + workerDrain time.Duration } type workerOpt struct { @@ -88,6 +89,15 @@ func WithMaxThreads(maxThreads int) Option { } } +// WithWorkerRequestDrainTimeout sets how long retired worker queues accept +// requests for handoff. Zero waits until the configured context is canceled. +func WithWorkerRequestDrainTimeout(timeout time.Duration) Option { + return func(o *opt) error { + o.workerDrain = timeout + return nil + } +} + func WithMetrics(m Metrics) Option { return func(o *opt) error { o.metrics = m diff --git a/scaling.go b/scaling.go index dd21a7e37c..16646912f8 100644 --- a/scaling.go +++ b/scaling.go @@ -93,6 +93,11 @@ func scaleWorkerThread(worker *worker, done chan struct{}, mstate *state.ThreadS return } + // Do not scale a worker from a previous PHP runtime. + if worker.done != done { + return + } + thread, err := addWorkerThread(worker) if err != nil { if globalLogger.Enabled(globalCtx, slog.LevelWarn) { diff --git a/testdata/reload.php b/testdata/reload.php new file mode 100644 index 0000000000..697f9a63c9 --- /dev/null +++ b/testdata/reload.php @@ -0,0 +1,8 @@ + $_SERVER['HTTP_HOST'], + 'tag' => $_SERVER['HTTP_X_TAG'], + 'uri' => $_SERVER['REQUEST_URI'], + 'body' => file_get_contents('php://input'), +]); diff --git a/threadregular_test.go b/threadregular_test.go index 7a1d0700f6..80a6c0a748 100644 --- a/threadregular_test.go +++ b/threadregular_test.go @@ -34,7 +34,7 @@ func TestRequestsQueuedBeforeThreadsAreReadyAreHandedOver(t *testing.T) { const requests = 5 errChans := make([]chan error, requests) for i := range errChans { - fc := &frankenPHPContext{done: make(chan any)} + fc := &frankenPHPContext{done: make(chan error)} errChan := make(chan error, 1) go func() { errChan <- handleRequestWithRegularPHPThreads(fc) diff --git a/worker.go b/worker.go index 4a596c2fb6..ab607ef422 100644 --- a/worker.go +++ b/worker.go @@ -3,6 +3,7 @@ package frankenphp // #include "frankenphp.h" import "C" import ( + "context" "fmt" "net/http" "os" @@ -27,6 +28,7 @@ type worker struct { maxThreads int requestOptions []RequestOption requestChan chan *frankenPHPContext + done <-chan struct{} threads []*phpThread threadMutex sync.RWMutex maxConsecutiveFailures int @@ -34,14 +36,16 @@ type worker struct { onThreadShutdown func(int) queuedRequests atomic.Int32 server *Server + metrics Metrics } var ( - workers []*worker - workersByName map[string]*worker - globalWorkersByPath map[string]*worker - watcherIsEnabled bool - startupFailChan chan error + workers []*worker + workersByName map[string]*worker + globalWorkersByPath map[string]*worker + watcherIsEnabled bool + startupFailChan chan error + workerRequestDrainTimeout time.Duration ) func initWorkers(opts []workerOpt) error { @@ -64,6 +68,8 @@ func initWorkers(opts []workerOpt) error { return err } + w.done = mainThread.done + totalThreadsToStart += w.num workers = append(workers, w) workersByName[w.name] = w @@ -205,6 +211,7 @@ func newWorker(o workerOpt) (*worker, error) { onThreadReady: o.onThreadReady, onThreadShutdown: o.onThreadShutdown, server: o.server, + metrics: metrics, } w.configureMercure(&o) @@ -273,7 +280,7 @@ func (worker *worker) isAtThreadLimit() bool { } func (worker *worker) handleRequest(fc *frankenPHPContext) error { - metrics.StartWorkerRequest(worker.name) + worker.metrics.StartWorkerRequest(worker.name) runtime.Gosched() @@ -285,7 +292,7 @@ func (worker *worker) handleRequest(fc *frankenPHPContext) error { case thread.requestChan <- fc: worker.threadMutex.RUnlock() <-fc.done - metrics.StopWorkerRequest(worker.name, time.Since(fc.startedAt)) + worker.metrics.StopWorkerRequest(worker.name, time.Since(fc.startedAt)) return nil default: @@ -297,7 +304,7 @@ func (worker *worker) handleRequest(fc *frankenPHPContext) error { // if no thread was available, mark the request as queued and apply the scaling strategy worker.queuedRequests.Add(1) - metrics.QueuedWorkerRequest(worker.name) + worker.metrics.QueuedWorkerRequest(worker.name) for { workerScaleChan := scaleChan @@ -308,18 +315,18 @@ func (worker *worker) handleRequest(fc *frankenPHPContext) error { select { case worker.requestChan <- fc: worker.queuedRequests.Add(-1) - metrics.DequeuedWorkerRequest(worker.name) - <-fc.done - metrics.StopWorkerRequest(worker.name, time.Since(fc.startedAt)) + worker.metrics.DequeuedWorkerRequest(worker.name) + err := <-fc.done + worker.metrics.StopWorkerRequest(worker.name, time.Since(fc.startedAt)) - return nil + return err case workerScaleChan <- fc: // the request has triggered scaling, continue to wait for a thread case <-timeoutChan(time.Duration(maxWaitTime.Load())): // the request has timed out stalling worker.queuedRequests.Add(-1) - metrics.DequeuedWorkerRequest(worker.name) - metrics.StopWorkerRequest(worker.name, time.Since(fc.startedAt)) + worker.metrics.DequeuedWorkerRequest(worker.name) + worker.metrics.StopWorkerRequest(worker.name, time.Since(fc.startedAt)) fc.reject(ErrMaxWaitTimeExceeded) @@ -327,3 +334,25 @@ func (worker *worker) handleRequest(fc *frankenPHPContext) error { } } } + +// drainRequests releases requests still arriving on a retired worker's queue. +// PHP threads must have stopped before it starts, so these requests were never executed. +func (worker *worker) drainRequests(ctx context.Context, timeout time.Duration) { + var expired <-chan time.Time + if timeout > 0 { + timer := time.NewTimer(timeout) + defer timer.Stop() + expired = timer.C + } + + for { + select { + case fc := <-worker.requestChan: + fc.done <- ErrNotRunning + case <-ctx.Done(): + return + case <-expired: + return + } + } +} diff --git a/worker_internal_test.go b/worker_internal_test.go index 99f96126a2..569488ae98 100644 --- a/worker_internal_test.go +++ b/worker_internal_test.go @@ -1,10 +1,14 @@ package frankenphp import ( + "context" + "io" "net/http/httptest" "os" "path/filepath" "runtime" + "strings" + "sync" "testing" "time" @@ -13,6 +17,68 @@ import ( "github.com/stretchr/testify/require" ) +func TestDrainWorkerRequests(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + stopped, resume := make(chan struct{}), make(chan struct{}) + release := sync.OnceFunc(func() { close(resume) }) + t.Cleanup(func() { release(); Shutdown() }) + require.NoError(t, Init( + WithContext(ctx), WithWorkerRequestDrainTimeout(0), + WithNumThreads(1), WithMaxThreads(2), WithMaxWaitTime(time.Second), + WithWorkers("retired", testDataPath+"/worker-with-counter.php", 1, + WithWorkerOnShutdown(func(int) { close(stopped); <-resume }), + ), + )) + w := workersByName["retired"] + request := httptest.NewRequest("POST", "http://localhost/worker-with-counter.php", strings.NewReader("payload")) + response := httptest.NewRecorder() + fc, err := newContextFromRequest(request, response, fallbackServer, WithRequestDocumentRoot(testDataPath, false)) + require.NoError(t, err) + require.Same(t, w, fc.worker) + shutdownDone := make(chan struct{}) + go func() { defer close(shutdownDone); Shutdown() }() + select { + case <-stopped: + case <-time.After(time.Second): + t.Fatal("worker did not reach shutdown") + } + result := make(chan error, 1) + go func() { result <- w.handleRequest(fc) }() + require.Eventually(t, func() bool { return w.queuedRequests.Load() == 1 }, time.Second, time.Millisecond) + release() + select { + case <-shutdownDone: + case <-time.After(time.Second): + t.Fatal("PHP shutdown did not finish") + } + select { + case err := <-result: + require.ErrorIs(t, err, ErrNotRunning) + case <-time.After(time.Second): + t.Fatal("retired worker request was not released") + } + require.Zero(t, w.queuedRequests.Load()) + require.Empty(t, response.Body.String()) + body, err := io.ReadAll(request.Body) + require.NoError(t, err) + require.Equal(t, "payload", string(body)) +} + +func TestWorkerRequestDrainerExpires(t *testing.T) { + w := &worker{requestChan: make(chan *frankenPHPContext)} + drained := make(chan struct{}) + go func() { + defer close(drained) + w.drainRequests(context.Background(), 10*time.Millisecond) + }() + select { + case <-drained: + case <-time.After(time.Second): + t.Fatal("worker queue drainer exceeded its configured timeout") + } +} + // TestRestartWorkersForceKillsStuckThread verifies the drain path does // not hang when a worker is stuck in a blocking PHP call (sleep, etc.). // macOS has no realtime signals so we can't unblock sleep() there; skip.