Skip to main content

magi/
conduct.rs

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