diff --git a/go/internal/api/routes/cascade.go b/go/internal/api/routes/cascade.go index a53eac1..ec4ee10 100644 --- a/go/internal/api/routes/cascade.go +++ b/go/internal/api/routes/cascade.go @@ -24,6 +24,11 @@ var CascadeR = &CascadeReviewer{ tracker: selfoptimize.NewCausalTracker(), } +// Tracker 暴露因果追踪器(供 server.go 注入 Redis 持久化) +func (cr *CascadeReviewer) Tracker() *selfoptimize.CausalTracker { + return cr.tracker +} + func (cr *CascadeReviewer) AddDependency(depID, dependsOnID string) { cr.mu.Lock() defer cr.mu.Unlock() diff --git a/go/internal/api/routes/core.go b/go/internal/api/routes/core.go index d36ccc8..3bedf39 100644 --- a/go/internal/api/routes/core.go +++ b/go/internal/api/routes/core.go @@ -133,6 +133,8 @@ func (a *API) Commit(w http.ResponseWriter, r *http.Request) { // 被动验证:对新记忆与已有记忆做 P1/P2/P3 匹配 go func() { + // 因果追踪:记录版本变更 + CascadeR.Tracker().RecordVersion(memID, req.Content, req.AgentID, "commit") // 搜索同 namespace 已有记忆 zeroVec := make([]float32, 1024) existing, _ := a.LanceDB.Search("memories", zeroVec, 50, req.Namespace) diff --git a/go/internal/api/routes/triggers.go b/go/internal/api/routes/triggers.go index aa0c2da..af87d75 100644 --- a/go/internal/api/routes/triggers.go +++ b/go/internal/api/routes/triggers.go @@ -240,6 +240,18 @@ func (sm *SkillManager) Trial(w http.ResponseWriter, r *http.Request) { sm.mu.Lock() defer sm.mu.Unlock() + skill := sm.recordTrialLocked(name, req.Success) + respond(w, 200, skill) +} + +// RecordTrial 程序化记录一次技能试验(无需 HTTP) +func (sm *SkillManager) RecordTrial(name string, success bool) { + sm.mu.Lock() + defer sm.mu.Unlock() + sm.recordTrialLocked(name, success) +} + +func (sm *SkillManager) recordTrialLocked(name string, success bool) *Skill { skill, exists := sm.skills[name] if !exists { skill = &Skill{ @@ -250,10 +262,10 @@ func (sm *SkillManager) Trial(w http.ResponseWriter, r *http.Request) { sm.skills[name] = skill } skill.Trials++ - if req.Success { + if success { skill.ETA = float64(skill.Trials-1) / float64(skill.Trials) } else { skill.ETA = float64(skill.Trials-1) / float64(skill.Trials) } - respond(w, 200, skill) + return skill } diff --git a/go/internal/api/server.go b/go/internal/api/server.go index c255b7d..b0b46a4 100644 --- a/go/internal/api/server.go +++ b/go/internal/api/server.go @@ -103,9 +103,12 @@ func NewServer() http.Handler { conflictDetector := governance.NewConflictDetector() conflictAPI := routes.NewConflictAPI(conflictDetector) - // 缺口 + // ─── 缺口 gapDetector := selfoptimize.NewGapDetector(emb, ldb) gapDetector.EnableGapRedisPersistence() + + // 因果追踪持久化 + routes.CascadeR.Tracker().EnableCausalRedisPersistence() gapAPI := routes.NewGapAPI(gapDetector) routes.InitGapRepair(gapDetector) // 共享同一个 GapDetector // G1a: 挂缺隙记录器(0 结果 → 自动记录 miss) @@ -184,6 +187,15 @@ func NewServer() http.Handler { } } selfoptimize.Validator.Validate(input.Content, validMems) + // 被动验证结果写回 quality_score + records := selfoptimize.Validator.GetRecords() + for _, r := range records { + if r.Confidence > 0 { + _ = ldb.Update("memories", r.MemoryID, map[string]any{ + "quality_score": r.Confidence, + }) + } + } } // ─── CO_OCCURS 追踪器(持久化到 Redis)─── @@ -198,6 +210,11 @@ func NewServer() http.Handler { // ─── V 值传播器(竞争性架构核心)───────────── vPropagator := selfoptimize.VProp + + // 注入 V 值查询函数到存储层(召回排序融合) + storage.VValueProvider = func(memoryID string) float64 { + return vPropagator.GetMemVValue(memoryID) + } // 在蒸馏完成回调中记录 V 值 originalOnComplete := distill.OnDistillComplete distill.OnDistillComplete = func(input distill.DistillInput, result distill.DistillResult) { @@ -215,6 +232,15 @@ func NewServer() http.Handler { } // 更新仪表盘指标 selfoptimize.Dash.RecordUseful() + // 记录蒸馏质量损失(avg_distill_loss = 1 - Overall) + loss := 1.0 - result.Overall + if loss < 0 { + loss = 0 + } + if loss > 1 { + loss = 1 + } + selfoptimize.Dash.RecordDistillLoss(loss) } // ─── 跨 Agent 缓存失效回调 ─────────────────── @@ -764,19 +790,31 @@ func NewServer() http.Handler { } } } - case "gap_scan": - // 检查已有缺口:过期 7 天的自动关闭 - gaps := gapDetector.List() - for _, g := range gaps { - if !g.Closed && time.Since(g.CreatedAt) > 7*24*time.Hour { - gapDetector.Close(g.Topic) + case "gap_scan": + // 检查已有缺口:过期 7 天的自动关闭 + gaps := gapDetector.List() + for _, g := range gaps { + if !g.Closed && time.Since(g.CreatedAt) > 7*24*time.Hour { + gapDetector.Close(g.Topic) + } + } + // 自动修复开放的缺口 + for _, g := range gaps { + if !g.Closed && routes.GapRepair != nil { + if routes.GapRepair.AutoRepair(g) { + selfoptimize.Dash.RecordGapClosed() } } } - if err != nil { - routes.Triggers.RecordFail(triggerID) - } - routes.WSBus.Broadcast("trigger.fired", map[string]string{ + } + if err != nil { + routes.Triggers.RecordFail(triggerID) + } + // Skill 结晶:每次蒸馏成功后记录 + if action == "distill" || action == "consolidation" { + routes.Skills.RecordTrial("auto_distill", true) + } + routes.WSBus.Broadcast("trigger.fired", map[string]string{ "trigger_id": triggerID, "action": action, }) @@ -785,10 +823,10 @@ func NewServer() http.Handler { } }() - // 自优化指标定时采集(每 6 小时) + // 自优化指标定时采集(每 30 分钟) go func() { for { - time.Sleep(6 * time.Hour) + time.Sleep(30 * time.Minute) m := selfoptimize.Dash.Metrics() date := time.Now().Format("2006-01-02") for k, v := range m { diff --git a/go/internal/selfoptimize/selfoptimize.go b/go/internal/selfoptimize/selfoptimize.go index 264d023..7286a33 100644 --- a/go/internal/selfoptimize/selfoptimize.go +++ b/go/internal/selfoptimize/selfoptimize.go @@ -408,6 +408,66 @@ type CausalTracker struct { mu sync.RWMutex entries map[string][]*TraceEntry // memory_id → version history deps map[string][]string // memory_id → depends_on[] + + // Redis 持久化 + redisClient *storage.RedisClient +} + +const causalRedisKey = "zhiyi:causal_entries" + +// EnableCausalRedisPersistence 启动因果追踪 Redis 持久化 +func (ct *CausalTracker) EnableCausalRedisPersistence() { + rc := storage.GetRedisClient() + if rc == nil { + return + } + ct.redisClient = rc + ct.loadCausalFromRedis() +} + +func (ct *CausalTracker) loadCausalFromRedis() { + if ct.redisClient == nil { + return + } + data, err := ct.redisClient.HGetAll(causalRedisKey) + if err != nil || len(data) == 0 { + return + } + // entries: memory_id → JSON数组 + // deps: memory_id → JSON数组 + ct.mu.Lock() + defer ct.mu.Unlock() + for key, jsonStr := range data { + if len(key) > 5 && key[:5] == "dep_:" { + // dep_:memory_id → [dependency_ids] + memID := key[5:] + var deps []string + if json.Unmarshal([]byte(jsonStr), &deps) == nil { + ct.deps[memID] = deps + } + } else { + // memory_id → [TraceEntry] + var entries []*TraceEntry + if json.Unmarshal([]byte(jsonStr), &entries) == nil { + ct.entries[key] = entries + } + } + } +} + +func (ct *CausalTracker) persistCausal() { + if ct.redisClient == nil { + return + } + for memID, entries := range ct.entries { + data, _ := json.Marshal(entries) + ct.redisClient.HSet(causalRedisKey, memID, string(data)) + } + for memID, deps := range ct.deps { + data, _ := json.Marshal(deps) + ct.redisClient.HSet(causalRedisKey, "dep_:"+memID, string(data)) + } + ct.redisClient.Expire(causalRedisKey, 90*24*time.Hour) } func NewCausalTracker() *CausalTracker { @@ -431,6 +491,7 @@ func (ct *CausalTracker) RecordVersion(memoryID, content, source, trigger string UpdatedAt: time.Now(), } ct.entries[memoryID] = append(ct.entries[memoryID], entry) + ct.persistCausal() } // AddDependency A depends_on B @@ -438,6 +499,7 @@ func (ct *CausalTracker) AddDependency(a, b string) { ct.mu.Lock() defer ct.mu.Unlock() ct.deps[a] = append(ct.deps[a], b) + ct.persistCausal() } // GetAffected 当 memoryID 被修正时,返回所有依赖它的记忆 @@ -560,6 +622,15 @@ func (d *Dashboard) RecordConflictResolved(auto bool) { d.persist() } +// RecordDistillLoss 记录蒸馏质量损失 +func (d *Dashboard) RecordDistillLoss(loss float64) { + d.mu.Lock() + defer d.mu.Unlock() + d.DistillLossSum += loss + d.DistillLossCount++ + d.persist() +} + func (pg *PrefetchGraph) GetPrefetch(query string) []string { pg.mu.RLock() defer pg.mu.RUnlock() diff --git a/go/internal/storage/lancedb.go b/go/internal/storage/lancedb.go index f3d58c0..3d8440a 100644 --- a/go/internal/storage/lancedb.go +++ b/go/internal/storage/lancedb.go @@ -5,6 +5,10 @@ package storage import "github.com/xiaoxue/memoryweave/internal/models" +// VValueProvider 外部注入的 V 值查询函数(由 server.go 设置 → selfoptimize.VProp) +// 用于召回排序时融合用户决策信号 +var VValueProvider func(memoryID string) float64 + // LanceDB 存储后端统一接口 // 实现者: SQLiteClient (生产), MemLanceClient (开发/降级) // 未来: Rust lancedb crate 通过 Unix Socket 提供 LanceDB 原生后端 diff --git a/go/internal/storage/lancedb_ipc.go b/go/internal/storage/lancedb_ipc.go index 57a0cf5..4a8db6e 100644 --- a/go/internal/storage/lancedb_ipc.go +++ b/go/internal/storage/lancedb_ipc.go @@ -204,6 +204,17 @@ func (rc *RustLanceDBClient) Search(table string, vector []float32, topK int, na // QualityScore from LLM (0-1) + computed importance, capped at 1.0 finalQuality = math.Min(0.4*computedImportance+0.6*r.QualityScore, 1.0) } + // V 值融合:用户决策信号(correct/useful/not-useful) + if VValueProvider != nil { + vVal := VValueProvider(r.ID) + if vVal > 0 { + // 30% V值 + 30% importance + 40% LLM quality + finalQuality = 0.3*vVal + 0.3*computedImportance + 0.4*r.QualityScore + if finalQuality > 1.0 { + finalQuality = 1.0 + } + } + } out = append(out, models.MemoryRecord{ ID: r.ID, Content: r.Content, Category: r.Category, Namespace: r.Namespace,