Files
homeagent-sdk/example/recoverydiag/plugin.go
JianFeeeee d57c5eaf3e fix(examples): 全插件安全审查修复(qq/a2a/memo/calendar/rss/browser/bili/recoverydiag)
审查发现并修复 7 项问题:
- P1 qq: downloadURL 裸 http.Get 无超时 → 120s client
- P2 a2a: inbound http.Server 零超时 → Read 30s/Write 120s/Idle 60s
- P3 bili: output_dir 配置项零校验 → 系统目录黑名单(/、/etc、/usr、/var 等)
- P4 recoverydiag: db_path LLM 可控任意 sqlite → 强制限制 data 目录内
- P5 memo/calendar/rss: os.WriteFile 直写 → atomicWriteJSON (temp+rename)
- P6 qq: 3 处后台 goroutine(已读/rcon转发/下载)加 panic recover
- P7 browser: dump-dom failback Kill 后补 wait 回收僵尸进程

recoverydiag 此前被 .gitignore 排除,但其 db_path 安全修复
属生产代码,故取消忽略并入库。

全部经 plugindev 重打包升版安装验证 config_kept=true。
2026-08-26 17:04:41 +08:00

1024 lines
30 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 (
"crypto/sha256"
"encoding/hex"
"encoding/json"
"fmt"
"log"
"os"
"os/exec"
"path/filepath"
"regexp"
"sort"
"strings"
"time"
sdk "gitcode.com/JianFeeeee/homeagent-sdk/sdk"
)
// Plugin 快速检查/崩溃取证工具集。全部确定性检出,返回结论而非原文,供 guard / failback 决策。
type Plugin struct {
name string
sdk *sdk.PluginSDK
muKey string
dataDir string
logDir string
cfgPath string
}
func (p *Plugin) Name() string { return p.name }
func (p *Plugin) Start(s *sdk.PluginSDK) error {
s.SetAutoRestart(true)
p.sdk = s
p.muKey = p.name + "_"
p.resolveDirs(s)
s.Settings().RegisterDef(sdk.ConfigDef{
Key: "db_check_cmd",
Default: "sqlite3",
Type: "string",
DisplayName: "sqlite3 CLI 路径",
Description: "diag_db 用到的 sqlite3 命令;留空则仅在可用时使用,缺失回退到内核 Settings 读取。留空=auto",
Category: "recoverydiag",
})
s.Settings().RegisterDef(sdk.ConfigDef{
Key: "recovery_kb_dir",
Default: "",
Type: "string",
DisplayName: "结论落盘目录",
Description: "diag_loc 结论 JSON 落盘目录,缺省用 <data_dir>/recovery_kb",
Category: "recoverydiag",
})
s.RegisterTool(p.muKey+"diag_triage", sdk.ToolDef{
Name: p.muKey + "diag_triage",
Description: "快速分诊:根据退出码/信号/存活状态粗分崩溃类别(进程死亡 vs 配置类不可达 vs 正常)。返回结论,不返回日志原文。",
Parameters: map[string]interface{}{
"type": "object",
"properties": map[string]interface{}{
"exit_code": map[string]interface{}{"type": "integer", "description": "进程退出码0=正常)"},
"signal": map[string]interface{}{"type": "string", "description": "终止信号名(如 SIGSEGV/SIGKILL/OOM可选"},
"uptime_ms": map[string]interface{}{"type": "integer", "description": "进程存活毫秒,可选"},
"still_alive": map[string]interface{}{"type": "boolean", "description": "主 agent 是否仍在运行/对心跳有响应,可选"},
"crash_reason": map[string]interface{}{"type": "string", "description": "守护方附带的已知原因描述,可选"},
},
},
NoMemory: true,
}, p.handleTriage)
s.RegisterTool(p.muKey+"diag_db", sdk.ToolDef{
Name: p.muKey + "diag_db",
Description: "config.db 完整性PRAGMA integrity_check+ LLM 源解析校验core.llm.sources.* 必备字段),逐项 ok/fail返回结论。",
Parameters: map[string]interface{}{
"type": "object",
"properties": map[string]interface{}{
"db_path": map[string]interface{}{"type": "string", "description": "config.db 路径,缺省用 <data_dir>/config.db"},
},
},
NoMemory: true,
}, p.handleDB)
s.RegisterTool(p.muKey+"diag_log_scan", sdk.ToolDef{
Name: p.muKey + "diag_log_scan",
Description: "在日志目录时间窗内统计已知错误签名panic/OOM/网络不可达/provider失败/sql/致命)出现次数,返回按类统计与主导结论。",
Parameters: map[string]interface{}{
"type": "object",
"properties": map[string]interface{}{
"log_dir": map[string]interface{}{"type": "string", "description": "日志目录,缺省用 <data_dir>/log"},
"since_minutes": map[string]interface{}{"type": "integer", "description": "只看最近 N 分钟,缺省看全部"},
"max_lines": map[string]interface{}{"type": "integer", "description": "最多扫描行数(防止读取过大文件),缺省 200000"},
},
},
NoMemory: true,
}, p.handleLogScan)
s.RegisterTool(p.muKey+"diag_delta", sdk.ToolDef{
Name: p.muKey + "diag_delta",
Description: "对比 baseline上次 good 快照/目录)与现状目录,输出 created/modified/deleted 文件清单与摘要,用于判定'改了什么'。",
Parameters: map[string]interface{}{
"type": "object",
"properties": map[string]interface{}{
"baseline_dir": map[string]interface{}{"type": "string", "description": "基线目录(快照解包目录),必传"},
"current_dir": map[string]interface{}{"type": "string", "description": "现状目录(如 agentfs merged/upper必传"},
"pattern": map[string]interface{}{"type": "string", "description": "只关注匹配该子串的相对路径,可选"},
"max_items": map[string]interface{}{"type": "integer", "description": "返回最多文件条数,缺省 500"},
},
},
NoMemory: true,
}, p.handleDelta)
s.RegisterTool(p.muKey+"diag_loc", sdk.ToolDef{
Name: p.muKey + "diag_loc",
Description: "综合分诊/DB/日志/快照四项结论,按因果强度正交排序定位根因并给出推荐恢复动作。调用前请先跑其余 diag_* 并把结论传入。",
Parameters: map[string]interface{}{
"type": "object",
"properties": map[string]interface{}{
"triage": map[string]interface{}{"type": "object", "description": "diag_triage 返回的结论对象"},
"db": map[string]interface{}{"type": "object", "description": "diag_db 返回的结论对象"},
"log_scan": map[string]interface{}{"type": "object", "description": "diag_log_scan 返回的结论对象"},
"delta": map[string]interface{}{"type": "object", "description": "diag_delta 返回的结论对象"},
"persist": map[string]interface{}{"type": "boolean", "description": "是否落盘结论到 recovery_kb 并回流知识库,缺省 true"},
},
},
NoMemory: true,
}, p.handleLoc)
log.Printf("[%s] started: data_dir=%s log_dir=%s", p.name, p.dataDir, p.logDir)
return nil
}
func (p *Plugin) Stop() error {
log.Printf("[%s] stopped", p.name)
return nil
}
func NewPluginFactory(name string, config map[string]interface{}) (sdk.Plugin, error) {
return &Plugin{name: name}, nil
}
// resolveDirs 从内核 Settings 解出数据目录与日志目录。
func (p *Plugin) resolveDirs(s *sdk.PluginSDK) {
if v, err := s.Settings().GetCore("core.daemon.data_dir"); err == nil && v != nil {
if sv, ok := v.(string); ok && sv != "" {
p.dataDir = sv
}
}
if v, err := s.Settings().GetCore("core.log.path"); err == nil && v != nil {
if sv, ok := v.(string); ok && sv != "" {
p.logDir = sv
}
}
if p.logDir == "" && p.dataDir != "" {
p.logDir = filepath.Join(p.dataDir, "log")
}
if p.dataDir != "" {
p.cfgPath = filepath.Join(p.dataDir, "config.db")
}
}
// ===== 参数解析辅助 =====
func argString(args map[string]interface{}, key string) string {
if v, ok := args[key]; ok {
switch x := v.(type) {
case string:
return x
case json.Number:
return x.String()
case float64:
return fmt.Sprintf("%.0f", x)
case int:
return fmt.Sprintf("%d", x)
case int64:
return fmt.Sprintf("%d", x)
case bool:
if x {
return "true"
}
return "false"
default:
return fmt.Sprint(x)
}
}
return ""
}
func argInt(args map[string]interface{}, key string, def int) int {
if v, ok := args[key]; ok {
switch x := v.(type) {
case float64:
return int(x)
case json.Number:
i, _ := x.Int64()
return int(i)
case int:
return x
case int64:
return int(x)
case string:
var i int
if _, err := fmt.Sscanf(x, "%d", &i); err == nil {
return i
}
}
}
return def
}
func argBool(args map[string]interface{}, key string) bool {
if v, ok := args[key]; ok {
switch x := v.(type) {
case bool:
return x
case string:
return x == "true" || x == "1" || x == "yes"
}
}
return false
}
func argMap(args map[string]interface{}, key string) map[string]interface{} {
if v, ok := args[key]; ok {
if m, ok := v.(map[string]interface{}); ok {
return m
}
if s, ok := v.(string); ok && s != "" {
var m map[string]interface{}
if json.Unmarshal([]byte(s), &m) == nil {
return m
}
}
}
return nil
}
func valueString(m map[string]interface{}, key string) string {
if m == nil {
return ""
}
return argString(m, key)
}
func valueInt(m map[string]interface{}, key string) int {
if m == nil {
return 0
}
return argInt(m, key, 0)
}
func content(v interface{}) map[string]interface{} {
out := map[string]interface{}{"content": v}
return out
}
func contentWith(m map[string]interface{}, c string) map[string]interface{} {
m["content"] = c
return m
}
// ---- diag_triage ----
func (p *Plugin) handleTriage(args map[string]interface{}) (interface{}, error) {
exitCode := argInt(args, "exit_code", 0)
signal := argString(args, "signal")
uptimeMS := argInt(args, "uptime_ms", 0)
stillAlive := argBool(args, "still_alive")
reason := argString(args, "crash_reason")
class := "normal_stop"
verdict := "healthy"
var detail []string
if reasons := strings.TrimSpace(reason); reasons != "" {
detail = append(detail, "守护方告知: "+reasons)
}
switch {
case stillAlive:
// 主 agent 进程仍在,但被判定需要检查 → 配置类/可达性问题优先(进程自身正常)
class = "config_unreachable"
verdict = "degraded"
detail = append(detail, "进程存活但健康检测触发,倾向配置/可达性类")
case signal != "":
s := strings.ToUpper(strings.ReplaceAll(signal, "-", ""))
class = "process_death"
verdict = "down"
detail = append(detail, fmt.Sprintf("被信号终止: %s", signal))
if s == "SIGKILL" || s == "KILL" || s == "OOM" || strings.Contains(strings.ToLower(signal), "oom") {
class = "process_starvation"
detail = append(detail, "疑似被强杀/OOM优先怀疑资源或失控")
} else if s == "SIGSEGV" || s == "SIGBUS" || s == "SIGABRT" || s == "SIGFPE" {
detail = append(detail, "疑似崩溃信号segv/abrt配合 diag_log_scan 的 panic/栈签名")
}
case exitCode != 0:
class = "process_death"
verdict = "down"
detail = append(detail, fmt.Sprintf("非零退出码: %d", exitCode))
if exitCode >= 128 {
detail = append(detail, "退出码>=128 通常是 128+signal配合信号判定")
}
default:
detail = append(detail, "退出码0且无信号正常停止")
}
c := map[string]interface{}{
"verdict": verdict,
"class": class,
"exit_code": exitCode,
"signal": signal,
"still_alive": stillAlive,
"uptime_ms": uptimeMS,
"detail": detail,
"recommendation": recForTriage(class),
}
return contentWith(c, describeTriage(c)), nil
}
func recForTriage(class string) string {
switch class {
case "process_death":
return "进程崩溃先查崩溃点diag_log_scan 栈/panic 签名);若无配置改动则重建 worker勿动配置"
case "process_starvation":
return "资源/强杀:检查内存/失控,勿动配置,重建 worker 并限制资源"
case "config_unreachable":
return "配置/可达性:查 LLM 源与网络配置diag_db + diag_delta on /etc必要时还原配置并 ReloadFromConfig"
default:
return "正常情况,无需恢复"
}
}
func describeTriage(c map[string]interface{}) string {
return fmt.Sprintf("判决: %s | 类别: %s | 建议: %s", c["verdict"], c["class"], c["recommendation"])
}
// ---- diag_db ----
type llmSourceData struct {
Name string `json:"name"`
BaseURL string `json:"base_url"`
Model string `json:"model"`
Adapter string `json:"adapter"`
APIKeySet bool `json:"api_key_present"`
Missing []string `json:"missing_fields"`
OK bool `json:"ok"`
}
func (p *Plugin) handleDB(args map[string]interface{}) (interface{}, error) {
dbPath := argString(args, "db_path")
if dbPath == "" {
dbPath = p.cfgPath
}
if dbPath == "" {
return content("无法定位 config.db未配置 data_dir请传入 db_path"), nil
}
// 安全校验db_path 仅允许 data 目录内的 sqlite 文件,防止被用作任意文件探测。
if p.dataDir != "" {
abs, err := filepath.Abs(dbPath)
if err != nil {
return content("db_path 解析失败: " + err.Error()), nil
}
base := filepath.Clean(p.dataDir)
if abs != base && !strings.HasPrefix(abs, base+string(filepath.Separator)) {
return content("db_path 必须位于数据目录内(" + base + ""), nil
}
}
res := map[string]interface{}{
"db_path": dbPath,
"exists": false,
}
if st, err := os.Stat(dbPath); err != nil || st.IsDir() {
res["integrity"] = "absent"
res["sources"] = []map[string]interface{}{}
res["summary"] = "config.db 缺失,属配置损坏类高危信号"
return contentWith(res, "config.db 缺失/不可访问"), nil
}
res["exists"] = true
res["size_bytes"] = func() int64 {
st, _ := os.Stat(dbPath)
if st != nil {
return st.Size()
}
return 0
}()
integrity, errTxt := p.dbIntegrity(dbPath)
res["integrity"] = integrity
if serr, ok := errTxt.(string); ok && serr != "" {
res["integrity_error"] = serr
}
sources, srcErr := p.dbSources(dbPath)
res["sources"] = sources
failed := 0
absent := 0
var missingFields []string
for _, s := range sources {
if !s.OK {
failed++
missingFields = append(missingFields, s.Name+":"+strings.Join(s.Missing, ","))
} else if s.BaseURL == "" {
absent++
}
}
res["source_count"] = len(sources)
res["source_failed"] = failed
res["missing_fields"] = missingFields
verdict := "ok"
summary := "config.db 完整LLM 源解析全部通过"
if integrity != "ok" {
verdict = "fail"
summary = "config.db 完整性校验失败,属配置损坏类,应还原配置快照并 ReloadFromConfig"
} else if failed > 0 {
verdict = "degraded"
summary = fmt.Sprintf("config.db 完整,但 %d 个 LLM 源缺必备字段(%s需修复源配置", failed, strings.Join(missingFields, ";"))
} else if sources == nil && srcErr != "" {
verdict = "unknown"
summary = "config.db 完整但无法解析 LLM 源:" + srcErr
}
res["verdict"] = verdict
res["summary"] = summary
return contentWith(res, summary), nil
}
// dbIntegrity 优先用 sqlite3 CLI 做 PRAGMA integrity_check缺失则用 Settings 兜底 + 头部魔法字节启发式。
func (p *Plugin) dbIntegrity(dbPath string) (string, interface{}) {
bin := p.sqliteBin()
if bin != "" {
out, err := exec.Command(bin, dbPath, "PRAGMA integrity_check;").CombinedOutput()
if err != nil {
return "error", fmt.Sprintf("sqlite3 运行失败: %v(%s)", err, strings.TrimSpace(string(out)))
}
trim := strings.TrimSpace(string(out))
if strings.Contains(trim, "ok") {
return "ok", ""
}
if trim != "" {
return "fail", hemlines(trim, 3)
}
return "unknown", "integrity_check 无输出"
}
// 无 CLI读头部魔法 + 是否 WAL 缺页pgno/心跳不必深析)作轻量启发式
hdr := make([]byte, 16)
f, err := os.Open(dbPath)
if err != nil {
return "error", "无法打开 config.db"
}
_, err = f.Read(hdr)
f.Close()
if err != nil || !strings.HasPrefix(string(hdr), "SQLite format 3\x00") {
return "fail", "非 SQLite 文件头,疑似损坏/截断"
}
return "ok", "" // 头部完好;深度一致性超出无 CLI 能力,标注降级
}
func hemlines(s string, n int) string {
lines := strings.Split(s, "\n")
lines = filterNonEmpty(lines)
if len(lines) > n {
return strings.Join(lines[:n], " | ")
}
return strings.Join(lines, " | ")
}
func filterNonEmpty(lines []string) []string {
var o []string
for _, l := range lines {
if strings.TrimSpace(l) != "" {
o = append(o, strings.TrimSpace(l))
}
}
return o
}
// dbSources 枚举 core.llm.sources.* 并校验必备字段。优先 sqlite3 CLI缺失回退内核 Settings。
func (p *Plugin) dbSources(dbPath string) ([]llmSourceData, string) {
bin := p.sqliteBin()
kv := map[string]string{}
if bin != "" {
out, err := exec.Command(bin,
dbPath,
"SELECT key, value FROM config WHERE key LIKE 'core.llm.sources.%';").CombinedOutput()
if err != nil {
return nil, fmt.Sprintf("sqlite3 查询失败: %v", err)
}
for _, line := range strings.Split(string(out), "\n") {
if idx := strings.IndexByte(line, '|'); idx >= 0 {
kv[line[:idx]] = line[idx+1:]
}
}
} else if p.sdk != nil {
keys, _ := p.sdk.Settings().ListCore("core.llm.sources")
for _, k := range keys {
if v, err := p.sdk.Settings().GetCore(k); err == nil && v != nil {
kv[k] = fmt.Sprint(v)
}
}
} else {
return nil, "既无 sqlite3 也无 Settings 可用"
}
byName := map[string]map[string]string{}
for k, v := range kv {
rest := strings.TrimPrefix(k, "core.llm.sources.")
parts := strings.SplitN(rest, ".", 2)
if len(parts) != 2 {
continue
}
if byName[parts[0]] == nil {
byName[parts[0]] = map[string]string{}
}
byName[parts[0]][parts[1]] = v
}
names := make([]string, 0, len(byName))
for n := range byName {
names = append(names, n)
}
sort.Strings(names)
var out []llmSourceData
for _, n := range names {
fields := byName[n]
var missing []string
for _, req := range []string{"base_url", "model", "adapter"} {
if strings.TrimSpace(fields[req]) == "" {
missing = append(missing, req)
}
}
out = append(out, llmSourceData{
Name: n,
BaseURL: fields["base_url"],
Model: fields["model"],
Adapter: fields["adapter"],
APIKeySet: strings.TrimSpace(fields["api_key"]) != "",
Missing: missing,
OK: len(missing) == 0,
})
}
return out, ""
}
func (p *Plugin) sqliteBin() string {
if p.sdk != nil {
if v, err := p.sdk.Settings().Get("db_check_cmd"); err == nil && v != nil {
if sv, ok := v.(string); ok && sv != "" && sv != "auto" {
if _, err := exec.LookPath(sv); err == nil {
return sv
}
return ""
}
}
}
if _, err := exec.LookPath("sqlite3"); err == nil {
return "sqlite3"
}
return ""
}
// ---- diag_log_scan ----
type sigRule struct {
Category string
Re *regexp.Regexp
}
var sigRules = []sigRule{
{"panic", regexp.MustCompile(`(?i)panic|nil pointer|invalid memory address|runtime error|SIGSEGV|coredump|stack overflow`)},
{"oom", regexp.MustCompile(`(?i)\boom\b|out of memory|memory allocation failed`)},
{"network", regexp.MustCompile(`(?i)no such host|connection refused|connection reset|timeout|unreachable|dns|lookup.*fail`)},
{"provider", regexp.MustCompile(`(?i)provider .* failed|marked unavailable|llm api unreachable|llmfallback|api key|401|403`)},
{"sql", regexp.MustCompile(`(?i)sql: |sqlite|database is locked|disk I/O error|no such table|constraint failed`)},
{"fatal", regexp.MustCompile(`(?i)\bfatal\b|\berror\b|failed`)},
}
func (p *Plugin) handleLogScan(args map[string]interface{}) (interface{}, error) {
logDir := argString(args, "log_dir")
if logDir == "" {
logDir = p.logDir
}
sinceMin := argInt(args, "since_minutes", 0)
maxLines := argInt(args, "max_lines", 200000)
if maxLines <= 0 {
maxLines = 200000
}
cutoff := time.Time{}
if sinceMin > 0 {
cutoff = time.Now().Add(-time.Duration(sinceMin) * time.Minute)
}
entries, err := os.ReadDir(logDir)
if err != nil {
return content(fmt.Sprintf("日志目录不可读: %v", err)), nil
}
// 只扫当前 raw 日志homed_YYYY-MM-DD_HH-MM-SS.log忽略已压缩归档
var files []string
for _, e := range entries {
if e.IsDir() || !strings.HasSuffix(e.Name(), ".log") {
continue
}
files = append(files, filepath.Join(logDir, e.Name()))
}
sort.Strings(files)
counts := map[string]int{}
total := 0
matched := 0
scannedLines := 0
filesRead := 0
for _, f := range files {
if scannedLines >= maxLines {
break
}
data, err := os.ReadFile(f)
if err != nil {
continue
}
filesRead++
for _, line := range strings.Split(string(data), "\n") {
if scannedLines >= maxLines {
break
}
scannedLines++
line = strings.TrimSpace(line)
if line == "" {
continue
}
// 时间窗过滤:行首时间戳形如 2026/08/03 10:03:52
if !cutoff.IsZero() {
ts := parseLogTS(line)
if !ts.IsZero() && ts.Before(cutoff) && (sinceMin > 0) {
continue
}
}
total++
for _, rule := range sigRules {
if rule.Re.MatchString(line) {
counts[rule.Category]++
matched++
break
}
}
}
}
// 排序取主导
type kv struct {
cat string
count int
}
var order []kv
for cat, n := range counts {
order = append(order, kv{cat, n})
}
sort.Slice(order, func(i, j int) bool {
if order[i].count != order[j].count {
return order[i].count > order[j].count
}
return order[i].cat < order[j].cat
})
dominant := ""
if len(order) > 0 {
dominant = order[0].cat
}
conclusion := "无已知错误签名命中"
switch dominant {
case "panic":
conclusion = "主导: panic/崩溃 → 配合栈定位,属进程死亡类"
case "oom":
conclusion = "主导: OOM/内存 → 进程失稳类,检查内存"
case "network", "provider":
conclusion = "主导: 网络/供应商不可达 → 配置或系统网络类"
case "sql":
conclusion = "主导: SQL/数据库错误 → 数据或配置损坏类"
case "fatal":
conclusion = "主导: 常规 error/failed → 需结合 DB/快照进一步定位"
}
res := map[string]interface{}{
"log_dir": logDir,
"files_read": filesRead,
"lines_scanned": scannedLines,
"lines_total_in_window": total,
"lines_matched": matched,
"counts": counts,
"dominant": dominant,
"conclusion": conclusion,
}
return contentWith(res, fmt.Sprintf("%s (命中 %d 行, 主导 %s)", conclusion, matched, dominant)), nil
}
// parseLogTS 解析 homed 时间戳前缀 2026/08/03 10:03:52。
var logTSRe = regexp.MustCompile(`^(\d{4})/(\d{2})/(\d{2}) (\d{2}):(\d{2}):(\d{2})`)
func parseLogTS(line string) time.Time {
m := logTSRe.FindStringSubmatch(line)
if m == nil {
return time.Time{}
}
ts, _ := time.ParseInLocation("2006-01-02 15:04:05",
fmt.Sprintf("%s-%s-%s %s:%s:%s", m[1], m[2], m[3], m[4], m[5], m[6]), time.Local)
return ts
}
// ---- diag_delta ----
type fileEntry struct {
Path string `json:"path"`
Type string `json:"type"` // created / modified / deleted
Size int64 `json:"size"`
NewHash string `json:"new_hash,omitempty"`
OldHash string `json:"old_hash,omitempty"`
}
func (p *Plugin) handleDelta(args map[string]interface{}) (interface{}, error) {
baseline := argString(args, "baseline_dir")
current := argString(args, "current_dir")
pattern := argString(args, "pattern")
maxItems := argInt(args, "max_items", 500)
if maxItems <= 0 {
maxItems = 500
}
if baseline == "" || current == "" {
return contentWith(map[string]interface{}{
"error": "baseline_dir 与 current_dir 均必填",
}, "缺少基线或现状目录:请先准备 last-good 快照解包目录"), nil
}
baseMissing := !dirExists(baseline)
currMissing := !dirExists(current)
if baseMissing {
return contentWith(map[string]interface{}{
"baseline_dir": baseline,
"current_dir": current,
"baseline_exists": false,
"summary": "基线不存在,无法差分(需先建立快照基线)",
}, "基线不存在,无法差分"), nil
}
if currMissing {
return contentWith(map[string]interface{}{
"baseline_dir": baseline,
"current_dir": current,
"current_exists": false,
"summary": "现状目录不存在",
}, "现状目录不存在"), nil
}
baseMap := walkHashes(baseline)
currMap := walkHashes(current)
var files []fileEntry
seen := map[string]bool{}
for path, ch := range currMap {
seen[path] = true
if pattern != "" && !strings.Contains(path, pattern) {
continue
}
if bh, ok := baseMap[path]; ok {
if bh.hash != ch.hash {
files = append(files, fileEntry{Path: path, Type: "modified", Size: ch.size, OldHash: bh.hash, NewHash: ch.hash})
}
} else {
files = append(files, fileEntry{Path: path, Type: "created", Size: ch.size, NewHash: ch.hash})
}
}
for path, bh := range baseMap {
if !seen[path] && (pattern == "" || strings.Contains(path, pattern)) {
files = append(files, fileEntry{Path: path, Type: "deleted", Size: bh.size, OldHash: bh.hash})
}
}
sort.Slice(files, func(i, j int) bool { return files[i].Path < files[j].Path })
summary := map[string]int{"created": 0, "modified": 0, "deleted": 0}
for _, f := range files {
summary[f.Type]++
}
shown := files
if len(shown) > maxItems {
shown = shown[:maxItems]
}
res := map[string]interface{}{
"baseline_dir": baseline,
"current_dir": current,
"summary": summary,
"total_diff": len(files),
"files": shown,
}
return contentWith(res, fmt.Sprintf("diff: %+v", res["summary"])), nil
}
type hashEnt struct {
hash string
size int64
}
func walkHashes(root string) map[string]hashEnt {
out := map[string]hashEnt{}
filepath.Walk(root, func(path string, info os.FileInfo, err error) error {
if err != nil || info.IsDir() {
return nil
}
rel, _ := filepath.Rel(root, path)
data, err := os.ReadFile(path)
if err != nil {
return nil
}
h := sha256.Sum256(data)
out[rel] = hashEnt{hash: hex.EncodeToString(h[:]), size: int64(len(data))}
return nil
})
return out
}
func dirExists(p string) bool {
st, err := os.Stat(p)
return err == nil && st.IsDir()
}
// ---- diag_loc ----
func (p *Plugin) handleLoc(args map[string]interface{}) (interface{}, error) {
triage := argMap(args, "triage")
db := argMap(args, "db")
logScan := argMap(args, "log_scan")
delta := argMap(args, "delta")
type hyp struct {
Cause string `json:"cause"`
Confidence int `json:"confidence"` // 0-100 因果强度
Evidence []string `json:"evidence"`
Recommend string `json:"recommendation"`
}
var hyps []hyp
tClass := valueString(triage, "class")
tVerdict := valueString(triage, "verdict")
dbVerdict := valueString(db, "verdict")
dom := valueString(logScan, "dominant")
deltaSummary := map[string]int{}
if delta != nil {
if s, ok := delta["summary"].(map[string]interface{}); ok {
for k, v := range s {
switch n := v.(type) {
case float64:
deltaSummary[k] = int(n)
case int:
deltaSummary[k] = n
case int64:
deltaSummary[k] = int(n)
}
}
}
}
dCreated := deltaSummary["created"]
dModified := deltaSummary["modified"]
dDeleted := deltaSummary["deleted"]
dTotal := dCreated + dModified + dDeleted
evidence := []string{}
if tVerdict != "" {
evidence = append(evidence, "triage="+tVerdict+"("+tClass+")")
}
if dbVerdict != "" {
evidence = append(evidence, "db="+dbVerdict)
}
if dom != "" {
evidence = append(evidence, "log_dominant="+dom)
}
if dTotal > 0 {
evidence = append(evidence, fmt.Sprintf("delta=%d 改动(%d改/%d增/%d删)", dTotal, dModified, dCreated, dDeleted))
} else {
evidence = append(evidence, "delta=无改动")
}
// 1) 进程失稳panic 主导时走更具体的 code_panic_loop 分支)
if (tClass == "process_death" || tClass == "process_starvation") &&
dbVerdict != "fail" && dbVerdict != "degraded" && dTotal == 0 &&
dom != "panic" {
hyps = append(hyps, hyp{
Cause: "process_instability",
Confidence: 75,
Evidence: in(evidence, "triage=down"),
Recommend: "重建 worker不动配置db 完好、无文件改动)",
})
}
// 2) 配置损坏
if dbVerdict == "fail" || dbVerdict == "degraded" {
hyps = append(hyps, hyp{
Cause: "config_corruption",
Confidence: 90,
Evidence: in(evidence, "db="+dbVerdict),
Recommend: "还原 core.llm.sources 配置快照 → ReloadFromConfig → 拉起主 agent",
})
}
// 3) 系统网络
if (dom == "network" || dom == "provider") && (dTotal > 0) {
hyps = append(hyps, hyp{
Cause: "system_network",
Confidence: 80,
Evidence: in(evidence, "log_dominant="+dom, "delta>0"),
Recommend: "还原 DNS/proxy/host 相关系统网络配置 → 重载主 agent",
})
}
// 4) 纯日志栈崩溃db 完好、无 delta
if tClass == "process_death" && dbVerdict == "ok" && dTotal == 0 && dom == "panic" {
hyps = append(hyps, hyp{
Cause: "code_panic_loop",
Confidence: 70,
Evidence: in(evidence, "log_dominant=panic", "db=ok", "delta=无改动"),
Recommend: "定位 panic 栈来源repeat+ 检查是否插件引起,必要时禁用对应插件后重建 worker",
})
}
// 未知/混合
if len(hyps) == 0 {
hyps = append(hyps, hyp{
Cause: "unknown_mixed",
Confidence: 20,
Evidence: evidence,
Recommend: "确定性命中不足,放开 webfetch/知识库,用 rescue 源 做最小 LLM 推理(依据 diag_* 结论摘要)",
})
}
sort.Slice(hyps, func(i, j int) bool { return hyps[i].Confidence > hyps[j].Confidence })
res := map[string]interface{}{
"evidence": evidence,
"ranked_hypotheses": hyps,
"final_recommendation": hyps[0].Recommend,
}
// 落盘 + 知识库回流(同类崩溃下次直接命中)
if argBool(args, "persist") {
p.persistConclusion(res, hyps[0].Cause, hyps[0].Recommend)
}
return contentWith(res, "定位: "+hyps[0].Cause+" | 建议: "+hyps[0].Recommend), nil
}
// persistConclusion 把定位结论写 recovery_kb/diag_<ts>.json并经知识库回流失败不阻塞
func (p *Plugin) persistConclusion(res map[string]interface{}, cause, recommend string) {
ts := time.Now()
entry := map[string]interface{}{
"ts": ts.Format(time.RFC3339),
"cause": cause,
"recommendation": recommend,
"evidence": valueFrom(res, "evidence"),
"ranked_hypotheses": res["ranked_hypotheses"],
"final_recommendation": recommend,
"tool": "diag_loc",
}
raw, _ := json.MarshalIndent(entry, "", " ")
dir := p.recoveryKBDir()
if dir != "" {
if err := os.MkdirAll(dir, 0755); err == nil {
path := filepath.Join(dir, fmt.Sprintf("diag_%s.json", ts.Format("2006-01-02_15-04-05")))
if err := os.WriteFile(path, raw, 0644); err == nil {
log.Printf("[%s] conclusion persisted to %s", p.name, path)
} else {
log.Printf("[%s] persist file %s: %v", p.name, path, err)
}
}
}
if p.sdk != nil && p.sdk.Knowledge() != nil {
kName := fmt.Sprintf("diag:%s:%s", cause, ts.Format("2006-01-02T15-04"))
content := fmt.Sprintf("恢复诊断结论(%s): %s。建议: %s。命中条件可复用。", ts.Format("2006-01-02 15:04:05"), cause, recommend)
if err := p.sdk.Knowledge().Add(kName, content); err != nil {
log.Printf("[%s] knowledge add %s: %v", p.name, kName, err)
}
}
}
func valueFrom(m map[string]interface{}, k string) interface{} {
if m == nil {
return nil
}
return m[k]
}
// recoveryKBDir 返回结论落盘目录,可配置,缺省 <data_dir>/recovery_kb。
func (p *Plugin) recoveryKBDir() string {
if p.sdk != nil {
if v, err := p.sdk.Settings().Get("recovery_kb_dir"); err == nil && v != nil {
if sv, ok := v.(string); ok && sv != "" {
return sv
}
}
}
if p.dataDir != "" {
return filepath.Join(p.dataDir, "recovery_kb")
}
return ""
}
// in 过滤 slice保留同时满足 items 中条件(简单子串匹配)的元素。
func in(src []string, items ...string) []string {
var o []string
for _, it := range items {
for _, s := range src {
if s == it {
o = append(o, it)
break
}
}
}
return o
}
var _ = json.Marshal