tick 执行单次推送循环。导出名仅用于测试。
(ctx context.Context)
| 101 | |
| 102 | // tick 执行单次推送循环。导出名仅用于测试。 |
| 103 | func (p *UsagePusher) tick(ctx context.Context) { |
| 104 | sender := p.currentSender() |
| 105 | if sender == nil { |
| 106 | return |
| 107 | } |
| 108 | |
| 109 | // 第一次启动:清零 xray 累计,建立 baseline,不 push。 |
| 110 | if !p.primed { |
| 111 | _ = p.api.DoUsage(true) |
| 112 | p.primed = true |
| 113 | return |
| 114 | } |
| 115 | |
| 116 | // 取一次 delta 快照(不 reset)。 |
| 117 | stats := p.api.DoUsage(false) |
| 118 | body, err := json.Marshal(stats) |
| 119 | if err != nil { |
| 120 | p.logger.Warn("nodeagent: marshal usage failed", "err", err) |
| 121 | return |
| 122 | } |
| 123 | |
| 124 | // 重发尚未 ack 的旧 seq(按 seq 升序,server 端去重)。 |
| 125 | var pendings []pendingUsage |
| 126 | p.pending.Range(func(k, v any) bool { |
| 127 | pendings = append(pendings, v.(pendingUsage)) |
| 128 | return true |
| 129 | }) |
| 130 | for _, pu := range pendings { |
| 131 | if err := sender.PushEvent("", "usage_push", pu.body, pu.seq); err != nil { |
| 132 | p.logger.Warn("nodeagent: re-push usage failed", "seq", pu.seq, "err", err) |
| 133 | return |
| 134 | } |
| 135 | } |
| 136 | |
| 137 | seq := p.nextSeq.Add(1) |
| 138 | pend := pendingUsage{seq: seq, body: body} |
| 139 | p.pending.Store(seq, pend) |
| 140 | |
| 141 | if err := sender.PushEvent("", "usage_push", body, seq); err != nil { |
| 142 | p.logger.Warn("nodeagent: push usage failed", "seq", seq, "err", err) |
| 143 | return |
| 144 | } |
| 145 | |
| 146 | // 异步等待 ack:成功 → DoUsage(true) 推进 baseline;超时 → 保留 pending 给下轮重发。 |
| 147 | go func(s Sender, seq uint64) { |
| 148 | waitCtx, cancel := context.WithTimeout(ctx, p.ackWait) |
| 149 | defer cancel() |
| 150 | if err := s.WaitAck(waitCtx, seq); err != nil { |
| 151 | p.logger.Debug("nodeagent: usage ack timeout", "seq", seq, "err", err) |
| 152 | return |
| 153 | } |
| 154 | _ = p.api.DoUsage(true) |
| 155 | p.pending.Delete(seq) |
| 156 | }(sender, seq) |
| 157 | } |
| 158 | |
| 159 | // PendingCount 返回当前未 ack 的 seq 数量(仅供测试/监控用)。 |
| 160 | func (p *UsagePusher) PendingCount() int { |