use std::sync::Arc;
use crate::batch::{BatchCensus, BatchItem, BatchReport, BatchStatus, ItemOutcome, ItemSource};
use crate::core::{BatchId, PlanIR, RunId, RuntimeError, Spend};
use super::executor::{RunOutcome, RunStatus, Runtime};
use super::{Mode, metrics};
#[derive(Debug)]
pub struct BatchSpec {
pub plan: PlanIR,
pub source: Arc<dyn ItemSource>,
pub page: usize,
pub max_items: Option<u64>,
}
impl BatchSpec {
pub fn new(plan: PlanIR, source: Arc<dyn ItemSource>) -> Self {
Self {
plan,
source,
page: 256,
max_items: None,
}
}
#[must_use]
pub const fn page(mut self, n: usize) -> Self {
self.page = n;
self
}
#[must_use]
pub const fn max_items(mut self, n: u64) -> Self {
self.max_items = Some(n);
self
}
}
impl Runtime {
pub async fn run_batch(
&self,
id: BatchId,
spec: &BatchSpec,
) -> Result<BatchReport, RuntimeError> {
let store = self.batches().ok_or_else(|| {
RuntimeError::PlanContract(
"batches need a batch store — build the runtime with `.batches(store)`".into(),
)
})?;
crate::plan::validate(&spec.plan, &self.contract())
.map_err(|e| RuntimeError::PlanContract(e.to_string()))?;
store
.open(id, &spec.plan.digest().to_hex())
.await
.map_err(RuntimeError::from_store)?;
let mut after = store.cursor(id).await.map_err(RuntimeError::from_store)?;
let mut processed = 0_u64;
loop {
if spec.max_items.is_some_and(|m| processed >= m) {
break;
}
let page = store_page(spec, after.as_deref()).await?;
if page.is_empty() {
store
.mark_exhausted(id)
.await
.map_err(RuntimeError::from_store)?;
break;
}
for item in page {
after = Some(item.key.clone());
let outcome = self.run_item(id, spec, &item).await?;
if outcome.is_terminal() {
processed += 1;
}
if spec.max_items.is_some_and(|m| processed >= m) {
break;
}
}
}
self.batch_report(id).await
}
async fn run_item(
&self,
id: BatchId,
spec: &BatchSpec,
item: &BatchItem,
) -> Result<ItemOutcome, RuntimeError> {
let store = self
.batches()
.ok_or_else(|| RuntimeError::PlanContract("no batch store".into()))?;
let reserved = store
.reserve(id, &item.key, RunId::generate())
.await
.map_err(RuntimeError::from_store)?;
if let Some(done) = reserved.outcome.as_ref().filter(|o| o.is_terminal()) {
return Ok(done.clone());
}
let outcome = if reserved.outcome.is_some() || self.has_journal(reserved.run).await? {
self.replay(reserved.run, Mode::Resume).await
} else {
self.admit_plan_as(
reserved.run,
spec.plan.clone(),
crate::core::Tainted::with_label(item.input.clone(), item.admission_label(id)),
super::executor::RunTerms::default(),
super::executor::Entry::Outside,
)
.await
};
let out = outcome?;
let (result, spend) = classify_item(out);
store
.record(id, &item.key, &result, spend)
.await
.map_err(RuntimeError::from_store)?;
self.meter().count(metrics::BATCH_ITEMS, result.as_str());
Ok(result)
}
async fn has_journal(&self, run: RunId) -> Result<bool, RuntimeError> {
match self.journal().head(run).await {
Ok(head) => Ok(head.seq > 0),
Err(e) => Err(RuntimeError::from_store(e)),
}
}
pub async fn batch_report(&self, id: BatchId) -> Result<BatchReport, RuntimeError> {
let store = self
.batches()
.ok_or_else(|| RuntimeError::PlanContract("batches need a batch store".into()))?;
if store
.plan_digest(id)
.await
.map_err(RuntimeError::from_store)?
.is_none()
{
return Err(RuntimeError::from_store(crate::core::StoreError::NotFound(
format!("batch {id}"),
)));
}
let census = store.census(id).await.map_err(RuntimeError::from_store)?;
let cursor = store.cursor(id).await.map_err(RuntimeError::from_store)?;
let exhausted = store
.is_exhausted(id)
.await
.map_err(RuntimeError::from_store)?;
Ok(BatchReport {
id,
status: status_of(&census, exhausted),
in_flight: census.in_flight + census.suspended,
exhausted: census.exhausted,
withheld: census.withheld,
spend: census.spend,
cursor,
})
}
}
fn status_of(c: &BatchCensus, source_exhausted: bool) -> BatchStatus {
if !source_exhausted || c.in_flight > 0 || c.suspended > 0 || c.exhausted > 0 || c.withheld > 0
{
return BatchStatus::Running;
}
BatchStatus::Completed {
succeeded: c.succeeded,
failed: c.failed,
quarantined: c.quarantined,
}
}
fn classify_item(out: RunOutcome) -> (ItemOutcome, Spend) {
let spend = out.spend();
let item = match out.status {
RunStatus::Succeeded => ItemOutcome::Succeeded,
RunStatus::Suspended(r) => ItemOutcome::Suspended(r.to_string()),
RunStatus::Quarantined(why) => ItemOutcome::Quarantined(why),
RunStatus::Abandoned { actor, reason } => {
ItemOutcome::Quarantined(format!("abandoned by '{actor}': {reason}"))
}
RunStatus::Failed(why) | RunStatus::Replanning(why) => ItemOutcome::Failed(why),
RunStatus::Exhausted(limit) => ItemOutcome::Exhausted(limit.to_string()),
RunStatus::Withheld { subject, reason } => {
ItemOutcome::Withheld(format!("authority '{subject}' withdrawn: {reason}"))
}
RunStatus::Swept | RunStatus::BrokeGlass { .. } | RunStatus::Observed => {
ItemOutcome::Quarantined(
"this item's run id resolves to one of the plane's own records, or to \
a session it only observed, rather than to the item's run — the \
reservation and the journal disagree"
.to_owned(),
)
}
RunStatus::Cancelled { actor, reason } => {
ItemOutcome::Failed(format!("cancelled by '{actor}': {reason}"))
}
};
(item, spend)
}
async fn store_page(spec: &BatchSpec, after: Option<&str>) -> Result<Vec<BatchItem>, RuntimeError> {
let page = spec
.source
.next(after, spec.page)
.await
.map_err(|e| RuntimeError::PlanContract(e.to_string()))?;
in_cursor_order(after, &page)?;
Ok(page)
}
fn in_cursor_order(after: Option<&str>, page: &[BatchItem]) -> Result<(), RuntimeError> {
let mut previous = after;
for item in page {
if let Some(p) = previous
&& item.key.as_bytes() <= p.as_bytes()
{
return Err(RuntimeError::PlanContract(format!(
"batch source returned key '{}' after '{p}': keys must strictly \
increase in byte order, which is the order the resume cursor \
compares — a source paging in any other order skips or repeats \
items on resume",
item.key
)));
}
previous = Some(item.key.as_str());
}
Ok(())
}