diff --git a/nodedb/src/control/crdt_admission.rs b/nodedb/src/control/crdt_admission.rs index 726333486..86f53d377 100644 --- a/nodedb/src/control/crdt_admission.rs +++ b/nodedb/src/control/crdt_admission.rs @@ -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 { @@ -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, @@ -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 { @@ -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(_), @@ -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, @@ -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(), @@ -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, diff --git a/nodedb/src/control/planner/calvin/dependent_recon.rs b/nodedb/src/control/planner/calvin/dependent_recon.rs index e29fa429d..08a1e90da 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) }; @@ -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 { + 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"); + } +} 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/src/control/planner/sql_plan_convert/dml/balanced_gate.rs b/nodedb/src/control/planner/sql_plan_convert/dml/balanced_gate.rs index c0284b647..12c488034 100644 --- a/nodedb/src/control/planner/sql_plan_convert/dml/balanced_gate.rs +++ b/nodedb/src/control/planner/sql_plan_convert/dml/balanced_gate.rs @@ -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(), diff --git a/nodedb/src/control/planner/sql_plan_convert/dml/crdt_gate.rs b/nodedb/src/control/planner/sql_plan_convert/dml/crdt_gate.rs index 9a6287856..f39153e12 100644 --- a/nodedb/src/control/planner/sql_plan_convert/dml/crdt_gate.rs +++ b/nodedb/src/control/planner/sql_plan_convert/dml/crdt_gate.rs @@ -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)) } diff --git a/nodedb/src/control/planner/sql_plan_convert/dml/insert/identity.rs b/nodedb/src/control/planner/sql_plan_convert/dml/insert/identity.rs index 3d89fdbe0..de5c321dd 100644 --- a/nodedb/src/control/planner/sql_plan_convert/dml/insert/identity.rs +++ b/nodedb/src/control/planner/sql_plan_convert/dml/insert/identity.rs @@ -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 diff --git a/nodedb/src/control/planner/sql_plan_convert/dml/update_delete/shared.rs b/nodedb/src/control/planner/sql_plan_convert/dml/update_delete/shared.rs index d077830a8..3da885547 100644 --- a/nodedb/src/control/planner/sql_plan_convert/dml/update_delete/shared.rs +++ b/nodedb/src/control/planner/sql_plan_convert/dml/update_delete/shared.rs @@ -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)) } 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. /// diff --git a/nodedb/src/control/server/shared/ddl/neutral/crdt_ops.rs b/nodedb/src/control/server/shared/ddl/neutral/crdt_ops.rs index 832adb071..8981db53e 100644 --- a/nodedb/src/control/server/shared/ddl/neutral/crdt_ops.rs +++ b/nodedb/src/control/server/shared/ddl/neutral/crdt_ops.rs @@ -17,6 +17,7 @@ use crate::control::crdt_post_image_policy::ExternalCrdtPostImagePolicy; use crate::control::planner::sql_plan_convert::convert::db_qualified; use crate::control::security::audit::ArcAuditEmitter; use crate::control::security::identity::{AuthenticatedIdentity, Permission}; +use crate::control::server::pgwire::types::error_map::error_to_sqlstate; use crate::control::server::response_shape::types::ShapedRows; use crate::control::server::shared::authorization::{authorize_collection, authorize_task_set}; use crate::control::server::shared::ddl::sql_parse::hex_decode; @@ -96,7 +97,10 @@ pub async fn crdt_state( }, ) .await - .map_err(|e| DdlError::new("XX000", e.to_string()))?; + .map_err(|e| { + let (_, sqlstate, message) = error_to_sqlstate(&e); + DdlError::new(sqlstate, message) + })?; let columns = vec!["crdt_state".to_string()]; @@ -163,7 +167,10 @@ pub async fn crdt_apply( let surrogate = state .surrogate_assigner .assign(database_id, tenant_id, collection, document_id.as_bytes()) - .map_err(|e| DdlError::new("XX000", e.to_string()))?; + .map_err(|e| { + let (_, sqlstate, message) = error_to_sqlstate(&e); + DdlError::new(sqlstate, message) + })?; let plan = PhysicalPlan::Crdt(CrdtOp::Apply { collection: nodedb_types::QualifiedCollection::new(database_id, collection), @@ -226,7 +233,10 @@ pub async fn crdt_apply( }, ) .await - .map_err(|e| DdlError::new("XX000", e.to_string()))?; + .map_err(|e| { + let (_, sqlstate, message) = error_to_sqlstate(&e); + DdlError::new(sqlstate, message) + })?; let columns = vec!["result".to_string()]; let mut row = Map::new(); diff --git a/nodedb/src/control/server/shared/ddl/neutral/dsl/crdt_merge.rs b/nodedb/src/control/server/shared/ddl/neutral/dsl/crdt_merge.rs index 6e2790798..5a0c395d3 100644 --- a/nodedb/src/control/server/shared/ddl/neutral/dsl/crdt_merge.rs +++ b/nodedb/src/control/server/shared/ddl/neutral/dsl/crdt_merge.rs @@ -8,6 +8,7 @@ use crate::bridge::envelope::PhysicalPlan; use crate::control::crdt_post_image_policy::ExternalCrdtPostImagePolicy; use crate::control::security::audit::ArcAuditEmitter; use crate::control::security::identity::{AuthenticatedIdentity, Permission}; +use crate::control::server::pgwire::types::error_map::error_to_sqlstate; use crate::control::server::shared::authorization::{authorize_collection, authorize_task_set}; use crate::control::state::SharedState; use crate::types::DatabaseId; @@ -87,7 +88,10 @@ pub async fn crdt_merge( }, ) .await - .map_err(|e| ddl_err("XX000", e.to_string()))?; + .map_err(|e| { + let (_, sqlstate, message) = error_to_sqlstate(&e); + ddl_err(sqlstate, message) + })?; if source_bytes.is_empty() { return Err(ddl_err( "02000", @@ -98,7 +102,10 @@ pub async fn crdt_merge( let target_surrogate = state .surrogate_assigner .assign(database_id, tenant_id, collection, target_id.as_bytes()) - .map_err(|e| ddl_err("XX000", e.to_string()))?; + .map_err(|e| { + let (_, sqlstate, message) = error_to_sqlstate(&e); + ddl_err(sqlstate, message) + })?; let apply_plan = PhysicalPlan::Crdt(CrdtOp::Apply { collection: nodedb_types::QualifiedCollection::new(database_id, collection), @@ -163,7 +170,10 @@ pub async fn crdt_merge( }, ) .await - .map_err(|e| ddl_err("XX000", e.to_string()))?; + .map_err(|e| { + let (_, sqlstate, message) = error_to_sqlstate(&e); + ddl_err(sqlstate, message) + })?; state.audit_record( crate::control::security::audit::AuditEvent::AdminAction, 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 a9a83a481..53e5cbeee 100644 --- a/nodedb/tests/wire/cases/crdt_write_rls_database_scope.rs +++ b/nodedb/tests/wire/cases/crdt_write_rls_database_scope.rs @@ -132,13 +132,13 @@ async fn crdt_merge_in_non_default_database_is_rls_enforced() { "a write policy on a non-default-database CRDT collection must reject \ a CRDT MERGE its predicate forbids", ); - // `CRDT MERGE`'s handler wraps every admission failure — RLS denial - // included — under this one SQLSTATE (`crdt_merge.rs`'s `dispatch_... - // .map_err(|e| ddl_err("XX000", ...))`); the substantive assertion is - // that an error surfaces here at all, since pre-fix the merge applied - // silently and no error, of any code, reached the client. + // `CRDT MERGE`'s handler classifies every admission failure — an RLS + // denial reaches the client as `42501`. Pre-fix the merge applied + // silently and no error, of any code, reached the client, so the + // substantive assertion stays "an error surfaces at all"; the code is + // pinned because a denial is the one case a client can act on. assert_eq!( - sqlstate, "XX000", + sqlstate, "42501", "expected the CRDT admission failure's SQLSTATE, got: {sqlstate}" ); assert_eq!( @@ -199,7 +199,7 @@ async fn crdt_merge_in_default_database_is_still_rls_enforced() { MERGE its predicate forbids", ); assert_eq!( - sqlstate, "XX000", + sqlstate, "42501", "expected the CRDT admission failure's SQLSTATE, got: {sqlstate}" ); assert_eq!( @@ -209,3 +209,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; +}