use std::{
fmt::{Debug, Formatter},
ops::Deref,
pin::{Pin, pin},
task::{Context, Poll},
};
use crate::AsyncIterator;
#[derive(Default, Clone, Debug)]
pub enum ProcessResultsStrategy {
#[default]
Partition,
BreakOnError,
}
pub struct ProcessResultsContainer<T, E> {
successes: Vec<T>,
errors: Vec<E>,
}
impl<T, E> Deref for ProcessResultsContainer<T, E> {
type Target = Vec<T>;
fn deref(&self) -> &Self::Target {
self.successes()
}
}
impl<T, E> ProcessResultsContainer<T, E> {
pub fn into_result(self) -> Result<Vec<T>, E> {
if !self.errors.is_empty() {
Err(self.into_errors().remove(0))
} else {
Ok(self.successes)
}
}
pub fn into_successes(self) -> Vec<T> {
self.successes
}
pub fn into_errors(self) -> Vec<E> {
self.errors
}
pub fn successes(&self) -> &Vec<T> {
self.successes.as_ref()
}
pub fn errors(&self) -> &Vec<E> {
self.errors.as_ref()
}
}
impl<T, E> Debug for ProcessResultsContainer<T, E>
where
T: Debug,
E: Debug,
{
fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
f.debug_struct("ProcessResultsContainer")
.field("successes", &self.successes)
.field("errors", &self.errors)
.finish()
}
}
impl<T, E> Clone for ProcessResultsContainer<T, E>
where
T: Clone,
E: Clone,
{
fn clone(&self) -> Self {
Self {
successes: self.successes.clone(),
errors: self.errors.clone(),
}
}
}
impl<T, E> From<(Vec<T>, Vec<E>)> for ProcessResultsContainer<T, E> {
fn from((successes, errors): (Vec<T>, Vec<E>)) -> Self {
Self { successes, errors }
}
}
pub struct ProcessResults<I, T, E>
where
I: AsyncIterator<Item = Result<T, E>>,
{
iter: I,
strategy: ProcessResultsStrategy,
}
impl<I, T, E> ProcessResults<I, T, E>
where
I: AsyncIterator<Item = Result<T, E>>,
{
pub fn new(iter: I) -> ProcessResults<I, T, E> {
Self {
iter,
strategy: ProcessResultsStrategy::default(),
}
}
pub fn with_process_strategy(mut self, strategy: ProcessResultsStrategy) -> Self {
self.strategy = strategy;
self
}
}
impl<I, T, E> Future for ProcessResults<I, T, E>
where
I: AsyncIterator<Item = Result<T, E>> + Unpin,
T: Unpin,
E: Unpin,
{
type Output = ProcessResultsContainer<T, E>;
fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
let strategy = self.strategy.clone();
let mut pinned_fut = pin!(self.get_mut().iter.sync_iter());
loop {
match pinned_fut.as_mut().poll(cx) {
Poll::Pending => {}
Poll::Ready(res) => {
let mut successes = vec![];
let mut errors = vec![];
for item in res {
match item {
Ok(item) => successes.push(item),
Err(error) => {
errors.push(error);
match strategy {
ProcessResultsStrategy::Partition => {}
ProcessResultsStrategy::BreakOnError => {
return Poll::Ready((vec![], errors).into());
}
}
}
}
}
return Poll::Ready((successes, errors).into());
}
}
}
}
}
impl<I, T, E> Debug for ProcessResults<I, T, E>
where
I: AsyncIterator<Item = Result<T, E>> + Debug,
{
fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
f.debug_struct("ProcessResults")
.field("iter", &self.iter)
.field("strategy", &self.strategy)
.finish()
}
}