From b9107fa925917c9dcd46cac077df1e05bec5e068 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=B0=8F=E5=94=AF?= Date: Sat, 12 Sep 2026 18:07:30 +0800 Subject: [PATCH] =?UTF-8?q?docs+test(recall):=20=E4=BF=AE=E6=AD=A3?= =?UTF-8?q?=E6=89=B9=E9=87=8F=E6=9B=B4=E6=96=B0=E6=B3=A8=E9=87=8A=EF=BC=88?= =?UTF-8?q?delta=20=E5=88=86=E7=BB=84=EF=BC=8C=E9=9D=9E=20CASE=20WHEN?= =?UTF-8?q?=EF=BC=89+=20=E8=A1=A5=20IPC=20=E5=8F=AA=E8=AF=BB=E6=A0=A1?= =?UTF-8?q?=E9=AA=8C=E7=94=A8=E4=BE=8B=20(t=5F8496e8b6)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - recall_write_buffer.go / lancedb_ipc.go: 注释与实现对齐 —— lance 不支持 CASE WHEN, Rust 侧按 delta 分组、每组一次 update 提交;整批版本数 = 不同 delta 的个数(通常 1 个) - batch_ipc_integration_test.go: +TestBatchIPCScanReadOnly(只读扫 socket 校验 recall_count 真落盘) 实测(临时 sidecar,/tmp/mw-verify-180517): VERSION_ACCOUNTING: 批量3条 → 1 个版本 | 逐条3次 → 3 个版本 BATCH_IPC_OK id=ep_1786982837511335732 recall_count 0 → 3 SCAN_READONLY: 总 1814 条, recall_count>0 的 742 条, 最大 828 --- .../storage/batch_ipc_integration_test.go | 35 +++++++++++++++++++ go/internal/storage/lancedb_ipc.go | 5 +-- go/internal/storage/recall_write_buffer.go | 6 ++-- 3 files changed, 42 insertions(+), 4 deletions(-) diff --git a/go/internal/storage/batch_ipc_integration_test.go b/go/internal/storage/batch_ipc_integration_test.go index bcf1dbf..627e808 100644 --- a/go/internal/storage/batch_ipc_integration_test.go +++ b/go/internal/storage/batch_ipc_integration_test.go @@ -124,6 +124,41 @@ func TestBatchIPCVersionAccounting(t *testing.T) { } } +// 只读校验:扫生产 socket,确认批量落盘真的写进了 recall_count(不产生任何写操作)。 +// +// ZHIYI_TEST_IPC_SOCK=/tmp/zhiyi-ipc.sock go test ./internal/storage/ -run TestBatchIPCScanReadOnly -v +func TestBatchIPCScanReadOnly(t *testing.T) { + sock := os.Getenv("ZHIYI_TEST_IPC_SOCK") + if sock == "" { + t.Skip("ZHIYI_TEST_IPC_SOCK 未设置") + } + resp := ipcCall(t, sock, map[string]any{"type": "lancedb_scan", "limit": 2000}) + var recs []struct { + ID string `json:"id"` + RecallCount int64 `json:"recall_count"` + LastRecalledAt string `json:"last_recalled_at"` + } + if err := json.Unmarshal([]byte(resp["report_json"].(string)), &recs); err != nil { + t.Fatalf("unmarshal: %v", err) + } + withCount, maxCount := 0, int64(0) + var sampleID, sampleTS string + for _, r := range recs { + if r.RecallCount > 0 { + withCount++ + if r.RecallCount > maxCount { + maxCount = r.RecallCount + sampleID, sampleTS = r.ID, r.LastRecalledAt + } + } + } + fmt.Printf("SCAN_READONLY: 总 %d 条, 其中 recall_count>0 的 %d 条, 最大 recall_count=%d (id=%s last_recalled_at=%s)\n", + len(recs), withCount, maxCount, sampleID, sampleTS) + if withCount == 0 { + t.Fatalf("scan 里没有任何 recall_count>0 —— 批量落盘可能没生效") + } +} + func TestBatchIPCUpdateRecallBatch(t *testing.T) { sock := os.Getenv("ZHIYI_TEST_IPC_SOCK") if sock == "" { diff --git a/go/internal/storage/lancedb_ipc.go b/go/internal/storage/lancedb_ipc.go index d8a5fa7..bce316e 100644 --- a/go/internal/storage/lancedb_ipc.go +++ b/go/internal/storage/lancedb_ipc.go @@ -453,8 +453,9 @@ func (rc *RustLanceDBClient) Update(table, id string, fields map[string]any) err } // UpdateRecallBatch 单事务批量更新召回元数据(2026-09-12 优化B)。 -// Rust 侧用一次 update 提交(CASE WHEN 表达式覆盖所有命中行)→ 整批只产生 1 个 LanceDB 版本。 -// recall_count 由 Rust 侧读「库现值 + delta」计算,避免 Go 缓存漂移。 +// Rust 侧把整批按 delta 分组,每组一次 update(`recall_count + delta` + id IN (...))提交 +// (lance 不支持 CASE WHEN)→ 整批版本数 = 不同 delta 的个数,通常 1 个(旧实现:每行 1 个)。 +// 注:Rust 侧固定操作 memories 表,table 参数暂未使用。 func (rc *RustLanceDBClient) UpdateRecallBatch(table string, items []RecallWriteItem, lastRecalledAt string) (int64, error) { if len(items) == 0 { return 0, nil diff --git a/go/internal/storage/recall_write_buffer.go b/go/internal/storage/recall_write_buffer.go index da57303..08d71d0 100644 --- a/go/internal/storage/recall_write_buffer.go +++ b/go/internal/storage/recall_write_buffer.go @@ -21,8 +21,10 @@ import ( // // Record() 读路径只把增量写进内存缓冲;**同一记忆在一个窗口内多次命中合并成 1 条 delta** // Flush() 后台按窗口(默认 300s)+ 阈值(默认 256 条不同记忆)批量提交一次; -// Rust 侧 lancedb_update_batch 用**一次** update 提交(CASE WHEN 表达式) -// → 每批只产生 1 个 LanceDB 版本(旧实现:每行 1 个版本) +// Rust 侧 lancedb_update_batch 把整批**按 delta 分组**,每组一次 +// update(`recall_count + delta` + id IN (...)) 提交(lance 不支持 CASE WHEN) +// → 每批版本数 = 不同 delta 的个数(delta 绝大多数为 1,故通常 1 个版本; +// 旧实现:每行 1 个版本) // // 语义取舍(明确记录,便于日后审计): // - last_recalled_at / freshness 最多延迟一个窗口(分钟级)。遗忘/衰减判定以「天」为单位,无影响。