Skip to main content

ag_session/
orchestration.rs

1//! Orchestration and orchestration-task lifecycle states.
2//!
3//! One orchestration groups the child sessions proposed by a single controller
4//! plan. The orchestration row tracks whether that plan is still awaiting the
5//! user's approval, actively fanning out, or settled; each task row tracks one
6//! child session through creation, execution, and settlement.
7
8use std::collections::HashSet;
9use std::fmt;
10use std::str::FromStr;
11
12use crate::SessionStatus;
13
14/// Maximum number of automatic focused-review remediation turns per managed
15/// worker settlement wave.
16pub const MAX_AUTOMATED_REVIEW_ITERATIONS: i64 = 3;
17
18/// Execution behavior for one persisted orchestration task.
19#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
20pub enum OrchestrationTaskKind {
21    /// Produces branch changes that proceed through review and integration.
22    #[default]
23    Implementation,
24    /// Produces a read-only report whose temporary worktree is discarded.
25    Research,
26}
27
28impl fmt::Display for OrchestrationTaskKind {
29    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
30        let value = match self {
31            OrchestrationTaskKind::Implementation => "Implementation",
32            OrchestrationTaskKind::Research => "Research",
33        };
34
35        formatter.write_str(value)
36    }
37}
38
39impl FromStr for OrchestrationTaskKind {
40    type Err = String;
41
42    fn from_str(value: &str) -> Result<Self, Self::Err> {
43        match value {
44            "Implementation" => Ok(OrchestrationTaskKind::Implementation),
45            "Research" => Ok(OrchestrationTaskKind::Research),
46            _ => Err(format!("Unknown orchestration task kind: {value}")),
47        }
48    }
49}
50
51/// Lifecycle state for one controller-owned orchestration.
52#[derive(Clone, Copy, Debug, Eq, PartialEq)]
53pub enum OrchestrationStatus {
54    /// The plan is persisted and parked on the campaign approval board.
55    AwaitingApproval,
56    /// The plan is approved and its tasks are fanning out.
57    Running,
58    /// Cancellation is blocking new fan-out while active children stop.
59    Canceling,
60    /// Every task settled and the controller is verifying the results.
61    Verifying,
62    /// Verification passed and user integration approval is required.
63    AwaitingIntegration,
64    /// Verified tasks are being merged or published in plan order.
65    Integrating,
66    /// Every verified task integrated and the campaign was archived.
67    Done,
68    /// The user canceled the orchestration or its controller session.
69    Canceled,
70}
71
72impl OrchestrationStatus {
73    /// Returns whether the orchestration is still open.
74    ///
75    /// An open plan blocks another plan from being persisted for the same
76    /// controller, including while it waits for approval.
77    pub fn is_active(self) -> bool {
78        matches!(
79            self,
80            OrchestrationStatus::AwaitingApproval
81                | OrchestrationStatus::Running
82                | OrchestrationStatus::Canceling
83                | OrchestrationStatus::Verifying
84                | OrchestrationStatus::AwaitingIntegration
85                | OrchestrationStatus::Integrating
86        )
87    }
88}
89
90impl fmt::Display for OrchestrationStatus {
91    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
92        let value = match self {
93            OrchestrationStatus::AwaitingApproval => "AwaitingApproval",
94            OrchestrationStatus::Running => "Running",
95            OrchestrationStatus::Canceling => "Canceling",
96            OrchestrationStatus::Verifying => "Verifying",
97            OrchestrationStatus::AwaitingIntegration => "AwaitingIntegration",
98            OrchestrationStatus::Integrating => "Integrating",
99            OrchestrationStatus::Done => "Done",
100            OrchestrationStatus::Canceled => "Canceled",
101        };
102
103        formatter.write_str(value)
104    }
105}
106
107impl FromStr for OrchestrationStatus {
108    type Err = String;
109
110    fn from_str(value: &str) -> Result<Self, Self::Err> {
111        match value {
112            "AwaitingApproval" => Ok(OrchestrationStatus::AwaitingApproval),
113            "Running" => Ok(OrchestrationStatus::Running),
114            "Canceling" => Ok(OrchestrationStatus::Canceling),
115            "Verifying" => Ok(OrchestrationStatus::Verifying),
116            "AwaitingIntegration" => Ok(OrchestrationStatus::AwaitingIntegration),
117            "Integrating" => Ok(OrchestrationStatus::Integrating),
118            "Done" => Ok(OrchestrationStatus::Done),
119            "Canceled" => Ok(OrchestrationStatus::Canceled),
120            _ => Err(format!("Unknown orchestration status: {value}")),
121        }
122    }
123}
124
125/// User-selected destination for verified orchestration task branches.
126#[derive(Clone, Copy, Debug, Eq, PartialEq)]
127pub enum IntegrationApproach {
128    /// Merge each verified child branch into the campaign base locally.
129    LocalMerge,
130    /// Publish each verified child branch as a forge review request.
131    ReviewRequest,
132}
133
134impl fmt::Display for IntegrationApproach {
135    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
136        let value = match self {
137            IntegrationApproach::LocalMerge => "LocalMerge",
138            IntegrationApproach::ReviewRequest => "ReviewRequest",
139        };
140
141        formatter.write_str(value)
142    }
143}
144
145impl FromStr for IntegrationApproach {
146    type Err = String;
147
148    fn from_str(value: &str) -> Result<Self, Self::Err> {
149        match value {
150            "LocalMerge" => Ok(IntegrationApproach::LocalMerge),
151            "ReviewRequest" => Ok(IntegrationApproach::ReviewRequest),
152            _ => Err(format!("Unknown integration approach: {value}")),
153        }
154    }
155}
156
157/// Lifecycle state for one orchestration task and its child session.
158#[derive(Clone, Copy, Debug, Eq, PartialEq)]
159pub enum OrchestrationTaskStatus {
160    /// Follow-up scope proposed by the controller and awaiting approval.
161    Proposed,
162    /// Persisted with the proposed plan, not yet approved or fanned out.
163    Planned,
164    /// The child session is being created and started.
165    Creating,
166    /// The child session is running its turn.
167    Running,
168    /// Focused review is running or awaiting a persisted result.
169    Reviewing,
170    /// The coordinator is applying one focused-review suggestion set.
171    ReviewApplying,
172    /// The child session parked on clarification questions.
173    WaitingForInput,
174    /// The child session finished and is ready for review or integration.
175    Ready,
176    /// A research child returned its report and its worktree was discarded.
177    Reported,
178    /// Approved feedback is being delivered to the existing managed child.
179    ContinuationPending,
180    /// Verification passed and this task awaits its integration gate.
181    AwaitingIntegration,
182    /// The coordinator has started branch merge or publication.
183    Merging,
184    /// Branch work was merged locally successfully.
185    Integrated,
186    /// The child branch was published and is waiting for its forge review
187    /// request to merge.
188    ReviewRequested,
189    /// Integration failed and requires attention.
190    IntegrationFailed,
191    /// Ownership permanently transferred from the coordinator to the user.
192    Detached,
193    /// The child session failed, or a straggler was canceled out of band.
194    Failed,
195    /// The task was canceled as part of a cascade cancel.
196    Canceled,
197}
198
199impl OrchestrationTaskStatus {
200    /// Maps one observed child-session status into the task state owned by
201    /// orchestration.
202    pub fn from_child_status(status: SessionStatus) -> Self {
203        match status {
204            SessionStatus::Draft
205            | SessionStatus::InProgress
206            | SessionStatus::Queued
207            | SessionStatus::Rebasing
208            | SessionStatus::Merging => Self::Running,
209            SessionStatus::Question => Self::WaitingForInput,
210            SessionStatus::Review | SessionStatus::AgentReview => Self::Reviewing,
211            SessionStatus::Merged | SessionStatus::Done => Self::Ready,
212            SessionStatus::Canceled => Self::Failed,
213        }
214    }
215
216    /// Returns whether the task reached a state that fan-in treats as settled.
217    ///
218    /// A canceled straggler counts as settled so out-of-band cancellation
219    /// unblocks the roll-up instead of stalling it.
220    pub fn is_settled(self) -> bool {
221        matches!(
222            self,
223            OrchestrationTaskStatus::Ready
224                | OrchestrationTaskStatus::Reported
225                | OrchestrationTaskStatus::Integrated
226                | OrchestrationTaskStatus::ReviewRequested
227                | OrchestrationTaskStatus::IntegrationFailed
228                | OrchestrationTaskStatus::Failed
229                | OrchestrationTaskStatus::Canceled
230                | OrchestrationTaskStatus::Detached
231        )
232    }
233
234    /// Returns the concise user-facing label shown in campaign status output.
235    pub fn campaign_label(self) -> &'static str {
236        self.labels().1
237    }
238
239    /// Returns whether the task currently occupies a parallelism slot.
240    ///
241    /// A task waiting for user input still holds its child session and
242    /// worktree, so it keeps consuming a slot until the user answers.
243    pub fn occupies_parallelism_slot(self) -> bool {
244        matches!(
245            self,
246            OrchestrationTaskStatus::Creating
247                | OrchestrationTaskStatus::Running
248                | OrchestrationTaskStatus::Reviewing
249                | OrchestrationTaskStatus::ReviewApplying
250                | OrchestrationTaskStatus::WaitingForInput
251                | OrchestrationTaskStatus::ContinuationPending
252                | OrchestrationTaskStatus::Merging
253        )
254    }
255
256    /// Returns whether a transition to `next` is valid.
257    ///
258    /// Retry re-enters `Creating` from a settled state with the same task key,
259    /// which is what makes replying "retry the failed tasks" a clean respawn
260    /// rather than a duplicate fan-out.
261    pub fn can_transition_to(self, next: OrchestrationTaskStatus) -> bool {
262        if self == next {
263            return true;
264        }
265
266        matches!(
267            (self, next),
268            (
269                OrchestrationTaskStatus::Proposed,
270                OrchestrationTaskStatus::Planned
271            ) | (
272                OrchestrationTaskStatus::Planned
273                    | OrchestrationTaskStatus::Ready
274                    | OrchestrationTaskStatus::Failed
275                    | OrchestrationTaskStatus::Canceled
276                    | OrchestrationTaskStatus::IntegrationFailed,
277                OrchestrationTaskStatus::Creating
278            ) | (
279                OrchestrationTaskStatus::Creating | OrchestrationTaskStatus::WaitingForInput,
280                OrchestrationTaskStatus::Running
281            ) | (
282                OrchestrationTaskStatus::Running,
283                OrchestrationTaskStatus::Reviewing | OrchestrationTaskStatus::WaitingForInput
284            ) | (
285                OrchestrationTaskStatus::Running
286                    | OrchestrationTaskStatus::Reviewing
287                    | OrchestrationTaskStatus::WaitingForInput,
288                OrchestrationTaskStatus::Ready | OrchestrationTaskStatus::Reported
289            ) | (
290                OrchestrationTaskStatus::Reviewing,
291                OrchestrationTaskStatus::ReviewApplying
292            ) | (
293                OrchestrationTaskStatus::ReviewApplying,
294                OrchestrationTaskStatus::Reviewing
295                    | OrchestrationTaskStatus::WaitingForInput
296                    | OrchestrationTaskStatus::Failed
297            ) | (
298                OrchestrationTaskStatus::Ready,
299                OrchestrationTaskStatus::AwaitingIntegration
300                    | OrchestrationTaskStatus::ContinuationPending
301                    | OrchestrationTaskStatus::Detached
302            ) | (
303                OrchestrationTaskStatus::AwaitingIntegration,
304                OrchestrationTaskStatus::Merging
305                    | OrchestrationTaskStatus::ContinuationPending
306                    | OrchestrationTaskStatus::Detached
307            ) | (
308                OrchestrationTaskStatus::IntegrationFailed,
309                OrchestrationTaskStatus::ContinuationPending | OrchestrationTaskStatus::Detached
310            ) | (
311                OrchestrationTaskStatus::ContinuationPending,
312                OrchestrationTaskStatus::Ready
313                    | OrchestrationTaskStatus::WaitingForInput
314                    | OrchestrationTaskStatus::Failed
315            ) | (
316                OrchestrationTaskStatus::Merging,
317                OrchestrationTaskStatus::Integrated
318                    | OrchestrationTaskStatus::ReviewRequested
319                    | OrchestrationTaskStatus::IntegrationFailed
320            ) | (
321                OrchestrationTaskStatus::ReviewRequested,
322                OrchestrationTaskStatus::Integrated | OrchestrationTaskStatus::IntegrationFailed
323            ) | (
324                OrchestrationTaskStatus::Planned
325                    | OrchestrationTaskStatus::Creating
326                    | OrchestrationTaskStatus::Running
327                    | OrchestrationTaskStatus::Reviewing
328                    | OrchestrationTaskStatus::ReviewApplying
329                    | OrchestrationTaskStatus::WaitingForInput,
330                OrchestrationTaskStatus::Failed | OrchestrationTaskStatus::Canceled
331            )
332        )
333    }
334
335    /// Returns whether this task no longer needs integration work.
336    pub fn is_integration_settled(self) -> bool {
337        matches!(
338            self,
339            OrchestrationTaskStatus::Integrated
340                | OrchestrationTaskStatus::Detached
341                | OrchestrationTaskStatus::Canceled
342                | OrchestrationTaskStatus::Failed
343        )
344    }
345
346    fn labels(self) -> (&'static str, &'static str) {
347        match self {
348            OrchestrationTaskStatus::Proposed => ("Proposed", "awaiting approval"),
349            OrchestrationTaskStatus::Planned => ("Planned", "waiting"),
350            OrchestrationTaskStatus::Creating => ("Creating", "starting"),
351            OrchestrationTaskStatus::Running => ("Running", "running"),
352            OrchestrationTaskStatus::Reviewing => ("Reviewing", "reviewing"),
353            OrchestrationTaskStatus::ReviewApplying => ("ReviewApplying", "applying review"),
354            OrchestrationTaskStatus::WaitingForInput => ("WaitingForInput", "waiting on you"),
355            OrchestrationTaskStatus::Ready => ("Ready", "ready"),
356            OrchestrationTaskStatus::Reported => ("Reported", "reported"),
357            OrchestrationTaskStatus::ContinuationPending => ("ContinuationPending", "continuing"),
358            OrchestrationTaskStatus::AwaitingIntegration => {
359                ("AwaitingIntegration", "awaiting integration")
360            }
361            OrchestrationTaskStatus::Merging => ("Merging", "integrating"),
362            OrchestrationTaskStatus::Integrated => ("Integrated", "integrated"),
363            OrchestrationTaskStatus::ReviewRequested => ("ReviewRequested", "review requested"),
364            OrchestrationTaskStatus::IntegrationFailed => {
365                ("IntegrationFailed", "integration failed")
366            }
367            OrchestrationTaskStatus::Detached => ("Detached", "detached"),
368            OrchestrationTaskStatus::Failed => ("Failed", "failed"),
369            OrchestrationTaskStatus::Canceled => ("Canceled", "canceled"),
370        }
371    }
372}
373
374/// Pure scheduling decision derived from one orchestration task snapshot.
375#[derive(Clone, Copy, Debug, Eq, PartialEq)]
376pub struct OrchestrationScheduleDecision {
377    /// Number of planned tasks that may claim a parallelism slot.
378    pub spawn_count: usize,
379    /// Whether every non-empty task has settled and roll-up can be claimed.
380    pub should_submit: bool,
381}
382
383/// Pure orchestration policy over typed task observations.
384pub struct OrchestrationPolicy;
385
386impl OrchestrationPolicy {
387    /// Decides fan-out capacity and roll-up readiness without persistence or
388    /// runtime dependencies.
389    pub fn schedule(
390        max_parallelism: usize,
391        task_statuses: &[Option<OrchestrationTaskStatus>],
392    ) -> OrchestrationScheduleDecision {
393        let occupied_slots = task_statuses
394            .iter()
395            .filter(|status| status.is_some_and(OrchestrationTaskStatus::occupies_parallelism_slot))
396            .count();
397        let planned_tasks = task_statuses
398            .iter()
399            .filter(|status| **status == Some(OrchestrationTaskStatus::Planned))
400            .count();
401        let spawn_count = max_parallelism
402            .saturating_sub(occupied_slots)
403            .min(planned_tasks);
404        let should_submit = !task_statuses.is_empty()
405            && task_statuses
406                .iter()
407                .all(|status| status.is_some_and(OrchestrationTaskStatus::is_settled));
408
409        OrchestrationScheduleDecision {
410            spawn_count,
411            should_submit,
412        }
413    }
414}
415
416/// Protocol-independent snapshot of one task proposed by an orchestration
417/// controller.
418#[derive(Clone)]
419pub struct OrchestrationPlanTask {
420    /// Observable conditions checked during settlement verification.
421    pub acceptance_criteria: Vec<String>,
422    /// Whether the task implements changes or returns research findings.
423    pub kind: OrchestrationTaskKind,
424    /// Standalone task prompt delivered to the child session.
425    pub prompt: String,
426    /// Stable kebab-case identity used for retries.
427    pub task_key: String,
428    /// Short user-facing task title.
429    pub title: String,
430    /// Best-effort repository-relative files or directories expected to change.
431    pub touched_areas: Vec<String>,
432}
433
434/// Validates one proposed subtask set before application code persists it.
435///
436/// # Errors
437///
438/// Returns a user-facing reason when the plan is too small, incomplete, or uses
439/// invalid task keys or planning paths.
440pub fn validate_subtasks(subtasks: &[OrchestrationPlanTask], is_retry: bool) -> Result<(), String> {
441    let is_research_wave = !subtasks.is_empty()
442        && subtasks
443            .iter()
444            .all(|subtask| subtask.kind == OrchestrationTaskKind::Research);
445    let has_research = subtasks
446        .iter()
447        .any(|subtask| subtask.kind == OrchestrationTaskKind::Research);
448    if has_research && !is_research_wave {
449        return Err(
450            "research and implementation tasks must be proposed in separate waves.".to_string(),
451        );
452    }
453    if subtasks.len() < 2 && !is_retry && !is_research_wave {
454        return Err("a meaningful orchestration requires at least two subtasks.".to_string());
455    }
456    let mut task_keys = HashSet::new();
457    for subtask in subtasks {
458        if !is_kebab_case_task_key(&subtask.task_key)
459            || !task_keys.insert(subtask.task_key.as_str())
460        {
461            return Err("every subtask needs a unique kebab-case task key.".to_string());
462        }
463        if subtask.prompt.trim().is_empty()
464            || subtask.title.trim().is_empty()
465            || subtask
466                .acceptance_criteria
467                .iter()
468                .all(|criterion| criterion.trim().is_empty())
469        {
470            return Err(format!(
471                "subtask `{}` needs a title, standalone prompt, and acceptance criteria.",
472                subtask.task_key
473            ));
474        }
475        for area in subtask
476            .touched_areas
477            .iter()
478            .filter(|_| subtask.kind == OrchestrationTaskKind::Implementation)
479        {
480            normalized_scope(area).map_err(|reason| {
481                format!(
482                    "subtask `{}` has invalid touched area `{area}`: {reason}.",
483                    subtask.task_key
484                )
485            })?;
486        }
487    }
488
489    Ok(())
490}
491
492fn is_kebab_case_task_key(task_key: &str) -> bool {
493    !task_key.is_empty()
494        && task_key.split('-').all(|segment| {
495            !segment.is_empty()
496                && segment
497                    .bytes()
498                    .all(|byte| byte.is_ascii_lowercase() || byte.is_ascii_digit())
499        })
500}
501
502fn normalized_scope(area: &str) -> Result<String, &'static str> {
503    let normalized = area.trim().trim_start_matches("./").trim_end_matches('/');
504    if normalized.is_empty()
505        || normalized.starts_with('/')
506        || normalized.split('/').any(|part| part == "..")
507    {
508        return Err("use a non-empty repository-relative path");
509    }
510    if normalized.contains(['*', '?', '[', ']', '{', '}']) {
511        return Err("use a literal file or directory path; wildcard patterns are not supported");
512    }
513
514    Ok(normalized.to_string())
515}
516
517impl fmt::Display for OrchestrationTaskStatus {
518    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
519        formatter.write_str(self.labels().0)
520    }
521}
522
523impl FromStr for OrchestrationTaskStatus {
524    type Err = String;
525
526    fn from_str(value: &str) -> Result<Self, Self::Err> {
527        match value {
528            "Proposed" => Ok(OrchestrationTaskStatus::Proposed),
529            "Planned" => Ok(OrchestrationTaskStatus::Planned),
530            "Creating" => Ok(OrchestrationTaskStatus::Creating),
531            "Running" => Ok(OrchestrationTaskStatus::Running),
532            "Reviewing" => Ok(OrchestrationTaskStatus::Reviewing),
533            "ReviewApplying" => Ok(OrchestrationTaskStatus::ReviewApplying),
534            "WaitingForInput" => Ok(OrchestrationTaskStatus::WaitingForInput),
535            "Ready" => Ok(OrchestrationTaskStatus::Ready),
536            "Reported" => Ok(OrchestrationTaskStatus::Reported),
537            "ContinuationPending" => Ok(OrchestrationTaskStatus::ContinuationPending),
538            "AwaitingIntegration" => Ok(OrchestrationTaskStatus::AwaitingIntegration),
539            "Merging" => Ok(OrchestrationTaskStatus::Merging),
540            "Integrated" => Ok(OrchestrationTaskStatus::Integrated),
541            "ReviewRequested" => Ok(OrchestrationTaskStatus::ReviewRequested),
542            "IntegrationFailed" => Ok(OrchestrationTaskStatus::IntegrationFailed),
543            "Detached" => Ok(OrchestrationTaskStatus::Detached),
544            "Failed" => Ok(OrchestrationTaskStatus::Failed),
545            "Canceled" => Ok(OrchestrationTaskStatus::Canceled),
546            _ => Err(format!("Unknown orchestration task status: {value}")),
547        }
548    }
549}
550
551#[cfg(test)]
552mod tests {
553    use super::*;
554
555    #[test]
556    /// Round-trips every orchestration status through its persisted form.
557    fn test_orchestration_status_round_trips_persisted_values() {
558        // Arrange
559        let statuses = [
560            OrchestrationStatus::AwaitingApproval,
561            OrchestrationStatus::Running,
562            OrchestrationStatus::Canceling,
563            OrchestrationStatus::Verifying,
564            OrchestrationStatus::AwaitingIntegration,
565            OrchestrationStatus::Integrating,
566            OrchestrationStatus::Done,
567            OrchestrationStatus::Canceled,
568        ];
569
570        // Act
571        let round_tripped = statuses.map(|status| {
572            status
573                .to_string()
574                .parse::<OrchestrationStatus>()
575                .expect("status should parse")
576        });
577
578        // Assert
579        assert_eq!(round_tripped, statuses);
580        assert!("Unknown".parse::<OrchestrationStatus>().is_err());
581    }
582
583    #[test]
584    fn integration_approach_round_trips_persisted_values() {
585        // Arrange
586        let approaches = [
587            IntegrationApproach::LocalMerge,
588            IntegrationApproach::ReviewRequest,
589        ];
590
591        // Act
592        let round_tripped = approaches.map(|approach| {
593            approach
594                .to_string()
595                .parse::<IntegrationApproach>()
596                .expect("approach should parse")
597        });
598
599        // Assert
600        assert_eq!(round_tripped, approaches);
601        assert!("Unknown".parse::<IntegrationApproach>().is_err());
602    }
603
604    #[test]
605    fn orchestration_task_kind_round_trips_persisted_values() {
606        // Arrange
607        let kinds = [
608            OrchestrationTaskKind::Implementation,
609            OrchestrationTaskKind::Research,
610        ];
611
612        // Act
613        let round_tripped = kinds.map(|kind| {
614            kind.to_string()
615                .parse::<OrchestrationTaskKind>()
616                .expect("kind should parse")
617        });
618
619        // Assert
620        assert_eq!(round_tripped, kinds);
621        assert_eq!(OrchestrationTaskKind::default(), kinds[0]);
622        assert!("Unknown".parse::<OrchestrationTaskKind>().is_err());
623    }
624
625    #[test]
626    /// Restricts restart re-linking to orchestrations that still need work.
627    fn test_only_unsettled_orchestrations_are_active() {
628        // Arrange / Act / Assert
629        assert!(OrchestrationStatus::AwaitingApproval.is_active());
630        assert!(OrchestrationStatus::Running.is_active());
631        assert!(OrchestrationStatus::Canceling.is_active());
632        assert!(OrchestrationStatus::Verifying.is_active());
633        assert!(OrchestrationStatus::AwaitingIntegration.is_active());
634        assert!(OrchestrationStatus::Integrating.is_active());
635        assert!(!OrchestrationStatus::Done.is_active());
636        assert!(!OrchestrationStatus::Canceled.is_active());
637    }
638
639    #[test]
640    /// Round-trips every task status through its persisted form.
641    fn test_orchestration_task_status_round_trips_persisted_values() {
642        // Arrange
643        let statuses = [
644            OrchestrationTaskStatus::Proposed,
645            OrchestrationTaskStatus::Planned,
646            OrchestrationTaskStatus::Creating,
647            OrchestrationTaskStatus::Running,
648            OrchestrationTaskStatus::Reviewing,
649            OrchestrationTaskStatus::ReviewApplying,
650            OrchestrationTaskStatus::WaitingForInput,
651            OrchestrationTaskStatus::Ready,
652            OrchestrationTaskStatus::Reported,
653            OrchestrationTaskStatus::ContinuationPending,
654            OrchestrationTaskStatus::AwaitingIntegration,
655            OrchestrationTaskStatus::Merging,
656            OrchestrationTaskStatus::Integrated,
657            OrchestrationTaskStatus::ReviewRequested,
658            OrchestrationTaskStatus::IntegrationFailed,
659            OrchestrationTaskStatus::Detached,
660            OrchestrationTaskStatus::Failed,
661            OrchestrationTaskStatus::Canceled,
662        ];
663
664        // Act
665        let round_tripped = statuses.map(|status| {
666            status
667                .to_string()
668                .parse::<OrchestrationTaskStatus>()
669                .expect("status should parse")
670        });
671
672        // Assert
673        assert_eq!(round_tripped, statuses);
674        assert!("Unknown".parse::<OrchestrationTaskStatus>().is_err());
675    }
676
677    #[test]
678    /// Treats a canceled straggler as settled so fan-in is not blocked by
679    /// out-of-band cancellation.
680    fn test_settled_task_statuses_include_cancellation() {
681        // Arrange / Act / Assert
682        assert!(OrchestrationTaskStatus::Ready.is_settled());
683        assert!(OrchestrationTaskStatus::Reported.is_settled());
684        assert!(OrchestrationTaskStatus::Integrated.is_settled());
685        assert!(OrchestrationTaskStatus::ReviewRequested.is_settled());
686        assert!(OrchestrationTaskStatus::IntegrationFailed.is_settled());
687        assert!(OrchestrationTaskStatus::Failed.is_settled());
688        assert!(OrchestrationTaskStatus::Canceled.is_settled());
689        assert!(!OrchestrationTaskStatus::Planned.is_settled());
690        assert!(!OrchestrationTaskStatus::Creating.is_settled());
691        assert!(!OrchestrationTaskStatus::Running.is_settled());
692        assert!(!OrchestrationTaskStatus::Reviewing.is_settled());
693        assert!(!OrchestrationTaskStatus::ReviewApplying.is_settled());
694        assert!(!OrchestrationTaskStatus::WaitingForInput.is_settled());
695    }
696
697    #[test]
698    /// Counts a task waiting for user input against the parallelism cap
699    /// because it still owns a live child session and worktree.
700    fn test_parallelism_slots_cover_every_live_child() {
701        // Arrange / Act / Assert
702        assert!(OrchestrationTaskStatus::Creating.occupies_parallelism_slot());
703        assert!(OrchestrationTaskStatus::Running.occupies_parallelism_slot());
704        assert!(OrchestrationTaskStatus::Reviewing.occupies_parallelism_slot());
705        assert!(OrchestrationTaskStatus::ReviewApplying.occupies_parallelism_slot());
706        assert!(OrchestrationTaskStatus::WaitingForInput.occupies_parallelism_slot());
707        assert!(!OrchestrationTaskStatus::Planned.occupies_parallelism_slot());
708        assert!(!OrchestrationTaskStatus::Ready.occupies_parallelism_slot());
709        assert!(!OrchestrationTaskStatus::Failed.occupies_parallelism_slot());
710        assert!(!OrchestrationTaskStatus::Canceled.occupies_parallelism_slot());
711    }
712
713    #[test]
714    /// Allows the fan-out, question, settle, and retry transitions the
715    /// coordinator drives, and rejects skipping creation.
716    fn test_task_status_transitions_cover_fan_out_and_retry() {
717        // Arrange / Act / Assert
718        assert!(
719            OrchestrationTaskStatus::Planned.can_transition_to(OrchestrationTaskStatus::Creating)
720        );
721        assert!(
722            OrchestrationTaskStatus::Creating.can_transition_to(OrchestrationTaskStatus::Running)
723        );
724        assert!(
725            OrchestrationTaskStatus::Running
726                .can_transition_to(OrchestrationTaskStatus::WaitingForInput)
727        );
728        assert!(
729            OrchestrationTaskStatus::WaitingForInput
730                .can_transition_to(OrchestrationTaskStatus::Running)
731        );
732        assert!(OrchestrationTaskStatus::Running.can_transition_to(OrchestrationTaskStatus::Ready));
733        assert!(
734            OrchestrationTaskStatus::Running.can_transition_to(OrchestrationTaskStatus::Reported)
735        );
736        assert!(
737            OrchestrationTaskStatus::Running.can_transition_to(OrchestrationTaskStatus::Reviewing)
738        );
739        assert!(
740            OrchestrationTaskStatus::Reviewing
741                .can_transition_to(OrchestrationTaskStatus::ReviewApplying)
742        );
743        assert!(
744            OrchestrationTaskStatus::ReviewApplying
745                .can_transition_to(OrchestrationTaskStatus::Reviewing)
746        );
747        assert!(
748            OrchestrationTaskStatus::Running.can_transition_to(OrchestrationTaskStatus::Canceled)
749        );
750        assert!(
751            OrchestrationTaskStatus::Failed.can_transition_to(OrchestrationTaskStatus::Creating)
752        );
753        assert!(OrchestrationTaskStatus::Ready.can_transition_to(OrchestrationTaskStatus::Ready));
754        assert!(
755            !OrchestrationTaskStatus::Planned.can_transition_to(OrchestrationTaskStatus::Running)
756        );
757        assert!(
758            !OrchestrationTaskStatus::Canceled.can_transition_to(OrchestrationTaskStatus::Ready)
759        );
760    }
761
762    #[test]
763    fn integration_settlement_waits_for_review_request_merge() {
764        // Arrange
765        let settled = [
766            OrchestrationTaskStatus::Integrated,
767            OrchestrationTaskStatus::Detached,
768            OrchestrationTaskStatus::Canceled,
769            OrchestrationTaskStatus::Failed,
770        ];
771        // Act / Assert
772        assert!(
773            settled
774                .into_iter()
775                .all(OrchestrationTaskStatus::is_integration_settled)
776        );
777        assert!(!OrchestrationTaskStatus::Reported.is_integration_settled());
778        assert!(!OrchestrationTaskStatus::AwaitingIntegration.is_integration_settled());
779        assert!(!OrchestrationTaskStatus::ReviewRequested.is_integration_settled());
780        assert!(
781            OrchestrationTaskStatus::Merging
782                .can_transition_to(OrchestrationTaskStatus::ReviewRequested)
783        );
784        assert!(
785            OrchestrationTaskStatus::ReviewRequested
786                .can_transition_to(OrchestrationTaskStatus::Integrated)
787        );
788        assert!(
789            OrchestrationTaskStatus::ReviewRequested
790                .can_transition_to(OrchestrationTaskStatus::IntegrationFailed)
791        );
792    }
793
794    #[test]
795    fn campaign_labels_cover_every_task_status() {
796        // Arrange
797        let statuses = [
798            OrchestrationTaskStatus::Proposed,
799            OrchestrationTaskStatus::Planned,
800            OrchestrationTaskStatus::Creating,
801            OrchestrationTaskStatus::Running,
802            OrchestrationTaskStatus::Reviewing,
803            OrchestrationTaskStatus::ReviewApplying,
804            OrchestrationTaskStatus::WaitingForInput,
805            OrchestrationTaskStatus::Ready,
806            OrchestrationTaskStatus::Reported,
807            OrchestrationTaskStatus::ContinuationPending,
808            OrchestrationTaskStatus::AwaitingIntegration,
809            OrchestrationTaskStatus::Merging,
810            OrchestrationTaskStatus::Integrated,
811            OrchestrationTaskStatus::ReviewRequested,
812            OrchestrationTaskStatus::IntegrationFailed,
813            OrchestrationTaskStatus::Detached,
814            OrchestrationTaskStatus::Failed,
815            OrchestrationTaskStatus::Canceled,
816        ];
817
818        // Act
819        let labels = statuses.map(OrchestrationTaskStatus::campaign_label);
820
821        // Assert
822        assert_eq!(
823            labels,
824            [
825                "awaiting approval",
826                "waiting",
827                "starting",
828                "running",
829                "reviewing",
830                "applying review",
831                "waiting on you",
832                "ready",
833                "reported",
834                "continuing",
835                "awaiting integration",
836                "integrating",
837                "integrated",
838                "review requested",
839                "integration failed",
840                "detached",
841                "failed",
842                "canceled",
843            ]
844        );
845    }
846
847    #[test]
848    fn validation_allows_one_research_task_and_ignores_its_touched_areas() {
849        // Arrange
850        let plan = [OrchestrationPlanTask {
851            acceptance_criteria: vec!["Architecture questions are answered".to_string()],
852            kind: OrchestrationTaskKind::Research,
853            prompt: "Inspect the architecture".to_string(),
854            task_key: "architecture".to_string(),
855            title: "Architecture research".to_string(),
856            touched_areas: vec!["**".to_string()],
857        }];
858
859        // Act
860        let result = validate_subtasks(&plan, false);
861
862        // Assert
863        assert_eq!(result, Ok(()));
864    }
865
866    #[test]
867    fn validation_rejects_one_implementation_task_and_invalid_implementation_scope() {
868        // Arrange
869        let single = [OrchestrationPlanTask {
870            acceptance_criteria: vec!["Feature is complete".to_string()],
871            kind: OrchestrationTaskKind::Implementation,
872            prompt: "Implement the feature".to_string(),
873            task_key: "feature".to_string(),
874            title: "Feature".to_string(),
875            touched_areas: Vec::new(),
876        }];
877        let invalid_scope = [
878            OrchestrationPlanTask {
879                touched_areas: vec!["crates/one/**".to_string()],
880                ..single[0].clone()
881            },
882            OrchestrationPlanTask {
883                task_key: "tests".to_string(),
884                touched_areas: Vec::new(),
885                ..single[0].clone()
886            },
887        ];
888        let mixed = [
889            OrchestrationPlanTask {
890                kind: OrchestrationTaskKind::Research,
891                ..single[0].clone()
892            },
893            OrchestrationPlanTask {
894                task_key: "implementation".to_string(),
895                ..single[0].clone()
896            },
897        ];
898
899        // Act
900        let single_result = validate_subtasks(&single, false);
901        let empty_result = validate_subtasks(&[], false);
902        let scope_result = validate_subtasks(&invalid_scope, false);
903        let mixed_result = validate_subtasks(&mixed, false);
904
905        // Assert
906        assert_eq!(
907            single_result,
908            Err("a meaningful orchestration requires at least two subtasks.".to_string())
909        );
910        assert_eq!(
911            empty_result,
912            Err("a meaningful orchestration requires at least two subtasks.".to_string())
913        );
914        assert!(scope_result.is_err_and(|reason| reason.contains("wildcard patterns")));
915        assert_eq!(
916            mixed_result,
917            Err(
918                "research and implementation tasks must be proposed in separate waves.".to_string()
919            )
920        );
921    }
922
923    #[test]
924    /// Derives fan-out capacity and roll-up readiness from typed task states.
925    fn test_orchestration_policy_schedules_available_slots_and_settlement() {
926        // Arrange
927        let active_statuses = [
928            Some(OrchestrationTaskStatus::Running),
929            Some(OrchestrationTaskStatus::WaitingForInput),
930            Some(OrchestrationTaskStatus::Planned),
931            Some(OrchestrationTaskStatus::Planned),
932        ];
933        let settled_statuses = [
934            Some(OrchestrationTaskStatus::Ready),
935            Some(OrchestrationTaskStatus::Failed),
936            Some(OrchestrationTaskStatus::Canceled),
937        ];
938        let invalid_statuses = [Some(OrchestrationTaskStatus::Ready), None];
939
940        // Act
941        let active_decision = OrchestrationPolicy::schedule(3, &active_statuses);
942        let settled_decision = OrchestrationPolicy::schedule(3, &settled_statuses);
943        let empty_decision = OrchestrationPolicy::schedule(3, &[]);
944        let invalid_decision = OrchestrationPolicy::schedule(3, &invalid_statuses);
945
946        // Assert
947        assert_eq!(
948            active_decision,
949            OrchestrationScheduleDecision {
950                spawn_count: 1,
951                should_submit: false,
952            }
953        );
954        assert_eq!(
955            settled_decision,
956            OrchestrationScheduleDecision {
957                spawn_count: 0,
958                should_submit: true,
959            }
960        );
961        assert_eq!(
962            empty_decision,
963            OrchestrationScheduleDecision {
964                spawn_count: 0,
965                should_submit: false,
966            }
967        );
968        assert!(!invalid_decision.should_submit);
969    }
970
971    #[test]
972    /// Maps every child-session lifecycle family into orchestration policy.
973    fn test_task_status_from_child_status_covers_session_lifecycle() {
974        // Arrange
975        let cases = [
976            (SessionStatus::Draft, OrchestrationTaskStatus::Running),
977            (SessionStatus::InProgress, OrchestrationTaskStatus::Running),
978            (SessionStatus::Queued, OrchestrationTaskStatus::Running),
979            (SessionStatus::Rebasing, OrchestrationTaskStatus::Running),
980            (SessionStatus::Merging, OrchestrationTaskStatus::Running),
981            (
982                SessionStatus::Question,
983                OrchestrationTaskStatus::WaitingForInput,
984            ),
985            (SessionStatus::Review, OrchestrationTaskStatus::Reviewing),
986            (
987                SessionStatus::AgentReview,
988                OrchestrationTaskStatus::Reviewing,
989            ),
990            (SessionStatus::Merged, OrchestrationTaskStatus::Ready),
991            (SessionStatus::Done, OrchestrationTaskStatus::Ready),
992            (SessionStatus::Canceled, OrchestrationTaskStatus::Failed),
993        ];
994
995        // Act / Assert
996        for (session_status, expected_task_status) in cases {
997            assert_eq!(
998                OrchestrationTaskStatus::from_child_status(session_status),
999                expected_task_status
1000            );
1001        }
1002    }
1003}