Skip to main content

magi/
triage.rs

1//! Held-task triage: walking `held` tasks so a hold left by an accident does
2//! not sit unread forever next to one a human placed on purpose.
3//!
4//! [`crate::queue::HoldSource`] already distinguishes "the daemon or
5//! conductor held this during its own recovery" ([`HoldSource::Machine`],
6//! documented as recoverable) from "an operator held this on purpose"
7//! ([`HoldSource::Manual`]). What was missing was anything that actually acts
8//! on that distinction: nothing walked the held list and asked whether a
9//! machine hold's cause was still true, and a record written before
10//! `hold_source` existed (schema < 3, `None`) was silently protected forever
11//! by [`Task::operator_held`]'s conservative default - never wrong, but also
12//! never looked at again by anything.
13//!
14//! [`run_once`] is that walk. For every `held` task it finds:
15//!
16//! - [`HoldSource::Machine`]: if [`machine_cause_resolved`] can tell the
17//!   cause is gone, the task goes straight back to `queued` - the same effect
18//!   as `magi task release`, just automatic. When it cannot tell, a
19//!   [`Question`] is filed once and the task stays held until answered.
20//! - `None` (a legacy record, or a hold nobody explained): always a question,
21//!   exactly once - the whole point being that "protected forever" must not
22//!   mean "never shown to anyone" either.
23//! - [`HoldSource::Manual`]: never touched automatically. Only once the hold
24//!   has sat untouched past [`MANUAL_STALE_AFTER`] does it earn a question of
25//!   its own, asking whether it is still wanted.
26//!
27//! Every question this module files carries [`NODE`] and uses the task id as
28//! [`Question::run`] - the same convention `crate::conduct`'s own questions
29//! use for a task rather than a run (see `conduct::apply_one`'s own comment
30//! on why the dedupe check there also filters on `node`, not `run` alone: an
31//! ordinary graph question's `run` is a real run id, and a coincidental
32//! equality with some task's id must not be read as "about this task").
33//! [`latest_triage_question`] follows the identical rule.
34//!
35//! # Why answers apply here rather than through `crate::queue::Task::block`
36//!
37//! `crate::conduct` blocks a task on its own question
38//! (`Task::block(vec![question_id], …)`), and `crate::daemon::resolve_blockers`
39//! unblocks it - unconditionally, back to `queued` - the moment that question
40//! is answered, whatever the answer actually said. That is correct for
41//! `conduct`: the *content* of the answer is meant for whoever reads
42//! `Task::answers` next, not for the resolver.
43//!
44//! A triage question's answer is different: "not yet" and "discard it" do two
45//! entirely different, non-resuming things, and only "resume it" may put the
46//! task back in line. Reusing the generic blocked/unblock path would resume
47//! every answer alike, so this module never calls [`Task::block`] and never
48//! leaves a triaged task anything but `held` while its question is open.
49//! [`interpret_answer`] reads [`Question::resolution`] itself and
50//! [`run_once`] acts on it directly: [`Task::release`] for an actual "resume
51//! it" choice, [`Queue::remove`] for "discard it" (捨ててよい really means
52//! "you may throw this away", not "leave it sitting held" - the English
53//! wording must say the same thing, not "leave it held"), and
54//! [`Task::hold_manual`] for anything else - which both keeps the task held
55//! and reclassifies it as a hold an operator has now actually seen, one
56//! `crate::conduct` and a later triage pass leave alone.
57//!
58//! The choice is read by its **position** in [`Question::choices`]
59//! ([`Wording::choices3`]/[`Wording::choices2`] always put "resume" first and
60//! "discard" third), never by comparing the answer text against [`Wording`]'s
61//! own strings picked from whatever config is in force *now* - the language a
62//! question was filed in and the language a later `run_once` call happens to
63//! read back are not guaranteed to be the same call's [`Config`], and a text
64//! comparison would silently misread a real "resume" answer as "keep held"
65//! the moment they disagree.
66//!
67//! # Idempotency
68//!
69//! [`run_once`] runs on every daemon idle tick (see `crate::daemon::poll`)
70//! and on every `magi task triage`, so an answered question must be applied to
71//! a task **at most once**, and a *fresh* question for the same task (once it
72//! is held again, or goes stale) must still be possible. The record is
73//! [`Task::triage_applied`], the question ids already applied. It cannot live
74//! in [`Task::hold_reason`]: [`Task::release`] clears that, so a "resume"
75//! answer left no trace, and a released task that failed back to `held` was
76//! released again by the same old answer with its attempts reset - forever.
77//! A task that comes back to `held` after an applied answer is therefore a
78//! new hold, handled per [`HoldSource`] (a fresh question for a machine hold).
79//! [`already_applied`] also still reads the `[triage:<short>]` marker
80//! [`keep_held_note`] appends to the hold reason, for records that pre-date
81//! the field.
82
83use std::path::{Path, PathBuf};
84use std::time::Duration;
85
86use jiff::Timestamp;
87
88use crate::ask::{Question, QuestionStatus, Questions};
89use crate::config::Config;
90use crate::disk;
91use crate::queue::{HoldSource, OperatorResume, Queue, Task, TaskStatus};
92
93/// Node recorded on every question this module files - `crate::conduct::NODE`
94/// for the same idea applied to a `crate::conduct` decision instead.
95pub const NODE: &str = "triage";
96
97/// Node on the question filed about a stuck dependency root - see
98/// [`ask_about_stuck_roots`]. Separate from [`NODE`] because the choices mean
99/// something else (position 2 is "detach the dependants", not "discard").
100pub const DEPS_NODE: &str = "triage-deps";
101
102/// Seat name on a filed question. Not a real agent seat - there is no model
103/// call anywhere in this module - but every [`Question`] needs one, and every
104/// other deterministic filer (`crate::land`'s merge approval) names itself
105/// the same way.
106const SEAT: &str = "triage";
107
108/// How long a [`HoldSource::Manual`] hold sits untouched before triage asks
109/// whether it is still wanted.
110///
111/// A judgement call, not a `magi.toml` setting - the same reasoning
112/// `ask::REPLY_QUIET_WINDOW` documents for itself: there is no operator
113/// preference for "how long is too long to ignore my own hold" that a
114/// per-repository config could be *right* about. Seven days is long enough
115/// that an ordinary multi-day hold (waiting on a dependency, waiting on the
116/// operator's own schedule) never gets nagged, and short enough that a hold
117/// nobody has looked at in a week surfaces again rather than aging into the
118/// kind of silent backlog this feature exists to prevent.
119const MANUAL_STALE_AFTER: Duration = Duration::from_secs(7 * 24 * 60 * 60);
120
121/// Localised strings for a filed question, the same idea as `crate::land`'s
122/// own `Words`/`words()` - only the languages magi can actually check are
123/// translated, and anything else falls back to English.
124struct Wording {
125    lang: &'static str,
126    resume: &'static str,
127    wait: &'static str,
128    discard: &'static str,
129    resume_now: &'static str,
130    keep_held: &'static str,
131}
132
133const EN: Wording = Wording {
134    lang: "en",
135    resume: "resume it",
136    wait: "not yet",
137    discard: "discard it",
138    resume_now: "resume it",
139    keep_held: "keep it held",
140};
141
142const JA: Wording = Wording {
143    lang: "ja",
144    resume: "再開してよい",
145    wait: "まだ待って",
146    discard: "捨ててよい",
147    resume_now: "再開する",
148    keep_held: "まだ止めておく",
149};
150
151/// Pick the wording. Codes and names both, the same acceptance
152/// `crate::land::words` gives `[graph] language`.
153fn wording(language: &str) -> &'static Wording {
154    if crate::lang::is_japanese(language) {
155        &JA
156    } else {
157        &EN
158    }
159}
160
161impl Wording {
162    fn choices3(&self) -> Vec<String> {
163        vec![
164            self.resume.to_owned(),
165            self.wait.to_owned(),
166            self.discard.to_owned(),
167        ]
168    }
169
170    fn choices2(&self) -> Vec<String> {
171        vec![self.resume_now.to_owned(), self.keep_held.to_owned()]
172    }
173
174    fn source_label(&self, source: Option<HoldSource>) -> &'static str {
175        match (self.lang, source) {
176            ("ja", Some(HoldSource::Machine)) => "machine(機械による自動保留)",
177            ("ja", Some(HoldSource::Manual)) => "manual(操作者による手動保留)",
178            ("ja", None) => "unknown(schema 3 未満の旧レコード、または理由未記録)",
179            (_, Some(HoldSource::Machine)) => "machine (automatic recovery hold)",
180            (_, Some(HoldSource::Manual)) => "manual (an operator held this)",
181            (_, None) => "unknown (pre-schema-3 record, or never recorded)",
182        }
183    }
184
185    /// The body under the summary: everything an operator needs to judge this
186    /// without opening a terminal - id, title, hold reason, hold source.
187    ///
188    /// Falls back to [`Task::last_error`] when [`Task::hold_reason`] is empty:
189    /// a record written before `Task::fail`/`Task::handed_off` started
190    /// copying `why` into `hold_reason` too, or one written by a still older
191    /// build, would otherwise leave the question blank for exactly the case
192    /// requirement 4 exists for.
193    fn detail(&self, task: &Task, why: &str) -> String {
194        let none = if self.lang == "ja" {
195            "(記録なし)"
196        } else {
197            "(none recorded)"
198        };
199        let reason = task
200            .hold_reason
201            .as_deref()
202            .or(task.last_error.as_deref())
203            .unwrap_or(none);
204        format!(
205            "task: {} ({})\ntitle: {}\nhold source: {}\nhold reason: {reason}\n\n{why}",
206            task.id,
207            task.short(),
208            task.title,
209            self.source_label(task.hold_source),
210        )
211    }
212
213    fn summary_machine_unknown(&self, task: &Task) -> String {
214        if self.lang == "ja" {
215            format!("保留タスク {} の再開可否を判断してください", task.short())
216        } else {
217            format!("decide whether to resume held task {}", task.short())
218        }
219    }
220
221    fn why_machine(&self) -> &'static str {
222        if self.lang == "ja" {
223            "機械的な保留(machine hold)ですが、原因がすでに解消しているかを自動では判断できませんでした。"
224        } else {
225            "This is a machine hold, but whether its cause has resolved could not be \
226             checked automatically."
227        }
228    }
229
230    fn summary_conductor_override(&self, task: &Task) -> String {
231        if self.lang == "ja" {
232            format!(
233                "再開と回答済みのタスク {} を conductor が再び保留しました",
234                task.short()
235            )
236        } else {
237            format!(
238                "task {} was resumed at your word, but the conductor held it again",
239                task.short()
240            )
241        }
242    }
243
244    fn why_conductor_override(&self, o: &OperatorResume) -> String {
245        let reason = o.conductor_rehold.as_deref().unwrap_or_default();
246        if self.lang == "ja" {
247            format!(
248                "{} に再開と回答済みですが、conductor が再び hold しました。conductor の理由: \
249                 {reason}\n\n強制再キューを選ぶと、以後 conductor はこのタスクを hold できません。",
250                o.at
251            )
252        } else {
253            format!(
254                "You answered \"resume\" at {}, but the conductor held the task again. \
255                 Its reason: {reason}\n\nForcing a requeue stops the conductor from \
256                 holding this task again.",
257                o.at
258            )
259        }
260    }
261
262    /// Positions match [`AnswerAction`]: 0 resume, 1 keep held, 2 discard.
263    fn choices_conductor_override(&self) -> Vec<String> {
264        if self.lang == "ja" {
265            vec![
266                "強制再キュー(conductor は再 hold 不可)".to_owned(),
267                "手動 hold のまま".to_owned(),
268                "捨ててよい".to_owned(),
269            ]
270        } else {
271            vec![
272                "force requeue (conductor must not hold again)".to_owned(),
273                "keep held (manual)".to_owned(),
274                "discard".to_owned(),
275            ]
276        }
277    }
278
279    fn summary_legacy(&self, task: &Task) -> String {
280        if self.lang == "ja" {
281            format!(
282                "hold_source が不明な保留タスク {} を確認してください",
283                task.short()
284            )
285        } else {
286            format!(
287                "held task {} has no recorded hold source - please take a look",
288                task.short()
289            )
290        }
291    }
292
293    fn why_legacy(&self) -> &'static str {
294        if self.lang == "ja" {
295            "hold_source が記録されていません。schema 3 より前のレコードか、理由が記録されなかった \
296             holdです。人が意図して止めたのか、クラッシュや強制再起動で宙に浮いただけなのか、\
297             このデータからは区別できません。"
298        } else {
299            "No hold_source was recorded - either a pre-schema-3 record, or a hold whose \
300             reason was never written down. Whether this was a deliberate hold or the \
301             leftover of a crash cannot be told from the data alone."
302        }
303    }
304
305    fn summary_manual_stale(&self, task: &Task, days: i64) -> String {
306        if self.lang == "ja" {
307            format!(
308                "{days}日間 保留されたままの手動保留タスク {} を確認してください",
309                task.short()
310            )
311        } else {
312            format!(
313                "held task {} has been on a manual hold for {days} day(s)",
314                task.short()
315            )
316        }
317    }
318
319    fn why_manual(&self) -> &'static str {
320        if self.lang == "ja" {
321            "操作者が明示的に止めた保留ですが、長期間そのままになっています。まだ止めておくか、\
322             再開するか教えてください。"
323        } else {
324            "An operator held this on purpose, but it has sat untouched for a while. Say \
325             whether to keep holding it or resume it."
326        }
327    }
328}
329
330/// Which of the three situations this module recognises a held task is in.
331#[derive(Debug, Clone, Copy, PartialEq, Eq)]
332enum Bucket {
333    /// `HoldSource::Machine`, cause not verifiably resolved.
334    MachineUnknown,
335    /// `HoldSource::Machine`, held by the conductor after the operator
336    /// answered "resume" (see [`Task::resume_override`]).
337    ConductorOverride,
338    /// `hold_source` is `None`.
339    Legacy,
340    /// `HoldSource::Manual`, held past [`MANUAL_STALE_AFTER`].
341    ManualStale,
342}
343
344/// What one [`run_once`] pass did, task ids in each list.
345#[derive(Debug, Clone, Default)]
346pub struct Report {
347    /// A `HoldSource::Machine` hold whose cause was found resolved, put back
348    /// in line automatically.
349    pub resumed: Vec<String>,
350    /// A fresh question was filed this pass.
351    pub asked: Vec<String>,
352    /// An operator's answer to an earlier triage question was applied.
353    pub answered: Vec<String>,
354    /// A `blocked` task whose `blocked_by` named an id that no longer exists
355    /// was moved to a machine hold this pass - see
356    /// [`quarantine_orphaned_blocked`]. Distinct from `asked`: the question
357    /// about it, if any, is filed in the same pass and only counted there.
358    pub quarantined: Vec<String>,
359}
360
361impl Report {
362    /// Is there nothing to report? Callers use this to skip logging an empty
363    /// pass rather than repeating "triaged 0 held task(s)" on every idle tick.
364    pub fn is_empty(&self) -> bool {
365        self.resumed.is_empty()
366            && self.asked.is_empty()
367            && self.answered.is_empty()
368            && self.quarantined.is_empty()
369    }
370}
371
372/// The repository a task's config and disk check should read from. Mirrors
373/// `crate::conduct::repo_for`'s own fallback, duplicated rather than shared
374/// because that one is private to its module and the two are one `if` each.
375fn repo_for(task: &Task) -> PathBuf {
376    if task.repo.as_os_str().is_empty() {
377        PathBuf::from(".")
378    } else {
379        task.repo.clone()
380    }
381}
382
383/// Recognise a `HoldSource::Machine` hold caused by the free-space gate
384/// (`crate::disk::gate`'s message, or `daemon::disk_gate`'s "could not
385/// measure" fallback) from `hold_reason` text alone.
386///
387/// There is no field recording *why* a machine hold happened - `hold_machine`
388/// takes only a reason string - so this is the one signal available, and disk
389/// pressure is the one cause this module can safely re-measure without
390/// touching git, a run, or an agent CLI. Keep the two prefixes here in sync
391/// with `crate::disk::gate`'s formatted string and `daemon::disk_gate`'s own
392/// message if either changes; nothing else ties them together.
393fn is_disk_hold(task: &Task) -> bool {
394    task.hold_reason
395        .as_deref()
396        .is_some_and(crate::disk::is_gate_reason)
397}
398
399/// Has a `HoldSource::Machine` hold's cause resolved? `Some(true)` means yes -
400/// safe to requeue. `Some(false)` means the same cause was checked and is
401/// still in force. `None` means this hold's cause is not one this module
402/// knows how to re-check at all, and a human has to look.
403fn machine_cause_resolved(task: &Task, cfg: &Config) -> Option<bool> {
404    if !is_disk_hold(task) {
405        return None;
406    }
407    let min = cfg.disk.min_free_bytes;
408    if min == 0 {
409        // The operator turned the gate off since this hold was placed - its
410        // one possible cause is gone by construction, no measurement needed.
411        return Some(true);
412    }
413    let free = disk::free_bytes(&repo_for(task)).ok()?;
414    Some(disk::gate(free, min).is_none())
415}
416
417/// Is a `HoldSource::Manual` hold old enough to earn a "still wanted?"
418/// question? Same comparison `crate::clean::due` uses for a run's fold grace,
419/// against [`MANUAL_STALE_AFTER`] instead of a configured one.
420fn manual_is_stale(task: &Task, now: Timestamp) -> bool {
421    now.as_second() - task.updated_at.as_second() > MANUAL_STALE_AFTER.as_secs() as i64
422}
423
424/// The marker [`apply_answer`] writes into [`Task::hold_reason`] and
425/// [`already_applied`] reads back - see this module's own doc on why.
426fn marker_for(q: &Question) -> String {
427    format!("[triage:{}]", q.short())
428}
429
430/// Has `q`'s answer already been applied to `task`? See this module's doc.
431/// Checks [`Task::triage_applied`] first, which survives [`Task::release`].
432fn already_applied(task: &Task, q: &Question) -> bool {
433    if task.triage_applied(&q.id) {
434        return true;
435    }
436    // Records written before `Task::triage_applied` existed carry only the
437    // "keep held" marker in the hold reason.
438    let marker = marker_for(q);
439    task.hold_reason
440        .as_deref()
441        .is_some_and(|r| r.contains(marker.as_str()))
442}
443
444/// The most recent question this module filed for `task_id`, any status -
445/// open (still waiting), answered (may need applying), or abandoned (settled
446/// with nothing decided). Filters on both `node` and `run`, never `run`
447/// alone - see this module's doc on why a bare `run` match is not safe.
448fn latest_triage_question(questions: &Questions, task_id: &str) -> Option<Question> {
449    latest_question(questions, NODE, task_id)
450}
451
452/// [`latest_triage_question`] for any of this module's nodes.
453fn latest_question(questions: &Questions, node: &str, task_id: &str) -> Option<Question> {
454    questions
455        .list()
456        .into_iter()
457        .filter(|q| q.node == node && q.run == task_id)
458        // `asked_at` first: ids carry only whole seconds plus a random
459        // suffix, so two questions filed in the same second order randomly.
460        .max_by(|a, b| a.asked_at.cmp(&b.asked_at).then_with(|| a.id.cmp(&b.id)))
461}
462
463/// What an answered triage question's choice means, independent of which
464/// language it was filed in.
465#[derive(Debug, Clone, Copy, PartialEq, Eq)]
466enum AnswerAction {
467    /// The first choice, always "resume it" / "再開してよい" / "再開する" -
468    /// see [`Wording::choices3`] and [`Wording::choices2`], whose first entry
469    /// is always the resume choice.
470    Resume,
471    /// The third choice, present only in [`Wording::choices3`]: "discard it"
472    /// / "捨ててよい".
473    Discard,
474    /// Anything else: the second choice ("not yet" / "keep it held"), or an
475    /// answer that does not match any offered choice at all (should not
476    /// happen for a multiple-choice question, but the safe default is to
477    /// keep holding rather than guess at "resume").
478    KeepHeld,
479}
480
481/// Read `q`'s answer as an [`AnswerAction`].
482///
483/// Matched by **position in `q.choices`**, never by comparing the answer text
484/// against [`Wording`]'s own strings: [`Wording`] is picked from the task's
485/// *current* repository config, which can differ from whatever language was
486/// in force when the question was filed (a `--config` override on one `magi
487/// task triage` call and not the next, or the repository's config edited in
488/// between). Comparing text would then silently misread an actual "resume"
489/// answer as "keep held" - the choices themselves are fixed at filing time,
490/// in [`file_question`], and never change afterwards, so their position is
491/// the one thing that stays true regardless of which language reads them
492/// back.
493fn interpret_answer(q: &Question) -> AnswerAction {
494    let resolution = q.resolution().unwrap_or_default();
495    match q.choices.iter().position(|c| *c == resolution) {
496        Some(0) => AnswerAction::Resume,
497        Some(2) => AnswerAction::Discard,
498        _ => AnswerAction::KeepHeld,
499    }
500}
501
502/// The note [`already_applied`] looks for, appended to (never replacing)
503/// whatever [`Task::hold_reason`] already said - the original cause is still
504/// worth reading in `magi task show` after the operator answers "not yet",
505/// and [`Task::hold_manual`] would otherwise overwrite it outright.
506fn keep_held_note(task: &Task, q: &Question, resolution: &str) -> String {
507    let marker = format!("{} operator: {resolution}", marker_for(q));
508    match task.hold_reason.as_deref() {
509        Some(existing) if !existing.is_empty() => format!("{existing}\n{marker}"),
510        _ => marker,
511    }
512}
513
514/// File a fresh triage question for `task` and return it. The caller is
515/// responsible for having already established there is no open one - see
516/// [`latest_triage_question`] - so this never checks again.
517fn file_question(
518    questions: &Questions,
519    task: &Task,
520    bucket: Bucket,
521    w: &Wording,
522    now: Timestamp,
523) -> Option<Question> {
524    let (summary, why, choices) = match bucket {
525        Bucket::MachineUnknown => (
526            w.summary_machine_unknown(task),
527            w.why_machine().to_owned(),
528            w.choices3(),
529        ),
530        Bucket::ConductorOverride => (
531            w.summary_conductor_override(task),
532            task.resume_override
533                .as_ref()
534                .map(|o| w.why_conductor_override(o))
535                .unwrap_or_default(),
536            w.choices_conductor_override(),
537        ),
538        Bucket::Legacy => (
539            w.summary_legacy(task),
540            w.why_legacy().to_owned(),
541            w.choices3(),
542        ),
543        Bucket::ManualStale => {
544            let days = (now.as_second() - task.updated_at.as_second()) / (24 * 60 * 60);
545            (
546                w.summary_manual_stale(task, days),
547                w.why_manual().to_owned(),
548                w.choices2(),
549            )
550        }
551    };
552    let mut q = Question::new(
553        task.id.clone(),
554        NODE.to_owned(),
555        SEAT.to_owned(),
556        summary,
557        w.detail(task, &why),
558        choices,
559    );
560    questions.put(&mut q).ok()?;
561    Some(q)
562}
563
564/// Move every `blocked` task whose `blocked_by` names a task or question id
565/// that no longer exists to a machine hold, before the per-`held` walk
566/// [`run_once`] does gets a look at it.
567///
568/// `crate::daemon::resolve_blockers` already catches the same situation on
569/// every idle poll, and [`Queue::remove`] already catches it the moment a
570/// dependency is deleted through `magi task rm` - both call the same
571/// [`crate::queue::missing_blockers`]/[`crate::queue::missing_blocker_hold_reason`]
572/// this does. This third copy exists because a dependency can also be deleted
573/// by hand (the file just removed from disk, not through either of those
574/// paths), and because a queue can carry a `blocked` task with a
575/// long-since-deleted dependency from *before* either catch above ever
576/// existed - and such a task is `blocked`, never `held`, so it is invisible
577/// to the rest of this module without this pass. Running it here, first, is
578/// also what makes `magi task triage` alone - with no daemon running at all -
579/// enough to fix one: the task lands `held` in this same call, and the
580/// ordinary loop below files its question in the very same pass.
581fn quarantine_orphaned_blocked(queue: &Queue, questions: &Questions) -> Vec<String> {
582    let mut quarantined = Vec::new();
583    for listed in queue.list() {
584        if listed.status != TaskStatus::Blocked || listed.blocked_by.is_empty() {
585            continue;
586        }
587        let Ok(_claim) = queue.claim(&listed.id) else {
588            continue;
589        };
590        let Ok(mut task) = queue.get(&listed.id) else {
591            continue;
592        };
593        if task.status != TaskStatus::Blocked {
594            continue;
595        }
596        let deleted = queue.apply_deleted_blockers(&mut task);
597        if !deleted.is_empty() && queue.put(&mut task).is_ok() {
598            for id in &deleted {
599                queue.note_dependency_deleted(&task, id);
600            }
601        }
602        if task.status != TaskStatus::Blocked {
603            continue;
604        }
605        let missing = crate::queue::missing_blockers(queue, questions, &task.blocked_by);
606        if missing.is_empty() {
607            continue;
608        }
609        let language = crate::lang::of_repo(&repo_for(&task));
610        task.hold_machine(Some(crate::queue::missing_blocker_hold_reason_in(
611            &task.blocked_by,
612            &missing,
613            &language,
614        )));
615        if queue.put(&mut task).is_ok() {
616            quarantined.push(task.id.clone());
617        }
618    }
619    quarantined
620}
621
622/// Choices of a [`DEPS_NODE`] question, by position: release the root, discard
623/// it, or detach the dependants from it.
624#[derive(Debug, Clone, Copy, PartialEq, Eq)]
625enum DepsAction {
626    Release,
627    Discard,
628    Detach,
629    /// An answer that matches no offered choice: nothing is changed.
630    Nothing,
631}
632
633fn interpret_deps_answer(q: &Question) -> DepsAction {
634    let resolution = q.resolution().unwrap_or_default();
635    match q.choices.iter().position(|c| *c == resolution) {
636        Some(0) => DepsAction::Release,
637        Some(1) => DepsAction::Discard,
638        Some(2) => DepsAction::Detach,
639        _ => DepsAction::Nothing,
640    }
641}
642
643/// Summary, body and choices of the question about one stuck root.
644fn deps_texts(w: &Wording, root: &Task, dependants: &[&Task]) -> (String, String, Vec<String>) {
645    let ja = w.lang == "ja";
646    let reason = root
647        .hold_reason
648        .as_deref()
649        .or(root.last_error.as_deref())
650        .unwrap_or(if ja {
651            "(記録なし)"
652        } else {
653            "(none recorded)"
654        });
655    let list: String = dependants
656        .iter()
657        .map(|t| format!("- {} {}\n", t.short(), t.title))
658        .collect();
659    if ja {
660        (
661            format!(
662                "{} ({}) が {} 件のタスクを止めています",
663                root.short(),
664                root.status.as_str(),
665                dependants.len()
666            ),
667            format!(
668                "task: {} ({})\ntitle: {}\n状態: {}\n理由: {reason}\n\n\
669                 このタスクは自動では実行されないため、依存している次のタスクは永遠に待ち続けます:\n{list}",
670                root.id,
671                root.short(),
672                root.title,
673                root.status.as_str()
674            ),
675            vec![
676                "依存先を再開する".to_owned(),
677                "依存先を捨てる(依存タスクは切り離して実行)".to_owned(),
678                "依存タスクを切り離す(依存先はそのまま)".to_owned(),
679            ],
680        )
681    } else {
682        (
683            format!(
684                "{} ({}) is freezing {} blocked task(s)",
685                root.short(),
686                root.status.as_str(),
687                dependants.len()
688            ),
689            format!(
690                "task: {} ({})\ntitle: {}\nstatus: {}\nreason: {reason}\n\n\
691                 Nothing in the loop will ever run this task, so these dependants wait \
692                 forever:\n{list}",
693                root.id,
694                root.short(),
695                root.title,
696                root.status.as_str()
697            ),
698            vec![
699                "release the dependency".to_owned(),
700                "discard the dependency (dependants are detached and run)".to_owned(),
701                "detach the dependants (the dependency stays as it is)".to_owned(),
702            ],
703        )
704    }
705}
706
707/// Detach every `blocked` task that names `root` directly: drop that one id
708/// from its `blocked_by` (`Task::unblock`), so a dependant with another
709/// unresolved blocker keeps waiting on it. Idempotent.
710fn detach_dependants(queue: &Queue, root: &str) {
711    for listed in queue.list() {
712        if listed.status != TaskStatus::Blocked || !listed.blocked_by.iter().any(|b| b == root) {
713            continue;
714        }
715        let Ok(_claim) = queue.claim(&listed.id) else {
716            continue;
717        };
718        let Ok(mut t) = queue.get(&listed.id) else {
719            continue;
720        };
721        if t.status != TaskStatus::Blocked {
722            continue;
723        }
724        t.unblock(root);
725        let _ = queue.put(&mut t);
726    }
727}
728
729/// Apply an answered [`DEPS_NODE`] question. Every step is idempotent, so a
730/// pass that dies half way is finished by the next one. Returns whether the
731/// answer was consumed.
732fn apply_deps_answer(queue: &Queue, questions: &Questions, q: &Question, now: Timestamp) -> bool {
733    let Ok(_claim) = queue.claim(&q.run) else {
734        return false;
735    };
736    let Ok(mut root) = queue.get(&q.run) else {
737        return false;
738    };
739    if root.triage_applied(&q.id) {
740        return false;
741    }
742    match interpret_deps_answer(q) {
743        DepsAction::Release => {
744            if root.status == TaskStatus::Running {
745                return false;
746            }
747            root.release();
748            root.resume_override = Some(OperatorResume {
749                question_id: q.id.clone(),
750                at: now,
751                conductor_rehold: None,
752                forced: false,
753                pinned_run: None,
754            });
755        }
756        DepsAction::Discard => {
757            detach_dependants(queue, &root.id);
758            return queue.remove(&root.id, false, questions).is_ok();
759        }
760        DepsAction::Detach => detach_dependants(queue, &root.id),
761        // Matches no offered choice, so there is nothing to apply. Left unmarked:
762        // marking it applied would read as "settled", and the next stuck check
763        // would file an identical question at once. Unmarked, it stays an
764        // answered question awaiting application, which suppresses re-asking.
765        DepsAction::Nothing => return false,
766    }
767    root.mark_triage_applied(&q.id);
768    queue.put(&mut root).is_ok()
769}
770
771/// One question per stuck dependency root, and apply the answers to earlier
772/// ones. See [`crate::blockers`] for what "stuck" means.
773///
774/// Not re-asked: a root with an open or answered-unapplied question; one whose
775/// own hold question ([`NODE`]) is still pending; and one whose question was
776/// abandoned unless the root or a task it freezes changed since. An *applied*
777/// answer leaves nothing stuck (released, discarded, or detached), so a root
778/// that is stuck again is a new situation and is asked about afresh.
779fn ask_about_stuck_roots(
780    queue: &Queue,
781    questions: &Questions,
782    config_override: Option<&Path>,
783    now: Timestamp,
784    report: &mut Report,
785) {
786    for q in questions.list() {
787        if q.node == DEPS_NODE
788            && q.status == QuestionStatus::Answered
789            && apply_deps_answer(queue, questions, &q, now)
790        {
791            report.answered.push(q.run.clone());
792        }
793    }
794
795    let all = questions.list();
796    let inv = crate::blockers::Inventory::new(queue.list(), &all);
797    let mut frozen: std::collections::BTreeMap<String, Vec<String>> = Default::default();
798    for (id, roots) in inv.stuck() {
799        for root in roots {
800            if root != id {
801                frozen.entry(root).or_default().push(id.clone());
802            }
803        }
804    }
805    // A cycle root freezes the others but may have no dependants of its own
806    // listed above when it is alone; a root with nothing to name is not asked.
807
808    for mut q in all
809        .iter()
810        .filter(|q| q.node == DEPS_NODE && q.status.open())
811        .cloned()
812    {
813        if !frozen.contains_key(&q.run) {
814            q.abandon("nothing waits on this task anymore");
815            let _ = questions.put(&mut q);
816        }
817    }
818
819    for (root_id, dependant_ids) in &frozen {
820        let Some(root) = inv.task(root_id) else {
821            continue;
822        };
823        let dependants: Vec<&Task> = dependant_ids.iter().filter_map(|d| inv.task(d)).collect();
824        let latest = all
825            .iter()
826            .filter(|q| q.node == DEPS_NODE && q.run == *root_id)
827            .max_by(|a, b| a.asked_at.cmp(&b.asked_at).then_with(|| a.id.cmp(&b.id)));
828        match latest {
829            Some(q) if q.status.open() => continue,
830            Some(q) if q.status == QuestionStatus::Answered && !root.triage_applied(&q.id) => {
831                continue;
832            }
833            Some(q) if q.status == QuestionStatus::Abandoned => {
834                let changed = root.updated_at > q.asked_at
835                    || dependants.iter().any(|t| t.updated_at > q.asked_at);
836                if !changed {
837                    continue;
838                }
839            }
840            _ => {}
841        }
842        if pending_for(questions, root) {
843            continue;
844        }
845        let cfg = Config::discover(&repo_for(root), config_override)
846            .ok()
847            .map(|(c, _)| c);
848        let w = wording(cfg.as_ref().map_or("en", |c| c.graph.language.as_str()));
849        let (summary, detail, choices) = deps_texts(w, root, &dependants);
850        let mut question = Question::new(
851            root.id.clone(),
852            DEPS_NODE.to_owned(),
853            SEAT.to_owned(),
854            summary,
855            detail,
856            choices,
857        );
858        if questions.put(&mut question).is_ok() {
859            report.asked.push(root.id.clone());
860        }
861    }
862}
863
864/// Run one deterministic triage pass over every `held` task in `queue`. No
865/// model call anywhere in this function - see this module's own doc for what
866/// each `HoldSource` gets instead.
867///
868/// `config_override` is threaded straight to [`Config::discover`], the same
869/// role `daemon::Opts::config` plays for `daemon::prepare` - an explicit
870/// `--config` from the caller, or `None` to let each task's own repository
871/// pick its layers.
872///
873/// Safe to call on every daemon idle tick and from `magi task triage` alike:
874/// a task already answered and applied is left alone (see
875/// [`already_applied`]), and a task with an open question is left alone too,
876/// so repeated calls with nothing new to say do nothing.
877///
878/// Also runs [`quarantine_orphaned_blocked`] first, so a `blocked` task whose
879/// dependency no longer exists is caught and turned into a fresh `held`
880/// question in this same pass, not left for a later call to notice.
881pub fn run_once(
882    queue: &Queue,
883    questions: &Questions,
884    config_override: Option<&Path>,
885    now: Timestamp,
886) -> Report {
887    let mut report = Report {
888        quarantined: quarantine_orphaned_blocked(queue, questions),
889        ..Report::default()
890    };
891    ask_about_stuck_roots(queue, questions, config_override, now, &mut report);
892    for listed in queue.list() {
893        if listed.status != TaskStatus::Held {
894            continue;
895        }
896        let Ok(_claim) = queue.claim(&listed.id) else {
897            continue;
898        };
899        let Ok(mut task) = queue.get(&listed.id) else {
900            continue;
901        };
902        // Re-read under the claim: a release or a re-hold landed by a human
903        // between the listing above and the claim just taken must not be
904        // clobbered by a decision based on the stale copy.
905        if task.status != TaskStatus::Held {
906            continue;
907        }
908        // A question about what this task freezes is already the one question
909        // it owes the operator; a second, about the hold itself, would ask two
910        // things at once.
911        if deps_pending(questions, &task) {
912            continue;
913        }
914
915        let cfg = Config::discover(&repo_for(&task), config_override)
916            .ok()
917            .map(|(c, _)| c);
918        let w = wording(cfg.as_ref().map_or("en", |c| c.graph.language.as_str()));
919
920        if let Some(q) = latest_triage_question(questions, &task.id) {
921            if q.status.open() {
922                // Already asked, still waiting - nothing to do this pass.
923                continue;
924            }
925            if q.status == QuestionStatus::Answered && !already_applied(&task, &q) {
926                match interpret_answer(&q) {
927                    AnswerAction::Resume => {
928                        // A second "resume", to the question about the
929                        // conductor's re-hold, forces it: the conductor may
930                        // not hold this task again. Any other resume records
931                        // the answer so a re-hold can be recognised.
932                        let contradicted = task
933                            .resume_override
934                            .as_ref()
935                            .is_some_and(|o| o.conductor_rehold.is_some());
936                        let record = match task.resume_override.take() {
937                            Some(mut o) if contradicted => {
938                                o.forced = true;
939                                o
940                            }
941                            _ => OperatorResume {
942                                question_id: q.id.clone(),
943                                at: now,
944                                conductor_rehold: None,
945                                forced: false,
946                                pinned_run: None,
947                            },
948                        };
949                        task.release();
950                        task.resume_override = Some(record);
951                        task.mark_triage_applied(&q.id);
952                        if queue.put(&mut task).is_ok() {
953                            report.answered.push(task.id.clone());
954                        }
955                    }
956                    AnswerAction::Discard => {
957                        if queue.remove(&task.id, false, questions).is_ok() {
958                            report.answered.push(task.id.clone());
959                        }
960                    }
961                    AnswerAction::KeepHeld => {
962                        let resolution = q.resolution().unwrap_or_default();
963                        let note = keep_held_note(&task, &q, &resolution);
964                        task.hold_manual(Some(note));
965                        task.mark_triage_applied(&q.id);
966                        if queue.put(&mut task).is_ok() {
967                            report.answered.push(task.id.clone());
968                        }
969                    }
970                }
971                continue;
972            }
973            // Abandoned, or an already-applied answer: fall through to the
974            // ordinary per-source handling below, which is how a stale
975            // `HoldSource::Manual` re-ask - or a fresh machine/legacy
976            // question, once a prior one settled the task back into a hold -
977            // gets filed.
978        }
979
980        match task.hold_source {
981            Some(HoldSource::Machine) => {
982                let overridden = task
983                    .resume_override
984                    .as_ref()
985                    .is_some_and(|o| o.conductor_rehold.is_some() && !o.forced);
986                if overridden {
987                    if file_question(questions, &task, Bucket::ConductorOverride, w, now).is_some()
988                    {
989                        report.asked.push(task.id.clone());
990                    }
991                } else if cfg.as_ref().and_then(|c| machine_cause_resolved(&task, c)) == Some(true)
992                {
993                    task.release();
994                    if queue.put(&mut task).is_ok() {
995                        report.resumed.push(task.id.clone());
996                    }
997                } else if file_question(questions, &task, Bucket::MachineUnknown, w, now).is_some()
998                {
999                    report.asked.push(task.id.clone());
1000                }
1001            }
1002            None => {
1003                if file_question(questions, &task, Bucket::Legacy, w, now).is_some() {
1004                    report.asked.push(task.id.clone());
1005                }
1006            }
1007            Some(HoldSource::Manual) => {
1008                if manual_is_stale(&task, now)
1009                    && file_question(questions, &task, Bucket::ManualStale, w, now).is_some()
1010                {
1011                    report.asked.push(task.id.clone());
1012                }
1013            }
1014        }
1015    }
1016    report
1017}
1018
1019/// The open triage question about `task_id`, if any - what `magi task show`
1020/// prints so a held task's card names the question waiting on it, not only
1021/// its hold reason. `None` once it is answered or abandoned: nothing is
1022/// waiting on it anymore.
1023pub fn open_question_for(questions: &Questions, task_id: &str) -> Option<Question> {
1024    latest_triage_question(questions, task_id).filter(|q| q.status.open())
1025}
1026
1027/// Does this module still have unfinished business with `task`?
1028///
1029/// True while its latest triage question is still open (waiting on an
1030/// answer), and true for a beat longer than [`open_question_for`] alone
1031/// would say: once answered, the question sits [`QuestionStatus::Answered`]
1032/// but unread until the next [`run_once`] pass actually applies it (see
1033/// [`already_applied`]), and [`run_once`] only ever runs on a fully idle
1034/// daemon tick - far less often than `crate::conduct` polls. A caller that
1035/// only checked "is a question open" would walk straight through that gap
1036/// the moment the operator answers, moving the task out of `held` before
1037/// [`run_once`] gets a turn - orphaning the very answer it was about to
1038/// apply, the same failure mode this function exists to keep `crate::conduct`
1039/// out of. `crate::conduct::apply_one` is exactly that caller.
1040pub fn pending_for(questions: &Questions, task: &Task) -> bool {
1041    let own = match latest_triage_question(questions, &task.id) {
1042        Some(q) if q.status.open() => true,
1043        Some(q) if q.status == QuestionStatus::Answered => !already_applied(task, &q),
1044        _ => false,
1045    };
1046    own || deps_pending(questions, task)
1047}
1048
1049/// Is the latest [`DEPS_NODE`] question about `task` still open, or answered
1050/// but not yet applied? Same reasoning as [`pending_for`]: the gap between an
1051/// answer and the next [`run_once`] pass must not be walked through.
1052fn deps_pending(questions: &Questions, task: &Task) -> bool {
1053    match latest_question(questions, DEPS_NODE, &task.id) {
1054        Some(q) if q.status.open() => true,
1055        Some(q) if q.status == QuestionStatus::Answered => !task.triage_applied(&q.id),
1056        _ => false,
1057    }
1058}
1059
1060/// Every task id with an open triage question right now - what `magi task
1061/// list` uses to mark a held task that is already waiting on an operator
1062/// decision, rather than have it read identically to one nobody has looked
1063/// at yet.
1064pub fn open_task_ids(questions: &Questions) -> std::collections::BTreeSet<String> {
1065    questions
1066        .list()
1067        .into_iter()
1068        .filter(|q| (q.node == NODE || q.node == DEPS_NODE) && q.status.open())
1069        .map(|q| q.run)
1070        .collect()
1071}
1072
1073#[cfg(test)]
1074mod tests {
1075    use super::*;
1076    use crate::ask::Answer;
1077    use crate::queue::Source;
1078    use jiff::SignedDuration;
1079
1080    fn store() -> (tempfile::TempDir, Queue, Questions) {
1081        let dir = tempfile::tempdir().unwrap();
1082        let q = Queue::at(dir.path().join("queue"));
1083        let s = Questions::at(dir.path().join("questions"));
1084        (dir, q, s)
1085    }
1086
1087    fn task(title: &str, repo: PathBuf) -> Task {
1088        Task::new(title.to_owned(), format!("do {title}"), repo, Source::Human)
1089    }
1090
1091    /// `[disk] min_free_bytes = 0` written next to a fictional repo, so
1092    /// `machine_cause_resolved` never has to ask the real disk anything - the
1093    /// same fixture pattern `daemon`'s own idle-loop tests use.
1094    fn gate_disabled_config(dir: &std::path::Path) -> PathBuf {
1095        let config = dir.join("magi.toml");
1096        std::fs::write(&config, "[disk]\nmin_free_bytes = 0\n").unwrap();
1097        config
1098    }
1099
1100    #[test]
1101    fn a_resolved_machine_hold_is_requeued_automatically() {
1102        let (dir, q, questions) = store();
1103        let config = gate_disabled_config(dir.path());
1104        let mut t = task("disk pressure", dir.path().join("repo"));
1105        t.hold_machine(Some(
1106            "not enough free space to start a run: 10 bytes free, 100 required by \
1107             `[disk] min_free_bytes`"
1108                .to_owned(),
1109        ));
1110        q.put(&mut t).unwrap();
1111
1112        let report = run_once(&q, &questions, Some(&config), Timestamp::now());
1113        assert_eq!(report.resumed, [t.id.clone()]);
1114        assert!(report.asked.is_empty());
1115
1116        let back = q.get(&t.id).unwrap();
1117        assert_eq!(back.status, TaskStatus::Queued);
1118        assert!(back.hold_source.is_none());
1119        assert!(questions.list().is_empty(), "nothing needed asking");
1120    }
1121
1122    #[test]
1123    fn a_machine_hold_with_no_recognised_cause_gets_one_question_not_two() {
1124        let (dir, q, questions) = store();
1125        let mut t = task("gate went red", dir.path().join("repo"));
1126        t.hold_machine(Some("gate red".to_owned()));
1127        q.put(&mut t).unwrap();
1128
1129        let first = run_once(&q, &questions, None, Timestamp::now());
1130        assert_eq!(first.asked, [t.id.clone()]);
1131        assert!(first.resumed.is_empty());
1132
1133        let open: Vec<_> = questions
1134            .list()
1135            .into_iter()
1136            .filter(|q| q.status.open())
1137            .collect();
1138        assert_eq!(open.len(), 1);
1139        assert_eq!(open[0].run, t.id);
1140        assert_eq!(open[0].node, NODE);
1141        assert_eq!(open[0].choices.len(), 3);
1142
1143        // A second pass with nothing new must not file a second question.
1144        let second = run_once(&q, &questions, None, Timestamp::now());
1145        assert!(second.asked.is_empty());
1146        assert_eq!(
1147            questions
1148                .list()
1149                .into_iter()
1150                .filter(|q| q.status.open())
1151                .count(),
1152            1
1153        );
1154    }
1155
1156    #[test]
1157    fn a_legacy_hold_with_no_recorded_source_gets_exactly_one_question() {
1158        let (dir, q, questions) = store();
1159        let mut t = task("schema 1 record", dir.path().join("repo"));
1160        t.status = TaskStatus::Held;
1161        assert!(t.hold_source.is_none(), "the case this test is about");
1162        q.put(&mut t).unwrap();
1163
1164        let first = run_once(&q, &questions, None, Timestamp::now());
1165        assert_eq!(first.asked, [t.id.clone()]);
1166
1167        let second = run_once(&q, &questions, None, Timestamp::now());
1168        assert!(
1169            second.asked.is_empty(),
1170            "the same legacy hold must not be asked about twice"
1171        );
1172        assert_eq!(
1173            questions
1174                .list()
1175                .into_iter()
1176                .filter(|q| q.status.open())
1177                .count(),
1178            1
1179        );
1180    }
1181
1182    #[test]
1183    fn a_manual_hold_is_never_auto_resumed() {
1184        let (dir, q, questions) = store();
1185        let mut t = task("operator stopped this", dir.path().join("repo"));
1186        t.hold_manual(Some("waiting on a decision".to_owned()));
1187        q.put(&mut t).unwrap();
1188
1189        let report = run_once(&q, &questions, None, Timestamp::now());
1190        assert!(report.resumed.is_empty());
1191        // Fresh, not stale yet - no question either.
1192        assert!(report.asked.is_empty());
1193
1194        let back = q.get(&t.id).unwrap();
1195        assert_eq!(back.status, TaskStatus::Held);
1196        assert_eq!(back.hold_source, Some(HoldSource::Manual));
1197        assert!(questions.list().is_empty());
1198    }
1199
1200    #[test]
1201    fn a_blocked_task_on_a_deleted_dependency_is_held_and_asked_about_in_one_pass() {
1202        // The five real tasks this whole change exists for are `blocked`, not
1203        // `held`, and no daemon has to be running for `magi task triage` alone
1204        // to reach them - `run_once` must both quarantine and ask in the same
1205        // call.
1206        let (dir, q, questions) = store();
1207        let mut still_going = task("still valid", dir.path().join("repo"));
1208        q.put(&mut still_going).unwrap();
1209
1210        let mut t = task("orphaned", dir.path().join("repo"));
1211        t.block(
1212            vec!["20260101-000000-gone".to_owned(), still_going.id.clone()],
1213            Some("waits on both".to_owned()),
1214        );
1215        q.put(&mut t).unwrap();
1216
1217        let report = run_once(&q, &questions, None, Timestamp::now());
1218        assert_eq!(report.quarantined, [t.id.clone()]);
1219        assert_eq!(
1220            report.asked,
1221            [t.id.clone()],
1222            "the fresh machine hold must earn a question in the same pass"
1223        );
1224
1225        let after = q.get(&t.id).unwrap();
1226        assert_eq!(after.status, TaskStatus::Held);
1227        assert_eq!(after.hold_source, Some(HoldSource::Machine));
1228        assert!(after.blocked_by.is_empty());
1229
1230        // The reason is not disk-pressure wording, so this must not be read
1231        // as a disk hold and silently auto-resumed.
1232        assert!(!is_disk_hold(&after));
1233
1234        let open: Vec<_> = questions
1235            .list()
1236            .into_iter()
1237            .filter(|q| q.status.open())
1238            .collect();
1239        assert_eq!(open.len(), 1);
1240        assert_eq!(open[0].run, t.id);
1241
1242        // A second pass with nothing new files no second question.
1243        let second = run_once(&q, &questions, None, Timestamp::now());
1244        assert!(second.quarantined.is_empty());
1245        assert!(second.asked.is_empty());
1246    }
1247
1248    /// A held root with a chain of dependants: 9db7 <- 4135 <- 6081, and a
1249    /// second direct dependant that also waits on a queued task.
1250    fn stuck_chain(q: &Queue, dir: &std::path::Path) -> (Task, Task, Task) {
1251        let mut root = task("root", dir.join("repo"));
1252        root.hold_manual(Some("waiting".to_owned()));
1253        q.put(&mut root).unwrap();
1254        let mut mid = task("mid", dir.join("repo"));
1255        mid.block(vec![root.id.clone()], None);
1256        q.put(&mut mid).unwrap();
1257        let mut leaf = task("leaf", dir.join("repo"));
1258        leaf.block(vec![mid.id.clone()], None);
1259        q.put(&mut leaf).unwrap();
1260        (root, mid, leaf)
1261    }
1262
1263    fn deps_questions(questions: &Questions) -> Vec<Question> {
1264        questions
1265            .list()
1266            .into_iter()
1267            .filter(|q| q.node == DEPS_NODE)
1268            .collect()
1269    }
1270
1271    fn answer_deps(questions: &Questions, choice: usize) {
1272        let mut asked = deps_questions(questions).remove(0);
1273        let c = asked.choices[choice].clone();
1274        asked.answer(Answer::Choice(c)).unwrap();
1275        questions.put(&mut asked).unwrap();
1276    }
1277
1278    #[test]
1279    fn a_stuck_chain_earns_one_question_naming_every_dependant_and_is_not_reasked() {
1280        let (dir, q, questions) = store();
1281        let (root, mid, leaf) = stuck_chain(&q, dir.path());
1282
1283        let first = run_once(&q, &questions, None, Timestamp::now());
1284        assert_eq!(first.asked, std::slice::from_ref(&root.id));
1285        let asked = deps_questions(&questions);
1286        assert_eq!(asked.len(), 1, "one question per root, not per dependant");
1287        assert_eq!(asked[0].run, root.id);
1288        assert_eq!(asked[0].choices.len(), 3);
1289        assert!(asked[0].detail.contains(mid.short()), "{}", asked[0].detail);
1290        assert!(
1291            asked[0].detail.contains(leaf.short()),
1292            "{}",
1293            asked[0].detail
1294        );
1295
1296        let second = run_once(&q, &questions, None, Timestamp::now());
1297        assert!(second.asked.is_empty());
1298        assert_eq!(deps_questions(&questions).len(), 1);
1299        assert!(pending_for(&questions, &q.get(&root.id).unwrap()));
1300    }
1301
1302    #[test]
1303    fn a_stuck_root_with_a_machine_hold_gets_only_the_dependency_question() {
1304        let (dir, q, questions) = store();
1305        let mut root = task("root", dir.path().join("repo"));
1306        root.hold_machine(Some("gate red".to_owned()));
1307        q.put(&mut root).unwrap();
1308        let mut dep = task("dep", dir.path().join("repo"));
1309        dep.block(vec![root.id.clone()], None);
1310        q.put(&mut dep).unwrap();
1311
1312        run_once(&q, &questions, None, Timestamp::now());
1313        run_once(&q, &questions, None, Timestamp::now());
1314        assert_eq!(deps_questions(&questions).len(), 1);
1315        assert_eq!(
1316            questions.list().len(),
1317            1,
1318            "no second question about the hold"
1319        );
1320    }
1321
1322    #[test]
1323    fn a_dependency_that_can_still_run_asks_nothing() {
1324        let (dir, q, questions) = store();
1325        let mut live = task("live", dir.path().join("repo"));
1326        q.put(&mut live).unwrap();
1327        let mut dep = task("dep", dir.path().join("repo"));
1328        dep.block(vec![live.id.clone()], None);
1329        q.put(&mut dep).unwrap();
1330        let report = run_once(&q, &questions, None, Timestamp::now());
1331        assert!(report.asked.is_empty());
1332        assert!(questions.list().is_empty());
1333    }
1334
1335    #[test]
1336    fn answering_release_requeues_the_root_and_is_not_reasked() {
1337        let (dir, q, questions) = store();
1338        let (root, _mid, _leaf) = stuck_chain(&q, dir.path());
1339        run_once(&q, &questions, None, Timestamp::now());
1340        answer_deps(&questions, 0);
1341
1342        let report = run_once(&q, &questions, None, Timestamp::now());
1343        assert_eq!(report.answered, std::slice::from_ref(&root.id));
1344        assert!(report.asked.is_empty(), "applying must not re-ask");
1345        assert_eq!(q.get(&root.id).unwrap().status, TaskStatus::Queued);
1346        let again = run_once(&q, &questions, None, Timestamp::now());
1347        assert!(again.asked.is_empty() && again.answered.is_empty());
1348        assert_eq!(deps_questions(&questions).len(), 1);
1349    }
1350
1351    #[test]
1352    fn answering_detach_frees_the_direct_dependant_and_keeps_the_root() {
1353        let (dir, q, questions) = store();
1354        let (root, mid, leaf) = stuck_chain(&q, dir.path());
1355        run_once(&q, &questions, None, Timestamp::now());
1356        answer_deps(&questions, 2);
1357
1358        let report = run_once(&q, &questions, None, Timestamp::now());
1359        assert!(report.asked.is_empty());
1360        assert_eq!(q.get(&root.id).unwrap().status, TaskStatus::Held);
1361        assert_eq!(q.get(&mid.id).unwrap().status, TaskStatus::Queued);
1362        let leaf = q.get(&leaf.id).unwrap();
1363        assert_eq!(leaf.status, TaskStatus::Blocked, "still waits on mid");
1364        assert_eq!(leaf.blocked_by, std::slice::from_ref(&mid.id));
1365    }
1366
1367    #[test]
1368    fn a_dependant_that_can_progress_through_another_blocker_is_not_stuck() {
1369        let (dir, q, questions) = store();
1370        let (root, mid, _leaf) = stuck_chain(&q, dir.path());
1371        let mut live = task("live", dir.path().join("repo"));
1372        q.put(&mut live).unwrap();
1373        let mut both = q.get(&mid.id).unwrap();
1374        both.block(vec![root.id.clone(), live.id.clone()], None);
1375        q.put(&mut both).unwrap();
1376        // `mid` can still progress through `live`, so nothing is stuck yet.
1377        assert!(
1378            run_once(&q, &questions, None, Timestamp::now())
1379                .asked
1380                .is_empty()
1381        );
1382    }
1383
1384    #[test]
1385    fn answering_discard_detaches_then_removes_the_root() {
1386        let (dir, q, questions) = store();
1387        let (root, mid, _leaf) = stuck_chain(&q, dir.path());
1388        run_once(&q, &questions, None, Timestamp::now());
1389        answer_deps(&questions, 1);
1390
1391        let report = run_once(&q, &questions, None, Timestamp::now());
1392        assert!(q.get(&root.id).is_err());
1393        assert_eq!(q.get(&mid.id).unwrap().status, TaskStatus::Queued);
1394        assert!(
1395            report.asked.is_empty(),
1396            "no per-dependant quarantine question"
1397        );
1398        assert!(report.quarantined.is_empty());
1399    }
1400
1401    #[test]
1402    fn an_unmatched_deps_answer_changes_nothing_and_is_not_reasked() {
1403        let (dir, q, questions) = store();
1404        let (root, mid, _leaf) = stuck_chain(&q, dir.path());
1405        run_once(&q, &questions, None, Timestamp::now());
1406        let mut asked = deps_questions(&questions).remove(0);
1407        asked.choices.clear();
1408        asked.answer(Answer::Text("dunno".to_owned())).unwrap();
1409        questions.put(&mut asked).unwrap();
1410
1411        for _ in 0..2 {
1412            let report = run_once(&q, &questions, None, Timestamp::now());
1413            assert!(report.asked.is_empty() && report.answered.is_empty());
1414        }
1415        assert_eq!(deps_questions(&questions).len(), 1);
1416        assert_eq!(q.get(&root.id).unwrap().status, TaskStatus::Held);
1417        assert_eq!(q.get(&mid.id).unwrap().status, TaskStatus::Blocked);
1418    }
1419
1420    #[test]
1421    fn a_dependency_cycle_terminates_and_is_asked_about_once() {
1422        let (dir, q, questions) = store();
1423        let mut a = task("a", dir.path().join("repo"));
1424        let mut b = task("b", dir.path().join("repo"));
1425        a.block(vec![b.id.clone()], None);
1426        b.block(vec![a.id.clone()], None);
1427        q.put(&mut a).unwrap();
1428        q.put(&mut b).unwrap();
1429
1430        let first = run_once(&q, &questions, None, Timestamp::now());
1431        assert_eq!(first.asked.len(), 1);
1432        let second = run_once(&q, &questions, None, Timestamp::now());
1433        assert!(second.asked.is_empty());
1434        assert_eq!(deps_questions(&questions).len(), 1);
1435    }
1436
1437    #[test]
1438    fn a_missing_dependency_is_still_quarantined_not_asked_about_as_stuck() {
1439        let (dir, q, questions) = store();
1440        let mut t = task("orphan", dir.path().join("repo"));
1441        t.block(vec!["20260101-000000-gone".to_owned()], None);
1442        q.put(&mut t).unwrap();
1443        let report = run_once(&q, &questions, None, Timestamp::now());
1444        assert_eq!(report.quarantined, [t.id.clone()]);
1445        assert!(deps_questions(&questions).is_empty());
1446    }
1447
1448    #[test]
1449    fn a_stale_manual_hold_earns_a_two_choice_question() {
1450        let (dir, q, questions) = store();
1451        let mut t = task("been sitting a while", dir.path().join("repo"));
1452        t.hold_manual(Some("waiting on a decision".to_owned()));
1453        q.put(&mut t).unwrap();
1454        // Back-date the hold past MANUAL_STALE_AFTER without waiting a week.
1455        let mut back = q.get(&t.id).unwrap();
1456        back.updated_at = Timestamp::now() - SignedDuration::new(8 * 24 * 60 * 60, 0);
1457        std::fs::write(
1458            q.path_of(&back.id),
1459            serde_json::to_string_pretty(&back).unwrap(),
1460        )
1461        .unwrap();
1462
1463        let report = run_once(&q, &questions, None, Timestamp::now());
1464        assert_eq!(report.asked, [t.id.clone()]);
1465        let open: Vec<_> = questions
1466            .list()
1467            .into_iter()
1468            .filter(|q| q.status.open())
1469            .collect();
1470        assert_eq!(open.len(), 1);
1471        assert_eq!(open[0].choices.len(), 2);
1472    }
1473
1474    #[test]
1475    fn answering_resume_releases_the_task() {
1476        let (dir, q, questions) = store();
1477        let mut t = task("gate went red", dir.path().join("repo"));
1478        t.hold_machine(Some("gate red".to_owned()));
1479        q.put(&mut t).unwrap();
1480        run_once(&q, &questions, None, Timestamp::now());
1481
1482        let mut asked = questions
1483            .list()
1484            .into_iter()
1485            .find(|q| q.run == t.id)
1486            .unwrap();
1487        asked.answer(Answer::Choice(EN.resume.to_owned())).unwrap();
1488        questions.put(&mut asked).unwrap();
1489
1490        let report = run_once(&q, &questions, None, Timestamp::now());
1491        assert_eq!(report.answered, [t.id.clone()]);
1492        let back = q.get(&t.id).unwrap();
1493        assert_eq!(back.status, TaskStatus::Queued);
1494        assert!(back.hold_source.is_none());
1495    }
1496
1497    /// A task released by a "resume" answer, run, and failed back to `held`.
1498    fn resumed_then_failed(q: &Queue, questions: &Questions, dir: &std::path::Path) -> Task {
1499        let mut t = task("gate went red", dir.join("repo"));
1500        t.hold_machine(Some("gate red".to_owned()));
1501        q.put(&mut t).unwrap();
1502        run_once(q, questions, None, Timestamp::now());
1503        let mut asked = questions
1504            .list()
1505            .into_iter()
1506            .find(|q| q.run == t.id)
1507            .unwrap();
1508        asked.answer(Answer::Choice(EN.resume.to_owned())).unwrap();
1509        questions.put(&mut asked).unwrap();
1510
1511        let report = run_once(q, questions, None, Timestamp::now());
1512        assert_eq!(report.answered, [t.id.clone()]);
1513        let mut back = q.get(&t.id).unwrap();
1514        assert_eq!(back.status, TaskStatus::Queued);
1515        assert_eq!(back.attempts, 0);
1516
1517        back.start("run-1".to_owned());
1518        back.fail("rebase conflict", 1);
1519        assert_eq!(back.status, TaskStatus::Held);
1520        q.put(&mut back).unwrap();
1521        back
1522    }
1523
1524    #[test]
1525    fn a_resume_answer_is_applied_once_and_a_new_machine_hold_is_asked_about() {
1526        let (dir, q, questions) = store();
1527        let t = resumed_then_failed(&q, &questions, dir.path());
1528
1529        let report = run_once(&q, &questions, None, Timestamp::now());
1530        assert!(report.answered.is_empty(), "the old answer must not replay");
1531        assert_eq!(report.asked, std::slice::from_ref(&t.id));
1532        let back = q.get(&t.id).unwrap();
1533        assert_eq!(back.status, TaskStatus::Held);
1534        assert_eq!(
1535            questions
1536                .list()
1537                .into_iter()
1538                .filter(|q| q.status.open())
1539                .count(),
1540            1
1541        );
1542    }
1543
1544    #[test]
1545    fn a_manual_hold_placed_after_a_resume_is_not_undone_by_the_old_answer() {
1546        let (dir, q, questions) = store();
1547        let mut t = resumed_then_failed(&q, &questions, dir.path());
1548        t.hold_manual(Some("operator stopped this".to_owned()));
1549        q.put(&mut t).unwrap();
1550        let before = questions.list().len();
1551
1552        let report = run_once(&q, &questions, None, Timestamp::now());
1553        assert!(report.answered.is_empty());
1554        assert!(report.asked.is_empty());
1555        let back = q.get(&t.id).unwrap();
1556        assert_eq!(back.status, TaskStatus::Held);
1557        assert_eq!(back.hold_source, Some(HoldSource::Manual));
1558        assert_eq!(questions.list().len(), before);
1559    }
1560
1561    #[test]
1562    fn answering_not_yet_keeps_it_held_as_a_manual_hold_and_does_not_reapply() {
1563        let (dir, q, questions) = store();
1564        let mut t = task("gate went red", dir.path().join("repo"));
1565        t.hold_machine(Some("gate red".to_owned()));
1566        q.put(&mut t).unwrap();
1567        run_once(&q, &questions, None, Timestamp::now());
1568
1569        let mut asked = questions
1570            .list()
1571            .into_iter()
1572            .find(|q| q.run == t.id)
1573            .unwrap();
1574        asked.answer(Answer::Choice(EN.wait.to_owned())).unwrap();
1575        questions.put(&mut asked).unwrap();
1576
1577        let report = run_once(&q, &questions, None, Timestamp::now());
1578        assert_eq!(report.answered, [t.id.clone()]);
1579        let back = q.get(&t.id).unwrap();
1580        assert_eq!(back.status, TaskStatus::Held);
1581        assert_eq!(back.hold_source, Some(HoldSource::Manual));
1582        assert!(
1583            back.hold_reason
1584                .as_deref()
1585                .is_some_and(|r| r.contains("gate red")),
1586            "the original cause must survive a \"not yet\" answer, not just the \
1587             triage marker: {:?}",
1588            back.hold_reason
1589        );
1590
1591        // A third pass, same instant: the manual hold is fresh (just
1592        // touched), so nothing more happens - in particular the already
1593        // answered question is not re-applied.
1594        let third = run_once(&q, &questions, None, Timestamp::now());
1595        assert!(third.answered.is_empty());
1596        assert!(third.asked.is_empty());
1597    }
1598
1599    #[test]
1600    fn answering_discard_removes_the_task_entirely() {
1601        let (dir, q, questions) = store();
1602        let mut t = task("gate went red", dir.path().join("repo"));
1603        t.hold_machine(Some("gate red".to_owned()));
1604        q.put(&mut t).unwrap();
1605        run_once(&q, &questions, None, Timestamp::now());
1606
1607        let mut asked = questions
1608            .list()
1609            .into_iter()
1610            .find(|q| q.run == t.id)
1611            .unwrap();
1612        asked.answer(Answer::Choice(EN.discard.to_owned())).unwrap();
1613        questions.put(&mut asked).unwrap();
1614
1615        let report = run_once(&q, &questions, None, Timestamp::now());
1616        assert_eq!(report.answered, [t.id.clone()]);
1617        assert!(
1618            q.get(&t.id).is_err(),
1619            "\"discard it\" (捨ててよい) must actually discard the task, not \
1620             just leave it sitting held forever"
1621        );
1622    }
1623
1624    #[test]
1625    fn an_answer_is_read_by_its_position_in_choices_not_by_the_callers_current_language() {
1626        // Filed while the repo's config reads Japanese...
1627        let (dir, q, questions) = store();
1628        let ja_config = dir.path().join("ja.toml");
1629        std::fs::write(&ja_config, "[graph]\nlanguage = \"ja\"\n").unwrap();
1630        let mut t = task("gate went red", dir.path().join("repo"));
1631        t.hold_machine(Some("gate red".to_owned()));
1632        q.put(&mut t).unwrap();
1633        run_once(&q, &questions, Some(&ja_config), Timestamp::now());
1634
1635        let mut asked = questions
1636            .list()
1637            .into_iter()
1638            .find(|q| q.run == t.id)
1639            .unwrap();
1640        assert_eq!(asked.choices[0], JA.resume, "filed in Japanese");
1641        asked.answer(Answer::Choice(JA.resume.to_owned())).unwrap();
1642        questions.put(&mut asked).unwrap();
1643
1644        // ...but applied against an English config (a later `--config`, or an
1645        // edited repository config). The Japanese "再開してよい" answer must
1646        // still be read as a resume, not silently misread as "keep held"
1647        // because it fails a text comparison against the English wording.
1648        let en_config = dir.path().join("en.toml");
1649        std::fs::write(&en_config, "[graph]\nlanguage = \"en\"\n").unwrap();
1650        let report = run_once(&q, &questions, Some(&en_config), Timestamp::now());
1651        assert_eq!(report.answered, [t.id.clone()]);
1652        let back = q.get(&t.id).unwrap();
1653        assert_eq!(
1654            back.status,
1655            TaskStatus::Queued,
1656            "a resume answer must resume the task regardless of which \
1657             language it is read back in"
1658        );
1659    }
1660
1661    #[test]
1662    fn a_question_falls_back_to_last_error_when_hold_reason_was_never_set() {
1663        // A record written before `Task::fail` started copying `why` into
1664        // `hold_reason` too (or one written by a still older build) carries
1665        // only `last_error`. The question detail must still name a cause
1666        // rather than reading "(none recorded)".
1667        let (dir, q, questions) = store();
1668        let mut t = task("kept failing the gate", dir.path().join("repo"));
1669        t.start("run-1".to_owned());
1670        t.fail("gate red three times running", 1);
1671        assert_eq!(t.status, TaskStatus::Held);
1672        t.hold_reason = None; // simulate a pre-fix or pre-schema record
1673        q.put(&mut t).unwrap();
1674
1675        run_once(&q, &questions, None, Timestamp::now());
1676        let asked = questions
1677            .list()
1678            .into_iter()
1679            .find(|q| q.run == t.id)
1680            .unwrap();
1681        assert!(
1682            asked.detail.contains("gate red three times running"),
1683            "the question must surface `last_error` when there is no \
1684             `hold_reason` to show instead: {}",
1685            asked.detail
1686        );
1687    }
1688
1689    #[test]
1690    fn open_question_for_and_open_task_ids_reflect_only_what_is_still_waiting() {
1691        let (dir, q, questions) = store();
1692        let mut t = task("gate went red", dir.path().join("repo"));
1693        t.hold_machine(Some("gate red".to_owned()));
1694        q.put(&mut t).unwrap();
1695
1696        assert!(open_question_for(&questions, &t.id).is_none());
1697        assert!(!open_task_ids(&questions).contains(&t.id));
1698
1699        run_once(&q, &questions, None, Timestamp::now());
1700        assert!(open_question_for(&questions, &t.id).is_some());
1701        assert!(open_task_ids(&questions).contains(&t.id));
1702
1703        let mut asked = questions
1704            .list()
1705            .into_iter()
1706            .find(|q| q.run == t.id)
1707            .unwrap();
1708        asked.answer(Answer::Choice(EN.resume.to_owned())).unwrap();
1709        questions.put(&mut asked).unwrap();
1710
1711        assert!(
1712            open_question_for(&questions, &t.id).is_none(),
1713            "an answered question is no longer open"
1714        );
1715        assert!(!open_task_ids(&questions).contains(&t.id));
1716    }
1717
1718    #[test]
1719    fn pending_for_stays_true_between_an_answer_and_the_next_run_once_pass() {
1720        // `open_question_for` alone goes `None` the instant the operator
1721        // answers, well before `run_once` - idle-tick only - gets a turn to
1722        // actually apply that answer (see `already_applied`). `pending_for`
1723        // exists so a caller polling far more often than `run_once` does -
1724        // `crate::conduct::apply_one` - does not walk through that gap.
1725        let (dir, q, questions) = store();
1726        let mut t = task("gate went red", dir.path().join("repo"));
1727        t.hold_machine(Some("gate red".to_owned()));
1728        q.put(&mut t).unwrap();
1729
1730        assert!(!pending_for(&questions, &t));
1731
1732        run_once(&q, &questions, None, Timestamp::now());
1733        let held = q.get(&t.id).unwrap();
1734        assert!(pending_for(&questions, &held), "still waiting on an answer");
1735
1736        let mut asked = questions
1737            .list()
1738            .into_iter()
1739            .find(|q| q.run == t.id)
1740            .unwrap();
1741        asked.answer(Answer::Choice(EN.wait.to_owned())).unwrap();
1742        questions.put(&mut asked).unwrap();
1743        assert!(!asked.status.open());
1744
1745        // The race window: answered, but `run_once` has not run again yet.
1746        let still_held = q.get(&t.id).unwrap();
1747        assert!(
1748            pending_for(&questions, &still_held),
1749            "answered but not yet applied is still pending"
1750        );
1751
1752        run_once(&q, &questions, None, Timestamp::now());
1753        let after = q.get(&t.id).unwrap();
1754        assert!(
1755            !pending_for(&questions, &after),
1756            "the answer is applied now, nothing left pending"
1757        );
1758    }
1759}