use std::fmt;
use std::future::Future;
use std::sync::{Arc, Mutex, MutexGuard};
use std::time::Duration;
use super::{FlowError, Scope, Tenant, at_most, each, within};
const TENANT: &str = "fan-out";
#[derive(Debug)]
#[non_exhaustive]
pub enum FanOutError<E> {
Item {
index: usize,
error: E,
},
TimedOut {
after: Duration,
},
Flow(FlowError),
}
impl<E: fmt::Display> fmt::Display for FanOutError<E> {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
match *self {
Self::Item { index, ref error } => write!(formatter, "item #{index} failed: {error}"),
Self::TimedOut { after } => write!(formatter, "timed out after {after:?}"),
Self::Flow(ref error) => write!(formatter, "{error}"),
}
}
}
impl<E: std::error::Error + 'static> std::error::Error for FanOutError<E> {
fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
match *self {
Self::Item { ref error, .. } => Some(error),
Self::Flow(ref error) => Some(error),
Self::TimedOut { .. } => None,
}
}
}
#[must_use = "a fan-out does nothing until `run` is awaited"]
#[derive(Debug)]
pub struct FanOut<I> {
items: I,
limit: Option<usize>,
deadline: Option<Duration>,
}
impl<I: IntoIterator> FanOut<I> {
pub fn new(items: I) -> Self {
Self {
items,
limit: None,
deadline: None,
}
}
pub fn at_most(mut self, limit: usize) -> Self {
self.limit = Some(limit);
self
}
pub fn within(mut self, deadline: Duration) -> Self {
self.deadline = Some(deadline);
self
}
pub async fn run<T, E, F, Fut>(self, body: F) -> Result<Vec<T>, FanOutError<E>>
where
F: Fn(I::Item) -> Fut,
Fut: Future<Output = Result<T, E>>,
{
let scope = Scope::root(Tenant::new(TENANT).map_err(FanOutError::Flow)?);
let limit = match self.limit {
Some(written) => Some(at_most(written).map_err(FanOutError::Flow)?),
None => None,
};
let first: Arc<Mutex<Option<(usize, E)>>> = Arc::new(Mutex::new(None));
let body = &body;
let work = each(
&scope,
"run",
limit,
self.items.into_iter().enumerate(),
|_step, (index, item)| {
let first = Arc::clone(&first);
let running = body(item);
async move {
match running.await {
Ok(value) => Ok(value),
Err(error) => {
lock(&first).get_or_insert((index, error));
Err(FlowError::failed("a fan-out item failed"))
}
}
}
},
);
let outcome = match self.deadline {
Some(deadline) => within(&scope, "deadline", deadline, work).await,
None => work.await,
};
match outcome {
Ok(values) => Ok(values),
Err(flow) => Err(lock(&first).take().map_or_else(
|| match flow {
FlowError::TimedOut { after, .. } => FanOutError::TimedOut { after },
other => FanOutError::Flow(other),
},
|(index, error)| FanOutError::Item { index, error },
)),
}
}
}
fn lock<E>(cell: &Mutex<Option<(usize, E)>>) -> MutexGuard<'_, Option<(usize, E)>> {
crate::journal::owner::lock(cell)
}