whlog 0.1.3

HybridLog allocator with three-region sliding window and scan iterator
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
use std::{
  slice,
  sync::{Arc, mpsc},
  thread,
  time::Duration,
};

use aok::{OK, Void};
use compio::runtime::Runtime;
use log::info;
use tempfile::tempdir;
use wdev::SegmentedDevice;
use wepoch::LightEpoch;
use whlog::{
  DEFAULT_INITIAL_ADDRESS, Error, HybridLog, HybridLogConfig, RecordOutput, SECTOR_ALIGNMENT,
};
use wrecord::{HEADER_SIZE, encode_to_slice};

/// 测试 1: 单页追加与内存直读
#[test]
fn test_append_and_memory_read() -> Void {
  let rt = Runtime::new()?;
  rt.block_on(async {
    let dir = tempdir()?;
    let db_path = dir.path().join("hlog_test1.db");
    let device = Arc::new(SegmentedDevice::single_file(&db_path)?);
    let epoch = Arc::new(LightEpoch::new(16));

    let config = HybridLogConfig::new(64 * 1024, 16, 0.5)?;
    let hlog = HybridLog::new(config, device, epoch)?;

    let k1 = b"user:1001";
    let v1 = b"alice_data";
    let addr1 = hlog.append(k1, v1, 0, false)?;
    assert_eq!(addr1, DEFAULT_INITIAL_ADDRESS);

    let k2 = b"user:1002";
    let v2 = b"bob_payload_string";
    let addr2 = hlog.append(k2, v2, addr1, false)?;
    assert!(addr2 > addr1);

    // 内存中直读
    assert!(hlog.is_in_memory(addr1));
    assert!(hlog.is_in_memory(addr2));
    assert!(hlog.is_mutable(addr1));
    assert!(hlog.is_mutable(addr2));

    let out1 = hlog.read_record(addr1).await?;
    assert!(matches!(out1, RecordOutput::Memory(_)));
    assert_eq!(out1.key()?, k1);
    assert_eq!(out1.value()?, v1);
    assert_eq!(out1.prev_address()?, 0);
    assert!(!out1.is_tombstone()?);

    let out2 = hlog.read_record(addr2).await?;
    assert!(matches!(out2, RecordOutput::Memory(_)));
    assert_eq!(out2.key()?, k2);
    assert_eq!(out2.value()?, v2);
    assert_eq!(out2.prev_address()?, addr1);

    info!("单页追加与内存直读测试通过");
    aok::Result::<()>::Ok(())
  })?;

  OK
}

/// 测试 3: 跨页换页(Page Turn)与 Padding 验证
#[test]
fn test_page_turn_and_padding() -> Void {
  let rt = Runtime::new()?;
  rt.block_on(async {
    let dir = tempdir()?;
    let db_path = dir.path().join("hlog_test3.db");
    let device = Arc::new(SegmentedDevice::single_file(&db_path)?);
    let epoch = Arc::new(LightEpoch::new(16));

    // 使用最小合法单页 4096 字节
    let page_size = SECTOR_ALIGNMENT; // 4096
    let config = HybridLogConfig::new(page_size, 16, 0.5)?;
    let hlog = HybridLog::new(config, device, epoch)?;

    // 初始 tail 从 64 开始,先写入若干记录填满大部分空间
    // 每条记录 16(header) + 8(key) + 400(val) = 424 字节
    let mut addrs = Vec::new();
    let val_400 = vec![b'A'; 400];
    for i in 0..9 {
      let key = format!("k:{i:06}");
      let addr = hlog.append(key.as_bytes(), &val_400, 0, false)?;
      addrs.push(addr);
    }

    // 此时第 0 页已消耗 64 + 9 * 424 = 3880 字节,剩余 4096 - 3880 = 216 字节
    // 写入一条大小为 16 + 8 + 300 = 324 字节的记录,必然无法容纳,触发换页!
    let overflow_key = b"k:overflow";
    let overflow_val = vec![b'B'; 300];
    let overflow_addr = hlog.append(overflow_key, &overflow_val, 0, false)?;

    // 验证新记录写入了第 1 页的起始位置(page_id=1, addr=4096)
    assert_eq!(
      overflow_addr, page_size as u64,
      "换页后新记录必须位于下一页开头"
    );

    // 验证原页末尾 3880 偏移处写入了 Pad 记录
    let pad_addr = 3880u64;
    let pad_res = hlog.read_record(pad_addr).await;
    assert!(
      matches!(pad_res, Err(Error::PadRecord(a)) if a == pad_addr),
      "读取填充位置应返回 PadRecord 错误"
    );

    // 验证新记录在第 1 页可正常读取
    let out = hlog.read_record(overflow_addr).await?;
    assert_eq!(out.key()?, overflow_key);
    assert_eq!(out.value()?, &overflow_val[..]);

    info!("跨页换页与 Padding 验证通过");
    aok::Result::<()>::Ok(())
  })?;

  OK
}

/// 测试 8: 多页连续 Scan 扫描与拉模式迭代器(自动跳过 PadRecord)
#[test]
fn test_scan_multipage_and_pull_iterator() -> Void {
  let rt = Runtime::new()?;
  rt.block_on(async {
    let dir = tempdir()?;
    let db_path = dir.path().join("hlog_scan_test.db");
    let device = Arc::new(SegmentedDevice::single_file(&db_path)?);
    let epoch = Arc::new(LightEpoch::new(16));

    let page_size = SECTOR_ALIGNMENT; // 4096
    let config = HybridLogConfig::new(page_size, 16, 0.5)?;
    let hlog = HybridLog::new(config, device, epoch)?;

    let mut expected_records = Vec::new();
    // 写入跨越至少 3 页的数据
    for i in 0..30 {
      let key = format!("scan_k:{i:04}");
      let val = vec![(i & 0xFF) as u8; 200];
      let addr = hlog.append(key.as_bytes(), &val, 0, false)?;
      expected_records.push((addr, key.into_bytes(), val));
    }

    // 1. 测试 push-based scan 全量扫描
    let mut scanned_records = Vec::new();
    hlog
      .scan(0, hlog.tail_address(), |addr, rec| {
        scanned_records.push((addr, rec.key().to_vec(), rec.value().to_vec()));
        Ok(true)
      })
      .await?;

    assert_eq!(scanned_records.len(), expected_records.len());
    for (actual, expected) in scanned_records.iter().zip(expected_records.iter()) {
      assert_eq!(actual.0, expected.0, "地址不一致");
      assert_eq!(actual.1, expected.1, "Key 不一致");
      assert_eq!(actual.2, expected.2, "Value 不一致");
    }

    // 2. 测试 pull-based ScanIterator
    let mut iter = hlog.scan_iter(0, hlog.tail_address());
    let mut pulled_records = Vec::new();
    while let Some((addr, out)) = iter.next().await? {
      pulled_records.push((addr, out.key()?.to_vec(), out.value()?.to_vec()));
    }

    assert_eq!(pulled_records.len(), expected_records.len());
    for (actual, expected) in pulled_records.iter().zip(expected_records.iter()) {
      assert_eq!(actual.0, expected.0);
      assert_eq!(actual.1, expected.1);
      assert_eq!(actual.2, expected.2);
    }

    info!("多页连续 Scan 扫描与拉模式迭代器测试通过");
    aok::Result::<()>::Ok(())
  })?;

  OK
}

/// 测试 9: 混合冷热数据穿透连续 Scan 扫描(磁盘区 + 内存只读区 + 内存可变区)
#[test]
fn test_scan_hybrid_disk_and_memory() -> Void {
  let rt = Runtime::new()?;
  rt.block_on(async {
    let dir = tempdir()?;
    let db_path = dir.path().join("hlog_scan_hybrid.db");
    let device = Arc::new(SegmentedDevice::single_file(&db_path)?);
    let epoch = Arc::new(LightEpoch::new(16));

    let page_size = SECTOR_ALIGNMENT; // 4096
    let config = HybridLogConfig::new(page_size, 16, 0.5)?;
    let hlog = HybridLog::new(config, device, epoch)?;

    // 写入第 0 页数据
    let mut all_records = Vec::new();
    for i in 0..5 {
      let key = format!("cold_k:{}", i);
      let val = vec![0x11; 100];
      let addr = hlog.append(key.as_bytes(), &val, 0, false)?;
      all_records.push((addr, key.into_bytes(), val));
    }

    // 填平第 0 页促使换页
    let fill_key = b"fill";
    let rem = page_size - hlog.config.page_offset(hlog.tail_address());
    let fill_val = vec![0x00; rem - HEADER_SIZE - fill_key.len()];
    let _ = hlog.append(fill_key, &fill_val, 0, false)?;

    // 写入第 1 页数据
    for i in 0..5 {
      let key = format!("hot_k:{}", i);
      let val = vec![0x22; 100];
      let addr = hlog.append(key.as_bytes(), &val, 0, false)?;
      all_records.push((addr, key.into_bytes(), val));
    }

    // 刷盘第 0 页并推进 HeadAddress 将其驱逐为磁盘冷数据
    hlog.flush_page(0).await?;
    hlog.shift_read_only_address(page_size as u64);
    hlog.shift_head_address(page_size as u64);

    assert!(hlog.is_on_disk(all_records[0].0));
    assert!(hlog.is_in_memory(all_records[5].0));

    // 执行跨三区扫描,应顺序读取冷数据与热数据
    let mut scanned = Vec::new();
    hlog
      .scan(0, hlog.tail_address(), |addr, rec| {
        if rec.key() != fill_key {
          scanned.push((addr, rec.key().to_vec(), rec.value().to_vec()));
        }
        Ok(true)
      })
      .await?;

    assert_eq!(scanned.len(), all_records.len());
    for (actual, expected) in scanned.iter().zip(all_records.iter()) {
      assert_eq!(actual.0, expected.0);
      assert_eq!(actual.1, expected.1);
      assert_eq!(actual.2, expected.2);
    }

    info!("混合冷热数据穿透连续 Scan 扫描测试通过");
    aok::Result::<()>::Ok(())
  })?;

  OK
}

/// 测试 10: Push-based Scan 提前终止(Early Termination)
#[test]
fn test_scan_early_termination() -> Void {
  let rt = Runtime::new()?;
  rt.block_on(async {
    let dir = tempdir()?;
    let db_path = dir.path().join("hlog_scan_early_stop.db");
    let device = Arc::new(SegmentedDevice::single_file(&db_path)?);
    let epoch = Arc::new(LightEpoch::new(16));

    let config = HybridLogConfig::new(SECTOR_ALIGNMENT, 16, 0.5)?;
    let hlog = HybridLog::new(config, device, epoch)?;

    for i in 0..20 {
      let key = format!("k:{i:02}");
      let val = b"data";
      let _ = hlog.append(key.as_bytes(), val, 0, false)?;
    }

    // 扫描并在第 5 条记录时提前终止
    let mut count = 0;
    hlog
      .scan(0, hlog.tail_address(), |_addr, _rec| {
        count += 1;
        if count == 5 {
          Ok(false) // 提前终止
        } else {
          Ok(true)
        }
      })
      .await?;

    assert_eq!(count, 5, "Scan 应在第 5 条记录处成功提前终止");

    info!("Push-based Scan 提前终止测试通过");
    aok::Result::<()>::Ok(())
  })?;

  OK
}

/// 测试 22: 页尾子头残片(0xFF 填充)的写入、读取拦截与扫描跳过
#[test]
fn test_subheader_fragment_pad() -> Void {
  let rt = Runtime::new()?;
  rt.block_on(async {
    let dir = tempdir()?;
    let db_path = dir.path().join("hlog_fragment.db");
    let device = Arc::new(SegmentedDevice::single_file(&db_path)?);
    let epoch = Arc::new(LightEpoch::new(16));

    let page_size = SECTOR_ALIGNMENT;
    let config = HybridLogConfig::new(page_size, 16, 0.5)?;
    let hlog = HybridLog::new(config, device, epoch)?;

    // R1: 16 + 10 + 4000 = 4026 → 页尾仅剩 6 字节(不足以容纳记录头)
    let v = vec![7u8; 4000];
    let addr1 = hlog.append(b"0123456789", &v, 0, false)?;
    assert_eq!(addr1, DEFAULT_INITIAL_ADDRESS);

    let addr2 = hlog.append(b"frag", b"tail", 0, false)?;
    assert_eq!(addr2, page_size as u64, "残片后新记录必须落于下一页开头");

    // 残片地址读取 → PadRecord;扫描跳过残片完整读出两条记录
    let fragment_addr = addr1 + 4026;
    assert_eq!(fragment_addr, page_size as u64 - 6);
    assert!(matches!(
      hlog.read_record(fragment_addr).await,
      Err(Error::PadRecord(_))
    ));

    let mut scanned = Vec::new();
    hlog
      .scan(0, hlog.tail_address(), |addr, rec| {
        scanned.push((addr, rec.key().to_vec()));
        Ok(true)
      })
      .await?;
    assert_eq!(
      scanned,
      vec![(addr1, b"0123456789".to_vec()), (addr2, b"frag".to_vec())]
    );

    info!("页尾子头残片处理测试通过");
    aok::Result::<()>::Ok(())
  })?;

  OK
}

/// 测试 23b: 在途预留零头不漏扫同页后续记录(对标 C# TsavoriteLogScanIterator.cs:779-790
/// scanUncommitted 模式的 SafeTailAddress + Thread.SpinWait(100) 复查语义)
///
/// append 协议「先 CAS 预占 tail 后编码」:本测试白盒清零中间记录字节,构造与
/// 在途预留槽位物理形态相同的零洞;扫描闭包读到前驱记录后经通道唤醒生产者线程,
/// 以 append encode_at 同款无锁裸指针路径补写该记录。扫描器须在零头上自旋等待
/// 编码完成后原址重试读出中间记录,且同页后续记录不漏。
#[test]
fn test_scan_inflight_zero_header_respin() -> Void {
  let rt = Runtime::new()?;
  rt.block_on(async {
    let dir = tempdir()?;
    let db_path = dir.path().join("hlog_inflight.db");
    let device = Arc::new(SegmentedDevice::single_file(&db_path)?);
    let epoch = Arc::new(LightEpoch::new(16));

    let config = HybridLogConfig::new(SECTOR_ALIGNMENT, 16, 1.0)?;
    let hlog = Arc::new(HybridLog::new(config, device, epoch)?);

    let addr0 = hlog.append(b"k0", b"v0", 0, false)?;
    let addr1 = hlog.append(b"k1", b"v1", 0, false)?;
    let addr2 = hlog.append(b"k2", b"v2", 0, false)?;

    // 白盒模拟在途预留窗口:中间记录(16B 头 + 2B 键 + 2B 值)整条清零,
    // 等价于 append 已发布 tail 但 encode_at 尚未落笔的物理形态
    let page1 = hlog.config.page_id(addr1);
    let off1 = hlog.config.page_offset(addr1);
    let rec1_len = HEADER_SIZE + 4;
    {
      let mut guard = hlog.buffer.write_page(page1);
      guard[off1..off1 + rec1_len].fill(0);
    }

    // 握手通道:扫描闭包读到 k0(零洞前驱)后唤醒生产者,生产者以与 encode_at
    // 完全一致的无锁裸指针路径补写 k1(零洞恰被扫描器触达时走自旋重试路径)
    let (encode_tx, encode_rx) = mpsc::channel::<()>();
    let producer = {
      let hlog = Arc::clone(&hlog);
      thread::spawn(move || {
        if encode_rx.recv().is_err() {
          return;
        }
        // 扫描器触达 k0 即将进入零洞:补写窗口足够覆盖自旋预算
        thread::sleep(Duration::from_micros(50));
        let slot = hlog.buffer.page_idx(page1);
        // SAFETY: 测试单扫描线程与生产者线程互斥于补写窗口,写入区间
        // [off1, off1 + rec1_len) 为已预留槽位,无其他并发访问者
        let ptr = unsafe { hlog.buffer.raw_page_ptr_mut(slot) };
        let dst = unsafe { slice::from_raw_parts_mut(ptr.add(off1), rec1_len) };
        let _ = encode_to_slice(dst, 0, b"k1", b"v1", false);
      })
    };

    let mut scanned = Vec::new();
    hlog
      .scan(0, hlog.tail_address(), |addr, rec| {
        scanned.push((addr, rec.key().to_vec()));
        if addr == addr0 {
          let _ = encode_tx.send(());
        }
        Ok(true)
      })
      .await?;
    producer.join().unwrap();

    assert_eq!(
      scanned,
      vec![
        (addr0, b"k0".to_vec()),
        (addr1, b"k1".to_vec()),
        (addr2, b"k2".to_vec())
      ],
      "在途零头自旋重试后必须原址读出 k1,且同页后续记录不得漏扫"
    );

    info!("在途预留零头自旋重试测试通过");
    aok::Result::<()>::Ok(())
  })?;

  OK
}