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