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 — 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}