diff --git a/internal/plugin/proc/corehandler.go b/internal/plugin/proc/corehandler.go index 7402ddd..4fb3fd6 100644 --- a/internal/plugin/proc/corehandler.go +++ b/internal/plugin/proc/corehandler.go @@ -121,552 +121,73 @@ type CoreSDK interface { // 权限梯度在此强制(§3.8):manifest 未声明的能力组被**明确拒绝**。 // 不静默忽略:C ABI 时代 case 23/24 返回成功但永远收不到事件 // (§1.3 的「给不了」而非「不给」),插件作者无从得知。 +// Handle 分派一次插件 → 内核的调用。 +// +// 权限梯度在此强制(§3.8):manifest 未声明的能力组被**明确拒绝**。 +// 不静默忽略:C ABI 时代 case 23/24 返回成功但永远收不到事件 +// (§1.3 的「给不了」而非「不给」),插件作者无从得知。 +// +// 体量上这里只做「method → 分部函数」一跳:原 51 个 case 的 550 行大 switch +// 已按协议面拆进同包 corehandler_register / _inject / _memory / +// _settings / _runtime。 func (h *coreHandler) Handle(method string, params json.RawMessage) (interface{}, error) { if ok, cap := h.caps.allows(method); !ok { return nil, errCapabilityDenied(h.name, method, cap) } switch method { + case MethodToolRegister, MethodStageRegister, MethodOutputRegister, MethodAPIRegister, + MethodInputRegister: + return h.handleRegister(method, params) - // ---- 注册面(原 case 1/2/3/4/46)---- - case MethodToolRegister: - return h.toolRegister(params) - case MethodStageRegister: - return h.stageRegister(params) - case MethodOutputRegister: - return h.outputRegister(params) - case MethodAPIRegister: - var p struct { - Name string `json:"name"` - } - if err := unmarshal(params, &p); err != nil { - return nil, err - } - return nil, h.sdk.RegisterPluginAPI(p.Name) - case MethodInputRegister: - var p struct { - Name string `json:"name"` - Def pubsdk.ChannelDef `json:"def"` - HasCleaner bool `json:"has_cleaner"` - } - if err := unmarshal(params, &p); err != nil { - return nil, err - } - if p.Name == "" { - return nil, fmt.Errorf("input.register: 缺少 name") - } - cleaner, err := h.cleanerProxy(CleanerScopeInput, p.Name, p.HasCleaner) - if err != nil { - return nil, fmt.Errorf("input.register: %w", err) - } - if err := validateContextPolicy("input.register", p.Def.ContextPolicy); err != nil { - return nil, err - } - // 整体传 p.Def(只是把函数型的 Cleaner 换成代理),不要手写字段白名单: - // 白名单会让新增字段静默丢失。 - def := p.Def - def.Cleaner = cleaner - return nil, h.sdk.RegisterInputChannel(p.Name, def) + case MethodIOInjectText, MethodIOInjectInterrupt, MethodIOInjectTextNoMem, + MethodIOInjectSync, MethodIOInjectMedia, MethodIOInjectMediaSync, + MethodIOInjectInterruptMedia, MethodIOSetToolBlocks: + return h.handleInject(method, params) - // ---- IO 注入(原 case 5/6/7/47)---- - // - // 注入标志位(no_memory / context_policy)由插件在调用点声明,默认 - // 记入记忆 + 不裁剪。策略值在入口校验:静默降级成 none 会让调用方 - // 以为自己声明的裁剪在生效。 - case MethodIOInjectText: - var p injectParams - if err := unmarshal(params, &p); err != nil { - return nil, err - } - if err := validateContextPolicy("io.injectText", p.ContextPolicy); err != nil { - return nil, err - } - h.sdk.InjectTextOpts(p.Source, p.Channel, h.resolveText(p), pubSdkInjectOpts(p.NoMemory, p.ContextPolicy, p.CleanerName, p.Priority)) - return nil, nil - case MethodIOInjectInterrupt: - var p injectParams - if err := unmarshal(params, &p); err != nil { - return nil, err - } - if err := validateContextPolicy("io.injectInterrupt", p.ContextPolicy); err != nil { - return nil, err - } - h.sdk.InjectInterruptTextOpts(p.Source, p.Channel, h.resolveText(p), pubSdkInjectOpts(p.NoMemory, p.ContextPolicy, p.CleanerName, p.Priority)) - return nil, nil - case MethodIOInjectTextNoMem: - var p injectParams - if err := unmarshal(params, &p); err != nil { - return nil, err - } - if err := validateContextPolicy("io.injectTextNoMem", p.ContextPolicy); err != nil { - return nil, err - } - // 旧 RPC 语义就是「不进记忆」,显式标志位只可能再叠上 context_policy。 - h.sdk.InjectTextOpts(p.Source, p.Channel, h.resolveText(p), pubSdkInjectOpts(true, p.ContextPolicy, p.CleanerName, p.Priority)) - return nil, nil - case MethodIOInjectSync: - var p injectParams - if err := unmarshal(params, &p); err != nil { - return nil, err - } - if err := validateContextPolicy("io.injectInputSync", p.ContextPolicy); err != nil { - return nil, err - } - reply := h.sdk.InjectInputSyncOpts(p.Source, p.Channel, h.resolveText(p), pubSdkInjectOpts(p.NoMemory, p.ContextPolicy, p.CleanerName, p.Priority)) - return map[string]interface{}{"reply": reply}, nil + case MethodMemoryRecall, MethodMemoryCommit, MethodMemoryIntrospect, MethodMemoryMerge, + MethodMemoryPurge: + return h.handleGraphMemory(method, params) - case MethodIOInjectMedia: - var p injectMediaParams - if err := unmarshal(params, &p); err != nil { - return nil, err - } - if err := validateContextPolicy("io.injectMedia", p.ContextPolicy); err != nil { - return nil, err - } - blocks, err := h.resolveBlocks(p) - if err != nil { - return nil, err - } - h.sdk.InjectInputMediaOpts(p.Source, p.Channel, p.Text, blocks, pubSdkInjectOpts(p.NoMemory, p.ContextPolicy, p.CleanerName, p.Priority)) - return nil, nil + case MethodDocQuery, MethodDocInsert, MethodDocInsertMedia, MethodDocRemove, + MethodDocStats: + return h.handleDocMemory(method, params) - case MethodIOInjectMediaSync: - var p injectMediaParams - if err := unmarshal(params, &p); err != nil { - return nil, err - } - if err := validateContextPolicy("io.injectMediaSync", p.ContextPolicy); err != nil { - return nil, err - } - blocks, err := h.resolveBlocks(p) - if err != nil { - return nil, err - } - reply := h.sdk.InjectInputMediaSyncOpts(p.Source, p.Channel, p.Text, blocks, pubSdkInjectOpts(p.NoMemory, p.ContextPolicy, p.CleanerName, p.Priority)) - return map[string]interface{}{"reply": reply}, nil + case MethodKnowledgeSearch, MethodKnowledgeAdd, MethodKnowledgeList: + return h.handleKnowledge(method, params) - case MethodIOInjectInterruptMedia: - var p injectMediaParams - if err := unmarshal(params, &p); err != nil { - return nil, err - } - if err := validateContextPolicy("io.injectInterruptMedia", p.ContextPolicy); err != nil { - return nil, err - } - blocks, err := h.resolveBlocks(p) - if err != nil { - return nil, err - } - h.sdk.InjectInterruptMediaOpts(p.Source, p.Channel, p.Text, blocks, pubSdkInjectOpts(p.NoMemory, p.ContextPolicy, p.CleanerName, p.Priority)) - return nil, nil - - // ---- 生命周期(原 case 8)---- - case MethodLifecycleAutoRestart: - var p struct { - Enabled bool `json:"enabled"` - } - if err := unmarshal(params, &p); err != nil { - return nil, err - } - h.sdk.SetAutoRestart(p.Enabled) - return nil, nil - - // ---- 图记忆(原 case 9/10/11/12/13)---- - case MethodMemoryRecall: - mem := h.sdk.Memory() - if mem == nil { - return nil, errUnavailable("memory") - } - var p struct { - Query []string `json:"query"` - Depth int `json:"depth"` - } - if err := unmarshal(params, &p); err != nil { - return nil, err - } - entities, relations, err := mem.Recall(p.Query, p.Depth) - if err != nil { - return nil, err - } - if entities == nil { - entities = []pubsdk.Entity{} - } - if relations == nil { - relations = []pubsdk.Relation{} - } - return map[string]interface{}{"entities": entities, "relations": relations}, nil - - case MethodMemoryCommit: - mem := h.sdk.Memory() - if mem == nil { - return nil, errUnavailable("memory") - } - var p struct { - Triples []pubsdk.Triple `json:"triples"` - } - if err := unmarshal(params, &p); err != nil { - return nil, err - } - return nil, mem.Commit(p.Triples) - - case MethodMemoryIntrospect: - mem := h.sdk.Memory() - if mem == nil { - return nil, errUnavailable("memory") - } - return mem.Introspect() - - case MethodMemoryMerge: - mem := h.sdk.Memory() - if mem == nil { - return nil, errUnavailable("memory") - } - var p struct { - Source string `json:"source"` - Target string `json:"target"` - } - if err := unmarshal(params, &p); err != nil { - return nil, err - } - n, err := mem.MergeEntities(p.Source, p.Target) - if err != nil { - return nil, err - } - return map[string]interface{}{"merged": n}, nil - - case MethodMemoryPurge: - mem := h.sdk.Memory() - if mem == nil { - return nil, errUnavailable("memory") - } - var p struct { - Criteria map[string]string `json:"criteria"` - Mode string `json:"mode"` - } - if err := unmarshal(params, &p); err != nil { - return nil, err - } - if p.Mode == "" { - p.Mode = "soft" - } - n, err := mem.Purge(p.Criteria, p.Mode) - if err != nil { - return nil, err - } - return map[string]interface{}{"purged": n}, nil - - // ---- 文档记忆(原 case 14/32/33/34)---- - case MethodDocQuery: - dm := h.sdk.DocMemory() - if dm == nil { - return nil, errUnavailable("doc memory") - } - var p struct { - Text string `json:"text"` - TopK int `json:"top_k"` - } - if err := unmarshal(params, &p); err != nil { - return nil, err - } - docs := dm.Query(p.Text, p.TopK) - if docs == nil { - docs = []*pubsdk.Doc{} - } - return map[string]interface{}{"docs": docs}, nil - - case MethodDocInsert: - dm := h.sdk.DocMemory() - if dm == nil { - return nil, errUnavailable("doc memory") - } - var p struct { - Doc *pubsdk.Doc `json:"doc,omitempty"` - DocRef SharedRef `json:"doc_ref,omitempty"` - } - if err := unmarshal(params, &p); err != nil { - return nil, err - } - // 文档全文可达几十 KB~数 MB,优先走共享内存。 - if err := h.resolveJSONRef(p.DocRef, &p.Doc); err != nil { - return nil, err - } - if p.Doc == nil { - return nil, fmt.Errorf("doc.insert: 缺少 doc 字段") - } - return nil, dm.Insert(p.Doc) - - case MethodDocInsertMedia: - dm := h.sdk.DocMemory() - if dm == nil { - return nil, errUnavailable("doc memory") - } - var p struct { - Doc *pubsdk.Doc `json:"doc,omitempty"` - Attachments []pubsdk.MediaAttachment `json:"attachments,omitempty"` - DocRef SharedRef `json:"doc_ref,omitempty"` - AttachRef SharedRef `json:"attachments_ref,omitempty"` - } - if err := unmarshal(params, &p); err != nil { - return nil, err - } - // 文档正文 + 附件(含媒体二进制/data URL)都优先走共享内存。 - if err := h.resolveJSONRef(p.DocRef, &p.Doc); err != nil { - return nil, err - } - if err := h.resolveJSONRef(p.AttachRef, &p.Attachments); err != nil { - return nil, err - } - if p.Doc == nil { - return nil, fmt.Errorf("doc.insertWithMedia: 缺少 doc 字段") - } - if err := dm.InsertWithMedia(p.Doc, p.Attachments); err != nil { - return nil, err - } - // 回传内核补过的字段:ID 新建时才生成,Content 含内核补的媒体标记, - // MediaDigests 是附件落盘后的完整 digest——插件靠它们后续引用同一份媒体。 - return map[string]interface{}{"doc": p.Doc}, nil - - case MethodDocRemove: - dm := h.sdk.DocMemory() - if dm == nil { - return nil, errUnavailable("doc memory") - } - var p struct { - ID string `json:"id"` - } - if err := unmarshal(params, &p); err != nil { - return nil, err - } - dm.Remove(p.ID) - return nil, nil - - case MethodDocStats: - dm := h.sdk.DocMemory() - if dm == nil { - return nil, errUnavailable("doc memory") - } - return dm.Stats(), nil - - // ---- 知识库(原 case 15/35/36)---- - case MethodKnowledgeSearch: - kn := h.sdk.Knowledge() - if kn == nil { - return nil, errUnavailable("knowledge") - } - var p struct { - Query string `json:"query"` - TopK int `json:"top_k"` - } - if err := unmarshal(params, &p); err != nil { - return nil, err - } - results, err := kn.Search(p.Query, p.TopK) - if err != nil { - return nil, err - } - if results == nil { - results = []*pubsdk.Knowledge{} - } - return map[string]interface{}{"results": results}, nil - - case MethodKnowledgeAdd: - kn := h.sdk.Knowledge() - if kn == nil { - return nil, errUnavailable("knowledge") - } - var p struct { - Name string `json:"name"` - Content string `json:"content,omitempty"` - ContentRef SharedRef `json:"content_ref,omitempty"` - } - if err := unmarshal(params, &p); err != nil { - return nil, err - } - // 知识正文可达数十 KB,优先走共享内存。内容是 JSON 字符串, - // 所以从 ref 读出后需再解一层。 - if err := h.resolveJSONRef(p.ContentRef, &p.Content); err != nil { - return nil, err - } - return nil, kn.Add(p.Name, p.Content) - - case MethodKnowledgeList: - kn := h.sdk.Knowledge() - if kn == nil { - return nil, errUnavailable("knowledge") - } - names, err := kn.List() - if err != nil { - return nil, err - } - if names == nil { - names = []string{} - } - return map[string]interface{}{"names": names}, nil - - // ---- 文本记忆(原 case 41)---- case MethodTextMemoryAppend: - tm := h.sdk.TextMemory() - if tm == nil { - return nil, errUnavailable("text memory") - } - var p struct { - Event pubsdk.TextEvent `json:"event"` - } - if err := unmarshal(params, &p); err != nil { - return nil, err - } - return nil, tm.Append(p.Event) + return h.handleTextMemory(method, params) - // ---- 设置(原 case 16/17/18/26~31/42~45/51)---- case MethodSettingsGet, MethodSettingsSet, MethodSettingsRegisterDef, MethodSettingsGetCore, MethodSettingsSetCore, MethodSettingsListCore, MethodSettingsGetPlugin, MethodSettingsSetPlugin, MethodSettingsListPlugin, - MethodSettingsList, MethodSettingsDefs, MethodSettingsDump, - MethodSettingsPlugins, MethodSettingsDataDir: - return h.settings(method, params) + MethodSettingsList, MethodSettingsDefs, MethodSettingsDump, MethodSettingsPlugins, + MethodSettingsDataDir: + return h.handleSettings(method, params) - // ---- LLM 源(原 case 19/20/37)---- - case MethodLLMListSources: - llm := h.sdk.LLM() - if llm == nil { - return nil, errUnavailable("llm") - } - sources := llm.ListSources() - if sources == nil { - sources = []string{} - } - return map[string]interface{}{"sources": sources}, nil - case MethodLLMSetSource: - llm := h.sdk.LLM() - if llm == nil { - return nil, errUnavailable("llm") - } - var p struct { - Name string `json:"name"` - } - if err := unmarshal(params, &p); err != nil { - return nil, err - } - return nil, llm.SetSource(p.Name) - case MethodLLMCurrentSource: - llm := h.sdk.LLM() - if llm == nil { - return nil, errUnavailable("llm") - } - return map[string]interface{}{"source": llm.CurrentSource()}, nil + case MethodLLMListSources, MethodLLMSetSource, MethodLLMCurrentSource: + return h.handleLLM(method, params) - // ---- 社交图(只读,原 case 21/22/38/39/40)---- case MethodSocialGetPerson, MethodSocialGetNetwork, MethodSocialGetTrait, MethodSocialGetRelation, MethodSocialListPersons: - return h.social(method, params) + return h.handleSocial(method, params) - // ---- 插件管理(原 case 48/49/50)---- - case MethodPluginReloadOne: - pm := h.sdk.PluginMgr() - if pm == nil { - return nil, errUnavailable("plugin manager") - } - var p struct { - Name string `json:"name"` - } - if err := unmarshal(params, &p); err != nil { - return nil, err - } - return nil, pm.ReloadOne(p.Name) - case MethodPluginListLoaded: - pm := h.sdk.PluginMgr() - if pm == nil { - return nil, errUnavailable("plugin manager") - } - list := pm.ListLoadedPlugins() - if list == nil { - list = []string{} - } - return map[string]interface{}{"plugins": list}, nil - case MethodPluginIsDisabled: - pm := h.sdk.PluginMgr() - if pm == nil { - return nil, errUnavailable("plugin manager") - } - var p struct { - Name string `json:"name"` - } - if err := unmarshal(params, &p); err != nil { - return nil, err - } - return map[string]interface{}{"disabled": pm.IsPluginDisabled(p.Name)}, nil + case MethodLifecycleAutoRestart: + return h.handleLifecycle(method, params) - // ---- 共享段锁仲裁(新增,§3.7)---- - case MethodStageLock: - if h.locks == nil { - return nil, fmt.Errorf("stage.lock: 锁仲裁未就绪") - } - return nil, h.locks.acquire(h.name) - case MethodStageUnlock: - if h.locks == nil { - return nil, fmt.Errorf("stage.unlock: 锁仲裁未就绪") - } - return nil, h.locks.release(h.name) + case MethodPluginReloadOne, MethodPluginListLoaded, MethodPluginIsDisabled: + return h.handlePluginMgr(method, params) - // ---- 事件订阅(原 case 23/24,子进程下首次真正可用,§3.6)---- - case MethodEventsSubscribe: - var p struct { - Types []pubsdk.EventType `json:"types"` - } - if err := unmarshal(params, &p); err != nil { - return nil, err - } - if h.evtRing == nil { - return nil, fmt.Errorf("%s: 事件环未就绪", method) - } - // 订阅请求来自子进程——handler 直接注册到 Bus, - // 事件经 EventRing 写入环后由子进程消费。 - h.evtRing.EvtRingSubscribe(p.Types) - return nil, nil + case MethodStageLock, MethodStageUnlock: + return h.handleStageLocks(method, params) - case MethodEventsUnsubscribe: - // 事件环的订阅没有持久化句柄(取消函数由 Subscribe 返回但子进程未保存)。 - // 当前设计:子进程 Stop 时由内核统一清理其订阅。 - return nil, nil + case MethodEventsSubscribe, MethodEventsUnsubscribe: + return h.handleEvents(method, params) - // ---- 共享槽池(内部传输层,见 protocol.go 注释)---- - case MethodArenaAlloc: - var p ArenaAllocParams - if err := unmarshal(params, &p); err != nil { - return nil, err - } - ref, err := h.arenaAlloc(p.Size) - if err != nil { - return nil, err - } - return ArenaAllocResult{Ref: ref}, nil + case MethodArenaAlloc, MethodArenaFree: + return h.handleArena(method, params) - case MethodArenaFree: - var p ArenaFreeParams - if err := unmarshal(params, &p); err != nil { - return nil, err - } - return nil, h.arenaFree(p.Ref) - - // ---- 多模态注入 ---- - // - // 之前这里是桩:返回“待共享段二进制通道落地”。后果是**子进程插件调 - // SetToolBlocks 必然失败**(模板只 log 一行),只有内置插件能用。 - // 现在媒体块经共享内存传递,该能力对两种插件形态等价。 - case MethodIOSetToolBlocks: - var p injectMediaParams - if err := unmarshal(params, &p); err != nil { - return nil, err - } - blocks, err := h.resolveBlocks(p) - if err != nil { - return nil, err - } - if len(blocks) == 0 { - return nil, fmt.Errorf("%s: blocks 为空", method) - } - h.sdk.SetToolBlocks(blocks) - return nil, nil } return nil, fmt.Errorf("未知 method: %s", method) diff --git a/internal/plugin/proc/corehandler_inject.go b/internal/plugin/proc/corehandler_inject.go new file mode 100644 index 0000000..8440bc6 --- /dev/null +++ b/internal/plugin/proc/corehandler_inject.go @@ -0,0 +1,130 @@ +package proc + +import ( + "encoding/json" + "fmt" +) + +// handleInject 处理 IO 注入面:文本 / 打断 / 同步 / 媒体块,以及 SetToolBlocks。 +// +// 本函数体是 corehandler.go 里 Handle 那一个大 switch 的**整块平移**: +// case 标签与 case 体逐字保留,只换了宿主函数(§3.2 的平移原则)。 +func (h *coreHandler) handleInject(method string, params json.RawMessage) (interface{}, error) { + switch method { + // ---- IO 注入(原 case 5/6/7/47)---- + // + // 注入标志位(no_memory / context_policy)由插件在调用点声明,默认 + // 记入记忆 + 不裁剪。策略值在入口校验:静默降级成 none 会让调用方 + // 以为自己声明的裁剪在生效。 + case MethodIOInjectText: + var p injectParams + if err := unmarshal(params, &p); err != nil { + return nil, err + } + if err := validateContextPolicy("io.injectText", p.ContextPolicy); err != nil { + return nil, err + } + h.sdk.InjectTextOpts(p.Source, p.Channel, h.resolveText(p), pubSdkInjectOpts(p.NoMemory, p.ContextPolicy, p.CleanerName, p.Priority)) + return nil, nil + case MethodIOInjectInterrupt: + var p injectParams + if err := unmarshal(params, &p); err != nil { + return nil, err + } + if err := validateContextPolicy("io.injectInterrupt", p.ContextPolicy); err != nil { + return nil, err + } + h.sdk.InjectInterruptTextOpts(p.Source, p.Channel, h.resolveText(p), pubSdkInjectOpts(p.NoMemory, p.ContextPolicy, p.CleanerName, p.Priority)) + return nil, nil + case MethodIOInjectTextNoMem: + var p injectParams + if err := unmarshal(params, &p); err != nil { + return nil, err + } + if err := validateContextPolicy("io.injectTextNoMem", p.ContextPolicy); err != nil { + return nil, err + } + // 旧 RPC 语义就是「不进记忆」,显式标志位只可能再叠上 context_policy。 + h.sdk.InjectTextOpts(p.Source, p.Channel, h.resolveText(p), pubSdkInjectOpts(true, p.ContextPolicy, p.CleanerName, p.Priority)) + return nil, nil + case MethodIOInjectSync: + var p injectParams + if err := unmarshal(params, &p); err != nil { + return nil, err + } + if err := validateContextPolicy("io.injectInputSync", p.ContextPolicy); err != nil { + return nil, err + } + reply := h.sdk.InjectInputSyncOpts(p.Source, p.Channel, h.resolveText(p), pubSdkInjectOpts(p.NoMemory, p.ContextPolicy, p.CleanerName, p.Priority)) + return map[string]interface{}{"reply": reply}, nil + + case MethodIOInjectMedia: + var p injectMediaParams + if err := unmarshal(params, &p); err != nil { + return nil, err + } + if err := validateContextPolicy("io.injectMedia", p.ContextPolicy); err != nil { + return nil, err + } + blocks, err := h.resolveBlocks(p) + if err != nil { + return nil, err + } + h.sdk.InjectInputMediaOpts(p.Source, p.Channel, p.Text, blocks, pubSdkInjectOpts(p.NoMemory, p.ContextPolicy, p.CleanerName, p.Priority)) + return nil, nil + + case MethodIOInjectMediaSync: + var p injectMediaParams + if err := unmarshal(params, &p); err != nil { + return nil, err + } + if err := validateContextPolicy("io.injectMediaSync", p.ContextPolicy); err != nil { + return nil, err + } + blocks, err := h.resolveBlocks(p) + if err != nil { + return nil, err + } + reply := h.sdk.InjectInputMediaSyncOpts(p.Source, p.Channel, p.Text, blocks, pubSdkInjectOpts(p.NoMemory, p.ContextPolicy, p.CleanerName, p.Priority)) + return map[string]interface{}{"reply": reply}, nil + + case MethodIOInjectInterruptMedia: + var p injectMediaParams + if err := unmarshal(params, &p); err != nil { + return nil, err + } + if err := validateContextPolicy("io.injectInterruptMedia", p.ContextPolicy); err != nil { + return nil, err + } + blocks, err := h.resolveBlocks(p) + if err != nil { + return nil, err + } + h.sdk.InjectInterruptMediaOpts(p.Source, p.Channel, p.Text, blocks, pubSdkInjectOpts(p.NoMemory, p.ContextPolicy, p.CleanerName, p.Priority)) + return nil, nil + + // ---- 多模态注入 ---- + // + // 之前这里是桩:返回“待共享段二进制通道落地”。后果是**子进程插件调 + // SetToolBlocks 必然失败**(模板只 log 一行),只有内置插件能用。 + // 现在媒体块经共享内存传递,该能力对两种插件形态等价。 + case MethodIOSetToolBlocks: + var p injectMediaParams + if err := unmarshal(params, &p); err != nil { + return nil, err + } + blocks, err := h.resolveBlocks(p) + if err != nil { + return nil, err + } + if len(blocks) == 0 { + return nil, fmt.Errorf("%s: blocks 为空", method) + } + h.sdk.SetToolBlocks(blocks) + return nil, nil + } + + // 组内不应到达:Handle 的分派表与本函数的 case 标签同源, + // 出现即说明两处不同步。 + return nil, fmt.Errorf("未知 method: %s", method) +} diff --git a/internal/plugin/proc/corehandler_memory.go b/internal/plugin/proc/corehandler_memory.go new file mode 100644 index 0000000..84a4379 --- /dev/null +++ b/internal/plugin/proc/corehandler_memory.go @@ -0,0 +1,306 @@ +package proc + +import ( + "encoding/json" + "fmt" + + pubsdk "gitcode.com/JianFeeeee/homeagent-sdk/sdk" +) + +// handleGraphMemory 处理图记忆(实体/关系):recall / commit / introspect / merge / purge。 +// +// 本函数体是 corehandler.go 里 Handle 那一个大 switch 的**整块平移**: +// case 标签与 case 体逐字保留,只换了宿主函数(§3.2 的平移原则)。 +func (h *coreHandler) handleGraphMemory(method string, params json.RawMessage) (interface{}, error) { + switch method { + // ---- 图记忆(原 case 9/10/11/12/13)---- + case MethodMemoryRecall: + mem := h.sdk.Memory() + if mem == nil { + return nil, errUnavailable("memory") + } + var p struct { + Query []string `json:"query"` + Depth int `json:"depth"` + } + if err := unmarshal(params, &p); err != nil { + return nil, err + } + entities, relations, err := mem.Recall(p.Query, p.Depth) + if err != nil { + return nil, err + } + if entities == nil { + entities = []pubsdk.Entity{} + } + if relations == nil { + relations = []pubsdk.Relation{} + } + return map[string]interface{}{"entities": entities, "relations": relations}, nil + + case MethodMemoryCommit: + mem := h.sdk.Memory() + if mem == nil { + return nil, errUnavailable("memory") + } + var p struct { + Triples []pubsdk.Triple `json:"triples"` + } + if err := unmarshal(params, &p); err != nil { + return nil, err + } + return nil, mem.Commit(p.Triples) + + case MethodMemoryIntrospect: + mem := h.sdk.Memory() + if mem == nil { + return nil, errUnavailable("memory") + } + return mem.Introspect() + + case MethodMemoryMerge: + mem := h.sdk.Memory() + if mem == nil { + return nil, errUnavailable("memory") + } + var p struct { + Source string `json:"source"` + Target string `json:"target"` + } + if err := unmarshal(params, &p); err != nil { + return nil, err + } + n, err := mem.MergeEntities(p.Source, p.Target) + if err != nil { + return nil, err + } + return map[string]interface{}{"merged": n}, nil + + case MethodMemoryPurge: + mem := h.sdk.Memory() + if mem == nil { + return nil, errUnavailable("memory") + } + var p struct { + Criteria map[string]string `json:"criteria"` + Mode string `json:"mode"` + } + if err := unmarshal(params, &p); err != nil { + return nil, err + } + if p.Mode == "" { + p.Mode = "soft" + } + n, err := mem.Purge(p.Criteria, p.Mode) + if err != nil { + return nil, err + } + return map[string]interface{}{"purged": n}, nil + + } + + // 组内不应到达:Handle 的分派表与本函数的 case 标签同源, + // 出现即说明两处不同步。 + return nil, fmt.Errorf("未知 method: %s", method) +} + +// handleDocMemory 处理文档记忆:query / insert / insertWithMedia / remove / stats。 +// +// 本函数体是 corehandler.go 里 Handle 那一个大 switch 的**整块平移**: +// case 标签与 case 体逐字保留,只换了宿主函数(§3.2 的平移原则)。 +func (h *coreHandler) handleDocMemory(method string, params json.RawMessage) (interface{}, error) { + switch method { + // ---- 文档记忆(原 case 14/32/33/34)---- + case MethodDocQuery: + dm := h.sdk.DocMemory() + if dm == nil { + return nil, errUnavailable("doc memory") + } + var p struct { + Text string `json:"text"` + TopK int `json:"top_k"` + } + if err := unmarshal(params, &p); err != nil { + return nil, err + } + docs := dm.Query(p.Text, p.TopK) + if docs == nil { + docs = []*pubsdk.Doc{} + } + return map[string]interface{}{"docs": docs}, nil + + case MethodDocInsert: + dm := h.sdk.DocMemory() + if dm == nil { + return nil, errUnavailable("doc memory") + } + var p struct { + Doc *pubsdk.Doc `json:"doc,omitempty"` + DocRef SharedRef `json:"doc_ref,omitempty"` + } + if err := unmarshal(params, &p); err != nil { + return nil, err + } + // 文档全文可达几十 KB~数 MB,优先走共享内存。 + if err := h.resolveJSONRef(p.DocRef, &p.Doc); err != nil { + return nil, err + } + if p.Doc == nil { + return nil, fmt.Errorf("doc.insert: 缺少 doc 字段") + } + return nil, dm.Insert(p.Doc) + + case MethodDocInsertMedia: + dm := h.sdk.DocMemory() + if dm == nil { + return nil, errUnavailable("doc memory") + } + var p struct { + Doc *pubsdk.Doc `json:"doc,omitempty"` + Attachments []pubsdk.MediaAttachment `json:"attachments,omitempty"` + DocRef SharedRef `json:"doc_ref,omitempty"` + AttachRef SharedRef `json:"attachments_ref,omitempty"` + } + if err := unmarshal(params, &p); err != nil { + return nil, err + } + // 文档正文 + 附件(含媒体二进制/data URL)都优先走共享内存。 + if err := h.resolveJSONRef(p.DocRef, &p.Doc); err != nil { + return nil, err + } + if err := h.resolveJSONRef(p.AttachRef, &p.Attachments); err != nil { + return nil, err + } + if p.Doc == nil { + return nil, fmt.Errorf("doc.insertWithMedia: 缺少 doc 字段") + } + if err := dm.InsertWithMedia(p.Doc, p.Attachments); err != nil { + return nil, err + } + // 回传内核补过的字段:ID 新建时才生成,Content 含内核补的媒体标记, + // MediaDigests 是附件落盘后的完整 digest——插件靠它们后续引用同一份媒体。 + return map[string]interface{}{"doc": p.Doc}, nil + + case MethodDocRemove: + dm := h.sdk.DocMemory() + if dm == nil { + return nil, errUnavailable("doc memory") + } + var p struct { + ID string `json:"id"` + } + if err := unmarshal(params, &p); err != nil { + return nil, err + } + dm.Remove(p.ID) + return nil, nil + + case MethodDocStats: + dm := h.sdk.DocMemory() + if dm == nil { + return nil, errUnavailable("doc memory") + } + return dm.Stats(), nil + + } + + // 组内不应到达:Handle 的分派表与本函数的 case 标签同源, + // 出现即说明两处不同步。 + return nil, fmt.Errorf("未知 method: %s", method) +} + +// handleKnowledge 处理知识库:search / add / list。 +// +// 本函数体是 corehandler.go 里 Handle 那一个大 switch 的**整块平移**: +// case 标签与 case 体逐字保留,只换了宿主函数(§3.2 的平移原则)。 +func (h *coreHandler) handleKnowledge(method string, params json.RawMessage) (interface{}, error) { + switch method { + // ---- 知识库(原 case 15/35/36)---- + case MethodKnowledgeSearch: + kn := h.sdk.Knowledge() + if kn == nil { + return nil, errUnavailable("knowledge") + } + var p struct { + Query string `json:"query"` + TopK int `json:"top_k"` + } + if err := unmarshal(params, &p); err != nil { + return nil, err + } + results, err := kn.Search(p.Query, p.TopK) + if err != nil { + return nil, err + } + if results == nil { + results = []*pubsdk.Knowledge{} + } + return map[string]interface{}{"results": results}, nil + + case MethodKnowledgeAdd: + kn := h.sdk.Knowledge() + if kn == nil { + return nil, errUnavailable("knowledge") + } + var p struct { + Name string `json:"name"` + Content string `json:"content,omitempty"` + ContentRef SharedRef `json:"content_ref,omitempty"` + } + if err := unmarshal(params, &p); err != nil { + return nil, err + } + // 知识正文可达数十 KB,优先走共享内存。内容是 JSON 字符串, + // 所以从 ref 读出后需再解一层。 + if err := h.resolveJSONRef(p.ContentRef, &p.Content); err != nil { + return nil, err + } + return nil, kn.Add(p.Name, p.Content) + + case MethodKnowledgeList: + kn := h.sdk.Knowledge() + if kn == nil { + return nil, errUnavailable("knowledge") + } + names, err := kn.List() + if err != nil { + return nil, err + } + if names == nil { + names = []string{} + } + return map[string]interface{}{"names": names}, nil + + } + + // 组内不应到达:Handle 的分派表与本函数的 case 标签同源, + // 出现即说明两处不同步。 + return nil, fmt.Errorf("未知 method: %s", method) +} + +// handleTextMemory 处理文本记忆追加。 +// +// 本函数体是 corehandler.go 里 Handle 那一个大 switch 的**整块平移**: +// case 标签与 case 体逐字保留,只换了宿主函数(§3.2 的平移原则)。 +func (h *coreHandler) handleTextMemory(method string, params json.RawMessage) (interface{}, error) { + switch method { + // ---- 文本记忆(原 case 41)---- + case MethodTextMemoryAppend: + tm := h.sdk.TextMemory() + if tm == nil { + return nil, errUnavailable("text memory") + } + var p struct { + Event pubsdk.TextEvent `json:"event"` + } + if err := unmarshal(params, &p); err != nil { + return nil, err + } + return nil, tm.Append(p.Event) + + } + + // 组内不应到达:Handle 的分派表与本函数的 case 标签同源, + // 出现即说明两处不同步。 + return nil, fmt.Errorf("未知 method: %s", method) +} diff --git a/internal/plugin/proc/corehandler_register.go b/internal/plugin/proc/corehandler_register.go new file mode 100644 index 0000000..9121a29 --- /dev/null +++ b/internal/plugin/proc/corehandler_register.go @@ -0,0 +1,61 @@ +package proc + +import ( + "encoding/json" + "fmt" + + pubsdk "gitcode.com/JianFeeeee/homeagent-sdk/sdk" +) + +// handleRegister 处理注册面:工具 / stage / 输出通道 / 插件 API / 入站通道。 +// +// 本函数体是 corehandler.go 里 Handle 那一个大 switch 的**整块平移**: +// case 标签与 case 体逐字保留,只换了宿主函数(§3.2 的平移原则)。 +func (h *coreHandler) handleRegister(method string, params json.RawMessage) (interface{}, error) { + switch method { + // ---- 注册面(原 case 1/2/3/4/46)---- + case MethodToolRegister: + return h.toolRegister(params) + case MethodStageRegister: + return h.stageRegister(params) + case MethodOutputRegister: + return h.outputRegister(params) + case MethodAPIRegister: + var p struct { + Name string `json:"name"` + } + if err := unmarshal(params, &p); err != nil { + return nil, err + } + return nil, h.sdk.RegisterPluginAPI(p.Name) + case MethodInputRegister: + var p struct { + Name string `json:"name"` + Def pubsdk.ChannelDef `json:"def"` + HasCleaner bool `json:"has_cleaner"` + } + if err := unmarshal(params, &p); err != nil { + return nil, err + } + if p.Name == "" { + return nil, fmt.Errorf("input.register: 缺少 name") + } + cleaner, err := h.cleanerProxy(CleanerScopeInput, p.Name, p.HasCleaner) + if err != nil { + return nil, fmt.Errorf("input.register: %w", err) + } + if err := validateContextPolicy("input.register", p.Def.ContextPolicy); err != nil { + return nil, err + } + // 整体传 p.Def(只是把函数型的 Cleaner 换成代理),不要手写字段白名单: + // 白名单会让新增字段静默丢失。 + def := p.Def + def.Cleaner = cleaner + return nil, h.sdk.RegisterInputChannel(p.Name, def) + + } + + // 组内不应到达:Handle 的分派表与本函数的 case 标签同源, + // 出现即说明两处不同步。 + return nil, fmt.Errorf("未知 method: %s", method) +} diff --git a/internal/plugin/proc/corehandler_runtime.go b/internal/plugin/proc/corehandler_runtime.go new file mode 100644 index 0000000..36f56d1 --- /dev/null +++ b/internal/plugin/proc/corehandler_runtime.go @@ -0,0 +1,148 @@ +package proc + +import ( + "encoding/json" + "fmt" + + pubsdk "gitcode.com/JianFeeeee/homeagent-sdk/sdk" +) + +// handlePluginMgr 处理插件管理:reloadOne / listLoaded / isDisabled。 +// +// 本函数体是 corehandler.go 里 Handle 那一个大 switch 的**整块平移**: +// case 标签与 case 体逐字保留,只换了宿主函数(§3.2 的平移原则)。 +func (h *coreHandler) handlePluginMgr(method string, params json.RawMessage) (interface{}, error) { + switch method { + // ---- 插件管理(原 case 48/49/50)---- + case MethodPluginReloadOne: + pm := h.sdk.PluginMgr() + if pm == nil { + return nil, errUnavailable("plugin manager") + } + var p struct { + Name string `json:"name"` + } + if err := unmarshal(params, &p); err != nil { + return nil, err + } + return nil, pm.ReloadOne(p.Name) + case MethodPluginListLoaded: + pm := h.sdk.PluginMgr() + if pm == nil { + return nil, errUnavailable("plugin manager") + } + list := pm.ListLoadedPlugins() + if list == nil { + list = []string{} + } + return map[string]interface{}{"plugins": list}, nil + case MethodPluginIsDisabled: + pm := h.sdk.PluginMgr() + if pm == nil { + return nil, errUnavailable("plugin manager") + } + var p struct { + Name string `json:"name"` + } + if err := unmarshal(params, &p); err != nil { + return nil, err + } + return map[string]interface{}{"disabled": pm.IsPluginDisabled(p.Name)}, nil + + } + + // 组内不应到达:Handle 的分派表与本函数的 case 标签同源, + // 出现即说明两处不同步。 + return nil, fmt.Errorf("未知 method: %s", method) +} + +// handleStageLocks 处理共享段锁仲裁(§3.7)。 +// +// 本函数体是 corehandler.go 里 Handle 那一个大 switch 的**整块平移**: +// case 标签与 case 体逐字保留,只换了宿主函数(§3.2 的平移原则)。 +func (h *coreHandler) handleStageLocks(method string, params json.RawMessage) (interface{}, error) { + switch method { + // ---- 共享段锁仲裁(新增,§3.7)---- + case MethodStageLock: + if h.locks == nil { + return nil, fmt.Errorf("stage.lock: 锁仲裁未就绪") + } + return nil, h.locks.acquire(h.name) + case MethodStageUnlock: + if h.locks == nil { + return nil, fmt.Errorf("stage.unlock: 锁仲裁未就绪") + } + return nil, h.locks.release(h.name) + + } + + // 组内不应到达:Handle 的分派表与本函数的 case 标签同源, + // 出现即说明两处不同步。 + return nil, fmt.Errorf("未知 method: %s", method) +} + +// handleEvents 处理事件订阅(§3.6,子进程下首次真正可用)。 +// +// 本函数体是 corehandler.go 里 Handle 那一个大 switch 的**整块平移**: +// case 标签与 case 体逐字保留,只换了宿主函数(§3.2 的平移原则)。 +func (h *coreHandler) handleEvents(method string, params json.RawMessage) (interface{}, error) { + switch method { + // ---- 事件订阅(原 case 23/24,子进程下首次真正可用,§3.6)---- + case MethodEventsSubscribe: + var p struct { + Types []pubsdk.EventType `json:"types"` + } + if err := unmarshal(params, &p); err != nil { + return nil, err + } + if h.evtRing == nil { + return nil, fmt.Errorf("%s: 事件环未就绪", method) + } + // 订阅请求来自子进程——handler 直接注册到 Bus, + // 事件经 EventRing 写入环后由子进程消费。 + h.evtRing.EvtRingSubscribe(p.Types) + return nil, nil + + case MethodEventsUnsubscribe: + // 事件环的订阅没有持久化句柄(取消函数由 Subscribe 返回但子进程未保存)。 + // 当前设计:子进程 Stop 时由内核统一清理其订阅。 + return nil, nil + + } + + // 组内不应到达:Handle 的分派表与本函数的 case 标签同源, + // 出现即说明两处不同步。 + return nil, fmt.Errorf("未知 method: %s", method) +} + +// handleArena 处理共享槽池分配/释放(内部传输层,见 protocol.go)。 +// +// 本函数体是 corehandler.go 里 Handle 那一个大 switch 的**整块平移**: +// case 标签与 case 体逐字保留,只换了宿主函数(§3.2 的平移原则)。 +func (h *coreHandler) handleArena(method string, params json.RawMessage) (interface{}, error) { + switch method { + // ---- 共享槽池(内部传输层,见 protocol.go 注释)---- + case MethodArenaAlloc: + var p ArenaAllocParams + if err := unmarshal(params, &p); err != nil { + return nil, err + } + ref, err := h.arenaAlloc(p.Size) + if err != nil { + return nil, err + } + return ArenaAllocResult{Ref: ref}, nil + + case MethodArenaFree: + var p ArenaFreeParams + if err := unmarshal(params, &p); err != nil { + return nil, err + } + return nil, h.arenaFree(p.Ref) + + } + + // 组内不应到达:Handle 的分派表与本函数的 case 标签同源, + // 出现即说明两处不同步。 + return nil, fmt.Errorf("未知 method: %s", method) +} diff --git a/internal/plugin/proc/corehandler_settings.go b/internal/plugin/proc/corehandler_settings.go new file mode 100644 index 0000000..d32540b --- /dev/null +++ b/internal/plugin/proc/corehandler_settings.go @@ -0,0 +1,112 @@ +package proc + +import ( + "encoding/json" + "fmt" +) + +// handleSettings 处理全部设置项读写(14 个 method 共用一条实现)。 +// +// 本函数体是 corehandler.go 里 Handle 那一个大 switch 的**整块平移**: +// case 标签与 case 体逐字保留,只换了宿主函数(§3.2 的平移原则)。 +func (h *coreHandler) handleSettings(method string, params json.RawMessage) (interface{}, error) { + switch method { + // ---- 设置(原 case 16/17/18/26~31/42~45/51)---- + case MethodSettingsGet, MethodSettingsSet, MethodSettingsRegisterDef, + MethodSettingsGetCore, MethodSettingsSetCore, MethodSettingsListCore, + MethodSettingsGetPlugin, MethodSettingsSetPlugin, MethodSettingsListPlugin, + MethodSettingsList, MethodSettingsDefs, MethodSettingsDump, + MethodSettingsPlugins, MethodSettingsDataDir: + return h.settings(method, params) + + } + + // 组内不应到达:Handle 的分派表与本函数的 case 标签同源, + // 出现即说明两处不同步。 + return nil, fmt.Errorf("未知 method: %s", method) +} + +// handleLLM 处理 LLM 源查询与切换。 +// +// 本函数体是 corehandler.go 里 Handle 那一个大 switch 的**整块平移**: +// case 标签与 case 体逐字保留,只换了宿主函数(§3.2 的平移原则)。 +func (h *coreHandler) handleLLM(method string, params json.RawMessage) (interface{}, error) { + switch method { + // ---- LLM 源(原 case 19/20/37)---- + case MethodLLMListSources: + llm := h.sdk.LLM() + if llm == nil { + return nil, errUnavailable("llm") + } + sources := llm.ListSources() + if sources == nil { + sources = []string{} + } + return map[string]interface{}{"sources": sources}, nil + case MethodLLMSetSource: + llm := h.sdk.LLM() + if llm == nil { + return nil, errUnavailable("llm") + } + var p struct { + Name string `json:"name"` + } + if err := unmarshal(params, &p); err != nil { + return nil, err + } + return nil, llm.SetSource(p.Name) + case MethodLLMCurrentSource: + llm := h.sdk.LLM() + if llm == nil { + return nil, errUnavailable("llm") + } + return map[string]interface{}{"source": llm.CurrentSource()}, nil + + } + + // 组内不应到达:Handle 的分派表与本函数的 case 标签同源, + // 出现即说明两处不同步。 + return nil, fmt.Errorf("未知 method: %s", method) +} + +// handleSocial 处理社交图只读查询(人物/网络/特质/关系)。 +// +// 本函数体是 corehandler.go 里 Handle 那一个大 switch 的**整块平移**: +// case 标签与 case 体逐字保留,只换了宿主函数(§3.2 的平移原则)。 +func (h *coreHandler) handleSocial(method string, params json.RawMessage) (interface{}, error) { + switch method { + // ---- 社交图(只读,原 case 21/22/38/39/40)---- + case MethodSocialGetPerson, MethodSocialGetNetwork, MethodSocialGetTrait, + MethodSocialGetRelation, MethodSocialListPersons: + return h.social(method, params) + + } + + // 组内不应到达:Handle 的分派表与本函数的 case 标签同源, + // 出现即说明两处不同步。 + return nil, fmt.Errorf("未知 method: %s", method) +} + +// handleLifecycle 处理生命周期开关(当前仅自动重启)。 +// +// 本函数体是 corehandler.go 里 Handle 那一个大 switch 的**整块平移**: +// case 标签与 case 体逐字保留,只换了宿主函数(§3.2 的平移原则)。 +func (h *coreHandler) handleLifecycle(method string, params json.RawMessage) (interface{}, error) { + switch method { + // ---- 生命周期(原 case 8)---- + case MethodLifecycleAutoRestart: + var p struct { + Enabled bool `json:"enabled"` + } + if err := unmarshal(params, &p); err != nil { + return nil, err + } + h.sdk.SetAutoRestart(p.Enabled) + return nil, nil + + } + + // 组内不应到达:Handle 的分派表与本函数的 case 标签同源, + // 出现即说明两处不同步。 + return nil, fmt.Errorf("未知 method: %s", method) +}