mirror of
https://gitcode.com/JianFeeeee/HomeAgent.git
synced 2026-10-03 15:53:56 +00:00
feat(vector): pluggable multimodal vector space
核心暴露 MultimodalEmbedder 接口,两条路径共享同一套 L0/L2/L3 向量缓存、media.Store 坐标、QueryMemoryMediaScored 检索: - onnx:内嵌 ONNX 模型(CLIP 等),通过 build tag 编译 - http:外部向量 API 服务(Jina v5 / OpenAI / 自建) 跨模态融合权重改为 CrossModalFusionConfig 可配置结构体, 移除所有模型特定硬编码(CLIP/Jina),版本切换只需改配置。 模型切换自动迁移: - StaleVecDigestsAll 支持全模态(image+audio+video) - 启动时并发重算(ONNX 4 workers / API 8 workers) - 修复 SQL 运算符优先级导致 kind 过滤失效的 bug 实测对比(492 篇生产文档 + 3 张真实图片): - TF-IDF:MRR 0.457(精确匹配快,语义差) - fastText:MRR 0.530(语义中等,延迟 8ms) - Jina v5-omni:MRR 0.900(全面领先,延迟 40ms) - 中文文本→图片:Jina MRR 0.833 vs CLIP 0.611 See docs/embedding-comparison.md for full benchmark.
This commit is contained in:
@ -3,6 +3,8 @@ package core
|
||||
import (
|
||||
"log"
|
||||
"runtime/debug"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"time"
|
||||
|
||||
"gitcode.com/JianFeeeee/HomeAgent/internal/memory/media"
|
||||
@ -180,63 +182,106 @@ func (a *Agent) describePendingMedia() {
|
||||
}
|
||||
log.Printf("[media] 已描述 %s (%s, %d 字, 源=%s)", shortDigest(it.Digest), kind, len([]rune(desc)), srcName)
|
||||
|
||||
// 描述成功后,若 CLIP 嵌入器可用且是图片,计算视觉向量。
|
||||
// 这是"描述 + 向量"两步同步完成的路径;对于历史已有描述但无向量的媒体,
|
||||
// 由启动时的 reembedStaleMedia 补算。
|
||||
if a.clipEmb != nil && kind == "image" && it.Kind == media.KindImage {
|
||||
vec, err := a.clipEmb.EmbedImageDense(data, mime)
|
||||
if err != nil {
|
||||
log.Printf("[media] 视觉嵌入失败 %s: %v", shortDigest(it.Digest), err)
|
||||
} else if err := a.mediaStore.SetVec(it.Digest, vec, a.clipEmb.Fingerprint()); err != nil {
|
||||
log.Printf("[media] 写向量失败 %s: %v", shortDigest(it.Digest), err)
|
||||
} else {
|
||||
log.Printf("[media] 已嵌入 %s (dim=%d)", shortDigest(it.Digest), len(vec))
|
||||
}
|
||||
}
|
||||
// 描述成功后无需再次做视觉嵌入:图片在进入 L0 记忆块时已由
|
||||
// embedMediaOnIngest 计算并写入 CAS,L0→L2→L3 只转移引用并复用坐标。
|
||||
// 历史已有图片或模型切换由启动时 reembedStaleMedia 一次性补算。
|
||||
}
|
||||
}
|
||||
|
||||
// reembedStaleMedia 在启动时为历史已有描述但无 CLIP 向量的图片补算视觉向量。
|
||||
// 避免安装 CLIP 后,旧图片永远只有描述文本、没有视觉向量,直到下次 Describe 才能写入。
|
||||
// reembedStaleMedia 在启动时批量迁移历史媒体向量到当前向量空间。
|
||||
//
|
||||
// 触发场景(任一变化都会导致旧向量无法参与查询):
|
||||
// - 切换模型(模型 A→模型 B,fp 变了)
|
||||
// - 切换向量维度(ONNX→HTTP dim 512→1024)
|
||||
// - 首次部署嵌入服务(历史无向量的媒体补算)
|
||||
// - 嵌入服务离线后重新上线(失败条目 vec_model 仍为空)
|
||||
//
|
||||
// 并发策略:启动时用 worker pool 并行迁移,避免上千张图片串行耗时过长。
|
||||
// 并发数在 ONNX 内嵌路径下不超 CPU 核心数(避免 ONNX 并发限流),
|
||||
// 外部 API 路径下不超 8(避免打爆外部服务)。
|
||||
func (a *Agent) reembedStaleMedia() {
|
||||
if a.clipEmb == nil || a.mediaStore == nil {
|
||||
if a.multimodalSpace == nil || a.mediaStore == nil {
|
||||
return
|
||||
}
|
||||
fp := a.clipEmb.Fingerprint()
|
||||
digests, err := a.mediaStore.StaleVecDigests(fp)
|
||||
fp := a.multimodalSpace.Fingerprint()
|
||||
digests, err := a.mediaStore.StaleVecDigestsAll(fp)
|
||||
if err != nil {
|
||||
log.Printf("[media] 查询需重算向量的媒体失败: %v", err)
|
||||
return
|
||||
}
|
||||
if len(digests) == 0 {
|
||||
log.Printf("[media] 无历史媒体需要补算视觉向量")
|
||||
log.Printf("[media] 无需迁移向量(所有媒体已与当前空间对齐 fp=%s)", shortFP(fp))
|
||||
return
|
||||
}
|
||||
log.Printf("[media] 启动补算视觉向量: %d 条 (fp=%s...)", len(digests), fp[:min(12, len(fp))])
|
||||
done := 0
|
||||
for _, d := range digests {
|
||||
it, err := a.mediaStore.Stat(d)
|
||||
if err != nil {
|
||||
continue
|
||||
}
|
||||
data, err := a.mediaStore.Get(d)
|
||||
if err != nil {
|
||||
continue
|
||||
}
|
||||
mime := it.MIME
|
||||
if mime == "" {
|
||||
mime = "image/png"
|
||||
}
|
||||
vec, err := a.clipEmb.EmbedImageDense(data, mime)
|
||||
if err != nil {
|
||||
log.Printf("[media] 启动补算失败 %s: %v", shortDigest(d), err)
|
||||
continue
|
||||
}
|
||||
if err := a.mediaStore.SetVec(d, vec, fp); err != nil {
|
||||
log.Printf("[media] 启动写入向量失败 %s: %v", shortDigest(d), err)
|
||||
continue
|
||||
}
|
||||
done++
|
||||
|
||||
// 并发度:ONNX 内嵌不超过 4,外部 API 不超过 8(由配置或实际环境动态定)
|
||||
workers := 4
|
||||
if fp[:min(4, len(fp))] == "http:" {
|
||||
workers = 8
|
||||
}
|
||||
log.Printf("[media] 启动补算完成: %d/%d", done, len(digests))
|
||||
log.Printf("[media] 启动向量迁移: %d 条 → 新空间 fp=%s dim=%d workers=%d",
|
||||
len(digests), shortFP(fp), a.multimodalSpace.Dim(), workers)
|
||||
|
||||
jobs := make(chan string, workers*2)
|
||||
var done, failed int64
|
||||
var failedMu sync.Mutex
|
||||
var wg sync.WaitGroup
|
||||
|
||||
for i := 0; i < workers; i++ {
|
||||
wg.Add(1)
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
for d := range jobs {
|
||||
if err := a.reembedOne(d, fp); err != nil {
|
||||
failedMu.Lock()
|
||||
failed++
|
||||
failedMu.Unlock()
|
||||
continue
|
||||
}
|
||||
atomic.AddInt64(&done, 1)
|
||||
}
|
||||
}()
|
||||
}
|
||||
|
||||
for i, d := range digests {
|
||||
jobs <- d
|
||||
// 每迁移 20 条输出进度日志,让用户看到迁移在推进
|
||||
if (i+1)%20 == 0 {
|
||||
log.Printf("[media] 向量迁移进度: %d/%d (done=%d failed=%d)", i+1, len(digests), atomic.LoadInt64(&done), failed)
|
||||
}
|
||||
}
|
||||
close(jobs)
|
||||
wg.Wait()
|
||||
log.Printf("[media] 向量迁移完成: 成功=%d 失败=%d 总计=%d fp=%s",
|
||||
done, failed, len(digests), shortFP(fp))
|
||||
}
|
||||
|
||||
// reembedOne 为单条媒体重新计算向量并写入。stat 错误时跳过(可能已被 GC 清除)。
|
||||
// Get 错误或 Embed 错误时静默跳过该条目(不影响迁移其他条目)。
|
||||
func (a *Agent) reembedOne(digest, fp string) error {
|
||||
it, err := a.mediaStore.Stat(digest)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
data, err := a.mediaStore.Get(digest)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
mime := it.MIME
|
||||
if mime == "" {
|
||||
mime = "image/png"
|
||||
}
|
||||
vec, err := a.multimodalSpace.EmbedImageDense(data, mime)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
return a.mediaStore.SetVec(digest, vec, fp)
|
||||
}
|
||||
|
||||
// shortFP 截断 fingerprint 为可读日志格式。
|
||||
func shortFP(fp string) string {
|
||||
if len(fp) > 12 {
|
||||
return fp[:12]
|
||||
}
|
||||
return fp
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user