Skip to main content

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