#[cfg(test)]
mod tests {
use std::sync::Arc;
use crate::bridge::envelope::{
Admission, ExemptReason, PhysicalPlan, Priority, Request, Response,
};
use crate::control::server::wal_dispatch::wal_append_if_write;
use crate::data::executor::core_loop::CoreLoop;
use crate::data::executor::handlers::kv::crud::KvWriteParams;
use crate::data::executor::task::ExecutionTask;
use crate::types::{DatabaseId, ReadConsistency, TenantId, TraceId, VShardId};
use crate::wal::manager::WalManager;
use nodedb_physical::physical_plan::KvOp;
use nodedb_types::Surrogate;
use nodedb_wal::TombstoneSet;
const TID: u64 = 1;
struct CoreHarness {
core: CoreLoop,
_req_tx: nodedb_bridge::buffer::Producer<crate::bridge::dispatch::BridgeRequest>,
_resp_rx: nodedb_bridge::buffer::Consumer<crate::bridge::dispatch::BridgeResponse>,
_dir: tempfile::TempDir,
}
fn make_core() -> CoreHarness {
use crate::bridge::dispatch::{BridgeRequest, BridgeResponse};
use nodedb_bridge::buffer::RingBuffer;
let dir = tempfile::tempdir().expect("tempdir");
let (req_tx, req_rx) = RingBuffer::channel::<BridgeRequest>(64);
let (resp_tx, resp_rx) = RingBuffer::channel::<BridgeResponse>(64);
let core = CoreLoop::open(
0,
req_rx,
resp_tx,
dir.path(),
Arc::new(nodedb_types::OrdinalClock::new()),
)
.expect("open core");
CoreHarness {
core,
_req_tx: req_tx,
_resp_rx: resp_rx,
_dir: dir,
}
}
fn installed_expire_at_ms(core: &CoreLoop, collection: &str, key: &[u8]) -> i64 {
core.kv_engine
.get_ttl_ms(DatabaseId::DEFAULT.as_u64(), TID, collection, key, 0)
.expect("key must exist with a TTL")
}
#[test]
fn kv_put_replay_installs_recorded_absolute_expiry_not_replay_time_clock() {
let entry = crate::control::server::wal_dispatch_kv::encode::encode_kv_put(
"sessions",
b"tok1",
b"payload",
5_000,
Some(6_000),
)
.expect("encode kv_put with absolute expiry");
let dir = tempfile::tempdir().expect("wal tempdir");
let wal = WalManager::open_for_testing(&dir.path().join("wal")).expect("open wal");
wal.append_put(
TenantId::new(TID),
VShardId::new(0),
DatabaseId::DEFAULT,
&entry,
)
.expect("append raw kv_put record");
wal.sync().expect("wal sync");
let records = wal.replay().expect("wal replay read");
let mut h = make_core();
h.core.replay_kv_wal(&records, 1, &TombstoneSet::new());
assert_eq!(
installed_expire_at_ms(&h.core, "sessions", b"tok1"),
6_000,
"replay must install the recorded absolute expiry verbatim, not \
recompute now_ms + ttl_ms at replay time"
);
}
#[test]
fn kv_batch_put_replay_installs_recorded_absolute_expiry_not_replay_time_clock() {
let entries = vec![
(b"k1".to_vec(), b"v1".to_vec()),
(b"k2".to_vec(), b"v2".to_vec()),
];
let entry = crate::control::server::wal_dispatch_kv::encode::encode_kv_batch_put(
"carts",
&entries,
5_000,
Some(6_000),
)
.expect("encode kv_batch_put with absolute expiry");
let dir = tempfile::tempdir().expect("wal tempdir");
let wal = WalManager::open_for_testing(&dir.path().join("wal")).expect("open wal");
wal.append_put(
TenantId::new(TID),
VShardId::new(0),
DatabaseId::DEFAULT,
&entry,
)
.expect("append raw kv_batch_put record");
wal.sync().expect("wal sync");
let records = wal.replay().expect("wal replay read");
let mut h = make_core();
h.core.replay_kv_wal(&records, 1, &TombstoneSet::new());
for key in [b"k1".as_slice(), b"k2".as_slice()] {
assert_eq!(
installed_expire_at_ms(&h.core, "carts", key),
6_000,
"batch_put replay must install the recorded absolute expiry \
verbatim for every entry, not recompute now_ms + ttl_ms"
);
}
}
#[test]
fn production_wal_append_emits_six_element_shape_for_ttl_put() {
let observed_now_ms = crate::engine::kv::current_ms();
let dir = tempfile::tempdir().expect("wal tempdir");
let wal = WalManager::open_for_testing(&dir.path().join("wal")).expect("open wal");
let plan = PhysicalPlan::Kv(KvOp::Put {
collection: "sessions".into(),
key: b"tok2".to_vec(),
value: b"payload".to_vec(),
ttl_ms: 5_000,
surrogate: Surrogate::new(1),
});
let outcome = wal_append_if_write(
&wal,
TenantId::new(TID),
VShardId::new(0),
DatabaseId::DEFAULT,
&plan,
)
.expect("wal append");
assert!(outcome.lsn.is_some());
let resolved = outcome
.resolved_now_ms
.expect("TTL-bearing Put must resolve a TTL instant");
assert!(
resolved >= observed_now_ms,
"resolved_now_ms ({resolved}) must be at or after the instant \
observed just before the call ({observed_now_ms})"
);
wal.sync().expect("wal sync");
let records = wal.replay().expect("wal replay read");
let record = records
.iter()
.find(|r| r.header.tenant_id == TID)
.expect("one record");
let (disc, _collection, _key, _value, ttl_ms, expire_at_ms) =
zerompk::from_msgpack::<(&str, String, Vec<u8>, Vec<u8>, u64, u64)>(&record.payload)
.expect("six-element kv_put shape");
assert_eq!(disc, "kv_put");
assert_eq!(ttl_ms, 5_000);
assert_eq!(expire_at_ms, resolved + 5_000);
}
#[test]
fn production_wal_append_emits_five_element_shape_for_no_ttl_put() {
let dir = tempfile::tempdir().expect("wal tempdir");
let wal = WalManager::open_for_testing(&dir.path().join("wal")).expect("open wal");
let plan = PhysicalPlan::Kv(KvOp::Put {
collection: "sessions".into(),
key: b"tok3".to_vec(),
value: b"payload".to_vec(),
ttl_ms: 0,
surrogate: Surrogate::new(1),
});
let outcome = wal_append_if_write(
&wal,
TenantId::new(TID),
VShardId::new(0),
DatabaseId::DEFAULT,
&plan,
)
.expect("wal append");
assert_eq!(
outcome.resolved_now_ms, None,
"a non-TTL Put must not resolve a TTL instant"
);
wal.sync().expect("wal sync");
let records = wal.replay().expect("wal replay read");
let record = &records[0];
zerompk::from_msgpack::<(&str, String, Vec<u8>, Vec<u8>, u64)>(&record.payload)
.expect("five-element kv_put shape for a non-TTL write");
}
fn task_with_resolved_now_ms(
plan: PhysicalPlan,
resolved_now_ms: Option<u64>,
) -> ExecutionTask {
ExecutionTask::new(Request {
request_id: crate::types::RequestId::new(1),
tenant_id: TenantId::new(TID),
database_id: DatabaseId::DEFAULT,
vshard_id: VShardId::new(0),
plan,
deadline: std::time::Instant::now() + std::time::Duration::from_secs(5),
priority: Priority::Normal,
trace_id: TraceId::ZERO,
consistency: ReadConsistency::Strong,
idempotency_key: None,
event_source: crate::event::EventSource::User,
user_roles: Vec::new(),
user_id: None,
statement_digest: None,
txn_id: None,
wal_lsn: None,
resolved_now_ms,
admission: Admission::Exempt(ExemptReason::Read),
})
}
#[test]
fn execute_kv_put_installs_task_resolved_now_ms_verbatim() {
let mut h = make_core();
let plan = PhysicalPlan::Kv(KvOp::Put {
collection: "sessions".into(),
key: b"tok4".to_vec(),
value: b"payload".to_vec(),
ttl_ms: 5_000,
surrogate: Surrogate::new(1),
});
let task = task_with_resolved_now_ms(plan, Some(1_000));
let resp: Response = h.core.execute_kv_put(
&task,
KvWriteParams {
did: DatabaseId::DEFAULT.as_u64(),
tid: TID,
collection: "sessions",
key: b"tok4",
value: b"payload",
ttl_ms: 5_000,
surrogate: Surrogate::new(1),
},
);
assert_eq!(resp.status, crate::bridge::envelope::Status::Ok);
assert_eq!(
installed_expire_at_ms(&h.core, "sessions", b"tok4"),
1_000 + 5_000,
"live apply must install resolved_now_ms + ttl_ms, not the wall clock"
);
}
#[test]
fn execute_kv_insert_installs_task_resolved_now_ms_verbatim() {
let mut h = make_core();
let plan = PhysicalPlan::Kv(KvOp::Insert {
collection: "sessions".into(),
key: b"tok5".to_vec(),
value: b"payload".to_vec(),
ttl_ms: 5_000,
surrogate: Surrogate::new(1),
});
let task = task_with_resolved_now_ms(plan, Some(1_000));
let resp: Response = h.core.execute_kv_insert(
&task,
KvWriteParams {
did: DatabaseId::DEFAULT.as_u64(),
tid: TID,
collection: "sessions",
key: b"tok5",
value: b"payload",
ttl_ms: 5_000,
surrogate: Surrogate::new(1),
},
);
assert_eq!(resp.status, crate::bridge::envelope::Status::Ok);
assert_eq!(
installed_expire_at_ms(&h.core, "sessions", b"tok5"),
1_000 + 5_000,
"execute_kv_insert must install resolved_now_ms + ttl_ms, not the wall clock"
);
}
}