#[cfg(feature = "logging")]
#[allow(dead_code)]
pub(crate) mod log_capture {
use std::collections::HashMap;
use std::fmt::Debug;
use std::sync::{Arc, Mutex};
use tracing::Subscriber;
use tracing::field::{Field, Visit};
use tracing::subscriber::DefaultGuard;
use tracing_subscriber::Layer;
use tracing_subscriber::layer::{Context, SubscriberExt as _};
pub(crate) type Events = Arc<Mutex<Vec<HashMap<String, String>>>>;
#[derive(Default)]
struct FieldGrab(HashMap<String, String>);
impl Visit for FieldGrab {
fn record_str(&mut self, field: &Field, value: &str) {
self.0.insert(field.name().to_owned(), value.to_owned());
}
fn record_u64(&mut self, field: &Field, value: u64) {
self.0.insert(field.name().to_owned(), value.to_string());
}
fn record_debug(&mut self, field: &Field, value: &dyn Debug) {
self.0
.entry(field.name().to_owned())
.or_insert_with(|| format!("{value:?}"));
}
}
struct Capture(Events);
impl<S: Subscriber> Layer<S> for Capture {
fn on_event(&self, event: &tracing::Event<'_>, _ctx: Context<'_, S>) {
let mut grab = FieldGrab::default();
event.record(&mut grab);
self.0.lock().unwrap().push(grab.0);
}
}
pub(crate) fn start() -> (Events, DefaultGuard) {
let events: Events = Arc::new(Mutex::new(Vec::new()));
let guard = tracing::subscriber::set_default(
tracing_subscriber::registry().with(Capture(Arc::clone(&events))),
);
(events, guard)
}
pub(crate) fn find(events: &Events, message: &str) -> HashMap<String, String> {
let captured = events.lock().unwrap();
captured
.iter()
.find(|fields| fields.get("message").is_some_and(|m| m == message))
.cloned()
.unwrap_or_else(|| panic!("no `{message}` event was emitted"))
}
}
#[cfg(all(feature = "memory", feature = "json"))]
pub(crate) mod batch {
use futures::StreamExt;
use crate::memory::{MemoryBroker, MemoryMessage, MemorySubscriber};
use crate::{BatchSubscriber, OutgoingMessage, Publisher};
pub(crate) async fn publish_numbers(broker: &MemoryBroker, name: &str, numbers: &[u32]) {
for n in numbers {
publish_payloads(broker, name, &[&serde_json::to_vec(n).unwrap()]).await;
}
}
pub(crate) async fn publish_payloads(broker: &MemoryBroker, name: &str, payloads: &[&[u8]]) {
let publisher = broker.publisher();
for payload in payloads {
publisher
.publish(OutgoingMessage::new(name, payload))
.await
.unwrap();
}
}
pub(crate) async fn pull_batch(sub: &mut MemorySubscriber) -> Vec<MemoryMessage> {
let mut stream =
std::pin::pin!(sub.batches(std::num::NonZeroUsize::new(64).expect("64 is nonzero")));
stream.next().await.unwrap().unwrap()
}
}