mirror of
https://gitcode.com/JianFeeeee/HomeAgent.git
synced 2026-10-03 15:53:56 +00:00
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:
@ -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)
|
||||
|
||||
Reference in New Issue
Block a user