use std::collections::{BTreeMap, HashMap, HashSet};
use std::path::PathBuf;
use std::sync::Arc;
use std::time::{Duration, Instant};
use anyhow::{Context as _, Result};
use jiff::Timestamp;
use tokio::task::JoinSet;
use crate::agent::{self, Invocation, SeatState};
use crate::ask::{ChoiceAction, Deputy, Question, Questions, Waiter as Note, WaiterKind, Who};
use crate::config::Config;
use crate::notices::{self, Notice};
use crate::prompt;
pub const NODE: &str = "deputy";
pub const TICK: Duration = Duration::from_secs(5);
pub const MAX_STARTS: u32 = 3;
const RESTART_AFTER: Duration = Duration::from_secs(30);
const SLACK_SECS: u64 = 120;
const CLAIM_STALE: Duration = Duration::from_secs(60);
const HANDOVER_TIMEOUT: Duration = Duration::from_secs(10 * 60);
pub type Halt = Arc<dyn Fn() -> bool + Send + Sync>;
pub fn brief(
task_id: &str,
reason: &str,
choices: &[String],
actions: &BTreeMap<String, ChoiceAction>,
) -> String {
let reason = reason.trim();
let mut s = format!(
"Task {task_id} (`magi task show {task_id}`) is blocked on this question. \
The conductor asked because: {}\n\n\
When the owner answers, magi records the question and the answer on the \
task and unblocks it, and the answer becomes part of the instructions of \
the task's next attempt.",
if reason.is_empty() {
"(it recorded no reasoning)"
} else {
reason
}
);
if !choices.is_empty() {
s.push_str("\n\nWhat each option does:");
for c in choices {
match actions.get(c) {
Some(a) => s.push_str(&format!("\n- `{c}`: also {}", a.describe())),
None => s.push_str(&format!(
"\n- `{c}`: recorded on the task as the answer and the task is unblocked"
)),
}
}
}
s
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Kind {
Conduct,
Land,
}
pub fn kind_of(q: &Question) -> Option<Kind> {
match q.node.as_str() {
crate::conduct::NODE => Some(Kind::Conduct),
crate::land::APPROVAL_NODE => Some(Kind::Land),
_ => None,
}
}
pub fn deadline(q: &Question, default_timeout: u64) -> i64 {
let secs = if q.answer_timeout > 0 {
q.answer_timeout
} else {
default_timeout
};
let from = match kind_of(q) {
Some(Kind::Land) => q.asked_at.as_second(),
_ => q.last_activity(),
};
from.saturating_add(secs as i64)
}
pub fn can_start(cfg: Option<&Config>, agent: &str) -> bool {
cfg.is_some_and(|c| {
c.daemon.max_deputies > 0
&& ((!agent.is_empty() && c.agent(agent).is_ok()) || c.resolve_roles().is_ok())
})
}
pub fn agent_of(q: &Question) -> &str {
q.deputy.as_ref().map_or("", |d| d.agent.as_str())
}
pub fn exhausted_past_deadline(
q: &Question,
startable: bool,
default_timeout: u64,
now: Timestamp,
) -> bool {
q.status.open()
&& q.deputy
.as_ref()
.is_some_and(|d| d.starts >= MAX_STARTS || !startable)
&& now.as_second() > deadline(q, default_timeout)
}
pub struct Deputies {
store: Questions,
home: PathBuf,
cfg: Option<Config>,
fallback_repo: PathBuf,
max: usize,
halt: Halt,
tasks: JoinSet<String>,
inflight: HashSet<String>,
memo: HashMap<String, Instant>,
}
impl Deputies {
pub fn new(
store: Questions,
home: PathBuf,
cfg: Option<Config>,
fallback_repo: PathBuf,
max: usize,
halt: Halt,
) -> Self {
Self {
store,
home,
cfg,
fallback_repo,
max,
halt,
tasks: JoinSet::new(),
inflight: HashSet::new(),
memo: HashMap::new(),
}
}
fn default_timeout(&self) -> u64 {
self.cfg
.as_ref()
.map_or(Config::default().graph.answer_timeout, |c| {
c.graph.answer_timeout
})
}
fn reap(&mut self) {
while let Some(done) = self.tasks.try_join_next() {
if let Ok(id) = done {
self.inflight.remove(&id);
}
}
if self.tasks.is_empty() {
self.inflight.clear();
}
}
fn attach(&self, q: &Question, kind: Kind) -> Option<Question> {
let default_timeout = self.default_timeout();
let repo = self.fallback_repo.to_string_lossy().into_owned();
let state = match kind {
Kind::Land => crate::run::RunState::load(&q.run).ok(),
Kind::Conduct => None,
};
let timeout = state
.as_ref()
.map_or(default_timeout, |s| s.config.graph.answer_timeout);
self.store
.update(&q.id, |r| {
if r.deputy.is_none() {
r.deputy = Some(Deputy::new(match kind {
Kind::Conduct => brief(&r.run, &r.detail, &r.choices, &r.actions),
Kind::Land => crate::land::deputy_brief(r, state.as_ref()),
}));
}
if kind == Kind::Conduct && r.cwd.is_none() {
r.cwd = Some(repo.clone());
}
if r.answer_timeout == 0 {
r.answer_timeout = match kind {
Kind::Conduct => default_timeout,
Kind::Land => timeout,
};
}
Ok(())
})
.map(|(r, ())| r)
.map_err(|e| tracing::warn!("question {}: cannot attach a deputy: {e:#}", q.short()))
.ok()
}
pub fn tick(&mut self, now: Timestamp) {
self.reap();
for q in self.store.list() {
if (self.halt)() {
return;
}
let Some(kind) = kind_of(&q) else {
continue;
};
if !q.status.open() {
continue;
}
let needs_cwd = kind == Kind::Conduct && q.cwd.is_none();
let q = if q.deputy.is_none() || needs_cwd || q.answer_timeout == 0 {
match self.attach(&q, kind) {
Some(q) => q,
None => continue,
}
} else {
q
};
let Some(dep) = q.deputy.as_ref() else {
continue;
};
if self.inflight.contains(&q.id)
|| self.store.read_lease(&q.id).is_some_and(|l| l.fresh(now))
{
continue;
}
if dep.starts >= MAX_STARTS {
self.give_up(&q);
continue;
}
if now.as_second() > deadline(&q, self.default_timeout())
&& q.unread_from_owner().is_none()
{
continue;
}
if self.inflight.len() >= self.max || !can_start(self.cfg.as_ref(), dep.agent.as_str())
{
continue;
}
if matches!(self.memo.get(&q.id), Some(until) if Instant::now() < *until) {
continue;
}
self.memo
.insert(q.id.clone(), Instant::now() + RESTART_AFTER);
self.inflight.insert(q.id.clone());
let job = Job {
store: self.store.clone(),
home: self.home.clone(),
cfg: self.cfg.clone(),
fallback_repo: self.fallback_repo.clone(),
halt: Arc::clone(&self.halt),
};
let id = q.id.clone();
self.tasks.spawn(async move {
if let Err(e) = job.turn(&id).await {
tracing::warn!("deputy for question {}: {e:#}", crate::ask::short_id(&id));
}
id
});
}
}
pub async fn drain(&mut self) {
while let Some(done) = self.tasks.join_next().await {
if let Ok(id) = done {
self.inflight.remove(&id);
}
}
self.inflight.clear();
}
fn give_up(&self, q: &Question) {
notices::raise_in(
&self.home,
Notice::warn(
&format!("deputy:{}", q.id),
format!(
"Question {} \"{}\": its follow-up agent ended {MAX_STARTS} times \
without an answer and is not restarted. What you say is recorded \
but nothing will read it; answer with one of the choices instead.",
q.short(),
q.summary
),
),
);
}
}
struct Job {
store: Questions,
home: PathBuf,
cfg: Option<Config>,
fallback_repo: PathBuf,
halt: Halt,
}
impl Job {
async fn drive(
&self,
spec: &crate::config::AgentSpec,
seat: &mut SeatState,
inv: &Invocation<'_>,
id: &str,
) -> Option<Result<agent::AgentOutput>> {
let fut = agent::invoke(spec, seat, inv);
tokio::pin!(fut);
let mut beat = tokio::time::interval(Duration::from_secs(1));
let mut beats = 0u32;
loop {
tokio::select! {
r = &mut fut => break Some(r),
_ = beat.tick() => {
if (self.halt)() {
break None;
}
beats += 1;
if beats % 20 == 0 {
self.store.beat(id, WaiterKind::Deputy);
}
}
}
}
}
fn park_quietly(&self, id: &str) {
let _ = self.store.update(id, |r| {
r.waiter = None;
Ok(())
});
self.store.drop_lease(id);
}
fn park(&self, id: &str) {
let _ = self.store.update(id, |r| {
if let Some(d) = r.deputy.as_mut() {
d.starts = d.starts.saturating_sub(1);
}
r.waiter = None;
Ok(())
});
self.store.drop_lease(id);
}
async fn turn(&self, id: &str) -> Result<()> {
let claim = self.store.root().join(format!("{id}.deputy-claim"));
if !crate::waiter::take_claim(&claim, CLAIM_STALE) {
return Ok(());
}
let release = crate::waiter::Release(claim);
let q = self.store.get(id)?;
let now = Timestamp::now();
let Some(dep) = q.deputy.clone() else {
return Ok(());
};
if !q.status.open()
|| dep.starts >= MAX_STARTS
|| self.store.read_lease(&q.id).is_some_and(|l| l.fresh(now))
{
return Ok(());
}
let cfg = self
.cfg
.as_ref()
.context("the deputy's configuration is not available")?;
let spec = match cfg.agent(&dep.agent) {
Ok(s) if !dep.agent.is_empty() => s.clone(),
_ => {
cfg.resolve_roles()
.context("resolving the deputy's agent")?
.conductor
}
};
let key = crate::ask::deputy_seat_key(&q.id);
let (mut seat, resumed) = match dep.seat.clone() {
Some(s)
if s.agent == spec.id && agent::has_session(spec.kind, &s, cfg.graph.sessions) =>
{
(s, true)
}
_ => (SeatState::new(&key, &spec.id, crate::rng::entropy()), false),
};
let cwd = q
.cwd
.as_deref()
.map(PathBuf::from)
.filter(|p| p.is_dir())
.unwrap_or_else(|| self.fallback_repo.clone());
let thread: Vec<(&str, &str)> = q
.thread
.iter()
.map(|t| {
(
if t.who == Who::Operator {
"operator"
} else {
"agent"
},
t.body.as_str(),
)
})
.collect();
let read = &thread[..q.delivered_turns.min(thread.len())];
let unread = q.unread_from_owner();
let snapshot = q.thread.len();
let body = prompt::deputy(&prompt::DeputyPrompt {
id: &q.id,
summary: &q.summary,
detail: &q.detail,
brief: &dep.brief,
choices: &q.choices,
thread: read,
unread: unread.as_deref(),
resumed,
handover: false,
land: kind_of(&q) == Some(Kind::Land),
language: &cfg.graph.language,
});
self.store.beat(&q.id, WaiterKind::Deputy);
let starts = dep.starts + 1;
let first = seat.clone();
self.store.update(&q.id, |r| {
if let Some(d) = r.deputy.as_mut() {
d.agent = spec.id.clone();
d.seat = Some(first);
d.starts = starts;
}
r.waiter = Some(Note {
kind: WaiterKind::Deputy,
since: now,
});
Ok(())
})?;
drop(release);
tracing::info!(
"question {}: deputy seat {} {} (start {starts}/{MAX_STARTS})",
q.short(),
seat.key,
if resumed { "resuming" } else { "starting" }
);
let left = (deadline(&q, cfg.graph.answer_timeout) - now.as_second()).max(0) as u64;
let artifacts = self.store.root().join(format!("{}.deputy", q.id));
let cache_dir = cfg.cache_dir();
let allow_write = true;
let writable = [self.store.root().to_path_buf()];
macro_rules! invocation {
($prompt:expr, $stem:expr, $timeout:expr) => {
Invocation {
cwd: &cwd,
prompt: $prompt,
timeout: $timeout,
allow_write,
sessions: cfg.graph.sessions,
artifacts: &artifacts,
stem: $stem,
run: &q.run,
node: NODE,
cache_dir: cache_dir.as_deref(),
attachments: &[],
writable: &writable,
}
};
}
let mut early = None;
if !resumed && cfg.graph.sessions {
let hbody = prompt::deputy(&prompt::DeputyPrompt {
id: &q.id,
summary: &q.summary,
detail: &q.detail,
brief: &dep.brief,
choices: &q.choices,
thread: read,
unread: None,
resumed: false,
handover: true,
land: kind_of(&q) == Some(Kind::Land),
language: &cfg.graph.language,
});
let hstem = format!("handover-{starts}");
let hlimit = HANDOVER_TIMEOUT.min(Duration::from_secs(left.max(1)));
let hinv = invocation!(&hbody, &hstem, hlimit);
let Some(done) = self.drive(&spec, &mut seat, &hinv, &q.id).await else {
self.park(&q.id);
return Ok(());
};
let kept = seat.clone();
self.store.update(&q.id, |r| {
if let Some(d) = r.deputy.as_mut() {
d.seat = Some(kept);
}
Ok(())
})?;
if !matches!(&done, Ok(o) if o.usable()) {
early = Some(done);
}
}
let left = if early.is_none() && !resumed && cfg.graph.sessions {
let now = Timestamp::now();
let again = self.store.get(&q.id)?;
let left = (deadline(&again, cfg.graph.answer_timeout) - now.as_second()).max(0) as u64;
if !again.status.open() || (left == 0 && again.unread_from_owner().is_none()) {
self.park_quietly(&q.id);
return Ok(());
}
left
} else {
left
};
let out = match early {
Some(done) => done,
None => {
let stem = format!("turn-{starts}");
let timeout = Duration::from_secs(left.max(60) + SLACK_SECS);
let inv = invocation!(&body, &stem, timeout);
match self.drive(&spec, &mut seat, &inv, &q.id).await {
Some(out) => out,
None => {
self.park(&q.id);
return Ok(());
}
}
}
};
let (text, why) = match out {
Ok(o) if o.usable() => (Some(o.text.trim().to_owned()), None),
Ok(o) if o.timed_out => (None, Some("its turn timed out".to_owned())),
Ok(o) if o.quota_exhausted() => (None, Some("the agent is out of quota".to_owned())),
Ok(o) => (
None,
Some(format!("its turn failed (exit {:?})", o.exit_code)),
),
Err(e) => (None, Some(format!("the agent could not be started: {e:#}"))),
};
let kept = seat.clone();
self.store.update(&q.id, |r| {
if let Some(d) = r.deputy.as_mut() {
d.seat = Some(kept);
}
r.waiter = None;
if let Some(text) = text
&& unread.is_some()
&& !text.is_empty()
&& r.status.open()
&& r.thread.len() == snapshot
{
r.delivered_turns = r.delivered_turns.max(snapshot);
let choices = r.choices.clone();
r.reply(text, choices)?;
}
Ok(())
})?;
self.store.drop_lease(&q.id);
if let Some(why) = why {
tracing::warn!("question {}: the deputy ended: {why}", q.short());
notices::raise_in(
&self.home,
Notice::warn(
&format!("deputy-turn:{}", q.id),
format!(
"Question {} \"{}\": its follow-up agent stopped ({why}).",
q.short(),
q.summary
),
),
);
}
Ok(())
}
}
pub async fn run(mut deputies: Deputies, stop: crate::daemon::Stop) {
while !stop.stopped() {
deputies.tick(Timestamp::now());
tokio::time::sleep(TICK).await;
}
}