docs+test(recall): 修正批量更新注释(delta 分组,非 CASE WHEN)+ 补 IPC 只读校验用例 (t_8496e8b6)
- 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
This commit is contained in:
parent
2e04adf054
commit
b9107fa925
|
|
@ -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 == "" {
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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 最多延迟一个窗口(分钟级)。遗忘/衰减判定以「天」为单位,无影响。
|
||||
|
|
|
|||
Loading…
Reference in New Issue