Compare commits
2
Commits
3e91c7d0a1
..
main
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
572665b3ac | ||
|
|
e480c5c55e |
@@ -65,6 +65,31 @@ func main() {
|
||||
|
||||
## 版本历史
|
||||
|
||||
### v1.7 — Session 会话机制 + 多工作区隔离 + Reporter 抽象
|
||||
|
||||
#### 变更
|
||||
|
||||
- **Session 会话机制** — 新增 `Session` 结构体,维护完整的对话历史(`[]schema.Message`),通过 `RWMutex` 实现线程安全的并发读写。全局 `SessionManager` 支持多会话隔离(`GetOrCreate`),基于 `map[string]*Session` 路由
|
||||
- **Working Memory 滑动窗口** — `GetWorkingMemory(limit)` 从后往前截取最近 N 条消息作为"短期工作记忆",并实现孤儿 ToolResult 防线:截断后若首条消息是孤立的工具响应(对应 ToolCall 已被丢弃),自动舍弃防止 API 400
|
||||
- **Reporter 输出抽象** — 新增 `Reporter` 接口(`OnThinking` / `OnToolCall` / `OnToolResult` / `OnMessage`),将引擎输出与展现层解耦。`TerminalReporter` 是首个实现,引擎不再直接 `fmt.Printf`
|
||||
- **WorkDir 从 Engine 下沉到 Session** — 引擎不再持有工作目录,WorkDir 跟随 Session 走。一个引擎实例可同时服务多个不同工作区的会话(多工作区复用单引擎)
|
||||
- **Run 签名重构** — `Run(ctx, userPrompt)` → `Run(ctx, session, reporter)`,会话成为一等公民
|
||||
- **引擎循环改造** — 每轮从 Session 的 Working Memory 构建上下文,工具执行结果实时 `Append` 回 Session,ReAct 循环结束后挂起等待人类下一条指令
|
||||
|
||||
#### 踩坑记录
|
||||
|
||||
| 问题 | 原因 | 解决 |
|
||||
|---|---|---|
|
||||
| 多会话并发操作同一目录文件冲突 | 两个 Session 的 WorkDir 相同,工具同时读写 | WorkDir 绑定 Session,不同会话指向不同工作区 |
|
||||
| 截断 Working Memory 后 API 报 400 | 丢弃了携带 ToolCall 的 Assistant 消息,但留下了对应的 ToolResult | `GetWorkingMemory` 检测并丢弃首部的孤儿 ToolResult |
|
||||
| `fmt.Printf` 无法适配飞书/钉钉/WebUI 等输出目标 | 引擎与终端输出硬耦合 | 抽象 Reporter 接口,`TerminalReporter` 仅为首个实现 |
|
||||
|
||||
#### 经验教训
|
||||
|
||||
1. **WorkDir 属于会话而非引擎** — 将工作目录从 Engine 移到 Session,一个引擎实例就能同时服务多个隔离的工作区(`project_front` / `project_back`),架构不变代码不变
|
||||
2. **Working Memory 不是全量历史** — 大模型 API 有 context window 上限,截取最近 N 条消息既控制成本又保持对话连贯。截断时必须保证 ToolCall / ToolResult 成对存在,否则 API 直接报错
|
||||
3. **Reporter 是引擎可移植的关键** — 引擎只负责"推理 + 调工具",不关心输出到哪里。CLI、飞书、WebUI 只需各自实现 Reporter 接口,引擎零改动
|
||||
|
||||
### v1.5 — Edit 工具 + 四级容错替换 + 多工具并发执行
|
||||
|
||||
#### 变更
|
||||
|
||||
+51
-26
@@ -4,16 +4,17 @@ import (
|
||||
"context"
|
||||
"log"
|
||||
"os"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"go-tiny-claw/internal/engine"
|
||||
"go-tiny-claw/internal/provider"
|
||||
"go-tiny-claw/internal/schema"
|
||||
"go-tiny-claw/internal/tools"
|
||||
)
|
||||
|
||||
func main() {
|
||||
|
||||
workDir, _ := os.Getwd()
|
||||
|
||||
// 1. 初始化真实的 Provider大脑
|
||||
// 这里你可以任意切换 NewZhipuClaudeProvider 或 NewZhipuOpenAIProvider,效果完全一致!
|
||||
llmProvider := provider.DeepseekOpenAIProvider("deepseek-v4-flash")
|
||||
@@ -21,34 +22,58 @@ func main() {
|
||||
registry := tools.NewRegistry()
|
||||
|
||||
// 挂载工具全家桶
|
||||
registry.Register(tools.NewReadFileTool(workDir))
|
||||
registry.Register(tools.NewWriteFileTool(workDir))
|
||||
registry.Register(tools.NewBashTool(workDir))
|
||||
registry.Register(tools.NewEditFileTool(workDir))
|
||||
registry.Register(tools.NewReadSkillTool(workDir))
|
||||
// registry.Register(tools.NewReadFileTool(workDir))
|
||||
// registry.Register(tools.NewWriteFileTool(workDir))
|
||||
// registry.Register(tools.NewBashTool(workDir))
|
||||
// registry.Register(tools.NewEditFileTool(workDir))
|
||||
// registry.Register(tools.NewReadSkillTool(workDir))
|
||||
|
||||
workDir, _ := os.Getwd()
|
||||
|
||||
registry.Register(tools.NewReadFileTool(workDir + "/tmp/project_front"))
|
||||
|
||||
// 实例化引擎,开启 EnableThinking = true
|
||||
eng := engine.NewAgentEngine(llmProvider, registry, workDir, false)
|
||||
eng := engine.NewAgentEngine(llmProvider, registry, false)
|
||||
reporter := engine.NewTerminalReporter()
|
||||
|
||||
// 发起一个需要局部修改的指令
|
||||
//prompt := `
|
||||
//我当前目录下有一个 server.go 文件。
|
||||
//请帮我把里面 "TODO: 增加鉴权逻辑" 下面的那个 if 语句,整个替换为:
|
||||
//if user == nil {
|
||||
// fmt.Println("Forbidden!")
|
||||
// return
|
||||
//}
|
||||
//`
|
||||
var wg sync.WaitGroup
|
||||
|
||||
//prompt := `
|
||||
//我当前目录下有 a.txt, b.txt, c.txt 三个文件。
|
||||
//为了节省时间,请你同时一次性读取这三个文件,并将它们的内容综合起来,告诉我它们分别记录了什么领域的信息。
|
||||
//`
|
||||
// ================= 模拟并发场景 1:飞书前端群 =================
|
||||
wg.Add(1)
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
sessionA := engine.GlobalSessionMgr.GetOrCreate("chat_front_001", workDir+"/tmp/project_front")
|
||||
|
||||
prompt := `我需要在当前目录下新建一个 ping.go,提供一个简单的 http ping 接口。写完之后,帮我把代码用 git 提交一下。`
|
||||
// 回合 1:获取机密
|
||||
log.Println("\n>>> 🙋♂️ [Session A / Turn 1]: 帮我看看 README.md 里记录了什么密钥?我的操作系统是windows")
|
||||
sessionA.Append(schema.Message{Role: schema.RoleUser, Content: "帮我看看 README.md 里记录了什么密钥?"})
|
||||
_ = eng.Run(context.Background(), sessionA, reporter)
|
||||
|
||||
err := eng.Run(context.Background(), prompt)
|
||||
if err != nil {
|
||||
log.Fatalf("引擎运行崩溃: %v", err)
|
||||
}
|
||||
// 故意制造大量“废话”对话,刷掉记忆 (假设 Working Memory Limit=6)
|
||||
for i := 0; i < 6; i++ {
|
||||
sessionA.Append(schema.Message{Role: schema.RoleUser, Content: "这只是一句闲聊占位符。"})
|
||||
sessionA.Append(schema.Message{Role: schema.RoleAssistant, Content: "好的,收到闲聊。"})
|
||||
}
|
||||
|
||||
// 回合 2:验证记忆截断 (此时第一轮的密钥已经被挤出 Working Memory 了!)
|
||||
log.Println("\n>>> 🙋♂️ [Session A / Turn 2]: 请直接告诉我,刚才第一轮你查到的那个密钥是什么?")
|
||||
sessionA.Append(schema.Message{Role: schema.RoleUser, Content: "请直接告诉我,刚才第一轮你查到的那个密钥是什么?不准调用工具!"})
|
||||
_ = eng.Run(context.Background(), sessionA, reporter)
|
||||
}()
|
||||
|
||||
// ================= 模拟并发场景 2:飞书后端群 =================
|
||||
wg.Add(1)
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
// 稍微错开一点时间发起请求
|
||||
time.Sleep(1 * time.Second)
|
||||
|
||||
sessionB := engine.GlobalSessionMgr.GetOrCreate("chat_back_002", workDir+"/tmp/project_back")
|
||||
|
||||
log.Println("\n>>> 🙋♂️ [Session B]: 别人查到了一个密钥,你这里能看到吗?")
|
||||
sessionB.Append(schema.Message{Role: schema.RoleUser, Content: "别人查到了一个密钥,你这里能看到吗?不准调用工具!"})
|
||||
_ = eng.Run(context.Background(), sessionB, reporter)
|
||||
}()
|
||||
|
||||
wg.Wait()
|
||||
}
|
||||
|
||||
+181
-94
@@ -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
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
@@ -0,0 +1,101 @@
|
||||
package prompt
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"log"
|
||||
|
||||
"go-tiny-claw/internal/schema"
|
||||
)
|
||||
|
||||
// Compactor 负责监控和压缩上下文内存,防止大模型发生 OOM
|
||||
type Compactor struct {
|
||||
MaxChars int // 触发压缩的最大字符数阈值 (水位线,可参考使用的大模型的token窗口大小)
|
||||
RetainLastMsgs int // Working Memory 保护区:最近的 N 条消息
|
||||
}
|
||||
|
||||
func NewCompactor(maxChars int, retainLastMsgs int) *Compactor {
|
||||
return &Compactor{
|
||||
MaxChars: maxChars,
|
||||
RetainLastMsgs: retainLastMsgs,
|
||||
}
|
||||
}
|
||||
|
||||
// Compact 接收准备发送给大模型的消息数组。
|
||||
// 如果总长度超标,对远期历史区进行全量掩码 (Masking),对短期保护区进行超长局部截断 (Truncation)。
|
||||
func (c *Compactor) Compact(msgs []schema.Message) []schema.Message {
|
||||
currentLength := c.estimateLength(msgs)
|
||||
|
||||
// 如果没有超过水位线,直接返回原数组 (大多数情况下的正常路径)
|
||||
if currentLength < c.MaxChars {
|
||||
return msgs
|
||||
}
|
||||
|
||||
log.Printf("[Compactor] ⚠️ 内存告警:当前上下文长度 (%d 字符) 超过阈值 (%d),触发压缩清理...\n", currentLength, c.MaxChars)
|
||||
|
||||
var compacted []schema.Message
|
||||
msgCount := len(msgs)
|
||||
|
||||
// 计算受保护的 Working Memory 起始索引
|
||||
protectStartIndex := msgCount - c.RetainLastMsgs
|
||||
if protectStartIndex < 0 {
|
||||
protectStartIndex = 0
|
||||
}
|
||||
|
||||
for i, msg := range msgs {
|
||||
// 1. 系统提示词 (System Prompt) 绝对不能动,直接保留
|
||||
if msg.Role == schema.RoleSystem {
|
||||
compacted = append(compacted, msg)
|
||||
continue
|
||||
}
|
||||
|
||||
// 我们必须拷贝一份新消息,因为在并发环境中直接修改原引用可能导致底层数据结构被污染
|
||||
newMsg := msg
|
||||
|
||||
isInWorkingMemory := i >= protectStartIndex
|
||||
|
||||
// 【核心驾驭逻辑】: 双重降级防线
|
||||
if msg.Role == schema.RoleUser && msg.ToolCallID != "" {
|
||||
// 对于工具的返回结果 (Observation/ToolResult)
|
||||
if !isInWorkingMemory {
|
||||
// 【第一道防线:远期历史】如果是早期对话,执行无情替换 (Full Masking)
|
||||
if len(msg.Content) > 200 {
|
||||
newMsg.Content = fmt.Sprintf("...[为了节省内存,早期的工具输出已被系统强制清理。原始长度: %d 字节]...", len(msg.Content))
|
||||
}
|
||||
} else {
|
||||
// 【第二道防线:短期记忆】即使处于近期保护区,只要单条内容过大,也必须截断防 OOM (Head-Tail Truncation)
|
||||
// 我们保留前 500 字符和后 500 字符(掐头去尾法,大模型通常只需要看开头报错和结尾总结)
|
||||
const maxKeep = 1000
|
||||
if len(msg.Content) > maxKeep {
|
||||
head := msg.Content[:500]
|
||||
tail := msg.Content[len(msg.Content)-500:]
|
||||
newMsg.Content = fmt.Sprintf("%s\n\n...[内容过长,中间 %d 字节已被系统截断]...\n\n%s", head, len(msg.Content)-maxKeep, tail)
|
||||
}
|
||||
}
|
||||
} else if msg.Role == schema.RoleAssistant && msg.Content != "" {
|
||||
// 对于大模型的冗长推理废话 (Thinking Trace)
|
||||
if !isInWorkingMemory && len(msg.Content) > 200 {
|
||||
newMsg.Content = "...[早期的推理思考过程已折叠]..."
|
||||
}
|
||||
}
|
||||
|
||||
// 注意:我们绝不会去动 msg.ToolCalls,因为这是模型行动的证据,是维系逻辑链的关键!
|
||||
compacted = append(compacted, newMsg)
|
||||
}
|
||||
|
||||
newLength := c.estimateLength(compacted)
|
||||
log.Printf("[Compactor] ✅ 压缩完成。上下文长度从 %d 降至 %d 字符。\n", currentLength, newLength)
|
||||
|
||||
return compacted
|
||||
}
|
||||
|
||||
// estimateLength 粗略计算当前上下文的总字符长度
|
||||
func (c *Compactor) estimateLength(msgs []schema.Message) int {
|
||||
length := 0
|
||||
for _, msg := range msgs {
|
||||
length += len(msg.Content)
|
||||
for _, tc := range msg.ToolCalls {
|
||||
length += len(tc.Name) + len(tc.Arguments)
|
||||
}
|
||||
}
|
||||
return length
|
||||
}
|
||||
Reference in New Issue
Block a user