From d4cb9cb8ccf80c901c22be6451de810e5c8c139f Mon Sep 17 00:00:00 2001 From: jianf <2198972886@qq.com> Date: Thu, 23 Jul 2026 16:34:26 +0800 Subject: [PATCH] refactor(clawhubadapter): merge manager/simulator, fix provider dispatch, add OC plugin load delegation - Manager now handles OC plugin loading via plugins/load JSON-RPC instead of launching a separate simulator process - Added 8 missing provider registration methods to manager - Added explicit registeredChannels declaration in manager - Provider dispatch now strips _provider suffix correctly (lookupType) - Skip .simulator directory in skills scanning - detect OC via package.json openclaw field (hasOCPackage) - All registration notifications include plugin name for correct tool prefix - Fix ToolRegistry/ProviderRegistry/ChannelRegistry to use plugin name from notification data instead of hardcoded 'manager' - test: add mockSettings, fix TestLoadOCPluginViaPluginStart panic --- .../plugins/clawhubadapter/manager/main.js | 69 +++++- .../clawhubadapter/manager/openclaw_cli.js | 22 -- internal/plugins/clawhubadapter/plugin.go | 222 +++++++++++------- internal/plugins/clawhubadapter/registry.go | 94 +++++++- .../plugins/clawhubadapter/sidecar_test.go | 19 ++ .../plugins/clawhubadapter/simulator/main.js | 53 ++++- .../clawhubadapter/simulator/openclaw_cli.js | 216 +++++++++++++++++ 7 files changed, 574 insertions(+), 121 deletions(-) delete mode 100644 internal/plugins/clawhubadapter/manager/openclaw_cli.js create mode 100644 internal/plugins/clawhubadapter/simulator/openclaw_cli.js diff --git a/internal/plugins/clawhubadapter/manager/main.js b/internal/plugins/clawhubadapter/manager/main.js index 59b4e01..1226b18 100644 --- a/internal/plugins/clawhubadapter/manager/main.js +++ b/internal/plugins/clawhubadapter/manager/main.js @@ -93,13 +93,14 @@ function tryDetectOCPackage(pkgDir, pkgName) { const loadedPlugins = {}; // name -> { entry, tools: [{name, execute, ...}] } const allTools = []; // flat list of all tools across all plugins const allProviders = {}; // type -> { name, instance } across all plugins +const registeredChannels = {}; // name -> { pluginName, channelPlugin, output, send, type } function registerPluginTools(name, tools, api) { for (const t of tools) { if (t && t.name) { t._plugin = name; allTools.push(t); - notify('register', { type: 'tool', data: { name: t.name, description: t.description, parameters: t.parameters } }); + notify('register', { type: 'tool', data: { name: t.name, description: t.description, parameters: t.parameters, plugin: name } }); } } loadedPlugins[name] = { tools, api }; @@ -174,7 +175,7 @@ function loadPlugin(pluginDir, name) { }, registerProvider: (p) => { if (p && p.id) allProviders['llm'] = { name: p.id, instance: p }; - notify('register', { type: 'provider', data: { name: p?.id || p?.name } }); + notify('register', { type: 'provider', data: { name: p?.id || p?.name, plugin: name } }); }, registerChannel: (ch) => { let chName = ch.name; @@ -191,7 +192,7 @@ function loadPlugin(pluginDir, name) { registeredChannels[chName] = { pluginName: name, output: ch.output || ch.send, type: chType }; } - notify('register', { type: 'channel', data: { name: chName, type: chType, id: chPlugin?.id } }); + notify('register', { type: 'channel', data: { name: chName, type: chType, id: chPlugin?.id, plugin: name } }); }, submitInput: (msg) => { notify('channel_input', { channel: name, payload: msg }); @@ -202,15 +203,47 @@ function loadPlugin(pluginDir, name) { registerService: (svc) => notify('register', { type: 'service', data: { name: svc.name } }), registerImageGenerationProvider: (p) => { if (p) allProviders['image_generation'] = { name: p.name, instance: p }; - notify('register', { type: 'image_generation_provider', data: { name: p?.name } }); + notify('register', { type: 'image_generation_provider', data: { name: p?.name, plugin: name } }); }, registerWebFetchProvider: (p) => { if (p) allProviders['web_fetch'] = { name: p.name, instance: p }; - notify('register', { type: 'web_fetch_provider', data: { name: p?.name } }); + notify('register', { type: 'web_fetch_provider', data: { name: p?.name, plugin: name } }); }, registerWebSearchProvider: (p) => { if (p) allProviders['web_search'] = { name: p.name, instance: p }; - notify('register', { type: 'web_search_provider', data: { name: p?.name } }); + notify('register', { type: 'web_search_provider', data: { name: p?.name, plugin: name } }); + }, + registerSpeechProvider: (p) => { + if (p) allProviders['speech'] = { name: p.name, instance: p }; + notify('register', { type: 'speech_provider', data: { name: p?.name, plugin: name } }); + }, + registerRealtimeTranscriptionProvider: (p) => { + if (p) allProviders['realtime_transcription'] = { name: p.name, instance: p }; + notify('register', { type: 'realtime_transcription_provider', data: { name: p?.name, plugin: name } }); + }, + registerRealtimeVoiceProvider: (p) => { + if (p) allProviders['realtime_voice'] = { name: p.name, instance: p }; + notify('register', { type: 'realtime_voice_provider', data: { name: p?.name, plugin: name } }); + }, + registerMediaUnderstandingProvider: (p) => { + if (p) allProviders['media_understanding'] = { name: p.name, instance: p }; + notify('register', { type: 'media_understanding_provider', data: { name: p?.name, plugin: name } }); + }, + registerMusicGenerationProvider: (p) => { + if (p) allProviders['music_generation'] = { name: p.name, instance: p }; + notify('register', { type: 'music_generation_provider', data: { name: p?.name, plugin: name } }); + }, + registerVideoGenerationProvider: (p) => { + if (p) allProviders['video_generation'] = { name: p.name, instance: p }; + notify('register', { type: 'video_generation_provider', data: { name: p?.name, plugin: name } }); + }, + registerEmbeddingProvider: (p) => { + if (p) allProviders['embedding'] = { name: p.name, instance: p }; + notify('register', { type: 'embedding_provider', data: { name: p?.name, plugin: name } }); + }, + registerMemoryEmbeddingProvider: (p) => { + if (p) allProviders['memory_embedding'] = { name: p.name, instance: p }; + notify('register', { type: 'memory_embedding_provider', data: { name: p?.name, plugin: name } }); }, start: (cb) => {}, stop: (cb) => {}, @@ -224,6 +257,8 @@ function loadPlugin(pluginDir, name) { // ---- Install npm package ---- function installNPMPackage(spec, skillsDir) { + // Strip npm: prefix if present + if (spec.startsWith('npm:')) spec = spec.slice(4); process.stderr.write(`[manager] installing: ${spec}\n`); const installDir = path.join(skillsDir, '.npm_install_' + Date.now()); fs.mkdirSync(installDir, { recursive: true }); @@ -342,6 +377,9 @@ async function runCLI(skillsDir, cliArgs) { const result = installNPMPackage(spec, skillsDir); if (result.error) throw new Error(result.error); console.log(`Installed: ${result.name}`); + // Write install result for CLI wrapper to pick up + const simDir = path.dirname(process.argv[1]); + fs.writeFileSync(path.join(simDir, '.install-result'), JSON.stringify({ name: result.name })); break; } @@ -640,6 +678,25 @@ rl.on('line', async (line) => { return; } + if (method === 'plugins/load') { + const dir = req.params?.dir; + const pluginName = req.params?.name || (dir ? path.basename(dir) : ''); + if (!dir) { sendError(id, -32602, 'dir required'); return; } + const resolvedDir = path.resolve(dir); + if (!fs.existsSync(resolvedDir)) { sendError(id, -32601, `directory not found: ${resolvedDir}`); return; } + + if (loadedPlugins[pluginName]) { + writeJSON({ jsonrpc: '2.0', id, result: { name: pluginName, tools: loadedPlugins[pluginName].tools.map(t => t.name), type: 'already_loaded' } }); + return; + } + + const ok = loadPlugin(resolvedDir, pluginName); + if (!ok) { sendError(id, -32603, `failed to load plugin: ${pluginName}`); return; } + + writeJSON({ jsonrpc: '2.0', id, result: { name: pluginName, tools: loadedPlugins[pluginName].tools.map(t => t.name), type: 'loaded' } }); + return; + } + if (method === 'tools/list') { const tools = allTools.map(t => ({ name: t.name, diff --git a/internal/plugins/clawhubadapter/manager/openclaw_cli.js b/internal/plugins/clawhubadapter/manager/openclaw_cli.js deleted file mode 100644 index f56f90a..0000000 --- a/internal/plugins/clawhubadapter/manager/openclaw_cli.js +++ /dev/null @@ -1,22 +0,0 @@ -#!/usr/bin/env node -const path = require('path'); -const fs = require('fs'); -const { fork } = require('child_process'); - -const binDir = __dirname; -const skillsdirPath = path.join(binDir, '..', '.skillsdir'); - -if (!fs.existsSync(skillsdirPath)) { - process.stderr.write('Error: OpenClaw manager not running (no .skillsdir found)\n'); - process.exit(1); -} - -const skillsDir = fs.readFileSync(skillsdirPath, 'utf8').trim(); -const managerPath = path.join(binDir, '..', 'manager.js'); - -const proc = fork(managerPath, [skillsDir, ...process.argv.slice(2)], { - stdio: 'inherit', - env: { ...process.env, OPENCLAW_CLI: '1' }, -}); - -proc.on('exit', (code) => process.exit(code)); diff --git a/internal/plugins/clawhubadapter/plugin.go b/internal/plugins/clawhubadapter/plugin.go index 6dd3d2a..9895ae1 100644 --- a/internal/plugins/clawhubadapter/plugin.go +++ b/internal/plugins/clawhubadapter/plugin.go @@ -16,21 +16,19 @@ import ( "path/filepath" "strings" "sync" + "time" "gitcode.com/JianFeeeee/HomeAgent/internal/plugin" sdk "gitcode.com/JianFeeeee/HomeAgent/internal/sdk" ) -//go:embed simulator/main.js -var simulatorSrc string - //go:embed manager/main.js var managerSrc string //go:embed pysimulator/main.py var pySimulatorSrc string -//go:embed manager/openclaw_cli.js +//go:embed simulator/openclaw_cli.js var openclawCliSrc string var SkillsDir string @@ -112,6 +110,9 @@ func (p *Plugin) Start(s *sdk.PluginSDK) error { // Load existing plugins from skills dir if entries, err := os.ReadDir(p.skillsDir); err == nil { for _, entry := range entries { + if strings.HasPrefix(entry.Name(), ".") { + continue + } skillPath := filepath.Join(p.skillsDir, entry.Name()) subs, err := os.ReadDir(skillPath) if err != nil { @@ -205,9 +206,58 @@ func (p *Plugin) Start(s *sdk.PluginSDK) error { }, }, p.handlePluginList) + // Start file-based IPC for CLI integration (settings sync + reload requests) + go p.ipcGoroutine(s) + return nil } +func (p *Plugin) ipcGoroutine(s *sdk.PluginSDK) { + ticker := time.NewTicker(3 * time.Second) + defer ticker.Stop() + + for range ticker.C { + p.mu.Lock() + simDir := p.simulatorDir + p.mu.Unlock() + if simDir == "" { + continue + } + + // 1. Write settings snapshot for CLI config:get/config:list + settings := s.Settings().Dump() + if data, err := json.MarshalIndent(settings, "", " "); err == nil { + os.WriteFile(filepath.Join(simDir, ".settings.json"), data, 0644) + } + + // 2. Process pending settings changes from CLI config:set + pendingPath := filepath.Join(simDir, ".settings-pending.json") + if data, err := os.ReadFile(pendingPath); err == nil { + var pending map[string]interface{} + if json.Unmarshal(data, &pending) == nil { + for k, v := range pending { + s.Settings().Set(k, v) + } + } + os.Remove(pendingPath) + } + + // 3. Process reload requests from CLI plugin:install + reloadPath := filepath.Join(simDir, ".reload-request") + if data, err := os.ReadFile(reloadPath); err == nil { + name := strings.TrimSpace(string(data)) + if name != "" { + if err := p.reloadPlugin(name); err != nil { + log.Printf("[clawhubadapter] reload from CLI: %v", err) + } else { + log.Printf("[clawhubadapter] reloaded plugin from CLI request: %s", name) + } + } + os.Remove(reloadPath) + } + } +} + // ─── Manager ──────────────────────────────────────────────── // launchManager 启动 OC 插件管理器(Node.js 进程),用于安装/卸载/加载 OC 格式插件 @@ -273,6 +323,10 @@ func (p *Plugin) handlePluginInstall(args map[string]interface{}) (interface{}, return p.installFromClawHub(pkg) } + if strings.HasPrefix(pkg, "npm:") { + pkg = strings.TrimPrefix(pkg, "npm:") + } + p.mu.Lock() mgr := p.manager p.mu.Unlock() @@ -549,62 +603,62 @@ func (p *Plugin) handlePluginList(args map[string]interface{}) (interface{}, err parts = append(parts, fmt.Sprintf(" %s v%s", sk.Name(), sk.Version())) } } - caps := p.dispatcher.Capabilities() - if len(caps) > 0 { - parts = append(parts, fmt.Sprintf("\nCapabilities (%d):", len(caps))) - for _, c := range caps { - parts = append(parts, fmt.Sprintf(" %s", c)) + caps := p.dispatcher.Capabilities() + if len(caps) > 0 { + parts = append(parts, fmt.Sprintf("\nCapabilities (%d):", len(caps))) + for _, c := range caps { + parts = append(parts, fmt.Sprintf(" %s", c)) + } } - } - if len(result.Plugins) == 0 && len(p.sidecars) <= 1 && len(p.skills) == 0 { - parts = append(parts, "没有已安装的插件。") - } - p.mu.Unlock() + if len(result.Plugins) == 0 && len(p.sidecars) <= 1 && len(p.skills) == 0 { + parts = append(parts, "没有已安装的插件。") + } + p.mu.Unlock() - return map[string]interface{}{ - "content": strings.Join(parts, "\n"), - }, nil + return map[string]interface{}{ + "content": strings.Join(parts, "\n"), + }, nil + } } } -} -// Fallback: list known plugins from Go side -p.mu.Lock() -defer p.mu.Unlock() + // Fallback: list known plugins from Go side + p.mu.Lock() + defer p.mu.Unlock() -var parts []string -parts = append(parts, fmt.Sprintf("Skills dir: %s\n", p.skillsDir)) + var parts []string + parts = append(parts, fmt.Sprintf("Skills dir: %s\n", p.skillsDir)) -if len(p.sidecars) > 0 { - parts = append(parts, fmt.Sprintf("\nSidecar/OC 插件 (%d):", len(p.sidecars))) - for _, sp := range p.sidecars { - tools, err := sp.ListTools() - toolList := "" - if err == nil { - var names []string - for _, t := range tools { - names = append(names, t.Name) + if len(p.sidecars) > 0 { + parts = append(parts, fmt.Sprintf("\nSidecar/OC 插件 (%d):", len(p.sidecars))) + for _, sp := range p.sidecars { + tools, err := sp.ListTools() + toolList := "" + if err == nil { + var names []string + for _, t := range tools { + names = append(names, t.Name) + } + toolList = strings.Join(names, ", ") } - toolList = strings.Join(names, ", ") + parts = append(parts, fmt.Sprintf(" %s: %s", sp.name, toolList)) } - parts = append(parts, fmt.Sprintf(" %s: %s", sp.name, toolList)) } -} -if len(p.skills) > 0 { - parts = append(parts, fmt.Sprintf("\nSKILL 插件 (%d):", len(p.skills))) - for _, sk := range p.skills { - parts = append(parts, fmt.Sprintf(" %s v%s", sk.Name(), sk.Version())) + if len(p.skills) > 0 { + parts = append(parts, fmt.Sprintf("\nSKILL 插件 (%d):", len(p.skills))) + for _, sk := range p.skills { + parts = append(parts, fmt.Sprintf(" %s v%s", sk.Name(), sk.Version())) + } } -} -caps := p.dispatcher.Capabilities() -if len(caps) > 0 { - parts = append(parts, fmt.Sprintf("\nCapabilities (%d):", len(caps))) - for _, c := range caps { - parts = append(parts, fmt.Sprintf(" %s", c)) + caps := p.dispatcher.Capabilities() + if len(caps) > 0 { + parts = append(parts, fmt.Sprintf("\nCapabilities (%d):", len(caps))) + for _, c := range caps { + parts = append(parts, fmt.Sprintf(" %s", c)) + } } -} if len(p.sidecars) == 0 && len(p.skills) == 0 { parts = append(parts, "没有已安装的插件。") @@ -616,48 +670,49 @@ if len(caps) > 0 { } func (p *Plugin) loadOCPlugin(s *sdk.PluginSDK, dir, name string) error { - simPath := filepath.Join(p.simulatorDir, "main.js") - if err := os.MkdirAll(p.simulatorDir, 0755); err != nil { - return fmt.Errorf("create simulator dir: %w", err) - } - if err := os.WriteFile(simPath, []byte(simulatorSrc), 0644); err != nil { - return fmt.Errorf("write simulator: %w", err) + // Delegate to manager's plugins/load instead of launching a separate simulator. + // The manager is the single Node.js process that handles all OC plugins. + p.mu.Lock() + mgr := p.manager + p.mu.Unlock() + if mgr == nil { + return fmt.Errorf("plugin manager not available") } - sp, err := launchProcess("node", simPath, dir, name) - if err != nil { - return fmt.Errorf("launch simulator: %w", err) - } - if sp == nil { - return nil - } - - // 通知是唯一注册路径。 - // 插件 register(api) 期间模拟器将所有 register* 调用以通知推送给 Go。 - // waitReady 之后所有初始化通知已缓冲在 notifyCh 中,同步排空处理。 - for done := false; !done; { - select { - case n := <-sp.NotifyChan(): - p.translateAndRegister(n, sp, s, name) - default: - done = true + // Check if manager already has this plugin loaded + data, listErr := mgr.call("plugins/list", nil) + if listErr == nil && data != nil { + var result struct { + Plugins []struct{ Name string `json:"name"` } `json:"plugins"` + } + if json.Unmarshal(data, &result) == nil { + for _, pl := range result.Plugins { + if pl.Name == name { + log.Printf("[clawhubadapter] plugin %s already loaded by manager, skipping", name) + return nil + } + } } } - // 启动持久通知协程:后续模拟器推送的注册通知持续转译注册到核心 - go p.notifyLoop(sp, s, name) - - // ListTools 仅验证日志,不参与注册 - tools, err := sp.ListTools() + // Ask manager to load this plugin (notifications flow through the manager's stdout → Go notifyLoop) + data, err := mgr.call("plugins/load", map[string]interface{}{ + "dir": dir, + "name": name, + }) if err != nil { - sp.Close() - return fmt.Errorf("list tools: %w", err) + return fmt.Errorf("manager plugins/load: %w", err) } - log.Printf("[clawhubadapter] ocplugin %s verified %d tools via ListTools", name, len(tools)) - p.mu.Lock() - p.sidecars = append(p.sidecars, sp) - p.mu.Unlock() + var result struct { + Name string `json:"name"` + Tools []string `json:"tools"` + } + if json.Unmarshal(data, &result) == nil { + log.Printf("[clawhubadapter] ocplugin %s loaded via manager with %d tools", name, len(result.Tools)) + } else { + log.Printf("[clawhubadapter] ocplugin %s loaded via manager, raw: %s", name, string(data)) + } return nil } @@ -678,7 +733,9 @@ func (p *Plugin) translateAndRegister(n OCNotification, sp *sidecarProcess, s *s if err := json.Unmarshal(n.Params, ¶ms); err != nil || params.Channel == "" { return } + channelInputBufMu.Lock() channelInputBuf[params.Channel] = append(channelInputBuf[params.Channel], params.Payload) + channelInputBufMu.Unlock() content, _ := params.Payload["content"].(string) if content == "" { data, _ := json.Marshal(params.Payload) @@ -848,6 +905,7 @@ func (p *Plugin) reloadPlugin(name string) error { hasMainJS := false hasMainPy := false hasOCManifest := false + hasOCPackage := false for _, f := range subs { switch f.Name() { case "main.js": @@ -856,6 +914,8 @@ func (p *Plugin) reloadPlugin(name string) error { hasMainPy = true case "openclaw.plugin.json": hasOCManifest = true + case "package.json": + hasOCPackage = hasOCExtensions(filepath.Join(pluginDir, "package.json")) } } @@ -866,7 +926,7 @@ func (p *Plugin) reloadPlugin(name string) error { return p.loadSidecar(p.sdk, pluginDir, name) case hasMainPy: return p.loadPySidecar(p.sdk, pluginDir, name) - case hasOCManifest: + case hasOCManifest || hasOCPackage: return p.loadOCPlugin(p.sdk, pluginDir, name) default: sk, err := plugin.LoadSKILL(pluginDir) diff --git a/internal/plugins/clawhubadapter/registry.go b/internal/plugins/clawhubadapter/registry.go index 5ffab8d..0ba13cc 100644 --- a/internal/plugins/clawhubadapter/registry.go +++ b/internal/plugins/clawhubadapter/registry.go @@ -10,7 +10,10 @@ import ( sdk "gitcode.com/JianFeeeee/HomeAgent/internal/sdk" ) -var channelInputBuf = map[string][]map[string]interface{}{} +var ( + channelInputBuf = map[string][]map[string]interface{}{} + channelInputBufMu sync.Mutex +) type ToolRegistry struct{} @@ -20,11 +23,16 @@ func (r *ToolRegistry) Dispatch(data json.RawMessage, pluginName string, sp *sid Label string `json:"label"` Description string `json:"description"` Parameters map[string]interface{} `json:"parameters"` + Plugin string `json:"plugin"` } if err := json.Unmarshal(data, &d); err != nil || d.Name == "" { return } - toolName := fmt.Sprintf("%s_%s", pluginName, d.Name) + pn := pluginName + if d.Plugin != "" { + pn = d.Plugin + } + toolName := fmt.Sprintf("%s_%s", pn, d.Name) tDef := sdk.ToolDef{ Name: toolName, Description: d.Description, @@ -64,29 +72,75 @@ func (r *ProviderRegistry) Dispatch(typeStr string, data json.RawMessage, plugin var d struct { Name string `json:"name"` Description string `json:"description"` + Plugin string `json:"plugin"` } json.Unmarshal(data, &d) + pn := pluginName + if d.Plugin != "" { + pn = d.Plugin + } + + // Strip _provider suffix: "image_generation_provider" -> "image_generation" + lookupType := typeStr + if lookupType != "provider" { + lookupType = strings.TrimSuffix(lookupType, "_provider") + } + if lookupType == "" { + return + } + for _, p := range providerMap { - if p.ocType == typeStr { - toolName := fmt.Sprintf("%s_%s", pluginName, p.toolSuffix) + if p.ocType == lookupType { + toolName := fmt.Sprintf("%s_%s", pn, p.toolSuffix) desc := p.desc if d.Name != "" { desc = fmt.Sprintf("[%s] %s", d.Name, desc) } + props := map[string]interface{}{ + "prompt": map[string]interface{}{"type": "string", "description": "Prompt for generation or search query"}, + } + switch p.ocType { + case "image_generation": + props["prompt"] = map[string]interface{}{"type": "string", "description": "Image description prompt"} + props["size"] = map[string]interface{}{"type": "string", "description": "Image size (e.g. 1024x1024)", "enum": []interface{}{"256x256", "512x512", "1024x1024", "1792x1024", "1024x1792"}} + case "web_search": + props["query"] = map[string]interface{}{"type": "string", "description": "Search query"} + delete(props, "prompt") + case "web_fetch": + props["url"] = map[string]interface{}{"type": "string", "description": "URL to fetch"} + delete(props, "prompt") + case "media_understanding": + props["url"] = map[string]interface{}{"type": "string", "description": "Media URL to analyze"} + props["media_type"] = map[string]interface{}{"type": "string", "description": "Media type", "enum": []interface{}{"image", "audio", "video"}} + case "speech": + props["text"] = map[string]interface{}{"type": "string", "description": "Text to synthesize"} + props["voice"] = map[string]interface{}{"type": "string", "description": "Voice identifier"} + case "music_generation": + props["prompt"] = map[string]interface{}{"type": "string", "description": "Music description prompt"} + props["duration"] = map[string]interface{}{"type": "number", "description": "Duration in seconds"} + case "video_generation": + props["prompt"] = map[string]interface{}{"type": "string", "description": "Video description prompt"} + props["duration"] = map[string]interface{}{"type": "number", "description": "Duration in seconds"} + case "realtime_transcription": + props["audio_url"] = map[string]interface{}{"type": "string", "description": "Audio URL to transcribe"} + case "realtime_voice": + props["text"] = map[string]interface{}{"type": "string", "description": "Text to speak"} + props["voice"] = map[string]interface{}{"type": "string", "description": "Voice identifier"} + } tDef := sdk.ToolDef{ Name: toolName, Description: desc, Parameters: map[string]interface{}{ "type": "object", - "properties": map[string]interface{}{}, + "properties": props, }, } handler := func(sp *sidecarProcess, ocType string) sdk.ToolHandler { return func(args map[string]interface{}) (interface{}, error) { return sp.CallProvider(ocType, args) } - }(sp, typeStr) + }(sp, lookupType) if err := s.RegisterTool(toolName, tDef, handler); err != nil { log.Printf("[clawhubadapter] register provider tool %s: %v", toolName, err) } @@ -94,20 +148,35 @@ func (r *ProviderRegistry) Dispatch(typeStr string, data json.RawMessage, plugin } } - log.Printf("[clawhubadapter] unknown provider type: %s (plugin: %s)", typeStr, pluginName) + if lookupType != typeStr { + log.Printf("[clawhubadapter] provider %s -> lookup %s (plugin: %s)", typeStr, lookupType, pluginName) + } + switch lookupType { + case "provider": + log.Printf("[clawhubadapter] generic LLM provider from %s handled natively by HomeAgent", pluginName) + case "embedding", "memory_embedding": + log.Printf("[clawhubadapter] embedding provider %s from %s ignored (HomeAgent uses native embeddings)", lookupType, pluginName) + default: + log.Printf("[clawhubadapter] unknown provider type: %s (plugin: %s)", typeStr, pluginName) + } } type ChannelRegistry struct{} func (r *ChannelRegistry) Dispatch(data json.RawMessage, pluginName string, sp *sidecarProcess, s *sdk.PluginSDK) { var d struct { - Name string `json:"name"` - Type string `json:"type"` + Name string `json:"name"` + Type string `json:"type"` + Plugin string `json:"plugin"` } if err := json.Unmarshal(data, &d); err != nil || d.Name == "" { return } chName := d.Name + pn := pluginName + if d.Plugin != "" { + pn = d.Plugin + } var caps int switch d.Type { @@ -120,12 +189,12 @@ func (r *ChannelRegistry) Dispatch(data json.RawMessage, pluginName string, sp * default: caps = 1 } - desc := fmt.Sprintf("OC channel %s (from %s)", chName, pluginName) + desc := fmt.Sprintf("OC channel %s (from %s)", chName, pn) s.RegisterOutputChannel(chName, caps, desc, func(args map[string]interface{}) (interface{}, error) { return sp.CallTool(chName, args) }) - readToolName := fmt.Sprintf("%s_read_%s_input", pluginName, strings.ReplaceAll(chName, "-", "_")) + readToolName := fmt.Sprintf("%s_read_%s_input", pn, strings.ReplaceAll(chName, "-", "_")) s.RegisterTool(readToolName, sdk.ToolDef{ Name: readToolName, Description: fmt.Sprintf("读取 %s 通道的待处理输入消息", chName), @@ -134,8 +203,10 @@ func (r *ChannelRegistry) Dispatch(data json.RawMessage, pluginName string, sp * "properties": map[string]interface{}{}, }, }, func(args map[string]interface{}) (interface{}, error) { + channelInputBufMu.Lock() buf := channelInputBuf[chName] if len(buf) == 0 { + channelInputBufMu.Unlock() return map[string]interface{}{"messages": []interface{}{}}, nil } msgs := make([]interface{}, len(buf)) @@ -143,6 +214,7 @@ func (r *ChannelRegistry) Dispatch(data json.RawMessage, pluginName string, sp * msgs[i] = m } channelInputBuf[chName] = nil + channelInputBufMu.Unlock() return map[string]interface{}{"messages": msgs}, nil }) } diff --git a/internal/plugins/clawhubadapter/sidecar_test.go b/internal/plugins/clawhubadapter/sidecar_test.go index a4df2ab..04fb02e 100644 --- a/internal/plugins/clawhubadapter/sidecar_test.go +++ b/internal/plugins/clawhubadapter/sidecar_test.go @@ -7,8 +7,26 @@ import ( "testing" sdk "gitcode.com/JianFeeeee/HomeAgent/internal/sdk" + pubsdk "gitcode.com/JianFeeeee/homeagent-sdk/sdk" ) +// mockSettings implements pubsdk.SettingsAPI for tests +type mockSettings struct{} + +func (m *mockSettings) Get(key string) (interface{}, error) { return nil, nil } +func (m *mockSettings) Set(key string, value interface{}) error { return nil } +func (m *mockSettings) List(prefix string) ([]string, error) { return nil, nil } +func (m *mockSettings) GetCore(key string) (interface{}, error) { return nil, nil } +func (m *mockSettings) SetCore(key string, value interface{}) error { return nil } +func (m *mockSettings) ListCore(prefix string) ([]string, error) { return nil, nil } +func (m *mockSettings) GetPlugin(plugin, key string) (interface{}, error) { return nil, nil } +func (m *mockSettings) SetPlugin(plugin, key string, value interface{}) error { return nil } +func (m *mockSettings) ListPlugin(plugin, prefix string) ([]string, error) { return nil, nil } +func (m *mockSettings) RegisterDef(def pubsdk.ConfigDef) {} +func (m *mockSettings) Defs(prefix string) []*pubsdk.ConfigDef { return nil } +func (m *mockSettings) Dump() map[string]interface{} { return nil } +func (m *mockSettings) Plugins() []string { return nil } + func TestLaunchSidecarNoMainJS(t *testing.T) { tmpDir := t.TempDir() sp, err := launchSidecar(tmpDir, "nonexistent") @@ -327,6 +345,7 @@ func TestLoadOCPluginViaPluginStart(t *testing.T) { var registeredTools []string registeredHandlers := make(map[string]sdk.ToolHandler) sdk := sdk.New("openclaw", sdk.SDKConfig{ + Settings: &mockSettings{}, RegTool: func(name string, def sdk.ToolDef, handler sdk.ToolHandler) error { registeredTools = append(registeredTools, name) registeredHandlers[name] = handler diff --git a/internal/plugins/clawhubadapter/simulator/main.js b/internal/plugins/clawhubadapter/simulator/main.js index 28bb948..c1b2f98 100644 --- a/internal/plugins/clawhubadapter/simulator/main.js +++ b/internal/plugins/clawhubadapter/simulator/main.js @@ -87,6 +87,9 @@ const registeredTools = []; // ---- Provider 存储(供 provider/call 用) ---- const registeredProviders = {}; // type -> { name, instance } +// ---- Channel 存储(供 tools/call 中通道输出路由用) ---- +const registeredChannels = {}; // name -> { output, send, channelPlugin, type } + function registerTool(defOrFactory, opts) { if (typeof defOrFactory === 'function') { const toolCtx = { @@ -192,7 +195,22 @@ const api = { }, // ---- Channel 注册 ---- - registerChannel: (ch) => notify('register', { type: 'channel', data: { name: ch.name, type: ch.type } }), + registerChannel: (ch) => { + let chName = ch.name; + let chType = ch.type || 'text'; + const chPlugin = ch.plugin; + if (chPlugin && typeof chPlugin === 'object') { + registeredChannels[chName] = { channelPlugin: chPlugin, type: chType }; + } else { + registeredChannels[chName] = { output: ch.output, send: ch.send, type: chType }; + } + notify('register', { type: 'channel', data: { name: chName, type: chType } }); + }, + + // ---- 输入提交(submitInput -> channel_input notification) ---- + submitInput: (msg) => { + notify('channel_input', { channel: api.name, payload: msg }); + }, // ---- Hook / 生命周期 ---- registerHook: (hook) => notify('register', { type: 'hook', data: { name: hook.name, event: hook.event } }), @@ -314,6 +332,39 @@ rl.on('line', async (line) => { const toolName = params.name; const args = params.arguments || {}; + // Channel output routing: if toolName matches a registered channel, use channel's output handler + const ch = registeredChannels[toolName]; + if (ch) { + try { + const channelPlugin = ch.channelPlugin; + if (channelPlugin && channelPlugin.outbound) { + const meta = args.meta || ''; + let metaObj = {}; + try { metaObj = typeof meta === 'string' ? JSON.parse(meta) : meta; } catch {} + const to = metaObj.user_id || metaObj.to || metaObj.group_id || ''; + const ctx = { to, text: args.payload || '', mediaUrl: metaObj.mediaUrl || '', cfg: {}, accountId: metaObj.accountId || null }; + let result; + if (ctx.mediaUrl && channelPlugin.outbound.sendMedia) { + result = await channelPlugin.outbound.sendMedia(ctx); + } else if (channelPlugin.outbound.sendText) { + result = await channelPlugin.outbound.sendText(ctx); + } else { + throw new Error(`channel ${toolName} has no sendText/sendMedia handler`); + } + writeJSON({ jsonrpc: '2.0', id, result: { status: 'sent', result } }); + } else if (typeof ch.output === 'function') { + const result = await ch.output(args.payload, args.type, args.meta); + writeJSON({ jsonrpc: '2.0', id, result: { status: 'sent', result } }); + } else if (typeof ch.send === 'function') { + const result = await ch.send(args.payload, args.meta); + writeJSON({ jsonrpc: '2.0', id, result: { status: 'sent', result } }); + } else { + sendError(id, -32601, `channel ${toolName} has no output handler`); + } + } catch (e) { sendError(id, -32603, e.message); } + return; + } + const tool = registeredTools.find(t => t.name === toolName); if (!tool) { sendError(id, -32601, `Tool not found: ${toolName}`); diff --git a/internal/plugins/clawhubadapter/simulator/openclaw_cli.js b/internal/plugins/clawhubadapter/simulator/openclaw_cli.js new file mode 100644 index 0000000..0383d83 --- /dev/null +++ b/internal/plugins/clawhubadapter/simulator/openclaw_cli.js @@ -0,0 +1,216 @@ +#!/usr/bin/env node +const path = require('path'); +const fs = require('fs'); +const { fork, execSync } = require('child_process'); + +const binDir = __dirname; +const simDir = path.resolve(binDir, '..'); + +function readJSON(file) { + try { return JSON.parse(fs.readFileSync(file, 'utf8')); } catch { return null; } +} + +// ---- Resolve skills directory ---- +function resolveSkillsDir() { + const skillsdirPath = path.join(simDir, '.skillsdir'); + if (fs.existsSync(skillsdirPath)) { + return fs.readFileSync(skillsdirPath, 'utf8').trim(); + } + const envDir = process.env.OPENCLAW_SKILLS_DIR; + if (envDir) return path.resolve(envDir); + const defaultDir = path.resolve(simDir, '..', 'skills'); + if (fs.existsSync(defaultDir)) return defaultDir; + return null; +} + +const skillsDir = resolveSkillsDir(); + +// ---- Detect system openclaw conflict ---- +if (process.env.OPENCLAW_CLI !== '1') { + const PATH = process.env.PATH || ''; + const pathEntries = PATH.split(path.delimiter); + for (const entry of pathEntries) { + if (entry === binDir) continue; + const otherPath = path.join(entry, 'openclaw'); + try { + if (fs.statSync(otherPath).isFile() && fs.realpathSync(otherPath) !== __filename) { + process.stderr.write(`[openclaw] warning: found another openclaw at ${otherPath}\n`); + process.stderr.write(`[openclaw] using: ${__filename}\n`); + } + } catch {} + } +} + +// ---- Settings helpers ---- +function getSettings() { + return readJSON(path.join(simDir, '.settings.json')) || {}; +} + +function writePendingSetting(key, value) { + const pendingPath = path.join(simDir, '.settings-pending.json'); + let pending = readJSON(pendingPath) || {}; + pending[key] = value; + fs.writeFileSync(pendingPath, JSON.stringify(pending, null, 2)); +} + +// ---- Permission check for plugin commands ---- +function requireSkillsDir() { + if (!skillsDir) { + process.stderr.write('Error: cannot determine skills directory. Start the manager first or set OPENCLAW_SKILLS_DIR.\n'); + process.exit(1); + } + if (!fs.existsSync(skillsDir)) { + fs.mkdirSync(skillsDir, { recursive: true }); + } +} + +// ---- Main ---- +const cmd = process.argv[2]; + +if (!cmd || cmd === 'help' || cmd === '--help') { + console.log(`HomeAgent OpenClaw Adapter 1.0.0 (openclaw-compatible) +Usage: openclaw [args] + +Commands: + plugin:install Install a plugin (npm:xxx, clawhub:xxx, or path) + plugin:uninstall Uninstall a plugin + plugin:list List installed plugins + auth:login [url] Show QR code for login + auth:status Check auth status + config:get Get config value + config:set Set config value + config:list List all config + doctor Run diagnostics + update Check for updates + env Show runtime info + --version Show version + help Show this help`); + process.exit(0); +} + +if (cmd === '--version' || cmd === 'version') { + console.log('HomeAgent OpenClaw Adapter 1.0.0 (openclaw-compatible)'); + process.exit(0); +} + +if (cmd === 'env') { + requireSkillsDir(); + const info = { + homeAgent: true, + openclawVersion: 'compatible', + platform: process.platform, + nodeVersion: process.version, + skillsDir: skillsDir, + simulatorDir: simDir, + openclawCli: __filename, + }; + console.log(JSON.stringify(info, null, 2)); + process.exit(0); +} + +if (cmd === 'doctor') { + console.log('OpenClaw Adapter diagnostics:'); + console.log(` Simulator dir: ${simDir}`); + console.log(` Skills dir: ${skillsDir || '(not set)'}`); + console.log(` Node: ${process.version}`); + console.log(` Platform: ${process.platform}`); + if (skillsDir) { + const entries = fs.readdirSync(skillsDir).filter(e => !e.startsWith('.')); + console.log(` Installed plugins: ${entries.length}`); + } + console.log(' Status: OK'); + process.exit(0); +} + +if (cmd === 'update') { + console.log('Update check: HomeAgent OpenClaw Adapter is bundled with HomeAgent.'); + console.log('To update, update HomeAgent itself.'); + process.exit(0); +} + +if (cmd === 'config:get') { + const key = process.argv[3]; + if (!key) { process.stderr.write('Usage: openclaw config:get \n'); process.exit(1); } + const settings = getSettings(); + if (key in settings) { + console.log(settings[key]); + } else { + process.stderr.write(`(not set: ${key})\n`); + process.exit(1); + } + process.exit(0); +} + +if (cmd === 'config:set') { + const key = process.argv[3]; + const value = process.argv[4]; + if (!key || value === undefined) { process.stderr.write('Usage: openclaw config:set \n'); process.exit(1); } + writePendingSetting(key, value); + console.log(`Set: ${key}=${value} (will apply on next sync)`); + process.exit(0); +} + +if (cmd === 'config:list') { + const settings = getSettings(); + const keys = Object.keys(settings); + if (keys.length === 0) { + console.log('(no settings)'); + } else { + for (const k of keys) { + console.log(` ${k}=${settings[k]}`); + } + } + process.exit(0); +} + +if (cmd === 'auth:login' || cmd === 'auth:qrcode') { + const url = process.argv[3] || 'openclaw://auth'; + try { + const qrcode = require('qrcode'); + qrcode.generate(url, { small: true }, (qr) => process.stdout.write(qr + '\n')); + } catch { + console.log('QR not available (install qrcode package). URL:'); + } + console.log(` ${url}`); + process.exit(0); +} + +if (cmd === 'auth:status') { + console.log('Auth status: running in HomeAgent mode (not applicable)'); + process.exit(0); +} + +// ---- Plugin lifecycle commands - needs manager.js fork ---- +if (cmd === 'plugin:install' || cmd === 'plugin:uninstall' || cmd === 'plugin:list') { + requireSkillsDir(); + const managerPath = path.join(simDir, 'manager.js'); + if (!fs.existsSync(managerPath)) { + process.stderr.write('Error: manager.js not found at ' + managerPath + '\n'); + process.exit(1); + } + const proc = fork(managerPath, [skillsDir, ...process.argv.slice(2)], { + stdio: 'inherit', + env: { ...process.env, OPENCLAW_CLI: '1' }, + }); + proc.on('exit', (code) => { + // After plugin:install, notify Go side via .reload-request + if (cmd === 'plugin:install' && code === 0) { + const spec = process.argv[3] || ''; + const installLog = path.join(simDir, '.install-result'); + if (fs.existsSync(installLog)) { + const result = readJSON(installLog); + if (result && result.name) { + fs.writeFileSync(path.join(simDir, '.reload-request'), result.name); + console.log(`[openclaw] queued reload for plugin: ${result.name}`); + } + fs.rmSync(installLog, { force: true }); + } + } + process.exit(code); + }); + return; +} + +process.stderr.write(`Unknown command: ${cmd}\n`); +process.stderr.write("Run 'openclaw help' for usage.\n"); +process.exit(1);