diff --git a/atsf_agent/internal/agent/runner.go b/atsf_agent/internal/agent/runner.go index ff733c1b..bc08413c 100644 --- a/atsf_agent/internal/agent/runner.go +++ b/atsf_agent/internal/agent/runner.go @@ -2,20 +2,29 @@ package agent import ( "context" + "log" "time" "atsflare-agent/internal/config" - "atsflare-agent/internal/heartbeat" "atsflare-agent/internal/protocol" "atsflare-agent/internal/state" - syncservice "atsflare-agent/internal/sync" ) +type HeartbeatService interface { + Register(ctx context.Context, payload protocol.NodePayload) error + Heartbeat(ctx context.Context, payload protocol.NodePayload) error +} + +type SyncService interface { + SyncOnStartup(ctx context.Context) error + SyncOnce(ctx context.Context) error +} + type Runner struct { Config *config.Config StateStore *state.Store - HeartbeatService *heartbeat.Service - SyncService *syncservice.Service + HeartbeatService HeartbeatService + SyncService SyncService } func (r *Runner) Run(ctx context.Context) error { @@ -23,11 +32,15 @@ func (r *Runner) Run(ctx context.Context) error { if err != nil { return err } - if err = r.SyncService.SyncOnStartup(ctx); err != nil { - return err - } if err = r.HeartbeatService.Register(ctx, r.nodePayload(nodeID)); err != nil { - return err + log.Printf("agent register failed: %v", err) + } + if err = r.SyncService.SyncOnStartup(ctx); err != nil { + r.recordSyncError(err) + log.Printf("agent startup sync failed: %v", err) + } + if err = r.HeartbeatService.Heartbeat(ctx, r.nodePayload(nodeID)); err != nil { + log.Printf("agent startup heartbeat failed: %v", err) } heartbeatTicker := time.NewTicker(r.Config.HeartbeatInterval) @@ -41,16 +54,32 @@ func (r *Runner) Run(ctx context.Context) error { return ctx.Err() case <-heartbeatTicker.C: if err = r.HeartbeatService.Heartbeat(ctx, r.nodePayload(nodeID)); err != nil { - return err + log.Printf("agent heartbeat failed: %v", err) } case <-syncTicker.C: if err = r.SyncService.SyncOnce(ctx); err != nil { - return err + r.recordSyncError(err) + log.Printf("agent sync failed: %v", err) } } } } +func (r *Runner) recordSyncError(err error) { + if err == nil || r.StateStore == nil { + return + } + snapshot, loadErr := r.StateStore.Load() + if loadErr != nil { + log.Printf("load state before recording sync error failed: %v", loadErr) + return + } + snapshot.LastError = err.Error() + if saveErr := r.StateStore.Save(snapshot); saveErr != nil { + log.Printf("save state after sync error failed: %v", saveErr) + } +} + func (r *Runner) nodePayload(nodeID string) protocol.NodePayload { snapshot, _ := r.StateStore.Load() return protocol.NodePayload{ diff --git a/atsf_agent/internal/agent/runner_test.go b/atsf_agent/internal/agent/runner_test.go new file mode 100644 index 00000000..0dc3a293 --- /dev/null +++ b/atsf_agent/internal/agent/runner_test.go @@ -0,0 +1,172 @@ +package agent + +import ( + "context" + "errors" + "path/filepath" + "sync" + "testing" + "time" + + "atsflare-agent/internal/config" + "atsflare-agent/internal/protocol" + "atsflare-agent/internal/state" +) + +type fakeHeartbeatService struct { + mu sync.Mutex + registerCalls int + heartbeatCalls int + registerErr error + heartbeatErrs []error + onHeartbeat func(int) +} + +func (f *fakeHeartbeatService) Register(ctx context.Context, payload protocol.NodePayload) error { + f.mu.Lock() + defer f.mu.Unlock() + f.registerCalls++ + return f.registerErr +} + +func (f *fakeHeartbeatService) Heartbeat(ctx context.Context, payload protocol.NodePayload) error { + f.mu.Lock() + f.heartbeatCalls++ + callIndex := f.heartbeatCalls + var err error + if len(f.heartbeatErrs) >= callIndex { + err = f.heartbeatErrs[callIndex-1] + } + onHeartbeat := f.onHeartbeat + f.mu.Unlock() + if onHeartbeat != nil { + onHeartbeat(callIndex) + } + return err +} + +type fakeSyncService struct { + mu sync.Mutex + startupErr error + syncOnceErr error + startupCalls int + syncOnceCalls int + onSyncOnceCall func(int) +} + +func (f *fakeSyncService) SyncOnStartup(ctx context.Context) error { + f.mu.Lock() + defer f.mu.Unlock() + f.startupCalls++ + return f.startupErr +} + +func (f *fakeSyncService) SyncOnce(ctx context.Context) error { + f.mu.Lock() + f.syncOnceCalls++ + callIndex := f.syncOnceCalls + callback := f.onSyncOnceCall + f.mu.Unlock() + if callback != nil { + callback(callIndex) + } + return f.syncOnceErr +} + +func TestRunnerKeepsHeartbeatWhenStartupSyncFails(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + + stateStore := state.NewStore(filepath.Join(t.TempDir(), "state.json")) + heartbeatService := &fakeHeartbeatService{ + onHeartbeat: func(callCount int) { + if callCount >= 2 { + cancel() + } + }, + } + syncService := &fakeSyncService{ + startupErr: errors.New("当前没有激活版本,保持当前 Nginx 配置"), + } + runner := &Runner{ + Config: &config.Config{ + NodeName: "edge-01", + NodeIP: "10.0.0.8", + AgentVersion: "0.1.0", + NginxVersion: "1.25.5", + HeartbeatInterval: 10 * time.Millisecond, + SyncInterval: 20 * time.Millisecond, + }, + StateStore: stateStore, + HeartbeatService: heartbeatService, + SyncService: syncService, + } + + err := runner.Run(ctx) + if !errors.Is(err, context.Canceled) { + t.Fatalf("expected context cancellation, got %v", err) + } + if heartbeatService.registerCalls != 1 { + t.Fatalf("expected 1 register call, got %d", heartbeatService.registerCalls) + } + if heartbeatService.heartbeatCalls < 2 { + t.Fatalf("expected heartbeat loop to continue, got %d heartbeat calls", heartbeatService.heartbeatCalls) + } + snapshot, loadErr := stateStore.Load() + if loadErr != nil { + t.Fatalf("failed to load state: %v", loadErr) + } + if snapshot.LastError != "当前没有激活版本,保持当前 Nginx 配置" { + t.Fatalf("expected startup sync error to be recorded, got %q", snapshot.LastError) + } +} + +func TestRunnerDoesNotExitOnHeartbeatOrSyncError(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + + stateStore := state.NewStore(filepath.Join(t.TempDir(), "state.json")) + heartbeatService := &fakeHeartbeatService{ + registerErr: errors.New("register timeout"), + heartbeatErrs: []error{errors.New("heartbeat timeout")}, + } + syncService := &fakeSyncService{ + syncOnceErr: errors.New("nginx reload failed"), + onSyncOnceCall: func(callCount int) { + if callCount >= 1 { + cancel() + } + }, + } + runner := &Runner{ + Config: &config.Config{ + NodeName: "edge-01", + NodeIP: "10.0.0.8", + AgentVersion: "0.1.0", + NginxVersion: "1.25.5", + HeartbeatInterval: 10 * time.Millisecond, + SyncInterval: 10 * time.Millisecond, + }, + StateStore: stateStore, + HeartbeatService: heartbeatService, + SyncService: syncService, + } + + err := runner.Run(ctx) + if !errors.Is(err, context.Canceled) { + t.Fatalf("expected context cancellation, got %v", err) + } + if heartbeatService.registerCalls != 1 { + t.Fatalf("expected register attempt, got %d", heartbeatService.registerCalls) + } + if syncService.syncOnceCalls == 0 { + t.Fatal("expected sync loop to continue after heartbeat/register errors") + } + snapshot, loadErr := stateStore.Load() + if loadErr != nil { + t.Fatalf("failed to load state: %v", loadErr) + } + if snapshot.LastError != "nginx reload failed" { + t.Fatalf("expected sync error to be recorded, got %q", snapshot.LastError) + } +} diff --git a/docs/development-plan.md b/docs/development-plan.md index 9c21dfce..d2fcd0d1 100644 --- a/docs/development-plan.md +++ b/docs/development-plan.md @@ -1,238 +1,238 @@ -# ATSFlare 开发计划 - -## 1. 目标 - -当前开发目标是完成 ATSFlare 的 MVP 闭环: - -```text -后台改规则 -> 点击发布 -> Agent 拉到新版本 -> 写入 Nginx 路由配置 -> nginx 校验并 reload -> 节点状态可见 -``` - -MVP 已于第一版完成。当前进入第二版迭代。 - -## 2. 第一版里程碑(已完成) - -### Phase 1: Server 数据层与发布闭环 ✅ - -交付: - -* 新模型(proxy_routes、config_versions) -* AutoMigrate -* 路由 CRUD API -* 发布 API -* 激活版本 API -* 渲染 service - -### Phase 2: Agent API 与节点状态 ✅ - -交付: - -* `nodes`、`apply_logs` -* Agent Token 鉴权(全局单 Token) -* 节点在线状态计算 -* 节点与应用日志查询接口 - -### Phase 3: Agent 本体 ✅ - -交付: - -* 本地配置文件读取 -* `node_id` 持久化 -* 心跳循环 -* 版本检查 -* 下载配置 -* 写入 Nginx 路由配置 -* 配置校验、reload 与失败回滚 -* 应用结果上报 - -### Phase 4: 管理端页面 ✅ - -交付: - -* 反代规则页 -* 版本页 -* 节点页 -* 应用记录页 - -### Phase 5: 联调与收尾 ✅ - -交付: - -* 手工部署说明 -* Agent 配置示例 -* 联调验证记录 - -## 3. 第二版里程碑 - -### V2 Phase 1: HTTPS/TLS 支持 - -目标: - -* 支持通过控制面配置 HTTPS 路由 -* 渲染出包含 443 端口的 `server` 块 -* 支持 HTTP → HTTPS 重定向块 -* 控制面支持托管证书 - -交付: - -* `proxy_routes` 新增字段:`enable_https`、`cert_id`、`redirect_http` -* 渲染器支持 HTTPS server 块生成 -* 前端反代规则页增加 HTTPS 配置表单 -* `tls_certificates` 表与模型 -* 证书导入能力:手动导入(粘贴 PEM)与文件导入 -* AutoMigrate 覆盖新字段 - -完成标准: - -* 创建含 HTTPS 字段的路由并发布,Agent 拉取后 Nginx 能以 HTTPS 正确转发 -* HTTP 重定向配置生效 -* 证书可通过控制面导入并被 HTTPS 路由引用 -* 未开启 HTTPS 的路由渲染行为与第一版保持一致 - -### V2 Phase 2: 域名管理与证书自动匹配 - -目标: - -* 控制面可管理域名并绑定证书 -* 反代规则编辑时按输入域名自动匹配证书 -* 支持 `*.example.com` 通配符证书匹配 - -交付: - -* `managed_domains` 表与模型 -* 域名管理 CRUD API -* 证书匹配 API(精确匹配 + 通配符匹配) -* 前端域名管理页 -* 前端反代规则页接入证书自动匹配 - -完成标准: - -* 创建域名并绑定证书后,反代规则输入域名可自动匹配证书 -* 通配符证书可匹配子域名(如 `api.example.com` 匹配 `*.example.com`) -* 无匹配证书时前端给出明确提示 - -### V2 Phase 3: Agent Token 管理 - -目标: - -* 支持多个命名 Token -* Token 可通过管理界面创建和撤销 -* 中间件改为查库验证 - -交付: - -* `agent_tokens` 表与模型 -* agent-auth 中间件改造(查表验证) -* 全局 Token 环境变量降级为 bootstrap 模式 -* Token CRUD API -* 前端 Token 管理页 - -完成标准: - -* 新建 Token 后 Agent 可用该 Token 访问 Agent API -* 撤销 Token 后 Agent 请求立即返回 401 -* 全局环境变量 Token 仅在数据库无有效记录时生效 - -### V2 Phase 4: 路由自定义头 - -目标: - -* 每条路由支持追加自定义 `proxy_set_header` 指令 - -交付: - -* `proxy_routes` 新增 `custom_headers` 字段(JSON,`[{"key":"X-My-Header","value":"foo"}]`) -* 渲染器按 `custom_headers` 在 `location /` 块中注入额外 header 指令 -* 前端反代规则页增加自定义头编辑器 - -完成标准: - -* 路由配置自定义头后发布,渲染结果包含对应 `proxy_set_header` 指令 -* 不配置自定义头的路由渲染行为与之前保持一致 - -### V2 Phase 5: 配置预览与变更摘要 - -目标: - -* 发布前可预览渲染结果 -* 发布前可查看与当前激活版本的变更摘要 - -交付: - -* `GET /api/config-versions/preview` 接口(返回渲染后的 Nginx 配置,不写库) -* `GET /api/config-versions/diff` 接口(返回新增/删除/修改的域名列表) -* 前端发布版本页增加"预览"按钮和变更摘要展示弹窗 - -完成标准: - -* 点击预览可查看即将生成的 Nginx 配置文本 -* 变更摘要正确列出相对于当前激活版本的域名变化 -* 预览和 diff 操作不产生版本记录 - -## 4. 第二版建议执行顺序 - -建议严格按以下顺序开发: - -1. HTTPS/TLS 支持(对现有渲染链路影响最小,独立可测) -2. 域名管理与证书自动匹配(与 HTTPS 强相关,尽早完成闭环) -3. Agent Token 管理(安全性改善,早做早稳) -4. 路由自定义头(纯增量,对现有结构影响小) -5. 配置预览与变更摘要(纯只读接口,最后补充) - -不要先做以下内容: - -* 节点分组与差异化下发 -* 灰度百分比发布 -* WAF、限流 -* Redis、Prometheus -* 多租户 - -## 5. 第二版每阶段验收检查 - -### V2 Phase 1 检查项 - -* `enable_https=true` 的路由发布后生成 443 端口 server 块 -* `redirect_http=true` 的路由生成 80 → 443 重定向块 -* 支持手动导入证书与文件导入证书 -* 未开启 HTTPS 的路由渲染结果不受影响 -* Agent 拉取后 Nginx reload 成功 - -### V2 Phase 2 检查项 - -* 可通过管理界面维护域名并绑定证书 -* 反代规则输入域名后可自动匹配证书 -* 通配符证书可匹配子域名(`*.example.com`) - -### V2 Phase 3 检查项 - -* 可通过管理界面创建 Token -* 新 Token 可被 Agent 使用 -* 撤销 Token 后访问立即失败 -* 全局 Token 环境变量在数据库有记录时失效 - -### V2 Phase 4 检查项 - -* 路由可添加自定义头 -* 渲染结果中正确包含自定义 header 指令 -* 无自定义头的路由渲染结果不受影响 - -### V2 Phase 5 检查项 - -* 预览接口返回正确的 Nginx 配置文本 -* diff 接口返回正确的域名变更列表 -* 两个接口均不产生数据库写入 - -## 6. 变更控制 - -开发中如果出现以下情况,需要先调整计划再继续编码: - -* V2 目标发生变化 -* 需要引入新的中间件(Redis、MQ 等) -* 需要新增核心数据模型 -* 需要把控制面扩展到 Nginx 全局配置 - -计划更新时,应同步修改: - -* `docs/design.md` -* `docs/development-guidelines.md` -* `docs/development-plan.md` +# ATSFlare 开发计划 + +## 1. 目标 + +当前开发目标是完成 ATSFlare 的 MVP 闭环: + +```text +后台改规则 -> 点击发布 -> Agent 拉到新版本 -> 写入 Nginx 路由配置 -> nginx 校验并 reload -> 节点状态可见 +``` + +MVP 已于第一版完成。当前进入第二版迭代。 + +## 2. 第一版里程碑(已完成) + +### Phase 1: Server 数据层与发布闭环 ✅ + +交付: + +* 新模型(proxy_routes、config_versions) +* AutoMigrate +* 路由 CRUD API +* 发布 API +* 激活版本 API +* 渲染 service + +### Phase 2: Agent API 与节点状态 ✅ + +交付: + +* `nodes`、`apply_logs` +* Agent Token 鉴权(全局单 Token) +* 节点在线状态计算 +* 节点与应用日志查询接口 + +### Phase 3: Agent 本体 ✅ + +交付: + +* 本地配置文件读取 +* `node_id` 持久化 +* 心跳循环 +* 版本检查 +* 下载配置 +* 写入 Nginx 路由配置 +* 配置校验、reload 与失败回滚 +* 应用结果上报 + +### Phase 4: 管理端页面 ✅ + +交付: + +* 反代规则页 +* 版本页 +* 节点页 +* 应用记录页 + +### Phase 5: 联调与收尾 ✅ + +交付: + +* 手工部署说明 +* Agent 配置示例 +* 联调验证记录 + +## 3. 第二版里程碑 + +### V2 Phase 1: HTTPS/TLS 支持 ✅ + +目标: + +* 支持通过控制面配置 HTTPS 路由 +* 渲染出包含 443 端口的 `server` 块 +* 支持 HTTP → HTTPS 重定向块 +* 控制面支持托管证书 + +交付: + +* `proxy_routes` 新增字段:`enable_https`、`cert_id`、`redirect_http` +* 渲染器支持 HTTPS server 块生成 +* 前端反代规则页增加 HTTPS 配置表单 +* `tls_certificates` 表与模型 +* 证书导入能力:手动导入(粘贴 PEM)与文件导入 +* AutoMigrate 覆盖新字段 + +完成标准: + +* 创建含 HTTPS 字段的路由并发布,Agent 拉取后 Nginx 能以 HTTPS 正确转发 +* HTTP 重定向配置生效 +* 证书可通过控制面导入并被 HTTPS 路由引用 +* 未开启 HTTPS 的路由渲染行为与第一版保持一致 + +### V2 Phase 2: 域名管理与证书自动匹配 ✅ + +目标: + +* 控制面可管理域名并绑定证书 +* 反代规则编辑时按输入域名自动匹配证书 +* 支持 `*.example.com` 通配符证书匹配 + +交付: + +* `managed_domains` 表与模型 +* 域名管理 CRUD API +* 证书匹配 API(精确匹配 + 通配符匹配) +* 前端域名管理页 +* 前端反代规则页接入证书自动匹配 + +完成标准: + +* 创建域名并绑定证书后,反代规则输入域名可自动匹配证书 +* 通配符证书可匹配子域名(如 `api.example.com` 匹配 `*.example.com`) +* 无匹配证书时前端给出明确提示 + +### V2 Phase 3: Agent Token 管理 + +目标: + +* 支持多个命名 Token +* Token 可通过管理界面创建和撤销 +* 中间件改为查库验证 + +交付: + +* `agent_tokens` 表与模型 +* agent-auth 中间件改造(查表验证) +* 全局 Token 环境变量降级为 bootstrap 模式 +* Token CRUD API +* 前端 Token 管理页 + +完成标准: + +* 新建 Token 后 Agent 可用该 Token 访问 Agent API +* 撤销 Token 后 Agent 请求立即返回 401 +* 全局环境变量 Token 仅在数据库无有效记录时生效 + +### V2 Phase 4: 路由自定义头 + +目标: + +* 每条路由支持追加自定义 `proxy_set_header` 指令 + +交付: + +* `proxy_routes` 新增 `custom_headers` 字段(JSON,`[{"key":"X-My-Header","value":"foo"}]`) +* 渲染器按 `custom_headers` 在 `location /` 块中注入额外 header 指令 +* 前端反代规则页增加自定义头编辑器 + +完成标准: + +* 路由配置自定义头后发布,渲染结果包含对应 `proxy_set_header` 指令 +* 不配置自定义头的路由渲染行为与之前保持一致 + +### V2 Phase 5: 配置预览与变更摘要 + +目标: + +* 发布前可预览渲染结果 +* 发布前可查看与当前激活版本的变更摘要 + +交付: + +* `GET /api/config-versions/preview` 接口(返回渲染后的 Nginx 配置,不写库) +* `GET /api/config-versions/diff` 接口(返回新增/删除/修改的域名列表) +* 前端发布版本页增加"预览"按钮和变更摘要展示弹窗 + +完成标准: + +* 点击预览可查看即将生成的 Nginx 配置文本 +* 变更摘要正确列出相对于当前激活版本的域名变化 +* 预览和 diff 操作不产生版本记录 + +## 4. 第二版建议执行顺序 + +建议严格按以下顺序开发: + +1. HTTPS/TLS 支持(对现有渲染链路影响最小,独立可测) +2. 域名管理与证书自动匹配(与 HTTPS 强相关,尽早完成闭环) +3. Agent Token 管理(安全性改善,早做早稳) +4. 路由自定义头(纯增量,对现有结构影响小) +5. 配置预览与变更摘要(纯只读接口,最后补充) + +不要先做以下内容: + +* 节点分组与差异化下发 +* 灰度百分比发布 +* WAF、限流 +* Redis、Prometheus +* 多租户 + +## 5. 第二版每阶段验收检查 + +### V2 Phase 1 检查项 + +* `enable_https=true` 的路由发布后生成 443 端口 server 块 +* `redirect_http=true` 的路由生成 80 → 443 重定向块 +* 支持手动导入证书与文件导入证书 +* 未开启 HTTPS 的路由渲染结果不受影响 +* Agent 拉取后 Nginx reload 成功 + +### V2 Phase 2 检查项 + +* 可通过管理界面维护域名并绑定证书 +* 反代规则输入域名后可自动匹配证书 +* 通配符证书可匹配子域名(`*.example.com`) + +### V2 Phase 3 检查项 + +* 可通过管理界面创建 Token +* 新 Token 可被 Agent 使用 +* 撤销 Token 后访问立即失败 +* 全局 Token 环境变量在数据库有记录时失效 + +### V2 Phase 4 检查项 + +* 路由可添加自定义头 +* 渲染结果中正确包含自定义 header 指令 +* 无自定义头的路由渲染结果不受影响 + +### V2 Phase 5 检查项 + +* 预览接口返回正确的 Nginx 配置文本 +* diff 接口返回正确的域名变更列表 +* 两个接口均不产生数据库写入 + +## 6. 变更控制 + +开发中如果出现以下情况,需要先调整计划再继续编码: + +* V2 目标发生变化 +* 需要引入新的中间件(Redis、MQ 等) +* 需要新增核心数据模型 +* 需要把控制面扩展到 Nginx 全局配置 + +计划更新时,应同步修改: + +* `docs/design.md` +* `docs/development-guidelines.md` +* `docs/development-plan.md`