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 missing = crate::queue::missing_blockers(queue, questions, &task.blocked_by);
597        if missing.is_empty() {
598            continue;
599        }
600        let language = crate::lang::of_repo(&repo_for(&task));
601        task.hold_machine(Some(crate::queue::missing_blocker_hold_reason_in(
602            &task.blocked_by,
603            &missing,
604            &language,
605        )));
606        if queue.put(&mut task).is_ok() {
607            quarantined.push(task.id.clone());
608        }
609    }
610    quarantined
611}
612
613/// Choices of a [`DEPS_NODE`] question, by position: release the root, discard
614/// it, or detach the dependants from it.
615#[derive(Debug, Clone, Copy, PartialEq, Eq)]
616enum DepsAction {
617    Release,
618    Discard,
619    Detach,
620    /// An answer that matches no offered choice: nothing is changed.
621    Nothing,
622}
623
624fn interpret_deps_answer(q: &Question) -> DepsAction {
625    let resolution = q.resolution().unwrap_or_default();
626    match q.choices.iter().position(|c| *c == resolution) {
627        Some(0) => DepsAction::Release,
628        Some(1) => DepsAction::Discard,
629        Some(2) => DepsAction::Detach,
630        _ => DepsAction::Nothing,
631    }
632}
633
634/// Summary, body and choices of the question about one stuck root.
635fn deps_texts(w: &Wording, root: &Task, dependants: &[&Task]) -> (String, String, Vec<String>) {
636    let ja = w.lang == "ja";
637    let reason = root
638        .hold_reason
639        .as_deref()
640        .or(root.last_error.as_deref())
641        .unwrap_or(if ja {
642            "(記録なし)"
643        } else {
644            "(none recorded)"
645        });
646    let list: String = dependants
647        .iter()
648        .map(|t| format!("- {} {}\n", t.short(), t.title))
649        .collect();
650    if ja {
651        (
652            format!(
653                "{} ({}) が {} 件のタスクを止めています",
654                root.short(),
655                root.status.as_str(),
656                dependants.len()
657            ),
658            format!(
659                "task: {} ({})\ntitle: {}\n状態: {}\n理由: {reason}\n\n\
660                 このタスクは自動では実行されないため、依存している次のタスクは永遠に待ち続けます:\n{list}",
661                root.id,
662                root.short(),
663                root.title,
664                root.status.as_str()
665            ),
666            vec![
667                "依存先を再開する".to_owned(),
668                "依存先を捨てる(依存タスクは切り離して実行)".to_owned(),
669                "依存タスクを切り離す(依存先はそのまま)".to_owned(),
670            ],
671        )
672    } else {
673        (
674            format!(
675                "{} ({}) is freezing {} blocked task(s)",
676                root.short(),
677                root.status.as_str(),
678                dependants.len()
679            ),
680            format!(
681                "task: {} ({})\ntitle: {}\nstatus: {}\nreason: {reason}\n\n\
682                 Nothing in the loop will ever run this task, so these dependants wait \
683                 forever:\n{list}",
684                root.id,
685                root.short(),
686                root.title,
687                root.status.as_str()
688            ),
689            vec![
690                "release the dependency".to_owned(),
691                "discard the dependency (dependants are detached and run)".to_owned(),
692                "detach the dependants (the dependency stays as it is)".to_owned(),
693            ],
694        )
695    }
696}
697
698/// Detach every `blocked` task that names `root` directly: drop that one id
699/// from its `blocked_by` (`Task::unblock`), so a dependant with another
700/// unresolved blocker keeps waiting on it. Idempotent.
701fn detach_dependants(queue: &Queue, root: &str) {
702    for listed in queue.list() {
703        if listed.status != TaskStatus::Blocked || !listed.blocked_by.iter().any(|b| b == root) {
704            continue;
705        }
706        let Ok(_claim) = queue.claim(&listed.id) else {
707            continue;
708        };
709        let Ok(mut t) = queue.get(&listed.id) else {
710            continue;
711        };
712        if t.status != TaskStatus::Blocked {
713            continue;
714        }
715        t.unblock(root);
716        let _ = queue.put(&mut t);
717    }
718}
719
720/// Apply an answered [`DEPS_NODE`] question. Every step is idempotent, so a
721/// pass that dies half way is finished by the next one. Returns whether the
722/// answer was consumed.
723fn apply_deps_answer(queue: &Queue, questions: &Questions, q: &Question, now: Timestamp) -> bool {
724    let Ok(_claim) = queue.claim(&q.run) else {
725        return false;
726    };
727    let Ok(mut root) = queue.get(&q.run) else {
728        return false;
729    };
730    if root.triage_applied(&q.id) {
731        return false;
732    }
733    match interpret_deps_answer(q) {
734        DepsAction::Release => {
735            if root.status == TaskStatus::Running {
736                return false;
737            }
738            root.release();
739            root.resume_override = Some(OperatorResume {
740                question_id: q.id.clone(),
741                at: now,
742                conductor_rehold: None,
743                forced: false,
744                pinned_run: None,
745            });
746        }
747        DepsAction::Discard => {
748            detach_dependants(queue, &root.id);
749            return queue.remove(&root.id, false, questions).is_ok();
750        }
751        DepsAction::Detach => detach_dependants(queue, &root.id),
752        // Matches no offered choice, so there is nothing to apply. Left unmarked:
753        // marking it applied would read as "settled", and the next stuck check
754        // would file an identical question at once. Unmarked, it stays an
755        // answered question awaiting application, which suppresses re-asking.
756        DepsAction::Nothing => return false,
757    }
758    root.mark_triage_applied(&q.id);
759    queue.put(&mut root).is_ok()
760}
761
762/// One question per stuck dependency root, and apply the answers to earlier
763/// ones. See [`crate::blockers`] for what "stuck" means.
764///
765/// Not re-asked: a root with an open or answered-unapplied question; one whose
766/// own hold question ([`NODE`]) is still pending; and one whose question was
767/// abandoned unless the root or a task it freezes changed since. An *applied*
768/// answer leaves nothing stuck (released, discarded, or detached), so a root
769/// that is stuck again is a new situation and is asked about afresh.
770fn ask_about_stuck_roots(
771    queue: &Queue,
772    questions: &Questions,
773    config_override: Option<&Path>,
774    now: Timestamp,
775    report: &mut Report,
776) {
777    for q in questions.list() {
778        if q.node == DEPS_NODE
779            && q.status == QuestionStatus::Answered
780            && apply_deps_answer(queue, questions, &q, now)
781        {
782            report.answered.push(q.run.clone());
783        }
784    }
785
786    let all = questions.list();
787    let inv = crate::blockers::Inventory::new(queue.list(), &all);
788    let mut frozen: std::collections::BTreeMap<String, Vec<String>> = Default::default();
789    for (id, roots) in inv.stuck() {
790        for root in roots {
791            if root != id {
792                frozen.entry(root).or_default().push(id.clone());
793            }
794        }
795    }
796    // A cycle root freezes the others but may have no dependants of its own
797    // listed above when it is alone; a root with nothing to name is not asked.
798
799    for mut q in all
800        .iter()
801        .filter(|q| q.node == DEPS_NODE && q.status.open())
802        .cloned()
803    {
804        if !frozen.contains_key(&q.run) {
805            q.abandon("nothing waits on this task anymore");
806            let _ = questions.put(&mut q);
807        }
808    }
809
810    for (root_id, dependant_ids) in &frozen {
811        let Some(root) = inv.task(root_id) else {
812            continue;
813        };
814        let dependants: Vec<&Task> = dependant_ids.iter().filter_map(|d| inv.task(d)).collect();
815        let latest = all
816            .iter()
817            .filter(|q| q.node == DEPS_NODE && q.run == *root_id)
818            .max_by(|a, b| a.asked_at.cmp(&b.asked_at).then_with(|| a.id.cmp(&b.id)));
819        match latest {
820            Some(q) if q.status.open() => continue,
821            Some(q) if q.status == QuestionStatus::Answered && !root.triage_applied(&q.id) => {
822                continue;
823            }
824            Some(q) if q.status == QuestionStatus::Abandoned => {
825                let changed = root.updated_at > q.asked_at
826                    || dependants.iter().any(|t| t.updated_at > q.asked_at);
827                if !changed {
828                    continue;
829                }
830            }
831            _ => {}
832        }
833        if pending_for(questions, root) {
834            continue;
835        }
836        let cfg = Config::discover(&repo_for(root), config_override)
837            .ok()
838            .map(|(c, _)| c);
839        let w = wording(cfg.as_ref().map_or("en", |c| c.graph.language.as_str()));
840        let (summary, detail, choices) = deps_texts(w, root, &dependants);
841        let mut question = Question::new(
842            root.id.clone(),
843            DEPS_NODE.to_owned(),
844            SEAT.to_owned(),
845            summary,
846            detail,
847            choices,
848        );
849        if questions.put(&mut question).is_ok() {
850            report.asked.push(root.id.clone());
851        }
852    }
853}
854
855/// Run one deterministic triage pass over every `held` task in `queue`. No
856/// model call anywhere in this function - see this module's own doc for what
857/// each `HoldSource` gets instead.
858///
859/// `config_override` is threaded straight to [`Config::discover`], the same
860/// role `daemon::Opts::config` plays for `daemon::prepare` - an explicit
861/// `--config` from the caller, or `None` to let each task's own repository
862/// pick its layers.
863///
864/// Safe to call on every daemon idle tick and from `magi task triage` alike:
865/// a task already answered and applied is left alone (see
866/// [`already_applied`]), and a task with an open question is left alone too,
867/// so repeated calls with nothing new to say do nothing.
868///
869/// Also runs [`quarantine_orphaned_blocked`] first, so a `blocked` task whose
870/// dependency no longer exists is caught and turned into a fresh `held`
871/// question in this same pass, not left for a later call to notice.
872pub fn run_once(
873    queue: &Queue,
874    questions: &Questions,
875    config_override: Option<&Path>,
876    now: Timestamp,
877) -> Report {
878    let mut report = Report {
879        quarantined: quarantine_orphaned_blocked(queue, questions),
880        ..Report::default()
881    };
882    ask_about_stuck_roots(queue, questions, config_override, now, &mut report);
883    for listed in queue.list() {
884        if listed.status != TaskStatus::Held {
885            continue;
886        }
887        let Ok(_claim) = queue.claim(&listed.id) else {
888            continue;
889        };
890        let Ok(mut task) = queue.get(&listed.id) else {
891            continue;
892        };
893        // Re-read under the claim: a release or a re-hold landed by a human
894        // between the listing above and the claim just taken must not be
895        // clobbered by a decision based on the stale copy.
896        if task.status != TaskStatus::Held {
897            continue;
898        }
899        // A question about what this task freezes is already the one question
900        // it owes the operator; a second, about the hold itself, would ask two
901        // things at once.
902        if deps_pending(questions, &task) {
903            continue;
904        }
905
906        let cfg = Config::discover(&repo_for(&task), config_override)
907            .ok()
908            .map(|(c, _)| c);
909        let w = wording(cfg.as_ref().map_or("en", |c| c.graph.language.as_str()));
910
911        if let Some(q) = latest_triage_question(questions, &task.id) {
912            if q.status.open() {
913                // Already asked, still waiting - nothing to do this pass.
914                continue;
915            }
916            if q.status == QuestionStatus::Answered && !already_applied(&task, &q) {
917                match interpret_answer(&q) {
918                    AnswerAction::Resume => {
919                        // A second "resume", to the question about the
920                        // conductor's re-hold, forces it: the conductor may
921                        // not hold this task again. Any other resume records
922                        // the answer so a re-hold can be recognised.
923                        let contradicted = task
924                            .resume_override
925                            .as_ref()
926                            .is_some_and(|o| o.conductor_rehold.is_some());
927                        let record = match task.resume_override.take() {
928                            Some(mut o) if contradicted => {
929                                o.forced = true;
930                                o
931                            }
932                            _ => OperatorResume {
933                                question_id: q.id.clone(),
934                                at: now,
935                                conductor_rehold: None,
936                                forced: false,
937                                pinned_run: None,
938                            },
939                        };
940                        task.release();
941                        task.resume_override = Some(record);
942                        task.mark_triage_applied(&q.id);
943                        if queue.put(&mut task).is_ok() {
944                            report.answered.push(task.id.clone());
945                        }
946                    }
947                    AnswerAction::Discard => {
948                        if queue.remove(&task.id, false, questions).is_ok() {
949                            report.answered.push(task.id.clone());
950                        }
951                    }
952                    AnswerAction::KeepHeld => {
953                        let resolution = q.resolution().unwrap_or_default();
954                        let note = keep_held_note(&task, &q, &resolution);
955                        task.hold_manual(Some(note));
956                        task.mark_triage_applied(&q.id);
957                        if queue.put(&mut task).is_ok() {
958                            report.answered.push(task.id.clone());
959                        }
960                    }
961                }
962                continue;
963            }
964            // Abandoned, or an already-applied answer: fall through to the
965            // ordinary per-source handling below, which is how a stale
966            // `HoldSource::Manual` re-ask - or a fresh machine/legacy
967            // question, once a prior one settled the task back into a hold -
968            // gets filed.
969        }
970
971        match task.hold_source {
972            Some(HoldSource::Machine) => {
973                let overridden = task
974                    .resume_override
975                    .as_ref()
976                    .is_some_and(|o| o.conductor_rehold.is_some() && !o.forced);
977                if overridden {
978                    if file_question(questions, &task, Bucket::ConductorOverride, w, now).is_some()
979                    {
980                        report.asked.push(task.id.clone());
981                    }
982                } else if cfg.as_ref().and_then(|c| machine_cause_resolved(&task, c)) == Some(true)
983                {
984                    task.release();
985                    if queue.put(&mut task).is_ok() {
986                        report.resumed.push(task.id.clone());
987                    }
988                } else if file_question(questions, &task, Bucket::MachineUnknown, w, now).is_some()
989                {
990                    report.asked.push(task.id.clone());
991                }
992            }
993            None => {
994                if file_question(questions, &task, Bucket::Legacy, w, now).is_some() {
995                    report.asked.push(task.id.clone());
996                }
997            }
998            Some(HoldSource::Manual) => {
999                if manual_is_stale(&task, now)
1000                    && file_question(questions, &task, Bucket::ManualStale, w, now).is_some()
1001                {
1002                    report.asked.push(task.id.clone());
1003                }
1004            }
1005        }
1006    }
1007    report
1008}
1009
1010/// The open triage question about `task_id`, if any - what `magi task show`
1011/// prints so a held task's card names the question waiting on it, not only
1012/// its hold reason. `None` once it is answered or abandoned: nothing is
1013/// waiting on it anymore.
1014pub fn open_question_for(questions: &Questions, task_id: &str) -> Option<Question> {
1015    latest_triage_question(questions, task_id).filter(|q| q.status.open())
1016}
1017
1018/// Does this module still have unfinished business with `task`?
1019///
1020/// True while its latest triage question is still open (waiting on an
1021/// answer), and true for a beat longer than [`open_question_for`] alone
1022/// would say: once answered, the question sits [`QuestionStatus::Answered`]
1023/// but unread until the next [`run_once`] pass actually applies it (see
1024/// [`already_applied`]), and [`run_once`] only ever runs on a fully idle
1025/// daemon tick - far less often than `crate::conduct` polls. A caller that
1026/// only checked "is a question open" would walk straight through that gap
1027/// the moment the operator answers, moving the task out of `held` before
1028/// [`run_once`] gets a turn - orphaning the very answer it was about to
1029/// apply, the same failure mode this function exists to keep `crate::conduct`
1030/// out of. `crate::conduct::apply_one` is exactly that caller.
1031pub fn pending_for(questions: &Questions, task: &Task) -> bool {
1032    let own = match latest_triage_question(questions, &task.id) {
1033        Some(q) if q.status.open() => true,
1034        Some(q) if q.status == QuestionStatus::Answered => !already_applied(task, &q),
1035        _ => false,
1036    };
1037    own || deps_pending(questions, task)
1038}
1039
1040/// Is the latest [`DEPS_NODE`] question about `task` still open, or answered
1041/// but not yet applied? Same reasoning as [`pending_for`]: the gap between an
1042/// answer and the next [`run_once`] pass must not be walked through.
1043fn deps_pending(questions: &Questions, task: &Task) -> bool {
1044    match latest_question(questions, DEPS_NODE, &task.id) {
1045        Some(q) if q.status.open() => true,
1046        Some(q) if q.status == QuestionStatus::Answered => !task.triage_applied(&q.id),
1047        _ => false,
1048    }
1049}
1050
1051/// Every task id with an open triage question right now - what `magi task
1052/// list` uses to mark a held task that is already waiting on an operator
1053/// decision, rather than have it read identically to one nobody has looked
1054/// at yet.
1055pub fn open_task_ids(questions: &Questions) -> std::collections::BTreeSet<String> {
1056    questions
1057        .list()
1058        .into_iter()
1059        .filter(|q| (q.node == NODE || q.node == DEPS_NODE) && q.status.open())
1060        .map(|q| q.run)
1061        .collect()
1062}
1063
1064#[cfg(test)]
1065mod tests {
1066    use super::*;
1067    use crate::ask::Answer;
1068    use crate::queue::Source;
1069    use jiff::SignedDuration;
1070
1071    fn store() -> (tempfile::TempDir, Queue, Questions) {
1072        let dir = tempfile::tempdir().unwrap();
1073        let q = Queue::at(dir.path().join("queue"));
1074        let s = Questions::at(dir.path().join("questions"));
1075        (dir, q, s)
1076    }
1077
1078    fn task(title: &str, repo: PathBuf) -> Task {
1079        Task::new(title.to_owned(), format!("do {title}"), repo, Source::Human)
1080    }
1081
1082    /// `[disk] min_free_bytes = 0` written next to a fictional repo, so
1083    /// `machine_cause_resolved` never has to ask the real disk anything - the
1084    /// same fixture pattern `daemon`'s own idle-loop tests use.
1085    fn gate_disabled_config(dir: &std::path::Path) -> PathBuf {
1086        let config = dir.join("magi.toml");
1087        std::fs::write(&config, "[disk]\nmin_free_bytes = 0\n").unwrap();
1088        config
1089    }
1090
1091    #[test]
1092    fn a_resolved_machine_hold_is_requeued_automatically() {
1093        let (dir, q, questions) = store();
1094        let config = gate_disabled_config(dir.path());
1095        let mut t = task("disk pressure", dir.path().join("repo"));
1096        t.hold_machine(Some(
1097            "not enough free space to start a run: 10 bytes free, 100 required by \
1098             `[disk] min_free_bytes`"
1099                .to_owned(),
1100        ));
1101        q.put(&mut t).unwrap();
1102
1103        let report = run_once(&q, &questions, Some(&config), Timestamp::now());
1104        assert_eq!(report.resumed, [t.id.clone()]);
1105        assert!(report.asked.is_empty());
1106
1107        let back = q.get(&t.id).unwrap();
1108        assert_eq!(back.status, TaskStatus::Queued);
1109        assert!(back.hold_source.is_none());
1110        assert!(questions.list().is_empty(), "nothing needed asking");
1111    }
1112
1113    #[test]
1114    fn a_machine_hold_with_no_recognised_cause_gets_one_question_not_two() {
1115        let (dir, q, questions) = store();
1116        let mut t = task("gate went red", dir.path().join("repo"));
1117        t.hold_machine(Some("gate red".to_owned()));
1118        q.put(&mut t).unwrap();
1119
1120        let first = run_once(&q, &questions, None, Timestamp::now());
1121        assert_eq!(first.asked, [t.id.clone()]);
1122        assert!(first.resumed.is_empty());
1123
1124        let open: Vec<_> = questions
1125            .list()
1126            .into_iter()
1127            .filter(|q| q.status.open())
1128            .collect();
1129        assert_eq!(open.len(), 1);
1130        assert_eq!(open[0].run, t.id);
1131        assert_eq!(open[0].node, NODE);
1132        assert_eq!(open[0].choices.len(), 3);
1133
1134        // A second pass with nothing new must not file a second question.
1135        let second = run_once(&q, &questions, None, Timestamp::now());
1136        assert!(second.asked.is_empty());
1137        assert_eq!(
1138            questions
1139                .list()
1140                .into_iter()
1141                .filter(|q| q.status.open())
1142                .count(),
1143            1
1144        );
1145    }
1146
1147    #[test]
1148    fn a_legacy_hold_with_no_recorded_source_gets_exactly_one_question() {
1149        let (dir, q, questions) = store();
1150        let mut t = task("schema 1 record", dir.path().join("repo"));
1151        t.status = TaskStatus::Held;
1152        assert!(t.hold_source.is_none(), "the case this test is about");
1153        q.put(&mut t).unwrap();
1154
1155        let first = run_once(&q, &questions, None, Timestamp::now());
1156        assert_eq!(first.asked, [t.id.clone()]);
1157
1158        let second = run_once(&q, &questions, None, Timestamp::now());
1159        assert!(
1160            second.asked.is_empty(),
1161            "the same legacy hold must not be asked about twice"
1162        );
1163        assert_eq!(
1164            questions
1165                .list()
1166                .into_iter()
1167                .filter(|q| q.status.open())
1168                .count(),
1169            1
1170        );
1171    }
1172
1173    #[test]
1174    fn a_manual_hold_is_never_auto_resumed() {
1175        let (dir, q, questions) = store();
1176        let mut t = task("operator stopped this", dir.path().join("repo"));
1177        t.hold_manual(Some("waiting on a decision".to_owned()));
1178        q.put(&mut t).unwrap();
1179
1180        let report = run_once(&q, &questions, None, Timestamp::now());
1181        assert!(report.resumed.is_empty());
1182        // Fresh, not stale yet - no question either.
1183        assert!(report.asked.is_empty());
1184
1185        let back = q.get(&t.id).unwrap();
1186        assert_eq!(back.status, TaskStatus::Held);
1187        assert_eq!(back.hold_source, Some(HoldSource::Manual));
1188        assert!(questions.list().is_empty());
1189    }
1190
1191    #[test]
1192    fn a_blocked_task_on_a_deleted_dependency_is_held_and_asked_about_in_one_pass() {
1193        // The five real tasks this whole change exists for are `blocked`, not
1194        // `held`, and no daemon has to be running for `magi task triage` alone
1195        // to reach them - `run_once` must both quarantine and ask in the same
1196        // call.
1197        let (dir, q, questions) = store();
1198        let mut still_going = task("still valid", dir.path().join("repo"));
1199        q.put(&mut still_going).unwrap();
1200
1201        let mut t = task("orphaned", dir.path().join("repo"));
1202        t.block(
1203            vec!["20260101-000000-gone".to_owned(), still_going.id.clone()],
1204            Some("waits on both".to_owned()),
1205        );
1206        q.put(&mut t).unwrap();
1207
1208        let report = run_once(&q, &questions, None, Timestamp::now());
1209        assert_eq!(report.quarantined, [t.id.clone()]);
1210        assert_eq!(
1211            report.asked,
1212            [t.id.clone()],
1213            "the fresh machine hold must earn a question in the same pass"
1214        );
1215
1216        let after = q.get(&t.id).unwrap();
1217        assert_eq!(after.status, TaskStatus::Held);
1218        assert_eq!(after.hold_source, Some(HoldSource::Machine));
1219        assert!(after.blocked_by.is_empty());
1220
1221        // The reason is not disk-pressure wording, so this must not be read
1222        // as a disk hold and silently auto-resumed.
1223        assert!(!is_disk_hold(&after));
1224
1225        let open: Vec<_> = questions
1226            .list()
1227            .into_iter()
1228            .filter(|q| q.status.open())
1229            .collect();
1230        assert_eq!(open.len(), 1);
1231        assert_eq!(open[0].run, t.id);
1232
1233        // A second pass with nothing new files no second question.
1234        let second = run_once(&q, &questions, None, Timestamp::now());
1235        assert!(second.quarantined.is_empty());
1236        assert!(second.asked.is_empty());
1237    }
1238
1239    /// A held root with a chain of dependants: 9db7 <- 4135 <- 6081, and a
1240    /// second direct dependant that also waits on a queued task.
1241    fn stuck_chain(q: &Queue, dir: &std::path::Path) -> (Task, Task, Task) {
1242        let mut root = task("root", dir.join("repo"));
1243        root.hold_manual(Some("waiting".to_owned()));
1244        q.put(&mut root).unwrap();
1245        let mut mid = task("mid", dir.join("repo"));
1246        mid.block(vec![root.id.clone()], None);
1247        q.put(&mut mid).unwrap();
1248        let mut leaf = task("leaf", dir.join("repo"));
1249        leaf.block(vec![mid.id.clone()], None);
1250        q.put(&mut leaf).unwrap();
1251        (root, mid, leaf)
1252    }
1253
1254    fn deps_questions(questions: &Questions) -> Vec<Question> {
1255        questions
1256            .list()
1257            .into_iter()
1258            .filter(|q| q.node == DEPS_NODE)
1259            .collect()
1260    }
1261
1262    fn answer_deps(questions: &Questions, choice: usize) {
1263        let mut asked = deps_questions(questions).remove(0);
1264        let c = asked.choices[choice].clone();
1265        asked.answer(Answer::Choice(c)).unwrap();
1266        questions.put(&mut asked).unwrap();
1267    }
1268
1269    #[test]
1270    fn a_stuck_chain_earns_one_question_naming_every_dependant_and_is_not_reasked() {
1271        let (dir, q, questions) = store();
1272        let (root, mid, leaf) = stuck_chain(&q, dir.path());
1273
1274        let first = run_once(&q, &questions, None, Timestamp::now());
1275        assert_eq!(first.asked, std::slice::from_ref(&root.id));
1276        let asked = deps_questions(&questions);
1277        assert_eq!(asked.len(), 1, "one question per root, not per dependant");
1278        assert_eq!(asked[0].run, root.id);
1279        assert_eq!(asked[0].choices.len(), 3);
1280        assert!(asked[0].detail.contains(mid.short()), "{}", asked[0].detail);
1281        assert!(
1282            asked[0].detail.contains(leaf.short()),
1283            "{}",
1284            asked[0].detail
1285        );
1286
1287        let second = run_once(&q, &questions, None, Timestamp::now());
1288        assert!(second.asked.is_empty());
1289        assert_eq!(deps_questions(&questions).len(), 1);
1290        assert!(pending_for(&questions, &q.get(&root.id).unwrap()));
1291    }
1292
1293    #[test]
1294    fn a_stuck_root_with_a_machine_hold_gets_only_the_dependency_question() {
1295        let (dir, q, questions) = store();
1296        let mut root = task("root", dir.path().join("repo"));
1297        root.hold_machine(Some("gate red".to_owned()));
1298        q.put(&mut root).unwrap();
1299        let mut dep = task("dep", dir.path().join("repo"));
1300        dep.block(vec![root.id.clone()], None);
1301        q.put(&mut dep).unwrap();
1302
1303        run_once(&q, &questions, None, Timestamp::now());
1304        run_once(&q, &questions, None, Timestamp::now());
1305        assert_eq!(deps_questions(&questions).len(), 1);
1306        assert_eq!(
1307            questions.list().len(),
1308            1,
1309            "no second question about the hold"
1310        );
1311    }
1312
1313    #[test]
1314    fn a_dependency_that_can_still_run_asks_nothing() {
1315        let (dir, q, questions) = store();
1316        let mut live = task("live", dir.path().join("repo"));
1317        q.put(&mut live).unwrap();
1318        let mut dep = task("dep", dir.path().join("repo"));
1319        dep.block(vec![live.id.clone()], None);
1320        q.put(&mut dep).unwrap();
1321        let report = run_once(&q, &questions, None, Timestamp::now());
1322        assert!(report.asked.is_empty());
1323        assert!(questions.list().is_empty());
1324    }
1325
1326    #[test]
1327    fn answering_release_requeues_the_root_and_is_not_reasked() {
1328        let (dir, q, questions) = store();
1329        let (root, _mid, _leaf) = stuck_chain(&q, dir.path());
1330        run_once(&q, &questions, None, Timestamp::now());
1331        answer_deps(&questions, 0);
1332
1333        let report = run_once(&q, &questions, None, Timestamp::now());
1334        assert_eq!(report.answered, std::slice::from_ref(&root.id));
1335        assert!(report.asked.is_empty(), "applying must not re-ask");
1336        assert_eq!(q.get(&root.id).unwrap().status, TaskStatus::Queued);
1337        let again = run_once(&q, &questions, None, Timestamp::now());
1338        assert!(again.asked.is_empty() && again.answered.is_empty());
1339        assert_eq!(deps_questions(&questions).len(), 1);
1340    }
1341
1342    #[test]
1343    fn answering_detach_frees_the_direct_dependant_and_keeps_the_root() {
1344        let (dir, q, questions) = store();
1345        let (root, mid, leaf) = stuck_chain(&q, dir.path());
1346        run_once(&q, &questions, None, Timestamp::now());
1347        answer_deps(&questions, 2);
1348
1349        let report = run_once(&q, &questions, None, Timestamp::now());
1350        assert!(report.asked.is_empty());
1351        assert_eq!(q.get(&root.id).unwrap().status, TaskStatus::Held);
1352        assert_eq!(q.get(&mid.id).unwrap().status, TaskStatus::Queued);
1353        let leaf = q.get(&leaf.id).unwrap();
1354        assert_eq!(leaf.status, TaskStatus::Blocked, "still waits on mid");
1355        assert_eq!(leaf.blocked_by, std::slice::from_ref(&mid.id));
1356    }
1357
1358    #[test]
1359    fn a_dependant_that_can_progress_through_another_blocker_is_not_stuck() {
1360        let (dir, q, questions) = store();
1361        let (root, mid, _leaf) = stuck_chain(&q, dir.path());
1362        let mut live = task("live", dir.path().join("repo"));
1363        q.put(&mut live).unwrap();
1364        let mut both = q.get(&mid.id).unwrap();
1365        both.block(vec![root.id.clone(), live.id.clone()], None);
1366        q.put(&mut both).unwrap();
1367        // `mid` can still progress through `live`, so nothing is stuck yet.
1368        assert!(
1369            run_once(&q, &questions, None, Timestamp::now())
1370                .asked
1371                .is_empty()
1372        );
1373    }
1374
1375    #[test]
1376    fn answering_discard_detaches_then_removes_the_root() {
1377        let (dir, q, questions) = store();
1378        let (root, mid, _leaf) = stuck_chain(&q, dir.path());
1379        run_once(&q, &questions, None, Timestamp::now());
1380        answer_deps(&questions, 1);
1381
1382        let report = run_once(&q, &questions, None, Timestamp::now());
1383        assert!(q.get(&root.id).is_err());
1384        assert_eq!(q.get(&mid.id).unwrap().status, TaskStatus::Queued);
1385        assert!(
1386            report.asked.is_empty(),
1387            "no per-dependant quarantine question"
1388        );
1389        assert!(report.quarantined.is_empty());
1390    }
1391
1392    #[test]
1393    fn an_unmatched_deps_answer_changes_nothing_and_is_not_reasked() {
1394        let (dir, q, questions) = store();
1395        let (root, mid, _leaf) = stuck_chain(&q, dir.path());
1396        run_once(&q, &questions, None, Timestamp::now());
1397        let mut asked = deps_questions(&questions).remove(0);
1398        asked.choices.clear();
1399        asked.answer(Answer::Text("dunno".to_owned())).unwrap();
1400        questions.put(&mut asked).unwrap();
1401
1402        for _ in 0..2 {
1403            let report = run_once(&q, &questions, None, Timestamp::now());
1404            assert!(report.asked.is_empty() && report.answered.is_empty());
1405        }
1406        assert_eq!(deps_questions(&questions).len(), 1);
1407        assert_eq!(q.get(&root.id).unwrap().status, TaskStatus::Held);
1408        assert_eq!(q.get(&mid.id).unwrap().status, TaskStatus::Blocked);
1409    }
1410
1411    #[test]
1412    fn a_dependency_cycle_terminates_and_is_asked_about_once() {
1413        let (dir, q, questions) = store();
1414        let mut a = task("a", dir.path().join("repo"));
1415        let mut b = task("b", dir.path().join("repo"));
1416        a.block(vec![b.id.clone()], None);
1417        b.block(vec![a.id.clone()], None);
1418        q.put(&mut a).unwrap();
1419        q.put(&mut b).unwrap();
1420
1421        let first = run_once(&q, &questions, None, Timestamp::now());
1422        assert_eq!(first.asked.len(), 1);
1423        let second = run_once(&q, &questions, None, Timestamp::now());
1424        assert!(second.asked.is_empty());
1425        assert_eq!(deps_questions(&questions).len(), 1);
1426    }
1427
1428    #[test]
1429    fn a_missing_dependency_is_still_quarantined_not_asked_about_as_stuck() {
1430        let (dir, q, questions) = store();
1431        let mut t = task("orphan", dir.path().join("repo"));
1432        t.block(vec!["20260101-000000-gone".to_owned()], None);
1433        q.put(&mut t).unwrap();
1434        let report = run_once(&q, &questions, None, Timestamp::now());
1435        assert_eq!(report.quarantined, [t.id.clone()]);
1436        assert!(deps_questions(&questions).is_empty());
1437    }
1438
1439    #[test]
1440    fn a_stale_manual_hold_earns_a_two_choice_question() {
1441        let (dir, q, questions) = store();
1442        let mut t = task("been sitting a while", dir.path().join("repo"));
1443        t.hold_manual(Some("waiting on a decision".to_owned()));
1444        q.put(&mut t).unwrap();
1445        // Back-date the hold past MANUAL_STALE_AFTER without waiting a week.
1446        let mut back = q.get(&t.id).unwrap();
1447        back.updated_at = Timestamp::now() - SignedDuration::new(8 * 24 * 60 * 60, 0);
1448        std::fs::write(
1449            q.path_of(&back.id),
1450            serde_json::to_string_pretty(&back).unwrap(),
1451        )
1452        .unwrap();
1453
1454        let report = run_once(&q, &questions, None, Timestamp::now());
1455        assert_eq!(report.asked, [t.id.clone()]);
1456        let open: Vec<_> = questions
1457            .list()
1458            .into_iter()
1459            .filter(|q| q.status.open())
1460            .collect();
1461        assert_eq!(open.len(), 1);
1462        assert_eq!(open[0].choices.len(), 2);
1463    }
1464
1465    #[test]
1466    fn answering_resume_releases_the_task() {
1467        let (dir, q, questions) = store();
1468        let mut t = task("gate went red", dir.path().join("repo"));
1469        t.hold_machine(Some("gate red".to_owned()));
1470        q.put(&mut t).unwrap();
1471        run_once(&q, &questions, None, Timestamp::now());
1472
1473        let mut asked = questions
1474            .list()
1475            .into_iter()
1476            .find(|q| q.run == t.id)
1477            .unwrap();
1478        asked.answer(Answer::Choice(EN.resume.to_owned())).unwrap();
1479        questions.put(&mut asked).unwrap();
1480
1481        let report = run_once(&q, &questions, None, Timestamp::now());
1482        assert_eq!(report.answered, [t.id.clone()]);
1483        let back = q.get(&t.id).unwrap();
1484        assert_eq!(back.status, TaskStatus::Queued);
1485        assert!(back.hold_source.is_none());
1486    }
1487
1488    /// A task released by a "resume" answer, run, and failed back to `held`.
1489    fn resumed_then_failed(q: &Queue, questions: &Questions, dir: &std::path::Path) -> Task {
1490        let mut t = task("gate went red", dir.join("repo"));
1491        t.hold_machine(Some("gate red".to_owned()));
1492        q.put(&mut t).unwrap();
1493        run_once(q, questions, None, Timestamp::now());
1494        let mut asked = questions
1495            .list()
1496            .into_iter()
1497            .find(|q| q.run == t.id)
1498            .unwrap();
1499        asked.answer(Answer::Choice(EN.resume.to_owned())).unwrap();
1500        questions.put(&mut asked).unwrap();
1501
1502        let report = run_once(q, questions, None, Timestamp::now());
1503        assert_eq!(report.answered, [t.id.clone()]);
1504        let mut back = q.get(&t.id).unwrap();
1505        assert_eq!(back.status, TaskStatus::Queued);
1506        assert_eq!(back.attempts, 0);
1507
1508        back.start("run-1".to_owned());
1509        back.fail("rebase conflict", 1);
1510        assert_eq!(back.status, TaskStatus::Held);
1511        q.put(&mut back).unwrap();
1512        back
1513    }
1514
1515    #[test]
1516    fn a_resume_answer_is_applied_once_and_a_new_machine_hold_is_asked_about() {
1517        let (dir, q, questions) = store();
1518        let t = resumed_then_failed(&q, &questions, dir.path());
1519
1520        let report = run_once(&q, &questions, None, Timestamp::now());
1521        assert!(report.answered.is_empty(), "the old answer must not replay");
1522        assert_eq!(report.asked, std::slice::from_ref(&t.id));
1523        let back = q.get(&t.id).unwrap();
1524        assert_eq!(back.status, TaskStatus::Held);
1525        assert_eq!(
1526            questions
1527                .list()
1528                .into_iter()
1529                .filter(|q| q.status.open())
1530                .count(),
1531            1
1532        );
1533    }
1534
1535    #[test]
1536    fn a_manual_hold_placed_after_a_resume_is_not_undone_by_the_old_answer() {
1537        let (dir, q, questions) = store();
1538        let mut t = resumed_then_failed(&q, &questions, dir.path());
1539        t.hold_manual(Some("operator stopped this".to_owned()));
1540        q.put(&mut t).unwrap();
1541        let before = questions.list().len();
1542
1543        let report = run_once(&q, &questions, None, Timestamp::now());
1544        assert!(report.answered.is_empty());
1545        assert!(report.asked.is_empty());
1546        let back = q.get(&t.id).unwrap();
1547        assert_eq!(back.status, TaskStatus::Held);
1548        assert_eq!(back.hold_source, Some(HoldSource::Manual));
1549        assert_eq!(questions.list().len(), before);
1550    }
1551
1552    #[test]
1553    fn answering_not_yet_keeps_it_held_as_a_manual_hold_and_does_not_reapply() {
1554        let (dir, q, questions) = store();
1555        let mut t = task("gate went red", dir.path().join("repo"));
1556        t.hold_machine(Some("gate red".to_owned()));
1557        q.put(&mut t).unwrap();
1558        run_once(&q, &questions, None, Timestamp::now());
1559
1560        let mut asked = questions
1561            .list()
1562            .into_iter()
1563            .find(|q| q.run == t.id)
1564            .unwrap();
1565        asked.answer(Answer::Choice(EN.wait.to_owned())).unwrap();
1566        questions.put(&mut asked).unwrap();
1567
1568        let report = run_once(&q, &questions, None, Timestamp::now());
1569        assert_eq!(report.answered, [t.id.clone()]);
1570        let back = q.get(&t.id).unwrap();
1571        assert_eq!(back.status, TaskStatus::Held);
1572        assert_eq!(back.hold_source, Some(HoldSource::Manual));
1573        assert!(
1574            back.hold_reason
1575                .as_deref()
1576                .is_some_and(|r| r.contains("gate red")),
1577            "the original cause must survive a \"not yet\" answer, not just the \
1578             triage marker: {:?}",
1579            back.hold_reason
1580        );
1581
1582        // A third pass, same instant: the manual hold is fresh (just
1583        // touched), so nothing more happens - in particular the already
1584        // answered question is not re-applied.
1585        let third = run_once(&q, &questions, None, Timestamp::now());
1586        assert!(third.answered.is_empty());
1587        assert!(third.asked.is_empty());
1588    }
1589
1590    #[test]
1591    fn answering_discard_removes_the_task_entirely() {
1592        let (dir, q, questions) = store();
1593        let mut t = task("gate went red", dir.path().join("repo"));
1594        t.hold_machine(Some("gate red".to_owned()));
1595        q.put(&mut t).unwrap();
1596        run_once(&q, &questions, None, Timestamp::now());
1597
1598        let mut asked = questions
1599            .list()
1600            .into_iter()
1601            .find(|q| q.run == t.id)
1602            .unwrap();
1603        asked.answer(Answer::Choice(EN.discard.to_owned())).unwrap();
1604        questions.put(&mut asked).unwrap();
1605
1606        let report = run_once(&q, &questions, None, Timestamp::now());
1607        assert_eq!(report.answered, [t.id.clone()]);
1608        assert!(
1609            q.get(&t.id).is_err(),
1610            "\"discard it\" (捨ててよい) must actually discard the task, not \
1611             just leave it sitting held forever"
1612        );
1613    }
1614
1615    #[test]
1616    fn an_answer_is_read_by_its_position_in_choices_not_by_the_callers_current_language() {
1617        // Filed while the repo's config reads Japanese...
1618        let (dir, q, questions) = store();
1619        let ja_config = dir.path().join("ja.toml");
1620        std::fs::write(&ja_config, "[graph]\nlanguage = \"ja\"\n").unwrap();
1621        let mut t = task("gate went red", dir.path().join("repo"));
1622        t.hold_machine(Some("gate red".to_owned()));
1623        q.put(&mut t).unwrap();
1624        run_once(&q, &questions, Some(&ja_config), Timestamp::now());
1625
1626        let mut asked = questions
1627            .list()
1628            .into_iter()
1629            .find(|q| q.run == t.id)
1630            .unwrap();
1631        assert_eq!(asked.choices[0], JA.resume, "filed in Japanese");
1632        asked.answer(Answer::Choice(JA.resume.to_owned())).unwrap();
1633        questions.put(&mut asked).unwrap();
1634
1635        // ...but applied against an English config (a later `--config`, or an
1636        // edited repository config). The Japanese "再開してよい" answer must
1637        // still be read as a resume, not silently misread as "keep held"
1638        // because it fails a text comparison against the English wording.
1639        let en_config = dir.path().join("en.toml");
1640        std::fs::write(&en_config, "[graph]\nlanguage = \"en\"\n").unwrap();
1641        let report = run_once(&q, &questions, Some(&en_config), Timestamp::now());
1642        assert_eq!(report.answered, [t.id.clone()]);
1643        let back = q.get(&t.id).unwrap();
1644        assert_eq!(
1645            back.status,
1646            TaskStatus::Queued,
1647            "a resume answer must resume the task regardless of which \
1648             language it is read back in"
1649        );
1650    }
1651
1652    #[test]
1653    fn a_question_falls_back_to_last_error_when_hold_reason_was_never_set() {
1654        // A record written before `Task::fail` started copying `why` into
1655        // `hold_reason` too (or one written by a still older build) carries
1656        // only `last_error`. The question detail must still name a cause
1657        // rather than reading "(none recorded)".
1658        let (dir, q, questions) = store();
1659        let mut t = task("kept failing the gate", dir.path().join("repo"));
1660        t.start("run-1".to_owned());
1661        t.fail("gate red three times running", 1);
1662        assert_eq!(t.status, TaskStatus::Held);
1663        t.hold_reason = None; // simulate a pre-fix or pre-schema record
1664        q.put(&mut t).unwrap();
1665
1666        run_once(&q, &questions, None, Timestamp::now());
1667        let asked = questions
1668            .list()
1669            .into_iter()
1670            .find(|q| q.run == t.id)
1671            .unwrap();
1672        assert!(
1673            asked.detail.contains("gate red three times running"),
1674            "the question must surface `last_error` when there is no \
1675             `hold_reason` to show instead: {}",
1676            asked.detail
1677        );
1678    }
1679
1680    #[test]
1681    fn open_question_for_and_open_task_ids_reflect_only_what_is_still_waiting() {
1682        let (dir, q, questions) = store();
1683        let mut t = task("gate went red", dir.path().join("repo"));
1684        t.hold_machine(Some("gate red".to_owned()));
1685        q.put(&mut t).unwrap();
1686
1687        assert!(open_question_for(&questions, &t.id).is_none());
1688        assert!(!open_task_ids(&questions).contains(&t.id));
1689
1690        run_once(&q, &questions, None, Timestamp::now());
1691        assert!(open_question_for(&questions, &t.id).is_some());
1692        assert!(open_task_ids(&questions).contains(&t.id));
1693
1694        let mut asked = questions
1695            .list()
1696            .into_iter()
1697            .find(|q| q.run == t.id)
1698            .unwrap();
1699        asked.answer(Answer::Choice(EN.resume.to_owned())).unwrap();
1700        questions.put(&mut asked).unwrap();
1701
1702        assert!(
1703            open_question_for(&questions, &t.id).is_none(),
1704            "an answered question is no longer open"
1705        );
1706        assert!(!open_task_ids(&questions).contains(&t.id));
1707    }
1708
1709    #[test]
1710    fn pending_for_stays_true_between_an_answer_and_the_next_run_once_pass() {
1711        // `open_question_for` alone goes `None` the instant the operator
1712        // answers, well before `run_once` - idle-tick only - gets a turn to
1713        // actually apply that answer (see `already_applied`). `pending_for`
1714        // exists so a caller polling far more often than `run_once` does -
1715        // `crate::conduct::apply_one` - does not walk through that gap.
1716        let (dir, q, questions) = store();
1717        let mut t = task("gate went red", dir.path().join("repo"));
1718        t.hold_machine(Some("gate red".to_owned()));
1719        q.put(&mut t).unwrap();
1720
1721        assert!(!pending_for(&questions, &t));
1722
1723        run_once(&q, &questions, None, Timestamp::now());
1724        let held = q.get(&t.id).unwrap();
1725        assert!(pending_for(&questions, &held), "still waiting on an answer");
1726
1727        let mut asked = questions
1728            .list()
1729            .into_iter()
1730            .find(|q| q.run == t.id)
1731            .unwrap();
1732        asked.answer(Answer::Choice(EN.wait.to_owned())).unwrap();
1733        questions.put(&mut asked).unwrap();
1734        assert!(!asked.status.open());
1735
1736        // The race window: answered, but `run_once` has not run again yet.
1737        let still_held = q.get(&t.id).unwrap();
1738        assert!(
1739            pending_for(&questions, &still_held),
1740            "answered but not yet applied is still pending"
1741        );
1742
1743        run_once(&q, &questions, None, Timestamp::now());
1744        let after = q.get(&t.id).unwrap();
1745        assert!(
1746            !pending_for(&questions, &after),
1747            "the answer is applied now, nothing left pending"
1748        );
1749    }
1750}