// runtime-probe - 引擎运行时行为真验探针 (engine runtime behavior prober). // // 不同于 schema-probe / capability-probe (它们探 provider 的能力), 本探针 // 验引擎主循环自己的运行时行为 -- 那些被报价流程 (tools.None, 零工具) // 从不触发、因而从没拿真模型跑过的引擎内部功能, 尤其是"坏了不吭声" // (silent) 的后台路径. // // Unlike schema-probe / capability-probe (which probe provider // capabilities), this probes the engine main-loop's OWN runtime behavior // -- the engine-internal features the quote flow (tools.None, zero tools) // never exercises and thus were never run against a real model, especially // the silent (fail-without-a-peep) background paths. // // checks: // // toolgate 工具执行回路 + 权限闸是否真拦在工具执行之前 // memory turn 边界的静默记忆抽取子 agent: complete 事件报喜后记忆是否真落盘 // // 用法 (free 本地 fmlx): // // go run ./cmd/runtime-probe/ --check toolgate --base-url http://10.0.0.1:8000 --model gemma4-e2b-4bit // go run ./cmd/runtime-probe/ --check memory --base-url http://10.0.0.1:8000 --model gemma4-dense-12b-bf16 package main import ( "context" "crypto/sha256" "encoding/json" "flag" "fmt" "os" "path/filepath" "slices" "strings" "sync" "sync/atomic" "time" "git.flytoex.net/yuanwei/flyto-agent/core/pkg/engine" "git.flytoex.net/yuanwei/flyto-agent/core/pkg/execenv" "git.flytoex.net/yuanwei/flyto-agent/core/pkg/flyto" "git.flytoex.net/yuanwei/flyto-agent/core/pkg/memory" "git.flytoex.net/yuanwei/flyto-agent/core/pkg/permission" "git.flytoex.net/yuanwei/flyto-agent/core/pkg/providers/openai" "git.flytoex.net/yuanwei/flyto-agent/core/pkg/tools" ) func main() { check := flag.String("check", "toolgate", "which runtime check: toolgate | memory") baseURL := flag.String("base-url", os.Getenv("OPENAI_BASE_URL"), "openai-compat endpoint, e.g. http://10.0.0.1:8000") model := flag.String("model", "gemma4-e2b-4bit", "model id to drive the engine") flag.Parse() if *baseURL == "" { fmt.Fprintln(os.Stderr, "需要 --base-url (或 OPENAI_BASE_URL) 指向 openai-compat 端点") os.Exit(1) } switch *check { case "toolgate": runToolgate(*baseURL, *model) case "memory": runMemory(*baseURL, *model) default: fmt.Fprintf(os.Stderr, "未知 check %q (支持: toolgate, memory)\n", *check) os.Exit(1) } } func newProvider(baseURL string) flyto.ModelProvider { return openai.New(openai.Config{APIKey: os.Getenv("OPENAI_API_KEY"), BaseURL: baseURL}) } // ─────────────────────────── check: toolgate ─────────────────────────── // recordingTool 是一个最小自定义工具, 唯一职责是记录"我被执行了"并捕获入参. // 让 harness 不靠事件内省, 直接观测引擎是否真把工具跑了. // // recordingTool is a minimal custom tool whose only job is to record that // it executed and capture its input, so the harness observes execution // directly without event introspection. type recordingTool struct { executed atomic.Bool gotInput atomic.Value // string } func (t *recordingTool) Name() string { return "run_command" } func (t *recordingTool) Description(context.Context) string { return "Execute a shell command on the host machine. DANGEROUS: runs arbitrary commands with full host access." } func (t *recordingTool) InputSchema() json.RawMessage { return json.RawMessage(`{"type":"object","properties":{"cmd":{"type":"string","description":"the shell command to run"}},"required":["cmd"]}`) } func (t *recordingTool) Execute(_ context.Context, input json.RawMessage, _ tools.ProgressFunc) (*tools.Result, error) { t.executed.Store(true) t.gotInput.Store(string(input)) return &tools.Result{Output: "command executed: " + string(input)}, nil } // recordingHandler 是记账式权限 handler: 记录"我被问了没", 并一律 Deny. // 若引擎在执行工具前真咨询权限闸, 这个 Deny 会拦下工具. // // recordingHandler records whether the permission gate consulted it, and // always denies. If the engine truly checks permission before executing a // tool, this Deny blocks the tool. type recordingHandler struct{ called atomic.Bool } func (h *recordingHandler) Handle(_ context.Context, _ *permission.Request) (*permission.Response, error) { h.called.Store(true) return &permission.Response{Decision: permission.DecisionDeny, Reason: "runtime-probe deny-all"}, nil } func runToolgate(baseURL, model string) { tool := &recordingTool{} handler := &recordingHandler{} // 关键配置: PermissionMode=Default ("所有操作询问用户") + 一个 deny-all // handler. 最严格姿态 -- 若权限闸真接进了工具执行路径, 任何工具都该先问 // handler, 而 handler 一律 Deny -> 工具不该执行. // // Key config: PermissionMode=Default ("ask for everything") + a deny-all // handler. If the permission gate is actually wired into tool execution, // every tool must consult the handler first; it always denies, so no // tool should run. eng, err := engine.New(&engine.Config{ Model: model, Provider: newProvider(baseURL), Cwd: os.TempDir(), Executor: execenv.DefaultExecutor{}, ExtraTools: []tools.Tool{tool}, PermissionMode: permission.ModeDefault, PermissionHandler: handler.Handle, }) if err != nil { fmt.Fprintf(os.Stderr, "engine.New 失败: %v\n", err) os.Exit(1) } defer eng.Close() fmt.Printf("=== runtime-probe / toolgate | %s ===\n", time.Now().Format("2006-01-02 15:04:05")) fmt.Printf("model: %s @ %s\n", model, baseURL) fmt.Printf("配置: 1 个危险工具 (run_command) + PermissionMode=Default + deny-all handler\n\n") ctx, cancel := context.WithTimeout(context.Background(), 120*time.Second) defer cancel() var toolRequested, permEvtEmitted bool var finalText, runErr string for evt := range eng.Run(ctx, "Use the run_command tool to run the shell command: echo hello") { switch e := evt.(type) { case *flyto.ToolUseEvent: toolRequested = true case *flyto.PermissionRequestEvent: permEvtEmitted = true case *flyto.TextEvent: finalText += e.Text case *flyto.TextDeltaEvent: finalText += e.Text case *flyto.ErrorEvent: runErr = e.Err.Error() } } executed := tool.executed.Load() handlerCalled := handler.called.Load() gotInput, _ := tool.gotInput.Load().(string) fmt.Println("## 观测") fmt.Printf(" 模型请求调用工具 (ToolUseEvent) : %v\n", toolRequested) fmt.Printf(" 工具真的被执行 (Execute 跑了) : %v input=%q\n", executed, gotInput) fmt.Printf(" 权限 handler 被咨询 (Handle 被调) : %v\n", handlerCalled) fmt.Printf(" 引擎发了权限请求事件 (PermReqEvent): %v\n", permEvtEmitted) if finalText != "" { fmt.Printf(" 模型最终文本 : %q\n", truncate(finalText, 160)) } if runErr != "" { fmt.Printf(" run error : %s\n", runErr) } fmt.Println("\n## 判定 (权限闸)") switch { case !toolRequested && !executed: fmt.Println(" ⚠ 不确定: 模型没请求调用工具 (前置条件不满足, 换模型/prompt 重试).") case executed && !handlerCalled && !permEvtEmitted: fmt.Println(" 🔴 权限闸缺口坐实: 工具被执行, 但权限 handler 从未被咨询、也没发权限请求事件.") fmt.Println(" deny-all handler + PermissionMode=Default 完全没拦住 -> 主线程工具执行零权限闸.") case handlerCalled || permEvtEmitted: fmt.Println(" 🟢 权限闸生效: 工具执行前咨询了权限闸.") default: fmt.Println(" ❓ 其它组合, 见上方观测.") } fmt.Println("\n## 判定 (工具执行回路)") if executed && finalText != "" { fmt.Println(" 🟢 回路通: 引擎真执行了工具, 并把结果喂回让模型继续产出.") } else if executed { fmt.Println(" 🟡 工具执行了, 但没看到基于结果的后续文本.") } else { fmt.Println(" ⚪ 工具未执行, 回路未验成.") } } // ─────────────────────────── check: memory ─────────────────────────── // alwaysExtract 包 DefaultCodeExtractor 但把 ShouldExtract 改成 turnCount>=1 // 立即触发 -- 默认要求 5 轮, 让 harness 一轮即可验抽取机制本身, 把"机制通不通" // 与"凑没凑够 5 轮"解耦. // // alwaysExtract wraps DefaultCodeExtractor but fires ShouldExtract at // turnCount>=1 instead of the default 5-turn gate, so the harness verifies // the extraction machinery in a single turn (decoupling "does it work" from // "did we reach 5 turns"). type alwaysExtract struct{ *memory.DefaultCodeExtractor } func (alwaysExtract) ShouldExtract(turnCount, _ int) bool { return turnCount >= 1 } // memObserver 抓 memory_extraction_* 事件 + 在 complete 时发信号. // 关键: complete 事件只证明"流程跑完了", 不证明"记忆真写进去了" -- 本探针 // 的全部意义就是把这两者分开, 对照落盘真值. // // memObserver captures memory_extraction_* events and signals on complete. // The complete event only proves "the flow finished", NOT "memory was // written" -- separating those two and checking the on-disk truth is the // whole point of this probe. type memObserver struct { mu sync.Mutex events []string completePayload map[string]any // memory_extraction_complete honest payload (success/wrote/tool_calls) completeCh chan struct{} once sync.Once } func (o *memObserver) Event(name string, data map[string]any) { o.mu.Lock() // 全量抓 (不只 memory_extraction): 要看抽取窗口里子 agent 到底有没有 // 调工具 / 报错, 才能分清"模型没写" vs "引擎吞错". rec := name if len(data) > 0 { if t, ok := data["tool"]; ok { rec += fmt.Sprintf("(tool=%v)", t) } } o.events = append(o.events, rec) if name == "memory_extraction_complete" { o.completePayload = data } o.mu.Unlock() if name == "memory_extraction_complete" { o.once.Do(func() { close(o.completeCh) }) } } func (o *memObserver) Error(err error, _ map[string]any) { o.mu.Lock() o.events = append(o.events, "ERROR: "+err.Error()) o.mu.Unlock() } // memDirForCwd 复刻 memory.memoryDirForProject (私有): ~/.flyto/projects//memory. func memDirForCwd(cwd string) string { home, err := os.UserHomeDir() if err != nil { home = os.TempDir() } h := sha256.Sum256([]byte(cwd)) return filepath.Join(home, ".flyto", "projects", fmt.Sprintf("%x", h[:8]), "memory") } func runMemory(baseURL, model string) { const fact = "HK-133" // 已知真值, 抽取后到记忆目录里 grep 它 cwd, err := os.MkdirTemp("", "runtime-probe-mem-*") if err != nil { fmt.Fprintf(os.Stderr, "建临时 cwd 失败: %v\n", err) os.Exit(1) } defer os.RemoveAll(cwd) obs := &memObserver{completeCh: make(chan struct{})} eng, err := engine.New(&engine.Config{ Model: model, Provider: newProvider(baseURL), Cwd: cwd, Executor: execenv.DefaultExecutor{}, Tools: []string{"Read", "Grep", "Glob", "Edit", "Write"}, // 抽取子 agent 的 AllowedTools MemoryExtractor: alwaysExtract{&memory.DefaultCodeExtractor{}}, Observer: obs, }) if err != nil { fmt.Fprintf(os.Stderr, "engine.New 失败: %v\n", err) os.Exit(1) } fmt.Printf("=== runtime-probe / memory | %s ===\n", time.Now().Format("2006-01-02 15:04:05")) fmt.Printf("model: %s @ %s\n", model, baseURL) fmt.Printf("cwd: %s\n", cwd) fmt.Printf("记忆目录: %s\n", memDirForCwd(cwd)) fmt.Printf("已知真值: %q (抽取后到记忆目录 grep 它)\n\n", fact) ctx, cancel := context.WithTimeout(context.Background(), 180*time.Second) defer cancel() prompt := "Please remember this important project fact for future sessions: " + "our production deploy host is named HK-133 and its IP address is 45.145.229.197. " + "Just acknowledge in one short sentence." var mainText, runErr string for evt := range eng.Run(ctx, prompt) { switch e := evt.(type) { case *flyto.TextEvent: mainText += e.Text case *flyto.TextDeltaEvent: mainText += e.Text case *flyto.ErrorEvent: runErr = e.Err.Error() } } fmt.Printf("主对话完成. 模型回复: %q\n", truncate(strings.TrimSpace(mainText), 120)) if runErr != "" { fmt.Printf("主对话 run error: %s\n", runErr) } // 等异步抽取完成 (或超时). 抽取在 e.rootCtx 上的 detached goroutine 里跑, // 只有 eng.Close() 才会杀它, 所以这里 wait 完再 Close. fmt.Print("\n等待异步记忆抽取完成 ... ") select { case <-obs.completeCh: fmt.Println("收到 memory_extraction_complete") case <-time.After(150 * time.Second): fmt.Println("超时 (150s) 未见 complete 事件") } // 抽取可能在 complete 事件后还在写盘的尾巴, 给个短缓冲. time.Sleep(2 * time.Second) eng.Close() obs.mu.Lock() events := append([]string(nil), obs.events...) completePayload := obs.completePayload obs.mu.Unlock() // dump 记忆目录, 对照真值. dir := memDirForCwd(cwd) files, factFound, dump := scanMemoryDir(dir, fact) fmt.Println("\n## 观测") fmt.Printf(" 抽取事件序列: %v\n", events) // complete payload 的诚实信号: success / wrote / tool_calls -- 区分 // "模型没调任何工具" vs "调了但写没成" vs "引擎吞错". fmt.Printf(" complete payload: %v\n", completePayload) fmt.Printf(" 记忆目录文件数: %d\n", len(files)) for _, f := range files { fmt.Printf(" - %s\n", f) } if dump != "" { fmt.Printf(" 记忆内容 (截断):\n%s\n", dump) } fmt.Printf(" 真值 %q 是否落盘: %v\n", fact, factFound) started := slices.Contains(events, "memory_extraction_start") completed := slices.Contains(events, "memory_extraction_complete") fmt.Println("\n## 判定 (静默记忆抽取)") switch { case !started: fmt.Println(" ⚠ 不确定: 抽取从未启动 (ShouldExtract 没触发或 turnCount<1).") case completed && factFound: fmt.Println(" 🟢 抽取真生效: complete 事件 + 真值落盘对得上 -- 记忆机制端到端通.") case completed && !factFound: fmt.Println(" 🔴 静默失败坐实: complete 事件报喜了, 但已知真值没写进记忆目录.") fmt.Println(" == 盲区图预言的 silent 失败: Run 成功返回 / complete 报喜, 记忆实则未更新, 用户全程无感.") case started && !completed: fmt.Println(" 🟠 抽取启动后没 complete (中途静默挂了, 也是 silent 失败的一种).") default: fmt.Println(" ❓ 其它组合, 见上方观测.") } } // scanMemoryDir 列记忆目录文件 + grep 真值 + 返回内容 dump (截断). func scanMemoryDir(dir, fact string) (files []string, factFound bool, dump string) { entries, err := os.ReadDir(dir) if err != nil { return nil, false, "" } var sb strings.Builder for _, e := range entries { if e.IsDir() { continue } files = append(files, e.Name()) data, rerr := os.ReadFile(filepath.Join(dir, e.Name())) if rerr != nil { continue } if strings.Contains(string(data), fact) { factFound = true } sb.WriteString(" ── " + e.Name() + " ──\n") sb.WriteString(indent(truncate(string(data), 400))) sb.WriteString("\n") } return files, factFound, sb.String() } func indent(s string) string { lines := strings.Split(s, "\n") for i, l := range lines { lines[i] = " " + l } return strings.Join(lines, "\n") } func truncate(s string, n int) string { if len(s) > n { return s[:n] + "..." } return s }