v1.7 Session 会话机制 + 多工作区隔离 + Reporter 抽象

This commit is contained in:
zhuyongxin
2026-05-29 14:09:12 +08:00
parent 3e91c7d0a1
commit e480c5c55e
6 changed files with 424 additions and 120 deletions
+181 -94
View File
@@ -18,23 +18,22 @@ type AgentEngine struct {
registry tools.Registry
// WorkDir (工作区): 借鉴 OpenClaw 的理念,Agent 必须有一个明确的物理边界
WorkDir string
// WorkDir string
EnableThinking bool // 【新增】慢思考模式开关
composer *prompt.PromptComposer
}
func NewAgentEngine(p provider.LLMProvider, r tools.Registry, workDir string, enableThinking bool) *AgentEngine {
// 移除了 Engine 层级的 WorkDir,因为 WorkDir 现在应该跟随 Session 走
func NewAgentEngine(p provider.LLMProvider, r tools.Registry, enableThinking bool) *AgentEngine {
return &AgentEngine{
provider: p,
registry: r,
WorkDir: workDir,
EnableThinking: enableThinking,
composer: prompt.NewPromptComposer(workDir),
}
}
func (e *AgentEngine) String() string {
return fmt.Sprintf("AgentEngine{workDir: %s, thinking: %v, registry: %s}", e.WorkDir, e.EnableThinking, e.registry)
return fmt.Sprintf("AgentEngine{thinking: %v, registry: %s}", e.EnableThinking, e.registry)
}
// dumpMessages 打印当前上下文中的所有消息 (调试用)
@@ -59,138 +58,226 @@ func dumpTools(tools []schema.ToolDefinition) {
}
}
func (e *AgentEngine) Run(ctx context.Context, userPrompt string) error {
log.Printf("[Engine] 引擎启动,锁定工作区: %s\n", e.WorkDir)
log.Printf("[Engine] 慢思考模式 (Thinking Phase): %v\n", e.EnableThinking)
func (e *AgentEngine) Run(ctx context.Context, session *Session, reporter Reporter) error {
log.Printf("[Engine] 唤醒会话 [%s],锁定工作区: %s\n", session.ID, session.WorkDir)
systemMsg := e.composer.Build()
contextHistory := []schema.Message{
systemMsg, // 注入动态组装的内核、AGENTS.md 与 Skills
{Role: schema.RoleUser, Content: userPrompt},
}
turnCount := 0
const maxTurns = 10
// 根据当前 Session 的工作区,动态组装最新的 System Prompt
composer := prompt.NewPromptComposer(session.WorkDir)
systemMsg := composer.Build()
for {
turnCount++
if turnCount > maxTurns {
log.Printf("[Engine] 已达最大轮数 (%d),强制终止。\n", maxTurns)
break
}
log.Printf("\n========== [Turn %d] 开始 ==========\n", turnCount)
dumpMessages(contextHistory)
// 获取当前挂载的所有工具定义
availableTools := e.registry.GetAvailableTools()
dumpTools(availableTools)
// ====================================================================
// Phase 1: 慢思考阶段 (Thinking) - 仅第一轮执行初始规划
// ====================================================================
if e.EnableThinking && turnCount == 1 {
log.Println("[Engine][Phase 1] 剥夺工具访问权,强制进入慢思考与规划阶段...")
// 1. 【上下文组装】: System Prompt + 截取最近的 6 条消息作为 Working Memory
// 在实际业务中,由于工具返回结果可能很长,短期工作记忆往往设为 6-10 条足以维系连贯对话
workingMemory := session.GetWorkingMemory(6)
var contextHistory []schema.Message
contextHistory = append(contextHistory, systemMsg)
contextHistory = append(contextHistory, workingMemory...)
// 2. ================= Phase 1: Thinking =================
if e.EnableThinking {
if reporter != nil {
reporter.OnThinking(ctx)
}
// 核心机制:传入的 availableTools 为 nil!
// 大模型看不到任何 JSON Schema,被迫只能输出纯文本的思考过程。
thinkResp, err := e.provider.Generate(ctx, contextHistory, nil)
if err != nil {
return fmt.Errorf("Thinking 阶段生成失败: %w", err)
return fmt.Errorf("Thinking 阶段失败: %w", err)
}
// 如果模型输出了思考过程,我们将其作为 Assistant 消息追加到上下文中
if thinkResp.Content != "" {
fmt.Printf("🧠 [内部思考 Trace]: %s\n", thinkResp.Content)
// 将思考过程持久化到 Session 中!
session.Append(*thinkResp)
// 把它追加到当前这一轮的临时上下文中,供 Action 阶段使用
contextHistory = append(contextHistory, *thinkResp)
}
// 插入过渡指令:让模型知道现在可以调用工具了
contextHistory = append(contextHistory, schema.Message{
Role: schema.RoleUser,
Content: "根据你的推理,现在请使用可用的工具来完成任务。执行具体行动。",
})
}
// ====================================================================
// Phase 2: 行动阶段 (Action) - 恢复工具,顺着规划执行
// ====================================================================
log.Println("[Engine][Phase 2] 恢复工具挂载,等待模型采取行动...")
// 此时的 contextHistory 中已经包含了上一阶段模型自己的 Thinking Trace + 过渡指令。
// 模型会顺着自己的逻辑,结合恢复的 availableTools 发起精准的工具调用。
// 3. ================= Phase 2: Action =================
actionResp, err := e.provider.Generate(ctx, contextHistory, availableTools)
if err != nil {
return fmt.Errorf("Action 阶段生成失败: %w", err)
return fmt.Errorf("Action 阶段失败: %w", err)
}
// 将大模型的行动响应持久化到 Session 中
session.Append(*actionResp)
contextHistory = append(contextHistory, *actionResp)
if actionResp.Content != "" {
fmt.Printf("🤖 [对外回复]: %s\n", actionResp.Content)
if actionResp.Content != "" && reporter != nil {
reporter.OnMessage(ctx, actionResp.Content)
}
// ====================================================================
// 退出与执行逻辑 (与上一讲保持一致)
// ====================================================================
if len(actionResp.ToolCalls) == 0 {
log.Println("[Engine] 模型未请求调用工具,任务宣告完成。")
// 如果没有工具调用,说明本次任务已完成,打破 ReAct 循环,挂起等待人类的下一条指令
break
}
log.Printf("[Engine] 模型请求并发调用 %d 个工具...\n", len(actionResp.ToolCalls))
// 【核心改造开始】: 从串行 (Sequential) 演进为并行 (Parallel)
// 1. 预分配一个固定长度的切片,用于安全地存放各个并发工具的执行结果(Observation)
// 长度与 ToolCalls 的数量完全一致
// 4. ================= 并发执行底层工具 =================
observationMsgs := make([]schema.Message, len(actionResp.ToolCalls))
// 2. 声明 WaitGroup 用于阻塞等待所有协程完成
var wg sync.WaitGroup
// 3. 遍历模型请求的所有工具,为每一个工具单独 Fork 出一个 Goroutine
for i, toolCall := range actionResp.ToolCalls {
wg.Add(1) // 增加计数器
wg.Add(1)
// 开启协程。注意:一定要将索引 i 和 toolCall 作为参数传入匿名函数,防止闭包变量捕获陷阱!
go func(idx int, call schema.ToolCall) {
defer wg.Done() // 协程结束时计数器减一
defer wg.Done()
log.Printf(" -> [Go-%d] 🛠️ 触发并行执行: %s\n", idx, call.Name)
// 调用底层 Registry 执行工具(物理操作)
result := e.registry.Execute(ctx, call)
if result.IsError {
log.Printf(" -> [Go-%d] ❌ 工具执行报错: %s\n", idx, result.Output)
} else {
log.Printf(" -> [Go-%d] ✅ 工具执行成功 (返回 %d 字节)\n", idx, len(result.Output))
if reporter != nil {
reporter.OnToolCall(ctx, call.Name, string(call.Arguments))
}
// 将执行结果封装为一条用户消息 (RoleUser)
obsMsg := schema.Message{
result := e.registry.Execute(ctx, call)
if reporter != nil {
displayOutput := result.Output
if len(displayOutput) > 200 {
displayOutput = displayOutput[:200] + "... (已截断)"
}
reporter.OnToolResult(ctx, call.Name, displayOutput, result.IsError)
}
observationMsgs[idx] = schema.Message{
Role: schema.RoleUser,
Content: result.Output,
ToolCallID: call.ID,
}
// 【线程安全】: 由于每个 Goroutine 操作的是预分配切片的不同索引,
// 这里不需要加锁 (Mutex),性能极高!
observationMsgs[idx] = obsMsg
}(i, toolCall) // 闭包传参
}(i, toolCall)
}
// 4. Join 阻塞等待:主循环挂起,直到所有的并发协程全部执行完毕
wg.Wait()
log.Println("[Engine] 所有并发工具执行完毕,开始聚合观察结果 (Observation)...")
// 5. 聚合装填:将并行的结果,按照原本的顺序,一次性追加到上下文时间线中
// 这等价于 contextHistory = append(contextHistory, observationMsgs...)
for _, obs := range observationMsgs {
contextHistory = append(contextHistory, obs)
}
// 将所有的工具执行结果(Observation)持久化到 Session 中,开启下一轮的复盘与推理
session.Append(observationMsgs...)
// systemMsg := e.composer.Build()
// contextHistory := []schema.Message{
// systemMsg, // 注入动态组装的内核、AGENTS.md 与 Skills
// {Role: schema.RoleUser, Content: userPrompt},
// }
// turnCount := 0
// const maxTurns = 10
// for {
// turnCount++
// if turnCount > maxTurns {
// log.Printf("[Engine] 已达最大轮数 (%d),强制终止。\n", maxTurns)
// break
// }
// log.Printf("\n========== [Turn %d] 开始 ==========\n", turnCount)
// dumpMessages(contextHistory)
// // 获取当前挂载的所有工具定义
// availableTools := e.registry.GetAvailableTools()
// dumpTools(availableTools)
// // ====================================================================
// // Phase 1: 慢思考阶段 (Thinking) - 仅第一轮执行初始规划
// // ====================================================================
// if e.EnableThinking && turnCount == 1 {
// log.Println("[Engine][Phase 1] 剥夺工具访问权,强制进入慢思考与规划阶段...")
// // 核心机制:传入的 availableTools 为 nil!
// // 大模型看不到任何 JSON Schema,被迫只能输出纯文本的思考过程。
// thinkResp, err := e.provider.Generate(ctx, contextHistory, nil)
// if err != nil {
// return fmt.Errorf("Thinking 阶段生成失败: %w", err)
// }
// // 如果模型输出了思考过程,我们将其作为 Assistant 消息追加到上下文中
// if thinkResp.Content != "" {
// fmt.Printf("🧠 [内部思考 Trace]: %s\n", thinkResp.Content)
// contextHistory = append(contextHistory, *thinkResp)
// }
// // 插入过渡指令:让模型知道现在可以调用工具了
// contextHistory = append(contextHistory, schema.Message{
// Role: schema.RoleUser,
// Content: "根据你的推理,现在请使用可用的工具来完成任务。执行具体行动。",
// })
// }
// // ====================================================================
// // Phase 2: 行动阶段 (Action) - 恢复工具,顺着规划执行
// // ====================================================================
// log.Println("[Engine][Phase 2] 恢复工具挂载,等待模型采取行动...")
// // 此时的 contextHistory 中已经包含了上一阶段模型自己的 Thinking Trace + 过渡指令。
// // 模型会顺着自己的逻辑,结合恢复的 availableTools 发起精准的工具调用。
// actionResp, err := e.provider.Generate(ctx, contextHistory, availableTools)
// if err != nil {
// return fmt.Errorf("Action 阶段生成失败: %w", err)
// }
// contextHistory = append(contextHistory, *actionResp)
// if actionResp.Content != "" {
// fmt.Printf("🤖 [对外回复]: %s\n", actionResp.Content)
// }
// // ====================================================================
// // 退出与执行逻辑 (与上一讲保持一致)
// // ====================================================================
// if len(actionResp.ToolCalls) == 0 {
// log.Println("[Engine] 模型未请求调用工具,任务宣告完成。")
// break
// }
// log.Printf("[Engine] 模型请求并发调用 %d 个工具...\n", len(actionResp.ToolCalls))
// // 【核心改造开始】: 从串行 (Sequential) 演进为并行 (Parallel)
// // 1. 预分配一个固定长度的切片,用于安全地存放各个并发工具的执行结果(Observation)
// // 长度与 ToolCalls 的数量完全一致
// observationMsgs := make([]schema.Message, len(actionResp.ToolCalls))
// // 2. 声明 WaitGroup 用于阻塞等待所有协程完成
// var wg sync.WaitGroup
// // 3. 遍历模型请求的所有工具,为每一个工具单独 Fork 出一个 Goroutine
// for i, toolCall := range actionResp.ToolCalls {
// wg.Add(1) // 增加计数器
// // 开启协程。注意:一定要将索引 i 和 toolCall 作为参数传入匿名函数,防止闭包变量捕获陷阱!
// go func(idx int, call schema.ToolCall) {
// defer wg.Done() // 协程结束时计数器减一
// log.Printf(" -> [Go-%d] 🛠️ 触发并行执行: %s\n", idx, call.Name)
// // 调用底层 Registry 执行工具(物理操作)
// result := e.registry.Execute(ctx, call)
// if result.IsError {
// log.Printf(" -> [Go-%d] ❌ 工具执行报错: %s\n", idx, result.Output)
// } else {
// log.Printf(" -> [Go-%d] ✅ 工具执行成功 (返回 %d 字节)\n", idx, len(result.Output))
// }
// // 将执行结果封装为一条用户消息 (RoleUser)
// obsMsg := schema.Message{
// Role: schema.RoleUser,
// Content: result.Output,
// ToolCallID: call.ID,
// }
// // 【线程安全】: 由于每个 Goroutine 操作的是预分配切片的不同索引,
// // 这里不需要加锁 (Mutex),性能极高!
// observationMsgs[idx] = obsMsg
// }(i, toolCall) // 闭包传参
// }
// // 4. Join 阻塞等待:主循环挂起,直到所有的并发协程全部执行完毕
// wg.Wait()
// log.Println("[Engine] 所有并发工具执行完毕,开始聚合观察结果 (Observation)...")
// // 5. 聚合装填:将并行的结果,按照原本的顺序,一次性追加到上下文时间线中
// // 这等价于 contextHistory = append(contextHistory, observationMsgs...)
// for _, obs := range observationMsgs {
// contextHistory = append(contextHistory, obs)
// }
}
return nil
+19
View File
@@ -0,0 +1,19 @@
package engine
import "context"
// Reporter 定义了 Agent 引擎向外界输出信息的规范。
// 这使得引擎可以无缝切换终端 (CLI)、飞书、钉钉甚至 WebUI 等不同的展现层。
type Reporter interface {
// OnThinking 当模型开始进行慢思考 (Reasoning) 时调用
OnThinking(ctx context.Context)
// OnToolCall 当模型决定并发调用工具时调用
OnToolCall(ctx context.Context, toolName string, args string)
// OnToolResult 当工具在底层执行完毕并返回结果时调用
OnToolResult(ctx context.Context, toolName string, result string, isError bool)
// OnMessage 当模型宣告任务完成,向用户输出最终纯文本回答时调用
OnMessage(ctx context.Context, content string)
}
+101
View File
@@ -0,0 +1,101 @@
package engine
import (
"sync"
"time"
"go-tiny-claw/internal/schema"
)
// Session 代表了一次持续的人机交互过程。它负责维护该会话的完整历史。
type Session struct {
ID string
WorkDir string // 该会话绑定的物理工作区
CreatedAt time.Time
UpdatedAt time.Time
// 存放此 Session 中所有的用户输入、大模型回复和工具调用结果
history []schema.Message
mu sync.RWMutex // 读写锁,防止并发读写历史时发生 Data Race
}
func NewSession(id string, workDir string) *Session {
return &Session{
ID: id,
WorkDir: workDir,
CreatedAt: time.Now(),
UpdatedAt: time.Now(),
history: make([]schema.Message, 0),
}
}
// Append 线程安全地向 Session 中追加消息
func (s *Session) Append(msgs ...schema.Message) {
s.mu.Lock()
defer s.mu.Unlock()
s.history = append(s.history, msgs...)
s.UpdatedAt = time.Now()
// 【持久化预留点】:在真实的工业级实现中(如 Claude Code),
// 我们会在这里将 s.history 以 JSONL 的格式 Append 到 workDir/.claw/sessions/xxx.jsonl 中。
// s.SaveToDisk()
}
// GetWorkingMemory 是驾驭工程的核心!
// 它不返回全量历史,而是从后往前截取最近的 N 条消息,形成 Agent 的“短期工作记忆”。
func (s *Session) GetWorkingMemory(limit int) []schema.Message {
s.mu.RLock()
defer s.mu.RUnlock()
total := len(s.history)
if total <= limit || limit <= 0 {
// 如果历史总量小于限制,或者不设限,全量返回 (需要深拷贝以防外部修改)
res := make([]schema.Message, total)
copy(res, s.history)
return res
}
// 截取最近的 limit 条消息
res := make([]schema.Message, limit)
copy(res, s.history[total-limit:])
// 【驾驭防线】:大模型 API 强制要求历史消息的连续性!
// 如果我们截断的第一条消息恰好是一个 ToolResult (RoleUser 且含有 ToolCallID),
// 但发出这个请求的 ToolCall 被我们截断抛弃了,大模型 API 会直接报 400 Bad Request。
// 因此,如果切片首条属于“孤儿”工具响应,我们必须将其强行舍弃,顺延到下一条正常的 User/Assistant 消息。
for len(res) > 0 {
if res[0].Role == schema.RoleUser && res[0].ToolCallID != "" {
res = res[1:]
} else {
break
}
}
return res
}
// ==========================================
// 全局 Session Manager: 用于多用户/多终端隔离
// ==========================================
type SessionManager struct {
sessions map[string]*Session
mu sync.RWMutex
}
var GlobalSessionMgr = &SessionManager{
sessions: make(map[string]*Session),
}
// GetOrCreate 获取或创建一个会话
func (sm *SessionManager) GetOrCreate(id string, workDir string) *Session {
sm.mu.Lock()
defer sm.mu.Unlock()
if sess, exists := sm.sessions[id]; exists {
return sess
}
sess := NewSession(id, workDir)
sm.sessions[id] = sess
return sess
}
+47
View File
@@ -0,0 +1,47 @@
package engine
import (
"context"
"fmt"
"strings"
)
type TerminalReporter struct{}
func NewTerminalReporter() *TerminalReporter {
return &TerminalReporter{}
}
func (r *TerminalReporter) OnThinking(ctx context.Context) {
fmt.Printf("\n[🤔 思考中] 模型正在推理...\n")
}
func (r *TerminalReporter) OnToolCall(ctx context.Context, toolName string, args string) {
fmt.Printf("[🛠️ 调用工具] %s\n", toolName)
// 清理参数中的换行符和特殊字符
displayArgs := strings.ReplaceAll(args, "\n", "\\n")
displayArgs = strings.ReplaceAll(displayArgs, "\r", "\\r")
if len(displayArgs) > 150 {
displayArgs = displayArgs[:150] + "... (已截断)"
}
fmt.Printf(" 参数: %s\n", displayArgs)
}
func (r *TerminalReporter) OnToolResult(ctx context.Context, toolName string, result string, isError bool) {
if isError {
fmt.Printf("[❌ 执行失败] %s\n", toolName)
// 显示错误信息
if result != "" {
fmt.Printf(" 错误: %s\n", result)
}
} else {
fmt.Printf("[✅ 执行成功] %s\n", toolName)
}
}
func (r *TerminalReporter) OnMessage(ctx context.Context, content string) {
if content == "" {
return
}
fmt.Printf("\n🤖 Agent 回复:\n%s\n\n", content)
}