Skip to main content

nmbrs_runtime/
resource_pool.rs

1// Copyright 2024-2026 Jonathan Shook
2// SPDX-License-Identifier: Apache-2.0
3
4//! Generic resource lifecycle and sharing pool — SRD-35 Push A.
5//!
6//! ## What this module provides
7//!
8//! Trait surface and runtime structure for sharing
9//! long-lived resources (CQL sessions, HTTP clients, …)
10//! across phases of one session, decoupled from per-phase
11//! adapter shells. The contract is intentionally generic —
12//! nothing here is CQL-specific.
13//!
14//! ### Core types
15//!
16//! - [`SharedResource`] — the trait every poolable resource
17//!   implements. Carries [`SharedResource::can_share`]
18//!   (capability declaration: thread-safe + designed for
19//!   sharing) and [`SharedResource::can_support_more_load`] (driver's
20//!   runtime judgement that another instance would relieve
21//!   knowable, substantial contention) — see the SRD-35
22//!   load-bearing rules.
23//! - [`ResourceKey`] — value-equality identity. Two keys
24//!   compare equal iff their adapter name and fields match
25//!   exactly. No derived hash function; the contract is
26//!   structural.
27//! - [`ShareCapability`] — strictest sharing the resource
28//!   *type* tolerates. Read by the pool at planning time
29//!   without an instance.
30//! - [`ResourceSharePolicy`] — user-elevatable isolation
31//!   policy. Must satisfy `policy >= capability_floor`.
32//! - [`ResourcePool`] — owns the `(ResourceKey, generation)
33//!   → Entry` map, tracks refcounts, emits lifecycle events.
34//!
35//! ## Push A scope
36//!
37//! Push A lays the trait foundation, the pool data
38//! structure, the lifecycle event emission, and a
39//! [`LegacyAdapterResource`] shim that wraps the existing
40//! `Arc<dyn DriverAdapter>` factories under `PerPhase`
41//! policy. Behaviour is byte-identical to today (a fresh
42//! resource per phase) but every phase boundary now emits
43//! the full `resource.{attach,init,detach,close}` event
44//! sequence.
45//!
46//! Push B will migrate the CQL adapter to `Shared` policy
47//! by splitting `CqlAdapter` into a per-phase shell + a
48//! `CassDriverInstance` that implements `SharedResource`
49//! directly with real async `close()`. The pool's
50//! multi-generation machinery is in place but exercised
51//! only by the synthetic mock resource in tests until then.
52
53use std::any::Any;
54use std::collections::{BTreeMap, HashMap};
55use std::future::Future;
56use std::pin::Pin;
57use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
58use std::sync::{Arc, Mutex, Weak};
59use std::time::Instant;
60
61use crate::adapter::DriverAdapter;
62use crate::observer::LogLevel;
63
64// =================================================================
65// Resource key — value-equality identity
66// =================================================================
67
68/// Structural value identity for a shared resource.
69///
70/// Two adapters that build equal `ResourceKey`s receive the
71/// same instance under `Shared` policy. The internal
72/// `BTreeMap` is sorted so equality is deterministic, and
73/// `Hash` is derived structurally for `HashMap` lookup. The
74/// pool never asks adapters to produce a "key digest"; the
75/// contract is value equality on the public type.
76///
77/// Adapters populate `fields` with **only** the params that
78/// distinguish one instance from another. Per-statement and
79/// per-phase shaping (timeouts, trace rates, dynamic
80/// controls) MUST NOT appear in the key — see SRD-35
81/// §"Instance-shaping vs shell-shaping params".
82#[derive(Clone, Debug, PartialEq, Eq, Hash, Default)]
83pub struct ResourceKey {
84    /// Adapter name (`"cql"`, `"http"`, …) — the
85    /// adapter-name surface, not the engine-specific driver
86    /// identifier (which the adapter folds into the fields
87    /// when meaningful).
88    pub adapter: String,
89    /// Identity-bearing param values, sorted so equality is
90    /// independent of insertion order.
91    pub fields: BTreeMap<String, String>,
92}
93
94impl ResourceKey {
95    /// Construct a key for the given adapter with no fields
96    /// yet. Chain [`Self::with`] to populate.
97    pub fn new(adapter: impl Into<String>) -> Self {
98        Self {
99            adapter: adapter.into(),
100            fields: BTreeMap::new(),
101        }
102    }
103
104    /// Set or replace the value for one identity-bearing
105    /// field, returning `self` for chaining.
106    pub fn with(mut self, key: impl Into<String>, value: impl Into<String>) -> Self {
107        self.fields.insert(key.into(), value.into());
108        self
109    }
110
111    /// Render the key for a single line of log output
112    /// (`adapter{k1=v1,k2=v2}`). Stable shape — the
113    /// lifecycle-event surface relies on this format being
114    /// consistent across log consumers.
115    pub fn fmt_for_log(&self) -> String {
116        let mut s = String::with_capacity(64);
117        s.push_str(&self.adapter);
118        s.push('{');
119        let mut first = true;
120        for (k, v) in &self.fields {
121            if !first {
122                s.push(',');
123            }
124            first = false;
125            s.push_str(k);
126            s.push('=');
127            // Don't log secrets. Adapters can prefix
128            // sensitive fields with `_secret_` to redact.
129            if k.starts_with("_secret_") || k == "password" {
130                s.push_str("***");
131            } else {
132                s.push_str(v);
133            }
134        }
135        s.push('}');
136        s
137    }
138
139    /// SRD-104 — the **canonical fingerprint rendering** of this key to a
140    /// stable string, used as the accessor lookup key
141    /// ([`polydat::ResourceAccessor::lookup`]). Unlike [`Self::fmt_for_log`]
142    /// this rendering is **identity-preserving**: two keys that differ
143    /// only in an identity-bearing field (e.g. `password`) MUST render to
144    /// different strings or they would collide in the accessor. A secret
145    /// field (`password`, `_secret_*`) contributes a SHA-256 digest of its
146    /// value rather than the value itself, so the fingerprint keeps its
147    /// identity without ever carrying the cleartext into a kernel binding
148    /// or a diagnostic. The `BTreeMap` makes the field order
149    /// deterministic, so equal keys always render identically.
150    ///
151    /// This is the single rendering shared by both sides of the SRD-104
152    /// bridge: the resource pool renders each live entry's key with it when
153    /// resolving a lookup, and a consumer node builds the same string from
154    /// the fingerprint in scope. Shape: `adapter{k1=v1,k2=v2}`.
155    pub fn render_key(&self) -> String {
156        let mut s = String::with_capacity(64);
157        s.push_str(&self.adapter);
158        s.push('{');
159        let mut first = true;
160        for (k, v) in &self.fields {
161            if !first {
162                s.push(',');
163            }
164            first = false;
165            s.push_str(k);
166            s.push('=');
167            if k.starts_with("_secret_") || k == "password" {
168                use sha2::{Digest, Sha256};
169                let digest = Sha256::digest(v.as_bytes());
170                s.push_str("sha256:");
171                for byte in &digest[..8] {
172                    s.push_str(&format!("{byte:02x}"));
173                }
174            } else {
175                s.push_str(v);
176            }
177        }
178        s.push('}');
179        s
180    }
181}
182
183// =================================================================
184// Capability and policy
185// =================================================================
186
187/// Strongest sharing the resource *type* can tolerate.
188/// Declared by the factory at registration time so the pool
189/// can plan the entry layout without configuring or
190/// instantiating any resource. The live instance's
191/// [`SharedResource::can_share`] is a runtime safety net
192/// that aborts the session if it disagrees with this
193/// type-level declaration.
194///
195/// The variants order from most-shared to most-isolated;
196/// `PartialOrd` semantics let the pool check policy ≥
197/// capability with a single comparison.
198#[derive(Copy, Clone, Debug, PartialEq, Eq, PartialOrd, Ord, Hash)]
199pub enum ShareCapability {
200    /// Safe to share one instance across every adapter
201    /// shell in the workload that produces an equal key.
202    /// CQL, HTTP-with-pool, OpenAPI typically.
203    Shared,
204    /// One instance per scenario subtree. Use when state
205    /// changes during the run that mustn't bleed across
206    /// scenarios (auth tokens, schema versions per
207    /// scenario).
208    PerScenario,
209    /// One instance per phase. Use when each phase
210    /// legitimately wants a fresh state object.
211    PerPhase,
212    /// One instance per fiber. Use only when the resource
213    /// type is `Send` but not `Sync`, or when the
214    /// underlying library mandates single-threaded use.
215    PerFiber,
216}
217
218/// User-selectable isolation. Must satisfy
219/// `policy >= capability_floor`.
220#[derive(Copy, Clone, Debug, PartialEq, Eq, PartialOrd, Ord, Hash)]
221pub enum ResourceSharePolicy {
222    Shared,
223    PerScenario,
224    PerPhase,
225    PerFiber,
226}
227
228impl ResourceSharePolicy {
229    /// Map a user-string (CLI / workload param) to the
230    /// enum. Used by the SRD-04 umbrella surface
231    /// (`--resource-share <adapter>:<policy>`).
232    pub fn parse(s: &str) -> Result<Self, String> {
233        match s.trim().to_ascii_lowercase().as_str() {
234            "shared" => Ok(Self::Shared),
235            "per-scenario" | "per_scenario" => Ok(Self::PerScenario),
236            "per-phase" | "per_phase" => Ok(Self::PerPhase),
237            "per-fiber" | "per_fiber" => Ok(Self::PerFiber),
238            other => Err(format!(
239                "unknown resource-share policy '{other}' \
240                 (expected: shared, per-scenario, per-phase, per-fiber)"
241            )),
242        }
243    }
244}
245
246impl std::fmt::Display for ResourceSharePolicy {
247    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
248        let s = match self {
249            Self::Shared => "shared",
250            Self::PerScenario => "per-scenario",
251            Self::PerPhase => "per-phase",
252            Self::PerFiber => "per-fiber",
253        };
254        f.write_str(s)
255    }
256}
257
258/// Default sharing policy — the strongest sharing the
259/// capability allows. Used when the user gives no override.
260pub fn default_policy_for(cap: ShareCapability) -> ResourceSharePolicy {
261    match cap {
262        ShareCapability::Shared => ResourceSharePolicy::Shared,
263        ShareCapability::PerScenario => ResourceSharePolicy::PerScenario,
264        ShareCapability::PerPhase => ResourceSharePolicy::PerPhase,
265        ShareCapability::PerFiber => ResourceSharePolicy::PerFiber,
266    }
267}
268
269/// Minimum policy compatible with the given capability.
270/// Used by the pool to validate user-supplied policies at
271/// session start.
272pub fn capability_floor(cap: ShareCapability) -> ResourceSharePolicy {
273    match cap {
274        ShareCapability::Shared => ResourceSharePolicy::Shared,
275        ShareCapability::PerScenario => ResourceSharePolicy::PerScenario,
276        ShareCapability::PerPhase => ResourceSharePolicy::PerPhase,
277        ShareCapability::PerFiber => ResourceSharePolicy::PerFiber,
278    }
279}
280
281// =================================================================
282// SharedResource trait
283// =================================================================
284
285/// One async-init-and-close return type used by the
286/// trait. Returns a `Send` future so the pool can hold it
287/// across `.await` points without runtime fuss.
288pub type ResourceFuture<'a, T> = Pin<Box<dyn Future<Output = T> + Send + 'a>>;
289
290/// The trait every poolable resource implements.
291///
292/// `Send + Sync + 'static` is required so the pool can hold
293/// `Arc<dyn SharedResource>` and clone it across fibers.
294///
295/// Default implementations make the trivial case
296/// (always-shareable, never-saturated, no-op init/close)
297/// boilerplate-free — only `resource_key()` is required to
298/// override.
299pub trait SharedResource: Send + Sync + 'static {
300    /// The structural identity of this resource.
301    fn resource_key(&self) -> &ResourceKey;
302
303    /// **Capability** — does this instance support being
304    /// shared by multiple adapter shells concurrently?
305    /// Default `true`. See SRD-35 §"Two trait methods on
306    /// the live resource".
307    fn can_share(&self) -> bool {
308        true
309    }
310
311    /// **Live capacity** — can this instance accept *another*
312    /// concurrent caller right now without substantial
313    /// contention?
314    ///
315    /// `true` (default) → "yes, route the next attach to me;
316    /// I have capacity." `false` → "no, I'm saturated; the
317    /// pool should spawn a sibling for the new attach."
318    /// Parallel naming to `can_share()`: `can_share()` says
319    /// *whether* sharing is structurally possible at all;
320    /// `can_support_more_load()` says *whether one more
321    /// caller is OK at this moment*.
322    ///
323    /// Drivers that override MUST document their decision
324    /// criterion in the type docstring (the operator reading
325    /// a `reason=capacity-declined` lifecycle event needs to
326    /// be able to interpret what triggered it).
327    /// MUST be cheap (atomic read or short metric query);
328    /// MUST NOT block. SRD-35 §"Validity rules" requires
329    /// the body to reflect *current* in-flight load — never
330    /// historical/peak/lifetime metrics. The pool's guard in
331    /// [`needs_sibling_spawn`] catches the historical-state
332    /// failure mode (driver returns `false` at quiescence)
333    /// and emits a `Warn` event so operators can spot the
334    /// driver bug.
335    fn can_support_more_load(&self) -> bool {
336        true
337    }
338
339    /// Optional async init beyond what construction
340    /// already did. Called by the pool on first attach for
341    /// a key, exactly once per `(key, generation)`.
342    fn init(&self) -> ResourceFuture<'_, Result<(), String>> {
343        Box::pin(async { Ok(()) })
344    }
345
346    /// Symmetric teardown. Called when the entry's
347    /// refcount hits zero. MUST block until network /
348    /// kernel resources are *actually released* — no
349    /// async-Drop races, no sockets in TIME_WAIT being
350    /// counted as closed, no libuv worker thread races.
351    fn close(self: Arc<Self>) -> ResourceFuture<'static, Result<(), String>> {
352        Box::pin(async { Ok(()) })
353    }
354
355    /// Bridge for the Push A legacy-adapter shim. Default
356    /// returns `None`; only [`LegacyAdapterResource`]
357    /// overrides it to surface its wrapped `DriverAdapter`
358    /// handle. Push B retires this entirely once each
359    /// adapter has its own `SharedResource` impl with no
360    /// hidden `DriverAdapter` inside.
361    ///
362    /// Real `SharedResource` implementations should leave
363    /// this at the default — the pool layer is the right
364    /// place to surface domain-specific handles, not the
365    /// trait surface.
366    fn as_legacy_adapter(&self) -> Option<Arc<dyn DriverAdapter>> {
367        None
368    }
369
370    /// SRD-104 — the resource's **accessor payload**: a type-erased handle
371    /// (`Arc<dyn Any + Send + Sync>`) a kernel node can obtain by
372    /// fingerprint through its kernel tree's resource scope. The pool stores it on
373    /// the entry right after a successful init, and its
374    /// [`polydat::ResourceAccessor`] impl hands out clones. Default `None`
375    /// — a resource opts in only when it wants kernels to reach a live
376    /// handle (the first consumer is a CQL session handle, SRD-103); no
377    /// existing adapter is affected.
378    fn accessor_payload(&self) -> Option<Arc<dyn Any + Send + Sync>> {
379        None
380    }
381}
382
383// =================================================================
384// Internal entry — one per (key, generation)
385// =================================================================
386
387/// Pool entry tracking one resource instance and its
388/// refcount state.
389struct Entry {
390    key: ResourceKey,
391    policy: ResourceSharePolicy,
392    generation: usize,
393
394    /// The lazily-constructed resource. `tokio::sync::OnceCell`
395    /// would work too; we use a `Mutex<Option<…>>` because
396    /// Push A doesn't yet need cross-fiber concurrent first-
397    /// attach contention (the pool's outer mutex serialises),
398    /// and the simpler shape avoids pulling tokio into the
399    /// type signature.
400    resource: Mutex<Option<Arc<dyn SharedResource>>>,
401
402    /// Set when init failed. Subsequent attaches return the
403    /// cached error rather than retrying.
404    poisoned: AtomicBool,
405    /// Cached init error message, populated alongside
406    /// `poisoned`.
407    init_error: Mutex<Option<String>>,
408
409    /// Predicted future attaches that haven't yet landed
410    /// on this entry. SRD-35 Push D's pre-map walker
411    /// (`pre_map_pending_uses`) seeds this from the
412    /// scenario tree at session bootstrap; each `attach`
413    /// decrements eagerly. The pool's close trigger fires
414    /// the moment `pending_uses == 0 && live_attaches ==
415    /// 0`, releasing the resource the instant its last
416    /// predicted user is done.
417    pending_uses: AtomicUsize,
418    /// Currently in-use shells. Decremented on detach.
419    live_attaches: AtomicUsize,
420
421    /// Wall-clock instant at which `init()` started — used
422    /// to compute `elapsed_ms` on `init.completed` /
423    /// `init.failed` events.
424    init_started_at: Mutex<Option<Instant>>,
425
426    /// SRD-104 — the resource's type-erased **accessor payload**, populated
427    /// right after a successful init from [`SharedResource::accessor_payload`]
428    /// and returned by the pool's [`polydat::ResourceAccessor`] impl. `None`
429    /// until init succeeds, and `None` for resources that don't opt in (the
430    /// trait default). The payload lives and dies with the entry — no second
431    /// store.
432    accessor_payload: Mutex<Option<Arc<dyn Any + Send + Sync>>>,
433}
434
435impl Entry {
436    fn new(key: ResourceKey, policy: ResourceSharePolicy, generation: usize) -> Self {
437        Self {
438            key,
439            policy,
440            generation,
441            resource: Mutex::new(None),
442            poisoned: AtomicBool::new(false),
443            init_error: Mutex::new(None),
444            pending_uses: AtomicUsize::new(0),
445            live_attaches: AtomicUsize::new(0),
446            init_started_at: Mutex::new(None),
447            accessor_payload: Mutex::new(None),
448        }
449    }
450}
451
452// =================================================================
453// `can_support_more_load()` validity guard
454// =================================================================
455
456/// Decide whether the pool should spawn a sibling
457/// generation for this attach.
458///
459/// Returns `true` only when the active resource genuinely
460/// can't take another caller (`can_support_more_load() ==
461/// false`) AND there's actual concurrent load
462/// (`live_attaches > 0`). Both halves are required:
463///
464/// - `can_support_more_load() == true` → existing instance
465///   has capacity; route the attach there. No spawn.
466/// - `can_support_more_load() == false` AND `live > 0` →
467///   genuine saturation; spawn a sibling.
468/// - `can_support_more_load() == false` AND `live == 0` →
469///   **driver bug**: the resource refused a new caller
470///   while sitting completely idle. SRD-35 §"Validity
471///   rules" forbids this — `can_support_more_load()` must
472///   reflect *current* in-flight load, never historical /
473///   peak / lifetime state. The classic failure mode is a
474///   ring buffer or moving-average metric that filled
475///   during an earlier phase and hasn't decayed; reading
476///   it at quiescence falsely reports saturation. The
477///   guard logs a `Warn` so the operator sees the driver
478///   bug, and routes the attach to the existing idle
479///   instance instead of wasting a sibling.
480///
481/// Push A/B do not yet wire the sibling-spawn path, but
482/// this helper lands now so Push D's spawn logic uses it
483/// from day one. Tests exercise it directly.
484#[allow(dead_code)]
485fn needs_sibling_spawn(entry: &Entry, resource: &dyn SharedResource) -> bool {
486    if resource.can_support_more_load() {
487        // Existing generation has capacity — route there.
488        return false;
489    }
490    let live = entry.live_attaches.load(Ordering::Acquire);
491    if live == 0 {
492        // Driver bug: refused a new attach at quiescence.
493        // Log once per occurrence with the resource key so
494        // the operator can attribute it. Don't spawn — the
495        // existing idle generation is the right answer.
496        crate::diag!(
497            LogLevel::Warn,
498            "{EVENT_FAMILY}.share.suppressed key={} generation={} \
499             reason=quiescent-decline \
500             note=can_support_more_load() returned false with live_attaches=0; \
501             driver may be reading historical state (filled ring buffer, \
502             moving-average that hasn't decayed) instead of in-flight load \
503             (SRD-35 §\"Validity rules\")",
504            entry.key.fmt_for_log(),
505            entry.generation,
506        );
507        return false;
508    }
509    // Genuine saturation: refusing under live load.
510    true
511}
512
513// =================================================================
514// Lifecycle event emission
515// =================================================================
516
517/// Family prefix on every event line emitted by the pool.
518/// Enables log-side filtering with one substring.
519const EVENT_FAMILY: &str = "resource";
520
521/// Emit one lifecycle event line. All events carry `key`,
522/// `generation`, and `policy`; per-event extra fields are
523/// appended to the formatted suffix.
524fn emit_event(level: LogLevel, name: &str, entry: &Entry, extra: &str) {
525    let suffix = if extra.is_empty() {
526        String::new()
527    } else {
528        format!(" {extra}")
529    };
530    crate::diag!(
531        level,
532        "{EVENT_FAMILY}.{name} key={} generation={} policy={}{suffix}",
533        entry.key.fmt_for_log(),
534        entry.generation,
535        entry.policy,
536    );
537}
538
539// =================================================================
540// Resource pool
541// =================================================================
542
543/// One `(ResourceKey, generation)` map for the lifetime of
544/// one session. Owns lazy init, refcount transitions, and
545/// the close trigger.
546///
547/// The pool is `Send + Sync` and intended to be held in an
548/// `Arc` from the session root (executor's `ExecCtx` or
549/// equivalent). Each phase activation calls
550/// [`ResourcePool::attach`] for the adapter names it needs;
551/// the returned [`AttachGuard`] holds a clone of the
552/// per-phase `DriverAdapter` shell and detaches on drop.
553pub struct ResourcePool {
554    inner: Mutex<PoolInner>,
555}
556
557struct PoolInner {
558    /// Active entries keyed by `(ResourceKey, generation)`.
559    /// We use a `Vec<Arc<Entry>>` for ordered iteration;
560    /// the `HashMap` is the lookup index.
561    entries_by_key: HashMap<(ResourceKey, usize), Arc<Entry>>,
562
563    /// Cumulative attach count — used for diagnostic
564    /// post-mortems and the `nmbrs describe drivers` output
565    /// that's deferred to a later SRD.
566    total_attaches: usize,
567}
568
569impl ResourcePool {
570    pub fn new() -> Self {
571        Self {
572            inner: Mutex::new(PoolInner {
573                entries_by_key: HashMap::new(),
574                total_attaches: 0,
575            }),
576        }
577    }
578
579    /// Total attach count since pool construction. Useful
580    /// for tests and post-mortem diagnostics.
581    pub fn total_attaches(&self) -> usize {
582        self.inner
583            .lock()
584            .unwrap_or_else(|e| e.into_inner())
585            .total_attaches
586    }
587
588    /// Number of distinct entries currently live in the
589    /// pool. One per `(key, generation)` that's been
590    /// attached at least once and not yet fully drained.
591    pub fn live_entries(&self) -> usize {
592        self.inner
593            .lock()
594            .unwrap_or_else(|e| e.into_inner())
595            .entries_by_key
596            .len()
597    }
598
599    /// Drain every live entry, awaiting each resource's
600    /// `close()` and removing the entry from the map. Call
601    /// at session end so `Shared`-policy entries — which
602    /// the per-attach detach intentionally keeps alive
603    /// across phases — release their network resources
604    /// before the process exits.
605    ///
606    /// Emits `resource.close.started reason=session-end`
607    /// per entry, then `resource.close.completed` /
608    /// `resource.close.failed` after each resource's
609    /// `close()` future resolves. Returns `Ok(())` even if
610    /// individual close calls fail — the failures are
611    /// logged and accounted, but the pool always finishes
612    /// the drain so the executor can move on.
613    pub async fn shutdown(self: &Arc<Self>) {
614        // Snapshot the entries, then drop the lock so the
615        // close futures can run unblocked.
616        let entries: Vec<Arc<Entry>> = {
617            let inner = self.inner.lock().unwrap_or_else(|e| e.into_inner());
618            inner.entries_by_key.values().cloned().collect()
619        };
620        for entry in entries {
621            // Take the resource out of the slot. If it's
622            // already gone (concurrent close, or never
623            // realised after init failure), skip.
624            let resource = {
625                let mut slot = entry.resource.lock().unwrap_or_else(|e| e.into_inner());
626                slot.take()
627            };
628            let Some(resource) = resource else {
629                self.remove_entry(&entry.key, entry.generation);
630                continue;
631            };
632            let live = entry.live_attaches.load(Ordering::Acquire);
633            let pending = entry.pending_uses.load(Ordering::Acquire);
634            // Note residual counts on the started event so
635            // operators can spot bugs in the pre-map walker
636            // when it lands (Push D): if `pending > 0` at
637            // session end, the walker over-predicted; if
638            // `live > 0`, an attach guard wasn't dropped.
639            emit_event(
640                LogLevel::Debug,
641                "close.started",
642                &entry,
643                &format!("reason=session-end live={live} pending={pending}"),
644            );
645            let started_at = Instant::now();
646            let result = resource.close().await;
647            let elapsed_ms = started_at.elapsed().as_millis() as u64;
648            match result {
649                Ok(()) => emit_event(
650                    LogLevel::Debug,
651                    "close.completed",
652                    &entry,
653                    &format!("elapsed_ms={elapsed_ms}"),
654                ),
655                Err(ref e) => emit_event(
656                    LogLevel::Warn,
657                    "close.failed",
658                    &entry,
659                    &format!("elapsed_ms={elapsed_ms} error={e:?}"),
660                ),
661            }
662            self.remove_entry(&entry.key, entry.generation);
663        }
664    }
665
666    /// SRD-35 Push D: declare a predicted future attach for
667    /// `key` under `policy`. Called by the pre-map walker
668    /// once per phase that will eventually attach this key
669    /// — the per-key counter accumulates so a `Shared`
670    /// entry attached by 50 phases starts the run with
671    /// `pending_uses = 50`. Each `attach` decrements
672    /// eagerly; the entry's close trigger fires the moment
673    /// `pending == 0 && live == 0`, releasing the resource
674    /// the instant its last user is done rather than holding
675    /// it until session end.
676    ///
677    /// Idempotent for a given `(key, policy)` — the entry is
678    /// created on first call, then increments on subsequent
679    /// calls. The pre-map walker is the only declared caller;
680    /// adapters never invoke this directly.
681    pub fn declare_pending_use(&self, key: ResourceKey, policy: ResourceSharePolicy) {
682        let entry = self.get_or_create_entry(key, policy, 0);
683        entry.pending_uses.fetch_add(1, Ordering::AcqRel);
684    }
685
686    /// SRD-35 Push D: explicit per-key `pending_uses`
687    /// decrement, mirroring [`Self::declare_pending_use`].
688    /// Most callers don't need this — the [`AttachGuard`]
689    /// detach path already decrements via the `attach`-time
690    /// "slot consumed" semantic, and the close trigger
691    /// observes the resulting `pending == 0 && live == 0`
692    /// invariant. Provided so external callers (tests, future
693    /// pre-map invalidation paths) can adjust the counter
694    /// explicitly. Saturating at 0 — calling this on a key
695    /// with no pending uses left is a no-op, not an
696    /// underflow.
697    ///
698    /// Returns `true` when the call drove `pending` to 0
699    /// AND `live` was already 0 (so the entry was eligible
700    /// for close). Doesn't itself trigger close — the close
701    /// path runs on detach. Tests that exercise the
702    /// counter-only lifecycle use the return to assert the
703    /// "eligible for close" gate without setting up a real
704    /// attach.
705    pub fn complete_pending_use(&self, key: &ResourceKey) -> bool {
706        let inner = self.inner.lock().unwrap_or_else(|e| e.into_inner());
707        let Some(entry) = inner.entries_by_key.get(&(key.clone(), 0)).cloned() else {
708            return false;
709        };
710        drop(inner);
711        let cur = entry.pending_uses.load(Ordering::Acquire);
712        if cur == 0 {
713            return entry.live_attaches.load(Ordering::Acquire) == 0;
714        }
715        let new_pending = entry.pending_uses.fetch_sub(1, Ordering::AcqRel) - 1;
716        new_pending == 0 && entry.live_attaches.load(Ordering::Acquire) == 0
717    }
718
719    /// SRD-35 Push D introspection helper: how many predicted
720    /// future attaches remain on the entry for `key` (gen 0).
721    /// Returns `None` when no entry exists for the key.
722    /// Tests use this to assert pre-map walker correctness
723    /// without poking the pool internals.
724    pub fn pending_uses_for(&self, key: &ResourceKey) -> Option<usize> {
725        let inner = self.inner.lock().unwrap_or_else(|e| e.into_inner());
726        inner
727            .entries_by_key
728            .get(&(key.clone(), 0))
729            .map(|e| e.pending_uses.load(Ordering::Acquire))
730    }
731
732    /// Look up or create the entry for the given key under
733    /// the given policy. Generation 0 by default; siblings
734    /// (Push D) will create higher generations on
735    /// `can_support_more_load()` recommendation.
736    fn get_or_create_entry(
737        &self,
738        key: ResourceKey,
739        policy: ResourceSharePolicy,
740        generation: usize,
741    ) -> Arc<Entry> {
742        let mut inner = self.inner.lock().unwrap_or_else(|e| e.into_inner());
743        let map_key = (key.clone(), generation);
744        if let Some(existing) = inner.entries_by_key.get(&map_key) {
745            return existing.clone();
746        }
747        let entry = Arc::new(Entry::new(key, policy, generation));
748        inner.entries_by_key.insert(map_key, entry.clone());
749        entry
750    }
751
752    /// Remove an entry from the active map. Called when its
753    /// refcount has fully drained.
754    fn remove_entry(&self, key: &ResourceKey, generation: usize) {
755        let mut inner = self.inner.lock().unwrap_or_else(|e| e.into_inner());
756        inner.entries_by_key.remove(&(key.clone(), generation));
757    }
758
759    /// SRD-104 — resolve the accessor payload for the live entry whose
760    /// [`ResourceKey`] renders (via [`ResourceKey::render_key`]) to `key`.
761    /// Returns a clone of the stored `Arc<dyn Any>`, or `None` when no live
762    /// entry matches the fingerprint or the matched entry has no payload
763    /// (not opted in, or not yet initialised). This is the pool's half of
764    /// the [`polydat::ResourceAccessor`] bridge; see [`install_accessor`].
765    fn lookup_accessor_payload(&self, key: &str) -> Option<Arc<dyn Any + Send + Sync>> {
766        let inner = self.inner.lock().unwrap_or_else(|e| e.into_inner());
767        for entry in inner.entries_by_key.values() {
768            if entry.key.render_key() == key {
769                return entry
770                    .accessor_payload
771                    .lock()
772                    .unwrap_or_else(|e| e.into_inner())
773                    .clone();
774            }
775        }
776        None
777    }
778
779    /// Lazily realise the resource for an entry, emitting
780    /// `resource.init.{started,completed,failed}` events
781    /// around the call. Subsequent calls observe the
782    /// already-populated slot and return the existing
783    /// resource (or the cached error if init was poisoned).
784    async fn ensure_initialized(
785        &self,
786        entry: &Arc<Entry>,
787        first_attach: bool,
788        factory: ResourceFactory<'_>,
789    ) -> Result<Arc<dyn SharedResource>, String> {
790        // Fast path — already realised.
791        if let Some(existing) = entry
792            .resource
793            .lock()
794            .unwrap_or_else(|e| e.into_inner())
795            .clone()
796        {
797            return Ok(existing);
798        }
799        if entry.poisoned.load(Ordering::Acquire) {
800            let err = entry
801                .init_error
802                .lock()
803                .unwrap_or_else(|e| e.into_inner())
804                .clone()
805                .unwrap_or_else(|| "init poisoned".into());
806            return Err(err);
807        }
808
809        let reason = if first_attach {
810            "first-attach"
811        } else {
812            "capacity-declined"
813        };
814        emit_event(
815            LogLevel::Debug,
816            "init.started",
817            entry,
818            &format!("reason={reason}"),
819        );
820        *entry
821            .init_started_at
822            .lock()
823            .unwrap_or_else(|e| e.into_inner()) = Some(Instant::now());
824
825        // Run the factory + the resource's own `init()`.
826        // Errors poison the entry so subsequent attaches
827        // surface the same error rather than retrying.
828        let outcome: Result<Arc<dyn SharedResource>, String> = async {
829            let resource = factory.build().await?;
830            resource.init().await?;
831            Ok(resource)
832        }
833        .await;
834
835        let elapsed_ms = entry
836            .init_started_at
837            .lock()
838            .unwrap_or_else(|e| e.into_inner())
839            .map(|t| t.elapsed().as_millis() as u64)
840            .unwrap_or(0);
841
842        match outcome {
843            Ok(resource) => {
844                *entry.resource.lock().unwrap_or_else(|e| e.into_inner()) = Some(resource.clone());
845                // SRD-104: capture the resource's accessor payload (if any)
846                // onto the entry so the pool's ResourceAccessor impl can hand
847                // it to kernel nodes by fingerprint. `None` for resources that
848                // don't opt in (the trait default) — no-op then.
849                *entry
850                    .accessor_payload
851                    .lock()
852                    .unwrap_or_else(|e| e.into_inner()) = resource.accessor_payload();
853                emit_event(
854                    LogLevel::Debug,
855                    "init.completed",
856                    entry,
857                    &format!("elapsed_ms={elapsed_ms}"),
858                );
859                Ok(resource)
860            }
861            Err(e) => {
862                entry.poisoned.store(true, Ordering::Release);
863                *entry.init_error.lock().unwrap_or_else(|e| e.into_inner()) = Some(e.clone());
864                emit_event(
865                    LogLevel::Error,
866                    "init.failed",
867                    entry,
868                    &format!("elapsed_ms={elapsed_ms} error={e:?}"),
869                );
870                Err(e)
871            }
872        }
873    }
874}
875
876impl Default for ResourcePool {
877    fn default() -> Self {
878        Self::new()
879    }
880}
881
882/// Boxed, `Send` one-shot closure that builds a shared resource,
883/// returning a [`ResourceFuture`] yielding the resource (or an
884/// error message).
885type ResourceFactoryFn<'a> =
886    Box<dyn FnOnce() -> ResourceFuture<'a, Result<Arc<dyn SharedResource>, String>> + Send + 'a>;
887
888/// Closure-shaped factory for building a resource the first
889/// time its key is attached. Boxed so the pool can hold it
890/// across `.await` points and so callers don't have to type
891/// the future signature inline.
892pub struct ResourceFactory<'a> {
893    inner: ResourceFactoryFn<'a>,
894}
895
896impl<'a> ResourceFactory<'a> {
897    pub fn new<F, Fut>(f: F) -> Self
898    where
899        F: FnOnce() -> Fut + Send + 'a,
900        Fut: Future<Output = Result<Arc<dyn SharedResource>, String>> + Send + 'a,
901    {
902        Self {
903            inner: Box::new(move || Box::pin(f()) as ResourceFuture<'a, _>),
904        }
905    }
906
907    fn build(self) -> ResourceFuture<'a, Result<Arc<dyn SharedResource>, String>> {
908        (self.inner)()
909    }
910}
911
912// =================================================================
913// Public attach / detach surface
914// =================================================================
915
916/// Public attach call. Builds (or reuses) the entry for
917/// `key` under `policy`, lazily realises the resource via
918/// `factory` if this is the first attach, emits
919/// `resource.attach`, and returns a guard that detaches on
920/// drop.
921///
922/// `phase` is a free-form label that lands on the
923/// `resource.attach` / `resource.detach` event lines so
924/// operators can correlate lifecycle events with the
925/// per-phase activation that drove them.
926pub async fn attach(
927    pool: &Arc<ResourcePool>,
928    key: ResourceKey,
929    policy: ResourceSharePolicy,
930    phase: impl Into<String>,
931    factory: ResourceFactory<'_>,
932) -> Result<AttachGuard, String> {
933    let phase = phase.into();
934    let entry = pool.get_or_create_entry(key.clone(), policy, 0);
935    let first_attach = entry
936        .resource
937        .lock()
938        .unwrap_or_else(|e| e.into_inner())
939        .is_none()
940        && !entry.poisoned.load(Ordering::Acquire);
941
942    let resource = pool
943        .ensure_initialized(&entry, first_attach, factory)
944        .await?;
945
946    // Capability check: the live resource's can_share()
947    // must agree with the policy. PerPhase / PerFiber
948    // policies don't require can_share()=true; only
949    // Shared / PerScenario do.
950    let needs_share = matches!(
951        policy,
952        ResourceSharePolicy::Shared | ResourceSharePolicy::PerScenario
953    );
954    if needs_share && !resource.can_share() {
955        return Err(format!(
956            "resource for {} declared can_share()=false but policy is {policy}; \
957             elevate isolation to per-phase or per-fiber, or fix the resource impl",
958            entry.key.fmt_for_log(),
959        ));
960    }
961
962    // Bookkeeping: increment live attach, emit attach event
963    let live = entry.live_attaches.fetch_add(1, Ordering::AcqRel) + 1;
964    let pending_dec = entry.pending_uses.load(Ordering::Acquire);
965    let pending = if pending_dec > 0 {
966        entry.pending_uses.fetch_sub(1, Ordering::AcqRel) - 1
967    } else {
968        0
969    };
970    {
971        let mut inner = pool.inner.lock().unwrap_or_else(|e| e.into_inner());
972        inner.total_attaches += 1;
973    }
974    emit_event(
975        LogLevel::Debug,
976        "attach",
977        &entry,
978        &format!("phase={phase:?} pending={pending} live={live}"),
979    );
980
981    Ok(AttachGuard {
982        pool: Arc::clone(pool),
983        entry,
984        resource,
985        phase,
986        detached: false,
987    })
988}
989
990/// Owns an attached resource for the duration of one phase
991/// activation. Drop emits `resource.detach` and, when the
992/// entry's refcount has fully drained, schedules an async
993/// `close()` (and emits `resource.close.*` events around it).
994///
995/// Holding the guard keeps the resource alive; releasing
996/// it lets the pool tear down the resource if no other
997/// shells reference it.
998///
999/// Implements `Debug` so test helpers like
1000/// `Result::expect_err` accept it as the success type when
1001/// the error path is the one under test.
1002pub struct AttachGuard {
1003    pool: Arc<ResourcePool>,
1004    entry: Arc<Entry>,
1005    resource: Arc<dyn SharedResource>,
1006    phase: String,
1007    detached: bool,
1008}
1009
1010impl std::fmt::Debug for AttachGuard {
1011    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1012        f.debug_struct("AttachGuard")
1013            .field("key", &self.entry.key)
1014            .field("generation", &self.entry.generation)
1015            .field("policy", &self.entry.policy)
1016            .field("phase", &self.phase)
1017            .field("detached", &self.detached)
1018            .finish()
1019    }
1020}
1021
1022impl AttachGuard {
1023    /// Borrow the underlying shared resource. The guard
1024    /// retains ownership; the `Arc` clone is cheap.
1025    pub fn resource(&self) -> Arc<dyn SharedResource> {
1026        Arc::clone(&self.resource)
1027    }
1028
1029    /// Resource key for the attached entry. Useful for log
1030    /// correlation in adapter shells.
1031    pub fn key(&self) -> &ResourceKey {
1032        &self.entry.key
1033    }
1034
1035    /// Explicit detach. Calling this drains the refcount
1036    /// and (if it hits zero) awaits the resource's
1037    /// `close()` synchronously, so the caller observes
1038    /// teardown completion at this point. The guard's
1039    /// `Drop` falls back to a best-effort spawn when
1040    /// `detach()` wasn't called explicitly.
1041    pub async fn detach(mut self) -> Result<(), String> {
1042        self.detach_inner(true).await
1043    }
1044
1045    async fn detach_inner(&mut self, await_close: bool) -> Result<(), String> {
1046        if self.detached {
1047            return Ok(());
1048        }
1049        self.detached = true;
1050
1051        let live = self.entry.live_attaches.fetch_sub(1, Ordering::AcqRel) - 1;
1052        let pending = self.entry.pending_uses.load(Ordering::Acquire);
1053        emit_event(
1054            LogLevel::Debug,
1055            "detach",
1056            &self.entry,
1057            &format!("phase={:?} pending={pending} live={live}", self.phase),
1058        );
1059
1060        // SRD-35 Push D: close trigger is unified across
1061        // every policy. The pool closes an entry the moment
1062        // its `(pending == 0, live == 0)` invariant is
1063        // satisfied — meaning "no more attaches predicted by
1064        // the pre-map AND nothing currently using it."
1065        //
1066        //   - `Shared` / `PerScenario`: pre-map seeds
1067        //     `pending_uses` to the predicted attach count
1068        //     (sum across all phases that produce this key).
1069        //     Each `attach` decrements pending eagerly; the
1070        //     last phase's detach drives both counters to
1071        //     zero and the entry closes immediately —
1072        //     releasing network resources before session end
1073        //     for keys whose users are all done. If the
1074        //     pre-map walker missed a phase (under-predict),
1075        //     the entry would close prematurely and be
1076        //     re-initialised on the next attach; if it
1077        //     over-predicts, `pending` stays > 0 and
1078        //     `pool.shutdown()` does the close at session
1079        //     end (the conservative default).
1080        //
1081        //   - `PerPhase` / `PerFiber`: pre-map intentionally
1082        //     skips these — each phase is its own key, so
1083        //     `pending` stays 0 throughout. Close fires on
1084        //     the first detach when `live` returns to 0.
1085        if live == 0 && pending == 0 {
1086            self.trigger_close(await_close, "refcount-zero").await?;
1087        }
1088        Ok(())
1089    }
1090
1091    async fn trigger_close(&self, await_close: bool, reason: &str) -> Result<(), String> {
1092        // Pull the resource out of the entry slot — only
1093        // the close path holds it after this point, and a
1094        // late attach would observe an empty slot and
1095        // re-init under a fresh generation.
1096        let resource_opt = {
1097            let mut slot = self
1098                .entry
1099                .resource
1100                .lock()
1101                .unwrap_or_else(|e| e.into_inner());
1102            slot.take()
1103        };
1104        let Some(resource) = resource_opt else {
1105            // Entry never realised (factory failed before
1106            // resource was set, or already closed).
1107            self.pool
1108                .remove_entry(&self.entry.key, self.entry.generation);
1109            return Ok(());
1110        };
1111
1112        emit_event(
1113            LogLevel::Debug,
1114            "close.started",
1115            &self.entry,
1116            &format!("reason={reason}"),
1117        );
1118        let started_at = Instant::now();
1119
1120        let entry_for_async = Arc::clone(&self.entry);
1121        let pool_for_async = Arc::clone(&self.pool);
1122        let close_future = resource.close();
1123
1124        if await_close {
1125            let result = close_future.await;
1126            let elapsed_ms = started_at.elapsed().as_millis() as u64;
1127            match result {
1128                Ok(()) => emit_event(
1129                    LogLevel::Debug,
1130                    "close.completed",
1131                    &self.entry,
1132                    &format!("elapsed_ms={elapsed_ms}"),
1133                ),
1134                Err(ref e) => emit_event(
1135                    LogLevel::Warn,
1136                    "close.failed",
1137                    &self.entry,
1138                    &format!("elapsed_ms={elapsed_ms} error={e:?}"),
1139                ),
1140            }
1141            self.pool
1142                .remove_entry(&self.entry.key, self.entry.generation);
1143            result
1144        } else {
1145            // Drop-path: spawn the close so we don't block
1146            // the dropping thread, but still emit the
1147            // events when it finishes. Errors get logged
1148            // via the events themselves.
1149            tokio::spawn(async move {
1150                let result = close_future.await;
1151                let elapsed_ms = started_at.elapsed().as_millis() as u64;
1152                match result {
1153                    Ok(()) => emit_event(
1154                        LogLevel::Debug,
1155                        "close.completed",
1156                        &entry_for_async,
1157                        &format!("elapsed_ms={elapsed_ms}"),
1158                    ),
1159                    Err(ref e) => emit_event(
1160                        LogLevel::Warn,
1161                        "close.failed",
1162                        &entry_for_async,
1163                        &format!("elapsed_ms={elapsed_ms} error={e:?}"),
1164                    ),
1165                }
1166                pool_for_async.remove_entry(&entry_for_async.key, entry_for_async.generation);
1167            });
1168            Ok(())
1169        }
1170    }
1171}
1172
1173impl Drop for AttachGuard {
1174    fn drop(&mut self) {
1175        if self.detached {
1176            return;
1177        }
1178        // Drop is sync; schedule the detach work onto the
1179        // current tokio runtime if there is one. If no
1180        // runtime is active (test fallback), do the
1181        // synchronous bookkeeping inline and skip the
1182        // close await.
1183        let live = self.entry.live_attaches.fetch_sub(1, Ordering::AcqRel) - 1;
1184        let pending = self.entry.pending_uses.load(Ordering::Acquire);
1185        self.detached = true;
1186        let phase = &self.phase;
1187        emit_event(
1188            LogLevel::Debug,
1189            "detach",
1190            &self.entry,
1191            &format!("phase={phase:?} pending={pending} live={live}"),
1192        );
1193
1194        // SRD-35 Push D unified close gate (mirrors
1195        // `detach_inner`): close as soon as both counters
1196        // hit zero, regardless of policy. Pre-map seeds
1197        // `pending_uses` so Shared/PerScenario entries close
1198        // when their last user drains — not at session end.
1199        if live == 0 && pending == 0 {
1200            // Take the resource out so the close path owns
1201            // it — even if no runtime is active to await
1202            // the close, the synchronous bookkeeping is
1203            // correct.
1204            let resource_opt = {
1205                let mut slot = self
1206                    .entry
1207                    .resource
1208                    .lock()
1209                    .unwrap_or_else(|e| e.into_inner());
1210                slot.take()
1211            };
1212            let Some(resource) = resource_opt else {
1213                self.pool
1214                    .remove_entry(&self.entry.key, self.entry.generation);
1215                return;
1216            };
1217
1218            // Try to spawn the close on the current
1219            // runtime. If we're outside a runtime
1220            // (tests, dryrun), the resource simply drops
1221            // synchronously here.
1222            let entry = Arc::clone(&self.entry);
1223            let pool = Arc::clone(&self.pool);
1224            emit_event(
1225                LogLevel::Debug,
1226                "close.started",
1227                &self.entry,
1228                "reason=refcount-zero",
1229            );
1230            let started_at = Instant::now();
1231
1232            if let Ok(handle) = tokio::runtime::Handle::try_current() {
1233                let close_future = resource.close();
1234                handle.spawn(async move {
1235                    let result = close_future.await;
1236                    let elapsed_ms = started_at.elapsed().as_millis() as u64;
1237                    match result {
1238                        Ok(()) => emit_event(
1239                            LogLevel::Debug,
1240                            "close.completed",
1241                            &entry,
1242                            &format!("elapsed_ms={elapsed_ms}"),
1243                        ),
1244                        Err(ref e) => emit_event(
1245                            LogLevel::Warn,
1246                            "close.failed",
1247                            &entry,
1248                            &format!("elapsed_ms={elapsed_ms} error={e:?}"),
1249                        ),
1250                    }
1251                    pool.remove_entry(&entry.key, entry.generation);
1252                });
1253            } else {
1254                // No runtime: drop synchronously, emit
1255                // a synthetic close.completed with
1256                // elapsed_ms=0 so log consumers still see
1257                // a paired close event.
1258                drop(resource);
1259                emit_event(
1260                    LogLevel::Debug,
1261                    "close.completed",
1262                    &self.entry,
1263                    "elapsed_ms=0",
1264                );
1265                self.pool
1266                    .remove_entry(&self.entry.key, self.entry.generation);
1267            }
1268        }
1269    }
1270}
1271
1272// =================================================================
1273// Legacy adapter wrapper — Push A bridge
1274// =================================================================
1275
1276/// Wraps an `Arc<dyn DriverAdapter>` from the existing
1277/// per-phase factory under `PerPhase` policy. Used until
1278/// each adapter migrates to a real `SharedResource` impl
1279/// (Push B for CQL, Push C for HTTP / OpenAPI / stdout).
1280///
1281/// `can_share()` returns `false` so the pool refuses to
1282/// give a `LegacyAdapterResource` to a second shell — it's
1283/// effectively single-use, matching today's behaviour.
1284///
1285/// Most callers shouldn't construct this directly — use
1286/// [`attach_legacy_adapter`] which handles the wrapping,
1287/// pool registration, and adapter unwrapping in one call.
1288pub struct LegacyAdapterResource {
1289    key: ResourceKey,
1290    adapter: Mutex<Option<Arc<dyn DriverAdapter>>>,
1291}
1292
1293impl LegacyAdapterResource {
1294    pub fn new(key: ResourceKey, adapter: Arc<dyn DriverAdapter>) -> Self {
1295        Self {
1296            key,
1297            adapter: Mutex::new(Some(adapter)),
1298        }
1299    }
1300
1301    /// Borrow the wrapped legacy adapter. Returns `None`
1302    /// after `close()` has consumed it.
1303    pub fn adapter(&self) -> Option<Arc<dyn DriverAdapter>> {
1304        self.adapter
1305            .lock()
1306            .unwrap_or_else(|e| e.into_inner())
1307            .clone()
1308    }
1309}
1310
1311impl SharedResource for LegacyAdapterResource {
1312    fn resource_key(&self) -> &ResourceKey {
1313        &self.key
1314    }
1315
1316    /// Legacy adapters are NOT yet declared shareable —
1317    /// each phase gets its own. Push B / C migrate
1318    /// adapters out of this shim and into real
1319    /// `can_share() = true` impls.
1320    fn can_share(&self) -> bool {
1321        false
1322    }
1323
1324    fn close(self: Arc<Self>) -> ResourceFuture<'static, Result<(), String>> {
1325        Box::pin(async move {
1326            // Drop the adapter Arc — the legacy factory
1327            // doesn't expose an async-close surface, so
1328            // this is best-effort sync teardown.
1329            let _ = self
1330                .adapter
1331                .lock()
1332                .unwrap_or_else(|e| e.into_inner())
1333                .take();
1334            Ok(())
1335        })
1336    }
1337
1338    fn as_legacy_adapter(&self) -> Option<Arc<dyn DriverAdapter>> {
1339        self.adapter()
1340    }
1341}
1342
1343/// Push B sibling of [`LegacyAdapterResource`] — wraps an
1344/// `Arc<dyn DriverAdapter>` but declares `can_share()=true`.
1345/// Used when an adapter has registered a
1346/// [`crate::adapter::SharedDriverRegistration`] opting it
1347/// into pool-shared semantics. The pool caches the
1348/// underlying `Arc<dyn DriverAdapter>` and hands the same
1349/// clone to every phase whose params produce the same
1350/// `ResourceKey`.
1351///
1352/// `close()` drops the wrapped adapter Arc — the actual
1353/// teardown timing is governed by the underlying
1354/// `DriverAdapter`'s `Drop`. Push D will add real async
1355/// `close()` to specific driver instances (e.g. awaiting
1356/// `cass_session_close()`'s `CassFuture` for the
1357/// cassandra-cpp engine).
1358pub struct SharedAdapterResource {
1359    key: ResourceKey,
1360    adapter: Mutex<Option<Arc<dyn DriverAdapter>>>,
1361}
1362
1363impl SharedAdapterResource {
1364    pub fn new(key: ResourceKey, adapter: Arc<dyn DriverAdapter>) -> Self {
1365        Self {
1366            key,
1367            adapter: Mutex::new(Some(adapter)),
1368        }
1369    }
1370
1371    pub fn adapter(&self) -> Option<Arc<dyn DriverAdapter>> {
1372        self.adapter
1373            .lock()
1374            .unwrap_or_else(|e| e.into_inner())
1375            .clone()
1376    }
1377}
1378
1379impl SharedResource for SharedAdapterResource {
1380    fn resource_key(&self) -> &ResourceKey {
1381        &self.key
1382    }
1383
1384    /// Push B: this wrapper declares the adapter shareable.
1385    /// The pool keeps one instance per `ResourceKey` and
1386    /// returns the same `Arc` to every matching attach.
1387    fn can_share(&self) -> bool {
1388        true
1389    }
1390
1391    fn close(self: Arc<Self>) -> ResourceFuture<'static, Result<(), String>> {
1392        Box::pin(async move {
1393            // Take the adapter Arc out of the slot so the
1394            // shutdown handshake runs on a stable owner.
1395            // Drop happens at the end of this scope; the
1396            // engine's Drop runs after `shutdown().await`
1397            // resolves.
1398            let adapter = self
1399                .adapter
1400                .lock()
1401                .unwrap_or_else(|e| e.into_inner())
1402                .take();
1403            if let Some(adapter) = adapter {
1404                adapter.shutdown().await;
1405            }
1406            Ok(())
1407        })
1408    }
1409
1410    fn as_legacy_adapter(&self) -> Option<Arc<dyn DriverAdapter>> {
1411        self.adapter()
1412    }
1413
1414    /// SRD-104 — surface the wrapped adapter's accessor payload (SRD-103's
1415    /// CQL session handle for the `cql` adapter). The pool captures this
1416    /// right after init and hands clones to kernel nodes by fingerprint. The
1417    /// pool-shared wrapper owns no payload itself; it delegates to the
1418    /// adapter, which builds one over its connected session. `None` for
1419    /// adapters that don't opt in (the trait default).
1420    fn accessor_payload(&self) -> Option<Arc<dyn Any + Send + Sync>> {
1421        self.adapter().and_then(|a| a.accessor_payload())
1422    }
1423}
1424
1425/// Push B high-level helper: attach a shareable adapter
1426/// through the pool under `Shared` policy. The first phase
1427/// whose params produce `key` triggers `factory`; every
1428/// subsequent matching phase reuses the same
1429/// `Arc<dyn DriverAdapter>`. Each call returns a fresh
1430/// guard; the pool's refcount drops the entry only when
1431/// every guard is released.
1432pub async fn attach_shared_adapter<F, Fut>(
1433    pool: &Arc<ResourcePool>,
1434    adapter_name: &str,
1435    phase: &str,
1436    key: ResourceKey,
1437    factory: F,
1438) -> Result<(Arc<dyn DriverAdapter>, AttachGuard), String>
1439where
1440    F: FnOnce() -> Fut + Send + 'static,
1441    Fut: Future<Output = Result<Arc<dyn DriverAdapter>, String>> + Send + 'static,
1442{
1443    let key_for_resource = key.clone();
1444    let factory = ResourceFactory::new(move || async move {
1445        let adapter = factory().await?;
1446        let res = SharedAdapterResource::new(key_for_resource, adapter);
1447        Ok(Arc::new(res) as Arc<dyn SharedResource>)
1448    });
1449    let guard = attach(pool, key, ResourceSharePolicy::Shared, phase, factory).await?;
1450    let adapter = guard.resource().as_legacy_adapter().ok_or_else(|| {
1451        format!(
1452            "internal: shared pool resource for adapter '{adapter_name}' did not surface \
1453         a DriverAdapter handle — should be unreachable"
1454        )
1455    })?;
1456    Ok((adapter, guard))
1457}
1458
1459/// Push A high-level helper for the executor: attach a
1460/// legacy adapter through the pool under `PerPhase` policy
1461/// and return both the unwrapped `Arc<dyn DriverAdapter>`
1462/// and the corresponding [`AttachGuard`].
1463///
1464/// The guard MUST outlive any use of the returned adapter
1465/// reference — typically the caller stores it alongside
1466/// the activity and drops it at phase teardown so the
1467/// pool sees the matching `resource.detach`.
1468///
1469/// `factory` builds the adapter the first time the entry
1470/// is realised. The pool only calls it once per
1471/// `(key, generation)` even though `PerPhase` policy
1472/// makes each phase its own key.
1473pub async fn attach_legacy_adapter<F, Fut>(
1474    pool: &Arc<ResourcePool>,
1475    adapter_name: &str,
1476    phase: &str,
1477    key_extras: &[(&str, &str)],
1478    factory: F,
1479) -> Result<(Arc<dyn DriverAdapter>, AttachGuard), String>
1480where
1481    F: FnOnce() -> Fut + Send + 'static,
1482    Fut: Future<Output = Result<Arc<dyn DriverAdapter>, String>> + Send + 'static,
1483{
1484    let mut key = ResourceKey::new(adapter_name);
1485    for (k, v) in key_extras {
1486        key = key.with(*k, *v);
1487    }
1488    let key_for_resource = key.clone();
1489    let factory = ResourceFactory::new(move || async move {
1490        let adapter = factory().await?;
1491        let res = LegacyAdapterResource::new(key_for_resource, adapter);
1492        Ok(Arc::new(res) as Arc<dyn SharedResource>)
1493    });
1494
1495    let guard = attach(pool, key, ResourceSharePolicy::PerPhase, phase, factory).await?;
1496    let resource = guard.resource();
1497    let adapter = resource.as_legacy_adapter().ok_or_else(|| {
1498        format!(
1499            "internal: pool resource for adapter '{adapter_name}' did not surface \
1500         a legacy DriverAdapter handle — Push A path should be unreachable"
1501        )
1502    })?;
1503    Ok((adapter, guard))
1504}
1505
1506// =================================================================
1507// Pre-map pending_uses walker — SRD-35 Push D
1508// =================================================================
1509
1510/// Walk the freshly-pre-mapped scenario tree and seed the
1511/// pool's `pending_uses` counter for every phase that will
1512/// attach a pool-shareable adapter. Called once at session
1513/// bootstrap, after [`crate::executor::pre_map_tree`] returns
1514/// and before the executor begins running any phase.
1515///
1516/// For each phase node:
1517///   1. Determine the adapter name — `phase.adapter` override
1518///      (when the phase declares one) wins over the
1519///      session-level `default_driver`.
1520///   2. Resolve the driver name without instantiating
1521///      anything (`resolve_driver_name`).
1522///   3. Look up the matching `SharedDriverRegistration`. If
1523///      none registered, the phase rides the legacy
1524///      `PerPhase` path — its key is per-phase-unique, so
1525///      pre-map contributes zero to any shared `pending_uses`
1526///      and the legacy close-on-detach path handles teardown.
1527///   4. Compute the resource key from the session's
1528///      `merged_params` via the registration's pure
1529///      `resource_key` function. `resource_key` failures are
1530///      hard errors at session bootstrap — surfacing the same
1531///      misconfiguration that `attach_shared_adapter` would
1532///      have hit at runtime, just earlier and in one place.
1533///   5. Increment the per-key `pending_uses` counter via
1534///      [`ResourcePool::declare_pending_use`].
1535///
1536/// This is a pure read of the scenario tree + phase params;
1537/// no adapters are instantiated. The walker can run before
1538/// `pool.shutdown()` would otherwise be a no-op — and must
1539/// run before the first phase, since pending counts are
1540/// the close trigger that releases shared resources promptly
1541/// when their last user finishes.
1542pub fn pre_map_pending_uses(
1543    pool: &Arc<ResourcePool>,
1544    tree: &crate::scene_tree::SceneTree,
1545    phases: &std::collections::HashMap<String, nmbrs_workload::model::WorkloadPhase>,
1546    default_driver: &str,
1547    merged_params: &std::collections::HashMap<String, String>,
1548) -> Result<(), String> {
1549    for node in tree.dfs_phases() {
1550        // Phase-level adapter override wins over session
1551        // default. Mirrors `executor::run_phase` (line 1903)
1552        // so the predicted key matches the runtime key
1553        // exactly.
1554        let adapter = phases
1555            .get(&node.name)
1556            .and_then(|p| p.adapter.clone())
1557            .unwrap_or_else(|| default_driver.to_string());
1558
1559        // Drive the same resolution the executor uses
1560        // (`<adapter>driver` selector param + DriverImpl
1561        // ranking). `None` here means the adapter has no
1562        // `DriverImpl` registered AND no fallback; skip.
1563        let selector = format!("{adapter}driver");
1564        let Some(driver_name) =
1565            crate::adapter::resolve_driver_name(&adapter, &selector, merged_params)
1566        else {
1567            continue;
1568        };
1569
1570        // Only adapters that opted into pool sharing have a
1571        // `SharedDriverRegistration`. Legacy adapters fall
1572        // through; their `PerPhase` close-on-detach path is
1573        // unaffected.
1574        let Some(reg) = crate::adapter::find_shared_driver(&adapter, driver_name) else {
1575            continue;
1576        };
1577
1578        let key = (reg.resource_key)(merged_params).map_err(|e| {
1579            format!(
1580                "resource pool pre-map: phase '{}' adapter '{}' driver '{}': {}",
1581                node.name, adapter, driver_name, e,
1582            )
1583        })?;
1584        pool.declare_pending_use(key, default_policy_for(reg.share_capability));
1585    }
1586    Ok(())
1587}
1588
1589// =================================================================
1590// SRD-104 — resource-accessor bridge (pool ⇄ polydat)
1591// =================================================================
1592
1593/// The pool the process-global [`polydat::ResourceAccessor`] currently
1594/// resolves against, held as a `Weak` so the process-lifetime bridge never
1595/// keeps a finished session's pool (and its live network resources) alive.
1596/// Re-pointed each session by [`install_accessor`]; empty until the first
1597/// install.
1598static ACTIVE_POOL: Mutex<Weak<ResourcePool>> = Mutex::new(Weak::new());
1599
1600/// The bridge every nmbrs kernel tree resolves resources through
1601/// ([`pool_resources`]). It carries no state itself — it resolves against
1602/// the swappable [`ACTIVE_POOL`] — so a host process that runs several
1603/// sessions (TUI, `metrics watch`) re-points the pool without touching the
1604/// trees built before.
1605struct PoolAccessorView;
1606
1607impl polydat::ResourceAccessor for PoolAccessorView {
1608    fn lookup(&self, key: &str) -> Option<Arc<dyn Any + Send + Sync>> {
1609        let pool = ACTIVE_POOL
1610            .lock()
1611            .unwrap_or_else(|e| e.into_inner())
1612            .upgrade()?;
1613        pool.lookup_accessor_payload(key)
1614    }
1615}
1616
1617/// SRD-104 — point the resource bridge at `pool`. Called once at session
1618/// start where the [`ResourcePool`] is created.
1619///
1620/// Every kernel tree resolves through the same bridge ([`pool_resources`]),
1621/// and the pool it resolves against is swappable ([`ACTIVE_POOL`]), so a
1622/// second session in the same host process re-points the bridge at its own
1623/// pool rather than stranding on the first session's (now-closed) one.
1624/// Idempotent per session.
1625pub fn install_accessor(pool: &Arc<ResourcePool>) {
1626    *ACTIVE_POOL.lock().unwrap_or_else(|e| e.into_inner()) = Arc::downgrade(pool);
1627}
1628
1629/// A resource scope with the pool bridge installed: what every nmbrs
1630/// kernel tree is compiled with at its root
1631/// ([`crate::bindings::compile_scope_kernel`]). A kernel bound under the
1632/// tree joins its scope, so a node anywhere in it that captured
1633/// `ctx.resources()` resolves pool resources by key.
1634pub fn pool_resources() -> polydat::ResourceScope {
1635    polydat::ResourceScope::with_accessor(Arc::new(PoolAccessorView))
1636}
1637
1638// =================================================================
1639// Tests
1640// =================================================================
1641
1642#[cfg(test)]
1643mod tests {
1644    use super::*;
1645    use std::sync::atomic::AtomicU32;
1646
1647    /// Synthetic resource for unit tests — records every
1648    /// lifecycle call so assertions can verify the pool's
1649    /// invariants without touching real network resources.
1650    struct MockResource {
1651        key: ResourceKey,
1652        can_share: bool,
1653        saturated_after_attaches: u32,
1654        attach_count: AtomicU32,
1655        init_calls: AtomicU32,
1656        close_calls: AtomicU32,
1657        init_should_fail: bool,
1658    }
1659
1660    impl MockResource {
1661        fn new(key: ResourceKey) -> Arc<Self> {
1662            Arc::new(Self {
1663                key,
1664                can_share: true,
1665                saturated_after_attaches: u32::MAX,
1666                attach_count: AtomicU32::new(0),
1667                init_calls: AtomicU32::new(0),
1668                close_calls: AtomicU32::new(0),
1669                init_should_fail: false,
1670            })
1671        }
1672
1673        fn with_can_share(key: ResourceKey, can_share: bool) -> Arc<Self> {
1674            Arc::new(Self {
1675                key,
1676                can_share,
1677                saturated_after_attaches: u32::MAX,
1678                attach_count: AtomicU32::new(0),
1679                init_calls: AtomicU32::new(0),
1680                close_calls: AtomicU32::new(0),
1681                init_should_fail: false,
1682            })
1683        }
1684
1685        fn with_init_failure(key: ResourceKey) -> Arc<Self> {
1686            Arc::new(Self {
1687                key,
1688                can_share: true,
1689                saturated_after_attaches: u32::MAX,
1690                attach_count: AtomicU32::new(0),
1691                init_calls: AtomicU32::new(0),
1692                close_calls: AtomicU32::new(0),
1693                init_should_fail: true,
1694            })
1695        }
1696    }
1697
1698    impl SharedResource for MockResource {
1699        fn resource_key(&self) -> &ResourceKey {
1700            &self.key
1701        }
1702        fn can_share(&self) -> bool {
1703            self.can_share
1704        }
1705        fn can_support_more_load(&self) -> bool {
1706            // True (has capacity) until `attach_count`
1707            // crosses the saturation threshold; then
1708            // false (saturated). `u32::MAX` as the
1709            // threshold ⇒ always has capacity, matching
1710            // the trait default.
1711            self.attach_count.load(Ordering::Acquire) < self.saturated_after_attaches
1712        }
1713
1714        fn init(&self) -> ResourceFuture<'_, Result<(), String>> {
1715            self.init_calls.fetch_add(1, Ordering::AcqRel);
1716            let fail = self.init_should_fail;
1717            Box::pin(async move {
1718                if fail {
1719                    Err("simulated init failure".into())
1720                } else {
1721                    Ok(())
1722                }
1723            })
1724        }
1725
1726        fn close(self: Arc<Self>) -> ResourceFuture<'static, Result<(), String>> {
1727            self.close_calls.fetch_add(1, Ordering::AcqRel);
1728            Box::pin(async { Ok(()) })
1729        }
1730    }
1731
1732    fn key(adapter: &str, vendor: &str) -> ResourceKey {
1733        ResourceKey::new(adapter).with("vendor", vendor)
1734    }
1735
1736    #[test]
1737    fn key_value_equality_is_field_order_independent() {
1738        let a = ResourceKey::new("cql")
1739            .with("hosts", "h1")
1740            .with("port", "9042");
1741        let b = ResourceKey::new("cql")
1742            .with("port", "9042")
1743            .with("hosts", "h1");
1744        assert_eq!(a, b);
1745        assert_eq!(a.fmt_for_log(), b.fmt_for_log());
1746    }
1747
1748    #[test]
1749    fn key_log_format_redacts_password() {
1750        let k = ResourceKey::new("cql")
1751            .with("hosts", "h1")
1752            .with("password", "hunter2");
1753        let s = k.fmt_for_log();
1754        assert!(s.contains("hosts=h1"));
1755        assert!(s.contains("password=***"), "got: {s}");
1756        assert!(!s.contains("hunter2"));
1757    }
1758
1759    #[test]
1760    fn policy_parse_accepts_canonical_and_underscore_forms() {
1761        assert_eq!(
1762            ResourceSharePolicy::parse("shared").unwrap(),
1763            ResourceSharePolicy::Shared
1764        );
1765        assert_eq!(
1766            ResourceSharePolicy::parse("per-phase").unwrap(),
1767            ResourceSharePolicy::PerPhase
1768        );
1769        assert_eq!(
1770            ResourceSharePolicy::parse("per_phase").unwrap(),
1771            ResourceSharePolicy::PerPhase
1772        );
1773        assert_eq!(
1774            ResourceSharePolicy::parse("PER-FIBER").unwrap(),
1775            ResourceSharePolicy::PerFiber
1776        );
1777        ResourceSharePolicy::parse("nope").unwrap_err();
1778    }
1779
1780    #[test]
1781    fn capability_floor_orders_isolation_correctly() {
1782        // Shared is the lowest (most-shared); PerFiber is
1783        // the highest (most-isolated). User-policy must be
1784        // ≥ floor for capability.
1785        assert!(capability_floor(ShareCapability::Shared) <= ResourceSharePolicy::Shared);
1786        assert!(capability_floor(ShareCapability::PerPhase) > ResourceSharePolicy::Shared);
1787    }
1788
1789    #[tokio::test]
1790    async fn shared_resource_survives_detach_until_more_phases_predicted() {
1791        // SRD-35 Push D: the pre-map walker seeds
1792        // `pending_uses` with the predicted total attach
1793        // count. While `pending > 0`, detach MUST NOT close
1794        // — there's still a future phase that will use the
1795        // resource. The 55-phase open/close storm in the
1796        // operator's session.log was a Push A regression
1797        // where pending wasn't seeded; the resource closed
1798        // on the first detach and re-initialised every time.
1799        let pool = Arc::new(ResourcePool::new());
1800        let mock = MockResource::new(key("cql", "cassandra-cpp"));
1801        let mock_for_factory = Arc::clone(&mock);
1802
1803        // Pre-map says: 2 phases will attach this key.
1804        pool.declare_pending_use(mock.resource_key().clone(), ResourceSharePolicy::Shared);
1805        pool.declare_pending_use(mock.resource_key().clone(), ResourceSharePolicy::Shared);
1806        assert_eq!(pool.pending_uses_for(mock.resource_key()), Some(2));
1807
1808        // Phase A: attach + detach. pending drops to 1, so
1809        // close must NOT fire — phase B is still expected.
1810        let guard = attach(
1811            &pool,
1812            mock.resource_key().clone(),
1813            ResourceSharePolicy::Shared,
1814            "phase_A",
1815            ResourceFactory::new(
1816                move || async move { Ok(mock_for_factory as Arc<dyn SharedResource>) },
1817            ),
1818        )
1819        .await
1820        .unwrap();
1821
1822        assert_eq!(mock.init_calls.load(Ordering::Acquire), 1);
1823        guard.detach().await.unwrap();
1824        assert_eq!(
1825            mock.close_calls.load(Ordering::Acquire),
1826            0,
1827            "Shared-policy detach MUST NOT close while another \
1828             phase is still predicted to attach"
1829        );
1830        assert_eq!(pool.pending_uses_for(mock.resource_key()), Some(1));
1831        assert_eq!(pool.live_entries(), 1);
1832
1833        // Phase B: attach + detach. pending drops to 0,
1834        // live drops to 0 — close fires NOW, not at session
1835        // shutdown. Promptly releases the resource the
1836        // moment its last user is done.
1837        let mock_for_factory2 = Arc::clone(&mock);
1838        let guard2 = attach(
1839            &pool,
1840            mock.resource_key().clone(),
1841            ResourceSharePolicy::Shared,
1842            "phase_B",
1843            ResourceFactory::new(move || async move {
1844                Ok(mock_for_factory2 as Arc<dyn SharedResource>)
1845            }),
1846        )
1847        .await
1848        .unwrap();
1849        guard2.detach().await.unwrap();
1850        assert_eq!(
1851            mock.close_calls.load(Ordering::Acquire),
1852            1,
1853            "close MUST fire when pending hits 0 AND live hits 0 — \
1854             that's the SRD-35 Push D close-on-zero contract"
1855        );
1856        assert_eq!(pool.live_entries(), 0, "entry removed once closed");
1857
1858        // Pool shutdown is now a no-op for this key.
1859        pool.shutdown().await;
1860        assert_eq!(
1861            mock.close_calls.load(Ordering::Acquire),
1862            1,
1863            "shutdown() must NOT double-close an entry"
1864        );
1865    }
1866
1867    #[tokio::test]
1868    async fn shared_resource_overpredicted_pending_holds_until_shutdown() {
1869        // Conservative fallback: if the pre-map walker
1870        // over-predicts (predicts N phases but only N-1
1871        // actually attach), the entry stays alive until
1872        // shutdown — better than closing prematurely. The
1873        // residual `pending > 0` at session end is also
1874        // visible on the `close.started reason=session-end`
1875        // line so an operator can spot pre-map bugs.
1876        let pool = Arc::new(ResourcePool::new());
1877        let mock = MockResource::new(key("cql", "cassandra-cpp"));
1878        let mock_for_factory = Arc::clone(&mock);
1879
1880        // Pre-map says 3 phases, only 1 actually runs.
1881        for _ in 0..3 {
1882            pool.declare_pending_use(mock.resource_key().clone(), ResourceSharePolicy::Shared);
1883        }
1884        let guard = attach(
1885            &pool,
1886            mock.resource_key().clone(),
1887            ResourceSharePolicy::Shared,
1888            "phase_A",
1889            ResourceFactory::new(
1890                move || async move { Ok(mock_for_factory as Arc<dyn SharedResource>) },
1891            ),
1892        )
1893        .await
1894        .unwrap();
1895        guard.detach().await.unwrap();
1896
1897        // pending = 2 after one attach — close MUST NOT
1898        // fire mid-run.
1899        assert_eq!(pool.pending_uses_for(mock.resource_key()), Some(2));
1900        assert_eq!(mock.close_calls.load(Ordering::Acquire), 0);
1901        assert_eq!(pool.live_entries(), 1);
1902
1903        pool.shutdown().await;
1904        assert_eq!(
1905            mock.close_calls.load(Ordering::Acquire),
1906            1,
1907            "shutdown() drains residual entries with pending > 0"
1908        );
1909    }
1910
1911    #[tokio::test]
1912    async fn per_phase_resource_closes_on_detach() {
1913        // Counter-test: under `PerPhase` policy, the
1914        // legacy-shim semantics still apply — close fires
1915        // on the detach that drains refcount, no waiting
1916        // for shutdown. The synthetic mock has
1917        // `can_share=true`, so we explicitly use PerPhase
1918        // to exercise this branch.
1919        let pool = Arc::new(ResourcePool::new());
1920        let mock = MockResource::new(key("legacy", "synthetic"));
1921        let mock_for_factory = Arc::clone(&mock);
1922
1923        let guard = attach(
1924            &pool,
1925            mock.resource_key().clone(),
1926            ResourceSharePolicy::PerPhase,
1927            "phase_A",
1928            ResourceFactory::new(
1929                move || async move { Ok(mock_for_factory as Arc<dyn SharedResource>) },
1930            ),
1931        )
1932        .await
1933        .unwrap();
1934
1935        guard.detach().await.unwrap();
1936        assert_eq!(
1937            mock.close_calls.load(Ordering::Acquire),
1938            1,
1939            "PerPhase detach MUST close immediately"
1940        );
1941        assert_eq!(pool.live_entries(), 0);
1942    }
1943
1944    #[tokio::test]
1945    async fn init_failure_poisons_the_entry() {
1946        let pool = Arc::new(ResourcePool::new());
1947        let mock = MockResource::with_init_failure(key("cql", "cassandra-cpp"));
1948        let mock_for_factory = Arc::clone(&mock);
1949
1950        let result = attach(
1951            &pool,
1952            mock.resource_key().clone(),
1953            ResourceSharePolicy::Shared,
1954            "phase_A",
1955            ResourceFactory::new(
1956                move || async move { Ok(mock_for_factory as Arc<dyn SharedResource>) },
1957            ),
1958        )
1959        .await;
1960
1961        let err = result.expect_err("init failure must propagate");
1962        assert!(err.contains("simulated init failure"), "got: {err}");
1963    }
1964
1965    #[tokio::test]
1966    async fn shared_policy_against_can_share_false_resource_errors() {
1967        let pool = Arc::new(ResourcePool::new());
1968        let mock = MockResource::with_can_share(key("legacy", "x"), false);
1969        let mock_for_factory = Arc::clone(&mock);
1970
1971        let result = attach(
1972            &pool,
1973            mock.resource_key().clone(),
1974            ResourceSharePolicy::Shared,
1975            "phase_A",
1976            ResourceFactory::new(
1977                move || async move { Ok(mock_for_factory as Arc<dyn SharedResource>) },
1978            ),
1979        )
1980        .await;
1981
1982        let err = result.expect_err("shared policy on can_share()=false resource must be rejected");
1983        assert!(err.contains("can_share()=false"), "got: {err}");
1984        assert!(
1985            err.contains("per-phase") || err.contains("per-fiber"),
1986            "error must point at the available policies, got: {err}"
1987        );
1988    }
1989
1990    #[tokio::test]
1991    async fn per_phase_policy_with_can_share_false_works() {
1992        let pool = Arc::new(ResourcePool::new());
1993        let mock = MockResource::with_can_share(key("legacy", "x"), false);
1994        let mock_for_factory = Arc::clone(&mock);
1995
1996        let guard = attach(
1997            &pool,
1998            mock.resource_key().clone(),
1999            ResourceSharePolicy::PerPhase,
2000            "phase_A",
2001            ResourceFactory::new(
2002                move || async move { Ok(mock_for_factory as Arc<dyn SharedResource>) },
2003            ),
2004        )
2005        .await
2006        .expect("PerPhase must accept can_share()=false");
2007
2008        // PerPhase + can_share=false is exactly the
2009        // legacy-adapter shape; works without elevation.
2010        guard.detach().await.unwrap();
2011    }
2012
2013    #[tokio::test]
2014    async fn shared_policy_reuses_one_instance_across_attaches() {
2015        let pool = Arc::new(ResourcePool::new());
2016        let mock = MockResource::new(key("cql", "cassandra-cpp"));
2017        let mock_for_factory = Arc::clone(&mock);
2018        let factory_called = Arc::new(AtomicU32::new(0));
2019        let factory_called_clone = Arc::clone(&factory_called);
2020
2021        // SRD-35 Push D: pre-map predicts 3 phases will
2022        // attach this key, but the test only exercises 2.
2023        // The residual pending = 1 keeps the entry alive
2024        // until shutdown — this lets the test assert "Shared
2025        // policy holds across phases" without timing-coupling
2026        // the close trigger to detach order.
2027        for _ in 0..3 {
2028            pool.declare_pending_use(mock.resource_key().clone(), ResourceSharePolicy::Shared);
2029        }
2030
2031        let g1 = attach(
2032            &pool,
2033            mock.resource_key().clone(),
2034            ResourceSharePolicy::Shared,
2035            "phase_A",
2036            ResourceFactory::new(move || {
2037                factory_called_clone.fetch_add(1, Ordering::AcqRel);
2038                let m = Arc::clone(&mock_for_factory);
2039                async move { Ok(m as Arc<dyn SharedResource>) }
2040            }),
2041        )
2042        .await
2043        .unwrap();
2044
2045        // Second attach for the SAME key MUST reuse the
2046        // existing instance — factory MUST NOT be called
2047        // again. Push A's pool reuses by key+generation
2048        // even though the legacy adapter shim doesn't
2049        // typically exercise this; the synthetic test
2050        // proves the contract end-to-end.
2051        let mock_for_factory2 = Arc::clone(&mock);
2052        let factory_called_clone2 = Arc::clone(&factory_called);
2053        let g2 = attach(
2054            &pool,
2055            mock.resource_key().clone(),
2056            ResourceSharePolicy::Shared,
2057            "phase_B",
2058            ResourceFactory::new(move || {
2059                factory_called_clone2.fetch_add(1, Ordering::AcqRel);
2060                let m = Arc::clone(&mock_for_factory2);
2061                async move { Ok(m as Arc<dyn SharedResource>) }
2062            }),
2063        )
2064        .await
2065        .unwrap();
2066
2067        assert_eq!(
2068            factory_called.load(Ordering::Acquire),
2069            1,
2070            "factory must not be called twice for the same key under Shared policy"
2071        );
2072        assert_eq!(
2073            mock.init_calls.load(Ordering::Acquire),
2074            1,
2075            "init() must fire exactly once"
2076        );
2077        assert_eq!(pool.total_attaches(), 2);
2078        assert_eq!(pool.live_entries(), 1);
2079
2080        // Detach both guards. Pre-map predicted 3 phases;
2081        // only 2 attached, so pending stays > 0 and the
2082        // entry survives until session shutdown. The
2083        // session.log bug this contract locks down was the
2084        // resource closing on first detach, defeating
2085        // sharing across the next phase's attach.
2086        g1.detach().await.unwrap();
2087        assert_eq!(
2088            mock.close_calls.load(Ordering::Acquire),
2089            0,
2090            "close MUST NOT fire while another guard is live"
2091        );
2092        g2.detach().await.unwrap();
2093        assert_eq!(
2094            mock.close_calls.load(Ordering::Acquire),
2095            0,
2096            "close MUST NOT fire while pending uses remain — \
2097             the resource stays alive for predicted future phases"
2098        );
2099        assert_eq!(pool.live_entries(), 1, "entry stays cached across phases");
2100
2101        // The session-end shutdown drains the residual.
2102        pool.shutdown().await;
2103        assert_eq!(
2104            mock.close_calls.load(Ordering::Acquire),
2105            1,
2106            "shutdown() drains the entry"
2107        );
2108        assert_eq!(pool.live_entries(), 0);
2109    }
2110
2111    /// Minimal `DriverAdapter` stub used by the legacy-
2112    /// wrapper test below. `map_op` returns an error
2113    /// because the test never exercises the dispenser
2114    /// path — we only need a `DriverAdapter` value to
2115    /// hand to `LegacyAdapterResource::new`.
2116    struct LegacyDummy;
2117    impl crate::adapter::DriverAdapter for LegacyDummy {
2118        fn name(&self) -> &str {
2119            "legacy_dummy"
2120        }
2121        fn map_op<'a>(
2122            &'a self,
2123            _template: &'a nmbrs_workload::model::ParsedOp,
2124            _parent: std::sync::Arc<dyn polydat::Kernel>,
2125        ) -> std::pin::Pin<
2126            Box<
2127                dyn std::future::Future<
2128                        Output = Result<Box<dyn crate::adapter::OpDispenser>, String>,
2129                    > + Send
2130                    + 'a,
2131            >,
2132        > {
2133            Box::pin(async move { Err("dummy".into()) })
2134        }
2135    }
2136
2137    #[test]
2138    fn capacity_decline_at_quiescence_is_caught_as_driver_bug() {
2139        // SRD-35 validity rule: a `can_support_more_load()`
2140        // body that returns FALSE while `live_attaches ==
2141        // 0` is reading historical state, not active load
2142        // (e.g., a ring buffer that filled during a
2143        // previous phase). The pool's guard MUST treat
2144        // this as a driver bug — log a Warn and route the
2145        // attach to the existing idle generation instead
2146        // of wasting a spawn. Real saturation (decline +
2147        // live > 0) MUST trigger the spawn unchanged.
2148        struct Reports {
2149            key: ResourceKey,
2150            has_capacity: bool,
2151        }
2152        impl SharedResource for Reports {
2153            fn resource_key(&self) -> &ResourceKey {
2154                &self.key
2155            }
2156            fn can_support_more_load(&self) -> bool {
2157                self.has_capacity
2158            }
2159        }
2160
2161        let entry = Entry::new(ResourceKey::new("buggy"), ResourceSharePolicy::Shared, 0);
2162
2163        // Has capacity → never spawn, regardless of load.
2164        let calm = Reports {
2165            key: ResourceKey::new("buggy"),
2166            has_capacity: true,
2167        };
2168        assert!(
2169            !needs_sibling_spawn(&entry, &calm),
2170            "no spawn when the resource has capacity"
2171        );
2172
2173        // No capacity at quiescence → driver bug, suppressed.
2174        let buggy = Reports {
2175            key: ResourceKey::new("buggy"),
2176            has_capacity: false,
2177        };
2178        assert_eq!(entry.live_attaches.load(Ordering::Acquire), 0);
2179        assert!(
2180            !needs_sibling_spawn(&entry, &buggy),
2181            "quiescent decline must be suppressed (driver bug)"
2182        );
2183
2184        // No capacity + actual load → genuine saturation, spawn.
2185        entry.live_attaches.fetch_add(3, Ordering::AcqRel);
2186        assert!(
2187            needs_sibling_spawn(&entry, &buggy),
2188            "decline under live load is genuine saturation; spawn sibling"
2189        );
2190
2191        // Has capacity + load → still no spawn.
2192        assert!(
2193            !needs_sibling_spawn(&entry, &calm),
2194            "no spawn when there is capacity, even under load"
2195        );
2196    }
2197
2198    #[tokio::test]
2199    async fn shared_adapter_resource_returns_one_instance_for_one_key() {
2200        // Push B contract: under `Shared` policy, two
2201        // attaches against the same `ResourceKey` must
2202        // produce two guards backed by the SAME wrapped
2203        // `Arc<dyn DriverAdapter>` — the pool's caching is
2204        // what stops 50 phases from opening 50 cluster
2205        // sessions.
2206        let pool = Arc::new(ResourcePool::new());
2207        let adapter: Arc<dyn DriverAdapter> = Arc::new(LegacyDummy);
2208        let key = ResourceKey::new("cql").with("driver", "synthetic");
2209        let factory_call_count = Arc::new(AtomicU32::new(0));
2210
2211        // SRD-35 Push D: simulate the pre-map walker
2212        // declaring 3 future phases. The two attaches below
2213        // consume 2 of those 3, leaving pending = 1 so the
2214        // entry survives until shutdown — the assertion
2215        // shape this test was designed to verify.
2216        for _ in 0..3 {
2217            pool.declare_pending_use(key.clone(), ResourceSharePolicy::Shared);
2218        }
2219
2220        let factory_count_a = Arc::clone(&factory_call_count);
2221        let adapter_a = Arc::clone(&adapter);
2222        let (got_a, guard_a) =
2223            attach_shared_adapter(&pool, "cql", "phase_A", key.clone(), move || {
2224                factory_count_a.fetch_add(1, Ordering::AcqRel);
2225                let a = Arc::clone(&adapter_a);
2226                async move { Ok(a) }
2227            })
2228            .await
2229            .unwrap();
2230
2231        let factory_count_b = Arc::clone(&factory_call_count);
2232        let adapter_b = Arc::clone(&adapter);
2233        let (got_b, guard_b) =
2234            attach_shared_adapter(&pool, "cql", "phase_B", key.clone(), move || {
2235                factory_count_b.fetch_add(1, Ordering::AcqRel);
2236                let a = Arc::clone(&adapter_b);
2237                async move { Ok(a) }
2238            })
2239            .await
2240            .unwrap();
2241
2242        // The factory MUST be called exactly once even
2243        // though two phases attached.
2244        assert_eq!(
2245            factory_call_count.load(Ordering::Acquire),
2246            1,
2247            "Shared policy MUST cache the factory output"
2248        );
2249        assert!(
2250            Arc::ptr_eq(&got_a, &got_b),
2251            "both attaches MUST surface the SAME Arc<dyn DriverAdapter>"
2252        );
2253        assert_eq!(pool.live_entries(), 1);
2254        assert_eq!(pool.total_attaches(), 2);
2255
2256        guard_a.detach().await.unwrap();
2257        assert_eq!(
2258            pool.live_entries(),
2259            1,
2260            "entry stays live while any guard holds it"
2261        );
2262        guard_b.detach().await.unwrap();
2263        assert_eq!(
2264            pool.live_entries(),
2265            1,
2266            "Shared-policy entry stays cached across phases — \
2267             session shutdown is the close trigger, not last detach"
2268        );
2269        pool.shutdown().await;
2270        assert_eq!(pool.live_entries(), 0, "shutdown drains the cached entry");
2271    }
2272
2273    #[test]
2274    fn declare_then_complete_pending_use_drains_to_zero() {
2275        // SRD-35 Push D counter lifecycle: each
2276        // `declare_pending_use` call increments per-key;
2277        // each `complete_pending_use` decrements (saturating
2278        // at 0) and reports whether the entry is now
2279        // eligible for close (`pending == 0 && live == 0`).
2280        // The pool's close trigger sits on `(pending, live)`
2281        // both being zero — this test exercises the
2282        // counter-only side without a real attach to keep
2283        // the assertion shape tight.
2284        let pool = Arc::new(ResourcePool::new());
2285        let k = ResourceKey::new("cql").with("hosts", "h1");
2286        assert_eq!(
2287            pool.pending_uses_for(&k),
2288            None,
2289            "no entry exists before declare"
2290        );
2291
2292        // Pre-map walker declares: 3 phases will use this key.
2293        pool.declare_pending_use(k.clone(), ResourceSharePolicy::Shared);
2294        pool.declare_pending_use(k.clone(), ResourceSharePolicy::Shared);
2295        pool.declare_pending_use(k.clone(), ResourceSharePolicy::Shared);
2296        assert_eq!(pool.pending_uses_for(&k), Some(3));
2297        assert_eq!(
2298            pool.live_entries(),
2299            1,
2300            "declare_pending_use creates the entry on first call"
2301        );
2302
2303        // Drain phases 1 and 2: not yet eligible for close.
2304        let eligible = pool.complete_pending_use(&k);
2305        assert!(!eligible, "pending == 2; not eligible yet");
2306        assert_eq!(pool.pending_uses_for(&k), Some(2));
2307
2308        let eligible = pool.complete_pending_use(&k);
2309        assert!(!eligible, "pending == 1; not eligible yet");
2310        assert_eq!(pool.pending_uses_for(&k), Some(1));
2311
2312        // Phase 3 drains: pending hits 0, live is 0 (no
2313        // real attach in this test) — eligible for close.
2314        let eligible = pool.complete_pending_use(&k);
2315        assert!(
2316            eligible,
2317            "pending == 0 && live == 0 — entry is eligible for close"
2318        );
2319        assert_eq!(pool.pending_uses_for(&k), Some(0));
2320
2321        // Saturation: extra completes are no-ops, not
2322        // underflow. Stays "eligible for close" because the
2323        // invariant still holds.
2324        let eligible = pool.complete_pending_use(&k);
2325        assert!(
2326            eligible,
2327            "saturating decrement: extra completes don't underflow"
2328        );
2329        assert_eq!(pool.pending_uses_for(&k), Some(0));
2330    }
2331
2332    #[tokio::test]
2333    async fn declare_pending_use_paired_with_attach_closes_promptly() {
2334        // SRD-35 Push D end-to-end: pre-map seeds pending,
2335        // attach decrements eagerly, detach observes the
2336        // (pending == 0, live == 0) invariant and closes.
2337        // Verifies the close trigger fires the moment the
2338        // last predicted phase finishes, not at session end.
2339        let pool = Arc::new(ResourcePool::new());
2340        let mock = MockResource::new(key("cql", "cassandra-cpp"));
2341
2342        // Pre-map: exactly 1 phase will use this key.
2343        pool.declare_pending_use(mock.resource_key().clone(), ResourceSharePolicy::Shared);
2344
2345        let mock_for_factory = Arc::clone(&mock);
2346        let guard = attach(
2347            &pool,
2348            mock.resource_key().clone(),
2349            ResourceSharePolicy::Shared,
2350            "phase_A",
2351            ResourceFactory::new(
2352                move || async move { Ok(mock_for_factory as Arc<dyn SharedResource>) },
2353            ),
2354        )
2355        .await
2356        .unwrap();
2357        // After attach, pending dropped from 1 to 0.
2358        assert_eq!(pool.pending_uses_for(mock.resource_key()), Some(0));
2359
2360        guard.detach().await.unwrap();
2361        // Detach drove live to 0; pending was already 0;
2362        // close fired without waiting for shutdown.
2363        assert_eq!(
2364            mock.close_calls.load(Ordering::Acquire),
2365            1,
2366            "Push D: close fires the moment the last predicted \
2367             phase detaches — not at session shutdown"
2368        );
2369        assert_eq!(pool.live_entries(), 0);
2370    }
2371
2372    #[test]
2373    fn render_key_is_identity_preserving_unlike_log_format() {
2374        // SRD-104: `render_key()` is the accessor fingerprint — it must be
2375        // identity-preserving (NOT redact secrets), because two keys that
2376        // differ only in an identity-bearing field would otherwise collide
2377        // in the accessor. `fmt_for_log()` redacts (safe for logs) and so
2378        // deliberately collides; `render_key()` must not.
2379        let a = ResourceKey::new("cql")
2380            .with("hosts", "h1")
2381            .with("password", "p1");
2382        let b = ResourceKey::new("cql")
2383            .with("hosts", "h1")
2384            .with("password", "p2");
2385        assert_eq!(
2386            a.fmt_for_log(),
2387            b.fmt_for_log(),
2388            "log format redacts password → collides (expected)"
2389        );
2390        assert_ne!(
2391            a.render_key(),
2392            b.render_key(),
2393            "render_key must distinguish identity-bearing password"
2394        );
2395        assert!(
2396            !a.render_key().contains("p1") && a.render_key().contains("password=sha256:"),
2397            "the fingerprint must carry a digest of the secret, never the cleartext; got: {}",
2398            a.render_key()
2399        );
2400        // BTreeMap makes field order irrelevant to the rendering.
2401        let c = ResourceKey::new("cql")
2402            .with("password", "p1")
2403            .with("hosts", "h1");
2404        assert_eq!(
2405            a.render_key(),
2406            c.render_key(),
2407            "render_key is field-order independent"
2408        );
2409    }
2410
2411    #[tokio::test]
2412    async fn accessor_payload_populates_and_looks_up_by_render_key() {
2413        // SRD-104 pool-side round-trip: a resource that opts into an
2414        // accessor payload has it captured onto its entry at init, and the
2415        // pool resolves it by the key's canonical `render_key()` rendering
2416        // — the same string a consumer node builds from the fingerprint in
2417        // scope. Exercises the pool half of the bridge directly, without the
2418        // process-global (which would race across parallel tests).
2419        struct WithPayload {
2420            key: ResourceKey,
2421        }
2422        impl SharedResource for WithPayload {
2423            fn resource_key(&self) -> &ResourceKey {
2424                &self.key
2425            }
2426            fn accessor_payload(&self) -> Option<Arc<dyn Any + Send + Sync>> {
2427                Some(Arc::new(4242u64) as Arc<dyn Any + Send + Sync>)
2428            }
2429        }
2430
2431        let pool = Arc::new(ResourcePool::new());
2432        let key = ResourceKey::new("cql")
2433            .with("hosts", "h1")
2434            .with("keyspace", "ks");
2435        let rendered = key.render_key();
2436
2437        // No entry yet → miss.
2438        assert!(
2439            pool.lookup_accessor_payload(&rendered).is_none(),
2440            "no live entry ⇒ None"
2441        );
2442
2443        let key_for_factory = key.clone();
2444        let guard = attach(
2445            &pool,
2446            key.clone(),
2447            ResourceSharePolicy::Shared,
2448            "phase_A",
2449            ResourceFactory::new(move || async move {
2450                Ok(Arc::new(WithPayload {
2451                    key: key_for_factory,
2452                }) as Arc<dyn SharedResource>)
2453            }),
2454        )
2455        .await
2456        .unwrap();
2457
2458        // Hit by the canonical rendering; downcast to the concrete payload.
2459        let payload = pool
2460            .lookup_accessor_payload(&rendered)
2461            .expect("payload present after successful init");
2462        let n = payload
2463            .downcast_ref::<u64>()
2464            .expect("payload downcasts to u64");
2465        assert_eq!(*n, 4242);
2466
2467        // Detach with no pending uses drains the entry → miss again.
2468        guard.detach().await.unwrap();
2469        assert!(
2470            pool.lookup_accessor_payload(&rendered).is_none(),
2471            "closed entry ⇒ None"
2472        );
2473    }
2474
2475    #[tokio::test]
2476    async fn default_resource_has_no_accessor_payload() {
2477        // A resource that doesn't override `accessor_payload` (the trait
2478        // default) contributes no payload — the miss policy of the bridge.
2479        let pool = Arc::new(ResourcePool::new());
2480        let mock = MockResource::new(key("cql", "cassandra-cpp"));
2481        let rendered = mock.resource_key().render_key();
2482        let mock_for_factory = Arc::clone(&mock);
2483        let guard = attach(
2484            &pool,
2485            mock.resource_key().clone(),
2486            ResourceSharePolicy::Shared,
2487            "phase_A",
2488            ResourceFactory::new(
2489                move || async move { Ok(mock_for_factory as Arc<dyn SharedResource>) },
2490            ),
2491        )
2492        .await
2493        .unwrap();
2494        assert!(
2495            pool.lookup_accessor_payload(&rendered).is_none(),
2496            "default accessor_payload() ⇒ None even for a live entry"
2497        );
2498        guard.detach().await.unwrap();
2499    }
2500
2501    #[tokio::test]
2502    async fn legacy_adapter_resource_declares_non_shareable() {
2503        // Sanity: the bridge type used by Push A's
2504        // executor wiring declares itself non-shareable,
2505        // forcing PerPhase semantics for any adapter that
2506        // hasn't migrated yet.
2507        let key = ResourceKey::new("legacy");
2508        let adapter: Arc<dyn DriverAdapter> = Arc::new(LegacyDummy);
2509        let wrapped = Arc::new(LegacyAdapterResource::new(key, adapter));
2510        assert!(
2511            !wrapped.can_share(),
2512            "LegacyAdapterResource MUST declare can_share()=false"
2513        );
2514    }
2515}