Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
35 changes: 25 additions & 10 deletions src/server/rpc/method.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ package rpc
import (
"fmt"
"os"
"sync"

"github.com/toolkits/pkg/logger"

Expand Down Expand Up @@ -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)
}
}

Expand All @@ -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++ {
Expand All @@ -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" {
Expand All @@ -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
}
5 changes: 5 additions & 0 deletions src/server/timer/timer_host_doing.go
Original file line number Diff line number Diff line change
Expand Up @@ -28,16 +28,21 @@ 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()

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)
Expand Down
4 changes: 3 additions & 1 deletion src/storage/redis.go
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand Down