From faa1871760f6121e9f5e75f43b2924e65e5ab939 Mon Sep 17 00:00:00 2001 From: Blizzard Date: Fri, 26 Jun 2026 10:00:51 +0800 Subject: [PATCH] =?UTF-8?q?feat(harness):=20=E6=88=90=E6=9C=AC/Token=20?= =?UTF-8?q?=E9=A2=84=E7=AE=97=E6=8A=A4=E6=A0=8F=20=E2=80=94=E2=80=94=20?= =?UTF-8?q?=E5=8D=95=E4=BB=BB=E5=8A=A1=E7=A1=AC=E4=B8=8A=E9=99=90=20+=20?= =?UTF-8?q?=E5=8D=95=E7=94=A8=E6=88=B7=E6=97=A5=E9=A2=84=E7=AE=97=EF=BC=88?= =?UTF-8?q?=E6=81=92=E6=B8=A9=E5=99=A8=E6=9C=80=E5=90=8E=E4=B8=80=E7=8E=AF?= =?UTF-8?q?=EF=BC=89?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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) --- project_analysis.md | 8 +- sundynix-dispatcher/cmd/dispatcher/main.go | 1 + .../internal/eino/compose_graph.go | 18 +++- sundynix-dispatcher/internal/eino/graph.go | 16 ++++ .../internal/eino/orchestrator.go | 95 +++++++++++++++++-- .../internal/eino/react_agent.go | 18 +++- sundynix-dispatcher/internal/eino/report.go | 12 +++ .../internal/harness/budget.go | 82 ++++++++++++++++ .../internal/harness/budget_test.go | 67 +++++++++++++ .../internal/nats/subscriber.go | 5 + sundynix-gateway/cmd/server/main.go | 14 +++ .../internal/handler/task_handler.go | 41 +++++++- sundynix-gateway/internal/nats/publisher.go | 5 + sundynix-gateway/internal/store/redis.go | 35 +++++++ sundynix-shared/bus/bus.go | 23 +++++ sundynix-shared/contract/task.go | 19 ++++ 16 files changed, 442 insertions(+), 17 deletions(-) create mode 100644 sundynix-dispatcher/internal/harness/budget.go create mode 100644 sundynix-dispatcher/internal/harness/budget_test.go diff --git a/project_analysis.md b/project_analysis.md index b226d9b..d45231f 100644 --- a/project_analysis.md +++ b/project_analysis.md @@ -107,8 +107,10 @@ Harness = 围绕 LLM 的可靠性 / 安全 / 质量治理层。4 个组件均为 | 输出脱敏 RedactSecrets | 26 行 + 39 测试 | ⭐⭐⭐ | 4 正则(sk-/AKIA/JWT/Bearer),流式逐片 | | 输入护栏 Guardrail | guardrail.go + 测试 | ⭐⭐½ | 注入正则(中英) + 体积限制 + 敏感词黑名单(默认空) | -**定性**:一个真·生产级熔断器 + 三个"能跑但偏启发式"的护栏/评测。当前更像**「测温计」(检测+记录)** -而非**「恒温器」(检测→决策→纠偏闭环)**——多数信号只打日志、不门控、不纠偏。 +**定性(更新)**:起初是**「测温计」(检测+记录)**——多数信号只打日志。经评测闭环、低分自动纠偏、 +忠实度评测、输出脱敏增强、输入护栏升级、成本预算护栏六项,现已成**「恒温器」(检测→决策→门控/纠偏闭环)**: +熔断挡后端雪崩、输入护栏两层挡注入/越狱、评测分级+低分重生成、输出脱敏挡密钥/PII、预算护栏挡失控成本。 +下一步重心转 P0/P1 生产硬骨头(优雅停机、HA、容量实测、本地模型)。 **已知短板 / 优化清单**(按性价比,逐项推进): @@ -123,7 +125,7 @@ Harness = 围绕 LLM 的可靠性 / 安全 / 质量治理层。4 个组件均为 - [x] **P3 坏输出自动纠偏** ✅:poor(<0.5) 触发评语驱动的重生成,重评后**仅采纳更优者(不退步)**, 采纳的修订版落会话历史 + 评测终值带 `corrected` 标记落库。`maxRefineRounds=1`、`canRefine`(模型就绪且熔断未开)门控; 单 goroutine 串 评测→纠偏→落历史避竞态。单测覆盖 采纳/不退步/非低分不触发,live 验证好答案不误触发。**至此 harness 由「测温计」迈入「恒温器」。** -- [ ] **P3 成本/Token 预算护栏**:单任务/单用户 token 上限与告警。 +- [x] **P3 成本/Token 预算护栏** ✅:用量估算计量(CJK/ASCII 启发式,无需分词器)。**单任务硬上限**:Budget 挂 ctx 沿图透传,各 LLM 节点(对话/ReAct/compose/报告)入口计输入、出口计输出,触顶即中止(报告优雅降级跳过剩余章节);防失控成本。**单用户日预算**:dispatcher 收尾经 NATS 回写 UsageEvent → 网关按用户按天累计 Redis → 提交前门控(USER_DAILY_TOKEN_BUDGET,超额 402),`/billing` 出当日用量/预算/余额。live:单任务 budget=30 → failed(已用约689);用户日 budget=200 → 402。 --- diff --git a/sundynix-dispatcher/cmd/dispatcher/main.go b/sundynix-dispatcher/cmd/dispatcher/main.go index cc8b383..42db8aa 100644 --- a/sundynix-dispatcher/cmd/dispatcher/main.go +++ b/sundynix-dispatcher/cmd/dispatcher/main.go @@ -55,6 +55,7 @@ func main() { log.Fatalf("[dispatcher] build eino graph: %v", err) } orch.SetGuardian(guardian) // 输入护栏 Tier2 + orch.SetUsageSink(sub) // 成本护栏:token 用量回写网关累计/计费 // 健康心跳:dispatcher 无 HTTP/工具端点,挂一个 NATS 应答让管理端「服务状态」探到它在线。 startedAt := time.Now() diff --git a/sundynix-dispatcher/internal/eino/compose_graph.go b/sundynix-dispatcher/internal/eino/compose_graph.go index 864c3e5..9cc1a78 100644 --- a/sundynix-dispatcher/internal/eino/compose_graph.go +++ b/sundynix-dispatcher/internal/eino/compose_graph.go @@ -63,6 +63,19 @@ func (o *Orchestrator) runComposeConversation(ctx context.Context, taskID string Upstream: append([]string{}, b.agentOut...), // 前序协作 agent 产出 → 接力 } msgs, _ := buildMessages(ctx, rc) + // 成本护栏:计入输入 token;触顶则中止整图。 + if bud := harness.BudgetFrom(ctx); bud != nil { + for _, m := range msgs { + bud.AddPrompt(m.Content) + } + if bud.Exceeded() { + if b.fatalErr == nil { + b.fatalErr = errBudget + } + tr.emit(node, "system", "error", "token 预算", "已达单任务预算上限,中止", 0) + return + } + } t0 := time.Now() // ChatModel 的 start/end 由 composeTracer(callbacks)落轨迹,这里不再手写 emit(归一)。 @@ -74,7 +87,7 @@ func (o *Orchestrator) runComposeConversation(ctx context.Context, taskID string defer sr.Close() chunks := 0 - var produced strings.Builder // 本节点产出(供下游 agent 接力) + var produced strings.Builder // 本节点产出(供下游 agent 接力) red := harness.NewStreamRedactor() // 输出护栏:跨分片脱敏,杜绝密钥被切断而漏检 emit := func(safe string) { if safe == "" { @@ -99,6 +112,9 @@ func (o *Orchestrator) runComposeConversation(ctx context.Context, taskID string emit(red.Push(chunk.Content)) } emit(red.Flush()) // 吐出暂留尾部 + if bud := harness.BudgetFrom(ctx); bud != nil { + bud.AddComplete(produced.String()) // 成本护栏:计入输出 token + } o.recordAgentOutput(b, produced.String()) tr.info(node, "system", "compose 图", fmt.Sprintf("%d 段输出 / %d 字(Eino compose 运行时)", chunks, len([]rune(produced.String())))) } diff --git a/sundynix-dispatcher/internal/eino/graph.go b/sundynix-dispatcher/internal/eino/graph.go index 1d6019d..1997431 100644 --- a/sundynix-dispatcher/internal/eino/graph.go +++ b/sundynix-dispatcher/internal/eino/graph.go @@ -273,6 +273,19 @@ func (o *Orchestrator) runAgent(ctx context.Context, taskID string, b *board, sy Upstream: append([]string{}, b.agentOut...), // 前序协作 agent 产出 → 接力 } msgs, _ := buildMessages(ctx, rc) + // 成本护栏:计入本节点输入 token;若已触顶则中止整图(防失控成本)。 + if bud := harness.BudgetFrom(ctx); bud != nil { + for _, m := range msgs { + bud.AddPrompt(m.Content) + } + if bud.Exceeded() { + if b.fatalErr == nil { + b.fatalErr = errBudget + } + tr.emit(node, "system", "error", "token 预算", "已达单任务预算上限,中止", 0) + return + } + } tr.emit(node, "model", "start", "模型流式推理", "", 0) t0 := time.Now() n := 0 @@ -304,6 +317,9 @@ func (o *Orchestrator) runAgent(ctx context.Context, taskID string, b *board, sy return } emit(red.Flush()) // 吐出暂留的尾部(最后一段疑似密钥的判定) + if bud := harness.BudgetFrom(ctx); bud != nil { + bud.AddComplete(produced.String()) // 成本护栏:计入本节点输出 token + } if red.Hits() > 0 { tr.info(node, "system", "输出护栏", fmt.Sprintf("已脱敏 %d 处疑似密钥/PII", red.Hits())) } diff --git a/sundynix-dispatcher/internal/eino/orchestrator.go b/sundynix-dispatcher/internal/eino/orchestrator.go index a626c0a..aa64d14 100644 --- a/sundynix-dispatcher/internal/eino/orchestrator.go +++ b/sundynix-dispatcher/internal/eino/orchestrator.go @@ -8,6 +8,8 @@ import ( "fmt" "log" "log/slog" + "os" + "strconv" "strings" "sync" "time" @@ -51,6 +53,14 @@ type EvalSink interface { PublishEval(ev *contract.EvalEvent) error } +// UsageSink 回写任务 token 用量供网关累计/计费(由 NATS bus 实现;可为 nil → 不回写)。 +type UsageSink interface { + PublishUsage(ev *contract.UsageEvent) error +} + +// errBudget 是单任务 token 预算触顶时的哨兵错误:任务以 failed 收尾并附明确原因(防失控成本)。 +var errBudget = errors.New("token 预算超限,已中止") + // errRejected 是审批节点拒绝(或超时)时图执行返回的哨兵错误:它是合法终态而非故障, // Handle 据此判 rejected 并优雅收尾(不计熔断失败)。 var errRejected = errors.New("approval rejected") @@ -79,16 +89,17 @@ const approvalTimeout = 5 * time.Minute // Orchestrator 把每个 DSL 任务动态编译为 Eino 图并执行(记忆召回 → 工具节点 → 注入 → 流式)。 type Orchestrator struct { - pool LLM - breaker *harness.CircuitBreaker - eval *harness.Evaluator - sink TokenSink - tools ToolCaller - exec ExecSink - status StatusSink // 任务生命周期状态回写(可为 nil) - approval ApprovalWaiter // HITL 审批等待(可为 nil → 审批节点自动放行) - evalSink EvalSink // 评测结果回写落库(可为 nil → 仅打日志) - guard *harness.Classifier // 输入护栏 Tier2:对网关标记的灰区任务做 LLM 裁决(可为 nil → 不做) + pool LLM + breaker *harness.CircuitBreaker + eval *harness.Evaluator + sink TokenSink + tools ToolCaller + exec ExecSink + status StatusSink // 任务生命周期状态回写(可为 nil) + approval ApprovalWaiter // HITL 审批等待(可为 nil → 审批节点自动放行) + evalSink EvalSink // 评测结果回写落库(可为 nil → 仅打日志) + guard *harness.Classifier // 输入护栏 Tier2:对网关标记的灰区任务做 LLM 裁决(可为 nil → 不做) + usageSink UsageSink // token 用量回写(可为 nil → 不回写) turnMu sync.Mutex // 保护 turns(攒批计数,多任务 goroutine 共享) turns map[string]int // sessionID → 累计轮次,用于每 N 轮触发 consolidate @@ -104,6 +115,52 @@ func NewOrchestrator(pool LLM, breaker *harness.CircuitBreaker, eval *harness.Ev // SetGuardian 注入输入护栏 Tier2 的 LLM 分类器(可选;不注入则灰区任务直接放行执行)。 func (o *Orchestrator) SetGuardian(c *harness.Classifier) { o.guard = c } +// SetUsageSink 注入 token 用量回写出口(可选;不注入则不上报用量)。 +func (o *Orchestrator) SetUsageSink(s UsageSink) { o.usageSink = s } + +// taskBudget 取本任务的 token 预算上限:优先 Meta(网关按用户/套餐下发),否则 env TASK_TOKEN_BUDGET(默认 20 万)。 +func (o *Orchestrator) taskBudget(t *contract.Task) int { + switch n := t.Meta[contract.MetaTokenBudget].(type) { + case float64: + if n > 0 { + return int(n) + } + case int: + if n > 0 { + return n + } + } + return envInt("TASK_TOKEN_BUDGET", 200000) +} + +// emitUsage 任务收尾回写本轮 token 用量(用量为 0 或无出口则跳过)。 +func (o *Orchestrator) emitUsage(t *contract.Task, b *harness.Budget) { + if o.usageSink == nil { + return + } + p, c, total := b.Snapshot() + if total == 0 { + return + } + uid, _ := t.Meta[contract.MetaUserID].(string) + if err := o.usageSink.PublishUsage(&contract.UsageEvent{ + TaskID: t.ID, UserID: uid, PromptTok: p, CompTok: c, TotalTok: total, + Exceeded: b.Exceeded(), TS: time.Now().UnixMilli(), + }); err != nil { + log.Printf("[usage] 回写用量失败 task=%s: %v", t.ID, err) + } +} + +// envInt 读正整数环境变量,缺省回退 def。 +func envInt(key string, def int) int { + if v := os.Getenv(key); v != "" { + if n, err := strconv.Atoi(v); err == nil && n > 0 { + return n + } + } + return def +} + // setStatus 回写一次任务状态流转(status 为 nil 时静默跳过)。 func (o *Orchestrator) setStatus(taskID, status, detail string) { if o.status == nil { @@ -180,6 +237,11 @@ func (o *Orchestrator) Handle(ctx context.Context, t *contract.Task) error { tctx, cancel := context.WithTimeout(ctx, taskExecTimeout) defer cancel() + // 成本护栏:单任务 token 预算挂到 ctx,沿图执行各 LLM 节点计量+封顶;收尾回写用量(计费/日预算)。 + budget := harness.NewBudget(o.taskBudget(t)) + tctx = harness.WithBudget(tctx, budget) + defer o.emitUsage(t, budget) + // 报告生成走专用多步编排(规划→分章并行检索撰写→汇聚→渲染 Word),而非通用对话图。 if intent, _ := t.Meta[contract.MetaIntent].(string); intent == contract.IntentReport { err := o.handleReport(tctx, t, tr) @@ -203,6 +265,19 @@ func (o *Orchestrator) Handle(ctx context.Context, t *contract.Task) error { o.setStatus(t.ID, contract.TaskRejected, truncate(answer, 120)) return nil } + if errors.Is(err, errBudget) { + // token 预算触顶:策略性中止,非后端故障。收尾流 + 置 failed(附明确原因),不计熔断、不重投。 + _, _, total := budget.Snapshot() + slog.WarnContext(ctx, "task aborted: token budget exceeded", "task_id", t.ID, "tokens", total) + if answer != "" { + _ = o.sink.PublishToken(t.ID, []byte(answer)) + } + _ = o.sink.PublishToken(t.ID, []byte("\n\n⚠️ 已达单任务 token 预算上限,自动中止。")) + _ = o.sink.CompleteStream(t.ID) + o.breaker.Report(true) + o.setStatus(t.ID, contract.TaskFailed, fmt.Sprintf("token 预算超限(已用约 %d)", total)) + return nil + } if err != nil { span.RecordError(err) span.SetStatus(codes.Error, err.Error()) diff --git a/sundynix-dispatcher/internal/eino/react_agent.go b/sundynix-dispatcher/internal/eino/react_agent.go index d9faa8a..bdca86e 100644 --- a/sundynix-dispatcher/internal/eino/react_agent.go +++ b/sundynix-dispatcher/internal/eino/react_agent.go @@ -202,6 +202,19 @@ func (o *Orchestrator) runReactAgent(ctx context.Context, taskID string, b *boar Upstream: append([]string{}, b.agentOut...), // 前序协作 agent 产出 → 接力 } msgs, _ := buildMessages(ctx, rc) + // 成本护栏:计入输入 token;触顶则中止整图。 + if bud := harness.BudgetFrom(ctx); bud != nil { + for _, m := range msgs { + bud.AddPrompt(m.Content) + } + if bud.Exceeded() { + if b.fatalErr == nil { + b.fatalErr = errBudget + } + tr.emit(node, "system", "error", "token 预算", "已达单任务预算上限,中止", 0) + return + } + } tr.emit(node, "model", "start", "ReAct 智能体(自主调工具)", fmt.Sprintf("%d 个工具可用", len(tools)), 0) t0 := time.Now() @@ -213,7 +226,7 @@ func (o *Orchestrator) runReactAgent(ctx context.Context, taskID string, b *boar defer sr.Close() chunks := 0 - var produced strings.Builder // 本节点产出(供下游 agent 接力) + var produced strings.Builder // 本节点产出(供下游 agent 接力) red := harness.NewStreamRedactor() // 输出护栏:跨分片脱敏,杜绝密钥被切断而漏检 emit := func(safe string) { if safe == "" { @@ -238,6 +251,9 @@ func (o *Orchestrator) runReactAgent(ctx context.Context, taskID string, b *boar emit(red.Push(chunk.Content)) } emit(red.Flush()) // 吐出暂留尾部 + if bud := harness.BudgetFrom(ctx); bud != nil { + bud.AddComplete(produced.String()) // 成本护栏:计入输出 token + } o.recordAgentOutput(b, produced.String()) tr.emit(node, "model", "end", "ReAct 智能体", fmt.Sprintf("%d 段输出 / %d 字", chunks, len([]rune(produced.String()))), time.Since(t0).Milliseconds()) diff --git a/sundynix-dispatcher/internal/eino/report.go b/sundynix-dispatcher/internal/eino/report.go index 679159f..a825082 100644 --- a/sundynix-dispatcher/internal/eino/report.go +++ b/sundynix-dispatcher/internal/eino/report.go @@ -10,6 +10,7 @@ import ( "time" "github.com/sundynix/sundynix-dispatcher/internal/dsl" + "github.com/sundynix/sundynix-dispatcher/internal/harness" "github.com/sundynix/sundynix-dispatcher/internal/llm" "github.com/sundynix/sundynix-shared/contract" ) @@ -204,6 +205,14 @@ func (o *Orchestrator) writeSection(ctx context.Context, topic, kb, heading stri ub.WriteString("\n") } ub.WriteString("请就「本章标题」撰写 200–400 字正文。只输出正文,不要重复标题、不要再列提纲。") + // 成本护栏:计入输入;若已触顶则跳过本章(报告优雅降级出部分稿,不整体失败)。 + if bud := harness.BudgetFrom(ctx); bud != nil { + bud.AddPrompt(sys + ub.String()) + if bud.Exceeded() { + tr.info(node, "system", "token 预算", "已达预算上限,跳过本章") + return "(已达 token 预算上限,本章自动跳过。)" + } + } cctx, cancel := llmCtx(ctx) defer cancel() txt, err := o.pool.Chat(cctx, []llm.ChatMessage{{Role: "system", Content: sys}, {Role: "user", Content: ub.String()}}) @@ -211,6 +220,9 @@ func (o *Orchestrator) writeSection(ctx context.Context, topic, kb, heading stri log.Printf("[report] 撰写「%s」失败: %v", heading, err) return "(本章撰写失败:" + err.Error() + ")" } + if bud := harness.BudgetFrom(ctx); bud != nil { + bud.AddComplete(txt) // 成本护栏:计入输出 + } return strings.TrimSpace(txt) } diff --git a/sundynix-dispatcher/internal/harness/budget.go b/sundynix-dispatcher/internal/harness/budget.go new file mode 100644 index 0000000..fbb526e --- /dev/null +++ b/sundynix-dispatcher/internal/harness/budget.go @@ -0,0 +1,82 @@ +package harness + +import ( + "context" + "sync" + "unicode" +) + +// EstimateTokens 粗估一段文本的 token 数(无需依赖模型/分词器,作预算护栏足够): +// CJK 字符约 1 token/字;其余(拉丁/数字/空白/标点)约 1 token/4 字符。偏保守(宁高勿低)。 +func EstimateTokens(s string) int { + cjk, other := 0, 0 + for _, r := range s { + if unicode.Is(unicode.Han, r) || unicode.Is(unicode.Hiragana, r) || + unicode.Is(unicode.Katakana, r) || unicode.Is(unicode.Hangul, r) { + cjk++ + } else { + other++ + } + } + return cjk + (other+3)/4 // 向上取整 +} + +// Budget 是单任务的 token 预算计量与封顶(并发安全:流式回调与节点循环可能并发计量)。 +// max ≤ 0 表示不限额(仅计量、不中止)。 +type Budget struct { + mu sync.Mutex + max int + prompt int + comp int +} + +// NewBudget 建一个上限为 max 的任务预算(max≤0 → 不限额,纯计量)。 +func NewBudget(max int) *Budget { return &Budget{max: max} } + +// AddPrompt / AddComplete 累计输入/输出 token(按文本估算)。 +func (b *Budget) AddPrompt(text string) { b.add(EstimateTokens(text), 0) } +func (b *Budget) AddComplete(text string) { b.add(0, EstimateTokens(text)) } + +func (b *Budget) add(p, c int) { + if b == nil { + return + } + b.mu.Lock() + b.prompt += p + b.comp += c + b.mu.Unlock() +} + +// Exceeded 报告是否已触顶(max≤0 恒为 false)。 +func (b *Budget) Exceeded() bool { + if b == nil { + return false + } + b.mu.Lock() + defer b.mu.Unlock() + return b.max > 0 && b.prompt+b.comp > b.max +} + +// Snapshot 返回当前用量快照(prompt / completion / total)。 +func (b *Budget) Snapshot() (prompt, comp, total int) { + if b == nil { + return 0, 0, 0 + } + b.mu.Lock() + defer b.mu.Unlock() + return b.prompt, b.comp, b.prompt + b.comp +} + +// budgetKey 是 context 里承载任务预算的私有键(避免跨包碰撞)。 +type budgetKey struct{} + +// WithBudget 把任务预算挂到 context,沿图执行透传(免改各节点函数签名)。 +func WithBudget(ctx context.Context, b *Budget) context.Context { + return context.WithValue(ctx, budgetKey{}, b) +} + +// BudgetFrom 取出 context 里的任务预算(无则返回 nil,调用方按"不限额"处理)。 +func BudgetFrom(ctx context.Context) *Budget { + b, _ := ctx.Value(budgetKey{}).(*Budget) + return b +} diff --git a/sundynix-dispatcher/internal/harness/budget_test.go b/sundynix-dispatcher/internal/harness/budget_test.go new file mode 100644 index 0000000..055e6ac --- /dev/null +++ b/sundynix-dispatcher/internal/harness/budget_test.go @@ -0,0 +1,67 @@ +package harness + +import ( + "context" + "testing" +) + +func TestEstimateTokens(t *testing.T) { + if got := EstimateTokens(""); got != 0 { + t.Errorf("空串应 0, got %d", got) + } + // 8 个 ASCII → (8+3)/4 = 2 + if got := EstimateTokens("abcdefgh"); got != 2 { + t.Errorf("ASCII 估算 want 2, got %d", got) + } + // 4 个中文 → 4 + if got := EstimateTokens("你好世界"); got != 4 { + t.Errorf("中文估算 want 4, got %d", got) + } +} + +func TestBudget_ExceedAndSnapshot(t *testing.T) { + b := NewBudget(10) + b.AddPrompt("你好世界") // 4 + b.AddComplete("你好") // 2 → total 6 + if b.Exceeded() { + t.Fatal("6/10 不应触顶") + } + b.AddComplete("一二三四五") // +5 → 11 + if !b.Exceeded() { + t.Fatal("11/10 应触顶") + } + p, c, total := b.Snapshot() + if p != 4 || c != 7 || total != 11 { + t.Errorf("快照 want 4/7/11, got %d/%d/%d", p, c, total) + } +} + +func TestBudget_Unlimited(t *testing.T) { + b := NewBudget(0) // 不限额 + b.AddComplete("非常非常非常长的一段中文内容反复堆叠占用大量预算额度") + if b.Exceeded() { + t.Error("max≤0 不应触顶") + } +} + +func TestBudget_NilSafe(t *testing.T) { + var b *Budget + b.AddPrompt("x") // 不应 panic + if b.Exceeded() { + t.Error("nil 预算视为不限额") + } + if _, _, total := b.Snapshot(); total != 0 { + t.Error("nil 快照应为 0") + } +} + +func TestBudget_Context(t *testing.T) { + if BudgetFrom(context.Background()) != nil { + t.Error("无预算的 ctx 应返回 nil") + } + b := NewBudget(100) + ctx := WithBudget(context.Background(), b) + if BudgetFrom(ctx) != b { + t.Error("应取回同一预算实例") + } +} diff --git a/sundynix-dispatcher/internal/nats/subscriber.go b/sundynix-dispatcher/internal/nats/subscriber.go index be5c40e..8a2a1c0 100644 --- a/sundynix-dispatcher/internal/nats/subscriber.go +++ b/sundynix-dispatcher/internal/nats/subscriber.go @@ -86,6 +86,11 @@ 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) diff --git a/sundynix-gateway/cmd/server/main.go b/sundynix-gateway/cmd/server/main.go index f7bedaa..9eda2e3 100644 --- a/sundynix-gateway/cmd/server/main.go +++ b/sundynix-gateway/cmd/server/main.go @@ -6,6 +6,7 @@ import ( "encoding/json" "log" "os" + "time" "github.com/sundynix/sundynix-gateway/internal/blob" "github.com/sundynix/sundynix-gateway/internal/nats" @@ -78,6 +79,19 @@ func main() { log.Printf("[gateway] subscribe eval: %v", err) } + // 成本护栏:订阅 dispatcher 回写的任务 token 用量,按用户按天累计到 Redis(供提交前日预算门控 / 计费)。 + if _, err := bus.SubscribeUsage(func(ev *contract.UsageEvent) { + if ev.UserID == "" || ev.TotalTok <= 0 { + return + } + day := time.UnixMilli(ev.TS).Format("20060102") + if _, err := cache.AddUsage(context.Background(), ev.UserID, day, ev.TotalTok); err != nil { + log.Printf("[gateway] 累计用量 user=%s 失败: %v", ev.UserID, err) + } + }); err != nil { + log.Printf("[gateway] subscribe usage: %v", err) + } + r := router.New(db, cache, bus, blobStore) addr := envOr("GATEWAY_ADDR", ":8080") log.Printf("[gateway] listening on %s", addr) diff --git a/sundynix-gateway/internal/handler/task_handler.go b/sundynix-gateway/internal/handler/task_handler.go index 30cddd7..10f4ef5 100644 --- a/sundynix-gateway/internal/handler/task_handler.go +++ b/sundynix-gateway/internal/handler/task_handler.go @@ -7,6 +7,8 @@ import ( "io" "log" "net/http" + "os" + "strconv" "time" "github.com/gin-contrib/sse" @@ -42,6 +44,17 @@ func (h *Handler) SubmitTask(c *gin.Context) { c.JSON(http.StatusUnprocessableEntity, gin.H{"error": err.Error()}) return } + // 成本护栏:单用户当日 token 日预算门控(USER_DAILY_TOKEN_BUDGET,0=不限)。超额则拒绝新任务。 + if budget := userDailyTokenBudget(); budget > 0 { + uid := userID(c) + used := h.cache.GetUsage(c.Request.Context(), uid, time.Now().Format("20060102")) + if used >= int64(budget) { + c.JSON(http.StatusPaymentRequired, gin.H{ + "error": "已达当日 token 预算上限", "used": used, "budget": budget, + }) + return + } + } // 附上用户标识(召回偏好记忆)与会话标识(召回短期多轮历史)。 // 真实场景由鉴权/会话中间件注入;此处用请求头,缺省匿名/默认会话。 task.Meta[contract.MetaUserID] = userID(c) @@ -362,12 +375,36 @@ func sessionID(c *gin.Context) string { return "default" } +// userDailyTokenBudget 读单用户当日 token 预算(env USER_DAILY_TOKEN_BUDGET,缺省/非法=0 即不限)。 +func userDailyTokenBudget() int { + if v := os.Getenv("USER_DAILY_TOKEN_BUDGET"); v != "" { + if n, err := strconv.Atoi(v); err == nil && n > 0 { + return n + } + } + return 0 +} + func (h *Handler) Billing(c *gin.Context) { - // TODO: 商业化与计费模块;暂以已提交任务计数演示真实读库。 n, err := h.db.CountTasks(c.Request.Context()) if err != nil { c.JSON(http.StatusBadGateway, gin.H{"error": err.Error()}) return } - c.JSON(http.StatusOK, gin.H{"status": "ok", "tasks_submitted": n, "persisted": h.db.Enabled()}) + // 成本护栏:当日 token 用量与预算(按用户隔离)。 + uid := userID(c) + used := h.cache.GetUsage(c.Request.Context(), uid, time.Now().Format("20060102")) + budget := userDailyTokenBudget() + remaining := -1 // -1 表示不限额 + if budget > 0 { + if r := budget - int(used); r > 0 { + remaining = r + } else { + remaining = 0 + } + } + c.JSON(http.StatusOK, gin.H{ + "status": "ok", "tasks_submitted": n, "persisted": h.db.Enabled(), + "token_used_today": used, "daily_budget": budget, "remaining": remaining, + }) } diff --git a/sundynix-gateway/internal/nats/publisher.go b/sundynix-gateway/internal/nats/publisher.go index 74f1eb2..93e31ce 100644 --- a/sundynix-gateway/internal/nats/publisher.go +++ b/sundynix-gateway/internal/nats/publisher.go @@ -74,6 +74,11 @@ 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) +} + // ServeConfig 让网关作为配置控制面,响应某 kind 的配置请求。 func (b *Bus) ServeConfig(kind string, provide func() *contract.ModelConfig) (func() error, error) { return b.inner.ServeConfig(kind, provide) diff --git a/sundynix-gateway/internal/store/redis.go b/sundynix-gateway/internal/store/redis.go index 03a06fd..e2d5434 100644 --- a/sundynix-gateway/internal/store/redis.go +++ b/sundynix-gateway/internal/store/redis.go @@ -50,6 +50,41 @@ func (r *Redis) Allow(ctx context.Context, key string, limit int64, window time. // ---- Token 流持久化(Redis Stream:可回放的追加日志,根治 SSE 连晚/重连丢 token)---- +// ---- Token 用量日预算(成本护栏:按用户按天累计 token,供提交前门控)---- + +// usageDailyKey 是某用户某天的累计 token key(day 形如 20260626,便于到期自然滚动)。 +func usageDailyKey(userID, day string) string { return "sundynix:usage:" + userID + ":" + day } + +// AddUsage 给某用户当天累计 token,并返回累计后总量(首次设 48h 过期自动清理)。降级返回 (0,nil)。 +func (r *Redis) AddUsage(ctx context.Context, userID, day string, tokens int) (int64, error) { + if r.rdb == nil || userID == "" { + return 0, nil + } + rk := usageDailyKey(userID, day) + n, err := r.rdb.IncrBy(ctx, rk, int64(tokens)).Result() + if err != nil { + return 0, err + } + if n == int64(tokens) { // 首次累计:设过期 + _ = r.rdb.Expire(ctx, rk, 48*time.Hour).Err() + } + return n, nil +} + +// GetUsage 取某用户当天已用 token(不存在/降级返回 0)。 +func (r *Redis) GetUsage(ctx context.Context, userID, day string) int64 { + if r.rdb == nil || userID == "" { + return 0 + } + n, err := r.rdb.Get(ctx, usageDailyKey(userID, day)).Int64() + if err != nil { + return 0 + } + return n +} + +// ---- Token 流持久化(Redis Stream:可回放的追加日志,根治 SSE 连晚/重连丢 token)---- + const streamTTL = 10 * time.Minute func streamKey(taskID string) string { return "sundynix:stream:" + taskID } diff --git a/sundynix-shared/bus/bus.go b/sundynix-shared/bus/bus.go index 3ab20df..a8c3052 100644 --- a/sundynix-shared/bus/bus.go +++ b/sundynix-shared/bus/bus.go @@ -363,6 +363,29 @@ func (b *Bus) SubscribeEval(onEvent func(*contract.EvalEvent)) (unsub func() err return sub.Unsubscribe, nil } +// PublishUsage 广播一次任务 token 用量(dispatcher 收尾调用)。 +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) +} + +// SubscribeUsage 订阅 token 用量(网关调用,累计到用户日预算 / 计费)。 +func (b *Bus) SubscribeUsage(onEvent func(*contract.UsageEvent)) (unsub func() error, err error) { + sub, err := b.nc.Subscribe(contract.SubjectUsage, func(m *nats.Msg) { + var ev contract.UsageEvent + if json.Unmarshal(m.Data, &ev) == nil { + onEvent(&ev) + } + }) + if err != nil { + return nil, fmt.Errorf("subscribe usage: %w", err) + } + return sub.Unsubscribe, nil +} + // ---- 人工审批(HITL,core NATS pub-sub)---- // PublishApproval 广播一次人工审批决定(网关在收到 UI 的批准/拒绝后调用)。 diff --git a/sundynix-shared/contract/task.go b/sundynix-shared/contract/task.go index e8e31d3..c66877c 100644 --- a/sundynix-shared/contract/task.go +++ b/sundynix-shared/contract/task.go @@ -40,8 +40,24 @@ const ( // 自动化评测结果回写:dispatcher 评完经此广播,网关订阅落 PG 并供 UI 查询。core NATS pub-sub。 SubjectEval = "sundynix.eval.task" + + // Token 用量回写:dispatcher 任务收尾经此广播本轮 token 用量,网关订阅累加到用户日预算并供计费。core NATS pub-sub。 + SubjectUsage = "sundynix.usage.task" ) +// UsageEvent 是一次任务的 token 用量(dispatcher 收尾经 SubjectUsage 回流给网关累计/计费)。 +// dispatcher 不持有计价,只报原始 token + 模型名;真钱成本由网关据 Pricing 折算。 +type UsageEvent struct { + TaskID string `json:"task_id"` + UserID string `json:"user_id,omitempty"` + Model string `json:"model,omitempty"` + PromptTok int `json:"prompt_tok"` // 输入 token(估算) + CompTok int `json:"comp_tok"` // 输出 token(估算) + TotalTok int `json:"total_tok"` // 合计 + Exceeded bool `json:"exceeded"` // 是否触顶单任务预算被中止 + TS int64 `json:"ts"` // unix 毫秒 +} + // 评测质量分级(据综合分 + 忠实度阈值,闭环门控/告警用)。 const ( EvalOK = "ok" // 综合 ≥ 0.75 且无忠实度风险 @@ -106,6 +122,9 @@ const ( // MetaSafetyCheck 是输入护栏「灰区升级」标志:网关 Tier1(归一化+正则)判为疑似但不确定时置 true, // Dispatcher 执行前据此调 LLM jailbreak 分类器(Tier2)裁决。明确干净/明确恶意的输入不带此标志,不付 LLM 成本。 MetaSafetyCheck = "safety_check" + // MetaTokenBudget 是单任务 token 预算上限(数字):网关按用户/套餐下发,Dispatcher 据此封顶单任务用量, + // 触顶即中止(防失控成本)。缺省用 Dispatcher 的 env TASK_TOKEN_BUDGET。 + MetaTokenBudget = "token_budget" // 配置控制面按 kind 寻址:sundynix.config..get / .updated。 // Gateway 持有配置,消费方(Dispatcher/mcp-go)经 NATS 取用/订阅变更。