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