Skip to content
Closed
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
53 changes: 50 additions & 3 deletions nodedb/src/control/crdt_admission.rs
Original file line number Diff line number Diff line change
Expand Up @@ -91,12 +91,33 @@ struct CrdtAdmissionWorkflow<'a> {
tenant_id: TenantId,
database_id: DatabaseId,
vshard_id: VShardId,
/// Bare collection name. Both doors that route this work hash it together
/// with the database, so it must not carry the qualified prefix here.
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 @@ -136,7 +157,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 @@ -201,6 +226,16 @@ 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.
//
// Routing is deliberately left on the caller's own form: each entry point
// derives its task vShard from the string it passes here, so re-deriving it
// would move work between cores on a path this change is not about.
let key = engine_key(database_id, collection);
let (document_id, delta) = match &plan {
PhysicalPlan::Crdt(
CrdtOp::Apply {
Expand All @@ -217,7 +252,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 @@ -243,6 +278,7 @@ pub(crate) async fn dispatch_crdt_apply_admitted_outcome(
database_id,
vshard_id,
collection,
engine_collection: &key,
timeout,
event_source,
policy,
Expand Down Expand Up @@ -301,7 +337,7 @@ async fn preview(
nodedb_types::CollectionKey::from_qualified_str(workflow.database_id, 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 @@ -435,6 +471,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 @@ -446,3 +460,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
22 changes: 18 additions & 4 deletions nodedb/src/control/server/http/routes/query_stream.rs
Original file line number Diff line number Diff line change
Expand Up @@ -169,8 +169,11 @@ pub(super) fn ndjson_body_stream(
None => break,
Some(Ok(b)) => b,
Some(Err(e)) => {
let (_status, msg) = GatewayErrorMap::to_http(&e);
let line = format!("{}\n", serde_json::json!({ "error": msg }));
let (status, msg) = GatewayErrorMap::to_http(&e);
let line = format!(
"{}\n",
serde_json::json!({ "error": msg, "status": status })
);
yield Ok(Bytes::from(line));
return;
}
Expand All @@ -182,9 +185,13 @@ pub(super) fn ndjson_body_stream(
// A malformed batch payload is surfaced as an in-band error
// line (matching the mid-stream dispatch-error path above)
// rather than silently dropping the batch.
let classified = crate::error_classify::classify(&e);
let line = format!(
"{}\n",
serde_json::json!({ "error": format!("malformed response batch: {e}") })
serde_json::json!({
"error": format!("malformed response batch: {}", classified.message()),
"code": classified.code().0,
})
);
yield Ok(Bytes::from(line));
return;
Expand All @@ -206,7 +213,14 @@ pub(super) fn ndjson_body_stream(
Err(e) => {
// In-band error line, matching the malformed-batch path
// above: the HTTP body itself never errors.
let line = format!("{}\n", serde_json::json!({ "error": format!("{e}") }));
let classified = crate::error_classify::classify(&e);
let line = format!(
"{}\n",
serde_json::json!({
"error": classified.message(),
"code": classified.code().0,
})
);
yield Ok(Bytes::from(line));
return;
}
Expand Down
8 changes: 4 additions & 4 deletions nodedb/src/control/server/response_shape/schema.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3,10 +3,10 @@
//! Planner-authoritative output schema types, plus the type mapping from the
//! planner's `SqlDataType` to the response shaper's wire-facing `DdlColType`.
//!
//! Nothing in this module is consumed by existing call sites yet; it is a
//! purely additive foundation for later threading the planner's resolved
//! output schema into response shaping (replacing the SQL-string re-parse
//! path).
//! The planner derives the schema from the compiled plan and the catalog, and
//! the session caches it with the physical tasks (`session/plan_cache.rs`).
//! Shaping receives it as `MaterializedShapeRequest::projection`, which drives
//! projection and the Control-Plane computed columns.

/// One output column of a resolved query, as known by the planner.
///
Expand Down
Loading
Loading