use crate::core::{RuntimeError, Timestamp};
pub const SUBJECTS_LISTED: usize = 10;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct Remedy {
pub cli: &'static str,
pub http: &'static str,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Condition {
pub kind: &'static str,
pub found: usize,
pub at_least: bool,
pub subjects: Vec<String>,
pub unlisted: usize,
pub remedy: Remedy,
}
#[derive(Debug, Clone, PartialEq, Eq, Default)]
pub struct Attention {
pub conditions: Vec<Condition>,
pub not_checked: Vec<&'static str>,
}
impl Attention {
#[must_use]
pub fn any(&self) -> bool {
!self.conditions.is_empty()
}
fn note(
&mut self,
kind: &'static str,
found: usize,
subjects: impl IntoIterator<Item = String>,
page: usize,
remedy: Remedy,
) {
if found > 0 {
let subjects: Vec<String> = subjects.into_iter().take(SUBJECTS_LISTED).collect();
self.conditions.push(Condition {
kind,
found,
at_least: found >= page,
unlisted: found.saturating_sub(subjects.len()),
subjects,
remedy,
});
}
}
}
impl super::Runtime {
pub async fn attention(&self, at: Timestamp, page: usize) -> Result<Attention, RuntimeError> {
let mut out = Attention::default();
let mut stopped = std::collections::BTreeSet::new();
for (outcome, kind, remedy) in [
(
"quarantined",
"run.quarantined",
Remedy {
cli: "establish an undecided effect with `reconcile`, or bring back the \
declaration or policy bundle the run was admitted under — then \
`quarantine` to reopen or abandon",
http: "establish an undecided effect with \
`POST /runs/{run}/reconcile`, or bring back the declaration or \
policy bundle the run was admitted under — then \
`POST /runs/{run}/reopen` or `POST /runs/{run}/abandon`",
},
),
(
"exhausted",
"run.exhausted",
Remedy {
cli: "raise the ceiling — or, stopped by a tool's rate ceiling, wait for \
its window to pass — and `replay`, or `cancel` the run",
http: "raise the ceiling — or, stopped by a tool's rate ceiling, wait \
for its window to pass — and resume the run with `agentplane \
replay` (this API does not resume runs), or \
`POST /runs/{run}/cancel`",
},
),
(
"withheld",
"run.withheld",
Remedy {
cli: "lift the withdrawal with `halt --lift`, then `replay` — or \
`cancel` the run",
http: "lift the withdrawal with `POST /halts/lift`, then resume the \
run with `agentplane replay` — or `POST /runs/{run}/cancel`",
},
),
] {
let found = self
.store()
.runs_by_outcome(outcome, page)
.await
.map_err(RuntimeError::from_store)?;
stopped.extend(found.iter().copied());
out.note(
kind,
found.len(),
found.iter().map(ToString::to_string),
page,
remedy,
);
}
stopped.extend(self.failed_on_landed_work(&mut out, page).await?);
self.held_by_stopped_runs(&mut out, &stopped, page).await?;
let abandoned = self
.store()
.abandoned_runs(page)
.await
.map_err(RuntimeError::from_store)?;
out.note(
"run.abandoned",
abandoned.len(),
abandoned.iter().map(ToString::to_string),
page,
Remedy {
cli: "no verb: the recovery sweep resumes these; one that keeps \
reappearing is a run nothing can drive",
http: "no verb: the recovery sweep resumes these; one that keeps \
reappearing is a run nothing can drive",
},
);
let overdue: Vec<String> = self
.store()
.waiting_runs(page)
.await
.map_err(RuntimeError::from_store)?
.into_iter()
.filter(|w| w.reason.until() <= at)
.map(|w| w.run.to_string())
.collect();
out.note(
"run.wait_expired",
overdue.len(),
overdue,
page,
Remedy {
cli: "`replay` the run: it reaches the announced wait and re-arms it",
http: "resume the run with `agentplane replay` (this API does not \
resume runs): it reaches the announced wait and re-arms it",
},
);
self.backlogs(&mut out, page, at).await?;
Ok(out)
}
async fn failed_on_landed_work(
&self,
out: &mut Attention,
page: usize,
) -> Result<Vec<crate::core::RunId>, RuntimeError> {
let failed = self
.store()
.runs_by_outcome("failed", page)
.await
.map_err(RuntimeError::from_store)?;
let mut landed = Vec::new();
for run in &failed {
if self.holds_landed_work(*run).await? {
landed.push(run.to_string());
}
}
let ceiling = if failed.len() >= page {
landed.len()
} else {
page
};
out.note(
"run.failed_with_landed_work",
landed.len(),
landed,
ceiling,
Remedy {
cli: "`replay` the run to finish its work, or `cancel` it to unwind what \
landed",
http: "resume the run with `agentplane replay` (this API does not resume \
runs) to finish its work, or `POST /runs/{run}/cancel` to unwind \
what landed",
},
);
Ok(failed)
}
async fn held_by_stopped_runs(
&self,
out: &mut Attention,
stopped: &std::collections::BTreeSet<crate::core::RunId>,
page: usize,
) -> Result<(), RuntimeError> {
let Some(quotas) = self.quota_store_if_wired() else {
out.not_checked.push(
"quota reservations — this plane holds no quota store, so whether a \
stopped run holds part of the tenant's period was not established",
);
return Ok(());
};
let held = quotas
.reservations(page)
.await
.map_err(RuntimeError::Store)?;
let named: Vec<String> = held
.iter()
.filter(|h| stopped.contains(&h.run))
.map(|h| h.run.to_string())
.collect();
let ceiling = if held.len() >= page {
named.len()
} else {
page
};
out.note(
"quota.held_by_stopped_run",
named.len(),
named,
ceiling,
Remedy {
cli: "each holds spend its tenant's period cannot admit against until it \
concludes: answer the run's own condition — `replay` it to finish, \
`cancel` it, or answer its `quarantine` — and what it did not spend \
is released",
http: "each holds spend its tenant's period cannot admit against until it \
concludes: resume it with `agentplane replay`, \
`POST /runs/{run}/cancel` it, or `POST /runs/{run}/abandon` a \
quarantine — and what it did not spend is released",
},
);
Ok(())
}
#[allow(clippy::too_many_lines)]
async fn backlogs(
&self,
out: &mut Attention,
page: usize,
now: Timestamp,
) -> Result<(), RuntimeError> {
match self.cases() {
Some(cases) => {
let breached = cases
.breached(page)
.await
.map_err(RuntimeError::from_store)?;
out.note(
"obligation.breached",
breached.len(),
breached.iter().map(|d| format!("{}/{}", d.case, d.name)),
page,
Remedy {
cli: "`acknowledge` the breach, saying what was done about it \
(each subject is <case>/<obligation>)",
http: "`POST /obligations/acknowledge` the breach, saying what \
was done about it (each subject is <case>/<obligation>)",
},
);
}
None => out.not_checked.push(
"obligations — this plane holds no case store, so whether any went \
unaccounted for was not established",
),
}
match self.cases() {
Some(cases) => {
let failed = cases
.last_drill()
.await
.map_err(RuntimeError::from_store)?
.is_some_and(|d| !d.sound);
out.note(
"drill.failed",
usize::from(failed),
std::iter::empty(),
usize::MAX,
Remedy {
cli: "the last recovery rehearsal found unrecoverable references — \
re-run `drill` and resolve what it names",
http: "the last recovery rehearsal found unrecoverable references \
(`GET /drill` reads the verdict) — re-run it with \
`agentplane drill` and resolve what it names",
},
);
}
None => out.not_checked.push(
"recovery rehearsal — this plane holds no case store, so there is \
nothing to drill and no verdict to read",
),
}
match self.tasks() {
Some(tasks) => {
let overdue = tasks
.overdue(now, page)
.await
.map_err(RuntimeError::from_store)?;
out.note(
"task.overdue",
overdue.len(),
overdue.iter().map(|t| t.id.to_string()),
page,
Remedy {
cli: "`decide` them. Escalation is not an operator act — the \
sweeper widens the audience itself when the window closes",
http: "`POST /tasks/{task}/decide` them. Escalation is not an \
operator act — the sweeper widens the audience itself when \
the window closes",
},
);
}
None => out.not_checked.push(
"worklist — this plane holds no task store, so whether any approval is \
overdue was not established",
),
}
match self.events() {
Some(events) => {
let dead = events
.dead_letters(page)
.await
.map_err(RuntimeError::from_store)?;
out.note(
"event.dead_lettered",
dead.len(),
dead.iter()
.map(|d| format!("{} {}", d.event.source, d.event.id)),
page,
Remedy {
cli: "no verb: the correlation key belongs to the emitter, so \
this is a diagnosis to carry to them (each subject is \
<source> <id>)",
http: "no verb: the correlation key belongs to the emitter, so \
this is a diagnosis to carry to them — `GET /dead-letters` \
lists what arrived",
},
);
}
None => out.not_checked.push(
"dead letters — this plane holds no event store, so whether an inbound \
message went unclaimed was not established",
),
}
#[cfg(feature = "push")]
match self.push() {
Some(push) => {
let parked = push.parked(page).await.map_err(RuntimeError::from_store)?;
out.note(
"push.parked",
parked.len(),
parked
.iter()
.map(|p| format!("{}/{}", p.config.task, p.config.id)),
page,
Remedy {
cli: "fix the endpoint, then `rearm` the registration (each \
subject is <run>/<id>)",
http: "fix the endpoint, then `POST /push/rearm` the registration \
(each subject is <run>/<id>)",
},
);
}
None => out.not_checked.push(
"push registrations — this plane holds no push store, so whether a \
worker gave up on one was not established",
),
}
#[cfg(not(feature = "push"))]
out.not_checked
.push("push registrations — this build has no outbound delivery");
Ok(())
}
}