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