MCPcopy Create free account
hub / github.com/0xUnixIO/pulse / tick

Method tick

internal/nodeagent/usage.go:103–157  ·  view source on GitHub ↗

tick 执行单次推送循环。导出名仅用于测试。

(ctx context.Context)

Source from the content-addressed store, hash-verified

101
102// tick 执行单次推送循环。导出名仅用于测试。
103func (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 数量(仅供测试/监控用)。
160func (p *UsagePusher) PendingCount() int {

Calls 6

currentSenderMethod · 0.95
AddMethod · 0.80
DoUsageMethod · 0.65
PushEventMethod · 0.65
WaitAckMethod · 0.65
DeleteMethod · 0.65