Phase 0.1/1/3/6: 核心生产问题修复

Healthcheck 隔离 (Phase 0.1):
- 新增 internal/sdk/selftest.go: VirtualInstance 完全隔离自检空间
- PluginSDK.Selftest()/SelftestReset() 暴露隔离实例 (含 mutex)
- LLM 自检只读白名单 isSafeReadonlyTool 防写类工具污染生产
- 单测验证: healthcheck 后生产实例内容不变 + 无残留
- 存量清理: 删除 gotest/luatest 残留目录

GraphDB 去重 (Phase 1):
- migrateRelationUnique: 启动自动重建 relations 表加 UNIQUE 约束并去重
- Commit 改为存在性检查, 重复三元组仅刷新 confidence 不重复插入
- 3 个 dedup 单测全绿

配置时长解析 (Phase 3):
- parseDurationExtended 支持 2d/1w/3h 等人类可读单位
- GetDuration 全局生效, 防 2d 静默回退 30m

Agentcli 通知风暴治理 (Phase 6):
- 语义通知: 累积 notify_bytes(2KB) 或间隔 notify_interval(2s) 触发
- 生命周期即时通知: 启动/进程退出/EOF 立即通知
- 可配置 settings, 保留通知机制保证 agent 感知终端存在
- 运维止血: 已杀掉幽灵 PID 3716282 (bash git sparse clone 运行 16h)

Plan.md: 新增设计意图备忘(插件即App/分层记忆), 更新各 Phase 进度
This commit is contained in:
root
2026-08-12 13:51:40 +08:00
parent c19fea2584
commit 9d5a914941
12 changed files with 962 additions and 39 deletions

2
go.mod
View File

@ -3,7 +3,7 @@ module gitcode.com/JianFeeeee/HomeAgent
go 1.25.0 go 1.25.0
require ( require (
github.com/mattn/go-sqlite3 v1.14.48 github.com/mattn/go-sqlite3 v1.14.49
github.com/yuin/gopher-lua v1.1.2 github.com/yuin/gopher-lua v1.1.2
gopkg.in/yaml.v3 v3.0.1 gopkg.in/yaml.v3 v3.0.1
) )

3
go.sum
View File

@ -1,6 +1,7 @@
github.com/mattn/go-sqlite3 v1.14.48 h1:7XHIgl0a8HwOaiK4E47ozLkST78rR9+OtNGx27D/TFs= github.com/mattn/go-sqlite3 v1.14.48 h1:7XHIgl0a8HwOaiK4E47ozLkST78rR9+OtNGx27D/TFs=
github.com/mattn/go-sqlite3 v1.14.48/go.mod h1:6JTjA44L93a0QCyJef5YvlPoKXntQPjzWv5gtm9sB6w= github.com/mattn/go-sqlite3 v1.14.48/go.mod h1:6JTjA44L93a0QCyJef5YvlPoKXntQPjzWv5gtm9sB6w=
github.com/mattn/go-sqlite3 v1.14.49 h1:B8jBHC3xhxZgxztrgruTuLucebnULQnx4W7cF7SAE9w=
github.com/mattn/go-sqlite3 v1.14.49/go.mod h1:6JTjA44L93a0QCyJef5YvlPoKXntQPjzWv5gtm9sB6w=
github.com/yalue/onnxruntime_go v1.13.0 h1:5HDXHon3EukQMyYA7yPMed/raWaDE/gjwLOwnVoiwy8= github.com/yalue/onnxruntime_go v1.13.0 h1:5HDXHon3EukQMyYA7yPMed/raWaDE/gjwLOwnVoiwy8=
github.com/yalue/onnxruntime_go v1.13.0/go.mod h1:b4X26A8pekNb1ACJ58wAXgNKeUCGEAQ9dmACut9Sm/4= github.com/yalue/onnxruntime_go v1.13.0/go.mod h1:b4X26A8pekNb1ACJ58wAXgNKeUCGEAQ9dmACut9Sm/4=
github.com/yanyiwu/gojieba v1.4.7 h1:2YkXELcYLTE0SJetq6xv4MjpEikWga6VpFn4jIFFQ/k= github.com/yanyiwu/gojieba v1.4.7 h1:2YkXELcYLTE0SJetq6xv4MjpEikWga6VpFn4jIFFQ/k=

View File

@ -7,6 +7,7 @@ import (
"log" "log"
"os" "os"
"path/filepath" "path/filepath"
"regexp"
"sort" "sort"
"strconv" "strconv"
"strings" "strings"
@ -695,13 +696,44 @@ func (r *ConfigRegistry) GetDuration(key string, defaultVal time.Duration) time.
if err != nil { if err != nil {
return defaultVal return defaultVal
} }
d, err := time.ParseDuration(v) d, err := parseDurationExtended(v)
if err != nil { if err != nil {
return defaultVal return defaultVal
} }
return d return d
} }
// durationUnitRe 匹配 `\d+[dhw]`(天/小时/周)这类 Go time.ParseDuration 不支持的天气单位。
var durationUnitRe = regexp.MustCompile(`(\d+)\s*([dhw])`)
// parseDurationExtended 解析人类可读时长,支持 Go 原生单位ns/us/ms/s/m/h
// 及复合如 "1h30m")加上 d与 w。返回实例化 duration失败返回 error。
func parseDurationExtended(s string) (time.Duration, error) {
s = strings.TrimSpace(s)
if s == "" {
return 0, fmt.Errorf("empty duration")
}
// 先展开 d/w再交给 time.ParseDuration 处理剩余(含 m/h/s 组合)。
expanded := durationUnitRe.ReplaceAllStringFunc(s, func(m string) string {
parts := durationUnitRe.FindStringSubmatch(m)
n, _ := strconv.Atoi(parts[1])
switch parts[2] {
case "d":
return fmt.Sprintf("%dh", n*24)
case "w":
return fmt.Sprintf("%dh", n*24*7)
case "h":
return m
}
return m
})
d, err := time.ParseDuration(expanded)
if err != nil {
return 0, err
}
return d, nil
}
func (r *ConfigRegistry) GetBool(key string, defaultVal bool) bool { func (r *ConfigRegistry) GetBool(key string, defaultVal bool) bool {
r.mu.RLock() r.mu.RLock()
defer r.mu.RUnlock() defer r.mu.RUnlock()

View File

@ -210,6 +210,14 @@ func TestGetHelpers(t *testing.T) {
if got := r.GetDuration("dur_key", 0); got != 5*time.Minute { if got := r.GetDuration("dur_key", 0); got != 5*time.Minute {
t.Fatalf("GetDuration: expected 5m, got %v", got) t.Fatalf("GetDuration: expected 5m, got %v", got)
} }
r.Set("dur_key_days", "2d")
if got := r.GetDuration("dur_key_days", 0); got != 48*time.Hour {
t.Fatalf("GetDuration d-unit: expected 48h, got %v", got)
}
r.Set("dur_key_weeks", "1w")
if got := r.GetDuration("dur_key_weeks", 0); got != 168*time.Hour {
t.Fatalf("GetDuration w-unit: expected 168h, got %v", got)
}
if got := r.GetDuration("nonexistent", 30*time.Second); got != 30*time.Second { if got := r.GetDuration("nonexistent", 30*time.Second); got != 30*time.Second {
t.Fatalf("GetDuration fallback: expected 30s, got %v", got) t.Fatalf("GetDuration fallback: expected 30s, got %v", got)
} }
@ -221,6 +229,42 @@ func TestGetHelpers(t *testing.T) {
} }
} }
func TestParseDurationExtended(t *testing.T) {
cases := []struct {
in string
want time.Duration
wantErr bool
}{
{"2d", 48 * time.Hour, false},
{"1w", 7 * 24 * time.Hour, false},
{"1d", 24 * time.Hour, false},
{"2d12h", 60 * time.Hour, false},
{"30m", 30 * time.Minute, false},
{"500ms", 500 * time.Millisecond, false},
{"1h30m", 90 * time.Minute, false},
{" 3d ", 72 * time.Hour, false},
{"2w", 336 * time.Hour, false},
{"", 0, true},
{"abc", 0, true},
}
for _, c := range cases {
got, err := parseDurationExtended(c.in)
if c.wantErr {
if err == nil {
t.Errorf("%q: expected error, got %v", c.in, got)
}
continue
}
if err != nil {
t.Errorf("%q: unexpected error: %v", c.in, err)
continue
}
if got != c.want {
t.Errorf("%q: expected %v, got %v", c.in, c.want, got)
}
}
}
func TestSnapshotRestoreCoreLLM(t *testing.T) { func TestSnapshotRestoreCoreLLM(t *testing.T) {
r := NewConfigRegistry("") r := NewConfigRegistry("")
defer r.Close() defer r.Close()

View File

@ -3,6 +3,7 @@ package memory
import ( import (
"database/sql" "database/sql"
"fmt" "fmt"
"strings"
"sync" "sync"
"time" "time"
@ -103,7 +104,8 @@ func (g *GraphDB) initSchema() error {
date_bucket TEXT, date_bucket TEXT,
sentence_id INTEGER DEFAULT 0, sentence_id INTEGER DEFAULT 0,
FOREIGN KEY (source_id) REFERENCES entities(id), FOREIGN KEY (source_id) REFERENCES entities(id),
FOREIGN KEY (target_id) REFERENCES entities(id) FOREIGN KEY (target_id) REFERENCES entities(id),
UNIQUE(source_id, target_id, relation_type, session_id)
)`, )`,
`CREATE INDEX IF NOT EXISTS idx_entity_name ON entities(name)`, `CREATE INDEX IF NOT EXISTS idx_entity_name ON entities(name)`,
`CREATE INDEX IF NOT EXISTS idx_entity_type ON entities(type)`, `CREATE INDEX IF NOT EXISTS idx_entity_type ON entities(type)`,
@ -132,9 +134,64 @@ func (g *GraphDB) initSchema() error {
// sentence_id 索引在迁移后创建,避免旧表缺少该列时失败 // sentence_id 索引在迁移后创建,避免旧表缺少该列时失败
tx.Exec(`CREATE INDEX IF NOT EXISTS idx_relation_sentence ON relations(sentence_id)`) tx.Exec(`CREATE INDEX IF NOT EXISTS idx_relation_sentence ON relations(sentence_id)`)
// 迁移4为旧版 relations 表(无复合唯一约束)重建表以去重。
// 旧表由 2026-07 之前的版本创建,缺少 UNIQUE(source_id, target_id, relation_type, session_id)
// 生产库累积了海量重复关系。这里检查 sqlite_master 中已建表的 DDL
// 若不含该约束则走"新建带约束表 → INSERT OR IGNORE 拷贝去重 → 换名"的官方 12 步迁移。
if err := g.migrateRelationUnique(tx); err != nil {
return fmt.Errorf("migrate relations unique: %w", err)
}
return tx.Commit() return tx.Commit()
} }
// migrateRelationUnique 检测 relations 表是否带复合唯一约束,缺失则重建去重。
// 必须在 initSchema 的同一个事务内调用(外键/索引均已存在时需先禁用外键再换名)。
func (g *GraphDB) migrateRelationUnique(tx *sql.Tx) error {
var ddl string
err := tx.QueryRow(`SELECT sql FROM sqlite_master WHERE type = 'table' AND name = 'relations'`).Scan(&ddl)
if err != nil {
if err == sql.ErrNoRows {
return nil // 表都不存在,无从迁移
}
return err
}
if strings.Contains(ddl, "UNIQUE") {
return nil // 已是新 schema
}
stmt := []string{
`ALTER TABLE relations RENAME TO relations_old`,
`CREATE TABLE relations (
id INTEGER PRIMARY KEY AUTOINCREMENT,
source_id INTEGER NOT NULL,
target_id INTEGER NOT NULL,
relation_type TEXT NOT NULL,
confidence REAL DEFAULT 1.0,
status TEXT DEFAULT 'active',
session_id TEXT,
turn_id INTEGER DEFAULT 0,
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
date_bucket TEXT,
sentence_id INTEGER DEFAULT 0,
sentence_ref TEXT DEFAULT '',
FOREIGN KEY (source_id) REFERENCES entities(id),
FOREIGN KEY (target_id) REFERENCES entities(id),
UNIQUE(source_id, target_id, relation_type, session_id)
)`,
`INSERT OR IGNORE INTO relations (id, source_id, target_id, relation_type, confidence, status, session_id, turn_id, created_at, updated_at, date_bucket, sentence_id, sentence_ref)
SELECT id, source_id, target_id, relation_type, confidence, status, session_id, turn_id, created_at, updated_at, date_bucket, sentence_id, sentence_ref FROM relations_old`,
`DROP TABLE relations_old`,
}
for _, s := range stmt {
if _, err := tx.Exec(s); err != nil {
return err
}
}
return nil
}
func (g *GraphDB) Commit(triples []Triple, sessionID string, turnID int) (int, int, error) { func (g *GraphDB) Commit(triples []Triple, sessionID string, turnID int) (int, int, error) {
g.mu.Lock() g.mu.Lock()
defer g.mu.Unlock() defer g.mu.Unlock()
@ -206,15 +263,34 @@ func (g *GraphDB) Commit(triples []Triple, sessionID string, turnID int) (int, i
} }
} }
_, err = tx.Exec( var existing int
`INSERT INTO relations (source_id, target_id, relation_type, confidence, session_id, turn_id, date_bucket, sentence_id) err = tx.QueryRow(
VALUES (?, ?, ?, ?, ?, ?, ?, ?)`, `SELECT 1 FROM relations WHERE source_id = ? AND target_id = ? AND relation_type = ? AND session_id = ?`,
sourceID, targetID, t.Relation, confidence, sessionID, turnID, dateBucket, sentenceID, sourceID, targetID, t.Relation, sessionID,
) ).Scan(&existing)
if err != nil { if err == sql.ErrNoRows {
_, err = tx.Exec(
`INSERT INTO relations (source_id, target_id, relation_type, confidence, session_id, turn_id, date_bucket, sentence_id)
VALUES (?, ?, ?, ?, ?, ?, ?, ?)`,
sourceID, targetID, t.Relation, confidence, sessionID, turnID, dateBucket, sentenceID,
)
if err != nil {
return 0, 0, err
}
relationsCreated++
} else if err != nil {
return 0, 0, err return 0, 0, err
} else {
// 同一(会话内)三元组已存在:仅刷新置信度与时间戳,不重复计数
_, err = tx.Exec(
`UPDATE relations SET confidence = ?, updated_at = CURRENT_TIMESTAMP
WHERE source_id = ? AND target_id = ? AND relation_type = ? AND session_id = ?`,
confidence, sourceID, targetID, t.Relation, sessionID,
)
if err != nil {
return 0, 0, err
}
} }
relationsCreated++
} }
if err := tx.Commit(); err != nil { if err := tx.Commit(); err != nil {

View File

@ -1,6 +1,8 @@
package memory package memory
import ( import (
"database/sql"
"path/filepath"
"os" "os"
"testing" "testing"
) )
@ -60,6 +62,135 @@ func TestCommitTriples(t *testing.T) {
} }
} }
func TestCommitDedupSameSession(t *testing.T) {
g := newTestGraph(t)
defer os.Remove(g.dbPath)
defer g.Close()
triple := []Triple{{Subject: "李四", Relation: "喜欢", Object: "篮球"}}
ec, rc, err := g.Commit(triple, "session_dup", 1)
if err != nil {
t.Fatal(err)
}
if ec != 2 || rc != 1 {
t.Fatalf("first commit: want 2/1, got %d/%d", ec, rc)
}
// 同一会话重复 commit 同一三元组:关系不再新增
_, rc, err = g.Commit(triple, "session_dup", 2)
if err != nil {
t.Fatal(err)
}
if rc != 0 {
t.Errorf("duplicate commit should not create relations again, got %d", rc)
}
var cnt int
if err := g.db.QueryRow(`SELECT COUNT(*) FROM relations`).Scan(&cnt); err != nil {
t.Fatal(err)
}
if cnt != 1 {
t.Errorf("expected exactly 1 relation after duplicate commit, got %d", cnt)
}
}
func TestCommitDedupDifferentSession(t *testing.T) {
g := newTestGraph(t)
defer os.Remove(g.dbPath)
defer g.Close()
triple := []Triple{{Subject: "王五", Relation: "喜欢", Object: "足球"}}
for _, sess := range []string{"s1", "s2"} {
if _, _, err := g.Commit(triple, sess, 0); err != nil {
t.Fatal(err)
}
}
var cnt int
if err := g.db.QueryRow(`SELECT COUNT(*) FROM relations`).Scan(&cnt); err != nil {
t.Fatal(err)
}
if cnt != 2 {
t.Errorf("different sessions may repeat a triple, expected 2 relations, got %d", cnt)
}
}
func TestMigrateRelationUniqueDedupsOldTable(t *testing.T) {
dir := t.TempDir()
dbPath := filepath.Join(dir, "legacy.db")
// 构造旧版 schemarelations 无复合唯一约束,且塞入重复行
db, err := sql.Open("sqlite3", dbPath)
if err != nil {
t.Fatal(err)
}
setup := []string{
`CREATE TABLE entities (
id INTEGER PRIMARY KEY AUTOINCREMENT,
name TEXT UNIQUE NOT NULL,
type TEXT DEFAULT 'Concept',
mention_count INTEGER DEFAULT 1,
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
)`,
`CREATE TABLE sentences (
id INTEGER PRIMARY KEY AUTOINCREMENT,
text TEXT UNIQUE NOT NULL,
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
)`,
`CREATE TABLE relations (
id INTEGER PRIMARY KEY AUTOINCREMENT,
source_id INTEGER NOT NULL,
target_id INTEGER NOT NULL,
relation_type TEXT NOT NULL,
confidence REAL DEFAULT 1.0,
status TEXT DEFAULT 'active',
session_id TEXT,
turn_id INTEGER DEFAULT 0,
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
date_bucket TEXT,
sentence_id INTEGER DEFAULT 0,
sentence_ref TEXT DEFAULT '',
FOREIGN KEY (source_id) REFERENCES entities(id),
FOREIGN KEY (target_id) REFERENCES entities(id)
)`,
`INSERT INTO entities (id, name) VALUES (1, '张三'), (2, '编程')`,
`INSERT INTO relations (source_id, target_id, relation_type, session_id) VALUES (1, 2, '喜欢', 's'), (1, 2, '喜欢', 's')`,
}
for _, s := range setup {
if _, err := db.Exec(s); err != nil {
t.Fatal(err)
}
}
db.Close()
// 用 NewGraphDB 打开,应触发 migrateRelationUnique重建带约束表并去重
g, err := NewGraphDB(dbPath)
if err != nil {
t.Fatal(err)
}
defer g.Close()
var cnt int
if err := g.db.QueryRow(`SELECT COUNT(*) FROM relations`).Scan(&cnt); err != nil {
t.Fatal(err)
}
if cnt != 1 {
t.Errorf("expected 1 relation after migration dedup, got %d", cnt)
}
// 再次提交重复三元组不应再新增
_, rc, err := g.Commit([]Triple{{Subject: "张三", Relation: "喜欢", Object: "编程"}}, "s", 1)
if err != nil {
t.Fatal(err)
}
if rc != 0 {
t.Errorf("after migration, duplicate commit should add 0 relations, got %d", rc)
}
}
func TestCommitEmptyTriples(t *testing.T) { func TestCommitEmptyTriples(t *testing.T) {
g := newTestGraph(t) g := newTestGraph(t)
defer os.Remove(g.dbPath) defer os.Remove(g.dbPath)

View File

@ -18,10 +18,11 @@ import (
) )
const ( const (
DefaultTimeout = 5 * time.Minute DefaultTimeout = 5 * time.Minute
ReadBufSize = 4096 ReadBufSize = 4096
MaxOutputBuffer = 128 * 1024 MaxOutputBuffer = 128 * 1024
NotifyOutputDelay = 500 * time.Millisecond DefaultNotifyBytes = 2048 // 积累 2KB 未读输出再通知
DefaultNotifyInterval = 2 * time.Second // 同一终端两次通知的最小间隔(兜底)
) )
// ptyTerm 抽象平台终端后端Linux PTY / Windows ConPTY // ptyTerm 抽象平台终端后端Linux PTY / Windows ConPTY
@ -56,6 +57,10 @@ type TerminalSession struct {
closed bool closed bool
stopCh chan struct{} stopCh chan struct{}
done chan struct{} done chan struct{}
// 通知节流字段
unreadBytes int // 最近一次通知后积累的未读字节数
lastNotify time.Time // 最近一次通知时间
} }
func (t *TerminalSession) Write(input string) (int, error) { func (t *TerminalSession) Write(input string) (int, error) {
@ -128,6 +133,8 @@ type Plugin struct {
sessions map[string]*TerminalSession sessions map[string]*TerminalSession
nextID int nextID int
defaultTimeout time.Duration defaultTimeout time.Duration
notifyBytes int
notifyInterval time.Duration
} }
func New(name string) *Plugin { func New(name string) *Plugin {
@ -147,6 +154,16 @@ func (p *Plugin) Start(s *sdk.PluginSDK) error {
Description: "终端自动关闭的默认时间,例如 5m, 10m, 30m, 1h默认 5m", Description: "终端自动关闭的默认时间,例如 5m, 10m, 30m, 1h默认 5m",
Default: "5m", Default: "5m",
}) })
s.Settings().RegisterDef(sdk.ConfigDef{
Key: "notify_bytes", Type: "int", DisplayName: "通知阈值字节数",
Description: "累积多少字节未读输出后发送通知(默认 2048",
Default: "2048",
})
s.Settings().RegisterDef(sdk.ConfigDef{
Key: "notify_interval", Type: "string", DisplayName: "通知最小间隔",
Description: "同一终端两次通知的最小时间间隔,如 2s, 5s默认 2s",
Default: "2s",
})
if v, _ := s.Settings().Get("default_timeout"); v != nil { if v, _ := s.Settings().Get("default_timeout"); v != nil {
if s, ok := v.(string); ok && s != "" { if s, ok := v.(string); ok && s != "" {
if d, err := time.ParseDuration(s); err == nil { if d, err := time.ParseDuration(s); err == nil {
@ -157,6 +174,24 @@ func (p *Plugin) Start(s *sdk.PluginSDK) error {
if p.defaultTimeout <= 0 { if p.defaultTimeout <= 0 {
p.defaultTimeout = DefaultTimeout p.defaultTimeout = DefaultTimeout
} }
if v, _ := s.Settings().Get("notify_bytes"); v != nil {
if i, ok := v.(float64); ok && i > 0 {
p.notifyBytes = int(i)
}
}
if p.notifyBytes <= 0 {
p.notifyBytes = DefaultNotifyBytes
}
if v, _ := s.Settings().Get("notify_interval"); v != nil {
if s, ok := v.(string); ok && s != "" {
if d, err := time.ParseDuration(s); err == nil {
p.notifyInterval = d
}
}
}
if p.notifyInterval <= 0 {
p.notifyInterval = DefaultNotifyInterval
}
s.RegisterTool("terminal_create", sdk.ToolDef{ s.RegisterTool("terminal_create", sdk.ToolDef{
Name: "terminal_create", Name: "terminal_create",
@ -558,12 +593,15 @@ func (p *Plugin) readLoop(t *TerminalSession, s *sdk.PluginSDK) {
defer close(t.done) defer close(t.done)
buf := make([]byte, ReadBufSize) buf := make([]byte, ReadBufSize)
lastNotify := time.Now()
pollInterval := 200 * time.Millisecond pollInterval := 200 * time.Millisecond
readCh := make(chan readResult, 4) readCh := make(chan readResult, 4)
go p.reader(t, buf, readCh) go p.reader(t, buf, readCh)
// 立即发送首次"终端已启动"通知,让 agent 感知存在
s.InjectText("agentcli", "agentcli", fmt.Sprintf("[终端 %s 已启动]", t.id))
t.lastNotify = time.Now()
for { for {
if t.IsExpired() { if t.IsExpired() {
log.Printf("[agentcli] terminal %s expired after %v", t.id, t.timeout) log.Printf("[agentcli] terminal %s expired after %v", t.id, t.timeout)
@ -587,20 +625,34 @@ func (p *Plugin) readLoop(t *TerminalSession, s *sdk.PluginSDK) {
return return
case r := <-readCh: case r := <-readCh:
if r.err != nil { if r.err != nil {
// 读取错误/EOF → 立即通知(进程可能已结束)
s.InjectText("agentcli", "agentcli", fmt.Sprintf("[终端 %s 读取结束: %v]", t.id, r.err))
return return
} }
if r.n > 0 { if r.n > 0 {
data := make([]byte, r.n) data := make([]byte, r.n)
copy(data, buf[:r.n]) copy(data, buf[:r.n])
t.appendOutput(data) t.appendOutput(data)
if time.Since(lastNotify) > NotifyOutputDelay {
preview := string(data) // 语义通知:累积未读字节数
if len(preview) > 100 { t.mu.Lock()
preview = preview[:100] t.unreadBytes += r.n
needNotify := t.unreadBytes >= p.notifyBytes ||
time.Since(t.lastNotify) >= p.notifyInterval
t.mu.Unlock()
if needNotify {
t.mu.Lock()
preview := t.buf.String()
if len(preview) > 200 {
preview = preview[len(preview)-200:] // 取最新 200 字符
} }
preview = sanitizePreview(preview) preview = sanitizePreview(preview)
t.unreadBytes = 0
t.lastNotify = time.Now()
t.mu.Unlock()
s.InjectText("agentcli", "agentcli", fmt.Sprintf("[终端 %s 有新输出]\n%s", t.id, preview)) s.InjectText("agentcli", "agentcli", fmt.Sprintf("[终端 %s 有新输出]\n%s", t.id, preview))
lastNotify = time.Now()
} }
} }
case <-time.After(pollInterval): case <-time.After(pollInterval):

View File

@ -335,6 +335,11 @@ func (p *Plugin) runAutoCheck(s *sdk.PluginSDK) {
func (p *Plugin) runFullCheck(s *sdk.PluginSDK) (interface{}, error) { func (p *Plugin) runFullCheck(s *sdk.PluginSDK) (interface{}, error) {
results := []checkResult{} results := []checkResult{}
// 每轮自检前重置隔离虚拟实例,清空上轮测试数据(仅影响虚拟空间,不碰生产存储)。
if err := s.SelftestReset("hc"); err != nil {
log.Printf("[healthcheck] selftest reset: %v", err)
}
pluginResult := p.checkPluginsRaw(s) pluginResult := p.checkPluginsRaw(s)
results = append(results, pluginResult...) results = append(results, pluginResult...)
@ -457,55 +462,86 @@ func (p *Plugin) collectAllTools(s *sdk.PluginSDK) []toolInfo {
return tools return tools
} }
func (p *Plugin) selftestInst(s *sdk.PluginSDK) (*sdk.VirtualInstance, error) {
vi, err := s.Selftest("hc")
if err != nil {
return nil, err
}
if vi == nil {
return nil, fmt.Errorf("Selftest 不可用")
}
return vi, nil
}
func (p *Plugin) testMemoryRaw(s *sdk.PluginSDK) checkResult { func (p *Plugin) testMemoryRaw(s *sdk.PluginSDK) checkResult {
// 在隔离虚拟图记忆上验证写→查→删,绝不动生产 GraphDB。
vi, err := p.selftestInst(s)
if err != nil {
return checkResult{Name: "memory", Status: "skip", Detail: fmt.Sprintf("虚拟实例不可用: %v", err), Pass: true}
}
marker := fmt.Sprintf("_hc_%d", time.Now().UnixNano()) marker := fmt.Sprintf("_hc_%d", time.Now().UnixNano())
triples := []sdk.Triple{ triples := []sdk.Triple{
{Subject: marker, Relation: "is", Object: "healthcheck_test", SubjectType: "System", ObjectType: "Flag"}, {Subject: marker, Relation: "is", Object: "healthcheck_test", SubjectType: "System", ObjectType: "Flag"},
} }
start := time.Now() start := time.Now()
if err := s.Memory().Commit(triples); err != nil { if err := vi.Memory.Commit(triples); err != nil {
return checkResult{Name: "memory_write", Status: "fail", Detail: fmt.Sprintf("写入失败: %v", err), Pass: false} return checkResult{Name: "memory", Status: "fail", Detail: fmt.Sprintf("写入失败: %v", err), Pass: false}
} }
n, err := s.Memory().Purge(map[string]string{"subject_contains": marker}, "hard") ents, rels, err := vi.Memory.Recall([]string{marker}, 1)
if err != nil { if err != nil {
return checkResult{Name: "memory_purge", Status: "fail", Detail: fmt.Sprintf("清理失败: %v", err), Pass: false} return checkResult{Name: "memory", Status: "fail", Detail: fmt.Sprintf("查询失败: %v", err), Pass: false}
}
if len(ents) == 0 && len(rels) == 0 {
return checkResult{Name: "memory", Status: "warn", Detail: "写入成功但查询未命中", Pass: true}
}
_, err = vi.Memory.Purge(map[string]string{"subject_contains": marker}, "hard")
if err != nil {
return checkResult{Name: "memory", Status: "fail", Detail: fmt.Sprintf("清理失败: %v", err), Pass: false}
} }
elapsed := time.Since(start) elapsed := time.Since(start)
return checkResult{ return checkResult{
Name: "memory", Name: "memory",
Status: "ok", Status: "ok",
Detail: fmt.Sprintf("写入+清理 %d 条, 耗时 %v", n, elapsed.Round(time.Millisecond)), Detail: fmt.Sprintf("隔离虚拟记忆写入+查询+清理正常, 耗时 %v", elapsed.Round(time.Millisecond)),
Pass: true, Pass: true,
} }
} }
func (p *Plugin) testKnowledgeRaw(s *sdk.PluginSDK) checkResult { func (p *Plugin) testKnowledgeRaw(s *sdk.PluginSDK) checkResult {
// 在隔离虚拟知识库上验证写→查→删,绝不动生产知识库。
vi, err := p.selftestInst(s)
if err != nil {
return checkResult{Name: "knowledge", Status: "skip", Detail: fmt.Sprintf("虚拟实例不可用: %v", err), Pass: true}
}
marker := fmt.Sprintf("_hc_knowledge_test_%d", time.Now().UnixNano()) marker := fmt.Sprintf("_hc_knowledge_test_%d", time.Now().UnixNano())
start := time.Now() start := time.Now()
if err := s.Knowledge().Add(marker, "健康检查测试标记,可忽略"); err != nil { if err := vi.Knowledge.Add(marker, "健康检查测试标记,可忽略"); err != nil {
return checkResult{Name: "knowledge", Status: "fail", Detail: fmt.Sprintf("写入失败: %v", err), Pass: false} return checkResult{Name: "knowledge", Status: "fail", Detail: fmt.Sprintf("写入失败: %v", err), Pass: false}
} }
results, err := s.Knowledge().Search("健康检查测试标记", 3) results, err := vi.Knowledge.Search("健康检查测试标记", 3)
if err != nil { if err != nil {
s.Knowledge().Remove(marker) vi.Knowledge.Remove(marker)
return checkResult{Name: "knowledge", Status: "fail", Detail: fmt.Sprintf("查询失败: %v", err), Pass: false} return checkResult{Name: "knowledge", Status: "fail", Detail: fmt.Sprintf("查询失败: %v", err), Pass: false}
} }
elapsed := time.Since(start) elapsed := time.Since(start)
// 清理测试条目,避免积累 // 清理测试条目,避免积累
s.Knowledge().Remove(marker) vi.Knowledge.Remove(marker)
if len(results) > 0 { if len(results) > 0 {
return checkResult{ return checkResult{
Name: "knowledge", Name: "knowledge",
Status: "ok", Status: "ok",
Detail: fmt.Sprintf("写入+查询正常, 耗时 %v", elapsed.Round(time.Millisecond)), Detail: fmt.Sprintf("隔离虚拟知识库写入+查询+清理正常, 耗时 %v", elapsed.Round(time.Millisecond)),
Pass: true, Pass: true,
} }
} }
@ -519,19 +555,25 @@ func (p *Plugin) testKnowledgeRaw(s *sdk.PluginSDK) checkResult {
} }
func (p *Plugin) testDocStoreRaw(s *sdk.PluginSDK) checkResult { func (p *Plugin) testDocStoreRaw(s *sdk.PluginSDK) checkResult {
// 在隔离虚拟文档记忆上验证写→查→删,绝不动生产 Document。
vi, err := p.selftestInst(s)
if err != nil {
return checkResult{Name: "documents", Status: "skip", Detail: fmt.Sprintf("虚拟实例不可用: %v", err), Pass: true}
}
start := time.Now() start := time.Now()
doc := &sdk.Doc{ doc := &sdk.Doc{
Title: fmt.Sprintf("健康检查测试文档 %d", time.Now().UnixNano()), Title: fmt.Sprintf("健康检查测试文档 %d", time.Now().UnixNano()),
Content: "这是一条由 healthcheck 插件创建的测试文档,用于验证文档记忆系统是否正常工作。", Content: "这是一条由 healthcheck 插件创建的测试文档,用于验证文档记忆系统是否正常工作。",
} }
if err := s.DocMemory().Insert(doc); err != nil { if err := vi.DocMemory.Insert(doc); err != nil {
return checkResult{Name: "documents", Status: "fail", Detail: fmt.Sprintf("写入失败: %v", err), Pass: false} return checkResult{Name: "documents", Status: "fail", Detail: fmt.Sprintf("写入失败: %v", err), Pass: false}
} }
// 清理测试文档避免积累SDK Insert 不回填 ID经 Query 按标题定位) // 清理测试文档避免积累SDK Insert 不回填 ID经 Query 按标题定位)
for _, d := range s.DocMemory().Query("健康检查测试文档", 10) { for _, d := range vi.DocMemory.Query("健康检查测试文档", 10) {
if d.ID != "" && strings.HasPrefix(d.Title, "健康检查测试文档") { if d.ID != "" && strings.HasPrefix(d.Title, "健康检查测试文档") {
s.DocMemory().Remove(d.ID) vi.DocMemory.Remove(d.ID)
} }
} }
@ -539,7 +581,7 @@ func (p *Plugin) testDocStoreRaw(s *sdk.PluginSDK) checkResult {
return checkResult{ return checkResult{
Name: "documents", Name: "documents",
Status: "ok", Status: "ok",
Detail: fmt.Sprintf("写入+删除正常, 耗时 %v", elapsed.Round(time.Millisecond)), Detail: fmt.Sprintf("隔离虚拟文档记忆写入+查询+清理正常, 耗时 %v", elapsed.Round(time.Millisecond)),
Pass: true, Pass: true,
} }
} }
@ -626,7 +668,8 @@ func (p *Plugin) testLLMDriven(s *sdk.PluginSDK) checkResult {
} }
// collectToolDefsForLLM 收集全部已注册的工具定义供 LLM 发现和测试。 // collectToolDefsForLLM 收集全部已注册的工具定义供 LLM 发现和测试。
// 动态排除本插件自身注册的工具(通过 selfToolNames避免 LLM 自我循环调用 // 动态排除本插件自身注册的工具(通过 selfToolNames避免 LLM 自我循环调用
// 且仅保留"只读/轻量验证"类工具(白名单语义),防止 LLM 自检污染生产数据或引发副作用。
func (p *Plugin) collectToolDefsForLLM(s *sdk.PluginSDK) []sdk.ToolDef { func (p *Plugin) collectToolDefsForLLM(s *sdk.PluginSDK) []sdk.ToolDef {
seen := map[string]bool{} seen := map[string]bool{}
var defs []sdk.ToolDef var defs []sdk.ToolDef
@ -635,6 +678,9 @@ func (p *Plugin) collectToolDefsForLLM(s *sdk.PluginSDK) []sdk.ToolDef {
if p.selfToolNames[d.Name] || seen[d.Name] { if p.selfToolNames[d.Name] || seen[d.Name] {
return return
} }
if !isSafeReadonlyTool(d.Name) {
return
}
seen[d.Name] = true seen[d.Name] = true
defs = append(defs, d) defs = append(defs, d)
} }
@ -651,10 +697,38 @@ func (p *Plugin) collectToolDefsForLLM(s *sdk.PluginSDK) []sdk.ToolDef {
return defs return defs
} }
// isSafeReadonlyTool 判断工具是否为"只读/无副作用、适合健康检查 LLM 自检"的工具。
// 仅白名单语义:不在白名单的工具一律不测(宁可少测,不可污染/引发副作用)。
func isSafeReadonlyTool(name string) bool {
// 明确只读的查询/列表类工具
readonlyExact := map[string]bool{
"memory_recall": true,
"memory_introspect": true,
"doc_query": true,
"knowledge_search": true,
"knowledge_list": true,
"person_query": true,
"person_network": true,
"llm_list_sources": true,
"output_list_channels": true,
"terminal_list": true,
}
if readonlyExact[name] {
return true
}
// 带 _list/_help 后缀的通常是只读展示
for _, sfx := range []string{"_list", "_help"} {
if strings.HasSuffix(name, sfx) {
return true
}
}
return false
}
// buildDiscoveryPrompt 为 LLM 构造工具探索 prompt。 // buildDiscoveryPrompt 为 LLM 构造工具探索 prompt。
func (p *Plugin) buildDiscoveryPrompt(toolDefs []sdk.ToolDef) string { func (p *Plugin) buildDiscoveryPrompt(toolDefs []sdk.ToolDef) string {
var b strings.Builder var b strings.Builder
b.WriteString(fmt.Sprintf(`你是一名系统健康检查专家。以下是系统中各插件提供的 %d 个工具(已自动排除健康检查插件自身工具): b.WriteString(fmt.Sprintf(`你是一名系统健康检查专家。以下是系统中各插件提供的 %d 个工具(已自动排除健康检查插件自身工具及所有会写/删/改生产数据或产生外部副作用的工具,以下均为只读/查询/列表类工具
你的任务是:逐一尝试调用这些工具,验证它们是否正常工作,并对于每个工具使用 healthcheck_report 工具上报测试结果。 你的任务是:逐一尝试调用这些工具,验证它们是否正常工作,并对于每个工具使用 healthcheck_report 工具上报测试结果。
@ -665,7 +739,7 @@ func (p *Plugin) buildDiscoveryPrompt(toolDefs []sdk.ToolDef) string {
4. 调用 healthcheck_report 工具上报tool_name, status=ok/fail/skip, detail=详情) 4. 调用 healthcheck_report 工具上报tool_name, status=ok/fail/skip, detail=详情)
注意: 注意:
- 有些工具有副作用(如写入数据),请使用安全参数,测试后应清理 - 所有工具均为只读、无副作用,可放心调用
- 尽可能覆盖所有工具 - 尽可能覆盖所有工具
- 每个工具只需测试一次 - 每个工具只需测试一次

View File

@ -3,6 +3,7 @@ package healthcheck
import ( import (
"encoding/json" "encoding/json"
"os" "os"
"strings"
"testing" "testing"
agentCore "gitcode.com/JianFeeeee/HomeAgent/internal/agent/core" agentCore "gitcode.com/JianFeeeee/HomeAgent/internal/agent/core"
@ -196,6 +197,20 @@ func TestHealthcheckWithMemory(t *testing.T) {
} }
defer memDB.Close() defer memDB.Close()
// 预置一条生产数据,验证 healthcheck 自检后不触碰它
_, _, err = memDB.Commit([]memory.Triple{{Subject: "用户", Relation: "喜欢", Object: "咖啡"}}, "test", 0)
if err != nil {
t.Fatal(err)
}
memSnapshot := func() (int, int) {
m, _ := memDB.Introspect()
ents, _ := m["entity_count"].(int)
rels, _ := m["relation_count"].(int)
return ents, rels
}
be, br := memSnapshot()
_, tc, err := setupPluginWith(sdk.SDKConfig{Memory: sdk.NewGraphMemory(memDB)}) _, tc, err := setupPluginWith(sdk.SDKConfig{Memory: sdk.NewGraphMemory(memDB)})
if err != nil { if err != nil {
t.Fatal(err) t.Fatal(err)
@ -212,11 +227,25 @@ func TestHealthcheckWithMemory(t *testing.T) {
json.Unmarshal(data, &resp) json.Unmarshal(data, &resp)
if resp["status"] != "ok" { if resp["status"] != "ok" {
t.Fatalf("expected status ok, got %v", resp["status"]) t.Fatalf("expected status ok, got %v (detail=%v)", resp["status"], resp["detail"])
} }
if pass, ok := resp["pass"].(bool); !ok || !pass { if pass, ok := resp["pass"].(bool); !ok || !pass {
t.Fatalf("expected pass=true, got pass=%v status=%v detail=%v", pass, resp["status"], resp["detail"]) t.Fatalf("expected pass=true, got pass=%v status=%v detail=%v", pass, resp["status"], resp["detail"])
} }
// 关键:生产记忆内容必须保持不变(未被 healthcheck 污染)
ae, ar := memSnapshot()
if ae != be || ar != br {
t.Fatalf("production memory polluted by healthcheck self-test: before=(%d,%d) after=(%d,%d)", be, br, ae, ar)
}
// 再次确认:注入实例中不应出现 _hc_ 测试实体
relResult, _ := memDB.Recall([]string{"_hc_"}, nil, 1, "")
for _, e := range relResult.Entities {
if len(e.Name) >= 4 && e.Name[:4] == "_hc_" {
t.Fatalf("healthcheck left _hc_ entity in production memory: %q", e.Name)
}
}
} }
func TestHealthcheckWithKnowledge(t *testing.T) { func TestHealthcheckWithKnowledge(t *testing.T) {
@ -232,6 +261,12 @@ func TestHealthcheckWithKnowledge(t *testing.T) {
} }
defer ks.Stop() defer ks.Stop()
// 预置一条真实知识,验证 healthcheck 自检后不触碰它
if err := ks.Add("生产知识点", "这是生产知识,不应被健康检查破坏"); err != nil {
t.Fatal(err)
}
before := len(ks.List())
_, tc, err := setupPluginWith(sdk.SDKConfig{Knowledge: sdk.NewKnowledge(ks)}) _, tc, err := setupPluginWith(sdk.SDKConfig{Knowledge: sdk.NewKnowledge(ks)})
if err != nil { if err != nil {
t.Fatal(err) t.Fatal(err)
@ -264,6 +299,20 @@ func TestHealthcheckWithKnowledge(t *testing.T) {
if knowledgeCheck == nil { if knowledgeCheck == nil {
t.Fatal("expected knowledge check in results") t.Fatal("expected knowledge check in results")
} }
if knowledgeCheck["status"] != "ok" {
t.Fatalf("expected knowledge check ok, got %v (detail=%v)", knowledgeCheck["status"], knowledgeCheck["detail"])
}
// 生产知识库内容必须保持不变(未被 healthcheck 污染)
after := len(ks.List())
if after != before {
t.Fatalf("production knowledge polluted: before=%d after=%d", before, after)
}
for _, name := range ks.List() {
if strings.HasPrefix(name, "_hc_knowledge_test_") {
t.Fatalf("healthcheck left test knowledge in production: %q", name)
}
}
} }
func TestHealthcheckWithDocStore(t *testing.T) { func TestHealthcheckWithDocStore(t *testing.T) {
@ -279,6 +328,12 @@ func TestHealthcheckWithDocStore(t *testing.T) {
} }
defer ds.Stop() defer ds.Stop()
// 预置一篇生产文档,验证 healthcheck 自检后不触碰它
if err := ds.Insert(&doc.Doc{ID: "prod_doc", Summary: "生产文档", Content: "这是生产文档,不应被健康检查破坏"}); err != nil {
t.Fatal(err)
}
before := len(ds.Query("生产文档", 10))
_, tc, err := setupPluginWith(sdk.SDKConfig{DocMemory: sdk.NewDocMemory(ds)}) _, tc, err := setupPluginWith(sdk.SDKConfig{DocMemory: sdk.NewDocMemory(ds)})
if err != nil { if err != nil {
t.Fatal(err) t.Fatal(err)
@ -311,6 +366,15 @@ func TestHealthcheckWithDocStore(t *testing.T) {
if docCheck == nil { if docCheck == nil {
t.Fatal("expected documents check in results") t.Fatal("expected documents check in results")
} }
if docCheck["status"] != "ok" {
t.Fatalf("expected documents check ok, got %v (detail=%v)", docCheck["status"], docCheck["detail"])
}
// 生产文档必须保持不变(未被 healthcheck 污染)
after := len(ds.Query("生产文档", 10))
if after != before {
t.Fatalf("production doc store polluted: before=%d after=%d", before, after)
}
} }
func TestLLMReportCollection(t *testing.T) { func TestLLMReportCollection(t *testing.T) {
@ -330,4 +394,88 @@ func TestLLMReportCollection(t *testing.T) {
} }
} }
// TestSafeReadonlyTool 验证 LLM 自检工具过滤:写类/副作用工具被拒绝,只读工具被放行。
func TestSafeReadonlyTool(t *testing.T) {
// 只读:应放行
readonly := []string{
"memory_recall", "memory_introspect", "doc_query",
"knowledge_search", "knowledge_list", "person_query",
"person_network", "llm_list_sources", "output_list_channels",
"terminal_list", "files_list",
}
for _, name := range readonly {
if !isSafeReadonlyTool(name) {
t.Errorf("expected readonly tool %q to be safe, but rejected", name)
}
}
// 写/删/改/副作用:应被拒绝
mutating := []string{
"memory_commit", "memory_edit", "memory_purge", "memory_delete_entity",
"memory_merge", "memory_block_merge",
"knowledge_create", "knowledge_delete",
"doc_commit", "doc_delete",
"person_set_trait", "person_relate",
"output_send__qq", "output_send__cli",
"llm_set_source", "config_set",
"timer_set", "plgreload", "plugin_disable",
"cmd_run", "files_write", "files_delete", "terminal_create", "terminal_write",
"terminal_close", "spawn_child",
}
for _, name := range mutating {
if isSafeReadonlyTool(name) {
t.Errorf("expected mutating tool %q to be rejected, but allowed", name)
}
}
}
// TestCollectToolDefsForLLMNoMutating 验证 collectToolDefsForLLM 不会把写类工具交给 LLM 自检。
func TestCollectToolDefsForLLMNoMutating(t *testing.T) {
p := &Plugin{name: "healthcheck", selfToolNames: map[string]bool{"healthcheck": true, "healthcheck_report": true}}
// 构造一个包含写类工具的 ToolDef 集合,注入工具注册表
stage := agentCore.NewStageHost()
defs := []sdk.ToolDef{
{Name: "memory_recall", Plugin: "memory", Description: "recall"},
{Name: "memory_commit", Plugin: "memory", Description: "commit"},
{Name: "doc_query", Plugin: "doc", Description: "query"},
{Name: "doc_commit", Plugin: "doc", Description: "commit doc"},
{Name: "knowledge_search", Plugin: "knowledge", Description: "search"},
{Name: "knowledge_create", Plugin: "knowledge", Description: "create"},
{Name: "cmd_run", Plugin: "cmd", Description: "run cmd"},
{Name: "files_list", Plugin: "files", Description: "list"},
}
for _, d := range defs {
d := d
stage.RegisterTool(d.Name, d, func(map[string]interface{}) (interface{}, error) { return nil, nil })
}
tc := newToolCapture()
s := newTestSDK(sdk.SDKConfig{
RegTool: tc.RegisterTool,
RegStage: tc.RegisterStage,
RegAPI: tc.RegisterAPI,
Tool: sdk.NewTool(stage, agentIO.NewIOManager()),
})
got := p.collectToolDefsForLLM(s)
allowed := map[string]bool{}
for _, d := range got {
allowed[d.Name] = true
}
// 只读工具应被包含
for _, name := range []string{"memory_recall", "doc_query", "knowledge_search", "files_list"} {
if !allowed[name] {
t.Errorf("expected readonly tool %q in LLM selftest set, missing", name)
}
}
// 写类/副作用工具绝不能被交给 LLM
for _, name := range []string{"memory_commit", "doc_commit", "knowledge_create", "cmd_run"} {
if allowed[name] {
t.Errorf("mutating tool %q must NOT be in LLM selftest set", name)
}
}
}

View File

@ -2,6 +2,7 @@ package sdk
import ( import (
"log" "log"
"sync"
pubsdk "gitcode.com/JianFeeeee/homeagent-sdk/sdk" pubsdk "gitcode.com/JianFeeeee/homeagent-sdk/sdk"
agentIO "gitcode.com/JianFeeeee/HomeAgent/internal/agent/io" agentIO "gitcode.com/JianFeeeee/HomeAgent/internal/agent/io"
@ -98,6 +99,9 @@ type PluginSDK struct {
config ConfigAPI config ConfigAPI
tool ToolAPI tool ToolAPI
indexer IndexerAPI indexer IndexerAPI
selftestMu sync.Mutex
selftest *VirtualInstance
} }
func (s *PluginSDK) PluginMgr() PluginManager { return s.pluginMgr } func (s *PluginSDK) PluginMgr() PluginManager { return s.pluginMgr }
@ -211,6 +215,42 @@ func New(name string, cfg SDKConfig) *PluginSDK {
} }
} }
// Selftest 返回一个隔离的虚拟自检实例healthcheck 等内置插件用),
// 完全独立于生产存储,不产生任何污染。首次调用创建,复用已存在实例;
// 每轮自检前调用 SelftestReset 重建以清空上轮测试数据。
func (s *PluginSDK) Selftest(scope string) (*VirtualInstance, error) {
s.selftestMu.Lock()
defer s.selftestMu.Unlock()
if s.selftest == nil {
vi, err := NewVirtualInstance(scope)
if err != nil {
return nil, err
}
s.selftest = vi
}
return s.selftest, nil
}
// SelftestReset 清理并重建隔离自检实例,用于每轮健康检查前重置状态。
func (s *PluginSDK) SelftestReset(scope string) error {
s.selftestMu.Lock()
defer s.selftestMu.Unlock()
if s.selftest == nil {
vi, err := NewVirtualInstance(scope)
if err != nil {
return err
}
s.selftest = vi
return nil
}
vi, err := s.selftest.Reset(scope)
if err != nil {
return err
}
s.selftest = vi
return nil
}
func (s *PluginSDK) Status() StatusAPI { return s.status } func (s *PluginSDK) Status() StatusAPI { return s.status }
func (s *PluginSDK) Supervisor() SupervisorAPI { return s.supervisor } func (s *PluginSDK) Supervisor() SupervisorAPI { return s.supervisor }
func (s *PluginSDK) Adapter() AdapterAPI { return s.adapter } func (s *PluginSDK) Adapter() AdapterAPI { return s.adapter }

108
internal/sdk/selftest.go Normal file
View File

@ -0,0 +1,108 @@
package sdk
import (
"fmt"
"os"
"path/filepath"
"sync"
"gitcode.com/JianFeeeee/HomeAgent/internal/knowledge"
"gitcode.com/JianFeeeee/HomeAgent/internal/memory"
doc "gitcode.com/JianFeeeee/HomeAgent/internal/memory/document"
"gitcode.com/JianFeeeee/HomeAgent/internal/memory/text"
)
// VirtualInstance 是完全隔离的虚拟存储集合,供内置插件(如 healthcheck
// 做不污染生产存储的"写→查→删"往返自检。所有写入发生在独立临时目录,
// 由 Cleanup 统一销毁。
type VirtualInstance struct {
Memory MemoryAPI
Knowledge KnowledgeAPI
DocMemory DocMemoryAPI
TextMemory TextMemoryAPI
dir string
mu sync.Mutex
created bool
}
// newBaseDir 返回一个隔离的临时根目录hot 前缀,避免与生产数据混淆)。
func newBaseDir(scope string) (string, error) {
return os.MkdirTemp("", "homeagent_selftest_"+scope+"_")
}
// NewVirtualInstance 创建完全隔离的测试实例,使用独立临时目录。
func NewVirtualInstance(scope string) (*VirtualInstance, error) {
if scope == "" {
scope = "hc"
}
base, err := newBaseDir(scope)
if err != nil {
return nil, err
}
v := &VirtualInstance{dir: base, created: true}
if err := v.initLocked(); err != nil {
os.RemoveAll(base)
return nil, err
}
return v, nil
}
// initLocked 初始化各隔离存储。调用方须持 v.mu。
func (v *VirtualInstance) initLocked() error {
memDB, err := memory.NewGraphDB(filepath.Join(v.dir, "graph.db"))
if err != nil {
return fmt.Errorf("virtual graph db: %w", err)
}
v.Memory = NewGraphMemory(memDB)
ks := knowledge.NewStore(filepath.Join(v.dir, "knowledge"))
if err := ks.Start(); err != nil {
return fmt.Errorf("virtual knowledge store: %w", err)
}
v.Knowledge = NewKnowledge(ks)
ds := doc.NewStore(filepath.Join(v.dir, "documents"))
if err := ds.Start(); err != nil {
return fmt.Errorf("virtual doc store: %w", err)
}
v.DocMemory = NewDocMemory(ds)
tm := text.New(filepath.Join(v.dir, "text"))
if err := tm.Start(); err != nil {
return fmt.Errorf("virtual text memory: %w", err)
}
v.TextMemory = NewTextMemory(tm)
return nil
}
// Cleanup 关闭并销毁该虚拟实例的全部临时存储。之后实例不可再用。
func (v *VirtualInstance) Cleanup() {
v.mu.Lock()
defer v.mu.Unlock()
if !v.created {
return
}
if tm, ok := v.TextMemory.(*textMemoryImpl); ok && tm.tm != nil {
tm.tm.Stop()
}
if ds, ok := v.DocMemory.(*docMemoryImpl); ok && ds.ds != nil {
ds.ds.Stop()
}
if ks, ok := v.Knowledge.(*knowledgeImpl); ok && ks.ks != nil {
ks.ks.Stop()
}
if gm, ok := v.Memory.(*graphMemory); ok && gm.db != nil {
gm.db.Close()
}
os.RemoveAll(v.dir)
v.created = false
}
// Reset 清理当前实例并重建一个全新的隔离实例,用于每轮自检前重置状态。
// 返回新实例;失败时旧实例已被清理、返回 err调用方需重新创建。
func (v *VirtualInstance) Reset(scope string) (*VirtualInstance, error) {
v.Cleanup()
return NewVirtualInstance(scope)
}

217
plan.md Normal file
View File

@ -0,0 +1,217 @@
# HomeAgent 生产问题修复计划
## 设计意图备忘(核心架构原则)
本框架的两大核心设计意图,贯穿所有插件/记忆/工具设计,**所有改动必须符合**
1. **插件即 Agent 的 App** —— Agent 像人用 App 一样用插件。
- QQ 插件应像 QQ 客户端:通知到来 → 看到预览/上下文 → 一键回复。
- **认知负荷最小**:中断/通知只给「发送者昵称 + 消息预览」,**元数据user_id/group_id/message_id完全走工具不进 prompt**,防提示词注入、防昵称欺诈、低认知负荷。
- **工具语义自解释**`qq_get_message``output_send__qq``qq_get_history` 的 Description 要让 Agent「读完即知怎么用」无需额外指令。
2. **基于相关性的分层记忆架构** —— L1 文本流 / L2 图谱 / L3 文档向量,按相关性蒸馏、检索、归档。
- GraphDB 去重、docToTriples 模板清理、Pipeline 增量蒸馏、嵌入模型内存优化,均服务于此。
---
## 0.1 紧急healthcheck 健康检查污染真实存储 ⚠️ 正在持续污染
**现象**2026-08-11 22:02 起,每 30 分钟一次):日志反复出现
`[knowledge] added: _hc_knowledge_test_<ts>`,且知识库出现 `gotest``luatest``_hc_knowledge_test_*` 等测试残留。
**根因**`internal/plugins/healthcheck/plugin.go` 的三个"写入通道"自检**全部在真实生产存储上写入再删除**
| 函数 | 写入 | 清理 | 固化问题 |
|------|------|------|----------|
| `testMemoryRaw` (:460) | `Memory().Commit(_hc_<ts> triple)` | `Purge(hard)` | GraphDB 空实体/AUTOINCREMENT id 膨胀 |
| `testKnowledgeRaw` (:485) | `Knowledge().Add(_hc_knowledge_test_<ts>)` | `Remove(marker)` | Add 异步 `go writeIndex` vs Remove 异步重写**竞态**`dirName`(用了 `/` 解析)与 `Remove``id=sanitize(name)` 计算**不一致**→ 删除可能失效 → **残留目录固化成文件** |
| `testDocStoreRaw` (:521) | `DocMemory().Insert(...)` | Query 后 Remove | 真实 docStore 写文件再删,抖动 |
**原则**:健康检查验证的是"写入通道是否可用"**结果不应固化进生产记忆**。改为**独立虚拟/影子空间**或**不落盘验证**。
**实现**
- [x] **核心实现**SDK 新增 `VirtualInstance``internal/sdk/selftest.go`healthcheck 的三个 raw 自检改为在完全隔离的虚拟空间(`os.MkdirTemp` 独立图库/知识库/文档/文本)上做真实"写→查→删",绝不碰生产存储。
- [x] `PluginSDK.Selftest()/SelftestReset()` 暴露隔离实例(含 mutex 防并发),每轮自检前 `SelftestReset` 重建清空上轮数据。
- [x] **LLM 驱动自检防护**`collectToolDefsForLLM` 改为只收集**只读白名单**工具(`isSafeReadonlyTool`),写/删/改生产数据及外部副作用工具memory_commit/doc_commit/knowledge_create/cmd_run/files_write/terminal_*/output_send/spawn_child 等)一律不交给 LLM 自检,防止 LLM 乱调污染生产。
- [x] 单测healthcheck 自检后注入的"生产"实例内容不变(快照对比 + `_hc_` 无残留);读/写工具白名单过滤测试通过。
- [x] 存量清理:删除生产残留的 `gotest/``luatest/``_hc_knowledge_test_*/` 目录(保留真实知识库);备份留于 `/tmp/opencode/knowledge_garbage_backup_20260812_122322`
- [ ] 部署验证:编译新 `homed` 部署后,`knowledge/``memory/graph.db``memory/documents/` 不再出现 `_hc_*` 残留。
---
## 0.2 紧急QQ 消息被无视agentcli 幽灵终端自喂送风暴)⚠️ 优先处理
**现象**2026-08-11 20:4x用户发 QQ 私聊消息agent 不回应。日志显示 agent 被 `agentcli` 终端 echo 洪水完全阻塞。
**根因(两个耦合缺陷,非记忆层)**
1. **`agentcli` 幽灵终端自喂送风暴(直接原因)**
- `internal/plugins/agentcli/plugin.go:602` — 终端一旦有输出就 `s.InjectText("agentcli","agentcli","[终端 X 有新输出]\n...")` 注入 agent 事件循环。
- 残留 `term_3`bash 的 `git sparse clone` 进度)持续吐数据 → 插件每 `NotifyOutputDelay`(几秒)注入一次新输入 → agent 每 6-14 秒跑完整 LLM 工具循环去 `terminal_read`/`terminal_list` → 读到新进度 → 再注入 → **无限自循环,独占整个 eventLoop**
- 日志2026-08-11 20:4820:52 期间约 20+ 次 `input from agentcli`,无一例外。
2. **QQ 消息无真正优先级(放大原因)**
- QQ 走 `interceptLoop``cancelLLM` + 塞入 `interceptCh``eventloop.go:84-85`)。
-`process()` 打断后 `continue` 回到**同一回合**`process.go:139`QQ 的 `[打断消息]` 只被追加到幽灵回合对话里夹带响应,**拿不到独立处理回合**。
- 面对 agentcli 持续自喂送QQ 永远排不到前面 → 20:51 之后的 QQ 消息一条未回。
### 止血(运维,立即执行)
- [x] 杀掉残留 `term_3` bash本机 PID 3716282→ 幽灵回声立停QQ 事件循环恢复。
- [x] **QQ 身份注入增强**`third_party/homeagent-sdk/example/qq/plugin.go` 中断模板新增 QQ号/群号(`evt.UserID`/`evt.GroupID`agent 无需先调 `get_message` 即可知发送者身份。
- 后续:`agentcli` 终端用后应 `terminal_close`,避免残留。
### 根治(改代码,入本计划)
- **Phase 6**agentcli 终端通知节流/去重,同源同 tag 的"有新输出"合并;避免对长输出逐段注入。
---
# HomeAgent 记忆层修复计划
> 基于生产实例诊断2026-08-11清理完成`graph.db` 从 6225 条 relations99.4% 垃圾) → **13 条真实关系**entities 从 29 → **17**。备份文件:`memory/graph.db.backup.20260811_163301`。
## 核心问题清单
| # | 问题 | 影响 | 位置 |
|---|------|------|------|
| 1 | GraphDB.Commit 对 relations **裸 INSERT 无去重** | 同一三元组每次归档无限重复,生产 5613 条 `文档→来源→context_archived` 重复垃圾 | `internal/memory/graph.go:209` |
| 2 | `docToTriples()` **对每篇冷文档永远生成固定模板三元组**`文档→来源→context_archived``文档→主题→{summary}` | 归档即写垃圾配合问题1指数级累积 | `internal/agent/core/distill.go:390-437` |
| 3 | `core.agent.distill_interval=2d` 配置 **静默失效**Go `time.ParseDuration` 不支持 `d` 单位),回退 30 分钟默认值 | archiveColdDocs **每小时跑**而非每 2 天,放大问题 1/2 | `internal/config/registry.go:698` |
| 4 | `Pipeline Distiller` 标称"10min 心跳"**实为 7 天一批回放** | 文档与行为严重不符,虽未直接制造垃圾但不可信 | `internal/memory/pipeline/pipeline.go:194` |
| 5 | 嵌入模型双模加载(中 200k + 英 378k300 维)直接导致 **2.4G 常驻** | OOM 风险、启动慢,文档未提内存代价 | 部署配置 `/data/cc.*.vec` |
---
## 修复计划
### Phase 1图记忆去重最小改动、最高收益✅ **进行中**
**目标**`Commit` 对 relations 加唯一约束 + 冲突即跳过,彻底阻断重复累积。
- [x] **Schema 迁移**`initSchema` 新增 `migrateRelationUnique`——检测旧 relations 表无复合唯一约束(旧 DD表自动重建为带 `UNIQUE(source_id, target_id, relation_type, session_id)` 的新表并 `INSERT OR IGNORE` 去重(官方 12 步迁移),无需人工干预。
- [x] **Commit 逻辑**:改为"查存在 → 不存在才 INSERT 并计数;已存在则仅刷新 confidence/updated_at",重复提交不新增、不重计。
- [x] **验证**`TestCommitDedupSameSession`(同会话重复 commit 不增行)、`TestCommitDedupDifferentSession`(跨会话允许重复)、`TestMigrateRelationUniqueDedupsOldTable`(旧表重建去重)全绿;`go build ./...` 通过。
- [ ] 生产部署后确认graph.db 36→ 去重2 组 `like/plugin` 重复消失),跑 1 周不再新增重复。
> entities 已有 `UNIQUE(name)` 保护,仅 relations 缺失。
---
### Phase 2归档三元组模板清理治本✅ **计划中**
**目标**`docToTriples` 不再把 `context_archived`/`Topic` 摘要当成实体写入图库。
- [ ] 重构 `docToTriples`:仅当 `doc.Source``context_archived` 且非空时写 `文档→来源``文档→主题` 仅当 summary 长度合理(<80 且非模板化时写否则跳过
- [ ] 引入 `doc.Meta["is_archived_context"]` 标记上下文归档文档 `docToTriples` 识别并跳过
- [ ] 单测验证构造冷文档 `archiveColdDocs` 无模板垃圾产出
---
### Phase 3配置持久化解析修复防配置失效✅ **计划中**
**目标**支持 `2d`/`1w` 等人类可读时长单位配置即时生效
- [x] 实现 `parseDurationExtended(string) time.Duration`正则识别 `\d+[dhw]` 换算为 `time.Hour*24` 再调 `time.ParseDuration`
- [x] 替换 `GetDuration` 加载点`registry.go` `GetDuration` 共用`main.go:423-426` 的四个间隔配置自动受益
- [x] 单测`TestParseDurationExtended``"2d"==48h``"1w"==168h``2d12h`复合/无单位错误输入+ `GetDuration` 集成用例全绿
- [x] 生产核实当前生产 `core.agent.distill_interval=30m`可解析非失效态修复为防御性未来 `2d`/`1w` 写入即可生效
---
### Phase 4Pipeline Distiller 行为对齐文档(可选,低优)✅ **计划中**
- [ ] 改为真正的增量蒸馏 tick 取最近 `RetentionDays` 未蒸馏记录 `extractKeyTriples` `Commit`标记 `Distilled=true`
- [ ] 移除 `CreatedAt.Before(cutoff)` 7 天门槛改为" tick 处理前 N "保持文档所述 10min 频率
- [ ] 单测验证启动即蒸馏 + 不重复蒸馏
---
### Phase 5嵌入模型内存优化运维侧✅ **计划中**
- [ ] 提供 **量化/裁剪** 选项`embedding_model_path` 支持 `top50k` 等规格或运行时 `mmap` 只加载词表需求词
- [ ] 文档补充内存预算说明双模 300 1.5G RAM/模型
- [ ] 生产可选降级仅保留中文模型主语言)。
---
### Phase 6agentcli 终端通知频率策略QQ 被无视的直接原因)✅ **进行中**
**背景**`readLoop``internal/plugins/agentcli/plugin.go:556-608`用写死常量 `NotifyOutputDelay = 500ms``plugin.go:24`**定时节流**只要终端持续输出 `git sparse clone` 进度就每 500ms 注入一条 `[终端 X 有新输出]` agent造成无限自喂送占满 eventLoop使 QQ 消息永远只能塞进回环通道且被 echo 上下文淹没
**原则(明确保留通知,不改中断机制)**
- **通知机制必须保留**——agent 需要感知"终端仍在运行可能有待读取的输出"否则会忘记终端的存在不知何时去 `terminal_read`
- 真正要改的是**写死的 500ms 定时节流**改为**基于输出语义/任务生命周期的通知策略**"多少输出一通知 / 命令执行结束再通知"由可配策略决定而非插件写死
**实现**
- [x] 通知从" 500ms 定时"改为**按终端生命周期事件触发**
- **进程结束 / 超时 / 读取错误 立即通知**已存在 + 增强 EOF 即时触发)。
- **持续运行仅产出进度 低频"有新输出"通知**累计 `notify_bytes` 默认 2KB 未读字节或距上次通知 `notify_interval` 默认 2s两条件满足任一即触发而非每 500ms
- **首次创建 立即通知"**已启动**"**确保 agent 感知终端存在
- [x] 通知频率可配置per-plugin settings`notify_bytes``notify_interval`把控制权交还 agent不写死
- [x] 纯进度输出仍吸入 `t.buf`agent 需要时用现有 `terminal_read` 主动拉全量保持 agent 可感知存在可自主决策取量)。
- [ ] 单测终端持续吐进度时消息注入频率显著低于 500ms/进程结束/出错/提示符时立即通知
- [ ] 运维止血杀掉残留 `term_3` bashPID 3716282验证 QQ 消息恢复响应
---
### Phase 7中断机制核验确认不需改动仅作记录✅ **已确认**
**结论**中断机制本身正确**无需改动**。文档ARCHITECTURE.md:9,368-384明确
> QQ 等外部插件提示走 `InjectInterrupt → interruptCh → interceptLoop → 塞 interceptCh内部回环通道→ process() 每轮前 drainInterrupts() 以 [打断消息] 注入当前对话流` —— **同一段 LLM 记忆内连贯处理**,不割裂、不另开新回合。
- [x] 验证回环通道存在`interceptCh``eventloop.go:85`+ `drainInterrupts()``eventloop.go:404` `process()` 每轮 LLM call 前非阻塞排空注入
- [x] 确认设计约束"不能开新回合保证 LLM 记忆连贯" —— 中断注入当前对话流QQ 在该语境下 `qq_get_message` 看消息回复接着干
- [ ] 仅作回归验证修复 Phase 6 QQ 消息在 agentcli 不泛滥时能正常经回环通道被响应20:49 已证明机制可达)。
---
## 验收标准
| 指标 | 当前 | 目标 | 验收方式 |
|------|------|------|----------|
| Graph relations 重复率 | ~90% (5613/6225) | 0% | `SELECT count(*), count(DISTINCT source_id||target_id||relation_type) FROM relations` |
| `context_archived` 关系残留 | 5613 | 0 | `grep` 关系表 |
| `distill_interval` 配置生效 | 失效(30m) | 2d | 日志 `heartbeat distill tick` 间隔 = 48h |
| Pipeline distiller 频率 | 7天一批 | 10min | 日志 `distilled N records` 10min |
| 启动内存占用 | 2.4G | <1.5G单模或可配 | `systemd` MemoryCurrent |
| agentcli 通知注入频率 | 500ms/ | 仅在生命周期事件/低频里程碑 | 日志 `input from agentcli` 密度显著下降 |
| healthcheck 污染生产存储 | 30min `_hc_*` | 0不落生产存储/虚拟空间 | `knowledge/``memory/` `_hc_*``gotest``luatest` 残留快照对比 |
---
## 实施顺序建议
1. **立即**Phase 0.1healthcheck 污染隔离)——正在持续污染生产最高优先顺手清理存量 `_hc_*`/`gotest`/`luatest`
2. **立即**Phase 1去重+ Phase 3配置解析)—— 互不依赖风险最低收益最大
3. **次日**Phase 2模板清理)—— 需确认 Phase 1 生效后防止旧垃圾再次写入
4. **终端洪水紧急项**Phase 6agentcli 通知频率策略)——解决 QQ 被无视的直接原因先运维止血杀残留终端再上代码
5. **后续**Phase 4/5 按需求排期Phase 7 确认中断机制无需改动
---
## 相关文件清单
```
internal/plugins/healthcheck/plugin.go # Phase 0.1:自检写入改为隔离空间/不落盘
internal/plugins/agentcli/plugin.go # Phase 6终端通知频率策略
internal/memory/graph.go # Commit 去重 + Schema 迁移
internal/agent/core/distill.go # docToTriples 重构
internal/config/registry.go # parseDurationExtended
internal/memory/pipeline/pipeline.go # distillOnce 增量化
internal/memory/static_embedder.go # 量化/裁剪入口(可选)
```
---
## 回滚预案
- Phase 1/2 修改数据库 Schema/写入逻辑保留 `graph.db.backup.*`出问题 `systemctl stop homeagent && cp backup graph.db && systemctl start`
- Phase 3 仅改配置解析回滚即改回 `time.ParseDuration`
- 所有改动需先跑 `make test`内存//文档/配置全绿再部署
---
*更新时间2026-08-11*
*生产实例`/home/newqqagent`systemd 托管二进制 `/usr/local/bin/homed` (v0.8.0, 2026-07-28 build)*