diff --git a/SAAS_P2_DESIGN.md b/SAAS_P2_DESIGN.md new file mode 100644 index 0000000..4441f22 --- /dev/null +++ b/SAAS_P2_DESIGN.md @@ -0,0 +1,142 @@ +# SaaS P2 —— 用量计量设计(Usage Metering) + +> 承接 P1(tenant_id 全链路已闭合,见 `SAAS_DESIGN.md` / DEPTH_ROADMAP T4.A)。 +> 本阶段目标:让 token 用量可**按租户持久归集**,成为**计费 / 配额**的事实源。 +> 纯后端、自包含;桌面端不改;不接支付(P5)。 + +--- + +## 1. 目标与非目标 + +**目标** +- token 用量按 **tenant** 维度持久落库(不再只有易失的 Redis 日计数)。 +- 计量单位为**积分(credits)**:可配 token↔积分汇率,扣量走扣积分; + **「按 token 直计」= 汇率 1:1 的退化情形**(同一条代码路径,配置切换,不做两套)。 +- 每条用量同时按 **Pricing 折算成金额**(内部账/毛利),可对账、可回溯重算。 +- 按 **租户 / 天** rollup,供 admin 查用量 + 积分消耗 + 成本趋势。 + +**非目标(划清边界,不过早做)** +- ❌ 支付 / 发票 / 订阅扣款 —— P5。 +- ❌ 套餐配额 enforcement(超额拒绝)—— P4,需先有 plan→quota 的产品定义。 +- ❌ 实时流式计量 —— 收尾一次事件已够,不做 per-token 流。 +- ❌ 前端计费页 —— 本阶段只出接口 + 数据,页面后续。 + +--- + +## 2. 现状(已有的一半) + +``` +提交(gateway) task.Meta[MetaUserID]=userID ──PublishTask──▶ dispatcher +dispatcher 收尾 emitUsage ── UsageEvent{task,user,model,tok...} ──▶ gateway +gateway SubscribeUsage ── AddUsage(user,day,tok) ──▶ Redis 计数(48h TTL) +门控 提交前 GetUsage(user,today) vs USER_DAILY_TOKEN_BUDGET +``` + +**够什么**:单用户当日预算门控。 +**缺什么**:① 无 tenant 维度;② Redis 48h 易失,计费要永久可审计;③ 有 `Pricing`(input/output per 1K)但没折算金额;④ 无聚合,查趋势要现算。 + +--- + +## 3. 设计 + +### 3.1 tenant_id 打通链路(接 P1) +- **提交**:`task.Meta[MetaTenantID] = tenantID(c)`(网关已有 `tenantID(c)` 助手)。 +- **契约**:`contract.UsageEvent` 加 `TenantID string`。 +- **dispatcher**:`emitUsage` 从 `t.Meta[MetaTenantID]` 读租户,带进事件。 +- 一处新增 Meta 键、一个契约字段、一行读取 —— 复用今天刚打通的 owner/tenant 提交链。 + +### 3.2 持久明细表 `sundynix_usage_event`(追加式,只增不改) +| 列 | 说明 | +|---|---| +| id / created_at | 雪花 + 时间(BaseModel) | +| tenant_id / owner | 租户 + 提交者 user.id(隔离 + 归属) | +| task_id | 关联任务 | +| model | 计费模型名 | +| prompt_tok / comp_tok / total_tok | 输入 / 输出 / 合计 | +| credits_micro | 折算积分,**int64 微积分**(×10⁻⁶),避取整损耗;展示时约整 | +| cost_micros | 折算金额,**int64 微单位**(币种最小单位 ×10⁻⁶),避浮点累积误差 | +| currency | 币种(跟 Pricing) | +| exceeded | 是否触顶单任务预算被中止 | + +- 每条 `UsageEvent` 落一行 —— **审计 / 对账 / 重算的事实源**。 +- 为何不复用 Redis:Redis 48h 就没了;计费要永久、可审计、可按新价重算历史。 +- 受租户表(标 `isTenantScoped()`);但消费在 background ctx(NATS 回写)→ tenant 从事件显式带,不靠插件自动填(与 Eval 同模式)。 + +### 3.3 成本折算 +- 落库时查 `Pricing[model]` → `cost = prompt/1000*input_per_1k + comp/1000*output_per_1k`,存 `cost_micros`(整数)。 +- **Pricing 缺失** → cost=0 + 照落事件(**计量不因缺价而丢量**;缺价可事后补价重算)。 +- tokens 是「估算」(契约原注释)→ 计费口径认这个近似;要精确则后续接 provider 真实 usage(已知项,非本轮)。 + +### 3.4 聚合表 `sundynix_usage_rollup`(租户 / 天 快照) +| 列 | 说明 | +|---|---| +| tenant_id + day | 联合唯一(`20060102`) | +| total_tok / credits_micro / cost_micros / task_count | 当天累加 | +| currency | 币种 | + +- 消费 `UsageEvent` 时**同步 upsert**(`ON CONFLICT(tenant_id,day) DO UPDATE ... = 原值 + 增量`)。 +- 明细表(event)给审计/重算,快照表(rollup)给快速读 —— 经典 append + materialized。 +- 配额/账单查 rollup **一行**,不扫明细。 + +### 3.6 计量单位:积分(credits)与 token 汇率 +**核心**:积分是租户面唯一计量/扣量单位;token 直计 = 汇率 1:1 的退化。一条路径,不做两套。 + +**汇率配置**(admin 可配) +- 全局默认 `TOKENS_PER_CREDIT`(如 1000 token = 1 积分;设 1 即 token 直计)。 +- 可选 **每模型权重** `credit_weight`(premium 模型每 token 烧更多积分)——复用 `Pricing` + 表加一列 `credit_weight float64`(缺省 1.0),不新建表。 +- 落库:`credits_micro = total_tok * 1e6 / TOKENS_PER_CREDIT * credit_weight`(整数运算,保精度)。 + +**积分账本 `sundynix_credit_ledger`(append-only,扣量的事实源)** +| 列 | 说明 | +|---|---| +| tenant_id | 租户 | +| kind | `grant`(充值/发放,+) / `usage`(消耗,−) / `adjust`(人工校正) | +| credits_micro | 带符号增量 | +| ref | 关联(usage→task_id,grant→订单/发放单号) | +| memo | 备注 | + +- 租户当前余额 = `SUM(credits_micro)`;为快取在 `sundynix_tenant.credit_balance_micro` + 维护物化余额(消费/充值时增量更新),查询读一列,不 SUM 全账本。 +- **一致性**:先写 ledger 明细,再增量改 balance;balance 可由 ledger 随时重算兜底。 +- 幂等:`usage` 分录对 `(tenant_id, ref=task_id, kind)` 唯一,防重投重复扣分。 + +**扣量 vs 拦截(分级,别过早强制)** +- **P2(本阶段,记账)**:每条 usage 记 `usage` 分录 + 扣物化余额。余额可为负(软扣,先记账不拦)。 +- **P4(enforcement)**:提交前查余额≤0 则拒绝新任务(硬闸)。需先有**套餐→发放积分**的产品定义,故 enforcement 留后。 + 开关 `CREDIT_ENFORCE`(默认 off) 控制软/硬,schema 现在就备好,改行为时不改表。 + +### 3.5 观测接口(admin,系统级) +- `GET /api/v1/admin/usage?tenant=&from=&to=` → 从 rollup 出某租户(或全平台)用量 + 积分消耗 + 成本趋势,附租户当前积分余额。 +- 走 `store.WithoutTenant`(系统级跨租户,与 admin/overview 同口径)。 +- 前端页后续;本阶段先接口 + 数据。 + +--- + +## 4. 落地顺序(增量) + +- **增量1 — 链路 + 明细 + 折算**:`MetaTenantID` + `UsageEvent.TenantID` + dispatcher 带租户; + `Pricing.credit_weight` 列 + `TOKENS_PER_CREDIT` 配置;`usage_event` 表 + 消费时折算 credits/cost 落明细。 + live 验证:跑一个任务 → 一行带 tenant + owner + credits + cost。 +- **增量2 — 账本 + 余额 + rollup**:`credit_ledger` 表(usage 分录,幂等)+ `tenant.credit_balance_micro` + 物化余额(软扣,可负);`usage_rollup` 表 + upsert。 + live 验证:多任务 → 余额递减正确、rollup 累加正确。 +- **增量3 — 观测接口**:`GET /admin/usage`(用量 + 积分 + 成本 + 余额)。 +- (前端用量页 / P4 硬拦截 `CREDIT_ENFORCE` / 充值发放 —— 后续,需求拉动) + +--- + +## 5. 权衡与已知风险(诚实标注,不一次做满) + +1. **投递可靠性**:`UsageEvent` 现走普通 NATS、收尾 best-effort,**丢事件=丢计费**。 + → 本轮先落表(已远强于 Redis)。若要**严格计费**,把 usage 升级到 **JetStream 持久投递**(像入库队列那样带重投/幂等)——列为已知项,等"要精确对账"的需求真来了再上,不提前做。 + → 幂等兜底:`usage_event` 可对 `task_id` 加唯一约束(一任务一计量),防重投重复计费。 +2. **成本依赖 Pricing 配全**:缺价记 0 不丢量,可事后补价重算明细 → 修 rollup。 +3. **token 估算**:计费认近似值;要精确接 provider usage 是后续。 +4. **rollup 与 event 一致性**:先写 event 再 upsert rollup;若 rollup 失败,event 仍在 → 可由 event 重建 rollup(提供一个重算命令兜底)。 + +--- + +## 6. 与 GPT 重构方案的关系 +- 采纳其正确诊断:**usage 先落事件、不接支付**(P2 与 P5 解耦)。 +- 不引入独立计费服务 / 消息中间件重构 —— 复用现有 NATS + PG,按需求拉动。 diff --git a/sundynix-dispatcher/internal/eino/orchestrator.go b/sundynix-dispatcher/internal/eino/orchestrator.go index 1a4d558..8700507 100644 --- a/sundynix-dispatcher/internal/eino/orchestrator.go +++ b/sundynix-dispatcher/internal/eino/orchestrator.go @@ -157,8 +157,9 @@ func (o *Orchestrator) emitUsage(t *contract.Task, b *harness.Budget) { return } uid, _ := t.Meta[contract.MetaUserID].(string) + tid, _ := t.Meta[contract.MetaTenantID].(string) if err := o.usageSink.PublishUsage(&contract.UsageEvent{ - TaskID: t.ID, UserID: uid, PromptTok: p, CompTok: c, TotalTok: total, + TaskID: t.ID, UserID: uid, TenantID: tid, PromptTok: p, CompTok: c, TotalTok: total, Exceeded: b.Exceeded(), TS: time.Now().UnixMilli(), }); err != nil { log.Printf("[usage] 回写用量失败 task=%s: %v", t.ID, err) diff --git a/sundynix-gateway/cmd/server/main.go b/sundynix-gateway/cmd/server/main.go index 8e35253..695bffb 100644 --- a/sundynix-gateway/cmd/server/main.go +++ b/sundynix-gateway/cmd/server/main.go @@ -104,9 +104,14 @@ func main() { return } 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 去重)。 + if err := db.SaveUsageEvent(context.Background(), ev); err != nil { + log.Printf("[gateway] 落库用量明细 task=%s 失败: %v", ev.TaskID, err) + } }); err != nil { log.Printf("[gateway] subscribe usage: %v", err) } diff --git a/sundynix-gateway/internal/handler/task_handler.go b/sundynix-gateway/internal/handler/task_handler.go index 7d6b399..669314e 100644 --- a/sundynix-gateway/internal/handler/task_handler.go +++ b/sundynix-gateway/internal/handler/task_handler.go @@ -59,6 +59,7 @@ func (h *Handler) SubmitTask(c *gin.Context) { // 附上用户标识(召回偏好记忆)与会话标识(召回短期多轮历史)。 // 真实场景由鉴权/会话中间件注入;此处用请求头,缺省匿名/默认会话。 task.Meta[contract.MetaUserID] = userID(c) + task.Meta[contract.MetaTenantID] = tenantID(c) task.Meta[contract.MetaSessionID] = sessionID(c) // 输入护栏灰区升级:Tier1(中间件)判为疑似的输入打标,Dispatcher 执行前调 LLM 分类器裁决。 if c.GetBool("guardrail_suspect") { diff --git a/sundynix-gateway/internal/store/model.go b/sundynix-gateway/internal/store/model.go index d529570..5698d42 100644 --- a/sundynix-gateway/internal/store/model.go +++ b/sundynix-gateway/internal/store/model.go @@ -249,10 +249,11 @@ func (p *Postgres) ListLinks(ctx context.Context, owner, kb string) ([]DocLink, // 表名 sundynix_pricing。ModelID 关联 sundynix_model.id,唯一。 type Pricing struct { BaseModel - ModelID string `gorm:"size:24;uniqueIndex"` // 关联 sundynix_model.id - InputPer1K float64 `gorm:"column:input_per_1k"` // 每 1K 输入 token 单价 - OutputPer1K float64 `gorm:"column:output_per_1k"` // 每 1K 输出 token 单价 - Currency string `gorm:"size:8"` // 币种(CNY / USD…) + ModelID string `gorm:"size:24;uniqueIndex"` // 关联 sundynix_model.id + InputPer1K float64 `gorm:"column:input_per_1k"` // 每 1K 输入 token 单价 + OutputPer1K float64 `gorm:"column:output_per_1k"` // 每 1K 输出 token 单价 + Currency string `gorm:"size:8"` // 币种(CNY / USD…) + CreditWeight float64 `gorm:"column:credit_weight"` // 积分权重(每 token 烧积分的倍率;0/缺省按 1.0 计) } func (Pricing) TableName() string { return "sundynix_pricing" } diff --git a/sundynix-gateway/internal/store/pgsql.go b/sundynix-gateway/internal/store/pgsql.go index 8378525..680f8f4 100644 --- a/sundynix-gateway/internal/store/pgsql.go +++ b/sundynix-gateway/internal/store/pgsql.go @@ -66,7 +66,7 @@ func OpenPostgres(dsn string) *Postgres { migrateLegacyIntIDs(db) migrateDocLinkToID(db) - if err := db.AutoMigrate(&User{}, &Task{}, &Eval{}, &LLMModel{}, &KB{}, &Doc{}, &Agent{}, &DocLink{}, &Pricing{}, &Prompt{}, &AuditLog{}, &GuardrailEvent{}, &Tenant{}, &TenantMember{}); err != nil { + if err := db.AutoMigrate(&User{}, &Task{}, &Eval{}, &LLMModel{}, &KB{}, &Doc{}, &Agent{}, &DocLink{}, &Pricing{}, &Prompt{}, &AuditLog{}, &GuardrailEvent{}, &Tenant{}, &TenantMember{}, &UsageEvent{}); err != nil { log.Printf("[store] postgres AutoMigrate 失败,降级运行: %v", err) return &Postgres{} } diff --git a/sundynix-gateway/internal/store/usage.go b/sundynix-gateway/internal/store/usage.go new file mode 100644 index 0000000..379db1c --- /dev/null +++ b/sundynix-gateway/internal/store/usage.go @@ -0,0 +1,100 @@ +package store + +import ( + "context" + "os" + "strconv" + + "gorm.io/gorm/clause" + + "github.com/sundynix/sundynix-shared/contract" +) + +// UsageEvent 是一条持久化的任务用量明细(追加式,只增不改)——计费 / 对账 / 重算的事实源。 +// 由网关消费 dispatcher 回写的 contract.UsageEvent 落库;tenant/owner 从事件显式带(background ctx 无请求租户)。 +// 表名 sundynix_usage_event。对 task_id 唯一 → 幂等(防重投重复计费)。 +type UsageEvent struct { + BaseModel + TenantID string `gorm:"size:64;index"` + Owner string `gorm:"size:64;index"` // 提交者 user.id(= 事件 UserID) + TaskID string `gorm:"size:64;uniqueIndex"` // 一任务一计量 → 幂等键 + Model string `gorm:"size:64"` // 计费模型名(空=按激活 chat 模型近似) + PromptTok int + CompTok int + TotalTok int + CreditsMicro int64 `gorm:"column:credits_micro"` // 折算积分 ×10⁻⁶ + CostMicros int64 `gorm:"column:cost_micros"` // 折算金额(币种最小单位 ×10⁻⁶) + Currency string `gorm:"size:8"` + Exceeded bool + TS int64 +} + +func (UsageEvent) TableName() string { return "sundynix_usage_event" } +func (UsageEvent) isTenantScoped() {} + +// TokensPerCredit 读 token→积分 汇率(env TOKENS_PER_CREDIT,缺省/非法=1000)。设 1 即 token 直计。 +func TokensPerCredit() float64 { + if v := os.Getenv("TOKENS_PER_CREDIT"); v != "" { + if n, err := strconv.ParseFloat(v, 64); err == nil && n > 0 { + return n + } + } + return 1000 +} + +// SaveUsageEvent 折算 credits + cost 并幂等落一条用量明细。 +// 折算模型:ev.Model 优先;为空则回退当前激活 chat 模型(近似——忽略 failover 到备用模型的情形)。 +// 缺 Pricing → cost=0、weight=1(计量不因缺价而丢量,可事后补价重算)。 +func (p *Postgres) SaveUsageEvent(ctx context.Context, ev *contract.UsageEvent) error { + if p.db == nil { + return nil + } + model := ev.Model + if model == "" { + if cfg := p.ActiveConfig(ctx, contract.ConfigKindChat); cfg != nil { + model = cfg.Model + } + } + + weight := 1.0 + currency := "" + var costMicros int64 + if pr := p.pricingForModelName(ctx, model); pr != nil { + if pr.CreditWeight > 0 { + weight = pr.CreditWeight + } + currency = pr.Currency + // cost(币种单位)= tok/1000 * per1k;×10⁶ 存微单位(整数)。 + cost := float64(ev.PromptTok)/1000*pr.InputPer1K + float64(ev.CompTok)/1000*pr.OutputPer1K + costMicros = int64(cost * 1e6) + } + // credits_micro = total_tok / tokensPerCredit * weight,×10⁶ 存微积分。 + creditsMicro := int64(float64(ev.TotalTok) / TokensPerCredit() * weight * 1e6) + + row := &UsageEvent{ + TenantID: ev.TenantID, Owner: ev.UserID, TaskID: ev.TaskID, Model: model, + PromptTok: ev.PromptTok, CompTok: ev.CompTok, TotalTok: ev.TotalTok, + CreditsMicro: creditsMicro, CostMicros: costMicros, Currency: currency, + Exceeded: ev.Exceeded, TS: ev.TS, + } + // 幂等:同一 task_id 已有明细则不重复插入(防 NATS 重投重复计费)。 + return p.db.WithContext(ctx).Clauses(clause.OnConflict{ + Columns: []clause.Column{{Name: "task_id"}}, + DoNothing: true, + }).Create(row).Error +} + +// pricingForModelName 按模型名查计价(join model 表,pricing 以 model_id 关联)。查不到返回 nil。 +func (p *Postgres) pricingForModelName(ctx context.Context, name string) *Pricing { + if p.db == nil || name == "" { + return nil + } + var pr Pricing + err := p.db.WithContext(ctx). + Joins("JOIN sundynix_model m ON m.id = sundynix_pricing.model_id"). + Where("m.model = ?", name).First(&pr).Error + if err != nil { + return nil + } + return &pr +} diff --git a/sundynix-shared/contract/task.go b/sundynix-shared/contract/task.go index ceb45cc..c7f60fd 100644 --- a/sundynix-shared/contract/task.go +++ b/sundynix-shared/contract/task.go @@ -62,6 +62,7 @@ const ( type UsageEvent struct { TaskID string `json:"task_id"` UserID string `json:"user_id,omitempty"` + TenantID string `json:"tenant_id,omitempty"` // 租户标识(按租户计量 / 计费) Model string `json:"model,omitempty"` PromptTok int `json:"prompt_tok"` // 输入 token(估算) CompTok int `json:"comp_tok"` // 输出 token(估算) @@ -129,6 +130,8 @@ type ApprovalDecision struct { const ( // MetaUserID 是 Task.Meta 中承载已登录用户标识的键(用于偏好记忆召回)。 MetaUserID = "user_id" + // MetaTenantID 是 Task.Meta 中承载租户标识的键(用于用量按租户计量 / 计费)。 + MetaTenantID = "tenant_id" // MetaSessionID 是 Task.Meta 中承载会话标识的键(用于短期多轮历史)。 MetaSessionID = "session_id" // MetaSafetyCheck 是输入护栏「灰区升级」标志:网关 Tier1(归一化+正则)判为疑似但不确定时置 true,