use tracing::info;
use super::format::{SYNC_HWM_CKPT_FORMAT_VERSION, SyncHwmCheckpointFile};
use super::paths::{sync_hwm_ckpt_dir, sync_hwm_ckpt_state_path};
use crate::data::executor::checkpoint_decode_error::CheckpointDecodeError;
use crate::data::executor::core_loop::CoreLoop;
use crate::types::Lsn;
impl CoreLoop {
pub fn load_sync_hwm_checkpoint(&mut self) -> crate::Result<()> {
let ckpt_dir = sync_hwm_ckpt_dir(&self.data_dir, self.core_id);
let path = sync_hwm_ckpt_state_path(&ckpt_dir);
if !path.exists() {
return Ok(());
}
let bytes = nodedb_wal::segment::read_checkpoint_framed(&path).map_err(|source| {
CheckpointDecodeError::ReadFile {
path: path.clone(),
source,
}
})?;
let file = zerompk::from_msgpack::<SyncHwmCheckpointFile>(&bytes).map_err(|source| {
CheckpointDecodeError::MsgpackDecode {
path: path.clone(),
source,
}
})?;
if file.format_version != SYNC_HWM_CKPT_FORMAT_VERSION {
return Err(CheckpointDecodeError::FormatVersion {
path: path.clone(),
found: file.format_version,
expected: SYNC_HWM_CKPT_FORMAT_VERSION,
}
.into());
}
let streams = file.hwm.len();
let producers = file.epoch_floor.len();
for (producer_id, stream_id, seq) in file.hwm {
self.sync_hwm.insert((producer_id, stream_id), seq);
}
for (producer_id, epoch) in file.epoch_floor {
self.producer_epoch_floor.insert(producer_id, epoch);
}
self.floors.sync_hwm_durable_lsn = Lsn::new(file.durable_through_lsn);
info!(
core = self.core_id,
streams,
producers,
durable_through_lsn = file.durable_through_lsn,
"sync HWM checkpoint restored"
);
Ok(())
}
}
#[cfg(test)]
mod tests {
use std::sync::Arc;
use nodedb_bridge::buffer::RingBuffer;
use nodedb_types::OrdinalClock;
use nodedb_types::sync::wire::SyncProvenance;
use tempfile::TempDir;
use super::*;
use crate::bridge::dispatch::{BridgeRequest, BridgeResponse};
use crate::data::executor::sync_gate::SyncAdmit;
fn make_prov(producer_id: u64, epoch: u64, stream_id: u64, seq: u64) -> SyncProvenance {
SyncProvenance {
producer_id,
epoch,
stream_id,
seq,
}
}
fn open_core_at(dir: &std::path::Path) -> CoreLoop {
let hlc = Arc::new(OrdinalClock::new());
let (req_tx, req_rx) = RingBuffer::channel::<BridgeRequest>(64);
let (resp_tx, _resp_rx) = RingBuffer::channel::<BridgeResponse>(64);
drop(req_tx); CoreLoop::open(0, req_rx, resp_tx, dir, hlc).expect("CoreLoop::open")
}
#[test]
fn restored_gate_still_rejects_a_duplicate() {
let dir = TempDir::new().expect("tempdir");
let mut before = open_core_at(dir.path());
let applied = make_prov(1, 3, 5, 1);
assert_eq!(before.sync_admit(&applied), SyncAdmit::Apply);
before.sync_commit(&applied);
let other_stream = make_prov(1, 3, 6, 1);
assert_eq!(before.sync_admit(&other_stream), SyncAdmit::Apply);
before.sync_commit(&other_stream);
before.advance_watermark(Lsn::new(900));
let reported = before
.checkpoint_sync_hwm()
.expect("flush to a writable dir must succeed");
assert_eq!(
reported,
Lsn::new(900),
"the flush must report exactly the LSN it made durable — the manager \
deletes WAL segments below whatever this returns"
);
drop(before);
let mut unrestored = open_core_at(dir.path());
assert!(
unrestored.sync_hwm.is_empty(),
"a fresh core's gate must start empty, or this test proves nothing"
);
assert_eq!(
unrestored.sync_admit(&make_prov(1, 3, 5, 1)),
SyncAdmit::Apply,
"before the restore the already-applied frame IS re-admitted — this is \
the duplicate the checkpoint exists to prevent"
);
drop(unrestored);
let mut after = open_core_at(dir.path());
after
.load_sync_hwm_checkpoint()
.expect("valid checkpoint must load");
assert_eq!(
after.sync_admit(&make_prov(1, 3, 5, 1)),
SyncAdmit::Duplicate,
"the frame at the restored high-watermark must not be applied again"
);
assert_eq!(
after.sync_admit(&make_prov(1, 3, 6, 1)),
SyncAdmit::Duplicate,
"every stream's high-watermark must restore, not just the first"
);
assert_eq!(
after.sync_admit(&make_prov(1, 3, 5, 2)),
SyncAdmit::Apply,
"the next frame in sequence must still be admitted"
);
assert_eq!(
after.sync_admit(&make_prov(1, 2, 5, 2)),
SyncAdmit::Fenced,
"the restored epoch floor must still fence a stale producer generation"
);
assert_eq!(
after.floors.sync_hwm_durable_lsn,
Lsn::new(900),
"the restored durable LSN is what a failed flush clamps to; losing it \
would pin truncation at zero for the rest of the process"
);
}
#[test]
fn replay_merges_over_the_restored_gate_max_wins() {
let dir = TempDir::new().expect("tempdir");
let mut before = open_core_at(dir.path());
before.sync_commit(&make_prov(1, 3, 5, 42));
before.checkpoint_sync_hwm().expect("flush");
drop(before);
let mut after = open_core_at(dir.path());
after
.load_sync_hwm_checkpoint()
.expect("valid checkpoint must load");
let mut maps = crate::wal::replay::SyncHwmReplayMaps::default();
maps.sync_hwm.insert((1, 5), 50);
maps.producer_epoch_floor.insert(1, 3);
after.install_sync_hwm_maps(maps);
assert_eq!(after.sync_hwm_value(1, 5), 50, "a later record must win");
let mut stale = crate::wal::replay::SyncHwmReplayMaps::default();
stale.sync_hwm.insert((1, 5), 10);
after.install_sync_hwm_maps(stale);
assert_eq!(
after.sync_hwm_value(1, 5),
50,
"an earlier record must never lower the high-watermark"
);
}
#[test]
fn absent_checkpoint_restores_nothing() {
let dir = TempDir::new().expect("tempdir");
let mut core = open_core_at(dir.path());
core.load_sync_hwm_checkpoint()
.expect("an absent checkpoint is a legitimate no-op, not an error");
assert!(core.sync_hwm.is_empty());
assert!(core.producer_epoch_floor.is_empty());
assert_eq!(core.floors.sync_hwm_durable_lsn, Lsn::ZERO);
}
#[test]
fn unknown_version_is_fail_stop() {
let dir = TempDir::new().expect("tempdir");
let ckpt_dir = sync_hwm_ckpt_dir(dir.path(), 0);
std::fs::create_dir_all(&ckpt_dir).expect("mkdir");
let file = SyncHwmCheckpointFile {
format_version: SYNC_HWM_CKPT_FORMAT_VERSION + 1,
durable_through_lsn: 5,
hwm: vec![(1, 5, 42)],
epoch_floor: vec![(1, 3)],
};
let bytes = zerompk::to_msgpack_vec(&file).expect("encode");
let path = sync_hwm_ckpt_state_path(&ckpt_dir);
let tmp = ckpt_dir.join("STATE.tmp");
nodedb_wal::segment::write_checkpoint_framed(&tmp, &path, &bytes).expect("write");
let mut core = open_core_at(dir.path());
assert!(
core.load_sync_hwm_checkpoint().is_err(),
"a file this build cannot read must abort the load, not restore nothing"
);
}
#[test]
fn corrupt_msgpack_body_is_fail_stop() {
let dir = TempDir::new().expect("tempdir");
let ckpt_dir = sync_hwm_ckpt_dir(dir.path(), 0);
std::fs::create_dir_all(&ckpt_dir).expect("mkdir");
let path = sync_hwm_ckpt_state_path(&ckpt_dir);
let tmp = ckpt_dir.join("STATE.tmp");
nodedb_wal::segment::write_checkpoint_framed(&tmp, &path, b"not valid msgpack")
.expect("write");
let mut core = open_core_at(dir.path());
assert!(
core.load_sync_hwm_checkpoint().is_err(),
"an undecodable checkpoint body must abort the load, not restore nothing"
);
}
}