diff --git a/sundynix-gateway/internal/handler/billing_pay.go b/sundynix-gateway/internal/handler/billing_pay.go index c57ed8c..0a8ecd2 100644 --- a/sundynix-gateway/internal/handler/billing_pay.go +++ b/sundynix-gateway/internal/handler/billing_pay.go @@ -39,10 +39,11 @@ const orderTTL = 30 * time.Minute func (h *Handler) BillingCreateOrder(c *gin.Context) { var b struct { PackID string `json:"pack_id"` + PlanID string `json:"plan_id"` // 传它=买订阅周期;与 pack_id 二选一 Channel string `json:"channel"` } - if err := c.ShouldBindJSON(&b); err != nil || strings.TrimSpace(b.PackID) == "" { - c.JSON(http.StatusBadRequest, gin.H{"error": "pack_id 必填"}) + if err := c.ShouldBindJSON(&b); err != nil || (strings.TrimSpace(b.PackID) == "" && strings.TrimSpace(b.PlanID) == "") { + c.JSON(http.StatusBadRequest, gin.H{"error": "pack_id 或 plan_id 必填"}) return } channel := strings.TrimSpace(b.Channel) @@ -61,21 +62,40 @@ func (h *Handler) BillingCreateOrder(c *gin.Context) { c.JSON(http.StatusBadRequest, gin.H{"error": "无计费租户上下文"}) return } - pk, err := h.db.GetPack(ctx, b.PackID) - if err != nil { - c.JSON(http.StatusBadRequest, gin.H{"error": "积分包不存在或已下架"}) - return - } - o := &store.PaymentOrder{ - TenantID: billing, UserID: uid, PackID: pk.ID, - AmountFen: pk.PriceFen, CreditsMicro: pk.CreditsMicro, - Channel: channel, Status: store.OrderPending, + // 订阅单与积分包单走同一条支付链路:只有订单内容不同,下单/回调/查单/掉单补偿全复用。 + var o *store.PaymentOrder + var desc string + if pid := strings.TrimSpace(b.PlanID); pid != "" { + pl := h.db.GetSubPlan(ctx, pid) + if pl == nil || !pl.Active { + c.JSON(http.StatusBadRequest, gin.H{"error": "订阅套餐不存在或已下架"}) + return + } + // 订阅单 credits_micro 恒为 0:积分不在付款时一次给,而是订阅期内按周期发放。 + o = &store.PaymentOrder{ + TenantID: billing, UserID: uid, Kind: store.OrderKindSub, PlanID: pl.ID, + AmountFen: pl.PriceFen, CreditsMicro: 0, + Channel: channel, Status: store.OrderPending, + } + desc = "sundynix 订阅 · " + pl.Name + } else { + pk, err := h.db.GetPack(ctx, b.PackID) + if err != nil { + c.JSON(http.StatusBadRequest, gin.H{"error": "积分包不存在或已下架"}) + return + } + o = &store.PaymentOrder{ + TenantID: billing, UserID: uid, PackID: pk.ID, Kind: store.OrderKindPack, + AmountFen: pk.PriceFen, CreditsMicro: pk.CreditsMicro, + Channel: channel, Status: store.OrderPending, + } + desc = "sundynix 积分充值 · " + pk.Name } if err := h.db.CreateOrder(ctx, o); err != nil { c.JSON(http.StatusBadGateway, gin.H{"error": err.Error()}) return } - intent, err := ch.CreatePay(ctx, o.ID, "sundynix 积分充值 · "+pk.Name, pk.PriceFen) + intent, err := ch.CreatePay(ctx, o.ID, desc, o.AmountFen) if err != nil { // 渠道下单失败的单直接作废,不留一堆永远付不了的 pending。 _ = h.db.ExpireOrder(ctx, o.ID) diff --git a/sundynix-gateway/internal/handler/subscription.go b/sundynix-gateway/internal/handler/subscription.go new file mode 100644 index 0000000..636dc50 --- /dev/null +++ b/sundynix-gateway/internal/handler/subscription.go @@ -0,0 +1,83 @@ +package handler + +import ( + "net/http" + "strings" + + "github.com/gin-gonic/gin" + + "github.com/sundynix/sundynix-gateway/internal/store" +) + +// 订阅 API。用户面只有「看套餐 / 看我的订阅」,下单复用现有 /billing/orders +// (多传 plan_id 即可),因为支付链路、幂等、掉单补偿都已经在那条路上验过了, +// 没必要为订阅再造一条支付路径。 + +// ---- 用户面 ---- + +// BillingSubPlans: GET /api/v1/billing/sub-plans —— 在售订阅套餐。 +func (h *Handler) BillingSubPlans(c *gin.Context) { + c.JSON(http.StatusOK, gin.H{"plans": h.db.ListSubPlans(c.Request.Context(), true)}) +} + +// MySubscription: GET /api/v1/billing/subscription —— 我的当前订阅(无则 null)。 +func (h *Handler) MySubscription(c *gin.Context) { + ctx := c.Request.Context() + billing := h.db.ResolveBillingTenantID(ctx, userID(c), tenantID(c)) + if billing == "" { + c.JSON(http.StatusOK, gin.H{"subscription": nil}) + return + } + sub := h.db.ActiveSubscription(ctx, billing) + if sub == nil { + c.JSON(http.StatusOK, gin.H{"subscription": nil}) + return + } + c.JSON(http.StatusOK, gin.H{"subscription": sub, "plan": h.db.GetSubPlan(ctx, sub.PlanID)}) +} + +// ---- 管理端 ---- + +// AdminSubPlans: GET /api/v1/admin/sub-plans —— 全部套餐(含下架)。 +func (h *Handler) AdminSubPlans(c *gin.Context) { + c.JSON(http.StatusOK, gin.H{"plans": h.db.ListSubPlans(c.Request.Context(), false)}) +} + +// AdminSaveSubPlan: PUT /api/v1/admin/sub-plans —— 新增/改套餐(id 空=新增)。 +// 积分以「积分」为单位收(面向人),服务端转 micro。 +func (h *Handler) AdminSaveSubPlan(c *gin.Context) { + var b struct { + ID string `json:"id"` + Name string `json:"name"` + PriceFen int64 `json:"price_fen"` + DurationDays int `json:"duration_days"` + RefillCredits float64 `json:"refill_credits"` + RefillInterval int `json:"refill_interval_days"` + Active bool `json:"active"` + Sort int `json:"sort"` + } + if err := c.ShouldBindJSON(&b); err != nil || strings.TrimSpace(b.Name) == "" { + c.JSON(http.StatusBadRequest, gin.H{"error": "name 必填"}) + return + } + pl := &store.SubscriptionPlan{ + BaseModel: store.BaseModel{ID: b.ID}, + Name: strings.TrimSpace(b.Name), + PriceFen: b.PriceFen, + DurationDays: b.DurationDays, + RefillCreditsMicro: int64(b.RefillCredits * 1e6), + RefillIntervalDays: b.RefillInterval, + Active: b.Active, + Sort: b.Sort, + } + if err := h.db.SaveSubPlan(c.Request.Context(), pl); err != nil { + c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()}) + return + } + c.JSON(http.StatusOK, gin.H{"id": pl.ID}) +} + +// AdminSubscriptions: GET /api/v1/admin/subscriptions —— 全平台订阅观测。 +func (h *Handler) AdminSubscriptions(c *gin.Context) { + c.JSON(http.StatusOK, gin.H{"subscriptions": h.db.AllSubscriptions(c.Request.Context(), 200)}) +} diff --git a/sundynix-gateway/internal/handler/subscription_tick.go b/sundynix-gateway/internal/handler/subscription_tick.go new file mode 100644 index 0000000..cbe1edf --- /dev/null +++ b/sundynix-gateway/internal/handler/subscription_tick.go @@ -0,0 +1,58 @@ +package handler + +import ( + "context" + "log" + "time" +) + +// 订阅推进定时器:周期扫 active 订阅 → 该发的发、该过期的置过期。 +// 与掉单补偿(payment_reconcile.go)同一范式:定时器只是"兜底触发器", +// 真正的语义与幂等都在 store.TickSubscription 里,两处不会漂移。 +// +// 为什么需要它:订阅是"有效期内每 N 天发一次积分",没有用户请求来驱动这个节拍。 +// 进程停机期间欠下的发放,由 TickSubscription 的补发逻辑一次性补齐。 +const subTickInterval = 10 * time.Minute + +// StartSubscriptionTicker 随进程生命周期运行;多实例并发也安全(发放靠 ledger 唯一索引幂等)。 +func (h *Handler) StartSubscriptionTicker(ctx context.Context) { + go func() { + t := time.NewTicker(subTickInterval) + defer t.Stop() + h.tickSubscriptions(ctx) // 启动即跑一次,把停机期间欠的补上 + for { + select { + case <-ctx.Done(): + return + case <-t.C: + h.tickSubscriptions(ctx) + } + } + }() + log.Printf("[sub] 订阅推进定时器已启动(每 %s 扫一次)", subTickInterval) +} + +func (h *Handler) tickSubscriptions(ctx context.Context) { + subs := h.db.DueSubscriptions(ctx, 200) + if len(subs) == 0 { + return + } + now := time.Now() + var granted, expired int + for i := range subs { + g, exp, err := h.db.TickSubscription(ctx, &subs[i], now) + if err != nil { + // 单条失败不影响其它订阅;下一轮会重试(幂等,不会重复发) + log.Printf("[sub] ⚠️ 推进订阅 %s 失败: %v", subs[i].ID, err) + continue + } + granted += g + if exp { + expired++ + } + } + // 只在有变化时记一行,避免空转刷屏 + if granted+expired > 0 { + log.Printf("[sub] 推进:发放 %d 笔、过期 %d 条(本轮 %d 条订阅)", granted, expired, len(subs)) + } +} diff --git a/sundynix-gateway/internal/router/router.go b/sundynix-gateway/internal/router/router.go index 0fba8b2..ef24fd1 100644 --- a/sundynix-gateway/internal/router/router.go +++ b/sundynix-gateway/internal/router/router.go @@ -13,12 +13,12 @@ import ( "github.com/prometheus/client_golang/prometheus/promhttp" "go.opentelemetry.io/contrib/instrumentation/github.com/gin-gonic/gin/otelgin" - "github.com/sundynix/sundynix-shared/blob" "github.com/sundynix/sundynix-gateway/internal/handler" "github.com/sundynix/sundynix-gateway/internal/middleware" "github.com/sundynix/sundynix-gateway/internal/nats" "github.com/sundynix/sundynix-gateway/internal/store" "github.com/sundynix/sundynix-gateway/internal/webui" + "github.com/sundynix/sundynix-shared/blob" ) // New 构建带有 Guardrail / 限流中间件的 Gin 引擎。 @@ -28,18 +28,19 @@ func New(db *store.Postgres, cache *store.Redis, bus *nats.Bus, blobStore *blob. r.Use(otelgin.Middleware("sundynix-gateway")) // OTel: HTTP server span(链路根 + 提取上游 traceparent) r.Use(middleware.RequestID()) // 生成/透传 X-Request-ID(日志关联) r.Use(middleware.Observe()) // Prometheus 指标 + 结构化访问日志(替代 gin 默认文本日志) - r.Use(cors()) // 桌面端/浏览器跨源访问 - r.Use(middleware.Auth()) // 解析 Bearer JWT,注入已验证 userID(非阻断)——须在限流前,供按用户限流 - r.Use(middleware.TenantContext(db)) // 多租户:注入当前 tenant_id(已登录才解析;须在 Auth 之后) - r.Use(middleware.SpaceContext(db)) // 共享工作区:注入当前 space_id(须在 TenantContext 之后) - r.Use(middleware.RateLimit(cache)) // 已认证按用户限流,否则按 IP(企业网多人共享 IP 不再互相拖累) - r.Use(middleware.Guardrail(db)) // Harness: Input Guardrail(命中落库 guardrail_event) + r.Use(cors()) // 桌面端/浏览器跨源访问 + r.Use(middleware.Auth()) // 解析 Bearer JWT,注入已验证 userID(非阻断)——须在限流前,供按用户限流 + r.Use(middleware.TenantContext(db)) // 多租户:注入当前 tenant_id(已登录才解析;须在 Auth 之后) + r.Use(middleware.SpaceContext(db)) // 共享工作区:注入当前 space_id(须在 TenantContext 之后) + r.Use(middleware.RateLimit(cache)) // 已认证按用户限流,否则按 IP(企业网多人共享 IP 不再互相拖累) + r.Use(middleware.Guardrail(db)) // Harness: Input Guardrail(命中落库 guardrail_event) h := handler.New(db, cache, bus, blobStore) // 微信支付渠道装配:DB 配置优先(admin 控制面热重载)→ env 兜底 → 隐藏。失败不阻断启动。 h.InitWechat(context.Background()) // 掉单补偿定时器:周期扫 pending 微信单确认到账(用户扫完码关页面也能补入账)。 h.StartReconcile(context.Background()) + h.StartSubscriptionTicker(context.Background()) // 订阅按周期发放积分 + 到期置失效 // 可观测性根端点:Prometheus 抓取 + k8s 存活/就绪探针(不挂业务中间件鉴权)。 r.GET("/metrics", gin.WrapH(promhttp.Handler())) @@ -49,73 +50,75 @@ func New(db *store.Postgres, cache *store.Redis, bus *nats.Bus, blobStore *blob. api := r.Group("/api/v1") { // —— 公开:鉴权端点 / 健康 / 按 task_id 寻址的 SSE 与导出(EventSource/下载无法带 Bearer)—— - api.POST("/auth/register", h.Register) // 注册 + 签发 JWT - api.POST("/auth/login", h.Login) // 登录 + 签发 JWT - api.GET("/auth/me", h.Me) // 当前登录用户(无效令牌 → 401) - api.GET("/health", h.Health) // 依赖健康聚合(顶栏五盏灯) - api.GET("/tasks/:id/stream", h.StreamTask) // SSE 回流 Token Stream(task_id 寻址) - api.GET("/tasks/:id/exec", h.StreamExec) // SSE 回流执行轨迹(task_id 寻址) - api.GET("/kb/ingest/:id/stream", h.KbIngestStream) // 入库进度 SSE(job_id 寻址) - api.GET("/reports/:id/export", h.ExportReport) // 按需导出(report_id 寻址) - api.GET("/reports/:id/download", h.ExportReport) // 兼容旧入口(默认 docx) + api.POST("/auth/register", h.Register) // 注册 + 签发 JWT + api.POST("/auth/login", h.Login) // 登录 + 签发 JWT + api.GET("/auth/me", h.Me) // 当前登录用户(无效令牌 → 401) + api.GET("/health", h.Health) // 依赖健康聚合(顶栏五盏灯) + api.GET("/tasks/:id/stream", h.StreamTask) // SSE 回流 Token Stream(task_id 寻址) + api.GET("/tasks/:id/exec", h.StreamExec) // SSE 回流执行轨迹(task_id 寻址) + api.GET("/kb/ingest/:id/stream", h.KbIngestStream) // 入库进度 SSE(job_id 寻址) + api.GET("/reports/:id/export", h.ExportReport) // 按需导出(report_id 寻址) + api.GET("/reports/:id/download", h.ExportReport) // 兼容旧入口(默认 docx) api.POST("/billing/callback/:channel", h.PaymentCallback) // 支付回调(渠道服务器带不了 Bearer;渠道验签是唯一的门) // —— 受保护:owner 作用域业务,必须携带有效 JWT —— p := api.Group("", middleware.RequireAuth()) { p.POST("/tasks", middleware.RequireTenantRole(db, store.RoleMember), h.SubmitTask) // 提交任务(烧租户积分):viewer 只读拦下 - p.GET("/tasks/:id", h.TaskStatus) // 任务生命周期状态(UI 轮询 submitted/running/done/failed/timeout/waiting/rejected) - p.POST("/tasks/:id/approve", middleware.Audit(db), h.ApproveTask) // HITL 人工审批决定(批准/拒绝,审计) - p.GET("/tenants/current", h.TenantCurrent) // 当前租户上下文 + 角色 + 可花余额(多租户) - p.GET("/me/tenants", h.MyTenantsList) // 我所属租户(供切换) - p.POST("/me/tenants", h.CreateMyTenant) // 自助建组织(创建者即 owner,建完切入) - p.POST("/me/tenant", h.SwitchTenant) // 切换当前活跃租户 + p.GET("/tasks/:id", h.TaskStatus) // 任务生命周期状态(UI 轮询 submitted/running/done/failed/timeout/waiting/rejected) + p.POST("/tasks/:id/approve", middleware.Audit(db), h.ApproveTask) // HITL 人工审批决定(批准/拒绝,审计) + p.GET("/tenants/current", h.TenantCurrent) // 当前租户上下文 + 角色 + 可花余额(多租户) + p.GET("/me/tenants", h.MyTenantsList) // 我所属租户(供切换) + p.POST("/me/tenants", h.CreateMyTenant) // 自助建组织(创建者即 owner,建完切入) + p.POST("/me/tenant", h.SwitchTenant) // 切换当前活跃租户 // 租户成员自助管理(薄 Web 面):作用于活跃租户,看 ≥viewer、写 ≥admin + 审计。 p.GET("/tenants/current/members", middleware.RequireTenantRole(db, store.RoleViewer), h.TenantMembers) p.POST("/tenants/current/members", middleware.RequireTenantRole(db, store.RoleAdmin), middleware.Audit(db), h.TenantAddMember) p.PUT("/tenants/current/members/:uid", middleware.RequireTenantRole(db, store.RoleAdmin), middleware.Audit(db), h.TenantSetMemberRole) p.DELETE("/tenants/current/members/:uid", middleware.RequireTenantRole(db, store.RoleAdmin), middleware.Audit(db), h.TenantRemoveMember) - p.GET("/me/usage", h.MyUsage) // 我的用量明细(余额 + 趋势 + 最近消耗) - p.GET("/tasks/:id/eval", h.TaskEval) // 自动化评测结果(综合/质量/忠实度/分级) - p.PUT("/memory", h.SetMemory) // 偏好记忆登记(→ mcp-go memory_upsert) - p.GET("/memory", h.ListMemory) // 列出当前用户偏好(记忆面板) - p.DELETE("/memory", h.DeleteMemory) // 软删一条偏好(?key=) - p.GET("/kb/list", h.KbList) // 当前空间的知识库列表(共享工作区) + p.GET("/me/usage", h.MyUsage) // 我的用量明细(余额 + 趋势 + 最近消耗) + p.GET("/tasks/:id/eval", h.TaskEval) // 自动化评测结果(综合/质量/忠实度/分级) + p.PUT("/memory", h.SetMemory) // 偏好记忆登记(→ mcp-go memory_upsert) + p.GET("/memory", h.ListMemory) // 列出当前用户偏好(记忆面板) + p.DELETE("/memory", h.DeleteMemory) // 软删一条偏好(?key=) + p.GET("/kb/list", h.KbList) // 当前空间的知识库列表(共享工作区) p.POST("/kb/create", middleware.RequireSpaceRole(db, store.RoleMember), h.KbCreate) // 新建知识库:空间只读 viewer 拦下 p.POST("/kb/ingest", middleware.RequireSpaceRole(db, store.RoleMember), h.KbIngest) // 文本入库:viewer 拦下 p.POST("/kb/ingest_file", middleware.RequireSpaceRole(db, store.RoleMember), h.KbIngestFile) // 文件入库:viewer 拦下 - p.POST("/kb/search", h.KbSearch) // 检索台(读,全员) - p.GET("/kb/vault", h.KbVault) // 文库列表 - p.GET("/kb/doc", h.KbDoc) // 取单篇文档 + p.POST("/kb/search", h.KbSearch) // 检索台(读,全员) + p.GET("/kb/vault", h.KbVault) // 文库列表 + p.GET("/kb/doc", h.KbDoc) // 取单篇文档 p.DELETE("/kb/doc", middleware.RequireSpaceRole(db, store.RoleMember), h.KbDeleteDoc) // 级联删文档:viewer 拦下 // Prompt 控制面(平台级配置:建版本 → 激活 → 控制面热下发各服务) - p.GET("/prompts", h.PromptList) // 列出全部版本 + 可配键 - p.POST("/prompts/version", h.PromptCreateVersion) // 建新版本(不自动激活) - p.POST("/prompts/activate", middleware.Audit(db), h.PromptActivate) // 激活某版本 → 广播热更新(审计) - p.POST("/prompts/deactivate", middleware.Audit(db), h.PromptDeactivate) // 撤销激活 → 回退代码默认(热,审计) - p.GET("/kb/links", h.KbLinks) // 某库双链 - p.POST("/kb/note", middleware.RequireSpaceRole(db, store.RoleMember), h.KbSaveNote) // 新建/编辑笔记:viewer 拦下 - p.GET("/kb/graph", h.KbGraph) // 知识图谱三元组 - p.GET("/agents", h.AgentList) // 当前空间的编排列表(共享工作区,含创建人) + p.GET("/prompts", h.PromptList) // 列出全部版本 + 可配键 + p.POST("/prompts/version", h.PromptCreateVersion) // 建新版本(不自动激活) + p.POST("/prompts/activate", middleware.Audit(db), h.PromptActivate) // 激活某版本 → 广播热更新(审计) + p.POST("/prompts/deactivate", middleware.Audit(db), h.PromptDeactivate) // 撤销激活 → 回退代码默认(热,审计) + p.GET("/kb/links", h.KbLinks) // 某库双链 + p.POST("/kb/note", middleware.RequireSpaceRole(db, store.RoleMember), h.KbSaveNote) // 新建/编辑笔记:viewer 拦下 + p.GET("/kb/graph", h.KbGraph) // 知识图谱三元组 + p.GET("/agents", h.AgentList) // 当前空间的编排列表(共享工作区,含创建人) p.POST("/agents", middleware.RequireSpaceRole(db, store.RoleMember), h.AgentSave) // 保存/更新编排:空间只读 viewer 拦下 p.DELETE("/agents", middleware.RequireSpaceRole(db, store.RoleMember), h.AgentDelete) // 删除编排:viewer 拦下(删他人另需 admin,见 handler) // 共享工作区(Space):切换 / 列表 / 建 / 成员管理(增量3) - p.GET("/me/spaces", h.SpacesList) // 活跃租户内我所属的空间(供切换) - p.POST("/me/space", h.SwitchSpace) // 切换活跃空间 - p.GET("/spaces/current", h.SpaceCurrent) // 当前空间上下文 + 我的角色 - p.POST("/spaces", middleware.RequireTenantRole(db, store.RoleMember), h.SpaceCreate) // 建空间:租户只读 viewer 拦下 - p.GET("/spaces/:id/members", h.SpaceMembers) // 空间成员列表 - p.POST("/spaces/:id/members", h.SpaceAddMember) // 拉人进空间(handler 内校验空间 admin) - p.PUT("/spaces/:id/members/:uid", h.SpaceSetMemberRole) // 改空间成员角色 - p.DELETE("/spaces/:id/members/:uid", h.SpaceRemoveMember) // 移除空间成员 - p.POST("/spaces/:id/archive", h.SpaceArchive) // 归档空间 + p.GET("/me/spaces", h.SpacesList) // 活跃租户内我所属的空间(供切换) + p.POST("/me/space", h.SwitchSpace) // 切换活跃空间 + p.GET("/spaces/current", h.SpaceCurrent) // 当前空间上下文 + 我的角色 + p.POST("/spaces", middleware.RequireTenantRole(db, store.RoleMember), h.SpaceCreate) // 建空间:租户只读 viewer 拦下 + p.GET("/spaces/:id/members", h.SpaceMembers) // 空间成员列表 + p.POST("/spaces/:id/members", h.SpaceAddMember) // 拉人进空间(handler 内校验空间 admin) + p.PUT("/spaces/:id/members/:uid", h.SpaceSetMemberRole) // 改空间成员角色 + p.DELETE("/spaces/:id/members/:uid", h.SpaceRemoveMember) // 移除空间成员 + p.POST("/spaces/:id/archive", h.SpaceArchive) // 归档空间 p.POST("/reports", middleware.RequireTenantRole(db, store.RoleMember), h.GenerateReport) // 报告生成(同样烧租户积分):viewer 只读拦下 p.GET("/billing", h.Billing) // 充值(P5.1 兑换码 + P5.2 微信 Native):动钱的 ≥member + 审计;查询全员可看。 p.GET("/billing/packs", h.BillingPacks) + p.GET("/billing/sub-plans", h.BillingSubPlans) // 在售订阅套餐 + p.GET("/billing/subscription", h.MySubscription) // 我的订阅(到期时间/发放次数) p.GET("/billing/orders", h.BillingOrders) p.GET("/billing/orders/:id", h.BillingOrderStatus) // 轮询单态(pending 时顺路主动查单确认) p.POST("/billing/redeem", middleware.RequireTenantRole(db, store.RoleMember), middleware.Audit(db), h.BillingRedeem) @@ -133,9 +136,9 @@ func New(db *store.Postgres, cache *store.Redis, bus *nats.Bus, blobStore *blob. admin.POST("/models/:id/active", h.SetActiveModel) admin.DELETE("/models/:id", h.DeleteModel) admin.POST("/models/test", h.TestModel) - admin.GET("/pricing", h.ListPricing) // 各模型计价(token↔真钱 + 积分权重) - admin.PUT("/pricing", h.SavePricing) // 设置某模型输入/输出单价 + 积分权重 - admin.GET("/billing-config", h.BillingConfig) // 全局计费规则(token→积分汇率 + 硬拦截开关) + admin.GET("/pricing", h.ListPricing) // 各模型计价(token↔真钱 + 积分权重) + admin.PUT("/pricing", h.SavePricing) // 设置某模型输入/输出单价 + 积分权重 + admin.GET("/billing-config", h.BillingConfig) // 全局计费规则(token→积分汇率 + 硬拦截开关) admin.PUT("/billing-config", h.SaveBillingConfig) admin.POST("/credits/grant", h.GrantCredits) // 给租户充值/发放积分 // 支付配置面(P5.1/P5.2):兑换码生成/查看 + 积分包配置 + 微信支付配置(DB 热生效) @@ -145,32 +148,35 @@ func New(db *store.Postgres, cache *store.Redis, bus *nats.Bus, blobStore *blob. admin.PUT("/packs", h.AdminSavePack) admin.GET("/payment/wechat", h.AdminGetWechatPay) admin.PUT("/payment/wechat", h.AdminSaveWechatPay) - admin.GET("/orders", h.AdminOrders) // 全平台充值订单流 + 状态计数 - admin.GET("/orders/reconcile", h.AdminReconcile) // 日终对账:paid 单 ↔ 账本 grant + admin.GET("/sub-plans", h.AdminSubPlans) // 订阅套餐(含下架) + admin.PUT("/sub-plans", h.AdminSaveSubPlan) // 配价格/时长/发放节奏 + admin.GET("/subscriptions", h.AdminSubscriptions) // 全平台订阅观测 + admin.GET("/orders", h.AdminOrders) // 全平台充值订单流 + 状态计数 + admin.GET("/orders/reconcile", h.AdminReconcile) // 日终对账:paid 单 ↔ 账本 grant admin.POST("/orders/:id/refund", h.AdminRefundOrder) // 人工退款:置 refunded + adjust 负分录 + 回退余额(审计) // 多租户成员管理(平台运维口径) - admin.GET("/tenants", h.AdminTenants) // 租户目录(成员数+余额) - admin.POST("/tenants", h.AdminCreateTenant) // 新建租户(可选指定 owner) + admin.GET("/tenants", h.AdminTenants) // 租户目录(成员数+余额) + admin.POST("/tenants", h.AdminCreateTenant) // 新建租户(可选指定 owner) admin.PUT("/tenants/:id/shared-billing", h.AdminSetSharedBilling) // 共享计费开关 admin.PUT("/tenants/:id/plan", h.AdminSetTenantPlan) // 方案等级 admin.PUT("/tenants/:id/status", h.AdminSetTenantStatus) // 租户状态 - admin.GET("/tenants/:id/members", h.AdminMembers) // 成员列表 - admin.POST("/tenants/:id/members", h.AdminAddMember) // 按邮箱加成员 - admin.PUT("/tenants/:id/members/:uid", h.AdminSetMemberRole) // 改角色 + admin.GET("/tenants/:id/members", h.AdminMembers) // 成员列表 + admin.POST("/tenants/:id/members", h.AdminAddMember) // 按邮箱加成员 + admin.PUT("/tenants/:id/members/:uid", h.AdminSetMemberRole) // 改角色 admin.DELETE("/tenants/:id/members/:uid", h.AdminRemoveMember) // 移除成员 - admin.GET("/status", h.AdminStatus) // 服务状态:基建/服务探活 + MCP 工具注册 - admin.GET("/overview", h.AdminOverview) // 系统级聚合:全平台用户/任务/评测/模型态/提示词态/健康 - admin.GET("/tasks", h.AdminTasks) // 全平台任务/运行观测(状态/租户筛 + HITL 待审批) - admin.GET("/tasks/:id", h.AdminTaskDetail) // 任务下钻:DSL/输出/轨迹/评测(跨租户,走 WithoutTenant) - admin.GET("/spaces", h.AdminSpaces) // 全平台空间观测(跨租户 Space + 成员数) - admin.GET("/usage", h.AdminUsage) // 用量/积分/成本:全平台按天趋势 + 租户排行 / 单租户余额 - admin.GET("/evals", h.AdminEvals) // 自动评测观测:质量趋势 + 计数 + 错题本(真数据) - admin.GET("/datasources", h.AdminDatasources) // 数据源清单:全平台知识库 + 文档数(真数据) - admin.POST("/kb/search", h.AdminKbSearch) // 检索试验台:按完整作用域键跨租户检索(支持单路 mode 对比) - admin.POST("/migrate-kb-storage", h.MigrateKBStorage) // 增量3:存量 KB 三库 owner/kb→space/kb 重灌(一次性) - admin.GET("/audit", h.AuditList) // 敏感操作审计流(倒序,翻页) - admin.GET("/guardrail-events", h.GuardrailEvents) // 护栏命中安全事件流(倒序,翻页) + admin.GET("/status", h.AdminStatus) // 服务状态:基建/服务探活 + MCP 工具注册 + admin.GET("/overview", h.AdminOverview) // 系统级聚合:全平台用户/任务/评测/模型态/提示词态/健康 + admin.GET("/tasks", h.AdminTasks) // 全平台任务/运行观测(状态/租户筛 + HITL 待审批) + admin.GET("/tasks/:id", h.AdminTaskDetail) // 任务下钻:DSL/输出/轨迹/评测(跨租户,走 WithoutTenant) + admin.GET("/spaces", h.AdminSpaces) // 全平台空间观测(跨租户 Space + 成员数) + admin.GET("/usage", h.AdminUsage) // 用量/积分/成本:全平台按天趋势 + 租户排行 / 单租户余额 + admin.GET("/evals", h.AdminEvals) // 自动评测观测:质量趋势 + 计数 + 错题本(真数据) + admin.GET("/datasources", h.AdminDatasources) // 数据源清单:全平台知识库 + 文档数(真数据) + admin.POST("/kb/search", h.AdminKbSearch) // 检索试验台:按完整作用域键跨租户检索(支持单路 mode 对比) + admin.POST("/migrate-kb-storage", h.MigrateKBStorage) // 增量3:存量 KB 三库 owner/kb→space/kb 重灌(一次性) + admin.GET("/audit", h.AuditList) // 敏感操作审计流(倒序,翻页) + admin.GET("/guardrail-events", h.GuardrailEvents) // 护栏命中安全事件流(倒序,翻页) } } diff --git a/sundynix-gateway/internal/store/payment.go b/sundynix-gateway/internal/store/payment.go index f373f93..e3054e3 100644 --- a/sundynix-gateway/internal/store/payment.go +++ b/sundynix-gateway/internal/store/payment.go @@ -4,6 +4,7 @@ import ( "context" "crypto/rand" "errors" + "log" "strings" "time" @@ -19,6 +20,12 @@ import ( // - CreditPack / RedeemCode 是平台级配置与凭证,不属于任何租户。 // 订单状态机:pending → paid | failed | expired;paid →(人工)refunded。 +const ( + // 订单类型:一次性积分包 vs 订阅周期 + OrderKindPack = "pack" + OrderKindSub = "sub" +) + const ( OrderPending = "pending" OrderPaid = "paid" @@ -49,9 +56,11 @@ func (CreditPack) TableName() string { return "sundynix_credit_pack" } // 兑换码入账也写一行(channel=redeem、amount_fen=0、即时 paid),全部充值一个查法。 type PaymentOrder struct { BaseModel - TenantID string `gorm:"size:64;index" json:"tenant_id"` // 计费租户(下单时解析并锁定) - UserID string `gorm:"size:64;index" json:"user_id"` // 操作人(审计) - PackID string `gorm:"size:24" json:"pack_id"` // redeem 渠道为空 + TenantID string `gorm:"size:64;index" json:"tenant_id"` // 计费租户(下单时解析并锁定) + UserID string `gorm:"size:64;index" json:"user_id"` // 操作人(审计) + PackID string `gorm:"size:24" json:"pack_id"` // redeem 渠道为空 + Kind string `gorm:"size:16;default:pack" json:"kind"` // pack=积分包(一次性) / sub=订阅 + PlanID string `gorm:"size:24" json:"plan_id"` // kind=sub 时的订阅套餐 AmountFen int64 `gorm:"column:amount_fen" json:"amount_fen"` CreditsMicro int64 `gorm:"column:credits_micro" json:"credits_micro"` Channel string `gorm:"size:16;index" json:"channel"` @@ -219,6 +228,7 @@ func (p *Postgres) MarkOrderPaid(ctx context.Context, orderID, channelTxn string return false, errStoreDisabled } changed := false + var paid *PaymentOrder err := p.db.WithContext(WithoutTenant(ctx)).Transaction(func(tx *gorm.DB) error { now := time.Now() res := tx.Model(&PaymentOrder{}). @@ -234,19 +244,35 @@ func (p *Postgres) MarkOrderPaid(ctx context.Context, orderID, channelTxn string if err := tx.First(&o, "id = ?", orderID).Error; err != nil { return err } - if err := tx.Create(&CreditLedger{ - TenantID: o.TenantID, Kind: LedgerGrant, CreditsMicro: o.CreditsMicro, Ref: o.ID, Memo: "充值 " + o.Channel, - }).Error; err != nil { - return err - } - if err := tx.Model(&Tenant{}).Where("id = ?", o.TenantID). - UpdateColumn("credit_balance_micro", gorm.Expr("credit_balance_micro + ?", o.CreditsMicro)).Error; err != nil { - return err + paid = &o + // 订阅单自身不带积分(积分由订阅按周期发放),跳过零额分录避免账本噪声。 + if o.CreditsMicro != 0 { + if err := tx.Create(&CreditLedger{ + TenantID: o.TenantID, Kind: LedgerGrant, CreditsMicro: o.CreditsMicro, Ref: o.ID, Memo: "充值 " + o.Channel, + }).Error; err != nil { + return err + } + if err := tx.Model(&Tenant{}).Where("id = ?", o.TenantID). + UpdateColumn("credit_balance_micro", gorm.Expr("credit_balance_micro + ?", o.CreditsMicro)).Error; err != nil { + return err + } } changed = true return nil }) - return changed, err + if err != nil { + return changed, err + } + // 订阅开通放在这里、而不是各调用方:回调与掉单补偿两条路都经过 MarkOrderPaid, + // 放在这一处才没人能漏掉。ActivateSubscription 按 orderID 幂等,重复调用无害。 + if paid != nil && paid.Kind == OrderKindSub && paid.PlanID != "" { + if _, aerr := p.ActivateSubscription(ctx, paid.TenantID, paid.PlanID, paid.ID); aerr != nil { + // 钱已收、订单已 paid:这里失败**不能**回滚订单(否则用户付了钱订单还回到 pending, + // 补偿定时器会再入账一次)。出声即可,可人工或下次回调补开通。 + log.Printf("[payment] ⚠️ 订单 %s 已入账但订阅开通失败: %v(需人工补开通)", paid.ID, aerr) + } + } + return changed, nil } // RefundOrder 人工退款(PAYMENT_DESIGN §5:admin 发起 → 订单置 refunded + 记 adjust 负分录 + 回退余额)。 diff --git a/sundynix-gateway/internal/store/pgsql.go b/sundynix-gateway/internal/store/pgsql.go index cf4dac6..9687d4e 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{}, &Space{}, &SpaceMember{}, &UsageEvent{}, &CreditLedger{}, &UsageRollup{}, &Setting{}, &CreditPack{}, &PaymentOrder{}, &RedeemCode{}); err != nil { + if err := db.AutoMigrate(&User{}, &Task{}, &Eval{}, &LLMModel{}, &KB{}, &Doc{}, &Agent{}, &DocLink{}, &Pricing{}, &Prompt{}, &AuditLog{}, &GuardrailEvent{}, &Tenant{}, &TenantMember{}, &Space{}, &SpaceMember{}, &UsageEvent{}, &CreditLedger{}, &UsageRollup{}, &Setting{}, &CreditPack{}, &PaymentOrder{}, &RedeemCode{}, &SubscriptionPlan{}, &Subscription{}); err != nil { log.Printf("[store] postgres AutoMigrate 失败,降级运行: %v", err) return &Postgres{} } diff --git a/sundynix-gateway/internal/store/subscription.go b/sundynix-gateway/internal/store/subscription.go new file mode 100644 index 0000000..a231c15 --- /dev/null +++ b/sundynix-gateway/internal/store/subscription.go @@ -0,0 +1,299 @@ +package store + +import ( + "context" + "errors" + "fmt" + "strings" + "time" + + "gorm.io/gorm" +) + +// 订阅(手动购买制)。刻意**不做自动续费**:微信 Native 扫码支付没有代扣能力, +// 真自动续费要走「委托代扣」——另一套产品与资质。这里的语义是: +// +// 用户扫码买一个订阅周期 → 有效期内每 N 天发一次积分 → 到期即失效,要续得再买一次。 +// +// N(间隔天数)与每次发多少积分都在套餐里配,后台可改。 +// +// **发放语义是「累加」而非「重置」**:每次刷新写一条 grant 分录、余额累加, +// 用不完的会留着,也绝不会清掉用户自己充值的积分。重置型(月度配额清零)会让 +// 「余额 = SUM(ledger)」这条对账不变量变复杂,且有误删用户已付费积分的风险, +// 故不采用。 + +// SubscriptionPlan 订阅套餐(价格 / 时长 / 发放节奏,全部后台可配)。 +type SubscriptionPlan struct { + BaseModel + Name string `gorm:"size:64" json:"name"` + PriceFen int64 `gorm:"column:price_fen" json:"price_fen"` // 售价(分) + DurationDays int `gorm:"column:duration_days" json:"duration_days"` // 一个订阅周期多少天 + RefillCreditsMicro int64 `gorm:"column:refill_credits_micro" json:"refill_credits_micro"` // 每次发放的积分 ×10⁻⁶ + RefillIntervalDays int `gorm:"column:refill_interval_days" json:"refill_interval_days"` // 每几天发一次 + Active bool `json:"active"` + Sort int `json:"sort"` +} + +func (SubscriptionPlan) TableName() string { return "sundynix_sub_plan" } + +// Subscription 一次已购订阅。到期即 expired,不自动续。 +type Subscription struct { + BaseModel + TenantID string `gorm:"size:64;index" json:"tenant_id"` + PlanID string `gorm:"size:24;index" json:"plan_id"` + OrderID string `gorm:"size:24" json:"order_id"` // 来源支付订单(人工发放为空) + Status string `gorm:"size:16;index" json:"status"` + StartedAt time.Time `json:"started_at"` + ExpiresAt time.Time `gorm:"index" json:"expires_at"` + RefillSeq int `gorm:"column:refill_seq" json:"refill_seq"` // 已发放次数;兼作幂等序号 + LastRefillAt *time.Time `json:"last_refill_at"` +} + +func (Subscription) TableName() string { return "sundynix_subscription" } +func (Subscription) isTenantScoped() {} // 用户面只看得到自己租户的订阅;系统级扫描须 WithoutTenant + +const ( + SubActive = "active" + SubExpired = "expired" +) + +// refillRef 是一次发放的幂等键,落到 credit_ledger.ref。 +// credit_ledger 上 (kind='grant', ref) 的唯一索引是最终闸门:定时器重跑、多实例并发、 +// 手动补发,撞到同一序号都只会成功一次。 +func refillRef(subID string, seq int) string { return fmt.Sprintf("sub:%s:%d", subID, seq) } + +// isDupKey 判断是否唯一索引冲突(= 这一笔已经发过了,幂等成功而非失败)。 +func isDupKey(err error) bool { + if err == nil { + return false + } + s := strings.ToLower(err.Error()) + return strings.Contains(s, "duplicate key") || strings.Contains(s, "unique constraint") || + strings.Contains(s, "unique violation") || strings.Contains(s, "constraint failed") +} + +// ---- 套餐配置 ---- + +func (p *Postgres) ListSubPlans(ctx context.Context, onlyActive bool) []SubscriptionPlan { + if p.db == nil { + return nil + } + q := p.db.WithContext(WithoutTenant(ctx)).Order("sort asc, price_fen asc") + if onlyActive { + q = q.Where("active = ?", true) + } + var out []SubscriptionPlan + q.Find(&out) + return out +} + +func (p *Postgres) GetSubPlan(ctx context.Context, id string) *SubscriptionPlan { + if p.db == nil || id == "" { + return nil + } + var pl SubscriptionPlan + if err := p.db.WithContext(WithoutTenant(ctx)).First(&pl, "id = ?", id).Error; err != nil { + return nil + } + return &pl +} + +// SaveSubPlan 新增或更新套餐(id 空=新增)。 +func (p *Postgres) SaveSubPlan(ctx context.Context, pl *SubscriptionPlan) error { + if p.db == nil { + return errStoreDisabled + } + if pl.DurationDays <= 0 { + return errors.New("订阅时长必须大于 0 天") + } + if pl.RefillIntervalDays <= 0 { + return errors.New("发放间隔必须大于 0 天") + } + // 间隔比时长还长 = 一个周期内一次都发不到第二回,多半是配错了,直接拦下。 + if pl.RefillIntervalDays > pl.DurationDays { + return errors.New("发放间隔不能大于订阅时长") + } + return p.db.WithContext(WithoutTenant(ctx)).Save(pl).Error +} + +// ---- 订阅生命周期 ---- + +// ActivateSubscription 支付成功后开通/续期,并立即发放第一笔积分。 +// 已有生效中的订阅则**顺延**到期时间(而不是新建一条),避免同租户多条 active 互相打架。 +// 幂等:同一 orderID 只会开通一次。 +func (p *Postgres) ActivateSubscription(ctx context.Context, tenantID, planID, orderID string) (*Subscription, error) { + if p.db == nil { + return nil, errStoreDisabled + } + pl := p.GetSubPlan(ctx, planID) + if pl == nil { + return nil, errors.New("订阅套餐不存在") + } + ctx = WithoutTenant(ctx) // 系统级:为目标租户开通,调用方可能是 admin + now := time.Now() + + var sub *Subscription + err := p.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error { + // 幂等闸:同一订单已开通过就直接返回,不重复延期 + if orderID != "" { + var exist Subscription + if err := tx.First(&exist, "order_id = ?", orderID).Error; err == nil { + sub = &exist + return nil + } + } + var cur Subscription + err := tx.Where("tenant_id = ? AND status = ?", tenantID, SubActive). + Order("expires_at desc").First(&cur).Error + switch { + case err == nil: // 续期:在原到期时间上顺延 + cur.ExpiresAt = cur.ExpiresAt.AddDate(0, 0, pl.DurationDays) + cur.PlanID, cur.OrderID = pl.ID, orderID + if err := tx.Save(&cur).Error; err != nil { + return err + } + sub = &cur + return nil + case errors.Is(err, gorm.ErrRecordNotFound): + s := &Subscription{ + TenantID: tenantID, PlanID: pl.ID, OrderID: orderID, Status: SubActive, + StartedAt: now, ExpiresAt: now.AddDate(0, 0, pl.DurationDays), + } + if err := tx.Create(s).Error; err != nil { + return err + } + sub = s + return nil + default: + return err + } + }) + if err != nil { + return nil, err + } + // 首笔发放走与定时器**同一套排期判断**(而不是无条件发一笔):否则同一订单重复 + // 开通时(回调重推、查单与回调赛跑)会各发一笔,序号递增绕过幂等索引 —— 白送钱。 + // 排期判断天然幂等:seq 已发过则下一笔的到期时间在未来,不会发。 + // 放在事务外:发放失败不该导致"已付款却没开通",定时器下一轮会补上。 + if _, _, err := p.TickSubscription(ctx, sub, now); err != nil { + return sub, nil // 开通已成功,发放失败交给定时器补 + } + return sub, nil +} + +// refillOnce 发放一次积分并推进序号。返回 granted=false 表示这一笔已发过(幂等)。 +func (p *Postgres) refillOnce(ctx context.Context, sub *Subscription, pl *SubscriptionPlan, now time.Time) (bool, error) { + if pl.RefillCreditsMicro <= 0 { + return false, nil + } + seq := sub.RefillSeq + 1 + err := p.GrantCredits(ctx, sub.TenantID, LedgerGrant, pl.RefillCreditsMicro, + refillRef(sub.ID, seq), "订阅发放 "+pl.Name) + if err != nil { + if isDupKey(err) { + return false, nil // 已发过:幂等成功 + } + return false, err + } + sub.RefillSeq = seq + sub.LastRefillAt = &now + return true, p.db.WithContext(WithoutTenant(ctx)).Model(&Subscription{}). + Where("id = ?", sub.ID). + Updates(map[string]any{"refill_seq": seq, "last_refill_at": now}).Error +} + +// DueSubscriptions 取到期需处理的订阅(系统级,跨租户)。 +func (p *Postgres) DueSubscriptions(ctx context.Context, limit int) []Subscription { + if p.db == nil { + return nil + } + if limit <= 0 || limit > 500 { + limit = 200 + } + var out []Subscription + p.db.WithContext(WithoutTenant(ctx)). + Where("status = ?", SubActive).Order("expires_at asc").Limit(limit).Find(&out) + return out +} + +// TickSubscription 推进一条订阅:先看是否到期,再看是否该发放。 +// 返回 (发放笔数, 是否刚过期)。定时器与手动触发共用这一份逻辑,避免两处行为漂移。 +func (p *Postgres) TickSubscription(ctx context.Context, sub *Subscription, now time.Time) (int, bool, error) { + pl := p.GetSubPlan(ctx, sub.PlanID) + if pl == nil { + return 0, false, errors.New("订阅套餐已不存在: " + sub.PlanID) + } + granted := 0 + // 补发:进程停机/定时器漏跑期间欠下的次数要一次性补齐,而不是只发最近一次。 + // 上限用到期时间卡住——过期之后的周期一律不发。 + for { + due := sub.StartedAt.AddDate(0, 0, pl.RefillIntervalDays*(sub.RefillSeq)) + if due.After(now) || !due.Before(sub.ExpiresAt) { + break + } + ok, err := p.refillOnce(ctx, sub, pl, now) + if err != nil { + return granted, false, err + } + if ok { + granted++ + } + if sub.RefillSeq > 1000 { // 防呆:配置异常(间隔 0)时不至于死循环 + break + } + } + if now.After(sub.ExpiresAt) { + if err := p.db.WithContext(WithoutTenant(ctx)).Model(&Subscription{}). + Where("id = ? AND status = ?", sub.ID, SubActive). + Update("status", SubExpired).Error; err != nil { + return granted, false, err + } + return granted, true, nil + } + return granted, false, nil +} + +// ActiveSubscription 取某租户当前生效的订阅(用户面:账单页展示到期时间)。 +func (p *Postgres) ActiveSubscription(ctx context.Context, tenantID string) *Subscription { + if p.db == nil || tenantID == "" { + return nil + } + var s Subscription + if err := p.db.WithContext(WithoutTenant(ctx)). + Where("tenant_id = ? AND status = ?", tenantID, SubActive). + Order("expires_at desc").First(&s).Error; err != nil { + return nil + } + return &s +} + +// AdminSubRow 是管理端订阅观测一行:订阅 + 租户名 + 套餐名。 +type AdminSubRow struct { + ID string `json:"id"` + TenantID string `json:"tenant_id"` + TenantName string `json:"tenant_name"` + PlanName string `json:"plan_name"` + Status string `json:"status"` + StartedAt time.Time `json:"started_at"` + ExpiresAt time.Time `json:"expires_at"` + RefillSeq int `json:"refill_seq"` +} + +// AllSubscriptions 全平台订阅(管理端观测,跨租户)。 +func (p *Postgres) AllSubscriptions(ctx context.Context, limit int) []AdminSubRow { + if p.db == nil { + return nil + } + if limit <= 0 || limit > 500 { + limit = 200 + } + var out []AdminSubRow + p.db.WithContext(WithoutTenant(ctx)).Table("sundynix_subscription as s"). + Select("s.id, s.tenant_id, s.status, s.started_at, s.expires_at, s.refill_seq, " + + "coalesce(t.name,'') as tenant_name, coalesce(pl.name,'') as plan_name"). + Joins("left join sundynix_tenant t on t.id = s.tenant_id"). + Joins("left join sundynix_sub_plan pl on pl.id = s.plan_id"). + Where("s.deleted_at is null"). + Order("s.expires_at desc").Limit(limit).Scan(&out) + return out +} diff --git a/sundynix-gateway/internal/store/subscription_test.go b/sundynix-gateway/internal/store/subscription_test.go new file mode 100644 index 0000000..4a2fa59 --- /dev/null +++ b/sundynix-gateway/internal/store/subscription_test.go @@ -0,0 +1,248 @@ +package store + +import ( + "context" + "testing" + "time" +) + +// 订阅是涉及钱的路径:发多了是白送,发少了是欠付费用户的。这组测试钉死三件事—— +// 幂等(重跑不重复发)、补发(漏跑要补齐)、到期边界(过期后一分不发)。 + +func seedPlan(t *testing.T, p *Postgres, durationDays, intervalDays int, credits int64) *SubscriptionPlan { + t.Helper() + pl := &SubscriptionPlan{ + Name: "测试套餐", PriceFen: 9900, DurationDays: durationDays, + RefillCreditsMicro: credits, RefillIntervalDays: intervalDays, Active: true, + } + if err := p.SaveSubPlan(context.Background(), pl); err != nil { + t.Fatalf("建套餐失败: %v", err) + } + return pl +} + +func balance(t *testing.T, p *Postgres, tenantID string) int64 { + t.Helper() + return p.TenantBalance(WithoutTenant(context.Background()), tenantID) +} + +func TestSubscription_ActivateGrantsFirstRefill(t *testing.T) { + p := newTestStore(t) + seedTenant(t, p, "t1") + pl := seedPlan(t, p, 30, 7, 100_000_000) + + sub, err := p.ActivateSubscription(context.Background(), "t1", pl.ID, "order-1") + if err != nil { + t.Fatalf("开通失败: %v", err) + } + if sub.Status != SubActive { + t.Fatalf("应为 active,得 %q", sub.Status) + } + if got := balance(t, p, "t1"); got != 100_000_000 { + t.Fatalf("开通即应发第一笔,余额应 100e6,得 %d", got) + } + assertBalanceInvariant(t, p, "t1") +} + +// 同一订单重复开通(回调重推 / 查单与回调赛跑)不能重复延期、不能重复发放。 +func TestSubscription_ActivateIsIdempotentPerOrder(t *testing.T) { + p := newTestStore(t) + seedTenant(t, p, "t1") + pl := seedPlan(t, p, 30, 7, 100_000_000) + ctx := context.Background() + + s1, err := p.ActivateSubscription(ctx, "t1", pl.ID, "order-1") + if err != nil { + t.Fatal(err) + } + s2, err := p.ActivateSubscription(ctx, "t1", pl.ID, "order-1") + if err != nil { + t.Fatal(err) + } + if !s1.ExpiresAt.Equal(s2.ExpiresAt) { + t.Fatalf("同一订单重复开通不该延期:%v → %v", s1.ExpiresAt, s2.ExpiresAt) + } + if got := balance(t, p, "t1"); got != 100_000_000 { + t.Fatalf("重复开通不该重复发放,余额应仍为 100e6,得 %d", got) + } + assertBalanceInvariant(t, p, "t1") +} + +// 续订(不同订单)应在原到期时间上顺延,而不是新建第二条 active。 +func TestSubscription_RenewExtendsInsteadOfDuplicating(t *testing.T) { + p := newTestStore(t) + seedTenant(t, p, "t1") + pl := seedPlan(t, p, 30, 7, 100_000_000) + ctx := context.Background() + + s1, _ := p.ActivateSubscription(ctx, "t1", pl.ID, "order-1") + s2, err := p.ActivateSubscription(ctx, "t1", pl.ID, "order-2") + if err != nil { + t.Fatal(err) + } + if s1.ID != s2.ID { + t.Fatalf("续订应复用同一条订阅,得两条:%s / %s", s1.ID, s2.ID) + } + want := s1.ExpiresAt.AddDate(0, 0, 30) + if !s2.ExpiresAt.Equal(want) { + t.Fatalf("续订应顺延 30 天:want %v got %v", want, s2.ExpiresAt) + } + var n int64 + p.db.WithContext(WithoutTenant(ctx)).Model(&Subscription{}). + Where("tenant_id = ? AND status = ?", "t1", SubActive).Count(&n) + if n != 1 { + t.Fatalf("同租户不应出现多条 active 订阅,得 %d 条", n) + } +} + +// 定时器漏跑(进程停机数周)后要把欠下的次数一次补齐,而不是只补最近一次。 +func TestSubscription_TickBackfillsMissedRefills(t *testing.T) { + p := newTestStore(t) + seedTenant(t, p, "t1") + pl := seedPlan(t, p, 30, 7, 100_000_000) // 30 天订阅,每 7 天发一次 + ctx := context.Background() + + sub, _ := p.ActivateSubscription(ctx, "t1", pl.ID, "order-1") // 已发第 1 笔 + // 快进 22 天:第 7/14/21 天各应发一次,共补 3 笔 + now := sub.StartedAt.AddDate(0, 0, 22) + granted, expired, err := p.TickSubscription(ctx, sub, now) + if err != nil { + t.Fatal(err) + } + if expired { + t.Fatal("22 天时不该过期(周期 30 天)") + } + if granted != 3 { + t.Fatalf("应补发 3 笔(第 7/14/21 天),得 %d", granted) + } + if got := balance(t, p, "t1"); got != 400_000_000 { + t.Fatalf("首笔 + 补 3 笔 = 400e6,得 %d", got) + } + assertBalanceInvariant(t, p, "t1") +} + +// 同一时刻重复 tick(多实例并发 / 定时器重叠)不能重复发放。 +func TestSubscription_TickIsIdempotent(t *testing.T) { + p := newTestStore(t) + seedTenant(t, p, "t1") + pl := seedPlan(t, p, 30, 7, 100_000_000) + ctx := context.Background() + + sub, _ := p.ActivateSubscription(ctx, "t1", pl.ID, "order-1") + now := sub.StartedAt.AddDate(0, 0, 8) + + if _, _, err := p.TickSubscription(ctx, sub, now); err != nil { + t.Fatal(err) + } + before := balance(t, p, "t1") + // 再 tick 两次,余额不能变 + for i := 0; i < 2; i++ { + if _, _, err := p.TickSubscription(ctx, sub, now); err != nil { + t.Fatal(err) + } + } + if after := balance(t, p, "t1"); after != before { + t.Fatalf("重复 tick 不该重复发放:%d → %d", before, after) + } + assertBalanceInvariant(t, p, "t1") +} + +// 过期后一分不发,且状态置 expired(到期即失效,无自动续费)。 +func TestSubscription_ExpiresAndStopsGranting(t *testing.T) { + p := newTestStore(t) + seedTenant(t, p, "t1") + pl := seedPlan(t, p, 14, 7, 100_000_000) // 14 天,7 天一发 → 最多发第 1、第 7 天两笔 + ctx := context.Background() + + sub, _ := p.ActivateSubscription(ctx, "t1", pl.ID, "order-1") + now := sub.StartedAt.AddDate(0, 0, 100) // 远超到期 + granted, expired, err := p.TickSubscription(ctx, sub, now) + if err != nil { + t.Fatal(err) + } + if !expired { + t.Fatal("早该过期了") + } + // 到期时间点(第 14 天)之后的周期不发:第 7 天那笔算,第 14 天正好等于到期不算 + if granted != 1 { + t.Fatalf("过期前只应补第 7 天那一笔,得 %d 笔", granted) + } + if got := balance(t, p, "t1"); got != 200_000_000 { + t.Fatalf("首笔 + 第 7 天 = 200e6,得 %d", got) + } + + var s Subscription + p.db.WithContext(WithoutTenant(ctx)).First(&s, "id = ?", sub.ID) + if s.Status != SubExpired { + t.Fatalf("状态应为 expired,得 %q", s.Status) + } + assertBalanceInvariant(t, p, "t1") +} + +// 配置校验:间隔比时长长 = 一个周期只发得到首笔,多半是配错了,直接拦。 +func TestSubPlan_RejectsBadConfig(t *testing.T) { + p := newTestStore(t) + ctx := context.Background() + for _, tc := range []struct { + name string + duration, ivl int + wantErrSubstring string + }{ + {"时长为 0", 0, 7, "订阅时长"}, + {"间隔为 0", 30, 0, "发放间隔"}, + {"间隔大于时长", 7, 30, "不能大于"}, + } { + t.Run(tc.name, func(t *testing.T) { + err := p.SaveSubPlan(ctx, &SubscriptionPlan{ + Name: "x", DurationDays: tc.duration, RefillIntervalDays: tc.ivl, RefillCreditsMicro: 1, + }) + if err == nil { + t.Fatal("应被拒绝") + } + }) + } +} + +var _ = time.Now + +// 订阅单经支付回调入账后必须真的开通订阅。这是"钱收了但订阅没生效"的高危点, +// 而且开通逻辑刻意放在 MarkOrderPaid 里(回调与掉单补偿两条路共用),这里一并钉死。 +func TestSubscription_ActivatedByOrderPayment(t *testing.T) { + p := newTestStore(t) + seedTenant(t, p, "t1") + pl := seedPlan(t, p, 30, 7, 100_000_000) + ctx := context.Background() + + o := &PaymentOrder{ + TenantID: "t1", UserID: "u1", Kind: OrderKindSub, PlanID: pl.ID, + AmountFen: pl.PriceFen, CreditsMicro: 0, // 订阅单自身不带积分 + Channel: ChannelWechat, Status: OrderPending, + } + if err := p.CreateOrder(ctx, o); err != nil { + t.Fatal(err) + } + changed, err := p.MarkOrderPaid(ctx, o.ID, "txn-1") + if err != nil || !changed { + t.Fatalf("入账应成功: changed=%v err=%v", changed, err) + } + + sub := p.ActiveSubscription(ctx, "t1") + if sub == nil { + t.Fatal("付款后应已开通订阅") + } + if sub.OrderID != o.ID { + t.Fatalf("订阅应关联来源订单 %s,得 %s", o.ID, sub.OrderID) + } + if got := balance(t, p, "t1"); got != 100_000_000 { + t.Fatalf("开通即发首笔,余额应 100e6,得 %d", got) + } + assertBalanceInvariant(t, p, "t1") + + // 回调重复推送:不能重复开通、不能重复发放 + if _, err := p.MarkOrderPaid(ctx, o.ID, "txn-1"); err != nil { + t.Fatal(err) + } + if got := balance(t, p, "t1"); got != 100_000_000 { + t.Fatalf("重复回调不该重复发放,得 %d", got) + } +} diff --git a/sundynix-gateway/internal/store/testdb_test.go b/sundynix-gateway/internal/store/testdb_test.go index 1666eca..62c63c5 100644 --- a/sundynix-gateway/internal/store/testdb_test.go +++ b/sundynix-gateway/internal/store/testdb_test.go @@ -34,7 +34,7 @@ func newTestStore(t *testing.T) *Postgres { if err := db.AutoMigrate( &User{}, &Tenant{}, &TenantMember{}, &CreditLedger{}, &PaymentOrder{}, &RedeemCode{}, &CreditPack{}, &UsageEvent{}, &UsageRollup{}, &Setting{}, &Pricing{}, &LLMModel{}, - &AuditLog{}, &Task{}, &Eval{}, + &AuditLog{}, &Task{}, &Eval{}, &SubscriptionPlan{}, &Subscription{}, &KB{}, // 租户作用域模型,验证隔离插件 ); err != nil { t.Fatalf("AutoMigrate 失败: %v", err)