feat: 接入 WSPrefetchAdapter 到 recall 管线

- 新增 WSPrefetchAdapter 实现 storage.PrefetchPusher 接口
- 当 agentID 为空时使用 WSBus.Broadcast(支持多 Agent)
- 在 NewAPI 中通过 SetPrefetchPusher 注入
- WebSocket prefetch.push 事件在 CO_OCCURS 权重 > 0.6 时触发
This commit is contained in:
xiaowei 2026-05-30 21:58:26 +08:00
parent 928468042f
commit 93671dc1d6
2 changed files with 16 additions and 1 deletions

View File

@ -24,11 +24,13 @@ type API struct {
}
func NewAPI(ldb storage.LanceDB, emb *storage.Embedder, rerank *storage.Reranker) *API {
pipeline := storage.NewRecallPipeline(emb, ldb, rerank)
pipeline.SetPrefetchPusher(&WSPrefetchAdapter{})
return &API{
LanceDB: ldb,
Embedder: emb,
Reranker: rerank,
Pipeline: storage.NewRecallPipeline(emb, ldb, rerank),
Pipeline: pipeline,
}
}

View File

@ -1,6 +1,19 @@
// 织忆 MemoryWeave — WebSocket 事件推送函数
package routes
import "github.com/xiaoxue/memoryweave/internal/models"
// WSPrefetchAdapter 实现 storage.PrefetchPusher 接口
type WSPrefetchAdapter struct{}
func (a *WSPrefetchAdapter) PushPrefetch(agentID string, memories []models.RecallResult) {
if agentID != "" {
WSBus.Push(agentID, "prefetch.push", memories)
} else {
WSBus.Broadcast("prefetch.push", memories)
}
}
// PushPrefetch 预取推送recall 管道调用)
func PushPrefetch(agentID string, memories interface{}) {
WSBus.Push(agentID, "prefetch.push", memories)