Skip to main content

magi/
conduct.rs

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