From b7d3ac8752a9b43649b25c2992befc7d35b373f1 Mon Sep 17 00:00:00 2001 From: EnRaiha <15997552+EnRaiha@users.noreply.github.com> Date: Sat, 26 Sep 2026 01:56:53 +0800 Subject: [PATCH 1/5] fix(control): key CRDT admission on the canonical collection `dispatch_crdt_apply_admitted_outcome` compared the plan's collection -- `QualifiedCollection::new(database_id, collection)`, so `{database_id}/{name}` on any database but the default -- against the bare name its caller typed, and `CrdtAdmissionInvalidPlan` came back as XX000 before any policy was consulted. A `CRDT MERGE` in a non-default database therefore never reached the RLS decision its own policy store already carried; in `default` the two forms coincide, which is why nothing surfaced. Reduce the request to the canonical form before comparing, and hand the same string to the preview that fences the apply: the CRDT engine is keyed by that name, so a preview built from the bare name read a different (empty) document than the apply it was fencing. Two places need the bare form and keep getting it: the catalog lookup (`get_collection` qualifies internally) and the vShard route. Routing stays on the form each caller passed, because every entry point derives its task vShard from that same string -- re-deriving it here would move work between cores on a path this change is not about. Callers disagree on the form by design: the SQL, HTTP, and sync entry points pass the name the client typed, while the native raw dispatch passes the stored, already-qualified one. De-qualification goes through `target_identity::bare_collection_name`, which strips the prefix only when it is really there, so an already-bare name and a collection whose own name contains `/` both survive untouched. --- nodedb/src/control/crdt_admission.rs | 45 ++++++++++++++++++++++++++-- 1 file changed, 42 insertions(+), 3 deletions(-) diff --git a/nodedb/src/control/crdt_admission.rs b/nodedb/src/control/crdt_admission.rs index 8dc7b0988..f287c48d2 100644 --- a/nodedb/src/control/crdt_admission.rs +++ b/nodedb/src/control/crdt_admission.rs @@ -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 { @@ -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, @@ -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 { @@ -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(_), @@ -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, @@ -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(), @@ -435,6 +471,9 @@ 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 keeps whatever form that caller used. + engine_collection: collection, timeout, event_source, policy, From 8b0b6f915b1bc784f14be0f06457058f514757f8 Mon Sep 17 00:00:00 2001 From: EnRaiha <15997552+EnRaiha@users.noreply.github.com> Date: Wed, 23 Sep 2026 01:05:00 +0800 Subject: [PATCH 2/5] ddl,http: classify the CRDT MERGE refusal and code the HTTP stream errors MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Two surfaces still flattened a classified failure while the routed pgwire paths already carried its class. - CRDT MERGE: the admission and apply failures now map through `error_to_sqlstate`, so a policy denial — `ExternalCrdtPostImagePolicy::deny` returns `RejectedAuthz` — reaches the client as `42501` (INSUFFICIENT_PRIVILEGE) instead of `XX000`. The "authorization returned no capability" site keeps `XX000`: that one is an internal invariant by decision, not a class a client can act on. - HTTP stream: the in-band error lines carry the numeric NodeDB code (shape and malformed-batch failures) and the status the gateway map already computed. Consumer trace for the class change: the two wire assertions that pinned the placeholder move with it — `nodedb/tests/wire/cases/crdt_write_rls_database_scope.rs` asserted `XX000` and now asserts `42501`, and the comment above the first one no longer describes the old wrapper. No other test, doc, or path keys on the CRDT merge SQLSTATE (`grep XX000` over `tests/wire/cases` finds no other crdt or merge site). The pgwire stream and DDL-dispatch sites are not touched here; the HTTP shaping surface above is the whole second half. Also refreshes the `response_shape/schema.rs` module comment, which still claimed nothing consumes the module — the session caches the schema with the physical tasks, and shaping receives it as `projection`. Verification: `cargo nextest run -p nodedb --test wire --all-features --cargo-profile ci --profile ci -E 'test(~crdt_write_rls_database_scope)'` fails on the pre-change tree (the denial surfaces as `XX000`) and passes with this change (`42501`). --- .../server/http/routes/query_stream.rs | 22 +++++++++++++++---- .../control/server/response_shape/schema.rs | 8 +++---- 2 files changed, 22 insertions(+), 8 deletions(-) diff --git a/nodedb/src/control/server/http/routes/query_stream.rs b/nodedb/src/control/server/http/routes/query_stream.rs index 9ca180aab..bd20dc10e 100644 --- a/nodedb/src/control/server/http/routes/query_stream.rs +++ b/nodedb/src/control/server/http/routes/query_stream.rs @@ -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; } @@ -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; @@ -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; } diff --git a/nodedb/src/control/server/response_shape/schema.rs b/nodedb/src/control/server/response_shape/schema.rs index 62e6bdb8c..ec0d5b963 100644 --- a/nodedb/src/control/server/response_shape/schema.rs +++ b/nodedb/src/control/server/response_shape/schema.rs @@ -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. /// From 684846ec251e4f4178dfe5d8e3a622473bb7a547 Mon Sep 17 00:00:00 2001 From: EnRaiha <15997552+EnRaiha@users.noreply.github.com> Date: Sat, 26 Sep 2026 03:40:43 +0800 Subject: [PATCH 3/5] fix(control,ddl): de-qualify the catalog gates and classify CRDT op errors The catalog gate reads in `crdt_gate`, `update_delete::shared` and `implicit_edges::catalog` still passed the collection the caller routed on, which is database-qualified, while the catalog keys collections by the bare name. In a non-default database every one of them missed, and each miss silently answered "no": a CRDT collection read as plain, so a predicate UPDATE/DELETE and an `ON CONFLICT DO UPDATE` skipped the refusal that keeps CRDT convergence authoritative; an edge-bearing collection read as plain, so a primary-key-equality UPDATE/DELETE skipped the mirrored-edge cleanup; and the edge-bearing marker was never set, so that cleanup could not find the collection at all. Reduce all three through `target_identity::bare_collection_name`, as the sibling gate reads already do; it strips the prefix only when it is really there and is identity for `DatabaseId::DEFAULT`. `crdt_state` and `crdt_apply` flattened every classified failure to `XX000` at three sites. Map them through `error_to_sqlstate`, so the policy denial an `ExternalCrdtPostImagePolicy::deny` produces reaches the client as `42501`, the code `CRDT MERGE` already reports, instead of a class a client cannot act on. Covered by four non-default-database variants mirroring the existing default-database gate tests (CRDT predicate UPDATE, CRDT predicate DELETE, CRDT upsert, and the reserved-edge-field expression UPDATE) and a wire assertion that `crdt_apply` reports `42501` for a post-image its policy forbids. --- .../control/planner/implicit_edges/catalog.rs | 10 ++- .../cases/crdt_write_rls_database_scope.rs | 69 +++++++++++++++++ .../cases/engine_surface_crdt_document.rs | 77 +++++++++++++++++++ .../tests/wire/cases/engine_surface_graph.rs | 21 +++++ 4 files changed, 175 insertions(+), 2 deletions(-) diff --git a/nodedb/src/control/planner/implicit_edges/catalog.rs b/nodedb/src/control/planner/implicit_edges/catalog.rs index 4691149a2..3e0fd1a99 100644 --- a/nodedb/src/control/planner/implicit_edges/catalog.rs +++ b/nodedb/src/control/planner/implicit_edges/catalog.rs @@ -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(()); }; diff --git a/nodedb/tests/wire/cases/crdt_write_rls_database_scope.rs b/nodedb/tests/wire/cases/crdt_write_rls_database_scope.rs index 8b755a49f..04335fc13 100644 --- a/nodedb/tests/wire/cases/crdt_write_rls_database_scope.rs +++ b/nodedb/tests/wire/cases/crdt_write_rls_database_scope.rs @@ -206,3 +206,72 @@ async fn crdt_merge_in_default_database_is_still_rls_enforced() { document's stored state untouched" ); } + +/// `SELECT crdt_apply(...)` must report an RLS write denial as `42501`, the +/// code `CRDT MERGE` already reports, rather than flattening every admission +/// failure to `XX000`. A client cannot act on a denial it cannot recognise. +/// +/// Runs in `default` on purpose: the flatten is database-independent, and the +/// engine key there is the bare collection name a wire test can construct. +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn crdt_apply_in_default_database_reports_rls_denial_as_42501() { + const COLL: &str = "crdt_rls_apply_notes"; + const DOC: &str = "doc1"; + + let server = TestServer::start().await; + let user = "crdt_rls_apply_default_user"; + + query_ok( + &server, + &format!("CREATE TABLE {COLL} (id TEXT PRIMARY KEY, title TEXT) WITH (crdt='true')"), + ) + .await; + create_scoped_user(&server, user, "default").await; + + // A predicate no real row ever satisfies: `id` is the primary key, so it + // holds the document's own id, never this sentinel. Every CRDT write on + // this collection is therefore denied. + query_ok( + &server, + &format!( + "CREATE RLS POLICY {COLL}_block ON {COLL} FOR WRITE \ + USING (id = 'sentinel_id_no_row_ever_has')" + ), + ) + .await; + + // A REAL Loro delta, shaped exactly as `CrdtState` models the document: + // collection = root map, row = `insert_container`, fields on the row map. + // A placeholder payload is refused by the preview for an unrelated reason, + // so the denial below must come from a genuine apply. + let delta_hex = { + let doc = loro::LoroDoc::new(); + let coll = doc.get_map(COLL); + let row = coll + .insert_container(DOC, loro::LoroMap::new()) + .expect("row container"); + row.insert("title", "t1").expect("field"); + doc.commit(); + let delta = doc + .export(loro::ExportMode::Snapshot) + .expect("export loro snapshot"); + hex::encode(delta) + }; + + let result = try_exec_as( + &server, + user, + "default", + &format!("SELECT crdt_apply('{COLL}', '{DOC}', '{delta_hex}')"), + ) + .await; + + let sqlstate = result.expect_err( + "a write policy on a CRDT collection must reject a crdt_apply whose \ + post-image its predicate forbids", + ); + assert_eq!( + sqlstate, "42501", + "expected the RLS denial's SQLSTATE, got: {sqlstate}" + ); +} diff --git a/nodedb/tests/wire/cases/engine_surface_crdt_document.rs b/nodedb/tests/wire/cases/engine_surface_crdt_document.rs index b031ef700..2fac427ca 100644 --- a/nodedb/tests/wire/cases/engine_surface_crdt_document.rs +++ b/nodedb/tests/wire/cases/engine_surface_crdt_document.rs @@ -292,3 +292,80 @@ async fn crdt_delete_returning_projects_deleted_row() { "DELETE ... RETURNING must still remove the row; got {remaining:?}" ); } + +/// The CRDT gate is read from a catalog keyed by the BARE collection name, so a +/// non-default database must refuse a predicate UPDATE exactly as `default` +/// does — mirrors `predicate_update_on_crdt_rejected`. +#[tokio::test] +async fn predicate_update_on_crdt_rejected_in_non_default_database() { + let (srv, _db) = TestServer::with_database("crdt_pred_upd_scope").await; + srv.exec( + "CREATE TABLE crdt_notes_nd (id TEXT PRIMARY KEY, title TEXT, body TEXT) \ + WITH (crdt='true')", + ) + .await + .unwrap(); + + srv.exec("INSERT INTO crdt_notes_nd (id, title, body) VALUES ('a', 't1', 'b1')") + .await + .unwrap(); + + srv.expect_error( + "UPDATE crdt_notes_nd SET title='x' WHERE title='t1'", + "predicate (non-primary-key) UPDATE on CRDT collection", + ) + .await; +} + +/// The DELETE side of the same gate — mirrors +/// `predicate_delete_on_crdt_rejected`. +#[tokio::test] +async fn predicate_delete_on_crdt_rejected_in_non_default_database() { + let (srv, _db) = TestServer::with_database("crdt_pred_del_scope").await; + srv.exec( + "CREATE TABLE crdt_notes_nd (id TEXT PRIMARY KEY, title TEXT, body TEXT) \ + WITH (crdt='true')", + ) + .await + .unwrap(); + + srv.exec("INSERT INTO crdt_notes_nd (id, title, body) VALUES ('a', 't1', 'b1')") + .await + .unwrap(); + + srv.expect_error( + "DELETE FROM crdt_notes_nd WHERE title='t1'", + "predicate (non-primary-key) DELETE on CRDT collection", + ) + .await; +} + +/// An explicit `ON CONFLICT DO UPDATE SET` on a CRDT collection is refused — +/// CRDT convergence IS the LWW full replace, so the caller's merge clause has +/// nowhere to run. The refusal rides the same catalog gate, so a non-default +/// database must refuse it too. +/// +/// `INSERT ... ON CONFLICT DO UPDATE SET` is the form that reaches the upsert +/// converter with the clause attached; the `UPSERT INTO` form is consumed and +/// rebuilt by the protocol-neutral collection DML parser before planning. +#[tokio::test] +async fn upsert_on_conflict_do_update_on_crdt_rejected_in_non_default_database() { + let (srv, _db) = TestServer::with_database("crdt_upsert_scope").await; + srv.exec( + "CREATE TABLE crdt_notes_nd (id TEXT PRIMARY KEY, title TEXT, body TEXT) \ + WITH (crdt='true')", + ) + .await + .unwrap(); + + srv.exec("INSERT INTO crdt_notes_nd (id, title, body) VALUES ('a', 't1', 'b1')") + .await + .unwrap(); + + srv.expect_error( + "INSERT INTO crdt_notes_nd (id, title, body) VALUES ('a', 't9', 'b9') \ + ON CONFLICT (id) DO UPDATE SET title = 't9'", + "UPSERT with ON CONFLICT DO UPDATE on CRDT collection", + ) + .await; +} diff --git a/nodedb/tests/wire/cases/engine_surface_graph.rs b/nodedb/tests/wire/cases/engine_surface_graph.rs index a4481e170..464029ebb 100644 --- a/nodedb/tests/wire/cases/engine_surface_graph.rs +++ b/nodedb/tests/wire/cases/engine_surface_graph.rs @@ -118,3 +118,24 @@ async fn engine_graph_flag_rejected_in_with_clause() { "expected graph-rejection error, got: {err}" ); } + +/// The edge-bearing gate is read from a catalog keyed by the BARE collection +/// name, so a non-default database must still refuse an expression update to a +/// reserved edge field — the mirrored edge could not be reconciled against it. +#[tokio::test] +async fn edge_field_expression_update_rejected_in_non_default_database() { + let (srv, _db) = TestServer::with_database("graph_edge_scope").await; + srv.exec("CREATE COLLECTION graph_edges_nd WITH (engine='document_schemaless')") + .await + .unwrap(); + + srv.exec("INSERT INTO graph_edges_nd { id: 'e1', _from: 'alice', _to: 'bob', _type: 'knows' }") + .await + .unwrap(); + + srv.expect_error( + "UPDATE graph_edges_nd SET _from = _to WHERE id = 'e1'", + "expression updates to reserved edge fields", + ) + .await; +} From d34722cd4ca0ffcb96ff94afeafb0da1b957c276 Mon Sep 17 00:00:00 2001 From: EnRaiha <15997552+EnRaiha@users.noreply.github.com> Date: Sun, 27 Sep 2026 13:40:33 +0800 Subject: [PATCH 4/5] control: read the edge-recon catalog with the bare collection name MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit `plan_needs_implicit_edge_recon` fed the database-qualified collection from the plan straight to a catalog keyed by the bare name, so in every non-default database the read missed, `has_implicit_edges` read false, this gate returned `None`, and the OLLP/Calvin dependent-edge reconnaissance never routed. The mirrored edges of a PK-equality UPDATE/DELETE were therefore never cleaned up — and the plan lowers exactly those writes to `Bulk*` so this gate picks them up, so the lowering had no effect outside the default database. This is the same defect the sibling planner gates had; the bare name is what the catalog stores. The returned collection still comes from the plan, because that is the routing key, and both call sites take only the `database_id`. The restore path's comment claimed its bare collection form was intended. It is not: the registry qualifies it before the engine is keyed, so restore still addresses a different document than an ordinary apply in a non-default database. That is pre-existing and left alone, but the comment now says so instead of asserting the opposite. --- nodedb/src/control/crdt_admission.rs | 9 ++++++++- .../control/planner/calvin/dependent_recon.rs | 16 +++++++++++++++- 2 files changed, 23 insertions(+), 2 deletions(-) diff --git a/nodedb/src/control/crdt_admission.rs b/nodedb/src/control/crdt_admission.rs index f287c48d2..7dadbdc28 100644 --- a/nodedb/src/control/crdt_admission.rs +++ b/nodedb/src/control/crdt_admission.rs @@ -472,7 +472,14 @@ pub(crate) async fn dispatch_crdt_restore_admitted( vshard_id, collection, // The restore path builds its own Apply from this same string, so the - // preview keeps whatever form that caller used. + // preview and that apply agree with each other. The string is the bare + // caller form, which the collection registry qualifies before the engine + // is keyed, so in a non-default database this path addresses a different + // document than every ordinary apply, which routes the database-qualified + // key (`engine_key` above). 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, diff --git a/nodedb/src/control/planner/calvin/dependent_recon.rs b/nodedb/src/control/planner/calvin/dependent_recon.rs index 4ff3ae86c..6d7bb8e58 100644 --- a/nodedb/src/control/planner/calvin/dependent_recon.rs +++ b/nodedb/src/control/planner/calvin/dependent_recon.rs @@ -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 @@ -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) }; From acd213880bb5e568db22ad305ee4e2ed409bdd7a Mon Sep 17 00:00:00 2001 From: EnRaiha <15997552+EnRaiha@users.noreply.github.com> Date: Sun, 27 Sep 2026 14:41:41 +0800 Subject: [PATCH 5/5] control: cover the edge-recon gate in a non-default database The fix had no test: `plan_needs_implicit_edge_recon` had zero call sites anywhere in the suite, and the existing non-default-database graph test sends an expression update that the planner rejects at plan time, so it never reaches this gate. Reverting the fix left every test green. Two tests now drive the gate directly: seed a catalog with an edge-bearing collection in a non-default database, hand the gate a `BulkUpdate` carrying the database-qualified collection the planner builds, and assert it fires and returns that qualified key. The second covers `DatabaseId::DEFAULT` as the identity case. Reverting the bare-name lookup fails the first and leaves the second passing, which is the red arm this needed. The restore-path comment also claimed the collection registry qualifies its bare string before the engine is keyed. It does not: the string reaches `from_stored` and the tenant engine's collection map verbatim, which is what makes restore disagree with an ordinary apply. The comment now says that. --- nodedb/src/control/crdt_admission.rs | 15 ++- .../control/planner/calvin/dependent_recon.rs | 124 ++++++++++++++++++ 2 files changed, 132 insertions(+), 7 deletions(-) diff --git a/nodedb/src/control/crdt_admission.rs b/nodedb/src/control/crdt_admission.rs index 7dadbdc28..1bc7ef6fa 100644 --- a/nodedb/src/control/crdt_admission.rs +++ b/nodedb/src/control/crdt_admission.rs @@ -473,13 +473,14 @@ pub(crate) async fn dispatch_crdt_restore_admitted( 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, which the collection registry qualifies before the engine - // is keyed, so in a non-default database this path addresses a different - // document than every ordinary apply, which routes the database-qualified - // key (`engine_key` above). 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. + // 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, diff --git a/nodedb/src/control/planner/calvin/dependent_recon.rs b/nodedb/src/control/planner/calvin/dependent_recon.rs index 6d7bb8e58..958ad7228 100644 --- a/nodedb/src/control/planner/calvin/dependent_recon.rs +++ b/nodedb/src/control/planner/calvin/dependent_recon.rs @@ -460,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 { + 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"); + } +}