Skip to main content

magi/
conduct.rs

1//! The conductor: a single agent seat that arranges the queue.
2//!
3//! [`crate::queue`] orders runnable tasks by `priority` alone; nothing in it
4//! can express that one task should wait for another, or that a task whose
5//! last run stopped short deserves a second look before the loop blindly
6//! retries it. Once per polling cycle, [`Conductor::maybe_run`] shows one
7//! agent seat three things - the runnable tasks, the tasks stuck `running`
8//! with no live daemon behind them, and the `failed`/`held` tasks nobody has
9//! decided about yet - and asks it to decide `blocked_by` for the first and a
10//! recovery for the other two. Everything else about the loop - which
11//! unblocked task runs next, one at a time, in `priority` order - is
12//! unchanged; see `crate::daemon`.
13//!
14//! # What the conductor may not do
15//!
16//! [`Decision`] has no field for `priority`, for deleting a task, or for
17//! touching git, a worktree, or a branch directly. [`Recovery::Review`] only
18//! ever reopens a branch `crate::daemon` itself resolved from the task's own
19//! run record ([`surviving_branch`]) - never a name the model wrote - through
20//! `crate::graph::Runner::review`, which reviews and verifies but never
21//! rewrites history.
22//!
23//! # Non-blocking by construction
24//!
25//! [`crate::ask::ask_and_wait`] is never called from here, and the prompt
26//! tells the model the same: that CLI command blocks until a human answers,
27//! and calling it from inside the conductor's own invocation would park the
28//! whole polling loop behind one task's question. Instead a decision that
29//! wants the operator's judgement carries a `question` field, and [`apply`]
30//! files it with [`Question::new`] and [`Questions::put`] and moves on in the
31//! same call. Something still has to read the owner's reply: the question
32//! carries a [`Deputy`] brief, and `magi serve` starts a short-lived seat of
33//! its own for it ([`crate::deputy`]) - never this loop.
34//!
35//! # Fails soft, always
36//!
37//! [`Conductor::maybe_run`] never returns an error: an unusable roster, a
38//! timed-out invocation, or a reply [`verdict::extract_json`] cannot parse are
39//! all logged and treated as "this cycle changes nothing." `crate::daemon`'s
40//! loop always falls through to its own `Queue::next_runnable` regardless of
41//! what happened here.
42
43use std::collections::BTreeSet;
44use std::path::{Path, PathBuf};
45use std::time::Duration;
46
47use anyhow::{Context as _, Result};
48use serde::Deserialize;
49
50use crate::agent::{self, Invocation, SeatState};
51use crate::ask::{Deputy, Question, Questions};
52use crate::config::Config;
53use crate::prompt;
54use crate::queue::{Queue, Task, TaskStatus};
55use crate::run::RunState;
56use crate::verdict;
57
58/// Seat name for the conductor's own CLI-side conversation, scoped away from
59/// every other seat magi ever opens - the same rule every other seat follows.
60const SEAT: &str = "conduct";
61
62/// Node name reported to the invoked agent (`MAGI_NODE`) and recorded on any
63/// question it files, so an operator reading the questions list can tell a
64/// conductor's question from one a run's own agent asked.
65pub const NODE: &str = "conduct";
66
67/// Wall-clock limit for one conductor turn. The conductor reads a queue
68/// listing and replies with json; it does not implement anything or run a
69/// build, so this is short - the same order of magnitude as
70/// `crate::chat`'s own single-turn, no-write invocations.
71const TURN_TIMEOUT: Duration = Duration::from_secs(300);
72
73/// What the conductor may choose for a `running`-but-stalled or a
74/// `failed`/`held` task.
75#[derive(Debug, Clone, Copy, PartialEq, Eq, Deserialize)]
76#[serde(rename_all = "lowercase")]
77pub enum Recovery {
78    /// Put it back in line, attempts reset - the same effect as
79    /// `magi task release`.
80    Requeue,
81    /// Leave it for a human, unchanged otherwise - the same effect as
82    /// `magi task hold`.
83    Hold,
84    /// Reopen the task's surviving branch as a review-only pass
85    /// (`crate::graph::Runner::review`) instead of competing from scratch.
86    /// Only takes effect when [`surviving_branch`] can actually name one;
87    /// otherwise `crate::daemon` falls back to [`Recovery::Requeue`].
88    Review,
89    /// Close the task outright, as [`Task::succeed`] - the same effect as
90    /// `magi task done`. For a `failed`/`held` task only, never `running`:
91    /// this is for a task whose own goal is already known to be met outside
92    /// the loop entirely (the branch was merged and the worktree cleaned up
93    /// by hand, say) and competing it again would only spend attempts on
94    /// work that has nothing left to do, not for a task that merely stopped
95    /// mid-competition and might still need to run.
96    Done,
97}
98
99/// The conductor's decision for one task. Deliberately has no `priority`
100/// field: see this module's doc.
101#[derive(Debug, Clone, Default, Deserialize)]
102pub struct Decision {
103    /// Task id, expected to be copied verbatim from what it was shown.
104    pub id: String,
105    /// Ids this task should be blocked on. Meaningful only for a runnable
106    /// (`queued`) task; ignored otherwise.
107    #[serde(default)]
108    pub blocked_by: Vec<String>,
109    /// One line explaining the block or the recovery.
110    #[serde(default)]
111    pub reason: Option<String>,
112    /// Recovery for a stalled or finished task. For a runnable (`queued`)
113    /// one, only [`Recovery::Hold`] has any effect - and only when
114    /// `blocked_by` is empty, since a `queued` task with something to wait on
115    /// is handled by that field instead - holding a task the operator has
116    /// already said should not compete again without filing another
117    /// `question` that only restates the same answer. `Requeue`, `Review`,
118    /// and `Done` have no meaning for a task that is already in line.
119    #[serde(default)]
120    pub recovery: Option<Recovery>,
121    /// A question for the operator. When present, [`apply`] files it (unless
122    /// one is already open for this task) and blocks the task on its id
123    /// instead of acting on `blocked_by` or `recovery`.
124    #[serde(default)]
125    pub question: Option<String>,
126    /// Fixed answers for `question`, if it has any. Empty means free text.
127    #[serde(default)]
128    pub choices: Vec<String>,
129}
130
131/// The conductor's whole reply for one cycle.
132///
133/// `decisions` is deliberately **not** `#[serde(default)]`, unlike every
134/// other field in this module. [`verdict::extract_json`] disambiguates
135/// between several balanced `{...}` spans in one reply by trying the type the
136/// caller wants against each of them, last first, and keeping the first that
137/// fits - which only works when a span that is not really the answer can
138/// fail to fit. A `Verdict` with no required field at all would make every
139/// span fit, including a `{}` left by stray trailing prose, and the reply's
140/// real `decisions` - earlier in the text - would never be reached. Requiring
141/// the key costs nothing: the prompt already asks for it on every reply, `[]`
142/// included.
143#[derive(Debug, Clone, Default, Deserialize)]
144pub struct Verdict {
145    /// One entry per task the conductor chose to say something about. A task
146    /// left out of this list is left exactly as it was.
147    pub decisions: Vec<Decision>,
148}
149
150/// A view of a task built for [`prompt::conduct`], shared by the runnable and
151/// stalled sections.
152fn view(t: &Task, max_attempts: usize) -> prompt::ConductTask {
153    prompt::ConductTask {
154        id: t.id.clone(),
155        title: t.title.clone(),
156        instruction: t.instruction.clone(),
157        repo: t.repo.display().to_string(),
158        priority: t.priority,
159        status: t.status.as_str().to_owned(),
160        attempts: t.attempts,
161        max_attempts,
162        last_error: t.last_error.clone(),
163        hold_reason: t.hold_reason.clone(),
164        hold_source: t.hold_source.map(|source| source.label().to_owned()),
165        blocked_by: t.blocked_by.clone(),
166        answers: t
167            .answers
168            .iter()
169            .map(|a| prompt::ConductAnswer {
170                question: a.question.clone(),
171                answer: a.answer.clone(),
172            })
173            .collect(),
174        operator_resume: t.resume_override.as_ref().map(|o| {
175            format!(
176                "the operator explicitly answered \"resume\" at {}; do not hold this \
177                 task again for the same reason unless there is new information",
178                o.at
179            )
180        }),
181    }
182}
183
184/// Did the operator answer a `magi ask` choice with a `resume` action that
185/// this task has not yet acted on? Then `Recovery::Requeue` must not turn the
186/// task into a fresh competition either.
187fn pinned_resume(task: &Task) -> bool {
188    task.resume_override
189        .as_ref()
190        .is_some_and(|o| o.pinned_run.is_some())
191}
192
193/// May the conductor's `Recovery::Hold` take effect on `task`, given an
194/// operator's recorded "resume" answer? The conductor is allowed to override
195/// that answer once - and the override is recorded so `crate::triage` can
196/// show the operator the contradiction - but never after the operator
197/// insisted (`forced`), and never a second time.
198///
199/// Returns `false` when the hold must not be applied (nothing to write).
200/// When it returns `true` the override, if any, has been recorded on `task`.
201fn may_hold(task: &mut Task, reason: &str) -> bool {
202    let Some(o) = task.resume_override.as_mut() else {
203        return true;
204    };
205    if o.forced || o.pinned_run.is_some() {
206        tracing::warn!(
207            "conductor tried to hold task {} after the operator forced a resume: {reason}",
208            task.id
209        );
210        return false;
211    }
212    if o.conductor_rehold.is_some() {
213        return false;
214    }
215    o.conductor_rehold = Some(reason.to_owned());
216    true
217}
218
219/// Severity as a lowercase word, matching how `crate::verdict::Severity` is
220/// spelled everywhere else an operator or a model reads it.
221fn severity_str(s: crate::verdict::Severity) -> &'static str {
222    match s {
223        crate::verdict::Severity::Nit => "nit",
224        crate::verdict::Severity::Minor => "minor",
225        crate::verdict::Severity::Major => "major",
226        crate::verdict::Severity::Blocker => "blocker",
227    }
228}
229
230/// The branch a task's last run left behind, if the tally ever ran on it -
231/// what [`Recovery::Review`] reopens. Derived from `crate::run::RunState`
232/// alone, never from anything the conductor wrote, so a hallucinated branch
233/// name can never reach `crate::graph::Runner::review`.
234fn surviving_branch(task: &Task) -> Option<String> {
235    let last = task.runs.last()?;
236    let state = RunState::load(last).ok()?;
237    state.winner().map(|c| c.branch.clone())
238}
239
240/// The reason to record when the conductor holds a `failed`/`held` task via
241/// [`Recovery::Hold`]. Always `Some`, never a bare `d.reason.clone()`.
242///
243/// [`Task::hold_machine`] only overwrites [`Task::hold_reason`] when given
244/// one, precisely so a stalled task's existing diagnosis survives a
245/// decision that has nothing new to add. That is the right default when the
246/// task was not already `held` - but a task that *is* already `held`, and
247/// is being held again here, is a different case: its current
248/// `hold_reason` may still read as a cause `crate::triage` knows how to
249/// re-check on its own (a disk-pressure message, say - see
250/// `triage::machine_cause_resolved`), and if the model gives no reason,
251/// leaving that text untouched would make this decision - the conductor
252/// choosing, informed by the operator's own answer, to keep the task held
253/// anyway - indistinguishable from the original, never-reconsidered hold.
254/// `crate::triage` would then read the stale text the next time its one
255/// recognised cause looks resolved and release the task straight through
256/// the decision this call was recording. Prepending (not appending) the new
257/// note keeps the old text as context without leaving the string starting
258/// with whatever pattern `crate::triage` matched before.
259fn reaffirmed_hold_reason(task: &Task, d: &Decision) -> String {
260    let note = match &d.reason {
261        Some(reason) => reason.clone(),
262        None => match task.answers.last() {
263            Some(a) => format!(
264                "conduct held this again with no new reason given; last operator \
265                 answer on record: {}",
266                a.answer
267            ),
268            None => "conduct held this again with no reason given".to_owned(),
269        },
270    };
271    match task.hold_reason.as_deref() {
272        Some(prior) if !prior.is_empty() => format!("{note}\n\n(previously: {prior})"),
273        _ => note,
274    }
275}
276
277/// The conductor's stated reason for a hold, or a stand-in when it gave none.
278fn hold_note(d: &Decision) -> String {
279    d.reason
280        .clone()
281        .unwrap_or_else(|| "(no reason given)".to_owned())
282}
283
284/// Everything the conductor is shown about a `failed`/`held` task's last run.
285async fn outcome_for(task: &Task, repo: &Path) -> prompt::ConductOutcome {
286    let Some(run_id) = task.runs.last().cloned() else {
287        return prompt::ConductOutcome {
288            run_id: "(none)".to_owned(),
289            unreadable: Some("this task has not produced a run yet".to_owned()),
290            run_status: None,
291            open_findings: Vec::new(),
292            rounds_used: 0,
293            rounds_max: 0,
294            rounds: Vec::new(),
295            branch: None,
296            branch_head: None,
297            references: None,
298            empty_candidate: false,
299        };
300    };
301    let state = match RunState::load(&run_id) {
302        Ok(s) => s,
303        Err(e) => {
304            // The exact failure this feature exists to stop hiding: a schema
305            // mismatch (or any other unreadable state) must never be treated
306            // as "nothing to recover" - it is surfaced here, verbatim, rather
307            // than swallowed into a quiet re-competition.
308            tracing::warn!(
309                "conductor: could not read run {run_id} for task {}: {e:#}",
310                task.short()
311            );
312            return prompt::ConductOutcome {
313                run_id,
314                unreadable: Some(format!("{e:#}")),
315                run_status: None,
316                open_findings: Vec::new(),
317                rounds_used: 0,
318                rounds_max: 0,
319                rounds: Vec::new(),
320                branch: None,
321                branch_head: None,
322                references: None,
323                empty_candidate: false,
324            };
325        }
326    };
327
328    let finding_view = |f: &crate::verdict::Finding| prompt::ConductFinding {
329        id: f.id.clone(),
330        title: f.title.clone(),
331        severity: severity_str(f.severity).to_owned(),
332    };
333    let open_findings = state
334        .open_findings()
335        .into_iter()
336        .map(finding_view)
337        .collect();
338    let rounds = state
339        .reviews
340        .iter()
341        .map(|r| prompt::ConductRound {
342            round: r.round,
343            findings: r
344                .reviews
345                .iter()
346                .flat_map(|rec| rec.findings.iter())
347                .map(finding_view)
348                .collect(),
349            addressed: r
350                .fix
351                .as_ref()
352                .map(|fx| fx.addressed.clone())
353                .unwrap_or_default(),
354            rejected: r
355                .fix
356                .as_ref()
357                .map(|fx| {
358                    fx.rejected
359                        .iter()
360                        .map(|rej| prompt::ConductRejection {
361                            id: rej.id.clone(),
362                            why: rej.why.clone(),
363                        })
364                        .collect()
365                })
366                .unwrap_or_default(),
367        })
368        .collect();
369    let branch = state.winner().map(|c| c.branch.clone());
370    let branch_head = match &branch {
371        Some(b) => crate::git::rev_parse(repo, b)
372            .await
373            .ok()
374            .map(|h| h.chars().take(8).collect()),
375        None => None,
376    };
377
378    prompt::ConductOutcome {
379        run_id,
380        unreadable: None,
381        run_status: Some(state.status.as_str().to_owned()),
382        open_findings,
383        rounds_used: state.reviews.len(),
384        rounds_max: state.config.graph.review_rounds,
385        rounds,
386        branch,
387        branch_head,
388        references: crate::refs::describe(&state.seeds),
389        empty_candidate: state.merge.as_ref().is_some_and(|m| m.empty),
390    }
391}
392
393/// Whether the task's latest run's winner branch already has an open pull
394/// request into the run's own base - a fact the repository can answer, so the
395/// operator is never asked it. `None` when there is no run or winner to check.
396async fn open_pr_fact(task: &Task, repo: &Path) -> Option<String> {
397    let id = task.runs.last()?;
398    let state = match RunState::load(id) {
399        Ok(s) => s,
400        Err(e) => return Some(format!("could not check open pull requests: {e:#}")),
401    };
402    let branch = state.winner()?.branch.clone();
403    let base = state.base_branch.clone();
404    match crate::land::find_open_pr(repo, &branch, &base).await {
405        Ok(crate::land::OpenPr::None) => None,
406        Ok(crate::land::OpenPr::One { url, .. }) => Some(format!(
407            "pull request {url} is already open for {branch} into {base}"
408        )),
409        Ok(crate::land::OpenPr::Many(urls)) => Some(format!(
410            "several pull requests are already open for {branch} into {base}: {}",
411            urls.join(" ")
412        )),
413        Err(e) => Some(format!("could not check open pull requests: {e:#}")),
414    }
415}
416
417/// Append what the repository says about the branches and commits a task
418/// names to every question the conductor is about to file, so the operator is
419/// never asked a fact git can answer ("is it already on main?"). Done here,
420/// in code, rather than trusted to the model: a question the model wrote
421/// without checking still reaches the operator with the answer attached.
422/// When it cannot be checked the question says so instead of guessing.
423async fn attach_facts(cfg: &Config, repo: &Path, queue: &Queue, verdict: &mut Verdict) {
424    for d in verdict
425        .decisions
426        .iter_mut()
427        .filter(|d| d.question.is_some())
428    {
429        let Ok(task) = queue.get(&d.id) else {
430            continue;
431        };
432        let text = format!("{}\n{}", task.title, task.instruction);
433        let pr_fact = open_pr_fact(&task, repo_for(&task, repo).as_path()).await;
434        if crate::refs::scan(&text).is_empty() {
435            if let (Some(fact), Some(q)) = (pr_fact, d.question.as_mut()) {
436                q.push_str(&format!(
437                    "\n\nChecked against the repository (magi did this, not the model):\n{fact}"
438                ));
439            }
440            continue;
441        }
442        let repo = repo_for(&task, repo);
443        let remote = &cfg.merge.remote;
444        let base = match cfg.merge.base.clone() {
445            Some(b) => Some(b),
446            None => crate::git::current_branch(&repo).await.ok().flatten(),
447        };
448        let facts = match base {
449            Some(base) => {
450                let base_name = base.clone();
451                let tracking = format!("{remote}/{base}");
452                // Best effort: a stale tracking ref would report work that
453                // has since landed as still missing from the base.
454                let refreshed = crate::git::fetch(&repo, remote, &base)
455                    .await
456                    .is_ok_and(|o| o.ok());
457                let against = if crate::git::rev_exists(&repo, &tracking).await {
458                    tracking
459                } else {
460                    base
461                };
462                match crate::git::rev_parse(&repo, &against).await {
463                    Ok(tip) => crate::refs::describe(
464                        &crate::refs::resolve(&repo, &tip, remote, &text).await,
465                    )
466                    .map(|facts| {
467                        if refreshed {
468                            facts
469                        } else {
470                            format!(
471                                "{facts}\n(could not fetch {remote}/{base_name}: this is \
472                                 against the local `{against}`, which may be behind the \
473                                 remote)"
474                            )
475                        }
476                    }),
477                    Err(e) => Some(format!("could not check the repository: {e:#}")),
478                }
479            }
480            None => Some("could not check the repository: no base branch known".to_owned()),
481        };
482        let facts = match (facts, pr_fact) {
483            (Some(f), Some(p)) => Some(format!("{f}\n{p}")),
484            (f, p) => f.or(p),
485        };
486        if let (Some(facts), Some(q)) = (facts, d.question.as_mut()) {
487            q.push_str(&format!(
488                "\n\nChecked against the repository (magi did this, not the model):\n{facts}"
489            ));
490        }
491    }
492}
493
494/// A `failed`/`held` task together with how its last run ended.
495async fn finished_view(t: &Task, repo: &Path, max_attempts: usize) -> prompt::ConductFinished {
496    prompt::ConductFinished {
497        task: view(t, max_attempts),
498        outcome: outcome_for(t, &repo_for(t, repo)).await,
499    }
500}
501
502/// The repository containing a task's branch. A task filed without a
503/// repository uses the daemon's repository, exactly as its later attempt does.
504fn repo_for(task: &Task, fallback: &Path) -> PathBuf {
505    if task.repo.as_os_str().is_empty() || task.repo == Path::new(".") {
506        fallback.to_path_buf()
507    } else {
508        task.repo.clone()
509    }
510}
511
512/// Apply one decision to the queue and the question store.
513///
514/// Takes the task's own claim before touching it: a model call spans a whole
515/// agent turn, and the queue can have moved on by the time its answer comes
516/// back. A claim that cannot be taken means something else owns this task
517/// right now - most often a live daemon mid-competition on it - so the
518/// conductor's now-stale view of it is dropped rather than raced against; see
519/// `crate::queue::Queue::claim`'s own doc on why a claim is proof, not a
520/// guess.
521///
522/// A decision is matched against the task's *current* status, re-read under
523/// the claim, not against whichever section of the prompt it came from: a
524/// `blocked_by` only takes effect on a `queued` task, and `recovery` only on
525/// one `running` (stalled) or `failed`/`held`, so a decision that no longer
526/// matches what the task actually is - it moved on between the read that
527/// built the prompt and this write - changes nothing.
528fn apply_one(queue: &Queue, questions: &Questions, d: &Decision) -> Result<()> {
529    let _claim = queue
530        .claim(&d.id)
531        .with_context(|| format!("task {} is claimed elsewhere right now", d.id))?;
532    let mut task = queue.get(&d.id).context("no such task")?;
533
534    // A conductor answer is never operator authorization.  In particular,
535    // do this before questions and blocking too: either would reclassify a
536    // manual hold and let a later deterministic resolver queue it.
537    if task.operator_held() {
538        return Ok(());
539    }
540
541    // `crate::triage` is already running its own question-and-answer cycle
542    // on this hold - see that module's doc on why it never calls
543    // `Task::block`, and `crate::triage::pending_for`'s own doc on why "open"
544    // alone is not enough here: an answered-but-not-yet-applied triage
545    // question is still triage's to finish, since `crate::triage::run_once`
546    // only runs on a fully idle tick and only ever looks at tasks still
547    // `held`. Blocking this task out from under it - even on an unrelated
548    // question - would move it to `Blocked`, and the pending triage answer
549    // would never be read back.
550    if task.status == TaskStatus::Held && crate::triage::pending_for(questions, &task) {
551        return Ok(());
552    }
553
554    // The operator's `resume` action is waiting for its turn: a question or a
555    // dependency would park the task before the daemon can run it.
556    let pinned = pinned_resume(&task);
557    if pinned && (d.question.is_some() || !d.blocked_by.is_empty()) {
558        return Ok(());
559    }
560
561    if let Some(text) = &d.question {
562        if task.status == TaskStatus::Done {
563            return Ok(());
564        }
565        // `Question::run` is a task id for conductor questions, but ordinary
566        // graph questions use it as a run id. A coincidental equality must
567        // not block this task on an answer meant for another node.
568        let question_id = match questions
569            .list()
570            .into_iter()
571            .find(|q| q.status.open() && q.node == NODE && q.run == task.id)
572        {
573            Some(existing) => existing.id,
574            // No cap on settled answers: a question the operator can answer
575            // is never refused. The task goes `Blocked` and at most one
576            // conductor question is open at a time (reused above), so each
577            // further question costs one human answer. Run attempts are
578            // bounded by `max_attempts` and the queue's attempt accounting.
579            // A model rewording the same question can still recur once per
580            // answer; that cost falls on the operator, by choice.
581            None => {
582                let mut q = Question::new(
583                    task.id.clone(),
584                    NODE.to_owned(),
585                    SEAT.to_owned(),
586                    text.clone(),
587                    d.reason.clone().unwrap_or_default(),
588                    d.choices.clone(),
589                );
590                // The wait is handed to a deputy (`crate::deputy`), so a
591                // free-text reply reaches an agent that knows why this was
592                // asked. The conductor still never blocks: this only records
593                // what the deputy will be told, and `serve` starts it.
594                q.deputy = Some(Deputy::new(crate::deputy::brief(
595                    &task.id,
596                    d.reason.as_deref().unwrap_or_default(),
597                    &q.choices,
598                    &q.actions,
599                )));
600                if task.repo.is_dir() && task.repo != Path::new(".") {
601                    q.cwd = Some(task.repo.to_string_lossy().into_owned());
602                }
603                questions.put(&mut q)?;
604                q.id
605            }
606        };
607        task.block(vec![question_id], d.reason.clone());
608        return queue.put(&mut task);
609    }
610
611    match task.status {
612        TaskStatus::Queued if !d.blocked_by.is_empty() => {
613            task.block(d.blocked_by.clone(), d.reason.clone());
614            queue.put(&mut task)?;
615        }
616        // See `Decision::recovery`'s doc: the only lever a runnable task has
617        // besides `blocked_by` is holding it outright, for a task whose
618        // answers already say it should not compete again.
619        TaskStatus::Queued if d.recovery == Some(Recovery::Hold) => {
620            if may_hold(&mut task, &hold_note(d)) {
621                task.hold_machine(d.reason.clone());
622                queue.put(&mut task)?;
623            }
624        }
625        TaskStatus::Running => match d.recovery {
626            Some(Recovery::Requeue) if pinned_resume(&task) => {}
627            Some(Recovery::Requeue) => {
628                task.requeue();
629                queue.put(&mut task)?;
630            }
631            Some(Recovery::Hold) if may_hold(&mut task, &hold_note(d)) => {
632                task.hold_machine(d.reason.clone());
633                queue.put(&mut task)?;
634            }
635            // `Review` reopens a branch, which only makes sense once a run
636            // has actually stopped; a task still `running` has nothing to
637            // reopen yet.
638            _ => {}
639        },
640        TaskStatus::Failed | TaskStatus::Held => match d.recovery {
641            // The operator picked a specific run to continue; a fresh
642            // competition would throw that run away.
643            Some(Recovery::Requeue) if pinned_resume(&task) => {}
644            Some(Recovery::Requeue) => {
645                task.requeue();
646                queue.put(&mut task)?;
647            }
648            Some(Recovery::Hold) => {
649                if may_hold(&mut task, &hold_note(d)) {
650                    task.hold_machine(Some(reaffirmed_hold_reason(&task, d)));
651                    queue.put(&mut task)?;
652                }
653            }
654            Some(Recovery::Review) => {
655                if let Some(branch) = surviving_branch(&task) {
656                    task.request_review(branch);
657                    queue.put(&mut task)?;
658                }
659                // No survivable branch: a decision naming `review` here is
660                // simply not actionable, and is dropped rather than guessed
661                // at - `crate::daemon` applies the same "no branch, no
662                // review" rule again, from its own read, right before it
663                // would actually start the run.
664            }
665            Some(Recovery::Done) => {
666                task.succeed();
667                // Same reason `magi task done` and the web UI's equivalent
668                // do this: the conductor closing a task by hand is just as
669                // much a finished story as the loop's own settle path.
670                crate::daemon::supersede_prior_runs(&task, &crate::run::home());
671                queue.put(&mut task)?;
672            }
673            None => {}
674        },
675        // `queued` with nothing to block on, `done`, or already `blocked`:
676        // nothing for this decision to do.
677        _ => {}
678    }
679    Ok(())
680}
681
682/// Apply every decision in `verdict`. A single bad decision - a task id that
683/// no longer exists, one already claimed elsewhere - is logged and skipped
684/// rather than losing every other decision in the same reply.
685pub fn apply(queue: &Queue, questions: &Questions, verdict: &Verdict) -> Result<()> {
686    for d in &verdict.decisions {
687        if let Err(e) = apply_one(queue, questions, d) {
688            tracing::warn!("conductor decision for task {}: {e:#}", d.id);
689        }
690    }
691    Ok(())
692}
693
694/// Where the conductor's seat is kept between turns, so the daemon waiter can
695/// resume it when the owner answers a question it asked.
696///
697/// A graph seat's session lives in its run's `RunState`; the conductor has no
698/// run, and until this file its seat existed only in memory.
699pub fn seat_path(home: &Path) -> PathBuf {
700    home.join("conduct").join("seat.json")
701}
702
703/// The seat the conductor last used, if it has taken a turn.
704pub fn load_seat(home: &Path) -> Option<SeatState> {
705    serde_json::from_str(&std::fs::read_to_string(seat_path(home)).ok()?).ok()
706}
707
708fn busy_path(home: &Path) -> PathBuf {
709    home.join("conduct").join("busy")
710}
711
712/// Is a conductor turn in flight right now?
713///
714/// A marker file rather than shared memory because the waiter reads it and the
715/// conductor is not necessarily in its process. A marker older than one turn's
716/// timeout (plus slack) was left by a turn that died and counts for nothing.
717pub fn busy(home: &Path) -> bool {
718    std::fs::metadata(busy_path(home))
719        .and_then(|m| m.modified())
720        .ok()
721        .and_then(|t| t.elapsed().ok())
722        .is_some_and(|age| age < TURN_TIMEOUT + Duration::from_secs(30))
723}
724
725/// Removes the busy marker when a turn ends, however it ends.
726struct Busy(PathBuf);
727
728impl Busy {
729    fn mark(home: &Path) -> Self {
730        let path = busy_path(home);
731        if let Some(dir) = path.parent() {
732            let _ = std::fs::create_dir_all(dir);
733        }
734        let _ = std::fs::write(&path, std::process::id().to_string());
735        Self(path)
736    }
737}
738
739impl Drop for Busy {
740    fn drop(&mut self) {
741        let _ = std::fs::remove_file(&self.0);
742    }
743}
744
745/// The conductor's state across polling cycles: its own CLI-side conversation
746/// and the last (revision, stalled ∪ finished ids) pair it actually acted on.
747#[derive(Debug, Default)]
748pub struct Conductor {
749    seat: Option<SeatState>,
750    last_seen: Option<(u64, BTreeSet<String>)>,
751}
752
753impl Conductor {
754    /// A conductor that has never run.
755    #[must_use]
756    pub fn new() -> Self {
757        Self::default()
758    }
759
760    fn snapshot(queue: &Queue, stalled: &[Task], finished: &[Task]) -> (u64, BTreeSet<String>) {
761        let ids = stalled
762            .iter()
763            .chain(finished)
764            .map(|t| t.id.clone())
765            .collect();
766        (queue.revision(), ids)
767    }
768
769    /// Whether calling the conductor could possibly do anything different
770    /// from last time: [`Queue::revision`] moved, or the set of stalled and
771    /// finished task ids changed.
772    ///
773    /// Deliberately **not** "stalled or finished is non-empty" - a task
774    /// sitting stalled or finished with nobody changing anything about it
775    /// must not be re-shown to the model every single poll forever; only a
776    /// change in *which* tasks are stalled or finished, or a queue write
777    /// changing something about a runnable one, is worth another look.
778    ///
779    /// Cheap and config-free on purpose, so `crate::daemon`'s poll loop can
780    /// skip `Config::discover`'s synchronous I/O entirely on a cycle where
781    /// this says no - which [`Conductor::maybe_run`] would otherwise only
782    /// discover after paying for that load. Both ask the identical question,
783    /// from the same [`Conductor::last_seen`], so they can never disagree
784    /// about whether there is anything to look at.
785    #[must_use]
786    pub fn worth_a_look(&self, queue: &Queue, stalled: &[Task], finished: &[Task]) -> bool {
787        self.last_seen.as_ref() != Some(&Self::snapshot(queue, stalled, finished))
788    }
789
790    /// Call the conductor once, unless nothing has changed since the last
791    /// time it was worth calling - see [`Conductor::worth_a_look`], the exact
792    /// same test. Never fatal - see this module's doc.
793    #[allow(clippy::too_many_arguments)]
794    pub async fn maybe_run(
795        &mut self,
796        cfg: &Config,
797        repo: &Path,
798        queue: &Queue,
799        questions: &Questions,
800        home: &Path,
801        queued: &[Task],
802        stalled: &[Task],
803        finished: &[Task],
804        max_attempts: usize,
805    ) {
806        let snapshot = Self::snapshot(queue, stalled, finished);
807        if self.last_seen.as_ref() == Some(&snapshot) {
808            return;
809        }
810        self.last_seen = Some(snapshot);
811        if let Err(e) = self
812            .run_once(
813                cfg,
814                repo,
815                queue,
816                questions,
817                home,
818                queued,
819                stalled,
820                finished,
821                max_attempts,
822            )
823            .await
824        {
825            tracing::warn!("conductor: {e:#}");
826        }
827    }
828
829    #[allow(clippy::too_many_arguments)]
830    async fn run_once(
831        &mut self,
832        cfg: &Config,
833        repo: &Path,
834        queue: &Queue,
835        questions: &Questions,
836        home: &Path,
837        queued: &[Task],
838        stalled: &[Task],
839        finished: &[Task],
840        max_attempts: usize,
841    ) -> Result<()> {
842        if queued.is_empty() && stalled.is_empty() && finished.is_empty() {
843            return Ok(());
844        }
845
846        let primary = cfg
847            .resolve_roles()
848            .context("resolving the conductor seat")?
849            .conductor;
850        // A string keeps today's resolution exactly (no CLI preflight here;
851        // an unavailable CLI fails at invocation). Only a multi-id array
852        // builds a chain, whose members are each tried at most once.
853        let chain = match cfg.roles.conductor.as_ref() {
854            Some(c) if c.ids().len() > 1 => {
855                agent::pick_chain(&cfg.agents, Some(c), &agent::installed, "conductor")?
856            }
857            _ => vec![primary],
858        };
859
860        let runnable_views: Vec<prompt::ConductTask> =
861            queued.iter().map(|t| view(t, max_attempts)).collect();
862        let stalled_views: Vec<prompt::ConductTask> =
863            stalled.iter().map(|t| view(t, max_attempts)).collect();
864        let mut finished_views = Vec::with_capacity(finished.len());
865        for t in finished {
866            finished_views.push(finished_view(t, repo, max_attempts).await);
867        }
868
869        let body = prompt::with_overlay(
870            prompt::conduct(
871                &runnable_views,
872                &stalled_views,
873                &finished_views,
874                &cfg.graph.language,
875            ),
876            cfg.prompts.overlay(NODE),
877        );
878
879        let artifacts = home.join("conduct").join("artifacts");
880        // Bound to a local: `Invocation` only borrows the cache path, and the
881        // `Option<PathBuf>` `cache_dir()` returns has to outlive that borrow.
882        let cache_dir = cfg.cache_dir();
883
884        // One marker for the whole chain: a fallback is still this turn.
885        let busy = Busy::mark(home);
886        let mut result: Result<Verdict> = Err(anyhow::anyhow!("no conductor agent ran"));
887        for spec in &chain {
888            // The prompt carries the whole queue state every cycle, so a
889            // fallback agent needs no history - but its seat is its own.
890            let needs_new_seat = !matches!(&self.seat, Some(s) if s.agent == spec.id);
891            if needs_new_seat {
892                self.seat = Some(SeatState::new(SEAT, &spec.id, crate::rng::entropy()));
893            }
894            let seat = self.seat.as_mut().expect("just ensured a seat exists");
895            // The first agent keeps the plain stem; a fallback's own artifacts
896            // must not overwrite it.
897            let stem = if std::ptr::eq(spec, &chain[0]) {
898                format!("turn-{}", seat.turns + 1)
899            } else {
900                format!("turn-{}-{}", seat.turns + 1, spec.id)
901            };
902            let inv = Invocation {
903                cwd: repo,
904                prompt: &body,
905                timeout: TURN_TIMEOUT,
906                // The conductor never edits anything - it only decides what
907                // blocks a task and what to do about one stuck or finished.
908                allow_write: false,
909                sessions: cfg.graph.sessions,
910                artifacts: &artifacts,
911                stem: &stem,
912                run: NODE,
913                node: NODE,
914                cache_dir: cache_dir.as_deref(),
915                attachments: &[],
916                writable: &[],
917            };
918            let out = agent::invoke(spec, seat, &inv).await;
919            // Kept after every turn, failed or not: the CLI-side conversation is
920            // what a resumed question needs, and it exists once the turn ran.
921            if let Ok(body) = serde_json::to_string(&*seat) {
922                let path = seat_path(home);
923                if let Some(dir) = path.parent() {
924                    let _ = std::fs::create_dir_all(dir);
925                }
926                let _ = std::fs::write(path, body);
927            }
928            let advance = agent::chain_advances(&out);
929            result = match out {
930                Err(e) => Err(e.context("invoking the conductor")),
931                Ok(out) if advance => Err(anyhow::anyhow!(
932                    "no usable reply (exit {:?}, timed out {})",
933                    out.exit_code,
934                    out.timed_out
935                )),
936                Ok(out) => verdict::extract_json(&out.text)
937                    .context("the conductor's reply could not be parsed"),
938            };
939            if result.is_ok() {
940                break;
941            }
942            if chain.len() > 1 {
943                tracing::warn!("conductor: `{}` failed, trying the next agent", spec.id);
944            }
945        }
946        drop(busy);
947        let mut verdict = result?;
948        attach_facts(cfg, repo, queue, &mut verdict).await;
949        apply(queue, questions, &verdict)
950    }
951}
952
953#[cfg(test)]
954mod tests {
955    use std::collections::BTreeMap;
956
957    use tempfile::tempdir;
958
959    use super::*;
960    use crate::ask::{Answer, QuestionStatus};
961    use crate::config::{AgentKind, AgentSpec, Graph};
962    use crate::queue::Source;
963
964    fn mock_agent(dir: &Path, script: &str, env: BTreeMap<String, String>) -> AgentSpec {
965        let path = dir.join("mock-conduct-agent.sh");
966        std::fs::write(&path, script).expect("write mock");
967        AgentSpec {
968            id: "mock".to_owned(),
969            kind: AgentKind::Command,
970            model: None,
971            command: vec!["sh".to_owned(), path.to_string_lossy().into_owned()],
972            extra_args: Vec::new(),
973            env,
974            prompt_delivery: None,
975        }
976    }
977
978    fn config(spec: AgentSpec) -> Config {
979        Config {
980            agents: vec![spec],
981            graph: Graph {
982                language: "en".to_owned(),
983                ..Graph::default()
984            },
985            ..Config::default()
986        }
987    }
988
989    fn task(title: &str) -> Task {
990        Task::new(
991            title.to_owned(),
992            format!("do {title}"),
993            std::path::PathBuf::from("."),
994            Source::Human,
995        )
996    }
997
998    const BROKEN: &str = "#!/bin/sh\ncat >/dev/null\nexit 3\n";
999    const GARBAGE: &str = "#!/bin/sh\ncat >/dev/null\nprintf 'not json at all\\n'\n";
1000
1001    fn env(reply: &str) -> BTreeMap<String, String> {
1002        BTreeMap::from([("MOCK_REPLY".to_owned(), reply.to_owned())])
1003    }
1004
1005    const REPLY: &str = "#!/bin/sh\ncat >/dev/null\nprintf '%s\\n' \"$MOCK_REPLY\"\n";
1006
1007    /// A throwaway repo with one commit on `main` and a second branch ahead
1008    /// of it, so `outcome_for`'s own `git::rev_parse` call has a real head to
1009    /// resolve.
1010    fn init_repo_with_branch(dir: &Path, branch: &str) {
1011        use crate::proc::Quiet as _;
1012        let run = |args: &[&str]| {
1013            let out = std::process::Command::new("git")
1014                .args(args)
1015                .current_dir(dir)
1016                .quiet()
1017                .output()
1018                .expect("spawn git");
1019            assert!(
1020                out.status.success(),
1021                "git {args:?} failed: {}",
1022                String::from_utf8_lossy(&out.stderr)
1023            );
1024        };
1025        run(&["init", "-b", "main"]);
1026        run(&["config", "user.name", "magi test"]);
1027        run(&["config", "user.email", "magi@example.com"]);
1028        std::fs::write(dir.join("README.md"), "# fixture\n").unwrap();
1029        run(&["add", "-A"]);
1030        run(&["commit", "-m", "init"]);
1031        run(&["checkout", "-b", branch]);
1032        std::fs::write(dir.join("change.txt"), "x\n").unwrap();
1033        run(&["add", "-A"]);
1034        run(&["commit", "-m", "candidate work"]);
1035    }
1036
1037    fn review_round_with_finding(
1038        round: usize,
1039        finding_id: &str,
1040        title: &str,
1041        addressed: &[&str],
1042        rejected: &[(&str, &str)],
1043    ) -> crate::run::ReviewRound {
1044        crate::run::ReviewRound {
1045            round,
1046            head: "deadbeef".to_owned(),
1047            verified_head: None,
1048            verified_at: None,
1049            reviews: vec![crate::run::ReviewRecord {
1050                attempts: 0,
1051                reviewer: 1,
1052                agent: "mock".to_owned(),
1053                summary: String::new(),
1054                findings: vec![crate::verdict::Finding {
1055                    id: finding_id.to_owned(),
1056                    severity: crate::verdict::Severity::Major,
1057                    file: None,
1058                    line: None,
1059                    title: title.to_owned(),
1060                    detail: String::new(),
1061                }],
1062                vote: None,
1063                failed: None,
1064                duration_ms: 0,
1065            }],
1066            e2e: Vec::new(),
1067            verify_retried: false,
1068            e2e_deferred: false,
1069            e2e_defer_reason: None,
1070            fix: Some(crate::run::FixRecord {
1071                agent: "mock".to_owned(),
1072                addressed: addressed.iter().map(|s| (*s).to_owned()).collect(),
1073                rejected: rejected
1074                    .iter()
1075                    .map(|(id, why)| crate::verdict::Rejection {
1076                        id: (*id).to_owned(),
1077                        why: (*why).to_owned(),
1078                    })
1079                    .collect(),
1080                notes: String::new(),
1081                committed: false,
1082                failed: None,
1083                duration_ms: 0,
1084                continuation: None,
1085            }),
1086            blocking: 1,
1087            answered: 1,
1088            expected: 1,
1089            clean: false,
1090            progressed: true,
1091            vote_split: false,
1092            reconsideration: Vec::new(),
1093            verdict: None,
1094        }
1095    }
1096
1097    #[test]
1098    fn outcome_for_carries_every_rounds_findings_and_the_branch_head() {
1099        crate::run::set_home(std::env::temp_dir().join("magi-conduct-tests-home"));
1100        let dir = tempdir().unwrap();
1101        let default_repo = dir.path().join("default");
1102        let task_repo = dir.path().join("task");
1103        std::fs::create_dir_all(&default_repo).unwrap();
1104        std::fs::create_dir_all(&task_repo).unwrap();
1105        init_repo_with_branch(&default_repo, "other-branch");
1106        init_repo_with_branch(&task_repo, "magi/f00d/A");
1107
1108        let mut config = Config::default();
1109        config.graph.review_rounds = 6;
1110        let mut state = crate::run::RunState::new(
1111            task_repo.clone(),
1112            "main".to_owned(),
1113            "deadbeef".to_owned(),
1114            "task".to_owned(),
1115            config,
1116        );
1117        state.status = crate::run::RunStatus::Blocked;
1118        state.candidates.push(crate::run::Candidate {
1119            index: 0,
1120            label: 'A',
1121            agent: "mock".to_owned(),
1122            branch: "magi/f00d/A".to_owned(),
1123            worktree: task_repo.clone(),
1124            summary: String::new(),
1125            stat: String::new(),
1126            files: 1,
1127            commits: 1,
1128            empty: false,
1129            failed: None,
1130            verified_noop: None,
1131            duration_ms: 0,
1132            folded: false,
1133        });
1134        state.tally = Some(crate::run::Tally {
1135            first_choice: std::collections::BTreeMap::new(),
1136            borda: std::collections::BTreeMap::new(),
1137            winner: 'A',
1138            rankings: 0,
1139            unanimous_initial: false,
1140            deliberated: false,
1141            changed_votes: 0,
1142            unanimous_final: false,
1143            tie_break: None,
1144            judges: 0,
1145            present: 0,
1146            quorum: 0,
1147            met_quorum: true,
1148            uncontested: Some("solo".to_owned()),
1149        });
1150        state.reviews = vec![
1151            review_round_with_finding(
1152                1,
1153                "R1-1-2",
1154                "answer content is dropped",
1155                &[],
1156                &[("R1-1-2", "the id leaving blocked_by is enough")],
1157            ),
1158            review_round_with_finding(2, "R2-1-3", "answer content is still dropped", &[], &[]),
1159        ];
1160        state.save().unwrap();
1161
1162        let mut t = task("outcome test");
1163        t.repo = task_repo;
1164        t.runs.push(state.id.clone());
1165
1166        let finished = tokio_test_block_on(finished_view(&t, &default_repo, 2));
1167        let outcome = finished.outcome;
1168
1169        assert!(outcome.unreadable.is_none());
1170        assert_eq!(outcome.run_status.as_deref(), Some("blocked"));
1171        assert_eq!(outcome.rounds_used, 2);
1172        assert_eq!(outcome.rounds_max, 6);
1173        assert_eq!(outcome.rounds.len(), 2);
1174        assert_eq!(outcome.rounds[0].findings[0].id, "R1-1-2");
1175        assert_eq!(outcome.rounds[0].rejected[0].id, "R1-1-2");
1176        assert!(outcome.rounds[1].addressed.is_empty());
1177        assert!(outcome.rounds[1].rejected.is_empty());
1178        assert_eq!(outcome.branch.as_deref(), Some("magi/f00d/A"));
1179        assert!(
1180            outcome.branch_head.is_some(),
1181            "a real branch must resolve a head commit: {outcome:?}"
1182        );
1183    }
1184
1185    /// A tiny single-threaded block-on, so an `async fn` can be exercised
1186    /// from a plain `#[test]` without pulling `tokio::test`'s multi-thread
1187    /// runtime into a test that does no other async work.
1188    fn tokio_test_block_on<F: std::future::Future>(f: F) -> F::Output {
1189        tokio::runtime::Builder::new_current_thread()
1190            .enable_all()
1191            .build()
1192            .unwrap()
1193            .block_on(f)
1194    }
1195
1196    #[test]
1197    fn view_carries_a_tasks_recorded_answers_into_the_conductor_prompt_input() {
1198        let mut t = task("answered");
1199        t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
1200        let v = view(&t, 2);
1201        assert_eq!(v.answers.len(), 1);
1202        assert_eq!(v.answers[0].question, "Which backend?");
1203        assert_eq!(v.answers[0].answer, "SQLite");
1204    }
1205
1206    #[test]
1207    fn a_dependency_decision_blocks_the_task_and_leaves_priority_alone() {
1208        let dir = tempdir().unwrap();
1209        let queue = Queue::at(dir.path().join("queue"));
1210        let questions = Questions::at(dir.path().join("questions"));
1211        let mut a = task("a");
1212        a.priority = 9;
1213        queue.put(&mut a).unwrap();
1214
1215        let verdict = Verdict {
1216            decisions: vec![Decision {
1217                id: a.id.clone(),
1218                blocked_by: vec!["20260101-000000-dead".to_owned()],
1219                reason: Some("waits on the other task".to_owned()),
1220                recovery: None,
1221                question: None,
1222                choices: Vec::new(),
1223            }],
1224        };
1225        apply(&queue, &questions, &verdict).unwrap();
1226
1227        let back = queue.get(&a.id).unwrap();
1228        assert_eq!(back.status, TaskStatus::Blocked);
1229        assert_eq!(back.blocked_by, ["20260101-000000-dead"]);
1230        assert_eq!(
1231            back.priority, 9,
1232            "the conductor's reply cannot carry priority"
1233        );
1234    }
1235
1236    #[test]
1237    fn a_question_decision_files_one_and_blocks_on_its_id() {
1238        let dir = tempdir().unwrap();
1239        let queue = Queue::at(dir.path().join("queue"));
1240        let questions = Questions::at(dir.path().join("questions"));
1241        let mut t = task("ambiguous");
1242        queue.put(&mut t).unwrap();
1243
1244        let verdict = Verdict {
1245            decisions: vec![Decision {
1246                id: t.id.clone(),
1247                blocked_by: Vec::new(),
1248                reason: Some("which backend?".to_owned()),
1249                recovery: None,
1250                question: Some("Which storage backend?".to_owned()),
1251                choices: vec!["SQLite".to_owned(), "Redis".to_owned()],
1252            }],
1253        };
1254        apply(&queue, &questions, &verdict).unwrap();
1255
1256        let back = queue.get(&t.id).unwrap();
1257        assert_eq!(back.status, TaskStatus::Blocked);
1258        assert_eq!(back.blocked_by.len(), 1);
1259        let q = questions.get(&back.blocked_by[0]).unwrap();
1260        assert_eq!(q.summary, "Which storage backend?");
1261        assert_eq!(q.node, NODE);
1262        assert!(q.status.open());
1263    }
1264
1265    #[test]
1266    fn a_task_with_an_open_question_already_reuses_it_rather_than_filing_a_second_one() {
1267        let dir = tempdir().unwrap();
1268        let queue = Queue::at(dir.path().join("queue"));
1269        let questions = Questions::at(dir.path().join("questions"));
1270        let mut t = task("asked once");
1271        queue.put(&mut t).unwrap();
1272
1273        let decision = Decision {
1274            id: t.id.clone(),
1275            reason: Some("still deciding".to_owned()),
1276            question: Some("Which backend?".to_owned()),
1277            ..Decision::default()
1278        };
1279        apply(
1280            &queue,
1281            &questions,
1282            &Verdict {
1283                decisions: vec![decision.clone()],
1284            },
1285        )
1286        .unwrap();
1287        assert_eq!(questions.list().len(), 1);
1288        let first_question_id = queue.get(&t.id).unwrap().blocked_by[0].clone();
1289
1290        // An operator releasing the blocked task by hand, without answering,
1291        // puts it back at `Queued` while the question stays open - exactly
1292        // the case the guard in `apply_one` exists for: a later cycle
1293        // proposing the very same question must reuse it, not file a second.
1294        let mut released = queue.get(&t.id).unwrap();
1295        released.release();
1296        queue.put(&mut released).unwrap();
1297
1298        apply(
1299            &queue,
1300            &questions,
1301            &Verdict {
1302                decisions: vec![decision],
1303            },
1304        )
1305        .unwrap();
1306        assert_eq!(questions.list().len(), 1, "no duplicate question was filed");
1307        let after = queue.get(&t.id).unwrap();
1308        assert_eq!(
1309            after.blocked_by,
1310            [first_question_id],
1311            "the existing open question is reused, not replaced"
1312        );
1313    }
1314
1315    #[test]
1316    fn a_same_id_question_from_another_node_is_not_reused() {
1317        let dir = tempdir().unwrap();
1318        let queue = Queue::at(dir.path().join("queue"));
1319        let questions = Questions::at(dir.path().join("questions"));
1320        let mut t = task("must ask the conductor");
1321        queue.put(&mut t).unwrap();
1322
1323        let mut unrelated = Question::new(
1324            t.id.clone(),
1325            "review".to_owned(),
1326            "reviewer-1".to_owned(),
1327            "An unrelated review question".to_owned(),
1328            String::new(),
1329            Vec::new(),
1330        );
1331        questions.put(&mut unrelated).unwrap();
1332
1333        apply(
1334            &queue,
1335            &questions,
1336            &Verdict {
1337                decisions: vec![Decision {
1338                    id: t.id.clone(),
1339                    question: Some("Which backend?".to_owned()),
1340                    ..Decision::default()
1341                }],
1342            },
1343        )
1344        .unwrap();
1345
1346        let blocked_by = &queue.get(&t.id).unwrap().blocked_by;
1347        assert_eq!(blocked_by.len(), 1);
1348        assert_ne!(blocked_by[0], unrelated.id);
1349        assert!(questions.get(&unrelated.id).unwrap().status.open());
1350        assert_eq!(questions.get(&blocked_by[0]).unwrap().node, NODE);
1351    }
1352
1353    #[test]
1354    fn answering_the_question_lets_the_resolver_clear_the_block_with_the_answer_kept() {
1355        let dir = tempdir().unwrap();
1356        let queue = Queue::at(dir.path().join("queue"));
1357        let questions = Questions::at(dir.path().join("questions"));
1358        let mut t = task("waits on an answer");
1359        queue.put(&mut t).unwrap();
1360
1361        apply(
1362            &queue,
1363            &questions,
1364            &Verdict {
1365                decisions: vec![Decision {
1366                    id: t.id.clone(),
1367                    blocked_by: Vec::new(),
1368                    reason: None,
1369                    recovery: None,
1370                    question: Some("Which backend?".to_owned()),
1371                    choices: Vec::new(),
1372                }],
1373            },
1374        )
1375        .unwrap();
1376        let blocked = queue.get(&t.id).unwrap();
1377        let question_id = blocked.blocked_by[0].clone();
1378
1379        let mut q = questions.get(&question_id).unwrap();
1380        q.answer(Answer::Text("SQLite".to_owned())).unwrap();
1381        questions.put(&mut q).unwrap();
1382        assert_eq!(q.status, QuestionStatus::Answered);
1383
1384        // `crate::daemon::resolve_blockers` is the deterministic resolver
1385        // that actually does this on the real queue; here it is enough to
1386        // prove the pure steps it is built from behave together.
1387        let mut task_after = queue.get(&t.id).unwrap();
1388        task_after.record_answer(q.summary.clone(), "SQLite".to_owned());
1389        task_after.unblock(&question_id);
1390        assert_eq!(task_after.status, TaskStatus::Queued);
1391        assert_eq!(task_after.answers[0].answer, "SQLite");
1392    }
1393
1394    #[test]
1395    fn a_stalled_task_can_be_requeued_or_held() {
1396        let dir = tempdir().unwrap();
1397        let queue = Queue::at(dir.path().join("queue"));
1398        let questions = Questions::at(dir.path().join("questions"));
1399
1400        let mut requeue_me = task("stuck a");
1401        requeue_me.start("run-1".to_owned());
1402        queue.put(&mut requeue_me).unwrap();
1403
1404        let mut hold_me = task("stuck b");
1405        hold_me.start("run-2".to_owned());
1406        queue.put(&mut hold_me).unwrap();
1407
1408        apply(
1409            &queue,
1410            &questions,
1411            &Verdict {
1412                decisions: vec![
1413                    Decision {
1414                        id: requeue_me.id.clone(),
1415                        recovery: Some(Recovery::Requeue),
1416                        ..Decision::default()
1417                    },
1418                    Decision {
1419                        id: hold_me.id.clone(),
1420                        recovery: Some(Recovery::Hold),
1421                        reason: Some("looks broken".to_owned()),
1422                        ..Decision::default()
1423                    },
1424                ],
1425            },
1426        )
1427        .unwrap();
1428
1429        let requeued = queue.get(&requeue_me.id).unwrap();
1430        assert_eq!(requeued.status, TaskStatus::Queued);
1431        assert_eq!(requeued.attempts, 0);
1432
1433        let held = queue.get(&hold_me.id).unwrap();
1434        assert_eq!(held.status, TaskStatus::Held);
1435        assert_eq!(held.hold_reason.as_deref(), Some("looks broken"));
1436    }
1437
1438    #[test]
1439    fn a_machine_held_task_asked_about_restores_to_held_once_answered() {
1440        // Reproduces the reported bug: a task held after burning its
1441        // attempts, once the conductor asks a follow-up question about it,
1442        // must come back `held` - not `queued` - once the question is
1443        // answered, whatever the answer said.
1444        let dir = tempdir().unwrap();
1445        let queue = Queue::at(dir.path().join("queue"));
1446        let questions = Questions::at(dir.path().join("questions"));
1447        let mut t = task("held out of attempts");
1448        t.hold_machine(Some("out of attempts".to_owned()));
1449        queue.put(&mut t).unwrap();
1450
1451        apply(
1452            &queue,
1453            &questions,
1454            &Verdict {
1455                decisions: vec![Decision {
1456                    id: t.id.clone(),
1457                    reason: Some("what should happen to this one?".to_owned()),
1458                    question: Some("Hold it, or try again?".to_owned()),
1459                    ..Decision::default()
1460                }],
1461            },
1462        )
1463        .unwrap();
1464        let blocked = queue.get(&t.id).unwrap();
1465        assert_eq!(blocked.status, TaskStatus::Blocked);
1466        let question_id = blocked.blocked_by[0].clone();
1467
1468        let mut q = questions.get(&question_id).unwrap();
1469        q.answer(Answer::Text("leave it held".to_owned())).unwrap();
1470        questions.put(&mut q).unwrap();
1471
1472        // What `crate::daemon::resolve_blockers` does on the real queue.
1473        let mut after = queue.get(&t.id).unwrap();
1474        after.record_answer(q.summary.clone(), "leave it held".to_owned());
1475        after.unblock(&question_id);
1476        assert_eq!(
1477            after.status,
1478            TaskStatus::Held,
1479            "must not fall back to queued"
1480        );
1481        assert_eq!(after.hold_reason.as_deref(), Some("out of attempts"));
1482    }
1483
1484    #[test]
1485    fn a_reaffirmed_hold_with_no_new_reason_is_not_silently_auto_released_by_triage() {
1486        // A different path to the same bug round 1 fixed: a task
1487        // machine-held for disk pressure, blocked on a conductor question,
1488        // answered "keep it held", and restored to `held`. If the
1489        // conductor's next decision reconfirms `recovery: hold` with no new
1490        // `reason` - allowed, `Decision::reason` is optional - `hold_machine`
1491        // must not silently leave the stale disk-pressure text in place, or
1492        // `crate::triage::run_once` reads it as an unexamined, auto-resolvable
1493        // hold and releases the task straight through the operator's answer
1494        // the instant disk space looks fine again.
1495        let dir = tempdir().unwrap();
1496        let queue = Queue::at(dir.path().join("queue"));
1497        let questions = Questions::at(dir.path().join("questions"));
1498        let mut t = task("disk pressure, then reconsidered");
1499        t.hold_machine(Some(
1500            "not enough free space to start a run: 10 bytes free, 100 required by \
1501             `[disk] min_free_bytes`"
1502                .to_owned(),
1503        ));
1504        t.record_answer(
1505            "How should this be handled?".to_owned(),
1506            "keep it held, a human will look at it later".to_owned(),
1507        );
1508        queue.put(&mut t).unwrap();
1509
1510        apply(
1511            &queue,
1512            &questions,
1513            &Verdict {
1514                decisions: vec![Decision {
1515                    id: t.id.clone(),
1516                    recovery: Some(Recovery::Hold),
1517                    ..Decision::default()
1518                }],
1519            },
1520        )
1521        .unwrap();
1522
1523        let after = queue.get(&t.id).unwrap();
1524        assert_eq!(after.status, TaskStatus::Held);
1525        assert!(
1526            !after
1527                .hold_reason
1528                .as_deref()
1529                .unwrap_or_default()
1530                .starts_with("not enough free space"),
1531            "the stale disk-pressure text must not survive a reconfirmed hold: {:?}",
1532            after.hold_reason
1533        );
1534
1535        // The disk gate would report space is fine now - `crate::triage`
1536        // must not read the old text and release the task through it.
1537        let cfg_dir = tempdir().unwrap();
1538        let config = cfg_dir.path().join("magi.toml");
1539        std::fs::write(&config, "[disk]\nmin_free_bytes = 0\n").unwrap();
1540        let report =
1541            crate::triage::run_once(&queue, &questions, Some(&config), jiff::Timestamp::now());
1542        assert!(report.resumed.is_empty(), "must not be auto-released");
1543        assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Held);
1544    }
1545
1546    #[test]
1547    fn a_conductor_rehold_after_a_resume_answer_is_not_asked_again_identically() {
1548        let dir = tempdir().unwrap();
1549        let queue = Queue::at(dir.path().join("queue"));
1550        let questions = Questions::at(dir.path().join("questions"));
1551        let config = dir.path().join("magi.toml");
1552        std::fs::write(&config, "[disk]\nmin_free_bytes = 0\n").unwrap();
1553        let now = jiff::Timestamp::now();
1554        let triage = || crate::triage::run_once(&queue, &questions, Some(&config), now);
1555        let open = || {
1556            questions
1557                .list()
1558                .into_iter()
1559                .filter(|q| q.node == "triage" && q.status.open())
1560                .collect::<Vec<_>>()
1561        };
1562        let hold = Decision {
1563            recovery: Some(Recovery::Hold),
1564            reason: Some("waiting on manual worktree cleanup".to_owned()),
1565            ..Decision::default()
1566        };
1567
1568        let mut t = task("looping hold");
1569        t.hold_machine(Some("waiting on manual worktree cleanup".to_owned()));
1570        queue.put(&mut t).unwrap();
1571        let hold = Decision {
1572            id: t.id.clone(),
1573            ..hold
1574        };
1575
1576        // 1. triage asks the ordinary machine question; the operator resumes.
1577        assert_eq!(triage().asked.len(), 1);
1578        let first = open().remove(0);
1579        let mut q = questions.get(&first.id).unwrap();
1580        let resume = q.choices[0].clone();
1581        q.answer(Answer::Choice(resume)).unwrap();
1582        questions.put(&mut q).unwrap();
1583        assert_eq!(triage().answered.len(), 1);
1584        assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Queued);
1585
1586        // 2. the conductor holds it again - allowed once, and recorded.
1587        apply(
1588            &queue,
1589            &questions,
1590            &Verdict {
1591                decisions: vec![hold.clone()],
1592            },
1593        )
1594        .unwrap();
1595        assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Held);
1596
1597        // 3. triage must say something different, once.
1598        assert_eq!(triage().asked.len(), 1);
1599        let second = open().remove(0);
1600        assert_ne!(second.summary, first.summary);
1601        assert_ne!(second.choices, first.choices);
1602        assert!(second.detail.contains("waiting on manual worktree cleanup"));
1603        assert!(triage().asked.is_empty(), "no duplicate question");
1604        assert_eq!(open().len(), 1);
1605
1606        // 4. forcing a requeue makes the conductor's hold a no-op.
1607        let mut q = questions.get(&second.id).unwrap();
1608        let force = q.choices[0].clone();
1609        q.answer(Answer::Choice(force)).unwrap();
1610        questions.put(&mut q).unwrap();
1611        assert_eq!(triage().answered.len(), 1);
1612        apply(
1613            &queue,
1614            &questions,
1615            &Verdict {
1616                decisions: vec![hold],
1617            },
1618        )
1619        .unwrap();
1620        assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Queued);
1621    }
1622
1623    #[test]
1624    fn a_runnable_task_can_be_held_directly_without_a_question() {
1625        let dir = tempdir().unwrap();
1626        let queue = Queue::at(dir.path().join("queue"));
1627        let questions = Questions::at(dir.path().join("questions"));
1628        let mut t = task("already answered, should stay put");
1629        queue.put(&mut t).unwrap();
1630
1631        apply(
1632            &queue,
1633            &questions,
1634            &Verdict {
1635                decisions: vec![Decision {
1636                    id: t.id.clone(),
1637                    recovery: Some(Recovery::Hold),
1638                    reason: Some("operator already said keep this held".to_owned()),
1639                    ..Decision::default()
1640                }],
1641            },
1642        )
1643        .unwrap();
1644
1645        let after = queue.get(&t.id).unwrap();
1646        assert_eq!(after.status, TaskStatus::Held);
1647        assert_eq!(after.hold_source, Some(crate::queue::HoldSource::Machine));
1648    }
1649
1650    #[test]
1651    fn done_recovery_closes_a_held_task_whose_goal_is_already_met() {
1652        // Reproduces the other half of the reported bug (task 6081): once
1653        // the operator's answer says the work already happened outside the
1654        // loop - PR merged, worktree cleaned up - the conductor needs an
1655        // actual terminal state to put the task in, not just a hold it will
1656        // keep being re-asked about.
1657        let dir = tempdir().unwrap();
1658        let queue = Queue::at(dir.path().join("queue"));
1659        let questions = Questions::at(dir.path().join("questions"));
1660        let mut t = task("already merged by hand");
1661        t.hold_machine(Some("branch survived, awaiting a decision".to_owned()));
1662        t.record_answer(
1663            "Handle this one?".to_owned(),
1664            "already merged and cleaned up, close it".to_owned(),
1665        );
1666        queue.put(&mut t).unwrap();
1667
1668        apply(
1669            &queue,
1670            &questions,
1671            &Verdict {
1672                decisions: vec![Decision {
1673                    id: t.id.clone(),
1674                    recovery: Some(Recovery::Done),
1675                    reason: Some("operator confirmed this already landed".to_owned()),
1676                    ..Decision::default()
1677                }],
1678            },
1679        )
1680        .unwrap();
1681
1682        let after = queue.get(&t.id).unwrap();
1683        assert_eq!(after.status, TaskStatus::Done);
1684        assert!(after.hold_reason.is_none());
1685        assert_eq!(after.answers.len(), 1, "the record of why is kept");
1686    }
1687
1688    #[test]
1689    fn done_recovery_is_ignored_for_a_runnable_or_running_task() {
1690        let dir = tempdir().unwrap();
1691        let queue = Queue::at(dir.path().join("queue"));
1692        let questions = Questions::at(dir.path().join("questions"));
1693
1694        let mut queued = task("never ran yet");
1695        queue.put(&mut queued).unwrap();
1696
1697        let mut running = task("mid-run");
1698        running.start("run-1".to_owned());
1699        queue.put(&mut running).unwrap();
1700
1701        for id in [queued.id.clone(), running.id.clone()] {
1702            apply(
1703                &queue,
1704                &questions,
1705                &Verdict {
1706                    decisions: vec![Decision {
1707                        id,
1708                        recovery: Some(Recovery::Done),
1709                        ..Decision::default()
1710                    }],
1711                },
1712            )
1713            .unwrap();
1714        }
1715
1716        assert_eq!(queue.get(&queued.id).unwrap().status, TaskStatus::Queued);
1717        assert_eq!(queue.get(&running.id).unwrap().status, TaskStatus::Running);
1718    }
1719
1720    #[test]
1721    fn a_question_after_two_settled_answers_is_still_filed_and_blocks() {
1722        // A genuine question must stay answerable however many were settled
1723        // before it: refusing it left `magi answer --list` empty.
1724        let dir = tempdir().unwrap();
1725        let queue = Queue::at(dir.path().join("queue"));
1726        let questions = Questions::at(dir.path().join("questions"));
1727        let mut t = task("asked about repeatedly");
1728        t.hold_machine(Some("out of attempts".to_owned()));
1729        t.record_answer("Handle this one? (1)".to_owned(), "not yet".to_owned());
1730        t.record_answer(
1731            "Handle this one? (2)".to_owned(),
1732            "still not yet".to_owned(),
1733        );
1734        queue.put(&mut t).unwrap();
1735        assert_eq!(questions.list().len(), 0);
1736
1737        let verdict = Verdict {
1738            decisions: vec![Decision {
1739                id: t.id.clone(),
1740                question: Some("Branch conflicts with origin/main, how do we proceed?".to_owned()),
1741                ..Decision::default()
1742            }],
1743        };
1744        apply(&queue, &questions, &verdict).unwrap();
1745
1746        let filed = questions.list();
1747        assert_eq!(filed.len(), 1, "the question was filed");
1748        assert_eq!(
1749            filed[0].summary,
1750            "Branch conflicts with origin/main, how do we proceed?"
1751        );
1752        let deputy = filed[0]
1753            .deputy
1754            .as_ref()
1755            .expect("the wait is handed to a deputy");
1756        assert!(deputy.brief.contains(&t.id), "{}", deputy.brief);
1757        assert_eq!(
1758            deputy.starts, 0,
1759            "filing never starts anything: the loop does not wait"
1760        );
1761        assert!(filed[0].status.open());
1762        assert_eq!(filed[0].node, NODE);
1763        let after = queue.get(&t.id).unwrap();
1764        assert_eq!(after.status, TaskStatus::Blocked);
1765        assert_eq!(after.blocked_by, vec![filed[0].id.clone()]);
1766        assert_eq!(after.answers.len(), 2, "the prior answers are untouched");
1767        assert!(
1768            !after
1769                .hold_reason
1770                .clone()
1771                .unwrap_or_default()
1772                .contains("conduct tried to ask"),
1773            "no hold was applied"
1774        );
1775
1776        apply(&queue, &questions, &verdict).unwrap();
1777        let again = questions.list();
1778        assert_eq!(again.len(), 1, "the open question is reused");
1779        assert_eq!(again[0].id, filed[0].id);
1780    }
1781
1782    #[test]
1783    fn a_second_conductor_question_is_still_allowed_after_one_settled_answer() {
1784        let dir = tempdir().unwrap();
1785        let queue = Queue::at(dir.path().join("queue"));
1786        let questions = Questions::at(dir.path().join("questions"));
1787        let mut t = task("asked about once already");
1788        t.hold_machine(Some("out of attempts".to_owned()));
1789        t.record_answer("Handle this one?".to_owned(), "not yet".to_owned());
1790        queue.put(&mut t).unwrap();
1791
1792        apply(
1793            &queue,
1794            &questions,
1795            &Verdict {
1796                decisions: vec![Decision {
1797                    id: t.id.clone(),
1798                    question: Some("Still not sure - now what?".to_owned()),
1799                    ..Decision::default()
1800                }],
1801            },
1802        )
1803        .unwrap();
1804
1805        assert_eq!(questions.list().len(), 1, "the second question was filed");
1806        assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Blocked);
1807    }
1808
1809    #[test]
1810    fn a_held_task_with_an_open_triage_question_is_left_to_triage() {
1811        // `crate::triage` never calls `Task::block`, and only ever looks at
1812        // tasks still `held` - so a conductor question that moved this task
1813        // to `blocked` would orphan triage's own open question. The guard in
1814        // `apply_one` must leave the task alone until triage's question
1815        // settles.
1816        let dir = tempdir().unwrap();
1817        let queue = Queue::at(dir.path().join("queue"));
1818        let questions = Questions::at(dir.path().join("questions"));
1819        let mut t = task("held, triage already asking about it");
1820        t.hold_machine(Some("cause unclear".to_owned()));
1821        queue.put(&mut t).unwrap();
1822
1823        let mut triage_q = Question::new(
1824            t.id.clone(),
1825            crate::triage::NODE.to_owned(),
1826            "triage".to_owned(),
1827            "Still needed?".to_owned(),
1828            String::new(),
1829            vec![
1830                "resume".to_owned(),
1831                "not yet".to_owned(),
1832                "discard".to_owned(),
1833            ],
1834        );
1835        questions.put(&mut triage_q).unwrap();
1836
1837        for decision in [
1838            Decision {
1839                id: t.id.clone(),
1840                question: Some("what now?".to_owned()),
1841                ..Decision::default()
1842            },
1843            Decision {
1844                id: t.id.clone(),
1845                recovery: Some(Recovery::Requeue),
1846                ..Decision::default()
1847            },
1848        ] {
1849            apply(
1850                &queue,
1851                &questions,
1852                &Verdict {
1853                    decisions: vec![decision],
1854                },
1855            )
1856            .unwrap();
1857        }
1858
1859        let after = queue.get(&t.id).unwrap();
1860        assert_eq!(
1861            after.status,
1862            TaskStatus::Held,
1863            "triage still owns this hold"
1864        );
1865        assert!(after.blocked_by.is_empty());
1866        assert_eq!(
1867            questions.list().len(),
1868            1,
1869            "no second, conductor-owned question was filed"
1870        );
1871    }
1872
1873    #[test]
1874    fn a_held_task_with_an_answered_but_unapplied_triage_question_is_still_left_alone() {
1875        // The race `crate::triage::pending_for` exists to close: the
1876        // operator has already answered the triage question (it is no
1877        // longer `open`), but `crate::triage::run_once` - which only runs on
1878        // a fully idle daemon tick, far less often than the conductor polls
1879        // - has not had a turn to apply it yet. `apply_one` must not treat
1880        // "not open" as "settled" here, or it would block the task out from
1881        // under an answer triage has not read back yet.
1882        let dir = tempdir().unwrap();
1883        let queue = Queue::at(dir.path().join("queue"));
1884        let questions = Questions::at(dir.path().join("questions"));
1885        let mut t = task("held, triage question answered but not yet applied");
1886        t.hold_machine(Some("cause unclear".to_owned()));
1887        queue.put(&mut t).unwrap();
1888
1889        let mut triage_q = Question::new(
1890            t.id.clone(),
1891            crate::triage::NODE.to_owned(),
1892            "triage".to_owned(),
1893            "Still needed?".to_owned(),
1894            String::new(),
1895            vec![
1896                "resume".to_owned(),
1897                "not yet".to_owned(),
1898                "discard".to_owned(),
1899            ],
1900        );
1901        questions.put(&mut triage_q).unwrap();
1902        triage_q
1903            .answer(Answer::Choice("not yet".to_owned()))
1904            .unwrap();
1905        questions.put(&mut triage_q).unwrap();
1906        assert!(!triage_q.status.open());
1907
1908        apply(
1909            &queue,
1910            &questions,
1911            &Verdict {
1912                decisions: vec![Decision {
1913                    id: t.id.clone(),
1914                    question: Some("what now?".to_owned()),
1915                    ..Decision::default()
1916                }],
1917            },
1918        )
1919        .unwrap();
1920
1921        let after = queue.get(&t.id).unwrap();
1922        assert_eq!(
1923            after.status,
1924            TaskStatus::Held,
1925            "triage's own answer is not yet applied - conduct must wait"
1926        );
1927        assert_eq!(
1928            questions.list().len(),
1929            1,
1930            "no conductor question was filed over the pending triage answer"
1931        );
1932    }
1933
1934    #[test]
1935    fn manual_hold_rejects_hostile_or_stale_conductor_recovery() {
1936        let dir = tempdir().unwrap();
1937        let queue = Queue::at(dir.path().join("queue"));
1938        let questions = Questions::at(dir.path().join("questions"));
1939        let mut held = task("manual recovery");
1940        held.priority = 300;
1941        held.runs.push("run20260912-224242-daf5".to_owned());
1942        held.hold_manual(Some(
1943            "active manual recovery run20260912-224242-daf5".to_owned(),
1944        ));
1945        queue.put(&mut held).unwrap();
1946
1947        // Every field a conductor may use to alter lifecycle state is ignored:
1948        // requeue/review would dispatch duplicate work, hold could overwrite
1949        // evidence, and a question would turn the hold into `blocked`.
1950        for decision in [
1951            Decision {
1952                id: held.id.clone(),
1953                recovery: Some(Recovery::Requeue),
1954                ..Decision::default()
1955            },
1956            Decision {
1957                id: held.id.clone(),
1958                recovery: Some(Recovery::Hold),
1959                reason: Some("stale replacement reason".to_owned()),
1960                ..Decision::default()
1961            },
1962            Decision {
1963                id: held.id.clone(),
1964                recovery: Some(Recovery::Review),
1965                ..Decision::default()
1966            },
1967            Decision {
1968                id: held.id.clone(),
1969                blocked_by: vec!["other-task".to_owned()],
1970                question: Some("retry now?".to_owned()),
1971                ..Decision::default()
1972            },
1973        ] {
1974            apply(
1975                &queue,
1976                &questions,
1977                &Verdict {
1978                    decisions: vec![decision],
1979                },
1980            )
1981            .unwrap();
1982        }
1983
1984        let after = queue.get(&held.id).unwrap();
1985        assert_eq!(after.status, TaskStatus::Held);
1986        assert!(after.operator_held());
1987        assert_eq!(after.priority, 300);
1988        assert_eq!(after.runs, ["run20260912-224242-daf5"]);
1989        assert_eq!(
1990            after.hold_reason.as_deref(),
1991            Some("active manual recovery run20260912-224242-daf5")
1992        );
1993        assert!(after.blocked_by.is_empty());
1994        assert!(questions.list().is_empty());
1995        assert!(
1996            queue.next_runnable().is_none(),
1997            "must not dispatch a duplicate"
1998        );
1999    }
2000
2001    #[test]
2002    fn machine_holds_remain_recoverable_and_manual_release_is_authorization() {
2003        let dir = tempdir().unwrap();
2004        let queue = Queue::at(dir.path().join("queue"));
2005        let questions = Questions::at(dir.path().join("questions"));
2006
2007        let mut automatic = task("disk gate");
2008        automatic.hold_machine(Some("disk full".to_owned()));
2009        queue.put(&mut automatic).unwrap();
2010        let requeue = || Verdict {
2011            decisions: vec![Decision {
2012                id: automatic.id.clone(),
2013                recovery: Some(Recovery::Requeue),
2014                ..Decision::default()
2015            }],
2016        };
2017        apply(&queue, &questions, &requeue()).unwrap();
2018        assert_eq!(queue.get(&automatic.id).unwrap().status, TaskStatus::Queued);
2019
2020        let mut manual = task("operator gate");
2021        manual.hold_manual(Some("wait for operator".to_owned()));
2022        queue.put(&mut manual).unwrap();
2023        apply(
2024            &queue,
2025            &questions,
2026            &Verdict {
2027                decisions: vec![Decision {
2028                    id: manual.id.clone(),
2029                    recovery: Some(Recovery::Requeue),
2030                    ..Decision::default()
2031                }],
2032            },
2033        )
2034        .unwrap();
2035        assert_eq!(queue.get(&manual.id).unwrap().status, TaskStatus::Held);
2036
2037        // This mirrors the CLI and web release routes: only an explicit
2038        // operator action clears the manual boundary.
2039        let mut released = queue.get(&manual.id).unwrap();
2040        released.release();
2041        queue.put(&mut released).unwrap();
2042        assert_eq!(queue.get(&manual.id).unwrap().status, TaskStatus::Queued);
2043    }
2044
2045    #[test]
2046    fn legacy_reasoned_hold_is_protected_without_losing_its_metadata() {
2047        let dir = tempdir().unwrap();
2048        let queue = Queue::at(dir.path().join("queue"));
2049        let questions = Questions::at(dir.path().join("questions"));
2050        let mut legacy = task("old explicit hold");
2051        legacy.status = TaskStatus::Held;
2052        legacy.hold_reason = Some("manual recovery already active".to_owned());
2053        legacy.hold_source = None;
2054        legacy.blocked_by = vec!["dependency".to_owned()];
2055        queue.put(&mut legacy).unwrap();
2056
2057        apply(
2058            &queue,
2059            &questions,
2060            &Verdict {
2061                decisions: vec![Decision {
2062                    id: legacy.id.clone(),
2063                    recovery: Some(Recovery::Requeue),
2064                    ..Decision::default()
2065                }],
2066            },
2067        )
2068        .unwrap();
2069
2070        let after = queue.get(&legacy.id).unwrap();
2071        assert_eq!(after.status, TaskStatus::Held);
2072        assert_eq!(after.hold_source, None);
2073        assert_eq!(after.hold_reason, legacy.hold_reason);
2074        assert_eq!(after.blocked_by, legacy.blocked_by);
2075    }
2076
2077    #[test]
2078    fn review_recovery_is_a_no_op_without_a_survivable_branch() {
2079        // `surviving_branch` reaches `RunState::load`, which reaches the
2080        // process-global `run::home()` - a `OnceLock`, so this only wins the
2081        // race the first time it runs in the binary; every other test still
2082        // reaches the same directory whichever call won, and this test's own
2083        // run id never collides with another test's.
2084        crate::run::set_home(std::env::temp_dir().join("magi-conduct-tests-home"));
2085        let dir = tempdir().unwrap();
2086        let queue = Queue::at(dir.path().join("queue"));
2087        let questions = Questions::at(dir.path().join("questions"));
2088        let mut t = task("blocked with no readable run");
2089        t.start("20260101-000000-dead".to_owned()); // no such run on disk
2090        t.fail("blocked", 5);
2091        queue.put(&mut t).unwrap();
2092
2093        apply(
2094            &queue,
2095            &questions,
2096            &Verdict {
2097                decisions: vec![Decision {
2098                    id: t.id.clone(),
2099                    recovery: Some(Recovery::Review),
2100                    ..Decision::default()
2101                }],
2102            },
2103        )
2104        .unwrap();
2105
2106        let after = queue.get(&t.id).unwrap();
2107        assert_eq!(
2108            after.status,
2109            TaskStatus::Failed,
2110            "with nothing to reopen, the decision is dropped rather than guessed at"
2111        );
2112        assert!(after.review_branch.is_none());
2113    }
2114
2115    #[test]
2116    fn requeue_and_review_recovery_are_ignored_for_a_runnable_task() {
2117        // `Hold` is the one exception - see `Decision::recovery`'s doc and
2118        // `a_runnable_task_can_be_held_directly_without_a_question` - but a
2119        // task already in line has nothing for `requeue` or `review` to do.
2120        let dir = tempdir().unwrap();
2121        let queue = Queue::at(dir.path().join("queue"));
2122        let questions = Questions::at(dir.path().join("questions"));
2123
2124        for recovery in [Recovery::Requeue, Recovery::Review] {
2125            let mut t = task("ordinary");
2126            queue.put(&mut t).unwrap();
2127
2128            apply(
2129                &queue,
2130                &questions,
2131                &Verdict {
2132                    decisions: vec![Decision {
2133                        id: t.id.clone(),
2134                        recovery: Some(recovery),
2135                        ..Decision::default()
2136                    }],
2137                },
2138            )
2139            .unwrap();
2140
2141            assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Queued);
2142        }
2143    }
2144
2145    #[tokio::test]
2146    async fn a_broken_agent_leaves_the_queue_untouched_and_does_not_error() {
2147        let dir = tempdir().unwrap();
2148        let cfg = config(mock_agent(dir.path(), BROKEN, BTreeMap::new()));
2149        let queue = Queue::at(dir.path().join("queue"));
2150        let questions = Questions::at(dir.path().join("questions"));
2151        let mut t = task("normal");
2152        queue.put(&mut t).unwrap();
2153
2154        let mut conductor = Conductor::new();
2155        conductor
2156            .maybe_run(
2157                &cfg,
2158                dir.path(),
2159                &queue,
2160                &questions,
2161                dir.path(),
2162                &[t.clone()],
2163                &[],
2164                &[],
2165                2,
2166            )
2167            .await;
2168
2169        assert_eq!(
2170            queue.get(&t.id).unwrap().status,
2171            TaskStatus::Queued,
2172            "a failed invocation must change nothing"
2173        );
2174        assert!(
2175            queue.next_runnable().is_some(),
2176            "the loop must still be able to take the next task"
2177        );
2178    }
2179
2180    #[tokio::test]
2181    async fn a_reply_with_no_json_leaves_the_queue_untouched() {
2182        let dir = tempdir().unwrap();
2183        let cfg = config(mock_agent(dir.path(), GARBAGE, BTreeMap::new()));
2184        let queue = Queue::at(dir.path().join("queue"));
2185        let questions = Questions::at(dir.path().join("questions"));
2186        let mut t = task("normal");
2187        queue.put(&mut t).unwrap();
2188
2189        let mut conductor = Conductor::new();
2190        conductor
2191            .maybe_run(
2192                &cfg,
2193                dir.path(),
2194                &queue,
2195                &questions,
2196                dir.path(),
2197                &[t.clone()],
2198                &[],
2199                &[],
2200                2,
2201            )
2202            .await;
2203
2204        assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Queued);
2205    }
2206
2207    #[tokio::test]
2208    async fn a_conductor_chain_falls_to_the_next_agent_once_each() {
2209        let dir = tempdir().unwrap();
2210        let mut t = task("chained");
2211        let reply = format!(
2212            "{{\"decisions\":[{{\"id\":\"{}\",\"blocked_by\":[\"x\"],\"reason\":\"why\"}}]}}",
2213            t.id
2214        );
2215        let counted = |id: &str, tail: &str| {
2216            let calls = dir.path().join(format!("{id}.calls"));
2217            let path = dir.path().join(format!("{id}.sh"));
2218            std::fs::write(
2219                &path,
2220                format!(
2221                    "#!/bin/sh\ncat >/dev/null\necho x >> '{}'\n{tail}\n",
2222                    calls.display()
2223                ),
2224            )
2225            .unwrap();
2226            AgentSpec {
2227                id: id.to_owned(),
2228                kind: AgentKind::Command,
2229                model: None,
2230                command: vec!["sh".to_owned(), path.to_string_lossy().into_owned()],
2231                extra_args: Vec::new(),
2232                env: BTreeMap::new(),
2233                prompt_delivery: None,
2234            }
2235        };
2236        let n_calls = |id: &str| {
2237            std::fs::read_to_string(dir.path().join(format!("{id}.calls")))
2238                .map_or(0, |s| s.lines().count())
2239        };
2240        let mut cfg = config(counted("a", "exit 3"));
2241        cfg.agents = vec![
2242            counted("a", "exit 3"),
2243            counted("b", &format!("printf '%s' '{reply}'")),
2244        ];
2245        cfg.roles.conductor = Some(crate::config::AgentChoice::Chain(vec![
2246            "a".into(),
2247            "b".into(),
2248            "a".into(),
2249        ]));
2250        let queue = Queue::at(dir.path().join("queue"));
2251        let questions = Questions::at(dir.path().join("questions"));
2252        queue.put(&mut t).unwrap();
2253
2254        let mut conductor = Conductor::new();
2255        conductor
2256            .maybe_run(
2257                &cfg,
2258                dir.path(),
2259                &queue,
2260                &questions,
2261                dir.path(),
2262                &[t.clone()],
2263                &[],
2264                &[],
2265                2,
2266            )
2267            .await;
2268
2269        assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Blocked);
2270        assert_eq!(n_calls("a"), 1, "each id is tried once");
2271        assert_eq!(n_calls("b"), 1);
2272        let saved = std::fs::read_to_string(seat_path(dir.path())).unwrap();
2273        assert!(
2274            saved.contains("\"b\""),
2275            "the seat on disk is the one that ran: {saved}"
2276        );
2277    }
2278
2279    #[tokio::test]
2280    async fn json_survives_code_fences_and_a_preamble() {
2281        let dir = tempdir().unwrap();
2282        let mut t = task("fenced");
2283        let reply = format!(
2284            "Sure, here is my decision.\n\n```json\n{{\"decisions\":[{{\"id\":\"{}\",\
2285             \"blocked_by\":[\"x\"],\"reason\":\"why\"}}]}}\n```\n",
2286            t.id
2287        );
2288        let cfg = config(mock_agent(dir.path(), REPLY, env(&reply)));
2289        let queue = Queue::at(dir.path().join("queue"));
2290        let questions = Questions::at(dir.path().join("questions"));
2291        queue.put(&mut t).unwrap();
2292
2293        let mut conductor = Conductor::new();
2294        conductor
2295            .maybe_run(
2296                &cfg,
2297                dir.path(),
2298                &queue,
2299                &questions,
2300                dir.path(),
2301                &[t.clone()],
2302                &[],
2303                &[],
2304                2,
2305            )
2306            .await;
2307
2308        let back = queue.get(&t.id).unwrap();
2309        assert_eq!(back.status, TaskStatus::Blocked);
2310        assert_eq!(back.blocked_by, ["x"]);
2311    }
2312
2313    #[tokio::test]
2314    async fn the_conductor_is_not_called_again_when_nothing_worth_looking_at_has_changed() {
2315        // Each real invocation writes its own artifact stem, `turn-<n>`, so
2316        // whether a second one happened is read off the artifacts directory.
2317        let dir = tempdir().unwrap();
2318        let cfg = config(mock_agent(dir.path(), REPLY, env("{\"decisions\":[]}")));
2319        let queue = Queue::at(dir.path().join("queue"));
2320        let questions = Questions::at(dir.path().join("questions"));
2321        let mut t = task("stable");
2322        queue.put(&mut t).unwrap();
2323        let artifacts = dir.path().join("conduct").join("artifacts");
2324        let turn = |n: usize| artifacts.join(format!("turn-{n}.out"));
2325
2326        let mut conductor = Conductor::new();
2327        conductor
2328            .maybe_run(
2329                &cfg,
2330                dir.path(),
2331                &queue,
2332                &questions,
2333                dir.path(),
2334                &[t.clone()],
2335                &[],
2336                &[],
2337                2,
2338            )
2339            .await;
2340        assert!(turn(1).is_file(), "the first cycle must call the conductor");
2341
2342        conductor
2343            .maybe_run(
2344                &cfg,
2345                dir.path(),
2346                &queue,
2347                &questions,
2348                dir.path(),
2349                &[t.clone()],
2350                &[],
2351                &[],
2352                2,
2353            )
2354            .await;
2355        assert!(
2356            !turn(2).is_file(),
2357            "an unchanged revision and an unchanged stalled/finished set must not call the \
2358             conductor twice"
2359        );
2360
2361        // Once the queue actually changes, the next `maybe_run` calls again.
2362        t.priority = 1;
2363        queue.put(&mut t).unwrap();
2364        conductor
2365            .maybe_run(
2366                &cfg,
2367                dir.path(),
2368                &queue,
2369                &questions,
2370                dir.path(),
2371                &[t.clone()],
2372                &[],
2373                &[],
2374                2,
2375            )
2376            .await;
2377        assert!(turn(2).is_file(), "a moved revision calls it again");
2378    }
2379
2380    #[tokio::test]
2381    async fn a_task_turning_stalled_calls_the_conductor_again_despite_an_unchanged_revision() {
2382        // The queue's own revision has not moved - nothing wrote to it - but
2383        // a task now looks stalled, purely because time passed. Calling
2384        // again here, and never again once this exact set has been shown
2385        // once, is the whole point of keying `worth_a_look` on the id set
2386        // rather than on "is it non-empty".
2387        let dir = tempdir().unwrap();
2388        let cfg = config(mock_agent(dir.path(), REPLY, env("{\"decisions\":[]}")));
2389        let queue = Queue::at(dir.path().join("queue"));
2390        let questions = Questions::at(dir.path().join("questions"));
2391        let mut t = task("quiet");
2392        queue.put(&mut t).unwrap();
2393        let artifacts = dir.path().join("conduct").join("artifacts");
2394        let turn = |n: usize| artifacts.join(format!("turn-{n}.out"));
2395
2396        let mut conductor = Conductor::new();
2397        conductor
2398            .maybe_run(
2399                &cfg,
2400                dir.path(),
2401                &queue,
2402                &questions,
2403                dir.path(),
2404                &[t.clone()],
2405                &[],
2406                &[],
2407                2,
2408            )
2409            .await;
2410        assert!(turn(1).is_file());
2411
2412        conductor
2413            .maybe_run(
2414                &cfg,
2415                dir.path(),
2416                &queue,
2417                &questions,
2418                dir.path(),
2419                &[],
2420                &[t.clone()],
2421                &[],
2422                2,
2423            )
2424            .await;
2425        assert!(
2426            turn(2).is_file(),
2427            "a task turning stalled must call the conductor again"
2428        );
2429
2430        // But once shown at this exact revision, showing the *same* stalled
2431        // set again must not call a third time.
2432        conductor
2433            .maybe_run(
2434                &cfg,
2435                dir.path(),
2436                &queue,
2437                &questions,
2438                dir.path(),
2439                &[],
2440                &[t.clone()],
2441                &[],
2442                2,
2443            )
2444            .await;
2445        assert!(
2446            !turn(3).is_file(),
2447            "the same stalled task lingering must not call the conductor every cycle"
2448        );
2449    }
2450
2451    #[test]
2452    fn worth_a_look_is_config_free_and_matches_maybe_runs_own_gate() {
2453        let dir = tempdir().unwrap();
2454        let queue = Queue::at(dir.path().join("queue"));
2455        let mut t = task("t");
2456        queue.put(&mut t).unwrap();
2457
2458        let mut conductor = Conductor::new();
2459        assert!(
2460            conductor.worth_a_look(&queue, &[], &[]),
2461            "a conductor that has never run has something to look at"
2462        );
2463
2464        conductor.last_seen = Some(Conductor::snapshot(&queue, &[], &[]));
2465        assert!(
2466            !conductor.worth_a_look(&queue, &[], &[]),
2467            "nothing changed and nothing is stalled or finished"
2468        );
2469        assert!(
2470            conductor.worth_a_look(&queue, &[t.clone()], &[]),
2471            "a stalled task is worth a look even at the same revision"
2472        );
2473        assert!(
2474            conductor.worth_a_look(&queue, &[], &[t.clone()]),
2475            "a finished task is worth a look even at the same revision"
2476        );
2477    }
2478
2479    #[tokio::test]
2480    async fn the_conduct_path_never_calls_ask_and_wait() {
2481        // Structural: grepping this module and `daemon.rs` for
2482        // `ask_and_wait` is the actual assertion this module's own doc
2483        // promises; this test exists so the promise has a name in the test
2484        // output too. `apply_one`'s question path uses `Questions::put`
2485        // exclusively.
2486        let dir = tempdir().unwrap();
2487        let queue = Queue::at(dir.path().join("queue"));
2488        let questions = Questions::at(dir.path().join("questions"));
2489        let mut t = task("asks without blocking");
2490        queue.put(&mut t).unwrap();
2491
2492        apply(
2493            &queue,
2494            &questions,
2495            &Verdict {
2496                decisions: vec![Decision {
2497                    id: t.id.clone(),
2498                    question: Some("ok?".to_owned()),
2499                    ..Decision::default()
2500                }],
2501            },
2502        )
2503        .unwrap();
2504        // Reaching here at all (no hang) is the assertion.
2505        assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Blocked);
2506    }
2507
2508    #[test]
2509    fn a_pinned_resume_is_not_held_or_requeued_by_the_conductor() {
2510        let mut t = Task::new("t".into(), "t".into(), PathBuf::new(), Source::Human);
2511        t.resume_override = Some(crate::queue::OperatorResume {
2512            question_id: "q".into(),
2513            at: jiff::Timestamp::now(),
2514            conductor_rehold: None,
2515            forced: false,
2516            pinned_run: Some("run-1".into()),
2517        });
2518        assert!(pinned_resume(&t));
2519        assert!(!may_hold(&mut t, "waiting for magi resume"));
2520        assert!(
2521            t.resume_override
2522                .as_ref()
2523                .unwrap()
2524                .conductor_rehold
2525                .is_none(),
2526            "a refused hold is not recorded as an override"
2527        );
2528        t.resume_override = None;
2529        assert!(!pinned_resume(&t));
2530        assert!(may_hold(&mut t, "no override, so a hold is allowed"));
2531    }
2532
2533    #[test]
2534    fn a_pinned_resume_is_not_blocked_by_a_conductor_question_or_dependency() {
2535        let dir = tempdir().unwrap();
2536        let queue = Queue::at(dir.path().join("queue"));
2537        let questions = Questions::at(dir.path().join("questions"));
2538        let mut t = task("resume me");
2539        t.resume_override = Some(crate::queue::OperatorResume {
2540            question_id: "q".into(),
2541            at: jiff::Timestamp::now(),
2542            conductor_rehold: None,
2543            forced: true,
2544            pinned_run: Some("run-1".into()),
2545        });
2546        queue.put(&mut t).unwrap();
2547
2548        apply(
2549            &queue,
2550            &questions,
2551            &Verdict {
2552                decisions: vec![
2553                    Decision {
2554                        id: t.id.clone(),
2555                        question: Some("really?".to_owned()),
2556                        ..Decision::default()
2557                    },
2558                    Decision {
2559                        id: t.id.clone(),
2560                        blocked_by: vec!["other".to_owned()],
2561                        ..Decision::default()
2562                    },
2563                ],
2564            },
2565        )
2566        .unwrap();
2567
2568        assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Queued);
2569        assert!(questions.list().is_empty());
2570    }
2571}