From b7cccf26bf36a464a14faa303e47185e0850835d Mon Sep 17 00:00:00 2001 From: henderkes Date: Tue, 29 Sep 2026 16:17:45 +0200 Subject: [PATCH 1/6] fix: hand off PHP requests during configuration reloads Requests that reach a PHP handler of a superseded runtime are handed to the matching module of the new app and dispatched on the new runtime. Hand-offs happen only before execution starts, so the body is unread, no response was written, and the request can be dispatched as is. --- caddy/app.go | 2 + caddy/module.go | 9 +++ caddy/reload.go | 35 +++++++++++ caddy/reload_test.go | 138 +++++++++++++++++++++++++++++++++++++++++++ caddy/servername.go | 7 +++ scaling.go | 5 ++ testdata/reload.php | 8 +++ worker.go | 10 ++++ 8 files changed, 214 insertions(+) create mode 100644 caddy/reload.go create mode 100644 caddy/reload_test.go create mode 100644 testdata/reload.php diff --git a/caddy/app.go b/caddy/app.go index fcee129180..f5bb578ca7 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 @@ -147,6 +148,7 @@ 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 { return err diff --git a/caddy/module.go b/caddy/module.go index 20dcec9ee2..7358f2b306 100644 --- a/caddy/module.go +++ b/caddy/module.go @@ -195,6 +195,10 @@ func (f *FrankenPHPModule) ServeHTTP(w http.ResponseWriter, r *http.Request, _ c } } + if app := activeApp.Load(); app != nil && app != f.app { + return f.requestReload(w, r) + } + ctx := r.Context() repl := ctx.Value(caddy.ReplacerCtxKey).(*caddy.Replacer) @@ -226,6 +230,11 @@ func (f *FrankenPHPModule) ServeHTTP(w http.ResponseWriter, r *http.Request, _ c } err := f.server.ServeHTTP(w, r, opts...) + if errors.Is(err, frankenphp.ErrNotRunning) { + if app := activeApp.Load(); app != nil && app != f.app { + return f.requestReload(w, r) + } + } if _, rejected := errors.AsType[frankenphp.ErrRejected](err); err != nil && !rejected { return caddyhttp.Error(http.StatusInternalServerError, err) diff --git a/caddy/reload.go b/caddy/reload.go new file mode 100644 index 0000000000..04ffb11cf1 --- /dev/null +++ b/caddy/reload.go @@ -0,0 +1,35 @@ +package caddy + +import ( + "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) + } + for _, m := range app.modules { + if m.server == nil || f.server == nil || m.server.Name() != f.server.Name() { + continue + } + select { + case <-m.app.started: + case <-r.Context().Done(): + return r.Context().Err() + case <-time.After(10 * time.Second): + return caddyhttp.Error(http.StatusServiceUnavailable, frankenphp.ErrNotRunning) + } + if !m.app.hasStarted.Load() { + return caddyhttp.Error(http.StatusServiceUnavailable, frankenphp.ErrNotRunning) + } + return m.ServeHTTP(w, r, nil) + } + return caddyhttp.Error(http.StatusServiceUnavailable, frankenphp.ErrNotRunning) +} diff --git a/caddy/reload_test.go b/caddy/reload_test.go new file mode 100644 index 0000000000..a11364e047 --- /dev/null +++ b/caddy/reload_test.go @@ -0,0 +1,138 @@ +package caddy + +import ( + "encoding/json" + "fmt" + "net/http" + "net/http/httptest" + "strings" + "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) + }) + } +} diff --git a/caddy/servername.go b/caddy/servername.go index 3ca25767fd..015830fd3d 100644 --- a/caddy/servername.go +++ b/caddy/servername.go @@ -53,6 +53,13 @@ func findHostInRoutes(routes caddyhttp.RouteList, target caddyhttp.MiddlewareHan return (*hp)[0] } } + for _, handler := range route.Handlers { + if subroute, ok := handler.(*caddyhttp.Subroute); ok { + if host := findHostInRoutes(subroute.Routes, target); host != "" { + return host + } + } + } } return "" 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/worker.go b/worker.go index 388dfbd031..4123b41739 100644 --- a/worker.go +++ b/worker.go @@ -27,6 +27,7 @@ type worker struct { maxThreads int requestOptions []RequestOption requestChan chan *frankenPHPContext + done <-chan struct{} threads []*phpThread threadMutex sync.RWMutex maxConsecutiveFailures int @@ -64,6 +65,8 @@ func initWorkers(opts []workerOpt) error { return err } + w.done = mainThread.done + totalThreadsToStart += w.num workers = append(workers, w) workersByName[w.name] = w @@ -277,6 +280,13 @@ func (worker *worker) handleRequest(fc *frankenPHPContext) error { return nil case workerScaleChan <- fc: // the request has triggered scaling, continue to wait for a thread + case <-worker.done: + worker.queuedRequests.Add(-1) + metrics.DequeuedWorkerRequest(worker.name) + metrics.StopWorkerRequest(worker.name, time.Since(fc.startedAt)) + + // No thread accepted the request; the caller can retry after a reload. + return ErrNotRunning case <-timeoutChan(time.Duration(maxWaitTime.Load())): // the request has timed out stalling worker.queuedRequests.Add(-1) From 9f05b14cb7e49f0f0a3fc8c5f782b14566b54b7c Mon Sep 17 00:00:00 2001 From: henderkes Date: Fri, 2 Oct 2026 11:37:08 +0200 Subject: [PATCH 2/6] fix: use stable server names for reload handoffs --- caddy/app.go | 21 +++++++++----------- caddy/module.go | 1 + caddy/reload.go | 44 ++++++++++++++++++++++++++++++----------- caddy/serveridx_test.go | 23 +++++++++++++++++++-- 4 files changed, 63 insertions(+), 26 deletions(-) diff --git a/caddy/app.go b/caddy/app.go index f5bb578ca7..4e359e851a 100644 --- a/caddy/app.go +++ b/caddy/app.go @@ -183,22 +183,19 @@ func (f *FrankenPHPApp) Stop() error { func (f *FrankenPHPApp) registerModules(repl *caddy.Replacer) error { modulesByIndex := make(map[int]*FrankenPHPModule, len(f.modules)) for _, module := range f.modules { - if module.ServerIndex == 0 { - if err := f.registerModule(repl, module); err != nil { - return err - } - continue - } - // 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 existingModule, ok := modulesByIndex[module.ServerIndex]; ok { - module.server = existingModule.server - continue + if module.ServerIndex != 0 { + if existingModule, ok := modulesByIndex[module.ServerIndex]; ok { + module.server = existingModule.server + module.reloadName = existingModule.reloadName + continue + } + modulesByIndex[module.ServerIndex] = module } - modulesByIndex[module.ServerIndex] = module + module.reloadName = f.resolveServerName(module) if err := f.registerModule(repl, module); err != nil { return err } @@ -209,7 +206,7 @@ func (f *FrankenPHPApp) registerModules(repl *caddy.Replacer) error { // register a server instance and its workers for a single Caddy module func (f *FrankenPHPApp) registerModule(repl *caddy.Replacer, module *FrankenPHPModule) error { - serverName := f.resolveServerName(module) + serverName := module.reloadName server, err := frankenphp.NewServer( module.resolvedDocumentRoot, frankenphp.WithServerName(serverName), diff --git a/caddy/module.go b/caddy/module.go index 7358f2b306..c76386e76d 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 } diff --git a/caddy/reload.go b/caddy/reload.go index 04ffb11cf1..bd60a2e18d 100644 --- a/caddy/reload.go +++ b/caddy/reload.go @@ -15,21 +15,41 @@ func (f *FrankenPHPModule) requestReload(w http.ResponseWriter, r *http.Request) if app == nil { return caddyhttp.Error(http.StatusServiceUnavailable, frankenphp.ErrNotRunning) } - for _, m := range app.modules { - if m.server == nil || f.server == nil || m.server.Name() != f.server.Name() { + select { + case <-app.started: + case <-r.Context().Done(): + return r.Context().Err() + case <-time.After(10 * time.Second): + return caddyhttp.Error(http.StatusServiceUnavailable, frankenphp.ErrNotRunning) + } + 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) reloadModule(name string) *FrankenPHPModule { + var target *FrankenPHPModule + for _, m := range f.modules { + if m.server == nil || m.reloadName != name { continue } - select { - case <-m.app.started: - case <-r.Context().Done(): - return r.Context().Err() - case <-time.After(10 * time.Second): - return caddyhttp.Error(http.StatusServiceUnavailable, frankenphp.ErrNotRunning) + // A php_server may embed several handlers sharing one server. + if target != nil && target.server != m.server { + return nil } - if !m.app.hasStarted.Load() { - return caddyhttp.Error(http.StatusServiceUnavailable, frankenphp.ErrNotRunning) + if target == nil { + target = m } - return m.ServeHTTP(w, r, nil) } - return caddyhttp.Error(http.StatusServiceUnavailable, frankenphp.ErrNotRunning) + return target } diff --git a/caddy/serveridx_test.go b/caddy/serveridx_test.go index 7414c07169..a76b5c65c4 100644 --- a/caddy/serveridx_test.go +++ b/caddy/serveridx_test.go @@ -55,12 +55,31 @@ func TestRegisterModulesWithoutServerIndexGetOwnServers(t *testing.T) { // and the first registered module defines the server configuration func TestRegisterModulesFirstModuleWinsPerIdx(t *testing.T) { app := &FrankenPHPApp{} - first := &FrankenPHPModule{ServerIndex: 2, resolvedDocumentRoot: "../testdata"} - second := &FrankenPHPModule{ServerIndex: 2, resolvedDocumentRoot: "../testdata/env"} + first := &FrankenPHPModule{Name: "first", ServerIndex: 2, resolvedDocumentRoot: "../testdata"} + second := &FrankenPHPModule{Name: "second", ServerIndex: 2, resolvedDocumentRoot: "../testdata/env"} app.modules = []*FrankenPHPModule{first, second} require.NoError(t, app.registerModules(caddy.NewReplacer())) require.Same(t, first.server, second.server) + require.Equal(t, first.reloadName, second.reloadName) + require.Same(t, first, app.reloadModule("first")) + require.Nil(t, app.reloadModule("second")) require.Len(t, app.opts, 1, "only one server must be registered for a shared index") } + +func TestReloadModuleUsesRegisteredLabels(t *testing.T) { + first := &FrankenPHPModule{Name: "site", ServerIndex: 1, resolvedDocumentRoot: "../testdata"} + shared := &FrankenPHPModule{Name: "site", ServerIndex: 1, resolvedDocumentRoot: "../testdata"} + other := &FrankenPHPModule{Name: "other", resolvedDocumentRoot: "../testdata"} + app := &FrankenPHPApp{modules: []*FrankenPHPModule{first, shared, other}} + require.NoError(t, app.registerModules(caddy.NewReplacer())) + + first.Name = "changed" + shared.Name = "changed" + require.Same(t, first, app.reloadModule("site")) + require.Same(t, other, app.reloadModule("other")) + + other.reloadName = "site" + require.Nil(t, app.reloadModule("site")) +} From fb93bc98cb4fc12152abfed5d496c44b1b9f107d Mon Sep 17 00:00:00 2001 From: henderkes Date: Fri, 2 Oct 2026 11:37:08 +0200 Subject: [PATCH 3/6] fix: reject new PHP requests before shutdown drains --- caddy/app.go | 4 ++ caddy/module.go | 5 +- caddy/reload.go | 21 +++++-- caddy/reload_test.go | 139 +++++++++++++++++++++++++++++++++++++++++++ frankenphp.go | 5 +- 5 files changed, 165 insertions(+), 9 deletions(-) diff --git a/caddy/app.go b/caddy/app.go index 4e359e851a..bc0b396b90 100644 --- a/caddy/app.go +++ b/caddy/app.go @@ -151,6 +151,7 @@ func (f *FrankenPHPApp) Start() error { activeApp.Store(f) frankenphp.Shutdown() if err := frankenphp.Init(f.opts...); err != nil { + activeApp.CompareAndSwap(f, nil) return err } @@ -160,6 +161,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 🐘") } diff --git a/caddy/module.go b/caddy/module.go index c76386e76d..5c56d31629 100644 --- a/caddy/module.go +++ b/caddy/module.go @@ -196,10 +196,6 @@ func (f *FrankenPHPModule) ServeHTTP(w http.ResponseWriter, r *http.Request, _ c } } - if app := activeApp.Load(); app != nil && app != f.app { - return f.requestReload(w, r) - } - ctx := r.Context() repl := ctx.Value(caddy.ReplacerCtxKey).(*caddy.Replacer) @@ -235,6 +231,7 @@ func (f *FrankenPHPModule) ServeHTTP(w http.ResponseWriter, r *http.Request, _ c 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 { diff --git a/caddy/reload.go b/caddy/reload.go index bd60a2e18d..9a67989691 100644 --- a/caddy/reload.go +++ b/caddy/reload.go @@ -17,10 +17,23 @@ func (f *FrankenPHPModule) requestReload(w http.ResponseWriter, r *http.Request) } select { case <-app.started: - case <-r.Context().Done(): - return r.Context().Err() - case <-time.After(10 * time.Second): - return caddyhttp.Error(http.StatusServiceUnavailable, frankenphp.ErrNotRunning) + default: + // A draining PHP request may be waiting on this request. + wait := 10 * time.Second + if app.MaxWaitTime != 0 && app.MaxWaitTime < wait { + wait = app.MaxWaitTime + } + timer := time.NewTimer(wait) + defer timer.Stop() + select { + case <-app.started: + case <-r.Context().Done(): + return r.Context().Err() + case <-app.ctx.Done(): + return caddyhttp.Error(http.StatusServiceUnavailable, frankenphp.ErrNotRunning) + case <-timer.C: + return caddyhttp.Error(http.StatusServiceUnavailable, frankenphp.ErrMaxWaitTimeExceeded) + } } if !app.hasStarted.Load() || f.server == nil { return caddyhttp.Error(http.StatusServiceUnavailable, frankenphp.ErrNotRunning) diff --git a/caddy/reload_test.go b/caddy/reload_test.go index a11364e047..e4e33bacbb 100644 --- a/caddy/reload_test.go +++ b/caddy/reload_test.go @@ -1,11 +1,16 @@ package caddy import ( + "context" "encoding/json" "fmt" + "log/slog" "net/http" "net/http/httptest" + "os" + "path/filepath" "strings" + "sync" "testing" "time" @@ -136,3 +141,137 @@ func TestReloadPreservesRequestAndResponse(t *testing.T) { }) } } + +func newReloadTestApp(root string, worker bool) *FrankenPHPApp { + app := &FrankenPHPApp{ + NumThreads: 2, + 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 := t.TempDir() + 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(" Date: Fri, 2 Oct 2026 12:00:59 +0200 Subject: [PATCH 4/6] fix: drain retired worker queues after reload --- caddy/app.go | 3 ++ caddy/module.go | 7 ++--- caddy/reload.go | 53 +++++++++++++++++++++------------ caddy/reload_test.go | 27 +++++++++++++++++ context.go | 8 ++--- frankenphp.go | 6 +++- options.go | 10 +++++++ threadregular_test.go | 2 +- worker.go | 61 ++++++++++++++++++++++++------------- worker_internal_test.go | 66 +++++++++++++++++++++++++++++++++++++++++ 10 files changed, 192 insertions(+), 51 deletions(-) diff --git a/caddy/app.go b/caddy/app.go index bc0b396b90..5a15b58341 100644 --- a/caddy/app.go +++ b/caddy/app.go @@ -130,6 +130,9 @@ func (f *FrankenPHPApp) Start() error { frankenphp.WithMaxIdleTime(f.MaxIdleTime), frankenphp.WithMaxRequests(f.MaxRequests), ) + if f.httpApp != nil { + f.opts = append(f.opts, frankenphp.WithWorkerRequestDrainTimeout(time.Duration(f.httpApp.GracePeriod))) + } // register global workers for _, w := range f.Workers { diff --git a/caddy/module.go b/caddy/module.go index 5c56d31629..ab75ee2dbb 100644 --- a/caddy/module.go +++ b/caddy/module.go @@ -188,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 } } diff --git a/caddy/reload.go b/caddy/reload.go index 9a67989691..fef44c1eae 100644 --- a/caddy/reload.go +++ b/caddy/reload.go @@ -1,6 +1,7 @@ package caddy import ( + "context" "net/http" "time" @@ -15,25 +16,8 @@ func (f *FrankenPHPModule) requestReload(w http.ResponseWriter, r *http.Request) if app == nil { return caddyhttp.Error(http.StatusServiceUnavailable, frankenphp.ErrNotRunning) } - select { - case <-app.started: - default: - // A draining PHP request may be waiting on this request. - wait := 10 * time.Second - if app.MaxWaitTime != 0 && app.MaxWaitTime < wait { - wait = app.MaxWaitTime - } - timer := time.NewTimer(wait) - defer timer.Stop() - select { - case <-app.started: - case <-r.Context().Done(): - return r.Context().Err() - case <-app.ctx.Done(): - return caddyhttp.Error(http.StatusServiceUnavailable, frankenphp.ErrNotRunning) - case <-timer.C: - return caddyhttp.Error(http.StatusServiceUnavailable, frankenphp.ErrMaxWaitTimeExceeded) - } + if err := app.waitForStartup(r.Context()); err != nil { + return err } if !app.hasStarted.Load() || f.server == nil { return caddyhttp.Error(http.StatusServiceUnavailable, frankenphp.ErrNotRunning) @@ -50,6 +34,37 @@ func (f *FrankenPHPModule) requestReload(w http.ResponseWriter, r *http.Request) 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 { diff --git a/caddy/reload_test.go b/caddy/reload_test.go index e4e33bacbb..9803fee67a 100644 --- a/caddy/reload_test.go +++ b/caddy/reload_test.go @@ -275,3 +275,30 @@ func TestFailedReloadReturnsServiceUnavailable(t *testing.T) { require.ErrorAs(t, err, &handlerErr) require.Equal(t, http.StatusServiceUnavailable, handlerErr.StatusCode) } + +func TestStartupWaitUsesConfiguredLimits(t *testing.T) { + for _, tc := range []struct { + name string + grace, maxWait time.Duration + }{ + {"grace_period", 10 * time.Millisecond, 0}, + {"max_wait_time", time.Hour, 10 * time.Millisecond}, + } { + t.Run(tc.name, func(t *testing.T) { + app := &FrankenPHPApp{ + ctx: context.Background(), started: make(chan any), MaxWaitTime: tc.maxWait, + httpApp: &caddyhttp.App{GracePeriod: caddy.Duration(tc.grace)}, + } + ctx, cancel := context.WithTimeout(context.Background(), time.Second) + defer cancel() + err := app.waitForStartup(ctx) + require.ErrorIs(t, err, frankenphp.ErrMaxWaitTimeExceeded) + }) + } + t.Run("no_limit", func(t *testing.T) { + app := &FrankenPHPApp{ctx: context.Background(), started: make(chan any), httpApp: &caddyhttp.App{}} + ctx, cancel := context.WithCancel(context.Background()) + cancel() + require.ErrorIs(t, app.waitForStartup(ctx), context.Canceled) + }) +} diff --git a/context.go b/context.go index bed3b36e14..10de5d817c 100644 --- a/context.go +++ b/context.go @@ -51,7 +51,7 @@ type frankenPHPContext struct { handlerParameters any handlerReturn any - done chan any + done chan error startedAt time.Time } @@ -76,7 +76,7 @@ func NewRequestWithContext(r *http.Request, opts ...RequestOption) (*http.Reques func newContextFromRequest(request *http.Request, responseWriter http.ResponseWriter, s *Server, opts ...RequestOption) (*frankenPHPContext, error) { fc := &frankenPHPContext{ ctx: request.Context(), - done: make(chan any), + done: make(chan error), startedAt: time.Now(), server: s, splitPath: s.splitPath, @@ -139,7 +139,7 @@ func newWorkerDummyContext(w *worker) (*frankenPHPContext, error) { } fc := &frankenPHPContext{ - done: make(chan any), + done: make(chan error), ctx: r.Context(), server: server, request: r, @@ -172,7 +172,7 @@ func newContextFromMessage(message any, rw http.ResponseWriter, ctx context.Cont } return &frankenPHPContext{ - done: make(chan any), + done: make(chan error), startedAt: time.Now(), server: server, worker: w, diff --git a/frankenphp.go b/frankenphp.go index 671b8d79ae..ae022ef404 100644 --- a/frankenphp.go +++ b/frankenphp.go @@ -313,7 +313,7 @@ func Init(options ...Option) error { registerExtensions() - opt := &opt{} + opt := &opt{workerDrain: shutDownGracePeriod} for _, o := range options { if err := o(opt); err != nil { shutdown() @@ -336,6 +336,7 @@ func Init(options ...Option) error { } maxWaitTime.Store(int64(opt.maxWaitTime)) + workerRequestDrainTimeout = opt.workerDrain maxRequestsPerThread = opt.maxRequests if opt.maxIdleTime > 0 { @@ -447,6 +448,9 @@ func shutdown() { drainWatchers() drainPHPThreads() + for _, worker := range workers { + go worker.drainRequests(globalCtx, workerRequestDrainTimeout) + } metrics.Shutdown() diff --git a/options.go b/options.go index e1eaeb7b55..a6342b7782 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 { @@ -85,6 +86,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/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 4123b41739..46d96829e1 100644 --- a/worker.go +++ b/worker.go @@ -3,6 +3,7 @@ package frankenphp // #include "frankenphp.h" import "C" import ( + "context" "fmt" "net/http" "os" @@ -35,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 { @@ -170,6 +173,7 @@ func newWorker(o workerOpt) (*worker, error) { onThreadReady: o.onThreadReady, onThreadShutdown: o.onThreadShutdown, server: o.server, + metrics: metrics, } w.configureMercure(&o) @@ -238,7 +242,7 @@ func (worker *worker) isAtThreadLimit() bool { } func (worker *worker) handleRequest(fc *frankenPHPContext) error { - metrics.StartWorkerRequest(worker.name) + worker.metrics.StartWorkerRequest(worker.name) runtime.Gosched() @@ -250,7 +254,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: @@ -262,7 +266,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 @@ -273,25 +277,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 <-worker.done: - worker.queuedRequests.Add(-1) - metrics.DequeuedWorkerRequest(worker.name) - metrics.StopWorkerRequest(worker.name, time.Since(fc.startedAt)) - - // No thread accepted the request; the caller can retry after a reload. - return ErrNotRunning 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) @@ -299,3 +296,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 f48fc16fb2..66ad684c94 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(2), 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. From f4cfc26f608c01688038f822b2be523daed6c3d4 Mon Sep 17 00:00:00 2001 From: henderkes Date: Fri, 2 Oct 2026 13:14:38 +0200 Subject: [PATCH 5/6] perf: keep reload checks off successful requests --- caddy/module.go | 16 +++++++++------- 1 file changed, 9 insertions(+), 7 deletions(-) diff --git a/caddy/module.go b/caddy/module.go index ab75ee2dbb..199b206cb6 100644 --- a/caddy/module.go +++ b/caddy/module.go @@ -224,15 +224,17 @@ func (f *FrankenPHPModule) ServeHTTP(w http.ResponseWriter, r *http.Request, _ c } err := f.server.ServeHTTP(w, r, opts...) - if errors.Is(err, frankenphp.ErrNotRunning) { - if app := activeApp.Load(); app != nil && app != f.app { - return f.requestReload(w, r) + 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) } - 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 From 7d1dd23881f571baf9ea92f59bce0eb8ea9c5c7a Mon Sep 17 00:00:00 2001 From: henderkes Date: Fri, 2 Oct 2026 13:41:21 +0200 Subject: [PATCH 6/6] test: resolve symlinked temporary roots in reload tests --- caddy/reload_test.go | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/caddy/reload_test.go b/caddy/reload_test.go index d7957aedde..c12d77f866 100644 --- a/caddy/reload_test.go +++ b/caddy/reload_test.go @@ -188,7 +188,8 @@ func (c *reloadTestWaitingContext) Done() <-chan struct{} { func TestReloadUsesNewRootWhileShutdownDrains(t *testing.T) { for _, worker := range []bool{false, true} { t.Run(fmt.Sprintf("worker=%t", worker), func(t *testing.T) { - root := t.TempDir() + 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))