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