magi/deputy.rs
1//! Deputies: the seat that waits on a conductor's question, so the conductor
2//! never has to - and on land's merge approval, whose free-text reply used to
3//! reach nobody either.
4//!
5//! [`crate::conduct`] files a question for the operator and moves on; it is
6//! non-blocking by construction, because one task's question must not park the
7//! whole polling loop. The price used to be that nothing at all was waiting on
8//! it: the owner's free-text reply landed in a thread no agent read, while the
9//! phone said "waiting for the agent". Question f1dc on task 684d sat open
10//! forever on a "setup done".
11//!
12//! A deputy is the answer. For each open conductor question, `magi serve`
13//! runs one short-lived agent seat that is handed what the conductor knew
14//! ([`brief`]: the task, why it asked, what each option leads to), blocks on
15//! the question with `magi ask --wait`, answers back in the thread with
16//! `magi ask --thread`, and - when the owner's own words clearly pick an
17//! option - records that with `magi ask --settle`. Applying the outcome is not
18//! the deputy's job: an answered conductor question goes through the daemon's
19//! existing `resolve_blockers` / action path, exactly as one tapped on the
20//! phone does.
21//!
22//! # Invariants
23//!
24//! - **One seat per question, keyed by seat name.** `deputy-<question id>`,
25//! never the conductor's shared seat (`conduct/seat.json`, overwritten by
26//! every cycle) and never an agent id.
27//! - **Persisted on the question** ([`crate::ask::Deputy`]), always through
28//! [`Questions::update`]: the owner's say and answer write the same file. A
29//! restarted daemon resumes the seat only when [`agent::has_session`] says
30//! it can; otherwise it starts a fresh seat and re-sends the whole context,
31//! never assuming memory.
32//! - **Bounded.** At most `daemon.max_deputies` at once; at most
33//! [`MAX_STARTS`] starts per question, and a restart does not reset that. A
34//! deputy never outlives the question's `answer_timeout`: the waiter retires
35//! the question at its deadline ([`crate::waiter`]) and the task goes to a
36//! machine hold (`daemon::resolve_blockers`).
37//! - **Never a second agent on one question.** A fresh lease or the claim file
38//! means somebody is already on it.
39//! - **A merge approval is served too, but stays land's.** Its deputy has no
40//! `cwd` (so the waiter never touches it), its deadline is `asked_at +
41//! answer_timeout` and never moves on a reply ([`deadline`]), and the only
42//! thing that retires it is `daemon::land_resume_state`. A say alone never
43//! merges: `--settle` accepts `merge` only for the owner's own word `merge`.
44//! - **A release-watch question is served too, and stays the watcher's.** The
45//! escalation, local-mode approval and failed-release questions
46//! ([`crate::release_watch`], node `release-bump`, seat `release-watch`) have
47//! no asker, so a say reached nobody. Its deputy has no `cwd` either, a
48//! fixed `asked_at + answer_timeout` deadline ([`deadline`]), and applies
49//! nothing: `--settle` records a choice and the watcher applies it on its next
50//! lap. A local approval's `merge` / `hold` go through the merge-approval
51//! rules ([`merge_gated`]). A deputy never closes or merges the pull request.
52//! - **Every open question the owner can say something to has a listener.**
53//! [`Kind::Triage`] serves the triage questions; [`Kind::Generic`] is the
54//! fallback for any other question with choices and no `cwd` (the divergence
55//! question today, a node nobody has written yet tomorrow). It is decided by
56//! the missing `cwd`, never by node name. Only a conductor question is ever
57//! given a `cwd` ([`Deputies::attach`]); the other kinds would otherwise be
58//! resumed and expired by the waiter beside the deputy. A test enumerates
59//! every node filed without an asker, so a new kind of question cannot ship
60//! without one.
61//! - **No new authority.** A deputy does not edit, merge or touch the queue,
62//! and `--settle` accepts only an offered label backed by a verbatim quote
63//! of the owner, never on a task the operator holds.
64
65use std::collections::{BTreeMap, HashMap, HashSet};
66use std::path::PathBuf;
67use std::sync::Arc;
68use std::time::{Duration, Instant};
69
70use anyhow::{Context as _, Result};
71use jiff::Timestamp;
72use tokio::task::JoinSet;
73
74use crate::agent::{self, Invocation, SeatState};
75use crate::ask::{ChoiceAction, Deputy, Question, Questions, Waiter as Note, WaiterKind, Who};
76use crate::config::Config;
77use crate::notices::{self, Notice};
78use crate::prompt;
79
80/// Node name a deputy runs under (`MAGI_NODE`); `magi ask --settle` accepts
81/// it and nothing else.
82pub const NODE: &str = "deputy";
83
84/// How often the runner looks at the store.
85pub const TICK: Duration = Duration::from_secs(5);
86
87/// Most turns ever started for one question. A deputy that keeps ending
88/// without an answer is a problem for the owner, not something to retry until
89/// the quota is gone.
90pub const MAX_STARTS: u32 = 3;
91
92/// Least gap between two starts for one question.
93const RESTART_AFTER: Duration = Duration::from_secs(30);
94
95/// Added to the question's remaining time for the invocation's wall clock, so
96/// the deputy's own `magi ask --wait` reaches the deadline first.
97const SLACK_SECS: u64 = 120;
98
99/// A claim file older than this was left by a process that died.
100const CLAIM_STALE: Duration = Duration::from_secs(60);
101
102/// How long the context-handover turn may take.
103const HANDOVER_TIMEOUT: Duration = Duration::from_secs(10 * 60);
104
105/// Stops the deputy when true: the daemon is parking.
106pub type Halt = Arc<dyn Fn() -> bool + Send + Sync>;
107
108/// What the conductor knew about one question, for its deputy.
109///
110/// `reason` is the conductor's own one-line reasoning (may be empty). Each
111/// choice is listed with what picking it does: an attached action, or - for
112/// the usual conductor question - that the answer is recorded on the task and
113/// the task is unblocked with it in its instructions.
114pub fn brief(
115 task_id: &str,
116 reason: &str,
117 choices: &[String],
118 actions: &BTreeMap<String, ChoiceAction>,
119) -> String {
120 let reason = reason.trim();
121 let mut s = format!(
122 "Task {task_id} (`magi task show {task_id}`) is blocked on this question. \
123 The conductor asked because: {}\n\n\
124 When the owner answers, magi records the question and the answer on the \
125 task and unblocks it, and the answer becomes part of the instructions of \
126 the task's next attempt.",
127 if reason.is_empty() {
128 "(it recorded no reasoning)"
129 } else {
130 reason
131 }
132 );
133 if !choices.is_empty() {
134 s.push_str("\n\nWhat each option does:");
135 for c in choices {
136 match actions.get(c) {
137 Some(a) => s.push_str(&format!("\n- `{c}`: also {}", a.describe())),
138 None => s.push_str(&format!(
139 "\n- `{c}`: recorded on the task as the answer and the task is unblocked"
140 )),
141 }
142 }
143 }
144 s
145}
146
147/// What a deputy serves: the question kinds it is attached to.
148#[derive(Debug, Clone, Copy, PartialEq, Eq)]
149pub enum Kind {
150 /// A conductor's question ([`crate::conduct::NODE`]).
151 Conduct,
152 /// The merge approval ([`crate::land::APPROVAL_NODE`]).
153 Land,
154 /// A question the release watcher filed ([`crate::release_watch`]).
155 Release,
156 /// A triage question about a held task or a stuck dependency root
157 /// ([`crate::triage::NODE`], [`crate::triage::DEPS_NODE`]).
158 Triage,
159 /// Any other question nobody asked from inside a run (no `cwd`) that offers
160 /// choices: the fallback, so a node nobody has met yet still has a listener.
161 Generic,
162}
163
164/// The kind of question a deputy serves for `q`, `None` for every other.
165///
166/// The fallback is decided by the absence of a `cwd`, never by the node name: a
167/// question filed with `magi ask` inside a run records its `cwd` and has an
168/// asker (and the waiter), and the node a seat asks from is the seat's own name
169/// (a reviewer's `review` is also the divergence question's node). A question
170/// with no choices cannot be settled and is a notice, not something to answer.
171pub fn kind_of(q: &Question) -> Option<Kind> {
172 match q.node.as_str() {
173 crate::conduct::NODE => Some(Kind::Conduct),
174 crate::land::APPROVAL_NODE => Some(Kind::Land),
175 // The seat matters: `bump` files choice-less notices on the same node,
176 // which nobody can answer and which must not cost a deputy.
177 crate::bump::NOTICE_NODE if q.seat == "release-watch" => Some(Kind::Release),
178 crate::bump::NOTICE_NODE => None,
179 crate::triage::NODE | crate::triage::DEPS_NODE => Some(Kind::Triage),
180 _ if q.cwd.is_none() && !q.choices.is_empty() => Some(Kind::Generic),
181 _ => None,
182 }
183}
184
185/// What a deputy is told about a question of no known kind: only what the
186/// question itself stored.
187pub fn generic_brief(q: &Question, actions: &BTreeMap<String, ChoiceAction>) -> String {
188 let mut s = format!(
189 "This question was filed by magi (node `{}`, seat `{}`) with no agent \
190 waiting on it. Magi records the owner's answer and the component that \
191 asked applies it; you apply nothing. You know only what the question \
192 itself says. A choice whose effect you cannot read from the question is \
193 not yours to guess: do not settle it, ask the owner what they mean with \
194 `--thread` instead.",
195 q.node, q.seat
196 );
197 if !q.choices.is_empty() {
198 s.push_str("\n\nWhat each option does:");
199 for c in &q.choices {
200 match actions.get(c) {
201 Some(a) => s.push_str(&format!("\n- `{c}`: also {}", a.describe())),
202 None => s.push_str(&format!(
203 "\n- `{c}`: recorded as the answer (nothing more is known)"
204 )),
205 }
206 }
207 }
208 s
209}
210
211/// Is picking `label` on `q` something that cannot be taken back, so that
212/// `--settle` must hold the owner's words to the same mechanical standard as a
213/// merge (`land::unhedged`)? Discarding a task deletes it; a divergence answer
214/// drops commits.
215pub fn destructive(q: &Question, label: &str) -> bool {
216 let at = q.choices.iter().position(|c| c == label);
217 match q.node.as_str() {
218 crate::triage::NODE => at == Some(2),
219 crate::triage::DEPS_NODE => at == Some(1),
220 n => n == crate::reconcile::NODE && q.seat == crate::reconcile::SEAT,
221 }
222}
223
224/// Does `--settle` hold `q` to the merge-approval rules (the verbatim quote of
225/// the latest message for `merge`, the whole message for `hold`)? A merge
226/// approval, and any other served question that offers `merge` (the release
227/// watcher's local-mode approval merges just as irreversibly).
228pub fn merge_gated(q: &Question) -> bool {
229 q.node == crate::land::APPROVAL_NODE
230 || (kind_of(q).is_some_and(|k| k != Kind::Conduct)
231 && q.choices.iter().any(|c| c == crate::land::APPROVE))
232}
233
234/// Does `q`'s clock run from `asked_at` and never move on a reply? True for the
235/// questions that something other than the waiter retires: a merge approval
236/// (land) and a release-watch question (the watcher, by silence being a hold).
237pub fn fixed_clock(q: &Question) -> bool {
238 matches!(kind_of(q), Some(Kind::Land | Kind::Release))
239}
240
241/// Second after which nobody is to be started or kept on `q`.
242///
243/// A conductor question runs from its last activity (`magi ask --thread`
244/// re-arms it, and the waiter retires it). A merge approval never moves: it
245/// runs from `asked_at`, exactly where `daemon::land_resume_state` abandons it,
246/// so a conversation cannot stretch the hold and land stays the only place that
247/// retires one.
248pub fn deadline(q: &Question, default_timeout: u64) -> i64 {
249 let secs = if q.answer_timeout > 0 {
250 q.answer_timeout
251 } else {
252 default_timeout
253 };
254 let from = if fixed_clock(q) {
255 q.asked_at.as_second()
256 } else {
257 q.last_activity()
258 };
259 from.saturating_add(secs as i64)
260}
261
262/// Can `magi serve` start a deputy at all under `cfg`? Not when deputies are
263/// switched off (`daemon.max_deputies = 0`) or the config could not be read.
264///
265/// Also false when the agent the deputy would run as cannot be resolved
266/// (`agent` is the deputy's recorded agent, empty when it has none): `turn`
267/// fails on that before it counts a start, so it would never reach
268/// [`MAX_STARTS`].
269pub fn can_start(cfg: Option<&Config>, agent: &str) -> bool {
270 cfg.is_some_and(|c| {
271 c.daemon.max_deputies > 0
272 && ((!agent.is_empty() && c.agent(agent).is_ok()) || c.resolve_roles().is_ok())
273 })
274}
275
276/// The agent a question's deputy records, empty when it has none yet.
277pub fn agent_of(q: &Question) -> &str {
278 q.deputy.as_ref().map_or("", |d| d.agent.as_str())
279}
280
281/// Has this question's deputy run out of starts, or can none ever start, with
282/// the deadline gone?
283///
284/// Then nothing will ever read an unread say, and the waiter must retire the
285/// question anyway instead of deferring to a deputy that no longer starts.
286/// `startable` is [`can_start`]; a fresh lease is the caller's to check.
287pub fn exhausted_past_deadline(
288 q: &Question,
289 startable: bool,
290 default_timeout: u64,
291 now: Timestamp,
292) -> bool {
293 q.status.open()
294 && q.deputy
295 .as_ref()
296 .is_some_and(|d| d.starts >= MAX_STARTS || !startable)
297 && now.as_second() > deadline(q, default_timeout)
298}
299
300/// The deputy runner: its own task inside `magi serve`, beside the waiter.
301pub struct Deputies {
302 store: Questions,
303 home: PathBuf,
304 cfg: Option<Config>,
305 /// Working directory for a question that recorded none.
306 fallback_repo: PathBuf,
307 max: usize,
308 halt: Halt,
309 tasks: JoinSet<String>,
310 inflight: HashSet<String>,
311 /// Earliest next start per question.
312 memo: HashMap<String, Instant>,
313}
314
315impl Deputies {
316 /// A runner over `store`, with at most `max` deputies at once.
317 pub fn new(
318 store: Questions,
319 home: PathBuf,
320 cfg: Option<Config>,
321 fallback_repo: PathBuf,
322 max: usize,
323 halt: Halt,
324 ) -> Self {
325 Self {
326 store,
327 home,
328 cfg,
329 fallback_repo,
330 max,
331 halt,
332 tasks: JoinSet::new(),
333 inflight: HashSet::new(),
334 memo: HashMap::new(),
335 }
336 }
337
338 fn default_timeout(&self) -> u64 {
339 self.cfg
340 .as_ref()
341 .map_or(Config::default().graph.answer_timeout, |c| {
342 c.graph.answer_timeout
343 })
344 }
345
346 fn reap(&mut self) {
347 while let Some(done) = self.tasks.try_join_next() {
348 if let Ok(id) = done {
349 self.inflight.remove(&id);
350 }
351 }
352 if self.tasks.is_empty() {
353 self.inflight.clear();
354 }
355 }
356
357 /// Give a question the record a deputy needs: the deputy itself (with a
358 /// brief rebuilt from what the question stored, when it was filed before
359 /// deputies existed - a lost reason is not invented) and the deadline.
360 ///
361 /// A conductor question also gets the working directory the waiter and
362 /// `magi ask --wait` use. Every other kind never does: `cwd` is what makes a
363 /// question the waiter's (it would resume and expire it beside the deputy),
364 /// and a release-watch question, a merge approval, a triage question and a
365 /// fallback one have no asker to resume.
366 fn attach(&self, q: &Question, kind: Kind) -> Option<Question> {
367 let default_timeout = self.default_timeout();
368 let repo = self.fallback_repo.to_string_lossy().into_owned();
369 let state = match kind {
370 Kind::Land => crate::run::RunState::load(&q.run).ok(),
371 Kind::Conduct | Kind::Release | Kind::Triage | Kind::Generic => None,
372 };
373 // The deadline land itself enforces for this run.
374 let timeout = state
375 .as_ref()
376 .map_or(default_timeout, |s| s.config.graph.answer_timeout);
377 self.store
378 .update(&q.id, |r| {
379 if r.deputy.is_none() {
380 r.deputy = Some(Deputy::new(match kind {
381 Kind::Conduct => brief(&r.run, &r.detail, &r.choices, &r.actions),
382 Kind::Land => crate::land::deputy_brief(r, state.as_ref()),
383 Kind::Release => crate::release_watch::deputy_brief(r, &self.home),
384 Kind::Triage => crate::triage::deputy_brief(
385 r,
386 &crate::queue::Queue::at(self.home.join("queue")),
387 ),
388 Kind::Generic => generic_brief(r, &r.actions),
389 }));
390 }
391 if kind == Kind::Conduct && r.cwd.is_none() {
392 r.cwd = Some(repo.clone());
393 }
394 if r.answer_timeout == 0 {
395 r.answer_timeout = match kind {
396 Kind::Conduct => default_timeout,
397 Kind::Land => timeout,
398 Kind::Release | Kind::Triage | Kind::Generic => default_timeout,
399 };
400 }
401 Ok(())
402 })
403 .map(|(r, ())| r)
404 .map_err(|e| tracing::warn!("question {}: cannot attach a deputy: {e:#}", q.short()))
405 .ok()
406 }
407
408 /// Look at every open conductor question once and start a deputy on each
409 /// that has none, within the limits. Returns without waiting for them.
410 pub fn tick(&mut self, now: Timestamp) {
411 self.reap();
412 for q in self.store.list() {
413 if (self.halt)() {
414 return;
415 }
416 let Some(kind) = kind_of(&q) else {
417 continue;
418 };
419 if !q.status.open() {
420 continue;
421 }
422 let needs_cwd = kind == Kind::Conduct && q.cwd.is_none();
423 let q = if q.deputy.is_none() || needs_cwd || q.answer_timeout == 0 {
424 match self.attach(&q, kind) {
425 Some(q) => q,
426 None => continue,
427 }
428 } else {
429 q
430 };
431 let Some(dep) = q.deputy.as_ref() else {
432 continue;
433 };
434 if self.inflight.contains(&q.id)
435 || self.store.read_lease(&q.id).is_some_and(|l| l.fresh(now))
436 {
437 continue;
438 }
439 if dep.starts >= MAX_STARTS {
440 self.give_up(&q);
441 continue;
442 }
443 // Past the deadline the waiter retires the question; a say that
444 // arrived in time is still read first.
445 if now.as_second() > deadline(&q, self.default_timeout())
446 && q.unread_from_owner().is_none()
447 {
448 continue;
449 }
450 if self.inflight.len() >= self.max || !can_start(self.cfg.as_ref(), dep.agent.as_str())
451 {
452 continue;
453 }
454 if matches!(self.memo.get(&q.id), Some(until) if Instant::now() < *until) {
455 continue;
456 }
457 self.memo
458 .insert(q.id.clone(), Instant::now() + RESTART_AFTER);
459 self.inflight.insert(q.id.clone());
460 let job = Job {
461 store: self.store.clone(),
462 home: self.home.clone(),
463 cfg: self.cfg.clone(),
464 fallback_repo: self.fallback_repo.clone(),
465 halt: Arc::clone(&self.halt),
466 };
467 let id = q.id.clone();
468 self.tasks.spawn(async move {
469 if let Err(e) = job.turn(&id).await {
470 tracing::warn!("deputy for question {}: {e:#}", crate::ask::short_id(&id));
471 }
472 id
473 });
474 }
475 }
476
477 /// Wait for every deputy turn in flight. The daemon never calls this; the
478 /// tests do, to look at the record once a turn is over.
479 pub async fn drain(&mut self) {
480 while let Some(done) = self.tasks.join_next().await {
481 if let Ok(id) = done {
482 self.inflight.remove(&id);
483 }
484 }
485 self.inflight.clear();
486 }
487
488 /// Say once that nobody is listening any more.
489 fn give_up(&self, q: &Question) {
490 notices::raise_in(
491 &self.home,
492 Notice::warn(
493 &format!("deputy:{}", q.id),
494 format!(
495 "Question {} \"{}\": its follow-up agent ended {MAX_STARTS} times \
496 without an answer and is not restarted. What you say is recorded \
497 but nothing will read it; answer with one of the choices instead.",
498 q.short(),
499 q.summary
500 ),
501 ),
502 );
503 }
504}
505
506/// Everything one deputy turn needs, owned so it can run detached.
507struct Job {
508 store: Questions,
509 home: PathBuf,
510 cfg: Option<Config>,
511 fallback_repo: PathBuf,
512 halt: Halt,
513}
514
515impl Job {
516 /// Run one invocation, beating the lease; `None` when the daemon is
517 /// parking and the turn was dropped.
518 async fn drive(
519 &self,
520 spec: &crate::config::AgentSpec,
521 seat: &mut SeatState,
522 inv: &Invocation<'_>,
523 id: &str,
524 ) -> Option<Result<agent::AgentOutput>> {
525 let fut = agent::invoke(spec, seat, inv);
526 tokio::pin!(fut);
527 let mut beat = tokio::time::interval(Duration::from_secs(1));
528 let mut beats = 0u32;
529 loop {
530 tokio::select! {
531 r = &mut fut => break Some(r),
532 _ = beat.tick() => {
533 if (self.halt)() {
534 break None;
535 }
536 beats += 1;
537 if beats % 20 == 0 {
538 self.store.beat(id, WaiterKind::Deputy);
539 }
540 }
541 }
542 }
543 }
544
545 /// The question settled or expired during the handover: release it without
546 /// starting the long turn. The start stays counted - the handover ran.
547 fn park_quietly(&self, id: &str) {
548 let _ = self.store.update(id, |r| {
549 r.waiter = None;
550 Ok(())
551 });
552 self.store.drop_lease(id);
553 }
554
555 /// Parking: the turn is dropped and the start refunded. The seat stays as
556 /// last persisted - after the handover turn, resumable.
557 fn park(&self, id: &str) {
558 let _ = self.store.update(id, |r| {
559 if let Some(d) = r.deputy.as_mut() {
560 d.starts = d.starts.saturating_sub(1);
561 }
562 r.waiter = None;
563 Ok(())
564 });
565 self.store.drop_lease(id);
566 }
567
568 async fn turn(&self, id: &str) -> Result<()> {
569 let claim = self.store.root().join(format!("{id}.deputy-claim"));
570 // The claim only covers the decision to start and the write that records
571 // it, so a daemon that dies holding it blocks a restart for a minute,
572 // not for a turn's length. The lease guards the turn itself.
573 if !crate::waiter::take_claim(&claim, CLAIM_STALE) {
574 return Ok(());
575 }
576 let release = crate::waiter::Release(claim);
577
578 // Decided again under the claim, on the record as it is now.
579 let q = self.store.get(id)?;
580 let now = Timestamp::now();
581 let Some(dep) = q.deputy.clone() else {
582 return Ok(());
583 };
584 if !q.status.open()
585 || dep.starts >= MAX_STARTS
586 || self.store.read_lease(&q.id).is_some_and(|l| l.fresh(now))
587 {
588 return Ok(());
589 }
590 let cfg = self
591 .cfg
592 .as_ref()
593 .context("the deputy's configuration is not available")?;
594 let spec = match cfg.agent(&dep.agent) {
595 Ok(s) if !dep.agent.is_empty() => s.clone(),
596 _ => {
597 cfg.resolve_roles()
598 .context("resolving the deputy's agent")?
599 .conductor
600 }
601 };
602 // The same CLI conversation when it can be resumed; otherwise a fresh
603 // seat - and either way the prompt carries the whole context.
604 let key = crate::ask::deputy_seat_key(&q.id);
605 let (mut seat, resumed) = match dep.seat.clone() {
606 Some(s)
607 if s.agent == spec.id && agent::has_session(spec.kind, &s, cfg.graph.sessions) =>
608 {
609 (s, true)
610 }
611 _ => (SeatState::new(&key, &spec.id, crate::rng::entropy()), false),
612 };
613 // A release-watch question has no `cwd`; its watch record names the
614 // checkout the pull request belongs to.
615 let recorded = q.cwd.clone().or_else(|| match kind_of(&q) {
616 Some(Kind::Release) => {
617 crate::release_watch::state_for_question(&self.home, &q.id).map(|st| st.repo)
618 }
619 // These record a task id in `run` (or a run id for a divergence's
620 // sibling kinds): the task's repository, when it can be found.
621 Some(Kind::Triage | Kind::Generic) => crate::queue::Queue::at(self.home.join("queue"))
622 .get(&q.run)
623 .ok()
624 .map(|t| t.repo.to_string_lossy().into_owned())
625 .filter(|r| !r.is_empty()),
626 _ => None,
627 });
628 let cwd = recorded
629 .as_deref()
630 .map(PathBuf::from)
631 .filter(|p| p.is_dir())
632 .unwrap_or_else(|| self.fallback_repo.clone());
633
634 // A settle's note is part of what the seat said, so a resumed seat
635 // reads its own report back.
636 let bodies: Vec<String> = q
637 .thread
638 .iter()
639 .map(|t| match &t.note {
640 Some(n) => format!("{}\n(note: {n})", t.body),
641 None => t.body.clone(),
642 })
643 .collect();
644 let thread: Vec<(&str, &str)> = q
645 .thread
646 .iter()
647 .zip(&bodies)
648 .map(|(t, body)| {
649 (
650 if t.who == Who::Operator {
651 "operator"
652 } else {
653 "agent"
654 },
655 body.as_str(),
656 )
657 })
658 .collect();
659 let read = &thread[..q.delivered_turns.min(thread.len())];
660 let unread = q.unread_from_owner();
661 let snapshot = q.thread.len();
662 let body = prompt::deputy(&prompt::DeputyPrompt {
663 id: &q.id,
664 summary: &q.summary,
665 detail: &q.detail,
666 brief: &dep.brief,
667 choices: &q.choices,
668 thread: read,
669 unread: unread.as_deref(),
670 resumed,
671 handover: false,
672 kind: kind_of(&q).unwrap_or(Kind::Conduct),
673 language: &cfg.graph.language,
674 });
675
676 // Recorded before the turn starts, so a daemon that dies mid-turn
677 // leaves a start counted and the seat on the record.
678 self.store.beat(&q.id, WaiterKind::Deputy);
679 let starts = dep.starts + 1;
680 let first = seat.clone();
681 self.store.update(&q.id, |r| {
682 if let Some(d) = r.deputy.as_mut() {
683 d.agent = spec.id.clone();
684 d.seat = Some(first);
685 d.starts = starts;
686 }
687 r.waiter = Some(Note {
688 kind: WaiterKind::Deputy,
689 since: now,
690 });
691 Ok(())
692 })?;
693 drop(release);
694 tracing::info!(
695 "question {}: deputy seat {} {} (start {starts}/{MAX_STARTS})",
696 q.short(),
697 seat.key,
698 if resumed { "resuming" } else { "starting" }
699 );
700
701 let left = (deadline(&q, cfg.graph.answer_timeout) - now.as_second()).max(0) as u64;
702 let artifacts = self.store.root().join(format!("{}.deputy", q.id));
703 let cache_dir = cfg.cache_dir();
704 // `magi ask` writes the question record and its lock, which a read-only
705 // sandbox refuses, so the seat cannot be read-only. It is told never to
706 // edit anything; the daemon, not the deputy, applies outcomes.
707 let allow_write = true;
708 // The question store is outside the repository, and `magi ask` writes it.
709 // A merge approval's deputy also files follow-up tasks with `magi task
710 // add`, which writes the queue; that is the one other place it writes.
711 let writable = [
712 self.store.root().to_path_buf(),
713 crate::queue::Queue::open().root().to_path_buf(),
714 ];
715 macro_rules! invocation {
716 ($prompt:expr, $stem:expr, $timeout:expr) => {
717 Invocation {
718 cwd: &cwd,
719 prompt: $prompt,
720 timeout: $timeout,
721 allow_write,
722 sessions: cfg.graph.sessions,
723 artifacts: &artifacts,
724 stem: $stem,
725 // The question's own run key (the task id) is what lets this
726 // seat's `magi ask` pass the ownership check on a conductor
727 // question.
728 run: &q.run,
729 node: NODE,
730 cache_dir: cache_dir.as_deref(),
731 attachments: &[],
732 writable: &writable,
733 }
734 };
735 }
736
737 // A fresh seat first takes a short turn that ends: `agent::invoke` only
738 // learns the CLI's session id when it returns, and the real turn blocks
739 // in `magi ask` for hours, so a daemon stopped mid-wait would otherwise
740 // leave nothing to resume. This turn persists the seat before the wait.
741 let mut early = None;
742 if !resumed && cfg.graph.sessions {
743 let hbody = prompt::deputy(&prompt::DeputyPrompt {
744 id: &q.id,
745 summary: &q.summary,
746 detail: &q.detail,
747 brief: &dep.brief,
748 choices: &q.choices,
749 thread: read,
750 unread: None,
751 resumed: false,
752 handover: true,
753 kind: kind_of(&q).unwrap_or(Kind::Conduct),
754 language: &cfg.graph.language,
755 });
756 let hstem = format!("handover-{starts}");
757 // Bounded by the question's own deadline, not only by its own cap: the
758 // lease it beats would otherwise keep an expired question alive.
759 let hlimit = HANDOVER_TIMEOUT.min(Duration::from_secs(left.max(1)));
760 let hinv = invocation!(&hbody, &hstem, hlimit);
761 let Some(done) = self.drive(&spec, &mut seat, &hinv, &q.id).await else {
762 self.park(&q.id);
763 return Ok(());
764 };
765 let kept = seat.clone();
766 self.store.update(&q.id, |r| {
767 if let Some(d) = r.deputy.as_mut() {
768 d.seat = Some(kept);
769 }
770 Ok(())
771 })?;
772 if !matches!(&done, Ok(o) if o.usable()) {
773 early = Some(done);
774 }
775 }
776
777 // The handover may have been slow: look at the question again, and
778 // measure the long turn from what is left now.
779 let left = if early.is_none() && !resumed && cfg.graph.sessions {
780 let now = Timestamp::now();
781 let again = self.store.get(&q.id)?;
782 let left = (deadline(&again, cfg.graph.answer_timeout) - now.as_second()).max(0) as u64;
783 if !again.status.open() || (left == 0 && again.unread_from_owner().is_none()) {
784 self.park_quietly(&q.id);
785 return Ok(());
786 }
787 left
788 } else {
789 left
790 };
791 let out = match early {
792 Some(done) => done,
793 None => {
794 let stem = format!("turn-{starts}");
795 let timeout = Duration::from_secs(left.max(60) + SLACK_SECS);
796 let inv = invocation!(&body, &stem, timeout);
797 match self.drive(&spec, &mut seat, &inv, &q.id).await {
798 Some(out) => out,
799 None => {
800 self.park(&q.id);
801 return Ok(());
802 }
803 }
804 }
805 };
806
807 let (text, why) = match out {
808 Ok(o) if o.usable() => (Some(o.text.trim().to_owned()), None),
809 Ok(o) if o.timed_out => (None, Some("its turn timed out".to_owned())),
810 Ok(o) if o.quota_exhausted() => (None, Some("the agent is out of quota".to_owned())),
811 Ok(o) => (
812 None,
813 Some(format!("its turn failed (exit {:?})", o.exit_code)),
814 ),
815 Err(e) => (None, Some(format!("the agent could not be started: {e:#}"))),
816 };
817 let kept = seat.clone();
818 self.store.update(&q.id, |r| {
819 if let Some(d) = r.deputy.as_mut() {
820 d.seat = Some(kept);
821 }
822 r.waiter = None;
823 // The owner's words were in the prompt. An agent that read them but
824 // answered in prose instead of `magi ask --thread` still spoke;
825 // keep it on the record so the owner reads it.
826 if let Some(text) = text
827 && unread.is_some()
828 && !text.is_empty()
829 && r.status.open()
830 && r.thread.len() == snapshot
831 {
832 r.delivered_turns = r.delivered_turns.max(snapshot);
833 let choices = r.choices.clone();
834 r.reply(text, choices)?;
835 }
836 Ok(())
837 })?;
838 self.store.drop_lease(&q.id);
839 if let Some(why) = why {
840 tracing::warn!("question {}: the deputy ended: {why}", q.short());
841 notices::raise_in(
842 &self.home,
843 Notice::warn(
844 &format!("deputy-turn:{}", q.id),
845 format!(
846 "Question {} \"{}\": its follow-up agent stopped ({why}).",
847 q.short(),
848 q.summary
849 ),
850 ),
851 );
852 }
853 Ok(())
854 }
855}
856
857/// The runner's loop, until `stop` is asked for.
858pub async fn run(mut deputies: Deputies, stop: crate::daemon::Stop) {
859 while !stop.stopped() {
860 deputies.tick(Timestamp::now());
861 tokio::time::sleep(TICK).await;
862 }
863}