feat(prod): DB 连接池上限 + LLM 失败暴露为 failed(生产级并发收尾)
为高并发生产做的三项收尾(配合已有的任务/工具并发消费): 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) <noreply@anthropic.com>
This commit is contained in:
@@ -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 {
|
||||
|
||||
Reference in New Issue
Block a user