diff --git a/meta/meta.go b/meta/meta.go index 883bbba..1e0619c 100644 --- a/meta/meta.go +++ b/meta/meta.go @@ -1,12 +1,15 @@ // Package meta 收集 HomeAgent SDK 的全部元数据。 // 版本号应与核心 meta.Version 保持一致。 -// ABI 版本与 Dispatch Method ID 应与核心仓 internal/meta/meta.go 保持一致。 package meta var ( // Version 是 HomeAgent SDK 版本号。 // 通过 `-ldflags="-X gitcode.com/JianFeeeee/homeagent-sdk/meta.Version=vX.Y.Z"` 注入。 - Version = "0.9.2" + // + // 1.0.0:插件运行模型从 C ABI 动态库改为子进程 + 共享内存。 + // 公开 SDK 接口(sdk/ 目录)**零改动**——插件业务代码不需要改一行, + // 但产物形态变了(plugin.so → plugin.bin),必须用新版 plugindev 重编。 + Version = "1.0.0" // Commit 是构建时的 Git commit hash。 Commit = "unknown" @@ -21,7 +24,10 @@ var ( CoreModule = "gitcode.com/JianFeeeee/HomeAgent" // CoreVersion 是此 SDK 所兼容的最低核心版本。 - CoreVersion = "0.9.2" + // + // 1.0.0 是硬下限而非建议值:0.9.x 内核只会 dlopen `.so`, + // 本版工具链产出的 `plugin.bin` 在旧内核上根本不会被识别。 + CoreVersion = "1.0.0" ) // FullVersion 返回完整的版本字符串。 @@ -29,79 +35,15 @@ func FullVersion() string { return SDKName + " v" + Version + " (" + Commit + ")" } -// ---- ABI 版本(与核心仓 internal/meta/meta.go 同步) ---- -// ABI 标识版本直接取内核版本号字符串(semver),与核心 Version 保持一致,不使用独立数字编码。 -// 协商层(C 结构体 int version 字段)使用 CABINum:由版本字符串派生的整数(major*100 + minor)。 -// 映射:v0.8.x → CABINum=800;v0.9.x → CABINum=900(invoke_stage 写回)。 -// 小版本(patch)演进不影响 ABI,CABINum 不变。version_min 保证旧 ABI 插件仍可加载。 - -var ( - // ABIVersion 是 ABI 标识版本(字符串 semver,与 SDK CoreVersion 对齐)。 - ABIVersion = CoreVersion - // ABIVersionMin 是兼容的最低 ABI 标识版本。 - ABIVersionMin = "0.8.0" -) - -const ( - // CABINum 是 C 层协商用的整数版本(major*100 + minor),随 ABIVersion 派生。 - CABINum = 900 - // CABINumMin 是 C 层兼容的最低整数版本。 - // 旧工具链(v0.8 之前)写入的整数 version=1,无写回能力但与新内核结构兼容, - // 因此最小值保持 1 以兼容全部旧插件(新插件 900 匹配,旧插件 1/2 通过); - // 仅当未来内核 ABI 破坏兼容时才提高该值。 - CABINumMin = 1 -) - -// ---- Dispatch Method IDs(与核心仓 internal/meta/meta.go 同步) ---- -const ( - CoreRegisterTool = 1 - CoreRegisterStage = 2 - CoreRegisterOutputCh = 3 - CoreRegisterPluginAPI = 4 - CoreInjectText = 5 - CoreInjectInterruptText = 6 - CoreInjectTextNoMemory = 7 - CoreSetAutoRestart = 8 - CoreMemoryRecall = 9 - CoreMemoryCommit = 10 - CoreMemoryIntrospect = 11 - CoreMemoryMerge = 12 - CoreMemoryPurge = 13 - CoreDocQuery = 14 - CoreKnowledgeSearch = 15 - CoreSettingsGet = 16 - CoreSettingsSet = 17 - CoreSettingsRegisterDef = 18 - CoreLLMListSources = 19 - CoreLLMSetSource = 20 - CoreSocialGetPerson = 21 - CoreSocialGetNetwork = 22 - CoreSubscribe = 23 - CoreUnsubscribe = 24 - CoreFreeString = 25 - CoreSettingsGetCore = 26 - CoreSettingsSetCore = 27 - CoreSettingsListCore = 28 - CoreSettingsGetPlugin = 29 - CoreSettingsSetPlugin = 30 - CoreSettingsListPlugin = 31 - CoreDocInsert = 32 - CoreDocRemove = 33 - CoreDocStats = 34 - CoreKnowledgeAdd = 35 - CoreKnowledgeList = 36 - CoreLLMCurrentSource = 37 - CoreSocialGetTrait = 38 - CoreSocialGetRelations = 39 - CoreSocialListPersons = 40 - CoreTextMemoryAppend = 41 - CoreSettingsList = 42 - CoreSettingsDefs = 43 - CoreSettingsDump = 44 - CoreSettingsPlugins = 45 - CoreRegisterInputCh = 46 - CoreInjectInputSync = 47 - CorePluginReloadOne = 48 - CorePluginListLoaded = 49 - CorePluginIsDisabled = 50 -) +// ---- 协议版本 ---- +// +// 子进程 RPC 的协议版本是一个独立的小整数,与 SDK/内核语义版本解耦: +// 语义版本变动频繁(修 bug、加字段),而 wire 协议只在**帧格式或握手语义** +// 变化时才升。当前值见核心仓 internal/plugin/proc/protocol.go 的 ProtocolVersion。 +// +// C ABI 时代的 ABIVersion / CABINum / 51 个 Core 整数 ID 已随 +// Part 6.2 删除 internal/plugin/cabi/ 一并退场: +// - 整数 method id 平移为 method 名字符串(proc/protocol.go 的 Method* 常量) +// - 版本协商改为握手帧里的 protocol 字段 +// +// 保留那些常量只会让人以为它们还在生效。 diff --git a/tools/plugindev/cmd_build.go b/tools/plugindev/cmd_build.go index 5cbfd70..2e974de 100644 --- a/tools/plugindev/cmd_build.go +++ b/tools/plugindev/cmd_build.go @@ -24,7 +24,8 @@ func cmdBuild(args []string) { // Read all config from plg.json first plg, err := readPlgJSON("plg.json") if err != nil { - fmt.Printf("error: read plg.json: %v\n", err); os.Exit(1) + fmt.Printf("error: read plg.json: %v\n", err) + os.Exit(1) } // Base config from plg.json @@ -39,11 +40,13 @@ func cmdBuild(args []string) { switch args[i] { case "--outdir": if i+1 < len(args) { - outDir = args[i+1]; i++ + outDir = args[i+1] + i++ } case "--target": if i+1 < len(args) { - targets = append(targets, args[i+1]); i++ + targets = append(targets, args[i+1]) + i++ } case "--bundle": bundle = true @@ -51,11 +54,13 @@ func cmdBuild(args []string) { bundle = false case "--sdk-path": if i+1 < len(args) { - sdkPath = args[i+1]; i++ + sdkPath = args[i+1] + i++ } case "--replace", "-R": if i+1 < len(args) { - cliReplaces = append(cliReplaces, args[i+1]); i++ + cliReplaces = append(cliReplaces, args[i+1]) + i++ } } } @@ -105,14 +110,17 @@ func cmdBuild(args []string) { } // allBundleTargets 是 --bundle 模式构建的全部平台。 -// 每个 OS 只有一个架构(amd64),避免二进制文件名冲突。 +// +// 子进程模式下各平台产物同名(plugin.bin)——进程边界即 ABI 边界, +// 不存在平台特有扩展名,故 zip 内按平台加后缀区分; +// 内核安装时按当前平台挑对应条目重命名为 plugin.bin。 var allBundleTargets = []struct { target string entry string // 二进制在 zip 中的文件名 }{ - {"linux/amd64", "plugin.so"}, - {"darwin/amd64", "plugin.dylib"}, - {"windows/amd64", "plugin.dll"}, + {"linux/amd64", "plugin.bin.linux.amd64"}, + {"darwin/amd64", "plugin.bin.darwin.amd64"}, + {"windows/amd64", "plugin.bin.windows.amd64"}, } func buildBundle(plg *PlgConfig, outDir string, sdkPath string) { @@ -120,9 +128,13 @@ func buildBundle(plg *PlgConfig, outDir string, sdkPath string) { buildDir := "build" os.MkdirAll(buildDir, 0755) - // Auto-generate C ABI bridge for non-Windows - bridgeCleanup := generateBridge("") - defer bridgeCleanup() + runtimeCleanup, err := generateProcRuntime() + if err != nil { + fmt.Printf(" error: %v\n", err) + return + } + defer runtimeCleanup() + thirdpartCleanup := linkThirdpart(plg, "linux/amd64") defer thirdpartCleanup() @@ -135,22 +147,18 @@ func buildBundle(plg *PlgConfig, outDir string, sdkPath string) { return } - outPath := filepath.Join(buildDir, cfg.entryFile) + // 每平台产物落到独立路径,避免相互覆盖 + outName := fmt.Sprintf("%s_%s_%s", cfg.entryFile, cfg.goos, cfg.goarch) + outPath := filepath.Join(buildDir, outName) - cmd := exec.Command("go", "build", "-buildmode=c-shared", "-o", outPath) + // 零 cgo:跨平台交叉编译不需目标平台 C 工具链 + cmd := exec.Command("go", "build", "-trimpath", "-o", outPath) cmd.Env = os.Environ() - cmd.Env = append(cmd.Env, "GOOS="+cfg.goos, "GOARCH="+cfg.goarch, "CGO_ENABLED=1") - - if cfg.goos == "windows" { - cc := detectWindowsCC() - if cc != "" { - cmd.Env = append(cmd.Env, "CC="+cc) - } - } - + cmd.Env = append(cmd.Env, "GOOS="+cfg.goos, "GOARCH="+cfg.goarch, "CGO_ENABLED=0") cmd.Stdout = os.Stdout cmd.Stderr = os.Stderr - fmt.Printf(" compiling %s/%s (-buildmode=c-shared)...\n", cfg.goos, cfg.goarch) + + fmt.Printf(" compiling %s/%s (子进程模式,CGO_ENABLED=0)...\n", cfg.goos, cfg.goarch) if err := cmd.Run(); err != nil { fmt.Printf(" error: build %s/%s: %v\n", cfg.goos, cfg.goarch, err) return @@ -168,7 +176,7 @@ func buildBundle(plg *PlgConfig, outDir string, sdkPath string) { for p := range platforms { plats = append(plats, p) } - writePluginJSON(plg, plats, "plugin.so") + writePluginJSON(plg, plats, procEntryFile) // package single .hmap with correctly named entries hmapPath := filepath.Join(outDir, fmt.Sprintf("%s_bundle.hmap", toSnake(plg.NameEn))) @@ -176,7 +184,10 @@ func buildBundle(plg *PlgConfig, outDir string, sdkPath string) { fmt.Printf(" packaged %s\n", filepath.Base(hmapPath)) } -func (p *PlgConfig) IsLua() bool { return p.Entry == "main.lua" } +// IsLua 判断是否为 Lua 插件(走解释器,不经过 Go 编译)。 +// +// 这是 entry 字段唯一仍在使用的用途:Go 插件不再看 entry 值,一律产出 plugin.bin。 +func (p *PlgConfig) IsLua() bool { return p.Entry == luaEntryFile } func readPlgJSON(path string) (*PlgConfig, error) { data, err := os.ReadFile(path) @@ -228,9 +239,14 @@ func writePluginJSON(plg *PlgConfig, platforms []string, entry string) { type buildConfig struct { goos string goarch string - entryFile string // "plugin.so" or "plugin.dll" + entryFile string // 一律为 plugin.bin(进程边界即 ABI 边界,无平台特有扩展名) } +// resolveBuild 解析目标平台。 +// +// 全平台统一产出 plugin.bin:子进程模式下不存在 .so/.dylib/.dll 的区分, +// 因为进程边界本身就是 ABI 边界——这正是三套独立 ABI 实现收敛为 +// 单一 RPC 实现的直接后果(§9.2:Windows 不再是能力退化的第三套实现)。 func resolveBuild(target string) (*buildConfig, string) { if target == "lua" || target == "" { return nil, "lua" @@ -245,14 +261,8 @@ func resolveBuild(target string) (*buildConfig, string) { } switch goos { - case "linux": - return &buildConfig{goos: goos, goarch: goarch, entryFile: "plugin.so"}, "" - case "darwin": - return &buildConfig{goos: goos, goarch: goarch, entryFile: "plugin.dylib"}, "" - case "freebsd": - return &buildConfig{goos: goos, goarch: goarch, entryFile: "plugin.so"}, "" - case "windows": - return &buildConfig{goos: goos, goarch: goarch, entryFile: "plugin.dll"}, "" + case "linux", "darwin", "freebsd", "windows": + return &buildConfig{goos: goos, goarch: goarch, entryFile: procEntryFile}, "" default: return nil, fmt.Sprintf("unsupported OS %q", goos) } @@ -403,7 +413,7 @@ func buildTarget(plg *PlgConfig, target, outDir, sdkPath string) { return } - // Resolve build config + // Resolve build config(全平台统一产出 plugin.bin) cfg, errMsg := resolveBuild(target) if cfg == nil { fmt.Printf(" error: %s\n", errMsg) @@ -414,9 +424,12 @@ func buildTarget(plg *PlgConfig, target, outDir, sdkPath string) { os.MkdirAll(buildDir, 0755) outPath := filepath.Join(buildDir, cfg.entryFile) - // Auto-generate C ABI bridge (all platforms use c-shared) - bridgeCleanup := generateBridge(cfg.goos) - defer bridgeCleanup() + runtimeCleanup, err := generateProcRuntime() + if err != nil { + fmt.Printf(" error: %v\n", err) + return + } + defer runtimeCleanup() // Auto-link thirdpart/ contents + source_dirs + replace targets thirdpartCleanup := linkThirdpart(plg, target) @@ -425,28 +438,15 @@ func buildTarget(plg *PlgConfig, target, outDir, sdkPath string) { // Write plugin.json with the correct entry for this target writePluginJSON(plg, nil, cfg.entryFile) - cmd := exec.Command("go", "build", "-buildmode=c-shared", "-o", outPath) + // 普通 go build + 零 cgo:交叉编译不再需要目标平台的 C 工具链 + // (旧路径靠 detectWindowsCC 找 MinGW,现在整个问题消失)。 + cmd := exec.Command("go", "build", "-trimpath", "-o", outPath) cmd.Env = os.Environ() - cmd.Env = append(cmd.Env, "GOOS="+cfg.goos, "GOARCH="+cfg.goarch, "CGO_ENABLED=1") - - // Auto-detect MinGW gcc on Windows - if cfg.goos == "windows" { - cc := detectWindowsCC() - if cc != "" { - cmd.Env = append(cmd.Env, "CC="+cc) - } - } - + cmd.Env = append(cmd.Env, "GOOS="+cfg.goos, "GOARCH="+cfg.goarch, "CGO_ENABLED=0") cmd.Stdout = os.Stdout cmd.Stderr = os.Stderr - // DEBUG: list files before building - entries, _ := os.ReadDir(".") - for _, e := range entries { - fmt.Printf(" [DEBUG] file: %s\n", e.Name()) - } - - fmt.Printf(" compiling %s/%s (-buildmode=c-shared)...\n", cfg.goos, cfg.goarch) + fmt.Printf(" compiling %s/%s (子进程模式,CGO_ENABLED=0)...\n", cfg.goos, cfg.goarch) if err := cmd.Run(); err != nil { fmt.Printf(" error: build %s/%s: %v\n", cfg.goos, cfg.goarch, err) return @@ -566,72 +566,8 @@ func toSnake(s string) string { return strings.ToLower(strings.ReplaceAll(s, " ", "_")) } -// detectWindowsCC looks for a MinGW-w64 gcc on Windows for c-shared builds. -func detectWindowsCC() string { - // Check CC from environment first - if cc := os.Getenv("CC"); cc != "" { - if _, err := exec.LookPath(cc); err == nil { - return cc - } - } - // Check common MinGW install paths - candidates := []string{ - "C:\\mingw64\\bin\\gcc.exe", - "C:\\MinGW\\bin\\gcc.exe", - "C:\\msys64\\mingw64\\bin\\gcc.exe", - "C:\\Users\\21989\\AppData\\Local\\Temp\\mingw64\\mingw64\\bin\\gcc.exe", - } - // Also search PATH for gcc - if path, err := exec.LookPath("gcc"); err == nil { - return path - } - for _, c := range candidates { - if _, err := os.Stat(c); err == nil { - return c - } - } - return "" -} - // stripIncludeGuard strips preprocessor guards and C++ comments from a C header, // since these can confuse cgo's type resolution. -// generateBridge generates the C ABI bridge files for non-Lua builds. -// Returns a cleanup function to remove generated files. -func generateBridge(goos string) func() { - const bridgeFile = "z_bridge_gen.go" - const cEntryFile = "z_entry.c" - os.Remove(bridgeFile) - os.Remove(cEntryFile) - - var files []string - - if goos == "windows" { - if err := os.WriteFile(bridgeFile, []byte(tmplBridge), 0644); err != nil { - fmt.Printf(" error: write bridge: %v\n", err) - return func() {} - } - files = append(files, bridgeFile) - } else { - if err := os.WriteFile(bridgeFile, []byte(tmplLinuxBridge), 0644); err != nil { - fmt.Printf(" error: write bridge: %v\n", err) - return func() {} - } - files = append(files, bridgeFile) - // Write C entry point file - if err := os.WriteFile(cEntryFile, []byte(tmplPluginInitC), 0644); err != nil { - fmt.Printf(" error: write C entry: %v\n", err) - return func() {} - } - files = append(files, cEntryFile) - } - - return func() { - for _, f := range files { - os.Remove(f) - } - } -} - // linkThirdpart scans thirdpart/, source_dirs from plg.json, and replace target dirs // for source files, generating auto-import stubs. Returns cleanup function. func linkThirdpart(plg *PlgConfig, target string) func() { diff --git a/tools/plugindev/cmd_init.go b/tools/plugindev/cmd_init.go index cbe2df0..5f03abe 100644 --- a/tools/plugindev/cmd_init.go +++ b/tools/plugindev/cmd_init.go @@ -7,8 +7,6 @@ import ( "sort" "strings" "text/template" - - "gitcode.com/JianFeeeee/homeagent-sdk/meta" ) func (p *PlgConfig) ReplacesToSlice() []string { @@ -61,10 +59,6 @@ type TemplateData struct { GoVersion string SDKModule string SDKVersion string - - // C ABI - CABIVersion int - CABIHeader string } func cmdInit(args []string) { @@ -139,9 +133,7 @@ func cmdInit(args []string) { Tags: []string{name}, Targets: targets, }, - IsLua: isLua, - CABIVersion: meta.CABINum, - CABIHeader: tmplCABIHeader, + IsLua: isLua, } // Detect SDK info for Go plugin go.mod. diff --git a/tools/plugindev/proc_runtime.go b/tools/plugindev/proc_runtime.go new file mode 100644 index 0000000..c745279 --- /dev/null +++ b/tools/plugindev/proc_runtime.go @@ -0,0 +1,91 @@ +package main + +import ( + "embed" + "fmt" + "os" +) + +// 子进程插件运行时(外部插件多进程化)。 +// +// 模板为何是**真实 .go 源文件** + //go:embed,而不是 raw string: +// 1100+ 行代码塞在字符串里写错只能等生成插件时才炸;作为源文件可被 +// gofmt / go vet / go/parser 直接检查(proc_runtime_test.go 的 16 项 +// 静态检查就以此为前提)。 +// +// 构建从 `-buildmode=c-shared` + CGO_ENABLED=1 变成普通 `go build` + +// CGO_ENABLED=0,交叉编译不再需要目标平台的 C 工具链(§3.1 连带消失项)。 +// +// 设计依据:docs/zh/架构迁移评估.md §3、docs/zh/plugin-migration-plan.md Part 3/6 + +//go:embed templates/proc_main.go.tmpl +//go:embed templates/proc_shm_unix.go.tmpl +//go:embed templates/proc_shm_windows.go.tmpl +var procTemplates embed.FS + +// procRuntimeFiles 列出生成到插件目录的运行时文件。 +// +// 共享段与事件通知的**传递机制**按平台不同(Unix 继承 fd, +// Windows 命名内核对象),故拆成带 build tag 的两个文件; +// 共享段**布局**与 RPC 逻辑完全平台无关,全在 proc_main 里。 +// +// 这正是三套独立 ABI 实现收敛为单一 RPC 实现的效果: +// 平台差异从「整套 stage 下发/写回逻辑各写一份」缩到「三个挂载函数」。 +var procRuntimeFiles = []struct { + tmpl string // 内嵌模板路径 + out string // 生成到插件目录的文件名 +}{ + {"templates/proc_main.go.tmpl", "z_proc_gen.go"}, + {"templates/proc_shm_unix.go.tmpl", "z_proc_shm_unix.go"}, + {"templates/proc_shm_windows.go.tmpl", "z_proc_shm_windows.go"}, +} + +// procEntryFile 是子进程插件的入口二进制名(与内核 internal/plugin/dynamic.go 的 binEntry 一致)。 +// +// 全平台同名:进程边界本身就是 ABI 边界,不存在平台特有的动态库扩展名 +// (对比 C ABI 时代的 .so/.dylib/.dll 三套产物 + 三套 ABI 实现)。 +const procEntryFile = "plugin.bin" + +// luaEntryFile 是 Lua 插件的入口。Lua 走解释器,不经过 Go 编译。 +const luaEntryFile = "main.lua" + +// procGenFile 是生成的主运行时文件名(兼容旧注释引用)。 +// 前缀 z_ 使其在目录列表中排在业务代码之后。 +const procGenFile = "z_proc_gen.go" + +// generateProcRuntime 把子进程运行时(平台无关主体 + 两个平台挂载实现) +// 写入插件目录,返回清理函数。 +func generateProcRuntime() (func(), error) { + // 清理历史 C ABI 产物:旧版 plugindev 生成过这两个文件,残留下来会与 + // 本模板的 main 冲突。无需人工清理就能从旧版升级。 + for _, stale := range []string{"z_bridge_gen.go", "z_entry.c"} { + os.Remove(stale) + } + + var written []string + cleanup := func() { + for _, f := range written { + os.Remove(f) + } + } + + for _, rf := range procRuntimeFiles { + data, err := procTemplates.ReadFile(rf.tmpl) + if err != nil { + cleanup() + return nil, fmt.Errorf("读取内嵌模板 %s: %w", rf.tmpl, err) + } + if err := os.WriteFile(rf.out, data, 0644); err != nil { + cleanup() + return nil, fmt.Errorf("写入 %s: %w", rf.out, err) + } + written = append(written, rf.out) + } + return cleanup, nil +} + +// isProcEntry 已删除:Go 插件一律产出 plugin.bin,不再看 plg.json 的 entry 值。 +// +// 为何忽略 entry:17 个存量插件的 plg.json 都写着 "plugin.so"。若把 entry 当作 +// 通道开关,迁移就得改 17 个文件——而「外部插件零改动」是本次迁移的硬约束。 +// entry 现在只用于区分 Lua(main.lua)与 Go 插件。 diff --git a/tools/plugindev/proc_runtime_test.go b/tools/plugindev/proc_runtime_test.go new file mode 100644 index 0000000..65f290f --- /dev/null +++ b/tools/plugindev/proc_runtime_test.go @@ -0,0 +1,394 @@ +package main + +import ( + "go/parser" + "go/token" + "os" + "regexp" + "strings" + "testing" +) + +// 子进程运行时模板的静态检查(Part 3)。 +// +// 为什么需要这些测试:模板是插件的运行时半身,它与内核 internal/plugin/proc/ +// 的协议名、共享段布局、字段索引必须逐一对齐。任一处漂移都会导致 +// 「插件编译通过但运行时读错字段」——比编译错误难查得多。 +// +// 模板改为真实 .go 源文件(而非 raw string)的直接收益就是这类检查可行。 + +func loadProcTemplate(t *testing.T) string { + t.Helper() + data, err := procTemplates.ReadFile("templates/proc_main.go.tmpl") + if err != nil { + t.Fatalf("读取内嵌模板: %v", err) + } + return string(data) +} + +// stripComments 去掉源码中的注释(用空白填充以保持偏移),只留可执行代码。 +func stripComments(t *testing.T, src string) string { + t.Helper() + fs := token.NewFileSet() + f, err := parser.ParseFile(fs, "proc_main.go", src, parser.ParseComments) + if err != nil { + t.Fatalf("解析模板: %v", err) + } + out := []byte(src) + for _, cg := range f.Comments { + s := fs.Position(cg.Pos()).Offset + e := fs.Position(cg.End()).Offset + for i := s; i < e && i < len(out); i++ { + if out[i] != '\n' { + out[i] = ' ' + } + } + } + return string(out) +} + +// 模板必须是合法 Go 源码。 +func TestProcTemplate_ParsesAsGo(t *testing.T) { + src := loadProcTemplate(t) + fs := token.NewFileSet() + if _, err := parser.ParseFile(fs, "proc_main.go", src, parser.AllErrors); err != nil { + t.Fatalf("模板不是合法 Go 源码: %v", err) + } +} + +// 模板必须提供 main(),且不得含 cgo 痕迹。 +// +// 零 cgo 是迁移的核心收益之一(§3.7 锁仲裁回内核后整个架构无 cgo); +// 一旦有人往模板里加 import "C",交叉编译立刻退回需要目标平台 C 工具链。 +func TestProcTemplate_HasMainAndNoCgo(t *testing.T) { + src := loadProcTemplate(t) + + if !strings.Contains(src, "func main()") { + t.Error("子进程模板必须有 main() 入口") + } + // 只检查代码,不检查注释——模板顶部的说明文字本身就提到了 C.CString/C.free + code := stripComments(t, src) + for _, forbidden := range []string{ + `import "C"`, + "//export ", + "C.CString", + "C.GoString", + "C.free", + } { + if strings.Contains(code, forbidden) { + t.Errorf("模板不应含 cgo 痕迹 %q(零 cgo 是迁移的核心收益)", forbidden) + } + } +} + +// 模板引用的 method 名必须与内核 internal/plugin/proc/protocol.go 一致。 +// +// 这里硬编码一份清单做对照:内核侧改了 method 名而模板没跟上时, +// 表现是插件调用返回「未知 method」,测试能提前拦住。 +func TestProcTemplate_CoversAllCoreMethods(t *testing.T) { + src := loadProcTemplate(t) + + // 51 个 C ABI method id 平移后的名字(§3.2),加 stage 锁仲裁 2 个 + required := []string{ + // 注册面 + "tool.register", "stage.register", "output.register", "api.register", "input.register", + // IO 注入 + "io.injectText", "io.injectInterrupt", "io.injectTextNoMem", "io.injectInputSync", + "io.setToolBlocks", + // 生命周期 + "lifecycle.autoRestart", + // 图记忆 + "memory.recall", "memory.commit", "memory.introspect", "memory.merge", "memory.purge", + // 文档记忆 + "doc.query", "doc.insert", "doc.remove", "doc.stats", + // 知识库 + "knowledge.search", "knowledge.add", "knowledge.list", + // 文本记忆 + "textmemory.append", + // 设置 + "settings.get", "settings.set", "settings.registerDef", + "settings.getCore", "settings.setCore", "settings.listCore", + "settings.getPlugin", "settings.setPlugin", "settings.listPlugin", + "settings.list", "settings.defs", "settings.dump", "settings.plugins", + "settings.dataDir", + // LLM + "llm.listSources", "llm.setSource", "llm.currentSource", + // 社交图 + "social.getPerson", "social.getNetwork", "social.getTrait", + "social.getRelations", "social.listPersons", + // 插件管理 + "plugin.reloadOne", "plugin.listLoaded", "plugin.isDisabled", + // 共享段锁仲裁(新增,C ABI 下不存在此概念) + "stage.lock", "stage.unlock", + } + for _, m := range required { + if !strings.Contains(src, `"`+m+`"`) { + t.Errorf("模板缺少 core method %q(内核已提供,插件侧未接线)", m) + } + } +} + +// 模板必须处理内核发来的全部 7 个调用(原 C ABI 的 7 个 //export)。 +func TestProcTemplate_HandlesAllKernelCalls(t *testing.T) { + src := loadProcTemplate(t) + for _, m := range []string{ + "handshake", + "plugin.init", "plugin.start", "plugin.stop", + "tool.invoke", "stage.invoke", "output.invoke", + } { + if !strings.Contains(src, `case "`+m+`"`) { + t.Errorf("模板未处理内核调用 %q", m) + } + } +} + +// 共享段布局常量必须与内核 internal/plugin/proc/shm.go 一致。 +// +// 字段索引错位是最危险的漂移:插件会读到相邻字段的数据, +// 而两边都不报错(同为 []byte)。 +func TestProcTemplate_ShmLayoutMatchesKernel(t *testing.T) { + src := loadProcTemplate(t) + + // 与内核 shm.go 的 offXxx 常量对齐(值比较,不依赖 gofmt 的对齐空白) + layout := map[string]string{ + "shmOffMagic": "0", + "shmOffVersion": "4", + "shmOffArenaBase": "8", + "shmOffArenaCap": "12", + "shmOffArenaUsed": "16", + "shmOffCtxBase": "20", + "shmOffSeq": "24", + // 与内核 stageFieldCount / sliceSize 对齐 + "shmStageFieldCount": "18", + "shmSliceSize": "8", + "shmVersion": "1", + } + constRe := func(name, want string) bool { + // gofmt 会对齐常量块,故容许 name 与 = 之间有任意空白 + re := regexp.MustCompile(`\b` + regexp.QuoteMeta(name) + `\s*=\s*` + regexp.QuoteMeta(want) + `\b`) + return re.MatchString(src) + } + for name, want := range layout { + if !constRe(name, want) { + t.Errorf("共享段常量 %s 应为 %s(须与内核 internal/plugin/proc/shm.go 一致)", name, want) + } + } + + // 字段枚举顺序:内核 stageField 的前若干项 + fieldOrder := []string{ + "fRawMessage = iota", "fUserID", "fGroupID", "fLLMText", + "fReasoningContent", "fFinalText", "fResponse", "fPhase", + "fContextMsgs", "fToolCalls", "fToolResults", "fMemory", + "fTokenUsage", "fErrors", + "fExtraMediaBlocks", "fExtraMediaType", "fExtraInputSource", "fExtraOutputChannel", + } + idx := -1 + for _, f := range fieldOrder { + at := strings.Index(src, f) + if at < 0 { + t.Fatalf("模板缺少字段常量 %s", f) + } + if at <= idx { + t.Errorf("字段常量 %s 的声明顺序与内核 stageField 枚举不一致", f) + } + idx = at + } +} + +// stage 处理必须「拿锁 → 读 → handler → 只写脏字段 → 放锁」。 +// +// 只写脏字段是消除 lost update 的核心:只读插件零写入, +// 不可能覆盖其他插件的改写(对照 C ABI 副本模型实测 35.8~36.8% 丢失)。 +func TestProcTemplate_StageFlowUsesLockAndDirtyWrite(t *testing.T) { + src := loadProcTemplate(t) + + for _, want := range []string{ + "func handleStageInvoke(", + "stage.lock", + "readStageContext()", + "takeStageSnapshot(", + "writeStageDirty(", + "stage.unlock", + } { + if !strings.Contains(src, want) { + t.Errorf("stage 处理链路缺少 %q", want) + } + } + + // 顺序检查:加锁必须在读取之前,写回必须在解锁之前 + iLock := strings.Index(src, `callCoreVoid("stage.lock"`) + iRead := strings.Index(src, "readStageContext()") + iWrite := strings.Index(src, "writeStageDirty(sc, snap)") + if iLock < 0 || iRead < 0 || iWrite < 0 { + t.Fatal("stage 链路关键调用缺失") + } + // readStageContext 的定义在前,调用在后;取 handleStageInvoke 内的位置 + stageFn := src[strings.Index(src, "func handleStageInvoke("):] + iLockFn := strings.Index(stageFn, `callCoreVoid("stage.lock"`) + iReadFn := strings.Index(stageFn, "readStageContext()") + iWriteFn := strings.Index(stageFn, "writeStageDirty(sc, snap)") + if !(iLockFn < iReadFn && iReadFn < iWriteFn) { + t.Error("stage 链路顺序应为 加锁 → 读取 → 写回") + } +} + +// 快照必须存序列化字符串而非 Go 值。 +// +// ❗ 这是修 C ABI 侧 11.3 时踩过的坑:StageContext 的切片字段与读出的值 +// 共享底层内容,handler 原地改元素(sc.ToolResults[0].Result = x)时, +// 直接持有 Go 值的快照会跟着变,脏字段计算失效、修复静默失效。 +func TestProcTemplate_SnapshotStoresSerializedStrings(t *testing.T) { + src := loadProcTemplate(t) + + if !strings.Contains(src, "strs map[int]string") || + !strings.Contains(src, "jsons map[int]string") { + t.Error("stageSnapshot 必须存序列化字符串(切片共享底层数组,存 Go 值会让脏字段计算失效)") + } + if !strings.Contains(src, "json.Marshal(v)") { + t.Error("takeStageSnapshot 应对容器字段做 json.Marshal") + } +} + +// arena 用尽必须显式报错,不得静默截断(§4.4 风险登记)。 +func TestProcTemplate_ArenaExhaustionErrors(t *testing.T) { + src := loadProcTemplate(t) + if !strings.Contains(src, "arena 空间不足") { + t.Error("shmWrite 在 arena 不足时必须报错,不得静默截断") + } +} + +// 日志必须走 stderr:stdout 是 RPC 通道,写日志会破坏 NDJSON 帧。 +func TestProcTemplate_LogsToStderr(t *testing.T) { + src := loadProcTemplate(t) + if !strings.Contains(src, "log.SetOutput(os.Stderr)") { + t.Error("日志必须走 stderr,否则会破坏 stdout 的 RPC 帧") + } +} + +// 请求必须在独立 goroutine 里处理。 +// +// handler 内会反向调用内核并等应答;若在读循环里同步处理, +// 就没人读应答帧 → 死锁。 +func TestProcTemplate_DispatchesRequestsConcurrently(t *testing.T) { + src := loadProcTemplate(t) + if !strings.Contains(src, "go handleKernelRequest(&req)") { + t.Error("请求须在独立 goroutine 处理(handler 内反向调用内核,同步处理会死锁)") + } +} + +// 协议与共享段版本不匹配必须拒绝,不得半兼容运行。 +func TestProcTemplate_RejectsVersionMismatch(t *testing.T) { + src := loadProcTemplate(t) + for _, want := range []string{"协议版本不匹配", "共享段版本不匹配", "共享段魔数不匹配"} { + if !strings.Contains(src, want) { + t.Errorf("握手应校验并拒绝 %q", want) + } + } +} + +// 全平台统一产出 plugin.bin。 +// +// 这是三套独立 ABI 实现(.so/.dylib/.dll)收敛为单一 RPC 实现的直接后果: +// 进程边界本身就是 ABI 边界,不存在平台特有的动态库扩展名。 +// §9.2 记录的「Windows DLL 路径只下发 3 字段、无写回」随之消失—— +// Windows 走的是与 Linux 完全相同的 RPC 实现。 +func TestResolveBuild_AllPlatformsProduceBin(t *testing.T) { + for _, target := range []string{ + "linux/amd64", "linux/arm64", + "darwin/amd64", "darwin/arm64", + "windows/amd64", + "freebsd/amd64", + } { + cfg, errMsg := resolveBuild(target) + if cfg == nil { + t.Fatalf("resolveBuild(%q) 失败: %s", target, errMsg) + } + if cfg.entryFile != procEntryFile { + t.Errorf("%s: 产物应为 %s,实际 %s", target, procEntryFile, cfg.entryFile) + } + } +} + +// lua 目标仍走解释器路径(entry 字段唯一仍在使用的用途)。 +func TestResolveBuild_LuaIsSeparatePath(t *testing.T) { + for _, target := range []string{"lua", ""} { + cfg, kind := resolveBuild(target) + if cfg != nil { + t.Errorf("%q 应返回 nil cfg(Lua 不经 Go 编译)", target) + } + if kind != "lua" { + t.Errorf("%q 应识别为 lua,实际 %q", target, kind) + } + } +} + +// 不支持的平台明确报错,不静默产出错误产物。 +func TestResolveBuild_UnsupportedOSErrors(t *testing.T) { + cfg, errMsg := resolveBuild("plan9/amd64") + if cfg != nil { + t.Error("不支持的平台应返回 nil cfg") + } + if !strings.Contains(errMsg, "unsupported") { + t.Errorf("应给出 unsupported 提示,实际 %q", errMsg) + } +} + +// bundle 产物在 zip 内按平台加后缀(全平台同名 plugin.bin 会相互覆盖)。 +func TestBundleTargets_HavePlatformSuffixedEntries(t *testing.T) { + seen := map[string]bool{} + for _, bt := range allBundleTargets { + if seen[bt.entry] { + t.Errorf("zip 条目名重复: %s(会相互覆盖)", bt.entry) + } + seen[bt.entry] = true + if !strings.HasPrefix(bt.entry, procEntryFile+".") { + t.Errorf("bundle 条目 %q 应以 %s. 为前缀", bt.entry, procEntryFile) + } + } + if len(allBundleTargets) == 0 { + t.Error("bundle 目标表不应为空") + } +} + +// C ABI 工具链残留必须彻底清除:不得再有 .so/.dylib/.dll 产物路径, +// 也不得再引用 c-shared 构建模式或 MinGW 探测。 +func TestToolchain_NoCABIResiduals(t *testing.T) { + for _, f := range []string{"cmd_build.go", "templates.go", "cmd_init.go", "proc_runtime.go"} { + data, err := os.ReadFile(f) + if err != nil { + t.Fatalf("读 %s: %v", f, err) + } + src := stripComments(t, string(data)) + for _, forbidden := range []string{ + "c-shared", + "CGO_ENABLED=1", + "detectWindowsCC", + "generateBridge", + "tmplLinuxBridge", + "tmplPluginInitC", + } { + if strings.Contains(src, forbidden) { + t.Errorf("%s 仍含 C ABI 残留 %q", f, forbidden) + } + } + } +} + +// Go 插件的构建不再读 plg.json 的 entry 值。 +// +// 这是「外部插件零改动」的关键:17 个存量插件的 plg.json 都写着 "plugin.so", +// 若把 entry 当通道开关,迁移就得改 17 个文件。 +func TestToolchain_IgnoresEntryForGoPlugins(t *testing.T) { + data, err := os.ReadFile("cmd_build.go") + if err != nil { + t.Fatalf("读 cmd_build.go: %v", err) + } + src := stripComments(t, string(data)) + if strings.Contains(src, "isProcEntry") { + t.Error("isProcEntry 应已删除——Go 插件一律产出 plugin.bin,不看 entry 值") + } + // entry 仅剩 Lua 判定这一处用途 + if !strings.Contains(src, "luaEntryFile") { + t.Error("IsLua 应改用 luaEntryFile 常量") + } +} diff --git a/tools/plugindev/stagediff_test.go b/tools/plugindev/stagediff_test.go new file mode 100644 index 0000000..bc5ae03 --- /dev/null +++ b/tools/plugindev/stagediff_test.go @@ -0,0 +1,247 @@ +package main + +import ( + "encoding/json" + "testing" + + sdk "gitcode.com/JianFeeeee/homeagent-sdk/sdk" +) + +// 本测试验证 tmplLinuxBridge 中 snapshotWritable + changedFieldsOnly 的语义(plan.md 11.3)。 +// 模板字符串本身无法直接单测,这里以同一份逻辑复刻,防止回归。 +// ❗ 模板与本文件须同步修改。 +// +// 关键陷阱(第一版实现踩过):stageContextWritable 返回的切片字段与 sc 共享底层数组, +// handler 原地改元素时"before 快照"会跟着变,diff 看不到变更 → 修复静默失效。 +// 故 before 必须是**序列化后的字符串快照**。 + +func writable(sc *sdk.StageContext) map[string]interface{} { + m := map[string]interface{}{ + "raw_message": sc.RawMessage, + "user_id": sc.UserID, + "group_id": sc.GroupID, + "phase": string(sc.Phase), + "llm_text": sc.LLMText, + "final_text": sc.FinalText, + "no_memory": sc.NoMemory, + } + if sc.Response != nil { + m["response"] = *sc.Response + } + if len(sc.ToolCalls) > 0 { + m["tool_calls"] = sc.ToolCalls + } + if len(sc.ToolResults) > 0 { + m["tool_results"] = sc.ToolResults + } + return m +} + +// snapshot 对应模板里的 snapshotWritable:逐字段序列化为不可变快照。 +func snapshot(sc *sdk.StageContext) map[string]string { + snap := map[string]string{} + for k, v := range writable(sc) { + b, err := json.Marshal(v) + if err != nil { + continue + } + snap[k] = string(b) + } + return snap +} + +// diffOnly 对应模板里的 changedFieldsOnly。 +func diffOnly(before map[string]string, after map[string]interface{}) map[string]interface{} { + diff := map[string]interface{}{} + keys := map[string]bool{} + for k := range before { + keys[k] = true + } + for k := range after { + keys[k] = true + } + for k := range keys { + bRaw, bHas := before[k] + a, aHas := after[k] + switch { + case aHas && !bHas: + diff[k] = a + case aHas && bHas: + ab, _ := json.Marshal(a) + if bRaw != string(ab) { + diff[k] = a + } + case bHas && !aHas: + switch k { + case "tool_calls": + diff[k] = []sdk.ToolCall{} + case "tool_results": + diff[k] = []sdk.ToolResult{} + } + } + } + return diff +} + +// 只读插件(如 weather 的 AfterToolcall)不改任何字段 → 零回传。 +// 这是修复 lost update 的关键:旧实现会回传它收到的旧快照,覆盖 sanitizer 的清洗结果。 +func TestChangedFieldsOnly_ReadOnlyPluginReturnsNothing(t *testing.T) { + sc := &sdk.StageContext{ + RawMessage: "hello", + LLMText: "world", + ToolResults: []sdk.ToolResult{ + {CallID: "c1", Name: "weather_query", Success: true, Result: "已清洗结果"}, + }, + } + before := snapshot(sc) + // 只读 handler:读了但没改 + _ = sc.ToolResults[0].Result + diff := diffOnly(before, writable(sc)) + + if len(diff) != 0 { + t.Fatalf("只读插件应零回传,实际回传 %d 个字段: %v", len(diff), diff) + } +} + +// 改写插件(如 sanitizer 改 ToolResults)→ 只回传被改的字段。 +// ⚠️ 这里是原地改切片元素,正是共享底层数组陷阱的触发场景。 +func TestChangedFieldsOnly_WriterReturnsOnlyChanged(t *testing.T) { + sc := &sdk.StageContext{ + RawMessage: "hello", + LLMText: "world", + ToolResults: []sdk.ToolResult{ + {CallID: "c1", Name: "weather_query", Success: true, Result: "带\x1b[31mANSI\x1b[0m脏数据"}, + }, + } + before := snapshot(sc) + // sanitizer handler:原地清洗 ToolResults + sc.ToolResults[0].Result = "带ANSI脏数据" + diff := diffOnly(before, writable(sc)) + + if len(diff) != 1 { + t.Fatalf("应只回传 tool_results 一个字段,实际 %d 个: %v", len(diff), diff) + } + if _, ok := diff["tool_results"]; !ok { + t.Fatalf("回传字段应为 tool_results,实际 %v", diff) + } + // raw_message / llm_text 未改,不应出现(否则会覆盖其他插件的改写) + if _, ok := diff["raw_message"]; ok { + t.Error("raw_message 未改却被回传(会覆盖其他插件的改写)") + } + if _, ok := diff["llm_text"]; ok { + t.Error("llm_text 未改却被回传") + } +} + +// 改写标量字段(如 before_output 改 FinalText)→ 只回传该字段。 +func TestChangedFieldsOnly_ScalarChange(t *testing.T) { + sc := &sdk.StageContext{ + RawMessage: "hi", + FinalText: " 带空白的回复 ", + LLMText: "原始", + } + before := snapshot(sc) + sc.FinalText = "带空白的回复" + diff := diffOnly(before, writable(sc)) + + if len(diff) != 1 || diff["final_text"] != "带空白的回复" { + t.Fatalf("应只回传 final_text,实际 %v", diff) + } +} + +// 首次设置 response(短路)→ 回传。 +func TestChangedFieldsOnly_NewResponseIsReturned(t *testing.T) { + sc := &sdk.StageContext{RawMessage: "hi"} + before := snapshot(sc) + resp := "被插件短路" + sc.Response = &resp + diff := diffOnly(before, writable(sc)) + + if v, ok := diff["response"]; !ok || v != "被插件短路" { + t.Fatalf("新设置的 response 应回传,实际 %v", diff) + } +} + +// 清空切片字段 → 显式回传空值让内核跟随。 +func TestChangedFieldsOnly_ClearedSliceIsReturnedAsEmpty(t *testing.T) { + sc := &sdk.StageContext{ + ToolCalls: []sdk.ToolCall{{ID: "t1", Name: "cmd_run"}}, + } + before := snapshot(sc) + sc.ToolCalls = nil // 插件拒绝了全部工具调用 + diff := diffOnly(before, writable(sc)) + + v, ok := diff["tool_calls"] + if !ok { + t.Fatalf("清空 tool_calls 应显式回传空值,实际 %v", diff) + } + if arr, _ := v.([]sdk.ToolCall); len(arr) != 0 { + t.Fatalf("应回传空切片,实际 %v", v) + } +} + +// 复刻现网场景(实验 13):sanitizer 清洗后 weather 只读回传,清洗结果不得被覆盖。 +// 旧实现下 weather 会回传自己收到的旧快照(含脏数据),覆盖 sanitizer 的清洗(丢失率 1.6~4.3%)。 +func TestChangedFieldsOnly_ProductionScenarioNoOverwrite(t *testing.T) { + dirty := "天气:晴 \x1b[31m28°C\x1b[0m" + clean := "天气:晴 28°C" + + // 内核下发的原始快照(两插件各拿到一份副本) + kernelSnapshot := map[string]interface{}{ + "raw_message": "查天气", + "llm_text": "", + "final_text": "", + "user_id": "u1", + "group_id": "", + "phase": "after_toolcall", + "no_memory": false, + "tool_results": []sdk.ToolResult{{CallID: "c1", Name: "weather_query", Result: dirty}}, + } + + // sanitizer 副本:清洗 + scSan := &sdk.StageContext{ + RawMessage: "查天气", + UserID: "u1", + Phase: sdk.StageAfterToolcall, + ToolResults: []sdk.ToolResult{{CallID: "c1", Name: "weather_query", Result: dirty}}, + } + beforeSan := snapshot(scSan) + scSan.ToolResults[0].Result = clean + diffSan := diffOnly(beforeSan, writable(scSan)) + + // weather 副本:只读,不改 + scWea := &sdk.StageContext{ + RawMessage: "查天气", + UserID: "u1", + Phase: sdk.StageAfterToolcall, + ToolResults: []sdk.ToolResult{{CallID: "c1", Name: "weather_query", Result: dirty}}, + } + beforeWea := snapshot(scWea) + diffWea := diffOnly(beforeWea, writable(scWea)) + + // weather 必须零回传,否则它的旧快照会覆盖 sanitizer 的清洗 + if len(diffWea) != 0 { + t.Fatalf("weather 只读却回传 %v —— 会覆盖 sanitizer 清洗结果", diffWea) + } + // sanitizer 必须回传 tool_results + if _, ok := diffSan["tool_results"]; !ok { + t.Fatalf("sanitizer 改写了 tool_results 却未回传:%v", diffSan) + } + + // 内核按 sanitizer → weather 顺序应用 diff(weather 后到,是最坏情形) + kernel := map[string]interface{}{} + for k, v := range kernelSnapshot { + kernel[k] = v + } + for k, v := range diffSan { + kernel[k] = v + } + for k, v := range diffWea { + kernel[k] = v + } + + res, _ := kernel["tool_results"].([]sdk.ToolResult) + if len(res) == 0 || res[0].Result != clean { + t.Fatalf("清洗结果被覆盖:期望 %q,实际 %v", clean, kernel["tool_results"]) + } +} diff --git a/tools/plugindev/templates.go b/tools/plugindev/templates.go index 9305ad5..bcc018c 100644 --- a/tools/plugindev/templates.go +++ b/tools/plugindev/templates.go @@ -152,721 +152,6 @@ function plugin.stop() sdk.log("info", "{{.Plg.Name}} stopped") end return plugin ` -// tmplBridge — Windows DLL C ABI bridge (unchanged) -const tmplBridge = `//go:build windows && cgo - -package main - -/* -#include -*/ -import "C" -import ( - "encoding/json" - "sync" - "unsafe" - sdk "gitcode.com/JianFeeeee/homeagent-sdk/sdk" -) - -var ( - mu sync.Mutex - handleMap = map[unsafe.Pointer]*bridgeState{} -) - -type bridgeState struct { - plugin sdk.Plugin - toolDefs map[string]sdk.ToolDef - handlers map[string]sdk.ToolHandler - stages map[string]sdk.StageHandler - settings map[string]interface{} - sdk *sdk.PluginSDK -} - -func newHandle(plg sdk.Plugin) unsafe.Pointer { - mu.Lock(); defer mu.Unlock() - h := C.malloc(C.size_t(1)) - handleMap[h] = &bridgeState{ - plugin: plg, toolDefs: make(map[string]sdk.ToolDef), - handlers: make(map[string]sdk.ToolHandler), stages: make(map[string]sdk.StageHandler), - settings: make(map[string]interface{}), - } - return h -} -func getState(h unsafe.Pointer) *bridgeState { mu.Lock(); defer mu.Unlock(); return handleMap[h] } -func delState(h unsafe.Pointer) { mu.Lock(); defer mu.Unlock(); delete(handleMap, h); C.free(h) } - -//export NewPlugin -func NewPlugin(name *C.char, configJSON *C.char) unsafe.Pointer { - goName := C.GoString(name) - var config map[string]interface{} - if configJSON != nil { - var wrapper map[string]interface{} - if err := json.Unmarshal([]byte(C.GoString(configJSON)), &wrapper); err == nil { - if c, ok := wrapper["config"].(map[string]interface{}); ok { config = c } - } - } - plg, err := NewPluginFactory(goName, config) - if err != nil { return nil } - return newHandle(plg) -} - -//export StartPlugin -func StartPlugin(handle unsafe.Pointer) C.int { - bs := getState(handle) - if bs == nil { return 1 } - mockSett := &bridgeSettings{data: bs.settings} - mockSDK := sdk.New(bs.plugin.Name(), mockSett, - func(name string, def sdk.ToolDef, handler sdk.ToolHandler) error { - bs.toolDefs[name] = def; bs.handlers[name] = handler; return nil - }, - func(stage sdk.Stage, handler sdk.StageHandler) { bs.stages[string(stage)] = handler }, - func(name string) error { return nil }, - func(name string, caps int, desc string, def sdk.ChannelDef, handler sdk.ToolHandler) error { return nil }, - ) - mockSDK.SetInputChannelRegistrar(func(name string, def sdk.ChannelDef) error { return nil }) - bs.sdk = mockSDK - if err := bs.plugin.Start(mockSDK); err != nil { return 1 } - return 0 -} - -//export StopPlugin -func StopPlugin(handle unsafe.Pointer) C.int { - bs := getState(handle) - if bs == nil { return 1 } - if bs.sdk != nil { - bs.sdk.RunStopHandlers() - } - if err := bs.plugin.Stop(); err != nil { return 1 } - return 0 -} - -//export DestroyPlugin -func DestroyPlugin(handle unsafe.Pointer) { - if bs := getState(handle); bs != nil { delState(handle) } -} - -//export GetToolDefsJSON -func GetToolDefsJSON(handle unsafe.Pointer) *C.char { - bs := getState(handle) - if bs == nil { return nil } - defs := make([]sdk.ToolDef, 0, len(bs.toolDefs)) - for _, def := range bs.toolDefs { defs = append(defs, def) } - b, _ := json.Marshal(defs) - return C.CString(string(b)) -} - -//export InvokeToolJSON -func InvokeToolJSON(handle unsafe.Pointer, toolName *C.char, argsJSON *C.char) *C.char { - bs := getState(handle) - if bs == nil || toolName == nil { return nil } - goName := C.GoString(toolName) - handler, ok := bs.handlers[goName] - if !ok { errMsg, _ := json.Marshal(map[string]interface{}{"error": "tool not found: " + goName}); return C.CString(string(errMsg)) } - var args map[string]interface{} - if argsJSON != nil { json.Unmarshal([]byte(C.GoString(argsJSON)), &args) } - r, err := handler(args) - if err != nil { errMsg, _ := json.Marshal(map[string]interface{}{"error": err.Error()}); return C.CString(string(errMsg)) } - b, _ := json.Marshal(r) - return C.CString(string(b)) -} - -//export GetStagesJSON -func GetStagesJSON(handle unsafe.Pointer) *C.char { - bs := getState(handle) - if bs == nil { return nil } - type se struct { Stage string ` + "`" + `json:"stage"` + "`" + ` } - var entries []se - for s := range bs.stages { entries = append(entries, se{s}) } - b, _ := json.Marshal(entries) - return C.CString(string(b)) -} - -//export InvokeStage -func InvokeStage(handle unsafe.Pointer, stage *C.char, contextJSON *C.char) C.int { - bs := getState(handle) - if bs == nil || stage == nil { return 1 } - goStage := C.GoString(stage) - handler, ok := bs.stages[goStage] - if !ok { return 1 } - var ctx map[string]interface{} - if contextJSON != nil { json.Unmarshal([]byte(C.GoString(contextJSON)), &ctx) } - sc := &sdk.StageContext{} - if ctx != nil { - if v, ok := ctx["raw_message"].(string); ok { sc.RawMessage = v } - if v, ok := ctx["user_id"].(string); ok { sc.UserID = v } - if v, ok := ctx["phase"].(string); ok { sc.Phase = sdk.Stage(v) } - } - if err := handler(sc); err != nil { return 1 } - return 0 -} - -//export FreeCString -func FreeCString(s *C.char) { C.free(unsafe.Pointer(s)) } - -type bridgeSettings struct{ data map[string]interface{} } -func (s *bridgeSettings) Get(key string) (interface{}, error) { v, ok := s.data[key]; if !ok { return nil, nil }; return v, nil } -func (s *bridgeSettings) Set(key string, value interface{}) error { s.data[key] = value; return nil } -func (s *bridgeSettings) List(prefix string) ([]string, error) { - var keys []string - for k := range s.data { if len(k) >= len(prefix) && k[:len(prefix)] == prefix { keys = append(keys, k) } } - return keys, nil -} -func (s *bridgeSettings) GetCore(key string) (interface{}, error) { return nil, nil } -func (s *bridgeSettings) SetCore(key string, value interface{}) error { return nil } -func (s *bridgeSettings) ListCore(prefix string) ([]string, error) { return nil, nil } -func (s *bridgeSettings) GetPlugin(plugin, key string) (interface{}, error) { return nil, nil } -func (s *bridgeSettings) SetPlugin(plugin, key string, value interface{}) error { return nil } -func (s *bridgeSettings) ListPlugin(plugin, prefix string) ([]string, error) { return nil, nil } -func (s *bridgeSettings) RegisterDef(def sdk.ConfigDef) {} -func (s *bridgeSettings) Defs(prefix string) []*sdk.ConfigDef { return nil } -func (s *bridgeSettings) Dump() map[string]interface{} { return s.data } -func (s *bridgeSettings) Plugins() []string { return nil } - -func main() {} -` - -// tmplCABIHeader — shared C ABI type definitions for both core and plugin -// 此模板中的常量应与 core/internal/meta/meta.go 保持一致(ABI 版本、dispatch method IDs)。 -const tmplCABIHeader = ` -#ifndef HOMEAGENT_CABI_H -#define HOMEAGENT_CABI_H -// HOMEAGENT_ABI_VERSION 与 sdk/meta/meta.go CABINum 同步(major*100+minor,v0.9.x→900) -#define HOMEAGENT_ABI_VERSION 900 -#ifdef __cplusplus -extern "C" { -#endif - -// PluginAPI — implemented by the plugin, called by the core -typedef struct { - int version; int version_min; - int (*init_plugin)(char*, char*, char**); - int (*start_plugin)(void*, int, char**); - int (*stop_plugin)(char**); - int (*invoke_tool)(char*, char*, char**, char**); - int (*invoke_stage)(char*, char*, char**, char**); - int (*invoke_output)(char*, char*, char*, char**); - void (*free_string)(char*); -} PluginAPI; - -// CoreAPI — implemented by the core, passed to plugin via start_plugin -// Uses single dispatch function to avoid function pointer ABI issues -typedef struct { - int version; int version_min; - int (*dispatch)(int method_id, void* ctx, char* s1, char* s2, char* s3, int i1, int i2, char** result, char** error); - void* ctx; -} CoreAPI; - -// Dispatch method IDs (plugin→core SDK calls) -enum { - CORE_REGISTER_TOOL = 1, - CORE_REGISTER_STAGE = 2, - CORE_REGISTER_OUTPUT_CH = 3, - CORE_REGISTER_PLUGIN_API = 4, - CORE_INJECT_TEXT = 5, - CORE_INJECT_INTERRUPT_TEXT = 6, - CORE_INJECT_TEXT_NO_MEMORY = 7, - CORE_INJECT_INPUT_SYNC = 47, - CORE_SET_AUTO_RESTART = 8, - CORE_MEMORY_RECALL = 9, - CORE_MEMORY_COMMIT = 10, - CORE_MEMORY_INTROSPECT = 11, - CORE_MEMORY_MERGE = 12, - CORE_MEMORY_PURGE = 13, - CORE_DOC_QUERY = 14, - CORE_KNOWLEDGE_SEARCH = 15, - CORE_SETTINGS_GET = 16, - CORE_SETTINGS_SET = 17, - CORE_SETTINGS_REGISTER_DEF = 18, - CORE_LLM_LIST_SOURCES = 19, - CORE_LLM_SET_SOURCE = 20, - CORE_SOCIAL_GET_PERSON = 21, - CORE_SOCIAL_GET_NETWORK = 22, - CORE_SUBSCRIBE = 23, - CORE_UNSUBSCRIBE = 24, - CORE_FREE_STRING = 25, - CORE_SETTINGS_GET_CORE = 26, - CORE_SETTINGS_SET_CORE = 27, - CORE_SETTINGS_LIST_CORE = 28, - CORE_SETTINGS_GET_PLUGIN = 29, - CORE_SETTINGS_SET_PLUGIN = 30, - CORE_SETTINGS_LIST_PLUGIN = 31, - CORE_DOC_INSERT = 32, - CORE_DOC_REMOVE = 33, - CORE_DOC_STATS = 34, - CORE_KNOWLEDGE_ADD = 35, - CORE_KNOWLEDGE_LIST = 36, - CORE_LLM_CURRENT_SOURCE = 37, - CORE_SOCIAL_GET_TRAIT = 38, - CORE_SOCIAL_GET_RELATIONS = 39, - CORE_SOCIAL_LIST_PERSONS = 40, - CORE_TEXT_MEMORY_APPEND = 41, - CORE_SETTINGS_LIST = 42, - CORE_SETTINGS_DEFS = 43, - CORE_SETTINGS_DUMP = 44, - CORE_SETTINGS_PLUGINS = 45, - CORE_REGISTER_INPUT_CH = 46, - CORE_INJECT_INPUT_SYNC = 47, - CORE_PLUGIN_RELOAD_ONE = 48, - CORE_PLUGIN_LIST_LOADED = 49, - CORE_PLUGIN_IS_DISABLED = 50, -}; - -#ifdef __cplusplus -} -#endif -#endif -` - -// tmplLinuxBridge — auto-generated Go bridge for Linux c-shared builds. -// Called by plugin's Start() with a PluginSDK that wraps CoreAPI dispatch. -// PluginSDK calls go through C ABI → CoreAPI dispatch → core's Go PluginSDK. -const tmplLinuxBridge = `package main - -/* -#include -int ha_dispatch(int method_id, void* core_api, char* s1, char* s2, char* s3, int i1, int i2, char** result, char** error); -*/ -import "C" -import ( - "encoding/json" - "fmt" - "sync" - "unsafe" - sdk "gitcode.com/JianFeeeee/homeagent-sdk/sdk" -) - -// ---- global state ---- - -var ( - mu sync.Mutex - currentPlg sdk.Plugin - currentSDK *sdk.PluginSDK - coreAPI unsafe.Pointer - - handlerMu sync.RWMutex - coreAPIMu sync.RWMutex - toolHandlers = map[string]sdk.ToolHandler{} - stageHandlers = map[string]sdk.StageHandler{} - outputHandlers = map[string]sdk.ToolHandler{} -) - -// ---- CoreAPI dispatch helpers ---- - -func callVoid(methodID int, s1, s2, s3 string, i1, i2 int) error { - coreAPIMu.RLock() - api := coreAPI - coreAPIMu.RUnlock() - var c1, c2, c3 *C.char - if s1 != "" { c1 = C.CString(s1); defer C.free(unsafe.Pointer(c1)) } - if s2 != "" { c2 = C.CString(s2); defer C.free(unsafe.Pointer(c2)) } - if s3 != "" { c3 = C.CString(s3); defer C.free(unsafe.Pointer(c3)) } - var cErr *C.char - if C.ha_dispatch(C.int(methodID), api, c1, c2, c3, C.int(i1), C.int(i2), nil, &cErr) != 0 && cErr != nil { - return fmt.Errorf("%s", C.GoString(cErr)) - } - return nil -} - -func callString(methodID int, s1, s2, s3 string, i1, i2 int) (string, error) { - coreAPIMu.RLock() - api := coreAPI - coreAPIMu.RUnlock() - var c1, c2, c3 *C.char - if s1 != "" { c1 = C.CString(s1); defer C.free(unsafe.Pointer(c1)) } - if s2 != "" { c2 = C.CString(s2); defer C.free(unsafe.Pointer(c2)) } - if s3 != "" { c3 = C.CString(s3); defer C.free(unsafe.Pointer(c3)) } - var strResult, cErr *C.char - if C.ha_dispatch(C.int(methodID), api, c1, c2, c3, C.int(i1), C.int(i2), &strResult, &cErr) != 0 && cErr != nil { - return "", fmt.Errorf("%s", C.GoString(cErr)) - } - if strResult != nil { - result := C.GoString(strResult) - C.ha_dispatch(C.int(25), api, strResult, nil, nil, 0, 0, nil, nil) - return result, nil - } - return "", nil -} - -// ---- buildPluginSDK: PluginSDK backed by CoreAPI dispatch ---- -// - ALL SDK methods route through C ABI → CoreAPI → core's PluginSDK -// - Handlers for tools/stages/output are stored locally AND registered via dispatch - -func buildPluginSDK(name string) *sdk.PluginSDK { - sett := &dispatchSettings{} - base := sdk.New(name, sett, - func(toolName string, def sdk.ToolDef, handler sdk.ToolHandler) error { - handlerMu.Lock() - toolHandlers[toolName] = handler - handlerMu.Unlock() - b, _ := json.Marshal(def) - return callVoid(1, toolName, string(b), "", 0, 0) - }, - func(stage sdk.Stage, handler sdk.StageHandler) { - handlerMu.Lock() - stageHandlers[string(stage)] = handler - handlerMu.Unlock() - callVoid(2, string(stage), "", "", 0, 0) - }, - func(name string) error { return callVoid(4, name, "", "", 0, 0) }, - func(name string, caps int, desc string, def sdk.ChannelDef, handler sdk.ToolHandler) error { - handlerMu.Lock() - outputHandlers[name] = handler - handlerMu.Unlock() - defJSON, _ := json.Marshal(def) - return callVoid(3, name, desc, string(defJSON), caps, 0) - }, - ) - base.SetIOInjector(dispatchIO{}) - base.SetMemoryAPI(dispatchMemory{}) - base.SetDocMemoryAPI(dispatchDocMemory{}) - base.SetKnowledgeAPI(dispatchKnowledge{}) - base.SetLLMAPI(dispatchLLM{}) - base.SetSocialAPI(dispatchSocial{}) - base.SetTextMemoryAPI(dispatchTextMemory{}) - base.SetPluginMgrAPI(dispatchPluginMgr{}) - base.SetInputChannelRegistrar( - func(name string, def sdk.ChannelDef) error { - defJSON, _ := json.Marshal(def) - return callVoid(46, name, string(defJSON), "", 0, 0) - }, - ) - return base -} - -// ---- dispatch IO (inline definitions) ---- - -type dispatchIO struct{} -func (dispatchIO) InjectInterruptText(s, c, t string) { callVoid(6, s, c, t, 0, 0) } -func (dispatchIO) InjectText(s, c, t string) { callVoid(5, s, c, t, 0, 0) } -func (dispatchIO) InjectTextNoMemory(s, c, t string) { callVoid(7, s, c, t, 0, 0) } -func (dispatchIO) InjectInputSync(s, c, t string) string { r, _ := callString(47, s, c, t, 0, 0); return r } -// SetToolBlocks 是 Go 原生(非 ABI)的多模态注入;跨 ABI 的外部插件无对应内核桥接, -// 故为空实现(满足接口即可)。需要多模态块时用插件内自持 SDK,不走 ABI。 -func (dispatchIO) SetToolBlocks([]sdk.ContentBlock) {} - -type dispatchMemory struct{} -func (dispatchMemory) Recall(q []string, d int) ([]sdk.Entity, []sdk.Relation, error) { - b, _ := json.Marshal(q); r, e := callString(9, string(b), "", "", d, 0) - if e != nil || r == "" { return nil, nil, e } - var v struct{ Entities []sdk.Entity; Relations []sdk.Relation } - if e = json.Unmarshal([]byte(r), &v); e != nil { return nil, nil, e } - if v.Entities == nil { v.Entities = []sdk.Entity{} } - if v.Relations == nil { v.Relations = []sdk.Relation{} } - return v.Entities, v.Relations, nil -} -func (dispatchMemory) Commit(t []sdk.Triple) error { b, _ := json.Marshal(t); return callVoid(10, string(b), "", "", 0, 0) } -func (dispatchMemory) Introspect() (map[string]interface{}, error) { r, e := callString(11, "", "", "", 0, 0); if e != nil || r == "" { return nil, e }; var m map[string]interface{}; return m, json.Unmarshal([]byte(r), &m) } -func (dispatchMemory) MergeEntities(s, t string) (int, error) { return 1, callVoid(12, s, t, "", 0, 0) } -func (dispatchMemory) Purge(c map[string]string, m string) (int, error) { b, _ := json.Marshal(c); i := 0; if m == "hard" { i = 1 }; return 1, callVoid(13, string(b), "", "", i, 0) } - -type dispatchDocMemory struct{} -func (dispatchDocMemory) Query(t string, k int) []*sdk.Doc { r, e := callString(14, t, "", "", k, 0); if e != nil || r == "" { return nil }; var d []*sdk.Doc; json.Unmarshal([]byte(r), &d); return d } -func (dispatchDocMemory) Insert(doc *sdk.Doc) error { b, _ := json.Marshal(doc); return callVoid(32, string(b), "", "", 0, 0) } -func (dispatchDocMemory) Remove(id string) { callVoid(33, id, "", "", 0, 0) } -func (dispatchDocMemory) Stats() map[string]interface{} { r, e := callString(34, "", "", "", 0, 0); if e != nil || r == "" { return nil }; var m map[string]interface{}; json.Unmarshal([]byte(r), &m); return m } - -type dispatchKnowledge struct{} -func (dispatchKnowledge) Search(q string, k int) ([]*sdk.Knowledge, error) { r, e := callString(15, q, "", "", k, 0); if e != nil || r == "" { return nil, e }; var v []*sdk.Knowledge; return v, json.Unmarshal([]byte(r), &v) } -func (dispatchKnowledge) Add(n, c string) error { return callVoid(35, n, c, "", 0, 0) } -func (dispatchKnowledge) List() ([]string, error) { r, e := callString(36, "", "", "", 0, 0); if e != nil || r == "" { return nil, e }; var v []string; return v, json.Unmarshal([]byte(r), &v) } - -type dispatchLLM struct{} -func (dispatchLLM) ListSources() []string { r, e := callString(19, "", "", "", 0, 0); if e != nil || r == "" { return nil }; var v []string; json.Unmarshal([]byte(r), &v); return v } -func (dispatchLLM) SetSource(n string) error { return callVoid(20, n, "", "", 0, 0) } -func (dispatchLLM) CurrentSource() string { r, e := callString(37, "", "", "", 0, 0); if e != nil || r == "" { return "" }; return r } - -type dispatchSocial struct{} -func (dispatchSocial) GetPerson(n string) (*sdk.PersonProfile, error) { r, e := callString(21, n, "", "", 0, 0); if e != nil || r == "" { return nil, e }; var v sdk.PersonProfile; return &v, json.Unmarshal([]byte(r), &v) } -func (dispatchSocial) GetTrait(n, t string) (string, bool) { r, e := callString(38, n, t, "", 0, 0); if e != nil || r == "" { return "", false }; var m map[string]interface{}; json.Unmarshal([]byte(r), &m); v, _ := m["value"].(string); ok, _ := m["found"].(bool); return v, ok } -func (dispatchSocial) GetRelations(name string) ([]sdk.SocialRelation, error) { r, e := callString(39, name, "", "", 0, 0); if e != nil || r == "" { return nil, e }; var v []sdk.SocialRelation; return v, json.Unmarshal([]byte(r), &v) } -func (dispatchSocial) GetNetwork(n string, d int) ([]*sdk.PersonProfile, error) { r, e := callString(22, n, "", "", d, 0); if e != nil || r == "" { return nil, e }; var v []*sdk.PersonProfile; return v, json.Unmarshal([]byte(r), &v) } -func (dispatchSocial) ListPersons() ([]string, error) { r, e := callString(40, "", "", "", 0, 0); if e != nil || r == "" { return nil, e }; var v []string; return v, json.Unmarshal([]byte(r), &v) } - -type dispatchTextMemory struct{} -func (dispatchTextMemory) Append(evt sdk.TextEvent) error { b, _ := json.Marshal(evt); return callVoid(41, string(b), "", "", 0, 0) } - -// ---- dispatchPluginMgr (CORE_PLUGIN_RELOAD_ONE = 48) ---- - -type dispatchPluginMgr struct{} - -func (dispatchPluginMgr) ReloadOne(name string) error { - return callVoid(48, name, "", "", 0, 0) -} - -func (dispatchPluginMgr) ListLoadedPlugins() []string { - r, e := callString(49, "", "", "", 0, 0) - if e != nil || r == "" { - return nil - } - var list []string - if json.Unmarshal([]byte(r), &list) != nil { - return nil - } - return list -} - -func (dispatchPluginMgr) IsPluginDisabled(name string) bool { - r, e := callString(50, name, "", "", 0, 0) - return e == nil && r == "1" -} - -// ---- dispatchSettings (inline) ---- - -type dispatchSettings struct{} -func (d *dispatchSettings) Get(key string) (interface{}, error) { - r, e := callString(16, key, "", "", 0, 0); if e != nil || r == "" { return nil, e }; var v interface{}; return v, json.Unmarshal([]byte(r), &v) -} -func (d *dispatchSettings) Set(key string, value interface{}) error { - b, _ := json.Marshal(value); return callVoid(17, key, string(b), "", 0, 0) -} -func (d *dispatchSettings) RegisterDef(def sdk.ConfigDef) { b, _ := json.Marshal(def); callVoid(18, string(b), "", "", 0, 0) } -func (d *dispatchSettings) List(prefix string) ([]string, error) { - r, e := callString(42, prefix, "", "", 0, 0); if e != nil || r == "" { return nil, e }; var v []string; return v, json.Unmarshal([]byte(r), &v) -} -func (d *dispatchSettings) GetCore(key string) (interface{}, error) { - r, e := callString(26, key, "", "", 0, 0); if e != nil || r == "" { return nil, e }; var v interface{}; return v, json.Unmarshal([]byte(r), &v) -} -func (d *dispatchSettings) SetCore(key string, value interface{}) error { - b, _ := json.Marshal(value); return callVoid(27, key, string(b), "", 0, 0) -} -func (d *dispatchSettings) ListCore(prefix string) ([]string, error) { - r, e := callString(28, prefix, "", "", 0, 0); if e != nil || r == "" { return nil, e }; var v []string; return v, json.Unmarshal([]byte(r), &v) -} -func (d *dispatchSettings) GetPlugin(plugin, key string) (interface{}, error) { - r, e := callString(29, plugin, key, "", 0, 0); if e != nil || r == "" { return nil, e }; var v interface{}; return v, json.Unmarshal([]byte(r), &v) -} -func (d *dispatchSettings) SetPlugin(plugin, key string, value interface{}) error { - b, _ := json.Marshal(value); return callVoid(30, plugin, key, string(b), 0, 0) -} -func (d *dispatchSettings) ListPlugin(plugin, prefix string) ([]string, error) { - r, e := callString(31, plugin, prefix, "", 0, 0); if e != nil || r == "" { return nil, e }; var v []string; return v, json.Unmarshal([]byte(r), &v) -} -func (d *dispatchSettings) Defs(prefix string) []*sdk.ConfigDef { - r, e := callString(43, prefix, "", "", 0, 0); if e != nil || r == "" { return nil }; var v []*sdk.ConfigDef; json.Unmarshal([]byte(r), &v); return v -} -func (d *dispatchSettings) Dump() map[string]interface{} { - r, e := callString(44, "", "", "", 0, 0); if e != nil || r == "" { return nil }; var m map[string]interface{}; json.Unmarshal([]byte(r), &m); return m -} -func (d *dispatchSettings) Plugins() []string { - r, e := callString(45, "", "", "", 0, 0); if e != nil || r == "" { return nil }; var v []string; json.Unmarshal([]byte(r), &v); return v -} -func (d *dispatchSettings) DataDir() string { - r, e := callString(51, "", "", "", 0, 0); if e != nil { return "" }; return r -} - -// ---- Go callbacks (called from z_entry.c via C) ---- - -//export go_init_plugin -func go_init_plugin(name *C.char, configJSON *C.char, errorOut **C.char) C.int { - plg, err := NewPluginFactory(C.GoString(name), nil) - if err != nil || plg == nil { - if err != nil { *errorOut = C.CString(err.Error()) } else { *errorOut = C.CString("NewPluginFactory returned nil") } - return 1 - } - mu.Lock(); currentPlg = plg; mu.Unlock() - _ = configJSON - return 0 -} - -//export go_start_plugin -func go_start_plugin(coreAPIptr unsafe.Pointer, coreVersion C.int, errorOut **C.char) C.int { - mu.Lock() - plg := currentPlg - coreAPIMu.Lock() - coreAPI = coreAPIptr - coreAPIMu.Unlock() - mu.Unlock() - _ = coreVersion - if plg == nil { *errorOut = C.CString("not initialized"); return 1 } - sdk := buildPluginSDK(plg.Name()) - mu.Lock(); currentSDK = sdk; mu.Unlock() - if err := plg.Start(sdk); err != nil { *errorOut = C.CString(err.Error()); return 1 } - return 0 -} - -//export go_stop_plugin -func go_stop_plugin(errorOut **C.char) C.int { - mu.Lock() - plg := currentPlg - sdk := currentSDK - currentPlg = nil - currentSDK = nil - coreAPIMu.Lock() - coreAPI = nil - coreAPIMu.Unlock() - mu.Unlock() - if sdk != nil { - sdk.RunStopHandlers() - } - if plg != nil { - if err := plg.Stop(); err != nil { *errorOut = C.CString(err.Error()); return 1 } - } - return 0 -} - -//export go_invoke_tool -func go_invoke_tool(name *C.char, argsJSON *C.char, resultOut **C.char, errorOut **C.char) C.int { - goName := C.GoString(name) - handlerMu.RLock() - h, ok := toolHandlers[goName] - handlerMu.RUnlock() - if !ok { *errorOut = C.CString("tool not found"); return 1 } - var args map[string]interface{} - if argsJSON != nil { json.Unmarshal([]byte(C.GoString(argsJSON)), &args) } - r, err := h(args) - if err != nil { *errorOut = C.CString(err.Error()); return 1 } - b, _ := json.Marshal(r) - *resultOut = C.CString(string(b)) - return 0 -} - -// fillStageContext 将内核传来的 ctx JSON 填充到插件侧 StageContext。 -func fillStageContext(sc *sdk.StageContext, ctxJSON string) { - var m map[string]interface{} - if err := json.Unmarshal([]byte(ctxJSON), &m); err != nil { - return - } - if v, _ := m["raw_message"].(string); v != "" { sc.RawMessage = v } - if v, _ := m["user_id"].(string); v != "" { sc.UserID = v } - if v, _ := m["group_id"].(string); v != "" { sc.GroupID = v } - if v, _ := m["phase"].(string); v != "" { sc.Phase = sdk.Stage(v) } - if v, _ := m["llm_text"].(string); v != "" { sc.LLMText = v } - if v, _ := m["final_text"].(string); v != "" { sc.FinalText = v } - if v, _ := m["no_memory"].(bool); v { sc.NoMemory = true } - if v, _ := m["response"].(string); v != "" { sc.Response = &v } - if v, _ := m["tool_calls"].([]interface{}); len(v) > 0 { - b, _ := json.Marshal(v); json.Unmarshal(b, &sc.ToolCalls) - } - if v, _ := m["tool_results"].([]interface{}); len(v) > 0 { - b, _ := json.Marshal(v); json.Unmarshal(b, &sc.ToolResults) - } -} - -// stageContextWritable 提取插件可写且内核会同步回去的字段。 -func stageContextWritable(sc *sdk.StageContext) map[string]interface{} { - m := map[string]interface{}{ - "raw_message": sc.RawMessage, - "user_id": sc.UserID, - "group_id": sc.GroupID, - "phase": string(sc.Phase), - "llm_text": sc.LLMText, - "final_text": sc.FinalText, - "no_memory": sc.NoMemory, - } - if sc.Response != nil { - m["response"] = *sc.Response - } - if len(sc.ToolCalls) > 0 { - m["tool_calls"] = sc.ToolCalls - } - if len(sc.ToolResults) > 0 { - m["tool_results"] = sc.ToolResults - } - return m -} - -//export go_invoke_stage -func go_invoke_stage(stage *C.char, ctxJSON *C.char, resultOut **C.char, errorOut **C.char) C.int { - goStage := C.GoString(stage) - handlerMu.RLock() - h, ok := stageHandlers[goStage] - handlerMu.RUnlock() - if !ok { return 0 } - sc := &sdk.StageContext{} - if ctxJSON != nil { - fillStageContext(sc, C.GoString(ctxJSON)) - } - if err := h(sc); err != nil { *errorOut = C.CString(err.Error()); return 1 } - // ABI v2: 回传插件修改后的上下文(若调用方要求) - if resultOut != nil { - if b, err := json.Marshal(stageContextWritable(sc)); err == nil { - *resultOut = C.CString(string(b)) - } - } - return 0 -} - -//export go_invoke_output -func go_invoke_output(channel *C.char, msgType *C.char, payloadJSON *C.char, errorOut **C.char) C.int { - goChan := C.GoString(channel) - handlerMu.RLock() - h, ok := outputHandlers[goChan] - handlerMu.RUnlock() - if !ok { return 0 } - // payloadJSON contains the full args JSON from output_send (e.g. {"content":"...","user_id":123}) - var args map[string]interface{} - if payloadJSON != nil { - json.Unmarshal([]byte(C.GoString(payloadJSON)), &args) - } - if _, err := h(args); err != nil { *errorOut = C.CString(err.Error()); return 1 } - return 0 -} - -//export go_free_string -func go_free_string(ptr *C.char) { C.free(unsafe.Pointer(ptr)) } - -func main() {} -` - -// tmplPluginInitC — C entry point for the plugin .so file. -// Contains PluginAPI, CoreAPI (single dispatch), and ha_dispatch bridge. -const tmplPluginInitC = `#include -#include - -// HOMEAGENT_ABI_VERSION 与 sdk/meta/meta.go CABINum 同步(major*100+minor,v0.9.x→900) -#define HOMEAGENT_ABI_VERSION 900 - -typedef struct { - int version; int version_min; - int (*init_plugin)(char*, char*, char**); - int (*start_plugin)(void*, int, char**); - int (*stop_plugin)(char**); - int (*invoke_tool)(char*, char*, char**, char**); - int (*invoke_stage)(char*, char*, char**, char**); - int (*invoke_output)(char*, char*, char*, char**); - void (*free_string)(char*); -} PluginAPI; - -typedef struct { - int version; int version_min; - int (*dispatch)(int, void*, char*, char*, char*, int, int, char**, char**); - void* ctx; -} CoreAPI; - -extern int go_init_plugin(char*, char*, char**); -extern int go_start_plugin(void*, int, char**); -extern int go_stop_plugin(char**); -extern int go_invoke_tool(char*, char*, char**, char**); -extern int go_invoke_stage(char*, char*, char**, char**); -extern int go_invoke_output(char*, char*, char*, char**); -extern void go_free_string(char*); - -int c_init_plugin(char* n, char* c, char** e) { return go_init_plugin(n, c, e); } -int c_start_plugin(void* a, int v, char** e) { return go_start_plugin(a, v, e); } -int c_stop_plugin(char** e) { return go_stop_plugin(e); } -int c_invoke_tool(char* n, char* a, char** r, char** e) { return go_invoke_tool(n, a, r, e); } -int c_invoke_stage(char* s, char* c, char** r, char** e) { return go_invoke_stage(s, c, r, e); } -int c_invoke_output(char* c, char* m, char* p, char** e) { return go_invoke_output(c, m, p, e); } -void c_free_string(char* p) { go_free_string(p); } - -// ha_dispatch — called by Go bridge, passes through to CoreAPI dispatch -int ha_dispatch(int id, void* api, char* s1, char* s2, char* s3, int i1, int i2, char** r, char** e) { - CoreAPI* a = (CoreAPI*)api; - if (!a || !a->dispatch) return 1; - return a->dispatch(id, a->ctx, s1, s2, s3, i1, i2, r, e); -} - -PluginAPI* plugin_init(void) { - static PluginAPI api; - memset(&api, 0, sizeof(api)); - api.version = HOMEAGENT_ABI_VERSION; api.version_min = HOMEAGENT_ABI_VERSION; - api.init_plugin = c_init_plugin; api.start_plugin = c_start_plugin; api.stop_plugin = c_stop_plugin; - api.invoke_tool = c_invoke_tool; api.invoke_stage = c_invoke_stage; api.invoke_output = c_invoke_output; - api.free_string = c_free_string; - return &api; -} -` - // ============================================================ // Remote Device Adapter Templates // ============================================================ diff --git a/tools/plugindev/templates/proc_main.go.tmpl b/tools/plugindev/templates/proc_main.go.tmpl new file mode 100644 index 0000000..537b76a --- /dev/null +++ b/tools/plugindev/templates/proc_main.go.tmpl @@ -0,0 +1,1291 @@ +package main + +// 子进程插件入口(由 plugindev 自动生成,请勿手工编辑)。 +// +// 与旧 C ABI bridge(z_bridge_gen.go)的关键差异: +// - **零 cgo**:没有 //export、没有 C.CString/C.free、不需要 -buildmode=c-shared +// - 51 个整数 method id 换成可读 method 名(内核侧 internal/plugin/proc/protocol.go) +// - StageContext 走共享内存(fd 3 传入的 memfd),插件在同一份状态上读改写, +// 消除副本模型的 lost update(实测 35.8~36.8% → 0) +// - **插件业务代码零改动**:仍是 NewPluginFactory + sdk.PluginSDK +// +// 设计依据:docs/zh/架构迁移评估.md 第三章 + +import ( + "bufio" + "encoding/binary" + "encoding/json" + "fmt" + "log" + "os" + "sync" + + sdk "gitcode.com/JianFeeeee/homeagent-sdk/sdk" +) + +// ---- 协议常量(须与内核 internal/plugin/proc/protocol.go 一致)---- + +const procProtocolVersion = 1 + +// ---- 共享段布局(须与内核 internal/plugin/proc/shm.go 一致)---- + +const ( + shmStageFieldCount = 18 + shmSliceSize = 8 + + shmOffMagic = 0 + shmOffVersion = 4 + shmOffArenaBase = 8 + shmOffArenaCap = 12 + shmOffArenaUsed = 16 + shmOffCtxBase = 20 + shmOffSeq = 24 + + shmMagic = 0x48415348 + shmVersion = 1 +) + +// 字段索引(顺序须与内核 stageField 枚举一致) +const ( + fRawMessage = iota + fUserID + fGroupID + fLLMText + fReasoningContent + fFinalText + fResponse + fPhase + fContextMsgs + fToolCalls + fToolResults + fMemory + fTokenUsage + fErrors + fExtraMediaBlocks + fExtraMediaType + fExtraInputSource + fExtraOutputChannel +) + +const ( + flagNoMemory = 0 + flagResponseSet = 1 +) + +// ---- 全局状态 ---- + +var ( + stdoutW = bufio.NewWriter(os.Stdout) + writeMu sync.Mutex + + nextID uint64 + pendMu sync.Mutex + pending = map[uint64]chan rpcResponse{} + + plg sdk.Plugin + pluginSDK *sdk.PluginSDK + pluginName string + + handlerMu sync.RWMutex + toolHandlers = map[string]sdk.ToolHandler{} + stageHandlers = map[string]sdk.StageHandler{} + outputHandlers = map[string]sdk.ToolHandler{} + + shm []byte + + // 事件环(§3.6):fd 4 = 事件环段 mmap,fd 5 = eventfd 读端 + evtRingData []byte + evtNotifier evtWaiter + evtHandlers = map[uint32]func(*sdk.Event){} + evtHandlerMu sync.RWMutex +) + +type rpcRequest struct { + ID uint64 `json:"id,omitempty"` + Method string `json:"method"` + Params json.RawMessage `json:"params,omitempty"` +} + +type rpcResponse struct { + ID uint64 `json:"id"` + Result json.RawMessage `json:"result,omitempty"` + Error string `json:"error,omitempty"` +} + +func writeFrame(v interface{}) { + b, err := json.Marshal(v) + if err != nil { + log.Printf("序列化帧失败: %v", err) + return + } + writeMu.Lock() + stdoutW.Write(b) + stdoutW.WriteByte('\n') + stdoutW.Flush() + writeMu.Unlock() +} + +func respond(id uint64, result interface{}) { + resp := rpcResponse{ID: id} + if result != nil { + if b, err := json.Marshal(result); err == nil { + resp.Result = b + } + } + writeFrame(&resp) +} + +func respondErr(id uint64, err error) { + writeFrame(&rpcResponse{ID: id, Error: err.Error()}) +} + +// callCore 反向调用内核(对应旧 bridge 的 callVoid/callString)。 +func callCore(method string, params interface{}) (json.RawMessage, error) { + pendMu.Lock() + nextID++ + id := nextID + ch := make(chan rpcResponse, 1) + pending[id] = ch + pendMu.Unlock() + + var raw json.RawMessage + if params != nil { + b, err := json.Marshal(params) + if err != nil { + return nil, err + } + raw = b + } + writeFrame(&rpcRequest{ID: id, Method: method, Params: raw}) + + resp := <-ch + if resp.Error != "" { + return nil, fmt.Errorf("%s", resp.Error) + } + return resp.Result, nil +} + +func callCoreVoid(method string, params interface{}) error { + _, err := callCore(method, params) + return err +} + +// ---- 共享段访问(插件作者永远不接触这些,§3.4)---- + +func shmU32(off int) uint32 { return binary.LittleEndian.Uint32(shm[off:]) } +func shmArenaBase() uint32 { return shmU32(shmOffArenaBase) } +func shmArenaCap() uint32 { return shmU32(shmOffArenaCap) } +func shmCtxBase() uint32 { return shmU32(shmOffCtxBase) } + +func shmDescOff(field int) uint32 { + return shmCtxBase() + uint32(field*shmSliceSize) +} + +func shmGetDesc(field int) (off, ln uint32) { + o := shmDescOff(field) + return binary.LittleEndian.Uint32(shm[o:]), binary.LittleEndian.Uint32(shm[o+4:]) +} + +func shmSetDesc(field int, off, ln uint32) { + o := shmDescOff(field) + binary.LittleEndian.PutUint32(shm[o:], off) + binary.LittleEndian.PutUint32(shm[o+4:], ln) +} + +func shmFlagsOff() uint32 { + return shmCtxBase() + uint32(shmStageFieldCount*shmSliceSize) +} + +func shmGetFlag(bit int) bool { return shm[shmFlagsOff()+uint32(bit)] != 0 } + +func shmSetFlag(bit int, v bool) { + b := byte(0) + if v { + b = 1 + } + shm[shmFlagsOff()+uint32(bit)] = b +} + +func shmRead(field int) []byte { + off, ln := shmGetDesc(field) + if off == 0 && ln == 0 { + return nil + } + if ln == 0 { + return []byte{} + } + base := shmArenaBase() + return shm[base+off : base+off+ln] +} + +// shmWrite 在 arena 上 append-only 分配并更新描述符。 +// arena 用尽显式报错,不静默截断(与内核侧同一约定)。 +func shmWrite(field int, data []byte) error { + if len(data) == 0 { + shmSetDesc(field, 1, 0) + return nil + } + used := shmU32(shmOffArenaUsed) + if used == 0 { + used = 1 + } + end := used + uint32(len(data)) + if end > shmArenaCap() { + return fmt.Errorf("共享段 arena 空间不足:需要 %d 字节,容量 %d,已用 %d", + len(data), shmArenaCap(), used) + } + base := shmArenaBase() + copy(shm[base+used:], data) + binary.LittleEndian.PutUint32(shm[shmOffArenaUsed:], end) + shmSetDesc(field, used, uint32(len(data))) + return nil +} + +func shmBumpSeq() { + v := binary.LittleEndian.Uint64(shm[shmOffSeq:]) + binary.LittleEndian.PutUint64(shm[shmOffSeq:], v+1) +} + +// readStageContext 从共享段构造插件侧原生 StageContext。 +// 全 16 个字段可见——C ABI 下只有 10 个(§8.3)。 +func readStageContext() (*sdk.StageContext, error) { + sc := &sdk.StageContext{} + + sc.RawMessage = string(shmRead(fRawMessage)) + sc.UserID = string(shmRead(fUserID)) + sc.GroupID = string(shmRead(fGroupID)) + sc.LLMText = string(shmRead(fLLMText)) + sc.ReasoningContent = string(shmRead(fReasoningContent)) + sc.FinalText = string(shmRead(fFinalText)) + sc.Phase = sdk.Stage(string(shmRead(fPhase))) + sc.NoMemory = shmGetFlag(flagNoMemory) + + if shmGetFlag(flagResponseSet) { + r := string(shmRead(fResponse)) + sc.Response = &r + } + + unmarshalField := func(field int, out interface{}) error { + b := shmRead(field) + if len(b) == 0 { + return nil + } + return json.Unmarshal(b, out) + } + if err := unmarshalField(fContextMsgs, &sc.ContextMsgs); err != nil { + return nil, err + } + if err := unmarshalField(fToolCalls, &sc.ToolCalls); err != nil { + return nil, err + } + if err := unmarshalField(fToolResults, &sc.ToolResults); err != nil { + return nil, err + } + if err := unmarshalField(fMemory, &sc.Memory); err != nil { + return nil, err + } + if err := unmarshalField(fTokenUsage, &sc.TokenUsage); err != nil { + return nil, err + } + if err := unmarshalField(fErrors, &sc.Errors); err != nil { + return nil, err + } + + extra := map[string]interface{}{} + for _, pair := range []struct { + field int + key string + }{ + {fExtraMediaBlocks, "media_blocks"}, + {fExtraMediaType, "media_type"}, + {fExtraInputSource, "input_source"}, + {fExtraOutputChannel, "output_channel"}, + } { + var v interface{} + if err := unmarshalField(pair.field, &v); err != nil { + return nil, err + } + if v != nil { + extra[pair.key] = v + } + } + if len(extra) > 0 { + sc.Extra = extra + } + return sc, nil +} + +// stageSnapshot 是 handler 运行前的序列化快照,用于计算脏字段。 +// +// ❗ 必须存序列化后的字符串:handler 原地改切片元素 +// (sc.ToolResults[0].Result = x)时,直接持有的 Go 值快照会跟着变, +// 脏字段计算失效——这个坑在修 C ABI 侧的 11.3 时已经踩过一次。 +type stageSnapshot struct { + strs map[int]string + jsons map[int]string + response string + responseSet bool + noMemory bool +} + +func takeStageSnapshot(sc *sdk.StageContext) *stageSnapshot { + sn := &stageSnapshot{strs: map[int]string{}, jsons: map[int]string{}} + sn.strs[fRawMessage] = sc.RawMessage + sn.strs[fUserID] = sc.UserID + sn.strs[fGroupID] = sc.GroupID + sn.strs[fLLMText] = sc.LLMText + sn.strs[fReasoningContent] = sc.ReasoningContent + sn.strs[fFinalText] = sc.FinalText + sn.strs[fPhase] = string(sc.Phase) + + marshal := func(v interface{}, n int) string { + if n == 0 { + return "" + } + b, err := json.Marshal(v) + if err != nil { + return "" + } + return string(b) + } + sn.jsons[fContextMsgs] = marshal(sc.ContextMsgs, len(sc.ContextMsgs)) + sn.jsons[fToolCalls] = marshal(sc.ToolCalls, len(sc.ToolCalls)) + sn.jsons[fToolResults] = marshal(sc.ToolResults, len(sc.ToolResults)) + sn.jsons[fMemory] = marshal(sc.Memory, len(sc.Memory)) + sn.jsons[fTokenUsage] = marshal(sc.TokenUsage, len(sc.TokenUsage)) + sn.jsons[fErrors] = marshal(sc.Errors, len(sc.Errors)) + + if sc.Response != nil { + sn.response = *sc.Response + sn.responseSet = true + } + sn.noMemory = sc.NoMemory + return sn +} + +// writeStageDirty 只把变更字段写回共享段,返回写回字段数。 +// +// **这是消除 lost update 的核心**:只读插件的脏字段集为空 → 零写入 → +// 不可能覆盖其他插件的改写(对照 C ABI 副本模型实测 35.8~36.8% 丢失)。 +func writeStageDirty(sc *sdk.StageContext, base *stageSnapshot) (int, error) { + now := takeStageSnapshot(sc) + changed := 0 + + for field, cur := range now.strs { + if base.strs[field] != cur { + if err := shmWrite(field, []byte(cur)); err != nil { + return changed, err + } + changed++ + } + } + for field, cur := range now.jsons { + if base.jsons[field] == cur { + continue + } + if err := shmWrite(field, []byte(cur)); err != nil { + return changed, err + } + changed++ + } + if base.responseSet != now.responseSet || base.response != now.response { + if now.responseSet { + if err := shmWrite(fResponse, []byte(now.response)); err != nil { + return changed, err + } + shmSetFlag(flagResponseSet, true) + changed++ + } + // Response 置回 nil 不清空内核已设的值:短路语义不应被撑销 + } + if base.noMemory != now.noMemory { + shmSetFlag(flagNoMemory, now.noMemory) + changed++ + } + + if changed > 0 { + shmBumpSeq() + } + return changed, nil +} + +// ---- SDK 装配:全部 API 经 RPC 打回内核(51 个 method 的插件侧一半)---- + +func buildPluginSDK(name string) *sdk.PluginSDK { + base := sdk.New(name, procSettings{}, + func(toolName string, def sdk.ToolDef, handler sdk.ToolHandler) error { + handlerMu.Lock() + toolHandlers[toolName] = handler + handlerMu.Unlock() + return callCoreVoid("tool.register", map[string]interface{}{ + "name": toolName, "def": def, + }) + }, + func(stage sdk.Stage, handler sdk.StageHandler) { + handlerMu.Lock() + stageHandlers[string(stage)] = handler + handlerMu.Unlock() + if err := callCoreVoid("stage.register", map[string]interface{}{ + "stage": string(stage), "scope": "global", + }); err != nil { + log.Printf("注册阶段 %s 失败: %v", stage, err) + } + }, + func(apiName string) error { + return callCoreVoid("api.register", map[string]interface{}{"name": apiName}) + }, + func(chName string, caps int, desc string, def sdk.ChannelDef, handler sdk.ToolHandler) error { + handlerMu.Lock() + outputHandlers[chName] = handler + handlerMu.Unlock() + return callCoreVoid("output.register", map[string]interface{}{ + "name": chName, "caps": caps, "desc": desc, + "def": map[string]interface{}{"NoMemory": def.NoMemory}, + }) + }, + ) + + base.SetIOInjector(procIO{}) + base.SetMemoryAPI(procMemory{}) + base.SetDocMemoryAPI(procDocMemory{}) + base.SetKnowledgeAPI(procKnowledge{}) + base.SetLLMAPI(procLLM{}) + base.SetSocialAPI(procSocial{}) + base.SetTextMemoryAPI(procTextMemory{}) + base.SetPluginMgrAPI(procPluginMgr{}) + base.SetInputChannelRegistrar(func(chName string, def sdk.ChannelDef) error { + return callCoreVoid("input.register", map[string]interface{}{ + "name": chName, + "def": map[string]interface{}{"NoMemory": def.NoMemory}, + }) + }) + return base +} + +type procIO struct{} + +func (procIO) InjectText(s, c, t string) { + callCoreVoid("io.injectText", map[string]string{"source": s, "channel": c, "text": t}) +} +func (procIO) InjectInterruptText(s, c, t string) { + callCoreVoid("io.injectInterrupt", map[string]string{"source": s, "channel": c, "text": t}) +} +func (procIO) InjectTextNoMemory(s, c, t string) { + callCoreVoid("io.injectTextNoMem", map[string]string{"source": s, "channel": c, "text": t}) +} +func (procIO) InjectInputSync(s, c, t string) string { + raw, err := callCore("io.injectInputSync", map[string]string{"source": s, "channel": c, "text": t}) + if err != nil { + return "" + } + var r struct { + Reply string `json:"reply"` + } + json.Unmarshal(raw, &r) + return r.Reply +} +func (procIO) SetToolBlocks(blocks []sdk.ContentBlock) { + if err := callCoreVoid("io.setToolBlocks", map[string]interface{}{"blocks": blocks}); err != nil { + log.Printf("SetToolBlocks: %v", err) + } +} + +type procMemory struct{} + +func (procMemory) Recall(q []string, d int) ([]sdk.Entity, []sdk.Relation, error) { + raw, err := callCore("memory.recall", map[string]interface{}{"query": q, "depth": d}) + if err != nil { + return nil, nil, err + } + var r struct { + Entities []sdk.Entity `json:"entities"` + Relations []sdk.Relation `json:"relations"` + } + if err := json.Unmarshal(raw, &r); err != nil { + return nil, nil, err + } + return r.Entities, r.Relations, nil +} +func (procMemory) Commit(t []sdk.Triple) error { + return callCoreVoid("memory.commit", map[string]interface{}{"triples": t}) +} +func (procMemory) Introspect() (map[string]interface{}, error) { + raw, err := callCore("memory.introspect", nil) + if err != nil { + return nil, err + } + var m map[string]interface{} + json.Unmarshal(raw, &m) + return m, nil +} +func (procMemory) MergeEntities(s, t string) (int, error) { + raw, err := callCore("memory.merge", map[string]string{"source": s, "target": t}) + if err != nil { + return 0, err + } + var r struct { + Merged int `json:"merged"` + } + json.Unmarshal(raw, &r) + return r.Merged, nil +} +func (procMemory) Purge(c map[string]string, mode string) (int, error) { + raw, err := callCore("memory.purge", map[string]interface{}{"criteria": c, "mode": mode}) + if err != nil { + return 0, err + } + var r struct { + Purged int `json:"purged"` + } + json.Unmarshal(raw, &r) + return r.Purged, nil +} + +type procDocMemory struct{} + +func (procDocMemory) Query(text string, topK int) []*sdk.Doc { + raw, err := callCore("doc.query", map[string]interface{}{"text": text, "top_k": topK}) + if err != nil { + return nil + } + var r struct { + Docs []*sdk.Doc `json:"docs"` + } + json.Unmarshal(raw, &r) + return r.Docs +} +func (procDocMemory) Insert(d *sdk.Doc) error { + return callCoreVoid("doc.insert", map[string]interface{}{"doc": d}) +} +func (procDocMemory) Remove(id string) { + callCoreVoid("doc.remove", map[string]string{"id": id}) +} +func (procDocMemory) Stats() map[string]interface{} { + raw, err := callCore("doc.stats", nil) + if err != nil { + return nil + } + var m map[string]interface{} + json.Unmarshal(raw, &m) + return m +} + +type procKnowledge struct{} + +func (procKnowledge) Search(q string, topK int) ([]*sdk.Knowledge, error) { + raw, err := callCore("knowledge.search", map[string]interface{}{"query": q, "top_k": topK}) + if err != nil { + return nil, err + } + var r struct { + Results []*sdk.Knowledge `json:"results"` + } + json.Unmarshal(raw, &r) + return r.Results, nil +} +func (procKnowledge) Add(name, content string) error { + return callCoreVoid("knowledge.add", map[string]string{"name": name, "content": content}) +} +func (procKnowledge) List() ([]string, error) { + raw, err := callCore("knowledge.list", nil) + if err != nil { + return nil, err + } + var r struct { + Names []string `json:"names"` + } + json.Unmarshal(raw, &r) + return r.Names, nil +} + +type procTextMemory struct{} + +func (procTextMemory) Append(evt sdk.TextEvent) error { + return callCoreVoid("textmemory.append", map[string]interface{}{"event": evt}) +} + +type procLLM struct{} + +func (procLLM) ListSources() []string { + raw, err := callCore("llm.listSources", nil) + if err != nil { + return nil + } + var r struct { + Sources []string `json:"sources"` + } + json.Unmarshal(raw, &r) + return r.Sources +} +func (procLLM) SetSource(name string) error { + return callCoreVoid("llm.setSource", map[string]string{"name": name}) +} +func (procLLM) CurrentSource() string { + raw, err := callCore("llm.currentSource", nil) + if err != nil { + return "" + } + var r struct { + Source string `json:"source"` + } + json.Unmarshal(raw, &r) + return r.Source +} + +type procSocial struct{} + +func (procSocial) GetPerson(name string) (*sdk.PersonProfile, error) { + raw, err := callCore("social.getPerson", map[string]string{"name": name}) + if err != nil { + return nil, err + } + var r struct { + Person *sdk.PersonProfile `json:"person"` + } + json.Unmarshal(raw, &r) + return r.Person, nil +} +func (procSocial) GetTrait(name, trait string) (string, bool) { + raw, err := callCore("social.getTrait", map[string]string{"name": name, "trait": trait}) + if err != nil { + return "", false + } + var r struct { + Value string `json:"value"` + Found bool `json:"found"` + } + json.Unmarshal(raw, &r) + return r.Value, r.Found +} +func (procSocial) GetRelations(name string) ([]sdk.SocialRelation, error) { + raw, err := callCore("social.getRelations", map[string]string{"name": name}) + if err != nil { + return nil, err + } + var r struct { + Relations []sdk.SocialRelation `json:"relations"` + } + json.Unmarshal(raw, &r) + return r.Relations, nil +} +func (procSocial) GetNetwork(name string, depth int) ([]*sdk.PersonProfile, error) { + raw, err := callCore("social.getNetwork", map[string]interface{}{"name": name, "depth": depth}) + if err != nil { + return nil, err + } + var r struct { + Network []*sdk.PersonProfile `json:"network"` + } + json.Unmarshal(raw, &r) + return r.Network, nil +} +func (procSocial) ListPersons() ([]string, error) { + raw, err := callCore("social.listPersons", nil) + if err != nil { + return nil, err + } + var r struct { + Persons []string `json:"persons"` + } + json.Unmarshal(raw, &r) + return r.Persons, nil +} + +type procPluginMgr struct{} + +func (procPluginMgr) ReloadOne(name string) error { + return callCoreVoid("plugin.reloadOne", map[string]string{"name": name}) +} +func (procPluginMgr) ListLoadedPlugins() []string { + raw, err := callCore("plugin.listLoaded", nil) + if err != nil { + return nil + } + var r struct { + Plugins []string `json:"plugins"` + } + json.Unmarshal(raw, &r) + return r.Plugins +} +func (procPluginMgr) IsPluginDisabled(name string) bool { + raw, err := callCore("plugin.isDisabled", map[string]string{"name": name}) + if err != nil { + return false + } + var r struct { + Disabled bool `json:"disabled"` + } + json.Unmarshal(raw, &r) + return r.Disabled +} + +type procSettings struct{} + +func (procSettings) Get(key string) (interface{}, error) { + return settingsValue("settings.get", map[string]string{"key": key}) +} +func (procSettings) Set(key string, v interface{}) error { + return callCoreVoid("settings.set", map[string]interface{}{"key": key, "value": v}) +} +func (procSettings) List(prefix string) ([]string, error) { + return settingsKeys("settings.list", map[string]string{"prefix": prefix}) +} +func (procSettings) GetCore(key string) (interface{}, error) { + return settingsValue("settings.getCore", map[string]string{"key": key}) +} +func (procSettings) SetCore(key string, v interface{}) error { + return callCoreVoid("settings.setCore", map[string]interface{}{"key": key, "value": v}) +} +func (procSettings) ListCore(prefix string) ([]string, error) { + return settingsKeys("settings.listCore", map[string]string{"prefix": prefix}) +} +func (procSettings) DataDir() string { + raw, err := callCore("settings.dataDir", nil) + if err != nil { + return "" + } + var r struct { + Dir string `json:"dir"` + } + json.Unmarshal(raw, &r) + return r.Dir +} +func (procSettings) GetPlugin(plugin, key string) (interface{}, error) { + return settingsValue("settings.getPlugin", map[string]string{"plugin": plugin, "key": key}) +} +func (procSettings) SetPlugin(plugin, key string, v interface{}) error { + return callCoreVoid("settings.setPlugin", map[string]interface{}{ + "plugin": plugin, "key": key, "value": v, + }) +} +func (procSettings) ListPlugin(plugin, prefix string) ([]string, error) { + return settingsKeys("settings.listPlugin", map[string]string{"plugin": plugin, "prefix": prefix}) +} +func (procSettings) RegisterDef(def sdk.ConfigDef) { + callCoreVoid("settings.registerDef", map[string]interface{}{"def": def}) +} +func (procSettings) Defs(prefix string) []*sdk.ConfigDef { + raw, err := callCore("settings.defs", map[string]string{"prefix": prefix}) + if err != nil { + return nil + } + var r struct { + Defs []*sdk.ConfigDef `json:"defs"` + } + json.Unmarshal(raw, &r) + return r.Defs +} +func (procSettings) Dump() map[string]interface{} { + raw, err := callCore("settings.dump", nil) + if err != nil { + return nil + } + var m map[string]interface{} + json.Unmarshal(raw, &m) + return m +} +func (procSettings) Plugins() []string { + raw, err := callCore("settings.plugins", nil) + if err != nil { + return nil + } + var r struct { + Plugins []string `json:"plugins"` + } + json.Unmarshal(raw, &r) + return r.Plugins +} + +func settingsValue(method string, params interface{}) (interface{}, error) { + raw, err := callCore(method, params) + if err != nil { + return nil, err + } + var r struct { + Value interface{} `json:"value"` + } + if err := json.Unmarshal(raw, &r); err != nil { + return nil, err + } + return r.Value, nil +} + +func settingsKeys(method string, params interface{}) ([]string, error) { + raw, err := callCore(method, params) + if err != nil { + return nil, err + } + var r struct { + Keys []string `json:"keys"` + } + if err := json.Unmarshal(raw, &r); err != nil { + return nil, err + } + return r.Keys, nil +} + +// ---- 内核 → 插件的调用处理 ---- + +func handleKernelRequest(req *rpcRequest) { + defer func() { + if r := recover(); r != nil { + // handler panic 只影响本次调用,不带崩进程; + // 真崩溃时进程退出,内核经 EOF 感知并按 recordCrash 处理。 + if req.ID != 0 { + respondErr(req.ID, fmt.Errorf("插件 handler panic: %v", r)) + } + log.Printf("handler panic (%s): %v", req.Method, r) + } + }() + + switch req.Method { + case "handshake": + handleHandshake(req) + + case "plugin.init": + var p struct { + Name string `json:"name"` + Config map[string]interface{} `json:"config"` + } + json.Unmarshal(req.Params, &p) + if p.Name != "" { + pluginName = p.Name + } + instance, err := NewPluginFactory(pluginName, p.Config) + if err != nil { + respondErr(req.ID, err) + return + } + plg = instance + respond(req.ID, nil) + + case "plugin.start": + if plg == nil { + respondErr(req.ID, fmt.Errorf("plugin.start 前未 init")) + return + } + pluginSDK = buildPluginSDK(pluginName) + if err := plg.Start(pluginSDK); err != nil { + respondErr(req.ID, err) + return + } + // 上报 AutoRestart:公开 SDK 的 SetAutoRestart 是纯 setter(无回调 hook), + // 插件在 Start() 里调它只改进程内副本。C ABI 路径下内核在 Start 返回后 + // 直接读 plgSDK.AutoRestart();子进程隔着进程边界读不到,故在此显式上报。 + // **不改公开 SDK 接口**(接口冻结约束)。 + if err := callCoreVoid("lifecycle.autoRestart", map[string]interface{}{ + "enabled": pluginSDK.AutoRestart(), + }); err != nil { + log.Printf("上报 autoRestart 失败: %v", err) + } + respond(req.ID, nil) + + case "plugin.stop": + if pluginSDK != nil { + pluginSDK.RunStopHandlers() + } + if plg != nil { + if err := plg.Stop(); err != nil { + log.Printf("Stop: %v", err) + } + } + respond(req.ID, nil) + stdoutW.Flush() + os.Exit(0) + + case "tool.invoke": + var p struct { + Name string `json:"name"` + Args map[string]interface{} `json:"args"` + } + json.Unmarshal(req.Params, &p) + handlerMu.RLock() + h, ok := toolHandlers[p.Name] + handlerMu.RUnlock() + if !ok { + respondErr(req.ID, fmt.Errorf("未注册的工具: %s", p.Name)) + return + } + res, err := h(p.Args) + if err != nil { + respondErr(req.ID, err) + return + } + respond(req.ID, map[string]interface{}{"result": res}) + + case "stage.invoke": + handleStageInvoke(req) + + case "output.invoke": + var p struct { + Channel string `json:"channel"` + Args map[string]interface{} `json:"args"` + } + json.Unmarshal(req.Params, &p) + handlerMu.RLock() + h, ok := outputHandlers[p.Channel] + handlerMu.RUnlock() + if !ok { + respondErr(req.ID, fmt.Errorf("未注册的输出通道: %s", p.Channel)) + return + } + // 同步返回真实结果——内核据此告知模型成功/失败,不再假成功(§9.4) + res, err := h(p.Args) + if err != nil { + respondErr(req.ID, err) + return + } + if m, ok := res.(map[string]interface{}); ok { + respond(req.ID, m) + return + } + respond(req.ID, map[string]interface{}{"status": "sent"}) + + case "events.subscribe": + var p struct { + Types []string `json:"types"` + } + json.Unmarshal(req.Params, &p) + // 订阅逻辑由内核 EventRing 处理,插件侧在此注册本地 handler。 + // 实际事件到达时由 evtConsumerLoop 分发。 + for _, t := range p.Types { + var idx uint32 + switch t { + case "raw_input": + idx = evtTypeRawInput + case "agent_output": + idx = evtTypeAgentOutput + case "agent_llm_chain": + idx = evtTypeAgentLLMChain + case "tool_call": + idx = evtTypeToolCall + case "reasoning": + idx = evtTypeReasoning + case "stage": + idx = evtTypeStage + case "system": + idx = evtTypeSystem + case "reasoning_delta": + idx = evtTypeReasoningDelta + case "content_delta": + idx = evtTypeContentDelta + case "skill_detected": + idx = evtTypeSkillDetected + default: + continue + } + evtHandlerMu.Lock() + evtHandlers[idx] = func(evt *sdk.Event) {} + evtHandlerMu.Unlock() + } + respond(req.ID, nil) + + case "events.unsubscribe": + // 清空全部 handler(子进程 Stop 时由内核统一清理订阅) + evtHandlerMu.Lock() + evtHandlers = map[uint32]func(*sdk.Event){} + evtHandlerMu.Unlock() + respond(req.ID, nil) + + default: + if req.ID != 0 { + respondErr(req.ID, fmt.Errorf("未实现的 method: %s", req.Method)) + } + } +} + +func handleHandshake(req *rpcRequest) { + var p struct { + Protocol int `json:"protocol"` + ShmVersion uint32 `json:"shm_version"` + ShmSize int `json:"shm_size"` + PluginName string `json:"plugin_name"` + EvtRingSize int `json:"evt_ring_size,omitempty"` + } + json.Unmarshal(req.Params, &p) + + if p.Protocol != procProtocolVersion { + respondErr(req.ID, fmt.Errorf("协议版本不匹配(内核 %d,插件 %d)——请用配套 plugindev 重编", + p.Protocol, procProtocolVersion)) + return + } + if p.ShmVersion != shmVersion { + respondErr(req.ID, fmt.Errorf("共享段版本不匹配(内核 %d,插件 %d)", p.ShmVersion, shmVersion)) + return + } + if p.PluginName != "" { + pluginName = p.PluginName + } + + // 挂载 StageContext 共享段。 + // 传递机制按平台不同(Unix 用继承的 fd,Windows 用命名段), + // 由 z_proc_shm_*.go 承担——本文件保持平台无关。 + if p.ShmSize > 0 { + m, err := attachStageShm(p.ShmSize) + if err != nil { + respondErr(req.ID, fmt.Errorf("挂载共享段失败: %w", err)) + return + } + if got := binary.LittleEndian.Uint32(m[shmOffMagic:]); got != shmMagic { + respondErr(req.ID, fmt.Errorf("共享段魔数不匹配(0x%x)", got)) + return + } + shm = m + } + + // 挂载事件环段 + 打开通知句柄(§3.6) + if p.EvtRingSize > 0 { + er, err := attachEvtRingShm(p.EvtRingSize) + if err != nil { + respondErr(req.ID, fmt.Errorf("挂载事件环段失败: %w", err)) + return + } + if got := binary.LittleEndian.Uint32(er[evtOffMagic : evtOffMagic+4]); got != evtRingMagic { + respondErr(req.ID, fmt.Errorf("事件环魔数不匹配(0x%x)", got)) + return + } + notifier, err := openEvtNotifier() + if err != nil { + respondErr(req.ID, fmt.Errorf("打开事件通知句柄失败: %w", err)) + return + } + evtRingData = er + evtNotifier = notifier + go evtConsumerLoop() + } + + respond(req.ID, map[string]interface{}{ + "protocol": procProtocolVersion, + "sdk_version": sdk.SDKVersion, + "plugin_name": pluginName, + "pid": os.Getpid(), + }) +} + +// handleStageInvoke 执行阶段处理器:拿锁 → 读共享段 → handler → 只写脏字段 → 放锁。 +// +// 插件作者的 handler 与 .so 时代完全一致(仍是 func(ctx *sdk.StageContext) error), +// 共享内存与锁的复杂度全部由本模板承担(§3.4)。 +func handleStageInvoke(req *rpcRequest) { + var p struct { + Stage string `json:"stage"` + Seq uint64 `json:"seq"` + } + json.Unmarshal(req.Params, &p) + + handlerMu.RLock() + h, ok := stageHandlers[p.Stage] + handlerMu.RUnlock() + if !ok { + respond(req.ID, map[string]interface{}{"dirty_fields": 0}) + return + } + if shm == nil { + respondErr(req.ID, fmt.Errorf("共享段未挂载")) + return + } + + // 跨进程写锁:内核仲裁(§3.7),持锁进程崩溃由内核代为释放 + if err := callCoreVoid("stage.lock", nil); err != nil { + respondErr(req.ID, fmt.Errorf("申请 stage 锁: %w", err)) + return + } + unlocked := false + unlock := func() { + if !unlocked { + unlocked = true + if err := callCoreVoid("stage.unlock", nil); err != nil { + log.Printf("释放 stage 锁: %v", err) + } + } + } + defer unlock() + + sc, err := readStageContext() + if err != nil { + respondErr(req.ID, fmt.Errorf("读共享段: %w", err)) + return + } + snap := takeStageSnapshot(sc) + + if err := h(sc); err != nil { + respondErr(req.ID, err) + return + } + + dirty, err := writeStageDirty(sc, snap) + if err != nil { + respondErr(req.ID, fmt.Errorf("写回共享段: %w", err)) + return + } + unlock() + respond(req.ID, map[string]interface{}{"dirty_fields": dirty, "seq": p.Seq}) +} + +// ---- 主循环 ---- + +func main() { + // 日志走 stderr:stdout 是 RPC 通道,写日志会破坏帧 + log.SetOutput(os.Stderr) + log.SetPrefix("[plugin] ") + + in := bufio.NewScanner(bufio.NewReader(os.Stdin)) + // 单帧上限 1MB:控制面帧本应很小,大 payload 走共享段 + in.Buffer(make([]byte, 0, 64*1024), 1024*1024) + + for in.Scan() { + line := make([]byte, len(in.Bytes())) + copy(line, in.Bytes()) + + var probe struct { + ID uint64 `json:"id"` + Method string `json:"method"` + } + if err := json.Unmarshal(line, &probe); err != nil { + log.Printf("非法 JSON 帧: %v", err) + continue + } + + // method 为空 = 内核对我们反向调用的应答 + if probe.Method == "" { + var resp rpcResponse + if err := json.Unmarshal(line, &resp); err != nil { + continue + } + pendMu.Lock() + ch, ok := pending[resp.ID] + delete(pending, resp.ID) + pendMu.Unlock() + if ok { + ch <- resp + } + continue + } + + var req rpcRequest + if err := json.Unmarshal(line, &req); err != nil { + continue + } + // 每个请求独立 goroutine:handler 内可能反向调用内核, + // 在读循环里同步处理会死锁(等应答但没人读)。 + go handleKernelRequest(&req) + } + + if err := in.Err(); err != nil { + log.Printf("读 stdin 出错: %v", err) + } + // stdin 关闭 = 内核结束了我们 + if pluginSDK != nil { + pluginSDK.RunStopHandlers() + } + if plg != nil { + plg.Stop() + } +} + +// ---- 事件环消费(§3.6,子进程侧)---- + +// 事件环布局常量(与内核 internal/plugin/proc/evtring.go 一致)。 +const ( + evtOffMagic = 0 + evtOffVersion = 4 + evtOffWriteSeq = 8 + evtOffCap = 16 + evtOffSlots = 20 + + evtRingMagic = 0x48455654 // "HEVT" + + // 事件类型位索引(与内核 encodeEvtType 一致) + evtTypeRawInput = 0 + evtTypeAgentOutput = 1 + evtTypeAgentLLMChain = 2 + evtTypeToolCall = 3 + evtTypeReasoning = 4 + evtTypeStage = 5 + evtTypeSystem = 6 + evtTypeReasoningDelta = 7 + evtTypeContentDelta = 8 + evtTypeSkillDetected = 9 + evtTypeMax = 10 + + evtRingSlotLen uint32 = 32 + evtRingCap uint32 = 8192 +) + +// evtTypeNames 把位索引还原成 pubsdk.EventType 字符串。 +var evtTypeNames = [evtTypeMax]string{ + "raw_input", "agent_output", "agent_llm_chain", "tool_call", + "reasoning", "stage", "system", "reasoning_delta", + "content_delta", "skill_detected", +} + +// evtConsumerLoop 是事件环消费主循环:eventfd.Read(阻塞走 netpoller)→ drain events → 分发。 +// 每个插件进程启动一个 goroutine,与 RPC 主循环并行。 +func evtConsumerLoop() { + if evtNotifier == nil || evtRingData == nil { + return + } + buf := make([]byte, 8) // eventfd uint64 计数 + var readSeq uint64 + + for { + if err := evtNotifier.Wait(buf); err != nil { + continue + } + // 循环 drain 直到无新事件(eventfd 计数合并,一次 Read 处理全部) + for { + writeSeq := binary.LittleEndian.Uint64(evtRingData[evtOffWriteSeq:]) + if readSeq >= writeSeq { + break + } + cap := uint64(evtRingCap) + if writeSeq-readSeq > cap { + readSeq = writeSeq - cap // 跳到最旧可读事件 + } + idx := readSeq % cap + slotOff := evtOffSlots + uint32(idx)*evtRingSlotLen + seq := binary.LittleEndian.Uint64(evtRingData[slotOff:]) + etype := binary.LittleEndian.Uint32(evtRingData[slotOff+8:]) + off := binary.LittleEndian.Uint32(evtRingData[slotOff+12:]) + slen := binary.LittleEndian.Uint32(evtRingData[slotOff+16:]) + + if seq != readSeq { + // slot 被覆盖,跳到最新 + if writeSeq > cap { + readSeq = writeSeq - cap + } else { + readSeq = writeSeq + } + continue + } + + // 读载荷并分发给注册的 handler + if off > 0 && slen > 0 && uint64(off)+uint64(slen) <= uint64(len(evtRingData)) { + evtTypeStr := "" + if int(etype) < len(evtTypeNames) { + evtTypeStr = evtTypeNames[etype] + } + evtHandlerMu.RLock() + handler, ok := evtHandlers[etype] + evtHandlerMu.RUnlock() + if ok && handler != nil && evtTypeStr != "" { + evt := &sdk.Event{Type: sdk.EventType(evtTypeStr)} + // payload 非核心(多数订阅者只看类型),简化为不解析 JSON + _ = evtRingData[off : off+slen] + handler(evt) + } + } + readSeq++ + } + } +} + +// evtWaiter 抽象事件通知等待。 +// +// Unix:eventfd(Linux)/ pipe(macOS)的读端,阻塞 Read 走 netpoller。 +// Windows:命名 Event 对象,WaitForSingleObject。 +// 平台实现在 z_proc_shm_unix.go / z_proc_shm_windows.go。 +type evtWaiter interface { + // Wait 阻塞直到有新事件;buf 供实现复用(Unix 读 8 字节计数)。 + Wait(buf []byte) error +} diff --git a/tools/plugindev/templates/proc_shm_unix.go.tmpl b/tools/plugindev/templates/proc_shm_unix.go.tmpl new file mode 100644 index 0000000..1976b10 --- /dev/null +++ b/tools/plugindev/templates/proc_shm_unix.go.tmpl @@ -0,0 +1,62 @@ +//go:build linux || darwin || freebsd + +package main + +import ( + "fmt" + "os" + "syscall" +) + +// Unix 侧共享段挂载:内核经 ExtraFiles 传入继承的 fd。 +// +// fd 布局(与内核 internal/plugin/proc/plugin.go 的 ExtraFiles 顺序一致): +// +// fd 3 = StageContext 段(memfd / 已 unlink 的临时文件) +// fd 4 = 事件环段 +// fd 5 = 事件通知(Linux eventfd / macOS pipe 读端) +// +// 继承的 fd 无需文件名,也不残留——这是选 memfd 而非 /dev/shm 的原因。 +const ( + fdStageShm = 3 + fdEvtRingShm = 4 + fdEvtNotifier = 5 +) + +// attachStageShm 挂载 StageContext 共享段。 +// +// 各进程 mmap 到不同虚拟地址,段内一律用相对偏移而非指针,故仍能正确解引用 +// (实验 2 已验证父子 mmap 基址不同时偏移解引用正确)。 +func attachStageShm(size int) ([]byte, error) { + return syscall.Mmap(fdStageShm, 0, size, + syscall.PROT_READ|syscall.PROT_WRITE, syscall.MAP_SHARED) +} + +// attachEvtRingShm 挂载事件环段。 +func attachEvtRingShm(size int) ([]byte, error) { + return syscall.Mmap(fdEvtRingShm, 0, size, + syscall.PROT_READ|syscall.PROT_WRITE, syscall.MAP_SHARED) +} + +// openEvtNotifier 打开事件通知读端。 +func openEvtNotifier() (evtWaiter, error) { + f := os.NewFile(fdEvtNotifier, "evtnotify") + if f == nil { + return nil, fmt.Errorf("fd %d 不是有效的通知句柄", fdEvtNotifier) + } + return &unixEvtWaiter{f: f}, nil +} + +// unixEvtWaiter 用 eventfd/pipe 的阻塞 Read 等待通知。 +// +// os.NewFile 把 fd 注册进 runtime netpoller,Read 阻塞时只 park goroutine, +// 不占 OS 线程(实验 1:200 个等待者仅增 1 个 OS 线程)。 +// 反面对照是经 cgo 调 sem_wait——那会阻塞整个 M。 +type unixEvtWaiter struct { + f *os.File +} + +func (w *unixEvtWaiter) Wait(buf []byte) error { + _, err := w.f.Read(buf) + return err +} diff --git a/tools/plugindev/templates/proc_shm_windows.go.tmpl b/tools/plugindev/templates/proc_shm_windows.go.tmpl new file mode 100644 index 0000000..eeef77e --- /dev/null +++ b/tools/plugindev/templates/proc_shm_windows.go.tmpl @@ -0,0 +1,153 @@ +//go:build windows + +package main + +import ( + "fmt" + "os" + "syscall" + "unsafe" +) + +// Windows 侧共享段挂载:走命名对象而非继承 fd。 +// +// 为何不能照抄 Unix:Windows 没有 fd 继承语义,`ExtraFiles` 在 os/exec 的 +// Windows 实现里不被支持。等价机制是命名内核对象——父进程用 +// CreateFileMapping / CreateEvent 建带名字的对象,子进程按同名 Open 拿到同一对象。 +// +// 名字经环境变量传入(内核 internal/plugin/proc/plugin_windows.go 设置), +// 而不是硬编码:多个 homed 实例并存时不能撞名。 +// +// **这是 §9.2 的正解**:C ABI 时代 Windows 是第三套独立 ABI 实现, +// stage 只下发 3 个字段且完全没有写回,sanitizer 这类改写型插件静默失效。 +// 现在 Windows 与 Unix 共用同一份 RPC 逻辑与同一份共享段布局, +// 差异被收敛到本文件的三个函数里。 +const ( + envStageShmName = "HOMEAGENT_SHM_STAGE" + envEvtRingName = "HOMEAGENT_SHM_EVTRING" + envEvtEventName = "HOMEAGENT_EVT_EVENT" +) + +// Windows API 绑定:用 LazyDLL 而非 golang.org/x/sys/windows。 +// +// 原因:OpenFileMappingW / OpenEventW 未被标准库 syscall 包导出。 +// 引入 x/sys 会给**每个插件的 go.mod 加一个新依赖**, +// 而「外部插件零改动」是本次迁移的硬约束(插件仅依赖公开 SDK)。 +// LazyDLL 属于标准库 syscall,零新增依赖。 +var ( + kernel32 = syscall.NewLazyDLL("kernel32.dll") + procOpenFileMappingW = kernel32.NewProc("OpenFileMappingW") + procOpenEventW = kernel32.NewProc("OpenEventW") +) + +const ( + winEventModifyState = 0x0002 + winSynchronize = 0x00100000 +) + +// openFileMappingW 封装 OpenFileMappingW。 +func openFileMappingW(access uint32, inherit bool, name *uint16) (syscall.Handle, error) { + var inheritFlag uintptr + if inherit { + inheritFlag = 1 + } + r, _, err := procOpenFileMappingW.Call( + uintptr(access), inheritFlag, uintptr(unsafe.Pointer(name))) + if r == 0 { + return 0, err + } + return syscall.Handle(r), nil +} + +// openEventW 封装 OpenEventW。 +func openEventW(access uint32, inherit bool, name *uint16) (syscall.Handle, error) { + var inheritFlag uintptr + if inherit { + inheritFlag = 1 + } + r, _, err := procOpenEventW.Call( + uintptr(access), inheritFlag, uintptr(unsafe.Pointer(name))) + if r == 0 { + return 0, err + } + return syscall.Handle(r), nil +} + +// attachStageShm 按名字打开 StageContext 段并映射。 +func attachStageShm(size int) ([]byte, error) { + return openNamedMapping(os.Getenv(envStageShmName), size, "StageContext 段") +} + +// attachEvtRingShm 按名字打开事件环段并映射。 +func attachEvtRingShm(size int) ([]byte, error) { + return openNamedMapping(os.Getenv(envEvtRingName), size, "事件环段") +} + +// openNamedMapping 打开命名共享段并映射为 []byte。 +// +// 与 Unix 的 mmap 语义对齐:MapViewOfFile 返回的地址在本进程虚拟空间, +// 段内偏移仍是相对的,故跨进程解引用正确。 +func openNamedMapping(name string, size int, what string) ([]byte, error) { + if name == "" { + return nil, fmt.Errorf("%s 名字未经环境变量传入", what) + } + namePtr, err := syscall.UTF16PtrFromString(name) + if err != nil { + return nil, fmt.Errorf("%s 名字非法: %w", what, err) + } + + h, err := openFileMappingW(syscall.FILE_MAP_WRITE, false, namePtr) + if err != nil { + return nil, fmt.Errorf("打开 %s(%s): %w", what, name, err) + } + + addr, err := syscall.MapViewOfFile(h, syscall.FILE_MAP_WRITE, 0, 0, uintptr(size)) + if err != nil { + syscall.CloseHandle(h) + return nil, fmt.Errorf("映射 %s: %w", what, err) + } + // 句柄不关:视图存活期间必须保持句柄有效,进程退出时由 OS 回收。 + + return unsafe.Slice((*byte)(unsafe.Pointer(addr)), size), nil +} + +// openEvtNotifier 按名字打开事件通知对象。 +func openEvtNotifier() (evtWaiter, error) { + name := os.Getenv(envEvtEventName) + if name == "" { + return nil, fmt.Errorf("事件通知对象名字未经环境变量传入") + } + namePtr, err := syscall.UTF16PtrFromString(name) + if err != nil { + return nil, fmt.Errorf("事件对象名字非法: %w", err) + } + h, err := openEventW(winSynchronize|winEventModifyState, false, namePtr) + if err != nil { + return nil, fmt.Errorf("打开事件对象(%s): %w", name, err) + } + return &windowsEvtWaiter{h: h}, nil +} + +// windowsEvtWaiter 用命名 Event 对象等待通知。 +// +// 与 eventfd 的差异:Event 是二元信号而非计数器,多次 SetEvent 只对应 +// 一次唤醒。这不影响正确性——消费者被唤醒后按 readSeq 追 writeSeq +// 批量 drain,一次唤醒能处理累积的全部事件。 +// +// WaitForSingleObject 阻塞的是 OS 线程而非仅 goroutine,故不如 eventfd +// 的 netpoller 路径省线程。每插件一个消费 goroutine,17 插件即 17 线程, +// 在可接受范围(实验 5 实测 17 子进程共 84 线程)。 +type windowsEvtWaiter struct { + h syscall.Handle +} + +func (w *windowsEvtWaiter) Wait(buf []byte) error { + ev, err := syscall.WaitForSingleObject(w.h, syscall.INFINITE) + if err != nil { + return err + } + if ev != syscall.WAIT_OBJECT_0 { + return fmt.Errorf("等待事件对象返回 0x%x", ev) + } + return nil +}