Skip to content

Commit 8ccecfb

Browse files
committed
feat(nodedb-wal): add TemporalPurge WAL record type and payload
Introduce `RecordType::TemporalPurge` (opcode 103 | required-replay flag) to distinguish bitemporal version purges from regular `Delete` records. Replay must apply purge records to avoid resurrectin superseded versions that the leader has already dropped. Add `TemporalPurgePayload` and `TemporalPurgeEngine` in a new `temporal_purge` module to carry engine tag, collection name, cutoff timestamp, and purged-row count for auditing and crash-recovery.
1 parent 1937199 commit 8ccecfb

3 files changed

Lines changed: 213 additions & 0 deletions

File tree

nodedb-wal/src/lib.rs

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -34,6 +34,7 @@ pub mod recovery;
3434
pub mod replay;
3535
pub mod segment;
3636
pub mod segmented;
37+
pub mod temporal_purge;
3738
pub mod tombstone;
3839
#[cfg(feature = "io-uring")]
3940
pub mod uring_writer;
@@ -47,5 +48,6 @@ pub use record::{RecordHeader, RecordType, WalRecord};
4748
pub use recovery::{RecoveryInfo, recover};
4849
pub use replay::{TombstoneSet, extract_tombstones};
4950
pub use segmented::{SegmentedWal, SegmentedWalConfig};
51+
pub use temporal_purge::{TemporalPurgeEngine, TemporalPurgePayload};
5052
pub use tombstone::{CollectionTombstonePayload, MAX_COLLECTION_NAME_LEN};
5153
pub use writer::WalWriter;

nodedb-wal/src/record/types.rs

Lines changed: 13 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -54,6 +54,16 @@ pub enum RecordType {
5454
/// Not required: a replay that skips these records produces a slightly
5555
/// coarser interpolation table but does not corrupt state.
5656
LsnMsAnchor = 102,
57+
58+
/// Bitemporal version purge — drops one or more *superseded* row
59+
/// versions (those with finite `_ts_valid_until`) once
60+
/// `audit_retain_ms` has elapsed. Distinct from `Delete`, which
61+
/// removes the current live row; replay must not conflate them
62+
/// because a `TemporalPurge` must never delete live state.
63+
///
64+
/// Required: a replay that skipped this record would leave purged
65+
/// versions resurrected and diverge from the leader's state.
66+
TemporalPurge = 103 | 0x8000,
5767
}
5868

5969
impl RecordType {
@@ -78,6 +88,7 @@ impl RecordType {
7888
x if x == 100 | 0x8000 => Some(Self::Checkpoint),
7989
x if x == 101 | 0x8000 => Some(Self::CollectionTombstoned),
8090
102 => Some(Self::LsnMsAnchor),
91+
x if x == 103 | 0x8000 => Some(Self::TemporalPurge),
8192
_ => None,
8293
}
8394
}
@@ -96,6 +107,7 @@ mod tests {
96107
assert!(!RecordType::is_required(RecordType::TimeseriesBatch as u16));
97108
assert!(!RecordType::is_required(RecordType::LogBatch as u16));
98109
assert!(!RecordType::is_required(RecordType::LsnMsAnchor as u16));
110+
assert!(RecordType::is_required(RecordType::TemporalPurge as u16));
99111
}
100112

101113
#[test]
@@ -114,6 +126,7 @@ mod tests {
114126
RecordType::Checkpoint,
115127
RecordType::CollectionTombstoned,
116128
RecordType::LsnMsAnchor,
129+
RecordType::TemporalPurge,
117130
] {
118131
assert_eq!(RecordType::from_raw(ty as u16), Some(ty));
119132
}

nodedb-wal/src/temporal_purge.rs

Lines changed: 198 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,198 @@
1+
//! `TemporalPurge` record payload codec.
2+
//!
3+
//! Emitted when the Control Plane's bitemporal-retention scheduler runs
4+
//! an audit-retention pass on a bitemporal collection and successfully
5+
//! drops some number of *superseded* versions below `cutoff_system_ms`.
6+
//!
7+
//! Distinct from `RecordType::Delete`: a `TemporalPurge` never removes
8+
//! the live / current state of a row — only history older than the
9+
//! audit-retain window. Crash recovery MUST treat these separately so a
10+
//! mid-flight purge that was interrupted does not resurface as a delete
11+
//! of the surviving latest version.
12+
//!
13+
//! Fixed little-endian wire format (no serde dep):
14+
//!
15+
//! ```text
16+
//! ┌────────────┬────────────┬───────────┬──────────────────┬────────────┐
17+
//! │engine_tag │name_len u32│name bytes │ cutoff_ms i64 │ count u64 │
18+
//! │ u8 │ │ │ │ │
19+
//! └────────────┴────────────┴───────────┴──────────────────┴────────────┘
20+
//! ```
21+
//!
22+
//! Tenant id lives on the record header, so it is not repeated here.
23+
24+
use crate::error::{Result, WalError};
25+
use crate::tombstone::MAX_COLLECTION_NAME_LEN;
26+
27+
/// Which engine produced the purge. Wire-stable — do not renumber.
28+
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
29+
#[repr(u8)]
30+
pub enum TemporalPurgeEngine {
31+
EdgeStore = 1,
32+
DocumentStrict = 2,
33+
Columnar = 3,
34+
}
35+
36+
impl TemporalPurgeEngine {
37+
pub fn from_raw(raw: u8) -> Option<Self> {
38+
match raw {
39+
1 => Some(Self::EdgeStore),
40+
2 => Some(Self::DocumentStrict),
41+
3 => Some(Self::Columnar),
42+
_ => None,
43+
}
44+
}
45+
}
46+
47+
/// Parsed temporal-purge payload.
48+
#[derive(Debug, Clone, PartialEq, Eq)]
49+
pub struct TemporalPurgePayload {
50+
pub engine: TemporalPurgeEngine,
51+
pub collection: String,
52+
pub cutoff_system_ms: i64,
53+
pub purged_count: u64,
54+
}
55+
56+
impl TemporalPurgePayload {
57+
pub fn new(
58+
engine: TemporalPurgeEngine,
59+
collection: impl Into<String>,
60+
cutoff_system_ms: i64,
61+
purged_count: u64,
62+
) -> Self {
63+
Self {
64+
engine,
65+
collection: collection.into(),
66+
cutoff_system_ms,
67+
purged_count,
68+
}
69+
}
70+
71+
pub fn wire_size(&self) -> usize {
72+
1 + 4 + self.collection.len() + 8 + 8
73+
}
74+
75+
pub fn to_bytes(&self) -> Result<Vec<u8>> {
76+
let name_bytes = self.collection.as_bytes();
77+
if name_bytes.len() > MAX_COLLECTION_NAME_LEN {
78+
return Err(WalError::PayloadTooLarge {
79+
size: name_bytes.len(),
80+
max: MAX_COLLECTION_NAME_LEN,
81+
});
82+
}
83+
let mut buf = Vec::with_capacity(self.wire_size());
84+
buf.push(self.engine as u8);
85+
buf.extend_from_slice(&(name_bytes.len() as u32).to_le_bytes());
86+
buf.extend_from_slice(name_bytes);
87+
buf.extend_from_slice(&self.cutoff_system_ms.to_le_bytes());
88+
buf.extend_from_slice(&self.purged_count.to_le_bytes());
89+
Ok(buf)
90+
}
91+
92+
pub fn from_bytes(buf: &[u8]) -> Result<Self> {
93+
if buf.len() < 1 + 4 {
94+
return Err(WalError::CorruptRecord {
95+
lsn: 0,
96+
detail: "temporal-purge payload shorter than engine_tag + name_len".into(),
97+
});
98+
}
99+
let engine =
100+
TemporalPurgeEngine::from_raw(buf[0]).ok_or_else(|| WalError::CorruptRecord {
101+
lsn: 0,
102+
detail: format!("temporal-purge unknown engine_tag {}", buf[0]),
103+
})?;
104+
let name_len = u32::from_le_bytes([buf[1], buf[2], buf[3], buf[4]]) as usize;
105+
if name_len > MAX_COLLECTION_NAME_LEN {
106+
return Err(WalError::CorruptRecord {
107+
lsn: 0,
108+
detail: format!("temporal-purge name_len {name_len} exceeds max"),
109+
});
110+
}
111+
let need = 1 + 4 + name_len + 8 + 8;
112+
if buf.len() < need {
113+
return Err(WalError::CorruptRecord {
114+
lsn: 0,
115+
detail: format!(
116+
"temporal-purge payload truncated: need {need} bytes, have {}",
117+
buf.len()
118+
),
119+
});
120+
}
121+
let name_end = 5 + name_len;
122+
let collection = std::str::from_utf8(&buf[5..name_end])
123+
.map_err(|e| WalError::CorruptRecord {
124+
lsn: 0,
125+
detail: format!("temporal-purge collection not utf8: {e}"),
126+
})?
127+
.to_string();
128+
let cutoff_system_ms = i64::from_le_bytes(
129+
buf[name_end..name_end + 8]
130+
.try_into()
131+
.expect("bounded above"),
132+
);
133+
let purged_count = u64::from_le_bytes(
134+
buf[name_end + 8..name_end + 16]
135+
.try_into()
136+
.expect("bounded above"),
137+
);
138+
Ok(Self {
139+
engine,
140+
collection,
141+
cutoff_system_ms,
142+
purged_count,
143+
})
144+
}
145+
}
146+
147+
#[cfg(test)]
148+
mod tests {
149+
use super::*;
150+
151+
#[test]
152+
fn roundtrip() {
153+
let p =
154+
TemporalPurgePayload::new(TemporalPurgeEngine::EdgeStore, "users", 1_000_000_000, 42);
155+
let bytes = p.to_bytes().unwrap();
156+
let decoded = TemporalPurgePayload::from_bytes(&bytes).unwrap();
157+
assert_eq!(p, decoded);
158+
}
159+
160+
#[test]
161+
fn all_engine_tags_roundtrip() {
162+
for e in [
163+
TemporalPurgeEngine::EdgeStore,
164+
TemporalPurgeEngine::DocumentStrict,
165+
TemporalPurgeEngine::Columnar,
166+
] {
167+
let p = TemporalPurgePayload::new(e, "c", 0, 0);
168+
let b = p.to_bytes().unwrap();
169+
assert_eq!(TemporalPurgePayload::from_bytes(&b).unwrap().engine, e);
170+
}
171+
}
172+
173+
#[test]
174+
fn rejects_unknown_engine_tag() {
175+
let mut buf = TemporalPurgePayload::new(TemporalPurgeEngine::EdgeStore, "c", 0, 0)
176+
.to_bytes()
177+
.unwrap();
178+
buf[0] = 99;
179+
assert!(TemporalPurgePayload::from_bytes(&buf).is_err());
180+
}
181+
182+
#[test]
183+
fn rejects_truncated() {
184+
let full = TemporalPurgePayload::new(TemporalPurgeEngine::Columnar, "users", 1, 1)
185+
.to_bytes()
186+
.unwrap();
187+
for cut in 0..full.len() {
188+
assert!(TemporalPurgePayload::from_bytes(&full[..cut]).is_err());
189+
}
190+
}
191+
192+
#[test]
193+
fn rejects_oversize_name() {
194+
let long = "a".repeat(MAX_COLLECTION_NAME_LEN + 1);
195+
let p = TemporalPurgePayload::new(TemporalPurgeEngine::EdgeStore, long, 0, 0);
196+
assert!(p.to_bytes().is_err());
197+
}
198+
}

0 commit comments

Comments
 (0)