use std::any::Any;
use std::collections::{BTreeMap, HashMap};
use std::future::Future;
use std::pin::Pin;
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use std::sync::{Arc, Mutex, Weak};
use std::time::Instant;
use crate::adapter::DriverAdapter;
use crate::observer::LogLevel;
#[derive(Clone, Debug, PartialEq, Eq, Hash, Default)]
pub struct ResourceKey {
pub adapter: String,
pub fields: BTreeMap<String, String>,
}
impl ResourceKey {
pub fn new(adapter: impl Into<String>) -> Self {
Self {
adapter: adapter.into(),
fields: BTreeMap::new(),
}
}
pub fn with(mut self, key: impl Into<String>, value: impl Into<String>) -> Self {
self.fields.insert(key.into(), value.into());
self
}
pub fn fmt_for_log(&self) -> String {
let mut s = String::with_capacity(64);
s.push_str(&self.adapter);
s.push('{');
let mut first = true;
for (k, v) in &self.fields {
if !first {
s.push(',');
}
first = false;
s.push_str(k);
s.push('=');
if k.starts_with("_secret_") || k == "password" {
s.push_str("***");
} else {
s.push_str(v);
}
}
s.push('}');
s
}
pub fn render_key(&self) -> String {
let mut s = String::with_capacity(64);
s.push_str(&self.adapter);
s.push('{');
let mut first = true;
for (k, v) in &self.fields {
if !first {
s.push(',');
}
first = false;
s.push_str(k);
s.push('=');
if k.starts_with("_secret_") || k == "password" {
use sha2::{Digest, Sha256};
let digest = Sha256::digest(v.as_bytes());
s.push_str("sha256:");
for byte in &digest[..8] {
s.push_str(&format!("{byte:02x}"));
}
} else {
s.push_str(v);
}
}
s.push('}');
s
}
}
#[derive(Copy, Clone, Debug, PartialEq, Eq, PartialOrd, Ord, Hash)]
pub enum ShareCapability {
Shared,
PerScenario,
PerPhase,
PerFiber,
}
#[derive(Copy, Clone, Debug, PartialEq, Eq, PartialOrd, Ord, Hash)]
pub enum ResourceSharePolicy {
Shared,
PerScenario,
PerPhase,
PerFiber,
}
impl ResourceSharePolicy {
pub fn parse(s: &str) -> Result<Self, String> {
match s.trim().to_ascii_lowercase().as_str() {
"shared" => Ok(Self::Shared),
"per-scenario" | "per_scenario" => Ok(Self::PerScenario),
"per-phase" | "per_phase" => Ok(Self::PerPhase),
"per-fiber" | "per_fiber" => Ok(Self::PerFiber),
other => Err(format!(
"unknown resource-share policy '{other}' \
(expected: shared, per-scenario, per-phase, per-fiber)"
)),
}
}
}
impl std::fmt::Display for ResourceSharePolicy {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
let s = match self {
Self::Shared => "shared",
Self::PerScenario => "per-scenario",
Self::PerPhase => "per-phase",
Self::PerFiber => "per-fiber",
};
f.write_str(s)
}
}
pub fn default_policy_for(cap: ShareCapability) -> ResourceSharePolicy {
match cap {
ShareCapability::Shared => ResourceSharePolicy::Shared,
ShareCapability::PerScenario => ResourceSharePolicy::PerScenario,
ShareCapability::PerPhase => ResourceSharePolicy::PerPhase,
ShareCapability::PerFiber => ResourceSharePolicy::PerFiber,
}
}
pub fn capability_floor(cap: ShareCapability) -> ResourceSharePolicy {
match cap {
ShareCapability::Shared => ResourceSharePolicy::Shared,
ShareCapability::PerScenario => ResourceSharePolicy::PerScenario,
ShareCapability::PerPhase => ResourceSharePolicy::PerPhase,
ShareCapability::PerFiber => ResourceSharePolicy::PerFiber,
}
}
pub type ResourceFuture<'a, T> = Pin<Box<dyn Future<Output = T> + Send + 'a>>;
pub trait SharedResource: Send + Sync + 'static {
fn resource_key(&self) -> &ResourceKey;
fn can_share(&self) -> bool {
true
}
fn can_support_more_load(&self) -> bool {
true
}
fn init(&self) -> ResourceFuture<'_, Result<(), String>> {
Box::pin(async { Ok(()) })
}
fn close(self: Arc<Self>) -> ResourceFuture<'static, Result<(), String>> {
Box::pin(async { Ok(()) })
}
fn as_legacy_adapter(&self) -> Option<Arc<dyn DriverAdapter>> {
None
}
fn accessor_payload(&self) -> Option<Arc<dyn Any + Send + Sync>> {
None
}
}
struct Entry {
key: ResourceKey,
policy: ResourceSharePolicy,
generation: usize,
resource: Mutex<Option<Arc<dyn SharedResource>>>,
poisoned: AtomicBool,
init_error: Mutex<Option<String>>,
pending_uses: AtomicUsize,
live_attaches: AtomicUsize,
init_started_at: Mutex<Option<Instant>>,
accessor_payload: Mutex<Option<Arc<dyn Any + Send + Sync>>>,
}
impl Entry {
fn new(key: ResourceKey, policy: ResourceSharePolicy, generation: usize) -> Self {
Self {
key,
policy,
generation,
resource: Mutex::new(None),
poisoned: AtomicBool::new(false),
init_error: Mutex::new(None),
pending_uses: AtomicUsize::new(0),
live_attaches: AtomicUsize::new(0),
init_started_at: Mutex::new(None),
accessor_payload: Mutex::new(None),
}
}
}
#[allow(dead_code)]
fn needs_sibling_spawn(entry: &Entry, resource: &dyn SharedResource) -> bool {
if resource.can_support_more_load() {
return false;
}
let live = entry.live_attaches.load(Ordering::Acquire);
if live == 0 {
crate::diag!(
LogLevel::Warn,
"{EVENT_FAMILY}.share.suppressed key={} generation={} \
reason=quiescent-decline \
note=can_support_more_load() returned false with live_attaches=0; \
driver may be reading historical state (filled ring buffer, \
moving-average that hasn't decayed) instead of in-flight load \
(SRD-35 §\"Validity rules\")",
entry.key.fmt_for_log(),
entry.generation,
);
return false;
}
true
}
const EVENT_FAMILY: &str = "resource";
fn emit_event(level: LogLevel, name: &str, entry: &Entry, extra: &str) {
let suffix = if extra.is_empty() {
String::new()
} else {
format!(" {extra}")
};
crate::diag!(
level,
"{EVENT_FAMILY}.{name} key={} generation={} policy={}{suffix}",
entry.key.fmt_for_log(),
entry.generation,
entry.policy,
);
}
pub struct ResourcePool {
inner: Mutex<PoolInner>,
}
struct PoolInner {
entries_by_key: HashMap<(ResourceKey, usize), Arc<Entry>>,
total_attaches: usize,
}
impl ResourcePool {
pub fn new() -> Self {
Self {
inner: Mutex::new(PoolInner {
entries_by_key: HashMap::new(),
total_attaches: 0,
}),
}
}
pub fn total_attaches(&self) -> usize {
self.inner
.lock()
.unwrap_or_else(|e| e.into_inner())
.total_attaches
}
pub fn live_entries(&self) -> usize {
self.inner
.lock()
.unwrap_or_else(|e| e.into_inner())
.entries_by_key
.len()
}
pub async fn shutdown(self: &Arc<Self>) {
let entries: Vec<Arc<Entry>> = {
let inner = self.inner.lock().unwrap_or_else(|e| e.into_inner());
inner.entries_by_key.values().cloned().collect()
};
for entry in entries {
let resource = {
let mut slot = entry.resource.lock().unwrap_or_else(|e| e.into_inner());
slot.take()
};
let Some(resource) = resource else {
self.remove_entry(&entry.key, entry.generation);
continue;
};
let live = entry.live_attaches.load(Ordering::Acquire);
let pending = entry.pending_uses.load(Ordering::Acquire);
emit_event(
LogLevel::Debug,
"close.started",
&entry,
&format!("reason=session-end live={live} pending={pending}"),
);
let started_at = Instant::now();
let result = resource.close().await;
let elapsed_ms = started_at.elapsed().as_millis() as u64;
match result {
Ok(()) => emit_event(
LogLevel::Debug,
"close.completed",
&entry,
&format!("elapsed_ms={elapsed_ms}"),
),
Err(ref e) => emit_event(
LogLevel::Warn,
"close.failed",
&entry,
&format!("elapsed_ms={elapsed_ms} error={e:?}"),
),
}
self.remove_entry(&entry.key, entry.generation);
}
}
pub fn declare_pending_use(&self, key: ResourceKey, policy: ResourceSharePolicy) {
let entry = self.get_or_create_entry(key, policy, 0);
entry.pending_uses.fetch_add(1, Ordering::AcqRel);
}
pub fn complete_pending_use(&self, key: &ResourceKey) -> bool {
let inner = self.inner.lock().unwrap_or_else(|e| e.into_inner());
let Some(entry) = inner.entries_by_key.get(&(key.clone(), 0)).cloned() else {
return false;
};
drop(inner);
let cur = entry.pending_uses.load(Ordering::Acquire);
if cur == 0 {
return entry.live_attaches.load(Ordering::Acquire) == 0;
}
let new_pending = entry.pending_uses.fetch_sub(1, Ordering::AcqRel) - 1;
new_pending == 0 && entry.live_attaches.load(Ordering::Acquire) == 0
}
pub fn pending_uses_for(&self, key: &ResourceKey) -> Option<usize> {
let inner = self.inner.lock().unwrap_or_else(|e| e.into_inner());
inner
.entries_by_key
.get(&(key.clone(), 0))
.map(|e| e.pending_uses.load(Ordering::Acquire))
}
fn get_or_create_entry(
&self,
key: ResourceKey,
policy: ResourceSharePolicy,
generation: usize,
) -> Arc<Entry> {
let mut inner = self.inner.lock().unwrap_or_else(|e| e.into_inner());
let map_key = (key.clone(), generation);
if let Some(existing) = inner.entries_by_key.get(&map_key) {
return existing.clone();
}
let entry = Arc::new(Entry::new(key, policy, generation));
inner.entries_by_key.insert(map_key, entry.clone());
entry
}
fn remove_entry(&self, key: &ResourceKey, generation: usize) {
let mut inner = self.inner.lock().unwrap_or_else(|e| e.into_inner());
inner.entries_by_key.remove(&(key.clone(), generation));
}
fn lookup_accessor_payload(&self, key: &str) -> Option<Arc<dyn Any + Send + Sync>> {
let inner = self.inner.lock().unwrap_or_else(|e| e.into_inner());
for entry in inner.entries_by_key.values() {
if entry.key.render_key() == key {
return entry
.accessor_payload
.lock()
.unwrap_or_else(|e| e.into_inner())
.clone();
}
}
None
}
async fn ensure_initialized(
&self,
entry: &Arc<Entry>,
first_attach: bool,
factory: ResourceFactory<'_>,
) -> Result<Arc<dyn SharedResource>, String> {
if let Some(existing) = entry
.resource
.lock()
.unwrap_or_else(|e| e.into_inner())
.clone()
{
return Ok(existing);
}
if entry.poisoned.load(Ordering::Acquire) {
let err = entry
.init_error
.lock()
.unwrap_or_else(|e| e.into_inner())
.clone()
.unwrap_or_else(|| "init poisoned".into());
return Err(err);
}
let reason = if first_attach {
"first-attach"
} else {
"capacity-declined"
};
emit_event(
LogLevel::Debug,
"init.started",
entry,
&format!("reason={reason}"),
);
*entry
.init_started_at
.lock()
.unwrap_or_else(|e| e.into_inner()) = Some(Instant::now());
let outcome: Result<Arc<dyn SharedResource>, String> = async {
let resource = factory.build().await?;
resource.init().await?;
Ok(resource)
}
.await;
let elapsed_ms = entry
.init_started_at
.lock()
.unwrap_or_else(|e| e.into_inner())
.map(|t| t.elapsed().as_millis() as u64)
.unwrap_or(0);
match outcome {
Ok(resource) => {
*entry.resource.lock().unwrap_or_else(|e| e.into_inner()) = Some(resource.clone());
*entry
.accessor_payload
.lock()
.unwrap_or_else(|e| e.into_inner()) = resource.accessor_payload();
emit_event(
LogLevel::Debug,
"init.completed",
entry,
&format!("elapsed_ms={elapsed_ms}"),
);
Ok(resource)
}
Err(e) => {
entry.poisoned.store(true, Ordering::Release);
*entry.init_error.lock().unwrap_or_else(|e| e.into_inner()) = Some(e.clone());
emit_event(
LogLevel::Error,
"init.failed",
entry,
&format!("elapsed_ms={elapsed_ms} error={e:?}"),
);
Err(e)
}
}
}
}
impl Default for ResourcePool {
fn default() -> Self {
Self::new()
}
}
type ResourceFactoryFn<'a> =
Box<dyn FnOnce() -> ResourceFuture<'a, Result<Arc<dyn SharedResource>, String>> + Send + 'a>;
pub struct ResourceFactory<'a> {
inner: ResourceFactoryFn<'a>,
}
impl<'a> ResourceFactory<'a> {
pub fn new<F, Fut>(f: F) -> Self
where
F: FnOnce() -> Fut + Send + 'a,
Fut: Future<Output = Result<Arc<dyn SharedResource>, String>> + Send + 'a,
{
Self {
inner: Box::new(move || Box::pin(f()) as ResourceFuture<'a, _>),
}
}
fn build(self) -> ResourceFuture<'a, Result<Arc<dyn SharedResource>, String>> {
(self.inner)()
}
}
pub async fn attach(
pool: &Arc<ResourcePool>,
key: ResourceKey,
policy: ResourceSharePolicy,
phase: impl Into<String>,
factory: ResourceFactory<'_>,
) -> Result<AttachGuard, String> {
let phase = phase.into();
let entry = pool.get_or_create_entry(key.clone(), policy, 0);
let first_attach = entry
.resource
.lock()
.unwrap_or_else(|e| e.into_inner())
.is_none()
&& !entry.poisoned.load(Ordering::Acquire);
let resource = pool
.ensure_initialized(&entry, first_attach, factory)
.await?;
let needs_share = matches!(
policy,
ResourceSharePolicy::Shared | ResourceSharePolicy::PerScenario
);
if needs_share && !resource.can_share() {
return Err(format!(
"resource for {} declared can_share()=false but policy is {policy}; \
elevate isolation to per-phase or per-fiber, or fix the resource impl",
entry.key.fmt_for_log(),
));
}
let live = entry.live_attaches.fetch_add(1, Ordering::AcqRel) + 1;
let pending_dec = entry.pending_uses.load(Ordering::Acquire);
let pending = if pending_dec > 0 {
entry.pending_uses.fetch_sub(1, Ordering::AcqRel) - 1
} else {
0
};
{
let mut inner = pool.inner.lock().unwrap_or_else(|e| e.into_inner());
inner.total_attaches += 1;
}
emit_event(
LogLevel::Debug,
"attach",
&entry,
&format!("phase={phase:?} pending={pending} live={live}"),
);
Ok(AttachGuard {
pool: Arc::clone(pool),
entry,
resource,
phase,
detached: false,
})
}
pub struct AttachGuard {
pool: Arc<ResourcePool>,
entry: Arc<Entry>,
resource: Arc<dyn SharedResource>,
phase: String,
detached: bool,
}
impl std::fmt::Debug for AttachGuard {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("AttachGuard")
.field("key", &self.entry.key)
.field("generation", &self.entry.generation)
.field("policy", &self.entry.policy)
.field("phase", &self.phase)
.field("detached", &self.detached)
.finish()
}
}
impl AttachGuard {
pub fn resource(&self) -> Arc<dyn SharedResource> {
Arc::clone(&self.resource)
}
pub fn key(&self) -> &ResourceKey {
&self.entry.key
}
pub async fn detach(mut self) -> Result<(), String> {
self.detach_inner(true).await
}
async fn detach_inner(&mut self, await_close: bool) -> Result<(), String> {
if self.detached {
return Ok(());
}
self.detached = true;
let live = self.entry.live_attaches.fetch_sub(1, Ordering::AcqRel) - 1;
let pending = self.entry.pending_uses.load(Ordering::Acquire);
emit_event(
LogLevel::Debug,
"detach",
&self.entry,
&format!("phase={:?} pending={pending} live={live}", self.phase),
);
if live == 0 && pending == 0 {
self.trigger_close(await_close, "refcount-zero").await?;
}
Ok(())
}
async fn trigger_close(&self, await_close: bool, reason: &str) -> Result<(), String> {
let resource_opt = {
let mut slot = self
.entry
.resource
.lock()
.unwrap_or_else(|e| e.into_inner());
slot.take()
};
let Some(resource) = resource_opt else {
self.pool
.remove_entry(&self.entry.key, self.entry.generation);
return Ok(());
};
emit_event(
LogLevel::Debug,
"close.started",
&self.entry,
&format!("reason={reason}"),
);
let started_at = Instant::now();
let entry_for_async = Arc::clone(&self.entry);
let pool_for_async = Arc::clone(&self.pool);
let close_future = resource.close();
if await_close {
let result = close_future.await;
let elapsed_ms = started_at.elapsed().as_millis() as u64;
match result {
Ok(()) => emit_event(
LogLevel::Debug,
"close.completed",
&self.entry,
&format!("elapsed_ms={elapsed_ms}"),
),
Err(ref e) => emit_event(
LogLevel::Warn,
"close.failed",
&self.entry,
&format!("elapsed_ms={elapsed_ms} error={e:?}"),
),
}
self.pool
.remove_entry(&self.entry.key, self.entry.generation);
result
} else {
tokio::spawn(async move {
let result = close_future.await;
let elapsed_ms = started_at.elapsed().as_millis() as u64;
match result {
Ok(()) => emit_event(
LogLevel::Debug,
"close.completed",
&entry_for_async,
&format!("elapsed_ms={elapsed_ms}"),
),
Err(ref e) => emit_event(
LogLevel::Warn,
"close.failed",
&entry_for_async,
&format!("elapsed_ms={elapsed_ms} error={e:?}"),
),
}
pool_for_async.remove_entry(&entry_for_async.key, entry_for_async.generation);
});
Ok(())
}
}
}
impl Drop for AttachGuard {
fn drop(&mut self) {
if self.detached {
return;
}
let live = self.entry.live_attaches.fetch_sub(1, Ordering::AcqRel) - 1;
let pending = self.entry.pending_uses.load(Ordering::Acquire);
self.detached = true;
let phase = &self.phase;
emit_event(
LogLevel::Debug,
"detach",
&self.entry,
&format!("phase={phase:?} pending={pending} live={live}"),
);
if live == 0 && pending == 0 {
let resource_opt = {
let mut slot = self
.entry
.resource
.lock()
.unwrap_or_else(|e| e.into_inner());
slot.take()
};
let Some(resource) = resource_opt else {
self.pool
.remove_entry(&self.entry.key, self.entry.generation);
return;
};
let entry = Arc::clone(&self.entry);
let pool = Arc::clone(&self.pool);
emit_event(
LogLevel::Debug,
"close.started",
&self.entry,
"reason=refcount-zero",
);
let started_at = Instant::now();
if let Ok(handle) = tokio::runtime::Handle::try_current() {
let close_future = resource.close();
handle.spawn(async move {
let result = close_future.await;
let elapsed_ms = started_at.elapsed().as_millis() as u64;
match result {
Ok(()) => emit_event(
LogLevel::Debug,
"close.completed",
&entry,
&format!("elapsed_ms={elapsed_ms}"),
),
Err(ref e) => emit_event(
LogLevel::Warn,
"close.failed",
&entry,
&format!("elapsed_ms={elapsed_ms} error={e:?}"),
),
}
pool.remove_entry(&entry.key, entry.generation);
});
} else {
drop(resource);
emit_event(
LogLevel::Debug,
"close.completed",
&self.entry,
"elapsed_ms=0",
);
self.pool
.remove_entry(&self.entry.key, self.entry.generation);
}
}
}
}
pub struct LegacyAdapterResource {
key: ResourceKey,
adapter: Mutex<Option<Arc<dyn DriverAdapter>>>,
}
impl LegacyAdapterResource {
pub fn new(key: ResourceKey, adapter: Arc<dyn DriverAdapter>) -> Self {
Self {
key,
adapter: Mutex::new(Some(adapter)),
}
}
pub fn adapter(&self) -> Option<Arc<dyn DriverAdapter>> {
self.adapter
.lock()
.unwrap_or_else(|e| e.into_inner())
.clone()
}
}
impl SharedResource for LegacyAdapterResource {
fn resource_key(&self) -> &ResourceKey {
&self.key
}
fn can_share(&self) -> bool {
false
}
fn close(self: Arc<Self>) -> ResourceFuture<'static, Result<(), String>> {
Box::pin(async move {
let _ = self
.adapter
.lock()
.unwrap_or_else(|e| e.into_inner())
.take();
Ok(())
})
}
fn as_legacy_adapter(&self) -> Option<Arc<dyn DriverAdapter>> {
self.adapter()
}
}
pub struct SharedAdapterResource {
key: ResourceKey,
adapter: Mutex<Option<Arc<dyn DriverAdapter>>>,
}
impl SharedAdapterResource {
pub fn new(key: ResourceKey, adapter: Arc<dyn DriverAdapter>) -> Self {
Self {
key,
adapter: Mutex::new(Some(adapter)),
}
}
pub fn adapter(&self) -> Option<Arc<dyn DriverAdapter>> {
self.adapter
.lock()
.unwrap_or_else(|e| e.into_inner())
.clone()
}
}
impl SharedResource for SharedAdapterResource {
fn resource_key(&self) -> &ResourceKey {
&self.key
}
fn can_share(&self) -> bool {
true
}
fn close(self: Arc<Self>) -> ResourceFuture<'static, Result<(), String>> {
Box::pin(async move {
let adapter = self
.adapter
.lock()
.unwrap_or_else(|e| e.into_inner())
.take();
if let Some(adapter) = adapter {
adapter.shutdown().await;
}
Ok(())
})
}
fn as_legacy_adapter(&self) -> Option<Arc<dyn DriverAdapter>> {
self.adapter()
}
fn accessor_payload(&self) -> Option<Arc<dyn Any + Send + Sync>> {
self.adapter().and_then(|a| a.accessor_payload())
}
}
pub async fn attach_shared_adapter<F, Fut>(
pool: &Arc<ResourcePool>,
adapter_name: &str,
phase: &str,
key: ResourceKey,
factory: F,
) -> Result<(Arc<dyn DriverAdapter>, AttachGuard), String>
where
F: FnOnce() -> Fut + Send + 'static,
Fut: Future<Output = Result<Arc<dyn DriverAdapter>, String>> + Send + 'static,
{
let key_for_resource = key.clone();
let factory = ResourceFactory::new(move || async move {
let adapter = factory().await?;
let res = SharedAdapterResource::new(key_for_resource, adapter);
Ok(Arc::new(res) as Arc<dyn SharedResource>)
});
let guard = attach(pool, key, ResourceSharePolicy::Shared, phase, factory).await?;
let adapter = guard.resource().as_legacy_adapter().ok_or_else(|| {
format!(
"internal: shared pool resource for adapter '{adapter_name}' did not surface \
a DriverAdapter handle — should be unreachable"
)
})?;
Ok((adapter, guard))
}
pub async fn attach_legacy_adapter<F, Fut>(
pool: &Arc<ResourcePool>,
adapter_name: &str,
phase: &str,
key_extras: &[(&str, &str)],
factory: F,
) -> Result<(Arc<dyn DriverAdapter>, AttachGuard), String>
where
F: FnOnce() -> Fut + Send + 'static,
Fut: Future<Output = Result<Arc<dyn DriverAdapter>, String>> + Send + 'static,
{
let mut key = ResourceKey::new(adapter_name);
for (k, v) in key_extras {
key = key.with(*k, *v);
}
let key_for_resource = key.clone();
let factory = ResourceFactory::new(move || async move {
let adapter = factory().await?;
let res = LegacyAdapterResource::new(key_for_resource, adapter);
Ok(Arc::new(res) as Arc<dyn SharedResource>)
});
let guard = attach(pool, key, ResourceSharePolicy::PerPhase, phase, factory).await?;
let resource = guard.resource();
let adapter = resource.as_legacy_adapter().ok_or_else(|| {
format!(
"internal: pool resource for adapter '{adapter_name}' did not surface \
a legacy DriverAdapter handle — Push A path should be unreachable"
)
})?;
Ok((adapter, guard))
}
pub fn pre_map_pending_uses(
pool: &Arc<ResourcePool>,
tree: &crate::scene_tree::SceneTree,
phases: &std::collections::HashMap<String, nmbrs_workload::model::WorkloadPhase>,
default_driver: &str,
merged_params: &std::collections::HashMap<String, String>,
) -> Result<(), String> {
for node in tree.dfs_phases() {
let adapter = phases
.get(&node.name)
.and_then(|p| p.adapter.clone())
.unwrap_or_else(|| default_driver.to_string());
let selector = format!("{adapter}driver");
let Some(driver_name) =
crate::adapter::resolve_driver_name(&adapter, &selector, merged_params)
else {
continue;
};
let Some(reg) = crate::adapter::find_shared_driver(&adapter, driver_name) else {
continue;
};
let key = (reg.resource_key)(merged_params).map_err(|e| {
format!(
"resource pool pre-map: phase '{}' adapter '{}' driver '{}': {}",
node.name, adapter, driver_name, e,
)
})?;
pool.declare_pending_use(key, default_policy_for(reg.share_capability));
}
Ok(())
}
static ACTIVE_POOL: Mutex<Weak<ResourcePool>> = Mutex::new(Weak::new());
struct PoolAccessorView;
impl polydat::ResourceAccessor for PoolAccessorView {
fn lookup(&self, key: &str) -> Option<Arc<dyn Any + Send + Sync>> {
let pool = ACTIVE_POOL
.lock()
.unwrap_or_else(|e| e.into_inner())
.upgrade()?;
pool.lookup_accessor_payload(key)
}
}
pub fn install_accessor(pool: &Arc<ResourcePool>) {
*ACTIVE_POOL.lock().unwrap_or_else(|e| e.into_inner()) = Arc::downgrade(pool);
}
pub fn pool_resources() -> polydat::ResourceScope {
polydat::ResourceScope::with_accessor(Arc::new(PoolAccessorView))
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::atomic::AtomicU32;
struct MockResource {
key: ResourceKey,
can_share: bool,
saturated_after_attaches: u32,
attach_count: AtomicU32,
init_calls: AtomicU32,
close_calls: AtomicU32,
init_should_fail: bool,
}
impl MockResource {
fn new(key: ResourceKey) -> Arc<Self> {
Arc::new(Self {
key,
can_share: true,
saturated_after_attaches: u32::MAX,
attach_count: AtomicU32::new(0),
init_calls: AtomicU32::new(0),
close_calls: AtomicU32::new(0),
init_should_fail: false,
})
}
fn with_can_share(key: ResourceKey, can_share: bool) -> Arc<Self> {
Arc::new(Self {
key,
can_share,
saturated_after_attaches: u32::MAX,
attach_count: AtomicU32::new(0),
init_calls: AtomicU32::new(0),
close_calls: AtomicU32::new(0),
init_should_fail: false,
})
}
fn with_init_failure(key: ResourceKey) -> Arc<Self> {
Arc::new(Self {
key,
can_share: true,
saturated_after_attaches: u32::MAX,
attach_count: AtomicU32::new(0),
init_calls: AtomicU32::new(0),
close_calls: AtomicU32::new(0),
init_should_fail: true,
})
}
}
impl SharedResource for MockResource {
fn resource_key(&self) -> &ResourceKey {
&self.key
}
fn can_share(&self) -> bool {
self.can_share
}
fn can_support_more_load(&self) -> bool {
self.attach_count.load(Ordering::Acquire) < self.saturated_after_attaches
}
fn init(&self) -> ResourceFuture<'_, Result<(), String>> {
self.init_calls.fetch_add(1, Ordering::AcqRel);
let fail = self.init_should_fail;
Box::pin(async move {
if fail {
Err("simulated init failure".into())
} else {
Ok(())
}
})
}
fn close(self: Arc<Self>) -> ResourceFuture<'static, Result<(), String>> {
self.close_calls.fetch_add(1, Ordering::AcqRel);
Box::pin(async { Ok(()) })
}
}
fn key(adapter: &str, vendor: &str) -> ResourceKey {
ResourceKey::new(adapter).with("vendor", vendor)
}
#[test]
fn key_value_equality_is_field_order_independent() {
let a = ResourceKey::new("cql")
.with("hosts", "h1")
.with("port", "9042");
let b = ResourceKey::new("cql")
.with("port", "9042")
.with("hosts", "h1");
assert_eq!(a, b);
assert_eq!(a.fmt_for_log(), b.fmt_for_log());
}
#[test]
fn key_log_format_redacts_password() {
let k = ResourceKey::new("cql")
.with("hosts", "h1")
.with("password", "hunter2");
let s = k.fmt_for_log();
assert!(s.contains("hosts=h1"));
assert!(s.contains("password=***"), "got: {s}");
assert!(!s.contains("hunter2"));
}
#[test]
fn policy_parse_accepts_canonical_and_underscore_forms() {
assert_eq!(
ResourceSharePolicy::parse("shared").unwrap(),
ResourceSharePolicy::Shared
);
assert_eq!(
ResourceSharePolicy::parse("per-phase").unwrap(),
ResourceSharePolicy::PerPhase
);
assert_eq!(
ResourceSharePolicy::parse("per_phase").unwrap(),
ResourceSharePolicy::PerPhase
);
assert_eq!(
ResourceSharePolicy::parse("PER-FIBER").unwrap(),
ResourceSharePolicy::PerFiber
);
ResourceSharePolicy::parse("nope").unwrap_err();
}
#[test]
fn capability_floor_orders_isolation_correctly() {
assert!(capability_floor(ShareCapability::Shared) <= ResourceSharePolicy::Shared);
assert!(capability_floor(ShareCapability::PerPhase) > ResourceSharePolicy::Shared);
}
#[tokio::test]
async fn shared_resource_survives_detach_until_more_phases_predicted() {
let pool = Arc::new(ResourcePool::new());
let mock = MockResource::new(key("cql", "cassandra-cpp"));
let mock_for_factory = Arc::clone(&mock);
pool.declare_pending_use(mock.resource_key().clone(), ResourceSharePolicy::Shared);
pool.declare_pending_use(mock.resource_key().clone(), ResourceSharePolicy::Shared);
assert_eq!(pool.pending_uses_for(mock.resource_key()), Some(2));
let guard = attach(
&pool,
mock.resource_key().clone(),
ResourceSharePolicy::Shared,
"phase_A",
ResourceFactory::new(
move || async move { Ok(mock_for_factory as Arc<dyn SharedResource>) },
),
)
.await
.unwrap();
assert_eq!(mock.init_calls.load(Ordering::Acquire), 1);
guard.detach().await.unwrap();
assert_eq!(
mock.close_calls.load(Ordering::Acquire),
0,
"Shared-policy detach MUST NOT close while another \
phase is still predicted to attach"
);
assert_eq!(pool.pending_uses_for(mock.resource_key()), Some(1));
assert_eq!(pool.live_entries(), 1);
let mock_for_factory2 = Arc::clone(&mock);
let guard2 = attach(
&pool,
mock.resource_key().clone(),
ResourceSharePolicy::Shared,
"phase_B",
ResourceFactory::new(move || async move {
Ok(mock_for_factory2 as Arc<dyn SharedResource>)
}),
)
.await
.unwrap();
guard2.detach().await.unwrap();
assert_eq!(
mock.close_calls.load(Ordering::Acquire),
1,
"close MUST fire when pending hits 0 AND live hits 0 — \
that's the SRD-35 Push D close-on-zero contract"
);
assert_eq!(pool.live_entries(), 0, "entry removed once closed");
pool.shutdown().await;
assert_eq!(
mock.close_calls.load(Ordering::Acquire),
1,
"shutdown() must NOT double-close an entry"
);
}
#[tokio::test]
async fn shared_resource_overpredicted_pending_holds_until_shutdown() {
let pool = Arc::new(ResourcePool::new());
let mock = MockResource::new(key("cql", "cassandra-cpp"));
let mock_for_factory = Arc::clone(&mock);
for _ in 0..3 {
pool.declare_pending_use(mock.resource_key().clone(), ResourceSharePolicy::Shared);
}
let guard = attach(
&pool,
mock.resource_key().clone(),
ResourceSharePolicy::Shared,
"phase_A",
ResourceFactory::new(
move || async move { Ok(mock_for_factory as Arc<dyn SharedResource>) },
),
)
.await
.unwrap();
guard.detach().await.unwrap();
assert_eq!(pool.pending_uses_for(mock.resource_key()), Some(2));
assert_eq!(mock.close_calls.load(Ordering::Acquire), 0);
assert_eq!(pool.live_entries(), 1);
pool.shutdown().await;
assert_eq!(
mock.close_calls.load(Ordering::Acquire),
1,
"shutdown() drains residual entries with pending > 0"
);
}
#[tokio::test]
async fn per_phase_resource_closes_on_detach() {
let pool = Arc::new(ResourcePool::new());
let mock = MockResource::new(key("legacy", "synthetic"));
let mock_for_factory = Arc::clone(&mock);
let guard = attach(
&pool,
mock.resource_key().clone(),
ResourceSharePolicy::PerPhase,
"phase_A",
ResourceFactory::new(
move || async move { Ok(mock_for_factory as Arc<dyn SharedResource>) },
),
)
.await
.unwrap();
guard.detach().await.unwrap();
assert_eq!(
mock.close_calls.load(Ordering::Acquire),
1,
"PerPhase detach MUST close immediately"
);
assert_eq!(pool.live_entries(), 0);
}
#[tokio::test]
async fn init_failure_poisons_the_entry() {
let pool = Arc::new(ResourcePool::new());
let mock = MockResource::with_init_failure(key("cql", "cassandra-cpp"));
let mock_for_factory = Arc::clone(&mock);
let result = attach(
&pool,
mock.resource_key().clone(),
ResourceSharePolicy::Shared,
"phase_A",
ResourceFactory::new(
move || async move { Ok(mock_for_factory as Arc<dyn SharedResource>) },
),
)
.await;
let err = result.expect_err("init failure must propagate");
assert!(err.contains("simulated init failure"), "got: {err}");
}
#[tokio::test]
async fn shared_policy_against_can_share_false_resource_errors() {
let pool = Arc::new(ResourcePool::new());
let mock = MockResource::with_can_share(key("legacy", "x"), false);
let mock_for_factory = Arc::clone(&mock);
let result = attach(
&pool,
mock.resource_key().clone(),
ResourceSharePolicy::Shared,
"phase_A",
ResourceFactory::new(
move || async move { Ok(mock_for_factory as Arc<dyn SharedResource>) },
),
)
.await;
let err = result.expect_err("shared policy on can_share()=false resource must be rejected");
assert!(err.contains("can_share()=false"), "got: {err}");
assert!(
err.contains("per-phase") || err.contains("per-fiber"),
"error must point at the available policies, got: {err}"
);
}
#[tokio::test]
async fn per_phase_policy_with_can_share_false_works() {
let pool = Arc::new(ResourcePool::new());
let mock = MockResource::with_can_share(key("legacy", "x"), false);
let mock_for_factory = Arc::clone(&mock);
let guard = attach(
&pool,
mock.resource_key().clone(),
ResourceSharePolicy::PerPhase,
"phase_A",
ResourceFactory::new(
move || async move { Ok(mock_for_factory as Arc<dyn SharedResource>) },
),
)
.await
.expect("PerPhase must accept can_share()=false");
guard.detach().await.unwrap();
}
#[tokio::test]
async fn shared_policy_reuses_one_instance_across_attaches() {
let pool = Arc::new(ResourcePool::new());
let mock = MockResource::new(key("cql", "cassandra-cpp"));
let mock_for_factory = Arc::clone(&mock);
let factory_called = Arc::new(AtomicU32::new(0));
let factory_called_clone = Arc::clone(&factory_called);
for _ in 0..3 {
pool.declare_pending_use(mock.resource_key().clone(), ResourceSharePolicy::Shared);
}
let g1 = attach(
&pool,
mock.resource_key().clone(),
ResourceSharePolicy::Shared,
"phase_A",
ResourceFactory::new(move || {
factory_called_clone.fetch_add(1, Ordering::AcqRel);
let m = Arc::clone(&mock_for_factory);
async move { Ok(m as Arc<dyn SharedResource>) }
}),
)
.await
.unwrap();
let mock_for_factory2 = Arc::clone(&mock);
let factory_called_clone2 = Arc::clone(&factory_called);
let g2 = attach(
&pool,
mock.resource_key().clone(),
ResourceSharePolicy::Shared,
"phase_B",
ResourceFactory::new(move || {
factory_called_clone2.fetch_add(1, Ordering::AcqRel);
let m = Arc::clone(&mock_for_factory2);
async move { Ok(m as Arc<dyn SharedResource>) }
}),
)
.await
.unwrap();
assert_eq!(
factory_called.load(Ordering::Acquire),
1,
"factory must not be called twice for the same key under Shared policy"
);
assert_eq!(
mock.init_calls.load(Ordering::Acquire),
1,
"init() must fire exactly once"
);
assert_eq!(pool.total_attaches(), 2);
assert_eq!(pool.live_entries(), 1);
g1.detach().await.unwrap();
assert_eq!(
mock.close_calls.load(Ordering::Acquire),
0,
"close MUST NOT fire while another guard is live"
);
g2.detach().await.unwrap();
assert_eq!(
mock.close_calls.load(Ordering::Acquire),
0,
"close MUST NOT fire while pending uses remain — \
the resource stays alive for predicted future phases"
);
assert_eq!(pool.live_entries(), 1, "entry stays cached across phases");
pool.shutdown().await;
assert_eq!(
mock.close_calls.load(Ordering::Acquire),
1,
"shutdown() drains the entry"
);
assert_eq!(pool.live_entries(), 0);
}
struct LegacyDummy;
impl crate::adapter::DriverAdapter for LegacyDummy {
fn name(&self) -> &str {
"legacy_dummy"
}
fn map_op<'a>(
&'a self,
_template: &'a nmbrs_workload::model::ParsedOp,
_parent: std::sync::Arc<dyn polydat::Kernel>,
) -> std::pin::Pin<
Box<
dyn std::future::Future<
Output = Result<Box<dyn crate::adapter::OpDispenser>, String>,
> + Send
+ 'a,
>,
> {
Box::pin(async move { Err("dummy".into()) })
}
}
#[test]
fn capacity_decline_at_quiescence_is_caught_as_driver_bug() {
struct Reports {
key: ResourceKey,
has_capacity: bool,
}
impl SharedResource for Reports {
fn resource_key(&self) -> &ResourceKey {
&self.key
}
fn can_support_more_load(&self) -> bool {
self.has_capacity
}
}
let entry = Entry::new(ResourceKey::new("buggy"), ResourceSharePolicy::Shared, 0);
let calm = Reports {
key: ResourceKey::new("buggy"),
has_capacity: true,
};
assert!(
!needs_sibling_spawn(&entry, &calm),
"no spawn when the resource has capacity"
);
let buggy = Reports {
key: ResourceKey::new("buggy"),
has_capacity: false,
};
assert_eq!(entry.live_attaches.load(Ordering::Acquire), 0);
assert!(
!needs_sibling_spawn(&entry, &buggy),
"quiescent decline must be suppressed (driver bug)"
);
entry.live_attaches.fetch_add(3, Ordering::AcqRel);
assert!(
needs_sibling_spawn(&entry, &buggy),
"decline under live load is genuine saturation; spawn sibling"
);
assert!(
!needs_sibling_spawn(&entry, &calm),
"no spawn when there is capacity, even under load"
);
}
#[tokio::test]
async fn shared_adapter_resource_returns_one_instance_for_one_key() {
let pool = Arc::new(ResourcePool::new());
let adapter: Arc<dyn DriverAdapter> = Arc::new(LegacyDummy);
let key = ResourceKey::new("cql").with("driver", "synthetic");
let factory_call_count = Arc::new(AtomicU32::new(0));
for _ in 0..3 {
pool.declare_pending_use(key.clone(), ResourceSharePolicy::Shared);
}
let factory_count_a = Arc::clone(&factory_call_count);
let adapter_a = Arc::clone(&adapter);
let (got_a, guard_a) =
attach_shared_adapter(&pool, "cql", "phase_A", key.clone(), move || {
factory_count_a.fetch_add(1, Ordering::AcqRel);
let a = Arc::clone(&adapter_a);
async move { Ok(a) }
})
.await
.unwrap();
let factory_count_b = Arc::clone(&factory_call_count);
let adapter_b = Arc::clone(&adapter);
let (got_b, guard_b) =
attach_shared_adapter(&pool, "cql", "phase_B", key.clone(), move || {
factory_count_b.fetch_add(1, Ordering::AcqRel);
let a = Arc::clone(&adapter_b);
async move { Ok(a) }
})
.await
.unwrap();
assert_eq!(
factory_call_count.load(Ordering::Acquire),
1,
"Shared policy MUST cache the factory output"
);
assert!(
Arc::ptr_eq(&got_a, &got_b),
"both attaches MUST surface the SAME Arc<dyn DriverAdapter>"
);
assert_eq!(pool.live_entries(), 1);
assert_eq!(pool.total_attaches(), 2);
guard_a.detach().await.unwrap();
assert_eq!(
pool.live_entries(),
1,
"entry stays live while any guard holds it"
);
guard_b.detach().await.unwrap();
assert_eq!(
pool.live_entries(),
1,
"Shared-policy entry stays cached across phases — \
session shutdown is the close trigger, not last detach"
);
pool.shutdown().await;
assert_eq!(pool.live_entries(), 0, "shutdown drains the cached entry");
}
#[test]
fn declare_then_complete_pending_use_drains_to_zero() {
let pool = Arc::new(ResourcePool::new());
let k = ResourceKey::new("cql").with("hosts", "h1");
assert_eq!(
pool.pending_uses_for(&k),
None,
"no entry exists before declare"
);
pool.declare_pending_use(k.clone(), ResourceSharePolicy::Shared);
pool.declare_pending_use(k.clone(), ResourceSharePolicy::Shared);
pool.declare_pending_use(k.clone(), ResourceSharePolicy::Shared);
assert_eq!(pool.pending_uses_for(&k), Some(3));
assert_eq!(
pool.live_entries(),
1,
"declare_pending_use creates the entry on first call"
);
let eligible = pool.complete_pending_use(&k);
assert!(!eligible, "pending == 2; not eligible yet");
assert_eq!(pool.pending_uses_for(&k), Some(2));
let eligible = pool.complete_pending_use(&k);
assert!(!eligible, "pending == 1; not eligible yet");
assert_eq!(pool.pending_uses_for(&k), Some(1));
let eligible = pool.complete_pending_use(&k);
assert!(
eligible,
"pending == 0 && live == 0 — entry is eligible for close"
);
assert_eq!(pool.pending_uses_for(&k), Some(0));
let eligible = pool.complete_pending_use(&k);
assert!(
eligible,
"saturating decrement: extra completes don't underflow"
);
assert_eq!(pool.pending_uses_for(&k), Some(0));
}
#[tokio::test]
async fn declare_pending_use_paired_with_attach_closes_promptly() {
let pool = Arc::new(ResourcePool::new());
let mock = MockResource::new(key("cql", "cassandra-cpp"));
pool.declare_pending_use(mock.resource_key().clone(), ResourceSharePolicy::Shared);
let mock_for_factory = Arc::clone(&mock);
let guard = attach(
&pool,
mock.resource_key().clone(),
ResourceSharePolicy::Shared,
"phase_A",
ResourceFactory::new(
move || async move { Ok(mock_for_factory as Arc<dyn SharedResource>) },
),
)
.await
.unwrap();
assert_eq!(pool.pending_uses_for(mock.resource_key()), Some(0));
guard.detach().await.unwrap();
assert_eq!(
mock.close_calls.load(Ordering::Acquire),
1,
"Push D: close fires the moment the last predicted \
phase detaches — not at session shutdown"
);
assert_eq!(pool.live_entries(), 0);
}
#[test]
fn render_key_is_identity_preserving_unlike_log_format() {
let a = ResourceKey::new("cql")
.with("hosts", "h1")
.with("password", "p1");
let b = ResourceKey::new("cql")
.with("hosts", "h1")
.with("password", "p2");
assert_eq!(
a.fmt_for_log(),
b.fmt_for_log(),
"log format redacts password → collides (expected)"
);
assert_ne!(
a.render_key(),
b.render_key(),
"render_key must distinguish identity-bearing password"
);
assert!(
!a.render_key().contains("p1") && a.render_key().contains("password=sha256:"),
"the fingerprint must carry a digest of the secret, never the cleartext; got: {}",
a.render_key()
);
let c = ResourceKey::new("cql")
.with("password", "p1")
.with("hosts", "h1");
assert_eq!(
a.render_key(),
c.render_key(),
"render_key is field-order independent"
);
}
#[tokio::test]
async fn accessor_payload_populates_and_looks_up_by_render_key() {
struct WithPayload {
key: ResourceKey,
}
impl SharedResource for WithPayload {
fn resource_key(&self) -> &ResourceKey {
&self.key
}
fn accessor_payload(&self) -> Option<Arc<dyn Any + Send + Sync>> {
Some(Arc::new(4242u64) as Arc<dyn Any + Send + Sync>)
}
}
let pool = Arc::new(ResourcePool::new());
let key = ResourceKey::new("cql")
.with("hosts", "h1")
.with("keyspace", "ks");
let rendered = key.render_key();
assert!(
pool.lookup_accessor_payload(&rendered).is_none(),
"no live entry ⇒ None"
);
let key_for_factory = key.clone();
let guard = attach(
&pool,
key.clone(),
ResourceSharePolicy::Shared,
"phase_A",
ResourceFactory::new(move || async move {
Ok(Arc::new(WithPayload {
key: key_for_factory,
}) as Arc<dyn SharedResource>)
}),
)
.await
.unwrap();
let payload = pool
.lookup_accessor_payload(&rendered)
.expect("payload present after successful init");
let n = payload
.downcast_ref::<u64>()
.expect("payload downcasts to u64");
assert_eq!(*n, 4242);
guard.detach().await.unwrap();
assert!(
pool.lookup_accessor_payload(&rendered).is_none(),
"closed entry ⇒ None"
);
}
#[tokio::test]
async fn default_resource_has_no_accessor_payload() {
let pool = Arc::new(ResourcePool::new());
let mock = MockResource::new(key("cql", "cassandra-cpp"));
let rendered = mock.resource_key().render_key();
let mock_for_factory = Arc::clone(&mock);
let guard = attach(
&pool,
mock.resource_key().clone(),
ResourceSharePolicy::Shared,
"phase_A",
ResourceFactory::new(
move || async move { Ok(mock_for_factory as Arc<dyn SharedResource>) },
),
)
.await
.unwrap();
assert!(
pool.lookup_accessor_payload(&rendered).is_none(),
"default accessor_payload() ⇒ None even for a live entry"
);
guard.detach().await.unwrap();
}
#[tokio::test]
async fn legacy_adapter_resource_declares_non_shareable() {
let key = ResourceKey::new("legacy");
let adapter: Arc<dyn DriverAdapter> = Arc::new(LegacyDummy);
let wrapped = Arc::new(LegacyAdapterResource::new(key, adapter));
assert!(
!wrapped.can_share(),
"LegacyAdapterResource MUST declare can_share()=false"
);
}
}