use std::fmt::Debug;
use async_trait::async_trait;
use serde::{Deserialize, Serialize};
use serde_json::Value;
use crate::core::{BatchId, Label, RunId, Spend, StoreError};
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct BatchItem {
pub key: String,
pub input: Value,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub label: Option<Label>,
}
impl BatchItem {
pub fn new(key: impl Into<String>, input: Value) -> Self {
Self {
key: key.into(),
input,
label: None,
}
}
#[must_use]
pub fn labelled(mut self, label: Label) -> Self {
self.label = Some(label);
self
}
#[must_use]
pub fn admission_label(&self, batch: BatchId) -> Label {
self.label.clone().unwrap_or_else(|| {
Label::untrusted(crate::core::SourceId::new(format!("batch:{batch}")))
})
}
}
#[derive(Debug, Clone, thiserror::Error)]
#[error("batch source: {0}")]
pub struct SourceError(String);
impl SourceError {
pub fn new(detail: impl Into<String>) -> Self {
Self(detail.into())
}
}
#[async_trait]
pub trait ItemSource: Send + Sync + Debug {
async fn next(&self, after: Option<&str>, limit: usize) -> Result<Vec<BatchItem>, SourceError>;
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum ItemOutcome {
Succeeded,
Failed(String),
Quarantined(String),
Suspended(String),
Exhausted(String),
Withheld(String),
}
impl ItemOutcome {
#[must_use]
pub const fn as_str(&self) -> &'static str {
match self {
Self::Succeeded => "succeeded",
Self::Failed(_) => "failed",
Self::Quarantined(_) => "quarantined",
Self::Suspended(_) => "suspended",
Self::Exhausted(_) => "exhausted",
Self::Withheld(_) => "withheld",
}
}
#[must_use]
pub fn all(detail: String) -> [Self; 6] {
[
Self::Succeeded,
Self::Failed(detail.clone()),
Self::Quarantined(detail.clone()),
Self::Suspended(detail.clone()),
Self::Exhausted(detail.clone()),
Self::Withheld(detail),
]
}
#[must_use]
pub fn parse(tag: &str, detail: String) -> Option<Self> {
Self::all(detail).into_iter().find(|o| o.as_str() == tag)
}
#[must_use]
pub const fn is_terminal(&self) -> bool {
!matches!(
self,
Self::Suspended(_) | Self::Exhausted(_) | Self::Withheld(_)
)
}
#[must_use]
pub fn terminal_tags() -> Vec<&'static str> {
Self::all(String::new())
.iter()
.filter(|o| o.is_terminal())
.map(Self::as_str)
.collect()
}
#[must_use]
pub const fn is_settled(&self) -> bool {
match self {
Self::Succeeded => true,
Self::Failed(_)
| Self::Quarantined(_)
| Self::Suspended(_)
| Self::Exhausted(_)
| Self::Withheld(_) => false,
}
}
#[must_use]
pub fn settled_tags() -> Vec<&'static str> {
Self::all(String::new())
.iter()
.filter(|o| o.is_settled())
.map(Self::as_str)
.collect()
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum BatchStatus {
Running,
Completed {
succeeded: u64,
failed: u64,
quarantined: u64,
},
}
impl BatchStatus {
#[must_use]
pub const fn everything_settled(&self) -> bool {
matches!(
self,
Self::Completed {
failed: 0,
quarantined: 0,
..
}
)
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ItemRecord {
pub key: String,
pub run: RunId,
pub outcome: Option<ItemOutcome>,
pub spend: Spend,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct BatchReport {
pub id: BatchId,
pub status: BatchStatus,
pub in_flight: u64,
pub exhausted: u64,
pub withheld: u64,
pub spend: Spend,
pub cursor: Option<String>,
}
impl BatchReport {
#[must_use]
pub const fn needs_attention(&self) -> bool {
self.in_flight > 0
|| self.exhausted > 0
|| self.withheld > 0
|| self.failed_or_quarantined() > 0
}
#[must_use]
pub const fn failed_or_quarantined(&self) -> u64 {
match self.status {
BatchStatus::Running => 0,
BatchStatus::Completed {
failed,
quarantined,
succeeded: _,
} => failed + quarantined,
}
}
}
#[async_trait]
pub trait BatchStore: Send + Sync + Debug {
fn tenant(&self) -> &str;
async fn open(&self, id: BatchId, plan_digest: &str) -> Result<(), StoreError>;
async fn plan_digest(&self, id: BatchId) -> Result<Option<String>, StoreError>;
async fn mark_exhausted(&self, id: BatchId) -> Result<(), StoreError>;
async fn is_exhausted(&self, id: BatchId) -> Result<bool, StoreError>;
async fn reserve(
&self,
batch: BatchId,
key: &str,
run: RunId,
) -> Result<ItemRecord, StoreError>;
async fn record(
&self,
batch: BatchId,
key: &str,
outcome: &ItemOutcome,
spend: Spend,
) -> Result<(), StoreError>;
async fn cursor(&self, batch: BatchId) -> Result<Option<String>, StoreError>;
async fn census(&self, batch: BatchId) -> Result<BatchCensus, StoreError>;
async fn items(&self, batch: BatchId, limit: usize) -> Result<Vec<ItemRecord>, StoreError>;
async fn items_needing_attention(
&self,
batch: BatchId,
limit: usize,
) -> Result<Vec<ItemRecord>, StoreError>;
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub struct BatchCensus {
pub succeeded: u64,
pub failed: u64,
pub quarantined: u64,
pub suspended: u64,
pub exhausted: u64,
pub withheld: u64,
pub in_flight: u64,
pub spend: Spend,
}
impl BatchCensus {
#[must_use]
pub const fn terminal(&self) -> u64 {
self.succeeded + self.failed + self.quarantined
}
}