use crate::pg::pool::DjogiPool;
use std::collections::HashMap;
use std::marker::PhantomData;
use std::str::FromStr;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex, OnceLock, Weak};
use tokio::sync::broadcast;
use tokio_postgres::AsyncMessage;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum EventKind {
Created,
Updated,
Deleted,
}
#[derive(Debug, Clone)]
pub struct ModelEvent<M: crate::model::Model> {
pub kind: EventKind,
pub id: M::Pk,
}
#[derive(Debug, thiserror::Error)]
pub enum NotifyError {
#[error("broadcast channel lagged by {skipped} events")]
ChannelLagged { skipped: u64 },
#[error("failed to start NOTIFY listener: {0}")]
ListenerStartFailed(String),
#[error("payload decode failed for {raw:?}: {source}")]
PayloadDecode {
raw: String,
#[source]
source: serde_json::Error,
},
#[error("invalid id {id:?} in payload (M::Pk::from_str rejected): {reason}")]
InvalidId { id: String, reason: String },
#[error("notify watcher terminated; re-subscribe to start a fresh listener")]
ListenerTerminated,
}
#[derive(Debug, Clone)]
struct RawEvent {
kind_str: String,
id_str: String,
}
struct PgListener {
senders: Arc<Mutex<HashMap<String, broadcast::Sender<RawEvent>>>>,
failed: Arc<AtomicBool>,
#[allow(dead_code)] client: tokio_postgres::Client,
}
impl PgListener {
fn is_failed(&self) -> bool {
self.failed.load(Ordering::Acquire)
}
}
impl Drop for PgListener {
fn drop(&mut self) {
tracing::debug!(
target: "djogi::notify",
"PgListener dropped — dedicated client torn down, watcher will exit"
);
}
}
fn close_all_senders(senders: &Mutex<HashMap<String, broadcast::Sender<RawEvent>>>) {
let mut guard = match senders.lock() {
Ok(guard) => guard,
Err(poisoned) => {
tracing::warn!(
target: "djogi::notify",
"senders mutex was poisoned; clearing on watcher exit anyway"
);
poisoned.into_inner()
}
};
guard.clear();
}
struct WatcherExitGuard {
senders: Arc<Mutex<HashMap<String, broadcast::Sender<RawEvent>>>>,
failed: Arc<AtomicBool>,
}
impl Drop for WatcherExitGuard {
fn drop(&mut self) {
close_all_senders(&self.senders);
self.failed.store(true, Ordering::Release);
tracing::debug!(
target: "djogi::notify",
"watcher exit guard fired — senders cleared, failed flag published"
);
}
}
fn registry() -> &'static Mutex<HashMap<u64, Weak<PgListener>>> {
static REGISTRY: OnceLock<Mutex<HashMap<u64, Weak<PgListener>>>> = OnceLock::new();
REGISTRY.get_or_init(|| Mutex::new(HashMap::new()))
}
fn pool_key(pool: &DjogiPool) -> u64 {
pool.pool_id
}
fn upgrade_existing<T>(map: &Mutex<HashMap<u64, Weak<T>>>, key: u64) -> Option<Arc<T>> {
let mut guard = map.lock().expect("notify registry mutex poisoned");
if let Some(weak) = guard.get(&key) {
if let Some(strong) = weak.upgrade() {
return Some(strong);
}
guard.remove(&key);
}
None
}
fn remove_if_current<T>(map: &Mutex<HashMap<u64, Weak<T>>>, key: u64, expected: &Arc<T>) -> bool {
let mut guard = map.lock().expect("notify registry mutex poisoned");
let should_remove = match guard.get(&key).and_then(Weak::upgrade) {
Some(current) => Arc::ptr_eq(¤t, expected),
None => guard.contains_key(&key),
};
if should_remove {
guard.remove(&key);
}
should_remove
}
fn install_or_lose<T>(map: &Mutex<HashMap<u64, Weak<T>>>, key: u64, candidate: Arc<T>) -> Arc<T> {
let mut guard = map.lock().expect("notify registry mutex poisoned");
match guard.get(&key).and_then(|w| w.upgrade()) {
Some(existing) => existing,
None => {
guard.insert(key, Arc::downgrade(&candidate));
candidate
}
}
}
fn decode_event<M: crate::model::Model>(raw: &RawEvent) -> Result<ModelEvent<M>, NotifyError>
where
M::Pk: FromStr,
<M::Pk as FromStr>::Err: std::fmt::Display,
{
let kind = match raw.kind_str.as_str() {
"create" => EventKind::Created,
"save" => EventKind::Updated,
"delete" => EventKind::Deleted,
other => {
return Err(NotifyError::PayloadDecode {
raw: serde_json::json!({ "kind": other, "id": &raw.id_str }).to_string(),
source: serde::de::Error::unknown_variant(other, &["create", "save", "delete"]),
});
}
};
let id = M::Pk::from_str(&raw.id_str).map_err(|e| NotifyError::InvalidId {
id: raw.id_str.clone(),
reason: e.to_string(),
})?;
Ok(ModelEvent { kind, id })
}
#[cfg(test)]
fn decode_payload<M: crate::model::Model>(payload: &str) -> Result<ModelEvent<M>, NotifyError>
where
M::Pk: FromStr,
<M::Pk as FromStr>::Err: std::fmt::Display,
{
let raw = parse_raw(payload).map_err(|source| NotifyError::PayloadDecode {
raw: payload.to_string(),
source,
})?;
decode_event::<M>(&raw)
}
async fn get_or_start_listener(pool: &DjogiPool) -> Result<Arc<PgListener>, NotifyError> {
let key = pool_key(pool);
loop {
if let Some(listener) = upgrade_existing(registry(), key) {
if !listener.is_failed() {
return Ok(listener);
}
if remove_if_current(registry(), key, &listener) {
break;
}
continue;
}
break;
}
let candidate = Arc::new(spawn_listener(pool).await?);
Ok(install_or_lose(registry(), key, candidate))
}
#[allow(clippy::disallowed_methods)]
async fn spawn_listener(pool: &DjogiPool) -> Result<PgListener, NotifyError> {
let url = pool.url.as_deref().ok_or_else(|| {
NotifyError::ListenerStartFailed(
"DjogiPool::url is None — pool was constructed via internal substrate \
without a URL, so the NOTIFY listener cannot spawn a dedicated \
connection. Use `DjogiPool::builder(url).build()` for adopter-facing \
pools."
.to_string(),
)
})?;
let (client, mut connection) = tokio_postgres::connect(url, tokio_postgres::NoTls)
.await
.map_err(|e| NotifyError::ListenerStartFailed(e.to_string()))?;
let senders: Arc<Mutex<HashMap<String, broadcast::Sender<RawEvent>>>> =
Arc::new(Mutex::new(HashMap::new()));
let failed: Arc<AtomicBool> = Arc::new(AtomicBool::new(false));
let senders_for_task = Arc::clone(&senders);
let exit_guard = WatcherExitGuard {
senders: Arc::clone(&senders),
failed: Arc::clone(&failed),
};
tokio::spawn(async move {
let _exit_guard = exit_guard;
use futures::StreamExt;
let mut stream = futures::stream::poll_fn(move |cx| connection.poll_message(cx));
while let Some(msg) = stream.next().await {
match msg {
Ok(AsyncMessage::Notification(n)) => {
let channel = n.channel();
let payload = n.payload();
let raw = match parse_raw(payload) {
Ok(r) => r,
Err(_) => {
tracing::warn!(
target: "djogi::notify",
channel = %channel,
payload = %payload,
"discarded malformed notify payload"
);
continue;
}
};
let senders_guard = match senders_for_task.lock() {
Ok(guard) => guard,
Err(_poisoned) => {
tracing::error!(
target: "djogi::notify",
"senders mutex poisoned mid-watch; exiting watcher \
(exit guard will clear and publish failed)"
);
break;
}
};
if let Some(tx) = senders_guard.get(channel) {
let _ = tx.send(raw);
}
}
Ok(AsyncMessage::Notice(n)) => {
tracing::debug!(target: "djogi::notify", "postgres notice: {n}");
}
Ok(_) => {}
Err(e) => {
tracing::error!(
target: "djogi::notify",
error = %e,
"notify connection terminated; subscribers will see \
ListenerTerminated on next recv()"
);
break;
}
}
}
});
Ok(PgListener {
senders,
failed,
client,
})
}
fn parse_raw(payload: &str) -> Result<RawEvent, serde_json::Error> {
let v: serde_json::Value = serde_json::from_str(payload)?;
let kind_str = v["kind"]
.as_str()
.ok_or_else(|| serde::de::Error::missing_field("kind"))?
.to_string();
let id_str = v["id"]
.as_str()
.ok_or_else(|| serde::de::Error::missing_field("id"))?
.to_string();
Ok(RawEvent { kind_str, id_str })
}
pub async fn subscribe<M>(pool: &DjogiPool) -> Result<TypedReceiver<M>, NotifyError>
where
M: crate::model::Model + 'static,
M::Pk: FromStr + Send + Sync + 'static,
<M::Pk as FromStr>::Err: std::fmt::Display + Send + Sync + 'static,
{
let listener = get_or_start_listener(pool).await?;
let channel = format!("djogi_{}", M::table_name());
crate::ident::check_plain_ident(&channel, false).map_err(|e| {
NotifyError::ListenerStartFailed(format!("subscribe: invalid channel {channel:?}: {e:?}"))
})?;
if listener.is_failed() {
return Err(NotifyError::ListenerTerminated);
}
let raw_rx = {
let mut senders = listener
.senders
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner());
let tx = senders
.entry(channel.clone())
.or_insert_with(|| broadcast::Sender::new(256));
tx.subscribe()
};
if listener.is_failed() {
let mut senders = listener
.senders
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner());
senders.remove(&channel);
return Err(NotifyError::ListenerTerminated);
}
let listen_sql = format!("LISTEN {channel}");
listener
.client
.simple_query(&listen_sql)
.await
.map_err(|e| NotifyError::ListenerStartFailed(format!("LISTEN failed: {e}")))?;
Ok(TypedReceiver {
raw: raw_rx,
_listener: listener,
_model: PhantomData,
})
}
pub struct TypedReceiver<M: crate::model::Model> {
raw: broadcast::Receiver<RawEvent>,
#[allow(dead_code)] _listener: Arc<PgListener>,
_model: PhantomData<M>,
}
impl<M> TypedReceiver<M>
where
M: crate::model::Model + 'static,
M::Pk: FromStr,
<M::Pk as FromStr>::Err: std::fmt::Display,
{
pub async fn recv(&mut self) -> Result<ModelEvent<M>, NotifyError> {
let raw = self.raw.recv().await.map_err(|e| match e {
broadcast::error::RecvError::Lagged(n) => NotifyError::ChannelLagged { skipped: n },
broadcast::error::RecvError::Closed => NotifyError::ListenerTerminated,
})?;
decode_event::<M>(&raw)
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::descriptor::ModelDescriptor;
use crate::model::Model;
use heeranjid::HeerId;
struct FakeModel;
impl crate::model::__sealed::Sealed for FakeModel {}
#[allow(clippy::manual_async_fn)]
impl Model for FakeModel {
type Pk = HeerId;
type Fields = ();
fn table_name() -> &'static str {
"fakes"
}
fn pk_value(&self) -> &Self::Pk {
unreachable!("decode tests don't call pk_value")
}
fn descriptor() -> &'static ModelDescriptor {
unreachable!("decode tests don't call descriptor")
}
fn get(
_ctx: &mut crate::context::DjogiContext,
_id: Self::Pk,
) -> impl std::future::Future<Output = Result<Self, crate::DjogiError>> + Send {
async { unreachable!() }
}
fn create(
_ctx: &mut crate::context::DjogiContext,
_v: Self,
) -> impl std::future::Future<Output = Result<Self, crate::DjogiError>> + Send {
async { unreachable!() }
}
fn save<'ctx>(
&'ctx mut self,
_ctx: &'ctx mut crate::context::DjogiContext,
) -> impl std::future::Future<Output = Result<(), crate::DjogiError>> + Send + 'ctx
{
async { unreachable!() }
}
fn delete(
self,
_ctx: &mut crate::context::DjogiContext,
) -> impl std::future::Future<Output = Result<(), crate::DjogiError>> + Send {
async { unreachable!() }
}
fn refresh_from_db<'ctx>(
&'ctx self,
_ctx: &'ctx mut crate::context::DjogiContext,
) -> impl std::future::Future<Output = Result<Self, crate::DjogiError>> + Send + 'ctx
{
async { unreachable!() }
}
}
#[test]
fn decode_create_round_trips() {
let payload = r#"{"kind":"create","id":"42"}"#;
let event = decode_payload::<FakeModel>(payload).expect("decode");
assert_eq!(event.kind, EventKind::Created);
assert_eq!(event.id, HeerId::from_i64(42).unwrap());
}
#[test]
fn decode_save_maps_to_updated() {
let payload = r#"{"kind":"save","id":"100"}"#;
let event = decode_payload::<FakeModel>(payload).expect("decode");
assert_eq!(event.kind, EventKind::Updated);
assert_eq!(event.id, HeerId::from_i64(100).unwrap());
}
#[test]
fn decode_delete_round_trips() {
let payload = r#"{"kind":"delete","id":"7"}"#;
let event = decode_payload::<FakeModel>(payload).expect("decode");
assert_eq!(event.kind, EventKind::Deleted);
assert_eq!(event.id, HeerId::from_i64(7).unwrap());
}
#[test]
fn decode_unknown_kind_errors() {
let payload = r#"{"kind":"unknown","id":"1"}"#;
let result = decode_payload::<FakeModel>(payload);
assert!(matches!(result, Err(NotifyError::PayloadDecode { .. })));
}
#[test]
fn decode_invalid_id_errors() {
let payload = r#"{"kind":"create","id":"not-a-number"}"#;
let result = decode_payload::<FakeModel>(payload);
assert!(matches!(result, Err(NotifyError::InvalidId { .. })));
}
#[test]
fn decode_missing_kind_errors() {
let payload = r#"{"id":"1"}"#;
let result = decode_payload::<FakeModel>(payload);
assert!(matches!(result, Err(NotifyError::PayloadDecode { .. })));
}
#[test]
fn decode_malformed_json_errors() {
let payload = "not json at all";
let result = decode_payload::<FakeModel>(payload);
assert!(matches!(result, Err(NotifyError::PayloadDecode { .. })));
}
#[test]
fn upgrade_existing_returns_strong_when_alive() {
let map: Mutex<HashMap<u64, Weak<u32>>> = Mutex::new(HashMap::new());
let live = Arc::new(123u32);
map.lock().unwrap().insert(1, Arc::downgrade(&live));
let upgraded = upgrade_existing(&map, 1).expect("upgrade should succeed while live");
assert_eq!(*upgraded, 123);
assert!(map.lock().unwrap().contains_key(&1));
}
#[test]
fn upgrade_existing_returns_none_and_reaps_after_drop() {
let map: Mutex<HashMap<u64, Weak<u32>>> = Mutex::new(HashMap::new());
{
let dying = Arc::new(456u32);
map.lock().unwrap().insert(2, Arc::downgrade(&dying));
}
assert!(map.lock().unwrap().contains_key(&2));
let result = upgrade_existing(&map, 2);
assert!(result.is_none(), "dangling Weak should fail to upgrade");
assert!(
!map.lock().unwrap().contains_key(&2),
"stale entry should be reaped on the failed-upgrade lookup"
);
}
#[test]
fn upgrade_existing_returns_none_for_missing_key() {
let map: Mutex<HashMap<u64, Weak<u32>>> = Mutex::new(HashMap::new());
assert!(upgrade_existing(&map, 99).is_none());
}
#[test]
fn remove_if_current_removes_matching_entry() {
let map: Mutex<HashMap<u64, Weak<u32>>> = Mutex::new(HashMap::new());
let failed = Arc::new(123u32);
map.lock().unwrap().insert(7, Arc::downgrade(&failed));
assert!(
remove_if_current(&map, 7, &failed),
"matching allocation should be removed"
);
assert!(
!map.lock().unwrap().contains_key(&7),
"registry slot must be empty after matching removal"
);
}
#[test]
fn remove_if_current_preserves_raced_in_replacement() {
let map: Mutex<HashMap<u64, Weak<u32>>> = Mutex::new(HashMap::new());
let failed = Arc::new(11u32);
let healthy = Arc::new(22u32);
map.lock().unwrap().insert(5, Arc::downgrade(&failed));
let observed_failed = upgrade_existing(&map, 5).expect("failed listener still strong");
assert!(Arc::ptr_eq(&failed, &observed_failed));
map.lock().unwrap().insert(5, Arc::downgrade(&healthy));
assert!(
!remove_if_current(&map, 5, &observed_failed),
"stale failed-listener cleanup must not erase a healthy replacement"
);
let remaining = map
.lock()
.unwrap()
.get(&5)
.and_then(Weak::upgrade)
.expect("healthy replacement should remain in registry");
assert!(
Arc::ptr_eq(&healthy, &remaining),
"registry must still point at the raced-in healthy listener"
);
}
#[test]
fn install_or_lose_first_caller_wins_slot() {
let map: Mutex<HashMap<u64, Weak<u32>>> = Mutex::new(HashMap::new());
let candidate = Arc::new(7u32);
let installed = install_or_lose(&map, 10, Arc::clone(&candidate));
assert!(Arc::ptr_eq(&candidate, &installed));
let weak = map.lock().unwrap().get(&10).cloned().unwrap();
let upgraded = weak.upgrade().expect("Weak should upgrade");
assert!(Arc::ptr_eq(&candidate, &upgraded));
}
#[test]
fn install_or_lose_second_caller_loses_to_existing() {
let map: Mutex<HashMap<u64, Weak<u32>>> = Mutex::new(HashMap::new());
let winner = Arc::new(11u32);
let loser = Arc::new(22u32);
let _winner_kept = install_or_lose(&map, 5, Arc::clone(&winner));
let resolved = install_or_lose(&map, 5, Arc::clone(&loser));
assert!(
Arc::ptr_eq(&winner, &resolved),
"racer should resolve to the canonical winner, not its own candidate"
);
assert_eq!(*loser, 22);
let weak = map.lock().unwrap().get(&5).cloned().unwrap();
let upgraded = weak.upgrade().expect("winner Weak should upgrade");
assert!(Arc::ptr_eq(&winner, &upgraded));
}
#[test]
fn install_or_lose_revives_dangling_slot() {
let map: Mutex<HashMap<u64, Weak<u32>>> = Mutex::new(HashMap::new());
{
let transient = Arc::new(33u32);
map.lock().unwrap().insert(8, Arc::downgrade(&transient));
}
assert!(map.lock().unwrap().contains_key(&8));
let fresh = Arc::new(44u32);
let resolved = install_or_lose(&map, 8, Arc::clone(&fresh));
assert!(
Arc::ptr_eq(&fresh, &resolved),
"fresh candidate should replace the dangling Weak"
);
let weak = map.lock().unwrap().get(&8).cloned().unwrap();
let upgraded = weak.upgrade().expect("fresh Weak should upgrade");
assert!(Arc::ptr_eq(&fresh, &upgraded));
}
#[test]
fn pool_id_is_unique_per_build() {
let a = crate::pg::pool::next_pool_id();
let b = crate::pg::pool::next_pool_id();
assert_ne!(
a, b,
"next_pool_id must allocate distinct ids on consecutive calls"
);
}
fn install_sender_and_subscribe(
senders: &Mutex<HashMap<String, broadcast::Sender<RawEvent>>>,
channel: &str,
capacity: usize,
) -> broadcast::Receiver<RawEvent> {
let mut guard = senders.lock().unwrap();
let tx = guard
.entry(channel.to_string())
.or_insert_with(|| broadcast::Sender::new(capacity));
tx.subscribe()
}
#[test]
fn close_all_senders_drops_each_sender_and_wakes_receivers() {
let senders: Mutex<HashMap<String, broadcast::Sender<RawEvent>>> =
Mutex::new(HashMap::new());
let mut rx_a = install_sender_and_subscribe(&senders, "djogi_a", 8);
let mut rx_b = install_sender_and_subscribe(&senders, "djogi_b", 8);
assert!(matches!(
rx_a.try_recv(),
Err(broadcast::error::TryRecvError::Empty)
));
assert!(matches!(
rx_b.try_recv(),
Err(broadcast::error::TryRecvError::Empty)
));
close_all_senders(&senders);
assert!(
senders.lock().unwrap().is_empty(),
"close_all_senders must clear the map"
);
assert!(
matches!(rx_a.try_recv(), Err(broadcast::error::TryRecvError::Closed)),
"receiver A must see Closed once its Sender is dropped from the map"
);
assert!(
matches!(rx_b.try_recv(), Err(broadcast::error::TryRecvError::Closed)),
"receiver B must see Closed once its Sender is dropped from the map"
);
}
#[test]
fn close_all_senders_recovers_from_poisoned_lock() {
let senders: Arc<Mutex<HashMap<String, broadcast::Sender<RawEvent>>>> =
Arc::new(Mutex::new(HashMap::new()));
let mut rx = install_sender_and_subscribe(&senders, "djogi_a", 4);
let poison_target = Arc::clone(&senders);
let _ = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
let _guard = poison_target.lock().unwrap();
panic!("inducing poison");
}));
assert!(
senders.is_poisoned(),
"lock should be poisoned by the panic-while-locked path"
);
close_all_senders(&senders);
let cleared = senders
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner());
assert!(
cleared.is_empty(),
"close_all_senders must clear despite a poisoned lock — \
without recover-via-into_inner the watcher exit guard \
would re-panic instead of cleaning up"
);
assert!(matches!(
rx.try_recv(),
Err(broadcast::error::TryRecvError::Closed)
));
}
#[test]
fn watcher_exit_guard_drop_clears_senders_and_marks_failed() {
let senders: Arc<Mutex<HashMap<String, broadcast::Sender<RawEvent>>>> =
Arc::new(Mutex::new(HashMap::new()));
let failed = Arc::new(AtomicBool::new(false));
let mut rx = install_sender_and_subscribe(&senders, "djogi_phase8_t11_evt", 16);
assert!(matches!(
rx.try_recv(),
Err(broadcast::error::TryRecvError::Empty)
));
assert!(!failed.load(Ordering::Acquire));
let guard = WatcherExitGuard {
senders: Arc::clone(&senders),
failed: Arc::clone(&failed),
};
drop(guard);
assert!(
failed.load(Ordering::Acquire),
"WatcherExitGuard::drop must publish failed=true so \
subscribers reap the dead listener"
);
assert!(
senders.lock().unwrap().is_empty(),
"WatcherExitGuard::drop must clear the senders map"
);
assert!(
matches!(rx.try_recv(), Err(broadcast::error::TryRecvError::Closed)),
"live receivers must see Closed after the watcher exit guard \
fires — pre-fix they hung forever (GH#131)"
);
}
#[test]
fn watcher_exit_guard_fires_on_panic_unwind() {
let senders: Arc<Mutex<HashMap<String, broadcast::Sender<RawEvent>>>> =
Arc::new(Mutex::new(HashMap::new()));
let failed = Arc::new(AtomicBool::new(false));
let mut rx = install_sender_and_subscribe(&senders, "djogi_a", 4);
let senders_for_panic = Arc::clone(&senders);
let failed_for_panic = Arc::clone(&failed);
let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
let _guard = WatcherExitGuard {
senders: senders_for_panic,
failed: failed_for_panic,
};
panic!("watcher task panicked mid-loop");
}));
assert!(
result.is_err(),
"panic should propagate to catch_unwind; if it doesn't, \
the test isn't actually exercising the panic-unwind path"
);
assert!(
failed.load(Ordering::Acquire),
"failed flag must be published on panic unwind"
);
assert!(
senders.lock().unwrap_or_else(|p| p.into_inner()).is_empty(),
"senders map must be cleared on panic unwind"
);
assert!(
matches!(rx.try_recv(), Err(broadcast::error::TryRecvError::Closed)),
"live receivers must see Closed even when the watcher \
panicked rather than exiting cleanly"
);
}
#[test]
fn watcher_exit_guard_with_no_subscribers_is_a_safe_no_op() {
let senders: Arc<Mutex<HashMap<String, broadcast::Sender<RawEvent>>>> =
Arc::new(Mutex::new(HashMap::new()));
let failed = Arc::new(AtomicBool::new(false));
let guard = WatcherExitGuard {
senders: Arc::clone(&senders),
failed: Arc::clone(&failed),
};
drop(guard);
assert!(failed.load(Ordering::Acquire));
assert!(senders.lock().unwrap().is_empty());
}
#[test]
fn pg_listener_is_failed_reflects_failed_flag() {
let failed = Arc::new(AtomicBool::new(false));
let read_is_failed = || failed.load(Ordering::Acquire);
assert!(!read_is_failed(), "fresh listener starts unfailed");
failed.store(true, Ordering::Release);
assert!(
read_is_failed(),
"Acquire load must observe the Release store published by the exit guard"
);
}
}