Files
HomeAgent/cmd/homed/bootstrap.go
JianFeeeee 69446a2649 feat(scheduler): 主 agent 忙时把积压任务自动转投给驻留子
问题(2026-09-19 线上实测):主 agent 被长任务占住时(现场:12 分 8 秒、69 次
工具调用),后来到达的消息全部以 level insufficient 排进中断队列干等 —— 同级
中断不能抢占同级运行任务(canPreempt),只能等前一个跑完。而内核本有驻留子
(独立 agent + 独立调度器)可并行干活。

行为(用户 2026-09-19 明确要求):
- 触发:运行任务持续 > offload_busy_after(5m) 且积压 >= offload_min_pending(3)
- 拉起/复用「转投专用」驻留子,把积压的纯排队输入转投过去
- 在原队列位置留下说明「[系统] N 条积压任务已转投给驻留子 agent X 处理…」

通道配置(按用户口径,与人工创建的子刻意不同):
- 不配 inputch(内核的干活 agent,不接收插件用户输入)
- 持有全部输出通道(结果要能发回 qq/webui 等正确通道)

三个设计要点(都是实测撞出来的,写进代码注释与设计文档 §7.1):
1. 检查必须在**独立 goroutine**:schedulerLoop 同步执行任务,放它里面在
   「正忙」期间根本回不到循环顶部 ⇒ 永不触发(我第一版就写错了,测试才发现)。
2. 只转投 TaskQueued 纯排队输入:中断任务带级别语义、self 任务与父的记忆面绑定。
3. 转投失败/关闭时必须把任务**放回队列前端**:吞一条输入比多处理一条更糟。

这是设计 §7「决策在父的模型手里」的**刻意例外**(父正忙、物理上无法决策,
而积压任务本来就是空的),已在文档中显式记录,且默认关闭、由部署方显式打开。

测试 11 条:只取排队输入 / 不足量不取 / 放回不丢任务 / 说明自解释 / 默认关闭 /
空闲不触发 / 端到端转投 / 上限不增殖 / 独立 goroutine 确实会触发。
2026-09-19 16:59:47 +08:00

877 lines
36 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

package main
import (
"context"
"flag"
"fmt"
"io"
"log"
"os"
"os/signal"
"path/filepath"
"strings"
"syscall"
"time"
agentPkg "gitcode.com/JianFeeeee/HomeAgent/internal/agent"
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/ipc"
"gitcode.com/JianFeeeee/HomeAgent/internal/knowledge"
logpkg "gitcode.com/JianFeeeee/HomeAgent/internal/log"
luapkg "gitcode.com/JianFeeeee/HomeAgent/internal/lua"
"gitcode.com/JianFeeeee/HomeAgent/internal/memory"
"gitcode.com/JianFeeeee/HomeAgent/internal/memory/document"
"gitcode.com/JianFeeeee/HomeAgent/internal/memory/media"
"gitcode.com/JianFeeeee/HomeAgent/internal/memory/pipeline"
"gitcode.com/JianFeeeee/HomeAgent/internal/memory/social"
"gitcode.com/JianFeeeee/HomeAgent/internal/memory/text"
"gitcode.com/JianFeeeee/HomeAgent/internal/memory/vector"
"gitcode.com/JianFeeeee/HomeAgent/internal/nlp"
"gitcode.com/JianFeeeee/HomeAgent/internal/plugin"
"gitcode.com/JianFeeeee/HomeAgent/internal/recovery"
sdk "gitcode.com/JianFeeeee/HomeAgent/internal/sdk"
"gitcode.com/JianFeeeee/HomeAgent/internal/supervisor"
"gitcode.com/JianFeeeee/HomeAgent/internal/tracker"
"gitcode.com/JianFeeeee/HomeAgent/pkg/embedding"
"gitcode.com/JianFeeeee/HomeAgent/pkg/types"
)
// defaultSystemPrompt 是内置默认人格模板:不含版本号字面量,
// 被问版本时以运行时快照为准(历史上写死版本号导致实例自称旧版本)。
const defaultSystemPrompt = `你是 HomeAgent一个持续运行的个人管家。
你的每次回复会自动发送到当前输出通道(默认=输入源),无需额外工具。
如需切换回复通道,使用 output_set_channel。
如需异步发送消息或通知,使用 output_send 指定通道和内容。
使用 output_list_channels 查看可用通道及其能力。
可用工具列表会由系统自动传入,按需使用即可。以下是你尤其需要关注的几类工具:
- memory_* — 图记忆(长期记忆,记录和查询个人信息/事实)
- knowledge_* — 知识库(查阅预设知识文档)
- doc_* — 文档记忆(近期对话的存档,查询后自动清除)
- person_* — 人物特质与社交关系网
- llm_* — LLM 源管理(列出/切换模型提供商)
- output_* — 输出通道管理(切换/发送消息)
- timer_set — 设置定时提醒
- plgreload — 热重载插件
- spawn_child — 生成子 Agent 异步执行独立任务(可传 max_turns 控制工具轮数,默认 5
并行策略:遇到多个互不依赖的子任务时,优先并行 spawn 多个子 Agent 而非自己串行逐个执行;
长耗时任务(批量处理、多轮搜索汇总)也应交给子 Agent避免阻塞当前对话。
- describe_image — 描述用户上传的图片
- transcribe_audio — 转写用户上传的音频
- ocr_image — 识别图片中的文字
命令与文件操作策略:
- cmd_run 经完整 shellbash执行支持管道、分号、&&、命令替换、heredoc、重定向。
- 多步交互式程序vim/top/ssh 会话、需要持续输入的进程)用 terminal_create 创建终端,
terminal_write 发送输入、terminal_read 读输出——不要用 cmd_run 硬等交互程序退出。
- 写文件优先 files_write原子+留档),生成多行内容时可用 heredoc 或 files_write
不要用 echo 拼接长文本。
- 读用户发来的文件用 files_read向 webui 回传图片/文件用 output_send__webui(type=image/file)。
当用户上传图片或音频时,系统会自动附着媒体内容。如果模型不支持直接处理多媒体,请使用上述工具。
回复你的真实想法,用自然语言与用户交流。不要在回复中使用 emoji 表情。`
// setupLogging 初始化日志:行号前缀 + 同时输出到控制台与 <data>/log/ 下的本次启动文件。
//
// 本函数体是 main() 里对应启动阶段的整块平移:语句、日志文本、错误语义不变,
// 只把「*dataDir」变成参数、把 defer 变成由调用点注册的 cleanup。
func setupLogging(dataDir string) string {
log.SetFlags(log.Ldate | log.Ltime | log.Lshortfile)
// 文件日志:同时输出到控制台和 data/log/ 目录
logDir := filepath.Join(dataDir, "log")
if err := os.MkdirAll(logDir, 0755); err != nil {
log.Printf("[homed] warning: cannot create log dir: %v", err)
} else {
logPath := filepath.Join(logDir, fmt.Sprintf("homed_%s.log", time.Now().Format("2006-01-02_15-04-05")))
logFile, err := os.OpenFile(logPath, os.O_CREATE|os.O_WRONLY|os.O_APPEND, 0644)
if err != nil {
log.Printf("[homed] warning: cannot open log file: %v", err)
} else {
log.SetOutput(io.MultiWriter(os.Stderr, logFile))
log.Printf("[homed] logging to %s", logPath)
}
}
return logDir
}
// ensureDataDirs 建好启动期需要的全部目录,返回 agent 的 overlayfs 工作目录。
//
// 本函数体是 main() 里对应启动阶段的整块平移:语句、日志文本、错误语义不变,
// 只把「*dataDir」变成参数、把 defer 变成由调用点注册的 cleanup。
func ensureDataDirs(dataDir string) string {
agentWorkDir := filepath.Join(dataDir, "agentfs")
dirs := []string{
dataDir,
filepath.Join(dataDir, "snapshots"),
filepath.Join(dataDir, "plugins"),
filepath.Join(dataDir, "changesets"),
filepath.Join(dataDir, "memory"),
filepath.Join(dataDir, "memory", "raw"),
filepath.Join(dataDir, "adapters"),
agentWorkDir,
}
for _, d := range dirs {
if err := os.MkdirAll(d, 0755); err != nil {
log.Fatalf("create dir %s: %v", d, err)
}
}
return agentWorkDir
}
// memoryStack 聚合记忆侧组件:图库、索引器、社交图、蒸馏器。
type memoryStack struct {
db *memory.GraphDB
indexer *memory.Indexer
social *social.SocialStore
distiller *pipeline.Distiller
}
// initMemoryStack 初始化图记忆 / 索引 / 社交图 / 蒸馏管线。
//
// 本函数体是 main() 里对应启动阶段的整块平移:语句、日志文本、错误语义不变,
// 只把「*dataDir」变成参数、把 defer 变成由调用点注册的 cleanup。
func initMemoryStack(dataDir string) (*memoryStack, func()) {
memDB, err := memory.NewGraphDB(filepath.Join(dataDir, "memory", "graph.db"))
if err != nil {
log.Printf("[homed] warning: memory init failed: %v", err)
memDB = nil
} else {
log.Printf("[homed] graph memory initialized")
}
memIdx := memory.NewIndexer(memDB)
memIdx.Sync() // 启动时立即同步避免前30分钟空窗
socialStore := social.New(memDB)
distiller := pipeline.NewDistiller(memDB, dataDir, pipeline.DistillerConfig{
Interval: 10 * time.Minute,
RetentionDays: 7,
BatchSize: 50,
})
if memDB != nil {
// 这里**故意不写 defer distiller.Stop()**:本函数在 return 时即触发
// defer而 Stop() → cancel() 会让刚启动的 distillLoop 立刻退出,
// 规则蒸馏管线启动即死、10min 心跳从不运行(旧 main() 拆分时的残留)。
// 停机由调用点注册的 cleanup 负责(见下方返回值)。
distiller.Start()
}
return &memoryStack{db: memDB, indexer: memIdx, social: socialStore, distiller: distiller},
func() {
// 与原 main 的两个 defer 同序LIFO先停蒸馏器再关图库。
if memDB != nil {
distiller.Stop()
}
if memDB != nil {
memDB.Close()
}
}
}
// startMemoryCandidateConsumer 起一个常驻 goroutine把 eventbus 上的 memory_candidate 事件写进文本记忆并喂给蒸馏管线。
//
// 本函数体是 main() 里对应启动阶段的整块平移:语句、日志文本、错误语义不变,
// 只把「*dataDir」变成参数、把 defer 变成由调用点注册的 cleanup。
func startMemoryCandidateConsumer(ctx context.Context, iom *agentIO.IOManager, textMem *text.Memory, memDB *memory.GraphDB, distiller *pipeline.Distiller) {
go func() {
for {
select {
case <-ctx.Done():
return
case evt, ok := <-iom.OutputChan():
if !ok {
return
}
if evt.Target == "memory" && evt.Type == "memory_candidate" {
source, _ := evt.Payload["source"].(string)
input, _ := evt.Payload["input"].(string)
response, _ := evt.Payload["response"].(string)
toolsUsed, _ := evt.Payload["tools_used"].([]string)
toolResults, _ := evt.Payload["tool_results"].([]interface{})
agentID, _ := evt.Payload["agent_id"].(string)
if input != "" && textMem != nil {
te := text.Event{
Timestamp: time.Now().Unix(),
Source: source,
Input: input,
Response: response,
ToolsUsed: toolsUsed,
AgentID: agentID,
}
if err := textMem.Append(te); err != nil {
log.Printf("[homed] text memory append: %v", err)
}
}
if input != "" && memDB != nil {
distiller.Append("agent", "user", input)
}
if response != "" && memDB != nil {
distiller.Append("agent", "assistant", response)
}
// 工具输出接入蒸馏管线
for _, tr := range toolResults {
if trMap, ok := tr.(map[string]interface{}); ok {
if text, ok := trMap["output"].(string); ok && text != "" && memDB != nil {
distiller.Append("agent", "tool", text)
}
}
}
}
}
}
}()
}
// initLLMProviders 按配置注册全部 LLM 源(每个源经 Lua 适配器协议转换),并把各适配器的并发额度汇总回 Lua VM。
//
// 本函数体是 main() 里对应启动阶段的整块平移:语句、日志文本、错误语义不变,
// 只把「*dataDir」变成参数、把 defer 变成由调用点注册的 cleanup。
func initLLMProviders(cfg *types.Config, luaVM *luapkg.VM, baseAPIKey string) *agentAPI.ProviderManager {
providerMgr := agentAPI.NewProviderManager()
adapterConcurrency := map[string]int{}
for _, src := range cfg.LLM.Sources {
if !agentAPI.IsValidSourceConfig(src.Name, src.BaseURL, src.Model, src.Adapter) {
log.Printf("[homed] skip invalid llm source %q (base_url=%q model=%q adapter=%q)", src.Name, src.BaseURL, src.Model, src.Adapter)
continue
}
key := src.APIKey
if key == "" {
key = baseAPIKey
}
luaProvider := agentAPI.NewLuaAdaptedProvider(agentAPI.BaseConfig{
Model: src.Model,
BaseURL: src.BaseURL,
APIKey: key,
Temperature: cfg.LLM.Temperature,
MaxTokens: cfg.LLM.MaxTokens,
ContextWindow: src.ContextWindow,
MaxConcurrent: src.MaxConcurrent,
Priority: src.Priority,
Vision: src.Vision,
Audio: src.Audio,
}, luaVM, src.Name, src.Adapter)
providerMgr.Register(src.Name, luaProvider)
if src.Adapter != "" {
adapterConcurrency[src.Adapter] += src.MaxConcurrent
}
}
luaVM.ConfigureConcurrency(adapterConcurrency)
if cfg.LLM.Provider != "" {
providerMgr.SetDefault(cfg.LLM.Provider)
}
return providerMgr
}
// initDocStore 启动文档记忆。flush 的唯一入口是 Stop(),所以关停时必须调用它。
//
// 本函数体是 main() 里对应启动阶段的整块平移:语句、日志文本、错误语义不变,
// 只把「*dataDir」变成参数、把 defer 变成由调用点注册的 cleanup。
func initDocStore(cfg *types.Config) (*document.Store, func()) {
docStore := document.NewStore(filepath.Join(cfg.Daemon.DataDir, "memory", "documents"), memory.TokenizeWords)
if err := docStore.Start(); err != nil {
log.Printf("[homed] warning: document store: %v", err)
}
// 关停时落盘。文档记忆的内存态变更(迁移结果、访问计数等)只在 flush
// 里写盘,而 flush 的唯一入口是 Stop()——此前全仓无人调用它,
// 于是迁移结果永不落盘、每次启动白算一遍。
return docStore, func() { docStore.Stop() }
}
// initMediaStore 按开关启动内容寻址的媒体存储;开不起来只告警(媒体记忆非对话必需品)。
//
// 本函数体是 main() 里对应启动阶段的整块平移:语句、日志文本、错误语义不变,
// 只把「*dataDir」变成参数、把 defer 变成由调用点注册的 cleanup。
func initMediaStore(cfgReg *internalConfig.ConfigRegistry, cfg *types.Config) (*media.Store, func()) {
var mediaStore *media.Store
if cfgReg.GetBool("core.memory.media.enabled", true) {
mediaDir := cfgReg.GetString("core.memory.media.dir",
filepath.Join(cfg.Daemon.DataDir, "memory", "media"))
ms, err := media.New(mediaDir)
if err != nil {
// 媒体存储开不起来不该阻止启动——它是记忆增强,不是对话必需品
log.Printf("[homed] warning: media store: %v媒体记忆已禁用", err)
} else {
mediaStore = ms
st := mediaStore.Stats()
log.Printf("[homed] media store active: %v 条 / %v 字节",
st["count"], st["total_bytes"])
}
}
return mediaStore, func() {
if mediaStore != nil {
mediaStore.Close()
}
}
}
// initMultimodalSpace 从公共注册表打开多模态向量 provider。返回 (空间, provider 名, 失败原因, cleanup):后两个值只用于状态报告。
//
// 本函数体是 main() 里对应启动阶段的整块平移:语句、日志文本、错误语义不变,
// 只把「*dataDir」变成参数、把 defer 变成由调用点注册的 cleanup。
func initMultimodalSpace(cfgReg *internalConfig.ConfigRegistry) (vector.MultimodalEmbedder, string, string, func()) {
var multimodalSpace vector.MultimodalEmbedder
// 这两个值只用于状态报告healthcheck_kernel 的 onnx 段):
// 「配了哪个 provider」与「为什么没启用」避免只能看到 false 却不知原因。
var mmProviderName, mmErr string
var closeAdapted func()
if mmProvider := cfgReg.GetString("core.memory.multimodal_space.provider", ""); mmProvider != "" {
mmProviderName = mmProvider
opts := map[string]string{}
const optPrefix = "core.memory.multimodal_space.options."
for _, key := range cfgReg.List("core.memory.multimodal_space.options.") {
opts[strings.TrimPrefix(key, optPrefix)] = cfgReg.GetString(key, "")
}
provider, err := embedding.Open(mmProvider, embedding.Config{Options: opts})
if err != nil {
mmErr = err.Error()
log.Printf("[homed] warning: 多模态向量 provider %q 打开失败: %v多模态向量检索已禁用已注册: %s",
mmProvider, err, strings.Join(embedding.Names(), ", "))
} else if adapted, err := vector.AdaptProvider(provider); err != nil {
provider.Close()
mmErr = err.Error()
log.Printf("[homed] warning: 多模态向量 provider %q 元数据不合法: %v多模态向量检索已禁用", mmProvider, err)
} else {
multimodalSpace = adapted
info := provider.Info()
// 指纹可能很长(模型文件哈希),日志里只取前 12 个字符便于对照。
shortFP := info.Fingerprint
if len(shortFP) > 12 {
shortFP = shortFP[:12]
}
log.Printf("[homed] multimodal space active: provider=%s dim=%d fp=%s modalities=%v",
mmProvider, info.Dimension, shortFP, info.Modalities)
}
}
if closeAdapted == nil {
closeAdapted = func() {}
}
return multimodalSpace, mmProviderName, mmErr, closeAdapted
}
// initKnowledgeStore 启动知识库。
//
// 本函数体是 main() 里对应启动阶段的整块平移:语句、日志文本、错误语义不变,
// 只把「*dataDir」变成参数、把 defer 变成由调用点注册的 cleanup。
func initKnowledgeStore(cfg *types.Config) *knowledge.Store {
ks := knowledge.NewStore(filepath.Join(cfg.Daemon.DataDir, "knowledge"))
if err := ks.Start(); err != nil {
log.Printf("[homed] warning: knowledge store: %v", err)
} else {
log.Printf("[homed] knowledge store active with %d items", len(ks.List()))
}
return ks
}
// loadPersonality 按「个人文件 > 配置项」的优先级解析人格内容,并对腐坏内容告警。
//
// 本函数体是 main() 里对应启动阶段的整块平移:语句、日志文本、错误语义不变,
// 只把「*dataDir」变成参数、把 defer 变成由调用点注册的 cleanup。
func loadPersonality(cfg *types.Config, cfgReg *internalConfig.ConfigRegistry) *agentPkg.Personality {
personalPath := filepath.Join(cfg.Daemon.DataDir, "personal", "personal.md")
personality, err := agentPkg.LoadPersonality(personalPath)
if err != nil {
log.Printf("[homed] warning: load personality: %v", err)
}
if personality != nil && personality.Content != "" {
log.Printf("[homed] 人格来源=文件 %s优先于配置项%d 字节", personalPath, len(personality.Content))
if hints := agentPkg.PersonaStaleHints(personality.Content); len(hints) > 0 {
log.Printf("[homed] warning: 人格文件含会腐坏的内容 %v — 建议迁到配置项 core.agent.personal_prompt"+
"(默认模板不含版本号,被问版本时以运行时快照为准)", hints)
}
} else if pv := cfgReg.GetString("core.agent.personal_prompt", internalConfig.DefaultPersonaPrompt); strings.TrimSpace(pv) != "" {
personality = &agentPkg.Personality{Content: pv, Path: "(core.agent.personal_prompt)"}
log.Printf("[homed] 人格来源=配置项 core.agent.personal_prompt%d 字节", len(pv))
} else {
log.Printf("[homed] 人格来源=无(配置项为空且无人格文件)")
}
return personality
}
// newStageAndRegistry 建阶段管道与插件注册表,把内核依赖接到注册表上。
//
// 本函数体是 main() 里对应启动阶段的整块平移:语句、日志文本、错误语义不变,
// 只把「*dataDir」变成参数、把 defer 变成由调用点注册的 cleanup。
func newStageAndRegistry(cfg *types.Config, cfgReg *internalConfig.ConfigRegistry, iom *agentIO.IOManager,
evBus *events.Bus, memDB *memory.GraphDB, textMem *text.Memory, docStore *document.Store,
mediaStore *media.Store, ks *knowledge.Store, providerMgr *agentAPI.ProviderManager,
dataDir string) (*agentCore.StageHost, *plugin.Registry) {
stageHost := agentCore.NewStageHost()
pluginReg := plugin.NewRegistry()
pluginReg.SetIOManager(iom)
pluginReg.SetEventBus(evBus)
pluginReg.SetMemory(memDB)
pluginReg.SetTextMemory(textMem)
pluginReg.SetDocStore(docStore)
pluginReg.SetMediaStore(mediaStore) // 插件写入的记忆也走媒体链路nil 时静默降级
pluginReg.SetKnowledge(ks)
pluginReg.SetProviderManager(providerMgr)
pluginReg.SetConfigRegistry(cfgReg)
pluginReg.SetPluginDir(cfg.Plugin.Dir)
pluginReg.SetDataDir(dataDir) // 插件 SettingsAPI.DataDir() 的数据根目录
// Wire registration callbacks: plugins' RegisterTool/RegisterStage → StageHost
pluginReg.SetToolRegistrar(func(name string, def sdk.ToolDef, handler sdk.ToolHandler) error {
log.Printf("[homed] SetToolRegistrar registering tool: %s (plugin=%s)", name, def.Plugin)
return stageHost.RegisterTool(name, def, handler)
})
pluginReg.SetStageRegistrar(func(stage sdk.Stage, handler sdk.StageHandler) {
stageHost.RegisterStage(stage, handler)
})
pluginReg.SetAPIRegistrar(func(name string) error {
return nil
})
pluginReg.SetToolCleaner(stageHost)
return stageHost, pluginReg
}
// newMainAgent 组装主 Agent把内核各面IO/记忆/文档/知识/媒体/社交/文本/插件/状态)接进 AgentConfig。
//
// 本函数体是 main() 里对应启动阶段的整块平移:语句、日志文本、错误语义不变,
// 只把「*dataDir」变成参数、把 defer 变成由调用点注册的 cleanup。
func newMainAgent(cfg *types.Config, cfgReg *internalConfig.ConfigRegistry, provider agentAPI.Provider,
providerMgr *agentAPI.ProviderManager, iom *agentIO.IOManager, memDB *memory.GraphDB,
memIdx *memory.Indexer, trk *tracker.Tracker, docStore *document.Store, ks *knowledge.Store,
socialStore *social.SocialStore, textMem *text.Memory, mediaStore *media.Store,
personality *agentPkg.Personality, pluginReg *plugin.Registry, embedder *memory.StaticEmbedder,
multimodalSpace vector.MultimodalEmbedder, mmProviderName, mmErr string,
stageHost *agentCore.StageHost, evBus *events.Bus) *agentCore.Agent {
sysPrompt := cfgReg.GetString("core.agent.system_prompt", defaultSystemPrompt)
if sysPrompt == "" {
sysPrompt = defaultSystemPrompt
}
agent := agentCore.New(agentCore.AgentConfig{
ID: "main",
SystemPrompt: sysPrompt,
Provider: provider,
ProviderManager: providerMgr,
IO: iom,
Memory: memDB,
Indexer: memIdx,
Tracker: trk,
DocStore: docStore,
Knowledge: ks,
SocialStore: socialStore,
TextMemory: textMem,
MediaStore: mediaStore,
Personality: personality,
// 人格落库面:首启门禁(任何通道都问一次)与 persona_set 工具用。
// 与 WebUI 向导共用 internal/config 的同一份落库逻辑。
PersonaStore: internalConfig.RegistryPersonaStore{Reg: cfgReg},
PluginReg: pluginReg,
PluginDir: cfg.Plugin.Dir,
// DataDir驻留子的 temp 图库锚点(<data>/residents/<id>/graph.db
// 漏接时的现象是"工具存在、可调用、但创建必失败"——只有真实二进制才看得出来。
DataDir: cfg.Daemon.DataDir,
DistillInterval: cfgReg.GetDuration("core.agent.distill_interval", 30*time.Minute),
ArchiveInterval: cfgReg.GetDuration("core.agent.archive_interval", 60*time.Minute),
ReviewInterval: cfgReg.GetDuration("core.agent.review_interval", 120*time.Minute),
MergeInterval: cfgReg.GetDuration("core.agent.merge_interval", 120*time.Minute),
MaxToolTurns: cfgReg.GetInt("core.agent.max_tool_turns", 10),
Offload: agentCore.OffloadOptions{
Enabled: cfgReg.GetBool("core.agent.offload_enabled", false),
BusyAfter: cfgReg.GetDuration("core.agent.offload_busy_after", 5*time.Minute),
MinPending: cfgReg.GetInt("core.agent.offload_min_pending", 3),
MaxResidents: cfgReg.GetInt("core.agent.offload_max_residents", 2),
},
ContextSavePath: filepath.Join(cfg.Daemon.DataDir, "memory", "context.json"),
EmbeddingModelPath: cfgReg.GetString("core.agent.embedding_model_path", ""),
Embedder: embedder,
MultimodalSpace: multimodalSpace,
EmbeddingProvider: mmProviderName,
EmbeddingError: mmErr,
StageHost: stageHost,
EventBus: evBus,
ThinkingEnabled: cfg.LLM.ThinkingEnabled,
InputProcessing: cfg.InputProcessing,
})
return agent
}
// initONNXParser 初始化依存句法分析器(内嵌 ONNX 模型,失败则退回规则引擎)。
//
// 本函数体是 main() 里对应启动阶段的整块平移:语句、日志文本、错误语义不变,
// 只把「*dataDir」变成参数、把 defer 变成由调用点注册的 cleanup。
func initONNXParser(cfg *types.Config, cfgReg *internalConfig.ConfigRegistry) {
modelPath := cfgReg.GetString("core.agent.onnx_model_path", "")
onnxParser, err := nlp.NewONNXParser(nlp.ONNXConfig{
ModelPath: modelPath,
DataDir: filepath.Join(cfg.Daemon.DataDir, "nlp"),
})
if err != nil {
log.Printf("[homed] warn: ONNX parser init: %v, using fallback", err)
} else {
nlp.SetDefaultParser(onnxParser)
log.Printf("[homed] dep parser initialized (model: %s)", modelPath)
}
}
// loadPlugins 建插件目录、按启动模式决定 allowlist然后加载全部插件。
//
// 本函数体是 main() 里对应启动阶段的整块平移:语句、日志文本、错误语义不变,
// 只把「*dataDir」变成参数、把 defer 变成由调用点注册的 cleanup。
func loadPlugins(cfg *types.Config, cfgReg *internalConfig.ConfigRegistry, pluginReg *plugin.Registry,
stageHost *agentCore.StageHost, bootMode, dataDir string) {
// Auto-create plugins directory (without hardcoding plugin names)
os.MkdirAll(cfg.Plugin.Dir, 0755)
// failback 受限启动:仅装载 failback 插件集webfetch/files/cmd 为内核内置,
// 此处仅控制外部插件,默认含 recoverydiag 以便直接在受限态产出恢复结论)
if bootMode == "failback" {
list := cfgReg.GetString("core.agent.failback_plugins", "webui,pluginmgr,recoverydiag")
// 优先使用 guard.yaml 经过 recovery 任务下发的插件集guard 是 failback 权威)
if task, terr := recovery.LoadTask(recovery.TaskPath(dataDir)); terr == nil && len(task.Plugins()) > 0 {
list = strings.Join(task.Plugins(), ",")
}
var names []string
for _, s := range strings.Split(list, ",") {
if s = strings.TrimSpace(s); s != "" {
names = append(names, s)
}
}
pluginReg.SetLoadAllowlist(names)
log.Printf("[homed] failback boot: plugin allowlist = %v", names)
}
// Load all plugins — each scans its own dir and is loaded via factory or .so
if err := pluginReg.Load(cfg.Plugin.Dir); err != nil {
log.Printf("[homed] warning: load plugins: %v", err)
}
log.Printf("[homed] stage host ready with %d registered tools", stageHost.ToolCount())
}
// startAgentRuntime 接线技能索引、起日志管理、启动 agent返回逆序关停的 cleanup。
//
// 本函数体是 main() 里对应启动阶段的整块平移:语句、日志文本、错误语义不变,
// 只把「*dataDir」变成参数、把 defer 变成由调用点注册的 cleanup。
func startAgentRuntime(cfgReg *internalConfig.ConfigRegistry, pluginReg *plugin.Registry,
agent *agentCore.Agent, logDir string, ctx context.Context) func() {
// 技能索引接线skillmgr 插件实现 SkillIndexProvider 时注入 agent方案B prompt 注入)
if sp := pluginReg.Get("skillmgr"); sp != nil {
if prov, ok := sp.(agentCore.SkillIndexProvider); ok {
agent.SetSkillIndexProvider(prov)
log.Printf("[homed] skill index wired from skillmgr plugin")
}
}
// 日志管理:层级压缩 + 保留策略
logManager := logpkg.NewManager(logDir, cfgReg)
go logManager.Start(ctx)
agent.Start()
return func() {
// 与原 main 的两个 defer 同序LIFO先停 agent再停日志管理。
agent.Stop()
logManager.Stop()
}
}
// startIPCServer 起 PING/ACK 心跳服务(含 kernel 状态快照),返回仅在启动成功后生效的 cleanup。
//
// 本函数体是 main() 里对应启动阶段的整块平移:语句、日志文本、错误语义不变,
// 只把「*dataDir」变成参数、把 defer 变成由调用点注册的 cleanup。
func startIPCServer(dataDir, bootMode string, agent *agentCore.Agent) (*ipc.Server, func()) {
started := false
ipcServer := ipc.NewServer(dataDir, func() *ipc.Status {
st := agent.GetKernelStatus()
llmOK := st != nil && st.LLM.Available
tools := 0
if st != nil {
tools = len(st.Tools)
}
uptime := int64(0)
if st != nil {
if d, err := time.ParseDuration(st.Uptime); err == nil {
uptime = int64(d.Seconds())
}
}
return &ipc.Status{
PID: os.Getpid(),
Boot: bootMode,
UptimeSec: uptime,
LLMOK: &llmOK,
Tools: tools,
LastDiag: lastDiagSummary(dataDir),
}
})
if err := ipcServer.Start(); err != nil {
log.Printf("[homed] warning: ipc heartbeat server: %v", err)
} else {
}
return ipcServer, func() {
if started {
ipcServer.Stop()
}
}
}
// startSupervisorRuntime 把真实存活源与重启通道接到 supervisor 上,返回重启请求通道。
//
// 本函数体是 main() 里对应启动阶段的整块平移:语句、日志文本、错误语义不变,
// 只把「*dataDir」变成参数、把 defer 变成由调用点注册的 cleanup。
func startSupervisorRuntime(sup *supervisor.Daemon, trk *tracker.Tracker, agent *agentCore.Agent) chan struct{} {
sup.SetTracker(trk)
sup.RegisterAgent("main")
// 真实存活源 + 重启通道daemon 心跳语义由此修正lastHB 只在确认存活时更新),
// 重启动作不再空转——清理后以特殊退出码交给 guard/systemd 重建。
restartCh := make(chan struct{}, 1)
sup.SetHeartbeatSource(func(id types.AgentID) (time.Time, types.HealthStatus, error) {
st := agent.GetKernelStatus()
if st == nil {
return time.Time{}, types.HealthDown, fmt.Errorf("no kernel status")
}
h := types.HealthHealthy
if !st.LLM.Available {
h = types.HealthDegraded
}
return time.Now(), h, nil
})
sup.SetRestartHandler(func(id types.AgentID) {
select {
case restartCh <- struct{}{}:
default:
}
})
return restartCh
}
// startHeartbeat 每 5s 触碰 <data>/heartbeatguard 据此判定 worker 存活/卡死),返回停止函数。
//
// 本函数体是 main() 里对应启动阶段的整块平移:语句、日志文本、错误语义不变,
// 只把「*dataDir」变成参数、把 defer 变成由调用点注册的 cleanup。
func startHeartbeat(dataDir string, ctx context.Context) func() {
// 心跳:每 5s 触碰 <data>/heartbeatguard 据此判定工作进程是否存活/卡死
hbPath := filepath.Join(dataDir, "heartbeat")
hbStop := make(chan struct{})
go func() {
t := time.NewTicker(5 * time.Second)
defer t.Stop()
writeHB := func() {
if f, err := os.OpenFile(hbPath, os.O_CREATE|os.O_WRONLY|os.O_APPEND, 0644); err == nil {
fmt.Fprintf(f, "t=%d\n", time.Now().Unix())
f.Close()
}
}
writeHB()
for {
select {
case <-t.C:
writeHB()
case <-hbStop:
return
case <-ctx.Done():
return
}
}
}()
return func() { close(hbStop) }
}
// waitForShutdown 阻塞至 SIGINT/SIGTERM 或 supervisor 请求重启,然后按原 main 的顺序清理,需要重建时以退出码交回 guard。
//
// 本函数体是 main() 里对应启动阶段的整块平移:语句、日志文本、错误语义不变,
// 只把「*dataDir」变成参数、把 defer 变成由调用点注册的 cleanup。
func waitForShutdown(ctx context.Context, dataDir string, restartCh chan struct{}, stopHeartbeat func(),
pluginReg *plugin.Registry, trk *tracker.Tracker, cfgReg *internalConfig.ConfigRegistry, sup *supervisor.Daemon) {
sigCh := make(chan os.Signal, 1)
signal.Notify(sigCh, syscall.SIGINT, syscall.SIGTERM)
restartRequested := false
select {
case <-sigCh:
log.Printf("[homed] shutting down...")
case <-restartCh:
restartRequested = true
log.Printf("[homed] restart requested, shutting down cleanly then exiting with code %d", exitRestartRequested)
}
stopHeartbeat()
pluginReg.StopAll()
if trk != nil {
trk.Stop()
}
if err := cfgReg.Flush(); err != nil {
log.Printf("[homed] flush config: %v", err)
}
sup.Shutdown()
log.Printf("[homed] stopped")
if restartRequested {
os.Exit(exitRestartRequested)
}
}
// initLuaVM 起 Lua VMLLM 协议适配启动失败只告警cleanup 为 no-op。
func initLuaVM(cfg *types.Config) (*luapkg.VM, func()) {
luaVM := luapkg.NewVM(filepath.Join(cfg.Daemon.DataDir, "adapters"))
if err := luaVM.Start(); err != nil {
log.Printf("[homed] warning: lua vm init failed: %v", err)
return luaVM, func() {}
}
return luaVM, luaVM.Stop
}
// initSupervisor 起守护管理(代理生命周期管理);起不来是致命错误。
func initSupervisor(cfg *types.Config) *supervisor.Daemon {
sup := supervisor.New(cfg)
if err := sup.Start(); err != nil {
log.Fatalf("start supervisor: %v", err)
}
return sup
}
// initTracker 起 overlayfs 变更追踪;无 overlayfs 支持时降级为非致命告警。
func initTracker(cfg *types.Config, agentWorkDir string) *tracker.Tracker {
trk := tracker.NewTracker(cfg.Daemon.DataDir, agentWorkDir,
tracker.WithKeepChangesets(100),
tracker.WithMaxChangesetAge(30*24*time.Hour),
)
if err := trk.Init(); err != nil {
log.Printf("[homed] warning: tracker init: %v", err)
} else {
if err := trk.Start(); err != nil {
log.Printf("[homed] warning: tracker mount overlay: %v (non-fatal: no overlayfs support?)", err)
} else {
log.Printf("[homed] change tracker active at %s", trk.MergeDir())
}
}
return trk
}
// initKernelAPI 建内核与插件之间的两个通道IOManagerIO 抽象层)+ EventBus事件总线
func initKernelAPI() (*agentIO.IOManager, *events.Bus) {
iom := agentIO.NewIOManager()
evBus := events.NewBus()
log.Printf("[homed] kernel API ready: IOManager + EventBus")
return iom, evBus
}
// initTextMemory 起文本记忆启动失败只告警cleanup 为 no-op。
func initTextMemory(cfg *types.Config) (*text.Memory, func()) {
textMem := text.New(filepath.Join(cfg.Daemon.DataDir, "memory", "text"))
if err := textMem.Start(); err != nil {
log.Printf("[homed] warning: text memory start: %v", err)
return textMem, func() {}
}
log.Printf("[homed] text memory active at %s", filepath.Join(cfg.Daemon.DataDir, "memory", "text"))
return textMem, textMem.Stop
}
// resolveBaseAPIKey 解析兜底 API key配置项 > LLM_API_KEY > DEEPSEEK_API_KEY。
func resolveBaseAPIKey(cfg *types.Config) string {
apiKey := cfg.LLM.APIKey
if apiKey == "" {
apiKey = os.Getenv("LLM_API_KEY")
}
if apiKey == "" {
apiKey = os.Getenv("DEEPSEEK_API_KEY")
}
return apiKey
}
// wirePluginSDK 把内核各面注入每个插件的 PluginSDK阶段6 将替换遗留的 util.Configure
func wirePluginSDK(pluginReg *plugin.Registry, luaVM *luapkg.VM, baseAPIKey string, sup *supervisor.Daemon,
trk *tracker.Tracker, cfg *types.Config, stageHost *agentCore.StageHost, memIdx *memory.Indexer,
agent *agentCore.Agent) {
pluginReg.SetLuaVM(luaVM)
pluginReg.SetBaseAPIKey(baseAPIKey)
pluginReg.SetSupervisor(supervisor.NewSDKAdapter(sup))
pluginReg.SetTracker(trk)
pluginReg.SetConfig(cfg)
pluginReg.SetStageHost(stageHost)
pluginReg.SetIndexer(memIdx)
pluginReg.SetStatusProvider(agent)
pluginReg.SetTerminalAPI(agent)
}
// resolveWebUIOverride 解析 webui 监听地址的覆盖值,空串表示不覆盖。
//
// 优先级CLI --webui > 核心配置 webui.listen_addr仅当它被改成非内置默认值
// 两者都不给时由 webui 插件自己的 settings["addr"] 决定。
//
// 为什么不写成“内核在插件加载前 Set 插件 settings['addr']”:那时
// config_webui 表还没建(表只在插件注册 def 时创建PluginSettings.Set 的
// INSERT 会失败而错误被忽略,随后插件 Start 里 RegisterDef 才建表并写入默认
// :8080 —— 于是 CLI --webui 与 webui.listen_addr **一直是死配置**
// 无论怎么传都监听 :8080。覆盖值改由插件自己接收webui.SetListenOverride
func resolveWebUIOverride(cfgReg *internalConfig.ConfigRegistry, httpAddr string) string {
if strings.TrimSpace(httpAddr) != "" {
return strings.TrimSpace(httpAddr)
}
// webui.listen_addr 的播种默认值就是 ":8080";与默认值相同视为“未配置”,
// 否则会把用户在设置页里改过的插件 addr 顶掉。
if v := strings.TrimSpace(cfgReg.GetString("webui.listen_addr", ":8080")); v != "" && v != ":8080" {
return v
}
return ""
}
// options 是 worker 的命令行参数。
type options struct {
dataDir string
httpAddr string
cliSocket string
role string
boot string
}
// parseFlags 解析命令行参数。
func parseFlags() options {
dataDir := flag.String("data", "", "data directory (default: auto-detect next to binary)")
httpAddr := flag.String("webui", "", "webui listen address (default: webui.listen_addr from config)")
cliSocket := flag.String("socket", "", "cli unix socket path (default: <data>/cli.sock)")
role := flag.String("role", "agent", "process role: guard (父守护) | agent (工作进程)")
boot := flag.String("boot", "normal", "agent boot mode: normal | failback (受限启动,仅 failback 插件集)")
flag.Parse()
return options{dataDir: *dataDir, httpAddr: *httpAddr, cliSocket: *cliSocket, role: *role, boot: *boot}
}
// compactConfigDB 在空闲页够多时压缩配置库;失败只告警(不影响启动)。
//
// 触发条件(见 internal/config.MaybeCompact空闲页 >= 1MB 且占页数 >= 25%。
// 放在插件加载之后调用——迁移/清理大值发生在插件 Start 里,之前调用没有意义。
func compactConfigDB(cfgReg *internalConfig.ConfigRegistry) {
before := int64(-1)
if st, err := os.Stat(cfgReg.DBPath()); err == nil {
before = st.Size()
}
done, err := cfgReg.MaybeCompact(1<<20, 0.25)
if err != nil {
log.Printf("[homed] warning: 配置库压缩失败: %v", err)
return
}
if !done {
return
}
after := before
if st, err := os.Stat(cfgReg.DBPath()); err == nil {
after = st.Size()
}
log.Printf("[homed] 配置库已压缩: %d -> %d 字节", before, after)
}