#![deny(clippy::wildcard_enum_match_arm)]
use nodedb_physical::physical_plan::ColumnarOp;
use crate::types::{DatabaseId, Lsn, TenantId, VShardId};
use crate::wal::manager::WalManager;
pub(super) fn wal_append_columnar_op(
wal: &WalManager,
tenant_id: TenantId,
vshard_id: VShardId,
database_id: DatabaseId,
op: &ColumnarOp,
) -> crate::Result<Option<Lsn>> {
let appended = match op {
ColumnarOp::Insert {
collection,
payload,
format: _,
intent: _,
on_conflict_updates: _,
surrogates,
schema_bytes: _,
provenance,
wal_lsn: _,
} => {
let wal_payload = super::timeseries::encode_columnar_batch_payload(
collection,
payload,
provenance.as_ref(),
surrogates,
)?;
Some(wal.append_timeseries_batch(tenant_id, vshard_id, database_id, &wal_payload)?)
}
ColumnarOp::Update {
collection,
filters,
updates,
} => {
let wal_payload =
super::timeseries::encode_columnar_dml_payload(collection, true, filters, updates)?;
Some(wal.append_timeseries_batch(tenant_id, vshard_id, database_id, &wal_payload)?)
}
ColumnarOp::Delete {
collection,
filters,
} => {
let wal_payload =
super::timeseries::encode_columnar_dml_payload(collection, false, filters, &[])?;
Some(wal.append_timeseries_batch(tenant_id, vshard_id, database_id, &wal_payload)?)
}
ColumnarOp::Scan { .. } | ColumnarOp::MaterializeScan { .. } => None,
};
Ok(appended)
}
#[cfg(test)]
mod tests {
use super::*;
use nodedb_physical::physical_plan::{ColumnarInsertIntent, PhysicalPlan};
fn open_wal(dir: &std::path::Path) -> WalManager {
WalManager::open_for_testing(&dir.join("test.wal")).expect("open wal")
}
fn has_record_of_type(wal: &WalManager, record_type: nodedb_wal::record::RecordType) -> bool {
wal.sync().expect("sync wal");
wal.replay().expect("read wal").into_iter().any(|r| {
nodedb_wal::record::RecordType::from_raw(r.logical_record_type()) == Some(record_type)
})
}
#[test]
fn insert_appends_timeseries_batch_record() {
let dir = tempfile::tempdir().expect("tempdir");
let wal = open_wal(dir.path());
let plan = PhysicalPlan::Columnar(ColumnarOp::Insert {
collection: "metrics".to_string(),
payload: vec![1, 2, 3],
format: "msgpack".to_string(),
intent: ColumnarInsertIntent::Insert,
on_conflict_updates: vec![],
surrogates: vec![],
schema_bytes: vec![],
provenance: None,
wal_lsn: None,
});
let outcome = super::super::wal_append_if_write(
&wal,
TenantId::new(1),
VShardId::new(0),
DatabaseId::DEFAULT,
&plan,
)
.expect("append");
assert!(outcome.lsn.is_some(), "Insert must produce a durable LSN");
assert!(has_record_of_type(
&wal,
nodedb_wal::record::RecordType::TimeseriesBatch
));
}
#[test]
fn scan_appends_nothing() {
let dir = tempfile::tempdir().expect("tempdir");
let wal = open_wal(dir.path());
let plan = PhysicalPlan::Columnar(ColumnarOp::Scan {
collection: "metrics".to_string(),
projection: vec![],
limit: 10,
filters: vec![],
rls_filters: vec![],
sort_keys: vec![],
system_time: Default::default(),
valid_at_ms: None,
prefilter: None,
computed_columns: vec![],
});
let outcome = super::super::wal_append_if_write(
&wal,
TenantId::new(1),
VShardId::new(0),
DatabaseId::DEFAULT,
&plan,
)
.expect("append");
assert!(outcome.lsn.is_none(), "Scan must produce no durable LSN");
}
}