Skip to content

Commit 2bc878c

Browse files
committed
ddl,http: classify the CRDT MERGE refusal and code the HTTP stream errors
Two surfaces still flattened a classified failure while the routed pgwire paths already carried its class. - CRDT MERGE: the admission/apply failures wrap a classified Error — ExternalCrdtPostImagePolicy::deny returns RejectedAuthz — so they map through error_to_sqlstate, which renders 42501 for a policy denial. The "authorization returned no capability" site stays XX000: issue #344 records it as an internal invariant by decision. - HTTP stream: the in-band error lines now 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 are updated with it — `nodedb/tests/wire/cases/crdt_write_rls_database_scope.rs:141,202` expected `XX000` and now expect `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/merge site). The pgwire stream and DDL-dispatch sites are NOT touched here: the local branch `fix/shaper-error-class` carried that work and its PR #341 was closed as stacked on #359, with the instruction to resubmit against main once the mapper PR lands. 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`. Refs #344. NOT RUN: cargo check -p nodedb --tests (the wire target).
1 parent 1ff3551 commit 2bc878c

4 files changed

Lines changed: 42 additions & 18 deletions

File tree

‎nodedb/src/control/server/http/routes/query_stream.rs‎

Lines changed: 18 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -169,8 +169,11 @@ pub(super) fn ndjson_body_stream(
169169
None => break,
170170
Some(Ok(b)) => b,
171171
Some(Err(e)) => {
172-
let (_status, msg) = GatewayErrorMap::to_http(&e);
173-
let line = format!("{}\n", serde_json::json!({ "error": msg }));
172+
let (status, msg) = GatewayErrorMap::to_http(&e);
173+
let line = format!(
174+
"{}\n",
175+
serde_json::json!({ "error": msg, "status": status })
176+
);
174177
yield Ok(Bytes::from(line));
175178
return;
176179
}
@@ -182,9 +185,13 @@ pub(super) fn ndjson_body_stream(
182185
// A malformed batch payload is surfaced as an in-band error
183186
// line (matching the mid-stream dispatch-error path above)
184187
// rather than silently dropping the batch.
188+
let classified = crate::error_classify::classify(&e);
185189
let line = format!(
186190
"{}\n",
187-
serde_json::json!({ "error": format!("malformed response batch: {e}") })
191+
serde_json::json!({
192+
"error": format!("malformed response batch: {}", classified.message()),
193+
"code": classified.code().0,
194+
})
188195
);
189196
yield Ok(Bytes::from(line));
190197
return;
@@ -206,7 +213,14 @@ pub(super) fn ndjson_body_stream(
206213
Err(e) => {
207214
// In-band error line, matching the malformed-batch path
208215
// above: the HTTP body itself never errors.
209-
let line = format!("{}\n", serde_json::json!({ "error": format!("{e}") }));
216+
let classified = crate::error_classify::classify(&e);
217+
let line = format!(
218+
"{}\n",
219+
serde_json::json!({
220+
"error": classified.message(),
221+
"code": classified.code().0,
222+
})
223+
);
210224
yield Ok(Bytes::from(line));
211225
return;
212226
}

‎nodedb/src/control/server/response_shape/schema.rs‎

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -3,10 +3,10 @@
33
//! Planner-authoritative output schema types, plus the type mapping from the
44
//! planner's `SqlDataType` to the response shaper's wire-facing `DdlColType`.
55
//!
6-
//! Nothing in this module is consumed by existing call sites yet; it is a
7-
//! purely additive foundation for later threading the planner's resolved
8-
//! output schema into response shaping (replacing the SQL-string re-parse
9-
//! path).
6+
//! The planner derives the schema from the compiled plan and the catalog, and
7+
//! the session caches it with the physical tasks (`session/plan_cache.rs`).
8+
//! Shaping receives it as `MaterializedShapeRequest::projection`, which drives
9+
//! projection and the Control-Plane computed columns.
1010
1111
/// One output column of a resolved query, as known by the planner.
1212
///

‎nodedb/src/control/server/shared/ddl/neutral/dsl/crdt_merge.rs‎

Lines changed: 13 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,7 @@ use crate::bridge::envelope::PhysicalPlan;
88
use crate::control::crdt_post_image_policy::ExternalCrdtPostImagePolicy;
99
use crate::control::security::audit::ArcAuditEmitter;
1010
use crate::control::security::identity::{AuthenticatedIdentity, Permission};
11+
use crate::control::server::pgwire::types::error_map::error_to_sqlstate;
1112
use crate::control::server::shared::authorization::{authorize_collection, authorize_task_set};
1213
use crate::control::state::SharedState;
1314
use crate::types::DatabaseId;
@@ -87,7 +88,10 @@ pub async fn crdt_merge(
8788
},
8889
)
8990
.await
90-
.map_err(|e| ddl_err("XX000", e.to_string()))?;
91+
.map_err(|e| {
92+
let (_, sqlstate, message) = error_to_sqlstate(&e);
93+
ddl_err(sqlstate, message)
94+
})?;
9195
if source_bytes.is_empty() {
9296
return Err(ddl_err(
9397
"02000",
@@ -98,7 +102,10 @@ pub async fn crdt_merge(
98102
let target_surrogate = state
99103
.surrogate_assigner
100104
.assign(database_id, tenant_id, collection, target_id.as_bytes())
101-
.map_err(|e| ddl_err("XX000", e.to_string()))?;
105+
.map_err(|e| {
106+
let (_, sqlstate, message) = error_to_sqlstate(&e);
107+
ddl_err(sqlstate, message)
108+
})?;
102109

103110
let apply_plan = PhysicalPlan::Crdt(CrdtOp::Apply {
104111
collection: nodedb_types::QualifiedCollection::new(database_id, collection),
@@ -163,7 +170,10 @@ pub async fn crdt_merge(
163170
},
164171
)
165172
.await
166-
.map_err(|e| ddl_err("XX000", e.to_string()))?;
173+
.map_err(|e| {
174+
let (_, sqlstate, message) = error_to_sqlstate(&e);
175+
ddl_err(sqlstate, message)
176+
})?;
167177

168178
state.audit_record(
169179
crate::control::security::audit::AuditEvent::AdminAction,

‎nodedb/tests/wire/cases/crdt_write_rls_database_scope.rs‎

Lines changed: 7 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -132,13 +132,13 @@ async fn crdt_merge_in_non_default_database_is_rls_enforced() {
132132
"a write policy on a non-default-database CRDT collection must reject \
133133
a CRDT MERGE its predicate forbids",
134134
);
135-
// `CRDT MERGE`'s handler wraps every admission failure — RLS denial
136-
// included — under this one SQLSTATE (`crdt_merge.rs`'s `dispatch_...
137-
// .map_err(|e| ddl_err("XX000", ...))`); the substantive assertion is
138-
// that an error surfaces here at all, since pre-fix the merge applied
139-
// silently and no error, of any code, reached the client.
135+
// `CRDT MERGE`'s handler classifies every admission failure — an RLS
136+
// denial reaches the client as `42501`. Pre-fix the merge applied
137+
// silently and no error, of any code, reached the client, so the
138+
// substantive assertion stays "an error surfaces at all"; the code is
139+
// pinned because a denial is the one case a client can act on.
140140
assert_eq!(
141-
sqlstate, "XX000",
141+
sqlstate, "42501",
142142
"expected the CRDT admission failure's SQLSTATE, got: {sqlstate}"
143143
);
144144
assert_eq!(
@@ -199,7 +199,7 @@ async fn crdt_merge_in_default_database_is_still_rls_enforced() {
199199
MERGE its predicate forbids",
200200
);
201201
assert_eq!(
202-
sqlstate, "XX000",
202+
sqlstate, "42501",
203203
"expected the CRDT admission failure's SQLSTATE, got: {sqlstate}"
204204
);
205205
assert_eq!(

0 commit comments

Comments
 (0)