From 29f19c5edd524df996dd694895934489983f5d6a Mon Sep 17 00:00:00 2001 From: ryan Date: Mon, 9 Mar 2026 23:50:55 +0800 Subject: [PATCH] =?UTF-8?q?=E9=87=8D=E6=9E=84=20Nginx=20=E7=AE=A1=E7=90=86?= =?UTF-8?q?=E5=99=A8=E5=92=8C=E5=90=8C=E6=AD=A5=E6=9C=8D=E5=8A=A1=EF=BC=8C?= =?UTF-8?q?=E6=B7=BB=E5=8A=A0=E8=BF=90=E8=A1=8C=E6=97=B6=E7=A1=AE=E4=BF=9D?= =?UTF-8?q?=E5=8A=9F=E8=83=BD=EF=BC=8C=E6=9B=B4=E6=96=B0=E7=9B=B8=E5=85=B3?= =?UTF-8?q?=E6=B5=8B=E8=AF=95=E7=94=A8=E4=BE=8B=E5=92=8C=E6=96=87=E6=A1=A3?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- atsf_agent/internal/agent/runner.go | 4 +- atsf_agent/internal/nginx/manager.go | 63 +++++++++++++++++++--- atsf_agent/internal/nginx/manager_test.go | 52 ++++++++++++++++-- atsf_agent/internal/sync/service.go | 36 +++++++++++-- atsf_agent/internal/sync/service_test.go | 64 +++++++++++++++++++++++ docs/deployment.md | 10 ++-- docs/design.md | 2 + docs/development-guidelines.md | 3 ++ 8 files changed, 213 insertions(+), 21 deletions(-) diff --git a/atsf_agent/internal/agent/runner.go b/atsf_agent/internal/agent/runner.go index 93ef9128..ff733c1b 100644 --- a/atsf_agent/internal/agent/runner.go +++ b/atsf_agent/internal/agent/runner.go @@ -23,10 +23,10 @@ func (r *Runner) Run(ctx context.Context) error { if err != nil { return err } - if err = r.HeartbeatService.Register(ctx, r.nodePayload(nodeID)); err != nil { + if err = r.SyncService.SyncOnStartup(ctx); err != nil { return err } - if err = r.SyncService.SyncOnce(ctx); err != nil { + if err = r.HeartbeatService.Register(ctx, r.nodePayload(nodeID)); err != nil { return err } diff --git a/atsf_agent/internal/nginx/manager.go b/atsf_agent/internal/nginx/manager.go index a5744fdc..9641a6d2 100644 --- a/atsf_agent/internal/nginx/manager.go +++ b/atsf_agent/internal/nginx/manager.go @@ -2,6 +2,8 @@ package nginx import ( "context" + "crypto/sha256" + "encoding/hex" "errors" "fmt" "os" @@ -13,6 +15,7 @@ import ( type Executor interface { Test(ctx context.Context) error Reload(ctx context.Context) error + EnsureRuntime(ctx context.Context, recreate bool) error } type CommandRunner interface { @@ -48,6 +51,10 @@ func (e *PathExecutor) Reload(ctx context.Context) error { return nil } +func (e *PathExecutor) EnsureRuntime(ctx context.Context, recreate bool) error { + return nil +} + type DockerExecutor struct { DockerBinary string ContainerName string @@ -57,7 +64,7 @@ type DockerExecutor struct { } func (e *DockerExecutor) Test(ctx context.Context) error { - if err := e.ensureContainer(ctx); err != nil { + if err := e.EnsureRuntime(ctx, false); err != nil { return err } output, err := e.Runner.Run(ctx, e.DockerBinary, "exec", e.ContainerName, "nginx", "-t") @@ -68,7 +75,7 @@ func (e *DockerExecutor) Test(ctx context.Context) error { } func (e *DockerExecutor) Reload(ctx context.Context) error { - if err := e.ensureContainer(ctx); err != nil { + if err := e.EnsureRuntime(ctx, false); err != nil { return err } output, err := e.Runner.Run(ctx, e.DockerBinary, "exec", e.ContainerName, "nginx", "-s", "reload") @@ -78,19 +85,39 @@ func (e *DockerExecutor) Reload(ctx context.Context) error { return nil } -func (e *DockerExecutor) ensureContainer(ctx context.Context) error { +func (e *DockerExecutor) EnsureRuntime(ctx context.Context, recreate bool) error { output, err := e.Runner.Run(ctx, e.DockerBinary, "inspect", "-f", "{{.State.Running}}", e.ContainerName) if err == nil { + if recreate { + if err := e.removeContainer(ctx); err != nil { + return err + } + return e.runContainer(ctx) + } if strings.TrimSpace(string(output)) == "true" { return nil } - startOutput, startErr := e.Runner.Run(ctx, e.DockerBinary, "start", e.ContainerName) - if startErr != nil { - return fmt.Errorf("docker start nginx failed: %w: %s", startErr, string(startOutput)) + if err := e.removeContainer(ctx); err != nil { + return err } - return nil + return e.runContainer(ctx) } + return e.runContainer(ctx) +} +func (e *DockerExecutor) removeContainer(ctx context.Context) error { + output, err := e.Runner.Run(ctx, e.DockerBinary, "rm", "-f", e.ContainerName) + if err != nil { + text := string(output) + if strings.Contains(text, "No such container") { + return nil + } + return fmt.Errorf("docker rm nginx failed: %w: %s", err, text) + } + return nil +} + +func (e *DockerExecutor) runContainer(ctx context.Context) error { runArgs := []string{ "run", "-d", "--name", e.ContainerName, @@ -133,6 +160,28 @@ func (m *Manager) Apply(ctx context.Context, content string) error { return nil } +func (m *Manager) EnsureRuntime(ctx context.Context, recreate bool) error { + if m.Executor == nil { + return errors.New("executor 未配置") + } + return m.Executor.EnsureRuntime(ctx, recreate) +} + +func (m *Manager) CurrentChecksum() (string, error) { + if m.RouteConfigPath == "" { + return "", errors.New("route config path 不能为空") + } + data, err := os.ReadFile(m.RouteConfigPath) + if err != nil { + if os.IsNotExist(err) { + return "", nil + } + return "", err + } + sum := sha256.Sum256(data) + return hex.EncodeToString(sum[:]), nil +} + type ExecutorOptions struct { NginxPath string DockerBinary string diff --git a/atsf_agent/internal/nginx/manager_test.go b/atsf_agent/internal/nginx/manager_test.go index 07b254d9..e466847a 100644 --- a/atsf_agent/internal/nginx/manager_test.go +++ b/atsf_agent/internal/nginx/manager_test.go @@ -50,6 +50,16 @@ func TestPathExecutorCommands(t *testing.T) { } } +func TestPathExecutorEnsureRuntimeNoop(t *testing.T) { + executor := &PathExecutor{ + Path: "/opt/nginx/sbin/nginx", + Runner: &fakeRunner{}, + } + if err := executor.EnsureRuntime(context.Background(), true); err != nil { + t.Fatalf("EnsureRuntime failed: %v", err) + } +} + func TestDockerExecutorStartsContainerWhenMissing(t *testing.T) { runner := &fakeRunner{ runFn: func(name string, args ...string) ([]byte, error) { @@ -103,14 +113,48 @@ func TestDockerExecutorStartsStoppedContainer(t *testing.T) { t.Fatalf("Reload failed: %v", err) } + if len(runner.calls) != 4 { + t.Fatalf("expected 4 calls, got %d", len(runner.calls)) + } + if runner.calls[1].args[0] != "rm" { + t.Fatalf("expected docker rm on second call, got %#v", runner.calls[1]) + } + if runner.calls[2].args[0] != "run" { + t.Fatalf("expected docker run on third call, got %#v", runner.calls[2]) + } + if runner.calls[3].args[0] != "exec" { + t.Fatalf("expected docker exec on fourth call, got %#v", runner.calls[3]) + } +} + +func TestDockerExecutorRecreatesContainerOnStartup(t *testing.T) { + runner := &fakeRunner{ + runFn: func(name string, args ...string) ([]byte, error) { + if len(args) >= 1 && args[0] == "inspect" { + return []byte("true"), nil + } + return []byte("ok"), nil + }, + } + executor := &DockerExecutor{ + DockerBinary: "docker", + ContainerName: "atsflare-nginx", + Image: "nginx:stable-alpine", + RouteConfigDir: filepath.Clean("/tmp/routes"), + Runner: runner, + } + + if err := executor.EnsureRuntime(context.Background(), true); err != nil { + t.Fatalf("EnsureRuntime failed: %v", err) + } if len(runner.calls) != 3 { t.Fatalf("expected 3 calls, got %d", len(runner.calls)) } - if runner.calls[1].args[0] != "start" { - t.Fatalf("expected docker start on second call, got %#v", runner.calls[1]) + if runner.calls[1].args[0] != "rm" { + t.Fatalf("expected docker rm on second call, got %#v", runner.calls[1]) } - if runner.calls[2].args[0] != "exec" { - t.Fatalf("expected docker exec on third call, got %#v", runner.calls[2]) + if runner.calls[2].args[0] != "run" { + t.Fatalf("expected docker run on third call, got %#v", runner.calls[2]) } } diff --git a/atsf_agent/internal/sync/service.go b/atsf_agent/internal/sync/service.go index bf257627..29708a89 100644 --- a/atsf_agent/internal/sync/service.go +++ b/atsf_agent/internal/sync/service.go @@ -3,7 +3,6 @@ package sync import ( "context" - "atsflare-agent/internal/nginx" "atsflare-agent/internal/protocol" "atsflare-agent/internal/state" ) @@ -18,13 +17,19 @@ type ConfigClient interface { ReportApplyLog(ctx context.Context, payload protocol.ApplyLogPayload) error } +type NginxManager interface { + Apply(ctx context.Context, content string) error + EnsureRuntime(ctx context.Context, recreate bool) error + CurrentChecksum() (string, error) +} + type Service struct { client ConfigClient - nginxManager *nginx.Manager + nginxManager NginxManager stateStore *state.Store } -func New(client ConfigClient, nginxManager *nginx.Manager, stateStore *state.Store) *Service { +func New(client ConfigClient, nginxManager NginxManager, stateStore *state.Store) *Service { return &Service{ client: client, nginxManager: nginxManager, @@ -33,6 +38,14 @@ func New(client ConfigClient, nginxManager *nginx.Manager, stateStore *state.Sto } func (s *Service) SyncOnce(ctx context.Context) error { + return s.sync(ctx, false) +} + +func (s *Service) SyncOnStartup(ctx context.Context) error { + return s.sync(ctx, true) +} + +func (s *Service) sync(ctx context.Context, startup bool) error { snapshot, err := s.stateStore.Load() if err != nil { return err @@ -41,7 +54,22 @@ func (s *Service) SyncOnce(ctx context.Context) error { if err != nil { return err } - if snapshot.CurrentVersion == config.Version && snapshot.CurrentChecksum == config.Checksum { + currentChecksum, err := s.nginxManager.CurrentChecksum() + if err != nil { + return err + } + if currentChecksum == config.Checksum { + if startup { + if err = s.nginxManager.EnsureRuntime(ctx, true); err != nil { + return err + } + } + snapshot.CurrentVersion = config.Version + snapshot.CurrentChecksum = config.Checksum + snapshot.LastError = "" + return s.stateStore.Save(snapshot) + } + if snapshot.CurrentVersion == config.Version && snapshot.CurrentChecksum == config.Checksum && !startup { return nil } if err = s.nginxManager.Apply(ctx, config.RenderedConfig); err != nil { diff --git a/atsf_agent/internal/sync/service_test.go b/atsf_agent/internal/sync/service_test.go index cdeddf06..98c12033 100644 --- a/atsf_agent/internal/sync/service_test.go +++ b/atsf_agent/internal/sync/service_test.go @@ -22,6 +22,14 @@ type fakeClient struct { reports []protocol.ApplyLogPayload } +type fakeManager struct { + applyErr error + currentChecksum string + currentChecksumErr error + ensureCalls []bool + applyContents []string +} + func (f *fakeExecutor) Test(ctx context.Context) error { return f.testErr } @@ -30,6 +38,10 @@ func (f *fakeExecutor) Reload(ctx context.Context) error { return f.reloadErr } +func (f *fakeExecutor) EnsureRuntime(ctx context.Context, recreate bool) error { + return nil +} + func (f *fakeClient) GetActiveConfig(ctx context.Context) (*protocol.ActiveConfigResponse, error) { return &f.config, nil } @@ -39,6 +51,20 @@ func (f *fakeClient) ReportApplyLog(ctx context.Context, payload protocol.ApplyL return nil } +func (m *fakeManager) Apply(ctx context.Context, content string) error { + m.applyContents = append(m.applyContents, content) + return m.applyErr +} + +func (m *fakeManager) EnsureRuntime(ctx context.Context, recreate bool) error { + m.ensureCalls = append(m.ensureCalls, recreate) + return nil +} + +func (m *fakeManager) CurrentChecksum() (string, error) { + return m.currentChecksum, m.currentChecksumErr +} + func TestSyncOnceSuccess(t *testing.T) { client := &fakeClient{ config: protocol.ActiveConfigResponse{ @@ -148,3 +174,41 @@ func TestSyncOnceRollbackOnNginxFailure(t *testing.T) { t.Fatal("expected failed apply report to be sent") } } + +func TestSyncOnStartupRecreatesRuntimeWhenChecksumMatches(t *testing.T) { + client := &fakeClient{ + config: protocol.ActiveConfigResponse{ + Version: "20260309-003", + Checksum: "checksum-3", + RenderedConfig: "server { listen 82; }", + CreatedAt: time.Now().Format(time.RFC3339), + }, + } + stateStore := state.NewStore(filepath.Join(t.TempDir(), "state.json")) + nodeID, err := stateStore.EnsureNodeID() + if err != nil { + t.Fatalf("EnsureNodeID failed: %v", err) + } + if err = stateStore.Save(&state.Snapshot{NodeID: nodeID}); err != nil { + t.Fatalf("failed to seed state: %v", err) + } + + manager := &fakeManager{currentChecksum: "checksum-3"} + service := New(client, manager, stateStore) + if err = service.SyncOnStartup(context.Background()); err != nil { + t.Fatalf("SyncOnStartup failed: %v", err) + } + if len(manager.ensureCalls) != 1 || !manager.ensureCalls[0] { + t.Fatal("expected startup sync to recreate runtime") + } + if len(client.reports) != 0 { + t.Fatal("expected no apply report when checksum already matches") + } + snapshot, err := stateStore.Load() + if err != nil { + t.Fatalf("failed to load state: %v", err) + } + if snapshot.CurrentChecksum != "checksum-3" || snapshot.CurrentVersion != "20260309-003" { + t.Fatal("expected snapshot to be refreshed from active config") + } +} diff --git a/docs/deployment.md b/docs/deployment.md index 6bed1cff..a020f2b8 100644 --- a/docs/deployment.md +++ b/docs/deployment.md @@ -133,9 +133,11 @@ go build -o atsflare-agent ./cmd/agent 2. Agent 拉取当前激活版本 3. Agent 写入 `route_config_path` 4. Agent 使用 `nginx_path` 指向的独立 Nginx,或自动准备 Docker Nginx 容器 -5. Agent 执行 `nginx -t` -6. Agent 执行 `nginx -s reload` -7. Agent 上报成功结果 +5. Agent 启动时先校验本地路由文件 checksum 与控制面激活版本是否一致 +6. Docker 模式下会重建容器,避免复用故障容器 +7. Agent 执行 `nginx -t` +8. Agent 执行 `nginx -s reload` +9. Agent 上报成功结果 ### 4.4 验证管理端状态 @@ -188,7 +190,7 @@ npm run build - Agent 运行器当前任一心跳或同步失败会直接退出,需要结合进程管理器拉起 - 尚未提供 systemd unit 文件 - 尚未提供 Docker Compose 或一键部署脚本 -- Docker 模式当前默认直接启动单容器 Nginx,挂载与端口策略仍是 MVP 水平 +- Docker 模式当前默认直接重建单容器 Nginx,挂载与端口策略仍是 MVP 水平 - 前端页面已可用,但交互和校验仍是 MVP 水平 - 目前联调说明以手工步骤为主,未内置完整自动化端到端脚本 diff --git a/docs/design.md b/docs/design.md index 31fe13de..6d188ab6 100644 --- a/docs/design.md +++ b/docs/design.md @@ -67,6 +67,8 @@ Agent 使用 Go 单体程序: * 未配置 `nginx_path` 时,默认通过 Docker 运行独立 Nginx 容器 * 管理本机 Nginx 路由配置文件和 reload * Agent 生成资源默认统一落在 `./data`,也允许通过单个基路径配置覆盖 +* Agent 启动时会校验本地路由文件哈希与控制面激活版本是否一致 +* Docker 模式启动时会重建独立 Nginx 容器,避免复用故障容器 ### Nginx 管理边界 diff --git a/docs/development-guidelines.md b/docs/development-guidelines.md index c81e2181..6eec653a 100644 --- a/docs/development-guidelines.md +++ b/docs/development-guidelines.md @@ -266,6 +266,8 @@ Agent 必须满足以下行为: * 上报最终应用结果 * 优先使用 `nginx_path` * 未配置 `nginx_path` 时自动准备并使用 Docker Nginx 容器 +* 启动时先校验本地路由文件 checksum 与控制面激活版本是否一致 +* Docker 模式启动时应重建容器,而不是继续复用异常停止的旧容器 ### 7.3 容错规范 @@ -275,6 +277,7 @@ Agent 必须满足以下行为: * 下载失败时,不修改本地配置 * 配置校验或 reload 失败时,自动尝试回滚 * 本地状态文件损坏时,允许重新初始化,但不能删除正在生效的 Nginx 配置 +* Docker 容器异常停止时,启动阶段应自动重建容器并重新校验配置 ### 7.4 外部命令规范