backbone_orm/org_scope.rs
1//! Org-tree request scope for the entitlement-union RLS fence (ADR-0028/0029).
2//!
3//! ADR-0028 replaces the single `company_id` scoping key with `org_unit_id`: one org tree per
4//! tenant database (root / company / branch nodes), and a session sees the UNION of the subtrees
5//! under every node it is entitled to — always including the root node, which owns tenant-wide
6//! shared rows. The database half is the policy
7//! `org_unit_id = ANY(string_to_array(current_setting('app.scope_unit_ids', true), ',')::uuid[])`;
8//! this module is the application half that resolves a session's entitled ids and carries them
9//! for the duration of a request.
10//!
11//! Session variables set by an org request scope, in full:
12//! - `app.scope_unit_ids` — the entitlement-union fence (read by org policies).
13//! - `app.company_id` — the legacy equality fence during the re-key transition, resolved from
14//! the acting node's company ancestry.
15//! - `app.acting_unit_id` — where new records land: the column DEFAULT
16//! `nullif(current_setting('app.acting_unit_id', true), '')::uuid` on decorated tables
17//! (ADR-0029) resolves INSERTs that omit `org_unit_id`. Unset/empty → NULL → NOT NULL
18//! violation: an insert outside a scope fails loud, never silently unscoped.
19//! - the six `app.*` audit variables of [`crate::audit_context`] (actor, correlation id,
20//! request facts) when the request carries a
21//! [`RequestAuditContext`](crate::audit_context::RequestAuditContext) — set by the audited
22//! twin [`with_org_request_scope_and_audit`], on the same request-dedicated connection, so
23//! the auditlog capture function's triggers read attribution off every write of the request.
24//!
25//! During the module-by-module re-key both fences are live at once: some tables still read
26//! `app.company_id` (ADR-0008 equality fence), org-re-keyed tables read `app.scope_unit_ids`.
27//! [`with_org_request_scope`] therefore sets ALL THREE session variables on one request-dedicated
28//! connection. When the last company-fenced table is re-keyed, the legacy bridge retires with
29//! the old fence.
30//!
31//! The scope binds the same `REQUEST_CONN` task-local as
32//! [`with_request_scope`](crate::company_scope::with_request_scope): every scoped execute helper
33//! in this crate — and therefore every generated repository — runs on that connection and
34//! inherits the fence variables without a single call-site change.
35//!
36//! **The task-local is not the fence.** RLS is. Unscoped statements see the variables unset and
37//! match zero rows — fail-closed, identical to the ADR-0008 contract.
38
39use crate::audit_context::{AUDIT_CONTEXT_VARS, RequestAuditContext};
40use sqlx::PgPool;
41use std::fmt;
42use std::future::Future;
43use std::sync::Arc;
44use tokio::sync::Mutex;
45use uuid::Uuid;
46
47/// The fence variables an org request scope sets — the reset inventory's first half. Private:
48/// callers reset through the scope wrappers, which clear fence AND audit variables together.
49const ORG_FENCE_VARS: [&str; 3] = ["app.scope_unit_ids", "app.company_id", "app.acting_unit_id"];
50
51tokio::task_local! {
52 /// Identity of the pool the ambient request connection was drawn from. Kept beside the scope so
53 /// a nested call can prove it is asking for the SAME database before reusing that connection —
54 /// with a database per tenant, reusing across pools would run the query on the wrong one.
55 static ORG_SCOPE_POOL: std::sync::Arc<sqlx::postgres::PgConnectOptions>;
56 /// The resolved org scope of the current request: entitled unit ids (union of subtrees,
57 /// root included) and the acting node. Unset for platform callers and non-request code.
58 static ORG_SCOPE: Arc<OrgScope>;
59}
60
61/// A resolved session scope over the org tree (ADR-0028).
62///
63/// Built by [`resolve_org_scope`]; carried by [`with_org_request_scope`].
64#[derive(Debug, Clone, PartialEq, Eq)]
65pub struct OrgScope {
66 scope_unit_ids: Vec<Uuid>,
67 acting_unit_id: Uuid,
68 /// The company node governing the acting node (itself for a company, nearest company
69 /// ancestor for a branch). `None` when the chain has no company node — then no legacy
70 /// `app.company_id` is set and company-fenced tables fail closed for this session.
71 legacy_company_id: Option<Uuid>,
72}
73
74impl OrgScope {
75 /// A single-company scope for paths that cannot run the resolver — chiefly composition
76 /// seams minting records on a node the request already pinned (ADR-0029's decorator makes
77 /// those inserts resolve their `org_unit_id` from `app.acting_unit_id`).
78 ///
79 /// **Precondition: `unit` must be a COMPANY node** (or whatever node kind the caller's
80 /// rows anchor on, with `legacy_company_id` semantics in mind). [`resolve_org_scope`]
81 /// derives the legacy `app.company_id` by walking ancestors from the acting node; this
82 /// constructor sets it verbatim, so a BRANCH handed here would bind a legacy variable
83 /// matching no company-fenced row and fail closed on every not-yet-stripped table.
84 ///
85 /// The scope ids are exactly `[unit]` — no root-shared rows, no sibling subtrees. That is
86 /// fail-narrow (the same property as [`execute_unit_scoped`]): inserts land on the unit
87 /// and reads see only the unit's own rows. Paths needing the session's entitlement union
88 /// must resolve a real scope instead.
89 pub fn for_company_unit(unit: Uuid) -> Self {
90 Self {
91 scope_unit_ids: vec![unit],
92 acting_unit_id: unit,
93 legacy_company_id: Some(unit),
94 }
95 }
96
97 /// Every unit id this session may see rows for: the union of the entitled subtrees plus the
98 /// root node. Order is stable (sorted) so the serialized form is deterministic.
99 pub fn scope_unit_ids(&self) -> &[Uuid] {
100 &self.scope_unit_ids
101 }
102
103 /// The node the session acts at — the default `org_unit_id` for new records and the UI's
104 /// context node.
105 pub fn acting_unit_id(&self) -> Uuid {
106 self.acting_unit_id
107 }
108
109 /// The company node for the legacy `app.company_id` fence during the re-key transition.
110 pub fn legacy_company_id(&self) -> Option<Uuid> {
111 self.legacy_company_id
112 }
113
114 /// The scope ids in the exact wire form the fence policy parses: comma-joined, no spaces.
115 fn scope_unit_ids_csv(&self) -> String {
116 self.scope_unit_ids
117 .iter()
118 .map(|id| id.to_string())
119 .collect::<Vec<_>>()
120 .join(",")
121 }
122}
123
124/// Failures of [`resolve_org_scope`] — all of them mean "do NOT open a scoped session".
125#[derive(Debug)]
126pub enum OrgScopeError {
127 /// The acting node does not exist in `organization.org_units` (wrong id, or the org spine
128 /// migration has not run in this database).
129 UnknownActingUnit(Uuid),
130 /// The database lacks the spine helpers (`organization.org_unit_subtree` / `org_unit_root`).
131 MissingSpineHelpers,
132 Database(sqlx::Error),
133}
134
135impl fmt::Display for OrgScopeError {
136 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
137 match self {
138 Self::UnknownActingUnit(id) => {
139 write!(f, "org scope: acting unit {id} is not an organization.org_units node")
140 }
141 Self::MissingSpineHelpers => write!(
142 f,
143 "org scope: organization.org_unit_subtree / org_unit_root are missing — \
144 has the org spine migration run in this database?"
145 ),
146 Self::Database(e) => write!(f, "org scope: resolution query failed: {e}"),
147 }
148 }
149}
150
151impl std::error::Error for OrgScopeError {}
152
153impl From<sqlx::Error> for OrgScopeError {
154 fn from(e: sqlx::Error) -> Self {
155 Self::Database(e)
156 }
157}
158
159/// Resolve a session's scope from the org tree.
160///
161/// `acting_unit` is the node the session acts at (from the signed token / tenant registry).
162/// `additional_entitled` are further nodes the session holds entitlements for (e.g. a group
163/// admin entitled to a sister company) — each contributes its whole subtree. The resolved scope
164/// is the union of those subtrees plus the root node.
165///
166/// `org_units` is unfenced by design (it IS the scoping substrate), so this runs on any
167/// connection of the tenant database regardless of fence state.
168pub async fn resolve_org_scope(
169 conn: &mut sqlx::PgConnection,
170 acting_unit: Uuid,
171 additional_entitled: &[Uuid],
172) -> Result<OrgScope, OrgScopeError> {
173 let mut roots = vec![acting_unit];
174 roots.extend_from_slice(additional_entitled);
175
176 let scope_unit_ids: Vec<Uuid> = sqlx::query_scalar(
177 "SELECT u.id FROM organization.org_unit_subtree($1::uuid[]) AS u(id) \
178 UNION SELECT organization.org_unit_root()",
179 )
180 .bind(&roots)
181 .fetch_all(&mut *conn)
182 .await
183 .map_err(|e| {
184 // 42883 undefined_function / 42P01 undefined_table both mean the spine is absent.
185 match &e {
186 sqlx::Error::Database(db) if db.code().as_deref() == Some("42883")
187 || db.code().as_deref() == Some("42P01") =>
188 {
189 OrgScopeError::MissingSpineHelpers
190 }
191 _ => OrgScopeError::Database(e),
192 }
193 })?;
194
195 if scope_unit_ids.is_empty() {
196 return Err(OrgScopeError::MissingSpineHelpers);
197 }
198 if !scope_unit_ids.contains(&acting_unit) {
199 // subtree() silently drops unknown roots; surface that instead of opening a
200 // narrower-than-asked scope (a wrong-but-working session is worse than a loud error).
201 return Err(OrgScopeError::UnknownActingUnit(acting_unit));
202 }
203
204 // Walk up from the acting node to the governing company node for the legacy fence.
205 let legacy_company_id: Option<Uuid> = sqlx::query_scalar(
206 "WITH RECURSIVE up AS ( \
207 SELECT id, parent_id, kind FROM organization.org_units WHERE id = $1 \
208 UNION ALL \
209 SELECT o.id, o.parent_id, o.kind FROM organization.org_units o \
210 JOIN up ON o.id = up.parent_id \
211 ) SELECT up.id FROM up WHERE up.kind::text = 'company' ORDER BY up.id LIMIT 1",
212 )
213 .bind(acting_unit)
214 .fetch_optional(&mut *conn)
215 .await?;
216
217 let mut sorted = scope_unit_ids;
218 sorted.sort_unstable();
219 sorted.dedup();
220
221 Ok(OrgScope {
222 scope_unit_ids: sorted,
223 acting_unit_id: acting_unit,
224 legacy_company_id,
225 })
226}
227
228/// Run `f` with a request-dedicated connection carrying the session's org scope.
229///
230/// Sets `app.scope_unit_ids` (the entitlement-union fence), the legacy `app.company_id`
231/// (equality fence, resolved from the acting node's company ancestry), and `app.acting_unit_id`
232/// (the acting-unit DEFAULT source for inserts on decorated tables, ADR-0029) at the session
233/// level, so org-re-keyed and not-yet-re-keyed tables are both fenced correctly for the whole
234/// request — including ID-only lookups, which ride the connection rather than the query text.
235///
236/// Mirrors [`with_request_scope`](crate::company_scope::with_request_scope)'s reset discipline:
237/// all three variables are cleared unconditionally before the connection returns to the pool,
238/// even when a `REQUEST_CONN` clone outlives the scope — a clone that runs queries after the
239/// reset does so unscoped (fail-closed), never with the previous session's scope.
240pub async fn with_org_request_scope<F, R>(pool: &PgPool, scope: OrgScope, f: F) -> Result<R, sqlx::Error>
241where
242 F: Future<Output = R>,
243{
244 with_org_request_scope_internal(pool, scope, None, f).await
245}
246
247/// [`with_org_request_scope`] plus the request's audit attribution (ADR-0025): the six
248/// [`RequestAuditContext`] variables are bound on the same request-dedicated connection, so
249/// every write of the request — including the ones that fire the auditlog capture function's
250/// triggers — reads the same actor and request facts off the connection it rides.
251///
252/// This is the wrapper a guarded route runs: the guard resolves the scope off the token, builds
253/// the audit context off the token's `sub` and the request itself, and the two channels travel
254/// together for the whole request. The reset discipline clears fence AND audit variables
255/// unconditionally: a pooled connection must never carry the previous request's attribution
256/// into the next one's audit rows.
257pub async fn with_org_request_scope_and_audit<F, R>(
258 pool: &PgPool,
259 scope: OrgScope,
260 audit: RequestAuditContext,
261 f: F,
262) -> Result<R, sqlx::Error>
263where
264 F: Future<Output = R>,
265{
266 with_org_request_scope_internal(pool, scope, Some(&audit), f).await
267}
268
269async fn with_org_request_scope_internal<F, R>(
270 pool: &PgPool,
271 scope: OrgScope,
272 audit: Option<&RequestAuditContext>,
273 f: F,
274) -> Result<R, sqlx::Error>
275where
276 F: Future<Output = R>,
277{
278 // Everything below the fast path lives in `open_org_request_scope`, behind a Box::pin. That is
279 // what keeps a nest cheap: this function's future would otherwise carry the whole slow path's
280 // locals — a pool connection, the Arc/Mutex holder, the bind statements and the task-local
281 // scope futures — and a three-level nest would stack three of those frames whether or not the
282 // slow path ever runs. Boxing moves them to the heap, so a nesting level costs a pointer.
283 // Reuse an ambient scope that already matches, instead of nesting a second one.
284 //
285 // A nested wrapper used to acquire a second connection, re-run the fence binds, and lay another
286 // prologue on the stack — roughly 600 KB per level in a debug build, so two levels plus the
287 // caller's own depth overflowed a default 2 MB worker thread and killed the task. It also held
288 // two pool connections for one logical request, which self-deadlocks a small pool.
289 //
290 // Three things must all hold before reuse is safe: a request connection is bound, the incoming
291 // scope is exactly the ambient one (a different scope needs its own binds), and the pool is the
292 // same (with a database per tenant, the ambient connection may be another tenant's entirely).
293 // A call carrying audit attribution always takes the slow path — its variables are not part of
294 // the scope comparison, so matching scopes say nothing about matching attribution.
295 if audit.is_none()
296 && crate::company_scope::current_request_conn().is_some()
297 && ORG_SCOPE.try_with(|ambient| **ambient == scope).unwrap_or(false)
298 && ORG_SCOPE_POOL
299 .try_with(|ambient| std::sync::Arc::ptr_eq(ambient, &pool.connect_options()))
300 .unwrap_or(false)
301 {
302 // The ambient scope owns the connection and will reset it; this call adds nothing to undo.
303 return Ok(f.await);
304 }
305
306 Box::pin(open_org_request_scope(pool, scope, audit, f)).await
307}
308
309/// The slow path: acquire a request-dedicated connection, bind the fence (and audit) variables on
310/// it, run `f` inside the task-locals, then reset every variable it may have set.
311async fn open_org_request_scope<F, R>(
312 pool: &PgPool,
313 scope: OrgScope,
314 audit: Option<&RequestAuditContext>,
315 f: F,
316) -> Result<R, sqlx::Error>
317where
318 F: Future<Output = R>,
319{
320 let mut conn = pool.acquire().await?;
321 sqlx::query("SELECT set_config('app.scope_unit_ids', $1, false)")
322 .bind(scope.scope_unit_ids_csv())
323 .execute(&mut *conn)
324 .await?;
325 sqlx::query("SELECT set_config('app.company_id', $1, false)")
326 .bind(scope.legacy_company_id.map(|id| id.to_string()).unwrap_or_default())
327 .execute(&mut *conn)
328 .await?;
329 sqlx::query("SELECT set_config('app.acting_unit_id', $1, false)")
330 .bind(scope.acting_unit_id.to_string())
331 .execute(&mut *conn)
332 .await?;
333 if let Some(audit) = audit {
334 audit.bind_on(&mut conn, false).await?;
335 }
336
337 let holder = Arc::new(Mutex::new(conn));
338 let scope_arc = Arc::new(scope);
339 let pool_identity = pool.connect_options();
340 // The audit attribution rides as a task-local TOO (not only as variables on
341 // the dedicated connection): write services below open their own pool
342 // transactions, whose connections the session-level binding never touches,
343 // and they relay it onto those transactions with
344 // `relay_ambient_audit_on` — the same discipline as the fence relay.
345 let audit_owned = audit.cloned();
346 let result = ORG_SCOPE_POOL
347 .scope(
348 pool_identity,
349 ORG_SCOPE.scope(scope_arc.clone(), async {
350 crate::company_scope::with_company_scope_internal(scope_arc.legacy_company_id, async {
351 let inner = crate::company_scope::with_request_conn_internal(holder.clone(), f);
352 match audit_owned {
353 Some(audit) => {
354 crate::audit_context::with_request_audit(audit, inner).await
355 }
356 None => inner.await,
357 }
358 })
359 .await
360 }),
361 )
362 .await;
363
364 // Unconditional reset of every variable the scope may have set — fence always, audit when
365 // the request carried a context (resetting anyway when it did not costs six cheap
366 // set_config calls and keeps the inventory one list; see with_request_scope for the
367 // lingering-clone reasoning — the same contract applies to every variable).
368 {
369 let mut guard = holder.lock().await;
370 for var in ORG_FENCE_VARS.into_iter().chain(AUDIT_CONTEXT_VARS) {
371 if let Err(e) = sqlx::query("SELECT set_config($1, '', false)")
372 .bind(var)
373 .execute(&mut **guard)
374 .await
375 {
376 tracing::error!(
377 target: "backbone_orm::org_scope",
378 error = %e,
379 var,
380 "failed to reset fence variable on org request connection; the pool \
381 connection may carry the previous session's scope — treat as a \
382 fence-hygiene incident",
383 );
384 }
385 }
386 }
387 Ok(result)
388}
389
390/// The resolved org scope of the current request, if one is bound.
391///
392/// Services use the acting node as the default `org_unit_id` for new records; `None` means no
393/// org scope (platform caller or non-request code) — do not guess a node then.
394pub fn current_org_scope() -> Option<OrgScope> {
395 ORG_SCOPE.try_with(|s| s.as_ref().clone()).ok()
396}
397
398/// Bind an explicit org scope onto an already-open transaction/connection, transaction-locally.
399///
400/// The transaction twin of [`with_org_request_scope`], for hand-written write services and jobs
401/// that manage their own transaction (mirrors
402/// [`bind_company_on`](crate::company_scope::bind_company_on)). Binds all three session
403/// variables, `app.acting_unit_id` included, so an insert inside the transaction can rely on
404/// the acting-unit DEFAULT — and relays the request's audit attribution, so the transaction's
405/// changes are recorded as the caller's, not `'system'`.
406pub async fn bind_org_scope_on(
407 conn: &mut sqlx::PgConnection,
408 scope: &OrgScope,
409) -> Result<(), sqlx::Error> {
410 sqlx::query("SELECT set_config('app.scope_unit_ids', $1, true)")
411 .bind(scope.scope_unit_ids_csv())
412 .execute(&mut *conn)
413 .await?;
414 sqlx::query("SELECT set_config('app.company_id', $1, true)")
415 .bind(scope.legacy_company_id.map(|id| id.to_string()).unwrap_or_default())
416 .execute(&mut *conn)
417 .await?;
418 sqlx::query("SELECT set_config('app.acting_unit_id', $1, true)")
419 .bind(scope.acting_unit_id.to_string())
420 .execute(&mut *conn)
421 .await?;
422 // The request's actor, so the transaction's audit rows name who made the change.
423 crate::audit_context::relay_ambient_audit_on(conn).await
424}
425
426/// Statement-level org fence for a write whose row belongs to one org node — the hand-written
427/// repository path for callers outside a request scope.
428///
429/// Inside a request scope ([`with_org_request_scope`] or the legacy
430/// [`with_request_scope`](crate::company_scope::with_request_scope)) the query runs on the
431/// request connection, which already carries both fence variables. Outside one, it opens a
432/// short transaction and sets `app.scope_unit_ids` to exactly `unit`: the INSERT's
433/// `WITH CHECK` only tests the row's own `org_unit_id`, so the row's node alone satisfies it —
434/// and a read that accidentally rides this helper sees only that node (fail-narrow, never
435/// leaky). Reads that need the session's whole entitlement union must run inside a resolved
436/// request scope, not this helper.
437pub async fn execute_unit_scoped<'q>(
438 pool: &PgPool,
439 unit: Uuid,
440 query: sqlx::query::Query<'q, sqlx::Postgres, sqlx::postgres::PgArguments>,
441) -> Result<sqlx::postgres::PgQueryResult, sqlx::Error> {
442 if let Some(conn) = crate::company_scope::current_request_conn() {
443 let mut g = conn.lock().await;
444 return query.execute(&mut **g).await;
445 }
446 let mut tx = pool.begin().await?;
447 sqlx::query("SELECT set_config('app.scope_unit_ids', $1, true)")
448 .bind(unit.to_string())
449 .execute(&mut *tx)
450 .await?;
451 crate::audit_context::relay_ambient_audit_on(&mut tx).await?;
452 let res = query.execute(&mut *tx).await?;
453 tx.commit().await?;
454 Ok(res)
455}
456
457/// Tenant-agnostic `execute` for hand-written module SQL (ADR-0029): ride the request-dedicated
458/// connection when one is bound — carrying whatever fence variables the COMPOSING service's scope
459/// set (`with_org_request_scope` / `with_request_scope`) — otherwise execute plainly on the pool.
460///
461/// Unlike [`execute_unit_scoped`] this invents no scope of its own: a module that knows nothing
462/// about tenancy must not fabricate a unit or a company. Under a composer's request scope the
463/// database fence owns isolation; with no scope bound this is a plain unfenced execute.
464pub async fn execute_scoped<'q>(
465 pool: &PgPool,
466 query: sqlx::query::Query<'q, sqlx::Postgres, sqlx::postgres::PgArguments>,
467) -> Result<sqlx::postgres::PgQueryResult, sqlx::Error> {
468 if let Some(conn) = crate::company_scope::current_request_conn() {
469 let mut g = conn.lock().await;
470 return query.execute(&mut **g).await;
471 }
472 query.execute(pool).await
473}
474
475/// Tenant-agnostic `fetch_optional` for an untyped row query — the read twin of
476/// [`execute_scoped`], same connection discipline: request-dedicated connection when bound,
477/// plain pool otherwise, no scope invented.
478pub async fn fetch_optional_row_scoped<'q>(
479 pool: &PgPool,
480 query: sqlx::query::Query<'q, sqlx::Postgres, sqlx::postgres::PgArguments>,
481) -> Result<Option<sqlx::postgres::PgRow>, sqlx::Error> {
482 if let Some(conn) = crate::company_scope::current_request_conn() {
483 let mut g = conn.lock().await;
484 return query.fetch_optional(&mut **g).await;
485 }
486 query.fetch_optional(pool).await
487}
488
489/// Tenant-agnostic `fetch_all` for an untyped row query — the set-returning sibling of
490/// [`fetch_optional_row_scoped`], same connection discipline: request-dedicated connection
491/// when bound, plain pool otherwise, no scope invented. Reads that must ride the request
492/// connection (a bare-pool read under a decorated deployment lands on a fresh connection
493/// with no fence variables and returns nothing) but owe no company predicate reach for
494/// this, not the company-scoped helpers.
495pub async fn fetch_all_rows_scoped<'q>(
496 pool: &PgPool,
497 query: sqlx::query::Query<'q, sqlx::Postgres, sqlx::postgres::PgArguments>,
498) -> Result<Vec<sqlx::postgres::PgRow>, sqlx::Error> {
499 if let Some(conn) = crate::company_scope::current_request_conn() {
500 let mut g = conn.lock().await;
501 return query.fetch_all(&mut **g).await;
502 }
503 query.fetch_all(pool).await
504}
505
506/// Tenant-agnostic `fetch_one` for an untyped row query — the single-row sibling of
507/// [`fetch_optional_row_scoped`], same connection discipline: request-dedicated connection
508/// when bound, plain pool otherwise, no scope invented.
509pub async fn fetch_one_row_scoped<'q>(
510 pool: &PgPool,
511 query: sqlx::query::Query<'q, sqlx::Postgres, sqlx::postgres::PgArguments>,
512) -> Result<sqlx::postgres::PgRow, sqlx::Error> {
513 if let Some(conn) = crate::company_scope::current_request_conn() {
514 let mut g = conn.lock().await;
515 return query.fetch_one(&mut **g).await;
516 }
517 query.fetch_one(pool).await
518}
519
520#[cfg(test)]
521mod tests {
522 //! Gated on `BACKBONE_ORM_RLS_DSN` (a superuser DSN). Self-contained: builds a minimal org
523 //! spine + an org-fenced table + an app role, then proves the resolver and the request
524 //! scope against them — the same shapes the organization/inventory migrations emit.
525 use super::{resolve_org_scope, with_org_request_scope, with_org_request_scope_and_audit};
526 use crate::audit_context::RequestAuditContext;
527 use sqlx::postgres::PgPoolOptions;
528 use sqlx::PgPool;
529 use sqlx::Row;
530 use uuid::Uuid;
531
532 fn dsn() -> Option<String> {
533 std::env::var("BACKBONE_ORM_RLS_DSN").ok()
534 }
535
536 async fn admin_pool(dsn: &str) -> PgPool {
537 PgPoolOptions::new().max_connections(4).connect(dsn).await.unwrap()
538 }
539
540 async fn app_pool(dsn: &str, role: &str) -> PgPool {
541 let after_at = dsn.rsplit('@').next().unwrap();
542 let url = format!("postgresql://{role}:orgpw@{after_at}");
543 PgPoolOptions::new().max_connections(1).connect(&url).await.unwrap()
544 }
545
546 /// Mint the per-run role name — a fixed name breaks on shared dev clusters (`DROP ROLE`
547 /// fails while the role holds grants in another database), a fresh one can never collide.
548 fn role_name() -> String {
549 format!("org_scope_app_{}", &Uuid::new_v4().simple().to_string()[..8])
550 }
551
552 async fn setup(admin: &PgPool, role: &str, root: Uuid, company: Uuid, branch: Uuid, other: Uuid) {
553 sqlx::raw_sql(&format!(
554 "DROP SCHEMA IF EXISTS organization CASCADE; \
555 DROP SCHEMA IF EXISTS org_scope_test CASCADE; \
556 CREATE SCHEMA organization; \
557 CREATE TABLE organization.org_units ( \
558 id uuid PRIMARY KEY, kind text NOT NULL, parent_id uuid, \
559 code text, name text NOT NULL, metadata jsonb NOT NULL DEFAULT '{{}}' \
560 ); \
561 CREATE UNIQUE INDEX one_root ON organization.org_units (kind) WHERE kind = 'root'; \
562 CREATE OR REPLACE FUNCTION organization.org_unit_subtree(p_roots uuid[]) \
563 RETURNS SETOF uuid LANGUAGE sql STABLE AS $$ \
564 WITH RECURSIVE tree AS ( \
565 SELECT o.id FROM organization.org_units o WHERE o.id = ANY(p_roots) \
566 UNION ALL \
567 SELECT o.id FROM organization.org_units o JOIN tree t ON o.parent_id = t.id \
568 ) SELECT id FROM tree $$; \
569 CREATE OR REPLACE FUNCTION organization.org_unit_root() RETURNS uuid \
570 LANGUAGE sql STABLE AS $$ \
571 SELECT id FROM organization.org_units WHERE kind = 'root' $$; \
572 CREATE SCHEMA org_scope_test; \
573 CREATE TABLE org_scope_test.t (id uuid PRIMARY KEY, org_unit_id uuid NOT NULL, code text); \
574 ALTER TABLE org_scope_test.t ENABLE ROW LEVEL SECURITY; \
575 ALTER TABLE org_scope_test.t FORCE ROW LEVEL SECURITY; \
576 CREATE POLICY t_org_isolation ON org_scope_test.t FOR ALL \
577 USING (org_unit_id = ANY(string_to_array(current_setting('app.scope_unit_ids', true), ',')::uuid[])) \
578 WITH CHECK (org_unit_id = ANY(string_to_array(current_setting('app.scope_unit_ids', true), ',')::uuid[])); \
579 CREATE ROLE {role} LOGIN PASSWORD 'orgpw'; \
580 GRANT USAGE ON SCHEMA organization, org_scope_test TO {role}; \
581 GRANT SELECT ON organization.org_units TO {role}; \
582 GRANT SELECT, INSERT, UPDATE, DELETE ON org_scope_test.t TO {role};",
583 ))
584 .execute(admin)
585 .await
586 .unwrap();
587
588 sqlx::query(
589 "INSERT INTO organization.org_units (id, kind, parent_id, code, name) VALUES \
590 ($1, 'root', NULL, 'root', 'Tenant root'), \
591 ($2, 'company', $1, 'CO', 'Company'), \
592 ($3, 'branch', $2, 'BR', 'Branch'), \
593 ($4, 'company', $1, 'OTHER', 'Other Company')",
594 )
595 .bind(root)
596 .bind(company)
597 .bind(branch)
598 .bind(other)
599 .execute(admin)
600 .await
601 .unwrap();
602 }
603
604 #[tokio::test]
605 async fn resolver_and_fence_shape() {
606 let Some(dsn) = dsn() else { eprintln!("skipping: set BACKBONE_ORM_RLS_DSN"); return };
607 let root = Uuid::new_v4();
608 let company = Uuid::new_v4();
609 let branch = Uuid::new_v4();
610 let other = Uuid::new_v4();
611 let role = role_name();
612 let admin = admin_pool(&dsn).await;
613 setup(&admin, &role, root, company, branch, other).await;
614
615 // Seed as superuser (bypasses RLS).
616 sqlx::query("INSERT INTO org_scope_test.t (id, org_unit_id, code) VALUES ($1,$2,'CO-WH'), ($3,$4,'BR-WH'), ($5,$6,'OTHER-WH')")
617 .bind(Uuid::new_v4()).bind(company)
618 .bind(Uuid::new_v4()).bind(branch)
619 .bind(Uuid::new_v4()).bind(other)
620 .execute(&admin).await.unwrap();
621
622 // Resolve from the branch node: subtree() DESCENDS only, so the scope is the branch's
623 // own subtree (itself) plus the root node — the parent company is NOT auto-included;
624 // seeing it requires entitlement. Legacy bridge = the governing company.
625 let mut conn = admin.acquire().await.unwrap();
626 let scope = resolve_org_scope(&mut conn, branch, &[]).await.unwrap();
627 let mut got = scope.scope_unit_ids().to_vec();
628 got.sort_unstable();
629 let mut want = vec![branch, root];
630 want.sort_unstable();
631 assert_eq!(got, want);
632 assert_eq!(scope.acting_unit_id(), branch);
633 assert_eq!(scope.legacy_company_id(), Some(company));
634
635 // Union entitlement: branch acting + sister company entitled.
636 let scope_union = resolve_org_scope(&mut conn, branch, &[other]).await.unwrap();
637 assert!(scope_union.scope_unit_ids().contains(&other));
638
639 // Unknown acting node is a loud error, not a narrower scope.
640 let ghost = resolve_org_scope(&mut conn, Uuid::new_v4(), &[]).await;
641 assert!(ghost.is_err());
642
643 // The request scope fences queries on the app role's pool. The scoped helper (not a raw
644 // pool fetch) is the interesting path: it must route onto the request connection the
645 // scope bound, which carries the fence variables.
646 let pool = app_pool(&dsn, &role).await;
647 let (codes, twin_codes): (Vec<String>, Vec<String>) =
648 with_org_request_scope(&pool, scope, async {
649 let codes: Vec<String> = crate::company_scope::fetch_all_scoped(
650 &pool,
651 sqlx::query_as::<_, (String,)>("SELECT code FROM org_scope_test.t ORDER BY code"),
652 )
653 .await
654 .unwrap()
655 .into_iter()
656 .map(|r| r.0)
657 .collect();
658 // The tenant-agnostic fetch-all twin rides the same request
659 // connection and therefore the same fence — the honest helper
660 // for reads that owe no company predicate.
661 let twin_codes: Vec<String> = super::fetch_all_rows_scoped(
662 &pool,
663 sqlx::query("SELECT code FROM org_scope_test.t ORDER BY code"),
664 )
665 .await
666 .unwrap()
667 .into_iter()
668 .map(|r| r.get::<String, &str>("code"))
669 .collect();
670 (codes, twin_codes)
671 })
672 .await
673 .unwrap();
674 assert_eq!(codes, ["BR-WH"], "branch session sees only its subtree (+ shared root rows), not the company or sister-company rows");
675 assert_eq!(
676 twin_codes, codes,
677 "the org-scope fetch-all twin fences identically to the scoped helper"
678 );
679
680 // A bare pool read (no request connection, no fence variables on the
681 // freshly acquired connection) sees nothing under the decorated fence —
682 // the failure mode that makes riding the request connection mandatory.
683 let bare: Vec<String> = sqlx::query("SELECT code FROM org_scope_test.t")
684 .fetch_all(&pool)
685 .await
686 .unwrap()
687 .into_iter()
688 .map(|r| r.get::<String, usize>(0))
689 .collect();
690 assert!(
691 bare.is_empty(),
692 "a bare pool read must not see fenced rows"
693 );
694
695 // After the scope, the pooled connection is clean (fail-closed for the next acquire).
696 let mut after = pool.acquire().await.unwrap();
697 let setting: String =
698 sqlx::query_scalar("SELECT current_setting('app.scope_unit_ids', true)")
699 .fetch_one(&mut *after)
700 .await
701 .unwrap();
702 assert_eq!(setting, "", "scope leaked onto the pooled connection");
703 }
704
705 /// The audit-context channel of the request scope (ADR-0025): the audited twin binds the
706 /// six `app.*` attribution variables on the request-dedicated connection, a row-level
707 /// trigger on a write through the scoped helpers reads them off that connection, and every
708 /// variable — fence AND audit — is cleared before the connection returns to the pool. The
709 /// unaudited twin leaves the audit variables unset (the capture function's `'system'`
710 /// fallback reads exactly that: empty).
711 ///
712 /// Runs in its own disposable database (dropped at the end), because the spine-building
713 /// sibling test above recreates the `organization` schema and would race this one's if they
714 /// shared a database — cargo runs test fns concurrently.
715 #[tokio::test]
716 async fn audit_context_binds_rides_and_clears_with_the_scope() {
717 let Some(maintenance_dsn) = dsn() else {
718 eprintln!("skipping: set BACKBONE_ORM_RLS_DSN");
719 return;
720 };
721 const PROOF_DB: &str = "backbone_orm_audit_ctx_probe";
722 let proof_dsn = {
723 let (before_db, ..) = maintenance_dsn.rsplit_once('/').unwrap();
724 format!("{before_db}/{PROOF_DB}")
725 };
726 let maintenance = admin_pool(&maintenance_dsn).await;
727 // Each statement rides its own simple-protocol query — CREATE/DROP DATABASE may not run
728 // inside a transaction block.
729 sqlx::raw_sql(&format!("DROP DATABASE IF EXISTS {PROOF_DB} WITH (FORCE)"))
730 .execute(&maintenance)
731 .await
732 .unwrap();
733 sqlx::raw_sql(&format!("CREATE DATABASE {PROOF_DB}"))
734 .execute(&maintenance)
735 .await
736 .unwrap();
737 let admin = admin_pool(&proof_dsn).await;
738
739 let root = Uuid::new_v4();
740 let company = Uuid::new_v4();
741 let role = role_name();
742 // The spine + fenced table of the sibling test, plus the audit probe: a table and an
743 // AFTER INSERT trigger capturing two of the attribution variables — the minimal shape
744 // of the auditlog module's capture function, whose contract this pins without the
745 // framework crate depending on a domain module.
746 sqlx::raw_sql(&format!(
747 "CREATE SCHEMA organization; \
748 CREATE TABLE organization.org_units ( \
749 id uuid PRIMARY KEY, kind text NOT NULL, parent_id uuid, \
750 code text, name text NOT NULL, metadata jsonb NOT NULL DEFAULT '{{}}' \
751 ); \
752 CREATE OR REPLACE FUNCTION organization.org_unit_subtree(p_roots uuid[]) \
753 RETURNS SETOF uuid LANGUAGE sql STABLE AS $$ \
754 WITH RECURSIVE tree AS ( \
755 SELECT o.id FROM organization.org_units o WHERE o.id = ANY(p_roots) \
756 UNION ALL \
757 SELECT o.id FROM organization.org_units o JOIN tree t ON o.parent_id = t.id \
758 ) SELECT id FROM tree $$; \
759 CREATE OR REPLACE FUNCTION organization.org_unit_root() RETURNS uuid \
760 LANGUAGE sql STABLE AS $$ \
761 SELECT id FROM organization.org_units WHERE kind = 'root' $$; \
762 CREATE SCHEMA org_scope_test; \
763 CREATE TABLE org_scope_test.t (id uuid PRIMARY KEY, org_unit_id uuid NOT NULL \
764 DEFAULT nullif(current_setting('app.acting_unit_id', true), '')::uuid, code text); \
765 ALTER TABLE org_scope_test.t ENABLE ROW LEVEL SECURITY; \
766 ALTER TABLE org_scope_test.t FORCE ROW LEVEL SECURITY; \
767 CREATE POLICY t_org_isolation ON org_scope_test.t FOR ALL \
768 USING (org_unit_id = ANY(string_to_array(current_setting('app.scope_unit_ids', true), ',')::uuid[])) \
769 WITH CHECK (org_unit_id = ANY(string_to_array(current_setting('app.scope_unit_ids', true), ',')::uuid[])); \
770 CREATE TABLE org_scope_test.audit_probe ( \
771 id bigserial PRIMARY KEY, actor text NOT NULL, correlation_id text NOT NULL \
772 ); \
773 CREATE FUNCTION org_scope_test.capture_probe() RETURNS trigger LANGUAGE plpgsql AS $$ \
774 BEGIN \
775 INSERT INTO org_scope_test.audit_probe (actor, correlation_id) \
776 VALUES (current_setting('app.actor', true), \
777 current_setting('app.correlation_id', true)); \
778 RETURN NULL; \
779 END $$; \
780 CREATE TRIGGER t_audit_probe AFTER INSERT ON org_scope_test.t \
781 FOR EACH ROW EXECUTE FUNCTION org_scope_test.capture_probe(); \
782 CREATE ROLE {role} LOGIN PASSWORD 'orgpw'; \
783 GRANT USAGE ON SCHEMA organization, org_scope_test TO {role}; \
784 GRANT SELECT ON organization.org_units TO {role}; \
785 GRANT SELECT, INSERT ON org_scope_test.t TO {role}; \
786 GRANT INSERT ON org_scope_test.audit_probe TO {role}; \
787 GRANT USAGE, SELECT ON SEQUENCE org_scope_test.audit_probe_id_seq TO {role};",
788 ))
789 .execute(&admin)
790 .await
791 .unwrap();
792 sqlx::query(
793 "INSERT INTO organization.org_units (id, kind, parent_id, code, name) VALUES \
794 ($1, 'root', NULL, 'ROOT', 'Tenant root'), \
795 ($2, 'company', $1, 'CO', 'Company')",
796 )
797 .bind(root)
798 .bind(company)
799 .execute(&admin)
800 .await
801 .unwrap();
802
803 let pool = app_pool(&proof_dsn, &role).await;
804 let scope = {
805 let mut conn = admin.acquire().await.unwrap();
806 resolve_org_scope(&mut conn, company, &[]).await.unwrap()
807 };
808
809 // 1. The audited twin: every attribution variable is set on the connection writes ride
810 // (inventory proof, same style as the fence-variable inventory above), and a trigger
811 // on a real insert reads them — the seam the auditlog capture function depends on.
812 let audit = RequestAuditContext {
813 actor: "user-77".to_string(),
814 correlation_id: "corr-77".to_string(),
815 client_ip: "203.0.113.7".to_string(),
816 user_agent: "probe-agent/1.0".to_string(),
817 http_method: "POST".to_string(),
818 resource_path: "/api/v1/probe/widgets".to_string(),
819 };
820 let inserted = Uuid::new_v4();
821 let settings: Vec<(String, String)> = with_org_request_scope_and_audit(
822 &pool,
823 scope.clone(),
824 audit.clone(),
825 async {
826 crate::company_scope::execute_scoped(
827 &pool,
828 sqlx::query("INSERT INTO org_scope_test.t (id, code) VALUES ($1, 'AUDITED')")
829 .bind(inserted),
830 )
831 .await
832 .unwrap();
833 crate::company_scope::fetch_all_scoped(
834 &pool,
835 sqlx::query_as::<_, (String, String)>(
836 "SELECT s.name, current_setting(s.name, true) FROM unnest(ARRAY[\
837 'app.actor','app.correlation_id','app.client_ip','app.user_agent',\
838 'app.http_method','app.resource_path']) AS s(name)",
839 ),
840 )
841 .await
842 .unwrap()
843 },
844 )
845 .await
846 .unwrap();
847 for (var, want) in audit.pairs() {
848 let got = settings
849 .iter()
850 .find(|(n, _)| n == var)
851 .map(|(_, v)| v.as_str())
852 .unwrap_or_else(|| panic!("{var} missing from the inventory query"));
853 assert_eq!(got, want, "{var} must ride the request connection");
854 }
855 // The fence still works under the audited twin: the acting-unit DEFAULT filled the row.
856 let landed: Uuid = sqlx::query_scalar("SELECT org_unit_id FROM org_scope_test.t WHERE id = $1")
857 .bind(inserted)
858 .fetch_one(&admin)
859 .await
860 .unwrap();
861 assert_eq!(landed, company, "audit channel must not disturb the fence");
862 let probe: (String, String) =
863 sqlx::query_as("SELECT actor, correlation_id FROM org_scope_test.audit_probe ORDER BY id DESC LIMIT 1")
864 .fetch_one(&admin)
865 .await
866 .unwrap();
867 assert_eq!(probe, ("user-77".to_string(), "corr-77".to_string()),
868 "a row-level trigger on the insert must read the bound attribution");
869
870 // 2. Pool hygiene: after the scope, no audit variable survives on the pooled connection.
871 let mut after = pool.acquire().await.unwrap();
872 for var in crate::audit_context::AUDIT_CONTEXT_VARS {
873 let v: String = sqlx::query_scalar("SELECT current_setting($1, true)")
874 .bind(var)
875 .fetch_one(&mut *after)
876 .await
877 .unwrap();
878 assert_eq!(v, "", "{var} leaked onto the pooled connection");
879 }
880 drop(after);
881
882 // 3. The unaudited twin leaves the channel unset — the capture function's `'system'`
883 // fallback and NULL columns read exactly this empty wire form.
884 let plain = Uuid::new_v4();
885 with_org_request_scope(&pool, scope, async {
886 crate::company_scope::execute_scoped(
887 &pool,
888 sqlx::query("INSERT INTO org_scope_test.t (id, code) VALUES ($1, 'PLAIN')")
889 .bind(plain),
890 )
891 .await
892 .unwrap();
893 })
894 .await
895 .unwrap();
896 let probe: (String, String) =
897 sqlx::query_as("SELECT actor, correlation_id FROM org_scope_test.audit_probe ORDER BY id DESC LIMIT 1")
898 .fetch_one(&admin)
899 .await
900 .unwrap();
901 assert_eq!(probe, (String::new(), String::new()),
902 "without an audit context the channel reads empty, never the previous request's");
903
904 // Zero residue: schemas + role in the proof database, then the database itself.
905 sqlx::raw_sql(&format!(
906 "DROP SCHEMA organization CASCADE; DROP SCHEMA org_scope_test CASCADE; DROP ROLE {role};"
907 ))
908 .execute(&admin)
909 .await
910 .unwrap();
911 pool.close().await;
912 admin.close().await;
913 sqlx::raw_sql(&format!("DROP DATABASE IF EXISTS {PROOF_DB} WITH (FORCE);"))
914 .execute(&maintenance)
915 .await
916 .unwrap();
917 maintenance.close().await;
918 }
919}