mirror of
https://gitcode.com/JianFeeeee/HomeAgent.git
synced 2026-09-22 01:48:11 +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;
|
||||
}
|
||||
|
||||
// 停止通道账号(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--) {
|
||||
|
||||
@ -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.
|
||||
|
||||
Reference in New Issue
Block a user