use std::fmt;
use serde::{Deserialize, Serialize};
use crate::event::EventId;
use crate::FlowId;
#[derive(Debug, Clone)]
pub struct FailedSubmission {
pub flow_id: Option<FlowId>,
pub failure: SubmissionFailure,
}
impl fmt::Display for FailedSubmission {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
if let Some(fid) = self.flow_id.as_ref() {
write!(f, "{}, flow id: {}", self.failure, fid)?;
} else {
write!(f, "{}", self.failure)?;
}
Ok(())
}
}
impl FailedSubmission {
pub fn is_unprocessable(&self) -> bool {
matches!(self.failure, SubmissionFailure::Unprocessable(_))
}
}
#[derive(Debug, Clone)]
pub enum SubmissionFailure {
Unprocessable(Vec<BatchItemResponse>),
NotAllSubmitted(Vec<BatchItemResponse>),
}
impl SubmissionFailure {
pub fn is_empty(&self) -> bool {
match self {
SubmissionFailure::Unprocessable(ref items) => items.is_empty(),
SubmissionFailure::NotAllSubmitted(ref items) => items.is_empty(),
}
}
pub fn len(&self) -> usize {
match self {
SubmissionFailure::Unprocessable(ref items) => items.len(),
SubmissionFailure::NotAllSubmitted(ref items) => items.len(),
}
}
pub fn iter(&self) -> impl Iterator<Item = &BatchItemResponse> {
match self {
SubmissionFailure::Unprocessable(ref items) => items.iter(),
SubmissionFailure::NotAllSubmitted(ref items) => items.iter(),
}
}
pub fn failed_response_items(&self) -> impl Iterator<Item = &BatchItemResponse> {
self.iter()
.filter(|item| item.publishing_status == PublishingStatus::Failed)
}
pub fn aborted_response_items(&self) -> impl Iterator<Item = &BatchItemResponse> {
self.iter()
.filter(|item| item.publishing_status == PublishingStatus::Aborted)
}
pub fn submitted_response_items(&self) -> impl Iterator<Item = &BatchItemResponse> {
self.iter()
.filter(|item| item.publishing_status == PublishingStatus::Submitted)
}
pub fn non_submitted_response_items(&self) -> impl Iterator<Item = &BatchItemResponse> {
self.iter()
.filter(|item| item.publishing_status != PublishingStatus::Submitted)
}
pub fn stats(&self) -> BatchStats {
let mut stats = BatchStats::default();
for item in self.iter() {
stats.n_items += 1;
match item.publishing_status {
PublishingStatus::Submitted => stats.n_submitted += 1,
PublishingStatus::Failed => {
stats.n_failed += 1;
stats.n_not_submitted += 1;
}
PublishingStatus::Aborted => {
stats.n_aborted += 1;
stats.n_not_submitted += 1;
}
}
}
stats
}
}
impl IntoIterator for SubmissionFailure {
type Item = BatchItemResponse;
type IntoIter = std::vec::IntoIter<BatchItemResponse>;
fn into_iter(self) -> Self::IntoIter {
match self {
SubmissionFailure::NotAllSubmitted(items) => items.into_iter(),
SubmissionFailure::Unprocessable(items) => items.into_iter(),
}
}
}
impl fmt::Display for SubmissionFailure {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
let num_submitted = self.submitted_response_items().count();
let num_failed = self.failed_response_items().count();
let num_aborted = self.aborted_response_items().count();
write!(
f,
"items:{}, submitted:{}, failed: {}, aborted: {}",
self.len(),
num_submitted,
num_failed,
num_aborted
)?;
Ok(())
}
}
#[derive(Debug, Default)]
pub struct BatchStats {
pub n_items: usize,
pub n_submitted: usize,
pub n_failed: usize,
pub n_aborted: usize,
pub n_not_submitted: usize,
}
impl BatchStats {
pub fn all_submitted(n: usize) -> Self {
Self {
n_items: n,
n_submitted: n,
..Self::default()
}
}
pub fn all_not_submitted(n: usize) -> Self {
Self {
n_items: n,
n_not_submitted: n,
..Self::default()
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct BatchItemResponse {
pub eid: Option<EventId>,
pub publishing_status: PublishingStatus,
pub step: Option<PublishingStep>,
#[serde(skip_serializing_if = "Option::is_none")]
pub detail: Option<String>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum PublishingStatus {
Submitted,
Failed,
Aborted,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum PublishingStep {
None,
Validating,
Partitioning,
Enriching,
Publishing,
}