1ee1e87371
问题: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>
141 lines
5.4 KiB
Go
141 lines
5.4 KiB
Go
// Package nats 是网关对共享 bus 的薄封装(发布任务 / 订阅 Token 回流)。
|
||
package nats
|
||
|
||
import (
|
||
"context"
|
||
"encoding/json"
|
||
"log"
|
||
|
||
sharedbus "github.com/sundynix/sundynix-shared/bus"
|
||
"github.com/sundynix/sundynix-shared/contract"
|
||
)
|
||
|
||
// Bus 包装共享 bus,向网关其余代码暴露发布能力。
|
||
type Bus struct {
|
||
inner *sharedbus.Bus
|
||
}
|
||
|
||
// MustConnect 接入 NATS 并确保任务流存在。
|
||
func MustConnect(url string) *Bus {
|
||
inner, err := sharedbus.Connect(url)
|
||
if err != nil {
|
||
log.Fatalf("[nats] connect: %v", err)
|
||
}
|
||
if err := inner.EnsureTaskStream(context.Background()); err != nil {
|
||
log.Fatalf("[nats] ensure stream: %v", err)
|
||
}
|
||
if err := inner.EnsureIngestStream(context.Background()); err != nil {
|
||
log.Fatalf("[nats] ensure ingest stream: %v", err)
|
||
}
|
||
if err := inner.EnsureStatusStream(context.Background()); err != nil {
|
||
log.Fatalf("[nats] ensure status stream: %v", err)
|
||
}
|
||
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)
|
||
return &Bus{inner: inner}
|
||
}
|
||
|
||
// PublishTask 把组装后的 Task 发布到 sundynix.tasks.<id>。
|
||
func (b *Bus) PublishTask(ctx context.Context, t *contract.Task) error {
|
||
seq, err := b.inner.PublishTask(ctx, t)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
log.Printf("[nats] published task %s (seq=%d)", t.ID, seq)
|
||
return nil
|
||
}
|
||
|
||
// SubscribeTokens 订阅 sundynix.streams.<taskID> 的 Token 回流,
|
||
// 每个 Token 触发 onToken,流结束触发 onDone,返回 unsub。
|
||
func (b *Bus) SubscribeTokens(taskID string, onToken func([]byte), onDone func()) (func() error, error) {
|
||
return b.inner.SubscribeTokens(taskID, onToken, onDone)
|
||
}
|
||
|
||
// SubscribeExec 订阅 sundynix.exec.<taskID> 的执行轨迹事件(用于"运行·观测"SSE)。
|
||
func (b *Bus) SubscribeExec(taskID string, onEvent func([]byte), onDone func()) (func() error, error) {
|
||
return b.inner.SubscribeExec(taskID, onEvent, onDone)
|
||
}
|
||
|
||
// CallTool 经 NATS 同步调用一个 MCP 工具(用于网关侧写偏好记忆等)。
|
||
func (b *Bus) CallTool(ctx context.Context, subject string, call *contract.ToolCall) (*contract.ToolResult, error) {
|
||
return b.inner.CallTool(ctx, subject, call)
|
||
}
|
||
|
||
// Ping 同步探测某节点健康(如 dispatcher 心跳主题)。无人应答 / 超时即返回错误(视为下线)。
|
||
func (b *Bus) Ping(ctx context.Context, subject string) ([]byte, error) {
|
||
return b.inner.Ping(ctx, subject)
|
||
}
|
||
|
||
// ConsumeTaskStatus 持久消费 dispatcher 回写的任务生命周期状态(落 PG,at-least-once + 幂等)。
|
||
func (b *Bus) ConsumeTaskStatus(ctx context.Context, h func(context.Context, *contract.TaskStatusEvent) error) (func(context.Context), error) {
|
||
return b.inner.ConsumeTaskStatus(ctx, h)
|
||
}
|
||
|
||
// PublishApproval 把一次人工审批决定发给 dispatcher(解除审批节点阻塞)。
|
||
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)
|
||
}
|
||
|
||
// ConsumeUsage 持久消费 dispatcher 回写的任务 token 用量(计费,at-least-once + 幂等)。
|
||
func (b *Bus) ConsumeUsage(ctx context.Context, h func(context.Context, *contract.UsageEvent) error) (func(context.Context), error) {
|
||
return b.inner.ConsumeUsage(ctx, h)
|
||
}
|
||
|
||
// ServeConfig 让网关作为配置控制面,响应某 kind 的配置请求。
|
||
func (b *Bus) ServeConfig(kind string, provide func() *contract.ModelConfig) (func() error, error) {
|
||
return b.inner.ServeConfig(kind, provide)
|
||
}
|
||
|
||
// PublishConfigUpdated 广播某 kind 的配置变更。
|
||
func (b *Bus) PublishConfigUpdated(kind string, cfg *contract.ModelConfig) error {
|
||
return b.inner.PublishConfigUpdated(kind, cfg)
|
||
}
|
||
|
||
// PublishIngest 把一条入库进度事件发到 sundynix.streams.<jobID>。
|
||
func (b *Bus) PublishIngest(jobID string, ev *contract.IngestEvent) error {
|
||
data, err := json.Marshal(ev)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
return b.inner.PublishToken(jobID, data)
|
||
}
|
||
|
||
// CompleteStream 发送入库流结束信号。
|
||
func (b *Bus) CompleteStream(jobID string) error { return b.inner.CompleteStream(jobID) }
|
||
|
||
// ---- 入库工作队列(JetStream 持久)----
|
||
|
||
// EnsureIngestStream 确保入库作业流存在(启动时调)。
|
||
func (b *Bus) EnsureIngestStream(ctx context.Context) error { return b.inner.EnsureIngestStream(ctx) }
|
||
|
||
// PublishIngestJob 把入库作业持久入队(崩溃不丢,由 worker 消费)。
|
||
func (b *Bus) PublishIngestJob(ctx context.Context, job *contract.IngestJob) error {
|
||
return b.inner.PublishIngestJob(ctx, job)
|
||
}
|
||
|
||
// ConsumeIngestJobs 启动入库 worker 池消费作业(有界并发 + 背压 + 崩溃重投)。
|
||
func (b *Bus) ConsumeIngestJobs(ctx context.Context, h sharedbus.IngestHandler) (func(context.Context), error) {
|
||
return b.inner.ConsumeIngestJobs(ctx, h)
|
||
}
|
||
|
||
// ---- Prompt 控制面 ----
|
||
|
||
// ServePrompts 让网关响应「取全部激活 prompt」请求。
|
||
func (b *Bus) ServePrompts(provide func() map[string]string) (func() error, error) {
|
||
return b.inner.ServePrompts(provide)
|
||
}
|
||
|
||
// PublishPromptsUpdated 广播最新激活集 → 各服务热更新。
|
||
func (b *Bus) PublishPromptsUpdated(m map[string]string) error {
|
||
return b.inner.PublishPromptsUpdated(m)
|
||
}
|
||
|
||
func (b *Bus) Close() { b.inner.Close() }
|