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