Files
HomeAgent/internal/memory/scene.go
JianFeeeee 49695c38f3 feat(memory): 场景双通道——主动声明与被动涌现并存,且互不吞噬
按「声明式的也要支持,相当于主动被动两条路」落实。此前两者只是恰好并存,
没有边界,实测会互相吃掉(下面的坑就是)。

- Triple.Scenes []string(多值):一轮写下的记忆**两条路都挂**。
  只挂一条会丢东西——只挂声明则细粒度唤起丢失,只挂涌现则首次交互
  (场景还没长出来)没有兜底。单值 Scene 保留兼容。
- TurnScene:Primary 用于写(优先涌现场景,首次退到声明场景兜底),
  Keys 是两条路的并集,用于召回(声明+涌动的场景一起进 RecallByScene)。
- EnterSceneWithHint:主动路 EnsureScene(声明即建场景,不等第二次),
  被动路 EnterScene(指纹聚类)。写侧由 executeToolCall 把本轮场景集合
  传给 memory_commit,模型不需要知道"场景"这回事。

踩到并修掉的坑(两条路互相吞噬):
  最初让声明场景也吸收**整轮指纹**,于是 chan:qq 的相似度永远是 1.0,
  把后续所有同类轮次全部吃掉 → 被动路再也长不出更细的场面,
  实测 turn2.Emergent=true 但 Primary 仍是 chan:qq、没有 auto: 场景。
  修法:给场景加 origin(declared/emergent):
  - 被动聚类只认 origin='emergent' 的场景(声明场景不进相似度空间);
  - 声明场景的特征**只从键自身解析**(chan:qq/peer:group_1 → {chan:qq, peer:group_1}),
    白名单 kind(chan/peer/peer_group/tool/topic/part),不猜——
    「老大2026-09-04_12:27_qq私聊图片」里的 12:27 也是 kind:value 形态,
    放进特征空间就是往相似度里灌垃圾(有测试钉住)。
  - 声明路的泛化靠**层级键前缀**(chan:qq 覆盖 chan:qq/peer:x),机制各归各。
- memgc -scene-stats 增加 [declared|emergent] 与 strength/features 两栏,
  可直接观察两条路各自在长什么。

新增/改写用例:
- TestDeclaredAndEmergentBothLearn:首次交互兜底到声明场景 → 第 2 轮长出
  细粒度涌现场景且**优先用于写入** → 声明场景不进相似度空间(防止压死被动路)
  但仍走声明键取回 → 两条路都进召回集合 → 声明键特征解析与白名单。
- TestEffectiveScenes:多值+单值合并去重保序。

go build/vet 干净,go test -count=1 ./... 全绿。
2026-09-15 09:37:29 +08:00

583 lines
19 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 memory
import (
"database/sql"
"fmt"
"sort"
"strings"
"time"
"unicode"
)
// ──────────────────────────────────────────────
// 场景式关联召回
//
// 词法与向量召回都建立在「字面/语义相似」上,而带**条件**的记忆天生不吃这一套:
//
// QQ回复禁用Markdown格式 --规定--> 纯文本不用Markdown
// 老大 --偏好--> QQ回复禁用Markdown格式
//
// 这些规则约束的是「回 QQ 消息这个场面」,而不是某个话题。用户措辞里没出现
// 「QQ」「Markdown」时它们召不回来而用户只说了「QQ」时词法路又会把上百个
// 含 QQ 的实体按建表顺序排前面把规则本体挤出注入预算实测输入「QQ回复格式」
// 命中 148 个实体,规则排第 32注入只取前 5 —— 规则根本没进去)。
//
// 场景引用把「触发条件」变成一等索引:节点记住自己属于哪个场面,
// 场面重现时按场景直接取回,与措辞无关。
//
// 场景键的形态是**分层字符串**,用 `/` 分隔,由宽到窄:
//
// chan:qq 通道级(在 QQ 上收发消息)
// chan:qq/peer:group_1027993713 再窄一层(具体群)
// tool:qq_get_message 工具级(取回消息正文这一步)
// chan:doc/src:qq 文档归档的来源
//
// 召回按**前缀**匹配:当前场景 `chan:qq` 会取回它自己以及所有更窄的场景
// `chan:qq/...`)——越窄的场景越具体,不该被漏掉;反向不成立。
// ──────────────────────────────────────────────
// maxSceneKeyLen 是场景键的长度上限。场景键要进索引、要参与前缀比较,
// 过长只说明有人把正文塞进了键里(那就该用句子/实体,而不是场景)。
const maxSceneKeyLen = 96
// NormalizeSceneKey 规范化场景键:按 `/` 分层、每层小写、层内空白与非法字符
// 归一成 `_`(连续多个只留一个)。
//
// 为什么要归一:场景键是**索引键**`chan:QQ` 与 `chan:qq` 必须是同一个场景,
// 否则同一条规则会因为写入时大小写不同而分裂成两个召不齐的场景。
// 为什么空白不算层级分隔符来源名里天然带空格如「老大2026-09-04 12:27
// QQ私聊图片」把它当层级会把一个平面名字拆成三层伪层级。
// 归一结果为空(全是非法字符)时返回空串,调用方应视为「没有场景」。
func NormalizeSceneKey(key string) string {
key = strings.TrimSpace(key)
if key == "" {
return ""
}
key = strings.ToLower(key)
parts := strings.Split(key, "/")
segs := make([]string, 0, len(parts))
for _, p := range parts {
seg := normalizeSceneSegment(p)
if seg != "" {
segs = append(segs, seg)
}
}
out := strings.Join(segs, "/")
if len(out) > maxSceneKeyLen {
out = strings.TrimRight(out[:maxSceneKeyLen], "/")
}
return out
}
// normalizeSceneSegment 归一单层场景:小写、空白与非法字符 → `_`(压缩连续)。
func normalizeSceneSegment(seg string) string {
var b strings.Builder
lastUnderscore := false
for _, r := range seg {
switch {
case r >= 'a' && r <= 'z', r >= '0' && r <= '9', r >= '\u4e00' && r <= '\u9fff',
r == ':' || r == '-' || r == '.':
b.WriteRune(r)
lastUnderscore = false
case r == '_' || unicode.IsSpace(r):
if !lastUnderscore && b.Len() > 0 {
b.WriteRune('_')
lastUnderscore = true
}
default:
// 其它字符(标点、表情)归一成 `_` 而不是静默丢弃:
// 「A/B」与「A_B」是两个不同的来源不能塌成一个场景。
if !lastUnderscore && b.Len() > 0 {
b.WriteRune('_')
lastUnderscore = true
}
}
}
return strings.Trim(b.String(), "_")
}
// ChannelScene 由输入/输出通道名构造场景键(`qq` → `chan:qq`)。
func ChannelScene(source string) string {
s := NormalizeSceneKey(source)
if s == "" {
return ""
}
return "chan:" + s
}
// ToolScene 由工具名构造场景键(`qq_get_message` → `tool:qq_get_message`)。
func ToolScene(tool string) string {
s := NormalizeSceneKey(tool)
if s == "" {
return ""
}
return "tool:" + s
}
// SceneStat 是单个场景的规模摘要(供 introspection / 运维观察)。
type SceneStat struct {
Key string `json:"key"`
Refs int `json:"refs"`
Relations int `json:"relations"`
Entities int `json:"entities"`
// Strength 是该场景被重现强化的次数Features 是它长出的特征数。
// 两者一起说明「这个场景是不是真的在涌现」,而不是被一次性写出来的。
Strength int `json:"strength"`
Features int `json:"features"`
// Origin 是这条场景来自哪条路declared主动声明或 emergent被动涌现
Origin string `json:"origin"`
UpdatedAt time.Time `json:"updated_at"`
}
// SceneRecall 是一次场景召回的产物。
type SceneRecall struct {
Scenes []string `json:"scenes"`
Relations []Relation `json:"relations"`
Entities []Entity `json:"entities"`
// Blocks 是该场景下的一等记忆块(图/音/文)。
//
// 为什么场景要能取回块:块是流水线里最细的子项目,而「那场对话里发过来的
// 那张图」只记住名字是没用的——场面重现时要把块本身带回来。
Blocks []MemoryBlock `json:"blocks,omitempty"`
// Documents 是该场景下的 L3 文档节点 id文档的场景由来源派生见 TagSceneDocument
Documents []string `json:"documents,omitempty"`
}
// tagSceneTx 在事务内把「关系 + 实体」挂到场景上(幂等 upsert
//
// weight 取关系的置信度:场景内的记忆也要能排序,置信度是目前唯一现成的
// 质量信号。重复写入同一节点只刷新 weight 与时间,不产生重复引用。
func tagSceneTx(tx *sql.Tx, sceneKey string, relationID int64, entityIDs []int64, weight float64) error {
key := NormalizeSceneKey(sceneKey)
if key == "" {
return nil
}
if relationID != 0 {
if err := tagSceneRefTx(tx, key, "relation", relationID, "", weight); err != nil {
return err
}
}
for _, eid := range entityIDs {
if eid != 0 {
if err := tagSceneRefTx(tx, key, "entity", eid, "", weight); err != nil {
return err
}
}
}
return nil
}
// tagSceneRefTx 在事务内把一个节点挂到场景上(幂等 upsert
//
// kind ∈ relation | entity | block | document。textID 供非数值主键的节点使用
// (块与文档的 id 是字符串),数值型节点传 0 并用 id。
//
// 为什么 weight 取 MAX 而不是覆盖:场景内的记忆也要能排序,置信度是目前唯一
// 现成的质量信号;同一节点被低置信度的重复写入命中的,不该把它从场景前排挤下去。
func tagSceneRefTx(tx *sql.Tx, sceneKey, kind string, id int64, textID string, weight float64) error {
key := NormalizeSceneKey(sceneKey)
if kind == "" || (id == 0 && textID == "") {
return nil
}
if key == "" {
return nil
}
if _, err := tx.Exec(
`INSERT INTO scenes (key) VALUES (?)
ON CONFLICT(key) DO UPDATE SET updated_at = CURRENT_TIMESTAMP`, key); err != nil {
return fmt.Errorf("upsert scene %q: %w", key, err)
}
var sceneID int64
if err := tx.QueryRow(`SELECT id FROM scenes WHERE key = ?`, key).Scan(&sceneID); err != nil {
return fmt.Errorf("select scene %q: %w", key, err)
}
if weight <= 0 {
weight = 1.0
}
if _, err := tx.Exec(
`INSERT INTO scene_refs (scene_id, kind, ref_id, ref_text, weight) VALUES (?, ?, ?, ?, ?)
ON CONFLICT(scene_id, kind, ref_id, ref_text)
DO UPDATE SET weight = MAX(weight, excluded.weight)`,
sceneID, kind, id, textID, weight); err != nil {
return fmt.Errorf("upsert scene ref %s/%d%s: %w", kind, id, textID, err)
}
return nil
}
// TagScene 给一批已有的关系补挂场景(存量记忆的场景标注入口)。
//
// 为什么需要「事后标注」:场景是后引入的维度,此前写下的规则(那批 QQ 规则
// 就是典型)没有任何场景引用,不补挂就永远吃不到场景召回。
func (g *GraphDB) TagScene(sceneKey string, relationIDs []int64) (int, error) {
g.mu.Lock()
defer g.mu.Unlock()
key := NormalizeSceneKey(sceneKey)
if key == "" || len(relationIDs) == 0 {
return 0, nil
}
tx, err := g.db.Begin()
if err != nil {
return 0, err
}
defer tx.Rollback()
n := 0
for _, rid := range relationIDs {
var sourceID, targetID int64
var confidence float64
if err := tx.QueryRow(
`SELECT source_id, target_id, confidence FROM relations WHERE id = ?`, rid,
).Scan(&sourceID, &targetID, &confidence); err != nil {
if err == sql.ErrNoRows {
continue
}
return n, err
}
if err := tagSceneTx(tx, key, rid, []int64{sourceID, targetID}, confidence); err != nil {
return n, err
}
n++
}
if err := tx.Commit(); err != nil {
return 0, err
}
return n, nil
}
// TagSceneByEntityGlob 把「任一端实体名匹配 GLOB pattern」的活跃关系标进场景
// 返回标注的关系数。dryRun 时只统计。
//
// 这是**存量引导**用的窄口子pattern 由调用方显式给出,不做任何自动猜测
// ——猜错的代价是把无关记忆钉死在某个场景上,之后每次进入该场景都会被注入,
// 比漏标更难发现。
//
// 为什么用 GLOB 而不是 LIKELIKE 对 ASCII **不区分大小写**,于是 `%QQ%`
// 会把对象里带 `/home/newqqagent` 的路径类记忆生产数据目录、email-mcp、
// dify-ops技能路径…实测 7 条一起卷进「QQ 场景」。GLOB 区分大小写,
// `*QQ*` 只命中真正写作 QQ 的那些名字。
func (g *GraphDB) TagSceneByEntityGlob(sceneKey, pattern string, dryRun bool) (int, error) {
key := NormalizeSceneKey(sceneKey)
if key == "" || pattern == "" {
return 0, nil
}
g.mu.Lock()
defer g.mu.Unlock()
rows, err := g.db.Query(
`SELECT r.id, r.source_id, r.target_id, r.confidence
FROM relations r
JOIN entities e1 ON r.source_id = e1.id
JOIN entities e2 ON r.target_id = e2.id
WHERE r.status = 'active' AND (e1.name GLOB ? OR e2.name GLOB ?)`,
pattern, pattern)
if err != nil {
return 0, err
}
type cand struct {
relID int64
sourceID, target int64
confidence float64
}
var cands []cand
for rows.Next() {
var c cand
if err := rows.Scan(&c.relID, &c.sourceID, &c.target, &c.confidence); err != nil {
rows.Close()
return 0, err
}
cands = append(cands, c)
}
rows.Close()
if err := rows.Err(); err != nil {
return 0, err
}
if dryRun {
return len(cands), nil
}
tx, err := g.db.Begin()
if err != nil {
return 0, err
}
defer tx.Rollback()
for _, c := range cands {
if err := tagSceneTx(tx, key, c.relID, []int64{c.sourceID, c.target}, c.confidence); err != nil {
return 0, err
}
}
if err := tx.Commit(); err != nil {
return 0, err
}
return len(cands), nil
}
// RecallByScene 按当前场景取回被钉在该场景上的记忆(前缀匹配,越窄越算命中)。
//
// 排序weight=写入时置信度)降序 → 关系时间降序。取回的是**关系全文**
// (含 relation_type 与 JOIN 出的原句),不只是实体名——带条件的规则本体
// 长在关系上,只给名字等于没召回。
//
// limit 同时约束关系数与实体数,避免一个场景把注入预算吃光。
func (g *GraphDB) RecallByScene(scenes []string, limit int) (*SceneRecall, error) {
out := &SceneRecall{}
var keys []string
seen := make(map[string]bool)
for _, s := range scenes {
k := NormalizeSceneKey(s)
if k == "" || seen[k] {
continue
}
seen[k] = true
keys = append(keys, k)
}
if len(keys) == 0 {
return out, nil
}
out.Scenes = keys
if limit <= 0 {
limit = 8
}
g.mu.RLock()
defer g.mu.RUnlock()
// 前缀条件:(key = ? OR key LIKE ? || '/%'),用 '/' 兜底防止
// `chan:qq` 误吞 `chan:qq2` 这种同前缀但不同层的场景。
conds := make([]string, 0, len(keys))
args := make([]interface{}, 0, len(keys)*2)
for _, k := range keys {
conds = append(conds, `(s.key = ? OR s.key LIKE ? || '/%')`)
args = append(args, k, k)
}
where := "(" + strings.Join(conds, " OR ") + ")"
relQuery := `SELECT r.id, r.source_id, r.target_id, e1.name, e2.name,
r.relation_type, r.confidence, r.status, r.session_id,
r.turn_id, r.created_at, COALESCE(r.date_bucket, ''),
COALESCE(r.sentence_id, 0), COALESCE(sn.text, ''),
MAX(sr.weight) AS w
FROM scene_refs sr
JOIN scenes s ON sr.scene_id = s.id
JOIN relations r ON sr.kind = 'relation' AND sr.ref_id = r.id
JOIN entities e1 ON r.source_id = e1.id
JOIN entities e2 ON r.target_id = e2.id
LEFT JOIN sentences sn ON r.sentence_id = sn.id
WHERE ` + where + ` AND r.status = 'active'
GROUP BY r.id
ORDER BY w DESC, r.updated_at DESC, r.id DESC
LIMIT ?`
relArgs := append(append([]interface{}{}, args...), limit)
rows, err := g.db.Query(relQuery, relArgs...)
if err != nil {
return nil, err
}
for rows.Next() {
var rel Relation
var w float64
if err := rows.Scan(&rel.ID, &rel.SourceID, &rel.TargetID, &rel.SourceName, &rel.TargetName,
&rel.RelationType, &rel.Confidence, &rel.Status, &rel.SessionID,
&rel.TurnID, &rel.CreatedAt, &rel.DateBucket, &rel.SentenceID, &rel.SentenceText, &w); err != nil {
rows.Close()
return nil, err
}
out.Relations = append(out.Relations, rel)
}
rows.Close()
if err := rows.Err(); err != nil {
return nil, err
}
entQuery := `SELECT e.id, e.name, e.type, e.mention_count, e.created_at, e.updated_at, MAX(sr.weight) AS w
FROM scene_refs sr
JOIN scenes s ON sr.scene_id = s.id
JOIN entities e ON sr.kind = 'entity' AND sr.ref_id = e.id
WHERE ` + where + `
GROUP BY e.id
ORDER BY w DESC, e.mention_count DESC, e.id DESC
LIMIT ?`
entArgs := append(append([]interface{}{}, args...), limit)
erows, err := g.db.Query(entQuery, entArgs...)
if err != nil {
return nil, err
}
for erows.Next() {
var e Entity
var w float64
if err := erows.Scan(&e.ID, &e.Name, &e.Type, &e.MentionCount, &e.CreatedAt, &e.UpdatedAt, &w); err != nil {
erows.Close()
return nil, err
}
out.Entities = append(out.Entities, e)
}
erows.Close()
if err := erows.Err(); err != nil {
return nil, err
}
blockQuery := `SELECT b.id, b.modality, b.text_content, b.payload_digest, b.mime,
b.size, b.width, b.height, b.fingerprint, b.source, b.tool,
COALESCE(b.scene, ''), b.created_at, b.updated_at, MAX(sr.weight) AS w
FROM scene_refs sr
JOIN scenes s ON sr.scene_id = s.id
JOIN memory_blocks b ON sr.kind = 'block' AND b.id = sr.ref_text
WHERE ` + where + `
GROUP BY b.id
ORDER BY w DESC, b.created_at DESC
LIMIT ?`
blockArgs := append(append([]interface{}{}, args...), limit)
brows, err := g.db.Query(blockQuery, blockArgs...)
if err != nil {
return nil, err
}
defer brows.Close()
for brows.Next() {
var b MemoryBlock
var w float64
if err := brows.Scan(&b.ID, &b.Modality, &b.Text, &b.PayloadDigest, &b.MIME,
&b.Size, &b.Width, &b.Height, &b.Fingerprint, &b.Source, &b.Tool,
&b.Scene, &b.CreatedAt, &b.UpdatedAt, &w); err != nil {
return nil, err
}
out.Blocks = append(out.Blocks, b)
}
if err := brows.Err(); err != nil {
return nil, err
}
docQuery := `SELECT sr.ref_text, MAX(sr.weight) AS w
FROM scene_refs sr JOIN scenes s ON sr.scene_id = s.id
WHERE ` + where + ` AND sr.kind = 'document'
GROUP BY sr.ref_text ORDER BY w DESC LIMIT ?`
docArgs := append(append([]interface{}{}, args...), limit)
drows, err := g.db.Query(docQuery, docArgs...)
if err != nil {
return nil, err
}
defer drows.Close()
for drows.Next() {
var id string
var w float64
if err := drows.Scan(&id, &w); err != nil {
return nil, err
}
out.Documents = append(out.Documents, id)
}
return out, drows.Err()
}
// SceneStats 返回各场景的规模,按引用数降序。
func (g *GraphDB) SceneStats() ([]SceneStat, error) {
g.mu.RLock()
defer g.mu.RUnlock()
rows, err := g.db.Query(
`SELECT s.key,
COUNT(sr.id),
SUM(CASE WHEN sr.kind = 'relation' THEN 1 ELSE 0 END),
SUM(CASE WHEN sr.kind = 'entity' THEN 1 ELSE 0 END),
COALESCE(s.strength, 1),
(SELECT COUNT(*) FROM scene_features f WHERE f.scene_id = s.id),
COALESCE(s.origin, 'emergent'),
s.updated_at
FROM scenes s LEFT JOIN scene_refs sr ON sr.scene_id = s.id
GROUP BY s.id ORDER BY COUNT(sr.id) DESC, s.key`)
if err != nil {
return nil, err
}
defer rows.Close()
var out []SceneStat
for rows.Next() {
var st SceneStat
var rels, ents sql.NullInt64
if err := rows.Scan(&st.Key, &st.Refs, &rels, &ents, &st.Strength, &st.Features, &st.Origin, &st.UpdatedAt); err != nil {
return nil, err
}
st.Relations = int(rels.Int64)
st.Entities = int(ents.Int64)
out = append(out, st)
}
if err := rows.Err(); err != nil {
return nil, err
}
sort.SliceStable(out, func(i, j int) bool { return out[i].Refs > out[j].Refs })
return out, nil
}
// PurgeStaleSceneRefs 清理指向已不存在节点的场景引用,返回删除数。
//
// 节点被清理PurgeNoise / PurgeOrphans / memory_delete_entity时不会级联
// 删 scene_refs见建表注释残留引用会让场景看起来很大却召回出空结果
// 也会让 SceneStats 说谎。这个函数把它们对齐。
func (g *GraphDB) PurgeStaleSceneRefs() (int, error) {
g.mu.Lock()
defer g.mu.Unlock()
return g.purgeStaleSceneRefsLocked()
}
func (g *GraphDB) purgeStaleSceneRefsLocked() (int, error) {
res, err := g.db.Exec(`DELETE FROM scene_refs WHERE
(kind = 'relation' AND ref_id NOT IN (SELECT id FROM relations))
OR (kind = 'entity' AND ref_id NOT IN (SELECT id FROM entities))
OR (kind = 'block' AND ref_text NOT IN (SELECT id FROM memory_blocks))
OR (kind = 'document' AND ref_text NOT IN (SELECT id FROM documents))`)
if err != nil {
return 0, err
}
n, _ := res.RowsAffected()
return int(n), nil
}
// ScenesOfRelation 取回一条关系当前所属的全部场景键。
//
// 用途memory_edit 是「删旧写新」——旧关系的 scene_refs 会随节点一起失效,
// 新关系若不重新挂上场景,这条记忆就**静默地脱离场景**,此后场面重现也召不回。
func (g *GraphDB) ScenesOfRelation(relationID int64) ([]string, error) {
g.mu.RLock()
defer g.mu.RUnlock()
rows, err := g.db.Query(
`SELECT s.key FROM scene_refs sr JOIN scenes s ON sr.scene_id = s.id
WHERE sr.kind = 'relation' AND sr.ref_id = ?`, relationID)
if err != nil {
return nil, err
}
defer rows.Close()
var keys []string
for rows.Next() {
var k string
if err := rows.Scan(&k); err != nil {
return nil, err
}
keys = append(keys, k)
}
return keys, rows.Err()
}
// TagSceneDocument 把 L3 文档节点挂到场景上document 层与场景模型兼容)。
//
// 文档的「场景」不由存储字段决定,而由**来源**派生chan:<source>)——存一份
// 冗余的 Doc.Scene 会随来源改名而说谎,是同一事实的第二份真相。
// 这里只登记「这份文档属于哪些场面」,供 scene 侧枚举与统计。
func (g *GraphDB) TagSceneDocument(sceneKey, docID string) error {
if docID == "" {
return nil
}
key := NormalizeSceneKey(sceneKey)
if key == "" {
return nil
}
g.mu.Lock()
defer g.mu.Unlock()
tx, err := g.db.Begin()
if err != nil {
return err
}
defer tx.Rollback()
if err := tagSceneRefTx(tx, key, "document", 0, docID, 1.0); err != nil {
return err
}
return tx.Commit()
}