Files
HomeAgent/internal/agent/core/distill.go
JianFeeeee d98bf512e1 feat(memory): 场景式关联召回——给记忆节点赋场景引用,场面重现即取回
背景(实测):带条件的记忆召不回来。生产库里明明有
「QQ回复禁用Markdown格式 --规定--> 纯文本不用Markdown」「老大 --偏好--> 同左」,
但输入「QQ回复格式」时命中 148 个实体、规则排第 32,注入只取前 5——规则根本没进去;
输入「在吗」这种零内容词的短消息,向量路反而灌进 17 个毫不相关的实体。

根因:词法/向量召回都建立在「字面或语义相似」上,而条件式记忆(在什么场合该怎么做)
约束的是**场面**不是话题。用户措辞不重合时它天然召不回;措辞太宽("QQ")时又被同形
命中淹没。另一处:自动注入只给实体名索引,而规则本体长在关系上(relation_type + object),
即使命中名字也拿不到「纯文本不用Markdown」这句正文。

改动:把「触发条件」升成一等索引维度。

- schema:新增 scenes(key) + scene_refs(scene_id, kind, ref_id, weight),
  kind ∈ relation|entity。刻意不建外键:节点可能先于引用被清理,
  悬空引用由读取侧 JOIN 过滤,级联删除会把清理变成跨表事务。
- 场景键是分层字符串(`/` 分隔,由宽到窄):chan:qq、chan:qq/peer:group_123、
  tool:qq_get_message。NormalizeSceneKey 归一(小写、空白/标点→_、按 `/` 分层),
  空白不算层级——否则「老大2026-09-04 12:27 QQ私聊图片」这种来源名会被拆成伪层级。
- 写入即挂场景:Triple 新增 Scene 字段,commit() 在同一事务里把「关系 + 两端实体」
  挂到场景上(同事务是必须的:关系进库但引用丢了 = 这条记忆永远无声地召不回来)。
- 召回:RecallByScene 前缀匹配(chan:qq 取回 chan:qq 及所有更窄场景;用 `/` 兜底
  防止 chan:qq 吞掉 chan:qq2),按 weight(=写入置信度)降序,返回**关系全文 + 原句**。
- 注入:BuildContextInScene 在词法/向量之外叠加场景路,FormatContext 把场景块排在
  最前(规则对行为的约束强于话题相关的实体名),上限 8 条 + 原句截断 60 字;
  场景实体不在【记忆索引】里重复占位。BuildContext(input) 保持原语义(无场景)。
- 当前场景推导:payload.scene 显式声明 > 通道(chan:qq)> 工具(tool:qq_get_message),
  并列命中不取交集。qq 通道本身 RecallPolicy=none(到达的是中断元文本),
  真正召回在 qq_get_message 工具上——现在那一步同时带上 chan:qq 与 tool:qq_get_message。
- 写入侧:memory_commit 新增 scene 参数(逐条 triples[].scene 优先,顶层 scene 作批次默认);
  docToTriples 按文档来源自动带 chan:<source>(QQ 归档的知识天然属于 QQ 场面)。
  不做自动猜测:猜错的场景会把无关记忆钉死,之后每次进入该场面都被注入。
- 存量引导:memgc -tag-scene <键> -entity-glob <GLOB>。用 GLOB 而非 LIKE——
  LIKE 对 ASCII 不区分大小写,`%QQ%` 会把对象带 /home/newqqagent 的路径类记忆
  (生产数据目录、email-mcp、dify-ops 路径…实测 7 条)一起卷进 QQ 场景。
- 清理对齐:PurgeNoise/PurgeOrphans 之后顺带删悬空场景引用,并提供
  PurgeStaleSceneRefs;memgc -scene-stats 看场景规模。

验证:go build/vet 干净,go test -count=1 ./... 全绿。
新增用例:场景键归一(含超长/分层/空白)、写入即挂场景(两端实体进、未标的实体不进)、
前缀语义(含 chan:qq2 反例)、weight 排序与 limit、GLOB 存量引导(dry-run 不写库)、
清理后无悬空引用、场景注入面(关系全文+原句+不在索引重复占位)、
agent 侧 sceneKeysFor 优先级(显式声明 > 通道 > 工具、数组形式、nil 安全)。

生产库实测(先 sqlite3 .backup 到 graph.db.bak-20260915-081043 再写):
把 22 条 QQ 相关关系标进 chan:qq(GLOB *QQ* 19 条 + *qq_* 3 条)。同一批输入前后对比:
- 「在吗」:改前注入 17 个无关实体;改后场景块直接给出「QQ回复禁用Markdown格式
  --规定--> 纯文本不用Markdown」等规则正文(零字面重合也能召回)。
- 「QQ回复格式」:改前规则排第 32 被截掉;改后排在场景块首位。
- 「帮我发个语音」:场景规则置顶,词法路的 qq通道语音输入 等仍在其后。
2026-09-15 08:13:52 +08:00

523 lines
16 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 core
import (
"fmt"
"log"
"runtime/debug"
"strings"
"time"
agentIO "gitcode.com/JianFeeeee/HomeAgent/internal/agent/io"
"gitcode.com/JianFeeeee/HomeAgent/internal/memory"
"gitcode.com/JianFeeeee/HomeAgent/internal/memory/document"
"gitcode.com/JianFeeeee/HomeAgent/internal/memory/vector"
"gitcode.com/JianFeeeee/HomeAgent/internal/nlp"
)
type ConsolidationTask struct {
Type string `json:"type"`
Reason string `json:"reason"`
Data interface{} `json:"data"`
}
func (a *Agent) enqueueConsolidationTask(task ConsolidationTask) {
msg := fmt.Sprintf(
"【记忆整理任务】\n类型: %s\n说明: %s\n\n注意:\n1. 仅使用 memory_merge 合并实体,或使用 memory_block_merge 标记不合并\n2. 不要使用 memory_commit 写入新的三元组\n3. 不要从这段任务文本中提取任何信息写入图库\n4. 只需要做出合并/不合并的判断并执行对应工具",
task.Type, task.Reason,
)
a.injectSelf(msg)
log.Printf("[agent] enqueued consolidation task: %s", task.Reason)
}
// ──────────────────────────────────────────────
// 四个独立心跳循环,各自拥有独立的 ticker 和配置
// ──────────────────────────────────────────────
// distillLoop 上下文裁剪(L1→L2),使用 distillInterval
func (a *Agent) distillLoop() {
defer func() {
if r := recover(); r != nil {
log.Printf("[agent] distillLoop panic recovered: %v\n%s", r, debug.Stack())
time.Sleep(time.Second)
go a.distillLoop()
}
}()
if a.docStore == nil {
return
}
ticker := time.NewTicker(a.distillInterval)
defer ticker.Stop()
for {
select {
case <-ticker.C:
log.Printf("[agent] heartbeat distill tick")
a.distillContext()
a.autoReloadPlugins()
case <-a.ctx.Done():
return
}
}
}
// archiveLoop 冷文档归档(L2→L3),使用 archiveInterval
func (a *Agent) archiveLoop() {
defer func() {
if r := recover(); r != nil {
log.Printf("[agent] archiveLoop panic recovered: %v\n%s", r, debug.Stack())
time.Sleep(time.Second)
go a.archiveLoop()
}
}()
if a.memory == nil {
return
}
ticker := time.NewTicker(a.archiveInterval)
defer ticker.Stop()
for {
select {
case <-ticker.C:
log.Printf("[agent] heartbeat archive tick")
a.archiveColdDocs()
case <-a.ctx.Done():
return
}
}
}
// mergeLoop 实体合并检测(GraphDB → LLM 裁决),使用 mergeInterval
func (a *Agent) mergeLoop() {
defer func() {
if r := recover(); r != nil {
log.Printf("[agent] mergeLoop panic recovered: %v\n%s", r, debug.Stack())
time.Sleep(time.Second)
go a.mergeLoop()
}
}()
if a.memory == nil {
return
}
ticker := time.NewTicker(a.mergeInterval)
defer ticker.Stop()
for {
select {
case <-ticker.C:
log.Printf("[agent] heartbeat merge tick")
a.detectEntityMerge()
case <-a.ctx.Done():
return
}
}
}
// reviewLoop 关系复审(GraphDB → ClearSentenceID → CleanupOrphanedSentences),使用 reviewInterval
func (a *Agent) reviewLoop() {
defer func() {
if r := recover(); r != nil {
log.Printf("[agent] reviewLoop panic recovered: %v\n%s", r, debug.Stack())
time.Sleep(time.Second)
go a.reviewLoop()
}
}()
if a.memory == nil {
return
}
ticker := time.NewTicker(a.reviewInterval)
defer ticker.Stop()
for {
select {
case <-ticker.C:
log.Printf("[agent] heartbeat review tick")
a.reviewRelations()
case <-a.ctx.Done():
return
}
}
}
// ──────────────────────────────────────────────
// 蒸馏逻辑
// ──────────────────────────────────────────────
func (a *Agent) distillContext() {
if a.docStore == nil {
return
}
n := a.context.Len()
if n > a.maxContextSize*2 {
archived := a.context.Prune("", a.maxContextSize, a.docStore)
if archived > 0 {
log.Printf("[agent] distill: pruned %d low-relevance events to document memory (total=%d)", archived, n)
}
}
}
// ──────────────────────────────────────────────
// 冷文档归档:docStore → GraphDB (L3→L4)
// ──────────────────────────────────────────────
func (a *Agent) archiveColdDocs() {
if a.memory == nil {
return
}
log.Printf("[agent] cold doc archival start")
if a.indexer != nil {
if err := a.indexer.Sync(); err != nil {
log.Printf("[agent] indexer sync error: %v", err)
}
}
if a.docStore != nil {
a.docStore.Reindex()
}
if a.docStore != nil {
coldDocs := a.docStore.FindColdDocs(72*time.Hour, 2)
for _, doc := range coldDocs {
triples := docToTriples(doc, a.embedder)
if len(triples) == 0 {
continue
}
ec, rc, blocks, err := a.commitTriplesWithMedia(triples, string(a.id)+"_doc_archival", 0, doc.Blocks)
if err != nil {
log.Printf("[agent] doc→graph archival error: %v", err)
continue
}
// 归档的实质是「信息从 L2 搬到 L3」。一条实体、一条关系都没写进
// 图库时,信息并没有搬过去,此时删文档等于直接丢数据。
//
// 这不是理论情形:Commit 会静默跳过实体名不合法的三元组
//(validEntityName 要求 2–50 字符),而 LLM 生成的长描述几乎
// 提不出合规实体名——实测 456 字图片描述得到 0 entities 0
// relations,随后文档被删、媒体引用被释放、blob 被 GC 清掉,
// 图片与描述彻底消失。保留文档,下一轮再试。
if ec == 0 && rc == 0 {
log.Printf("[agent] doc→graph: %s 未写入任何实体/关系,保留文档待下轮重试"+
"(三元组 %d 条全被实体名校验拒绝)", doc.ID, len(triples))
continue
}
log.Printf("[agent] doc→graph: %s → %d entities, %d relations, %d blocks", doc.ID, ec, rc, blocks)
// 文档持有的一等块写入 L3,并以 document --contains--> block 边关联;
// 块 ID 原样保留(迁移而非重建)。块迁走后删除文档即完成迁移。
if len(doc.Blocks) > 0 {
if bound := a.linkBlocksToDocument(doc.ID, doc.Blocks); bound != len(doc.Blocks) {
log.Printf("[agent] doc→graph: %s 块迁移不完整 (%d/%d),保留文档待下轮重试",
doc.ID, bound, len(doc.Blocks))
continue
}
}
a.docStore.Remove(doc.ID)
}
}
}
// ──────────────────────────────────────────────
// 实体合并检测:GraphDB → LLM 裁决
// ──────────────────────────────────────────────
func (a *Agent) detectEntityMerge() {
if a.memory == nil {
return
}
log.Printf("[agent] entity merge detection start")
result, err := a.memory.Recall(nil, nil, 1, "")
if err != nil || result == nil || len(result.Entities) < 2 {
return
}
llmCandidates := 0
maxCandidates := 5
for i := 0; i < len(result.Entities) && llmCandidates < maxCandidates; i++ {
for j := i + 1; j < len(result.Entities) && llmCandidates < maxCandidates; j++ {
ea, eb := result.Entities[i].Name, result.Entities[j].Name
if ea > eb {
ea, eb = eb, ea
}
key := ea + "||" + eb
// 跳过已标记"不合并"的实体对
a.noMergeMu.Lock()
rounds, ok := a.noMergeMarkers[key]
if ok {
rounds--
if rounds <= 0 {
delete(a.noMergeMarkers, key)
} else {
a.noMergeMarkers[key] = rounds
}
}
a.noMergeMu.Unlock()
if ok {
continue
}
// 复合相似度:字符二元组 + 语义向量(仅增强检测,不做自动合并)
sim := entitySimilarity(result.Entities[i].Name, result.Entities[j].Name)
semSim := entitySemanticSimilarity(result.Entities[i].Name, result.Entities[j].Name, a.embedder)
if semSim > sim {
sim = semSim
}
if sim > 0.75 {
llmCandidates++
a.enqueueConsolidationTask(ConsolidationTask{
Type: "entity_merge",
Reason: fmt.Sprintf(
"实体「%s」(类型:%s, 提及%d次) 与「%s」(类型:%s, 提及%d次) 相似度 %.0f%%,可能指代同一事物,请判断是否需要合并",
result.Entities[i].Name, result.Entities[i].Type, result.Entities[i].MentionCount,
result.Entities[j].Name, result.Entities[j].Type, result.Entities[j].MentionCount,
sim*100,
),
Data: map[string]interface{}{
"entity_a": result.Entities[i].Name,
"entity_a_type": result.Entities[i].Type,
"entity_a_mentions": result.Entities[i].MentionCount,
"entity_b": result.Entities[j].Name,
"entity_b_type": result.Entities[j].Type,
"entity_b_mentions": result.Entities[j].MentionCount,
"similarity": sim,
},
})
}
}
}
if llmCandidates > 0 {
log.Printf("[agent] entity merge: %d merge candidates sent for LLM decision", llmCandidates)
} else {
log.Printf("[agent] entity merge: no similar entities found")
}
}
// ──────────────────────────────────────────────
// 关系复审:GraphDB → ClearSentenceID → CleanupOrphanedSentences
// ──────────────────────────────────────────────
func (a *Agent) reviewRelations() {
if a.memory == nil {
return
}
log.Printf("[agent] relation review start")
reviewCount := 0
const maxReviewBatch = 5
relResult, err := a.memory.Recall(nil, nil, 1, "")
if err != nil || relResult == nil {
return
}
for _, rel := range relResult.Relations {
if reviewCount >= maxReviewBatch {
break
}
if rel.SentenceID == 0 || rel.SentenceText == "" {
continue
}
a.enqueueConsolidationTask(ConsolidationTask{
Type: "relation_review",
Reason: fmt.Sprintf(
"【关系复审】原始句子: '%s'\n当前三元组: (%s → %s → %s) 置信度 %.2f\n请判断是否需要修正(如相对引用未解析、主宾颠倒、噪音三元组等),如需修正请用 memory_edit 工具",
rel.SentenceText, rel.SourceName, rel.RelationType, rel.TargetName, rel.Confidence,
),
Data: map[string]interface{}{
"relation_id": rel.ID,
"source": rel.SourceName,
"relation_type": rel.RelationType,
"target": rel.TargetName,
"confidence": rel.Confidence,
"sentence": rel.SentenceText,
},
})
// 清除句子引用(复审后解除关联)
if err := a.memory.ClearSentenceID(rel.ID); err != nil {
log.Printf("[agent] clear sentence_id for relation %d: %v", rel.ID, err)
}
reviewCount++
}
if reviewCount > 0 {
// 清理无引用的句子
if deleted, err := a.memory.CleanupOrphanedSentences(); err != nil {
log.Printf("[agent] cleanup orphaned sentences: %v", err)
} else if deleted > 0 {
log.Printf("[agent] cleanup %d orphaned sentences", deleted)
}
log.Printf("[agent] relation review: %d relations sent for review", reviewCount)
}
}
// entitySemanticSimilarity 使用词嵌入向量余弦相似度计算实体名语义相似度
func entitySemanticSimilarity(a, b string, embedder *memory.StaticEmbedder) float64 {
if a == "" || b == "" || embedder == nil || !embedder.Loaded() {
return 0
}
va := embedder.Vectorize(a)
vb := embedder.Vectorize(b)
if len(va) == 0 || len(vb) == 0 {
return 0
}
return vector.CosineSimilarity(va, vb)
}
func entitySimilarity(a, b string) float64 {
if a == "" || b == "" {
return 0
}
if a == b {
return 1.0
}
runesA, runesB := []rune(a), []rune(b)
if len(runesA) < 2 || len(runesB) < 2 {
if len(runesA) == len(runesB) && len(runesA) == 1 {
if runesA[0] == runesB[0] {
return 1.0
}
}
return 0
}
setA := make(map[string]bool)
for i := 0; i < len(runesA)-1; i++ {
setA[string(runesA[i:i+2])] = true
}
setB := make(map[string]bool)
for i := 0; i < len(runesB)-1; i++ {
setB[string(runesB[i:i+2])] = true
}
intersect := 0
for bg := range setA {
if setB[bg] {
intersect++
}
}
union := len(setA) + len(setB) - intersect
if union <= 0 {
return 0
}
return float64(intersect) / float64(union)
}
func docToTriples(doc *document.Doc, embedder nlp.Vectorizer) []memory.Triple {
var triples []memory.Triple
if doc == nil {
return triples
}
if doc.Source == "graph" || doc.Source == "" {
return nil
}
// 文档归档的知识是有**来源场面**的:来自 QQ 的对话归档,其三元组就该
// 钉在 chan:qq 上。这样「又来一条 QQ 消息」时,这批知识靠场景就能取回,
// 不必指望本轮措辞与它们字面重合。
docScene := memory.ChannelScene(doc.Source)
isArchivedContext := doc.Meta != nil && doc.Meta["is_archived_context"] == "true"
// 文档元数据:仅当 summary 合理(非空、非模板化、长度适中)时才写「主题」
if !isArchivedContext && doc.Summary != "" && len([]rune(doc.Summary)) < 80 && !isTemplateSummary(doc.Summary) {
triples = append(triples, memory.Triple{
Subject: "文档",
SubjectType: "Concept",
Relation: "主题",
Object: doc.Summary,
ObjectType: "Topic",
Confidence: 1.0,
Scene: docScene,
})
}
// 媒体不再参与三元组:它作为一等块由 linkBlocksToDocument
// 写入 L3 并以 document --contains--> block 边关联,
// 不经过文本描述与 NLP 提取器。
// NLP 通用提取
e := nlp.NewExtractor(nil)
if embedder != nil {
e.SetEmbedder(embedder)
}
result := e.Extract(doc.Content)
if result != nil {
for _, nt := range result.Triples {
mt := nlp.ToMemoryTriple(nt)
if mt.Subject != "" && mt.Relation != "" && mt.Object != "" {
mt.Scene = docScene
triples = append(triples, mt)
}
}
}
// 仅当来源非归档上下文且非空时写「来源」——归档文档写死模板三元组属于垃圾
if doc.Source != "" && doc.Source != "context_archived" {
triples = append(triples, memory.Triple{
Subject: "文档",
SubjectType: "Concept",
Relation: "来源",
Object: doc.Source,
ObjectType: "Source",
Confidence: 1.0,
Scene: docScene,
})
}
// 噪音闸门:NLP 提取器不认常用词(「结果 / 什么 / 待命」都能当主语),
// 而落库闸门 validEntityName 只管名字像不像名字。这一层是防止
// 「每个文档的常用词都变成实体」的唯一防线(CutExact 时代的那层已随
// 提取器换代丢失,见 memory.IsNoiseEntity 的说明)。
return memory.FilterNoiseTriples(triples)
}
// isTemplateSummary 识别 summarizeEntries 生成的模板化摘要
// (形如「来自 N 个来源的 M 条对话 (src1, src2) 涉及: kw1, kw2」),
// 这类摘要无独立信息量,不应作为「主题」实体写入图库。
func isTemplateSummary(s string) bool {
if s == "" {
return true
}
return strings.HasPrefix(s, "来自 ") && strings.Contains(s, "条对话")
}
func (a *Agent) emitMemoryCandidate(source, input, response string, toolResults []ToolResultItem, toolsUsed []string) {
a.io.EmitOutput("memory", "memory_candidate", map[string]interface{}{
"source": source,
"input": input,
"response": response,
"tool_results": toolResults,
"tools_used": toolsUsed,
"agent_id": string(a.id),
"timestamp": time.Now().Unix(),
})
}
func (a *Agent) processConsolidation(evt *agentIO.InputEvent, input string) {
start := time.Now()
stageCtx := a.stageCtxFromInput(input, evt.Source, "")
stageCtx.Extra["output_channel"] = evt.OutputChannel
a.injectSourceContext(stageCtx, evt)
_, toolsUsed, _, err := a.process(input, stageCtx)
if err != nil {
log.Printf("[agent] consolidation error: %v", err)
return
}
log.Printf("[agent] consolidation done (%dms, tools=%v)", time.Since(start).Milliseconds(), toolsUsed)
}