use std::fmt::Debug;
use async_trait::async_trait;
use serde::{Deserialize, Serialize};
use serde_json::Value;
use crate::core::{BatchId, RunId, Spend, StoreError};
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct BatchItem {
pub key: String,
pub input: Value,
}
impl BatchItem {
pub fn new(key: impl Into<String>, input: Value) -> Self {
Self {
key: key.into(),
input,
}
}
}
#[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),
}
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",
}
}
#[must_use]
pub const fn is_terminal(&self) -> bool {
!matches!(self, Self::Suspended(_) | Self::Exhausted(_))
}
}
#[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 spend: Spend,
pub cursor: Option<String>,
}
#[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>;
}
#[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 in_flight: u64,
pub spend: Spend,
}
impl BatchCensus {
#[must_use]
pub const fn terminal(&self) -> u64 {
self.succeeded + self.failed + self.quarantined
}
}