mirror of
https://gitcode.com/JianFeeeee/HomeAgent.git
synced 2026-10-04 00:03:59 +00:00
clawhubadapter: agent 全量管理 OpenClaw 插件与通道(管理接口补齐)
新工具(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。
This commit is contained in:
@ -758,6 +758,42 @@ rl.on('line', async (line) => {
|
|||||||
return;
|
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') {
|
if (method === 'plugins/list') {
|
||||||
const list = Object.entries(loadedPlugins).map(([name, p]) => ({
|
const list = Object.entries(loadedPlugins).map(([name, p]) => ({
|
||||||
name,
|
name,
|
||||||
@ -767,6 +803,18 @@ rl.on('line', async (line) => {
|
|||||||
return;
|
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') {
|
if (method === 'plugins/install') {
|
||||||
const pkg = req.params?.package;
|
const pkg = req.params?.package;
|
||||||
if (!pkg) { sendError(id, -32602, 'package required'); return; }
|
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; }
|
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
|
// Remove tools
|
||||||
const idxs = [];
|
const idxs = [];
|
||||||
for (let i = allTools.length - 1; i >= 0; i--) {
|
for (let i = allTools.length - 1; i >= 0; i--) {
|
||||||
|
|||||||
@ -217,6 +217,79 @@ func (p *Plugin) Start(s *sdk.PluginSDK) error {
|
|||||||
},
|
},
|
||||||
}, p.handlePluginList)
|
}, 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)
|
// Start file-based IPC for CLI integration (settings sync + reload requests)
|
||||||
go p.ipcGoroutine(s)
|
go p.ipcGoroutine(s)
|
||||||
|
|
||||||
@ -680,6 +753,245 @@ func (p *Plugin) handlePluginList(args map[string]interface{}) (interface{}, err
|
|||||||
}, nil
|
}, 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 {
|
func (p *Plugin) loadOCPlugin(s *sdk.PluginSDK, dir, name string) error {
|
||||||
// Delegate to manager's plugins/load instead of launching a separate simulator.
|
// 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.
|
// The manager is the single Node.js process that handles all OC plugins.
|
||||||
|
|||||||
Reference in New Issue
Block a user