Skip to main content

loopflow/task/
mod.rs

1//! Durable execution of one Linear Task.
2//!
3//! A Task owns one durable worktree and provider transcript. Ordered PRs
4//! own the serial branches that advance the Task. Publication intent is recorded
5//! before GitHub is called, then the GitHub PR record is attached to that intent.
6
7use std::path::PathBuf;
8use std::str::FromStr;
9
10use serde::{Deserialize, Serialize};
11use time::OffsetDateTime;
12
13use crate::child::{prefixed_uuid_id, AbandonIntent};
14use crate::durable::RunId;
15pub use crate::durable::TaskId;
16use crate::id::WaveId;
17use crate::planning::TaskPlan;
18use crate::project::ProjectId;
19
20pub mod actions;
21pub mod runner;
22
23#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
24pub enum TaskDataError {
25    #[error("invalid Task id: {0}")]
26    InvalidId(String),
27    #[error("invalid Task: {0}")]
28    InvalidInvariant(String),
29}
30
31prefixed_uuid_id!(TaskPrId, "pr_", TaskDataError, TaskDataError::InvalidId);
32
33#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
34#[serde(rename_all = "snake_case")]
35#[non_exhaustive]
36pub enum TaskLifecyclePhase {
37    First,
38    Loop,
39    Finally,
40}
41
42impl TaskLifecyclePhase {
43    pub fn as_str(self) -> &'static str {
44        match self {
45            Self::First => "first",
46            Self::Loop => "loop",
47            Self::Finally => "finally",
48        }
49    }
50
51    pub(crate) fn storage_str(self) -> &'static str {
52        match self {
53            Self::First => "kickoff",
54            Self::Loop => "iterate",
55            Self::Finally => "gate",
56        }
57    }
58
59    pub(crate) fn from_storage_str(value: &str) -> Result<Self, TaskDataError> {
60        match value {
61            "kickoff" => Ok(Self::First),
62            "iterate" => Ok(Self::Loop),
63            "gate" => Ok(Self::Finally),
64            _ => Err(TaskDataError::InvalidInvariant(format!(
65                "invalid stored Task lifecycle phase: {value}"
66            ))),
67        }
68    }
69}
70
71#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
72pub struct TaskPhasePlan {
73    pub flow: String,
74}
75
76impl TaskPhasePlan {
77    fn validate(&self, phase: TaskLifecyclePhase) -> Result<(), TaskDataError> {
78        if self.flow.trim().is_empty() {
79            return Err(TaskDataError::InvalidInvariant(format!(
80                "{} flow cannot be empty",
81                phase.as_str()
82            )));
83        }
84        Ok(())
85    }
86}
87
88#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
89pub struct TaskLifecyclePlan {
90    pub first: TaskPhasePlan,
91    #[serde(rename = "loop")]
92    pub loop_: TaskPhasePlan,
93    pub finally: TaskPhasePlan,
94}
95
96impl TaskLifecyclePlan {
97    pub fn standard(
98        first_flow: impl Into<String>,
99        loop_flow: impl Into<String>,
100        finally_flow: impl Into<String>,
101    ) -> Self {
102        Self {
103            first: TaskPhasePlan {
104                flow: first_flow.into(),
105            },
106            loop_: TaskPhasePlan {
107                flow: loop_flow.into(),
108            },
109            finally: TaskPhasePlan {
110                flow: finally_flow.into(),
111            },
112        }
113    }
114
115    pub fn defaults() -> Self {
116        Self::standard("task-design", "slice", "ship")
117    }
118    pub fn phase(&self, phase: TaskLifecyclePhase) -> &TaskPhasePlan {
119        match phase {
120            TaskLifecyclePhase::First => &self.first,
121            TaskLifecyclePhase::Loop => &self.loop_,
122            TaskLifecyclePhase::Finally => &self.finally,
123        }
124    }
125
126    fn validate(&self) -> Result<(), TaskDataError> {
127        self.first.validate(TaskLifecyclePhase::First)?;
128        self.loop_.validate(TaskLifecyclePhase::Loop)?;
129        self.finally.validate(TaskLifecyclePhase::Finally)
130    }
131}
132
133#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
134pub struct TaskGateProposal {
135    pub done: bool,
136    pub reason: String,
137}
138
139impl TaskGateProposal {
140    fn validate(&self) -> Result<(), TaskDataError> {
141        if self.reason.trim().is_empty() {
142            return Err(TaskDataError::InvalidInvariant(
143                "gate proposal reason cannot be empty".to_string(),
144            ));
145        }
146        Ok(())
147    }
148}
149
150#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
151pub struct GithubPr {
152    pub number: u32,
153    pub url: String,
154    /// The PR's current head commit, from `gh pr list --json headRefOid`. CI
155    /// evidence is authoritative only for this head. `None` on rows written
156    /// before head SHAs were recorded.
157    pub head_sha: Option<String>,
158}
159
160/// The state of a PR head's required checks, as last observed by a reconcile.
161/// Failure dominates: any failing required check makes the head `Failing`
162/// regardless of what else is still pending.
163#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
164#[serde(rename_all = "snake_case")]
165pub enum CiState {
166    Pending,
167    Passing,
168    Failing,
169}
170
171impl CiState {
172    pub fn as_str(self) -> &'static str {
173        match self {
174            Self::Pending => "pending",
175            Self::Passing => "passing",
176            Self::Failing => "failing",
177        }
178    }
179}
180
181/// The required check `lf pr land` greens itself, by clearing `scratch/`.
182///
183/// See [`CiCheck::land_time_precondition`] for what this class means and the one
184/// rule for admitting a name to it. `crate::ops::land::clear_scratch` is the step
185/// that resolves this check; `land_time_precondition_names_a_real_ci_job` pins
186/// this literal against `.github/workflows/ci.yml`.
187const LAND_TIME_PRECONDITION_CHECK: &str = "scratch-clear";
188
189/// One required check that is not passing, named so the `ci-fix` skill can
190/// resolve the exact failure from its logs.
191#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
192pub struct CiCheck {
193    pub name: String,
194    pub url: Option<String>,
195}
196
197impl CiCheck {
198    /// Whether this check asserts a *land-time precondition* — a condition
199    /// `lf pr land` establishes itself, so no repair turn can green it.
200    ///
201    /// This is the one class of required check a Task body cannot act on.
202    /// `scratch-clear` fails whenever `scratch/` holds anything but `.gitkeep`,
203    /// which is true of every PR carrying its own design doc — i.e. every Task PR
204    /// during first and loop, by construction — and
205    /// `crate::ops::land::clear_scratch` is what greens it, not a code change. A
206    /// body woken to "repair" it could only delete the Task's design artifact,
207    /// to green a check land greens anyway.
208    ///
209    /// A name belongs here only when an `lf pr land` step is what resolves it.
210    /// This is not a catalogue of CI jobs, and it is not a mute button for checks
211    /// that are merely hard to fix: a check anyone *could* fix by changing the
212    /// tree does not belong here, however annoying it is.
213    pub fn land_time_precondition(&self) -> bool {
214        self.name == LAND_TIME_PRECONDITION_CHECK
215    }
216}
217
218/// The required-check reading for one PR head. `head_sha` pins it: a reading is
219/// stale — and never wakes work — once the PR's head moves past it.
220#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
221pub struct CiObservation {
222    pub head_sha: String,
223    pub state: CiState,
224    pub failing_checks: Vec<CiCheck>,
225    pub observed_at: OffsetDateTime,
226}
227
228impl CiObservation {
229    /// The current failing required checks by name, sorted — the content half of
230    /// a [`CiIncident`]'s identity. Empty unless `state` is `Failing`.
231    pub fn failure_set(&self) -> Vec<String> {
232        let mut names: Vec<String> = self
233            .failing_checks
234            .iter()
235            .map(|check| check.name.clone())
236            .collect();
237        names.sort();
238        names.dedup();
239        names
240    }
241
242    /// Whether this reading makes a `ci-fix` wake *legal*: the current head is
243    /// failing a required check that a repair turn could actually act on.
244    ///
245    /// This asks only about legality. Whether a wake has already fired for this
246    /// exact failure is a separate question with a separate owner — the durable
247    /// CI incident, keyed on the incident identity and claimed by one Run. Those
248    /// two questions used to be conflated in one mutable JSON marker on this
249    /// struct, which meant the wake was deduplicated by a value re-derived on
250    /// every reconcile and committed only once a body had already been born.
251    ///
252    /// A head whose failures are *all* land-time preconditions is red and not
253    /// repairable ([`CiCheck::land_time_precondition`]): waking a body there
254    /// spends a full turn on work whose only successful action is destructive.
255    /// The reading still reports the failure — this refuses the wake, it does not
256    /// deny the red — so `lf ci` and `lf task status` are unchanged.
257    pub fn wake_legal(&self) -> bool {
258        if self.state != CiState::Failing {
259            return false;
260        }
261        // An unnamed failure is one we could not classify, not one proven
262        // harmless, so it still wakes: a filter that swallows unknown failures is
263        // a mute button. Only a non-empty set whose every member land resolves is
264        // provably not a repair.
265        self.failing_checks.is_empty()
266            || self
267                .failing_checks
268                .iter()
269                .any(|check| !check.land_time_precondition())
270    }
271
272    /// Whether this head is red *only* on land-time preconditions — failing, with
273    /// at least one named failure, and every named failure one that
274    /// [`CiCheck::land_time_precondition`] resolves at land.
275    ///
276    /// This is the dual of [`CiObservation::wake_legal`] within the failing
277    /// state: such a head holds nothing a Task body could repair (`lf pr land`
278    /// greens it by clearing `scratch/`), so it does not belong to CI repair.
279    /// The action model and Waves supervision read it to stop recommending a
280    /// doomed Resume or labelling settlement preparation as "fixing CI".
281    ///
282    /// False for a passing or pending head, and false the moment any failure is a
283    /// real leaf or an unclassified one — the same anti-mute-button rule as
284    /// `wake_legal`: an unnamed failure keeps the head owned by CI.
285    pub fn only_land_time_preconditions(&self) -> bool {
286        self.state == CiState::Failing
287            && !self.failing_checks.is_empty()
288            && self
289                .failing_checks
290                .iter()
291                .all(|check| check.land_time_precondition())
292    }
293}
294
295/// One failed CI head carried forward after the PR's current observation moves
296/// on. This is evidence about the recovery loop, never a wake queue.
297#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
298pub struct CiIncident {
299    pub identity: String,
300    pub task_id: TaskId,
301    pub pr_id: TaskPrId,
302    pub repo: String,
303    pub pr_number: u32,
304    pub failed_head_sha: String,
305    /// The first authoritative post-turn remote head that settlement observed to
306    /// differ from `failed_head_sha` — the head the repair body actually shipped
307    /// for this incident. `None` until a ci-fix body advances the head; written
308    /// once and never overwritten by a later push.
309    pub repaired_head_sha: Option<String>,
310    pub failure_set: Vec<String>,
311    pub provider_completed_at: Option<OffsetDateTime>,
312    pub poll_observed_at: Option<OffsetDateTime>,
313    pub webhook_received_at: Option<OffsetDateTime>,
314    /// Active or most recent Run selected to repair this incident. A successor
315    /// may replace it only after the prior Run lost execution authority.
316    pub claimed_run_id: Option<RunId>,
317    pub responded_at: Option<OffsetDateTime>,
318    pub green_at: Option<OffsetDateTime>,
319    pub merged_at: Option<OffsetDateTime>,
320    pub blocked_at: Option<OffsetDateTime>,
321    pub blocked_reason: Option<String>,
322    pub created_at: OffsetDateTime,
323    pub updated_at: OffsetDateTime,
324}
325
326/// The last attempt to refresh one persisted GitHub PR. This metadata lives on
327/// `TaskPr` beside the cached PR fields: it bounds repeated reads across `lf`
328/// processes without making GitHub the source of truth.
329#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
330pub struct GithubObservation {
331    pub checked_at: OffsetDateTime,
332    pub result: GithubObservationResult,
333}
334
335#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
336#[serde(tag = "state", rename_all = "snake_case")]
337pub enum GithubObservationResult {
338    Fresh,
339    Partial { reason: String },
340    Degraded { reason: String },
341}
342
343#[derive(Debug, Clone, Copy, PartialEq, Eq)]
344pub(crate) enum TaskPrRepairKind {
345    AvoidableRebaseAgent,
346    ManualGitRepair,
347}
348
349impl TaskPrRepairKind {
350    pub(crate) fn as_str(self) -> &'static str {
351        match self {
352            Self::AvoidableRebaseAgent => "avoidable_rebase_agent",
353            Self::ManualGitRepair => "manual_git_repair",
354        }
355    }
356}
357
358#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
359#[serde(rename_all = "snake_case")]
360#[non_exhaustive]
361pub enum PrPhase {
362    Working,
363    Publishing,
364    Open,
365    Merged,
366    Abandoned,
367}
368
369impl PrPhase {
370    pub fn as_str(self) -> &'static str {
371        match self {
372            Self::Working => "working",
373            Self::Publishing => "publishing",
374            Self::Open => "open",
375            Self::Merged => "merged",
376            Self::Abandoned => "abandoned",
377        }
378    }
379
380    pub fn is_active(self) -> bool {
381        matches!(self, Self::Working | Self::Publishing | Self::Open)
382    }
383
384    pub fn is_settled(self) -> bool {
385        matches!(self, Self::Merged | Self::Abandoned)
386    }
387}
388
389#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
390#[serde(rename_all = "snake_case")]
391pub enum AfterMerge {
392    ContinueTask,
393    CompleteTask,
394}
395
396impl AfterMerge {
397    pub fn as_str(self) -> &'static str {
398        match self {
399            Self::ContinueTask => "continue_task",
400            Self::CompleteTask => "complete_task",
401        }
402    }
403}
404
405impl FromStr for AfterMerge {
406    type Err = TaskDataError;
407
408    fn from_str(value: &str) -> Result<Self, Self::Err> {
409        match value {
410            "continue_task" => Ok(Self::ContinueTask),
411            "complete_task" => Ok(Self::CompleteTask),
412            _ => Err(TaskDataError::InvalidInvariant(format!(
413                "invalid after-merge disposition: {value}"
414            ))),
415        }
416    }
417}
418
419#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
420#[serde(rename_all = "snake_case")]
421pub enum PrMergeMode {
422    User,
423    Auto,
424}
425
426impl PrMergeMode {
427    pub fn as_str(self) -> &'static str {
428        match self {
429            Self::User => "user",
430            Self::Auto => "auto",
431        }
432    }
433}
434
435impl FromStr for PrMergeMode {
436    type Err = TaskDataError;
437
438    fn from_str(value: &str) -> Result<Self, Self::Err> {
439        match value {
440            "user" => Ok(Self::User),
441            "auto" => Ok(Self::Auto),
442            _ => Err(TaskDataError::InvalidInvariant(format!(
443                "invalid PR merge mode: {value}"
444            ))),
445        }
446    }
447}
448
449#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
450pub struct PrMergeRequest {
451    pub mode: PrMergeMode,
452    pub requested_at: OffsetDateTime,
453    pub head_sha: String,
454    pub after_merge: AfterMerge,
455    pub next_slug: Option<String>,
456}
457
458#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
459pub struct PrPublication {
460    pub requested_at: OffsetDateTime,
461    pub github: Option<GithubPr>,
462    pub merge: Option<PrMergeRequest>,
463}
464
465#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
466pub struct TaskPr {
467    pub id: TaskPrId,
468    pub task_id: TaskId,
469    pub sequence: u32,
470    pub slug: String,
471    pub branch: String,
472    pub base_commit: String,
473    /// Another Task's PR this worktree was placed on, or `None` when rooted on
474    /// the default branch. `base_commit` is that parent's exact fork commit; the
475    /// link clears after the parent merges and this PR collapses onto main.
476    pub parent_pr_id: Option<TaskPrId>,
477    pub publication: Option<PrPublication>,
478    pub merge_commit: Option<String>,
479    pub abandoned_at: Option<OffsetDateTime>,
480    /// The most recent required-check reading for this PR's open head. `None`
481    /// until the head has been observed; ignored once the head moves past
482    /// `CiObservation::head_sha`.
483    pub ci_observation: Option<CiObservation>,
484    /// Last GitHub refresh attempt. A fresh result coalesces reads briefly; a
485    /// degraded result opens a longer circuit while the durable PR fields stand.
486    pub github_observation: Option<GithubObservation>,
487    /// Id of the first-class Linear attachment linking this PR on its owning
488    /// issue. `None` until the PR is first published; carried forward so later
489    /// publishes update the same attachment in place.
490    pub linear_attachment_id: Option<String>,
491    /// Id of the loopflow-managed Linear comment carrying the PR URL and state.
492    /// Its presence switches the writeback from `commentCreate` to `commentUpdate`.
493    pub linear_comment_id: Option<String>,
494    /// `None` when the last Linear linkage writeback succeeded; the last error
495    /// string when it degraded. The GitHub publication still succeeded — this only
496    /// records that the Linear side is behind, and is cleared on the next success.
497    pub linear_link_error: Option<String>,
498    pub created_at: OffsetDateTime,
499    pub updated_at: OffsetDateTime,
500}
501
502impl TaskPr {
503    pub fn phase(&self) -> PrPhase {
504        if self.abandoned_at.is_some() {
505            PrPhase::Abandoned
506        } else if self.merge_commit.is_some() {
507            PrPhase::Merged
508        } else if self
509            .publication
510            .as_ref()
511            .is_some_and(|publication| publication.github.is_some())
512        {
513            PrPhase::Open
514        } else if self.publication.is_some() {
515            PrPhase::Publishing
516        } else {
517            PrPhase::Working
518        }
519    }
520
521    pub fn github(&self) -> Option<&GithubPr> {
522        self.publication
523            .as_ref()
524            .and_then(|publication| publication.github.as_ref())
525    }
526
527    /// The current head SHA of the open PR, when recorded.
528    pub fn head_sha(&self) -> Option<&str> {
529        self.github().and_then(|github| github.head_sha.as_deref())
530    }
531
532    /// The explicit merge request, only while it names the current PR head.
533    pub fn merge_request(&self) -> Option<&PrMergeRequest> {
534        let request = self.publication.as_ref()?.merge.as_ref()?;
535        (self.head_sha() == Some(request.head_sha.as_str())).then_some(request)
536    }
537
538    /// Settlement disposition. A PR merged without an explicit request safely
539    /// continues the Task; only a head-pinned request may complete it.
540    pub fn after_merge(&self) -> AfterMerge {
541        self.merge_request()
542            .map_or(AfterMerge::ContinueTask, |request| request.after_merge)
543    }
544
545    pub fn next_slug(&self) -> Option<&str> {
546        self.merge_request()
547            .and_then(|request| request.next_slug.as_deref())
548    }
549
550    /// The CI reading, but only while it still describes the PR's current head.
551    /// Once the head moves, the reading is stale and this returns `None` — the
552    /// same freshness rule that keeps stale failures from waking work.
553    pub fn fresh_ci(&self) -> Option<&CiObservation> {
554        let observation = self.ci_observation.as_ref()?;
555        match self.head_sha() {
556            Some(head) if head == observation.head_sha => Some(observation),
557            _ => None,
558        }
559    }
560
561    /// Whether the current published head satisfies its merge checks.
562    pub fn merge_checks_passed(&self) -> bool {
563        self.phase() == PrPhase::Open
564            && self
565                .fresh_ci()
566                .is_some_and(|observation| observation.state == CiState::Passing)
567    }
568
569    pub fn is_active(&self) -> bool {
570        self.phase().is_active()
571    }
572
573    pub fn is_settled(&self) -> bool {
574        self.phase().is_settled()
575    }
576
577    pub fn validate(&self) -> Result<(), TaskDataError> {
578        if self.sequence == 0 {
579            return Err(TaskDataError::InvalidInvariant(
580                "task pull request sequence starts at 1".to_string(),
581            ));
582        }
583        if self.slug.trim().is_empty() {
584            return Err(TaskDataError::InvalidInvariant(
585                "task PR slug cannot be empty".to_string(),
586            ));
587        }
588        if self.branch.trim().is_empty() || self.base_commit.trim().is_empty() {
589            return Err(TaskDataError::InvalidInvariant(
590                "task PR requires a branch and base commit".to_string(),
591            ));
592        }
593        if let Some(publication) = &self.publication {
594            if let Some(github) = &publication.github {
595                if github.number == 0 || github.url.trim().is_empty() {
596                    return Err(TaskDataError::InvalidInvariant(
597                        "GitHub PR number and URL cannot be empty".to_string(),
598                    ));
599                }
600            }
601            if let Some(request) = &publication.merge {
602                if request.next_slug.as_deref().is_some_and(|slug| {
603                    slug.split('-').any(|word| {
604                        word.is_empty()
605                            || !word
606                                .bytes()
607                                .all(|byte| byte.is_ascii_lowercase() || byte.is_ascii_digit())
608                    })
609                }) {
610                    return Err(TaskDataError::InvalidInvariant(
611                        "next branch slug must be lowercase kebab-case".to_string(),
612                    ));
613                }
614                if request.after_merge == AfterMerge::CompleteTask && request.next_slug.is_some() {
615                    return Err(TaskDataError::InvalidInvariant(
616                        "a completing pull request cannot name a next branch".to_string(),
617                    ));
618                }
619                let github = publication.github.as_ref().ok_or_else(|| {
620                    TaskDataError::InvalidInvariant(
621                        "merge request requires a GitHub PR".to_string(),
622                    )
623                })?;
624                if request.head_sha.trim().is_empty()
625                    || github.head_sha.as_deref() != Some(request.head_sha.as_str())
626                {
627                    return Err(TaskDataError::InvalidInvariant(
628                        "merge request must name the current GitHub PR head".to_string(),
629                    ));
630                }
631            }
632        }
633        if self.merge_commit.is_some() && self.github().is_none() {
634            return Err(TaskDataError::InvalidInvariant(
635                "merged PR requires a GitHub PR".to_string(),
636            ));
637        }
638        if self.merge_commit.is_some() && self.abandoned_at.is_some() {
639            return Err(TaskDataError::InvalidInvariant(
640                "a PR cannot be both merged and abandoned".to_string(),
641            ));
642        }
643        Ok(())
644    }
645}
646
647#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
648#[serde(rename_all = "snake_case")]
649#[non_exhaustive]
650pub enum PmWritebackOperation {
651    CompleteTask,
652}
653
654#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
655#[serde(tag = "state", rename_all = "snake_case")]
656pub enum PmWritebackState {
657    Current,
658    Pending {
659        operation: PmWritebackOperation,
660        error: String,
661    },
662}
663
664/// Freshness of the GitHub observation behind the Task's returned PR state.
665/// The durable attempt metadata lives on `TaskPr`; this derived view tells one
666/// caller whether it read GitHub, reused a recent reading, or opened a degraded
667/// circuit while preserving the cached PR fields.
668#[derive(Debug, Clone, PartialEq, Eq, Serialize, Default)]
669#[serde(tag = "freshness", rename_all = "snake_case")]
670pub enum Observation {
671    /// No remote read applies, as for an unpublished working PR.
672    #[default]
673    NotRequired,
674    /// GitHub answered during this reconcile.
675    Fresh { observed_at: OffsetDateTime },
676    /// A recent successful reading was reused without spending another request.
677    Cached { observed_at: OffsetDateTime },
678    /// A bounded read failed. The reason and retry boundary are durable, so
679    /// later local controls reuse the cached state without hammering GitHub.
680    Degraded {
681        reason: String,
682        cached_as_of: OffsetDateTime,
683        retry_at: OffsetDateTime,
684    },
685}
686
687#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
688pub struct Task {
689    pub id: TaskId,
690    /// Current planning facts from the PM system.
691    pub plan: TaskPlan,
692    pub pm_writeback: PmWritebackState,
693    /// Root ownership. Wave name and checkout are resolved from this id.
694    pub wave_id: WaveId,
695    /// Required runtime parent. Every Task reports through one durable Project
696    /// Work; its Wave retains root inspection and override authority.
697    pub project_id: ProjectId,
698    pub worktree: PathBuf,
699    pub workspace_slug: String,
700    /// Three pinned phase flows.
701    pub lifecycle: TaskLifecyclePlan,
702    /// Current phase entry. `phase_epoch` advances on every transition,
703    /// including Finally → Loop, so stale bodies cannot rewind the Task.
704    pub lifecycle_phase: TaskLifecyclePhase,
705    pub phase_epoch: u32,
706    pub phase_cursor: u32,
707    pub phase_iteration: u32,
708    /// Number of Gate entries attempted by this Task.
709    pub gate_cycle: u32,
710    /// Outcome proposed by Iterate and awaiting Gate approval.
711    pub gate_proposal: Option<TaskGateProposal>,
712    /// Provider/model selection for the next body generation. This is mutable
713    /// lease state, not Task identity.
714    pub agent: String,
715    /// Harness family for the next/current body generation.
716    pub provider: String,
717    /// Transcript handle reusable only by a compatible provider generation.
718    pub provider_session_id: Option<String>,
719    /// Set when abandonment is *requested*, not when it is applied. No launch
720    /// path may start a Run for Task Work carrying this.
721    pub abandon_intent: Option<AbandonIntent>,
722    pub created_at: OffsetDateTime,
723    pub updated_at: OffsetDateTime,
724    /// Per-command view derived from the active PR's durable observation cache.
725    #[serde(skip)]
726    pub observation: Observation,
727}
728
729impl Task {
730    /// Why a supervisor must not start another process generation, if it must not.
731    ///
732    /// Two intents bar an automatic restart:
733    ///
734    /// - abandonment has been *requested* — the runner has not consumed the
735    ///   command yet, but the decision is made;
736    /// - publication or merge was explicitly requested. A merely published PR
737    ///   remains ordinary Task continuity and does not bar the supervisor.
738    ///
739    /// The third was the 2026-07-14 W2-129 failure: `Open` is not terminal
740    /// and carries no live process, so it reads exactly like Work that
741    /// merely stopped. A wake therefore launched generation 2, which reopened
742    /// the flow at `task/clarify` and began re-doing work whose PR (#878) was
743    /// already awaiting a merge. An explicit merge request is not an invitation
744    /// to start over.
745    ///
746    /// A User may still `lf task resume` a submitted Task explicitly;
747    /// this bars the supervisor, not the operator.
748    pub fn supervisor_restart_bar(&self, active_pr: Option<&TaskPr>) -> Option<String> {
749        if let Some(bar) = self.terminal_or_abandon_bar() {
750            return Some(bar);
751        }
752        if let Some(pr) = active_pr {
753            match pr.phase() {
754                PrPhase::Publishing => return Some(self.publishing_bar()),
755                PrPhase::Open if pr.merge_request().is_some() => {
756                    return Some(self.open_pr_bar(pr));
757                }
758                PrPhase::Open => {}
759                PrPhase::Working | PrPhase::Merged | PrPhase::Abandoned => {}
760            }
761        }
762        None
763    }
764
765    /// The restart bar for an automated `ci-fix` wake. Identical to the supervisor
766    /// bar, except an `Open` PR is *permitted* when its current head carries a
767    /// failing required check ([`CiObservation::wake_legal`]). This is the one
768    /// automated path allowed to restart a submitted Task, and only on fresh
769    /// current-head failure evidence — never a blind wake over passing or pending
770    /// work, and never past the terminal, abandon, or publishing bars.
771    ///
772    /// The bar answers legality only. The incident's Run claim answers whether
773    /// this failure already owns execution.
774    pub fn ci_fix_restart_bar(&self, active_pr: Option<&TaskPr>) -> Option<String> {
775        if let Some(bar) = self.terminal_or_abandon_bar() {
776            return Some(bar);
777        }
778        if let Some(pr) = active_pr {
779            match pr.phase() {
780                PrPhase::Publishing => return Some(self.publishing_bar()),
781                // An open PR restarts only on a current-head required-check
782                // failure; otherwise it stays barred exactly as the supervisor
783                // bar leaves it.
784                PrPhase::Open
785                    if pr.merge_request().is_some()
786                        && !pr.fresh_ci().is_some_and(CiObservation::wake_legal) =>
787                {
788                    return Some(self.open_pr_bar(pr));
789                }
790                PrPhase::Open | PrPhase::Working | PrPhase::Merged | PrPhase::Abandoned => {}
791            }
792        }
793        None
794    }
795
796    /// The product-level bar that holds for every automatic restart intent.
797    /// Terminal execution state belongs to Work and is checked by the caller.
798    pub(crate) fn terminal_or_abandon_bar(&self) -> Option<String> {
799        if let Some(intent) = &self.abandon_intent {
800            return Some(format!(
801                "Task {} is being abandoned: {}",
802                self.plan.identifier, intent.reason
803            ));
804        }
805        None
806    }
807
808    fn publishing_bar(&self) -> String {
809        format!(
810            "Task {} requested PR publication but has no GitHub PR record; \
811             resume it explicitly with `lf task resume {}` to retry publication",
812            self.plan.identifier, self.plan.identifier,
813        )
814    }
815
816    /// The open-PR restart refusal. Only a *supervisor* (`supervisor_restart_bar`
817    /// / a non-wake-legal `ci_fix_restart_bar`) ever reads this — an operator
818    /// resume takes the abandon-only `ExplicitResume` bar. The text names the
819    /// explicit settlement owner instead of recommending `lf task resume`, which
820    /// a supervisor re-running would only self-loop on.
821    fn open_pr_bar(&self, pr: &TaskPr) -> String {
822        let number = pr.github().expect("open Task PR passed validation").number;
823        let request = pr
824            .merge_request()
825            .expect("open PR restart bar requires a current merge request");
826        let short = request.head_sha.chars().take(12).collect::<String>();
827        match request.mode {
828            PrMergeMode::User => format!(
829                "Task {} requested a user merge of pull request #{} at head {}. \
830                 The supervisor will not restart it until that explicit merge \
831                 request settles or the head changes.",
832                self.plan.identifier, number, short,
833            ),
834            PrMergeMode::Auto => format!(
835                "Task {} requested GitHub auto-merge of pull request #{} at head {}. \
836                 The supervisor will not restart it until that explicit merge \
837                 request settles or the head changes.",
838                self.plan.identifier, number, short,
839            ),
840        }
841    }
842
843    pub fn validate(&self) -> Result<(), TaskDataError> {
844        if self.workspace_slug.trim().is_empty() {
845            return Err(TaskDataError::InvalidInvariant(format!(
846                "Task {} requires a workspace slug",
847                self.id
848            )));
849        }
850        self.lifecycle.validate()?;
851        if self.phase_epoch == 0 {
852            return Err(TaskDataError::InvalidInvariant(
853                "Task lifecycle phase epoch must be positive".to_string(),
854            ));
855        }
856        if self.lifecycle_phase == TaskLifecyclePhase::Finally && self.gate_proposal.is_none() {
857            return Err(TaskDataError::InvalidInvariant(
858                "Task finally phase requires a proposed outcome".to_string(),
859            ));
860        }
861        if self.lifecycle_phase != TaskLifecyclePhase::Finally && self.gate_proposal.is_some() {
862            return Err(TaskDataError::InvalidInvariant(
863                "Task gate proposal is valid only during finally phase".to_string(),
864            ));
865        }
866        if let Some(proposal) = &self.gate_proposal {
867            proposal.validate()?;
868        }
869        if matches!(self.pm_writeback, PmWritebackState::Pending { .. }) && self.gate_cycle == 0 {
870            return Err(TaskDataError::InvalidInvariant(
871                "pending PM completion requires an active gate cycle".to_string(),
872            ));
873        }
874        Ok(())
875    }
876
877    pub fn phase_plan(&self) -> &TaskPhasePlan {
878        self.lifecycle.phase(self.lifecycle_phase)
879    }
880
881    pub fn lifecycle_cycle(&self) -> u32 {
882        match self.lifecycle_phase {
883            TaskLifecyclePhase::First => 0,
884            TaskLifecyclePhase::Loop => self.gate_cycle + 1,
885            TaskLifecyclePhase::Finally => self.gate_cycle,
886        }
887    }
888
889    pub fn enter_loop(&mut self) -> Result<(), TaskDataError> {
890        if self.lifecycle_phase != TaskLifecyclePhase::First
891            && self.lifecycle_phase != TaskLifecyclePhase::Finally
892        {
893            return Err(TaskDataError::InvalidInvariant(
894                "only first or finally may enter loop".to_string(),
895            ));
896        }
897        self.lifecycle_phase = TaskLifecyclePhase::Loop;
898        self.phase_epoch += 1;
899        self.phase_cursor = 0;
900        self.phase_iteration = 0;
901        self.gate_proposal = None;
902        self.updated_at = OffsetDateTime::now_utc();
903        Ok(())
904    }
905
906    pub fn enter_finally(&mut self, proposal: TaskGateProposal) -> Result<(), TaskDataError> {
907        if self.lifecycle_phase != TaskLifecyclePhase::Loop {
908            return Err(TaskDataError::InvalidInvariant(
909                "only loop may enter finally".to_string(),
910            ));
911        }
912        proposal.validate()?;
913        self.lifecycle_phase = TaskLifecyclePhase::Finally;
914        self.phase_epoch += 1;
915        self.phase_cursor = 0;
916        self.phase_iteration = 0;
917        self.gate_cycle += 1;
918        self.gate_proposal = Some(proposal);
919        self.updated_at = OffsetDateTime::now_utc();
920        Ok(())
921    }
922
923    pub fn approved_gate_proposal(&self) -> Result<TaskGateProposal, TaskDataError> {
924        if self.lifecycle_phase != TaskLifecyclePhase::Finally {
925            return Err(TaskDataError::InvalidInvariant(
926                "only finally may approve a proposed outcome".to_string(),
927            ));
928        }
929        self.gate_proposal.clone().ok_or_else(|| {
930            TaskDataError::InvalidInvariant("Task gate has no proposed outcome".to_string())
931        })
932    }
933}
934
935#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
936#[serde(tag = "kind", rename_all = "snake_case")]
937pub enum TaskEventKind {
938    Started,
939    BodyHandedOff {
940        handoff: crate::child::ChildBodyHandoff,
941    },
942    Progress {
943        summary: String,
944    },
945    PrStarted {
946        pr_id: TaskPrId,
947        sequence: u32,
948        branch: String,
949        base_commit: String,
950    },
951    PrOpened {
952        pr_id: TaskPrId,
953        sequence: u32,
954        number: u32,
955        url: String,
956    },
957    PrMerged {
958        pr_id: TaskPrId,
959        sequence: u32,
960        number: u32,
961        url: String,
962        merge_commit: String,
963    },
964    Completed {
965        summary: String,
966    },
967    Failed {
968        error: String,
969        resumable: bool,
970    },
971}
972
973impl TaskEventKind {
974    /// Whether the event crosses the required Task → Project boundary.
975    pub fn is_project_observable(&self) -> bool {
976        !matches!(self, Self::Started | Self::Progress { .. })
977    }
978
979    /// Whether a Project-observable Task event also belongs in the root Wave.
980    /// This currently mirrors the Project boundary; the server-topology design
981    /// must decide whether the duplicate delivery remains necessary.
982    pub fn is_root_wave_observable(&self) -> bool {
983        self.is_project_observable()
984    }
985}
986
987#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
988pub struct TaskEvent {
989    pub id: i64,
990    pub task_id: TaskId,
991    pub kind: TaskEventKind,
992    pub created_at: OffsetDateTime,
993}
994
995#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
996pub struct TaskObservation {
997    pub task_id: TaskId,
998    pub issue_identifier: String,
999    pub event_id: i64,
1000    pub event: TaskEventKind,
1001}
1002
1003impl TaskObservation {
1004    pub fn inbox_id(&self) -> String {
1005        format!("task-{}-{}", self.task_id, self.event_id)
1006    }
1007
1008    pub fn prompt(&self) -> String {
1009        let payload = serde_json::to_string(&self.event)
1010            .expect("Task observation always serializes to structured JSON");
1011        format!(
1012            "<task_observation task_id=\"{}\" issue=\"{}\" event_id=\"{}\">\n{}\n</task_observation>",
1013            self.task_id, self.issue_identifier, self.event_id, payload
1014        )
1015    }
1016}
1017
1018/// The durable cursor for streaming human Linear edits into one Task.
1019/// It is the exactly-once ledger — what issue revision and comments have already
1020/// become Task direction — plus the health of the last observation, so
1021/// `lf task status` can show stale reads and their degraded reason.
1022#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1023pub struct TaskLinearObservation {
1024    pub task_id: TaskId,
1025    /// Linear issue `updatedAt` last folded in. Monotonic: a read whose revision
1026    /// is older is dropped as a stale/out-of-order response.
1027    pub last_revision: String,
1028    /// Content basis for the next title/description diff.
1029    pub last_title: String,
1030    pub last_description: String,
1031    pub last_success_at: OffsetDateTime,
1032    /// `Some` after a failed observation (auth/quota/network); `None` is healthy.
1033    pub degraded_reason: Option<String>,
1034    pub updated_at: OffsetDateTime,
1035}
1036
1037/// One Linear observation, ready to persist atomically as Task direction. The
1038/// directive is applied only if the stored title/description still differ
1039/// (compare-and-set), and each follow-up becomes a command only on its first
1040/// entry into the ledger — so overlapping polls, restarts, and out-of-order
1041/// responses never duplicate direction.
1042#[derive(Debug, Clone)]
1043pub struct LinearObservationApply {
1044    pub task_id: TaskId,
1045    pub revision: String,
1046    pub title: String,
1047    pub description: String,
1048    pub observed_at: OffsetDateTime,
1049    /// A title/description edit to persist as one authored Steer.
1050    pub content_steer: Option<String>,
1051    /// Human comments observed this pass, oldest first.
1052    pub follow_ups: Vec<LinearFollowUp>,
1053}
1054
1055#[derive(Debug, Clone)]
1056pub struct LinearFollowUp {
1057    pub comment_id: String,
1058    pub text: String,
1059}
1060
1061/// What one [`LinearObservationApply`] actually wrote — enough for the caller to
1062/// report receipts without re-reading the store.
1063#[derive(Debug, Clone, PartialEq, Eq)]
1064pub struct LinearObservationOutcome {
1065    /// The Task had no cursor yet: this observation seeded the baseline and
1066    /// emitted no direction (existing comments are marked seen, not replayed).
1067    pub baselined: bool,
1068    pub content_steer_applied: bool,
1069    pub follow_ups_created: Vec<crate::durable::SteerId>,
1070}
1071
1072#[cfg(test)]
1073mod tests {
1074    use super::{
1075        AfterMerge, GithubPr, PmWritebackOperation, PmWritebackState, PrPhase, PrPublication, Task,
1076        TaskGateProposal, TaskId, TaskLifecyclePhase, TaskLifecyclePlan, TaskObservation, TaskPr,
1077        TaskPrId,
1078    };
1079    use crate::planning::{LinearIssueId, TaskPlan};
1080
1081    fn task() -> Task {
1082        let now = time::OffsetDateTime::now_utc();
1083        Task {
1084            id: TaskId::new(),
1085            plan: TaskPlan {
1086                id: LinearIssueId::new("issue-1").unwrap(),
1087                identifier: "INF-123".to_string(),
1088                title: "Ship it".to_string(),
1089                description: String::new(),
1090                pm_snapshot_synced_at: 1,
1091            },
1092            pm_writeback: PmWritebackState::Current,
1093            wave_id: crate::id::WaveId::new(),
1094            project_id: crate::project::ProjectId::new(),
1095            worktree: "/tmp/task".into(),
1096            workspace_slug: "ship-it".to_string(),
1097            lifecycle: TaskLifecyclePlan::defaults(),
1098            lifecycle_phase: TaskLifecyclePhase::Loop,
1099            phase_epoch: 1,
1100            phase_cursor: 0,
1101            phase_iteration: 0,
1102            gate_cycle: 0,
1103            gate_proposal: None,
1104            agent: "codex".to_string(),
1105            provider: "codex".to_string(),
1106            provider_session_id: None,
1107            abandon_intent: None,
1108            created_at: now,
1109            updated_at: now,
1110            observation: crate::task::Observation::NotRequired,
1111        }
1112    }
1113
1114    #[test]
1115    fn task_ids_are_prefixed_and_round_trip() {
1116        let task = TaskId::new();
1117        assert_eq!(TaskId::parse(task.as_str()).unwrap(), task);
1118        let pr = TaskPrId::new();
1119        assert_eq!(TaskPrId::parse(pr.as_str()).unwrap(), pr);
1120    }
1121
1122    #[test]
1123    fn task_observation_has_a_stable_structured_inbox_identity() {
1124        let observation = TaskObservation {
1125            task_id: TaskId::from_raw("ts_example"),
1126            issue_identifier: "INF-123".to_string(),
1127            event_id: 42,
1128            event: super::TaskEventKind::Failed {
1129                error: "provider stopped".to_string(),
1130                resumable: true,
1131            },
1132        };
1133
1134        assert_eq!(observation.inbox_id(), "task-ts_example-42");
1135        assert!(observation.prompt().contains("<task_observation"));
1136        assert!(observation.prompt().contains("\"kind\":\"failed\""));
1137    }
1138
1139    #[test]
1140    fn pending_pm_writeback_has_a_stable_json_shape() {
1141        let state = PmWritebackState::Pending {
1142            operation: PmWritebackOperation::CompleteTask,
1143            error: "offline".to_string(),
1144        };
1145
1146        assert_eq!(
1147            serde_json::to_value(state).unwrap(),
1148            serde_json::json!({
1149                "state": "pending",
1150                "operation": "complete_task",
1151                "error": "offline"
1152            })
1153        );
1154    }
1155
1156    #[test]
1157    fn pr_phase_is_derived_from_durable_evidence() {
1158        let now = time::OffsetDateTime::now_utc();
1159        let mut pr = TaskPr {
1160            id: TaskPrId::new(),
1161            task_id: TaskId::new(),
1162            sequence: 1,
1163            slug: "ship-it".to_string(),
1164            branch: "jack/ship-it".to_string(),
1165            base_commit: "abc".to_string(),
1166            parent_pr_id: None,
1167            publication: None,
1168            merge_commit: None,
1169            abandoned_at: None,
1170            created_at: now,
1171            updated_at: now,
1172            ci_observation: None,
1173            github_observation: None,
1174            linear_attachment_id: None,
1175            linear_comment_id: None,
1176            linear_link_error: None,
1177        };
1178        assert_eq!(pr.phase(), PrPhase::Working);
1179
1180        pr.publication = Some(PrPublication {
1181            requested_at: now,
1182            github: None,
1183            merge: None,
1184        });
1185        assert_eq!(pr.phase(), PrPhase::Publishing);
1186
1187        pr.publication.as_mut().unwrap().github = Some(GithubPr {
1188            number: 872,
1189            url: "https://github.com/loopflowstudio/loopflow/pull/872".to_string(),
1190            head_sha: None,
1191        });
1192        assert_eq!(pr.phase(), PrPhase::Open);
1193
1194        pr.merge_commit = Some("def".to_string());
1195        assert_eq!(pr.phase(), PrPhase::Merged);
1196        assert!(pr.validate().is_ok());
1197
1198        pr.abandoned_at = Some(now);
1199        assert_eq!(pr.phase(), PrPhase::Abandoned);
1200        assert!(pr.validate().is_err());
1201    }
1202
1203    #[test]
1204    fn merge_request_contains_its_disposition() {
1205        let now = time::OffsetDateTime::now_utc();
1206        let mut pr = TaskPr {
1207            id: TaskPrId::new(),
1208            task_id: TaskId::new(),
1209            sequence: 1,
1210            slug: "ship-it".to_string(),
1211            branch: "jack/ship-it".to_string(),
1212            base_commit: "abc".to_string(),
1213            parent_pr_id: None,
1214            publication: Some(PrPublication {
1215                requested_at: now,
1216                github: Some(GithubPr {
1217                    number: 872,
1218                    url: "https://github.com/loopflowstudio/loopflow/pull/872".to_string(),
1219                    head_sha: Some("head".to_string()),
1220                }),
1221                merge: Some(super::PrMergeRequest {
1222                    mode: super::PrMergeMode::User,
1223                    requested_at: now,
1224                    head_sha: "head".to_string(),
1225                    after_merge: AfterMerge::ContinueTask,
1226                    next_slug: Some("released_upgrade".to_string()),
1227                }),
1228            }),
1229            merge_commit: None,
1230            abandoned_at: None,
1231            created_at: now,
1232            updated_at: now,
1233            ci_observation: None,
1234            github_observation: None,
1235            linear_attachment_id: None,
1236            linear_comment_id: None,
1237            linear_link_error: None,
1238        };
1239        assert!(pr.validate().is_err());
1240
1241        let merge = pr.publication.as_mut().unwrap().merge.as_mut().unwrap();
1242        merge.next_slug = Some("released-upgrade".to_string());
1243        assert!(pr.validate().is_ok());
1244
1245        pr.publication
1246            .as_mut()
1247            .unwrap()
1248            .merge
1249            .as_mut()
1250            .unwrap()
1251            .after_merge = AfterMerge::CompleteTask;
1252        assert!(pr.validate().is_err());
1253    }
1254
1255    fn open_pr(head_sha: &str, observation: Option<super::CiObservation>) -> TaskPr {
1256        let now = time::OffsetDateTime::now_utc();
1257        TaskPr {
1258            id: TaskPrId::new(),
1259            task_id: TaskId::new(),
1260            sequence: 1,
1261            slug: "ship-it".to_string(),
1262            branch: "jack/ship-it".to_string(),
1263            base_commit: "abc".to_string(),
1264            parent_pr_id: None,
1265            publication: Some(PrPublication {
1266                requested_at: now,
1267                github: Some(GithubPr {
1268                    number: 900,
1269                    url: "https://github.com/loopflow/loopflow/pull/900".to_string(),
1270                    head_sha: Some(head_sha.to_string()),
1271                }),
1272                merge: None,
1273            }),
1274            merge_commit: None,
1275            abandoned_at: None,
1276            ci_observation: observation,
1277            github_observation: None,
1278            linear_attachment_id: None,
1279            linear_comment_id: None,
1280            linear_link_error: None,
1281            created_at: now,
1282            updated_at: now,
1283        }
1284    }
1285
1286    #[test]
1287    fn fresh_ci_ignores_a_reading_for_a_past_head() {
1288        let now = time::OffsetDateTime::now_utc();
1289        let observation = super::CiObservation {
1290            head_sha: "old-head".to_string(),
1291            state: super::CiState::Failing,
1292            failing_checks: vec![super::CiCheck {
1293                name: "build".to_string(),
1294                url: None,
1295            }],
1296            observed_at: now,
1297        };
1298        // Reading matches the current head: fresh.
1299        let current = open_pr("old-head", Some(observation.clone()));
1300        assert_eq!(
1301            current.fresh_ci().map(|ci| ci.state),
1302            Some(super::CiState::Failing)
1303        );
1304        // Head has moved on: the stale reading never surfaces (and never wakes work).
1305        let moved = open_pr("new-head", Some(observation));
1306        assert!(moved.fresh_ci().is_none());
1307    }
1308
1309    #[test]
1310    fn merge_checks_require_current_head_passing_checks() {
1311        let observation = |head: &str, state| super::CiObservation {
1312            head_sha: head.to_string(),
1313            state,
1314            failing_checks: Vec::new(),
1315            observed_at: time::OffsetDateTime::now_utc(),
1316        };
1317
1318        assert!(open_pr(
1319            "current",
1320            Some(observation("current", super::CiState::Passing))
1321        )
1322        .merge_checks_passed());
1323        assert!(!open_pr(
1324            "current",
1325            Some(observation("current", super::CiState::Pending))
1326        )
1327        .merge_checks_passed());
1328        assert!(
1329            !open_pr("current", Some(observation("old", super::CiState::Passing)))
1330                .merge_checks_passed()
1331        );
1332        assert!(!open_pr("current", None).merge_checks_passed());
1333    }
1334
1335    fn failing(head: &str, checks: &[&str]) -> super::CiObservation {
1336        super::CiObservation {
1337            head_sha: head.to_string(),
1338            state: super::CiState::Failing,
1339            failing_checks: checks
1340                .iter()
1341                .map(|name| super::CiCheck {
1342                    name: name.to_string(),
1343                    url: None,
1344                })
1345                .collect(),
1346            observed_at: time::OffsetDateTime::now_utc(),
1347        }
1348    }
1349
1350    fn with_merge_request(mut pr: TaskPr, mode: super::PrMergeMode) -> TaskPr {
1351        let head_sha = pr.head_sha().expect("test PR has a head").to_string();
1352        pr.publication.as_mut().unwrap().merge = Some(super::PrMergeRequest {
1353            mode,
1354            requested_at: time::OffsetDateTime::now_utc(),
1355            head_sha,
1356            after_merge: AfterMerge::ContinueTask,
1357            next_slug: None,
1358        });
1359        pr
1360    }
1361
1362    /// The observation answers legality — is this head failing *now* — and nothing
1363    /// else. Whether a wake already fired for a failure is the command ledger's
1364    /// question, keyed on the incident identity; it used to be a mutable marker on
1365    /// this struct, which is what let a repeat poll race a body's birth.
1366    #[test]
1367    fn ci_wake_legality_follows_only_the_current_reading() {
1368        let obs = failing("h1", &["build", "lint"]);
1369        assert!(obs.wake_legal());
1370
1371        // The failing set is order-independent and deduplicated: it is the content
1372        // half of the incident identity, so two readings of one failure must hash
1373        // the same however GitHub ordered them.
1374        assert_eq!(
1375            obs.failure_set(),
1376            vec!["build".to_string(), "lint".to_string()]
1377        );
1378        assert_eq!(
1379            failing("h1", &["lint", "build"]).failure_set(),
1380            obs.failure_set()
1381        );
1382        assert_eq!(
1383            failing("h1", &["build", "build"]).failure_set(),
1384            vec!["build".to_string()]
1385        );
1386
1387        // A passing or pending reading is never legal to wake on.
1388        let mut green = obs.clone();
1389        green.state = super::CiState::Passing;
1390        assert!(!green.wake_legal());
1391        let mut pending = obs.clone();
1392        pending.state = super::CiState::Pending;
1393        assert!(!pending.wake_legal());
1394    }
1395
1396    /// A head red *only* on a land-time precondition is not a repair a body can
1397    /// perform: `lf pr land` clears `scratch/`, and the only action a woken body
1398    /// could take is deleting the Task's design doc.
1399    ///
1400    /// The direction that matters is the second half. Suppression fires only when
1401    /// every named failure is land-resolved — a real leaf alongside it still
1402    /// arms, and an *unnamed* failure still arms, because a filter that swallows
1403    /// failures it cannot classify is a mute button rather than a classifier.
1404    #[test]
1405    fn a_head_red_only_on_a_land_time_precondition_is_not_wakeable() {
1406        assert!(!failing("h1", &["scratch-clear"]).wake_legal());
1407
1408        // A real leaf alongside it is still a repair worth waking for.
1409        assert!(failing("h1", &["scratch-clear", "rust-test"]).wake_legal());
1410        assert!(failing("h1", &["rust-test"]).wake_legal());
1411
1412        // Failing with nothing named: unclassified, not proven harmless.
1413        assert!(failing("h1", &[]).wake_legal());
1414
1415        // The reading stays honest — this refuses the wake, it does not deny the
1416        // red. Status and `lf ci` still name the failure.
1417        let obs = failing("h1", &["scratch-clear"]);
1418        assert_eq!(obs.state, super::CiState::Failing);
1419        assert_eq!(obs.failure_set(), vec!["scratch-clear".to_string()]);
1420    }
1421
1422    /// The land-resolved predicate is the exact dual of `wake_legal` within the
1423    /// failing state, so the two must never both refuse the same head.
1424    #[test]
1425    fn only_land_time_preconditions_is_the_dual_of_wake_legal() {
1426        // Red only on scratch-clear: land-resolved, and no wake.
1427        let scratch = failing("h1", &["scratch-clear"]);
1428        assert!(scratch.only_land_time_preconditions());
1429        assert!(!scratch.wake_legal());
1430
1431        // A real leaf beside it (or alone): not land-resolved, wake arms.
1432        for obs in [
1433            failing("h1", &["scratch-clear", "rust-test"]),
1434            failing("h1", &["rust-test"]),
1435        ] {
1436            assert!(!obs.only_land_time_preconditions());
1437            assert!(obs.wake_legal());
1438        }
1439
1440        // Unclassified (empty) failure: not "only preconditions", still wakes.
1441        let empty = failing("h1", &[]);
1442        assert!(!empty.only_land_time_preconditions());
1443        assert!(empty.wake_legal());
1444
1445        // Passing/pending is never "only preconditions".
1446        let mut green = scratch.clone();
1447        green.state = super::CiState::Passing;
1448        assert!(!green.only_land_time_preconditions());
1449        let mut pending = scratch.clone();
1450        pending.state = super::CiState::Pending;
1451        assert!(!pending.only_land_time_preconditions());
1452    }
1453
1454    /// The one literal, pinned against the workflow that defines it. A rename of
1455    /// the CI job turns this red *at the const* instead of silently re-arming
1456    /// wakes in production — the anti-rot guard for a name that cannot be derived.
1457    #[test]
1458    fn land_time_precondition_names_a_real_ci_job() {
1459        let workflow =
1460            std::path::Path::new(env!("CARGO_MANIFEST_DIR")).join("../../.github/workflows/ci.yml");
1461        let yaml = std::fs::read_to_string(&workflow)
1462            .unwrap_or_else(|e| panic!("read {}: {e}", workflow.display()));
1463        let job = format!("\n  {}:", super::LAND_TIME_PRECONDITION_CHECK);
1464        assert!(
1465            yaml.contains(&job),
1466            "no job named `{}` in {} — the const no longer names a real check, so \
1467             every design-carrying PR is arming ci-fix wakes again",
1468            super::LAND_TIME_PRECONDITION_CHECK,
1469            workflow.display()
1470        );
1471    }
1472
1473    #[test]
1474    fn ci_fix_restart_bar_permits_only_a_failing_open_pr_wake() {
1475        let task = task(); // status Waiting, no abandon intent
1476
1477        // Publication alone is not a restart bar.
1478        let published = open_pr("h1", None);
1479        assert!(task.supervisor_restart_bar(Some(&published)).is_none());
1480        assert!(task.ci_fix_restart_bar(Some(&published)).is_none());
1481
1482        // Open PR, fresh failing head: the ci-fix wake is permitted where the plain
1483        // supervisor restart stays barred.
1484        let legal = with_merge_request(
1485            open_pr("h1", Some(failing("h1", &["build"]))),
1486            super::PrMergeMode::Auto,
1487        );
1488        assert!(task.supervisor_restart_bar(Some(&legal)).is_some());
1489        assert!(task.ci_fix_restart_bar(Some(&legal)).is_none());
1490
1491        // Passing head → not legal → barred.
1492        let mut green_obs = failing("h1", &[]);
1493        green_obs.state = super::CiState::Passing;
1494        let green = with_merge_request(open_pr("h1", Some(green_obs)), super::PrMergeMode::Auto);
1495        assert!(task.ci_fix_restart_bar(Some(&green)).is_some());
1496
1497        // Stale reading (observation head != PR head) → fresh_ci None → barred.
1498        let stale = with_merge_request(
1499            open_pr("h2", Some(failing("h1", &["build"]))),
1500            super::PrMergeMode::Auto,
1501        );
1502        assert!(task.ci_fix_restart_bar(Some(&stale)).is_some());
1503
1504        // Red only on a land-time precondition → no repair exists → barred, and
1505        // the automated restart is the one path that could have overridden the
1506        // open-PR bar. A real leaf beside it still permits the wake.
1507        let land_only = with_merge_request(
1508            open_pr("h1", Some(failing("h1", &["scratch-clear"]))),
1509            super::PrMergeMode::Auto,
1510        );
1511        assert!(task.ci_fix_restart_bar(Some(&land_only)).is_some());
1512        let mixed = with_merge_request(
1513            open_pr("h1", Some(failing("h1", &["scratch-clear", "rust-test"]))),
1514            super::PrMergeMode::Auto,
1515        );
1516        assert!(task.ci_fix_restart_bar(Some(&mixed)).is_none());
1517
1518        // The bar does not deduplicate. A head that already woke a body still reads
1519        // as legal here — refusing the second launch is the ledger's job, and
1520        // asking the question twice is what let the two answers drift.
1521        assert!(task.ci_fix_restart_bar(Some(&legal)).is_none());
1522    }
1523
1524    #[test]
1525    fn task_rejects_impossible_lifecycle_and_writeback_state() {
1526        let mut task = task();
1527        task.pm_writeback = PmWritebackState::Pending {
1528            operation: PmWritebackOperation::CompleteTask,
1529            error: "too early".to_string(),
1530        };
1531        assert!(task.validate().is_err());
1532
1533        task.gate_cycle = 1;
1534        assert!(task.validate().is_ok());
1535
1536        task.lifecycle.loop_.flow.clear();
1537        assert!(task.validate().is_err());
1538    }
1539
1540    #[test]
1541    fn task_lifecycle_repeats_loop_and_finally_until_approval() {
1542        let mut task = task();
1543        task.lifecycle_phase = TaskLifecyclePhase::First;
1544
1545        assert_eq!(task.lifecycle_cycle(), 0);
1546        task.enter_loop().unwrap();
1547        assert_eq!(task.lifecycle_phase, TaskLifecyclePhase::Loop);
1548        assert_eq!(task.lifecycle_cycle(), 1);
1549        assert_eq!(task.phase_epoch, 2);
1550
1551        let proposal = TaskGateProposal {
1552            done: false,
1553            reason: "iteration needs another pass".to_string(),
1554        };
1555        task.phase_cursor = 2;
1556        task.phase_iteration = 3;
1557        task.enter_finally(proposal.clone()).unwrap();
1558        assert_eq!(task.lifecycle_phase, TaskLifecyclePhase::Finally);
1559        assert_eq!(task.lifecycle_cycle(), 1);
1560        assert_eq!(task.gate_cycle, 1);
1561        assert_eq!(task.approved_gate_proposal().unwrap(), proposal);
1562        assert_eq!((task.phase_cursor, task.phase_iteration), (0, 0));
1563
1564        task.enter_loop().unwrap();
1565        assert_eq!(task.lifecycle_phase, TaskLifecyclePhase::Loop);
1566        assert_eq!(task.lifecycle_cycle(), 2);
1567        assert_eq!(task.gate_proposal, None);
1568        assert_eq!(task.phase_epoch, 4);
1569    }
1570
1571    #[test]
1572    fn task_lifecycle_uses_public_names_without_rewriting_storage() {
1573        for (phase, public, stored) in [
1574            (TaskLifecyclePhase::First, "first", "kickoff"),
1575            (TaskLifecyclePhase::Loop, "loop", "iterate"),
1576            (TaskLifecyclePhase::Finally, "finally", "gate"),
1577        ] {
1578            assert_eq!(phase.as_str(), public);
1579            assert_eq!(phase.storage_str(), stored);
1580            assert_eq!(TaskLifecyclePhase::from_storage_str(stored).unwrap(), phase);
1581        }
1582    }
1583
1584    #[test]
1585    fn standard_lifecycle_pins_each_phase_flow() {
1586        let plan = TaskLifecyclePlan::standard("task-design", "code", "ship");
1587        assert_eq!(plan.first.flow, "task-design");
1588        assert_eq!(plan.loop_.flow, "code");
1589        assert_eq!(plan.finally.flow, "ship");
1590    }
1591}