mirror of
https://gitcode.com/JianFeeeee/HomeAgent.git
synced 2026-09-23 02:18:06 +00:00
refactor(proc): coreHandler.Handle 550 行按 method 组拆成 14 个分部函数
原 Handle 是一个 550 行的巨型 switch(C ABI 51 个 case 的整块平移),
按协议面拆进同包 5 个新文件、14 个小函数:
corehandler_register.go handleRegister 注册面(tool/stage/output/api/input)
corehandler_inject.go handleInject IO 注入 + SetToolBlocks
corehandler_memory.go handleGraphMemory 图记忆
handleDocMemory 文档记忆
handleKnowledge 知识库
handleTextMemory 文本记忆
corehandler_settings.go handleSettings 设置(14 个 method 共用一条实现)
handleLLM LLM 源
handleSocial 社交图只读
handleLifecycle 生命周期开关
corehandler_runtime.go handlePluginMgr 插件管理
handleStageLocks 段锁仲裁
handleEvents 事件订阅
handleArena 共享槽池
Handle 保留 capability 强制检查,只做「method → 分部函数」一跳。
零漂移保证:case 标签由脚本从原文提取(不手抄常量名),case 体逐字搬迁,
逐函数比对确认 59 个标签 / 510 行 case 体与原文件完全一致(仅行首缩进经 gofmt 重排)。
验证:go build ./... / go vet / go test ./internal/plugin/... / -race 全绿。
This commit is contained in:
@ -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)
|
||||
|
||||
130
internal/plugin/proc/corehandler_inject.go
Normal file
130
internal/plugin/proc/corehandler_inject.go
Normal file
@ -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)
|
||||
}
|
||||
306
internal/plugin/proc/corehandler_memory.go
Normal file
306
internal/plugin/proc/corehandler_memory.go
Normal file
@ -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)
|
||||
}
|
||||
61
internal/plugin/proc/corehandler_register.go
Normal file
61
internal/plugin/proc/corehandler_register.go
Normal file
@ -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)
|
||||
}
|
||||
148
internal/plugin/proc/corehandler_runtime.go
Normal file
148
internal/plugin/proc/corehandler_runtime.go
Normal file
@ -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)
|
||||
}
|
||||
112
internal/plugin/proc/corehandler_settings.go
Normal file
112
internal/plugin/proc/corehandler_settings.go
Normal file
@ -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)
|
||||
}
|
||||
Reference in New Issue
Block a user