mirror of
https://gitcode.com/JianFeeeee/HomeAgent.git
synced 2026-10-03 15:53:56 +00:00
feat: complete OpenClaw plugin API bridge with async notification architecture
- simulator/main.js: All 30+ register* methods now send JSON-RPC notifications to Go instead of being silent no-ops. Each plugin capability registration (tool, provider, channel, hook, http_route, command, service, etc.) is forwarded to Go via the notify channel. - sidecar.go: Refactored to async reader goroutine. Single goroutine reads all stdout lines, routes responses to pending callers via channel by ID, and dispatches notifications (no-ID messages) to notifyCh for handling. - plugin.go: drainNotify() collects all capabilities registered during plugin init. handleNotify() logs each registered capability for visibility. - All 9 openclaw tests pass including simulator + full pipeline tests.
This commit is contained in:
@ -4,6 +4,7 @@ import (
|
||||
"bufio"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"io"
|
||||
"log"
|
||||
"os"
|
||||
"os/exec"
|
||||
@ -20,8 +21,8 @@ type sidecarRequest struct {
|
||||
}
|
||||
|
||||
type sidecarResponse struct {
|
||||
JSONRPC string `json:"jsonrpc"`
|
||||
ID int `json:"id"`
|
||||
JSONRPC string `json:"jsonrpc"`
|
||||
ID int `json:"id"`
|
||||
Result *json.RawMessage `json:"result,omitempty"`
|
||||
Error *struct {
|
||||
Code int `json:"code"`
|
||||
@ -42,16 +43,111 @@ type OCCallResult struct {
|
||||
} `json:"content"`
|
||||
}
|
||||
|
||||
// OCNotification 是模拟器主动推送的通知
|
||||
type OCNotification struct {
|
||||
JSONRPC string `json:"jsonrpc"`
|
||||
Method string `json:"method"`
|
||||
Params json.RawMessage `json:"params"`
|
||||
}
|
||||
|
||||
// OCNamedParam 通知参数中至少包含 name 的结构
|
||||
type OCNamedParam struct {
|
||||
Type string `json:"type"`
|
||||
Data *struct {
|
||||
Name string `json:"name,omitempty"`
|
||||
} `json:"data"`
|
||||
}
|
||||
|
||||
type sidecarProcess struct {
|
||||
name string
|
||||
dir string
|
||||
cmd *exec.Cmd
|
||||
stdin *bufio.Writer
|
||||
stdout *bufio.Scanner
|
||||
mu sync.Mutex
|
||||
nextID int
|
||||
closed bool
|
||||
name string
|
||||
dir string
|
||||
cmd *exec.Cmd
|
||||
stdin *bufio.Writer
|
||||
mu sync.Mutex
|
||||
nextID int
|
||||
closed bool
|
||||
stopped bool
|
||||
|
||||
// 异步 reader
|
||||
pending map[int]chan<- []byte
|
||||
notifyCh chan OCNotification
|
||||
readerStop chan struct{}
|
||||
readerWg sync.WaitGroup
|
||||
readerReady chan struct{}
|
||||
}
|
||||
|
||||
func newSidecarProcess(name, dir string, cmd *exec.Cmd, stdin *bufio.Writer, stdout io.Reader) *sidecarProcess {
|
||||
sp := &sidecarProcess{
|
||||
name: name,
|
||||
dir: dir,
|
||||
cmd: cmd,
|
||||
stdin: stdin,
|
||||
pending: make(map[int]chan<- []byte),
|
||||
notifyCh: make(chan OCNotification, 1024),
|
||||
readerStop: make(chan struct{}),
|
||||
readerReady: make(chan struct{}),
|
||||
}
|
||||
sp.readerWg.Add(1)
|
||||
go sp.readLoop(stdout)
|
||||
<-sp.readerReady
|
||||
return sp
|
||||
}
|
||||
|
||||
func (s *sidecarProcess) readLoop(r io.Reader) {
|
||||
defer s.readerWg.Done()
|
||||
scanner := bufio.NewScanner(bufio.NewReader(r))
|
||||
// 加大 scanner buffer 防止长行截断
|
||||
scanner.Buffer(make([]byte, 0, 1024*64), 1024*64)
|
||||
close(s.readerReady)
|
||||
|
||||
for {
|
||||
select {
|
||||
case <-s.readerStop:
|
||||
return
|
||||
default:
|
||||
}
|
||||
|
||||
if !scanner.Scan() {
|
||||
if scanner.Err() != nil {
|
||||
log.Printf("[openclaw] sidecar %s read error: %v", s.name, scanner.Err())
|
||||
}
|
||||
return
|
||||
}
|
||||
line := scanner.Text()
|
||||
|
||||
var base struct {
|
||||
ID *int `json:"id"`
|
||||
Method string `json:"method,omitempty"`
|
||||
Error *json.RawMessage `json:"error,omitempty"`
|
||||
}
|
||||
if err := json.Unmarshal([]byte(line), &base); err != nil {
|
||||
continue
|
||||
}
|
||||
|
||||
if base.ID != nil {
|
||||
s.mu.Lock()
|
||||
ch, ok := s.pending[*base.ID]
|
||||
delete(s.pending, *base.ID)
|
||||
s.mu.Unlock()
|
||||
if ok {
|
||||
ch <- []byte(line)
|
||||
close(ch)
|
||||
}
|
||||
} else if base.Method != "" {
|
||||
var notif OCNotification
|
||||
if err := json.Unmarshal([]byte(line), ¬if); err == nil {
|
||||
select {
|
||||
case s.notifyCh <- notif:
|
||||
default:
|
||||
log.Printf("[openclaw] sidecar %s notify channel full, dropping: %s", s.name, notif.Method)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (s *sidecarProcess) NotifyChan() <-chan OCNotification {
|
||||
return s.notifyCh
|
||||
}
|
||||
|
||||
func launchSidecar(dir, name string) (*sidecarProcess, error) {
|
||||
@ -70,9 +166,6 @@ func launchProcess(bin, arg, dir, name string) (*sidecarProcess, error) {
|
||||
}
|
||||
}
|
||||
|
||||
// 将 dir(插件目录)作为最后一个参数传给 Node.js 进程
|
||||
// 这样: node <script> <plugin-dir>
|
||||
// echoplugin 的 main.js 忽略它, 模拟器用它加载真实插件
|
||||
cmd := exec.Command(nodePath, arg, dir)
|
||||
cmd.Dir = dir
|
||||
cmd.Stderr = os.Stderr
|
||||
@ -90,13 +183,7 @@ func launchProcess(bin, arg, dir, name string) (*sidecarProcess, error) {
|
||||
return nil, fmt.Errorf("start %s: %w", name, err)
|
||||
}
|
||||
|
||||
sp := &sidecarProcess{
|
||||
name: name,
|
||||
dir: dir,
|
||||
cmd: cmd,
|
||||
stdin: bufio.NewWriter(stdin),
|
||||
stdout: bufio.NewScanner(bufio.NewReader(stdout)),
|
||||
}
|
||||
sp := newSidecarProcess(name, dir, cmd, bufio.NewWriter(stdin), stdout)
|
||||
|
||||
if err := sp.waitReady(); err != nil {
|
||||
sp.Close()
|
||||
@ -121,55 +208,68 @@ func (s *sidecarProcess) waitReady() error {
|
||||
}
|
||||
|
||||
func (s *sidecarProcess) call(method string, params interface{}) ([]byte, error) {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
ch := make(chan []byte, 1)
|
||||
|
||||
s.mu.Lock()
|
||||
if s.closed || s.stopped {
|
||||
s.mu.Unlock()
|
||||
return nil, fmt.Errorf("sidecar %s closed", s.name)
|
||||
}
|
||||
s.nextID++
|
||||
id := s.nextID
|
||||
s.pending[id] = ch
|
||||
|
||||
req := sidecarRequest{
|
||||
JSONRPC: "2.0",
|
||||
ID: id,
|
||||
Method: method,
|
||||
Params: params,
|
||||
}
|
||||
|
||||
data, err := json.Marshal(req)
|
||||
if err != nil {
|
||||
delete(s.pending, id)
|
||||
s.mu.Unlock()
|
||||
return nil, err
|
||||
}
|
||||
|
||||
if _, err := s.stdin.Write(data); err != nil {
|
||||
delete(s.pending, id)
|
||||
s.mu.Unlock()
|
||||
return nil, err
|
||||
}
|
||||
if _, err := s.stdin.Write([]byte("\n")); err != nil {
|
||||
delete(s.pending, id)
|
||||
s.mu.Unlock()
|
||||
return nil, err
|
||||
}
|
||||
if err := s.stdin.Flush(); err != nil {
|
||||
delete(s.pending, id)
|
||||
s.mu.Unlock()
|
||||
return nil, err
|
||||
}
|
||||
s.mu.Unlock()
|
||||
|
||||
if !s.stdout.Scan() {
|
||||
if s.stdout.Err() != nil {
|
||||
return nil, fmt.Errorf("sidecar %s read: %w", s.name, s.stdout.Err())
|
||||
select {
|
||||
case raw := <-ch:
|
||||
if raw == nil {
|
||||
return nil, fmt.Errorf("sidecar %s error: nil response", s.name)
|
||||
}
|
||||
return nil, fmt.Errorf("sidecar %s closed unexpectedly", s.name)
|
||||
var resp sidecarResponse
|
||||
if err := json.Unmarshal(raw, &resp); err != nil {
|
||||
return nil, fmt.Errorf("sidecar %s unmarshal: %w", s.name, err)
|
||||
}
|
||||
if resp.Error != nil {
|
||||
return nil, fmt.Errorf("sidecar %s error: %s", s.name, resp.Error.Message)
|
||||
}
|
||||
if resp.Result == nil {
|
||||
return nil, nil
|
||||
}
|
||||
return []byte(*resp.Result), nil
|
||||
case <-time.After(30 * time.Second):
|
||||
s.mu.Lock()
|
||||
delete(s.pending, id)
|
||||
s.mu.Unlock()
|
||||
return nil, fmt.Errorf("sidecar %s call %s timeout", s.name, method)
|
||||
}
|
||||
line := s.stdout.Text()
|
||||
|
||||
var resp sidecarResponse
|
||||
if err := json.Unmarshal([]byte(line), &resp); err != nil {
|
||||
return nil, fmt.Errorf("sidecar %s unmarshal: %w", s.name, err)
|
||||
}
|
||||
if resp.Error != nil {
|
||||
return nil, fmt.Errorf("sidecar %s error: %s", s.name, resp.Error.Message)
|
||||
}
|
||||
if resp.Result == nil {
|
||||
return nil, nil
|
||||
}
|
||||
return []byte(*resp.Result), nil
|
||||
}
|
||||
|
||||
func (s *sidecarProcess) ListTools() ([]OCPTool, error) {
|
||||
@ -213,16 +313,41 @@ func (s *sidecarProcess) CallTool(name string, args map[string]interface{}) (str
|
||||
return sb, nil
|
||||
}
|
||||
|
||||
// DrainNotify 排空通知通道
|
||||
func (s *sidecarProcess) DrainNotify() {
|
||||
for {
|
||||
select {
|
||||
case <-s.notifyCh:
|
||||
default:
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (s *sidecarProcess) Close() {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
if s.closed || s.stopped {
|
||||
s.mu.Unlock()
|
||||
return
|
||||
}
|
||||
s.stopped = true
|
||||
s.mu.Unlock()
|
||||
|
||||
// 先杀进程(关闭 stdout pipe),然后 reader 的 Scan() 会退出
|
||||
if s.cmd != nil && s.cmd.Process != nil {
|
||||
s.cmd.Process.Kill()
|
||||
s.cmd.Wait()
|
||||
}
|
||||
|
||||
// 等 reader 循环结束
|
||||
s.readerWg.Wait()
|
||||
|
||||
s.mu.Lock()
|
||||
for id, ch := range s.pending {
|
||||
close(ch)
|
||||
delete(s.pending, id)
|
||||
}
|
||||
s.mu.Unlock()
|
||||
|
||||
log.Printf("[openclaw] sidecar %s stopped", s.name)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user