#![deny(clippy::wildcard_enum_match_arm)]
use nodedb_physical::physical_plan::CrdtOp;
use crate::types::{DatabaseId, Lsn, TenantId, VShardId};
use crate::wal::manager::WalManager;
pub(super) fn wal_append_crdt_op(
wal: &WalManager,
tenant_id: TenantId,
vshard_id: VShardId,
database_id: DatabaseId,
op: &CrdtOp,
) -> crate::Result<Option<Lsn>> {
let appended = match op {
CrdtOp::Apply {
collection,
delta,
provenance,
..
} => {
let payload = crate::wal::CrdtDeltaWalPayload {
bytes: delta.clone(),
collection: Some(collection.clone()),
provenance: provenance.clone(),
};
let crdt_payload =
zerompk::to_msgpack_vec(&payload).map_err(|e| crate::Error::Serialization {
format: "msgpack".into(),
detail: format!("wal crdt delta: {e}"),
})?;
Some(wal.append_crdt_delta(tenant_id, vshard_id, database_id, &crdt_payload)?)
}
CrdtOp::ImportSnapshot {
collection, bytes, ..
} => {
let payload = crate::wal::CrdtDeltaWalPayload {
bytes: bytes.clone(),
collection: Some(collection.clone()),
provenance: None,
};
let crdt_payload =
zerompk::to_msgpack_vec(&payload).map_err(|e| crate::Error::Serialization {
format: "msgpack".into(),
detail: format!("wal crdt snapshot import: {e}"),
})?;
Some(wal.append_crdt_delta(tenant_id, vshard_id, database_id, &crdt_payload)?)
}
CrdtOp::ListInsert {
collection,
document_id,
list_path,
index,
fields_json,
surrogate: _,
} => {
let payload = crate::wal::CrdtListOpWalRecord::Insert {
collection: collection.clone(),
document_id: document_id.clone(),
list_path: list_path.clone(),
index: *index as u64,
fields_json: fields_json.clone(),
};
let bytes = encode_crdt_list_op_payload(payload)?;
Some(wal.append_crdt_list_op(tenant_id, vshard_id, database_id, &bytes)?)
}
CrdtOp::ListDelete {
collection,
document_id,
list_path,
index,
surrogate: _,
} => {
let payload = crate::wal::CrdtListOpWalRecord::Delete {
collection: collection.clone(),
document_id: document_id.clone(),
list_path: list_path.clone(),
index: *index as u64,
};
let bytes = encode_crdt_list_op_payload(payload)?;
Some(wal.append_crdt_list_op(tenant_id, vshard_id, database_id, &bytes)?)
}
CrdtOp::ListMove {
collection,
document_id,
list_path,
from_index,
to_index,
surrogate: _,
} => {
let payload = crate::wal::CrdtListOpWalRecord::Move {
collection: collection.clone(),
document_id: document_id.clone(),
list_path: list_path.clone(),
from_index: *from_index as u64,
to_index: *to_index as u64,
};
let bytes = encode_crdt_list_op_payload(payload)?;
Some(wal.append_crdt_list_op(tenant_id, vshard_id, database_id, &bytes)?)
}
CrdtOp::DocUpsert {
collection,
document_id,
fields_json,
surrogate,
partial,
returning: _,
} => {
let payload = crate::wal::CrdtDocOpWalRecord::Upsert {
collection: collection.clone(),
document_id: document_id.clone(),
surrogate: surrogate.as_u32(),
fields_json: fields_json.clone(),
partial: *partial,
};
let bytes = encode_crdt_doc_op_payload(payload)?;
Some(wal.append_crdt_doc_op(tenant_id, vshard_id, database_id, &bytes)?)
}
CrdtOp::DocDelete {
collection,
document_id,
surrogate,
returning: _,
} => {
let payload = crate::wal::CrdtDocOpWalRecord::Delete {
collection: collection.clone(),
document_id: document_id.clone(),
surrogate: surrogate.as_u32(),
};
let bytes = encode_crdt_doc_op_payload(payload)?;
Some(wal.append_crdt_doc_op(tenant_id, vshard_id, database_id, &bytes)?)
}
CrdtOp::Read { .. }
| CrdtOp::ReadConstraints { .. }
| CrdtOp::GetPolicy { .. }
| CrdtOp::SetPolicy { .. }
| CrdtOp::ReadAtVersion { .. }
| CrdtOp::GetVersionVector { .. }
| CrdtOp::ExportDelta { .. }
| CrdtOp::CompactAtVersion { .. } => None,
CrdtOp::SetConstraints { .. } | CrdtOp::DropConstraints { .. } => None,
CrdtOp::RestoreToVersion { .. } => None,
};
Ok(appended)
}
fn encode_crdt_list_op_payload(payload: crate::wal::CrdtListOpWalRecord) -> crate::Result<Vec<u8>> {
zerompk::to_msgpack_vec(&payload).map_err(|e| crate::Error::Serialization {
format: "msgpack".into(),
detail: format!("wal crdt list op: {e}"),
})
}
fn encode_crdt_doc_op_payload(payload: crate::wal::CrdtDocOpWalRecord) -> crate::Result<Vec<u8>> {
zerompk::to_msgpack_vec(&payload).map_err(|e| crate::Error::Serialization {
format: "msgpack".into(),
detail: format!("wal crdt doc op: {e}"),
})
}
#[cfg(test)]
mod tests {
use super::*;
use nodedb_physical::physical_plan::PhysicalPlan;
use nodedb_types::Surrogate;
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 apply_appends_crdt_delta_record() {
let dir = tempfile::tempdir().expect("tempdir");
let wal = open_wal(dir.path());
let plan = PhysicalPlan::Crdt(CrdtOp::Apply {
collection: "docs".to_string(),
document_id: "d1".to_string(),
delta: vec![9, 9, 9],
peer_id: 1,
mutation_id: 1,
surrogate: Surrogate::new(3),
provenance: None,
constraint_version_required: 0,
});
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(), "Apply must produce a durable LSN");
assert!(has_record_of_type(
&wal,
nodedb_wal::record::RecordType::CrdtDelta
));
}
#[test]
fn list_insert_appends_crdt_list_op_record() {
let dir = tempfile::tempdir().expect("tempdir");
let wal = open_wal(dir.path());
let plan = PhysicalPlan::Crdt(CrdtOp::ListInsert {
collection: "docs".to_string(),
document_id: "d1".to_string(),
list_path: "blocks".to_string(),
index: 0,
fields_json: "{}".to_string(),
surrogate: Surrogate::new(3),
});
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(),
"ListInsert must produce a durable LSN"
);
assert!(has_record_of_type(
&wal,
nodedb_wal::record::RecordType::CrdtListOp
));
}
#[test]
fn doc_upsert_appends_crdt_doc_op_record() {
let dir = tempfile::tempdir().expect("tempdir");
let wal = open_wal(dir.path());
let plan = PhysicalPlan::Crdt(CrdtOp::DocUpsert {
collection: "docs".to_string(),
document_id: "d1".to_string(),
fields_json: r#"{"a":1}"#.to_string(),
surrogate: Surrogate::new(3),
partial: false,
returning: 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(),
"DocUpsert must produce a durable LSN"
);
assert!(has_record_of_type(
&wal,
nodedb_wal::record::RecordType::CrdtDocOp
));
}
#[test]
fn read_op_appends_nothing() {
let dir = tempfile::tempdir().expect("tempdir");
let wal = open_wal(dir.path());
let plan = PhysicalPlan::Crdt(CrdtOp::Read {
collection: "docs".to_string(),
document_id: "d1".to_string(),
});
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(), "read op must produce no durable LSN");
}
}