fix: 状态/用量回写升级 JetStream 持久 —— 堵住漏账(core NATS fire-and-forget)
问题:usage(计费)/status 回写走 core NATS,网关离线/慢消费者/NATS 抖动期间 dispatcher 发的事件直接丢——任务照跑照烧 token,但这次计费凭空消失(漏账),且零重试零对账。 (对比:提交/审批/入库本就 JetStream 持久,唯独回写是 best-effort。) 修复(照 tasks/approvals 套路): - 新增 JetStream 流 SUNDYNIX_USAGE(MaxAge 72h) / SUNDYNIX_STATUS(24h),捕获 usage.task/status.task。 - PublishUsage/PublishTaskStatus 改 js.Publish(同步等 stream ack);dispatcher+gateway 启动各自 ensure 流。 - ConsumeUsage/ConsumeTaskStatus 持久消费者 + 显式 ack:落库成功 Ack、失败 Nak 重投自愈、脏数据 Term。 - 幂等保证 at-least-once 安全:usage_event.task_id 唯一 + 门控;SaveUsageEvent 返回 inserted, 仅新插入才累计 Redis 日计数(非幂等旁路,防重投重复累加);UpdateTaskStatus 按 task_id 覆盖幂等。 live 验证(复现原漏账场景):提交任务→立刻杀网关→dispatcher 跑完把 usage 发进持久流 (网关离线,usage_event=0 但流积压 1 条=钱没丢)→重启网关→自动补消费:任务 done、 usage_event 补上、公司A 余额扣 0.098、消费者 num_pending/ack_pending 归零。旧设计下这笔会永久丢失。 eval 回写仍 core NATS(仅观测,低价值,暂不改)。 Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
+85
-25
@@ -333,30 +333,60 @@ func (b *Bus) Ping(ctx context.Context, subject string) ([]byte, error) {
|
||||
return msg.Data, nil
|
||||
}
|
||||
|
||||
// ---- 任务生命周期状态回写(core NATS pub-sub)----
|
||||
// ---- 任务生命周期状态回写(JetStream 持久,at-least-once + 幂等落库)----
|
||||
|
||||
// PublishTaskStatus 广播一次任务状态流转(dispatcher 调用)。
|
||||
// EnsureStatusStream 幂等地创建/更新状态回写流,持久捕获 sundynix.status.task。
|
||||
func (b *Bus) EnsureStatusStream(ctx context.Context) error {
|
||||
_, err := b.js.CreateOrUpdateStream(ctx, jetstream.StreamConfig{
|
||||
Name: contract.StreamStatus,
|
||||
Subjects: []string{contract.SubjectTaskStatus},
|
||||
Storage: jetstream.FileStorage,
|
||||
MaxAge: 24 * time.Hour, // 状态一天内必被消费,过期回收防无限增长
|
||||
})
|
||||
return err
|
||||
}
|
||||
|
||||
// PublishTaskStatus 发布一次任务状态流转到持久流(dispatcher 调用);同步等 stream ack,失败即返回。
|
||||
func (b *Bus) PublishTaskStatus(ev *contract.TaskStatusEvent) error {
|
||||
data, err := json.Marshal(ev)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
return b.nc.Publish(contract.SubjectTaskStatus, data)
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
|
||||
defer cancel()
|
||||
_, err = b.js.Publish(ctx, contract.SubjectTaskStatus, data)
|
||||
return err
|
||||
}
|
||||
|
||||
// SubscribeTaskStatus 订阅任务状态流转(网关调用,落 PG + 推 UI)。
|
||||
// 队列组订阅:多网关副本下每条状态只由一个副本落库,避免重复写(HA)。
|
||||
func (b *Bus) SubscribeTaskStatus(onEvent func(*contract.TaskStatusEvent)) (unsub func() error, err error) {
|
||||
sub, err := b.nc.QueueSubscribe(contract.SubjectTaskStatus, contract.QueueGateway, func(m *nats.Msg) {
|
||||
var ev contract.TaskStatusEvent
|
||||
if json.Unmarshal(m.Data, &ev) == nil {
|
||||
onEvent(&ev)
|
||||
}
|
||||
// ConsumeTaskStatus 持久消费任务状态并落库(网关调用)。h 返回 error → Nak 重投(落库失败自愈);
|
||||
// nil → Ack。多网关副本共用同一 durable,队列式分摊(HA)。落库幂等,重投不致错。
|
||||
func (b *Bus) ConsumeTaskStatus(ctx context.Context, h func(context.Context, *contract.TaskStatusEvent) error) (drain func(context.Context), err error) {
|
||||
cons, err := b.js.CreateOrUpdateConsumer(ctx, contract.StreamStatus, jetstream.ConsumerConfig{
|
||||
Durable: contract.ConsumerStatus,
|
||||
AckPolicy: jetstream.AckExplicitPolicy,
|
||||
FilterSubject: contract.SubjectTaskStatus,
|
||||
AckWait: time.Minute,
|
||||
MaxAckPending: 512,
|
||||
})
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("subscribe task status: %w", err)
|
||||
return nil, fmt.Errorf("create status consumer: %w", err)
|
||||
}
|
||||
return sub.Unsubscribe, nil
|
||||
cc, err := cons.Consume(func(msg jetstream.Msg) {
|
||||
var ev contract.TaskStatusEvent
|
||||
if json.Unmarshal(msg.Data(), &ev) != nil {
|
||||
_ = msg.Term() // 脏数据,丢弃不重投
|
||||
return
|
||||
}
|
||||
if herr := h(extractTrace(context.Background(), nats.Header(msg.Headers())), &ev); herr != nil {
|
||||
_ = msg.Nak() // 落库失败 → 重投兜底
|
||||
return
|
||||
}
|
||||
_ = msg.Ack()
|
||||
})
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("consume task status: %w", err)
|
||||
}
|
||||
return func(context.Context) { cc.Stop() }, nil
|
||||
}
|
||||
|
||||
// ---- 自动化评测结果回写(core NATS pub-sub)----
|
||||
@@ -384,28 +414,58 @@ func (b *Bus) SubscribeEval(onEvent func(*contract.EvalEvent)) (unsub func() err
|
||||
return sub.Unsubscribe, nil
|
||||
}
|
||||
|
||||
// PublishUsage 广播一次任务 token 用量(dispatcher 收尾调用)。
|
||||
// EnsureUsageStream 幂等地创建/更新用量回写流,持久捕获 sundynix.usage.task(计费凭据不丢)。
|
||||
func (b *Bus) EnsureUsageStream(ctx context.Context) error {
|
||||
_, err := b.js.CreateOrUpdateStream(ctx, jetstream.StreamConfig{
|
||||
Name: contract.StreamUsage,
|
||||
Subjects: []string{contract.SubjectUsage},
|
||||
Storage: jetstream.FileStorage,
|
||||
MaxAge: 72 * time.Hour, // 计费留足追赶窗口(网关长时间离线后仍能补齐)
|
||||
})
|
||||
return err
|
||||
}
|
||||
|
||||
// PublishUsage 发布一次任务 token 用量到持久流(dispatcher 收尾调用);同步等 stream ack。
|
||||
func (b *Bus) PublishUsage(ev *contract.UsageEvent) error {
|
||||
data, err := json.Marshal(ev)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
return b.nc.Publish(contract.SubjectUsage, data)
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
|
||||
defer cancel()
|
||||
_, err = b.js.Publish(ctx, contract.SubjectUsage, data)
|
||||
return err
|
||||
}
|
||||
|
||||
// SubscribeUsage 订阅 token 用量(网关调用,累计到用户日预算 / 计费)。
|
||||
// 队列组:多副本下每条用量只累加一次,避免日预算被重复计(HA)。
|
||||
func (b *Bus) SubscribeUsage(onEvent func(*contract.UsageEvent)) (unsub func() error, err error) {
|
||||
sub, err := b.nc.QueueSubscribe(contract.SubjectUsage, contract.QueueGateway, func(m *nats.Msg) {
|
||||
var ev contract.UsageEvent
|
||||
if json.Unmarshal(m.Data, &ev) == nil {
|
||||
onEvent(&ev)
|
||||
}
|
||||
// ConsumeUsage 持久消费用量并计费(网关调用)。h 返回 error → Nak 重投;nil → Ack。
|
||||
// 落库幂等(usage_event.task_id 唯一 + 门控)→ 重投绝不重复扣费。多网关副本队列式分摊(HA)。
|
||||
func (b *Bus) ConsumeUsage(ctx context.Context, h func(context.Context, *contract.UsageEvent) error) (drain func(context.Context), err error) {
|
||||
cons, err := b.js.CreateOrUpdateConsumer(ctx, contract.StreamUsage, jetstream.ConsumerConfig{
|
||||
Durable: contract.ConsumerUsage,
|
||||
AckPolicy: jetstream.AckExplicitPolicy,
|
||||
FilterSubject: contract.SubjectUsage,
|
||||
AckWait: time.Minute,
|
||||
MaxAckPending: 512,
|
||||
})
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("subscribe usage: %w", err)
|
||||
return nil, fmt.Errorf("create usage consumer: %w", err)
|
||||
}
|
||||
return sub.Unsubscribe, nil
|
||||
cc, err := cons.Consume(func(msg jetstream.Msg) {
|
||||
var ev contract.UsageEvent
|
||||
if json.Unmarshal(msg.Data(), &ev) != nil {
|
||||
_ = msg.Term()
|
||||
return
|
||||
}
|
||||
if herr := h(extractTrace(context.Background(), nats.Header(msg.Headers())), &ev); herr != nil {
|
||||
_ = msg.Nak()
|
||||
return
|
||||
}
|
||||
_ = msg.Ack()
|
||||
})
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("consume usage: %w", err)
|
||||
}
|
||||
return func(context.Context) { cc.Stop() }, nil
|
||||
}
|
||||
|
||||
// ---- 人工审批(HITL,core NATS pub-sub)----
|
||||
|
||||
@@ -46,6 +46,13 @@ const (
|
||||
StreamApprovals = "SUNDYNIX_APPROVALS" // 审批决定 JetStream 流(持久,决定不因 dispatcher 离线而丢)
|
||||
ConsumerApprovals = "approval-resumers" // 审批决定持久消费者(队列组:多副本下每条决定只一个副本处理 resume)
|
||||
|
||||
// 状态 / 用量回写升级为 JetStream 持久(此前 core NATS:网关离线/慢消费者会丢——尤其 usage 丢=漏账)。
|
||||
// 落库幂等(task_id 唯一 + 门控),故 at-least-once 重投安全,不会重复扣费。subject 均在 tasks.> 之外。
|
||||
StreamStatus = "SUNDYNIX_STATUS" // 任务状态回写流(持久,done/failed 不因网关离线而丢)
|
||||
ConsumerStatus = "gateway-status" // 状态回写持久消费者(队列组:多网关副本每条只落一次)
|
||||
StreamUsage = "SUNDYNIX_USAGE" // 用量回写流(持久,计费凭据不丢)
|
||||
ConsumerUsage = "gateway-usage" // 用量回写持久消费者
|
||||
|
||||
// BucketCheckpoints 是 HITL 持久化中断的 JetStream KV 桶名:存 compose 图 checkpoint
|
||||
// (键=task_id)与 resume 记录(键=pending:task_id),dispatcher 重启后可据此恢复在途审批。
|
||||
BucketCheckpoints = "SUNDYNIX_CHECKPOINTS"
|
||||
|
||||
Reference in New Issue
Block a user