Files
HomeAgent/internal/agent/core/context.go
JianFeeeee f855893d1c feat(memory): 媒体接入 L0/L2——digest 挂到对话事件,归档时引用随之转移
a822674 的 CAS 层之上把媒体真正接进记忆链路。此前 CAS 只是个孤立的
存储包,没有任何写入方。

## 媒体进入对话有两条路,两条都只把文字留给记忆

  1. 用户直接发图 → processMediaInput → mediaToBlocks
     ContextEvent.Input 只存 alt 文本("[从 qq 收到了 image]"),
     base64 随 message 数组发给模型后就丢了。
  2. 插件注入 → SetToolBlocks → process.go 的 mediaMsg
     ToolResultItem.Output 只存那句 "[已将图片注入后续对话] /tmp/x.png"。

于是下一轮起,模型能看到的只剩一句路径或一句 alt。那个文件被删、被覆盖,
或者本来就是 /tmp 下的临时产物,连线索都断了。

现在两条路在同一处收口(captureBlockMedia):从 ContentBlock 的 data URL
取出字节存进 CAS,digest 挂到当轮 ContextEvent。

## 改动

internal/agent/core/mediaref.go(新)
  - captureBlockMedia:ContentBlock → CAS。只处理 data URL——http(s) URL
    拿不到字节就无法内容寻址,而「下载它再存」会把一次对话变成一次网络
    请求(超时、鉴权、SSRF 全来了),不在本层解决。
  - stage/drainMediaDigests:媒体在 process() 期间被捕获,而承载它的
    ContextEvent 要等 process() 返回后才 Append——此刻还没有 owner_id,
    故先缓存。与既有 pendingMedia 同一手法,同受 a.mu 保护。
  - bindEventMedia:双向落地。evt.Media 让事件记得引了什么(随
    context.json 持久化),media_refs 让 CAS 知道谁在引用(GC 的判断依据)。
    只写一边的话,要么 GC 误删仍被引用的内容,要么孤儿永远清不掉。
  - mediaSummaryForEvent:把已有描述拼成一行写进 Input。这是方案 C 的
    落点——**描述文本才是持久语义记忆,blob 只是缓存**。blob 可能被容量
    GC 淘汰,但描述会一直留在 L0/L2/L3 的文本里,让「那张紫蓝红三色带图」
    几个月后仍可被检索。

ContextEvent 新增 ID 与 Media 两个字段,都是 omitempty:
  - ID 懒生成,只有真要挂媒体时才赋值。绝大多数对话没有媒体,全量生成
    会让每条事件都多一个字段进 context.json。
  - 存量 context.json 读回来两字段皆空,不影响任何既有行为(有测试)。

RelevanceContext.Prune 归档时转移引用(transferMediaRefs):
  **先挂到归档文档、再注销原事件引用**。顺序不能反——先销后挂会让引用
  计数瞬时归零,若此刻后台 GC 正在跑就会把仍被记忆引用的内容当孤儿清掉。
  为此把 Prune 内的局部类型 scored 提为包级 scoredEvent(局部类型无法
  出现在方法签名上)。

media 包新增 OwnerContext/OwnerDocument/OwnerGraphSentence 常量:
  owner_kind 进了主键,拼错一个字符就是一条永远对不上的孤立引用——
  AddRef 不报错,DropOwner 也永远匹配不到。

## 配置

core.memory.media.enabled(默认 true)、.dir、.max_mb(2048)、
.gc_interval(6h)、.gc_min_age(1h)。

关闭后全链路静默跳过,对话行为与本特性上线前完全一致(有测试)。
mediaStore 为 nil 时同理——它是记忆增强,不是对话必需品,开不起来
只记一条 warning 不阻止启动。

## 测试(11 例)

入库与 MIME 归类、http URL 跳过、nil store 全链路 no-op、音视频混合、
stage/drain 清空语义、懒生成 ID、描述作为持久记忆、**归档转移期间内容
始终可读且 refcount 不归零**、无媒体存储时归档照常、context.json
向后兼容往返。

全仓 go build / go vet / go test 通过,SDK 冻结 diff = 0。

## 尚未接入

L3 图库的 graph_sentence owner(常量已备好,无写入方)、
描述生成的后台任务(Pending() 已就绪,尚无消费者)、
媒体 GC 的定时触发(配置项已注册,尚未接 ticker)。
2026-09-04 20:53:32 +08:00

447 lines
12 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 (
"encoding/json"
"fmt"
"log"
"os"
"path/filepath"
"sort"
"strings"
"sync"
"time"
"gitcode.com/JianFeeeee/HomeAgent/internal/memory"
"gitcode.com/JianFeeeee/HomeAgent/internal/memory/document"
"gitcode.com/JianFeeeee/HomeAgent/internal/memory/media"
"gitcode.com/JianFeeeee/HomeAgent/internal/memory/vector"
sdk "gitcode.com/JianFeeeee/HomeAgent/internal/sdk"
)
type ToolResultItem struct {
Name string `json:"name"`
Output string `json:"output"`
}
type ContextEvent struct {
// ID 是事件的稳定标识媒体引用media_refs.owner_id挂在它上面。
//
// 惰性生成:只有真的要挂媒体时才赋值(见 bindEventMedia
// 全量生成会让每条事件都多一个字段进 context.json而绝大多数对话没有媒体。
// omitempty 保证存量 context.json 读回来时该字段为空,不影响任何既有行为。
ID string `json:"id,omitempty"`
Timestamp time.Time `json:"timestamp"`
Source string `json:"source"`
Input string `json:"input"`
Response string `json:"response,omitempty"`
ToolsUsed []string `json:"tools_used,omitempty"`
ToolResults []ToolResultItem `json:"tool_results,omitempty"`
// Media 是本轮对话涉及的媒体 digestsha256 十六进制)。
//
// 存 digest 而不存路径:路径会失效(/tmp 探针图、下载缓存、别的进程的
// 临时产物digest 是内容本身的身份,配合 internal/memory/media 的 CAS
// 永远能取回原始字节——只要它还没被容量 GC 淘汰。
Media []string `json:"media,omitempty"`
Vector vector.Vector `json:"-"`
}
const contextFlushInterval = 5 * time.Second
type RelevanceContext struct {
mu sync.Mutex
events []*ContextEvent
embedder *memory.StaticEmbedder
savePath string
saveTimer *time.Timer
dirty bool
toolDefLookup func(name string) *sdk.ToolDef
channelDefLookup func(name string) (sdk.ChannelDef, bool)
// mediaStore 只用于 Prune 时把媒体引用从事件转给归档文档。
// 为 nil 时引用转移静默跳过(媒体存储未启用)。
mediaStore *media.Store
}
// SetMediaStore 注入媒体存储,供 L0→L2 归档时转移媒体引用。
func (c *RelevanceContext) SetMediaStore(s *media.Store) {
c.mu.Lock()
defer c.mu.Unlock()
c.mediaStore = s
}
// transferMediaRefs 把被归档事件的媒体引用转给目标文档(调用方已持 c.mu
//
// 先挂后销:若反序,引用计数会瞬时归零,此时若后台 GC 正在跑
// 就会把仍被记忆引用的内容当孤儿清掉。
func (c *RelevanceContext) transferMediaRefs(archive []scoredEvent, docID string) {
if c.mediaStore == nil || docID == "" {
return
}
for _, s := range archive {
evt := s.event
if evt == nil || evt.ID == "" || len(evt.Media) == 0 {
continue
}
for _, d := range evt.Media {
if err := c.mediaStore.AddRef(d, media.OwnerDocument, docID); err != nil {
log.Printf("[media] 归档转移 AddRef 失败 (%s → doc %s): %v", shortDigest(d), docID, err)
}
}
if _, err := c.mediaStore.DropOwner(media.OwnerContext, evt.ID); err != nil {
log.Printf("[media] 归档转移 DropOwner 失败 (evt %s): %v", evt.ID, err)
}
}
}
func NewRelevanceContext(savePath string, embedder *memory.StaticEmbedder) *RelevanceContext {
rc := &RelevanceContext{
embedder: embedder,
savePath: savePath,
}
if savePath != "" {
rc.load()
}
return rc
}
func (c *RelevanceContext) SetToolDefLookup(fn func(name string) *sdk.ToolDef) {
c.mu.Lock()
defer c.mu.Unlock()
c.toolDefLookup = fn
}
func (c *RelevanceContext) SetChannelDefLookup(fn func(name string) (sdk.ChannelDef, bool)) {
c.mu.Lock()
defer c.mu.Unlock()
c.channelDefLookup = fn
}
func (c *RelevanceContext) load() {
data, err := os.ReadFile(c.savePath)
if err != nil {
return
}
var events []*ContextEvent
if err := json.Unmarshal(data, &events); err != nil {
return
}
for _, evt := range events {
evt.Vector = c.computeVector(evt)
}
c.events = events
}
func textForVector(evt *ContextEvent, toolDefLookup func(name string) *sdk.ToolDef, channelDefLookup func(name string) (sdk.ChannelDef, bool)) string {
var text string
switch {
case evt.Source == "agent" && evt.Response != "":
text = evt.Response
case evt.Source == "cold_storage":
text = evt.Input + " " + evt.Response
default:
text = evt.Input
}
// 计算层:应用输入通道的 Cleaner不改原文仅在计算层清洗
if channelDefLookup != nil {
if chDef, ok := channelDefLookup(evt.Source); ok && chDef.Cleaner != nil {
text = chDef.Cleaner(text)
}
if chDef, ok := channelDefLookup(evt.Source); ok && chDef.NoMemory {
return ""
}
}
// 计算层附加工具输出NoMemory 跳过,其余经 Cleaner 过滤
if toolDefLookup != nil {
noMemory := make(map[string]bool)
for _, tr := range evt.ToolResults {
def := toolDefLookup(tr.Name)
if def != nil && def.NoMemory {
noMemory[tr.Name] = true
}
}
for _, tr := range evt.ToolResults {
if noMemory[tr.Name] {
continue
}
cleaned := tr.Output
def := toolDefLookup(tr.Name)
if def != nil && def.Cleaner != nil {
cleaned = def.Cleaner(cleaned)
}
text += " " + cleaned
}
}
return memory.CleanText(text)
}
// toolOutputClean 根据工具定义的 NoMemory/Cleaner 清洗输出,用于计算层。
// 返回 "" 表示跳过NoMemory否则返回清洗后文本Cleaner 或原文)。
func (c *RelevanceContext) toolOutputClean(name, output string) string {
if c.toolDefLookup == nil {
return output
}
def := c.toolDefLookup(name)
if def == nil {
return output
}
if def.NoMemory {
return ""
}
if def.Cleaner != nil {
return def.Cleaner(output)
}
return output
}
// inputChannelClean 根据输入通道的 Def 清洗输入文本,用于计算层。
func (c *RelevanceContext) inputChannelClean(source, input string) string {
if c.channelDefLookup == nil {
return input
}
chDef, ok := c.channelDefLookup(source)
if !ok {
return input
}
if chDef.Cleaner != nil {
return chDef.Cleaner(input)
}
return input
}
// channelCleanerForDoc 返回 ChannelCleaner使 document 包在存档时能按来源查找 Cleaner。
func (c *RelevanceContext) channelCleanerForDoc() document.ChannelCleaner {
if c.channelDefLookup == nil {
return nil
}
return func(source string) func(string) string {
chDef, ok := c.channelDefLookup(source)
if !ok {
return nil
}
return chDef.Cleaner
}
}
func (c *RelevanceContext) computeVector(evt *ContextEvent) vector.Vector {
text := textForVector(evt, c.toolDefLookup, c.channelDefLookup)
if text == "" {
return nil
}
return c.embedder.Vectorize(text)
}
func (c *RelevanceContext) Save() error {
if c.savePath == "" {
return nil
}
if err := os.MkdirAll(filepath.Dir(c.savePath), 0755); err != nil {
return err
}
data, err := json.Marshal(c.events)
if err != nil {
return err
}
return os.WriteFile(c.savePath, data, 0644)
}
func (c *RelevanceContext) Append(evt ContextEvent) {
c.mu.Lock()
defer c.mu.Unlock()
evt.Vector = c.computeVector(&evt)
c.events = append(c.events, &evt)
c.save()
}
func (c *RelevanceContext) InsertByTimestamp(evt ContextEvent) {
c.mu.Lock()
defer c.mu.Unlock()
evt.Vector = c.computeVector(&evt)
idx := sort.Search(len(c.events), func(i int) bool {
return c.events[i].Timestamp.After(evt.Timestamp)
})
c.events = append(c.events, nil)
copy(c.events[idx+1:], c.events[idx:])
c.events[idx] = &evt
c.save()
}
func (c *RelevanceContext) save() error {
if c.savePath == "" {
return nil
}
if !c.dirty {
c.dirty = true
if c.saveTimer == nil {
c.saveTimer = time.AfterFunc(contextFlushInterval, c.flush)
} else {
c.saveTimer.Reset(contextFlushInterval)
}
}
return nil
}
func (c *RelevanceContext) flush() {
c.mu.Lock()
defer c.mu.Unlock()
if !c.dirty {
return
}
data, err := json.Marshal(c.events)
if err != nil {
return
}
if err := os.WriteFile(c.savePath, data, 0644); err != nil {
return
}
c.dirty = false
}
// scoredEvent 是 Prune 里按相关度排序的事件。
//
// 提为包级类型(原先是 Prune 内的局部类型transferMediaRefs 需要
// 把待归档列表传进去,局部类型无法出现在方法签名上。
type scoredEvent struct {
event *ContextEvent
score float64
idx int
}
func (c *RelevanceContext) Prune(currentInput string, topK int, docStore *document.Store) int {
c.mu.Lock()
defer c.mu.Unlock()
if len(c.events) <= topK {
return 0
}
pCount := 10
if pCount > len(c.events) {
pCount = len(c.events)
}
protected := c.events[len(c.events)-pCount:]
candidates := c.events[:len(c.events)-pCount]
if len(candidates) == 0 {
return 0
}
queryVec := c.embedder.VectorizeClean(currentInput)
scoredEvents := make([]scoredEvent, len(candidates))
for i, evt := range candidates {
score := vector.CosineSimilarity(queryVec, evt.Vector)
scoredEvents[i] = scoredEvent{event: evt, score: score, idx: i}
}
sort.Slice(scoredEvents, func(i, j int) bool {
return scoredEvents[i].score > scoredEvents[j].score
})
keepCount := topK - len(protected)
if keepCount < 0 {
keepCount = 0
}
keep := scoredEvents
if len(keep) > keepCount {
keep = keep[:keepCount]
}
archive := scoredEvents[keepCount:]
c.events = make([]*ContextEvent, 0, len(keep)+len(protected))
for _, s := range keep {
c.events = append(c.events, s.event)
}
c.events = append(c.events, protected...)
sort.Slice(c.events, func(i, j int) bool {
return c.events[i].Timestamp.Before(c.events[j].Timestamp)
})
archived := 0
if docStore != nil && len(archive) > 0 {
entries := make([]document.ContextEntry, len(archive))
for i, s := range archive {
entries[i] = document.ContextEntry{
Timestamp: s.event.Timestamp,
Source: s.event.Source,
Content: s.event.Input,
Response: s.event.Response,
ToolResults: convertToolResults(s.event.ToolResults),
}
}
doc, err := docStore.ContextToDoc("context_archived", entries, c.embedder, nil, c.toolOutputClean, c.channelCleanerForDoc())
if err == nil && doc != nil {
archived = len(entries)
// 媒体引用随事件一起从 L0 转到 L2先把引用挂到归档文档上
// 再注销原事件的引用。顺序不能反——先销后挂会让引用计数
// 瞬时归零,若此时 GC 正在跑(后台任务)就会把仍被记忆引用的
// 内容当孤儿清掉。
c.transferMediaRefs(archive, doc.ID)
}
}
c.save()
return archived
}
func (c *RelevanceContext) Format() string {
c.mu.Lock()
defer c.mu.Unlock()
if len(c.events) == 0 {
return ""
}
var sb strings.Builder
sb.WriteString("【近期事件】\n")
for _, e := range c.events {
sb.WriteString(fmt.Sprintf("[%s] %s: %s", e.Timestamp.Format("15:04:05"), e.Source, e.Input))
if e.Response != "" {
sb.WriteString(fmt.Sprintf(" → %s", truncateStr(e.Response, 80)))
}
sb.WriteString("\n")
}
return sb.String()
}
func (c *RelevanceContext) Recent(n int) []ContextEvent {
c.mu.Lock()
defer c.mu.Unlock()
if n <= 0 || n > len(c.events) {
n = len(c.events)
}
result := make([]ContextEvent, n)
for i, evt := range c.events[len(c.events)-n:] {
result[i] = *evt
}
return result
}
func (c *RelevanceContext) Len() int {
c.mu.Lock()
defer c.mu.Unlock()
return len(c.events)
}
func convertToolResults(items []ToolResultItem) []document.ToolResultItem {
if items == nil {
return nil
}
result := make([]document.ToolResultItem, len(items))
for i, item := range items {
result[i] = document.ToolResultItem{Name: item.Name, Output: item.Output}
}
return result
}