feat(eval): 评测结果回写升 JetStream 持久 —— 消灭最后一处 core NATS 回写 (P1)
完成度审计 P1 + 记忆 nats-durability 的既定规矩「计费/需落库的回写一律 JetStream+幂等,别fire-and-forget」。此前 eval 是全仓最后一处 core NATS pub-sub 回写:网关离线/慢消费者时评测结果直接丢——而质量趋势/门控都依赖它。 照抄已升级的 status/usage 范式(同为 dispatcher→gateway→PG 回写): - contract 加 StreamEval/ConsumerEval;bus 加 EnsureEvalStream + ConsumeEval (durable consumer,AckExplicit,落库失败 Nak 重投自愈),PublishEval 改 js.Publish 同步等 ack。删 core NATS 的 SubscribeEval。 - gateway/dispatcher 两个 wrapper 在 connect 时 EnsureEvalStream;gateway main 的评测订阅从 SubscribeEval(fire-and-forget)换 ConsumeEval(handler 返 error→ Nak),接入优雅停机 drain。 - 幂等前提已满足:SaveEval 按 task_id upsert,at-least-once 重投只覆盖不重复。 验证:e2e 去重测试(TestGatewayQueueDedup)升级到 ConsumeEval,50 条两副本合计 处理50次零重复;live 真端到端:提交任务→跑完→评测经新 eval 流落 PG(level=ok) 一次成功;两服务启动 eval 流 ensure 无报错。go build/vet/test 全绿。 Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
+45
-12
@@ -389,29 +389,62 @@ func (b *Bus) ConsumeTaskStatus(ctx context.Context, h func(context.Context, *co
|
||||
return func(context.Context) { cc.Stop() }, nil
|
||||
}
|
||||
|
||||
// ---- 自动化评测结果回写(core NATS pub-sub)----
|
||||
// ---- 自动化评测结果回写(JetStream 持久,at-least-once + 幂等落库)----
|
||||
// 此前是 core NATS pub-sub:网关离线/慢消费者会丢评测结果(质量趋势/门控依赖它)。
|
||||
// 升级为持久流,与 status/usage 同级。SaveEval 按 task_id upsert 幂等,重投安全。
|
||||
|
||||
// PublishEval 广播一次评测结果(dispatcher 调用)。
|
||||
// EnsureEvalStream 幂等地创建/更新评测回写流,持久捕获 sundynix.eval.task。
|
||||
func (b *Bus) EnsureEvalStream(ctx context.Context) error {
|
||||
_, err := b.js.CreateOrUpdateStream(ctx, jetstream.StreamConfig{
|
||||
Name: contract.StreamEval,
|
||||
Subjects: []string{contract.SubjectEval},
|
||||
Storage: jetstream.FileStorage,
|
||||
MaxAge: 24 * time.Hour, // 评测一天内必被消费,过期回收防无限增长
|
||||
})
|
||||
return err
|
||||
}
|
||||
|
||||
// PublishEval 发布一次评测结果到持久流(dispatcher 调用);同步等 stream ack,失败即返回。
|
||||
func (b *Bus) PublishEval(ev *contract.EvalEvent) error {
|
||||
data, err := json.Marshal(ev)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
return b.nc.Publish(contract.SubjectEval, data)
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
|
||||
defer cancel()
|
||||
_, err = b.js.Publish(ctx, contract.SubjectEval, data)
|
||||
return err
|
||||
}
|
||||
|
||||
// SubscribeEval 订阅评测结果(网关调用,落 PG)。队列组:多副本下每条只落一次(HA)。
|
||||
func (b *Bus) SubscribeEval(onEvent func(*contract.EvalEvent)) (unsub func() error, err error) {
|
||||
sub, err := b.nc.QueueSubscribe(contract.SubjectEval, contract.QueueGateway, func(m *nats.Msg) {
|
||||
var ev contract.EvalEvent
|
||||
if json.Unmarshal(m.Data, &ev) == nil {
|
||||
onEvent(&ev)
|
||||
}
|
||||
// ConsumeEval 持久消费评测结果并落库(网关调用)。h 返回 error → Nak 重投(落库失败自愈);
|
||||
// nil → Ack。多网关副本共用同一 durable,队列式分摊(HA)。落库幂等,重投不致错。
|
||||
func (b *Bus) ConsumeEval(ctx context.Context, h func(context.Context, *contract.EvalEvent) error) (drain func(context.Context), err error) {
|
||||
cons, err := b.js.CreateOrUpdateConsumer(ctx, contract.StreamEval, jetstream.ConsumerConfig{
|
||||
Durable: contract.ConsumerEval,
|
||||
AckPolicy: jetstream.AckExplicitPolicy,
|
||||
FilterSubject: contract.SubjectEval,
|
||||
AckWait: time.Minute,
|
||||
MaxAckPending: 512,
|
||||
})
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("subscribe eval: %w", err)
|
||||
return nil, fmt.Errorf("create eval consumer: %w", err)
|
||||
}
|
||||
return sub.Unsubscribe, nil
|
||||
cc, err := cons.Consume(func(msg jetstream.Msg) {
|
||||
var ev contract.EvalEvent
|
||||
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 eval: %w", err)
|
||||
}
|
||||
return func(context.Context) { cc.Stop() }, nil
|
||||
}
|
||||
|
||||
// EnsureUsageStream 幂等地创建/更新用量回写流,持久捕获 sundynix.usage.task(计费凭据不丢)。
|
||||
|
||||
@@ -382,15 +382,20 @@ func TestGatewayQueueDedup(t *testing.T) {
|
||||
}
|
||||
defer gwB.Close()
|
||||
|
||||
// eval 流已升 JetStream;基础 bus.Connect 不 ensure(那是网关/dispatcher wrapper 的活),测试手动建。
|
||||
ctx := context.Background()
|
||||
if err := pub.EnsureEvalStream(ctx); err != nil {
|
||||
t.Fatalf("ensure eval stream: %v", err)
|
||||
}
|
||||
var total int64
|
||||
count := func(_ *contract.EvalEvent) { atomic.AddInt64(&total, 1) }
|
||||
if _, err := gwA.SubscribeEval(count); err != nil {
|
||||
t.Fatalf("gwA sub: %v", err)
|
||||
count := func(_ context.Context, _ *contract.EvalEvent) error { atomic.AddInt64(&total, 1); return nil }
|
||||
if _, err := gwA.ConsumeEval(ctx, count); err != nil {
|
||||
t.Fatalf("gwA consume: %v", err)
|
||||
}
|
||||
if _, err := gwB.SubscribeEval(count); err != nil {
|
||||
t.Fatalf("gwB sub: %v", err)
|
||||
if _, err := gwB.ConsumeEval(ctx, count); err != nil {
|
||||
t.Fatalf("gwB consume: %v", err)
|
||||
}
|
||||
time.Sleep(100 * time.Millisecond) // 等订阅就绪
|
||||
time.Sleep(100 * time.Millisecond) // 等消费者就绪
|
||||
|
||||
const n = 50
|
||||
for i := 0; i < n; i++ {
|
||||
|
||||
@@ -52,12 +52,15 @@ const (
|
||||
ConsumerStatus = "gateway-status" // 状态回写持久消费者(队列组:多网关副本每条只落一次)
|
||||
StreamUsage = "SUNDYNIX_USAGE" // 用量回写流(持久,计费凭据不丢)
|
||||
ConsumerUsage = "gateway-usage" // 用量回写持久消费者
|
||||
StreamEval = "SUNDYNIX_EVAL" // 评测结果回写流(持久,网关离线不丢评测——SaveEval 按 task_id upsert 幂等)
|
||||
ConsumerEval = "gateway-eval" // 评测回写持久消费者
|
||||
|
||||
// BucketCheckpoints 是 HITL 持久化中断的 JetStream KV 桶名:存 compose 图 checkpoint
|
||||
// (键=task_id)与 resume 记录(键=pending:task_id),dispatcher 重启后可据此恢复在途审批。
|
||||
BucketCheckpoints = "SUNDYNIX_CHECKPOINTS"
|
||||
|
||||
// 自动化评测结果回写:dispatcher 评完经此广播,网关订阅落 PG 并供 UI 查询。core NATS pub-sub。
|
||||
// 自动化评测结果回写:dispatcher 评完经此发到持久流,网关消费落 PG 供 UI 查询。
|
||||
// 已升 JetStream(此前 core NATS:网关离线/慢消费者会丢评测);SaveEval upsert 幂等,重投安全。
|
||||
SubjectEval = "sundynix.eval.task"
|
||||
|
||||
// Token 用量回写:dispatcher 任务收尾经此广播本轮 token 用量,网关订阅累加到用户日预算并供计费。core NATS pub-sub。
|
||||
|
||||
Reference in New Issue
Block a user