#![deny(clippy::wildcard_enum_match_arm)]
use nodedb_physical::physical_plan::TimeseriesOp;
use crate::control::security::credential::CredentialStore;
use crate::types::{DatabaseId, Lsn, TenantId, VShardId};
use crate::wal::manager::WalManager;
pub(super) fn wal_append_timeseries_op(
wal: &WalManager,
tenant_id: TenantId,
vshard_id: VShardId,
database_id: DatabaseId,
op: &TimeseriesOp,
credentials: Option<&CredentialStore>,
) -> crate::Result<Option<Lsn>> {
let appended = match op {
TimeseriesOp::Ingest {
collection,
payload,
format: _,
provenance,
..
} => {
if let Some(creds) = credentials
&& let Ok(Some(coll)) = creds.catalog().get_collection(
DatabaseId::DEFAULT,
tenant_id.as_u64(),
collection,
)
&& let Some(config) = coll.get_timeseries_config()
&& config.get("wal").and_then(|v| v.as_str()) == Some("false")
{
None
} else {
let wal_payload =
encode_timeseries_batch_payload(collection, payload, provenance.as_ref())?;
Some(wal.append_timeseries_batch(
tenant_id,
vshard_id,
database_id,
&wal_payload,
)?)
}
}
TimeseriesOp::Scan { .. } => None,
};
Ok(appended)
}
pub(crate) fn encode_timeseries_batch_payload(
collection: &str,
payload: &[u8],
provenance: Option<&nodedb_types::sync::wire::SyncProvenance>,
) -> crate::Result<Vec<u8>> {
zerompk::to_msgpack_vec(&("timeseries", collection, payload, provenance)).map_err(|e| {
crate::Error::Serialization {
format: "msgpack".into(),
detail: format!("wal timeseries batch: {e}"),
}
})
}
pub(crate) fn encode_columnar_batch_payload(
collection: &str,
payload: &[u8],
provenance: Option<&nodedb_types::sync::wire::SyncProvenance>,
surrogates: &[nodedb_types::Surrogate],
) -> crate::Result<Vec<u8>> {
let record = nodedb_types::columnar::ColumnarWalRecord {
kind: "columnar".to_string(),
collection: collection.to_string(),
payload: payload.to_vec(),
provenance: provenance.cloned(),
surrogates: surrogates.to_vec(),
};
zerompk::to_msgpack_vec(&record).map_err(|e| crate::Error::Serialization {
format: "msgpack".into(),
detail: format!("wal columnar batch: {e}"),
})
}
pub fn wal_append_timeseries(
wal: &WalManager,
tenant_id: TenantId,
vshard_id: VShardId,
collection: &str,
payload: &[u8],
provenance: Option<&nodedb_types::sync::wire::SyncProvenance>,
credentials: Option<&CredentialStore>,
) -> crate::Result<Option<nodedb_types::Lsn>> {
let database_id = DatabaseId::DEFAULT;
if let Some(creds) = credentials
&& let Ok(Some(coll)) =
creds
.catalog()
.get_collection(database_id, tenant_id.as_u64(), collection)
&& let Some(config) = coll.get_timeseries_config()
&& config.get("wal").and_then(|v| v.as_str()) == Some("false")
{
return Ok(None);
}
let wal_payload = encode_timeseries_batch_payload(collection, payload, provenance)?;
let lsn = wal.append_timeseries_batch(tenant_id, vshard_id, database_id, &wal_payload)?;
Ok(Some(lsn))
}
pub(crate) fn encode_columnar_dml_payload(
collection: &str,
is_update: bool,
filters: &[u8],
updates: &[(String, Vec<u8>)],
) -> crate::Result<Vec<u8>> {
let record = nodedb_types::columnar::ColumnarDmlWalRecord {
kind: "columnar_dml".to_string(),
collection: collection.to_string(),
is_update,
filters: filters.to_vec(),
updates: updates.to_vec(),
};
zerompk::to_msgpack_vec(&record).map_err(|e| crate::Error::Serialization {
format: "msgpack".into(),
detail: format!("wal columnar dml: {e}"),
})
}
pub struct ColumnarWalAppendArgs<'a> {
pub collection: &'a str,
pub payload: &'a [u8],
pub provenance: Option<&'a nodedb_types::sync::wire::SyncProvenance>,
pub surrogates: &'a [nodedb_types::Surrogate],
}
pub fn wal_append_columnar(
wal: &WalManager,
tenant_id: TenantId,
vshard_id: VShardId,
database_id: DatabaseId,
args: ColumnarWalAppendArgs<'_>,
) -> crate::Result<Option<nodedb_types::Lsn>> {
let ColumnarWalAppendArgs {
collection,
payload,
provenance,
surrogates,
} = args;
let wal_payload = encode_columnar_batch_payload(collection, payload, provenance, surrogates)?;
let lsn = wal.append_timeseries_batch(tenant_id, vshard_id, database_id, &wal_payload)?;
Ok(Some(lsn))
}
#[cfg(test)]
mod tests {
use super::*;
use nodedb_physical::physical_plan::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 ingest_appends_timeseries_batch_record() {
let dir = tempfile::tempdir().expect("tempdir");
let wal = open_wal(dir.path());
let plan = PhysicalPlan::Timeseries(TimeseriesOp::Ingest {
collection: "metrics".to_string(),
payload: vec![1, 2, 3],
format: "samples".to_string(),
wal_lsn: None,
surrogates: vec![],
provenance: 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(), "Ingest 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::Timeseries(TimeseriesOp::Scan {
collection: "metrics".to_string(),
time_range: (0, i64::MAX),
projection: vec![],
limit: 10,
filters: vec![],
bucket_interval_ms: 0,
group_by: vec![],
aggregates: vec![],
gap_fill: String::new(),
computed_columns: vec![],
rls_filters: vec![],
system_time: Default::default(),
valid_at_ms: 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_none(), "Scan must produce no durable LSN");
}
}