#![deny(clippy::wildcard_enum_match_arm)]
use crate::bridge::envelope::PhysicalPlan;
use crate::control::security::credential::CredentialStore;
use crate::types::{DatabaseId, TenantId, VShardId};
use crate::wal::manager::WalManager;
use super::super::wal_dispatch_kv;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct WalAppendOutcome {
pub lsn: Option<crate::types::Lsn>,
pub resolved_now_ms: Option<u64>,
}
pub struct WalAppendRequest<'a> {
pub wal: &'a WalManager,
pub tenant_id: TenantId,
pub vshard_id: VShardId,
pub database_id: DatabaseId,
pub plan: &'a PhysicalPlan,
pub credentials: Option<&'a CredentialStore>,
pub now_override: Option<u64>,
}
pub fn wal_append_if_write(
wal: &WalManager,
tenant_id: TenantId,
vshard_id: VShardId,
database_id: DatabaseId,
plan: &PhysicalPlan,
) -> crate::Result<WalAppendOutcome> {
wal_append_if_write_with_creds(wal, tenant_id, vshard_id, database_id, plan, None)
}
pub fn wal_append_if_write_with_creds(
wal: &WalManager,
tenant_id: TenantId,
vshard_id: VShardId,
database_id: DatabaseId,
plan: &PhysicalPlan,
credentials: Option<&CredentialStore>,
) -> crate::Result<WalAppendOutcome> {
wal_append(WalAppendRequest {
wal,
tenant_id,
vshard_id,
database_id,
plan,
credentials,
now_override: None,
})
}
pub fn wal_append(req: WalAppendRequest<'_>) -> crate::Result<WalAppendOutcome> {
let WalAppendRequest {
wal,
tenant_id,
vshard_id,
database_id,
plan,
credentials,
now_override,
} = req;
let mut resolved_now_ms: Option<u64> = None;
let appended: Option<crate::types::Lsn> = match plan {
PhysicalPlan::Document(op) => {
super::document::wal_append_document_op(wal, tenant_id, vshard_id, database_id, op)?
}
PhysicalPlan::Vector(op) => {
super::vector::wal_append_vector_op(wal, tenant_id, vshard_id, database_id, op)?
}
PhysicalPlan::Crdt(op) => {
super::crdt::wal_append_crdt_op(wal, tenant_id, vshard_id, database_id, op)?
}
PhysicalPlan::Graph(op) => {
super::graph::wal_append_graph_op(wal, tenant_id, vshard_id, database_id, op)?
}
PhysicalPlan::Columnar(op) => {
super::columnar::wal_append_columnar_op(wal, tenant_id, vshard_id, database_id, op)?
}
PhysicalPlan::Timeseries(op) => super::timeseries::wal_append_timeseries_op(
wal,
tenant_id,
vshard_id,
database_id,
op,
credentials,
)?,
PhysicalPlan::Kv(kv_op) => {
let outcome = wal_dispatch_kv::wal_append_kv_op(
wal,
tenant_id,
vshard_id,
database_id,
kv_op,
now_override,
)?;
resolved_now_ms = outcome.resolved_now_ms;
outcome.lsn
}
PhysicalPlan::Array(op) => {
super::array::wal_append_array_op(wal, tenant_id, vshard_id, database_id, op)?
}
PhysicalPlan::Text(op) => {
super::text::wal_append_text_op(wal, tenant_id, vshard_id, database_id, op)?
}
PhysicalPlan::Spatial(op) => {
super::spatial::wal_append_spatial_op(wal, tenant_id, vshard_id, database_id, op)?
}
PhysicalPlan::Meta(_) | PhysicalPlan::Query(_) | PhysicalPlan::ClusterArray(_) => None,
};
Ok(WalAppendOutcome {
lsn: appended,
resolved_now_ms,
})
}
#[cfg(test)]
mod tests {
use super::*;
use nodedb_physical::physical_plan::{SpatialOp, TextOp};
use nodedb_types::Surrogate;
use nodedb_types::geometry::Geometry;
fn open_wal(dir: &std::path::Path) -> WalManager {
WalManager::open_for_testing(&dir.join("test.wal")).expect("open wal")
}
fn last_record_of_type(
wal: &WalManager,
record_type: nodedb_wal::record::RecordType,
) -> nodedb_wal::WalRecord {
wal.sync().expect("sync wal");
wal.replay()
.expect("read wal")
.into_iter()
.rfind(|r| {
nodedb_wal::record::RecordType::from_raw(r.logical_record_type())
== Some(record_type)
})
.expect("expected record of this type")
}
#[test]
fn fts_index_doc_appends_and_decodes() {
let dir = tempfile::tempdir().expect("tempdir");
let wal = open_wal(dir.path());
let plan = PhysicalPlan::Text(TextOp::FtsIndexDoc {
collection: "docs".to_string(),
surrogate: Surrogate::new(7),
text: "hello world".to_string(),
provenance: None,
});
let outcome = wal_append_if_write(
&wal,
TenantId::new(1),
VShardId::new(0),
DatabaseId::DEFAULT,
&plan,
)
.expect("append");
assert!(
outcome.lsn.is_some(),
"FtsIndexDoc must produce a durable LSN"
);
let record = last_record_of_type(&wal, nodedb_wal::record::RecordType::FtsIndex);
let decoded =
nodedb_wal::record::FtsIndexPayload::from_bytes(&record.payload).expect("decode");
assert_eq!(decoded.collection, "docs");
assert_eq!(decoded.text, "hello world");
assert_eq!(
decoded.doc_id,
crate::engine::document::store::surrogate_to_doc_id(Surrogate::new(7))
);
}
#[test]
fn fts_delete_doc_appends_and_decodes() {
let dir = tempfile::tempdir().expect("tempdir");
let wal = open_wal(dir.path());
let plan = PhysicalPlan::Text(TextOp::FtsDeleteDoc {
collection: "docs".to_string(),
surrogate: Surrogate::new(7),
provenance: None,
});
let outcome = wal_append_if_write(
&wal,
TenantId::new(1),
VShardId::new(0),
DatabaseId::DEFAULT,
&plan,
)
.expect("append");
assert!(
outcome.lsn.is_some(),
"FtsDeleteDoc must produce a durable LSN"
);
let record = last_record_of_type(&wal, nodedb_wal::record::RecordType::FtsDelete);
let decoded =
nodedb_wal::record::FtsDeletePayload::from_bytes(&record.payload).expect("decode");
assert_eq!(decoded.collection, "docs");
assert_eq!(
decoded.doc_id,
crate::engine::document::store::surrogate_to_doc_id(Surrogate::new(7))
);
}
#[test]
fn spatial_insert_appends_and_decodes() {
let dir = tempfile::tempdir().expect("tempdir");
let wal = open_wal(dir.path());
let plan = PhysicalPlan::Spatial(SpatialOp::Insert {
collection: "places".to_string(),
field: "loc".to_string(),
surrogate: Surrogate::new(9),
geometry: Geometry::point(10.0, 20.0),
provenance: None,
});
let outcome = wal_append_if_write(
&wal,
TenantId::new(1),
VShardId::new(0),
DatabaseId::DEFAULT,
&plan,
)
.expect("append");
assert!(
outcome.lsn.is_some(),
"SpatialOp::Insert must produce a durable LSN"
);
let record = last_record_of_type(&wal, nodedb_wal::record::RecordType::SpatialPut);
let decoded =
nodedb_wal::record::SpatialPutPayload::from_bytes(&record.payload).expect("decode");
assert_eq!(decoded.collection, "places");
assert_eq!(decoded.field, "loc");
let geometry: Geometry =
zerompk::from_msgpack(&decoded.geometry_bytes).expect("decode geometry");
assert_eq!(geometry, Geometry::point(10.0, 20.0));
}
#[test]
fn spatial_delete_appends_and_decodes() {
let dir = tempfile::tempdir().expect("tempdir");
let wal = open_wal(dir.path());
let plan = PhysicalPlan::Spatial(SpatialOp::Delete {
collection: "places".to_string(),
field: "loc".to_string(),
surrogate: Surrogate::new(9),
provenance: None,
});
let outcome = wal_append_if_write(
&wal,
TenantId::new(1),
VShardId::new(0),
DatabaseId::DEFAULT,
&plan,
)
.expect("append");
assert!(
outcome.lsn.is_some(),
"SpatialOp::Delete must produce a durable LSN"
);
let record = last_record_of_type(&wal, nodedb_wal::record::RecordType::SpatialDelete);
let decoded =
nodedb_wal::record::SpatialDeletePayload::from_bytes(&record.payload).expect("decode");
assert_eq!(decoded.collection, "places");
assert_eq!(decoded.field, "loc");
}
}