mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-09-29 05:56:38 +08:00
重构 Nginx 管理器和同步服务,添加运行时确保功能,更新相关测试用例和文档
This commit is contained in:
@@ -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
|
||||
}
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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])
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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")
|
||||
}
|
||||
}
|
||||
|
||||
+6
-4
@@ -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 水平
|
||||
- 目前联调说明以手工步骤为主,未内置完整自动化端到端脚本
|
||||
|
||||
|
||||
@@ -67,6 +67,8 @@ Agent 使用 Go 单体程序:
|
||||
* 未配置 `nginx_path` 时,默认通过 Docker 运行独立 Nginx 容器
|
||||
* 管理本机 Nginx 路由配置文件和 reload
|
||||
* Agent 生成资源默认统一落在 `./data`,也允许通过单个基路径配置覆盖
|
||||
* Agent 启动时会校验本地路由文件哈希与控制面激活版本是否一致
|
||||
* Docker 模式启动时会重建独立 Nginx 容器,避免复用故障容器
|
||||
|
||||
### Nginx 管理边界
|
||||
|
||||
|
||||
@@ -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 外部命令规范
|
||||
|
||||
|
||||
Reference in New Issue
Block a user