diff --git a/docs/changelog/index.md b/docs/changelog/index.md index f30a1399..aea222bb 100644 --- a/docs/changelog/index.md +++ b/docs/changelog/index.md @@ -12,6 +12,7 @@ sidebar: false ### 🛠 修复 +- 修复节点配置自动重试时,OpenResty PID 文件异常导致重复启动和端口占用的问题;自动识别并恢复现有 master,串行执行配置应用与重启,仅跟踪本实例的旧 master 及子进程,等待其退出后再启动,避免其他 nginx 实例阻塞启动与重启。配置回滚后继续按指数退避重试目标版本。 - 默认信任 Cloudflare 官方 IPv4/IPv6 网段,并允许管理员追加或覆盖可信代理 CIDR;显式 `[]` 可关闭信任。客户端地址恢复后统一用于访问日志、WAF、`X-Real-IP` 和 `X-Forwarded-For` 追加项,旧版自定义 OpenResty 主模板也会自动补入相关指令。 - 修复 Casdoor 等认证源返回字符串 `id` 时 OIDC 登录回调失败的问题;自动注册独立生成本地用户 ID,避免第三方数字 ID 与已有用户或其他认证源冲突,已有用户及绑定不受影响。 diff --git a/docs/deployment/agent.md b/docs/deployment/agent.md index fc1ace69..6071f22e 100644 --- a/docs/deployment/agent.md +++ b/docs/deployment/agent.md @@ -156,4 +156,6 @@ curl -fsSL https://raw.githubusercontent.com/Rain-kl/OpenFlare/main/scripts/unin | --- |---------------------------------------------------------------------------------------------------------| | `agent_token 和 discovery_token 不能同时为空` | 检查 `agent.json` 至少配置了一个 Token | | 节点一直离线 | 在 Agent 节点执行 `curl -I http://your-server:3000`,确认 Server 地址可达 | -| 发布后重复失败 | Agent 会阻断同一 `version + checksum` 的重复应用;在节点详情页点击「强制同步」,或重新发布新版本 | +| 发布后重复失败 | Agent 会按指数退避自动强制同步,初始间隔 10 秒,最长间隔 5 分钟,直到目标版本应用成功;回滚至旧配置也会继续重试。可在节点详情页点击「强制同步」立即触发 | +| PID 文件为空或失效,但端口仍被占用 | Agent 会核对 OpenResty master 的进程身份,恢复 PID 文件并重载现有进程;按配置路径和进程父子关系识别本实例,其他 nginx 实例不会阻塞启动。若出现多个匹配的 master,会报告错误并拒绝重复启动 | +| OpenResty 重启等待超时 | 重启会先发送优雅退出信号,最多等待 10 秒让本实例旧 master 和 worker 退出(持续跟踪已识别的 worker,即使其被重新托管);超时会保留错误,不会启动第二个实例。检查仍未结束的连接及进程后再重试 | diff --git a/internal/apps/agent/nginx/manager.go b/internal/apps/agent/nginx/manager.go index 43ac6951..f04f0092 100644 --- a/internal/apps/agent/nginx/manager.go +++ b/internal/apps/agent/nginx/manager.go @@ -112,9 +112,11 @@ func (r *OSCommandRunner) Run(ctx context.Context, name string, args ...string) // PathExecutor runs OpenResty using a configured binary and config path. type PathExecutor struct { - Path string - ConfigPath string - Runner CommandRunner + Path string + ConfigPath string + Runner CommandRunner + inspectProcesses func() ([]runtimeProcess, error) + signalProcess func(runtimeProcess, bool) error } // Test validates the current OpenResty configuration. @@ -130,21 +132,7 @@ func (e *PathExecutor) Test(ctx context.Context) error { // Reload reloads OpenResty or starts it when no runtime process is running. func (e *PathExecutor) Reload(ctx context.Context) error { - slog.Debug("running openresty reload with binary", "path", e.Path, "config", e.ConfigPath) - output, err := e.Runner.Run(ctx, e.Path, "-s", "reload", "-c", e.ConfigPath) - if err != nil { - if isOpenrestyNotRunningError(string(output)) { - slog.Warn("openresty reload reported runtime is not running, starting binary", "path", e.Path) - startOutput, startErr := e.Runner.Run(ctx, e.Path, "-c", e.ConfigPath) - if startErr != nil { - return fmt.Errorf("openresty reload failed: %w: %s; start failed: %w: %s", err, string(output), startErr, string(startOutput)) - } - return nil - } - return fmt.Errorf("openresty reload failed: %w: %s", err, string(output)) - } - slog.Debug("openresty reload succeeded with binary", "path", e.Path) - return nil + return e.recoverRuntime(ctx) } // EnsureRuntime validates configuration and reloads the OpenResty runtime. @@ -162,20 +150,7 @@ func (e *PathExecutor) CheckHealth(ctx context.Context) error { // Restart stops and starts the OpenResty runtime process. func (e *PathExecutor) Restart(ctx context.Context) error { - slog.Info("restarting openresty with binary", "path", e.Path, "config", e.ConfigPath) - output, err := e.Runner.Run(ctx, e.Path, "-s", "quit", "-c", e.ConfigPath) - if err != nil { - text := string(output) - if !isIgnorableOpenrestyStopError(text) { - return fmt.Errorf("openresty stop failed: %w: %s", err, text) - } - } - output, err = e.Runner.Run(ctx, e.Path, "-c", e.ConfigPath) - if err != nil { - return fmt.Errorf("openresty start failed: %w: %s", err, string(output)) - } - slog.Info("openresty restart succeeded with binary", "path", e.Path) - return nil + return e.restartRuntime(ctx) } // Manager applies OpenResty configuration and manages runtime assets. @@ -197,6 +172,7 @@ type Manager struct { Executor Executor atomicFileWriter func(path string, data []byte, perm os.FileMode) error wafIPGroupsMu sync.Mutex + runtimeMu sync.Mutex } // ApplyStatus reports the outcome of an OpenResty configuration apply. @@ -261,6 +237,8 @@ type wafIPGroupsRuntimeConfig struct { // Apply writes, validates, and activates new OpenResty configuration files. func (m *Manager) Apply(ctx context.Context, mainConfig string, routeConfig string, supportFiles []protocol.SupportFile) ApplyOutcome { + m.runtimeMu.Lock() + defer m.runtimeMu.Unlock() slog.Info("openresty apply started", "main_config", m.MainConfigPath, "route_config", m.RouteConfigPath, "cert_files", len(supportFiles)) backup, err := m.backup() if err != nil { @@ -364,7 +342,7 @@ func (m *Manager) rollbackAfterFailedApply(ctx context.Context, backup *backupSt if backup != nil && backup.MainExisted { return fatalApplyOutcome(fmt.Errorf("apply failed: %w; rollback recovery failed: %w", applyErr, err)) } - if fallbackErr := m.EnsureSafeFallbackRuntime(ctx, fmt.Sprintf("apply failed: %v; rollback recovery failed: %v", applyErr, err)); fallbackErr != nil { + if fallbackErr := m.ensureSafeFallbackRuntime(ctx, fmt.Sprintf("apply failed: %v; rollback recovery failed: %v", applyErr, err)); fallbackErr != nil { return fatalApplyOutcome(fmt.Errorf("apply failed: %w; rollback recovery failed: %w; fallback recovery failed: %w", applyErr, err, fallbackErr)) } message := fmt.Sprintf("apply failed, but fallback runtime started: %v; rollback recovery failed: %v", applyErr, err) @@ -426,6 +404,8 @@ func (m *Manager) EnsureLuaAssets() error { // EnsureRuntime validates and reloads the current OpenResty runtime configuration. func (m *Manager) EnsureRuntime(ctx context.Context, recreate bool) error { + m.runtimeMu.Lock() + defer m.runtimeMu.Unlock() if m.Executor == nil { return errors.New("executor 未配置") } @@ -435,6 +415,12 @@ func (m *Manager) EnsureRuntime(ctx context.Context, recreate bool) error { // EnsureSafeFallbackRuntime starts a minimal safe default OpenResty runtime. func (m *Manager) EnsureSafeFallbackRuntime(ctx context.Context, reason string) error { + m.runtimeMu.Lock() + defer m.runtimeMu.Unlock() + return m.ensureSafeFallbackRuntime(ctx, reason) +} + +func (m *Manager) ensureSafeFallbackRuntime(ctx context.Context, reason string) error { if m.Executor == nil { return errors.New("executor 未配置") } @@ -455,6 +441,8 @@ func (m *Manager) EnsureSafeFallbackRuntime(ctx context.Context, reason string) // CheckHealth verifies that OpenResty configuration and health endpoints are available. func (m *Manager) CheckHealth(ctx context.Context) error { + m.runtimeMu.Lock() + defer m.runtimeMu.Unlock() if m.Executor == nil { return errors.New("executor 未配置") } @@ -471,6 +459,8 @@ func (m *Manager) CheckHealth(ctx context.Context) error { // Restart restarts the OpenResty runtime process. func (m *Manager) Restart(ctx context.Context) error { + m.runtimeMu.Lock() + defer m.runtimeMu.Unlock() if m.Executor == nil { return errors.New("executor 未配置") } @@ -480,6 +470,8 @@ func (m *Manager) Restart(ctx context.Context) error { // CurrentChecksum returns a stable checksum for the active OpenResty configuration bundle. func (m *Manager) CurrentChecksum() (string, error) { + m.runtimeMu.Lock() + defer m.runtimeMu.Unlock() if m.RouteConfigPath == "" { return "", errors.New("route config path 不能为空") } @@ -710,7 +702,7 @@ func writeAtomicFile(path string, data []byte, perm os.FileMode) (resultErr erro resultErr = closeErr } } - _ = os.Remove(tempPath) + _ = os.Remove(tempPath) //nolint:gosec // Temp path is created by os.CreateTemp in the managed destination directory. }() if err = tempFile.Chmod(perm); err != nil { return err @@ -725,7 +717,7 @@ func writeAtomicFile(path string, data []byte, perm os.FileMode) (resultErr erro return err } closed = true - if err = os.Rename(tempPath, path); err != nil { + if err = os.Rename(tempPath, path); err != nil { //nolint:gosec // Destination is an Agent-managed runtime file (including the validated pid directive). return err } return nil @@ -817,24 +809,6 @@ func parseExtVersion(output string) string { var nginxVersionPattern = regexp.MustCompile(`(?im)(?:nginx|openresty) version:\s*(?:nginx|openresty)/(\S+)`) -func isIgnorableOpenrestyStopError(output string) bool { - text := strings.ToLower(strings.TrimSpace(output)) - if text == "" { - return false - } - return strings.Contains(text, "invalid pid") || strings.Contains(text, "no such process") -} - -func isOpenrestyNotRunningError(output string) bool { - text := strings.ToLower(strings.TrimSpace(output)) - if text == "" { - return false - } - return strings.Contains(text, "invalid pid") || - strings.Contains(text, "no such process") || - strings.Contains(text, "open()") && strings.Contains(text, "nginx.pid") && strings.Contains(text, "failed") -} - type backupState struct { MainExisted bool MainData []byte diff --git a/internal/apps/agent/nginx/manager_test.go b/internal/apps/agent/nginx/manager_test.go index a55d3061..99a331d8 100644 --- a/internal/apps/agent/nginx/manager_test.go +++ b/internal/apps/agent/nginx/manager_test.go @@ -104,9 +104,10 @@ func (e *scriptedExecutor) Restart(ctx context.Context) error { func TestPathExecutorCommands(t *testing.T) { runner := &fakeRunner{} executor := &PathExecutor{ - Path: "/usr/local/openresty/nginx/sbin/openresty", - ConfigPath: "/data/etc/nginx/nginx.conf", - Runner: runner, + Path: "/usr/local/openresty/nginx/sbin/openresty", + ConfigPath: "/data/etc/nginx/nginx.conf", + Runner: runner, + inspectProcesses: func() ([]runtimeProcess, error) { return nil, nil }, } if err := executor.Test(context.Background()); err != nil { @@ -118,7 +119,7 @@ func TestPathExecutorCommands(t *testing.T) { expected := []runCall{ {name: "/usr/local/openresty/nginx/sbin/openresty", args: []string{"-t", "-c", "/data/etc/nginx/nginx.conf"}}, - {name: "/usr/local/openresty/nginx/sbin/openresty", args: []string{"-s", "reload", "-c", "/data/etc/nginx/nginx.conf"}}, + {name: "/usr/local/openresty/nginx/sbin/openresty", args: []string{"-c", "/data/etc/nginx/nginx.conf"}}, } if !reflect.DeepEqual(runner.calls, expected) { t.Fatalf("unexpected calls: %#v", runner.calls) @@ -128,9 +129,10 @@ func TestPathExecutorCommands(t *testing.T) { func TestPathExecutorEnsureRuntimeNoop(t *testing.T) { runner := &fakeRunner{} executor := &PathExecutor{ - Path: "/usr/local/openresty/nginx/sbin/openresty", - ConfigPath: "/data/etc/nginx/nginx.conf", - Runner: runner, + Path: "/usr/local/openresty/nginx/sbin/openresty", + ConfigPath: "/data/etc/nginx/nginx.conf", + Runner: runner, + inspectProcesses: func() ([]runtimeProcess, error) { return nil, nil }, } if err := executor.EnsureRuntime(context.Background(), true); err != nil { t.Fatalf("EnsureRuntime failed: %v", err) @@ -140,47 +142,34 @@ func TestPathExecutorEnsureRuntimeNoop(t *testing.T) { } } -func TestPathExecutorRestartIgnoresMissingPID(t *testing.T) { - runner := &fakeRunner{ - runFn: func(name string, args ...string) ([]byte, error) { - if len(args) == 2 && args[0] == "-s" && args[1] == "quit" { - return []byte("openresty: [error] invalid PID number \"\" in \"/usr/local/openresty/nginx/logs/nginx.pid\""), errors.New("exit status 1") - } - return []byte(""), nil - }, - } +func TestPathExecutorRestartStartsWhenNoProcesses(t *testing.T) { + runner := &fakeRunner{} executor := &PathExecutor{ - Path: "/usr/local/openresty/nginx/sbin/openresty", - ConfigPath: "/data/etc/nginx/nginx.conf", - Runner: runner, + Path: "/usr/local/openresty/nginx/sbin/openresty", + ConfigPath: "/data/etc/nginx/nginx.conf", + Runner: runner, + inspectProcesses: func() ([]runtimeProcess, error) { return nil, nil }, } if err := executor.Restart(context.Background()); err != nil { t.Fatalf("Restart failed: %v", err) } - if len(runner.calls) != 2 { - t.Fatalf("expected 2 restart calls, got %d", len(runner.calls)) + if len(runner.calls) != 1 { + t.Fatalf("expected 1 start call, got %d", len(runner.calls)) } } func TestPathExecutorReloadStartsWhenRuntimeIsNotRunning(t *testing.T) { - runner := &fakeRunner{ - runFn: func(name string, args ...string) ([]byte, error) { - if len(args) >= 2 && args[0] == "-s" && args[1] == "reload" { - return []byte("openresty: [error] invalid PID number \"\" in \"/usr/local/openresty/nginx/logs/nginx.pid\""), errors.New("exit status 1") - } - return []byte(""), nil - }, - } + runner := &fakeRunner{} executor := &PathExecutor{ - Path: "/usr/local/openresty/nginx/sbin/openresty", - ConfigPath: "/data/etc/nginx/nginx.conf", - Runner: runner, + Path: "/usr/local/openresty/nginx/sbin/openresty", + ConfigPath: "/data/etc/nginx/nginx.conf", + Runner: runner, + inspectProcesses: func() ([]runtimeProcess, error) { return nil, nil }, } if err := executor.Reload(context.Background()); err != nil { t.Fatalf("Reload failed: %v", err) } expected := []runCall{ - {name: "/usr/local/openresty/nginx/sbin/openresty", args: []string{"-s", "reload", "-c", "/data/etc/nginx/nginx.conf"}}, {name: "/usr/local/openresty/nginx/sbin/openresty", args: []string{"-c", "/data/etc/nginx/nginx.conf"}}, } if !reflect.DeepEqual(runner.calls, expected) { diff --git a/internal/apps/agent/nginx/runtime_process.go b/internal/apps/agent/nginx/runtime_process.go new file mode 100644 index 00000000..e437a690 --- /dev/null +++ b/internal/apps/agent/nginx/runtime_process.go @@ -0,0 +1,219 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package nginx + +import ( + "context" + "errors" + "fmt" + "log/slog" + "os" + "path/filepath" + "regexp" + "strconv" + "strings" + "time" +) + +const ( + runtimeShutdownTimeout = 10 * time.Second + runtimeProcessPollInterval = 50 * time.Millisecond +) + +// runtimeProcess records an inspected process identity, including its start time +// to reject a PID that has since been reused. +type runtimeProcess struct { + PID int + ParentPID int + Master bool + StartTime string +} + +// runtimeProcesses returns candidates; only verified masters and their tracked +// descendants establish membership of this instance. +func (e *PathExecutor) runtimeProcesses(ctx context.Context) ([]runtimeProcess, error) { + if e.inspectProcesses != nil { + return e.inspectProcesses() + } + return inspectOpenrestyProcesses(ctx, e.Path, e.ConfigPath) +} + +func (e *PathExecutor) signalRuntime(ctx context.Context, process runtimeProcess, quit bool) error { + if e.signalProcess != nil { + return e.signalProcess(process, quit) + } + return signalOpenrestyProcess(ctx, e.Path, e.ConfigPath, process, quit) +} + +func findRuntimeMaster(processes []runtimeProcess) (runtimeProcess, error) { + var master runtimeProcess + for _, process := range processes { + if !process.Master { + continue + } + if master.PID != 0 { + return runtimeProcess{}, errors.New("multiple matching openresty masters; refusing ambiguous runtime recovery") + } + master = process + } + return master, nil +} + +// runtimeMasterIdentities seeds instance membership from verified masters. +func runtimeMasterIdentities(processes []runtimeProcess) map[int]string { + identities := make(map[int]string) + for _, process := range processes { + if process.Master { + identities[process.PID] = process.StartTime + } + } + return identities +} + +// runtimeDescendants selects known identities and their children. Keeping the +// identities between shutdown polls also recognizes workers reparented to init. +func runtimeDescendants(processes []runtimeProcess, identities map[int]string) []runtimeProcess { + selected := make(map[int]bool) + for _, process := range processes { + if start, ok := identities[process.PID]; ok && start == process.StartTime { + selected[process.PID] = true + } + } + for changed := true; changed; { + changed = false + for _, process := range processes { + if !selected[process.PID] && selected[process.ParentPID] { + selected[process.PID] = true + changed = true + } + } + } + var result []runtimeProcess + for _, process := range processes { + if selected[process.PID] { + identities[process.PID] = process.StartTime + result = append(result, process) + } + } + return result +} + +func (e *PathExecutor) recoverRuntime(ctx context.Context) error { + if err := ctx.Err(); err != nil { + return err + } + processes, err := e.runtimeProcesses(ctx) + if err != nil { + return fmt.Errorf("inspect openresty runtime: %w", err) + } + master, err := findRuntimeMaster(processes) + if err != nil { + return err + } + if master.PID != 0 { + repaired, err := e.repairPIDFile(master.PID) + if err != nil { + return err + } + if err := e.signalRuntime(ctx, master, false); err != nil { + return fmt.Errorf("reload recovered openresty master: %w", err) + } + if repaired { + slog.WarnContext(ctx, "recovered openresty master with invalid pid file", "pid", master.PID, "config", e.ConfigPath) + } + return nil + } + return e.startRuntime(ctx) +} + +func (e *PathExecutor) startRuntime(ctx context.Context) error { + output, err := e.Runner.Run(ctx, e.Path, "-c", e.ConfigPath) + if err != nil { + return fmt.Errorf("openresty start failed: %w: %s", err, output) + } + return nil +} + +var pidDirectivePattern = regexp.MustCompile(`(?m)^\s*pid\s+([^;\n]+);`) + +func (e *PathExecutor) repairPIDFile(pid int) (bool, error) { + config, err := os.ReadFile(e.ConfigPath) + if err != nil { + return false, fmt.Errorf("read config for pid recovery: %w", err) + } + matches := pidDirectivePattern.FindAllSubmatch(config, -1) + if len(matches) != 1 { + return false, errors.New("pid recovery requires one explicit pid directive") + } + path := strings.Trim(string(matches[0][1]), " \t\"'") + if !filepath.IsAbs(path) { + return false, errors.New("pid recovery requires an absolute pid path") + } + expected := strconv.Itoa(pid) + "\n" + //nolint:gosec // PID path comes from the validated, Agent-managed main configuration. + if data, err := os.ReadFile(path); err == nil && string(data) == expected { + return false, nil + } + if err := writeAtomicFile(path, []byte(strconv.Itoa(pid)+"\n"), nginxConfigFilePerm); err != nil { + return false, fmt.Errorf("repair openresty pid file: %w", err) + } + return true, nil +} + +func (e *PathExecutor) restartRuntime(ctx context.Context) error { + if err := ctx.Err(); err != nil { + return err + } + processes, err := e.runtimeProcesses(ctx) + if err != nil { + return err + } + master, err := findRuntimeMaster(processes) + if err != nil { + return err + } + if master.PID == 0 { + return e.startRuntime(ctx) + } + identities := runtimeMasterIdentities(processes) + runtimeDescendants(processes, identities) + if err := e.signalRuntime(ctx, master, true); err != nil { + return fmt.Errorf("stop openresty master: %w", err) + } + // A successful signal only acknowledges delivery, not process termination. + // Never create a second master while the old master or workers remain. + waitCtx, cancel := context.WithTimeout(ctx, runtimeShutdownTimeout) + defer cancel() + ticker := time.NewTicker(runtimeProcessPollInterval) + defer ticker.Stop() + for { + remaining, err := e.runtimeProcesses(waitCtx) + if err != nil { + return err + } + if len(runtimeDescendants(remaining, identities)) == 0 { + break + } + select { + case <-waitCtx.Done(): + return fmt.Errorf("waiting for openresty processes to exit before restart: %w", waitCtx.Err()) + case <-ticker.C: + } + } + return e.startRuntime(ctx) +} + +func matchesRuntimeMaster(command, configPath string) bool { + command = strings.ReplaceAll(command, "\x00", " ") + if !strings.HasPrefix(command, "nginx: master process ") && !strings.HasPrefix(command, "openresty: master process ") { + return false + } + fields := strings.Fields(command) + for i := 0; i+1 < len(fields); i++ { + if fields[i] == "-c" { + return filepath.IsAbs(fields[i+1]) && filepath.Clean(fields[i+1]) == filepath.Clean(configPath) + } + } + return false +} diff --git a/internal/apps/agent/nginx/runtime_process_linux.go b/internal/apps/agent/nginx/runtime_process_linux.go new file mode 100644 index 00000000..d78138f0 --- /dev/null +++ b/internal/apps/agent/nginx/runtime_process_linux.go @@ -0,0 +1,44 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package nginx + +import ( + "context" + "fmt" + "os" + "os/exec" + "syscall" +) + +func inspectOpenrestyProcesses(ctx context.Context, binary, configPath string) ([]runtimeProcess, error) { + if err := ctx.Err(); err != nil { + return nil, err + } + path, err := exec.LookPath(binary) + if err != nil { + return nil, err + } + executable, err := os.Stat(path) + if err != nil { + return nil, err + } + return inspectProcRuntime("/proc", executable, binary, configPath, os.Geteuid()) +} + +func signalOpenrestyProcess(ctx context.Context, binary, config string, process runtimeProcess, quit bool) error { + processes, err := inspectOpenrestyProcesses(ctx, binary, config) + if err != nil { + return err + } + for _, current := range processes { + if current.PID == process.PID && current.Master && current.StartTime == process.StartTime { + signal := syscall.SIGHUP + if quit { + signal = syscall.SIGQUIT + } + return syscall.Kill(current.PID, signal) + } + } + return fmt.Errorf("openresty master identity changed before signalling pid %d", process.PID) +} diff --git a/internal/apps/agent/nginx/runtime_process_other.go b/internal/apps/agent/nginx/runtime_process_other.go new file mode 100644 index 00000000..67fec8a3 --- /dev/null +++ b/internal/apps/agent/nginx/runtime_process_other.go @@ -0,0 +1,19 @@ +//go:build !unix + +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package nginx + +import ( + "context" + "errors" +) + +func inspectOpenrestyProcesses(_ context.Context, _, _ string) ([]runtimeProcess, error) { + return nil, errors.New("safe openresty process recovery is unsupported on this platform; refusing an unverified start") +} + +func signalOpenrestyProcess(_ context.Context, _, _ string, _ runtimeProcess, _ bool) error { + return errors.New("safe openresty process signalling is unsupported on this platform") +} diff --git a/internal/apps/agent/nginx/runtime_process_proc.go b/internal/apps/agent/nginx/runtime_process_proc.go new file mode 100644 index 00000000..def23fa7 --- /dev/null +++ b/internal/apps/agent/nginx/runtime_process_proc.go @@ -0,0 +1,118 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package nginx + +import ( + "fmt" + "os" + "path/filepath" + "strconv" + "strings" +) + +//nolint:unused // Used by the Linux process inspector. +const procParentPIDIndex = 1 + +//nolint:unused // Used by the Linux process inspector. +const procStartTimeIndex = 19 // Field 22, counting from state after the process name. + +//nolint:unused // Used by the Linux executor; portable for process fixture tests. +func inspectProcRuntime(root string, executable os.FileInfo, binary, configPath string, uid int) ([]runtimeProcess, error) { + entries, err := os.ReadDir(root) + if err != nil { + return nil, err + } + var processes []runtimeProcess + for _, entry := range entries { + pid, err := strconv.Atoi(entry.Name()) + if err != nil || pid <= 0 || !entry.IsDir() { + continue + } + process, found, err := inspectProcProcess(filepath.Join(root, entry.Name()), pid, executable, binary, configPath, uid) + if err != nil { + if os.IsNotExist(err) { + continue + } + return nil, err + } + if found { + processes = append(processes, process) + } + } + return processes, nil +} + +//nolint:unused // Used by the Linux process inspector. +func inspectProcProcess(dir string, pid int, executable os.FileInfo, binary, configPath string, uid int) (runtimeProcess, bool, error) { + //nolint:gosec // Read a fixed proc entry under a numeric PID directory. + status, err := os.ReadFile(filepath.Join(dir, "status")) + if err != nil { + return runtimeProcess{}, false, err + } + owned := effectiveProcessUID(string(status)) == uid + //nolint:gosec // Read a fixed proc entry under a numeric PID directory. + command, err := os.ReadFile(filepath.Join(dir, "cmdline")) + if err != nil { + return runtimeProcess{}, false, err + } + title := strings.ReplaceAll(string(command), "\x00", " ") + if !strings.HasPrefix(title, "nginx: ") && !strings.HasPrefix(title, "openresty: ") { + return runtimeProcess{}, false, nil + } + master := owned && matchesRuntimeMaster(title, configPath) + exe, err := os.Stat(filepath.Join(dir, "exe")) + if err != nil { + if !os.IsPermission(err) { + return runtimeProcess{}, false, err + } + // Non-master candidates are only retained by runtimeDescendants when + // their parent belongs to this instance. A capability-enabled non-root + // master may be non-dumpable, making /proc/PID/exe inaccessible; verify + // its invocation and UID instead. + fields := strings.Fields(title) + if master && (len(fields) < 4 || fields[3] != binary) { + return runtimeProcess{}, false, fmt.Errorf("cannot verify executable of openresty master %d", pid) + } + } else if !os.SameFile(executable, exe) { + if master { + return runtimeProcess{}, false, fmt.Errorf("executable changed for openresty master %d; refusing duplicate start", pid) + } + return runtimeProcess{}, false, nil + } + //nolint:gosec // Read a fixed proc entry under a numeric PID directory. + stat, err := os.ReadFile(filepath.Join(dir, "stat")) + if err != nil { + return runtimeProcess{}, false, err + } + closeParen := strings.LastIndexByte(string(stat), ')') + if closeParen < 0 { + return runtimeProcess{}, false, fmt.Errorf("invalid stat for process %d", pid) + } + fields := strings.Fields(string(stat)[closeParen+1:]) + if len(fields) <= procStartTimeIndex { + return runtimeProcess{}, false, fmt.Errorf("incomplete stat for process %d", pid) + } + if fields[0] == "Z" { + return runtimeProcess{}, false, nil + } + parentPID, err := strconv.Atoi(fields[procParentPIDIndex]) + if err != nil { + return runtimeProcess{}, false, fmt.Errorf("invalid parent pid for process %d: %w", pid, err) + } + return runtimeProcess{PID: pid, ParentPID: parentPID, Master: master, StartTime: fields[procStartTimeIndex]}, true, nil +} + +//nolint:unused // Used by the Linux process inspector. +func effectiveProcessUID(status string) int { + for _, line := range strings.Split(status, "\n") { + fields := strings.Fields(line) + if len(fields) >= 3 && fields[0] == "Uid:" { + uid, err := strconv.Atoi(fields[2]) + if err == nil { + return uid + } + } + } + return -1 +} diff --git a/internal/apps/agent/nginx/runtime_process_proc_test.go b/internal/apps/agent/nginx/runtime_process_proc_test.go new file mode 100644 index 00000000..a8fa84a0 --- /dev/null +++ b/internal/apps/agent/nginx/runtime_process_proc_test.go @@ -0,0 +1,172 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package nginx + +import ( + "fmt" + "os" + "path/filepath" + "strconv" + "strings" + "testing" +) + +func TestInspectProcRuntime(t *testing.T) { + root := t.TempDir() + binary := filepath.Join(root, "openresty") + if err := os.WriteFile(binary, nil, 0o755); err != nil { + t.Fatal(err) + } + executable, err := os.Stat(binary) + if err != nil { + t.Fatal(err) + } + config := filepath.Join(root, "nginx.conf") + writeProcFixture(t, root, 42, 1, 1000, "nginx: master process "+binary+" -c "+config, "S", binary) + // Root masters may run workers as a different user. Track them during shutdown. + writeProcFixture(t, root, 43, 42, 1001, "nginx: worker process", "S", binary) + writeProcFixture(t, root, 44, 42, 1000, "nginx: worker process", "Z", binary) + writeProcFixture(t, root, 45, 1, 1000, "unrelated command", "S", binary) + // Another instance using the same executable must not block this instance. + writeProcFixture(t, root, 48, 1, 1000, "nginx: master process "+binary+" -c /other/nginx.conf", "S", binary) + writeProcFixture(t, root, 49, 48, 1000, "nginx: worker process", "S", binary) + otherBinary := filepath.Join(root, "other") + if err := os.WriteFile(otherBinary, nil, 0o755); err != nil { + t.Fatal(err) + } + writeProcFixture(t, root, 46, 1, 1000, "nginx: worker process", "S", otherBinary) + processes, err := inspectProcRuntime(root, executable, binary, config, 1000) + if err != nil { + t.Fatal(err) + } + processes = runtimeDescendants(processes, runtimeMasterIdentities(processes)) + if len(processes) != 2 { + t.Fatalf("inspected processes = %v, want live master and worker", processes) + } + if processes[0] != (runtimeProcess{PID: 42, ParentPID: 1, Master: true, StartTime: "142"}) { + t.Errorf("master identity = %v, want PID 42 and start time 142", processes[0]) + } + if processes[1].PID != 43 || processes[1].Master { + t.Errorf("worker identity = %v, want worker 43", processes[1]) + } + if err := os.RemoveAll(filepath.Join(root, "42")); err != nil { + t.Fatal(err) + } + remaining, err := inspectProcRuntime(root, executable, binary, config, 1000) + if err != nil { + t.Fatal(err) + } + if master, err := findRuntimeMaster(remaining); err != nil || master.PID != 0 { + t.Errorf("unattributed worker blocks cold start: master=%v err=%v", master, err) + } + writeProcFixture(t, root, 47, 1, 1000, "nginx: master process "+binary+" -c "+config, "S", otherBinary) + if _, err := inspectProcRuntime(root, executable, binary, config, 1000); err == nil { + t.Error("inspectProcRuntime(replaced executable) = nil, want refusal") + } +} + +func TestInspectProcRuntimeIgnoresDifferentConfigOrOwner(t *testing.T) { + for _, tc := range []struct { + name string + uid int + config string + }{ + {name: "different config", uid: 1000, config: "/other/nginx.conf"}, + {name: "different owner", uid: 1001, config: "/data/nginx.conf"}, + } { + t.Run(tc.name, func(t *testing.T) { + root := t.TempDir() + binary := filepath.Join(root, "openresty") + if err := os.WriteFile(binary, nil, 0o755); err != nil { + t.Fatal(err) + } + executable, err := os.Stat(binary) + if err != nil { + t.Fatal(err) + } + writeProcFixture(t, root, 42, 1, tc.uid, "nginx: master process "+binary+" -c "+tc.config, "S", binary) + processes, err := inspectProcRuntime(root, executable, binary, "/data/nginx.conf", 1000) + if err != nil { + t.Fatal(err) + } + if master, err := findRuntimeMaster(processes); err != nil || master.PID != 0 { + t.Errorf("unrelated master blocks cold start: master=%v err=%v", master, err) + } + }) + } +} + +func writeProcFixture(t *testing.T, root string, pid, parentPID, uid int, title, state, binary string) { + t.Helper() + dir := filepath.Join(root, strconv.Itoa(pid)) + if err := os.Mkdir(dir, 0o755); err != nil { + t.Fatal(err) + } + fields := make([]string, 20) + for i := range fields { + fields[i] = "0" + } + fields[0] = state + fields[procParentPIDIndex] = strconv.Itoa(parentPID) + fields[19] = strconv.Itoa(pid + 100) + for name, content := range map[string]string{ + "status": fmt.Sprintf("Uid:\t%d\t%d\t%d\t%d\n", uid, uid, uid, uid), + "cmdline": title + "\x00", + "stat": fmt.Sprintf("%d (nginx worker) %s\n", pid, strings.Join(fields, " ")), + } { + if err := os.WriteFile(filepath.Join(dir, name), []byte(content), 0o644); err != nil { + t.Fatal(err) + } + } + if err := os.Symlink(binary, filepath.Join(dir, "exe")); err != nil { + t.Fatal(err) + } +} + +func TestInspectProcRuntimeScopesUnreadableExecutables(t *testing.T) { + if os.Geteuid() == 0 { + t.Skip("root can read the permission-denied fixture") + } + root := t.TempDir() + binary := filepath.Join(root, "openresty") + if err := os.WriteFile(binary, nil, 0o755); err != nil { + t.Fatal(err) + } + executable, err := os.Stat(binary) + if err != nil { + t.Fatal(err) + } + denied := filepath.Join(root, "denied") + if err := os.Mkdir(denied, 0o700); err != nil { + t.Fatal(err) + } + inaccessible := filepath.Join(denied, "exe") + if err := os.Symlink(binary, inaccessible); err != nil { + t.Fatal(err) + } + config := "/data/nginx.conf" + writeProcFixture(t, root, 42, 1, 1000, "nginx: master process "+binary+" -c "+config, "S", inaccessible) + writeProcFixture(t, root, 43, 42, 1001, "nginx: worker process", "S", inaccessible) + writeProcFixture(t, root, 50, 1, 1000, "nginx: master process "+binary+" -c /other/nginx.conf", "S", inaccessible) + writeProcFixture(t, root, 51, 50, 1001, "nginx: worker process", "S", inaccessible) + if err := os.Chmod(denied, 0); err != nil { + t.Fatal(err) + } + t.Cleanup(func() { + if err := os.Chmod(denied, 0o700); err != nil { + t.Error(err) + } + }) + if _, err := os.Stat(filepath.Join(root, "51", "exe")); !os.IsPermission(err) { + t.Fatalf("fixture exe read = %v, want permission denied", err) + } + processes, err := inspectProcRuntime(root, executable, binary, config, 1000) + if err != nil { + t.Fatal(err) + } + processes = runtimeDescendants(processes, runtimeMasterIdentities(processes)) + if len(processes) != 2 || processes[0].PID != 42 || processes[1].PID != 43 { + t.Fatalf("instance processes = %v, want own master and worker", processes) + } +} diff --git a/internal/apps/agent/nginx/runtime_process_ps.go b/internal/apps/agent/nginx/runtime_process_ps.go new file mode 100644 index 00000000..7f5ca345 --- /dev/null +++ b/internal/apps/agent/nginx/runtime_process_ps.go @@ -0,0 +1,49 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package nginx + +import ( + "fmt" + "strconv" + "strings" +) + +//nolint:unused // Used by the non-Linux Unix inspector. +const ( + psCommandIndex = 8 + psMasterBinaryIndex = 11 +) + +//nolint:unused // Used by the non-Linux Unix inspector; portable for fixture tests. +func parsePSRuntime(output, binary, configPath string, owner int) ([]runtimeProcess, error) { + var processes []runtimeProcess + for _, line := range strings.Split(output, "\n") { + fields := strings.Fields(line) + if len(fields) <= psCommandIndex { + continue + } + uid, err := strconv.Atoi(fields[2]) + if err != nil { + continue + } + title := strings.Join(fields[psCommandIndex:], " ") + if !strings.HasPrefix(title, "nginx: ") && !strings.HasPrefix(title, "openresty: ") { + continue + } + master := uid == owner && matchesRuntimeMaster(title, configPath) + if master && (len(fields) <= psMasterBinaryIndex || fields[psMasterBinaryIndex] != binary) { + return nil, fmt.Errorf("cannot verify openresty master invocation %s", fields[0]) + } + pid, err := strconv.Atoi(fields[0]) + if err != nil { + return nil, err + } + parentPID, err := strconv.Atoi(fields[1]) + if err != nil { + return nil, err + } + processes = append(processes, runtimeProcess{PID: pid, ParentPID: parentPID, Master: master, StartTime: strings.Join(fields[3:psCommandIndex], " ")}) + } + return processes, nil +} diff --git a/internal/apps/agent/nginx/runtime_process_ps_test.go b/internal/apps/agent/nginx/runtime_process_ps_test.go new file mode 100644 index 00000000..9ee04ce7 --- /dev/null +++ b/internal/apps/agent/nginx/runtime_process_ps_test.go @@ -0,0 +1,24 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package nginx + +import "testing" + +func TestParsePSRuntimeInstanceMembership(t *testing.T) { + output := `42 1 1000 Wed Oct 7 16:00:00 2026 nginx: master process openresty -c /data/nginx.conf +43 42 1001 Wed Oct 7 16:00:01 2026 nginx: worker process +50 1 1000 Wed Oct 7 16:00:00 2026 nginx: master process openresty -c /other/nginx.conf +51 50 1001 Wed Oct 7 16:00:01 2026 nginx: worker process` + processes, err := parsePSRuntime(output, "openresty", "/data/nginx.conf", 1000) + if err != nil { + t.Fatal(err) + } + processes = runtimeDescendants(processes, runtimeMasterIdentities(processes)) + if len(processes) != 2 || processes[0].PID != 42 || processes[1].PID != 43 || processes[1].ParentPID != 42 { + t.Fatalf("instance processes = %v, want master 42 and worker 43", processes) + } + if processes[0].StartTime != "Wed Oct 7 16:00:00 2026" { + t.Errorf("master start time = %q", processes[0].StartTime) + } +} diff --git a/internal/apps/agent/nginx/runtime_process_test.go b/internal/apps/agent/nginx/runtime_process_test.go new file mode 100644 index 00000000..272bb50c --- /dev/null +++ b/internal/apps/agent/nginx/runtime_process_test.go @@ -0,0 +1,258 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package nginx + +import ( + "context" + "errors" + "os" + "path/filepath" + "strings" + "testing" + "time" +) + +func TestReloadRecoversLiveMasterWithEmptyPID(t *testing.T) { + dir := t.TempDir() + config := filepath.Join(dir, "nginx.conf") + pidPath := filepath.Join(dir, "nginx.pid") + if err := os.WriteFile(config, []byte("pid "+pidPath+";\n"), 0o644); err != nil { + t.Fatal(err) + } + if err := os.WriteFile(pidPath, nil, 0o644); err != nil { + t.Fatal(err) + } + runner := &fakeRunner{runFn: func(_ string, _ ...string) ([]byte, error) { + return []byte(`invalid PID number ""`), errors.New("exit status 1") + }} + master := runtimeProcess{PID: 42, Master: true, StartTime: "123"} + var signalled runtimeProcess + executor := &PathExecutor{Path: "openresty", ConfigPath: config, Runner: runner, + inspectProcesses: func() ([]runtimeProcess, error) { return []runtimeProcess{master}, nil }, + signalProcess: func(process runtimeProcess, quit bool) error { + if quit { + t.Error("Reload requested quit, want reload") + } + signalled = process + return nil + }, + } + if err := executor.Reload(context.Background()); err != nil { + t.Fatalf("Reload(empty PID, live master) = %v, want nil", err) + } + if signalled != master { + t.Errorf("Reload signalled %v, want %v", signalled, master) + } + if len(runner.calls) != 0 { + t.Errorf("Reload commands = %v, want direct master signal and no start", runner.calls) + } + data, err := os.ReadFile(pidPath) + if err != nil { + t.Fatal(err) + } + if string(data) != "42\n" { + t.Errorf("recovered PID = %q, want 42 newline", data) + } +} + +func TestReloadStartsWithUnrelatedWorkers(t *testing.T) { + runner := &fakeRunner{runFn: func(_ string, _ ...string) ([]byte, error) { return []byte("invalid pid"), errors.New("failed") }} + executor := &PathExecutor{Path: "openresty", Runner: runner, + inspectProcesses: func() ([]runtimeProcess, error) { return []runtimeProcess{{PID: 99}}, nil }, + } + if err := executor.Reload(context.Background()); err == nil || !strings.Contains(err.Error(), "start failed") { + t.Errorf("Reload(unrelated workers) = %v, want attempted start", err) + } + if len(runner.calls) != 1 { + t.Errorf("Reload commands = %v, want one start", runner.calls) + } +} + +func TestRestartWaitsForOldWorkersBeforeStart(t *testing.T) { + inspections := 0 + quit := false + runner := &fakeRunner{runFn: func(_ string, _ ...string) ([]byte, error) { + if !quit || inspections < 3 { + t.Errorf("start before shutdown completed: quit=%v inspections=%d", quit, inspections) + } + return nil, nil + }} + executor := &PathExecutor{Path: "openresty", Runner: runner, + inspectProcesses: func() ([]runtimeProcess, error) { + inspections++ + if inspections == 1 { + return []runtimeProcess{{PID: 42, Master: true}, {PID: 43, ParentPID: 42}}, nil + } + if inspections == 2 { + return []runtimeProcess{{PID: 43, ParentPID: 1}}, nil + } + return nil, nil + }, + signalProcess: func(_ runtimeProcess, gracefulQuit bool) error { quit = gracefulQuit; return nil }, + } + if err := executor.Restart(context.Background()); err != nil { + t.Fatalf("Restart(draining workers) = %v, want nil", err) + } + if len(runner.calls) != 1 { + t.Errorf("Restart commands=%v, want one start", runner.calls) + } +} + +func TestRestartCancellationNeverStartsSecondMaster(t *testing.T) { + runner := &fakeRunner{} + executor := &PathExecutor{Path: "openresty", Runner: runner, + inspectProcesses: func() ([]runtimeProcess, error) { return []runtimeProcess{{PID: 42, Master: true}}, nil }, + signalProcess: func(_ runtimeProcess, _ bool) error { return nil }, + } + ctx, cancel := context.WithCancel(context.Background()) + cancel() + if err := executor.Restart(ctx); !errors.Is(err, context.Canceled) { + t.Errorf("Restart(cancelled)=%v, want context.Canceled", err) + } + if len(runner.calls) != 0 { + t.Errorf("Restart(cancelled) commands=%v, want no start", runner.calls) + } +} + +func TestRestartTimeoutNeverStartsWhileWorkersRemain(t *testing.T) { + runner := &fakeRunner{} + quits := 0 + executor := &PathExecutor{Path: "openresty", Runner: runner, + inspectProcesses: func() ([]runtimeProcess, error) { + if quits == 0 { + return []runtimeProcess{{PID: 42, Master: true}, {PID: 43, ParentPID: 42}}, nil + } + return []runtimeProcess{{PID: 43, ParentPID: 1}}, nil + }, + signalProcess: func(_ runtimeProcess, quit bool) error { + if quit { + quits++ + } + return nil + }, + } + ctx, cancel := context.WithTimeout(context.Background(), 20*time.Millisecond) + defer cancel() + if err := executor.Restart(ctx); !errors.Is(err, context.DeadlineExceeded) { + t.Errorf("Restart(draining worker, deadline) = %v, want deadline exceeded", err) + } + if quits != 1 || len(runner.calls) != 0 { + t.Errorf("Restart quit count=%d commands=%v, want one quit and no start", quits, runner.calls) + } +} + +func TestReloadRefusesAmbiguousMasters(t *testing.T) { + runner := &fakeRunner{} + executor := &PathExecutor{Runner: runner, + inspectProcesses: func() ([]runtimeProcess, error) { + return []runtimeProcess{{PID: 42, Master: true}, {PID: 43, Master: true}}, nil + }, + } + if err := executor.Reload(context.Background()); err == nil { + t.Error("Reload(two masters) = nil, want refusal") + } + if len(runner.calls) != 0 { + t.Errorf("Reload commands = %v, want no start", runner.calls) + } +} + +type blockingRuntimeExecutor struct { + fakeExecutor + entered chan struct{} + release chan struct{} + restarted chan struct{} +} + +func (e *blockingRuntimeExecutor) Reload(context.Context) error { + close(e.entered) + <-e.release + return nil +} +func (e *blockingRuntimeExecutor) Restart(context.Context) error { close(e.restarted); return nil } + +func TestManagerApplySerializesRestart(t *testing.T) { + dir := t.TempDir() + executor := &blockingRuntimeExecutor{entered: make(chan struct{}), release: make(chan struct{}), restarted: make(chan struct{})} + manager := &Manager{MainConfigPath: filepath.Join(dir, "etc", "nginx.conf"), RouteConfigPath: filepath.Join(dir, "etc", "routes.conf"), Executor: executor} + applied := make(chan ApplyOutcome, 1) + go func() { applied <- manager.Apply(context.Background(), "events {}\nhttp {}", "", nil) }() + select { + case <-executor.entered: + case <-time.After(3 * time.Second): + t.Fatal("Apply did not reach reload") + } + restarted := make(chan error, 1) + go func() { restarted <- manager.Restart(context.Background()) }() + select { + case <-executor.restarted: + t.Error("Restart overlapped Apply, want serialized lifecycle") + case <-time.After(30 * time.Millisecond): + } + close(executor.release) + if result := <-applied; result.Status != ApplyStatusSuccess { + t.Errorf("Apply status=%v message=%s, want success", result.Status, result.Message) + } + if err := <-restarted; err != nil { + t.Errorf("Restart(after Apply)=%v, want nil", err) + } +} + +func TestRepairPIDRefusesRelativeOrMissingPath(t *testing.T) { + config := filepath.Join(t.TempDir(), "nginx.conf") + if err := os.WriteFile(config, []byte("pid logs/nginx.pid;"), 0o644); err != nil { + t.Fatal(err) + } + executor := &PathExecutor{ConfigPath: config} + if _, err := executor.repairPIDFile(42); err == nil || !strings.Contains(err.Error(), "absolute") { + t.Errorf("repairPIDFile(relative path)=%v, want absolute path error", err) + } +} + +func TestRuntimeDescendantsTracksOnlyInstanceIdentities(t *testing.T) { + processes := []runtimeProcess{ + {PID: 40, ParentPID: 43, StartTime: "140"}, // Descendant precedes its parent. + {PID: 42, Master: true, StartTime: "142"}, + {PID: 43, ParentPID: 42, StartTime: "143"}, + {PID: 50, StartTime: "150"}, // Unrelated master and worker. + {PID: 51, ParentPID: 50, StartTime: "151"}, + } + identities := runtimeMasterIdentities(processes) + if got := runtimeDescendants(processes, identities); len(got) != 3 { + t.Fatalf("instance processes = %v, want master and two descendants", got) + } + remaining := []runtimeProcess{ + {PID: 43, ParentPID: 1, StartTime: "143"}, // Reparented old worker. + {PID: 40, ParentPID: 1, StartTime: "240"}, // Reused PID. + {PID: 50, StartTime: "150"}, + {PID: 51, ParentPID: 50, StartTime: "151"}, + } + got := runtimeDescendants(remaining, identities) + if len(got) != 1 || got[0].PID != 43 { + t.Fatalf("remaining instance processes = %v, want old worker 43", got) + } +} + +func TestRestartIgnoresUnrelatedProcessesAfterShutdown(t *testing.T) { + inspections := 0 + runner := &fakeRunner{} + executor := &PathExecutor{Path: "openresty", Runner: runner, + inspectProcesses: func() ([]runtimeProcess, error) { + inspections++ + processes := []runtimeProcess{{PID: 50}, {PID: 51, ParentPID: 50}} + if inspections == 1 { + processes = append(processes, runtimeProcess{PID: 42, Master: true}) + } + return processes, nil + }, + signalProcess: func(_ runtimeProcess, _ bool) error { return nil }, + } + ctx, cancel := context.WithTimeout(context.Background(), time.Second) + defer cancel() + if err := executor.Restart(ctx); err != nil { + t.Fatalf("Restart(other nginx instance) = %v, want nil", err) + } + if len(runner.calls) != 1 { + t.Errorf("Restart commands = %v, want one start", runner.calls) + } +} diff --git a/internal/apps/agent/nginx/runtime_process_unix.go b/internal/apps/agent/nginx/runtime_process_unix.go new file mode 100644 index 00000000..59eb49a0 --- /dev/null +++ b/internal/apps/agent/nginx/runtime_process_unix.go @@ -0,0 +1,40 @@ +//go:build unix && !linux + +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package nginx + +import ( + "context" + "fmt" + "os" + "os/exec" + "syscall" +) + +// Non-Linux Unix systems expose process ownership and start time through ps. +func inspectOpenrestyProcesses(ctx context.Context, binary, configPath string) ([]runtimeProcess, error) { + output, err := exec.CommandContext(ctx, "ps", "-axo", "pid=,ppid=,uid=,lstart=,command=").Output() + if err != nil { + return nil, fmt.Errorf("inspect openresty processes: %w", err) + } + return parsePSRuntime(string(output), binary, configPath, os.Geteuid()) +} + +func signalOpenrestyProcess(ctx context.Context, binary, config string, process runtimeProcess, quit bool) error { + processes, err := inspectOpenrestyProcesses(ctx, binary, config) + if err != nil { + return err + } + for _, current := range processes { + if current == process && current.Master { + signal := syscall.SIGHUP + if quit { + signal = syscall.SIGQUIT + } + return syscall.Kill(current.PID, signal) + } + } + return fmt.Errorf("openresty master identity changed before signalling pid %d", process.PID) +} diff --git a/internal/apps/agent/sync/service.go b/internal/apps/agent/sync/service.go index ee5d4ce1..5cfd1477 100644 --- a/internal/apps/agent/sync/service.go +++ b/internal/apps/agent/sync/service.go @@ -193,7 +193,7 @@ func (s *Service) applyRenderedConfig(ctx context.Context, mode string, snapshot } if !shouldReportApplyLog(alreadySynced, applyResult.reportResult) { slog.Debug("skipping duplicate apply log report", "version", config.Version, "checksum", config.Checksum, "result", applyResult.reportResult) - if applyResult.reportResult == ApplyResultFailed { + if applyResult.reportResult != ApplyResultSuccess { return outcomeError(config.Version, applyResult.message) } if err := s.syncReferencedWAFIPGroups(ctx, rendered.supportFiles); err != nil { @@ -215,7 +215,7 @@ func (s *Service) applyRenderedConfig(ctx context.Context, mode string, snapshot slog.Error("report apply log failed", "version", config.Version, "result", applyResult.reportResult, "error", err) return err } - if applyResult.reportResult == ApplyResultFailed { + if applyResult.reportResult != ApplyResultSuccess { slog.Warn("failed apply log reported", "version", config.Version) return outcomeError(config.Version, applyResult.message) } diff --git a/internal/apps/agent/sync/service_test.go b/internal/apps/agent/sync/service_test.go index d8d02745..8022400a 100644 --- a/internal/apps/agent/sync/service_test.go +++ b/internal/apps/agent/sync/service_test.go @@ -732,8 +732,8 @@ func TestSyncOnceReportsWarningWhenRollbackKeepsOpenrestyHealthy(t *testing.T) { if err = service.SyncOnce(context.Background(), &protocol.ActiveConfigMeta{ Version: client.config.Version, Checksum: client.config.Checksum, - }); err != nil { - t.Fatalf("expected warning outcome to keep sync successful, got %v", err) + }); err == nil { + t.Fatal("SyncOnce(rolled-back target) returned nil, want an error so automatic retry continues") } snapshot, err := stateStore.Load()