Skip to main content

magi/
triage.rs

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