mirror of
https://gitcode.com/JianFeeeee/HomeAgent.git
synced 2026-10-03 15:53:56 +00:00
refactor: migrate built-in plugins to SDK-only interface
- Six-phase plan complete: webui/cli/healthcheck/pluginmgr/clawhubadapter now interact with the kernel exclusively via internal/sdk interfaces; all Configure() calls and package-level global injection removed - buildSDK in internal/plugin/registry.go is the single assembly point - Add internal/sdk/events.go exporting event types/constants - Fix ProviderManager cooldown sharing: LuaAdaptedProvider.Name() now returns the source name instead of lua_<adapter>, so multiple sources sharing an adapter (single script load via shared VM AdapterCache) no longer share failure-cooldown state - Verified: build/vet/tests green, deployed to homeagent.service with full plugin capability testing via local OpenAI-compatible mock
This commit is contained in:
@ -31,7 +31,12 @@ func (tc *toolCapture) RegisterAPI(name string) error {
|
||||
func setupPlugin() (*Plugin, *toolCapture, error) {
|
||||
p := New("agentcli")
|
||||
tc := newToolCapture()
|
||||
sdk := sdk.New("agentcli", sdk.SDKConfig{RegTool: tc.RegisterTool, RegStage: tc.RegisterStage, RegAPI: tc.RegisterAPI})
|
||||
sdk := sdk.New("agentcli", sdk.SDKConfig{
|
||||
RegTool: tc.RegisterTool,
|
||||
RegStage: tc.RegisterStage,
|
||||
RegAPI: tc.RegisterAPI,
|
||||
Settings: sdk.NewSettings("agentcli", nil),
|
||||
})
|
||||
if err := p.Start(sdk); err != nil {
|
||||
return nil, nil, err
|
||||
}
|
||||
|
||||
@ -31,18 +31,11 @@ var pySimulatorSrc string
|
||||
//go:embed simulator/openclaw_cli.js
|
||||
var openclawCliSrc string
|
||||
|
||||
var SkillsDir string
|
||||
var SimulatorDir string
|
||||
|
||||
func init() {
|
||||
plugin.RegisterPluginMeta("clawhubadapter", "ClawHub 适配器", "ClawHub Adapter")
|
||||
plugin.RegisterFactory("clawhubadapter", func(name string, config map[string]interface{}) (sdk.Plugin, error) {
|
||||
dir := SkillsDir
|
||||
if dir == "" {
|
||||
dataDir, ok := config["data_dir"].(string)
|
||||
if !ok {
|
||||
return nil, fmt.Errorf("clawhubadapter plugin: config missing 'data_dir' or not a string")
|
||||
}
|
||||
dir := ""
|
||||
if dataDir, ok := config["data_dir"].(string); ok && dataDir != "" {
|
||||
dir = filepath.Join(dataDir, "skills")
|
||||
}
|
||||
return New(name, dir), nil
|
||||
@ -63,8 +56,8 @@ type Plugin struct {
|
||||
}
|
||||
|
||||
func New(name, skillsDir string) *Plugin {
|
||||
sd := SimulatorDir
|
||||
if sd == "" {
|
||||
sd := ""
|
||||
if skillsDir != "" {
|
||||
sd = filepath.Join(skillsDir, ".simulator")
|
||||
}
|
||||
return &Plugin{
|
||||
@ -100,6 +93,20 @@ func (p *Plugin) Start(s *sdk.PluginSDK) error {
|
||||
p.simulatorDir = s
|
||||
}
|
||||
}
|
||||
// 默认目录:内核 data_dir 下 skills 目录(与配置 core.daemon.data_dir 对齐)
|
||||
if p.skillsDir == "" {
|
||||
if v, _ := s.Settings().GetCore("daemon.data_dir"); v != nil {
|
||||
if dir, ok := v.(string); ok && dir != "" {
|
||||
p.skillsDir = filepath.Join(dir, "skills")
|
||||
}
|
||||
}
|
||||
}
|
||||
if p.skillsDir == "" {
|
||||
p.skillsDir = filepath.Join("data", "skills")
|
||||
}
|
||||
if p.simulatorDir == "" {
|
||||
p.simulatorDir = filepath.Join(p.skillsDir, ".simulator")
|
||||
}
|
||||
|
||||
// Launch OC plugin manager first (handles OC-format plugin installation and lifecycle)
|
||||
os.MkdirAll(p.skillsDir, 0755)
|
||||
@ -283,7 +290,7 @@ func (p *Plugin) launchManager(s *sdk.PluginSDK) error {
|
||||
// Ensure skills dir exists for the manager to scan
|
||||
os.MkdirAll(p.skillsDir, 0755)
|
||||
|
||||
sp, err := launchProcess("node", managerPath, p.skillsDir, "manager")
|
||||
sp, err := launchProcess("node", managerPath, p.skillsDir, "manager", p.simulatorDir)
|
||||
if err != nil {
|
||||
return fmt.Errorf("launch manager: %w", err)
|
||||
}
|
||||
@ -775,7 +782,7 @@ func (p *Plugin) loadPySidecar(s *sdk.PluginSDK, dir, name string) error {
|
||||
}
|
||||
}
|
||||
|
||||
sp, err := launchProcess(pythonBin, simPath, dir, name)
|
||||
sp, err := launchProcess(pythonBin, simPath, dir, name, p.simulatorDir)
|
||||
if err != nil {
|
||||
return fmt.Errorf("launch pysimulator: %w", err)
|
||||
}
|
||||
@ -808,7 +815,7 @@ func (p *Plugin) loadPySidecar(s *sdk.PluginSDK, dir, name string) error {
|
||||
}
|
||||
|
||||
func (p *Plugin) loadSidecar(s *sdk.PluginSDK, dir, name string) error {
|
||||
sp, err := launchSidecar(dir, name)
|
||||
sp, err := launchSidecar(dir, name, p.simulatorDir)
|
||||
if err != nil {
|
||||
return fmt.Errorf("launch: %w", err)
|
||||
}
|
||||
|
||||
@ -143,15 +143,15 @@ func (s *sidecarProcess) NotifyChan() <-chan OCNotification {
|
||||
return s.notifyCh
|
||||
}
|
||||
|
||||
func launchSidecar(dir, name string) (*sidecarProcess, error) {
|
||||
func launchSidecar(dir, name, simDir string) (*sidecarProcess, error) {
|
||||
mainJS := filepath.Join(dir, "main.js")
|
||||
if _, err := os.Stat(mainJS); os.IsNotExist(err) {
|
||||
return nil, nil
|
||||
}
|
||||
return launchProcess("node", mainJS, dir, name)
|
||||
return launchProcess("node", mainJS, dir, name, simDir)
|
||||
}
|
||||
|
||||
func launchProcess(bin, arg, dir, name string) (*sidecarProcess, error) {
|
||||
func launchProcess(bin, arg, dir, name, simDir string) (*sidecarProcess, error) {
|
||||
nodePath := bin
|
||||
if bin == "node" {
|
||||
if p := os.Getenv("NODE_PATH"); p != "" {
|
||||
@ -164,8 +164,8 @@ func launchProcess(bin, arg, dir, name string) (*sidecarProcess, error) {
|
||||
cmd.Stderr = os.Stderr
|
||||
|
||||
// Add openclaw CLI bin dir to PATH so subprocesses can exec 'openclaw' command
|
||||
if SimulatorDir != "" {
|
||||
binDir := filepath.Join(SimulatorDir, "bin")
|
||||
if simDir != "" {
|
||||
binDir := filepath.Join(simDir, "bin")
|
||||
if info, err := os.Stat(binDir); err == nil && info.IsDir() {
|
||||
env := os.Environ()
|
||||
binDirPath := binDir + string(os.PathListSeparator)
|
||||
|
||||
@ -26,10 +26,12 @@ 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 (m *mockSettings) DefsCore(prefix string) []*sdk.ConfigDef { return nil }
|
||||
func (m *mockSettings) DefsPlugin(plugin, prefix string) []*sdk.ConfigDef { return nil }
|
||||
|
||||
func TestLaunchSidecarNoMainJS(t *testing.T) {
|
||||
tmpDir := t.TempDir()
|
||||
sp, err := launchSidecar(tmpDir, "nonexistent")
|
||||
sp, err := launchSidecar(tmpDir, "nonexistent", "")
|
||||
if err != nil {
|
||||
t.Fatalf("unexpected error: %v", err)
|
||||
}
|
||||
@ -50,7 +52,7 @@ func TestLaunchSidecarAndListTools(t *testing.T) {
|
||||
t.Fatalf("write test plugin: %v", err)
|
||||
}
|
||||
|
||||
sp, err := launchSidecar(tmpDir, "echoplugin")
|
||||
sp, err := launchSidecar(tmpDir, "echoplugin", "")
|
||||
if err != nil {
|
||||
t.Fatalf("launch sidecar: %v", err)
|
||||
}
|
||||
@ -92,7 +94,7 @@ func TestCallEchoTool(t *testing.T) {
|
||||
t.Fatalf("write test plugin: %v", err)
|
||||
}
|
||||
|
||||
sp, err := launchSidecar(tmpDir, "echoplugin")
|
||||
sp, err := launchSidecar(tmpDir, "echoplugin", "")
|
||||
if err != nil {
|
||||
t.Fatalf("launch sidecar: %v", err)
|
||||
}
|
||||
@ -124,7 +126,7 @@ func TestCallAddTool(t *testing.T) {
|
||||
t.Fatalf("write test plugin: %v", err)
|
||||
}
|
||||
|
||||
sp, err := launchSidecar(tmpDir, "echoplugin")
|
||||
sp, err := launchSidecar(tmpDir, "echoplugin", "")
|
||||
if err != nil {
|
||||
t.Fatalf("launch sidecar: %v", err)
|
||||
}
|
||||
@ -157,7 +159,7 @@ func TestCallNonexistentTool(t *testing.T) {
|
||||
t.Fatalf("write test plugin: %v", err)
|
||||
}
|
||||
|
||||
sp, err := launchSidecar(tmpDir, "echoplugin")
|
||||
sp, err := launchSidecar(tmpDir, "echoplugin", "")
|
||||
if err != nil {
|
||||
t.Fatalf("launch sidecar: %v", err)
|
||||
}
|
||||
@ -182,7 +184,7 @@ func TestConcurrentCalls(t *testing.T) {
|
||||
t.Fatalf("write test plugin: %v", err)
|
||||
}
|
||||
|
||||
sp, err := launchSidecar(tmpDir, "echoplugin")
|
||||
sp, err := launchSidecar(tmpDir, "echoplugin", "")
|
||||
if err != nil {
|
||||
t.Fatalf("launch sidecar: %v", err)
|
||||
}
|
||||
@ -217,7 +219,7 @@ func launchSimulator(t *testing.T, pluginDir, name string) *sidecarProcess {
|
||||
if err != nil {
|
||||
t.Fatalf("abs simulator path: %v", err)
|
||||
}
|
||||
sp, err := launchProcess("node", simPath, pluginDir, name)
|
||||
sp, err := launchProcess("node", simPath, pluginDir, name, "")
|
||||
if err != nil {
|
||||
t.Fatalf("launch simulator for %s: %v", name, err)
|
||||
}
|
||||
@ -338,8 +340,6 @@ func TestLoadOCPluginViaPluginStart(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
SimulatorDir = filepath.Join(t.TempDir(), ".simulator")
|
||||
|
||||
p := New("openclaw", skillsDir)
|
||||
|
||||
var registeredTools []string
|
||||
|
||||
@ -12,8 +12,6 @@ import (
|
||||
"strings"
|
||||
"sync"
|
||||
|
||||
agentCore "gitcode.com/JianFeeeee/HomeAgent/internal/agent/core"
|
||||
internalConfig "gitcode.com/JianFeeeee/HomeAgent/internal/config"
|
||||
"gitcode.com/JianFeeeee/HomeAgent/internal/plugin"
|
||||
sdk "gitcode.com/JianFeeeee/HomeAgent/internal/sdk"
|
||||
)
|
||||
@ -21,22 +19,6 @@ import (
|
||||
// DefaultSocket 由 main.go 在 Load() 前设置,覆盖默认 socket 路径。
|
||||
var DefaultSocket string
|
||||
|
||||
// 以下通过 Configure() 注入内核依赖
|
||||
var (
|
||||
pluginReg *plugin.Registry
|
||||
cfgReg *internalConfig.ConfigRegistry
|
||||
statusProv agentCore.StatusProvider
|
||||
pluginDir string
|
||||
)
|
||||
|
||||
// Configure 由 main.go 在 Load() 前调用,注入内核依赖供结构化命令使用。
|
||||
func Configure(pr *plugin.Registry, cr *internalConfig.ConfigRegistry, sp agentCore.StatusProvider, pDir string) {
|
||||
pluginReg = pr
|
||||
cfgReg = cr
|
||||
statusProv = sp
|
||||
pluginDir = pDir
|
||||
}
|
||||
|
||||
func init() {
|
||||
plugin.RegisterPluginMeta("cli", "CLI", "CLI")
|
||||
plugin.RegisterFactory("cli", func(name string, config map[string]interface{}) (sdk.Plugin, error) {
|
||||
@ -189,18 +171,16 @@ func (p *Plugin) cliAPIKey(s *sdk.PluginSDK) string {
|
||||
}
|
||||
}
|
||||
}
|
||||
return p.webuiAPIKey()
|
||||
return p.webuiAPIKey(s)
|
||||
}
|
||||
|
||||
func (p *Plugin) webuiAPIKey() string {
|
||||
if cfgReg == nil {
|
||||
func (p *Plugin) webuiAPIKey(s *sdk.PluginSDK) string {
|
||||
if s == nil {
|
||||
return ""
|
||||
}
|
||||
ps := cfgReg.PluginConfig("webui")
|
||||
if v, _ := ps.Get("api_key"); v != nil {
|
||||
if s, ok := v.(string); ok {
|
||||
return s
|
||||
}
|
||||
v, _ := s.Settings().GetPlugin("webui", "api_key")
|
||||
if k, ok := v.(string); ok {
|
||||
return k
|
||||
}
|
||||
return ""
|
||||
}
|
||||
@ -215,11 +195,11 @@ func (p *Plugin) handleBuiltin(conn net.Conn, line string, s *sdk.PluginSDK) boo
|
||||
case "/help":
|
||||
p.cmdHelp(conn)
|
||||
case "/status":
|
||||
p.cmdStatus(conn)
|
||||
p.cmdStatus(conn, s)
|
||||
case "/kernel":
|
||||
p.cmdKernel(conn)
|
||||
p.cmdKernel(conn, s)
|
||||
case "/settings":
|
||||
p.cmdSettings(conn, parts)
|
||||
p.cmdSettings(conn, parts, s)
|
||||
case "/plugin":
|
||||
p.cmdPlugin(conn, parts, s)
|
||||
case "/memory":
|
||||
@ -227,7 +207,7 @@ func (p *Plugin) handleBuiltin(conn net.Conn, line string, s *sdk.PluginSDK) boo
|
||||
case "/knowledge":
|
||||
p.cmdKnowledge(conn, s)
|
||||
case "/agents":
|
||||
p.cmdAgents(conn)
|
||||
p.cmdAgents(conn, s)
|
||||
default:
|
||||
return false
|
||||
}
|
||||
@ -260,12 +240,13 @@ func (p *Plugin) cmdHelp(conn net.Conn) {
|
||||
|
||||
// ======== /status ========
|
||||
|
||||
func (p *Plugin) cmdStatus(conn net.Conn) {
|
||||
if statusProv == nil {
|
||||
func (p *Plugin) cmdStatus(conn net.Conn, s *sdk.PluginSDK) {
|
||||
st := s.Status()
|
||||
if st == nil {
|
||||
writeLine(conn, map[string]interface{}{"type": "error", "error": "status provider not available"})
|
||||
return
|
||||
}
|
||||
ks := statusProv.GetKernelStatus()
|
||||
ks := st.GetKernelStatus()
|
||||
|
||||
llmStatus := "不可用"
|
||||
if ks.LLM.Available {
|
||||
@ -294,19 +275,21 @@ func (p *Plugin) cmdStatus(conn net.Conn) {
|
||||
|
||||
// ======== /kernel ========
|
||||
|
||||
func (p *Plugin) cmdKernel(conn net.Conn) {
|
||||
if statusProv == nil {
|
||||
func (p *Plugin) cmdKernel(conn net.Conn, s *sdk.PluginSDK) {
|
||||
st := s.Status()
|
||||
if st == nil {
|
||||
writeLine(conn, map[string]interface{}{"type": "error", "error": "status provider not available"})
|
||||
return
|
||||
}
|
||||
data, _ := json.MarshalIndent(statusProv.GetKernelStatus(), "", " ")
|
||||
data, _ := json.MarshalIndent(st.GetKernelStatus(), "", " ")
|
||||
writeLine(conn, map[string]interface{}{"type": "response", "content": string(data)})
|
||||
}
|
||||
|
||||
// ======== /settings ========
|
||||
|
||||
func (p *Plugin) cmdSettings(conn net.Conn, parts []string) {
|
||||
if cfgReg == nil {
|
||||
func (p *Plugin) cmdSettings(conn net.Conn, parts []string, s *sdk.PluginSDK) {
|
||||
sett := s.Settings()
|
||||
if sett == nil {
|
||||
writeLine(conn, map[string]interface{}{"type": "error", "error": "config registry not available"})
|
||||
return
|
||||
}
|
||||
@ -318,7 +301,7 @@ func (p *Plugin) cmdSettings(conn net.Conn, parts []string) {
|
||||
}
|
||||
key := parts[2]
|
||||
val := strings.Join(parts[3:], " ")
|
||||
if err := cfgReg.Set(key, val); err != nil {
|
||||
if err := sett.SetCore(key, val); err != nil {
|
||||
writeLine(conn, map[string]interface{}{"type": "error", "error": err.Error()})
|
||||
return
|
||||
}
|
||||
@ -330,7 +313,7 @@ func (p *Plugin) cmdSettings(conn net.Conn, parts []string) {
|
||||
if len(parts) >= 2 {
|
||||
prefix = parts[1]
|
||||
}
|
||||
keys := cfgReg.List(prefix)
|
||||
keys, _ := sett.ListCore(prefix)
|
||||
sort.Strings(keys)
|
||||
if len(keys) == 0 {
|
||||
writeLine(conn, map[string]interface{}{"type": "response", "content": "无匹配配置项"})
|
||||
@ -338,7 +321,7 @@ func (p *Plugin) cmdSettings(conn net.Conn, parts []string) {
|
||||
}
|
||||
var lines []string
|
||||
for _, k := range keys {
|
||||
v, _ := cfgReg.Get(k)
|
||||
v, _ := sett.GetCore(k)
|
||||
lines = append(lines, fmt.Sprintf(" %s = %v", k, v))
|
||||
}
|
||||
writeLine(conn, map[string]interface{}{
|
||||
@ -409,20 +392,22 @@ func (p *Plugin) cmdPlugin(conn net.Conn, parts []string, s *sdk.PluginSDK) {
|
||||
writeLine(conn, map[string]interface{}{"type": "response", "content": "用法: /plugin remove <name>"})
|
||||
return
|
||||
}
|
||||
if pmgr == nil {
|
||||
writeLine(conn, map[string]interface{}{"type": "error", "error": "plugin manager not available"})
|
||||
return
|
||||
}
|
||||
name := parts[2]
|
||||
if pluginDir == "" {
|
||||
dir := pmgr.PluginDir()
|
||||
if dir == "" {
|
||||
writeLine(conn, map[string]interface{}{"type": "error", "error": "plugin dir not configured"})
|
||||
return
|
||||
}
|
||||
dir := filepath.Join(pluginDir, name)
|
||||
if err := os.RemoveAll(dir); err != nil {
|
||||
if err := os.RemoveAll(filepath.Join(dir, name)); err != nil {
|
||||
writeLine(conn, map[string]interface{}{"type": "error", "error": err.Error()})
|
||||
return
|
||||
}
|
||||
// 同步清理禁用表
|
||||
if pmgr != nil {
|
||||
pmgr.EnablePlugin(name)
|
||||
}
|
||||
_ = pmgr.EnablePlugin(name)
|
||||
writeLine(conn, map[string]interface{}{"type": "response", "content": fmt.Sprintf("插件 %s 已删除,执行 /plugin reload 生效", name)})
|
||||
|
||||
case "info":
|
||||
@ -430,20 +415,33 @@ func (p *Plugin) cmdPlugin(conn net.Conn, parts []string, s *sdk.PluginSDK) {
|
||||
writeLine(conn, map[string]interface{}{"type": "response", "content": "用法: /plugin info <name>"})
|
||||
return
|
||||
}
|
||||
if pluginReg == nil {
|
||||
writeLine(conn, map[string]interface{}{"type": "error", "error": "plugin registry not available"})
|
||||
if pmgr == nil {
|
||||
writeLine(conn, map[string]interface{}{"type": "error", "error": "plugin manager not available"})
|
||||
return
|
||||
}
|
||||
plg := pluginReg.Get(parts[2])
|
||||
if plg == nil {
|
||||
if pmgr != nil && pmgr.IsPluginDisabled(parts[2]) {
|
||||
writeLine(conn, map[string]interface{}{"type": "response", "content": fmt.Sprintf("插件 %q 已禁用", parts[2])})
|
||||
return
|
||||
name := parts[2]
|
||||
metas := pmgr.PluginMetas()
|
||||
meta, hasMeta := metas[name]
|
||||
loaded := false
|
||||
for _, n := range pmgr.ListLoadedPlugins() {
|
||||
if n == name {
|
||||
loaded = true
|
||||
break
|
||||
}
|
||||
writeLine(conn, map[string]interface{}{"type": "response", "content": fmt.Sprintf("插件 %q 未安装", parts[2])})
|
||||
}
|
||||
if loaded {
|
||||
display := name
|
||||
if hasMeta && meta.NameZh != "" {
|
||||
display = meta.NameZh
|
||||
}
|
||||
writeLine(conn, map[string]interface{}{"type": "response", "content": fmt.Sprintf("名称: %s (%s)\n状态: 已加载", display, name)})
|
||||
return
|
||||
}
|
||||
writeLine(conn, map[string]interface{}{"type": "response", "content": fmt.Sprintf("名称: %s\n状态: 已加载", plg.Name())})
|
||||
if pmgr.IsPluginDisabled(name) {
|
||||
writeLine(conn, map[string]interface{}{"type": "response", "content": fmt.Sprintf("插件 %q 已禁用", name)})
|
||||
return
|
||||
}
|
||||
writeLine(conn, map[string]interface{}{"type": "response", "content": fmt.Sprintf("插件 %q 未安装", name)})
|
||||
|
||||
case "disable":
|
||||
if len(parts) < 3 {
|
||||
@ -541,12 +539,13 @@ func (p *Plugin) cmdKnowledge(conn net.Conn, s *sdk.PluginSDK) {
|
||||
|
||||
// ======== /agents ========
|
||||
|
||||
func (p *Plugin) cmdAgents(conn net.Conn) {
|
||||
if statusProv == nil {
|
||||
func (p *Plugin) cmdAgents(conn net.Conn, s *sdk.PluginSDK) {
|
||||
st := s.Status()
|
||||
if st == nil {
|
||||
writeLine(conn, map[string]interface{}{"type": "error", "error": "status provider not available"})
|
||||
return
|
||||
}
|
||||
ks := statusProv.GetKernelStatus()
|
||||
ks := st.GetKernelStatus()
|
||||
data, _ := json.MarshalIndent(map[string]string{"agent_id": ks.AgentID}, "", " ")
|
||||
writeLine(conn, map[string]interface{}{"type": "response", "content": string(data)})
|
||||
}
|
||||
|
||||
@ -34,7 +34,12 @@ func (tc *toolCapture) RegisterAPI(name string) error { return nil }
|
||||
func setupPlugin() (*Plugin, *toolCapture, error) {
|
||||
p := New("cmd")
|
||||
tc := newToolCapture()
|
||||
sdk := sdk.New("cmd", sdk.SDKConfig{RegTool: tc.RegisterTool, RegStage: tc.RegisterStage, RegAPI: tc.RegisterAPI})
|
||||
sdk := sdk.New("cmd", sdk.SDKConfig{
|
||||
RegTool: tc.RegisterTool,
|
||||
RegStage: tc.RegisterStage,
|
||||
RegAPI: tc.RegisterAPI,
|
||||
Settings: sdk.NewSettings("cmd", nil),
|
||||
})
|
||||
if err := p.Start(sdk); err != nil {
|
||||
return nil, nil, err
|
||||
}
|
||||
|
||||
@ -9,27 +9,10 @@ import (
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
agentAPI "gitcode.com/JianFeeeee/HomeAgent/internal/agent/api"
|
||||
agentCore "gitcode.com/JianFeeeee/HomeAgent/internal/agent/core"
|
||||
agentIO "gitcode.com/JianFeeeee/HomeAgent/internal/agent/io"
|
||||
"gitcode.com/JianFeeeee/HomeAgent/internal/knowledge"
|
||||
"gitcode.com/JianFeeeee/HomeAgent/internal/memory"
|
||||
doc "gitcode.com/JianFeeeee/HomeAgent/internal/memory/document"
|
||||
"gitcode.com/JianFeeeee/HomeAgent/internal/plugin"
|
||||
sdk "gitcode.com/JianFeeeee/HomeAgent/internal/sdk"
|
||||
)
|
||||
|
||||
var (
|
||||
hcStageHost *agentCore.StageHost
|
||||
hcIOMgr *agentIO.IOManager
|
||||
hcPluginReg *plugin.Registry
|
||||
hcMemory *memory.GraphDB
|
||||
hcKnowledge *knowledge.Store
|
||||
hcDocStore *doc.Store
|
||||
hcProviderMgr *agentAPI.ProviderManager
|
||||
hcStatusProvider agentCore.StatusProvider
|
||||
)
|
||||
|
||||
type toolInfo struct {
|
||||
Name string `json:"name"`
|
||||
Source string `json:"source"`
|
||||
@ -50,24 +33,9 @@ type llmReport struct {
|
||||
Detail string `json:"detail,omitempty"`
|
||||
}
|
||||
|
||||
func Configure(sh *agentCore.StageHost, iom *agentIO.IOManager, pr *plugin.Registry,
|
||||
mem *memory.GraphDB, ks *knowledge.Store, ds *doc.Store, pm *agentAPI.ProviderManager, sp agentCore.StatusProvider) {
|
||||
hcStageHost = sh
|
||||
hcIOMgr = iom
|
||||
hcPluginReg = pr
|
||||
hcMemory = mem
|
||||
hcKnowledge = ks
|
||||
hcDocStore = ds
|
||||
hcProviderMgr = pm
|
||||
hcStatusProvider = sp
|
||||
}
|
||||
|
||||
func init() {
|
||||
plugin.RegisterPluginMeta("healthcheck", "健康检查", "Health Check")
|
||||
plugin.RegisterFactory("healthcheck", func(name string, config map[string]interface{}) (sdk.Plugin, error) {
|
||||
if hcStageHost == nil {
|
||||
return nil, nil
|
||||
}
|
||||
return New(name), nil
|
||||
})
|
||||
}
|
||||
@ -215,7 +183,7 @@ func (p *Plugin) Start(s *sdk.PluginSDK) error {
|
||||
"properties": map[string]interface{}{},
|
||||
},
|
||||
}, func(args map[string]interface{}) (interface{}, error) {
|
||||
return p.listAllTools()
|
||||
return p.listAllTools(s)
|
||||
})
|
||||
|
||||
p.selfToolNames["healthcheck_memory"] = true
|
||||
@ -227,7 +195,7 @@ func (p *Plugin) Start(s *sdk.PluginSDK) error {
|
||||
"properties": map[string]interface{}{},
|
||||
},
|
||||
}, func(args map[string]interface{}) (interface{}, error) {
|
||||
return p.checkMemory()
|
||||
return p.checkMemory(s)
|
||||
})
|
||||
|
||||
p.selfToolNames["healthcheck_report"] = true
|
||||
@ -255,7 +223,7 @@ func (p *Plugin) Start(s *sdk.PluginSDK) error {
|
||||
return map[string]interface{}{"ok": true, "received": count}, nil
|
||||
})
|
||||
|
||||
if hcStatusProvider != nil {
|
||||
if s.Status() != nil {
|
||||
p.selfToolNames["healthcheck_kernel"] = true
|
||||
s.RegisterTool("healthcheck_kernel", sdk.ToolDef{
|
||||
Name: "healthcheck_kernel",
|
||||
@ -265,7 +233,7 @@ func (p *Plugin) Start(s *sdk.PluginSDK) error {
|
||||
"properties": map[string]interface{}{},
|
||||
},
|
||||
}, func(args map[string]interface{}) (interface{}, error) {
|
||||
return hcStatusProvider.GetKernelStatus(), nil
|
||||
return s.Status().GetKernelStatus(), nil
|
||||
})
|
||||
}
|
||||
|
||||
@ -300,9 +268,9 @@ func (p *Plugin) Start(s *sdk.PluginSDK) error {
|
||||
p.startAutoCheck(s, p.autoInterval)
|
||||
}
|
||||
|
||||
log.Printf("[healthcheck] ready (stageHost=%v iom=%v reg=%v mem=%v ks=%v ds=%v pm=%v sp=%v)",
|
||||
hcStageHost != nil, hcIOMgr != nil, hcPluginReg != nil,
|
||||
hcMemory != nil, hcKnowledge != nil, hcDocStore != nil, hcProviderMgr != nil, hcStatusProvider != nil)
|
||||
log.Printf("[healthcheck] ready (tool=%v mem=%v ks=%v ds=%v llm=%v plugins=%v status=%v)",
|
||||
s.Tool() != nil, s.Memory() != nil, s.Knowledge() != nil,
|
||||
s.DocMemory() != nil, s.LLM() != nil, s.PluginMgr() != nil, s.Status() != nil)
|
||||
return nil
|
||||
}
|
||||
|
||||
@ -367,36 +335,32 @@ func (p *Plugin) runAutoCheck(s *sdk.PluginSDK) {
|
||||
func (p *Plugin) runFullCheck(s *sdk.PluginSDK) (interface{}, error) {
|
||||
results := []checkResult{}
|
||||
|
||||
pluginResult := p.checkPluginsRaw()
|
||||
pluginResult := p.checkPluginsRaw(s)
|
||||
results = append(results, pluginResult...)
|
||||
|
||||
toolResult := p.checkToolsRaw()
|
||||
toolResult := p.checkToolsRaw(s)
|
||||
results = append(results, toolResult...)
|
||||
|
||||
if hcMemory != nil {
|
||||
r := p.testMemoryRaw()
|
||||
results = append(results, r)
|
||||
if s.Memory() != nil {
|
||||
results = append(results, p.testMemoryRaw(s))
|
||||
} else {
|
||||
results = append(results, checkResult{Name: "memory", Status: "skip", Detail: "图记忆未初始化", Pass: true})
|
||||
}
|
||||
|
||||
if hcKnowledge != nil {
|
||||
r := p.testKnowledgeRaw()
|
||||
results = append(results, r)
|
||||
if s.Knowledge() != nil {
|
||||
results = append(results, p.testKnowledgeRaw(s))
|
||||
} else {
|
||||
results = append(results, checkResult{Name: "knowledge", Status: "skip", Detail: "知识库未初始化", Pass: true})
|
||||
}
|
||||
|
||||
if hcDocStore != nil {
|
||||
r := p.testDocStoreRaw()
|
||||
results = append(results, r)
|
||||
if s.DocMemory() != nil {
|
||||
results = append(results, p.testDocStoreRaw(s))
|
||||
} else {
|
||||
results = append(results, checkResult{Name: "documents", Status: "skip", Detail: "文档记忆未初始化", Pass: true})
|
||||
}
|
||||
|
||||
if hcProviderMgr != nil {
|
||||
r := p.testLLMDriven()
|
||||
results = append(results, r)
|
||||
if s.LLM() != nil {
|
||||
results = append(results, p.testLLMDriven(s))
|
||||
} else {
|
||||
results = append(results, checkResult{Name: "llm_discovery", Status: "skip", Detail: "LLM Provider 未初始化", Pass: true})
|
||||
}
|
||||
@ -424,7 +388,7 @@ func (p *Plugin) runFullCheck(s *sdk.PluginSDK) (interface{}, error) {
|
||||
}
|
||||
|
||||
func (p *Plugin) checkPlugins(s *sdk.PluginSDK) (interface{}, error) {
|
||||
results := p.checkPluginsRaw()
|
||||
results := p.checkPluginsRaw(s)
|
||||
return map[string]interface{}{
|
||||
"status": "ok",
|
||||
"plugins": results,
|
||||
@ -432,12 +396,12 @@ func (p *Plugin) checkPlugins(s *sdk.PluginSDK) (interface{}, error) {
|
||||
}, nil
|
||||
}
|
||||
|
||||
func (p *Plugin) checkPluginsRaw() []checkResult {
|
||||
if hcPluginReg == nil {
|
||||
func (p *Plugin) checkPluginsRaw(s *sdk.PluginSDK) []checkResult {
|
||||
if s.PluginMgr() == nil {
|
||||
return []checkResult{{Name: "plugins", Status: "skip", Detail: "插件注册表未初始化", Pass: true}}
|
||||
}
|
||||
|
||||
names := hcPluginReg.List()
|
||||
names := s.PluginMgr().ListLoadedPlugins()
|
||||
if names == nil {
|
||||
names = []string{}
|
||||
}
|
||||
@ -449,8 +413,8 @@ func (p *Plugin) checkPluginsRaw() []checkResult {
|
||||
}}
|
||||
}
|
||||
|
||||
func (p *Plugin) listAllTools() (interface{}, error) {
|
||||
tools := p.collectAllTools()
|
||||
func (p *Plugin) listAllTools(s *sdk.PluginSDK) (interface{}, error) {
|
||||
tools := p.collectAllTools(s)
|
||||
return map[string]interface{}{
|
||||
"status": "ok",
|
||||
"count": len(tools),
|
||||
@ -458,8 +422,8 @@ func (p *Plugin) listAllTools() (interface{}, error) {
|
||||
}, nil
|
||||
}
|
||||
|
||||
func (p *Plugin) checkToolsRaw() []checkResult {
|
||||
tools := p.collectAllTools()
|
||||
func (p *Plugin) checkToolsRaw(s *sdk.PluginSDK) []checkResult {
|
||||
tools := p.collectAllTools(s)
|
||||
return []checkResult{{
|
||||
Name: "tools",
|
||||
Status: "ok",
|
||||
@ -468,7 +432,7 @@ func (p *Plugin) checkToolsRaw() []checkResult {
|
||||
}}
|
||||
}
|
||||
|
||||
func (p *Plugin) collectAllTools() []toolInfo {
|
||||
func (p *Plugin) collectAllTools(s *sdk.PluginSDK) []toolInfo {
|
||||
seen := map[string]bool{}
|
||||
var tools []toolInfo
|
||||
|
||||
@ -480,14 +444,12 @@ func (p *Plugin) collectAllTools() []toolInfo {
|
||||
tools = append(tools, toolInfo{Name: name, Source: source, Description: desc})
|
||||
}
|
||||
|
||||
if hcStageHost != nil {
|
||||
for _, def := range hcStageHost.GetToolDefs() {
|
||||
if s.Tool() != nil {
|
||||
for _, def := range s.Tool().GetToolDefs() {
|
||||
addTool(def.Name, "plugin", def.Description)
|
||||
}
|
||||
}
|
||||
|
||||
if hcIOMgr != nil {
|
||||
for _, def := range hcIOMgr.GetAllTools() {
|
||||
for _, def := range s.Tool().GetAllTools() {
|
||||
addTool(def.Name, "device", def.Description)
|
||||
}
|
||||
}
|
||||
@ -495,23 +457,18 @@ func (p *Plugin) collectAllTools() []toolInfo {
|
||||
return tools
|
||||
}
|
||||
|
||||
func (p *Plugin) testMemoryRaw() checkResult {
|
||||
func (p *Plugin) testMemoryRaw(s *sdk.PluginSDK) checkResult {
|
||||
marker := fmt.Sprintf("_hc_%d", time.Now().UnixNano())
|
||||
triples := []memory.Triple{
|
||||
triples := []sdk.Triple{
|
||||
{Subject: marker, Relation: "is", Object: "healthcheck_test", SubjectType: "System", ObjectType: "Flag"},
|
||||
}
|
||||
|
||||
start := time.Now()
|
||||
ec, rc, err := hcMemory.Commit(triples, "healthcheck", 0)
|
||||
if err != nil {
|
||||
if err := s.Memory().Commit(triples); err != nil {
|
||||
return checkResult{Name: "memory_write", Status: "fail", Detail: fmt.Sprintf("写入失败: %v", err), Pass: false}
|
||||
}
|
||||
|
||||
if _, _, err := hcMemory.Commit(triples, "healthcheck_cleanup", 0); err != nil {
|
||||
log.Printf("[healthcheck] memory cleanup error: %v", err)
|
||||
}
|
||||
|
||||
n, err := hcMemory.Purge(map[string]string{"subject_contains": marker}, "hard")
|
||||
n, err := s.Memory().Purge(map[string]string{"subject_contains": marker}, "hard")
|
||||
if err != nil {
|
||||
return checkResult{Name: "memory_purge", Status: "fail", Detail: fmt.Sprintf("清理失败: %v", err), Pass: false}
|
||||
}
|
||||
@ -520,25 +477,29 @@ func (p *Plugin) testMemoryRaw() checkResult {
|
||||
return checkResult{
|
||||
Name: "memory",
|
||||
Status: "ok",
|
||||
Detail: fmt.Sprintf("写入 %d 实体/%d 关系, 清理 %d 条, 耗时 %v", ec, rc, n, elapsed.Round(time.Millisecond)),
|
||||
Detail: fmt.Sprintf("写入+清理 %d 条, 耗时 %v", n, elapsed.Round(time.Millisecond)),
|
||||
Pass: true,
|
||||
}
|
||||
}
|
||||
|
||||
func (p *Plugin) testKnowledgeRaw() checkResult {
|
||||
func (p *Plugin) testKnowledgeRaw(s *sdk.PluginSDK) checkResult {
|
||||
marker := fmt.Sprintf("_hc_knowledge_test_%d", time.Now().UnixNano())
|
||||
start := time.Now()
|
||||
|
||||
if err := hcKnowledge.Add(marker, "健康检查测试标记,可忽略"); err != nil {
|
||||
if err := s.Knowledge().Add(marker, "健康检查测试标记,可忽略"); err != nil {
|
||||
return checkResult{Name: "knowledge", Status: "fail", Detail: fmt.Sprintf("写入失败: %v", err), Pass: false}
|
||||
}
|
||||
|
||||
results := hcKnowledge.Search("健康检查测试标记", 3)
|
||||
results, err := s.Knowledge().Search("健康检查测试标记", 3)
|
||||
if err != nil {
|
||||
s.Knowledge().Remove(marker)
|
||||
return checkResult{Name: "knowledge", Status: "fail", Detail: fmt.Sprintf("查询失败: %v", err), Pass: false}
|
||||
}
|
||||
|
||||
elapsed := time.Since(start)
|
||||
|
||||
// 清理测试条目,避免积累
|
||||
hcKnowledge.Remove(marker)
|
||||
s.Knowledge().Remove(marker)
|
||||
|
||||
if len(results) > 0 {
|
||||
return checkResult{
|
||||
@ -557,20 +518,21 @@ func (p *Plugin) testKnowledgeRaw() checkResult {
|
||||
}
|
||||
}
|
||||
|
||||
func (p *Plugin) testDocStoreRaw() checkResult {
|
||||
func (p *Plugin) testDocStoreRaw(s *sdk.PluginSDK) checkResult {
|
||||
start := time.Now()
|
||||
doc := &doc.Doc{
|
||||
Summary: "健康检查测试文档",
|
||||
doc := &sdk.Doc{
|
||||
Title: fmt.Sprintf("健康检查测试文档 %d", time.Now().UnixNano()),
|
||||
Content: "这是一条由 healthcheck 插件创建的测试文档,用于验证文档记忆系统是否正常工作。",
|
||||
Tags: []string{"healthcheck", "test"},
|
||||
Source: "healthcheck",
|
||||
}
|
||||
if err := hcDocStore.Insert(doc); err != nil {
|
||||
if err := s.DocMemory().Insert(doc); err != nil {
|
||||
return checkResult{Name: "documents", Status: "fail", Detail: fmt.Sprintf("写入失败: %v", err), Pass: false}
|
||||
}
|
||||
|
||||
if doc.ID != "" {
|
||||
hcDocStore.Remove(doc.ID)
|
||||
// 清理测试文档,避免积累(SDK Insert 不回填 ID,经 Query 按标题定位)
|
||||
for _, d := range s.DocMemory().Query("健康检查测试文档", 10) {
|
||||
if d.ID != "" && strings.HasPrefix(d.Title, "健康检查测试文档") {
|
||||
s.DocMemory().Remove(d.ID)
|
||||
}
|
||||
}
|
||||
|
||||
elapsed := time.Since(start)
|
||||
@ -582,9 +544,9 @@ func (p *Plugin) testDocStoreRaw() checkResult {
|
||||
}
|
||||
}
|
||||
|
||||
func (p *Plugin) testLLMDriven() checkResult {
|
||||
provider := hcProviderMgr.Default()
|
||||
if provider == nil {
|
||||
func (p *Plugin) testLLMDriven(s *sdk.PluginSDK) checkResult {
|
||||
llmName := s.LLM().CurrentSource()
|
||||
if llmName == "" {
|
||||
return checkResult{Name: "llm_discovery", Status: "skip", Detail: "无可用 LLM Provider", Pass: true}
|
||||
}
|
||||
|
||||
@ -593,7 +555,7 @@ func (p *Plugin) testLLMDriven() checkResult {
|
||||
defer cancel()
|
||||
|
||||
// 收集所有工具定义(排除健康检查自身的工具以避免循环测试)
|
||||
toolDefs := p.collectToolDefsForLLM()
|
||||
toolDefs := p.collectToolDefsForLLM(s)
|
||||
|
||||
if len(toolDefs) == 0 {
|
||||
return checkResult{Name: "llm_discovery", Status: "skip", Detail: "没有可测试的工具", Pass: true}
|
||||
@ -608,15 +570,14 @@ func (p *Plugin) testLLMDriven() checkResult {
|
||||
// 构建 prompt
|
||||
prompt := p.buildDiscoveryPrompt(toolDefs)
|
||||
|
||||
msgs := []agentAPI.Message{{Role: "user", Content: prompt}}
|
||||
msgs := []sdk.LLMMessage{{Role: "user", Content: prompt}}
|
||||
tools := convertToolDefs(toolDefs)
|
||||
|
||||
llmName := provider.Name()
|
||||
turnCount := 0
|
||||
toolCallCount := 0
|
||||
|
||||
for turn := 0; turn < p.llmMaxTurns; turn++ {
|
||||
resp, err := provider.Chat(ctx, &agentAPI.CompletionRequest{
|
||||
resp, err := s.LLM().Chat(ctx, &sdk.LLMCompletionRequest{
|
||||
Messages: msgs,
|
||||
MaxTokens: p.llmMaxTokens,
|
||||
Tools: tools,
|
||||
@ -638,12 +599,12 @@ func (p *Plugin) testLLMDriven() checkResult {
|
||||
break
|
||||
}
|
||||
|
||||
msgs = append(msgs, agentAPI.Message{Role: "assistant", Content: resp.Content, ToolCalls: resp.ToolCalls})
|
||||
msgs = append(msgs, sdk.LLMMessage{Role: "assistant", Content: resp.Content, ToolCalls: resp.ToolCalls})
|
||||
|
||||
for _, tc := range resp.ToolCalls {
|
||||
toolCallCount++
|
||||
content := p.executeToolForLLM(tc)
|
||||
msgs = append(msgs, agentAPI.Message{Role: "tool", ToolCallID: tc.ID, Content: content})
|
||||
content := p.executeToolForLLM(s, tc)
|
||||
msgs = append(msgs, sdk.LLMMessage{Role: "tool", ToolCallID: tc.ID, Content: content})
|
||||
}
|
||||
}
|
||||
|
||||
@ -666,7 +627,7 @@ func (p *Plugin) testLLMDriven() checkResult {
|
||||
|
||||
// collectToolDefsForLLM 收集全部已注册的工具定义供 LLM 发现和测试。
|
||||
// 动态排除本插件自身注册的工具(通过 selfToolNames),避免 LLM 自我循环调用。
|
||||
func (p *Plugin) collectToolDefsForLLM() []sdk.ToolDef {
|
||||
func (p *Plugin) collectToolDefsForLLM(s *sdk.PluginSDK) []sdk.ToolDef {
|
||||
seen := map[string]bool{}
|
||||
var defs []sdk.ToolDef
|
||||
|
||||
@ -678,14 +639,12 @@ func (p *Plugin) collectToolDefsForLLM() []sdk.ToolDef {
|
||||
defs = append(defs, d)
|
||||
}
|
||||
|
||||
if hcStageHost != nil {
|
||||
for _, d := range hcStageHost.GetToolDefs() {
|
||||
if s.Tool() != nil {
|
||||
for _, d := range s.Tool().GetToolDefs() {
|
||||
addDef(d)
|
||||
}
|
||||
}
|
||||
if hcIOMgr != nil {
|
||||
for _, d := range hcIOMgr.GetAllTools() {
|
||||
addDef(sdk.ToolDef{Name: d.Name, Description: d.Description, Parameters: d.Parameters})
|
||||
for _, d := range s.Tool().GetAllTools() {
|
||||
addDef(d)
|
||||
}
|
||||
}
|
||||
|
||||
@ -715,10 +674,10 @@ func (p *Plugin) buildDiscoveryPrompt(toolDefs []sdk.ToolDef) string {
|
||||
}
|
||||
|
||||
// executeToolForLLM 在 LLM 工具循环中执行工具调用。
|
||||
// healthcheck_report 通过 StageHost 路由到自身注册的 handler,负责收集 LLM 上报。
|
||||
func (p *Plugin) executeToolForLLM(tc agentAPI.ToolCall) string {
|
||||
if hcStageHost != nil {
|
||||
result, err := hcStageHost.ExecuteTool(tc.Name, tc.Arguments)
|
||||
// healthcheck_report 经 SDK ToolAPI 路由到自身注册的 handler,负责收集 LLM 上报。
|
||||
func (p *Plugin) executeToolForLLM(s *sdk.PluginSDK, tc sdk.LLMToolCall) string {
|
||||
if s.Tool() != nil {
|
||||
result, err := s.Tool().ExecuteTool(tc.Name, tc.Arguments)
|
||||
if err != nil {
|
||||
return fmt.Sprintf("调用工具 %s 失败: %v", tc.Name, err)
|
||||
}
|
||||
@ -726,7 +685,7 @@ func (p *Plugin) executeToolForLLM(tc agentAPI.ToolCall) string {
|
||||
return string(data)
|
||||
}
|
||||
|
||||
return fmt.Sprintf("工具 %s 不可执行(StageHost 未初始化)", tc.Name)
|
||||
return fmt.Sprintf("工具 %s 不可执行(工具注册表未初始化)", tc.Name)
|
||||
}
|
||||
|
||||
func convertToolDefs(defs []sdk.ToolDef) []interface{} {
|
||||
@ -744,11 +703,11 @@ func convertToolDefs(defs []sdk.ToolDef) []interface{} {
|
||||
return tools
|
||||
}
|
||||
|
||||
func (p *Plugin) checkMemory() (interface{}, error) {
|
||||
if hcMemory == nil {
|
||||
func (p *Plugin) checkMemory(s *sdk.PluginSDK) (interface{}, error) {
|
||||
if s.Memory() == nil {
|
||||
return map[string]interface{}{"status": "skip", "pass": true, "detail": "图记忆未初始化"}, nil
|
||||
}
|
||||
r := p.testMemoryRaw()
|
||||
r := p.testMemoryRaw(s)
|
||||
c := map[string]interface{}{
|
||||
"status": r.Status,
|
||||
"pass": r.Pass,
|
||||
|
||||
@ -10,7 +10,6 @@ import (
|
||||
"gitcode.com/JianFeeeee/HomeAgent/internal/knowledge"
|
||||
"gitcode.com/JianFeeeee/HomeAgent/internal/memory"
|
||||
doc "gitcode.com/JianFeeeee/HomeAgent/internal/memory/document"
|
||||
"gitcode.com/JianFeeeee/HomeAgent/internal/plugin"
|
||||
sdk "gitcode.com/JianFeeeee/HomeAgent/internal/sdk"
|
||||
)
|
||||
|
||||
@ -34,16 +33,28 @@ func (tc *toolCapture) RegisterTool(name string, def sdk.ToolDef, handler sdk.To
|
||||
func (tc *toolCapture) RegisterStage(stage sdk.Stage, handler sdk.StageHandler) {}
|
||||
func (tc *toolCapture) RegisterAPI(name string) error { return nil }
|
||||
|
||||
func setupPlugin() (*Plugin, *toolCapture, error) {
|
||||
sh := agentCore.NewStageHost()
|
||||
iom := agentIO.NewIOManager()
|
||||
pr := plugin.NewRegistry()
|
||||
func newTestSDK(cfg sdk.SDKConfig) *sdk.PluginSDK {
|
||||
if cfg.Settings == nil {
|
||||
cfg.Settings = sdk.NewSettings("healthcheck", nil)
|
||||
}
|
||||
if cfg.Tool == nil {
|
||||
cfg.Tool = sdk.NewTool(agentCore.NewStageHost(), agentIO.NewIOManager())
|
||||
}
|
||||
return sdk.New("healthcheck", cfg)
|
||||
}
|
||||
|
||||
Configure(sh, iom, pr, nil, nil, nil, nil, nil)
|
||||
p := New("healthcheck")
|
||||
func setupPlugin() (*Plugin, *toolCapture, error) {
|
||||
return setupPluginWith(sdk.SDKConfig{})
|
||||
}
|
||||
|
||||
func setupPluginWith(cfg sdk.SDKConfig) (*Plugin, *toolCapture, error) {
|
||||
tc := newToolCapture()
|
||||
sdk := sdk.New("healthcheck", sdk.SDKConfig{RegTool: tc.RegisterTool, RegStage: tc.RegisterStage, RegAPI: tc.RegisterAPI})
|
||||
if err := p.Start(sdk); err != nil {
|
||||
cfg.RegTool = tc.RegisterTool
|
||||
cfg.RegStage = tc.RegisterStage
|
||||
cfg.RegAPI = tc.RegisterAPI
|
||||
p := New("healthcheck")
|
||||
s := newTestSDK(cfg)
|
||||
if err := p.Start(s); err != nil {
|
||||
return nil, nil, err
|
||||
}
|
||||
return p, tc, nil
|
||||
@ -185,15 +196,8 @@ func TestHealthcheckWithMemory(t *testing.T) {
|
||||
}
|
||||
defer memDB.Close()
|
||||
|
||||
sh := agentCore.NewStageHost()
|
||||
iom := agentIO.NewIOManager()
|
||||
pr := plugin.NewRegistry()
|
||||
|
||||
Configure(sh, iom, pr, memDB, nil, nil, nil, nil)
|
||||
p := New("healthcheck")
|
||||
tc := newToolCapture()
|
||||
sdk := sdk.New("healthcheck", sdk.SDKConfig{RegTool: tc.RegisterTool, RegStage: tc.RegisterStage, RegAPI: tc.RegisterAPI})
|
||||
if err := p.Start(sdk); err != nil {
|
||||
_, tc, err := setupPluginWith(sdk.SDKConfig{Memory: sdk.NewGraphMemory(memDB)})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
@ -228,15 +232,8 @@ func TestHealthcheckWithKnowledge(t *testing.T) {
|
||||
}
|
||||
defer ks.Stop()
|
||||
|
||||
sh := agentCore.NewStageHost()
|
||||
iom := agentIO.NewIOManager()
|
||||
pr := plugin.NewRegistry()
|
||||
|
||||
Configure(sh, iom, pr, nil, ks, nil, nil, nil)
|
||||
p := New("healthcheck")
|
||||
tc := newToolCapture()
|
||||
sdk := sdk.New("healthcheck", sdk.SDKConfig{RegTool: tc.RegisterTool, RegStage: tc.RegisterStage, RegAPI: tc.RegisterAPI})
|
||||
if err := p.Start(sdk); err != nil {
|
||||
_, tc, err := setupPluginWith(sdk.SDKConfig{Knowledge: sdk.NewKnowledge(ks)})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
@ -282,15 +279,8 @@ func TestHealthcheckWithDocStore(t *testing.T) {
|
||||
}
|
||||
defer ds.Stop()
|
||||
|
||||
sh := agentCore.NewStageHost()
|
||||
iom := agentIO.NewIOManager()
|
||||
pr := plugin.NewRegistry()
|
||||
|
||||
Configure(sh, iom, pr, nil, nil, ds, nil, nil)
|
||||
p := New("healthcheck")
|
||||
tc := newToolCapture()
|
||||
sdk := sdk.New("healthcheck", sdk.SDKConfig{RegTool: tc.RegisterTool, RegStage: tc.RegisterStage, RegAPI: tc.RegisterAPI})
|
||||
if err := p.Start(sdk); err != nil {
|
||||
_, tc, err := setupPluginWith(sdk.SDKConfig{DocMemory: sdk.NewDocMemory(ds)})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
@ -340,11 +330,4 @@ func TestLLMReportCollection(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestConfigureNilStageHost(t *testing.T) {
|
||||
Configure(nil, nil, nil, nil, nil, nil, nil, nil)
|
||||
if hcStageHost != nil {
|
||||
t.Fatal("expected hcStageHost to be nil")
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
|
||||
@ -10,15 +10,13 @@ import (
|
||||
agentAPI "gitcode.com/JianFeeeee/HomeAgent/internal/agent/api"
|
||||
agentCore "gitcode.com/JianFeeeee/HomeAgent/internal/agent/core"
|
||||
agentIO "gitcode.com/JianFeeeee/HomeAgent/internal/agent/io"
|
||||
internalConfig "gitcode.com/JianFeeeee/HomeAgent/internal/config"
|
||||
"gitcode.com/JianFeeeee/HomeAgent/internal/knowledge"
|
||||
luaVM "gitcode.com/JianFeeeee/HomeAgent/internal/lua"
|
||||
"gitcode.com/JianFeeeee/HomeAgent/internal/memory"
|
||||
doc "gitcode.com/JianFeeeee/HomeAgent/internal/memory/document"
|
||||
"gitcode.com/JianFeeeee/HomeAgent/internal/plugin"
|
||||
cli "gitcode.com/JianFeeeee/HomeAgent/internal/plugins/cli"
|
||||
healthcheck "gitcode.com/JianFeeeee/HomeAgent/internal/plugins/healthcheck"
|
||||
openclaw "gitcode.com/JianFeeeee/HomeAgent/internal/plugins/clawhubadapter"
|
||||
webui "gitcode.com/JianFeeeee/HomeAgent/internal/plugins/webui"
|
||||
sdk "gitcode.com/JianFeeeee/HomeAgent/internal/sdk"
|
||||
)
|
||||
|
||||
@ -77,11 +75,17 @@ func setupIntegrationWithProvider(t *testing.T, pm *agentAPI.ProviderManager) *t
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
pluginReg.SetStageHost(stageHost)
|
||||
pluginReg.SetProviderManager(pm)
|
||||
pluginReg.SetKnowledge(ks)
|
||||
pluginReg.SetDocStore(docStore)
|
||||
|
||||
cli.DefaultSocket = filepath.Join(tmpDir, "cli.sock")
|
||||
openclaw.SkillsDir = filepath.Join(tmpDir, "skills")
|
||||
os.MkdirAll(openclaw.SkillsDir, 0755)
|
||||
webui.Configure(":0", nil, memDB, nil, nil, nil, iom, nil, ks, nil, nil, pluginReg, nil, nil)
|
||||
healthcheck.Configure(stageHost, iom, pluginReg, memDB, ks, docStore, pm, nil)
|
||||
|
||||
// 经 ConfigRegistry 装配内核路径配置(clawhubadapter/pluginmgr 等经 SDK settings 读取)
|
||||
cfgReg := internalConfig.NewConfigRegistry("")
|
||||
cfgReg.SeedDefaults(tmpDir)
|
||||
pluginReg.SetConfigRegistry(cfgReg)
|
||||
|
||||
plgDir := filepath.Join(tmpDir, "plugins")
|
||||
os.MkdirAll(plgDir, 0755)
|
||||
@ -229,7 +233,7 @@ func TestIntegrationCmdRunStderr(t *testing.T) {
|
||||
defer env.cleanup()
|
||||
|
||||
result, err := env.stageHost.ExecuteTool("cmd_run", map[string]interface{}{
|
||||
"command": "echo stderr_test >&2",
|
||||
"command": "sh -c \"echo stderr_test >&2\"",
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
@ -520,7 +524,7 @@ func TestIntegrationLLMDrivenDiscoveryWithRealKey(t *testing.T) {
|
||||
Model: "deepseek-v4-flash",
|
||||
BaseURL: "https://api.deepseek.com",
|
||||
APIKey: apiKey,
|
||||
}, vm, "deepseek"))
|
||||
}, vm, "deepseek", "deepseek"))
|
||||
|
||||
// Setup — 加载所有真实内置插件
|
||||
env := setupIntegrationWithProvider(t, pm)
|
||||
|
||||
@ -66,11 +66,7 @@ var downloadClient = &http.Client{
|
||||
},
|
||||
}
|
||||
|
||||
var (
|
||||
PluginDir string // 由 main.go 设置
|
||||
Reg *plugin.Registry // 由 main.go 设置
|
||||
HTTPAddr = "127.0.0.1:9876" // 监听地址,可被 main.go 覆写或 settings 配置
|
||||
)
|
||||
var HTTPAddr = "127.0.0.1:9876" // 监听地址,可被 settings 配置
|
||||
|
||||
func init() {
|
||||
plugin.RegisterPluginMeta("pluginmgr", "插件管理", "Plugin Manager")
|
||||
@ -80,12 +76,14 @@ func init() {
|
||||
}
|
||||
|
||||
type Plugin struct {
|
||||
name string
|
||||
mu sync.Mutex
|
||||
server *http.Server
|
||||
mux *http.ServeMux
|
||||
listen net.Listener
|
||||
httpURL string
|
||||
name string
|
||||
mu sync.Mutex
|
||||
server *http.Server
|
||||
mux *http.ServeMux
|
||||
listen net.Listener
|
||||
httpURL string
|
||||
sdk *sdk.PluginSDK
|
||||
pluginDir string
|
||||
}
|
||||
|
||||
func New(name string) *Plugin {
|
||||
@ -96,6 +94,7 @@ func (p *Plugin) Name() string { return p.name }
|
||||
|
||||
func (p *Plugin) Start(s *sdk.PluginSDK) error {
|
||||
s.SetAutoRestart(true)
|
||||
p.sdk = s
|
||||
s.Settings().RegisterDef(sdk.ConfigDef{
|
||||
Key: "http_addr",
|
||||
Default: HTTPAddr,
|
||||
@ -111,6 +110,12 @@ func (p *Plugin) Start(s *sdk.PluginSDK) error {
|
||||
}
|
||||
}
|
||||
|
||||
if v, _ := s.Settings().GetCore("plugin.dir"); v != nil {
|
||||
if dir, ok := v.(string); ok && dir != "" {
|
||||
p.pluginDir = dir
|
||||
}
|
||||
}
|
||||
|
||||
p.registerTools(s)
|
||||
|
||||
if HTTPAddr != "" {
|
||||
@ -374,7 +379,7 @@ func (p *Plugin) installFromData(data []byte) (interface{}, error) {
|
||||
}, nil
|
||||
}
|
||||
|
||||
dir := PluginDir
|
||||
dir := p.pluginDir
|
||||
if dir == "" {
|
||||
return map[string]interface{}{"error": "plugin dir not configured"}, nil
|
||||
}
|
||||
@ -409,7 +414,7 @@ func (p *Plugin) installFromData(data []byte) (interface{}, error) {
|
||||
}
|
||||
|
||||
func (p *Plugin) listPlugins() (interface{}, error) {
|
||||
dir := PluginDir
|
||||
dir := p.pluginDir
|
||||
if dir == "" {
|
||||
return []map[string]interface{}{}, nil
|
||||
}
|
||||
@ -446,7 +451,7 @@ func (p *Plugin) listPlugins() (interface{}, error) {
|
||||
}
|
||||
|
||||
func (p *Plugin) removePlugin(name string) (interface{}, error) {
|
||||
dir := filepath.Join(PluginDir, name)
|
||||
dir := filepath.Join(p.pluginDir, name)
|
||||
if _, err := os.Stat(dir); os.IsNotExist(err) {
|
||||
return map[string]interface{}{"error": "plugin not found", "name": name}, nil
|
||||
}
|
||||
@ -456,8 +461,10 @@ func (p *Plugin) removePlugin(name string) (interface{}, error) {
|
||||
}
|
||||
|
||||
// 同步清理禁用表
|
||||
if Reg != nil {
|
||||
Reg.EnablePlugin(name)
|
||||
if p.sdk != nil && p.sdk.PluginMgr() != nil {
|
||||
if err := p.sdk.PluginMgr().EnablePlugin(name); err != nil {
|
||||
log.Printf("[pluginmgr] enable %s after remove: %v", name, err)
|
||||
}
|
||||
}
|
||||
|
||||
return map[string]interface{}{
|
||||
@ -468,7 +475,7 @@ func (p *Plugin) removePlugin(name string) (interface{}, error) {
|
||||
}
|
||||
|
||||
func (p *Plugin) pluginInfo(name string) (interface{}, error) {
|
||||
dir := filepath.Join(PluginDir, name)
|
||||
dir := filepath.Join(p.pluginDir, name)
|
||||
m, err := plugin.ReadManifest(dir)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("plugin %q not found", name)
|
||||
|
||||
@ -9,28 +9,13 @@ import (
|
||||
"io"
|
||||
"log"
|
||||
"net/http"
|
||||
"os"
|
||||
"sort"
|
||||
"strconv"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
agentAPI "gitcode.com/JianFeeeee/HomeAgent/internal/agent/api"
|
||||
agentCore "gitcode.com/JianFeeeee/HomeAgent/internal/agent/core"
|
||||
agentIO "gitcode.com/JianFeeeee/HomeAgent/internal/agent/io"
|
||||
internalConfig "gitcode.com/JianFeeeee/HomeAgent/internal/config"
|
||||
"gitcode.com/JianFeeeee/HomeAgent/internal/events"
|
||||
"gitcode.com/JianFeeeee/HomeAgent/internal/knowledge"
|
||||
"gitcode.com/JianFeeeee/HomeAgent/internal/meta"
|
||||
luaVM "gitcode.com/JianFeeeee/HomeAgent/internal/lua"
|
||||
"gitcode.com/JianFeeeee/HomeAgent/internal/memory"
|
||||
"gitcode.com/JianFeeeee/HomeAgent/internal/memory/text"
|
||||
"gitcode.com/JianFeeeee/HomeAgent/internal/plugin"
|
||||
sdk "gitcode.com/JianFeeeee/HomeAgent/internal/sdk"
|
||||
"gitcode.com/JianFeeeee/HomeAgent/internal/skill"
|
||||
"gitcode.com/JianFeeeee/HomeAgent/internal/supervisor"
|
||||
"gitcode.com/JianFeeeee/HomeAgent/internal/tracker"
|
||||
"gitcode.com/JianFeeeee/HomeAgent/pkg/types"
|
||||
)
|
||||
|
||||
@ -49,26 +34,24 @@ func init() {
|
||||
}
|
||||
|
||||
type Handler struct {
|
||||
supervisor *supervisor.Daemon
|
||||
memory *memory.GraphDB
|
||||
indexer *memory.Indexer
|
||||
skills *skill.Manager
|
||||
lua *luaVM.VM
|
||||
config *types.Config
|
||||
startTime time.Time
|
||||
iom *agentIO.IOManager
|
||||
textMem *text.Memory
|
||||
knowledge *knowledge.Store
|
||||
tracker *tracker.Tracker
|
||||
cfgReg *internalConfig.ConfigRegistry
|
||||
pluginReg *plugin.Registry
|
||||
pluginMgr sdk.PluginManager
|
||||
eventBus *events.Bus
|
||||
statusProvider agentCore.StatusProvider
|
||||
providerMgr *agentAPI.ProviderManager
|
||||
baseAPIKey string
|
||||
sessionMu sync.Mutex
|
||||
sessions map[string]time.Time
|
||||
sdk *sdk.PluginSDK
|
||||
supervisor sdk.SupervisorAPI
|
||||
memory sdk.MemoryAPI
|
||||
indexer sdk.IndexerAPI
|
||||
skills sdk.SkillAPI
|
||||
adapter sdk.AdapterAPI
|
||||
config sdk.ConfigAPI
|
||||
startTime time.Time
|
||||
textMem sdk.TextMemoryAPI
|
||||
knowledge sdk.KnowledgeAPI
|
||||
tracker sdk.TrackerAPI
|
||||
settings sdk.SettingsAPI
|
||||
pluginMgr sdk.PluginManager
|
||||
status sdk.StatusAPI
|
||||
llm sdk.LLMAPI
|
||||
|
||||
sessionMu sync.Mutex
|
||||
sessions map[string]time.Time
|
||||
|
||||
chatMu sync.Mutex
|
||||
chatHistory []ChatMsg
|
||||
@ -104,51 +87,64 @@ type termState struct {
|
||||
created time.Time
|
||||
}
|
||||
|
||||
func (h *Handler) SetPluginMgr(mgr sdk.PluginManager) { h.pluginMgr = mgr }
|
||||
|
||||
const maxChatHistory = 200
|
||||
const maxCmdHistory = 100
|
||||
const maxTerminals = 50
|
||||
|
||||
func NewHandler(sup *supervisor.Daemon, mem *memory.GraphDB, sk *skill.Manager, lua *luaVM.VM, cfg *types.Config, iom *agentIO.IOManager, tm *text.Memory, ks *knowledge.Store, tr *tracker.Tracker, cr *internalConfig.ConfigRegistry, pr *plugin.Registry, evBus *events.Bus, sp agentCore.StatusProvider, pm *agentAPI.ProviderManager, baseKey string) *Handler {
|
||||
var idx *memory.Indexer
|
||||
if mem != nil {
|
||||
idx = memory.NewIndexer(mem)
|
||||
func NewHandler(s *sdk.PluginSDK) *Handler {
|
||||
var (
|
||||
sup sdk.SupervisorAPI
|
||||
mem sdk.MemoryAPI
|
||||
idx sdk.IndexerAPI
|
||||
sk sdk.SkillAPI
|
||||
ad sdk.AdapterAPI
|
||||
cfg sdk.ConfigAPI
|
||||
tm sdk.TextMemoryAPI
|
||||
ks sdk.KnowledgeAPI
|
||||
tr sdk.TrackerAPI
|
||||
se sdk.SettingsAPI
|
||||
pm sdk.PluginManager
|
||||
st sdk.StatusAPI
|
||||
llm sdk.LLMAPI
|
||||
)
|
||||
if s != nil {
|
||||
sup, mem, idx = s.Supervisor(), s.Memory(), s.Indexer()
|
||||
sk, ad, cfg = s.Skill(), s.Adapter(), s.Config()
|
||||
tm, ks, tr = s.TextMemory(), s.Knowledge(), s.Tracker()
|
||||
se, pm = s.Settings(), s.PluginMgr()
|
||||
st, llm = s.Status(), s.LLM()
|
||||
}
|
||||
h := &Handler{
|
||||
supervisor: sup,
|
||||
memory: mem,
|
||||
indexer: idx,
|
||||
skills: sk,
|
||||
lua: lua,
|
||||
config: cfg,
|
||||
startTime: time.Now(),
|
||||
iom: iom,
|
||||
textMem: tm,
|
||||
knowledge: ks,
|
||||
tracker: tr,
|
||||
cfgReg: cr,
|
||||
pluginReg: pr,
|
||||
eventBus: evBus,
|
||||
statusProvider: sp,
|
||||
providerMgr: pm,
|
||||
baseAPIKey: baseKey,
|
||||
sessions: make(map[string]time.Time),
|
||||
termStates: make(map[string]*termState),
|
||||
sdk: s,
|
||||
supervisor: sup,
|
||||
memory: mem,
|
||||
indexer: idx,
|
||||
skills: sk,
|
||||
adapter: ad,
|
||||
config: cfg,
|
||||
startTime: time.Now(),
|
||||
textMem: tm,
|
||||
knowledge: ks,
|
||||
tracker: tr,
|
||||
settings: se,
|
||||
pluginMgr: pm,
|
||||
status: st,
|
||||
llm: llm,
|
||||
sessions: make(map[string]time.Time),
|
||||
termStates: make(map[string]*termState),
|
||||
}
|
||||
h.loadChatHistory()
|
||||
if evBus != nil {
|
||||
if s != nil {
|
||||
go h.trackToolEvents()
|
||||
}
|
||||
return h
|
||||
}
|
||||
|
||||
func (h *Handler) loadChatHistory() {
|
||||
if h.cfgReg == nil {
|
||||
if h.settings == nil {
|
||||
return
|
||||
}
|
||||
ps := h.cfgReg.PluginConfig("webui")
|
||||
v, err := ps.Get("chathistory")
|
||||
v, err := h.settings.Get("chathistory")
|
||||
if err != nil || v == nil {
|
||||
return
|
||||
}
|
||||
@ -166,12 +162,15 @@ func (h *Handler) loadChatHistory() {
|
||||
}
|
||||
|
||||
func (h *Handler) trackToolEvents() {
|
||||
h.eventBus.Subscribe(events.EventToolCall, func(ev *events.Event) {
|
||||
if h.sdk == nil {
|
||||
return
|
||||
}
|
||||
h.sdk.Subscribe(sdk.EventToolCall, func(ev *sdk.Event) {
|
||||
h.handleToolEvent(ev)
|
||||
})
|
||||
}
|
||||
|
||||
func (h *Handler) handleToolEvent(ev *events.Event) {
|
||||
func (h *Handler) handleToolEvent(ev *sdk.Event) {
|
||||
payload := ev.Payload
|
||||
tool, _ := payload["tool"].(string)
|
||||
args, _ := payload["args"].(map[string]interface{})
|
||||
@ -235,20 +234,19 @@ func getStr(m map[string]interface{}, key string) string {
|
||||
|
||||
func (h *Handler) getWebUIConfig() (apiKey, username, password string, ttl time.Duration) {
|
||||
ttl = 24 * time.Hour
|
||||
if h.cfgReg == nil {
|
||||
if h.settings == nil {
|
||||
return
|
||||
}
|
||||
ps := h.cfgReg.PluginConfig("webui")
|
||||
if v, _ := ps.Get("api_key"); v != nil {
|
||||
if v, _ := h.settings.Get("api_key"); v != nil {
|
||||
apiKey, _ = v.(string)
|
||||
}
|
||||
if v, _ := ps.Get("username"); v != nil {
|
||||
if v, _ := h.settings.Get("username"); v != nil {
|
||||
username, _ = v.(string)
|
||||
}
|
||||
if v, _ := ps.Get("password"); v != nil {
|
||||
if v, _ := h.settings.Get("password"); v != nil {
|
||||
password, _ = v.(string)
|
||||
}
|
||||
if v, _ := ps.Get("session_ttl_hours"); v != nil {
|
||||
if v, _ := h.settings.Get("session_ttl_hours"); v != nil {
|
||||
switch n := v.(type) {
|
||||
case float64:
|
||||
if n > 0 { ttl = time.Duration(n) * time.Hour }
|
||||
@ -433,12 +431,15 @@ func (h *Handler) handleStatus(w http.ResponseWriter, r *http.Request) {
|
||||
http.Error(w, "method not allowed", http.StatusMethodNotAllowed)
|
||||
return
|
||||
}
|
||||
agents := h.supervisor.ListAgents()
|
||||
agentCount := 0
|
||||
if h.supervisor != nil {
|
||||
agentCount = len(h.supervisor.ListAgents())
|
||||
}
|
||||
writeJSON(w, http.StatusOK, map[string]interface{}{
|
||||
"status": "running",
|
||||
"uptime": time.Since(h.startTime).Round(time.Second).String(),
|
||||
"agents": len(agents),
|
||||
"version": meta.Version,
|
||||
"agents": agentCount,
|
||||
"version": sdk.SDKVersion,
|
||||
"startedAt": h.startTime,
|
||||
})
|
||||
}
|
||||
@ -448,16 +449,20 @@ func (h *Handler) handleKernel(w http.ResponseWriter, r *http.Request) {
|
||||
http.Error(w, "method not allowed", http.StatusMethodNotAllowed)
|
||||
return
|
||||
}
|
||||
if h.statusProvider == nil {
|
||||
if h.status == nil {
|
||||
writeJSON(w, http.StatusServiceUnavailable, map[string]string{"error": "kernel status provider not available"})
|
||||
return
|
||||
}
|
||||
writeJSON(w, http.StatusOK, h.statusProvider.GetKernelStatus())
|
||||
writeJSON(w, http.StatusOK, h.status.GetKernelStatus())
|
||||
}
|
||||
|
||||
func (h *Handler) handleAgents(w http.ResponseWriter, r *http.Request) {
|
||||
switch r.Method {
|
||||
case http.MethodGet:
|
||||
if h.supervisor == nil {
|
||||
writeJSON(w, http.StatusServiceUnavailable, map[string]string{"error": "supervisor not available"})
|
||||
return
|
||||
}
|
||||
agents := h.supervisor.ListAgents()
|
||||
writeJSON(w, http.StatusOK, map[string]interface{}{"agents": agents})
|
||||
case http.MethodPost:
|
||||
@ -470,7 +475,10 @@ func (h *Handler) handleAgents(w http.ResponseWriter, r *http.Request) {
|
||||
writeJSON(w, http.StatusBadRequest, map[string]string{"error": "agent id is required"})
|
||||
return
|
||||
}
|
||||
h.config.Agents = append(h.config.Agents, cfg)
|
||||
if h.config != nil {
|
||||
kcfg := h.config.Get()
|
||||
kcfg.Agents = append(kcfg.Agents, cfg)
|
||||
}
|
||||
writeJSON(w, http.StatusCreated, map[string]string{"id": string(cfg.ID)})
|
||||
default:
|
||||
http.Error(w, "method not allowed", http.StatusMethodNotAllowed)
|
||||
@ -484,7 +492,11 @@ func (h *Handler) handleAgentByID(w http.ResponseWriter, r *http.Request) {
|
||||
if len(parts) == 1 {
|
||||
switch r.Method {
|
||||
case http.MethodGet:
|
||||
status, err := h.supervisor.GetAgentStatus(agentID)
|
||||
if h.supervisor == nil {
|
||||
writeJSON(w, http.StatusServiceUnavailable, map[string]string{"error": "supervisor not available"})
|
||||
return
|
||||
}
|
||||
status, err := h.supervisor.GetAgentStatus(string(agentID))
|
||||
if err != nil {
|
||||
writeJSON(w, http.StatusNotFound, map[string]string{"error": err.Error()})
|
||||
return
|
||||
@ -515,7 +527,11 @@ func (h *Handler) handleSnapshots(w http.ResponseWriter, r *http.Request, agentI
|
||||
case http.MethodGet:
|
||||
writeJSON(w, http.StatusOK, map[string]interface{}{"agent_id": agentID, "snapshots": []map[string]interface{}{}})
|
||||
case http.MethodPost:
|
||||
snap, err := h.supervisor.PreActionSnapshot(agentID)
|
||||
if h.supervisor == nil {
|
||||
writeJSON(w, http.StatusServiceUnavailable, map[string]string{"error": "supervisor not available"})
|
||||
return
|
||||
}
|
||||
snap, err := h.supervisor.PreActionSnapshot(string(agentID))
|
||||
if err != nil {
|
||||
writeJSON(w, http.StatusInternalServerError, map[string]string{"error": err.Error()})
|
||||
return
|
||||
@ -531,8 +547,12 @@ func (h *Handler) handleRollback(w http.ResponseWriter, r *http.Request, agentID
|
||||
http.Error(w, "method not allowed", http.StatusMethodNotAllowed)
|
||||
return
|
||||
}
|
||||
if h.supervisor == nil {
|
||||
writeJSON(w, http.StatusServiceUnavailable, map[string]string{"error": "supervisor not available"})
|
||||
return
|
||||
}
|
||||
snapID := types.SnapshotID(parts[2])
|
||||
if err := h.supervisor.RollbackAgent(agentID, snapID); err != nil {
|
||||
if err := h.supervisor.RollbackAgent(string(agentID), string(snapID)); err != nil {
|
||||
writeJSON(w, http.StatusInternalServerError, map[string]string{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
@ -598,28 +618,25 @@ func (h *Handler) handleMemory(w http.ResponseWriter, r *http.Request) {
|
||||
if depth <= 0 {
|
||||
depth = 2
|
||||
}
|
||||
result, err := h.memory.Recall(keywords, nil, depth, "")
|
||||
entities, relations, err := h.memory.Recall(keywords, depth)
|
||||
if err != nil {
|
||||
writeJSON(w, http.StatusInternalServerError, map[string]string{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
writeJSON(w, http.StatusOK, result)
|
||||
writeJSON(w, http.StatusOK, map[string]interface{}{"entities": entities, "relations": relations})
|
||||
case http.MethodPost:
|
||||
var req struct {
|
||||
Triples []memory.Triple `json:"triples"`
|
||||
SessionID string `json:"session_id"`
|
||||
TurnID int `json:"turn_id"`
|
||||
Triples []sdk.Triple `json:"triples"`
|
||||
}
|
||||
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
|
||||
writeJSON(w, http.StatusBadRequest, map[string]string{"error": "invalid request"})
|
||||
return
|
||||
}
|
||||
ec, rc, err := h.memory.Commit(req.Triples, req.SessionID, req.TurnID)
|
||||
if err != nil {
|
||||
if err := h.memory.Commit(req.Triples); err != nil {
|
||||
writeJSON(w, http.StatusInternalServerError, map[string]string{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
writeJSON(w, http.StatusCreated, map[string]int{"entities_created": ec, "relations_created": rc})
|
||||
writeJSON(w, http.StatusCreated, map[string]interface{}{"status": "committed", "committed": len(req.Triples)})
|
||||
case http.MethodDelete:
|
||||
var req struct {
|
||||
Criteria map[string]string `json:"criteria"`
|
||||
@ -650,7 +667,11 @@ func (h *Handler) handleMemoryContext(w http.ResponseWriter, r *http.Request) {
|
||||
return
|
||||
}
|
||||
userInput := r.URL.Query().Get("q")
|
||||
injected := h.indexer.BuildContext(userInput)
|
||||
injected, err := h.indexer.BuildContext(userInput)
|
||||
if err != nil {
|
||||
writeJSON(w, http.StatusInternalServerError, map[string]string{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
writeJSON(w, http.StatusOK, map[string]interface{}{
|
||||
"context": h.indexer.FormatContext(injected),
|
||||
"summary": injected.Summary,
|
||||
@ -701,12 +722,13 @@ func (h *Handler) handleKnowledge(w http.ResponseWriter, r *http.Request) {
|
||||
}
|
||||
query := r.URL.Query().Get("q")
|
||||
if query != "" {
|
||||
results := h.knowledge.Search(query, 10)
|
||||
results, _ := h.knowledge.Search(query, 10)
|
||||
writeJSON(w, http.StatusOK, map[string]interface{}{"results": results})
|
||||
return
|
||||
}
|
||||
categories, _ := h.knowledge.List()
|
||||
writeJSON(w, http.StatusOK, map[string]interface{}{
|
||||
"categories": h.knowledge.List(),
|
||||
"categories": categories,
|
||||
"stats": h.knowledge.Stats(),
|
||||
})
|
||||
|
||||
@ -795,13 +817,13 @@ func (h *Handler) handleTextMemory(w http.ResponseWriter, r *http.Request) {
|
||||
}
|
||||
|
||||
func (h *Handler) handleAdapters(w http.ResponseWriter, r *http.Request) {
|
||||
if h.lua == nil {
|
||||
if h.adapter == nil {
|
||||
writeJSON(w, http.StatusServiceUnavailable, map[string]string{"error": "lua vm not available"})
|
||||
return
|
||||
}
|
||||
switch r.Method {
|
||||
case http.MethodGet:
|
||||
writeJSON(w, http.StatusOK, map[string]interface{}{"adapters": h.lua.ListAdapters()})
|
||||
writeJSON(w, http.StatusOK, map[string]interface{}{"adapters": h.adapter.List()})
|
||||
case http.MethodPost:
|
||||
var req struct {
|
||||
Name string `json:"name"`
|
||||
@ -811,12 +833,7 @@ func (h *Handler) handleAdapters(w http.ResponseWriter, r *http.Request) {
|
||||
writeJSON(w, http.StatusBadRequest, map[string]string{"error": "invalid request"})
|
||||
return
|
||||
}
|
||||
path := fmt.Sprintf("%s/%s.lua", h.lua.AdapterDir(), req.Name)
|
||||
if err := os.WriteFile(path, []byte(req.Code), 0644); err != nil {
|
||||
writeJSON(w, http.StatusInternalServerError, map[string]string{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
if err := h.lua.LoadAdapter(path); err != nil {
|
||||
if err := h.adapter.Load(req.Name, req.Code); err != nil {
|
||||
writeJSON(w, http.StatusInternalServerError, map[string]string{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
@ -827,7 +844,7 @@ func (h *Handler) handleAdapters(w http.ResponseWriter, r *http.Request) {
|
||||
}
|
||||
|
||||
func (h *Handler) handleAdapterByID(w http.ResponseWriter, r *http.Request) {
|
||||
if h.lua == nil {
|
||||
if h.adapter == nil {
|
||||
writeJSON(w, http.StatusServiceUnavailable, map[string]string{"error": "lua vm not available"})
|
||||
return
|
||||
}
|
||||
@ -838,7 +855,7 @@ func (h *Handler) handleAdapterByID(w http.ResponseWriter, r *http.Request) {
|
||||
}
|
||||
switch r.Method {
|
||||
case http.MethodGet:
|
||||
for _, a := range h.lua.ListAdapters() {
|
||||
for _, a := range h.adapter.List() {
|
||||
if a.Name == name {
|
||||
writeJSON(w, http.StatusOK, a)
|
||||
return
|
||||
@ -846,12 +863,10 @@ func (h *Handler) handleAdapterByID(w http.ResponseWriter, r *http.Request) {
|
||||
}
|
||||
http.NotFound(w, r)
|
||||
case http.MethodDelete:
|
||||
path := fmt.Sprintf("%s/%s.lua", h.lua.AdapterDir(), name)
|
||||
if err := os.Remove(path); err != nil {
|
||||
if err := h.adapter.Remove(name); err != nil {
|
||||
writeJSON(w, http.StatusNotFound, map[string]string{"error": "adapter not found"})
|
||||
return
|
||||
}
|
||||
h.lua.RemoveAdapter(name)
|
||||
writeJSON(w, http.StatusOK, map[string]string{"status": "deleted", "name": name})
|
||||
default:
|
||||
http.Error(w, "method not allowed", http.StatusMethodNotAllowed)
|
||||
@ -865,7 +880,7 @@ func (h *Handler) handleNetwork(w http.ResponseWriter, r *http.Request) {
|
||||
}
|
||||
writeJSON(w, http.StatusOK, map[string]interface{}{
|
||||
"network_status": "monitoring",
|
||||
"endpoints": h.config.Defaults.LLMEndpoints,
|
||||
"endpoints": h.config.Get().Defaults.LLMEndpoints,
|
||||
})
|
||||
}
|
||||
|
||||
@ -876,10 +891,9 @@ func (h *Handler) addChatMsg(msg ChatMsg) {
|
||||
h.chatHistory = h.chatHistory[len(h.chatHistory)-maxChatHistory:]
|
||||
}
|
||||
// persist to webui config table as compact JSON
|
||||
if h.cfgReg != nil {
|
||||
ps := h.cfgReg.PluginConfig("webui")
|
||||
if h.settings != nil {
|
||||
b, _ := json.Marshal(h.chatHistory)
|
||||
ps.Set("chathistory", string(b))
|
||||
_ = h.settings.Set("chathistory", string(b))
|
||||
}
|
||||
h.chatMu.Unlock()
|
||||
}
|
||||
@ -929,7 +943,11 @@ func (h *Handler) handleChat(w http.ResponseWriter, r *http.Request) {
|
||||
}
|
||||
|
||||
h.addChatMsg(ChatMsg{Role: "user", Content: body.Message, Time: time.Now().Format(time.RFC3339)})
|
||||
resp := h.iom.InjectTextSync("cli", body.Message)
|
||||
if h.sdk == nil {
|
||||
writeJSON(w, http.StatusServiceUnavailable, map[string]string{"error": "agent unavailable"})
|
||||
return
|
||||
}
|
||||
resp := h.sdk.InjectTextSync("cli", "cli", body.Message)
|
||||
if resp == nil {
|
||||
writeJSON(w, http.StatusServiceUnavailable, map[string]string{"error": "agent unavailable"})
|
||||
return
|
||||
@ -965,7 +983,7 @@ func (h *Handler) handleChatEvents(w http.ResponseWriter, r *http.Request) {
|
||||
flusher.Flush()
|
||||
|
||||
done := r.Context().Done()
|
||||
if h.eventBus == nil {
|
||||
if h.sdk == nil {
|
||||
fmt.Fprintf(w, "event: error\ndata: {\"msg\":\"event bus unavailable\"}\n\n")
|
||||
flusher.Flush()
|
||||
return
|
||||
@ -994,15 +1012,15 @@ func (h *Handler) handleChatEvents(w http.ResponseWriter, r *http.Request) {
|
||||
var unsubs []func()
|
||||
for _, t := range subTypes {
|
||||
t2 := t
|
||||
unsub := h.eventBus.Subscribe(events.EventType(t2), func(evt *events.Event) {
|
||||
if evt.Type == events.EventToolCall {
|
||||
unsub := h.sdk.Subscribe(sdk.EventType(t2), func(evt *sdk.Event) {
|
||||
if evt.Type == sdk.EventToolCall {
|
||||
toolName, _ := evt.Payload["tool"].(string)
|
||||
log.Printf("[SSE] received tool_call event: tool=%s", toolName)
|
||||
}
|
||||
data, _ := json.Marshal(evt)
|
||||
select {
|
||||
case writeCh <- fmt.Sprintf("event: %s\ndata: %s\n", evt.Type, string(data)):
|
||||
if evt.Type == events.EventToolCall {
|
||||
if evt.Type == sdk.EventToolCall {
|
||||
toolName, _ := evt.Payload["tool"].(string)
|
||||
log.Printf("[SSE] wrote tool_call to writeCh: tool=%s", toolName)
|
||||
}
|
||||
@ -1031,16 +1049,20 @@ func (h *Handler) handleChatEvents(w http.ResponseWriter, r *http.Request) {
|
||||
}
|
||||
|
||||
func (h *Handler) handleConfig(w http.ResponseWriter, r *http.Request) {
|
||||
if h.config == nil {
|
||||
writeJSON(w, http.StatusServiceUnavailable, map[string]string{"error": "config not available"})
|
||||
return
|
||||
}
|
||||
switch r.Method {
|
||||
case http.MethodGet:
|
||||
writeJSON(w, http.StatusOK, h.config)
|
||||
writeJSON(w, http.StatusOK, h.config.Get())
|
||||
case http.MethodPut:
|
||||
var cfg types.Config
|
||||
if err := json.NewDecoder(r.Body).Decode(&cfg); err != nil {
|
||||
writeJSON(w, http.StatusBadRequest, map[string]string{"error": "invalid config"})
|
||||
return
|
||||
}
|
||||
h.config = &cfg
|
||||
h.config.Put(&cfg)
|
||||
writeJSON(w, http.StatusOK, map[string]string{"status": "config_updated"})
|
||||
default:
|
||||
http.Error(w, "method not allowed", http.StatusMethodNotAllowed)
|
||||
@ -1048,7 +1070,7 @@ func (h *Handler) handleConfig(w http.ResponseWriter, r *http.Request) {
|
||||
}
|
||||
|
||||
func (h *Handler) handleSettings(w http.ResponseWriter, r *http.Request) {
|
||||
if h.cfgReg == nil {
|
||||
if h.settings == nil {
|
||||
writeJSON(w, http.StatusNotFound, map[string]string{"error": "config registry not available"})
|
||||
return
|
||||
}
|
||||
@ -1056,30 +1078,35 @@ func (h *Handler) handleSettings(w http.ResponseWriter, r *http.Request) {
|
||||
case http.MethodGet:
|
||||
prefix := r.URL.Query().Get("prefix")
|
||||
values := make(map[string]interface{})
|
||||
meta := make(map[string]*internalConfig.ConfigDef)
|
||||
meta := make(map[string]*sdk.ConfigDef)
|
||||
|
||||
if strings.HasPrefix(prefix, "plugin.") {
|
||||
// 插件配置:从插件自身 config_<name> 表读取
|
||||
pluginName := prefix[7:]
|
||||
ps := h.cfgReg.PluginConfig(pluginName)
|
||||
keys, _ := ps.List("")
|
||||
keys, _ := h.settings.ListPlugin(pluginName, "")
|
||||
for _, k := range keys {
|
||||
v, _ := ps.Get(k)
|
||||
v, _ := h.settings.GetPlugin(pluginName, k)
|
||||
fullKey := prefix + "." + k
|
||||
values[fullKey] = v
|
||||
if def := h.cfgReg.GetDef(fullKey); def != nil {
|
||||
meta[fullKey] = def
|
||||
}
|
||||
}
|
||||
for _, def := range h.settings.DefsPlugin(pluginName, "") {
|
||||
fullKey := prefix + "." + def.Key
|
||||
meta[fullKey] = def
|
||||
}
|
||||
} else {
|
||||
// 核心配置:从 core config 表读取
|
||||
keys := h.cfgReg.List(prefix)
|
||||
for _, k := range keys {
|
||||
v, _ := h.cfgReg.Get(k)
|
||||
values[k] = v
|
||||
// 核心配置:从 core config 表读取(键可为任意前缀,如 core.llm.*、webui.*)
|
||||
all := h.settings.Dump()
|
||||
var keys []string
|
||||
for k := range all {
|
||||
if strings.HasPrefix(k, prefix) {
|
||||
keys = append(keys, k)
|
||||
}
|
||||
}
|
||||
defs := h.cfgReg.ListDefs(prefix)
|
||||
for _, d := range defs {
|
||||
sort.Strings(keys)
|
||||
for _, k := range keys {
|
||||
values[k] = all[k]
|
||||
}
|
||||
for _, d := range h.settings.DefsCore(prefix) {
|
||||
meta[d.Key] = d
|
||||
// 有 def 但 DB 中尚无值的 key,用 default 填充以便在 WebUI 中显示和编辑
|
||||
if _, exists := values[d.Key]; !exists {
|
||||
@ -1087,29 +1114,35 @@ func (h *Handler) handleSettings(w http.ResponseWriter, r *http.Request) {
|
||||
}
|
||||
}
|
||||
// 无前缀时同时加载所有插件配置
|
||||
if prefix == "" && h.pluginReg != nil {
|
||||
for _, p := range h.pluginReg.List() {
|
||||
ps := h.cfgReg.PluginConfig(p)
|
||||
pkeys, _ := ps.List("")
|
||||
if prefix == "" {
|
||||
for _, p := range h.settings.Plugins() {
|
||||
if p == "core" {
|
||||
continue
|
||||
}
|
||||
pkeys, _ := h.settings.ListPlugin(p, "")
|
||||
for _, k := range pkeys {
|
||||
v, _ := ps.Get(k)
|
||||
v, _ := h.settings.GetPlugin(p, k)
|
||||
fullKey := "plugin." + p + "." + k
|
||||
values[fullKey] = v
|
||||
if def := h.cfgReg.GetDef(fullKey); def != nil {
|
||||
meta[fullKey] = def
|
||||
}
|
||||
}
|
||||
for _, def := range h.settings.DefsPlugin(p, "") {
|
||||
fullKey := "plugin." + p + "." + def.Key
|
||||
meta[fullKey] = def
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
plugins := []string{"core"}
|
||||
pm := h.pluginReg.PluginMetas()
|
||||
if h.pluginReg != nil {
|
||||
for _, p := range h.pluginReg.List() {
|
||||
for _, p := range h.settings.Plugins() {
|
||||
if p != "core" {
|
||||
plugins = append(plugins, "plugin."+p)
|
||||
}
|
||||
}
|
||||
var pm map[string]sdk.PluginMeta
|
||||
if h.pluginMgr != nil {
|
||||
pm = h.pluginMgr.PluginMetas()
|
||||
}
|
||||
var disabledPlugins []sdk.DisabledPluginInfo
|
||||
if h.pluginMgr != nil {
|
||||
disabledPlugins = h.pluginMgr.ListDisabledPlugins()
|
||||
@ -1134,20 +1167,21 @@ func (h *Handler) handleSettings(w http.ResponseWriter, r *http.Request) {
|
||||
if strings.HasPrefix(body.Key, "plugin.") {
|
||||
parts := strings.SplitN(body.Key, ".", 3)
|
||||
if len(parts) >= 3 {
|
||||
ps := h.cfgReg.PluginConfig(parts[1])
|
||||
if err := ps.Set(parts[2], body.Value); err != nil {
|
||||
if err := h.settings.SetPlugin(parts[1], parts[2], body.Value); err != nil {
|
||||
writeJSON(w, http.StatusInternalServerError, map[string]string{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
}
|
||||
} else {
|
||||
if err := h.cfgReg.Set(body.Key, body.Value); err != nil {
|
||||
if err := h.settings.SetCore(body.Key, body.Value); err != nil {
|
||||
writeJSON(w, http.StatusInternalServerError, map[string]string{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
}
|
||||
if strings.HasPrefix(body.Key, "core.llm.") && h.providerMgr != nil && h.lua != nil {
|
||||
h.reloadLLMProviders()
|
||||
if strings.HasPrefix(body.Key, "core.llm.") && h.llm != nil {
|
||||
if err := h.llm.ReloadFromConfig(); err != nil {
|
||||
log.Printf("[webui] failed to reload LLM providers: %v", err)
|
||||
}
|
||||
}
|
||||
writeJSON(w, http.StatusOK, map[string]string{"status": "ok"})
|
||||
default:
|
||||
@ -1155,36 +1189,13 @@ func (h *Handler) handleSettings(w http.ResponseWriter, r *http.Request) {
|
||||
}
|
||||
}
|
||||
|
||||
func (h *Handler) reloadLLMProviders() {
|
||||
cfg := h.cfgReg.ToConfig()
|
||||
h.providerMgr.Reset()
|
||||
for _, src := range cfg.LLM.Sources {
|
||||
key := src.APIKey
|
||||
if key == "" {
|
||||
key = h.baseAPIKey
|
||||
}
|
||||
provider := agentAPI.NewLuaAdaptedProvider(agentAPI.BaseConfig{
|
||||
Model: src.Model,
|
||||
BaseURL: src.BaseURL,
|
||||
APIKey: key,
|
||||
Temperature: cfg.LLM.Temperature,
|
||||
MaxTokens: cfg.LLM.MaxTokens,
|
||||
ContextWindow: src.ContextWindow,
|
||||
}, h.lua, src.Adapter)
|
||||
h.providerMgr.Register(src.Name, provider)
|
||||
}
|
||||
if cfg.LLM.Provider != "" {
|
||||
_ = h.providerMgr.SetDefault(cfg.LLM.Provider)
|
||||
}
|
||||
}
|
||||
|
||||
func (h *Handler) handleOpenAICompletions(w http.ResponseWriter, r *http.Request) {
|
||||
if r.Method != http.MethodPost {
|
||||
http.Error(w, "method not allowed", http.StatusMethodNotAllowed)
|
||||
return
|
||||
}
|
||||
|
||||
if h.iom == nil {
|
||||
if h.sdk == nil {
|
||||
writeJSON(w, http.StatusServiceUnavailable, map[string]string{"error": "IO manager not available"})
|
||||
return
|
||||
}
|
||||
@ -1211,9 +1222,9 @@ func (h *Handler) handleOpenAICompletions(w http.ResponseWriter, r *http.Request
|
||||
return
|
||||
}
|
||||
|
||||
response := h.iom.InjectTextSync("http", lastMsg.Content)
|
||||
response := h.sdk.InjectTextSync("http", "http", lastMsg.Content)
|
||||
if response == nil {
|
||||
writeJSON(w, http.StatusInternalServerError, map[string]string{"error": "no response from agent"})
|
||||
writeJSON(w, http.StatusServiceUnavailable, map[string]string{"error": "no response from agent"})
|
||||
return
|
||||
}
|
||||
|
||||
@ -1375,11 +1386,10 @@ func (h *Handler) handleTracker(w http.ResponseWriter, r *http.Request) {
|
||||
|
||||
func (h *Handler) pluginmgrAddr() string {
|
||||
addr := "127.0.0.1:9876"
|
||||
if h.cfgReg == nil {
|
||||
if h.settings == nil {
|
||||
return addr
|
||||
}
|
||||
ps := h.cfgReg.PluginConfig("pluginmgr")
|
||||
if v, err := ps.Get("http_addr"); err == nil {
|
||||
if v, err := h.settings.GetPlugin("pluginmgr", "http_addr"); err == nil {
|
||||
if s, ok := v.(string); ok && s != "" {
|
||||
addr = s
|
||||
}
|
||||
@ -1436,17 +1446,11 @@ func (h *Handler) handlePluginByID(w http.ResponseWriter, r *http.Request) {
|
||||
}
|
||||
|
||||
if path == "reload" && r.Method == http.MethodPost {
|
||||
if h.pluginReg == nil {
|
||||
if h.pluginMgr == nil {
|
||||
writeJSON(w, http.StatusInternalServerError, map[string]string{"error": "plugin registry not available"})
|
||||
return
|
||||
}
|
||||
dir := ""
|
||||
if h.cfgReg != nil {
|
||||
if v, _ := h.cfgReg.Get("core.plugin.dir"); v != nil {
|
||||
dir, _ = v.(string)
|
||||
}
|
||||
}
|
||||
if _, err := h.pluginReg.Reload(dir); err != nil {
|
||||
if _, err := h.pluginMgr.ReloadPlugins(); err != nil {
|
||||
writeJSON(w, http.StatusInternalServerError, map[string]string{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
|
||||
@ -17,12 +17,19 @@ import (
|
||||
"gitcode.com/JianFeeeee/HomeAgent/internal/events"
|
||||
"gitcode.com/JianFeeeee/HomeAgent/internal/knowledge"
|
||||
"gitcode.com/JianFeeeee/HomeAgent/internal/memory"
|
||||
"gitcode.com/JianFeeeee/HomeAgent/internal/plugin"
|
||||
sdk "gitcode.com/JianFeeeee/HomeAgent/internal/sdk"
|
||||
"gitcode.com/JianFeeeee/HomeAgent/internal/supervisor"
|
||||
"gitcode.com/JianFeeeee/HomeAgent/internal/tracker"
|
||||
"gitcode.com/JianFeeeee/HomeAgent/pkg/types"
|
||||
)
|
||||
|
||||
func testSDK(cfg sdk.SDKConfig) *sdk.PluginSDK {
|
||||
if cfg.EventBus == nil {
|
||||
cfg.EventBus = events.NewBus()
|
||||
}
|
||||
return sdk.New("webui", cfg)
|
||||
}
|
||||
|
||||
func newTestHandler(t *testing.T) (*Handler, *supervisor.Daemon) {
|
||||
t.Helper()
|
||||
cfg := &types.Config{
|
||||
@ -34,7 +41,11 @@ func newTestHandler(t *testing.T) (*Handler, *supervisor.Daemon) {
|
||||
sup := supervisor.New(cfg)
|
||||
sup.Start()
|
||||
|
||||
return NewHandler(sup, nil, nil, nil, cfg, nil, nil, nil, nil, nil, nil, events.NewBus(), nil, nil, ""), sup
|
||||
s := testSDK(sdk.SDKConfig{
|
||||
Supervisor: supervisor.NewSDKAdapter(sup),
|
||||
Config: sdk.NewConfig(cfg),
|
||||
})
|
||||
return NewHandler(s), sup
|
||||
}
|
||||
|
||||
func TestAuthMiddleware(t *testing.T) {
|
||||
@ -44,9 +55,21 @@ func TestAuthMiddleware(t *testing.T) {
|
||||
cfgReg.PluginConfig("webui").Set("password", "secret-pass")
|
||||
cfgReg.PluginConfig("webui").Set("session_ttl_hours", "24")
|
||||
|
||||
h, sup := newTestHandler(t)
|
||||
sup := supervisor.New(&types.Config{
|
||||
Daemon: types.DaemonConfig{
|
||||
CheckInterval: time.Minute,
|
||||
HeartbeatInterval: 30 * time.Second,
|
||||
},
|
||||
})
|
||||
sup.Start()
|
||||
defer sup.Shutdown()
|
||||
h.cfgReg = cfgReg
|
||||
|
||||
s := testSDK(sdk.SDKConfig{
|
||||
Supervisor: supervisor.NewSDKAdapter(sup),
|
||||
Settings: sdk.NewSettings("webui", cfgReg),
|
||||
Config: sdk.NewConfig(&types.Config{}),
|
||||
})
|
||||
h := NewHandler(s)
|
||||
mux := http.NewServeMux()
|
||||
h.RegisterRoutes(mux)
|
||||
|
||||
@ -194,7 +217,12 @@ func TestHandleKnowledgeSearch(t *testing.T) {
|
||||
sup.Start()
|
||||
defer sup.Shutdown()
|
||||
|
||||
h := NewHandler(sup, nil, nil, nil, cfg, nil, nil, ks, nil, nil, nil, events.NewBus(), nil, nil, "")
|
||||
s := testSDK(sdk.SDKConfig{
|
||||
Supervisor: supervisor.NewSDKAdapter(sup),
|
||||
Knowledge: sdk.NewKnowledge(ks),
|
||||
Config: sdk.NewConfig(cfg),
|
||||
})
|
||||
h := NewHandler(s)
|
||||
|
||||
req := httptest.NewRequest(http.MethodGet, "/api/v1/knowledge?q=test", nil)
|
||||
w := httptest.NewRecorder()
|
||||
@ -226,7 +254,12 @@ func TestHandleKnowledgeCreate(t *testing.T) {
|
||||
sup.Start()
|
||||
defer sup.Shutdown()
|
||||
|
||||
h := NewHandler(sup, nil, nil, nil, cfg, nil, nil, ks, nil, nil, nil, events.NewBus(), nil, nil, "")
|
||||
s := testSDK(sdk.SDKConfig{
|
||||
Supervisor: supervisor.NewSDKAdapter(sup),
|
||||
Knowledge: sdk.NewKnowledge(ks),
|
||||
Config: sdk.NewConfig(cfg),
|
||||
})
|
||||
h := NewHandler(s)
|
||||
|
||||
body := `{"name":"new_doc","content":"fresh content"}`
|
||||
req := httptest.NewRequest(http.MethodPost, "/api/v1/knowledge", strings.NewReader(body))
|
||||
@ -291,7 +324,12 @@ func TestHandleTrackerStats(t *testing.T) {
|
||||
sup.Start()
|
||||
defer sup.Shutdown()
|
||||
|
||||
h := NewHandler(sup, nil, nil, nil, cfg, nil, nil, nil, tr, nil, nil, events.NewBus(), nil, nil, "")
|
||||
s := testSDK(sdk.SDKConfig{
|
||||
Supervisor: supervisor.NewSDKAdapter(sup),
|
||||
Tracker: tr,
|
||||
Config: sdk.NewConfig(cfg),
|
||||
})
|
||||
h := NewHandler(s)
|
||||
|
||||
req := httptest.NewRequest(http.MethodGet, "/api/v1/tracker", nil)
|
||||
w := httptest.NewRecorder()
|
||||
@ -305,8 +343,6 @@ func TestHandleTrackerStats(t *testing.T) {
|
||||
func TestHandleOpenAICompletionsNoMessages(t *testing.T) {
|
||||
h, sup := newTestHandler(t)
|
||||
defer sup.Shutdown()
|
||||
// 给 handler 一个 IOManager,才能通过 nil 检查到达消息校验
|
||||
h.iom = agentIO.NewIOManager()
|
||||
|
||||
body := `{"model":"test"}`
|
||||
req := httptest.NewRequest(http.MethodPost, "/v1/chat/completions", strings.NewReader(body))
|
||||
@ -322,8 +358,6 @@ func TestHandleOpenAICompletionsNoMessages(t *testing.T) {
|
||||
func TestHandleOpenAICompletionsLastMsgNotUser(t *testing.T) {
|
||||
h, sup := newTestHandler(t)
|
||||
defer sup.Shutdown()
|
||||
// 给 handler 一个 IOManager,才能通过 nil 检查到达消息校验
|
||||
h.iom = agentIO.NewIOManager()
|
||||
|
||||
body := `{"messages":[{"role":"assistant","content":"hi"}]}`
|
||||
req := httptest.NewRequest(http.MethodPost, "/v1/chat/completions", strings.NewReader(body))
|
||||
@ -392,9 +426,27 @@ func TestHandleConfigGet(t *testing.T) {
|
||||
}
|
||||
|
||||
func TestRegisterRoutes(t *testing.T) {
|
||||
h, sup := newTestHandler(t)
|
||||
cfgReg := internalConfig.NewConfigRegistry("")
|
||||
cfgReg.PluginConfig("webui").Set("api_key", "test-api-key")
|
||||
cfgReg.PluginConfig("webui").Set("username", "admin")
|
||||
cfgReg.PluginConfig("webui").Set("password", "secret-pass")
|
||||
|
||||
sup := supervisor.New(&types.Config{
|
||||
Daemon: types.DaemonConfig{
|
||||
CheckInterval: time.Minute,
|
||||
HeartbeatInterval: 30 * time.Second,
|
||||
},
|
||||
})
|
||||
sup.Start()
|
||||
defer sup.Shutdown()
|
||||
|
||||
s := testSDK(sdk.SDKConfig{
|
||||
Supervisor: supervisor.NewSDKAdapter(sup),
|
||||
Settings: sdk.NewSettings("webui", cfgReg),
|
||||
Config: sdk.NewConfig(&types.Config{}),
|
||||
})
|
||||
h := NewHandler(s)
|
||||
|
||||
mux := http.NewServeMux()
|
||||
h.RegisterRoutes(mux)
|
||||
|
||||
@ -407,7 +459,7 @@ func TestRegisterRoutes(t *testing.T) {
|
||||
{"/api/v1/agents", http.MethodGet, http.StatusOK},
|
||||
{"/api/v1/config", http.MethodGet, http.StatusOK},
|
||||
{"/api/v1/network", http.MethodGet, http.StatusOK},
|
||||
{"/", http.MethodGet, http.StatusOK},
|
||||
{"/", http.MethodGet, http.StatusFound},
|
||||
{"/api/v1/memory", http.MethodGet, http.StatusServiceUnavailable},
|
||||
{"/api/v1/knowledge", http.MethodGet, http.StatusServiceUnavailable},
|
||||
{"/api/v1/tracker", http.MethodGet, http.StatusServiceUnavailable},
|
||||
@ -416,6 +468,9 @@ func TestRegisterRoutes(t *testing.T) {
|
||||
|
||||
for _, tt := range tests {
|
||||
req := httptest.NewRequest(tt.method, tt.path, nil)
|
||||
if strings.HasPrefix(tt.path, "/api/v1/") {
|
||||
req.Header.Set("X-API-Key", "test-api-key")
|
||||
}
|
||||
w := httptest.NewRecorder()
|
||||
mux.ServeHTTP(w, req)
|
||||
|
||||
@ -470,8 +525,12 @@ func TestSettingsAPIFlow(t *testing.T) {
|
||||
sup.Start()
|
||||
defer sup.Shutdown()
|
||||
|
||||
pluginReg := plugin.NewRegistry()
|
||||
h := NewHandler(sup, nil, nil, nil, &types.Config{}, nil, nil, nil, nil, cfgReg, pluginReg, events.NewBus(), nil, nil, "")
|
||||
s := testSDK(sdk.SDKConfig{
|
||||
Supervisor: supervisor.NewSDKAdapter(sup),
|
||||
Settings: sdk.NewSettings("webui", cfgReg),
|
||||
Config: sdk.NewConfig(&types.Config{}),
|
||||
})
|
||||
h := NewHandler(s)
|
||||
|
||||
t.Run("GET_settings_lists_keys_and_plugins", func(t *testing.T) {
|
||||
req := httptest.NewRequest(http.MethodGet, "/api/v1/settings", nil)
|
||||
@ -578,7 +637,11 @@ func TestSettingsAPIFlow(t *testing.T) {
|
||||
})
|
||||
|
||||
t.Run("settings_not_available_without_registry", func(t *testing.T) {
|
||||
h2 := NewHandler(sup, nil, nil, nil, &types.Config{}, nil, nil, nil, nil, nil, nil, events.NewBus(), nil, nil, "")
|
||||
s2 := testSDK(sdk.SDKConfig{
|
||||
Supervisor: supervisor.NewSDKAdapter(sup),
|
||||
Config: sdk.NewConfig(&types.Config{}),
|
||||
})
|
||||
h2 := NewHandler(s2)
|
||||
req := httptest.NewRequest(http.MethodGet, "/api/v1/settings", nil)
|
||||
w := httptest.NewRecorder()
|
||||
h2.handleSettings(w, req)
|
||||
@ -603,8 +666,12 @@ func TestSettingsWithPluginRegistry(t *testing.T) {
|
||||
sup.Start()
|
||||
defer sup.Shutdown()
|
||||
|
||||
pluginReg := plugin.NewRegistry()
|
||||
h := NewHandler(sup, nil, nil, nil, &types.Config{}, nil, nil, nil, nil, cfgReg, pluginReg, events.NewBus(), nil, nil, "")
|
||||
s := testSDK(sdk.SDKConfig{
|
||||
Supervisor: supervisor.NewSDKAdapter(sup),
|
||||
Settings: sdk.NewSettings("webui", cfgReg),
|
||||
Config: sdk.NewConfig(&types.Config{}),
|
||||
})
|
||||
h := NewHandler(s)
|
||||
|
||||
req := httptest.NewRequest(http.MethodGet, "/api/v1/settings", nil)
|
||||
w := httptest.NewRecorder()
|
||||
@ -633,7 +700,7 @@ type echoProvider struct{ name string }
|
||||
func (p *echoProvider) Name() string { return p.name }
|
||||
func (p *echoProvider) MaxContextTokens() int { return 8192 }
|
||||
func (p *echoProvider) Chat(ctx context.Context, req *agentAPI.CompletionRequest) (*agentAPI.CompletionResponse, error) {
|
||||
content := "echo: " + req.Messages[len(req.Messages)-1].Content
|
||||
content := "echo: " + lastUserContent(req.Messages)
|
||||
return &agentAPI.CompletionResponse{Content: content, FinishReason: "stop"}, nil
|
||||
}
|
||||
func (p *echoProvider) ChatStream(ctx context.Context, req *agentAPI.CompletionRequest) (<-chan agentAPI.StreamChunk, error) {
|
||||
@ -642,6 +709,15 @@ func (p *echoProvider) ChatStream(ctx context.Context, req *agentAPI.CompletionR
|
||||
return ch, nil
|
||||
}
|
||||
|
||||
func lastUserContent(msgs []agentAPI.Message) string {
|
||||
for i := len(msgs) - 1; i >= 0; i-- {
|
||||
if msgs[i].Role == "user" {
|
||||
return msgs[i].Content
|
||||
}
|
||||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
func init() {
|
||||
// 避免测试时自动输出
|
||||
}
|
||||
@ -656,19 +732,22 @@ func TestHandleCompletionsEndToEnd(t *testing.T) {
|
||||
}
|
||||
defer memDB.Close()
|
||||
|
||||
pm := agentAPI.NewProviderManager()
|
||||
pm.Register("echo", &echoProvider{name: "echo"})
|
||||
|
||||
agent := agentCore.New(agentCore.AgentConfig{
|
||||
ID: "test",
|
||||
SystemPrompt: "你是测试助手",
|
||||
Provider: &echoProvider{name: "echo"},
|
||||
IO: iom,
|
||||
Memory: memDB,
|
||||
Indexer: nil,
|
||||
ID: "test",
|
||||
SystemPrompt: "你是测试助手",
|
||||
Provider: &echoProvider{name: "echo"},
|
||||
ProviderManager: pm,
|
||||
IO: iom,
|
||||
Memory: memDB,
|
||||
Indexer: nil,
|
||||
ContextSavePath: "",
|
||||
})
|
||||
agent.Start()
|
||||
defer agent.Stop()
|
||||
|
||||
// Handler 需要 iom
|
||||
sup := supervisor.New(&types.Config{
|
||||
Daemon: types.DaemonConfig{
|
||||
CheckInterval: time.Minute,
|
||||
@ -678,7 +757,12 @@ func TestHandleCompletionsEndToEnd(t *testing.T) {
|
||||
sup.Start()
|
||||
defer sup.Shutdown()
|
||||
|
||||
h := NewHandler(sup, nil, nil, nil, &types.Config{}, iom, nil, nil, nil, nil, nil, events.NewBus(), nil, nil, "")
|
||||
s := testSDK(sdk.SDKConfig{
|
||||
Supervisor: supervisor.NewSDKAdapter(sup),
|
||||
IOManager: iom,
|
||||
Config: sdk.NewConfig(&types.Config{}),
|
||||
})
|
||||
h := NewHandler(s)
|
||||
|
||||
t.Run("POST_chat_completions_returns_echo", func(t *testing.T) {
|
||||
body := `{"model":"test","messages":[{"role":"user","content":"你好"}]}`
|
||||
@ -711,7 +795,10 @@ func TestHandleCompletionsEndToEnd(t *testing.T) {
|
||||
})
|
||||
|
||||
t.Run("POST_chat_completions_no_iom_returns_503", func(t *testing.T) {
|
||||
h2 := NewHandler(sup, nil, nil, nil, &types.Config{}, nil, nil, nil, nil, nil, nil, events.NewBus(), nil, nil, "")
|
||||
s2 := testSDK(sdk.SDKConfig{
|
||||
Supervisor: supervisor.NewSDKAdapter(sup),
|
||||
})
|
||||
h2 := NewHandler(s2)
|
||||
body := `{"messages":[{"role":"user","content":"hi"}]}`
|
||||
req := httptest.NewRequest(http.MethodPost, "/v1/chat/completions", strings.NewReader(body))
|
||||
req.Header.Set("Content-Type", "application/json")
|
||||
|
||||
@ -7,117 +7,28 @@ import (
|
||||
"log"
|
||||
"net/http"
|
||||
|
||||
agentAPI "gitcode.com/JianFeeeee/HomeAgent/internal/agent/api"
|
||||
agentCore "gitcode.com/JianFeeeee/HomeAgent/internal/agent/core"
|
||||
agentIO "gitcode.com/JianFeeeee/HomeAgent/internal/agent/io"
|
||||
internalConfig "gitcode.com/JianFeeeee/HomeAgent/internal/config"
|
||||
"gitcode.com/JianFeeeee/HomeAgent/internal/events"
|
||||
"gitcode.com/JianFeeeee/HomeAgent/internal/knowledge"
|
||||
luaVM "gitcode.com/JianFeeeee/HomeAgent/internal/lua"
|
||||
"gitcode.com/JianFeeeee/HomeAgent/internal/memory"
|
||||
"gitcode.com/JianFeeeee/HomeAgent/internal/memory/text"
|
||||
"gitcode.com/JianFeeeee/HomeAgent/internal/plugin"
|
||||
"gitcode.com/JianFeeeee/HomeAgent/internal/skill"
|
||||
sdk "gitcode.com/JianFeeeee/HomeAgent/internal/sdk"
|
||||
"gitcode.com/JianFeeeee/HomeAgent/internal/supervisor"
|
||||
"gitcode.com/JianFeeeee/HomeAgent/internal/tracker"
|
||||
"gitcode.com/JianFeeeee/HomeAgent/pkg/types"
|
||||
)
|
||||
|
||||
// 包级依赖注入 — 由 main.go 在 Load() 前调用 Configure() 设置。
|
||||
var (
|
||||
webuiAddr string
|
||||
webuiSup *supervisor.Daemon
|
||||
webuiMem *memory.GraphDB
|
||||
webuiSK *skill.Manager
|
||||
webuiLua *luaVM.VM
|
||||
webuiCfg *types.Config
|
||||
webuiIOM *agentIO.IOManager
|
||||
webuiTM *text.Memory
|
||||
webuiKS *knowledge.Store
|
||||
webuiTR *tracker.Tracker
|
||||
webuiCR *internalConfig.ConfigRegistry
|
||||
webuiPR *plugin.Registry
|
||||
webuiEvBus *events.Bus
|
||||
webuiStatusProvider agentCore.StatusProvider
|
||||
webuiProviderMgr *agentAPI.ProviderManager
|
||||
webuiBaseAPIKey string
|
||||
)
|
||||
|
||||
// Configure 注入 WebUI 插件需要的内核依赖。必须在 Load() 之前调用。
|
||||
func Configure(addr string,
|
||||
sup *supervisor.Daemon, mem *memory.GraphDB, sk *skill.Manager,
|
||||
lua *luaVM.VM, cfg *types.Config, iom *agentIO.IOManager,
|
||||
tm *text.Memory, ks *knowledge.Store, tr *tracker.Tracker,
|
||||
cr *internalConfig.ConfigRegistry, pr *plugin.Registry, evBus *events.Bus,
|
||||
sp agentCore.StatusProvider, pm *agentAPI.ProviderManager, baseKey string,
|
||||
) {
|
||||
webuiAddr = addr
|
||||
webuiSup, webuiMem, webuiSK, webuiLua = sup, mem, sk, lua
|
||||
webuiCfg, webuiIOM, webuiTM, webuiKS = cfg, iom, tm, ks
|
||||
webuiTR, webuiCR, webuiPR, webuiEvBus = tr, cr, pr, evBus
|
||||
webuiStatusProvider = sp
|
||||
webuiProviderMgr = pm
|
||||
webuiBaseAPIKey = baseKey
|
||||
}
|
||||
|
||||
func init() {
|
||||
plugin.RegisterPluginMeta("webui", "Web 控制台", "WebUI")
|
||||
plugin.RegisterFactory("webui", func(name string, config map[string]interface{}) (sdk.Plugin, error) {
|
||||
if webuiSup == nil {
|
||||
return nil, nil // 未 Configure 则跳过(不给日志警告)
|
||||
}
|
||||
addr := webuiAddr
|
||||
if a, ok := config["addr"].(string); ok {
|
||||
addr = a
|
||||
}
|
||||
return New(name, addr,
|
||||
webuiSup, webuiMem, webuiSK, webuiLua,
|
||||
webuiCfg, webuiIOM, webuiTM, webuiKS,
|
||||
webuiTR, webuiCR, webuiPR, webuiEvBus, webuiStatusProvider,
|
||||
webuiProviderMgr, webuiBaseAPIKey,
|
||||
), nil
|
||||
return New(name), nil
|
||||
})
|
||||
}
|
||||
|
||||
type Plugin struct {
|
||||
name string
|
||||
addr string
|
||||
handler *Handler
|
||||
server *http.Server
|
||||
mux *http.ServeMux
|
||||
|
||||
sup *supervisor.Daemon
|
||||
mem *memory.GraphDB
|
||||
sk *skill.Manager
|
||||
lua *luaVM.VM
|
||||
cfg *types.Config
|
||||
iom *agentIO.IOManager
|
||||
tm *text.Memory
|
||||
ks *knowledge.Store
|
||||
tr *tracker.Tracker
|
||||
cr *internalConfig.ConfigRegistry
|
||||
pr *plugin.Registry
|
||||
evBus *events.Bus
|
||||
statusProvider agentCore.StatusProvider
|
||||
providerMgr *agentAPI.ProviderManager
|
||||
baseAPIKey string
|
||||
}
|
||||
|
||||
func New(name, addr string,
|
||||
sup *supervisor.Daemon, mem *memory.GraphDB, sk *skill.Manager,
|
||||
lua *luaVM.VM, cfg *types.Config, iom *agentIO.IOManager,
|
||||
tm *text.Memory, ks *knowledge.Store, tr *tracker.Tracker,
|
||||
cr *internalConfig.ConfigRegistry, pr *plugin.Registry, evBus *events.Bus,
|
||||
sp agentCore.StatusProvider, pm *agentAPI.ProviderManager, baseKey string,
|
||||
) *Plugin {
|
||||
func New(name string) *Plugin {
|
||||
return &Plugin{
|
||||
name: name,
|
||||
addr: addr,
|
||||
mux: http.NewServeMux(),
|
||||
sup: sup, mem: mem, sk: sk, lua: lua, cfg: cfg,
|
||||
iom: iom, tm: tm, ks: ks, tr: tr, cr: cr, pr: pr, evBus: evBus,
|
||||
statusProvider: sp, providerMgr: pm, baseAPIKey: baseKey,
|
||||
}
|
||||
}
|
||||
|
||||
@ -157,11 +68,18 @@ func (p *Plugin) Name() string { return p.name }
|
||||
func (p *Plugin) Start(s *sdk.PluginSDK) error {
|
||||
s.SetAutoRestart(true)
|
||||
|
||||
addr := ":8080"
|
||||
if v, _ := s.Settings().Get("addr"); v != nil {
|
||||
if s2, ok := v.(string); ok && s2 != "" {
|
||||
addr = s2
|
||||
}
|
||||
}
|
||||
|
||||
s.RegisterOutputChannel("webui", 1, "Web 控制台", sdk.ChannelDef{}, func(args map[string]interface{}) (interface{}, error) {
|
||||
payload, _ := args["payload"].(string)
|
||||
if payload != "" {
|
||||
p.evBus.Publish(&events.Event{
|
||||
Type: events.EventAgentOutput,
|
||||
s.Publish(&sdk.Event{
|
||||
Type: sdk.EventAgentOutput,
|
||||
Payload: map[string]interface{}{
|
||||
"content": payload,
|
||||
"channel": "webui",
|
||||
@ -171,6 +89,7 @@ func (p *Plugin) Start(s *sdk.PluginSDK) error {
|
||||
return map[string]interface{}{"status": "ok"}, nil
|
||||
})
|
||||
|
||||
s.Settings().RegisterDef(sdk.ConfigDef{Key: "addr", Default: ":8080", Type: "string", DisplayName: "监听地址", Description: "Web 控制台监听地址", Category: "webui"})
|
||||
s.Settings().RegisterDef(sdk.ConfigDef{Key: "api_key", Default: "", Type: "password", DisplayName: "API 密钥", Description: "访问 API 时需要的密钥", Category: "webui"})
|
||||
s.Settings().RegisterDef(sdk.ConfigDef{Key: "username", Default: "admin", Type: "string", DisplayName: "登录用户名", Description: "Web 控制台登录用户名", Category: "webui"})
|
||||
s.Settings().RegisterDef(sdk.ConfigDef{Key: "password", Default: "", Type: "password", DisplayName: "Web 控制台登录密码", Description: "Web 控制台登录密码", Category: "webui"})
|
||||
@ -178,7 +97,7 @@ func (p *Plugin) Start(s *sdk.PluginSDK) error {
|
||||
p.ensureAuthBootstrap(s)
|
||||
|
||||
s.RegisterStage(sdk.StagePreAction, func(ctx *sdk.StageContext) error {
|
||||
p.evBus.Publish(&events.Event{Type: events.EventStage, Payload: map[string]interface{}{"phase": "pre_action", "message": "thinking"}})
|
||||
s.Publish(&sdk.Event{Type: sdk.EventStage, Payload: map[string]interface{}{"phase": "pre_action", "message": "thinking"}})
|
||||
return nil
|
||||
})
|
||||
s.RegisterStage(sdk.StageBeforeToolcall, func(ctx *sdk.StageContext) error {
|
||||
@ -186,22 +105,20 @@ func (p *Plugin) Start(s *sdk.PluginSDK) error {
|
||||
if len(ctx.ToolCalls) > 0 {
|
||||
tool = ctx.ToolCalls[0].Name
|
||||
}
|
||||
p.evBus.Publish(&events.Event{Type: events.EventStage, Payload: map[string]interface{}{"phase": "before_toolcall", "tool": tool, "message": "tool:" + tool}})
|
||||
s.Publish(&sdk.Event{Type: sdk.EventStage, Payload: map[string]interface{}{"phase": "before_toolcall", "tool": tool, "message": "tool:" + tool}})
|
||||
return nil
|
||||
})
|
||||
s.RegisterStage(sdk.StageBeforeOutput, func(ctx *sdk.StageContext) error {
|
||||
p.evBus.Publish(&events.Event{Type: events.EventStage, Payload: map[string]interface{}{"phase": "before_output", "message": "output"}})
|
||||
s.Publish(&sdk.Event{Type: sdk.EventStage, Payload: map[string]interface{}{"phase": "before_output", "message": "output"}})
|
||||
return nil
|
||||
})
|
||||
|
||||
h := NewHandler(p.sup, p.mem, p.sk, p.lua, p.cfg, p.iom, p.tm, p.ks, p.tr, p.cr, p.pr, p.evBus, p.statusProvider, p.providerMgr, p.baseAPIKey)
|
||||
h.SetPluginMgr(s.PluginMgr())
|
||||
p.handler = h
|
||||
h.RegisterRoutes(p.mux)
|
||||
p.handler = NewHandler(s)
|
||||
p.handler.RegisterRoutes(p.mux)
|
||||
|
||||
p.server = &http.Server{Addr: p.addr, Handler: p.mux}
|
||||
p.server = &http.Server{Addr: addr, Handler: p.mux}
|
||||
go func() {
|
||||
log.Printf("[webui] HTTP server listening on %s", p.addr)
|
||||
log.Printf("[webui] HTTP server listening on %s", addr)
|
||||
if err := p.server.ListenAndServe(); err != nil && err != http.ErrServerClosed {
|
||||
log.Printf("[webui] server error: %v", err)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user