use std::{
collections::HashMap,
marker::PhantomData,
pin::Pin,
sync::{
Arc, Mutex,
atomic::{AtomicBool, Ordering},
},
task::{Context, Poll},
};
type RegistrySender = Sender<Result<PgTaskId, Error>>;
use apalis_codec::json::JsonCodec;
use apalis_core::{backend::shared::MakeShared, worker::context::WorkerContext};
use diesel::RunQueryDsl;
use futures::{
Stream,
channel::mpsc::{self, Receiver, Sender},
};
use ulid::Ulid;
use crate::{
CompactType, Config, Error, PgPool, PgTask, PgTaskId, PostgresStorage, queries, sink::PgSink,
};
type RegistryEntry = (Ulid, RegistrySender);
type RegistryMap = HashMap<String, Vec<RegistryEntry>>;
type SharedRegistry = Arc<Mutex<RegistryMap>>;
pub struct SharedPostgresStorage<Codec = JsonCodec<CompactType>> {
pool: PgPool,
registry: SharedRegistry,
listener_alive: Arc<AtomicBool>,
_marker: PhantomData<Codec>,
}
impl<Codec> SharedPostgresStorage<Codec> {
#[must_use]
pub fn new(pool: PgPool) -> Self {
let registry: SharedRegistry = Arc::new(Mutex::new(HashMap::new()));
Self {
pool,
registry,
listener_alive: Arc::new(AtomicBool::new(false)),
_marker: PhantomData,
}
}
fn spawn_registry_listener(&self) {
let pool = self.pool.clone();
let registry = self.registry.clone();
let listener_alive = self.listener_alive.clone();
if let Err(error) = std::thread::Builder::new()
.name("apalis-postgres-shared-listener".to_owned())
.spawn(move || {
let mut conn = match pool.get() {
Ok(conn) => conn,
Err(error) => {
exit_listener(
®istry,
&listener_alive,
Some(format!(
"failed to get pooled connection for shared LISTEN: {error}"
)),
);
return;
}
};
if let Err(error) =
diesel::sql_query("LISTEN \"apalis::job::insert\"").execute(&mut conn)
{
exit_listener(
®istry,
&listener_alive,
Some(format!("failed to start shared LISTEN listener: {error}")),
);
return;
}
run_listener_loop(&mut conn, ®istry, &listener_alive);
let _ = diesel::sql_query("UNLISTEN \"apalis::job::insert\"").execute(&mut conn);
})
{
exit_listener(
&self.registry,
&self.listener_alive,
Some(format!("failed to spawn listener: {error}")),
);
}
}
}
fn run_listener_loop(
conn: &mut diesel::r2d2::PooledConnection<
diesel::r2d2::ConnectionManager<diesel::PgConnection>,
>,
registry: &SharedRegistry,
listener_alive: &AtomicBool,
) {
loop {
for notification in conn.notifications_iter() {
let notification = match notification {
Ok(notification) => notification,
Err(error) => {
exit_listener(
registry,
listener_alive,
Some(format!("failed to receive shared notification: {error}")),
);
return;
}
};
let Ok(event) = serde_json::from_str::<crate::InsertEvent>(¬ification.payload)
else {
continue;
};
let (event_queue, ids) = event.into_ids();
let Ok(mut registry) = registry.lock() else {
listener_alive.store(false, Ordering::Release);
return;
};
deliver_to_queue(&mut registry, &event_queue, &ids);
}
match registry.lock() {
Ok(registry) => {
if listener_should_exit(®istry) {
listener_alive.store(false, Ordering::Release);
drop(registry);
return;
}
}
Err(_) => {
listener_alive.store(false, Ordering::Release);
return;
}
}
std::thread::sleep(queries::NOTIFY_LISTENER_POLL_INTERVAL);
}
}
fn exit_listener(registry: &SharedRegistry, listener_alive: &AtomicBool, error: Option<String>) {
match registry.lock() {
Ok(mut guard) => {
if let Some(message) = error {
broadcast_notify_error_locked(&mut guard, message);
}
listener_alive.store(false, Ordering::Release);
drop(guard);
}
Err(_) => {
listener_alive.store(false, Ordering::Release);
}
}
}
fn deliver_to_queue(registry: &mut RegistryMap, queue: &str, ids: &[PgTaskId]) {
if let Some(senders) = registry.get_mut(queue) {
for &id in ids {
senders.retain_mut(|(_, sender)| match sender.try_send(Ok(id)) {
Ok(()) => true,
Err(error) if error.is_disconnected() => false,
Err(_) => true,
});
}
if senders.is_empty() {
registry.remove(queue);
}
}
}
fn listener_should_exit(registry: &RegistryMap) -> bool {
registry.is_empty()
}
fn claim_listener_spawn(listener_alive: &AtomicBool) -> bool {
!listener_alive.swap(true, Ordering::AcqRel)
}
#[cfg(test)]
fn broadcast_notify_error(registry: &SharedRegistry, message: String) {
let Ok(mut guard) = registry.lock() else {
return;
};
broadcast_notify_error_locked(&mut guard, message);
}
fn broadcast_notify_error_locked(registry: &mut RegistryMap, message: String) {
registry.retain(|_, senders| {
senders.retain_mut(|(_, sender)| {
match sender.try_send(Err(Error::NotifyListener(message.clone()))) {
Ok(()) => true,
Err(error) => !error.is_disconnected(),
}
});
!senders.is_empty()
});
}
impl<Codec> std::fmt::Debug for SharedPostgresStorage<Codec> {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("SharedPostgresStorage")
.finish_non_exhaustive()
}
}
#[derive(Debug, thiserror::Error)]
#[non_exhaustive]
pub enum SharedPostgresError {
#[error("registry lock poisoned")]
RegistryLocked,
}
impl<Args, Codec> MakeShared<Args> for SharedPostgresStorage<Codec> {
type Backend = PostgresStorage<Args, Codec, SharedFetcher>;
type Config = Config;
type MakeError = SharedPostgresError;
fn make_shared(&mut self) -> Result<Self::Backend, Self::MakeError>
where
Self::Config: Default,
{
self.make_shared_with_config(Config::new(std::any::type_name::<Args>()))
}
fn make_shared_with_config(
&mut self,
config: Self::Config,
) -> Result<Self::Backend, Self::MakeError> {
let (sender, receiver) =
mpsc::channel(crate::queries::clamp_notify_capacity(config.buffer_size()));
let mut registry = self
.registry
.lock()
.map_err(|_| SharedPostgresError::RegistryLocked)?;
let queue = config.queue().to_string();
let registration_id = Ulid::new();
registry
.entry(queue)
.or_default()
.push((registration_id, sender));
let should_spawn_listener = claim_listener_spawn(&self.listener_alive);
drop(registry);
if should_spawn_listener {
self.spawn_registry_listener();
}
let registration = Arc::new(SharedRegistration {
id: registration_id,
queue: config.queue().to_string(),
registry: self.registry.clone(),
pool: self.pool.clone(),
});
Ok(PostgresStorage {
_marker: PhantomData,
sink: PgSink::new(&self.pool, &config),
pool: self.pool.clone(),
config,
fetcher: SharedFetcher {
receiver,
_registration: registration,
},
lease_token: crate::queries::worker::mint_lease_token().into(),
})
}
}
struct SharedRegistration {
id: Ulid,
queue: String,
registry: SharedRegistry,
pool: PgPool,
}
impl std::fmt::Debug for SharedRegistration {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("SharedRegistration")
.field("queue", &self.queue)
.finish_non_exhaustive()
}
}
impl Drop for SharedRegistration {
fn drop(&mut self) {
let became_empty = match self.registry.lock() {
Ok(mut registry) => {
if let Some(senders) = registry.get_mut(&self.queue) {
senders.retain(|(id, _)| *id != self.id);
if senders.is_empty() {
registry.remove(&self.queue);
}
}
registry.is_empty()
}
Err(_) => false,
};
if became_empty {
let pool = self.pool.clone();
let _ = std::thread::Builder::new()
.name("apalis-postgres-shared-drop".to_owned())
.spawn(move || {
if let Ok(mut conn) = pool.get() {
let _ = diesel::sql_query("SELECT pg_notify('apalis::job::insert', '')")
.execute(&mut conn);
}
});
}
}
}
pub struct SharedFetcher {
receiver: Receiver<Result<PgTaskId, Error>>,
_registration: Arc<SharedRegistration>,
}
impl std::fmt::Debug for SharedFetcher {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("SharedFetcher").finish_non_exhaustive()
}
}
impl Stream for SharedFetcher {
type Item = Result<PgTaskId, Error>;
fn poll_next(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
Pin::new(&mut self.get_mut().receiver).poll_next(cx)
}
}
impl crate::fetcher::PgFetcherSource for SharedFetcher {
const STORAGE_NAME: &'static str = "SharedPostgresStorage";
fn into_compact_stream(
self,
pool: PgPool,
config: Config,
worker: WorkerContext,
lease_token: std::sync::Arc<str>,
) -> apalis_core::backend::TaskStream<PgTask<CompactType>, Error> {
crate::fetcher::notify_backed_compact_stream(
Self::STORAGE_NAME,
self,
pool,
config,
worker,
lease_token,
)
}
}
#[cfg(test)]
mod tests {
use apalis_core::backend::{Backend, BackendExt, shared::MakeShared};
use diesel::{
PgConnection,
r2d2::{ConnectionManager, Pool},
};
use lets_expect::{AssertionError, AssertionResult, *};
use super::*;
struct SharedObservation {
queue: String,
buffer_size: usize,
debug: String,
}
fn unchecked_pool() -> PgPool {
let manager = ConnectionManager::<PgConnection>::new("postgres://127.0.0.1:1/not-used");
Pool::builder()
.max_size(1)
.connection_timeout(std::time::Duration::from_millis(10))
.build_unchecked(manager)
}
fn shared_debug() -> String {
let shared: SharedPostgresStorage = SharedPostgresStorage::new(unchecked_pool());
format!("{shared:?}")
}
fn make_default_shared() -> Result<SharedObservation, SharedPostgresError> {
let mut shared: SharedPostgresStorage = SharedPostgresStorage::new(unchecked_pool());
let storage = <SharedPostgresStorage as MakeShared<String>>::make_shared(&mut shared)?;
Ok(SharedObservation {
queue: storage.config.queue().to_string(),
buffer_size: storage.config.buffer_size(),
debug: format!("{storage:?}"),
})
}
fn make_configured_shared() -> Result<SharedObservation, SharedPostgresError> {
let mut shared: SharedPostgresStorage = SharedPostgresStorage::new(unchecked_pool());
let config = Config::new("shared-unit").set_buffer_size(3);
let storage = <SharedPostgresStorage as MakeShared<String>>::make_shared_with_config(
&mut shared,
config,
)?;
Ok(SharedObservation {
queue: storage.get_queue().to_string(),
buffer_size: storage.config.buffer_size(),
debug: format!("{:?}", storage.fetcher),
})
}
fn shared_trait_surfaces() -> Result<(String, String), SharedPostgresError> {
let mut shared: SharedPostgresStorage = SharedPostgresStorage::new(unchecked_pool());
let config = Config::new("shared-traits");
let storage = <SharedPostgresStorage as MakeShared<String>>::make_shared_with_config(
&mut shared,
config,
)?;
let worker = WorkerContext::new::<()>("shared-trait-worker");
let middleware_name = std::any::type_name_of_val(&storage.middleware()).to_owned();
let stream_name = std::any::type_name_of_val(&storage.poll_compact(&worker)).to_owned();
Ok((middleware_name, stream_name))
}
fn registration_debug_and_drop() -> (String, bool) {
let registry: SharedRegistry = Arc::new(Mutex::new(HashMap::new()));
let (sender, _receiver) = mpsc::channel(1);
let id = Ulid::new();
registry
.lock()
.expect("fresh shared registry is not poisoned")
.insert("shared-registration".to_owned(), vec![(id, sender)]);
let debug = {
let registration = SharedRegistration {
id,
queue: "shared-registration".to_owned(),
registry: registry.clone(),
pool: unchecked_pool(),
};
format!("{registration:?}")
};
let removed = registry
.lock()
.expect("fresh shared registry is not poisoned")
.is_empty();
(debug, removed)
}
fn drop_leaves_remaining(target_queue: &str, sibling_queues: &[&str]) -> usize {
let registry: SharedRegistry = Arc::new(Mutex::new(HashMap::new()));
let target_id = Ulid::new();
{
let mut reg = registry
.lock()
.expect("fresh shared registry is not poisoned");
let (sender, _r) = mpsc::channel(1);
reg.insert(target_queue.to_owned(), vec![(target_id, sender)]);
for sibling in sibling_queues {
let (sender, _r) = mpsc::channel(1);
reg.insert((*sibling).to_owned(), vec![(Ulid::new(), sender)]);
}
}
{
let registration = SharedRegistration {
id: target_id,
queue: target_queue.to_owned(),
registry: registry.clone(),
pool: unchecked_pool(),
};
drop(registration);
}
registry
.lock()
.expect("fresh shared registry is not poisoned")
.len()
}
fn drop_when_registry_empties() -> usize {
drop_leaves_remaining("shared-only", &[])
}
fn drop_when_registry_has_siblings() -> usize {
drop_leaves_remaining("shared-target", &["shared-other-a", "shared-other-b"])
}
fn drop_siblings_keeps_their_identities() -> (usize, bool) {
let registry: SharedRegistry = Arc::new(Mutex::new(HashMap::new()));
let target_id = Ulid::new();
{
let mut reg = registry
.lock()
.expect("fresh shared registry is not poisoned");
let (sender, _r) = mpsc::channel(1);
reg.insert("shared-target".to_owned(), vec![(target_id, sender)]);
for sibling in ["shared-other-a", "shared-other-b"] {
let (sender, _r) = mpsc::channel(1);
reg.insert(sibling.to_owned(), vec![(Ulid::new(), sender)]);
}
}
drop(SharedRegistration {
id: target_id,
queue: "shared-target".to_owned(),
registry: registry.clone(),
pool: unchecked_pool(),
});
let reg = registry
.lock()
.expect("fresh shared registry is not poisoned");
let only_target_removed = !reg.contains_key("shared-target")
&& reg.contains_key("shared-other-a")
&& reg.contains_key("shared-other-b");
(reg.len(), only_target_removed)
}
fn drop_one_of_two_keeps_sibling_sender() -> usize {
let registry: SharedRegistry = Arc::new(Mutex::new(HashMap::new()));
let queue = "shared-coexist".to_owned();
let first_id = Ulid::new();
let second_id = Ulid::new();
let (first_sender, _first_rx) = mpsc::channel(1);
let (second_sender, _second_rx) = mpsc::channel(1);
registry
.lock()
.expect("fresh registry is not poisoned")
.insert(
queue.clone(),
vec![(first_id, first_sender), (second_id, second_sender)],
);
drop(SharedRegistration {
id: first_id,
queue: queue.clone(),
registry: registry.clone(),
pool: unchecked_pool(),
});
let guard = registry.lock().expect("registry is not poisoned");
guard.get(&queue).map(Vec::len).unwrap_or(0)
}
fn drop_with_poisoned_registry() -> usize {
let registry: SharedRegistry = Arc::new(Mutex::new(HashMap::new()));
let queue = "shared-poisoned-drop".to_owned();
let id = Ulid::new();
let (sender, _receiver) = mpsc::channel(1);
registry
.lock()
.expect("fresh registry is not poisoned")
.insert(queue.clone(), vec![(id, sender)]);
let poison_target = registry.clone();
let join = std::thread::spawn(move || {
let _guard = poison_target
.lock()
.expect("fresh registry lock is not poisoned");
panic!("synthetic poisoning panic");
});
let _ = join.join();
drop(SharedRegistration {
id,
queue: queue.clone(),
registry: registry.clone(),
pool: unchecked_pool(),
});
registry
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.get(&queue)
.map(Vec::len)
.unwrap_or(0)
}
fn double_make_shared_same_queue() -> Result<(), SharedPostgresError> {
let mut shared: SharedPostgresStorage = SharedPostgresStorage::new(unchecked_pool());
let config = Config::new("double-make-shared");
let _first = <SharedPostgresStorage as MakeShared<String>>::make_shared_with_config(
&mut shared,
config.clone(),
)?;
let _second = <SharedPostgresStorage as MakeShared<String>>::make_shared_with_config(
&mut shared,
config,
)?;
Ok(())
}
fn broadcast_notify_error_observation() -> (usize, bool) {
let registry: SharedRegistry = Arc::new(Mutex::new(HashMap::new()));
let (alive_sender, _alive_receiver) = mpsc::channel(1);
let (dead_sender, dead_receiver) = mpsc::channel::<Result<PgTaskId, Error>>(1);
drop(dead_receiver);
{
let mut reg = registry.lock().expect("fresh registry is not poisoned");
reg.insert("alive".to_owned(), vec![(Ulid::new(), alive_sender)]);
reg.insert("dead".to_owned(), vec![(Ulid::new(), dead_sender)]);
}
broadcast_notify_error(®istry, "synthetic listener failure".to_owned());
let reg = registry.lock().expect("registry is not poisoned");
let retained = reg.len();
let alive_is_the_sole_survivor = reg.contains_key("alive") && !reg.contains_key("dead");
(retained, alive_is_the_sole_survivor)
}
fn new_task_id() -> PgTaskId {
PgTaskId::new(Ulid::new())
}
fn deliver_prunes_disconnected_sender() -> usize {
let mut registry: RegistryMap = HashMap::new();
let (dead_sender, dead_receiver) = mpsc::channel::<Result<PgTaskId, Error>>(1);
drop(dead_receiver);
registry.insert(
"shared-deliver-dead".to_owned(),
vec![(Ulid::new(), dead_sender)],
);
deliver_to_queue(&mut registry, "shared-deliver-dead", &[new_task_id()]);
registry.len()
}
fn deliver_keeps_full_sender() -> usize {
let mut registry: RegistryMap = HashMap::new();
let (full_sender, _full_receiver) = mpsc::channel::<Result<PgTaskId, Error>>(1);
registry.insert(
"shared-deliver-full".to_owned(),
vec![(Ulid::new(), full_sender)],
);
let ids = [new_task_id(), new_task_id(), new_task_id(), new_task_id()];
deliver_to_queue(&mut registry, "shared-deliver-full", &ids);
registry
.get("shared-deliver-full")
.map(Vec::len)
.unwrap_or(0)
}
fn deliver_broadcasts_id_to_every_live_sender() -> (bool, bool) {
let mut registry: RegistryMap = HashMap::new();
let (first_sender, mut first_receiver) = mpsc::channel::<Result<PgTaskId, Error>>(1);
let (second_sender, mut second_receiver) = mpsc::channel::<Result<PgTaskId, Error>>(1);
registry.insert(
"shared-deliver-fanout".to_owned(),
vec![(Ulid::new(), first_sender), (Ulid::new(), second_sender)],
);
let id = new_task_id();
deliver_to_queue(&mut registry, "shared-deliver-fanout", &[id]);
let got_id = |receiver: &mut Receiver<Result<PgTaskId, Error>>| matches!(receiver.try_recv(), Ok(Ok(got)) if got == id);
(got_id(&mut first_receiver), got_id(&mut second_receiver))
}
fn deliver_to_absent_queue_leaves_others_intact() -> (usize, bool) {
let mut registry: RegistryMap = HashMap::new();
let (sender, _receiver) = mpsc::channel::<Result<PgTaskId, Error>>(1);
registry.insert("shared-other".to_owned(), vec![(Ulid::new(), sender)]);
deliver_to_queue(&mut registry, "shared-absent", &[new_task_id()]);
let unrelated_len = registry.get("shared-other").map(Vec::len).unwrap_or(0);
let absent_created = registry.contains_key("shared-absent");
(unrelated_len, absent_created)
}
fn broadcast_notify_error_keeps_full_sender() -> usize {
let registry: SharedRegistry = Arc::new(Mutex::new(HashMap::new()));
let (mut full_sender, _full_receiver) = mpsc::channel::<Result<PgTaskId, Error>>(1);
while full_sender.try_send(Ok(new_task_id())).is_ok() {}
registry
.lock()
.expect("fresh registry is not poisoned")
.insert(
"shared-error-full".to_owned(),
vec![(Ulid::new(), full_sender)],
);
broadcast_notify_error(®istry, "synthetic listener failure".to_owned());
let reg = registry.lock().expect("registry is not poisoned");
reg.get("shared-error-full").map(Vec::len).unwrap_or(0)
}
fn empty_registry() -> RegistryMap {
HashMap::new()
}
fn registry_with_one_consumer() -> RegistryMap {
let mut registry: RegistryMap = HashMap::new();
let (sender, _receiver) = mpsc::channel::<Result<PgTaskId, Error>>(1);
registry.insert(
"shared-still-active".to_owned(),
vec![(Ulid::new(), sender)],
);
registry
}
fn exit_listener_observation() -> (bool, bool) {
let registry: SharedRegistry = Arc::new(Mutex::new(HashMap::new()));
let (sender, mut receiver) = mpsc::channel::<Result<PgTaskId, Error>>(1);
registry
.lock()
.expect("fresh registry is not poisoned")
.insert("shared-exit".to_owned(), vec![(Ulid::new(), sender)]);
let listener_alive = AtomicBool::new(true);
exit_listener(
®istry,
&listener_alive,
Some("synthetic listener spawn failure".to_owned()),
);
let alive_after = listener_alive.load(Ordering::Acquire);
let error_delivered = matches!(receiver.try_recv(), Ok(Err(_)));
(alive_after, error_delivered)
}
fn listener_spawn_claims() -> (bool, bool) {
let listener_alive = AtomicBool::new(false);
let first = claim_listener_spawn(&listener_alive);
let second = claim_listener_spawn(&listener_alive);
(first, second)
}
fn debug_mentions_type(expected: &'static str) -> impl Fn(&String) -> AssertionResult {
move |debug| {
if debug.contains(expected) {
Ok(())
} else {
Err(AssertionError::new(vec![format!(
"expected debug output containing {expected:?}, got {debug}"
)]))
}
}
}
fn uses_default_queue(result: &SharedObservation) -> AssertionResult {
if result.queue == std::any::type_name::<String>()
&& result.buffer_size == 10
&& result.debug.contains("SharedFetcher")
{
Ok(())
} else {
Err(AssertionError::new(vec![format!(
"unexpected default shared storage: queue={:?}, buffer={}, debug={}",
result.queue, result.buffer_size, result.debug
)]))
}
}
fn uses_configured_queue(result: &SharedObservation) -> AssertionResult {
if result.queue == "shared-unit"
&& result.buffer_size == 3
&& result.debug.contains("SharedFetcher")
{
Ok(())
} else {
Err(AssertionError::new(vec![format!(
"unexpected configured shared storage: queue={:?}, buffer={}, debug={}",
result.queue, result.buffer_size, result.debug
)]))
}
}
fn constructs_backend_traits(result: &(String, String)) -> AssertionResult {
if result.0.contains("PgMiddleware") && result.1.contains("Stream") {
Ok(())
} else {
Err(AssertionError::new(vec![format!(
"unexpected shared trait surfaces: {result:?}"
)]))
}
}
fn removes_registration(result: &(String, bool)) -> AssertionResult {
if result.0.contains("SharedRegistration") && result.1 {
Ok(())
} else {
Err(AssertionError::new(vec![format!(
"expected registration debug and drop cleanup, got {result:?}"
)]))
}
}
fn make_shared_with_poisoned_registry() -> Result<(), SharedPostgresError> {
let mut shared: SharedPostgresStorage = SharedPostgresStorage::new(unchecked_pool());
let registry = shared.registry.clone();
let join = std::thread::spawn(move || {
let _guard = registry
.lock()
.expect("fresh registry lock is not poisoned");
panic!("synthetic poisoning panic");
});
let _ = join.join();
let config = Config::new("poisoned-registry");
<SharedPostgresStorage as MakeShared<String>>::make_shared_with_config(&mut shared, config)
.map(|_| ())
}
fn is_registry_locked(error: &SharedPostgresError) -> AssertionResult {
match error {
SharedPostgresError::RegistryLocked => Ok(()),
}
}
lets_expect! {
expect(shared_debug()) {
to describes_the_shared_factory { debug_mentions_type("SharedPostgresStorage") }
}
expect(make_default_shared()) {
when no_config_is_supplied {
to uses_the_task_type_as_the_namespace { be_ok_and uses_default_queue }
}
}
expect(make_configured_shared()) {
when config_is_supplied {
to exposes_the_queue_and_fetcher { be_ok_and uses_configured_queue }
}
}
expect(shared_trait_surfaces()) {
when backend_traits_are_requested {
to builds_middleware_and_compact_stream { be_ok_and constructs_backend_traits }
}
}
expect(registration_debug_and_drop()) {
when registration_is_dropped {
to removes_the_namespace_from_the_registry { removes_registration }
}
}
expect(drop_when_registry_empties()) {
when dropping_the_last_registration_empties_the_registry {
to leaves_no_remaining_registrations { equal(0) }
}
}
expect(drop_when_registry_has_siblings()) {
when dropping_one_of_several_registrations {
to keeps_sibling_registrations_intact { equal(2) }
}
}
expect(drop_siblings_keeps_their_identities()) {
when dropping_one_of_several_registrations_with_siblings_present {
to removes_only_the_dropped_queue_and_keeps_both_named_siblings {
equal((2_usize, true))
}
}
}
expect(drop_one_of_two_keeps_sibling_sender()) {
when dropping_one_of_two_consumers_on_the_same_queue {
to leaves_the_other_senders_sender_in_place { equal(1) }
}
}
expect(drop_with_poisoned_registry()) {
when the_registry_mutex_is_poisoned {
to leaves_the_registration_in_place_without_waking_the_listener {
equal(1)
}
}
}
expect(double_make_shared_same_queue()) {
when the_same_queue_is_registered_twice {
to accepts_the_second_registration { be_ok }
}
}
expect(broadcast_notify_error_observation()) {
when listener_broadcasts_an_error_to_a_mixed_registry {
to keeps_the_live_queue_and_drops_the_disconnected_one { equal((1_usize, true)) }
}
}
expect(broadcast_notify_error_keeps_full_sender()) {
when a_senders_channel_is_full_but_still_connected {
to keeps_the_back_pressured_sender_registered { equal(1) }
}
}
expect(deliver_prunes_disconnected_sender()) {
when a_queues_only_sender_has_a_dropped_receiver {
to prunes_the_disconnected_sender_and_removes_the_empty_queue { equal(0) }
}
}
expect(deliver_keeps_full_sender()) {
when a_queues_sender_channel_is_full_but_still_connected {
to keeps_the_back_pressured_sender_registered { equal(1) }
}
}
expect(deliver_broadcasts_id_to_every_live_sender()) {
when a_queue_has_multiple_live_consumers {
to broadcasts_the_wakeup_id_to_every_sender { equal((true, true)) }
}
}
expect(deliver_to_absent_queue_leaves_others_intact()) {
when a_notification_targets_a_queue_with_no_registered_consumers {
to leaves_other_queues_untouched_and_creates_no_entry {
equal((1_usize, false))
}
}
}
expect(listener_should_exit(®istry)) {
let registry = empty_registry();
when no_consumers_remain_registered {
to signals_the_listener_to_exit { be_true }
}
when a_consumer_is_still_registered {
let registry = registry_with_one_consumer();
to keeps_the_listener_running { be_false }
}
}
expect(exit_listener_observation()) {
when the_listener_exits_after_a_spawn_or_connection_failure {
to clears_the_alive_flag_and_broadcasts_the_error { equal((false, true)) }
}
}
expect(listener_spawn_claims()) {
when two_registrations_race_to_claim_the_listener_spawn {
to spawns_on_the_first_registration_only { equal((true, false)) }
}
}
expect(make_shared_with_poisoned_registry()) {
when the_registry_mutex_is_poisoned_by_a_panic_in_another_thread {
to surfaces_registry_locked_rather_than_panicking_or_succeeding {
be_err_and is_registry_locked
}
}
}
}
}