Skip to main content

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.
405pub async fn bind_org_scope_on(
406    conn: &mut sqlx::PgConnection,
407    scope: &OrgScope,
408) -> Result<(), sqlx::Error> {
409    sqlx::query("SELECT set_config('app.scope_unit_ids', $1, true)")
410        .bind(scope.scope_unit_ids_csv())
411        .execute(&mut *conn)
412        .await?;
413    sqlx::query("SELECT set_config('app.company_id', $1, true)")
414        .bind(scope.legacy_company_id.map(|id| id.to_string()).unwrap_or_default())
415        .execute(&mut *conn)
416        .await?;
417    sqlx::query("SELECT set_config('app.acting_unit_id', $1, true)")
418        .bind(scope.acting_unit_id.to_string())
419        .execute(&mut *conn)
420        .await?;
421    Ok(())
422}
423
424/// Statement-level org fence for a write whose row belongs to one org node — the hand-written
425/// repository path for callers outside a request scope.
426///
427/// Inside a request scope ([`with_org_request_scope`] or the legacy
428/// [`with_request_scope`](crate::company_scope::with_request_scope)) the query runs on the
429/// request connection, which already carries both fence variables. Outside one, it opens a
430/// short transaction and sets `app.scope_unit_ids` to exactly `unit`: the INSERT's
431/// `WITH CHECK` only tests the row's own `org_unit_id`, so the row's node alone satisfies it —
432/// and a read that accidentally rides this helper sees only that node (fail-narrow, never
433/// leaky). Reads that need the session's whole entitlement union must run inside a resolved
434/// request scope, not this helper.
435pub async fn execute_unit_scoped<'q>(
436    pool: &PgPool,
437    unit: Uuid,
438    query: sqlx::query::Query<'q, sqlx::Postgres, sqlx::postgres::PgArguments>,
439) -> Result<sqlx::postgres::PgQueryResult, sqlx::Error> {
440    if let Some(conn) = crate::company_scope::current_request_conn() {
441        let mut g = conn.lock().await;
442        return query.execute(&mut **g).await;
443    }
444    let mut tx = pool.begin().await?;
445    sqlx::query("SELECT set_config('app.scope_unit_ids', $1, true)")
446        .bind(unit.to_string())
447        .execute(&mut *tx)
448        .await?;
449    let res = query.execute(&mut *tx).await?;
450    tx.commit().await?;
451    Ok(res)
452}
453
454/// Tenant-agnostic `execute` for hand-written module SQL (ADR-0029): ride the request-dedicated
455/// connection when one is bound — carrying whatever fence variables the COMPOSING service's scope
456/// set (`with_org_request_scope` / `with_request_scope`) — otherwise execute plainly on the pool.
457///
458/// Unlike [`execute_unit_scoped`] this invents no scope of its own: a module that knows nothing
459/// about tenancy must not fabricate a unit or a company. Under a composer's request scope the
460/// database fence owns isolation; with no scope bound this is a plain unfenced execute.
461pub async fn execute_scoped<'q>(
462    pool: &PgPool,
463    query: sqlx::query::Query<'q, sqlx::Postgres, sqlx::postgres::PgArguments>,
464) -> Result<sqlx::postgres::PgQueryResult, sqlx::Error> {
465    if let Some(conn) = crate::company_scope::current_request_conn() {
466        let mut g = conn.lock().await;
467        return query.execute(&mut **g).await;
468    }
469    query.execute(pool).await
470}
471
472/// Tenant-agnostic `fetch_optional` for an untyped row query — the read twin of
473/// [`execute_scoped`], same connection discipline: request-dedicated connection when bound,
474/// plain pool otherwise, no scope invented.
475pub async fn fetch_optional_row_scoped<'q>(
476    pool: &PgPool,
477    query: sqlx::query::Query<'q, sqlx::Postgres, sqlx::postgres::PgArguments>,
478) -> Result<Option<sqlx::postgres::PgRow>, sqlx::Error> {
479    if let Some(conn) = crate::company_scope::current_request_conn() {
480        let mut g = conn.lock().await;
481        return query.fetch_optional(&mut **g).await;
482    }
483    query.fetch_optional(pool).await
484}
485
486/// Tenant-agnostic `fetch_all` for an untyped row query — the set-returning sibling of
487/// [`fetch_optional_row_scoped`], same connection discipline: request-dedicated connection
488/// when bound, plain pool otherwise, no scope invented. Reads that must ride the request
489/// connection (a bare-pool read under a decorated deployment lands on a fresh connection
490/// with no fence variables and returns nothing) but owe no company predicate reach for
491/// this, not the company-scoped helpers.
492pub async fn fetch_all_rows_scoped<'q>(
493    pool: &PgPool,
494    query: sqlx::query::Query<'q, sqlx::Postgres, sqlx::postgres::PgArguments>,
495) -> Result<Vec<sqlx::postgres::PgRow>, sqlx::Error> {
496    if let Some(conn) = crate::company_scope::current_request_conn() {
497        let mut g = conn.lock().await;
498        return query.fetch_all(&mut **g).await;
499    }
500    query.fetch_all(pool).await
501}
502
503/// Tenant-agnostic `fetch_one` for an untyped row query — the single-row sibling of
504/// [`fetch_optional_row_scoped`], same connection discipline: request-dedicated connection
505/// when bound, plain pool otherwise, no scope invented.
506pub async fn fetch_one_row_scoped<'q>(
507    pool: &PgPool,
508    query: sqlx::query::Query<'q, sqlx::Postgres, sqlx::postgres::PgArguments>,
509) -> Result<sqlx::postgres::PgRow, sqlx::Error> {
510    if let Some(conn) = crate::company_scope::current_request_conn() {
511        let mut g = conn.lock().await;
512        return query.fetch_one(&mut **g).await;
513    }
514    query.fetch_one(pool).await
515}
516
517#[cfg(test)]
518mod tests {
519    //! Gated on `BACKBONE_ORM_RLS_DSN` (a superuser DSN). Self-contained: builds a minimal org
520    //! spine + an org-fenced table + an app role, then proves the resolver and the request
521    //! scope against them — the same shapes the organization/inventory migrations emit.
522    use super::{resolve_org_scope, with_org_request_scope, with_org_request_scope_and_audit};
523    use crate::audit_context::RequestAuditContext;
524    use sqlx::postgres::PgPoolOptions;
525    use sqlx::PgPool;
526    use sqlx::Row;
527    use uuid::Uuid;
528
529    fn dsn() -> Option<String> {
530        std::env::var("BACKBONE_ORM_RLS_DSN").ok()
531    }
532
533    async fn admin_pool(dsn: &str) -> PgPool {
534        PgPoolOptions::new().max_connections(4).connect(dsn).await.unwrap()
535    }
536
537    async fn app_pool(dsn: &str, role: &str) -> PgPool {
538        let after_at = dsn.rsplit('@').next().unwrap();
539        let url = format!("postgresql://{role}:orgpw@{after_at}");
540        PgPoolOptions::new().max_connections(1).connect(&url).await.unwrap()
541    }
542
543    /// Mint the per-run role name — a fixed name breaks on shared dev clusters (`DROP ROLE`
544    /// fails while the role holds grants in another database), a fresh one can never collide.
545    fn role_name() -> String {
546        format!("org_scope_app_{}", &Uuid::new_v4().simple().to_string()[..8])
547    }
548
549    async fn setup(admin: &PgPool, role: &str, root: Uuid, company: Uuid, branch: Uuid, other: Uuid) {
550        sqlx::raw_sql(&format!(
551            "DROP SCHEMA IF EXISTS organization CASCADE; \
552             DROP SCHEMA IF EXISTS org_scope_test CASCADE; \
553             CREATE SCHEMA organization; \
554             CREATE TABLE organization.org_units ( \
555                 id uuid PRIMARY KEY, kind text NOT NULL, parent_id uuid, \
556                 code text, name text NOT NULL, metadata jsonb NOT NULL DEFAULT '{{}}' \
557             ); \
558             CREATE UNIQUE INDEX one_root ON organization.org_units (kind) WHERE kind = 'root'; \
559             CREATE OR REPLACE FUNCTION organization.org_unit_subtree(p_roots uuid[]) \
560                 RETURNS SETOF uuid LANGUAGE sql STABLE AS $$ \
561                 WITH RECURSIVE tree AS ( \
562                     SELECT o.id FROM organization.org_units o WHERE o.id = ANY(p_roots) \
563                     UNION ALL \
564                     SELECT o.id FROM organization.org_units o JOIN tree t ON o.parent_id = t.id \
565                 ) SELECT id FROM tree $$; \
566             CREATE OR REPLACE FUNCTION organization.org_unit_root() RETURNS uuid \
567                 LANGUAGE sql STABLE AS $$ \
568                 SELECT id FROM organization.org_units WHERE kind = 'root' $$; \
569             CREATE SCHEMA org_scope_test; \
570             CREATE TABLE org_scope_test.t (id uuid PRIMARY KEY, org_unit_id uuid NOT NULL, code text); \
571             ALTER TABLE org_scope_test.t ENABLE ROW LEVEL SECURITY; \
572             ALTER TABLE org_scope_test.t FORCE ROW LEVEL SECURITY; \
573             CREATE POLICY t_org_isolation ON org_scope_test.t FOR ALL \
574                 USING (org_unit_id = ANY(string_to_array(current_setting('app.scope_unit_ids', true), ',')::uuid[])) \
575                 WITH CHECK (org_unit_id = ANY(string_to_array(current_setting('app.scope_unit_ids', true), ',')::uuid[])); \
576             CREATE ROLE {role} LOGIN PASSWORD 'orgpw'; \
577             GRANT USAGE ON SCHEMA organization, org_scope_test TO {role}; \
578             GRANT SELECT ON organization.org_units TO {role}; \
579             GRANT SELECT, INSERT, UPDATE, DELETE ON org_scope_test.t TO {role};",
580        ))
581        .execute(admin)
582        .await
583        .unwrap();
584
585        sqlx::query(
586            "INSERT INTO organization.org_units (id, kind, parent_id, code, name) VALUES \
587                 ($1, 'root',   NULL, 'root',  'Tenant root'), \
588                 ($2, 'company', $1,   'CO',    'Company'), \
589                 ($3, 'branch',  $2,   'BR',    'Branch'), \
590                 ($4, 'company', $1,   'OTHER', 'Other Company')",
591        )
592        .bind(root)
593        .bind(company)
594        .bind(branch)
595        .bind(other)
596        .execute(admin)
597        .await
598        .unwrap();
599    }
600
601    #[tokio::test]
602    async fn resolver_and_fence_shape() {
603        let Some(dsn) = dsn() else { eprintln!("skipping: set BACKBONE_ORM_RLS_DSN"); return };
604        let root = Uuid::new_v4();
605        let company = Uuid::new_v4();
606        let branch = Uuid::new_v4();
607        let other = Uuid::new_v4();
608        let role = role_name();
609        let admin = admin_pool(&dsn).await;
610        setup(&admin, &role, root, company, branch, other).await;
611
612        // Seed as superuser (bypasses RLS).
613        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')")
614            .bind(Uuid::new_v4()).bind(company)
615            .bind(Uuid::new_v4()).bind(branch)
616            .bind(Uuid::new_v4()).bind(other)
617            .execute(&admin).await.unwrap();
618
619        // Resolve from the branch node: subtree() DESCENDS only, so the scope is the branch's
620        // own subtree (itself) plus the root node — the parent company is NOT auto-included;
621        // seeing it requires entitlement. Legacy bridge = the governing company.
622        let mut conn = admin.acquire().await.unwrap();
623        let scope = resolve_org_scope(&mut conn, branch, &[]).await.unwrap();
624        let mut got = scope.scope_unit_ids().to_vec();
625        got.sort_unstable();
626        let mut want = vec![branch, root];
627        want.sort_unstable();
628        assert_eq!(got, want);
629        assert_eq!(scope.acting_unit_id(), branch);
630        assert_eq!(scope.legacy_company_id(), Some(company));
631
632        // Union entitlement: branch acting + sister company entitled.
633        let scope_union = resolve_org_scope(&mut conn, branch, &[other]).await.unwrap();
634        assert!(scope_union.scope_unit_ids().contains(&other));
635
636        // Unknown acting node is a loud error, not a narrower scope.
637        let ghost = resolve_org_scope(&mut conn, Uuid::new_v4(), &[]).await;
638        assert!(ghost.is_err());
639
640        // The request scope fences queries on the app role's pool. The scoped helper (not a raw
641        // pool fetch) is the interesting path: it must route onto the request connection the
642        // scope bound, which carries the fence variables.
643        let pool = app_pool(&dsn, &role).await;
644        let (codes, twin_codes): (Vec<String>, Vec<String>) =
645            with_org_request_scope(&pool, scope, async {
646                let codes: Vec<String> = crate::company_scope::fetch_all_scoped(
647                    &pool,
648                    sqlx::query_as::<_, (String,)>("SELECT code FROM org_scope_test.t ORDER BY code"),
649                )
650                .await
651                .unwrap()
652                .into_iter()
653                .map(|r| r.0)
654                .collect();
655                // The tenant-agnostic fetch-all twin rides the same request
656                // connection and therefore the same fence — the honest helper
657                // for reads that owe no company predicate.
658                let twin_codes: Vec<String> = super::fetch_all_rows_scoped(
659                    &pool,
660                    sqlx::query("SELECT code FROM org_scope_test.t ORDER BY code"),
661                )
662                .await
663                .unwrap()
664                .into_iter()
665                .map(|r| r.get::<String, &str>("code"))
666                .collect();
667                (codes, twin_codes)
668            })
669            .await
670            .unwrap();
671        assert_eq!(codes, ["BR-WH"], "branch session sees only its subtree (+ shared root rows), not the company or sister-company rows");
672        assert_eq!(
673            twin_codes, codes,
674            "the org-scope fetch-all twin fences identically to the scoped helper"
675        );
676
677        // A bare pool read (no request connection, no fence variables on the
678        // freshly acquired connection) sees nothing under the decorated fence —
679        // the failure mode that makes riding the request connection mandatory.
680        let bare: Vec<String> = sqlx::query("SELECT code FROM org_scope_test.t")
681            .fetch_all(&pool)
682            .await
683            .unwrap()
684            .into_iter()
685            .map(|r| r.get::<String, usize>(0))
686            .collect();
687        assert!(
688            bare.is_empty(),
689            "a bare pool read must not see fenced rows"
690        );
691
692        // After the scope, the pooled connection is clean (fail-closed for the next acquire).
693        let mut after = pool.acquire().await.unwrap();
694        let setting: String =
695            sqlx::query_scalar("SELECT current_setting('app.scope_unit_ids', true)")
696                .fetch_one(&mut *after)
697                .await
698                .unwrap();
699        assert_eq!(setting, "", "scope leaked onto the pooled connection");
700    }
701
702    /// The audit-context channel of the request scope (ADR-0025): the audited twin binds the
703    /// six `app.*` attribution variables on the request-dedicated connection, a row-level
704    /// trigger on a write through the scoped helpers reads them off that connection, and every
705    /// variable — fence AND audit — is cleared before the connection returns to the pool. The
706    /// unaudited twin leaves the audit variables unset (the capture function's `'system'`
707    /// fallback reads exactly that: empty).
708    ///
709    /// Runs in its own disposable database (dropped at the end), because the spine-building
710    /// sibling test above recreates the `organization` schema and would race this one's if they
711    /// shared a database — cargo runs test fns concurrently.
712    #[tokio::test]
713    async fn audit_context_binds_rides_and_clears_with_the_scope() {
714        let Some(maintenance_dsn) = dsn() else {
715            eprintln!("skipping: set BACKBONE_ORM_RLS_DSN");
716            return;
717        };
718        const PROOF_DB: &str = "backbone_orm_audit_ctx_probe";
719        let proof_dsn = {
720            let (before_db, ..) = maintenance_dsn.rsplit_once('/').unwrap();
721            format!("{before_db}/{PROOF_DB}")
722        };
723        let maintenance = admin_pool(&maintenance_dsn).await;
724        // Each statement rides its own simple-protocol query — CREATE/DROP DATABASE may not run
725        // inside a transaction block.
726        sqlx::raw_sql(&format!("DROP DATABASE IF EXISTS {PROOF_DB} WITH (FORCE)"))
727            .execute(&maintenance)
728            .await
729            .unwrap();
730        sqlx::raw_sql(&format!("CREATE DATABASE {PROOF_DB}"))
731            .execute(&maintenance)
732            .await
733            .unwrap();
734        let admin = admin_pool(&proof_dsn).await;
735
736        let root = Uuid::new_v4();
737        let company = Uuid::new_v4();
738        let role = role_name();
739        // The spine + fenced table of the sibling test, plus the audit probe: a table and an
740        // AFTER INSERT trigger capturing two of the attribution variables — the minimal shape
741        // of the auditlog module's capture function, whose contract this pins without the
742        // framework crate depending on a domain module.
743        sqlx::raw_sql(&format!(
744            "CREATE SCHEMA organization; \
745             CREATE TABLE organization.org_units ( \
746                 id uuid PRIMARY KEY, kind text NOT NULL, parent_id uuid, \
747                 code text, name text NOT NULL, metadata jsonb NOT NULL DEFAULT '{{}}' \
748             ); \
749             CREATE OR REPLACE FUNCTION organization.org_unit_subtree(p_roots uuid[]) \
750                 RETURNS SETOF uuid LANGUAGE sql STABLE AS $$ \
751                 WITH RECURSIVE tree AS ( \
752                     SELECT o.id FROM organization.org_units o WHERE o.id = ANY(p_roots) \
753                     UNION ALL \
754                     SELECT o.id FROM organization.org_units o JOIN tree t ON o.parent_id = t.id \
755                 ) SELECT id FROM tree $$; \
756             CREATE OR REPLACE FUNCTION organization.org_unit_root() RETURNS uuid \
757                 LANGUAGE sql STABLE AS $$ \
758                 SELECT id FROM organization.org_units WHERE kind = 'root' $$; \
759             CREATE SCHEMA org_scope_test; \
760             CREATE TABLE org_scope_test.t (id uuid PRIMARY KEY, org_unit_id uuid NOT NULL \
761                 DEFAULT nullif(current_setting('app.acting_unit_id', true), '')::uuid, code text); \
762             ALTER TABLE org_scope_test.t ENABLE ROW LEVEL SECURITY; \
763             ALTER TABLE org_scope_test.t FORCE ROW LEVEL SECURITY; \
764             CREATE POLICY t_org_isolation ON org_scope_test.t FOR ALL \
765                 USING (org_unit_id = ANY(string_to_array(current_setting('app.scope_unit_ids', true), ',')::uuid[])) \
766                 WITH CHECK (org_unit_id = ANY(string_to_array(current_setting('app.scope_unit_ids', true), ',')::uuid[])); \
767             CREATE TABLE org_scope_test.audit_probe ( \
768                 id bigserial PRIMARY KEY, actor text NOT NULL, correlation_id text NOT NULL \
769             ); \
770             CREATE FUNCTION org_scope_test.capture_probe() RETURNS trigger LANGUAGE plpgsql AS $$ \
771                 BEGIN \
772                     INSERT INTO org_scope_test.audit_probe (actor, correlation_id) \
773                     VALUES (current_setting('app.actor', true), \
774                             current_setting('app.correlation_id', true)); \
775                     RETURN NULL; \
776                 END $$; \
777             CREATE TRIGGER t_audit_probe AFTER INSERT ON org_scope_test.t \
778                 FOR EACH ROW EXECUTE FUNCTION org_scope_test.capture_probe(); \
779             CREATE ROLE {role} LOGIN PASSWORD 'orgpw'; \
780             GRANT USAGE ON SCHEMA organization, org_scope_test TO {role}; \
781             GRANT SELECT ON organization.org_units TO {role}; \
782             GRANT SELECT, INSERT ON org_scope_test.t TO {role}; \
783             GRANT INSERT ON org_scope_test.audit_probe TO {role}; \
784             GRANT USAGE, SELECT ON SEQUENCE org_scope_test.audit_probe_id_seq TO {role};",
785        ))
786        .execute(&admin)
787        .await
788        .unwrap();
789        sqlx::query(
790            "INSERT INTO organization.org_units (id, kind, parent_id, code, name) VALUES \
791                 ($1, 'root',    NULL, 'ROOT', 'Tenant root'), \
792                 ($2, 'company', $1,   'CO',   'Company')",
793        )
794        .bind(root)
795        .bind(company)
796        .execute(&admin)
797        .await
798        .unwrap();
799
800        let pool = app_pool(&proof_dsn, &role).await;
801        let scope = {
802            let mut conn = admin.acquire().await.unwrap();
803            resolve_org_scope(&mut conn, company, &[]).await.unwrap()
804        };
805
806        // 1. The audited twin: every attribution variable is set on the connection writes ride
807        //    (inventory proof, same style as the fence-variable inventory above), and a trigger
808        //    on a real insert reads them — the seam the auditlog capture function depends on.
809        let audit = RequestAuditContext {
810            actor: "user-77".to_string(),
811            correlation_id: "corr-77".to_string(),
812            client_ip: "203.0.113.7".to_string(),
813            user_agent: "probe-agent/1.0".to_string(),
814            http_method: "POST".to_string(),
815            resource_path: "/api/v1/probe/widgets".to_string(),
816        };
817        let inserted = Uuid::new_v4();
818        let settings: Vec<(String, String)> = with_org_request_scope_and_audit(
819            &pool,
820            scope.clone(),
821            audit.clone(),
822            async {
823                crate::company_scope::execute_scoped(
824                    &pool,
825                    sqlx::query("INSERT INTO org_scope_test.t (id, code) VALUES ($1, 'AUDITED')")
826                        .bind(inserted),
827                )
828                .await
829                .unwrap();
830                crate::company_scope::fetch_all_scoped(
831                    &pool,
832                    sqlx::query_as::<_, (String, String)>(
833                        "SELECT s.name, current_setting(s.name, true) FROM unnest(ARRAY[\
834                         'app.actor','app.correlation_id','app.client_ip','app.user_agent',\
835                         'app.http_method','app.resource_path']) AS s(name)",
836                    ),
837                )
838                .await
839                .unwrap()
840            },
841        )
842        .await
843        .unwrap();
844        for (var, want) in audit.pairs() {
845            let got = settings
846                .iter()
847                .find(|(n, _)| n == var)
848                .map(|(_, v)| v.as_str())
849                .unwrap_or_else(|| panic!("{var} missing from the inventory query"));
850            assert_eq!(got, want, "{var} must ride the request connection");
851        }
852        // The fence still works under the audited twin: the acting-unit DEFAULT filled the row.
853        let landed: Uuid = sqlx::query_scalar("SELECT org_unit_id FROM org_scope_test.t WHERE id = $1")
854            .bind(inserted)
855            .fetch_one(&admin)
856            .await
857            .unwrap();
858        assert_eq!(landed, company, "audit channel must not disturb the fence");
859        let probe: (String, String) =
860            sqlx::query_as("SELECT actor, correlation_id FROM org_scope_test.audit_probe ORDER BY id DESC LIMIT 1")
861                .fetch_one(&admin)
862                .await
863                .unwrap();
864        assert_eq!(probe, ("user-77".to_string(), "corr-77".to_string()),
865            "a row-level trigger on the insert must read the bound attribution");
866
867        // 2. Pool hygiene: after the scope, no audit variable survives on the pooled connection.
868        let mut after = pool.acquire().await.unwrap();
869        for var in crate::audit_context::AUDIT_CONTEXT_VARS {
870            let v: String = sqlx::query_scalar("SELECT current_setting($1, true)")
871                .bind(var)
872                .fetch_one(&mut *after)
873                .await
874                .unwrap();
875            assert_eq!(v, "", "{var} leaked onto the pooled connection");
876        }
877        drop(after);
878
879        // 3. The unaudited twin leaves the channel unset — the capture function's `'system'`
880        //    fallback and NULL columns read exactly this empty wire form.
881        let plain = Uuid::new_v4();
882        with_org_request_scope(&pool, scope, async {
883            crate::company_scope::execute_scoped(
884                &pool,
885                sqlx::query("INSERT INTO org_scope_test.t (id, code) VALUES ($1, 'PLAIN')")
886                    .bind(plain),
887            )
888            .await
889            .unwrap();
890        })
891        .await
892        .unwrap();
893        let probe: (String, String) =
894            sqlx::query_as("SELECT actor, correlation_id FROM org_scope_test.audit_probe ORDER BY id DESC LIMIT 1")
895                .fetch_one(&admin)
896                .await
897                .unwrap();
898        assert_eq!(probe, (String::new(), String::new()),
899            "without an audit context the channel reads empty, never the previous request's");
900
901        // Zero residue: schemas + role in the proof database, then the database itself.
902        sqlx::raw_sql(&format!(
903            "DROP SCHEMA organization CASCADE; DROP SCHEMA org_scope_test CASCADE; DROP ROLE {role};"
904        ))
905        .execute(&admin)
906        .await
907        .unwrap();
908        pool.close().await;
909        admin.close().await;
910        sqlx::raw_sql(&format!("DROP DATABASE IF EXISTS {PROOF_DB} WITH (FORCE);"))
911            .execute(&maintenance)
912            .await
913            .unwrap();
914        maintenance.close().await;
915    }
916}