diff --git a/sundynix-dispatcher/internal/nats/subscriber.go b/sundynix-dispatcher/internal/nats/subscriber.go index e35e2d3..ac2d591 100644 --- a/sundynix-dispatcher/internal/nats/subscriber.go +++ b/sundynix-dispatcher/internal/nats/subscriber.go @@ -27,6 +27,13 @@ func MustConnect(url string) *Subscriber { 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) + } 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 695bffb..313df25 100644 --- a/sundynix-gateway/cmd/server/main.go +++ b/sundynix-gateway/cmd/server/main.go @@ -76,13 +76,13 @@ func main() { log.Printf("[gateway] serve prompts: %v", err) } - // 任务生命周期:订阅 dispatcher 回写的状态流转(running/done/failed/timeout),落 PG 供 UI 查询。 - if _, err := bus.SubscribeTaskStatus(func(ev *contract.TaskStatusEvent) { - if err := db.UpdateTaskStatus(context.Background(), ev.TaskID, ev.Status, ev.Detail); err != nil { - log.Printf("[gateway] 更新任务状态 %s=%s 失败: %v", ev.TaskID, ev.Status, err) - } - }); err != nil { - log.Printf("[gateway] subscribe task status: %v", err) + // 任务生命周期:持久消费 dispatcher 回写的状态流转(running/done/failed/timeout),落 PG 供 UI 查询。 + // JetStream at-least-once:落库失败 → Nak 重投自愈;UpdateTaskStatus 幂等(按 task_id 覆盖),重投无害。 + statusDrain, serr := bus.ConsumeTaskStatus(context.Background(), func(ctx context.Context, ev *contract.TaskStatusEvent) error { + return db.UpdateTaskStatus(ctx, ev.TaskID, ev.Status, ev.Detail) + }) + if serr != nil { + log.Printf("[gateway] consume task status: %v", serr) } // 评测闭环:订阅 dispatcher 回写的自动化评测结果,落 PG 供 UI 查询 / 质量趋势 / 门控。 @@ -98,22 +98,27 @@ func main() { log.Printf("[gateway] subscribe eval: %v", err) } - // 成本护栏:订阅 dispatcher 回写的任务 token 用量,按用户按天累计到 Redis(供提交前日预算门控 / 计费)。 - if _, err := bus.SubscribeUsage(func(ev *contract.UsageEvent) { + // 成本护栏:持久消费 dispatcher 回写的任务 token 用量(计费事实源,JetStream at-least-once + 幂等)。 + usageDrain, uerr := bus.ConsumeUsage(context.Background(), func(ctx context.Context, ev *contract.UsageEvent) error { if ev.UserID == "" || ev.TotalTok <= 0 { - return + return nil // 无效事件:Ack 丢弃,不重投 } - day := time.UnixMilli(ev.TS).Format("20060102") - // 快速配额校验:Redis 按用户按天累计(48h TTL,供提交前门控)。 - if _, err := cache.AddUsage(context.Background(), ev.UserID, day, ev.TotalTok); err != nil { - log.Printf("[gateway] 累计用量 user=%s 失败: %v", ev.UserID, err) + // 计费事实源:折算 credits + cost 幂等落持久明细(按 task_id 去重)。失败 → Nak 重投自愈。 + inserted, err := db.SaveUsageEvent(ctx, ev) + if err != nil { + return err } - // 计费事实源:折算 credits + cost 幂等落持久明细(按 task_id 去重)。 - if err := db.SaveUsageEvent(context.Background(), ev); err != nil { - log.Printf("[gateway] 落库用量明细 task=%s 失败: %v", ev.TaskID, err) + // 仅"新插入"时累计 Redis 日计数(非幂等,重投不能重复累加;供提交前日预算门控)。 + if inserted { + day := time.UnixMilli(ev.TS).Format("20060102") + if _, err := cache.AddUsage(ctx, ev.UserID, day, ev.TotalTok); err != nil { + log.Printf("[gateway] 累计用量 user=%s 失败: %v", ev.UserID, err) // best-effort,不影响已落库的计费 + } } - }); err != nil { - log.Printf("[gateway] subscribe usage: %v", err) + return nil + }) + if uerr != nil { + log.Printf("[gateway] consume usage: %v", uerr) } // 入库工作队列:启动有界并发 worker 池消费 JetStream 入库作业(背压防 OOM、崩溃重投兜底)。 @@ -155,6 +160,13 @@ func main() { ingestDrain(ictx) icancel() } + // 停状态/用量持久消费者(未 ack 的重启后由 JetStream 重投,幂等落库不丢不重)。 + if statusDrain != nil { + statusDrain(context.Background()) + } + if usageDrain != nil { + usageDrain(context.Background()) + } log.Println("[gateway] 已优雅停机") } diff --git a/sundynix-gateway/internal/nats/publisher.go b/sundynix-gateway/internal/nats/publisher.go index 300a4d5..690a18f 100644 --- a/sundynix-gateway/internal/nats/publisher.go +++ b/sundynix-gateway/internal/nats/publisher.go @@ -27,7 +27,13 @@ func MustConnect(url string) *Bus { if err := inner.EnsureIngestStream(context.Background()); err != nil { log.Fatalf("[nats] ensure ingest stream: %v", err) } - log.Printf("[nats] connected %s, task + ingest streams ready", url) + 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} } @@ -62,9 +68,9 @@ func (b *Bus) Ping(ctx context.Context, subject string) ([]byte, error) { return b.inner.Ping(ctx, subject) } -// SubscribeTaskStatus 订阅 dispatcher 回写的任务生命周期状态(落 PG)。 -func (b *Bus) SubscribeTaskStatus(onEvent func(*contract.TaskStatusEvent)) (func() error, error) { - return b.inner.SubscribeTaskStatus(onEvent) +// 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(解除审批节点阻塞)。 @@ -77,9 +83,9 @@ func (b *Bus) SubscribeEval(onEvent func(*contract.EvalEvent)) (func() error, er return b.inner.SubscribeEval(onEvent) } -// SubscribeUsage 订阅 dispatcher 回写的任务 token 用量(累计到用户日预算 / 计费)。 -func (b *Bus) SubscribeUsage(onEvent func(*contract.UsageEvent)) (func() error, error) { - return b.inner.SubscribeUsage(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 的配置请求。 diff --git a/sundynix-gateway/internal/store/usage.go b/sundynix-gateway/internal/store/usage.go index aaf1cb8..ad23e9e 100644 --- a/sundynix-gateway/internal/store/usage.go +++ b/sundynix-gateway/internal/store/usage.go @@ -57,12 +57,13 @@ func tokensPerCreditEnv() float64 { return 1000 } -// SaveUsageEvent 折算 credits + cost 并幂等落一条用量明细。 +// SaveUsageEvent 折算 credits + cost 并幂等落一条用量明细。返回 inserted(是否本次新插入): +// 供调用方据此只在"新插入"时更新非幂等的旁路(如 Redis 日计数),保证 at-least-once 重投下不重复累计。 // 折算模型:ev.Model 优先;为空则回退当前激活 chat 模型(近似——忽略 failover 到备用模型的情形)。 // 缺 Pricing → cost=0、weight=1(计量不因缺价而丢量,可事后补价重算)。 -func (p *Postgres) SaveUsageEvent(ctx context.Context, ev *contract.UsageEvent) error { +func (p *Postgres) SaveUsageEvent(ctx context.Context, ev *contract.UsageEvent) (inserted bool, err error) { if p.db == nil { - return nil + return false, nil } model := ev.Model if model == "" { @@ -96,7 +97,7 @@ func (p *Postgres) SaveUsageEvent(ctx context.Context, ev *contract.UsageEvent) // 一个事务内:落明细 → 扣积分(账本+余额) → 累加 rollup。 // 幂等锚点:usage_event 的 task_id 唯一;插入若被冲突吞掉(RowsAffected==0)→重投,跳过后续,绝不重复计费。 - return p.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error { + err = p.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error { res := tx.Clauses(clause.OnConflict{ Columns: []clause.Column{{Name: "task_id"}}, DoNothing: true, @@ -107,11 +108,13 @@ func (p *Postgres) SaveUsageEvent(ctx context.Context, ev *contract.UsageEvent) if res.RowsAffected == 0 { return nil // 重投:明细已存在,账本/余额/rollup 均不再动 } + inserted = true if err := applyUsageCredit(tx, ev.TenantID, ev.TaskID, creditsMicro); err != nil { return err } return upsertRollup(tx, ev.TenantID, day, int64(ev.TotalTok), creditsMicro, costMicros, currency) }) + return inserted, err } // RecentUsage 返回某租户最近 n 条用量明细(供用户看"最近消耗")。受租户表→用户 ctx 自动过滤, diff --git a/sundynix-shared/bus/bus.go b/sundynix-shared/bus/bus.go index 0bc412d..bfe0519 100644 --- a/sundynix-shared/bus/bus.go +++ b/sundynix-shared/bus/bus.go @@ -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)---- diff --git a/sundynix-shared/contract/task.go b/sundynix-shared/contract/task.go index c7f60fd..45b71c0 100644 --- a/sundynix-shared/contract/task.go +++ b/sundynix-shared/contract/task.go @@ -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"