From 4c3002a60db4206918ff955a1d00ca705fe07cc7 Mon Sep 17 00:00:00 2001 From: ning <710leo@gmail.com> Date: Fri, 4 Sep 2026 00:32:45 +0800 Subject: [PATCH] =?UTF-8?q?fix:=20=E7=BB=93=E6=9E=9C=E5=9B=9E=E5=86=99?= =?UTF-8?q?=E4=B8=8D=E5=86=8D=E9=98=BB=E5=A1=9E=E4=BB=BB=E5=8A=A1=E4=B8=8B?= =?UTF-8?q?=E5=8F=91=EF=BC=8C=E4=BF=AE=E5=A4=8D=E8=B7=A8=E6=9C=BA=E6=88=BF?= =?UTF-8?q?=20edge=20=E4=BB=BB=E5=8A=A1=E9=80=9A=E9=81=93=E6=AD=BB?= =?UTF-8?q?=E9=94=81?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 线上现象:edge 上的 categraf 收不到任何新任务,日志里每 6 秒一条 `AI ibex rpc call Server.Report fail: timeout`,持续数小时不自愈。 成因是 Report 里把「回写执行结果」和「下发待办任务」放在了同一次同步调用里。 回写是逐条落库的,edge 上每条还要发一次 HTTP 到 center,跨机房往返约 200ms。 一旦 agent 手里攒了一批结果(线上是 1356 条),单次 Report 需要约 141 秒, 而 agent 的 RPC 超时是 5 秒。于是: agent 超时断开 -> Locals.Clean 不执行 -> 本地结果永远清不掉 -> 下一轮原样重发 -> 服务端再跑一遍 -> 永远超时 清理逻辑被它自己要清理的东西挡在门外,而且 AssignTasks 一直发不出去, 这台机器的任务通道就被永久堵死了。跨机房部署下能承受的结果数上限只有 5s / 200ms ≈ 25 条,对一个专为跨地域设计的组件来说过低。 本次改动: 1. src/server/rpc/method.go 回写改为后台 goroutine,Report 立即返回 AssignTasks,两件事解耦。 用 sync.Map 保证每个 ident 同时只有一个回写在跑,防止 agent 重发时 goroutine 堆积(线上观察到 edge 到 center 的连接数因此从 126 涨到 160)。 跳过的批次打 Warning 留痕,因为 agent 侧 Clean 之后就取不回来了。 2. src/server/rpc/method.go handleDoneTask 单条回写失败改为 continue,不再 return 中断整批。 原来一条写不进去,它后面的所有结果都没有机会落库。 3. src/server/timer/timer_host_doing.go 取数失败时直接返回,保留上一轮缓存。原来错误只记日志不返回, 随后照常用残缺数据覆盖,导致 center 不通时 agent 静默收不到任务, 排查时唯一线索只有一行 models.TableRecordGets fail。 4. src/storage/redis.go IdInit 改用 SetNX。这个 id 是 edge 与 center 不通时给任务兜底发号用的, 要求跨重启不重复;无条件 Set 会让每次重启都从 IDINITIAL 重新发号, 发出与历史任务相同的 id。线上 edge 的 id 键正是精确的 2^32。 验证:gofmt 无差异;改动的三个包 go build / go vet 通过; 在 n9e-plus 里用 replace 指向本仓库全量 go build 通过。 (src/server/router 的编译错误是既有问题,未修改的树上同样存在。) --- src/server/rpc/method.go | 35 ++++++++++++++++++++-------- src/server/timer/timer_host_doing.go | 5 ++++ src/storage/redis.go | 4 +++- 3 files changed, 33 insertions(+), 11 deletions(-) diff --git a/src/server/rpc/method.go b/src/server/rpc/method.go index 1bf38a6..cc44e14 100644 --- a/src/server/rpc/method.go +++ b/src/server/rpc/method.go @@ -3,6 +3,7 @@ package rpc import ( "fmt" "os" + "sync" "github.com/toolkits/pkg/logger" @@ -37,12 +38,26 @@ func (*Server) GetTaskMeta(id int64, resp *types.TaskMetaResponse) error { return nil } +// reporting 保证同一个 ident 同时只有一个回写 goroutine 在跑。 +// agent 在拿到成功响应前会反复重发同一批结果,没有这层保护会堆积 goroutine。 +var reporting sync.Map // ident -> struct{} + func (*Server) Report(req types.ReportRequest, resp *types.ReportResponse) error { - if req.ReportTasks != nil && len(req.ReportTasks) > 0 { - err := handleDoneTask(req) - if err != nil { - resp.Message = err.Error() - return nil + // 结果回写必须与任务下发解耦。回写是逐条落库的,在 edge 上每条还要发一次 HTTP 到 + // center,跨机房部署时总耗时远超 agent 的 RPC 超时(5s)。如果同步做,一批回写不掉的 + // 结果会让 Report 迟迟不返回:agent 每次超时重连,本地结果因此永远清不掉,下一轮又原样 + // 重发,而 AssignTasks 一直发不出去——这台机器的任务通道就被永久堵死了。 + if len(req.ReportTasks) > 0 { + if _, busy := reporting.LoadOrStore(req.Ident, struct{}{}); busy { + // 上一批还在回写。这一批本次不落库,agent 侧 Clean 之后就取不回来了, + // 所以这里必须留痕。正常情况下回写是毫秒级的,只有积压补写时才会走到。 + logger.Warningf("skip report of %d task(s) from %s: previous writeback still in flight", + len(req.ReportTasks), req.Ident) + } else { + go func(r types.ReportRequest) { + defer reporting.Delete(r.Ident) + handleDoneTask(r) + }(req) } } @@ -61,7 +76,9 @@ func (*Server) Report(req types.ReportRequest, resp *types.ReportResponse) error return nil } -func handleDoneTask(req types.ReportRequest) error { +// handleDoneTask 逐条回写 agent 上报的执行结果。单条失败只记录日志并继续处理后面的, +// 不能中断整批:一条写不进去的结果会让它后面的所有结果都没有机会落库。 +func handleDoneTask(req types.ReportRequest) { count := len(req.ReportTasks) val, ok := os.LookupEnv("CONTINUOUS_OUTPUT") for i := 0; i < count; i++ { @@ -70,7 +87,7 @@ func handleDoneTask(req types.ReportRequest) error { err := models.RealTimeUpdateOutput(t.Id, req.Ident, t.Stdout, t.Stderr) if err != nil { logger.Errorf("cannot update output, id:%d, hostname:%s, clock:%d, status:%s, err: %v", t.Id, req.Ident, t.Clock, t.Status, err) - return err + continue } } else { if t.Status == "success" || t.Status == "failed" { @@ -87,12 +104,10 @@ func handleDoneTask(req types.ReportRequest) error { err := models.MarkDoneStatus(t.Id, t.Clock, req.Ident, t.Status, t.Stdout, t.Stderr, isEdgeAlertTriggered) if err != nil { logger.Errorf("cannot mark task done, id:%d, hostname:%s, clock:%d, status:%s, err: %v", t.Id, req.Ident, t.Clock, t.Status, err) - return err + continue } } } } - - return nil } diff --git a/src/server/timer/timer_host_doing.go b/src/server/timer/timer_host_doing.go index 813c1bc..00895d8 100644 --- a/src/server/timer/timer_host_doing.go +++ b/src/server/timer/timer_host_doing.go @@ -28,9 +28,13 @@ func loopCacheHostDoing() { } func cacheHostDoing() error { + // 任何一路取数失败都直接返回,保留上一轮的缓存。 + // 用残缺的数据覆盖缓存会让 agent 静默地收不到任务:Report 照常返回,只是 AssignTasks + // 空了,排查时唯一的线索只有这里的一行日志。宁可短暂用旧数据,也不要下发一个空集合。 doingsFromDb, err := models.TableRecordGets[[]models.TaskHostDoing](models.TaskHostDoing{}.TableName(), "") if err != nil { logger.Errorf("models.TableRecordGets fail: %v", err) + return err } ctx := context.Background() @@ -38,6 +42,7 @@ func cacheHostDoing() error { doingsFromRedis, err := models.CacheRecordGets[models.TaskHostDoing](ctx) if err != nil { logger.Errorf("models.CacheRecordGets fail: %v", err) + return err } set := make(map[string][]models.TaskHostDoing) diff --git a/src/storage/redis.go b/src/storage/redis.go index f691fad..2592480 100644 --- a/src/storage/redis.go +++ b/src/storage/redis.go @@ -29,8 +29,10 @@ func CacheMGet(ctx context.Context, keys []string) [][]byte { const IDINITIAL = 1 << 32 +// IdInit 初始化发号器。必须用 SetNX:这个 id 是 edge 与 center 网络不通时给任务兜底用的, +// 要求跨重启不重复。无条件 Set 会让每次重启都从 IDINITIAL 重新发号,发出与历史任务相同的 id。 func IdInit() error { - return Cache.Set(context.Background(), "id", IDINITIAL, 0).Err() + return Cache.SetNX(context.Background(), "id", IDINITIAL, 0).Err() } func IdGet() (int64, error) {