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//! - **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::{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(
188    q: &Question,
189    actions: &BTreeMap<String, ChoiceAction>,
190    home: &Path,
191) -> String {
192    // `Question::run` may also hold a task id; only a run id can be shown.
193    let run = crate::run::is_run_id(&q.run).then_some(q.run.as_str());
194    let knows = if run.is_some() {
195        "You know what the question itself says and what you can learn from \
196         its run by the read-only investigation described below."
197    } else {
198        "You know only what the question itself says."
199    };
200    let mut s = format!(
201        "This question was filed by magi (node `{}`, seat `{}`) with no agent \
202         waiting on it. Magi records the owner's answer and the component that \
203         asked applies it; you apply nothing. {knows} A choice whose effect you \
204         cannot read is not yours to guess: do not settle it, ask the owner \
205         what they mean with `--thread` instead.",
206        q.node, q.seat
207    );
208    if let Some(run) = run {
209        let artifacts = home.join("runs").join(run).join("artifacts");
210        s.push_str(&format!(
211            "\n\nThis question concerns run `{run}`. You may investigate it on \
212             your own, read-only: run `magi show {run}` and read files under \
213             `{}`. When the owner asks about details (for example which part \
214             was wrong), look it up yourself and answer; do not reply that you \
215             cannot see the details, and do not tell the owner to run `magi show` \
216             themselves. Investigating never changes anything: edit, commit, push \
217             and apply nothing. If the run cannot be read, say exactly that and do \
218             not guess at the cause or at what an option does.",
219            artifacts.display()
220        ));
221    }
222    if !q.choices.is_empty() {
223        s.push_str("\n\nWhat each option does:");
224        for c in &q.choices {
225            match actions.get(c) {
226                Some(a) => s.push_str(&format!("\n- `{c}`: also {}", a.describe())),
227                None => s.push_str(&format!(
228                    "\n- `{c}`: recorded as the answer (nothing more is known)"
229                )),
230            }
231        }
232    }
233    s
234}
235
236/// Is picking `label` on `q` something that cannot be taken back, so that
237/// `--settle` must hold the owner's words to the same mechanical standard as a
238/// merge (`land::unhedged`)? Discarding a task deletes it; a divergence answer
239/// drops commits.
240pub fn destructive(q: &Question, label: &str) -> bool {
241    let at = q.choices.iter().position(|c| c == label);
242    match q.node.as_str() {
243        crate::triage::NODE => at == Some(2),
244        crate::triage::DEPS_NODE => at == Some(1),
245        n => n == crate::reconcile::NODE && q.seat == crate::reconcile::SEAT,
246    }
247}
248
249/// Does `--settle` hold `q` to the merge-approval rules (the verbatim quote of
250/// the latest message for `merge`, the whole message for `hold`)? A merge
251/// approval, and any other served question that offers `merge` (the release
252/// watcher's local-mode approval merges just as irreversibly).
253pub fn merge_gated(q: &Question) -> bool {
254    q.node == crate::land::APPROVAL_NODE
255        || (kind_of(q).is_some_and(|k| k != Kind::Conduct)
256            && q.choices.iter().any(|c| c == crate::land::APPROVE))
257}
258
259/// Does `q`'s clock run from `asked_at` and never move on a reply? True for the
260/// questions that something other than the waiter retires: a merge approval
261/// (land) and a release-watch question (the watcher, by silence being a hold).
262pub fn fixed_clock(q: &Question) -> bool {
263    matches!(kind_of(q), Some(Kind::Land | Kind::Release)) || q.node == crate::github_text::ASK_NODE
264}
265
266/// Second after which nobody is to be started or kept on `q`.
267///
268/// A conductor question runs from its last activity (`magi ask --thread`
269/// re-arms it, and the waiter retires it). A merge approval never moves: it
270/// runs from `asked_at`, exactly where `daemon::land_resume_state` abandons it,
271/// so a conversation cannot stretch the hold and land stays the only place that
272/// retires one.
273pub fn deadline(q: &Question, default_timeout: u64) -> i64 {
274    let secs = if q.answer_timeout > 0 {
275        q.answer_timeout
276    } else {
277        default_timeout
278    };
279    let from = if fixed_clock(q) {
280        q.asked_at.as_second()
281    } else {
282        q.last_activity()
283    };
284    from.saturating_add(secs as i64)
285}
286
287/// Can `magi serve` start a deputy at all under `cfg`? Not when deputies are
288/// switched off (`daemon.max_deputies = 0`) or the config could not be read.
289///
290/// Also false when the agent the deputy would run as cannot be resolved
291/// (`agent` is the deputy's recorded agent, empty when it has none): `turn`
292/// fails on that before it counts a start, so it would never reach
293/// [`MAX_STARTS`].
294pub fn can_start(cfg: Option<&Config>, agent: &str) -> bool {
295    cfg.is_some_and(|c| {
296        c.daemon.max_deputies > 0
297            && ((!agent.is_empty() && c.agent(agent).is_ok()) || c.resolve_roles().is_ok())
298    })
299}
300
301/// The agent a question's deputy records, empty when it has none yet.
302pub fn agent_of(q: &Question) -> &str {
303    q.deputy.as_ref().map_or("", |d| d.agent.as_str())
304}
305
306/// Has this question's deputy run out of starts, or can none ever start, with
307/// the deadline gone?
308///
309/// Then nothing will ever read an unread say, and the waiter must retire the
310/// question anyway instead of deferring to a deputy that no longer starts.
311/// `startable` is [`can_start`]; a fresh lease is the caller's to check.
312pub fn exhausted_past_deadline(
313    q: &Question,
314    startable: bool,
315    default_timeout: u64,
316    now: Timestamp,
317) -> bool {
318    q.status.open()
319        && q.deputy
320            .as_ref()
321            .is_some_and(|d| d.starts >= MAX_STARTS || !startable)
322        && now.as_second() > deadline(q, default_timeout)
323}
324
325/// The deputy runner: its own task inside `magi serve`, beside the waiter.
326pub struct Deputies {
327    store: Questions,
328    home: PathBuf,
329    cfg: Option<Config>,
330    /// Working directory for a question that recorded none.
331    fallback_repo: PathBuf,
332    max: usize,
333    halt: Halt,
334    tasks: JoinSet<String>,
335    inflight: HashSet<String>,
336    /// Earliest next start per question.
337    memo: HashMap<String, Instant>,
338}
339
340impl Deputies {
341    /// A runner over `store`, with at most `max` deputies at once.
342    pub fn new(
343        store: Questions,
344        home: PathBuf,
345        cfg: Option<Config>,
346        fallback_repo: PathBuf,
347        max: usize,
348        halt: Halt,
349    ) -> Self {
350        Self {
351            store,
352            home,
353            cfg,
354            fallback_repo,
355            max,
356            halt,
357            tasks: JoinSet::new(),
358            inflight: HashSet::new(),
359            memo: HashMap::new(),
360        }
361    }
362
363    fn default_timeout(&self) -> u64 {
364        self.cfg
365            .as_ref()
366            .map_or(Config::default().graph.answer_timeout, |c| {
367                c.graph.answer_timeout
368            })
369    }
370
371    fn reap(&mut self) {
372        while let Some(done) = self.tasks.try_join_next() {
373            if let Ok(id) = done {
374                self.inflight.remove(&id);
375            }
376        }
377        if self.tasks.is_empty() {
378            self.inflight.clear();
379        }
380    }
381
382    /// Give a question the record a deputy needs: the deputy itself (with a
383    /// brief rebuilt from what the question stored, when it was filed before
384    /// deputies existed - a lost reason is not invented) and the deadline.
385    ///
386    /// A conductor question also gets the working directory the waiter and
387    /// `magi ask --wait` use. Every other kind never does: `cwd` is what makes a
388    /// question the waiter's (it would resume and expire it beside the deputy),
389    /// and a release-watch question, a merge approval, a triage question and a
390    /// fallback one have no asker to resume.
391    fn attach(&self, q: &Question, kind: Kind) -> Option<Question> {
392        let default_timeout = self.default_timeout();
393        let repo = self.fallback_repo.to_string_lossy().into_owned();
394        let state = match kind {
395            Kind::Land => crate::run::RunState::load(&q.run).ok(),
396            Kind::Conduct | Kind::Release | Kind::Triage | Kind::Generic => None,
397        };
398        // The deadline land itself enforces for this run.
399        let timeout = state
400            .as_ref()
401            .map_or(default_timeout, |s| s.config.graph.answer_timeout);
402        self.store
403            .update(&q.id, |r| {
404                if r.deputy.is_none() {
405                    r.deputy = Some(Deputy::new(match kind {
406                        Kind::Conduct => brief(&r.run, &r.detail, &r.choices, &r.actions),
407                        Kind::Land => crate::land::deputy_brief(r, state.as_ref()),
408                        Kind::Release => crate::release_watch::deputy_brief(r, &self.home),
409                        Kind::Triage => crate::triage::deputy_brief(
410                            r,
411                            &crate::queue::Queue::at(self.home.join("queue")),
412                        ),
413                        Kind::Generic => generic_brief(r, &r.actions, &self.home),
414                    }));
415                }
416                if kind == Kind::Conduct && r.cwd.is_none() {
417                    r.cwd = Some(repo.clone());
418                }
419                if r.answer_timeout == 0 {
420                    r.answer_timeout = match kind {
421                        Kind::Conduct => default_timeout,
422                        Kind::Land => timeout,
423                        Kind::Release | Kind::Triage | Kind::Generic => default_timeout,
424                    };
425                }
426                Ok(())
427            })
428            .map(|(r, ())| r)
429            .map_err(|e| tracing::warn!("question {}: cannot attach a deputy: {e:#}", q.short()))
430            .ok()
431    }
432
433    /// Look at every open conductor question once and start a deputy on each
434    /// that has none, within the limits. Returns without waiting for them.
435    pub fn tick(&mut self, now: Timestamp) {
436        self.reap();
437        for q in self.store.list() {
438            if (self.halt)() {
439                return;
440            }
441            let Some(kind) = kind_of(&q) else {
442                continue;
443            };
444            if !q.status.open() {
445                continue;
446            }
447            let needs_cwd = kind == Kind::Conduct && q.cwd.is_none();
448            let q = if q.deputy.is_none() || needs_cwd || q.answer_timeout == 0 {
449                match self.attach(&q, kind) {
450                    Some(q) => q,
451                    None => continue,
452                }
453            } else {
454                q
455            };
456            let Some(dep) = q.deputy.as_ref() else {
457                continue;
458            };
459            if self.inflight.contains(&q.id)
460                || self.store.read_lease(&q.id).is_some_and(|l| l.fresh(now))
461            {
462                continue;
463            }
464            if dep.starts >= MAX_STARTS {
465                self.give_up(&q);
466                continue;
467            }
468            // Past the deadline the waiter retires the question; a say that
469            // arrived in time is still read first.
470            if now.as_second() > deadline(&q, self.default_timeout())
471                && q.unread_from_owner().is_none()
472            {
473                continue;
474            }
475            if self.inflight.len() >= self.max || !can_start(self.cfg.as_ref(), dep.agent.as_str())
476            {
477                continue;
478            }
479            if matches!(self.memo.get(&q.id), Some(until) if Instant::now() < *until) {
480                continue;
481            }
482            self.memo
483                .insert(q.id.clone(), Instant::now() + RESTART_AFTER);
484            self.inflight.insert(q.id.clone());
485            let job = Job {
486                store: self.store.clone(),
487                home: self.home.clone(),
488                cfg: self.cfg.clone(),
489                fallback_repo: self.fallback_repo.clone(),
490                halt: Arc::clone(&self.halt),
491            };
492            let id = q.id.clone();
493            self.tasks.spawn(async move {
494                if let Err(e) = job.turn(&id).await {
495                    tracing::warn!("deputy for question {}: {e:#}", crate::ask::short_id(&id));
496                }
497                id
498            });
499        }
500    }
501
502    /// Wait for every deputy turn in flight. The daemon never calls this; the
503    /// tests do, to look at the record once a turn is over.
504    pub async fn drain(&mut self) {
505        while let Some(done) = self.tasks.join_next().await {
506            if let Ok(id) = done {
507                self.inflight.remove(&id);
508            }
509        }
510        self.inflight.clear();
511    }
512
513    /// Say once that nobody is listening any more.
514    fn give_up(&self, q: &Question) {
515        notices::raise_in(
516            &self.home,
517            Notice::warn(
518                &format!("deputy:{}", q.id),
519                format!(
520                    "Question {} \"{}\": its follow-up agent ended {MAX_STARTS} times \
521                     without an answer and is not restarted. What you say is recorded \
522                     but nothing will read it; answer with one of the choices instead.",
523                    q.short(),
524                    q.summary
525                ),
526            ),
527        );
528    }
529}
530
531/// Everything one deputy turn needs, owned so it can run detached.
532struct Job {
533    store: Questions,
534    home: PathBuf,
535    cfg: Option<Config>,
536    fallback_repo: PathBuf,
537    halt: Halt,
538}
539
540impl Job {
541    /// Run one invocation, beating the lease; `None` when the daemon is
542    /// parking and the turn was dropped.
543    async fn drive(
544        &self,
545        spec: &crate::config::AgentSpec,
546        seat: &mut SeatState,
547        inv: &Invocation<'_>,
548        id: &str,
549    ) -> Option<Result<agent::AgentOutput>> {
550        let fut = agent::invoke(spec, seat, inv);
551        tokio::pin!(fut);
552        let mut beat = tokio::time::interval(Duration::from_secs(1));
553        let mut beats = 0u32;
554        loop {
555            tokio::select! {
556                r = &mut fut => break Some(r),
557                _ = beat.tick() => {
558                    if (self.halt)() {
559                        break None;
560                    }
561                    beats += 1;
562                    if beats % 20 == 0 {
563                        self.store.beat(id, WaiterKind::Deputy);
564                    }
565                }
566            }
567        }
568    }
569
570    /// The question settled or expired during the handover: release it without
571    /// starting the long turn. The start stays counted - the handover ran.
572    fn park_quietly(&self, id: &str) {
573        let _ = self.store.update(id, |r| {
574            r.waiter = None;
575            Ok(())
576        });
577        self.store.drop_lease(id);
578    }
579
580    /// Parking: the turn is dropped and the start refunded. The seat stays as
581    /// last persisted - after the handover turn, resumable.
582    fn park(&self, id: &str) {
583        let _ = self.store.update(id, |r| {
584            if let Some(d) = r.deputy.as_mut() {
585                d.starts = d.starts.saturating_sub(1);
586            }
587            r.waiter = None;
588            Ok(())
589        });
590        self.store.drop_lease(id);
591    }
592
593    async fn turn(&self, id: &str) -> Result<()> {
594        let claim = self.store.root().join(format!("{id}.deputy-claim"));
595        // The claim only covers the decision to start and the write that records
596        // it, so a daemon that dies holding it blocks a restart for a minute,
597        // not for a turn's length. The lease guards the turn itself.
598        if !crate::waiter::take_claim(&claim, CLAIM_STALE) {
599            return Ok(());
600        }
601        let release = crate::waiter::Release(claim);
602
603        // Decided again under the claim, on the record as it is now.
604        let q = self.store.get(id)?;
605        let now = Timestamp::now();
606        let Some(dep) = q.deputy.clone() else {
607            return Ok(());
608        };
609        if !q.status.open()
610            || dep.starts >= MAX_STARTS
611            || self.store.read_lease(&q.id).is_some_and(|l| l.fresh(now))
612        {
613            return Ok(());
614        }
615        let cfg = self
616            .cfg
617            .as_ref()
618            .context("the deputy's configuration is not available")?;
619        let spec = match cfg.agent(&dep.agent) {
620            Ok(s) if !dep.agent.is_empty() => s.clone(),
621            _ => {
622                cfg.resolve_roles()
623                    .context("resolving the deputy's agent")?
624                    .conductor
625            }
626        };
627        // The same CLI conversation when it can be resumed; otherwise a fresh
628        // seat - and either way the prompt carries the whole context.
629        let key = crate::ask::deputy_seat_key(&q.id);
630        let (mut seat, resumed) = match dep.seat.clone() {
631            Some(s)
632                if s.agent == spec.id && agent::has_session(spec.kind, &s, cfg.graph.sessions) =>
633            {
634                (s, true)
635            }
636            _ => (SeatState::new(&key, &spec.id, crate::rng::entropy()), false),
637        };
638        // A release-watch question has no `cwd`; its watch record names the
639        // checkout the pull request belongs to.
640        let recorded = q.cwd.clone().or_else(|| match kind_of(&q) {
641            Some(Kind::Release) => {
642                crate::release_watch::state_for_question(&self.home, &q.id).map(|st| st.repo)
643            }
644            // These record a task id in `run` (or a run id for a divergence's
645            // sibling kinds): the task's repository, when it can be found.
646            Some(Kind::Triage | Kind::Generic) => crate::queue::Queue::at(self.home.join("queue"))
647                .get(&q.run)
648                .ok()
649                .map(|t| t.repo.to_string_lossy().into_owned())
650                .filter(|r| !r.is_empty()),
651            _ => None,
652        });
653        let cwd = recorded
654            .as_deref()
655            .map(PathBuf::from)
656            .filter(|p| p.is_dir())
657            .unwrap_or_else(|| self.fallback_repo.clone());
658
659        // A settle's note is part of what the seat said, so a resumed seat
660        // reads its own report back.
661        let bodies: Vec<String> = q
662            .thread
663            .iter()
664            .map(|t| match &t.note {
665                Some(n) => format!("{}\n(note: {n})", t.body),
666                None => t.body.clone(),
667            })
668            .collect();
669        let thread: Vec<(&str, &str)> = q
670            .thread
671            .iter()
672            .zip(&bodies)
673            .map(|(t, body)| {
674                (
675                    if t.who == Who::Operator {
676                        "operator"
677                    } else {
678                        "agent"
679                    },
680                    body.as_str(),
681                )
682            })
683            .collect();
684        let read = &thread[..q.delivered_turns.min(thread.len())];
685        let unread = q.unread_from_owner();
686        let snapshot = q.thread.len();
687        let body = prompt::deputy(&prompt::DeputyPrompt {
688            id: &q.id,
689            summary: &q.summary,
690            detail: &q.detail,
691            brief: &dep.brief,
692            choices: &q.choices,
693            thread: read,
694            unread: unread.as_deref(),
695            resumed,
696            handover: false,
697            kind: kind_of(&q).unwrap_or(Kind::Conduct),
698            language: &cfg.graph.language,
699        });
700
701        // Recorded before the turn starts, so a daemon that dies mid-turn
702        // leaves a start counted and the seat on the record.
703        self.store.beat(&q.id, WaiterKind::Deputy);
704        let starts = dep.starts + 1;
705        let first = seat.clone();
706        self.store.update(&q.id, |r| {
707            if let Some(d) = r.deputy.as_mut() {
708                d.agent = spec.id.clone();
709                d.seat = Some(first);
710                d.starts = starts;
711            }
712            r.waiter = Some(Note {
713                kind: WaiterKind::Deputy,
714                since: now,
715            });
716            Ok(())
717        })?;
718        drop(release);
719        tracing::info!(
720            "question {}: deputy seat {} {} (start {starts}/{MAX_STARTS})",
721            q.short(),
722            seat.key,
723            if resumed { "resuming" } else { "starting" }
724        );
725
726        let left = (deadline(&q, cfg.graph.answer_timeout) - now.as_second()).max(0) as u64;
727        let artifacts = self.store.root().join(format!("{}.deputy", q.id));
728        let cache_dir = cfg.cache_dir();
729        // `magi ask` writes the question record and its lock, which a read-only
730        // sandbox refuses, so the seat cannot be read-only. It is told never to
731        // edit anything; the daemon, not the deputy, applies outcomes.
732        let allow_write = true;
733        // The question store is outside the repository, and `magi ask` writes it.
734        // A merge approval's deputy also files follow-up tasks with `magi task
735        // add`, which writes the queue; that is the one other place it writes.
736        let writable = [
737            self.store.root().to_path_buf(),
738            crate::queue::Queue::open().root().to_path_buf(),
739        ];
740        macro_rules! invocation {
741            ($prompt:expr, $stem:expr, $timeout:expr) => {
742                Invocation {
743                    cwd: &cwd,
744                    prompt: $prompt,
745                    timeout: $timeout,
746                    allow_write,
747                    unsandboxed: false,
748                    sessions: cfg.graph.sessions,
749                    artifacts: &artifacts,
750                    stem: $stem,
751                    // The question's own run key (the task id) is what lets this
752                    // seat's `magi ask` pass the ownership check on a conductor
753                    // question.
754                    run: &q.run,
755                    node: NODE,
756                    cache_dir: cache_dir.as_deref(),
757                    attachments: &[],
758                    writable: &writable,
759                }
760            };
761        }
762
763        // A fresh seat first takes a short turn that ends: `agent::invoke` only
764        // learns the CLI's session id when it returns, and the real turn blocks
765        // in `magi ask` for hours, so a daemon stopped mid-wait would otherwise
766        // leave nothing to resume. This turn persists the seat before the wait.
767        let mut early = None;
768        if !resumed && cfg.graph.sessions {
769            let hbody = prompt::deputy(&prompt::DeputyPrompt {
770                id: &q.id,
771                summary: &q.summary,
772                detail: &q.detail,
773                brief: &dep.brief,
774                choices: &q.choices,
775                thread: read,
776                unread: None,
777                resumed: false,
778                handover: true,
779                kind: kind_of(&q).unwrap_or(Kind::Conduct),
780                language: &cfg.graph.language,
781            });
782            let hstem = format!("handover-{starts}");
783            // Bounded by the question's own deadline, not only by its own cap: the
784            // lease it beats would otherwise keep an expired question alive.
785            let hlimit = HANDOVER_TIMEOUT.min(Duration::from_secs(left.max(1)));
786            let hinv = invocation!(&hbody, &hstem, hlimit);
787            let Some(done) = self.drive(&spec, &mut seat, &hinv, &q.id).await else {
788                self.park(&q.id);
789                return Ok(());
790            };
791            let kept = seat.clone();
792            self.store.update(&q.id, |r| {
793                if let Some(d) = r.deputy.as_mut() {
794                    d.seat = Some(kept);
795                }
796                Ok(())
797            })?;
798            if !matches!(&done, Ok(o) if o.usable()) {
799                early = Some(done);
800            }
801        }
802
803        // The handover may have been slow: look at the question again, and
804        // measure the long turn from what is left now.
805        let left = if early.is_none() && !resumed && cfg.graph.sessions {
806            let now = Timestamp::now();
807            let again = self.store.get(&q.id)?;
808            let left = (deadline(&again, cfg.graph.answer_timeout) - now.as_second()).max(0) as u64;
809            if !again.status.open() || (left == 0 && again.unread_from_owner().is_none()) {
810                self.park_quietly(&q.id);
811                return Ok(());
812            }
813            left
814        } else {
815            left
816        };
817        let out = match early {
818            Some(done) => done,
819            None => {
820                let stem = format!("turn-{starts}");
821                let timeout = Duration::from_secs(left.max(60) + SLACK_SECS);
822                let inv = invocation!(&body, &stem, timeout);
823                match self.drive(&spec, &mut seat, &inv, &q.id).await {
824                    Some(out) => out,
825                    None => {
826                        self.park(&q.id);
827                        return Ok(());
828                    }
829                }
830            }
831        };
832
833        let (text, why) = match out {
834            Ok(o) if o.usable() => (Some(o.text.trim().to_owned()), None),
835            Ok(o) if o.timed_out => (None, Some("its turn timed out".to_owned())),
836            Ok(o) if o.quota_exhausted() => (None, Some("the agent is out of quota".to_owned())),
837            Ok(o) => (
838                None,
839                Some(format!("its turn failed (exit {:?})", o.exit_code)),
840            ),
841            Err(e) => (None, Some(format!("the agent could not be started: {e:#}"))),
842        };
843        let kept = seat.clone();
844        self.store.update(&q.id, |r| {
845            if let Some(d) = r.deputy.as_mut() {
846                d.seat = Some(kept);
847            }
848            r.waiter = None;
849            // The owner's words were in the prompt. An agent that read them but
850            // answered in prose instead of `magi ask --thread` still spoke;
851            // keep it on the record so the owner reads it.
852            if let Some(text) = text
853                && unread.is_some()
854                && !text.is_empty()
855                && r.status.open()
856                && r.thread.len() == snapshot
857            {
858                r.delivered_turns = r.delivered_turns.max(snapshot);
859                let choices = r.choices.clone();
860                r.reply(text, choices)?;
861            }
862            Ok(())
863        })?;
864        self.store.drop_lease(&q.id);
865        if let Some(why) = why {
866            tracing::warn!("question {}: the deputy ended: {why}", q.short());
867            notices::raise_in(
868                &self.home,
869                Notice::warn(
870                    &format!("deputy-turn:{}", q.id),
871                    format!(
872                        "Question {} \"{}\": its follow-up agent stopped ({why}).",
873                        q.short(),
874                        q.summary
875                    ),
876                ),
877            );
878        }
879        Ok(())
880    }
881}
882
883/// The runner's loop, until `stop` is asked for.
884pub async fn run(mut deputies: Deputies, stop: crate::daemon::Stop) {
885    while !stop.stopped() {
886        deputies.tick(Timestamp::now());
887        tokio::time::sleep(TICK).await;
888    }
889}