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