Files
sundynix-agentix/sundynix-dispatcher/internal/nats/subscriber.go
T
Blizzard faa1871760 feat(harness): 成本/Token 预算护栏 —— 单任务硬上限 + 单用户日预算(恒温器最后一环)
token 用量估算计量(CJK≈1/字、ASCII≈1/4字,无需分词器,护栏够用)。

单任务硬上限(dispatcher):Budget 挂 ctx 沿图透传,各 LLM 节点(对话/ReAct/compose/
报告)入口计输入 token、出口计输出,触顶即中止整图——防失控成本(死循环/超大报告)。
报告路径优雅降级:触顶跳过剩余章节出部分稿,不整体失败。预算来源 Meta.token_budget
或 env TASK_TOKEN_BUDGET(默认 20 万)。

单用户日预算(gateway):dispatcher 收尾经 NATS 回写 UsageEvent → 网关按用户按天累计
Redis(48h 过期自滚动)→ 提交前门控 USER_DAILY_TOKEN_BUDGET(0=不限,超额 402)。
/billing 升级为真实用量:当日已用 / 日预算 / 余额。

契约 UsageEvent + MetaTokenBudget + SubjectUsage;bus Publish/SubscribeUsage;
orchestrator SetUsageSink + 预算触顶 failed(不计熔断)。harness budget 5 单测,三模块全绿。
live:单任务 budget=30 → failed(已用约689);用户日 budget=200 → /billing remaining=0 → 402。

至此 harness 由「测温计」完成向「恒温器」的演进(评测闭环/纠偏/忠实度/脱敏/输入护栏/预算六项)。

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-26 10:00:51 +08:00

115 lines
4.3 KiB
Go
Raw 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)
}
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 {
stop, err := s.inner.ConsumeTasks(ctx, func(c context.Context, t *contract.Task) error {
return h(c, t)
})
if err != nil {
return err
}
defer stop()
<-ctx.Done()
return ctx.Err()
}
// 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)
}
// 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)
}
func (s *Subscriber) Close() { s.inner.Close() }