1use 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 pub head_sha: Option<String>,
158}
159
160#[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
181const LAND_TIME_PRECONDITION_CHECK: &str = "scratch-clear";
188
189#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
192pub struct CiCheck {
193 pub name: String,
194 pub url: Option<String>,
195}
196
197impl CiCheck {
198 pub fn land_time_precondition(&self) -> bool {
214 self.name == LAND_TIME_PRECONDITION_CHECK
215 }
216}
217
218#[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 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 pub fn wake_legal(&self) -> bool {
258 if self.state != CiState::Failing {
259 return false;
260 }
261 self.failing_checks.is_empty()
266 || self
267 .failing_checks
268 .iter()
269 .any(|check| !check.land_time_precondition())
270 }
271
272 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#[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 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 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#[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 pub parent_pr_id: Option<TaskPrId>,
477 pub publication: Option<PrPublication>,
478 pub merge_commit: Option<String>,
479 pub abandoned_at: Option<OffsetDateTime>,
480 pub ci_observation: Option<CiObservation>,
484 pub github_observation: Option<GithubObservation>,
487 pub linear_attachment_id: Option<String>,
491 pub linear_comment_id: Option<String>,
494 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 pub fn head_sha(&self) -> Option<&str> {
529 self.github().and_then(|github| github.head_sha.as_deref())
530 }
531
532 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 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 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 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#[derive(Debug, Clone, PartialEq, Eq, Serialize, Default)]
669#[serde(tag = "freshness", rename_all = "snake_case")]
670pub enum Observation {
671 #[default]
673 NotRequired,
674 Fresh { observed_at: OffsetDateTime },
676 Cached { observed_at: OffsetDateTime },
678 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 pub plan: TaskPlan,
692 pub pm_writeback: PmWritebackState,
693 pub wave_id: WaveId,
695 pub project_id: ProjectId,
698 pub worktree: PathBuf,
699 pub workspace_slug: String,
700 pub lifecycle: TaskLifecyclePlan,
702 pub lifecycle_phase: TaskLifecyclePhase,
705 pub phase_epoch: u32,
706 pub phase_cursor: u32,
707 pub phase_iteration: u32,
708 pub gate_cycle: u32,
710 pub gate_proposal: Option<TaskGateProposal>,
712 pub agent: String,
715 pub provider: String,
717 pub provider_session_id: Option<String>,
719 pub abandon_intent: Option<AbandonIntent>,
722 pub created_at: OffsetDateTime,
723 pub updated_at: OffsetDateTime,
724 #[serde(skip)]
726 pub observation: Observation,
727}
728
729impl Task {
730 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 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 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 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 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 pub fn is_project_observable(&self) -> bool {
976 !matches!(self, Self::Started | Self::Progress { .. })
977 }
978
979 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#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1023pub struct TaskLinearObservation {
1024 pub task_id: TaskId,
1025 pub last_revision: String,
1028 pub last_title: String,
1030 pub last_description: String,
1031 pub last_success_at: OffsetDateTime,
1032 pub degraded_reason: Option<String>,
1034 pub updated_at: OffsetDateTime,
1035}
1036
1037#[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 pub content_steer: Option<String>,
1051 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#[derive(Debug, Clone, PartialEq, Eq)]
1064pub struct LinearObservationOutcome {
1065 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 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 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 #[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 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 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 #[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 assert!(failing("h1", &["scratch-clear", "rust-test"]).wake_legal());
1410 assert!(failing("h1", &["rust-test"]).wake_legal());
1411
1412 assert!(failing("h1", &[]).wake_legal());
1414
1415 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 #[test]
1425 fn only_land_time_preconditions_is_the_dual_of_wake_legal() {
1426 let scratch = failing("h1", &["scratch-clear"]);
1428 assert!(scratch.only_land_time_preconditions());
1429 assert!(!scratch.wake_legal());
1430
1431 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 let empty = failing("h1", &[]);
1442 assert!(!empty.only_land_time_preconditions());
1443 assert!(empty.wake_legal());
1444
1445 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 #[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(); 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 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 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 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 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 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}