From a965a74d2792a5a708195d6ea12a25f145852d8d Mon Sep 17 00:00:00 2001 From: root Date: Sun, 2 Aug 2026 13:32:51 +0800 Subject: [PATCH] =?UTF-8?q?clawhubadapter:=20agent=20=E5=85=A8=E9=87=8F?= =?UTF-8?q?=E7=AE=A1=E7=90=86=20OpenClaw=20=E6=8F=92=E4=BB=B6=E4=B8=8E?= =?UTF-8?q?=E9=80=9A=E9=81=93=EF=BC=88=E7=AE=A1=E7=90=86=E6=8E=A5=E5=8F=A3?= =?UTF-8?q?=E8=A1=A5=E9=BD=90=EF=BC=89?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 新工具(agent 可直接调用): - plugin_info:插件详情(OC/sidecar/SKILL 类型、工具列表、关联通道及状态) - plugin_reload:重新加载插件(reloadPlugin) - channel_list:全部注册通道 + 实时运行状态 - channel_send:向通道注入消息(SendToChannel → channel/send → pollQueues) - channel_start / channel_stop:通道账号启停 manager 新 RPC: - plugins/channels:registeredChannels 摘要(type/status/accounts) - channel/start:fire-and-forget startChannels 恢复账号(startAccount 永久挂起语义) - channel/stop:stopAccount 停止心跳/轮询,支持按 channel 全停或按 accountId 单停 Go 侧 channelSummary 合并 manager 注册信息与 channel_status 实时缓存(缓存优先)。 验证:微信端到端(本地 mock LLM 源)channel_list 返回 wechat running=true connected=true、 channel_send 投递成功;standalone manager 验证 channel/start 心跳恢复与 channel/stop。 --- .../plugins/clawhubadapter/manager/main.js | 55 +++ internal/plugins/clawhubadapter/plugin.go | 312 ++++++++++++++++++ 2 files changed, 367 insertions(+) diff --git a/internal/plugins/clawhubadapter/manager/main.js b/internal/plugins/clawhubadapter/manager/main.js index ea1c148..719605b 100644 --- a/internal/plugins/clawhubadapter/manager/main.js +++ b/internal/plugins/clawhubadapter/manager/main.js @@ -758,6 +758,42 @@ rl.on('line', async (line) => { return; } + // 停止通道账号(stopAccount:停止心跳/轮询)。channel 可省略(停该插件全部通道) + if (method === 'channel/stop') { + const chName = req.params?.channel; + const accountId = req.params?.accountId; + if (chName) { + const ch = registeredChannels[chName]; + if (!ch) { sendError(id, -32601, `channel not found: ${chName}`); return; } + if (accountId) { + const gateway = ch.channelPlugin && ch.channelPlugin.gateway; + if (gateway && typeof gateway.stopAccount === 'function') { + gateway.stopAccount({ account: { accountId }, channelRuntime: ch.runtime, cfg: ch.pluginConfig || {} }).catch((e) => + process.stderr.write(`[manager] ${chName}: stopAccount(${accountId}) failed: ${e.message}\n`)); + } + if (ch.accounts) delete ch.accounts[accountId]; + } else { + stopChannels(chName); + } + writeJSON({ jsonrpc: '2.0', id, result: { status: 'stopped', channel: chName, accountId: accountId || 'all' } }); + } else { + for (const name of Object.keys(registeredChannels)) stopChannels(name); + writeJSON({ jsonrpc: '2.0', id, result: { status: 'stopped', channel: 'all' } }); + } + return; + } + + // 启动通道账号(gateway.startAccount fire-and-forget) + if (method === 'channel/start') { + const chName = req.params?.channel; + const ch = registeredChannels[chName]; + if (!ch) { sendError(id, -32601, `channel not found: ${chName}`); return; } + ch.accounts = {}; + startChannels(chName).catch((e) => process.stderr.write(`[manager] ${chName}: startChannels failed: ${e.message}\n`)); + writeJSON({ jsonrpc: '2.0', id, result: { status: 'started', channel: chName } }); + return; + } + if (method === 'plugins/list') { const list = Object.entries(loadedPlugins).map(([name, p]) => ({ name, @@ -767,6 +803,18 @@ rl.on('line', async (line) => { return; } + if (method === 'plugins/channels') { + const list = Object.entries(registeredChannels).map(([name, ch]) => ({ + name, + plugin: ch.pluginName, + type: ch.type || 'text', + status: ch.status || {}, + accounts: ch.accounts ? Object.keys(ch.accounts) : [], + })); + writeJSON({ jsonrpc: '2.0', id, result: { channels: list } }); + return; + } + if (method === 'plugins/install') { const pkg = req.params?.package; if (!pkg) { sendError(id, -32602, 'package required'); return; } @@ -794,6 +842,13 @@ rl.on('line', async (line) => { if (!loadedPlugins[name]) { sendError(id, -32601, `plugin not found: ${name}`); return; } + // 先优雅停靠通道账号(stopAccount 停止心跳/轮询),再卸载 + for (const chName of Object.keys(registeredChannels)) { + if (registeredChannels[chName].pluginName === name) { + stopChannels(chName); + } + } + // Remove tools const idxs = []; for (let i = allTools.length - 1; i >= 0; i--) { diff --git a/internal/plugins/clawhubadapter/plugin.go b/internal/plugins/clawhubadapter/plugin.go index 3903768..8b97431 100644 --- a/internal/plugins/clawhubadapter/plugin.go +++ b/internal/plugins/clawhubadapter/plugin.go @@ -217,6 +217,79 @@ func (p *Plugin) Start(s *sdk.PluginSDK) error { }, }, p.handlePluginList) + s.RegisterTool(tp+"plugin_info", sdk.ToolDef{ + Name: tp + "plugin_info", + Description: "查看单个已安装插件的详细信息:类型(OC/sidecar/SKILL)、工具列表、关联通道及运行状态。", + Parameters: map[string]interface{}{ + "type": "object", + "properties": map[string]interface{}{ + "name": map[string]interface{}{"type": "string", "description": "插件名称(目录名)"}, + }, + "required": []string{"name"}, + }, + }, p.handlePluginInfo) + + s.RegisterTool(tp+"plugin_reload", sdk.ToolDef{ + Name: tp + "plugin_reload", + Description: "重新加载已安装的插件(代码或配置变更后生效,如 ClawHub 更新)。", + Parameters: map[string]interface{}{ + "type": "object", + "properties": map[string]interface{}{ + "name": map[string]interface{}{"type": "string", "description": "插件名称(目录名)"}, + }, + "required": []string{"name"}, + }, + }, p.handlePluginReload) + + s.RegisterTool(tp+"channel_list", sdk.ToolDef{ + Name: tp + "channel_list", + Description: "列出所有已注册的消息通道及其运行状态(running/connected/账号列表)。", + Parameters: map[string]interface{}{ + "type": "object", + "properties": map[string]interface{}{}, + }, + }, p.handleChannelList) + + s.RegisterTool(tp+"channel_send", sdk.ToolDef{ + Name: tp + "channel_send", + Description: "向指定通道注入一条消息(经通道插件的 chatPolls/getPolls 轮询取走,如微信/钉钉通用通道)。用于主动向通道投递内容。", + Parameters: map[string]interface{}{ + "type": "object", + "properties": map[string]interface{}{ + "channel": map[string]interface{}{"type": "string", "description": "通道名称(channel_list 查询)"}, + "content": map[string]interface{}{"type": "string", "description": "消息内容"}, + "from": map[string]interface{}{"type": "string", "description": "发送方标识(可选)"}, + "accountId": map[string]interface{}{"type": "string", "description": "账号 ID(可选,默认 default)"}, + }, + "required": []string{"channel", "content"}, + }, + }, p.handleChannelSend) + + s.RegisterTool(tp+"channel_start", sdk.ToolDef{ + Name: tp + "channel_start", + Description: "启动指定通道的账号(重新执行 startAccount,恢复心跳/收消息轮询)。", + Parameters: map[string]interface{}{ + "type": "object", + "properties": map[string]interface{}{ + "channel": map[string]interface{}{"type": "string", "description": "通道名称(channel_list 查询)"}, + }, + "required": []string{"channel"}, + }, + }, p.handleChannelStart) + + s.RegisterTool(tp+"channel_stop", sdk.ToolDef{ + Name: tp + "channel_stop", + Description: "停止指定通道(执行 stopAccount,停止心跳/收消息轮询)。channel 省略则停止全部通道。", + Parameters: map[string]interface{}{ + "type": "object", + "properties": map[string]interface{}{ + "channel": map[string]interface{}{"type": "string", "description": "通道名称(可选,省略则停止全部)"}, + "accountId": map[string]interface{}{"type": "string", "description": "账号 ID(可选,默认停止该通道全部账号)"}, + }, + "required": []string{}, + }, + }, p.handleChannelStop) + // Start file-based IPC for CLI integration (settings sync + reload requests) go p.ipcGoroutine(s) @@ -680,6 +753,245 @@ func (p *Plugin) handlePluginList(args map[string]interface{}) (interface{}, err }, nil } +func (p *Plugin) handlePluginInfo(args map[string]interface{}) (interface{}, error) { + name, _ := args["name"].(string) + if name == "" { + return errorResult("name is required"), nil + } + p.mu.Lock() + mgr := p.manager + p.mu.Unlock() + + typ := "" + var tools []string + + if mgr != nil { + if data, err := mgr.call("plugins/list", nil); err == nil && data != nil { + var result struct { + Plugins []struct { + Name string `json:"name"` + Tools []struct { + Name string `json:"name"` + Description string `json:"description"` + } `json:"tools"` + } `json:"plugins"` + } + if json.Unmarshal(data, &result) == nil { + for _, pl := range result.Plugins { + if pl.Name == name { + typ = "OC (manager)" + for _, t := range pl.Tools { + tools = append(tools, t.Name) + } + } + } + } + } + } + + if typ == "" { + p.mu.Lock() + for _, sp := range p.sidecars { + if sp == mgr || sp.name != name { + continue + } + typ = "sidecar" + if tl, err := sp.ListTools(); err == nil { + for _, t := range tl { + tools = append(tools, t.Name) + } + } + } + for _, sk := range p.skills { + if sk.Name() == name { + typ = "SKILL" + } + } + p.mu.Unlock() + } + + if typ == "" { + return errorResult(fmt.Sprintf("未找到插件 '%s'", name)), nil + } + + var parts []string + parts = append(parts, fmt.Sprintf("插件: %s\n类型: %s", name, typ)) + if len(tools) > 0 { + parts = append(parts, fmt.Sprintf("工具 (%d):\n - %s", len(tools), strings.Join(tools, "\n - "))) + } + + var own []string + for _, c := range p.channelSummary(mgr) { + if c["plugin"] == name { + own = append(own, fmt.Sprintf(" %s | type=%v running=%v connected=%v accounts=%v", + c["name"], c["type"], c["running"], c["connected"], c["accounts"])) + } + } + if len(own) > 0 { + parts = append(parts, "通道:\n"+strings.Join(own, "\n")) + } else { + parts = append(parts, "通道: 无") + } + + return map[string]interface{}{ + "content": strings.Join(parts, "\n"), + }, nil +} + +func (p *Plugin) handlePluginReload(args map[string]interface{}) (interface{}, error) { + name, _ := args["name"].(string) + if name == "" { + return errorResult("name is required"), nil + } + if err := p.reloadPlugin(name); err != nil { + return errorResult(fmt.Sprintf("reload failed: %v", err)), nil + } + return map[string]interface{}{ + "content": fmt.Sprintf("插件 %s 已重新加载", name), + }, nil +} + +// channelSummary 汇总 manager 通道注册信息与 Go 侧 channel_status 实时缓存 +func (p *Plugin) channelSummary(mgr *sidecarProcess) []map[string]interface{} { + var out []map[string]interface{} + if mgr == nil { + return out + } + data, err := mgr.call("plugins/channels", nil) + if err != nil || data == nil { + return out + } + var result struct { + Channels []struct { + Name string `json:"name"` + Plugin string `json:"plugin"` + Type string `json:"type"` + Status map[string]interface{} `json:"status"` + Accounts []string `json:"accounts"` + } `json:"channels"` + } + if json.Unmarshal(data, &result) != nil { + return out + } + for _, c := range result.Channels { + entry := map[string]interface{}{ + "name": c.Name, "plugin": c.Plugin, "type": c.Type, "accounts": c.Accounts, + } + channelStatusMu.Lock() + if st, ok := channelStatus[c.Name]; ok { + entry["running"] = st["running"] + entry["connected"] = st["connected"] + } else { + entry["running"] = c.Status["running"] + entry["connected"] = c.Status["connected"] + } + channelStatusMu.Unlock() + out = append(out, entry) + } + return out +} + +func (p *Plugin) handleChannelList(args map[string]interface{}) (interface{}, error) { + p.mu.Lock() + mgr := p.manager + p.mu.Unlock() + if mgr == nil { + return errorResult("plugin manager not available"), nil + } + chans := p.channelSummary(mgr) + if len(chans) == 0 { + return map[string]interface{}{ + "content": "没有已注册的通道。", + }, nil + } + var lines []string + for _, c := range chans { + lines = append(lines, fmt.Sprintf("- %s | plugin=%v type=%v running=%v connected=%v accounts=%v", + c["name"], c["plugin"], c["type"], c["running"], c["connected"], c["accounts"])) + } + return map[string]interface{}{ + "content": fmt.Sprintf("已注册通道 (%d):\n%s", len(chans), strings.Join(lines, "\n")), + }, nil +} + +func (p *Plugin) handleChannelSend(args map[string]interface{}) (interface{}, error) { + channel, _ := args["channel"].(string) + content, _ := args["content"].(string) + if channel == "" || content == "" { + return errorResult("channel and content are required"), nil + } + payload := map[string]interface{}{"content": content} + if from, _ := args["from"].(string); from != "" { + payload["from"] = from + } + if acc, _ := args["accountId"].(string); acc != "" { + payload["accountId"] = acc + } + if err := p.SendToChannel(channel, payload); err != nil { + return errorResult(fmt.Sprintf("send failed: %v", err)), nil + } + return map[string]interface{}{ + "content": fmt.Sprintf("消息已投递到通道 %s(等待插件轮询取走)", channel), + }, nil +} + +func (p *Plugin) handleChannelStart(args map[string]interface{}) (interface{}, error) { + channel, _ := args["channel"].(string) + if channel == "" { + return errorResult("channel is required"), nil + } + p.mu.Lock() + mgr := p.manager + p.mu.Unlock() + if mgr == nil { + return errorResult("plugin manager not available"), nil + } + data, err := mgr.call("channel/start", map[string]interface{}{"channel": channel}) + if err != nil { + return errorResult(fmt.Sprintf("start failed: %v", err)), nil + } + var result struct { + Status string `json:"status"` + } + json.Unmarshal(data, &result) + return map[string]interface{}{ + "content": fmt.Sprintf("通道 %s 启动请求已发出(%v)", channel, result.Status), + }, nil +} + +func (p *Plugin) handleChannelStop(args map[string]interface{}) (interface{}, error) { + p.mu.Lock() + mgr := p.manager + p.mu.Unlock() + if mgr == nil { + return errorResult("plugin manager not available"), nil + } + params := map[string]interface{}{} + channel, _ := args["channel"].(string) + accountId, _ := args["accountId"].(string) + if channel != "" { + params["channel"] = channel + } + if accountId != "" { + params["accountId"] = accountId + } + data, err := mgr.call("channel/stop", params) + if err != nil { + return errorResult(fmt.Sprintf("stop failed: %v", err)), nil + } + var result struct { + Status string `json:"status"` + } + json.Unmarshal(data, &result) + target := channel + if target == "" { + target = "全部" + } + return map[string]interface{}{ + "content": fmt.Sprintf("通道 %s 停止请求已发出(%v)", target, result.Status), + }, nil +} + func (p *Plugin) loadOCPlugin(s *sdk.PluginSDK, dir, name string) error { // 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.