use crate::Error;
use crate::control::server::shared::session::read_set::ReadSetEntry;
use crate::types::VShardId;
use nodedb_cluster::calvin::types::{EngineKeySet, ReadWriteSet, SortedVec, TxClass};
use nodedb_physical::physical_plan::{GraphOp, PhysicalPlan};
use nodedb_physical::physical_task::PhysicalTask;
use nodedb_types::{DatabaseId, TenantId};
use super::shared::{
collection_name_from_plan, read_set_from, surrogate_from_plan, versioned_reads_from,
};
pub fn build_dependent_tx_class(
tasks: &[PhysicalTask],
tenant_id: TenantId,
collection: &str,
predicted_surrogates: &[u32],
reads: &[ReadSetEntry],
) -> crate::Result<TxClass> {
build_dependent_tx_class_impl(
tasks,
tenant_id,
collection,
predicted_surrogates,
reads,
false,
)
}
pub fn build_single_vshard_dependent_tx_class(
tasks: &[PhysicalTask],
tenant_id: TenantId,
collection: &str,
predicted_surrogates: &[u32],
reads: &[ReadSetEntry],
) -> crate::Result<TxClass> {
build_dependent_tx_class_impl(
tasks,
tenant_id,
collection,
predicted_surrogates,
reads,
true,
)
}
fn build_dependent_tx_class_impl(
tasks: &[PhysicalTask],
tenant_id: TenantId,
collection: &str,
predicted_surrogates: &[u32],
reads: &[ReadSetEntry],
allow_single_vshard: bool,
) -> crate::Result<TxClass> {
use std::collections::BTreeMap;
let database_id = tasks
.first()
.map_or(DatabaseId::DEFAULT, |task| task.database_id);
if tasks.iter().any(|task| task.database_id != database_id)
|| reads.iter().any(|read| read.database_id != database_id)
{
return Err(Error::BadRequest {
detail: "Calvin transaction spans multiple databases".to_owned(),
});
}
let mut doc_surrogates: BTreeMap<String, Vec<u32>> = BTreeMap::new();
let mut edge_pairs: BTreeMap<String, Vec<(u32, u32)>> = BTreeMap::new();
let mut edge_homes: BTreeMap<String, Vec<u32>> = BTreeMap::new();
doc_surrogates
.entry(collection.to_owned())
.or_default()
.extend_from_slice(predicted_surrogates);
for task in tasks {
if let PhysicalPlan::Graph(
GraphOp::EdgePut {
collection: edge_coll,
src_id,
dst_id,
src_surrogate,
dst_surrogate,
..
}
| GraphOp::EdgeDelete {
collection: edge_coll,
src_id,
dst_id,
src_surrogate,
dst_surrogate,
..
},
) = &task.plan
{
edge_pairs
.entry(edge_coll.clone())
.or_default()
.push((src_surrogate.as_u32(), dst_surrogate.as_u32()));
let homes = edge_homes.entry(edge_coll.clone()).or_default();
homes.push(VShardId::from_key(src_id.as_bytes()).as_u32());
homes.push(VShardId::from_key(dst_id.as_bytes()).as_u32());
continue;
}
let coll = collection_name_from_plan(&task.plan);
if coll.is_empty() || coll == collection {
continue;
}
let surrogate = surrogate_from_plan(&task.plan);
doc_surrogates.entry(coll).or_default().push(surrogate);
}
let mut write_sets: Vec<EngineKeySet> = doc_surrogates
.into_iter()
.map(|(coll, surrogates)| EngineKeySet::Document {
collection: coll,
surrogates: SortedVec::new(surrogates),
})
.collect();
for (edge_coll, pairs) in edge_pairs {
let homes = edge_homes.remove(&edge_coll).ok_or_else(|| Error::Internal {
detail: format!(
"build_dependent_tx_class invariant violated: no edge_homes for collection {edge_coll}"
),
})?;
write_sets.push(EngineKeySet::Edge {
collection: edge_coll,
edges: SortedVec::new(pairs),
home_vshards: SortedVec::new(homes),
});
}
write_sets.sort_by(|a, b| a.collection().cmp(b.collection()));
let write_set = ReadWriteSet::new(write_sets);
let read_set = read_set_from(reads);
let plans: Vec<&PhysicalPlan> = tasks.iter().map(|t| &t.plan).collect();
let plans_bytes = zerompk::to_msgpack_vec(&plans).map_err(|e| Error::Serialization {
format: "msgpack".to_owned(),
detail: format!("failed to encode PhysicalPlan vec for Calvin dependent TxClass: {e}"),
})?;
let versioned_reads = versioned_reads_from(reads);
let result = if allow_single_vshard {
TxClass::new_single_vshard_in_database(
read_set,
write_set,
plans_bytes,
tenant_id,
database_id,
None,
versioned_reads,
)
} else {
TxClass::new_in_database(
read_set,
write_set,
plans_bytes,
tenant_id,
database_id,
None,
versioned_reads,
)
};
result.map_err(|e| Error::BadRequest {
detail: format!("invalid dependent TxClass: {e}"),
})
}
#[cfg(test)]
mod tests {
use super::*;
use crate::types::DatabaseId;
use nodedb_physical::physical_plan::DocumentOp;
fn bulk_delete_task(collection: &str) -> PhysicalTask {
PhysicalTask {
tenant_id: TenantId::new(1),
vshard_id: VShardId::new(0),
database_id: DatabaseId::DEFAULT,
plan: PhysicalPlan::Document(DocumentOp::BulkDelete {
collection: collection.to_owned(),
filters: vec![],
returning: None,
ollp_predicted_surrogates: None,
ollp_predicted_edges: None,
}),
post_set_op: nodedb_physical::physical_task::PostSetOp::None,
txn_id: None,
}
}
#[test]
fn single_collection_predicate_strict_rejects_but_single_vshard_builder_accepts() {
let tasks = vec![bulk_delete_task("users")];
let want_vshard =
VShardId::from_collection_in_database(DatabaseId::DEFAULT, "users").as_u32();
let strict = build_dependent_tx_class(&tasks, TenantId::new(1), "users", &[7, 8], &[]);
assert!(
matches!(strict, Err(crate::Error::BadRequest { .. })),
"strict dependent builder must reject single-vshard write set"
);
let tx =
build_single_vshard_dependent_tx_class(&tasks, TenantId::new(1), "users", &[7, 8], &[])
.expect("single-vshard dependent TxClass accepted");
assert_eq!(tx.participating_vshards().len(), 1);
assert_eq!(tx.participating_vshards()[0].as_u32(), want_vshard);
}
}