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) {