use crate::bridge::envelope::{PhysicalPlan, Response, Status, WriteSetEntry};
use crate::engine::document::store::surrogate_to_doc_id;
use crate::types::{DatabaseId, Lsn, TenantId, VShardId};
use crate::wal::manager::WalManager;
use nodedb_physical::physical_plan::DocumentOp;
use nodedb_types::Surrogate;
use super::document::{encode_document_delete_record, encode_document_put_record};
pub fn plan_post_apply_redo(plan: &PhysicalPlan) -> Option<String> {
if let PhysicalPlan::Document(DocumentOp::PointUpdate { collection, .. }) = plan {
Some(collection.clone())
} else if let PhysicalPlan::Document(DocumentOp::Upsert { collection, .. }) = plan {
Some(collection.clone())
} else if let PhysicalPlan::Document(DocumentOp::BulkUpdate { collection, .. }) = plan {
Some(collection.clone())
} else if let PhysicalPlan::Document(DocumentOp::BulkDelete { collection, .. }) = plan {
Some(collection.clone())
} else if let PhysicalPlan::Document(DocumentOp::BatchInsert { collection, .. }) = plan {
Some(collection.clone())
} else if let PhysicalPlan::Document(DocumentOp::Truncate { collection, .. }) = plan {
Some(collection.clone())
} else if let PhysicalPlan::Document(DocumentOp::UpdateFromJoin {
target_collection, ..
}) = plan
{
Some(target_collection.clone())
} else {
None
}
}
pub fn append_write_set_redo(
wal: &WalManager,
tenant_id: TenantId,
vshard_id: VShardId,
database_id: DatabaseId,
collection: &str,
write_set: &[WriteSetEntry],
) -> crate::Result<Option<Lsn>> {
let mut last: Option<Lsn> = None;
for entry in write_set {
let doc_id = surrogate_to_doc_id(Surrogate::new(entry.surrogate));
let lsn = if entry.is_delete {
let record = encode_document_delete_record(collection, &doc_id, entry.surrogate)?;
wal.append_delete(tenant_id, vshard_id, database_id, &record)?
} else {
let record =
encode_document_put_record(collection, &doc_id, &entry.value, entry.surrogate)?;
wal.append_put(tenant_id, vshard_id, database_id, &record)?
};
last = Some(lsn);
}
Ok(last)
}
pub fn mint_dispatch_local_redo(
wal: &WalManager,
tenant_id: TenantId,
database_id: DatabaseId,
collection: &str,
resp: &Response,
) -> crate::Result<()> {
if resp.status != Status::Ok || resp.write_set.is_empty() {
return Ok(());
}
let vshard_id = VShardId::from_collection_in_database(database_id, collection);
append_write_set_redo(
wal,
tenant_id,
vshard_id,
database_id,
collection,
&resp.write_set,
)?;
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
use nodedb_physical::physical_plan::ReturningSpec;
use nodedb_types::sync::wire::SyncProvenance;
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 point_update_is_post_apply_redo() {
let plan = PhysicalPlan::Document(DocumentOp::PointUpdate {
collection: "docs".to_string(),
document_id: "d1".to_string(),
surrogate: Surrogate::new(1),
pk_bytes: Vec::new(),
updates: Vec::new(),
returning: None::<ReturningSpec>,
});
assert_eq!(plan_post_apply_redo(&plan).as_deref(), Some("docs"));
}
#[test]
fn bulk_update_is_post_apply_redo() {
let plan = PhysicalPlan::Document(DocumentOp::BulkUpdate {
collection: "docs".to_string(),
filters: Vec::new(),
updates: Vec::new(),
returning: None::<ReturningSpec>,
ollp_predicted_surrogates: None,
ollp_predicted_edges: None,
});
assert_eq!(plan_post_apply_redo(&plan).as_deref(), Some("docs"));
}
#[test]
fn update_from_join_is_post_apply_redo() {
let plan = PhysicalPlan::Document(DocumentOp::UpdateFromJoin {
target_collection: "docs".to_string(),
source_collection: "src".to_string(),
source_alias: "s".to_string(),
target_join_col: "sku".to_string(),
source_join_col: "sku".to_string(),
updates: Vec::new(),
target_filters: Vec::new(),
returning: None::<ReturningSpec>,
resolve_only: false,
source_rows: None,
});
assert_eq!(plan_post_apply_redo(&plan).as_deref(), Some("docs"));
}
#[test]
fn bulk_delete_is_post_apply_redo() {
let plan = PhysicalPlan::Document(DocumentOp::BulkDelete {
collection: "docs".to_string(),
filters: Vec::new(),
returning: None::<ReturningSpec>,
ollp_predicted_surrogates: None,
ollp_predicted_edges: None,
});
assert_eq!(plan_post_apply_redo(&plan).as_deref(), Some("docs"));
}
#[test]
fn batch_insert_is_post_apply_redo() {
let plan = PhysicalPlan::Document(DocumentOp::BatchInsert {
collection: "docs".to_string(),
documents: Vec::new(),
surrogates: Vec::new(),
});
assert_eq!(plan_post_apply_redo(&plan).as_deref(), Some("docs"));
}
#[test]
fn point_get_is_not_post_apply_redo() {
let plan = PhysicalPlan::Document(DocumentOp::PointGet {
collection: "docs".to_string(),
document_id: "d1".to_string(),
surrogate: Surrogate::ZERO,
pk_bytes: Vec::new(),
rls_filters: Vec::new(),
system_time: nodedb_types::SystemTimeScope::Current,
valid_at_ms: None,
});
assert!(plan_post_apply_redo(&plan).is_none());
}
#[test]
fn write_set_put_appends_replayable_put_record() {
let dir = tempfile::tempdir().expect("tempdir");
let wal = open_wal(dir.path());
let entries = vec![WriteSetEntry {
surrogate: 9,
is_delete: false,
value: vec![1, 2, 3],
}];
let lsn = append_write_set_redo(
&wal,
TenantId::new(1),
VShardId::new(0),
DatabaseId::DEFAULT,
"docs",
&entries,
)
.expect("append");
assert!(
lsn.is_some(),
"a put write-set entry must append a redo LSN"
);
let record = last_record_of_type(&wal, nodedb_wal::record::RecordType::Put);
let (collection, document_id, value, _prov, surrogate) =
zerompk::from_msgpack::<(String, String, Vec<u8>, Option<SyncProvenance>, u32)>(
&record.payload,
)
.expect("decode put payload");
assert_eq!(collection, "docs");
assert_eq!(document_id, surrogate_to_doc_id(Surrogate::new(9)));
assert_eq!(value, vec![1, 2, 3]);
assert_eq!(surrogate, 9);
}
#[test]
fn write_set_delete_appends_replayable_delete_record() {
let dir = tempfile::tempdir().expect("tempdir");
let wal = open_wal(dir.path());
let entries = vec![WriteSetEntry {
surrogate: 9,
is_delete: true,
value: Vec::new(),
}];
append_write_set_redo(
&wal,
TenantId::new(1),
VShardId::new(0),
DatabaseId::DEFAULT,
"docs",
&entries,
)
.expect("append");
let record = last_record_of_type(&wal, nodedb_wal::record::RecordType::Delete);
let (collection, document_id, _prov, surrogate) =
zerompk::from_msgpack::<(String, String, Option<SyncProvenance>, u32)>(&record.payload)
.expect("decode delete payload");
assert_eq!(collection, "docs");
assert_eq!(document_id, surrogate_to_doc_id(Surrogate::new(9)));
assert_eq!(surrogate, 9);
}
#[test]
fn empty_write_set_appends_nothing() {
let dir = tempfile::tempdir().expect("tempdir");
let wal = open_wal(dir.path());
let lsn = append_write_set_redo(
&wal,
TenantId::new(1),
VShardId::new(0),
DatabaseId::DEFAULT,
"docs",
&[],
)
.expect("append");
assert!(lsn.is_none(), "empty write-set must append no record");
}
}