From 940330cdb7716b8795d28860fbea365727b8b87b Mon Sep 17 00:00:00 2001 From: Blizzard Date: Sat, 18 Jul 2026 11:52:40 +0800 Subject: [PATCH] =?UTF-8?q?feat(eval):=20=E8=AF=84=E6=B5=8B=E7=BB=93?= =?UTF-8?q?=E6=9E=9C=E5=9B=9E=E5=86=99=E5=8D=87=20JetStream=20=E6=8C=81?= =?UTF-8?q?=E4=B9=85=20=E2=80=94=E2=80=94=20=E6=B6=88=E7=81=AD=E6=9C=80?= =?UTF-8?q?=E5=90=8E=E4=B8=80=E5=A4=84=20core=20NATS=20=E5=9B=9E=E5=86=99?= =?UTF-8?q?=20(P1)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 完成度审计 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 --- .../internal/nats/subscriber.go | 3 + sundynix-gateway/cmd/server/main.go | 19 ++++--- sundynix-gateway/internal/nats/publisher.go | 11 ++-- sundynix-shared/bus/bus.go | 57 +++++++++++++++---- sundynix-shared/bus/bus_e2e_test.go | 17 ++++-- sundynix-shared/contract/task.go | 5 +- 6 files changed, 81 insertions(+), 31 deletions(-) diff --git a/sundynix-dispatcher/internal/nats/subscriber.go b/sundynix-dispatcher/internal/nats/subscriber.go index ac2d591..d0f81cd 100644 --- a/sundynix-dispatcher/internal/nats/subscriber.go +++ b/sundynix-dispatcher/internal/nats/subscriber.go @@ -34,6 +34,9 @@ func MustConnect(url string) *Subscriber { if err := inner.EnsureUsageStream(context.Background()); err != nil { log.Fatalf("[dispatcher/nats] ensure usage stream: %v", err) } + if err := inner.EnsureEvalStream(context.Background()); err != nil { + log.Fatalf("[dispatcher/nats] ensure eval stream: %v", err) + } log.Printf("[dispatcher/nats] connected %s", url) return &Subscriber{inner: inner} } diff --git a/sundynix-gateway/cmd/server/main.go b/sundynix-gateway/cmd/server/main.go index d81cb58..a2cc71e 100644 --- a/sundynix-gateway/cmd/server/main.go +++ b/sundynix-gateway/cmd/server/main.go @@ -100,17 +100,17 @@ func main() { log.Printf("[gateway] consume task status: %v", serr) } - // 评测闭环:订阅 dispatcher 回写的自动化评测结果,落 PG 供 UI 查询 / 质量趋势 / 门控。 - if _, err := bus.SubscribeEval(func(ev *contract.EvalEvent) { + // 评测闭环:持久消费 dispatcher 回写的自动化评测结果,落 PG 供 UI 查询 / 质量趋势 / 门控。 + // JetStream at-least-once + SaveEval 按 task_id upsert 幂等;落库失败返 error → Nak 重投自愈。 + evalDrain, everr := bus.ConsumeEval(context.Background(), func(ctx context.Context, ev *contract.EvalEvent) error { flags, _ := json.Marshal(ev.Flags) - if err := db.SaveEval(context.Background(), &store.Eval{ + return db.SaveEval(ctx, &store.Eval{ TaskID: ev.TaskID, Overall: ev.Overall, Rule: ev.Rule, LLM: ev.LLM, Faithful: ev.Faithful, Level: ev.Level, Flags: string(flags), Reason: ev.Reason, Sources: ev.Sources, Corrected: ev.Corrected, - }); err != nil { - log.Printf("[gateway] 落库评测 %s 失败: %v", ev.TaskID, err) - } - }); err != nil { - log.Printf("[gateway] subscribe eval: %v", err) + }) + }) + if everr != nil { + log.Printf("[gateway] consume eval: %v", everr) } // 成本护栏:持久消费 dispatcher 回写的任务 token 用量(计费事实源,JetStream at-least-once + 幂等)。 @@ -182,6 +182,9 @@ func main() { if usageDrain != nil { usageDrain(context.Background()) } + if evalDrain != nil { + evalDrain(context.Background()) + } log.Println("[gateway] 已优雅停机") } diff --git a/sundynix-gateway/internal/nats/publisher.go b/sundynix-gateway/internal/nats/publisher.go index 690a18f..1223461 100644 --- a/sundynix-gateway/internal/nats/publisher.go +++ b/sundynix-gateway/internal/nats/publisher.go @@ -33,7 +33,10 @@ func MustConnect(url string) *Bus { if err := inner.EnsureUsageStream(context.Background()); err != nil { log.Fatalf("[nats] ensure usage stream: %v", err) } - log.Printf("[nats] connected %s, task + ingest + status + usage streams ready", url) + if err := inner.EnsureEvalStream(context.Background()); err != nil { + log.Fatalf("[nats] ensure eval stream: %v", err) + } + log.Printf("[nats] connected %s, task + ingest + status + usage + eval streams ready", url) return &Bus{inner: inner} } @@ -78,9 +81,9 @@ func (b *Bus) PublishApproval(dec *contract.ApprovalDecision) error { return b.inner.PublishApproval(dec) } -// SubscribeEval 订阅 dispatcher 回写的自动化评测结果(落 PG)。 -func (b *Bus) SubscribeEval(onEvent func(*contract.EvalEvent)) (func() error, error) { - return b.inner.SubscribeEval(onEvent) +// ConsumeEval 持久消费 dispatcher 回写的自动化评测结果(落 PG,at-least-once + 幂等)。 +func (b *Bus) ConsumeEval(ctx context.Context, h func(context.Context, *contract.EvalEvent) error) (func(context.Context), error) { + return b.inner.ConsumeEval(ctx, h) } // ConsumeUsage 持久消费 dispatcher 回写的任务 token 用量(计费,at-least-once + 幂等)。 diff --git a/sundynix-shared/bus/bus.go b/sundynix-shared/bus/bus.go index bfe0519..94cf78a 100644 --- a/sundynix-shared/bus/bus.go +++ b/sundynix-shared/bus/bus.go @@ -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(计费凭据不丢)。 diff --git a/sundynix-shared/bus/bus_e2e_test.go b/sundynix-shared/bus/bus_e2e_test.go index 815161b..5be76a2 100644 --- a/sundynix-shared/bus/bus_e2e_test.go +++ b/sundynix-shared/bus/bus_e2e_test.go @@ -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++ { diff --git a/sundynix-shared/contract/task.go b/sundynix-shared/contract/task.go index 45b71c0..5fcfedc 100644 --- a/sundynix-shared/contract/task.go +++ b/sundynix-shared/contract/task.go @@ -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。