use tracing::warn;
use super::core_loop::CoreLoop;
use crate::bridge::envelope::{PhysicalPlan, Status};
use crate::types::{DatabaseId, Lsn, TenantId, VShardId};
use nodedb_physical::physical_plan::ColumnarOp;
use nodedb_types::columnar::ColumnarDmlWalRecord;
impl CoreLoop {
pub(in crate::data::executor) fn try_replay_columnar_predicate_dml(
&mut self,
payload: &[u8],
tenant_id: u64,
database_id: DatabaseId,
record_lsn: u64,
tombstones: &nodedb_wal::TombstoneSet,
) -> Option<usize> {
let record: ColumnarDmlWalRecord = zerompk::from_msgpack(payload).ok()?;
if record.kind != "columnar_dml" {
return None;
}
if tombstones.is_tombstoned(
database_id.as_u64(),
tenant_id,
&record.collection,
record_lsn,
) {
return Some(0);
}
if self.floors.replay_floors.columnar.covers(record_lsn) {
return Some(0);
}
let tid = TenantId::new(tenant_id);
let vshard_id = VShardId::from_collection_in_database(database_id, &record.collection);
let plan = if record.is_update {
PhysicalPlan::Columnar(ColumnarOp::Update {
collection: record.collection.clone(),
filters: record.filters.clone(),
updates: record.updates.clone(),
})
} else {
PhysicalPlan::Columnar(ColumnarOp::Delete {
collection: record.collection.clone(),
filters: record.filters.clone(),
})
};
let task = Self::replay_task(
tid,
database_id,
vshard_id,
plan,
Some(Lsn::new(record_lsn)),
);
let response = if record.is_update {
self.execute_columnar_update(
&task,
&record.collection,
&record.filters,
&record.updates,
None,
)
} else {
self.execute_columnar_delete(&task, &record.collection, &record.filters, None)
};
if response.status != Status::Ok {
warn!(
core = self.core_id,
collection = %record.collection,
lsn = record_lsn,
is_update = record.is_update,
error = ?response.error_code,
"columnar predicate DML WAL replay failed; skipping record"
);
return Some(0);
}
Some(1)
}
}
#[cfg(test)]
mod tests {
use std::sync::Arc;
use crate::bridge::envelope::PhysicalPlan;
use crate::control::server::wal_dispatch::wal_append_if_write;
use crate::types::{DatabaseId, TenantId, VShardId};
use crate::wal::manager::WalManager;
use nodedb_physical::physical_plan::{ColumnarInsertIntent, ColumnarOp};
use nodedb_query::scan_filter::{FilterOp, ScanFilter};
use nodedb_types::Value;
use nodedb_wal::TombstoneSet;
use super::CoreLoop;
const TID: u64 = 1;
const COLLECTION: &str = "m";
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 insert_plan(rows: Vec<Value>) -> PhysicalPlan {
let payload = nodedb_types::value_to_msgpack(&Value::Array(rows))
.expect("encode columnar insert payload");
PhysicalPlan::Columnar(ColumnarOp::Insert {
collection: COLLECTION.into(),
payload,
format: "msgpack".into(),
intent: ColumnarInsertIntent::Insert,
on_conflict_updates: Vec::new(),
surrogates: Vec::new(),
schema_bytes: Vec::new(),
provenance: None,
wal_lsn: None,
})
}
fn row(id: i64, v: i64) -> Value {
Value::Object(std::collections::HashMap::from([
("id".to_string(), Value::Integer(id)),
("v".to_string(), Value::Integer(v)),
]))
}
fn eq_filter_bytes(field: &str, value: Value) -> Vec<u8> {
let filters = vec![ScanFilter {
field: field.to_string(),
op: FilterOp::Eq,
value,
clauses: Vec::new(),
expr: None,
}];
zerompk::to_msgpack_vec(&filters).expect("encode filters")
}
fn append_via_autocommit(plans: &[PhysicalPlan]) -> Vec<nodedb_wal::WalRecord> {
let dir = tempfile::tempdir().expect("wal tempdir");
let wal = WalManager::open_for_testing(&dir.path().join("wal")).expect("open wal");
for plan in plans {
let outcome = wal_append_if_write(
&wal,
TenantId::new(TID),
VShardId::new(0),
DatabaseId::DEFAULT,
plan,
)
.expect("wal append");
assert!(
outcome.lsn.is_some(),
"columnar predicate DML autocommit writes must produce a durable WAL record"
);
}
wal.sync().expect("wal sync");
wal.replay().expect("wal replay read")
}
fn scan_ids(core: &mut CoreLoop) -> Vec<(i64, i64)> {
let key = (
DatabaseId::DEFAULT,
TenantId::new(TID),
COLLECTION.to_string(),
);
let engine = core
.columnar_engines
.get(&key)
.expect("columnar engine present after replay");
let schema = engine.schema();
let id_idx = schema
.columns
.iter()
.position(|c| c.name == "id")
.expect("schema has 'id' column");
let v_idx = schema
.columns
.iter()
.position(|c| c.name == "v")
.expect("schema has 'v' column");
engine
.scan_memtable_rows()
.map(|row| {
let id = match &row[id_idx] {
Value::Integer(n) => *n,
other => panic!("expected integer id, got {other:?}"),
};
let v = match &row[v_idx] {
Value::Integer(n) => *n,
other => panic!("expected integer v, got {other:?}"),
};
(id, v)
})
.collect()
}
#[test]
fn autocommit_delete_produces_durable_lsn() {
let dir = tempfile::tempdir().expect("wal tempdir");
let wal = WalManager::open_for_testing(&dir.path().join("wal")).expect("open wal");
let delete = PhysicalPlan::Columnar(ColumnarOp::Delete {
collection: COLLECTION.into(),
filters: eq_filter_bytes("id", Value::Integer(2)),
});
let outcome = wal_append_if_write(
&wal,
TenantId::new(TID),
VShardId::new(0),
DatabaseId::DEFAULT,
&delete,
)
.expect("wal append delete");
assert!(
outcome.lsn.is_some(),
"autocommit columnar predicate DELETE must be durably WAL-appended \
(pre-fix: ColumnarOp::Delete fell through the catch-all and was never logged)"
);
}
#[test]
fn autocommit_update_produces_durable_lsn() {
let dir = tempfile::tempdir().expect("wal tempdir");
let wal = WalManager::open_for_testing(&dir.path().join("wal")).expect("open wal");
let update = PhysicalPlan::Columnar(ColumnarOp::Update {
collection: COLLECTION.into(),
filters: eq_filter_bytes("id", Value::Integer(1)),
updates: vec![(
"v".to_string(),
nodedb_types::value_to_msgpack(&Value::Integer(999)).expect("encode"),
)],
});
let outcome = wal_append_if_write(
&wal,
TenantId::new(TID),
VShardId::new(0),
DatabaseId::DEFAULT,
&update,
)
.expect("wal append update");
assert!(
outcome.lsn.is_some(),
"autocommit columnar predicate UPDATE must be durably WAL-appended \
(pre-fix: ColumnarOp::Update fell through the catch-all and was never logged)"
);
}
#[test]
fn deleted_rows_do_not_reappear_after_replay_from_empty() {
let insert = insert_plan(vec![row(1, 10), row(2, 20), row(3, 30)]);
let delete = PhysicalPlan::Columnar(ColumnarOp::Delete {
collection: COLLECTION.into(),
filters: eq_filter_bytes("id", Value::Integer(2)),
});
let records = append_via_autocommit(&[insert, delete]);
let mut h = make_core();
h.core
.replay_timeseries_wal(&records, 1, &TombstoneSet::new());
let mut rows = scan_ids(&mut h.core);
rows.sort();
assert_eq!(
rows,
vec![(1, 10), (3, 30)],
"deleted row (id=2) must NOT reappear after replay from empty \
(pre-fix: the Delete was never WAL-logged, so replay only re-ran \
the Insert and the deleted row resurrected)"
);
}
#[test]
fn updated_row_present_exactly_once_after_replay_from_empty() {
let insert = insert_plan(vec![row(1, 10), row(2, 20)]);
let update = PhysicalPlan::Columnar(ColumnarOp::Update {
collection: COLLECTION.into(),
filters: eq_filter_bytes("id", Value::Integer(1)),
updates: vec![(
"v".to_string(),
nodedb_types::value_to_msgpack(&Value::Integer(999)).expect("encode"),
)],
});
let records = append_via_autocommit(&[insert, update]);
let mut h = make_core();
h.core
.replay_timeseries_wal(&records, 1, &TombstoneSet::new());
let mut rows = scan_ids(&mut h.core);
rows.sort();
assert_eq!(
rows,
vec![(1, 999), (2, 20)],
"updated row must carry the new value exactly once after replay \
(pre-fix: the Update was never WAL-logged, so the value reverted; \
a non-idempotent double-apply would instead duplicate the row)"
);
}
}