mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-10-08 16:46:37 +08:00
fix(agent): recover OpenResty runtime after invalid PID (#45)
* fix(agent): recover OpenResty runtime after invalid PID * fix(agent): scope OpenResty shutdown to instance processes --------- Co-authored-by: Noru Wyrms <wyrmsnoru@gmail.com>
This commit is contained in:
@@ -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 与已有用户或其他认证源冲突,已有用户及绑定不受影响。
|
||||
|
||||
|
||||
@@ -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,即使其被重新托管);超时会保留错误,不会启动第二个实例。检查仍未结束的连接及进程后再重试 |
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
@@ -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")
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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()
|
||||
|
||||
Reference in New Issue
Block a user