use std::collections::BTreeSet;
use std::time::Duration;
use meerkat_mob::{MobError, MobRunStatus, RunId};
use crate::mob_handle_runtime::MobRuntimeError;
pub const MOB_STOP_FLOW_SETTLE_BUDGET: Duration = Duration::from_secs(10);
#[derive(Debug)]
pub struct MobStopFlowRunsUnsettled {
pub runs: Vec<RunId>,
pub last_refusal: Option<MobError>,
pub budget: Duration,
}
impl std::fmt::Display for MobStopFlowRunsUnsettled {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(
f,
"mob stop did not settle its flow runs within {:?}: unsettled runs [",
self.budget
)?;
for (index, run) in self.runs.iter().enumerate() {
if index > 0 {
write!(f, ", ")?;
}
write!(f, "{run}")?;
}
write!(f, "]")?;
match &self.last_refusal {
Some(refusal) => write!(f, "; last stop refusal: {refusal}"),
None => write!(f, "; no stop refusal recorded"),
}
}
}
fn flow_terminal_run(kind: &meerkat_mob::MobEventKind) -> Option<&RunId> {
match kind {
meerkat_mob::MobEventKind::FlowCompleted { run_id, .. }
| meerkat_mob::MobEventKind::FlowFailed { run_id, .. }
| meerkat_mob::MobEventKind::FlowCanceled { run_id, .. } => Some(run_id),
_ => None,
}
}
pub(crate) trait FlowTerminals {
async fn next_terminal(&mut self) -> Option<RunId>;
fn try_next_terminal(&mut self) -> Option<RunId>;
}
pub(crate) trait MachineCommits {
async fn changed(&mut self) -> Result<(), ()>;
}
pub(crate) trait MobStopTarget {
type Terminals: FlowTerminals;
type Commits: MachineCommits;
async fn subscribe_flow_terminals(&self) -> Result<Self::Terminals, MobError>;
async fn active_flow_runs(&self) -> Result<Vec<RunId>, MobError>;
async fn cancel_flow(&self, run_id: RunId) -> Result<(), MobError>;
fn machine_commits(&self) -> Self::Commits;
async fn stop(&self) -> Result<(), MobError>;
}
pub(crate) async fn settle_flow_runs_then_stop<T: MobStopTarget>(
target: &T,
budget: Duration,
) -> Result<(), MobRuntimeError> {
let mut unsettled = BTreeSet::new();
let mut last_refusal = None;
let settle_and_stop = async {
let mut terminals = target.subscribe_flow_terminals().await?;
let mut commits = target.machine_commits();
let refusal = match target.stop().await {
Ok(()) => return Ok(()),
Err(refusal @ MobError::InvalidTransition { .. }) => refusal,
Err(error) => return Err(error),
};
let mut settled = 0usize;
while terminals.try_next_terminal().is_some() {
settled += 1;
}
for run_id in target.active_flow_runs().await? {
unsettled.insert(run_id);
}
if unsettled.is_empty() && settled == 0 {
return Err(refusal);
}
last_refusal = Some(refusal);
for run_id in unsettled.clone() {
target.cancel_flow(run_id).await?;
}
while !unsettled.is_empty() {
let Some(run_id) = terminals.next_terminal().await else {
return Err(MobError::ActorCommandChannelClosed);
};
unsettled.remove(&run_id);
}
loop {
if commits.changed().await.is_err() {
return Err(MobError::ActorCommandChannelClosed);
}
commits = target.machine_commits();
match target.stop().await {
Ok(()) => return Ok(()),
Err(refusal @ MobError::InvalidTransition { .. }) => {
last_refusal = Some(refusal);
}
Err(error) => return Err(error),
}
}
};
match tokio::time::timeout(budget, settle_and_stop).await {
Ok(result) => result.map_err(MobRuntimeError::from),
Err(_) => Err(MobRuntimeError::MobStopFlowRunsUnsettled(Box::new(
MobStopFlowRunsUnsettled {
runs: unsettled.into_iter().collect(),
last_refusal,
budget,
},
))),
}
}
pub(crate) struct HandleStopTarget<'a>(pub(crate) &'a meerkat_mob::MobHandle);
pub(crate) struct LedgerFlowTerminals(meerkat_mob::MobEventsSubscription);
impl FlowTerminals for LedgerFlowTerminals {
async fn next_terminal(&mut self) -> Option<RunId> {
loop {
let event = self.0.event_rx.recv().await?;
if let Some(run_id) = flow_terminal_run(&event.kind) {
return Some(run_id.clone());
}
}
}
fn try_next_terminal(&mut self) -> Option<RunId> {
while let Ok(event) = self.0.event_rx.try_recv() {
if let Some(run_id) = flow_terminal_run(&event.kind) {
return Some(run_id.clone());
}
}
None
}
}
impl MachineCommits for meerkat_mob::MobMachineStateChanges {
async fn changed(&mut self) -> Result<(), ()> {
meerkat_mob::MobMachineStateChanges::changed(self)
.await
.map_err(|_| ())
}
}
impl MobStopTarget for HandleStopTarget<'_> {
type Terminals = LedgerFlowTerminals;
type Commits = meerkat_mob::MobMachineStateChanges;
async fn subscribe_flow_terminals(&self) -> Result<Self::Terminals, MobError> {
Ok(LedgerFlowTerminals(self.0.events().subscribe().await?))
}
async fn active_flow_runs(&self) -> Result<Vec<RunId>, MobError> {
Ok(self
.0
.list_runs(None)
.await?
.into_iter()
.filter(|run| match run.status {
MobRunStatus::Pending | MobRunStatus::Running => true,
MobRunStatus::Completed | MobRunStatus::Failed | MobRunStatus::Canceled => false,
})
.map(|run| run.run_id)
.collect())
}
async fn cancel_flow(&self, run_id: RunId) -> Result<(), MobError> {
self.0.cancel_flow(run_id).await
}
fn machine_commits(&self) -> Self::Commits {
self.0.machine_state_changes()
}
async fn stop(&self) -> Result<(), MobError> {
let report = self.0.stop().await?;
for (identity, reason) in report.not_holdable() {
tracing::warn!(
agent_identity = %identity,
?reason,
"mob stop could not hold this member's run starts"
);
}
Ok(())
}
}
#[cfg(test)]
#[allow(clippy::expect_used, clippy::unwrap_used, clippy::panic)]
mod tests {
use std::collections::VecDeque;
use std::sync::Mutex;
use std::sync::atomic::{AtomicUsize, Ordering};
use meerkat_mob::MobState;
use tokio::sync::{mpsc, watch};
use super::*;
fn refusal() -> MobError {
MobError::InvalidTransition {
from: MobState::Running,
to: MobState::Stopped,
}
}
struct FakeTarget {
runs: Vec<RunId>,
cancel_ends_run: bool,
refusal_commits: AtomicUsize,
stop_answers: Mutex<VecDeque<Result<(), MobError>>>,
stop_calls: AtomicUsize,
cancels: AtomicUsize,
terminal_tx: mpsc::UnboundedSender<RunId>,
terminal_rx: Mutex<Option<mpsc::UnboundedReceiver<RunId>>>,
commits: watch::Sender<u64>,
}
impl FakeTarget {
fn new(
runs: Vec<RunId>,
cancel_ends_run: bool,
stop_answers: Vec<Result<(), MobError>>,
) -> Self {
let (terminal_tx, terminal_rx) = mpsc::unbounded_channel();
Self {
runs,
cancel_ends_run,
refusal_commits: AtomicUsize::new(usize::MAX),
stop_answers: Mutex::new(stop_answers.into()),
stop_calls: AtomicUsize::new(0),
cancels: AtomicUsize::new(0),
terminal_tx,
terminal_rx: Mutex::new(Some(terminal_rx)),
commits: watch::channel(0).0,
}
}
}
struct FakeTerminals(mpsc::UnboundedReceiver<RunId>);
impl FlowTerminals for FakeTerminals {
async fn next_terminal(&mut self) -> Option<RunId> {
self.0.recv().await
}
fn try_next_terminal(&mut self) -> Option<RunId> {
self.0.try_recv().ok()
}
}
struct FakeCommits(watch::Receiver<u64>);
impl MachineCommits for FakeCommits {
async fn changed(&mut self) -> Result<(), ()> {
self.0.changed().await.map_err(|_| ())
}
}
impl MobStopTarget for FakeTarget {
type Terminals = FakeTerminals;
type Commits = FakeCommits;
async fn subscribe_flow_terminals(&self) -> Result<Self::Terminals, MobError> {
Ok(FakeTerminals(
self.terminal_rx
.lock()
.unwrap()
.take()
.expect("one subscription"),
))
}
async fn active_flow_runs(&self) -> Result<Vec<RunId>, MobError> {
Ok(self.runs.clone())
}
async fn cancel_flow(&self, run_id: RunId) -> Result<(), MobError> {
self.cancels.fetch_add(1, Ordering::SeqCst);
if self.cancel_ends_run {
self.terminal_tx.send(run_id).unwrap();
}
Ok(())
}
fn machine_commits(&self) -> Self::Commits {
FakeCommits(self.commits.subscribe())
}
async fn stop(&self) -> Result<(), MobError> {
self.stop_calls.fetch_add(1, Ordering::SeqCst);
let answer = self
.stop_answers
.lock()
.unwrap()
.pop_front()
.unwrap_or_else(|| Err(refusal()));
if answer.is_err()
&& self
.refusal_commits
.fetch_update(Ordering::SeqCst, Ordering::SeqCst, |left| {
left.checked_sub(1)
})
.is_ok()
{
self.commits.send_modify(|commit| *commit += 1);
}
answer
}
}
#[tokio::test(start_paused = true)]
async fn a_mob_that_stops_is_stopped_with_one_call() {
let target = FakeTarget::new(vec![RunId::new()], true, vec![Ok(())]);
settle_flow_runs_then_stop(&target, MOB_STOP_FLOW_SETTLE_BUDGET)
.await
.expect("stopped");
assert_eq!(target.stop_calls.load(Ordering::SeqCst), 1);
assert_eq!(target.cancels.load(Ordering::SeqCst), 0, "nothing settled");
}
#[tokio::test(start_paused = true)]
async fn a_refusal_with_no_flow_runs_is_returned_as_is() {
let target = FakeTarget::new(Vec::new(), true, vec![Err(refusal())]);
let error = settle_flow_runs_then_stop(&target, MOB_STOP_FLOW_SETTLE_BUDGET)
.await
.expect_err("refused");
assert!(
matches!(
error,
MobRuntimeError::Mob(MobError::InvalidTransition { .. })
),
"{error:?}"
);
assert_eq!(target.stop_calls.load(Ordering::SeqCst), 1);
}
#[tokio::test(start_paused = true)]
async fn a_stop_refused_before_the_runs_finish_run_is_reissued_on_the_commit() {
let run = RunId::new();
let target = FakeTarget::new(
vec![run],
true,
vec![Err(refusal()), Err(refusal()), Ok(())],
);
settle_flow_runs_then_stop(&target, MOB_STOP_FLOW_SETTLE_BUDGET)
.await
.expect("stopped once the run left the machine");
assert_eq!(target.cancels.load(Ordering::SeqCst), 1);
assert_eq!(
target.stop_calls.load(Ordering::SeqCst),
3,
"the first Stop, the one refused before FinishRun, the one after"
);
}
#[tokio::test(start_paused = true)]
async fn a_finish_run_after_unrelated_commits_still_admits_the_stop() {
let target = FakeTarget::new(
vec![RunId::new()],
true,
vec![
Err(refusal()),
Err(refusal()),
Err(refusal()),
Err(refusal()),
Err(refusal()),
Ok(()),
],
);
settle_flow_runs_then_stop(&target, MOB_STOP_FLOW_SETTLE_BUDGET)
.await
.expect("stopped once the run left the machine");
assert_eq!(target.stop_calls.load(Ordering::SeqCst), 6);
}
#[tokio::test(start_paused = true)]
async fn an_unrelated_refusal_after_the_settled_runs_surfaces_at_the_guard() {
let target = FakeTarget::new(vec![RunId::new()], true, Vec::new());
target.refusal_commits.store(3, Ordering::SeqCst);
let error = settle_flow_runs_then_stop(&target, MOB_STOP_FLOW_SETTLE_BUDGET)
.await
.expect_err("still refused");
let MobRuntimeError::MobStopFlowRunsUnsettled(unsettled) = &error else {
panic!("the typed hang-guard error: {error:?}");
};
assert!(unsettled.runs.is_empty(), "the run itself settled");
assert!(matches!(
unsettled.last_refusal,
Some(MobError::InvalidTransition { .. })
));
assert_eq!(
target.stop_calls.load(Ordering::SeqCst),
4,
"the first Stop and one re-issue per commit, no more"
);
}
#[tokio::test(start_paused = true)]
async fn a_flow_run_that_never_ends_trips_the_hang_guard() {
let run = RunId::new();
let target = FakeTarget::new(vec![run.clone()], false, vec![Err(refusal())]);
let error = settle_flow_runs_then_stop(&target, MOB_STOP_FLOW_SETTLE_BUDGET)
.await
.expect_err("unsettled");
let MobRuntimeError::MobStopFlowRunsUnsettled(unsettled) = &error else {
panic!("the typed hang-guard error: {error:?}");
};
assert_eq!(unsettled.runs, vec![run.clone()]);
assert!(matches!(
unsettled.last_refusal,
Some(MobError::InvalidTransition { .. })
));
assert_eq!(unsettled.budget, MOB_STOP_FLOW_SETTLE_BUDGET);
let rendered = error.to_string();
assert!(
rendered.contains(&run.to_string()) && rendered.contains("last stop refusal"),
"{rendered}"
);
}
}