1674252d81
把"逐轮盲写抽取"升级为 Mem0 式对账(方案见 memory_industry_analysis.md 落地节):
mcp-go:
- Profile 加 Importance(1~10, poignancy) + LastSeenAt(为 Generative Agents 读路径
Score=w1·Relevance+w2·Recency+w3·Importance 铺路)。
- Upsert 收 importance + 每次置 last_seen(印证);新增 Delete(软删,BaseModel.DeletedAt
已具备,失效不物删可审计)+ Touch;memory_upsert 透传 importance、新增 memory_delete 工具。
dispatcher:
- extractMemory → consolidateMemory:一次 LLM 调用同时做 抽取+对账,输出
[{op:ADD|UPDATE|DELETE|NOOP,key,value,importance}];ADD/UPDATE→upsert、DELETE→软删;
sanitizeOps 防幻删(DELETE 须命中已有)/夹 importance[1,10]/同key保末个/丢 NOOP。
- 攒批:每 3 轮(per-session 计数)才 consolidate 一次,省成本,对齐 ChatGPT 周期整理。
从根上解决 exact-key 盲写的记忆腐烂。
验证:parseOps/sanitizeOps/parseProfile 纯逻辑单测;store 集成测试(真 PG)覆盖
importance/last_seen 写入 + 软删(live 0 / 物理 1);dispatcher -race 全过。
(注:完整多轮 LLM consolidate 未做实跑,属构造性验证 + 沿用已证 pool.Chat 模式。)
P2 待做:读路径按 Score(Recency+Importance) 排序/衰减/截断 + 桌面端记忆面板。
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
215 lines
8.1 KiB
Go
215 lines
8.1 KiB
Go
// Package eino 封装基于 CloudWeGo Eino 的 Agent 图编排引擎。
|
||
package eino
|
||
|
||
import (
|
||
"context"
|
||
"encoding/json"
|
||
"fmt"
|
||
"log"
|
||
"sync"
|
||
"time"
|
||
|
||
"github.com/cloudwego/eino/schema"
|
||
|
||
"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"
|
||
)
|
||
|
||
// 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)
|
||
}
|
||
|
||
// 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)
|
||
}
|
||
|
||
// 工具调用超时;超时即降级(不带工具上下文继续推理)。
|
||
const toolCallTimeout = 3 * time.Second
|
||
|
||
// Orchestrator 把每个 DSL 任务动态编译为 Eino 图并执行(记忆召回 → 工具节点 → 注入 → 流式)。
|
||
type Orchestrator struct {
|
||
pool LLM
|
||
breaker *harness.CircuitBreaker
|
||
eval *harness.Evaluator
|
||
sink TokenSink
|
||
tools ToolCaller
|
||
exec ExecSink
|
||
|
||
turnMu sync.Mutex // 保护 turns(攒批计数,多任务 goroutine 共享)
|
||
turns map[string]int // sessionID → 累计轮次,用于每 N 轮触发 consolidate
|
||
}
|
||
|
||
// NewOrchestrator 持有依赖;图按任务的 DSL 在 Handle 内动态编译。
|
||
// exec 为执行可视化事件出口(可为 nil,则不发轨迹事件);eval 为自动化评测(可为 nil)。
|
||
func NewOrchestrator(pool LLM, breaker *harness.CircuitBreaker, eval *harness.Evaluator, sink TokenSink, tools ToolCaller, exec ExecSink) (*Orchestrator, error) {
|
||
return &Orchestrator{pool: pool, breaker: breaker, eval: eval, sink: sink, tools: tools, exec: exec}, nil
|
||
}
|
||
|
||
// Handle 消费一个任务:按 DSL 编译 Eino 图并执行,把 Token 流回流到 sundynix.streams.<id>。
|
||
func (o *Orchestrator) Handle(ctx context.Context, t *contract.Task) error {
|
||
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)
|
||
return nil
|
||
}
|
||
|
||
// 报告生成走专用多步编排(规划→分章并行检索撰写→汇聚→渲染 Word),而非通用对话图。
|
||
if intent, _ := t.Meta[contract.MetaIntent].(string); intent == contract.IntentReport {
|
||
return o.handleReport(ctx, t, tr)
|
||
}
|
||
log.Printf("[eino] task %s received (graph=%d bytes), 按图执行(拓扑+连线+分支)...", t.ID, len(t.Graph))
|
||
tr.info("task", "system", "任务受理", fmt.Sprintf("DSL %d 字节,按图执行", len(t.Graph)))
|
||
|
||
// 按 DSL 图的真实拓扑/连线/分支执行(graph.go 解释器),agent 节点流式回流 token。
|
||
answer, err := o.runGraph(ctx, t, tr)
|
||
if err != nil {
|
||
log.Printf("[eino] task %s graph error: %v", t.ID, err)
|
||
_ = o.sink.CompleteStream(t.ID)
|
||
o.breaker.Report(false)
|
||
return err
|
||
}
|
||
|
||
if cerr := o.sink.CompleteStream(t.ID); cerr != nil {
|
||
log.Printf("[eino] complete stream failed: %v", cerr)
|
||
}
|
||
log.Printf("[eino] task %s done (%d 字答复)", t.ID, len([]rune(answer)))
|
||
o.breaker.Report(true)
|
||
|
||
// 写回阶段:离开热路径、异步落历史 + (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 %.2f)flags=%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) {
|
||
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)
|
||
}
|
||
}
|