Files
sundynix-agentix/sundynix-dispatcher/internal/eino/orchestrator.go
T
Blizzard 38a74e3911 fix(history): 空答复不落历史 + 发送前过滤空消息(根治会话毒化)
现象:某轮 LLM 返回空(余额不足/失败/降级)→ 空 assistant 消息被写进会话历史
→ 此后该 session 每次请求都带一条空 content 的 assistant 消息 → DeepSeek 400
"Invalid assistant message: content or tool_calls must be set" → 又空 → 又写空 → 自我循环毒化。

- Fix A(断源):orchestrator.memorize 对空答复直接 return,不落历史(失败的一轮不留痕)。
- Fix B(兜底):buildMessages 发送前剔除 content 与 tool_calls 均空的历史消息,
  防御任何来源的脏历史(含存量)。

验证:清理被污染的 default 会话后恢复正常;新逻辑下空答复不再写入。dispatcher
build+vet+test 全绿。

注:诊断中发现激活模型 deepseek-v4-pro 是推理模型(答案在 reasoning_content,content 为空),
已在管理端切回 deepseek-chat。仍待办:LLM 调用失败应暴露为 failed 而非 done-空输出。

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-24 15:53:30 +08:00

313 lines
13 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
// Package eino 封装基于 CloudWeGo Eino 的 Agent 图编排引擎。
package eino
import (
"context"
"encoding/json"
"errors"
"fmt"
"log"
"log/slog"
"strings"
"sync"
"time"
"github.com/cloudwego/eino/components/model"
"github.com/cloudwego/eino/schema"
"go.opentelemetry.io/otel/attribute"
"go.opentelemetry.io/otel/codes"
"go.opentelemetry.io/otel/trace"
"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"
"github.com/sundynix/sundynix-shared/otelx"
)
// TokenSink 是 Token 流回流出口(由 NATS bus 实现)。
type TokenSink interface {
PublishToken(taskID string, token []byte) error
CompleteStream(taskID string) error
}
// ToolCaller 经 NATS 调起第 5 层 MCP 工具(由 NATS bus 实现)。
type ToolCaller interface {
CallTool(ctx context.Context, subject string, call *contract.ToolCall) (*contract.ToolResult, error)
}
// StatusSink 回写任务生命周期状态(由 NATS bus 实现;可为 nil → 不回写)。
type StatusSink interface {
PublishTaskStatus(taskID, status, detail string) error
}
// ApprovalWaiter 阻塞等待审批节点的人工决定(由 NATS bus 实现;可为 nil → 审批节点自动放行)。
type ApprovalWaiter interface {
WaitApproval(ctx context.Context, taskID string, timeout time.Duration) (*contract.ApprovalDecision, error)
}
// errRejected 是审批节点拒绝(或超时)时图执行返回的哨兵错误:它是合法终态而非故障,
// Handle 据此判 rejected 并优雅收尾(不计熔断失败)。
var errRejected = errors.New("approval rejected")
// LLM 是编排所需的语言模型能力(生产由 *llm.Pool 实现)。抽成接口便于测试注入假模型。
type LLM interface {
Ready() bool
ChatStream(ctx context.Context, msgs []llm.ChatMessage, onToken func(string)) error
StreamText(ctx context.Context, text string, onToken func([]byte)) error
Chat(ctx context.Context, msgs []llm.ChatMessage) (string, error)
// ToolCallingModel 返回支持函数调用的模型(ReAct agent 用);不支持则返回 nil。
ToolCallingModel() model.ToolCallingChatModel
// ChatModel 返回 Eino ChatModel 组件(compose.Graph 编排用);未就绪则 nil。
ChatModel() model.BaseChatModel
}
// 工具调用超时;超时即降级(不带工具上下文继续推理)。
const toolCallTimeout = 3 * time.Second
// taskExecTimeout 是单个任务整体执行上限;超时即判 timeout(状态机),避免无限期"运行中"。
// 含 HITL 审批等待预算(approvalTimeout+ 常规图执行;须 < bus 消费者 AckWait(15min) 以免重投。
const taskExecTimeout = 10 * time.Minute
// approvalTimeout 是单个审批节点等待人工决定的上限;超时安全默认拒绝(fail-safe)。
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 → 审批节点自动放行)
turnMu sync.Mutex // 保护 turns(攒批计数,多任务 goroutine 共享)
turns map[string]int // sessionID → 累计轮次,用于每 N 轮触发 consolidate
}
// NewOrchestrator 持有依赖;图按任务的 DSL 在 Handle 内动态编译。
// exec 为执行可视化事件出口(可为 nil,则不发轨迹事件);eval 为自动化评测(可为 nil);
// status 为任务生命周期状态回写出口(可为 nil);approval 为 HITL 审批等待(可为 nil)。
func NewOrchestrator(pool LLM, breaker *harness.CircuitBreaker, eval *harness.Evaluator, sink TokenSink, tools ToolCaller, exec ExecSink, status StatusSink, approval ApprovalWaiter) (*Orchestrator, error) {
return &Orchestrator{pool: pool, breaker: breaker, eval: eval, sink: sink, tools: tools, exec: exec, status: status, approval: approval}, nil
}
// setStatus 回写一次任务状态流转(status 为 nil 时静默跳过)。
func (o *Orchestrator) setStatus(taskID, status, detail string) {
if o.status == nil {
return
}
if err := o.status.PublishTaskStatus(taskID, status, detail); err != nil {
log.Printf("[eino] 回写任务状态 %s=%s 失败: %v", taskID, status, err)
}
}
// finishStatus 据收尾错误把任务置为 done / timeout / failed。
func (o *Orchestrator) finishStatus(taskID string, err error) {
switch {
case err == nil:
o.setStatus(taskID, contract.TaskDone, "")
case errors.Is(err, context.DeadlineExceeded):
o.setStatus(taskID, contract.TaskTimeout, "执行超时")
default:
o.setStatus(taskID, contract.TaskFailed, truncate(err.Error(), 200))
}
}
// Handle 消费一个任务:按 DSL 编译 Eino 图并执行,把 Token 流回流到 sundynix.streams.<id>。
func (o *Orchestrator) Handle(ctx context.Context, t *contract.Task) error {
// 护栏:丢弃空任务(无 id),避免误投/历史脏数据被当真任务处理并触发状态回写放大。
if t.ID == "" {
log.Printf("[eino] 跳过空任务(无 id")
return nil
}
// 链路根(dispatcher 侧):续上 gateway 经 NATS 传来的 trace,覆盖整个图执行。
ctx, span := otelx.Tracer().Start(ctx, "task.execute",
trace.WithAttributes(attribute.String("sundynix.task_id", t.ID)))
defer span.End()
tr := o.tracer(t.ID)
defer tr.done()
// 熔断开启:快速拒绝,但要让客户端解阻(回流提示 + 收尾流),不静默丢弃。
if !o.breaker.Allow() {
log.Printf("[eino] 熔断开启,拒绝任务 %s", t.ID)
tr.info("task", "system", "服务熔断", "后端连续失败,暂时拒绝新任务,请稍后重试")
_ = o.sink.PublishToken(t.ID, []byte("⚠️ 服务繁忙(已触发熔断保护),请稍后重试。"))
_ = o.sink.CompleteStream(t.ID)
o.setStatus(t.ID, contract.TaskFailed, "服务熔断")
return nil
}
// 任务状态机:进入执行 → running;整体加超时上限,超时判 timeout(杜绝无限期"运行中")。
o.setStatus(t.ID, contract.TaskRunning, "")
tctx, cancel := context.WithTimeout(ctx, taskExecTimeout)
defer cancel()
// 报告生成走专用多步编排(规划→分章并行检索撰写→汇聚→渲染 Word),而非通用对话图。
if intent, _ := t.Meta[contract.MetaIntent].(string); intent == contract.IntentReport {
err := o.handleReport(tctx, t, tr)
o.finishStatus(t.ID, err)
return err
}
// ctx 携带 task.execute span → 这些任务生命周期日志自动带 trace_id,可与 Jaeger 链路互跳。
slog.InfoContext(ctx, "task received", "task_id", t.ID, "graph_bytes", len(t.Graph))
tr.info("task", "system", "任务受理", fmt.Sprintf("DSL %d 字节,按图执行", len(t.Graph)))
// 按 DSL 图执行:compose.GraphEINO_COMPOSE=1)或自研 graph.go(默认);agent 节点流式回流 token。
answer, err := o.executeGraph(tctx, t, tr)
if errors.Is(err, errRejected) {
// HITL 拒绝:合法终态,非故障。收尾流 + 置 rejected,不计熔断、不重投。
slog.InfoContext(ctx, "task rejected by approval", "task_id", t.ID)
if answer != "" {
_ = o.sink.PublishToken(t.ID, []byte(answer))
}
_ = o.sink.CompleteStream(t.ID)
o.breaker.Report(true) // 拒绝是人为决策,不算后端失败
o.setStatus(t.ID, contract.TaskRejected, truncate(answer, 120))
return nil
}
if err != nil {
span.RecordError(err)
span.SetStatus(codes.Error, err.Error())
slog.ErrorContext(ctx, "task graph error", "task_id", t.ID, "err", err.Error())
_ = o.sink.CompleteStream(t.ID)
o.breaker.Report(false)
o.finishStatus(t.ID, err)
return err
}
if cerr := o.sink.CompleteStream(t.ID); cerr != nil {
log.Printf("[eino] complete stream failed: %v", cerr)
}
slog.InfoContext(ctx, "task done", "task_id", t.ID, "answer_runes", len([]rune(answer)))
o.breaker.Report(true)
o.finishStatus(t.ID, nil)
// 写回阶段:离开热路径、异步落历史 + (TODO)抽取记忆。
go o.memorize(t, answer)
// 自动化评测:离开热路径,对本轮输出打分并记录(规则 + LLM-as-judge)。
go o.evaluate(t, dsl.Compile(t.Graph).Query, answer)
return nil
}
// evaluate 异步对一次输出做自动化评测并记录评分(off 热路径,不影响响应)。
func (o *Orchestrator) evaluate(t *contract.Task, input, output string) {
if o.eval == nil {
return
}
ctx, cancel := context.WithTimeout(context.Background(), 15*time.Second)
defer cancel()
r := o.eval.Score(ctx, input, output)
log.Printf("[eval] task %s 综合 %.2f(规则 %.2f / LLM %.2fflags=%v %s",
t.ID, r.Overall, r.Rule, r.LLM, r.Flags, r.Reason)
}
// fetchMemory 经 MCP memory_get 工具召回用户常驻画像。
// 工具不可用/超时/无 user_id 时返回空串,降级为无记忆推理(不阻断主流程)。
func (o *Orchestrator) fetchMemory(ctx context.Context, userID, _ string) string {
if o.tools == nil || userID == "" {
return ""
}
cctx, cancel := context.WithTimeout(ctx, toolCallTimeout)
defer cancel()
res, err := o.tools.CallTool(cctx, contract.ToolSubjectGo("memory_get"), &contract.ToolCall{
Tool: "memory_get",
Args: map[string]any{"user_id": userID},
})
if err != nil {
log.Printf("[eino] memory_get unavailable for %s, degrade: %v", userID, err)
return ""
}
if !res.OK {
log.Printf("[eino] memory_get error for %s: %s", userID, res.Error)
return ""
}
log.Printf("[eino] memory_get ok for %s: %s", userID, res.Content)
return res.Content
}
// fetchHistory 经 MCP history_get 工具召回会话短期多轮历史,转为 Eino 消息。
// 工具不可用/无 session 时返回空,降级为无历史(不阻断主流程)。
func (o *Orchestrator) fetchHistory(ctx context.Context, sessionID string) []*schema.Message {
if o.tools == nil || sessionID == "" {
return nil
}
cctx, cancel := context.WithTimeout(ctx, toolCallTimeout)
defer cancel()
res, err := o.tools.CallTool(cctx, contract.ToolSubjectGo("history_get"), &contract.ToolCall{
Tool: "history_get",
Args: map[string]any{"session_id": sessionID},
})
if err != nil || res == nil || !res.OK || res.Content == "" {
return nil
}
var turns []struct {
Role string `json:"role"`
Content string `json:"content"`
}
if json.Unmarshal([]byte(res.Content), &turns) != nil {
return nil
}
msgs := make([]*schema.Message, 0, len(turns))
for _, tn := range turns {
if tn.Role == "assistant" {
msgs = append(msgs, schema.AssistantMessage(tn.Content, nil))
} else {
msgs = append(msgs, schema.UserMessage(tn.Content))
}
}
if len(msgs) > 0 {
log.Printf("[eino] history_get ok for %s: %d 条历史", sessionID, len(msgs))
}
return msgs
}
// consolidateEveryTurns 控制记忆对账的攒批节奏:每 N 轮 consolidate 一次(不逐轮,省成本)。
const consolidateEveryTurns = 3
// memorize 写回阶段(异步、离热路径):落短期历史;每 N 轮做一次记忆对账(consolidate)。
func (o *Orchestrator) memorize(t *contract.Task, answer string) {
// 空答复(LLM 失败/降级/拒绝)不落历史:空 assistant 消息会被 LLM API 拒绝(400),
// 一旦写入会毒化该会话后续所有请求。失败的一轮干脆不留痕。
if strings.TrimSpace(answer) == "" {
return
}
uid, _ := t.Meta[contract.MetaUserID].(string)
sid, _ := t.Meta[contract.MetaSessionID].(string)
if sid != "" && o.tools != nil {
o.appendHistory(sid, "user", dsl.Compile(t.Graph).Query) // 落真实用户输入,而非 DSL 原文
o.appendHistory(sid, "assistant", answer)
log.Printf("[eino] (writeback) task %s 已落会话历史 session=%s", t.ID, sid)
}
if uid == "" || sid == "" || o.tools == nil {
return
}
// 攒批:累计轮次,每 N 轮才把近期对话与已有画像交给 LLM 对账一次。
o.turnMu.Lock()
if o.turns == nil {
o.turns = map[string]int{}
}
o.turns[sid]++
n := o.turns[sid]
o.turnMu.Unlock()
if n%consolidateEveryTurns == 0 {
ctx := context.Background()
o.consolidateMemory(ctx, uid, o.fetchHistory(ctx, sid))
}
}
func (o *Orchestrator) appendHistory(sessionID, role, content string) {
cctx, cancel := context.WithTimeout(context.Background(), toolCallTimeout)
defer cancel()
if _, err := o.tools.CallTool(cctx, contract.ToolSubjectGo("history_append"), &contract.ToolCall{
Tool: "history_append",
Args: map[string]any{"session_id": sessionID, "role": role, "content": content},
}); err != nil {
log.Printf("[eino] history_append failed: %v", err)
}
}