Files
Blizzard 940330cdb7 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>
2026-07-18 11:52:40 +08:00

157 lines
6.5 KiB
Go
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
// Package nats 是调度器对共享 bus 的薄封装(消费任务 / 回写 Token)。
package nats
import (
"context"
"log"
"time"
sharedbus "github.com/sundynix/sundynix-shared/bus"
"github.com/sundynix/sundynix-shared/contract"
)
// TaskHandler 处理单个任务。
type TaskHandler func(ctx context.Context, t *contract.Task) error
// Subscriber 包装共享 bus,向调度器暴露消费能力。
type Subscriber struct {
inner *sharedbus.Bus
}
// MustConnect 接入 NATS 并确保任务流存在(消费者声明在 Consume 时完成)。
func MustConnect(url string) *Subscriber {
inner, err := sharedbus.Connect(url)
if err != nil {
log.Fatalf("[dispatcher/nats] connect: %v", err)
}
if err := inner.EnsureTaskStream(context.Background()); err != nil {
log.Fatalf("[dispatcher/nats] ensure stream: %v", err)
}
// 状态/用量回写流:dispatcher 是发布方,js.Publish 需流已存在(幂等,网关也会 ensure)。
if err := inner.EnsureStatusStream(context.Background()); err != nil {
log.Fatalf("[dispatcher/nats] ensure status stream: %v", err)
}
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}
}
// ConsumeTasks 从 sundynix.tasks.* 持续消费任务(队列组负载均衡),阻塞至 ctx 取消。
func (s *Subscriber) ConsumeTasks(ctx context.Context, h TaskHandler) error {
drain, err := s.inner.ConsumeTasks(ctx, func(c context.Context, t *contract.Task) error {
return h(c, t)
})
if err != nil {
return err
}
<-ctx.Done() // 收到停机信号
// 优雅停机:停止消费新任务 + 等在途任务跑完(至多 DrainTimeout),再退出。
to := sharedbus.DrainTimeout()
log.Printf("[dispatcher] 收到停机信号,drain 在途任务(≤%s)…", to)
dctx, cancel := context.WithTimeout(context.Background(), to)
defer cancel()
drain(dctx)
log.Printf("[dispatcher] drain 完成,退出")
return nil
}
// PublishToken / CompleteStream 让 Subscriber 满足 eino.TokenSink
// 把推理 Token 回流到 sundynix.streams.<taskID>。
func (s *Subscriber) PublishToken(taskID string, token []byte) error {
return s.inner.PublishToken(taskID, token)
}
func (s *Subscriber) CompleteStream(taskID string) error {
return s.inner.CompleteStream(taskID)
}
// PublishExec / CompleteExec 让 Subscriber 满足 eino.ExecSink
// 把执行轨迹事件回流到 sundynix.exec.<taskID>(与 Token 流分开)。
func (s *Subscriber) PublishExec(taskID string, data []byte) error {
return s.inner.PublishExec(taskID, data)
}
func (s *Subscriber) CompleteExec(taskID string) error {
return s.inner.CompleteExec(taskID)
}
// CallTool 让 Subscriber 满足 eino.ToolCaller,经 NATS request-reply 调起第 5 层 MCP 工具。
func (s *Subscriber) CallTool(ctx context.Context, subject string, call *contract.ToolCall) (*contract.ToolResult, error) {
return s.inner.CallTool(ctx, subject, call)
}
// ServeHealth 在 dispatcher 心跳主题上应答探活,让管理端「服务状态」判定其在线。
func (s *Subscriber) ServeHealth(provide func() []byte) (func() error, error) {
return s.inner.ServeHealth(contract.SubjectHealthDispatcher, provide)
}
// PublishTaskStatus 让 Subscriber 满足 eino.StatusSink,把任务状态流转回写给网关。
func (s *Subscriber) PublishTaskStatus(taskID, status, detail string) error {
return s.inner.PublishTaskStatus(&contract.TaskStatusEvent{
TaskID: taskID, Status: status, Detail: detail, TS: time.Now().UnixMilli(),
})
}
// PublishEval 让 Subscriber 满足 eino.EvalSink,把自动化评测结果回写给网关落库。
func (s *Subscriber) PublishEval(ev *contract.EvalEvent) error {
return s.inner.PublishEval(ev)
}
// PublishUsage 让 Subscriber 满足 eino.UsageSink,把任务 token 用量回写给网关累计/计费。
func (s *Subscriber) PublishUsage(ev *contract.UsageEvent) error {
return s.inner.PublishUsage(ev)
}
// WaitApproval 让 Subscriber 满足 eino.ApprovalWaiter,阻塞等待审批节点的人工决定。
func (s *Subscriber) WaitApproval(ctx context.Context, taskID string, timeout time.Duration) (*contract.ApprovalDecision, error) {
return s.inner.WaitApproval(ctx, taskID, timeout)
}
// Checkpoints 打开 HITL 持久化中断的 KV 桶(compose checkpoint + resume 记录)。
// 返回的 *KVHandle 结构化满足 eino.CheckpointKVGet/Put/Delete),可直接 SetCheckpoints。
func (s *Subscriber) Checkpoints(ctx context.Context, bucket string, ttl time.Duration) (*sharedbus.KVHandle, error) {
return s.inner.Checkpoints(ctx, bucket, ttl)
}
// EnsureApprovalStream 确保审批决定流存在(中断/恢复模型用,决定持久抗 dispatcher 离线)。
func (s *Subscriber) EnsureApprovalStream(ctx context.Context) error {
return s.inner.EnsureApprovalStream(ctx)
}
// ConsumeApprovals 持久消费审批决定,每条交给 h 驱动 resume(队列组多副本安全)。
func (s *Subscriber) ConsumeApprovals(ctx context.Context, h func(context.Context, *contract.ApprovalDecision)) (func(context.Context), error) {
return s.inner.ConsumeApprovals(ctx, h)
}
// RequestModelConfig 向控制面(Gateway)取当前激活的对话模型配置。
func (s *Subscriber) RequestModelConfig(ctx context.Context) (*contract.ModelConfig, error) {
return s.inner.RequestConfig(ctx, contract.ConfigKindChat)
}
// SubscribeModelConfigUpdated 订阅对话模型配置热更新。
func (s *Subscriber) SubscribeModelConfigUpdated(onUpdate func(*contract.ModelConfig)) (func() error, error) {
return s.inner.SubscribeConfigUpdated(contract.ConfigKindChat, onUpdate)
}
// FetchModelConfigWithRetry 后台重试拉取初始对话模型配置(容忍 gateway 晚于 dispatcher 启动)。
func (s *Subscriber) FetchModelConfigWithRetry(ctx context.Context, apply func(*contract.ModelConfig)) {
s.inner.RequestConfigWithRetry(ctx, contract.ConfigKindChat, apply)
}
// FetchPromptsWithRetry 后台重试拉取初始激活 prompt 集(覆盖内置默认)。
func (s *Subscriber) FetchPromptsWithRetry(ctx context.Context, apply func(map[string]string)) {
s.inner.RequestPromptsWithRetry(ctx, apply)
}
// SubscribePromptsUpdated 订阅 prompt 激活集热更新。
func (s *Subscriber) SubscribePromptsUpdated(onUpdate func(map[string]string)) (func() error, error) {
return s.inner.SubscribePromptsUpdated(onUpdate)
}
func (s *Subscriber) Close() { s.inner.Close() }