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
This commit is contained in:
2026-07-23 16:34:26 +08:00
parent 4b6caae605
commit d4cb9cb8cc
7 changed files with 574 additions and 121 deletions

View File

@ -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,

View File

@ -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));

View File

@ -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, &params); 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)

View File

@ -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
})
}

View File

@ -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

View File

@ -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}`);

View File

@ -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 <command> [args]
Commands:
plugin:install <spec> Install a plugin (npm:xxx, clawhub:xxx, or path)
plugin:uninstall <name> Uninstall a plugin
plugin:list List installed plugins
auth:login [url] Show QR code for login
auth:status Check auth status
config:get <key> Get config value
config:set <key> <val> 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 <key>\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 <key> <value>\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);