#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
pub enum FailurePolicy {
FailFast,
#[default]
CollectAll,
Quorum(usize),
BestEffort,
}
#[derive(Clone, Debug, Default)]
pub struct ParallelOptions {
pub max_concurrency: usize,
pub failure_policy: FailurePolicy,
pub item_timeout: Option<std::time::Duration>,
pub total_timeout: Option<std::time::Duration>,
pub cancellation: Option<crate::harness::cancel::CancellationToken>,
}
impl ParallelOptions {
pub fn with_max_concurrency(mut self, n: usize) -> Self {
self.max_concurrency = n;
self
}
pub fn with_failure_policy(mut self, policy: FailurePolicy) -> Self {
self.failure_policy = policy;
self
}
pub fn with_item_timeout(mut self, timeout: std::time::Duration) -> Self {
self.item_timeout = Some(timeout);
self
}
pub fn with_total_timeout(mut self, timeout: std::time::Duration) -> Self {
self.total_timeout = Some(timeout);
self
}
pub fn with_cancellation(mut self, token: crate::harness::cancel::CancellationToken) -> Self {
self.cancellation = Some(token);
self
}
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct ItemOutcome<T> {
pub index: usize,
pub result: std::result::Result<T, String>,
}
impl<T> ItemOutcome<T> {
pub fn is_ok(&self) -> bool {
self.result.is_ok()
}
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct ParallelOutcome<T> {
pub outcomes: Vec<ItemOutcome<T>>,
}
impl<T> ParallelOutcome<T> {
pub fn success_count(&self) -> usize {
self.outcomes.iter().filter(|o| o.is_ok()).count()
}
pub fn failure_count(&self) -> usize {
self.outcomes.iter().filter(|o| !o.is_ok()).count()
}
pub fn successes(&self) -> Vec<&T> {
self.outcomes
.iter()
.filter_map(|o| o.result.as_ref().ok())
.collect()
}
pub fn into_successes(self) -> Vec<T> {
self.outcomes
.into_iter()
.filter_map(|o| o.result.ok())
.collect()
}
}