Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
61 changes: 56 additions & 5 deletions nodedb/src/control/crdt_admission.rs
Original file line number Diff line number Diff line change
Expand Up @@ -85,12 +85,36 @@ struct CrdtAdmissionWorkflow<'a> {
tenant_id: TenantId,
database_id: DatabaseId,
vshard_id: VShardId,
/// The collection the doors that route this work hash together with the
/// database to pick a vShard. Apply admission binds the canonical,
/// database-qualified key so the sequencer slot, the preview dispatch and
/// the raft entry all ride the vShard the planner derived; the restore path
/// binds its caller's own form, which its plan and apply are both built from.
collection: &'a str,
/// Canonical, database-qualified key the CRDT engine stores the collection
/// under. The plan under admission carries this form, so the preview that
/// fences it has to as well or the two address different documents.
engine_collection: &'a str,
timeout: Duration,
event_source: EventSource,
policy: &'a dyn CrdtPostImagePolicy,
}

/// The canonical, database-qualified key for `collection`, whatever form the
/// caller passed.
///
/// `QualifiedCollection::new` is the constructor every plan uses, so reducing
/// the request through it is what makes the request's collection comparable
/// with the plan's — and it is the string the CRDT engine is keyed by.
fn engine_key(database_id: DatabaseId, collection: &str) -> String {
nodedb_types::QualifiedCollection::new(
database_id,
&crate::control::target_identity::bare_collection_name(database_id, collection),
)
.as_str()
.to_owned()
}

/// Whether an operation changes the Loro frontier and must serialize with an
/// admission preview when executed directly on a single-node Data Plane.
pub fn changes_crdt_frontier(op: &CrdtOp) -> bool {
Expand Down Expand Up @@ -130,7 +154,11 @@ pub async fn dispatch_authorized_crdt_apply_admitted_outcome(
event_source,
policy,
} = request;
enforce_external_signing_policy(state, &authorized, collection)?;
let database_id = authorized.database_id();
// The catalog qualifies the name itself, so this lookup is the one place
// that needs the bare form back.
let bare = crate::control::target_identity::bare_collection_name(database_id, collection);
enforce_external_signing_policy(state, &authorized, &bare)?;
let task = authorized.into_physical_task();
dispatch_crdt_apply_admitted_outcome(
state,
Expand Down Expand Up @@ -191,6 +219,12 @@ pub(crate) async fn dispatch_crdt_apply_admitted_outcome(
event_source,
policy,
} = request;
// The plan carries the canonical engine key; the request may carry either
// form. Reducing the request to the canonical form is what lets a plan built
// for a non-default database match the bare name its caller typed -- and it
// is the string the engine is keyed by, so the preview below has to use it
// too or it reads a different (empty) document than the apply writes.
let key = engine_key(database_id, collection);
let (document_id, delta) = match &plan {
PhysicalPlan::Crdt(
CrdtOp::Apply {
Expand All @@ -207,7 +241,7 @@ pub(crate) async fn dispatch_crdt_apply_admitted_outcome(
expected_frontier_digest: None,
..
},
) if plan_collection.as_str() == collection => (document_id.clone(), delta.clone()),
) if plan_collection.as_str() == key.as_str() => (document_id.clone(), delta.clone()),
PhysicalPlan::Crdt(
CrdtOp::Apply {
expected_frontier_digest: Some(_),
Expand All @@ -224,13 +258,19 @@ pub(crate) async fn dispatch_crdt_apply_admitted_outcome(
});
}
};
let vshard_id = VShardId::from_collection_in_database(database_id, collection);
// Route on the canonical key, not on the caller's own form: the planner
// derives its task vShard from the database-qualified collection, so the
// sequencer slot, the preview dispatch and the raft entry have to key on
// that same string or admission fences a different vShard than the one the
// plan's own writes land on.
let vshard_id = VShardId::from_collection_in_database(database_id, &key);
let workflow = CrdtAdmissionWorkflow {
state,
tenant_id,
database_id,
vshard_id,
collection,
collection: &key,
engine_collection: &key,
timeout,
event_source,
policy,
Expand Down Expand Up @@ -290,7 +330,7 @@ async fn preview(
workflow.collection,
PhysicalPlan::Crdt(CrdtOp::PreviewApply {
collection: nodedb_types::QualifiedCollection::from_stored(
workflow.collection.to_owned(),
workflow.engine_collection.to_owned(),
),
document_id: document_id.to_owned(),
delta: delta.to_vec(),
Expand Down Expand Up @@ -422,6 +462,17 @@ pub(crate) async fn dispatch_crdt_restore_admitted(
database_id,
vshard_id,
collection,
// The restore path builds its own Apply from this same string, so the
// preview and that apply agree with each other. The string is the bare
// caller form, and nothing qualifies it on the way in: it is handed to
// `from_stored` verbatim and the tenant engine keys its collections by
// exactly that string. Every ordinary apply instead routes the
// database-qualified key (`engine_key` above), so in a non-default
// database restore addresses a different document than an apply of the
// same collection. Pre-existing and deliberately not changed here:
// canonicalizing it means changing the form the restore caller passes,
// and `CrdtOp::Apply.collection` below is rebuilt from this same string.
engine_collection: collection,
timeout,
event_source,
policy,
Expand Down
140 changes: 139 additions & 1 deletion nodedb/src/control/planner/calvin/dependent_recon.rs
Original file line number Diff line number Diff line change
Expand Up @@ -63,6 +63,11 @@ pub struct DependentReconOutcome {
/// task (`BulkUpdate`/`BulkDelete`) whose target collection has
/// `has_implicit_edges` set in the catalog, else `None`.
///
/// The returned collection is the form the plan carries — database-qualified
/// outside `DatabaseId::DEFAULT` — because that is the routing key. The catalog
/// lookup underneath reduces it to the bare name the catalog is keyed by; both
/// current call sites take only the `database_id`.
///
/// A genuine catalog READ error propagates as a typed [`crate::Error`]:
/// misrouting a delete on a real I/O fault would silently skip edge cleanup
/// (dangling edges). An ABSENT catalog (`None`) or absent collection row
Expand All @@ -83,8 +88,17 @@ pub fn plan_needs_implicit_edge_recon(
let db = dep_task.database_id;
let edge_bearing = {
let catalog = state.credentials.catalog();
// The plan carries the database-qualified collection (the router keys
// its vShard on that form), but the catalog stores collections under
// the bare name. Reading it qualified misses in every non-default
// database, so `has_implicit_edges` reads false, this gate returns
// `None`, and the OLLP/Calvin recon that cleans up mirrored edges never
// routes — the plan lowers a PK-equality UPDATE/DELETE to `Bulk*`
// precisely so this gate picks it up. Identity for
// `DatabaseId::DEFAULT`.
let bare = crate::control::target_identity::bare_collection_name(db, &coll);
catalog
.get_collection(db, tenant_id.as_u64(), &coll)?
.get_collection(db, tenant_id.as_u64(), &bare)?
.map(|c| c.has_implicit_edges)
.unwrap_or(false)
};
Expand Down Expand Up @@ -426,3 +440,127 @@ async fn dispatch_dependent_edge_recon_inner(
apply_result,
})
}

#[cfg(test)]
mod tests {
use std::sync::Arc;

use super::*;
use crate::control::security::catalog::StoredCollection;
use crate::types::VShardId;
use nodedb_physical::physical_plan::{DocumentOp, PhysicalPlan};
use nodedb_physical::physical_task::{PhysicalTask, PostSetOp};
use nodedb_types::QualifiedCollection;

/// A `SharedState` with a real on-disk catalog and no Data Plane: this test
/// only reads the catalog, so nothing else has to be live.
fn state_with_edge_bearing_collection(
dir: &tempfile::TempDir,
database_id: DatabaseId,
tenant_id: u64,
) -> Arc<SharedState> {
let wal_dir = dir.path().join("wal");
std::fs::create_dir_all(&wal_dir).unwrap();
let wal = Arc::new(crate::wal::WalManager::open_for_testing(&wal_dir).unwrap());
let (dispatcher, _) = crate::bridge::dispatch::Dispatcher::new(1, 16);
let state = SharedState::open(
crate::control::state::DataPlaneHandles {
dispatcher,
quiesce: crate::bridge::quiesce::CollectionQuiesce::new(),
array_catalog: crate::control::array_catalog::ArrayCatalog::handle(),
system_metrics: Arc::new(crate::control::metrics::SystemMetrics::new()),
},
wal,
&dir.path().join("catalog.redb"),
&crate::config::auth::AuthConfig::default(),
Default::default(),
false,
crate::data::executor::core_loop::test_governor(),
)
.unwrap();

let mut coll = StoredCollection::new(tenant_id, "edges_nd", "admin");
coll.collection_type = nodedb_types::CollectionType::document();
coll.has_implicit_edges = true;
state
.credentials
.catalog()
.put_collection(database_id, &coll)
.unwrap();
state
}

/// A `BulkUpdate` on `collection`, the shape the planner lowers a
/// PK-equality `DELETE`/`UPDATE` into so this gate picks it up.
fn bulk_update_task(database_id: DatabaseId, collection: QualifiedCollection) -> PhysicalTask {
PhysicalTask {
tenant_id: TenantId::new(1),
vshard_id: VShardId::new(0),
database_id,
plan: PhysicalPlan::Document(DocumentOp::BulkUpdate {
collection,
filters: vec![],
updates: vec![],
returning: None,
ollp_predicted_surrogates: None,
ollp_predicted_edges: None,
rls_filters: vec![],
rls_write_check: nodedb_types::RlsWriteCheck::pending_injection(),
resolved_sum_targets: Vec::new(),
declared_primary_key: None,
}),
post_set_op: PostSetOp::None,
txn_id: None,
}
}

/// The gate must fire for an edge-bearing collection in a NON-DEFAULT
/// database.
///
/// The plan carries the database-qualified collection, but the catalog is
/// keyed by the bare name, so a lookup with the qualified form misses, this
/// gate returns `None`, and the OLLP/Calvin reconnaissance that cleans up
/// mirrored edges never routes. That failure is silent — the write still
/// succeeds — so only a test that asserts the gate's own answer catches it.
#[test]
fn gate_fires_for_an_edge_bearing_collection_in_a_non_default_database() {
let dir = tempfile::tempdir().unwrap();
let database_id = DatabaseId::new(1024);
let state = state_with_edge_bearing_collection(&dir, database_id, 1);
let task = bulk_update_task(
database_id,
QualifiedCollection::new(database_id, "edges_nd"),
);

let fired = plan_needs_implicit_edge_recon(&state, &[task], TenantId::new(1)).unwrap();
assert!(
fired.is_some(),
"the gate must fire for an edge-bearing collection in a non-default database; \
a qualified catalog lookup misses and the mirrored edges leak"
);
let (collection, db) = fired.unwrap();
assert_eq!(db, database_id);
assert_eq!(
collection, "1024/edges_nd",
"the returned collection is the plan's routing key, not the bare catalog name"
);
}

/// The default database is the identity case and must keep firing.
#[test]
fn gate_fires_for_an_edge_bearing_collection_in_the_default_database() {
let dir = tempfile::tempdir().unwrap();
let state = state_with_edge_bearing_collection(&dir, DatabaseId::DEFAULT, 1);
let task = bulk_update_task(
DatabaseId::DEFAULT,
QualifiedCollection::new(DatabaseId::DEFAULT, "edges_nd"),
);

let fired = plan_needs_implicit_edge_recon(&state, &[task], TenantId::new(1)).unwrap();
assert!(
fired.is_some(),
"the default-database path must be unchanged"
);
assert_eq!(fired.unwrap().0, "edges_nd");
}
}
10 changes: 8 additions & 2 deletions nodedb/src/control/planner/implicit_edges/catalog.rs
Original file line number Diff line number Diff line change
Expand Up @@ -33,8 +33,14 @@ pub async fn mark_collection_edge_bearing(
collection: &str,
) -> crate::Result<()> {
let catalog = state.credentials.catalog();
let Some(mut coll) = catalog.get_collection(database_id, tenant_id.as_u64(), collection)?
else {
// Callers route the plan on the database-qualified collection, but the
// catalog keys collections by the bare name, so strip the qualifier back off
// here or an edge-bearing collection in a non-default database misses the
// read, the flag is never set, and implicit-edge UPDATE/DELETE cleanup is
// silently skipped, leaking stale mirrored edges. Identity for
// `DatabaseId::DEFAULT`.
let bare = crate::control::target_identity::bare_collection_name(database_id, collection);
let Some(mut coll) = catalog.get_collection(database_id, tenant_id.as_u64(), &bare)? else {
// Collection row absent — don't fail the write over flag bookkeeping.
return Ok(());
};
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -70,8 +70,14 @@ pub(in crate::control::planner::sql_plan_convert::dml) fn document_collection_wr
return Ok(WriteGates::default());
};
let catalog = credentials.catalog();
// The caller routes the plan on the database-qualified collection, but the
// catalog keys collections by the bare name, so strip the qualifier back off
// here or an INSERT into a non-default database misses its row and silently
// declares neither gate — losing CRDT convergence, and the BALANCED boundary
// with it. Identity for `DatabaseId::DEFAULT`.
let bare = crate::control::target_identity::bare_collection_name(ctx.database_id, collection);
Ok(catalog
.get_collection(ctx.database_id, ctx.tenant_id.as_u64(), collection)?
.get_collection(ctx.database_id, ctx.tenant_id.as_u64(), &bare)?
.map(|c| WriteGates {
crdt: c.crdt,
balanced: c.balanced.is_some(),
Expand Down
8 changes: 7 additions & 1 deletion nodedb/src/control/planner/sql_plan_convert/dml/crdt_gate.rs
Original file line number Diff line number Diff line change
Expand Up @@ -36,8 +36,14 @@ pub(in crate::control::planner::sql_plan_convert::dml) fn document_collection_is
return Ok(false);
};
let catalog = credentials.catalog();
// Callers route the plan on the database-qualified collection, but the
// catalog keys collections by the bare name, so strip the qualifier back off
// here or a CRDT collection in a non-default database reads as non-CRDT and
// its predicate write silently bypasses CRDT convergence. Identity for
// `DatabaseId::DEFAULT`.
let bare = crate::control::target_identity::bare_collection_name(ctx.database_id, collection);
Ok(catalog
.get_collection(ctx.database_id, ctx.tenant_id.as_u64(), collection)?
.get_collection(ctx.database_id, ctx.tenant_id.as_u64(), &bare)?
.map(|c| c.crdt)
.unwrap_or(false))
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -38,9 +38,14 @@ pub(in super::super::super) fn declared_primary_key_name(
let Some(credentials) = ctx.credentials.as_ref() else {
return Ok(None);
};
// Callers route the plan on the database-qualified collection, but the
// catalog keys collections by the bare name, so strip the qualifier back off
// here or a declared PRIMARY KEY reads as undeclared in a non-default
// database and its NOT NULL goes unenforced. Identity for `DEFAULT`.
let bare = crate::control::target_identity::bare_collection_name(ctx.database_id, collection);
credentials
.catalog()
.declared_primary_key(ctx.database_id, ctx.tenant_id.as_u64(), collection)
.declared_primary_key(ctx.database_id, ctx.tenant_id.as_u64(), &bare)
}

/// Resolve a row's document id and surrogate, refusing a NULL or omitted
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -26,8 +26,14 @@ pub(super) fn document_collection_is_edge_bearing(
return Ok(false);
};
let catalog = credentials.catalog();
// Callers route the plan on the database-qualified collection, but the
// catalog keys collections by the bare name, so strip the qualifier back off
// here or an edge-bearing collection in a non-default database reads as
// non-edge-bearing and its PK-equality UPDATE/DELETE skips the mirrored-edge
// cleanup, leaking stale edges. Identity for `DatabaseId::DEFAULT`.
let bare = crate::control::target_identity::bare_collection_name(ctx.database_id, collection);
Ok(catalog
.get_collection(ctx.database_id, ctx.tenant_id.as_u64(), collection)?
.get_collection(ctx.database_id, ctx.tenant_id.as_u64(), &bare)?
.map(|c| c.has_implicit_edges)
.unwrap_or(false))
}
Expand Down
Loading
Loading