use aion_core::{
Event, HealthSample, HealthStatus, Payload, RunId, WorkloopSpec, current_lease_terminal,
};
use aion_store::workloop::InvariantStateRecord;
use chrono::{DateTime, Utc};
use super::error::WorkloopError;
use crate::durability::{Recorder, WorkflowStartRecord};
use crate::lifecycle::continue_as_new::guard_no_pending_work;
use crate::time::retire_run_deadline;
#[derive(Clone, Debug)]
pub struct WorkloopIterationClose {
pub routes: Vec<String>,
pub carry: Payload,
pub invariant_states: Vec<(String, Payload)>,
}
#[derive(Clone, Debug)]
pub struct IterationOutcome {
pub samples: Vec<HealthSample>,
pub next_run_id: RunId,
pub records: Vec<InvariantStateRecord>,
}
#[must_use]
pub fn derive_health_samples(
spec: &WorkloopSpec,
routes: &[String],
window_seq: Option<u64>,
) -> Vec<HealthSample> {
spec.invariants()
.iter()
.map(|invariant| {
let confirmed = routes
.iter()
.any(|route| invariant.confirms.contains(route));
HealthSample {
invariant: invariant.name.clone(),
status: if confirmed {
HealthStatus::Confirmed
} else {
HealthStatus::Unconfirmed
},
window_seq,
}
})
.collect()
}
#[must_use]
pub fn failed_iteration_samples(spec: &WorkloopSpec, window_seq: Option<u64>) -> Vec<HealthSample> {
spec.invariants()
.iter()
.map(|invariant| HealthSample {
invariant: invariant.name.clone(),
status: HealthStatus::Unconfirmed,
window_seq,
})
.collect()
}
fn invariant_records(
spec: &WorkloopSpec,
close: &WorkloopIterationClose,
loop_id: &aion_core::WorkflowId,
window_seq: Option<u64>,
recorded_at: DateTime<Utc>,
) -> Result<Vec<InvariantStateRecord>, WorkloopError> {
close
.invariant_states
.iter()
.map(|(name, payload)| {
let invariant = spec
.invariants()
.iter()
.find(|invariant| &invariant.name == name)
.ok_or_else(|| WorkloopError::UndeclaredInvariant {
loop_id: loop_id.clone(),
invariant: name.clone(),
})?;
Ok(InvariantStateRecord {
loop_id: loop_id.clone(),
invariant: name.clone(),
payload: payload.clone(),
record_type: invariant.record_type.clone(),
window_seq,
recorded_at,
})
})
.collect()
}
fn current_generation_identity(
history: &[Event],
) -> Result<(String, aion_core::PackageVersion), WorkloopError> {
history
.iter()
.rev()
.find_map(|event| {
if let Event::WorkflowStarted {
workflow_type,
package_version,
..
} = event
{
Some((workflow_type.clone(), package_version.clone()))
} else {
None
}
})
.ok_or_else(|| WorkloopError::Engine {
reason: "workloop history has no WorkflowStarted".to_owned(),
})
}
#[derive(Clone, Copy, Debug)]
pub struct IterationContext<'a> {
pub history: &'a [Event],
pub run_id: &'a RunId,
pub loop_id: &'a aion_core::WorkflowId,
pub spec: &'a WorkloopSpec,
pub window_seq: Option<u64>,
pub recorded_at: DateTime<Utc>,
}
pub async fn close_iteration(
recorder: &mut Recorder,
context: IterationContext<'_>,
close: WorkloopIterationClose,
) -> Result<IterationOutcome, WorkloopError> {
let IterationContext {
history,
run_id,
loop_id,
spec,
window_seq,
recorded_at,
} = context;
if current_lease_terminal(history).is_some() {
return Err(WorkloopError::Engine {
reason: format!("workloop {loop_id} run {run_id} already recorded a terminal"),
});
}
guard_no_pending_work(aion_core::run_segment(history, run_id)).map_err(|error| {
WorkloopError::Engine {
reason: error.to_string(),
}
})?;
let samples = derive_health_samples(spec, &close.routes, window_seq);
let records = invariant_records(spec, &close, loop_id, window_seq, recorded_at)?;
let (workflow_type, package_version) = current_generation_identity(history)?;
let next_run_id = RunId::new_v4();
recorder
.record_workloop_iteration_boundary(
recorded_at,
close.routes.clone(),
samples.clone(),
close.carry.clone(),
run_id.clone(),
WorkflowStartRecord {
workflow_type,
input: close.carry,
run_id: next_run_id.clone(),
parent_run_id: Some(run_id.clone()),
parent_workflow_id: None,
package_version,
},
)
.await?;
retire_run_deadline(recorder, history, run_id).await?;
Ok(IterationOutcome {
samples,
next_run_id,
records,
})
}
#[cfg(test)]
mod tests {
use std::time::Duration;
use aion_core::{InvariantSpec, ToleranceSpec, WorkloopArming};
use super::*;
fn spec() -> Result<WorkloopSpec, Box<dyn std::error::Error>> {
Ok(WorkloopSpec::new(
WorkloopArming::every(Duration::from_secs(1500))?,
vec![
InvariantSpec {
name: String::from("serving"),
record_type: String::from("ServeState"),
tolerance: ToleranceSpec::count(3),
confirms: vec![String::from("sweep"), String::from("dispatch")],
},
InvariantSpec {
name: String::from("drained"),
record_type: String::from("DrainState"),
tolerance: ToleranceSpec::count(0),
confirms: vec![String::from("drain")],
},
],
Duration::from_secs(86_400),
)?)
}
#[test]
fn every_invariant_is_sampled_on_the_same_tick() -> Result<(), Box<dyn std::error::Error>> {
let samples = derive_health_samples(
&spec()?,
&[String::from("sweep"), String::from("start")],
Some(4),
);
assert_eq!(samples.len(), 2);
assert_eq!(samples[0].invariant, "serving");
assert_eq!(samples[0].status, HealthStatus::Confirmed);
assert_eq!(samples[0].window_seq, Some(4));
assert_eq!(samples[1].invariant, "drained");
assert_eq!(samples[1].status, HealthStatus::Unconfirmed);
Ok(())
}
#[test]
fn a_failed_iteration_reds_every_invariant() -> Result<(), Box<dyn std::error::Error>> {
let samples = failed_iteration_samples(&spec()?, None);
assert!(
samples
.iter()
.all(|sample| sample.status == HealthStatus::Unconfirmed)
);
assert_eq!(samples.len(), 2);
Ok(())
}
}