From 700845d64ad071d316a171ea5e8cab4459345f01 Mon Sep 17 00:00:00 2001 From: Blizzard Date: Wed, 24 Jun 2026 16:14:19 +0800 Subject: [PATCH] =?UTF-8?q?feat(prod):=20DB=20=E8=BF=9E=E6=8E=A5=E6=B1=A0?= =?UTF-8?q?=E4=B8=8A=E9=99=90=20+=20LLM=20=E5=A4=B1=E8=B4=A5=E6=9A=B4?= =?UTF-8?q?=E9=9C=B2=E4=B8=BA=20failed=EF=BC=88=E7=94=9F=E4=BA=A7=E7=BA=A7?= =?UTF-8?q?=E5=B9=B6=E5=8F=91=E6=94=B6=E5=B0=BE=EF=BC=89?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 为高并发生产做的三项收尾(配合已有的任务/工具并发消费): 1. DB 连接池上限(pgsql.go / memory/store.go):SetMaxOpenConns(默认 25, DB_MAX_OPEN_CONNS 可调)+ MaxIdleConns 5 + ConnMaxLifetime 1h。 防高并发无限开连接打爆 PG(max_connections 默认 100)。 2. LLM 失败暴露为 failed(graph.go):board.fatalErr —— agent 模型调用出错即上抛, runGraph 中止后续节点并返回错误 → Handle 判 failed(带原因),不再静默 done-空。 可观测/可告警,生产排障必需。 压测验证(dispatcher 并发=50, 池=25, deepseek-v4-pro 推理): - 平台同一秒并发收下 40 任务,全程零 DB/连接错误,平台开销≈0(裸 LLM 1.8s vs 平台 P50 1.7s)。 - 并发 10 健康 4.7/s;20+ 延迟暴涨 = DeepSeek 开发账号并发限流(外部),非平台。 - 失败注入(错误模型名)→ 任务正确判 failed 并回传原因。 结论:平台并发机制达生产级;真实吞吐上限 = 自托管模型容量(生产 Qwen,可加卡线性扩)。 Co-Authored-By: Claude Opus 4.8 (1M context) --- sundynix-dispatcher/internal/eino/graph.go | 18 +++++++++++++--- sundynix-gateway/internal/store/pgsql.go | 24 ++++++++++++++++++++++ sundynix-mcp-go/internal/memory/store.go | 14 +++++++++++++ 3 files changed, 53 insertions(+), 3 deletions(-) diff --git a/sundynix-dispatcher/internal/eino/graph.go b/sundynix-dispatcher/internal/eino/graph.go index 3729e0f..97533f0 100644 --- a/sundynix-dispatcher/internal/eino/graph.go +++ b/sundynix-dispatcher/internal/eino/graph.go @@ -33,6 +33,7 @@ type board struct { answer string // 当前成稿(多 agent 协作时 = 最近一个 agent 的产出 = 成品) agentOut []string // 各上游 agent 的产出(按序),注入下游 agent 上下文以实现接力协作 rejected bool // HITL 审批节点拒绝/超时 → 置位,runGraph 中止并返回 errRejected + fatalErr error // agent 节点 LLM 调用失败 → 置位,runGraph 中止并上抛 → 任务判 failed(而非 done-空) } // runGraph 按 DSL 图的真实拓扑与连线执行(替代旧的线性拍平 compileFlow)。 @@ -56,7 +57,7 @@ func (o *Orchestrator) runGraph(ctx context.Context, t *contract.Task, tr *execT b.profile = o.fetchMemory(ctx, b.uid, b.query) b.history = o.fetchHistory(ctx, b.sid) o.runConversation(ctx, t.ID, b, plan.System, tr, "agent") - return b.answer, nil + return b.answer, b.fatalErr // 模型失败 → 上抛判 failed } // 建邻接与入度(只认两端都存在的边)。保留整条边以便 branch 按 true/false 标签选路。 @@ -155,8 +156,8 @@ func (o *Orchestrator) runGraph(ctx context.Context, t *contract.Task, tr *execT tr.info(n.Kind+":"+n.ID, "system", labelOf(n, n.Kind), "未识别节点,跳过") } nspan.End() - if b.rejected { - break // 审批拒绝/超时:中止后续节点,整图判 rejected + if b.rejected || b.fatalErr != nil { + break // 审批拒绝/超时 或 模型失败:中止后续节点 } for _, tgt := range propagate { active[tgt] = true @@ -166,11 +167,17 @@ func (o *Orchestrator) runGraph(ctx context.Context, t *contract.Task, tr *execT if b.rejected { return b.answer, errRejected // 合法终态,Handle 据此判 rejected 并优雅收尾 } + if b.fatalErr != nil { + return b.answer, b.fatalErr // 上抛 → Handle 判 failed(带原因) + } // 图里无 agent 节点(纯工具/检索图)也要出一段模型答复,否则没有输出。 if b.answer == "" { o.runConversation(ctx, t.ID, b, plan.System, tr, "agent") } + if b.fatalErr != nil { + return b.answer, b.fatalErr + } return b.answer, nil } @@ -283,6 +290,11 @@ func (o *Orchestrator) runAgent(ctx context.Context, taskID string, b *board, sy } if err != nil { tr.emit(node, "model", "error", "模型流式推理", err.Error(), time.Since(t0).Milliseconds()) + // 未产出任何 token 即失败 → 标记致命错,让任务判 failed(暴露原因,便于监控告警), + // 而非静默 done-空。已流出部分 token 的中断也算失败(结果不完整)。 + if b.fatalErr == nil { + b.fatalErr = fmt.Errorf("agent 模型推理失败: %w", err) + } return } if redacted > 0 { diff --git a/sundynix-gateway/internal/store/pgsql.go b/sundynix-gateway/internal/store/pgsql.go index 8b011f7..c4344cc 100644 --- a/sundynix-gateway/internal/store/pgsql.go +++ b/sundynix-gateway/internal/store/pgsql.go @@ -5,6 +5,9 @@ import ( "context" "errors" "log" + "os" + "strconv" + "time" "gorm.io/driver/postgres" "gorm.io/gorm" @@ -13,6 +16,26 @@ import ( "github.com/sundynix/sundynix-shared/contract" ) +// 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 +} + +// tunePool 给连接池设上限:高并发下不至于无限开连接打爆 PG(max_connections 默认 100)。 +// 各服务默认 25,可经 DB_MAX_OPEN_CONNS / DB_MAX_IDLE_CONNS 调整。 +func tunePool(db *gorm.DB) { + if sqlDB, err := db.DB(); err == nil { + sqlDB.SetMaxOpenConns(envInt("DB_MAX_OPEN_CONNS", 25)) + sqlDB.SetMaxIdleConns(envInt("DB_MAX_IDLE_CONNS", 5)) + sqlDB.SetConnMaxLifetime(time.Hour) + } +} + // errStoreDisabled 表示 Postgres 处于降级(未连接)模式,写操作无法进行。 var errStoreDisabled = errors.New("postgres store disabled") @@ -36,6 +59,7 @@ func OpenPostgres(dsn string) *Postgres { log.Printf("[store] postgres 不可用,降级运行(不持久化): %v", err) return &Postgres{} } + tunePool(db) // 连接池上限,防高并发打爆 PG // 一次性迁移:旧表用整型自增 id,与新雪花字符串 id 不兼容(AutoMigrate 不改主键类型)。 // 备份模型密钥(唯一不可再生的数据) → 重建全部表 → 回灌模型。其余为可重建的测试数据。 migrateLegacyIntIDs(db) diff --git a/sundynix-mcp-go/internal/memory/store.go b/sundynix-mcp-go/internal/memory/store.go index 944bd7a..05d8ceb 100644 --- a/sundynix-mcp-go/internal/memory/store.go +++ b/sundynix-mcp-go/internal/memory/store.go @@ -7,7 +7,9 @@ import ( "fmt" "log" "math" + "os" "sort" + "strconv" "strings" "time" @@ -69,6 +71,18 @@ func Open(dsn string) *Store { log.Printf("[memory] postgres 不可用,记忆降级(召回为空): %v", err) return &Store{} } + // 连接池上限,防高并发打爆 PG(max_connections 默认 100);默认 25,可经 DB_MAX_OPEN_CONNS 调。 + if sqlDB, derr := db.DB(); derr == nil { + maxOpen := 25 + if v := os.Getenv("DB_MAX_OPEN_CONNS"); v != "" { + if n, perr := strconv.Atoi(v); perr == nil && n > 0 { + maxOpen = n + } + } + sqlDB.SetMaxOpenConns(maxOpen) + sqlDB.SetMaxIdleConns(5) + sqlDB.SetConnMaxLifetime(time.Hour) + } // 一次性迁移:旧表用复合主键 (user_id,key) 无 id/时间戳,与雪花规约不兼容。 migrateLegacyProfile(db) if err := db.AutoMigrate(&Profile{}); err != nil {