use std::fmt;
use serde::{Deserialize, Serialize};
use crate::event::EventId;
use crate::FlowId;
#[derive(Debug, Default, Clone, Serialize, Deserialize)]
pub struct BatchResponse {
#[serde(skip_serializing_if = "Option::is_none")]
pub flow_id: Option<FlowId>,
pub batch_items: Vec<BatchItemResponse>,
}
impl BatchResponse {
pub fn is_empty(&self) -> bool {
self.batch_items.is_empty()
}
pub fn len(&self) -> usize {
self.batch_items.len()
}
pub fn failed_response_items(&self) -> impl Iterator<Item = &BatchItemResponse> {
self.batch_items
.iter()
.filter(|item| item.publishing_status == PublishingStatus::Failed)
}
pub fn aborted_response_items(&self) -> impl Iterator<Item = &BatchItemResponse> {
self.batch_items
.iter()
.filter(|item| item.publishing_status == PublishingStatus::Aborted)
}
pub fn submitted_response_items(&self) -> impl Iterator<Item = &BatchItemResponse> {
self.batch_items
.iter()
.filter(|item| item.publishing_status == PublishingStatus::Submitted)
}
pub fn non_submitted_response_items(&self) -> impl Iterator<Item = &BatchItemResponse> {
self.batch_items
.iter()
.filter(|item| item.publishing_status != PublishingStatus::Submitted)
}
pub fn iter(&self) -> impl Iterator<Item = &BatchItemResponse> {
self.batch_items.iter()
}
pub fn stats(&self) -> BatchStats {
let mut stats = BatchStats::default();
for item in &self.batch_items {
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 BatchResponse {
type Item = BatchItemResponse;
type IntoIter = std::vec::IntoIter<BatchItemResponse>;
fn into_iter(self) -> Self::IntoIter {
self.batch_items.into_iter()
}
}
impl fmt::Display for BatchResponse {
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,
"BatchResponse(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 {
let mut me = Self::default();
me.n_items = n;
me.n_submitted = n;
me
}
pub fn all_not_submitted(n: usize) -> Self {
let mut me = Self::default();
me.n_items = n;
me.n_not_submitted = n;
me
}
}
#[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,
}