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}