use std::collections::{HashSet, VecDeque};
use std::io::BufRead;
use std::path::Path;
use std::sync::Arc;
use std::time::Duration;
use anyhow::{anyhow, Result};
use sha2::{Digest, Sha256};
use tokio::sync::mpsc;
use crate::chat::types::{ConversationEvent, ConversationItem, Lifecycle};
use crate::child_control::{
absorb_commands as absorb_child_commands, apply_input as apply_child_input, input_is_current,
reconcile_stale_deliveries, ChildTarget, CommandStop, DecisionResolution, PendingInput,
};
use crate::child_session::{
task_write_lease_from_env, unincorporated_directive_version, BoundaryResult,
ChildBodyHandoffRequest, ChildBodyOutcome, ChildCommand, ChildCommandEffect, ChildCommandId,
ChildCommandKind, ChildCommandSource, ChildCommandState, ChildDirective, ChildLeaseState,
ChildRef, ChildWriteLease,
};
use crate::engine::wave_config::read_wave_config;
use crate::engine::InteractionPolicy;
use crate::harness::{
classify_disconnect_recovery, drain_turn_failure_reason, ApprovalPolicy, Harness,
RecoveryDecision,
};
use crate::interaction_review::{
InteractionReview, InteractionReviewDisposition, InteractionReviewEvidence,
InteractionReviewId, InteractionReviewPr, InteractionReviewStatus, InteractionReviewer,
};
use crate::interactive_handoff::{InteractiveHandoffOutcome, InteractiveHandoffParent};
use crate::store::{open_existing_store, SharedStore};
use crate::task::interactive_rendezvous::{self, Rendezvous};
use crate::task::{
CiCheck, Observation, PrPhase, TaskEventKind, TaskGateProposal, TaskLifecyclePhase,
TaskSession, TaskSessionId, TaskSessionStatus,
};
use crate::wave::playhead::{
BodyProvenance, Playhead, PlayheadEvent, QueuedInvocation, StepKind, StepOutcome,
};
use crate::wave::Wave;
#[derive(Debug)]
struct PreparedTaskStep {
turn: crate::lf::commands::run::PreparedHarnessTurn,
review: Option<InteractionReview>,
}
#[derive(Debug)]
struct StartedTaskStep {
review: Option<InteractionReviewId>,
provider_turn_active: bool,
}
pub async fn run_task_session(session_id: TaskSessionId, generation: u32) -> Result<()> {
let lease = task_write_lease_from_env().map_err(|error| anyhow!(error))?;
if lease.generation != generation {
anyhow::bail!(
"Task generation {generation} does not match its ambient write lease generation {}",
lease.generation
);
}
let result = run_task_session_inner(
session_id.clone(),
&lease,
Box::new(crate::harness::default_create_harness),
)
.await;
if let Err(error) = &result {
record_unhandled_failure(&session_id, &lease, error).await;
}
result
}
async fn owning_wave(store: &SharedStore, session: &TaskSession) -> Result<Wave> {
store
.get_wave(&session.wave_id)
.await?
.ok_or_else(|| anyhow!("owning Wave {} is not registered", session.wave_id))
}
async fn run_task_session_inner(
session_id: TaskSessionId,
lease: &ChildWriteLease,
create_harness: crate::harness::CreateHarness,
) -> Result<()> {
let store: SharedStore = Arc::new(
open_existing_store()
.await
.ok_or_else(|| anyhow!("no Loopflow registry on this machine"))?,
);
run_task_session_with(store, session_id, lease, create_harness).await
}
async fn run_task_session_with(
store: SharedStore,
session_id: TaskSessionId,
lease: &ChildWriteLease,
create_harness: crate::harness::CreateHarness,
) -> Result<()> {
let generation = lease.generation;
let mut session = store
.get_task_session(&session_id)
.await?
.ok_or_else(|| anyhow!("Task Session {session_id} not found"))?;
let wave = owning_wave(&store, &session).await?;
let recorded_generation = session
.latest_process
.as_ref()
.map(|process| process.generation);
if recorded_generation != Some(generation) {
anyhow::bail!(
"Task Session {session_id} generation mismatch: expected {:?}, got {generation}",
recorded_generation
);
}
if let Some(process) = &mut session.latest_process {
process.mark_booted();
}
let from = session.status;
session.set_status(TaskSessionStatus::Running, "provider turn is active");
store.activate_task_process(&session, lease).await?;
store
.append_task_event_for_lease(
&session.id,
lease,
&TaskEventKind::StatusChanged {
from,
to: TaskSessionStatus::Running,
reason: session.status_reason.clone(),
},
)
.await?;
store
.append_task_event_for_lease(&session.id, lease, &TaskEventKind::Started)
.await?;
reconcile_stale_deliveries(&store, ChildTarget::Task(&session.id, lease)).await?;
// Claim before choosing the flow: a durable `CiFix` command is what decides
// whether this generation is a ci-fix body, and that choice happens before a
// harness exists. The claim also reassigns a predecessor's still-`Claimed`
// wake to this generation, which is how a crashed repair resumes on the same
// command rather than silently reverting to the lifecycle phase.
let mut seen_commands = HashSet::new();
let claimed = claim_commands(&store, &session, lease, &mut seen_commands).await?;
// A woken open-PR Task whose claimed wake still names a failing current head
// runs the single-step `ci-fix` flow; every other launch resumes the standard
// task lifecycle phase. The wake stays `Claimed` for the whole turn — settling
// it belongs to the repair's exit, not its entry.
let (mut ci_fix_wake, commands) = arm_ci_fix_wake(&store, &session, lease, claimed).await?;
let mut flow = if ci_fix_wake.is_some() {
Playhead::new(QueuedInvocation::load(&session.worktree, "ci-fix")?).0
} else {
let mut flow = resume_task_phase(&session)?;
// Reconcile any interactive handoff the parent's prior body opened before
// this body runs a provider turn. A completed outcome advances past work the
// human finished, a hand-back resumes the same step, and a still-open handoff
// parks the parent without ever starting the provider.
if reconcile_interactive_rendezvous_at_birth(&store, &mut session, lease, &mut flow).await?
{
let outcome = ChildBodyOutcome::Interrupted {
reason: session.status_reason.clone(),
};
return finish_parked(&store, &mut session, lease, None, outcome, None).await;
}
flow
};
let mut prepared = prepare_task_flow_step(
&store,
&mut session,
lease,
wave.name(),
&flow,
ci_fix_wake.as_ref(),
)
.await?;
let (harness_name, _) = crate::engine::config::parse_agent(&session.agent);
let (event_tx, mut event_rx) = mpsc::unbounded_channel();
let mut harness = create_harness(&harness_name, ApprovalPolicy::AutoApprove, event_tx)?;
harness.set_provider_session_id(session.provider_session_id.clone());
store
.validate_child_write_lease(&ChildRef::Task(session.id.clone()), lease)
.await?;
harness.start(&prepared.turn.config).await?;
session.provider = harness_name;
session.provider_session_id = harness.provider_session_id();
if let Some(process) = &mut session.latest_process {
process.observe_provider(
&session.provider,
session.provider_session_id.clone(),
harness.process_group_id(),
);
}
if let Err(error) = store.update_task_session_for_lease(&session, lease).await {
let _ = harness.stop().await;
return Err(error.into());
}
let mut state_fingerprint = task_state_fingerprint(&session)?;
let mut iteration_start_head = pr_head_for_session(&store, &session).await?;
let mut gate_fingerprint = if session.lifecycle_phase == TaskLifecyclePhase::Gate {
Some(task_gate_fingerprint(&session)?)
} else {
None
};
let mut pending = VecDeque::new();
// Claimed above, before the flow choice. `commands` is that same batch minus
// the ci-fix wake this body was born for, if any — absorb must never see the
// command it is already servicing.
if let Some(stop) = absorb_commands(
&store,
&session,
lease,
commands,
harness.as_mut(),
false,
&mut pending,
)
.await?
{
return finish_command_stop(&store, &mut session, lease, harness.as_mut(), stop, None)
.await;
}
let mut review_start = None;
let mut review_recovery = None;
let mut interaction_review = if let Some(review) = prepared.review.take() {
open_interaction_review_body(&store, &session, lease, &mut flow, &review).await?;
match (review.status, &review.reviewer) {
(InteractionReviewStatus::Requested, InteractionReviewer::Human) => {
store
.activate_human_interaction_review(&session, &review.id, lease)
.await?;
review_start = Some(PendingInput::system(prepared.turn.input.clone()));
}
(InteractionReviewStatus::Requested, _) => {}
(InteractionReviewStatus::Active, InteractionReviewer::Human) => {
review_start = Some(PendingInput::system(format!(
"Resume human interaction review {} after a Task process restart. Continue \
the `{}` exercise in this existing provider transcript. Human messages arrive as FIFO follow-ups; \
answer them here and record each answer with `lf task review reply {} \"<answer and evidence>\"`. \
The human finishes the checkpoint with `lf task review complete {0} --disposition \
approved|changes-requested --outcome \"<findings and evidence>\"`. The complete Task and skill \
context follows so recovery does not depend on the interrupted turn having reached the provider.\n\n{}",
review.id, review.step, review.id, prepared.turn.input
)));
}
(InteractionReviewStatus::Active, _) => {
review_recovery = Some(PendingInput::system(format!(
"Resume interaction review {} after a Task process restart. Inspect the \
existing provider transcript for the latest reviewer question, then answer it with \
`lf task review reply {} \"<answer and evidence>\"`.",
review.id, review.id
)));
}
(InteractionReviewStatus::Completed, _) => {
let disposition = review.disposition.ok_or_else(|| {
anyhow!(
"completed interaction review {} has no disposition",
review.id
)
})?;
let outcome = review.outcome.as_deref().ok_or_else(|| {
anyhow!("completed interaction review {} has no outcome", review.id)
})?;
review_recovery = Some(PendingInput::system(format!(
"Recover the already-completed interaction review {} with disposition `{}`. \
The durable reviewer outcome is:\n{}",
review.id,
disposition.as_str(),
outcome
)));
}
}
Some(review.id)
} else {
None
};
if let Some(review_start) = review_start {
pending.push_front(review_start);
}
if pending.is_empty() {
pending.extend(review_recovery);
}
// Record this body's turns the way `flowloop/wave.rs` does. Without it a
// Task Session's spend reaches no store at all: the provider runs in this
// process, so no child `lf` records on its behalf.
let capture = flow.current().and_then(|step| {
let context = crate::journal::trace_capture_context(
Path::new(&session.worktree),
Some(step.flow.clone()),
Some(step.step.clone()),
)?;
match crate::trace::CaptureHandle::begin(
context,
prepared.turn.context.clone(),
crate::trace::CaptureStart {
provider: prepared.turn.harness.clone(),
model: prepared.turn.model.clone(),
surface: "headless".to_string(),
input_op: "initial".to_string(),
gather_ms: prepared.turn.context_gather_ms,
render_ms: prepared.turn.context_render_ms,
raw_provider: true,
},
) {
Ok(capture) => Some(capture),
Err(error) => {
// Spend telemetry must never take a Task body down.
tracing::warn!(%error, "failed to establish Task trace capture");
None
}
}
});
if let Some(capture) = &capture {
capture.set_provider_session_id(session.provider_session_id.clone());
}
let mut flow_turn_active = false;
let mut provider_turn_active =
apply_next_pending(&store, &session, lease, harness.as_mut(), &mut pending).await?;
if !provider_turn_active && interaction_review.is_none() {
start_task_flow_turn(
&store,
&mut session,
lease,
harness.as_mut(),
&mut flow,
prepared.turn,
)
.await?;
flow_turn_active = true;
provider_turn_active = true;
}
let (attachment_tx, mut attachment_rx) = mpsc::unbounded_channel();
std::thread::spawn(move || {
for line in std::io::stdin().lock().lines() {
let Ok(line) = line else { break };
if attachment_tx.send(line).is_err() {
break;
}
}
});
println!(
"task {}> attached; /status, /interrupt [message], /detach, or type a message/instruction",
session.launch.issue.identifier
);
let mut command_poll = tokio::time::interval(Duration::from_millis(200));
command_poll.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
let mut last_text = String::new();
let mut turn_had_durable_side_effect = false;
// One preempt per provider turn; cleared with `provider_turn_active`.
let mut review_preempted = false;
'runner: loop {
tokio::select! {
line = attachment_rx.recv() => {
if let Some(line) = line {
handle_attachment(&store, &session, lease, line).await?;
}
}
_ = command_poll.tick() => {
let (wake, commands) = if provider_turn_active {
let claimed = claim_commands(
&store,
&session,
lease,
&mut seen_commands,
).await?;
let (ci_fix, commands): (Vec<_>, Vec<_>) = claimed
.into_iter()
.partition(|command| {
matches!(&command.kind, ChildCommandKind::CiFix { .. })
});
for command in &ci_fix {
seen_commands.remove(&command.id);
}
// A review turn is the agent waiting, so no `TurnCompleted`
// is coming to release the repair. Take that boundary once,
// and only for a wake the PR's current reading still names —
// a failure no repair can green is not `wake_legal`, so it
// yields no current incident and never reaches here.
if !review_preempted
&& interaction_review.is_some()
&& !ci_fix.is_empty()
&& holds_current_ci_fix_wake(&store, &session, &ci_fix).await?
{
harness.interrupt().await?;
review_preempted = true;
}
(None, commands)
} else if ci_fix_wake.is_none() {
claim_and_arm_ci_fix(
&store,
&session,
lease,
&mut seen_commands,
).await?
} else {
(
None,
claim_commands(&store, &session, lease, &mut seen_commands).await?,
)
};
if let Some(wake) = wake {
start_ci_fix_flow(
&store,
&mut session,
lease,
wave.name(),
harness.as_mut(),
&mut flow,
&wake,
).await?;
ci_fix_wake = Some(wake);
// The bounded repair owns this body's exit. The durable Gate
// review stays open for the next Task generation.
interaction_review = None;
flow_turn_active = true;
provider_turn_active = true;
last_text.clear();
}
if let Some(stop) = absorb_commands(
&store,
&session,
lease,
commands,
harness.as_mut(),
provider_turn_active,
&mut pending,
).await? {
return finish_command_stop(&store, &mut session, lease, harness.as_mut(), stop, capture.as_ref()).await;
}
if !provider_turn_active {
provider_turn_active = apply_next_pending(
&store,
&session,
lease,
harness.as_mut(),
&mut pending,
).await?;
}
}
event = event_rx.recv() => {
let Some(event) = event else {
return finish_failed(
&store,
&mut session,
lease,
harness.as_mut(),
"provider event stream closed",
capture.as_ref(),
).await;
};
if let Some(capture) = &capture {
capture.record_conversation(event.clone());
}
let provider_session_id = harness.provider_session_id();
if provider_session_id != session.provider_session_id {
session.provider_session_id = provider_session_id;
if let Some(process) = &mut session.latest_process {
process.observe_provider(
&session.provider,
session.provider_session_id.clone(),
harness.process_group_id(),
);
}
store.update_task_session_for_lease(&session, lease).await?;
}
match event {
ConversationEvent::TextDelta { content, .. } => last_text.push_str(&content),
ConversationEvent::TurnStarted { .. } => {
turn_had_durable_side_effect = false;
}
ConversationEvent::ItemCompleted { item, .. } => {
if matches!(
item,
ConversationItem::Command { .. } | ConversationItem::File { .. }
) {
turn_had_durable_side_effect = true;
}
}
ConversationEvent::TurnCompleted { status, .. } => {
provider_turn_active = false;
review_preempted = false;
if status == Lifecycle::Failed {
let reason = drain_turn_failure_reason(
&mut event_rx,
"provider turn failed",
);
return handle_body_failure(
&store,
&mut session,
lease,
harness.as_mut(),
&wave,
&reason,
turn_had_durable_side_effect,
capture.as_ref(),
)
.await;
}
if ci_fix_wake.is_none() {
let (wake, commands) = claim_and_arm_ci_fix(
&store,
&session,
lease,
&mut seen_commands,
).await?;
if let Some(wake) = wake {
start_ci_fix_flow(
&store,
&mut session,
lease,
wave.name(),
harness.as_mut(),
&mut flow,
&wake,
).await?;
ci_fix_wake = Some(wake);
// The repair takes the just-released provider
// boundary before Gate or lifecycle progression.
// Its durable review remains open for recovery.
interaction_review = None;
flow_turn_active = true;
provider_turn_active = true;
last_text.clear();
}
if let Some(stop) = absorb_commands(
&store,
&session,
lease,
commands,
harness.as_mut(),
provider_turn_active,
&mut pending,
).await? {
return finish_command_stop(
&store,
&mut session,
lease,
harness.as_mut(),
stop,
capture.as_ref(),
).await;
}
if provider_turn_active {
continue 'runner;
}
provider_turn_active = apply_next_pending(
&store,
&session,
lease,
harness.as_mut(),
&mut pending,
).await?;
if provider_turn_active {
continue 'runner;
}
}
// A repair never parks on the parent's rendezvous: ci-fix
// boot skips the birth reconcile, so a prior body's
// handoff is still pending here, and parking would write
// the `ci-fix` playhead into the phase's cursor — the
// same rejection by a second route. It stays pending for
// the next parent body, and the exit below parks the Task
// anyway, which is all this park wanted.
if ci_fix_wake.is_none()
&& flow_turn_active
&& status == Lifecycle::Completed
&& parked_on_interactive_handoff(&store, &session).await?
{
// The agent opened an interactive handoff this turn.
// Park without advancing: clear the active body as
// interrupted so the interactive step stays current,
// then end this body waiting on a human.
finish_task_flow_turn(&mut flow, Lifecycle::Interrupted)?;
record_task_flow_position(&mut session, &flow)?;
set_and_record_status(
&store,
&mut session,
lease,
TaskSessionStatus::Waiting,
"interactive handoff open; waiting for a human",
)
.await?;
let outcome = ChildBodyOutcome::Interrupted {
reason: session.status_reason.clone(),
};
return finish_parked(
&store,
&mut session,
lease,
Some(harness.as_mut()),
outcome,
capture.as_ref(),
)
.await;
}
let resume_interrupted_flow =
flow_turn_active && status == Lifecycle::Interrupted;
let completed_review = if flow_turn_active {
None
} else if let Some(review_id) = interaction_review.as_ref() {
completed_interaction_review(&store, review_id).await?
} else {
None
};
if completed_review.is_some() {
// Completion and its final FollowUp commit atomically. Claim after
// observing completion so that message cannot leak into the next phase.
let commands = claim_commands(
&store,
&session,
lease,
&mut seen_commands,
).await?;
if let Some(stop) = absorb_commands(
&store,
&session,
lease,
commands,
harness.as_mut(),
false,
&mut pending,
).await? {
return finish_command_stop(
&store,
&mut session,
lease,
harness.as_mut(),
stop,
capture.as_ref(),
).await;
}
if apply_next_pending(
&store,
&session,
lease,
harness.as_mut(),
&mut pending,
).await? {
provider_turn_active = true;
continue 'runner;
}
}
let review_body_completed = completed_review.is_some();
let mut flow_iteration_completed = if flow_turn_active {
finish_task_flow_turn(&mut flow, status)?
} else if review_body_completed {
interaction_review = None;
finish_task_flow_turn(&mut flow, Lifecycle::Completed)?
} else {
false
};
if flow_turn_active || review_body_completed {
let latest = store
.get_task_session(&session.id)
.await?
.ok_or_else(|| {
anyhow!("Task Session {} disappeared", session.id)
})?;
sync_terminal_task_state(&mut session, &latest);
if ci_fix_wake.is_none() {
record_task_flow_position(&mut session, &flow)?;
}
store.update_task_session_for_lease(&session, lease).await?;
}
if session.status == TaskSessionStatus::Abandoned {
let _ = harness.stop().await;
if let Some(process) = &mut session.latest_process {
process.state = ChildLeaseState::Finished;
process.outcome = Some(ChildBodyOutcome::Interrupted {
reason: session.status_reason.clone(),
});
}
store.finish_task_process(&session, lease).await?;
return Ok(());
}
// A ci-fix body is one bounded turn, and this is its exit.
// Standing above the lifecycle loop rather than at its
// tail is the whole fix, and every path below may assume
// no repair body is live — `settle_ci_fix_turn` documents
// why, and why a new lifecycle path belongs below it.
if let Some(wake) = ci_fix_wake.as_ref() {
// The flow ended, or the turn was cut short. Anything
// else is a repair still mid-flow.
if flow_iteration_completed || status == Lifecycle::Interrupted {
// Settlement judges head advancement against the
// authoritative remote head: a `Fresh` reconcile
// bypasses both the store observation cache and
// gh's HTTP cache, so the head the repair body
// just pushed is what settlement reads — never a
// warm pre-turn observation.
let observed_pr =
crate::ops::task::reconcile_task_pr_fresh_for_lease(
&store,
&mut session,
lease,
)
.await
.map_err(|error| anyhow!(error.to_string()))?;
let _ = harness.stop().await;
return settle_ci_fix_turn(
&store,
&mut session,
lease,
wake,
observed_pr.as_ref(),
iteration_start_head.as_deref(),
status,
capture.as_ref(),
)
.await;
}
}
flow_turn_active = false;
loop {
if flow_iteration_completed
&& session.lifecycle_phase == TaskLifecyclePhase::Kickoff
{
let reason = if matches!(
completed_review
.as_ref()
.map(|(disposition, _)| disposition),
Some(InteractionReviewDisposition::ChangesRequested)
) {
"Task kickoff requested changes; iteration is starting"
} else {
"Task kickoff approved; autonomous iteration is starting"
};
session.enter_iterate()?;
session.set_status(TaskSessionStatus::Running, reason);
store.update_task_session_for_lease(&session, lease).await?;
flow = resume_task_phase(&session)?;
flow_iteration_completed = false;
state_fingerprint = task_state_fingerprint(&session)?;
gate_fingerprint = None;
last_text.clear();
}
if matches!(
completed_review.as_ref().map(|(disposition, _)| disposition),
Some(InteractionReviewDisposition::ChangesRequested)
) && session.lifecycle_phase == TaskLifecyclePhase::Gate
{
state_fingerprint = task_state_fingerprint(&session)?;
gate_fingerprint = None;
session.enter_iterate()?;
session.set_status(
TaskSessionStatus::Running,
"Task gate requested changes; returning to iteration",
);
store.update_task_session_for_lease(&session, lease).await?;
let started = start_resumed_task_phase(
&store,
&mut session,
lease,
harness.as_mut(),
&mut flow,
wave.name(),
)
.await?;
interaction_review = started.review;
flow_turn_active = interaction_review.is_none();
provider_turn_active = started.provider_turn_active;
last_text.clear();
continue 'runner;
}
while let Some(input) = pending.pop_front() {
if !pending_input_is_current(&store, &session, lease, &input).await? {
continue;
}
if resume_interrupted_flow {
open_task_flow_body(&mut flow, &session)?;
flow_turn_active = true;
}
let command = input.command_id.map(|id| (id, input.effect));
apply_input(
&store,
&session,
lease,
harness.as_mut(),
&input.text,
command,
input.decision,
).await?;
provider_turn_active = true;
continue 'runner;
}
if interaction_review.is_some() {
last_text.clear();
continue 'runner;
}
let approved_gate = if flow_iteration_completed
&& session.lifecycle_phase == TaskLifecyclePhase::Gate
{
let next_gate_fingerprint = task_gate_fingerprint(&session)?;
if gate_fingerprint.as_ref() != Some(&next_gate_fingerprint) {
state_fingerprint = task_state_fingerprint(&session)?;
gate_fingerprint = None;
session.enter_iterate()?;
session.set_status(
TaskSessionStatus::Running,
"Task gate requested changes; returning to iteration",
);
store.update_task_session_for_lease(&session, lease).await?;
let started = start_resumed_task_phase(
&store,
&mut session,
lease,
harness.as_mut(),
&mut flow,
wave.name(),
)
.await?;
interaction_review = started.review;
flow_turn_active = interaction_review.is_none();
provider_turn_active = started.provider_turn_active;
last_text.clear();
continue 'runner;
}
Some(session.approved_gate_proposal()?)
} else {
None
};
if !flow_iteration_completed && status != Lifecycle::Interrupted {
let prepared = prepare_task_flow_step(
&store,
&mut session,
lease,
wave.name(),
&flow,
ci_fix_wake.as_ref(),
)
.await?;
let started = start_prepared_task_step(
&store,
&mut session,
lease,
harness.as_mut(),
&mut flow,
prepared,
)
.await?;
interaction_review = started.review;
flow_turn_active = interaction_review.is_none();
provider_turn_active = started.provider_turn_active;
continue 'runner;
}
let summary = progress_summary(&last_text);
let latest = store
.get_task_session(&session.id)
.await?
.ok_or_else(|| anyhow!("Task Session {} disappeared", session.id))?;
sync_terminal_task_state(&mut session, &latest);
session.current_directive_version = latest.current_directive_version;
session.incorporated_directive_version =
latest.incorporated_directive_version;
let pending_directive = unincorporated_directive_version(
session.current_directive_version,
session.incorporated_directive_version,
);
let observed_pr = crate::ops::task::reconcile_task_pr_for_lease(
&store,
&mut session,
lease,
)
.await
.map_err(|error| anyhow!(error.to_string()))?;
// Reconcile keeps the cached PR row through a GitHub
// outage and names the failure on the session. For a
// turn that just ran, that degraded reading is an
// infrastructure blocker: it could not have verified
// a repair.
let github_degraded = match &session.observation {
Observation::Degraded { reason, .. } => Some(reason.clone()),
_ => None,
};
// The head before this turn is the baseline for the
// no-change check; the head we just observed becomes
// the baseline for the next turn.
let head_before_turn = iteration_start_head;
iteration_start_head = observed_pr
.as_ref()
.and_then(|pr| pr.head_sha().map(str::to_string));
// Merged, not merely settled: the branch below reports
// this PR as merged and waits on its review gate, which
// is false of an abandoned one.
let merged_completing_pr = observed_pr.as_ref().is_some_and(|pr| {
pr.phase() == crate::task::PrPhase::Merged
&& pr
.publication
.as_ref()
.is_some_and(|publication| {
publication.after_merge
== crate::task::AfterMerge::CompleteTask
})
});
let needs_rotation = if merged_completing_pr {
// A completing PR settles the Task, never rotates to a next PR.
false
} else if observed_pr
.as_ref()
.is_some_and(|pr| pr.is_settled())
{
true
} else if observed_pr.is_none() {
store.active_task_pr(&session.id).await?.is_none()
} else {
false
};
let (stopped_status, stopped_reason) = if let Some(proposal) = approved_gate {
(proposal.status, proposal.reason)
} else if session.status == TaskSessionStatus::Completed {
(
TaskSessionStatus::Completed,
session.status_reason.clone(),
)
} else if let Some(version) = pending_directive {
(
TaskSessionStatus::Blocked,
format!(
"current directive v{version} was applied but not incorporated; resume the Task flow and acknowledge it before settling"
),
)
} else if status == Lifecycle::Interrupted {
(
TaskSessionStatus::Waiting,
"Task flow step interrupted; waiting for resume or another instruction".to_string(),
)
} else if merged_completing_pr {
// The PR merged to complete the Task, but a required review is
// still open. Wait for the gate to close before completion; do
// not rotate to another PR.
let number = observed_pr
.as_ref()
.and_then(|pr| pr.github())
.map(|github| github.number);
let reason = match number {
Some(number) => format!(
"pull request #{number} merged; awaiting required review before completion"
),
None => "pull request merged; awaiting required review before completion"
.to_string(),
};
(TaskSessionStatus::Waiting, reason)
} else if needs_rotation {
crate::ops::task::ensure_working_pr_for_lease(
&store,
&mut session,
lease,
)
.await
.map_err(|error| anyhow!(error.to_string()))?;
session.status_reason =
"Task PR settled; starting the next PR".to_string();
store.update_task_session_for_lease(&session, lease).await?;
let prepared = prepare_task_flow_step(
&store,
&mut session,
lease,
wave.name(),
&flow,
ci_fix_wake.as_ref(),
)
.await?;
let started = start_prepared_task_step(
&store,
&mut session,
lease,
harness.as_mut(),
&mut flow,
prepared,
)
.await?;
interaction_review = started.review;
flow_turn_active = interaction_review.is_none();
provider_turn_active = started.provider_turn_active;
last_text.clear();
continue 'runner;
} else if let Some(pr) = observed_pr
.as_ref()
.filter(|pr| pr.phase() == PrPhase::Open)
{
let head_advanced =
match (head_before_turn.as_deref(), pr.head_sha()) {
// No baseline (the PR was opened during
// this turn): opening it is progress.
(None, _) => true,
(Some(start), Some(current)) => start != current,
(Some(_), None) => false,
};
crate::ops::task::decide_open_pr_status(
pr,
github_degraded.as_deref(),
head_advanced,
)
} else {
let next_fingerprint = task_state_fingerprint(&session)?;
if next_fingerprint != state_fingerprint {
state_fingerprint = next_fingerprint;
session.status_reason =
"Task flow changed the worktree; starting another iteration"
.to_string();
store.update_task_session_for_lease(&session, lease).await?;
let prepared = prepare_task_flow_step(
&store,
&mut session,
lease,
wave.name(),
&flow,
ci_fix_wake.as_ref(),
)
.await?;
let started = start_prepared_task_step(
&store,
&mut session,
lease,
harness.as_mut(),
&mut flow,
prepared,
)
.await?;
interaction_review = started.review;
flow_turn_active = interaction_review.is_none();
provider_turn_active = started.provider_turn_active;
last_text.clear();
continue 'runner;
}
(
TaskSessionStatus::Blocked,
"Task flow completed without a PR or any worktree change; another automatic iteration would spin".to_string(),
)
};
if session.lifecycle_phase == TaskLifecyclePhase::Iterate
&& status != Lifecycle::Interrupted
{
let waiting_for_ci = observed_pr.as_ref().is_some_and(|pr| {
pr.phase() == PrPhase::Open && !pr.review_ready()
});
session.enter_gate(TaskGateProposal {
status: stopped_status,
reason: stopped_reason,
})?;
if waiting_for_ci {
let number = observed_pr
.as_ref()
.and_then(|pr| pr.github())
.map(|github| github.number);
let reason = match number {
Some(number) => format!(
"pull request #{number} is waiting for fresh passing required checks before Task review"
),
None => "pull request is waiting for fresh passing required checks before Task review"
.to_string(),
};
set_and_record_status(
&store,
&mut session,
lease,
TaskSessionStatus::Waiting,
reason,
)
.await?;
return finish_parked(
&store,
&mut session,
lease,
Some(harness.as_mut()),
ChildBodyOutcome::Completed,
capture.as_ref(),
)
.await;
}
session.set_status(
TaskSessionStatus::Running,
format!(
"Task outcome is awaiting gate cycle {}",
session.gate_cycle
),
);
gate_fingerprint = Some(task_gate_fingerprint(&session)?);
store.update_task_session_for_lease(&session, lease).await?;
let started = start_resumed_task_phase(
&store,
&mut session,
lease,
harness.as_mut(),
&mut flow,
wave.name(),
)
.await?;
interaction_review = started.review;
flow_turn_active = interaction_review.is_none();
provider_turn_active = started.provider_turn_active;
last_text.clear();
continue 'runner;
}
// Persist non-status fields while the generation is still active.
// The following transaction alone chooses commands or inactivity.
store.update_task_session_for_lease(&session, lease).await?;
let boundary = store
.claim_task_commands_or_stop_for_lease(
&session.id,
lease,
stopped_status,
&stopped_reason,
)
.await?;
let boundary_commands = match boundary {
BoundaryResult::Commands(commands) => {
filter_new_commands(commands, &mut seen_commands)
}
BoundaryResult::Stopped(stopped) => {
let _ = harness.stop().await;
let from = session.status;
session = stopped;
if !summary.is_empty() {
store.append_task_event_for_lease(
&session.id,
lease,
&TaskEventKind::Progress {
summary: summary.clone(),
},
).await?;
}
if session.status == TaskSessionStatus::Completed {
store.append_task_event_for_lease(
&session.id,
lease,
&TaskEventKind::Completed { summary },
).await?;
}
store.append_task_event_for_lease(
&session.id,
lease,
&TaskEventKind::StatusChanged {
from,
to: session.status,
reason: session.status_reason.clone(),
},
).await?;
if let Some(process) = &mut session.latest_process {
process.state = ChildLeaseState::Finished;
process.outcome = Some(if session.status == TaskSessionStatus::Completed {
ChildBodyOutcome::Completed
} else {
ChildBodyOutcome::Interrupted {
reason: session.status_reason.clone(),
}
});
}
store.finish_task_process(&session, lease).await?;
return Ok(());
}
};
let resume_requested = boundary_commands.iter().any(|command| {
matches!(&command.kind, ChildCommandKind::Resume { .. })
});
if let Some(stop) = absorb_commands(
&store,
&session,
lease,
boundary_commands,
harness.as_mut(),
false,
&mut pending,
).await? {
return finish_command_stop(
&store,
&mut session,
lease,
harness.as_mut(),
stop,
capture.as_ref(),
)
.await;
}
if resume_requested && pending.is_empty() {
let prepared = prepare_task_flow_step(
&store,
&mut session,
lease,
wave.name(),
&flow,
ci_fix_wake.as_ref(),
)
.await?;
let started = start_prepared_task_step(
&store,
&mut session,
lease,
harness.as_mut(),
&mut flow,
prepared,
)
.await?;
interaction_review = started.review;
flow_turn_active = interaction_review.is_none();
provider_turn_active = started.provider_turn_active;
continue 'runner;
}
}
}
ConversationEvent::Error { code, message } => {
let reason = format!("{code}: {message}");
return handle_body_failure(
&store,
&mut session,
lease,
harness.as_mut(),
&wave,
&reason,
turn_had_durable_side_effect,
capture.as_ref(),
)
.await;
}
ConversationEvent::ItemStarted { .. }
| ConversationEvent::ItemUpdated { .. }
| ConversationEvent::ReasoningDelta { .. }
| ConversationEvent::DiffUpdated { .. }
| ConversationEvent::TurnUsage { .. }
| ConversationEvent::SuggestedActions { .. }
| ConversationEvent::StatusChanged { .. } => {}
}
}
}
}
}
async fn prepare_task_flow_step(
store: &SharedStore,
session: &mut TaskSession,
lease: &ChildWriteLease,
wave_name: &str,
flow: &Playhead,
ci_fix: Option<&CiFixWake>,
) -> Result<PreparedTaskStep> {
let latest = store
.get_task_session(&session.id)
.await?
.ok_or_else(|| anyhow!("Task Session {} disappeared", session.id))?;
session.current_directive_version = latest.current_directive_version;
session.incorporated_directive_version = latest.incorporated_directive_version;
let directives = store
.child_directives(&ChildRef::Task(session.id.clone()))
.await?;
let directive = directives
.iter()
.find(|directive| directive.version == session.current_directive_version)
.ok_or_else(|| {
anyhow!(
"Task Session {} has no current directive v{}",
session.id,
session.current_directive_version
)
})?;
let step = flow
.current()
.ok_or_else(|| anyhow!("Task flow has no current step"))?;
if step.kind != StepKind::Skill {
anyhow::bail!(
"Task flow step {} is {:?}; durable Task flows currently require skills",
step.step,
step.kind
);
}
session.status_reason = format!(
"Task {} cycle {}, iteration {}, step {}/{}: {}",
session.lifecycle_phase.as_str(),
session.lifecycle_cycle(),
step.iteration + 1,
step.index + 1,
step.total,
step.step
);
store.update_task_session_for_lease(session, lease).await?;
let pr = store
.active_task_pr(&session.id)
.await?
.ok_or_else(|| anyhow!("Task Session {} has no active PR", session.id))?;
// The `ci-fix` step gets the failure seed from the wake command that selected
// it; every other Task-flow step gets the standard task seed. The flow and the
// wake are chosen together at birth, so a `ci-fix` step without a wake would
// mean the runner selected the flow from something other than the ledger.
let seed = match (step.step.as_str(), ci_fix) {
("ci-fix", Some(wake)) => ci_fix_seed(session, &pr, wake, wave_name),
("ci-fix", None) => {
anyhow::bail!(
"Task Session {} is running the ci-fix flow with no claimed ci-fix wake",
session.id
)
}
_ => task_seed(session, &pr, wave_name, directive),
};
let mut prepared =
crate::lf::commands::run::prepare_harness_turn(&step.step, &seed, wave_name, None)?;
prepared.config.agent = Some(session.agent.clone());
let skill = crate::engine::load_skill(&step.step, Path::new(&session.worktree))?;
let review = if skill.interactive.unwrap_or(false) {
let id = InteractionReviewId::new();
let policy = session.phase_plan().interaction_policy;
let (reviewer, prompt, reviewer_name) = match policy {
InteractionPolicy::Require => {
let protocol = human_interaction_review_protocol(&id, &step.step);
(
InteractionReviewer::Human,
format!(
"{protocol}\n\n{}",
skill
.content
.as_deref()
.unwrap_or("Follow the named skill.")
),
"Human",
)
}
InteractionPolicy::Defer => (
InteractionReviewer::Project(session.project_session_id.clone()),
interaction_review_prompt(
&id,
&step.step,
skill
.content
.as_deref()
.unwrap_or("Follow the named skill."),
),
"Project",
),
};
let request = InteractionReview {
id: id.clone(),
wave_id: session.wave_id.clone(),
project_session_id: session.project_session_id.clone(),
task_session_id: session.id.clone(),
phase: session.lifecycle_phase,
phase_epoch: session.phase_epoch,
flow: session.phase_plan().flow.clone(),
step: step.step.clone(),
step_index: step.index,
phase_iteration: step.iteration,
policy,
reviewer,
status: InteractionReviewStatus::Requested,
reason: session
.gate_proposal
.as_ref()
.map(|proposal| proposal.reason.clone())
.unwrap_or_else(|| session.status_reason.clone()),
prompt,
evidence: InteractionReviewEvidence {
worktree: session.worktree.clone(),
branch: pr.branch.clone(),
base_commit: pr.base_commit.clone(),
head_commit: crate::engine::git::rev_parse(&session.worktree, "HEAD")?,
worktree_fingerprint: task_state_fingerprint(session)?,
pr: pr.github().map(|github| InteractionReviewPr {
number: github.number,
url: github.url.clone(),
}),
},
requested_by_generation: lease.generation,
reviewer_generation: None,
disposition: None,
outcome: None,
requested_at: time::OffsetDateTime::now_utc(),
completed_at: None,
};
let review = store
.open_interaction_review(session, &request, lease)
.await?
.0;
if review.reviewer == InteractionReviewer::Human {
prepared.input.push_str("\n\n");
prepared
.input
.push_str(&human_interaction_review_protocol(&review.id, &step.step));
}
session.status_reason = format!(
"Task {} cycle {}, interactive step {} is {} in {reviewer_name} review {}",
session.lifecycle_phase.as_str(),
session.lifecycle_cycle(),
step.step,
review.status.as_str(),
review.id
);
store.update_task_session_for_lease(session, lease).await?;
Some(review)
} else {
None
};
Ok(PreparedTaskStep {
turn: prepared,
review,
})
}
fn human_interaction_review_protocol(review_id: &InteractionReviewId, skill: &str) -> String {
format!(
"Conduct the interactive `{skill}` exercise with the human in this existing Task provider \
session. Ask bounded questions and wait for their FIFO follow-up messages. Respond in this \
transcript, then record each answer with `lf task review reply {review_id} \
\"<answer and evidence>\"`. Do not approve yourself. The human finishes the checkpoint with \
`lf task review complete {review_id} --disposition approved|changes-requested --outcome \
\"<findings and evidence>\"`. Approval lets the lifecycle advance; requested changes return \
the same Task to Iterate."
)
}
fn interaction_review_prompt(
review_id: &InteractionReviewId,
skill: &str,
instructions: &str,
) -> String {
format!(
"Conduct the interactive `{skill}` exercise as the parent reviewer for this Task. \
Do not implement the child work yourself. Inspect the supplied evidence and apply the skill \
instructions from the reviewer role. Ask the Task a FIFO question with \
`lf project review message {review_id} \"<question>\"`. Finish with \
`lf project review complete {review_id} --disposition approved|changes-requested \
--outcome \"<findings and evidence>\"`.\n\n{instructions}"
)
}
fn open_task_flow_body(flow: &mut Playhead, session: &TaskSession) -> Result<()> {
let step = flow
.current()
.ok_or_else(|| anyhow!("Task flow has no current step"))?;
if step.kind != StepKind::Skill {
anyhow::bail!("Task flow step {} is not a skill", step.step);
}
flow.start_body(BodyProvenance::for_step(&step, &session.worktree))?;
Ok(())
}
async fn start_task_flow_turn(
store: &SharedStore,
session: &mut TaskSession,
lease: &ChildWriteLease,
harness: &mut dyn Harness,
flow: &mut Playhead,
prepared: crate::lf::commands::run::PreparedHarnessTurn,
) -> Result<()> {
open_task_flow_body(flow, session)?;
apply_input(store, session, lease, harness, &prepared.input, None, None).await?;
store
.mark_child_directive_applied_for_lease(
&ChildRef::Task(session.id.clone()),
lease,
session.current_directive_version,
)
.await?;
Ok(())
}
async fn open_interaction_review_body(
store: &SharedStore,
session: &TaskSession,
lease: &ChildWriteLease,
flow: &mut Playhead,
review: &InteractionReview,
) -> Result<()> {
if review.task_session_id != session.id
|| review.phase_epoch != session.phase_epoch
|| review.step_index != session.phase_cursor
|| review.phase_iteration != session.phase_iteration
{
anyhow::bail!(
"interaction review {} is stale for this Task step",
review.id
);
}
open_task_flow_body(flow, session)?;
store
.mark_child_directive_applied_for_lease(
&ChildRef::Task(session.id.clone()),
lease,
session.current_directive_version,
)
.await?;
Ok(())
}
async fn start_prepared_task_step(
store: &SharedStore,
session: &mut TaskSession,
lease: &ChildWriteLease,
harness: &mut dyn Harness,
flow: &mut Playhead,
mut prepared: PreparedTaskStep,
) -> Result<StartedTaskStep> {
if let Some(review) = prepared.review.take() {
open_interaction_review_body(store, session, lease, flow, &review).await?;
let provider_turn_active = if review.reviewer == InteractionReviewer::Human {
store
.activate_human_interaction_review(session, &review.id, lease)
.await?;
apply_input(
store,
session,
lease,
harness,
&prepared.turn.input,
None,
None,
)
.await?;
true
} else {
false
};
Ok(StartedTaskStep {
review: Some(review.id),
provider_turn_active,
})
} else {
start_task_flow_turn(store, session, lease, harness, flow, prepared.turn).await?;
Ok(StartedTaskStep {
review: None,
provider_turn_active: true,
})
}
}
async fn start_resumed_task_phase(
store: &SharedStore,
session: &mut TaskSession,
lease: &ChildWriteLease,
harness: &mut dyn Harness,
flow: &mut Playhead,
wave_name: &str,
) -> Result<StartedTaskStep> {
*flow = resume_task_phase(session)?;
let prepared = prepare_task_flow_step(store, session, lease, wave_name, flow, None).await?;
start_prepared_task_step(store, session, lease, harness, flow, prepared).await
}
fn finish_task_flow_turn(flow: &mut Playhead, status: Lifecycle) -> Result<bool> {
let body_id = flow
.active
.as_ref()
.map(|body| body.body_id.clone())
.ok_or_else(|| anyhow!("Task flow turn completed without an active body"))?;
let outcome = match status {
Lifecycle::Completed => StepOutcome::Completed,
Lifecycle::Interrupted => StepOutcome::Interrupted,
_ => anyhow::bail!("Task flow turn ended with unexpected status {status:?}"),
};
let events = flow.finish_body(&body_id, outcome, status.name())?;
Ok(events
.iter()
.any(|event| matches!(event, PlayheadEvent::InvocationCompleted { .. })))
}
async fn completed_interaction_review(
store: &SharedStore,
review_id: &InteractionReviewId,
) -> Result<Option<(InteractionReviewDisposition, String)>> {
let review = store
.get_interaction_review(review_id)
.await?
.ok_or_else(|| anyhow!("interaction review {review_id} disappeared"))?;
if review.status != InteractionReviewStatus::Completed {
return Ok(None);
}
let disposition = review
.disposition
.ok_or_else(|| anyhow!("completed interaction review {review_id} has no disposition"))?;
let outcome = review
.outcome
.ok_or_else(|| anyhow!("completed interaction review {review_id} has no outcome"))?;
Ok(Some((disposition, outcome)))
}
/// Reconcile the parent against any interactive handoff before this body runs a
/// provider turn. Returns `true` when the parent is parked on a human and the
/// body must end without starting a turn. A completed outcome advances the flow,
/// hand-back resumes the same step, and a failed handoff blocks the parent.
async fn reconcile_interactive_rendezvous_at_birth(
store: &SharedStore,
session: &mut TaskSession,
lease: &ChildWriteLease,
flow: &mut Playhead,
) -> Result<bool> {
let review_owns_current_step = store
.interaction_review_at(
&session.id,
session.phase_epoch,
session.phase_iteration,
session.phase_cursor,
)
.await?
.is_some();
let parent = InteractiveHandoffParent::Task(session.id.clone());
match interactive_rendezvous::resolve(store, &parent, lease.generation).await? {
Rendezvous::None => Ok(false),
Rendezvous::Waiting => {
set_and_record_status(
store,
session,
lease,
TaskSessionStatus::Waiting,
"interactive handoff open; waiting for a human",
)
.await?;
Ok(true)
}
Rendezvous::Resume { outcome, fresh } => {
resume_interactive_step(
store,
session,
lease,
flow,
outcome,
fresh,
review_owns_current_step,
)
.await
}
}
}
/// Resolve a terminal interactive handoff at body birth. Completion advances the
/// flow past work the human finished unless an InteractionReview owns the same
/// step; a handoff can wake that review but cannot decide it. Hand-back resumes
/// the same step, and failure blocks the parent for an operator. Evidence is
/// recorded once, on the generation that wins the wake claim (`fresh`).
async fn resume_interactive_step(
store: &SharedStore,
session: &mut TaskSession,
lease: &ChildWriteLease,
flow: &mut Playhead,
outcome: InteractiveHandoffOutcome,
fresh: bool,
review_owns_current_step: bool,
) -> Result<bool> {
if fresh {
let detail = match &outcome {
InteractiveHandoffOutcome::Completed { summary }
| InteractiveHandoffOutcome::HandedBack { summary } => summary.clone(),
InteractiveHandoffOutcome::Failed { reason } => reason.clone(),
};
store
.append_task_event_for_lease(
&session.id,
lease,
&TaskEventKind::Progress {
summary: format!(
"interactive handoff {}: {detail}",
outcome.status().as_str()
),
},
)
.await?;
}
match outcome {
InteractiveHandoffOutcome::Failed { reason } => {
set_and_record_status(
store,
session,
lease,
TaskSessionStatus::Blocked,
format!("interactive handoff failed: {reason}"),
)
.await?;
Ok(true)
}
InteractiveHandoffOutcome::Completed { .. } if !review_owns_current_step => {
advance_past_interactive_step(flow, session)?;
record_task_flow_position(session, flow)?;
store.update_task_session_for_lease(session, lease).await?;
Ok(false)
}
InteractiveHandoffOutcome::Completed { .. }
| InteractiveHandoffOutcome::HandedBack { .. } => Ok(false),
}
}
/// Advance the flow cursor one step past the resolved interactive step, reusing
/// the ordinary body start/finish path so the playhead settles exactly as it
/// would after a completed provider turn.
fn advance_past_interactive_step(flow: &mut Playhead, session: &TaskSession) -> Result<()> {
open_task_flow_body(flow, session)?;
finish_task_flow_turn(flow, Lifecycle::Completed)?;
Ok(())
}
/// True when the parent has an unresolved interactive handoff — open, or terminal
/// but not yet woken. The agent opened one this turn, so the parent must park
/// rather than advance: a still-open handoff waits on a human, and a
/// completed-this-turn handoff must be woken exactly once by the next body's birth
/// reconcile, not advanced here (which would skip the following step).
async fn parked_on_interactive_handoff(store: &SharedStore, session: &TaskSession) -> Result<bool> {
let parent = InteractiveHandoffParent::Task(session.id.clone());
let handoffs = store.list_interactive_handoffs(Some(&parent)).await?;
Ok(interactive_rendezvous::pending(&handoffs).is_some())
}
/// End a parked body: the session status is already `Waiting` or `Blocked`, so only
/// the process is settled and the parent stays non-terminal. The caller supplies
/// `outcome` — only it knows whether the turn finished or was cut short.
/// Close the body's trace capture so its last turn's usage is persisted.
/// A capture left open leaves a `running` turn with NULL usage -- the orphan
/// row `lf doctor` reports -- so this runs on every terminal path.
fn finish_capture(capture: Option<&crate::trace::CaptureHandle>, outcome: &str) {
let Some(capture) = capture else { return };
if let Err(error) = capture.finish(outcome, false) {
tracing::warn!(%error, "failed to finish Task trace capture");
}
}
async fn finish_parked(
store: &SharedStore,
session: &mut TaskSession,
lease: &ChildWriteLease,
harness: Option<&mut dyn Harness>,
outcome: ChildBodyOutcome,
capture: Option<&crate::trace::CaptureHandle>,
) -> Result<()> {
finish_capture(capture, "completed");
if let Some(harness) = harness {
let _ = harness.stop().await;
}
if let Some(process) = &mut session.latest_process {
process.state = ChildLeaseState::Finished;
process.outcome = Some(outcome);
}
store.finish_task_process(session, lease).await?;
Ok(())
}
fn record_task_flow_position(session: &mut TaskSession, flow: &Playhead) -> Result<()> {
let root = flow
.stack
.first()
.ok_or_else(|| anyhow!("Task flow has no root invocation"))?;
if root.flow != session.phase_plan().flow {
anyhow::bail!(
"Task Session {} {} flow is {:?}, but its playhead is {:?}",
session.id,
session.lifecycle_phase.as_str(),
session.phase_plan().flow,
root.flow
);
}
session.phase_cursor = root.cursor;
session.phase_iteration = root.iteration;
session.updated_at = time::OffsetDateTime::now_utc();
Ok(())
}
fn resume_task_phase(session: &TaskSession) -> Result<Playhead> {
let (flow, _) = Playhead::resume_root(
QueuedInvocation::load(&session.worktree, &session.phase_plan().flow)?,
session.phase_cursor,
session.phase_iteration,
)?;
Ok(flow)
}
fn sync_terminal_task_state(session: &mut TaskSession, latest: &TaskSession) {
if latest.status.is_terminal() {
session.status = latest.status;
session.status_reason = latest.status_reason.clone();
session.status_at = latest.status_at;
session.pm_writeback = latest.pm_writeback.clone();
}
}
fn task_state_fingerprint(session: &TaskSession) -> Result<String> {
let state = crate::engine::git::worktree_state(Path::new(&session.worktree))?;
Ok(hex::encode(Sha256::digest(state.as_bytes())))
}
/// The active PR's current head SHA, or `None` when there is no active PR.
/// Captured at iteration boundaries as a GitHub-side progress baseline so the
/// runner can tell a no-change ci-fix (head unchanged) from a push (head
/// advanced) without relying on worktree churn.
async fn pr_head_for_session(store: &SharedStore, session: &TaskSession) -> Result<Option<String>> {
Ok(store
.active_task_pr(&session.id)
.await?
.and_then(|pr| pr.github().map(|g| g.head_sha.clone()))
.flatten())
}
fn task_gate_fingerprint(session: &TaskSession) -> Result<String> {
let state = crate::engine::git::material_worktree_state(Path::new(&session.worktree))?;
Ok(hex::encode(Sha256::digest(state.as_bytes())))
}
async fn pending_input_is_current(
store: &SharedStore,
session: &TaskSession,
lease: &ChildWriteLease,
input: &PendingInput,
) -> Result<bool> {
input_is_current(store, ChildTarget::Task(&session.id, lease), input).await
}
async fn apply_next_pending(
store: &SharedStore,
session: &TaskSession,
lease: &ChildWriteLease,
harness: &mut dyn Harness,
pending: &mut VecDeque<PendingInput>,
) -> Result<bool> {
while let Some(input) = pending.pop_front() {
if !pending_input_is_current(store, session, lease, &input).await? {
continue;
}
let command = input.command_id.map(|id| (id, input.effect));
apply_input(
store,
session,
lease,
harness,
&input.text,
command,
input.decision,
)
.await?;
return Ok(true);
}
Ok(false)
}
async fn handle_attachment(
store: &SharedStore,
session: &TaskSession,
lease: &ChildWriteLease,
line: String,
) -> Result<()> {
let line = line.trim();
if line.is_empty() {
return Ok(());
}
if line == "/status" {
println!(
"{} {} {}",
session.launch.issue.identifier,
session.status.as_str(),
session.status_reason
);
return Ok(());
}
if line == "/detach" {
let _ = std::process::Command::new("tmux")
.args(["detach-client"])
.status();
return Ok(());
}
if !line.starts_with("/interrupt") {
let review = store
.interaction_review_at(
&session.id,
session.phase_epoch,
session.phase_iteration,
session.phase_cursor,
)
.await?;
if let Some(review) = review.filter(|review| {
review.reviewer == InteractionReviewer::Human && !review.status.is_terminal()
}) {
let command = store
.send_human_interaction_review_message(
&review.id,
ChildCommandSource::Attachment,
line,
)
.await?;
println!("queued {} for human review {}", command.id, review.id);
return Ok(());
}
}
let kind = if let Some(message) = line.strip_prefix("/interrupt") {
let message = message.trim();
ChildCommandKind::Interrupt {
replacement: (!message.is_empty()).then(|| message.to_string()),
}
} else {
ChildCommandKind::Steer {
text: line.to_string(),
}
};
let command = ChildCommand::new(
ChildRef::Task(session.id.clone()),
ChildCommandSource::Attachment,
kind,
);
let replacement = match &command.kind {
ChildCommandKind::Steer { text } => Some(text.clone()),
ChildCommandKind::Interrupt {
replacement: Some(text),
} => Some(text.clone()),
_ => None,
};
let (superseded, directive_event) = if let Some(text) = replacement {
let latest = store
.get_task_session(&session.id)
.await?
.ok_or_else(|| anyhow!("Task Session {} disappeared", session.id))?;
let directive = ChildDirective::replacement(
ChildRef::Task(session.id.clone()),
latest.current_directive_version + 1,
text,
command.source.clone(),
command.id.clone(),
);
let superseded = store
.create_child_command_with_directive(&command, &directive)
.await?;
(
superseded,
Some((directive.id, directive.version, directive.kind)),
)
} else if matches!(&command.kind, ChildCommandKind::Interrupt { .. }) {
(
store.supersede_and_create_child_command(&command).await?,
None,
)
} else {
store.create_child_command(&command).await?;
(Vec::new(), None)
};
for command_id in superseded {
store
.append_task_event_for_lease(
&session.id,
lease,
&TaskEventKind::CommandChanged {
command_id,
state: ChildCommandState::Superseded,
effect: None,
error: None,
},
)
.await?;
}
if let Some((directive_id, version, directive_kind)) = directive_event {
store
.append_task_event_for_lease(
&session.id,
lease,
&TaskEventKind::DirectiveChanged {
directive_id,
version,
directive_kind,
},
)
.await?;
}
store
.append_task_event_for_lease(
&session.id,
lease,
&TaskEventKind::CommandChanged {
command_id: command.id.clone(),
state: ChildCommandState::Persisted,
effect: command.effect,
error: None,
},
)
.await?;
println!("queued {}", command.id);
Ok(())
}
async fn absorb_commands(
store: &SharedStore,
session: &TaskSession,
lease: &ChildWriteLease,
commands: Vec<ChildCommand>,
harness: &mut dyn Harness,
turn_active: bool,
pending: &mut VecDeque<PendingInput>,
) -> Result<Option<CommandStop>> {
absorb_child_commands(
store,
ChildTarget::Task(&session.id, lease),
commands,
harness,
turn_active,
pending,
)
.await
}
async fn claim_commands(
store: &SharedStore,
session: &TaskSession,
lease: &ChildWriteLease,
seen: &mut HashSet<ChildCommandId>,
) -> Result<Vec<ChildCommand>> {
let commands = store
.claim_child_commands_for_lease(&ChildRef::Task(session.id.clone()), lease)
.await?;
Ok(filter_new_commands(commands, seen))
}
async fn claim_and_arm_ci_fix(
store: &SharedStore,
session: &TaskSession,
lease: &ChildWriteLease,
seen: &mut HashSet<ChildCommandId>,
) -> Result<(Option<CiFixWake>, Vec<ChildCommand>)> {
let commands = claim_commands(store, session, lease, seen).await?;
arm_ci_fix_wake(store, session, lease, commands).await
}
async fn start_ci_fix_flow(
store: &SharedStore,
session: &mut TaskSession,
lease: &ChildWriteLease,
wave_name: &str,
harness: &mut dyn Harness,
flow: &mut Playhead,
wake: &CiFixWake,
) -> Result<()> {
*flow = Playhead::new(QueuedInvocation::load(&session.worktree, "ci-fix")?).0;
let prepared =
prepare_task_flow_step(store, session, lease, wave_name, flow, Some(wake)).await?;
start_task_flow_turn(store, session, lease, harness, flow, prepared.turn).await
}
fn filter_new_commands(
commands: Vec<ChildCommand>,
seen: &mut HashSet<ChildCommandId>,
) -> Vec<ChildCommand> {
commands
.into_iter()
.filter(|command| seen.insert(command.id.clone()))
.collect()
}
async fn record_unhandled_failure(
session_id: &TaskSessionId,
lease: &ChildWriteLease,
error: &anyhow::Error,
) {
let Some(store) = open_existing_store().await.map(Arc::new) else {
return;
};
let Ok(Some(mut session)) = store.get_task_session(session_id).await else {
return;
};
if !session.status.is_process_active()
|| session
.latest_process
.as_ref()
.map(|process| process.generation)
!= Some(lease.generation)
{
return;
}
let from = session.status;
let message = format!("task process failed: {error}");
session.set_status(TaskSessionStatus::Failed, &message);
if store
.update_task_session_for_lease(&session, lease)
.await
.is_err()
{
return;
}
let _ = store
.append_task_event_for_lease(
&session.id,
lease,
&TaskEventKind::StatusChanged {
from,
to: TaskSessionStatus::Failed,
reason: message.clone(),
},
)
.await;
let _ = store
.append_task_event_for_lease(
&session.id,
lease,
&TaskEventKind::Failed {
error: message.clone(),
resumable: true,
},
)
.await;
if let Some(process) = &mut session.latest_process {
process.state = ChildLeaseState::Finished;
process.outcome = Some(ChildBodyOutcome::Failed { reason: message });
}
let _ = store.finish_task_process(&session, lease).await;
}
/// Send `text` to the harness and record the driving command's fate: accepted on
/// success, failed (with the error propagated) otherwise. `command` is `None`
/// for the task seed, which has no command to reconcile.
async fn apply_input(
store: &SharedStore,
session: &TaskSession,
lease: &ChildWriteLease,
harness: &mut dyn Harness,
text: &str,
command: Option<(ChildCommandId, ChildCommandEffect)>,
decision: Option<DecisionResolution>,
) -> Result<()> {
let (command_id, effect) = command
.map(|(command_id, effect)| (Some(command_id), effect))
.unwrap_or((None, ChildCommandEffect::NextTurn));
apply_child_input(
store,
ChildTarget::Task(&session.id, lease),
harness,
PendingInput {
command_id,
text: text.to_string(),
effect,
decision,
},
)
.await
}
/// Apply a status transition and persist it: set the status, update the row, and
/// append the paired `StatusChanged` event.
async fn set_and_record_status(
store: &SharedStore,
session: &mut TaskSession,
lease: &ChildWriteLease,
status: TaskSessionStatus,
reason: impl Into<String>,
) -> Result<()> {
let from = session.status;
session.set_status(status, reason);
store.update_task_session_for_lease(session, lease).await?;
store
.append_task_event_for_lease(
&session.id,
lease,
&TaskEventKind::StatusChanged {
from,
to: status,
reason: session.status_reason.clone(),
},
)
.await?;
if status == TaskSessionStatus::Blocked {
if let Some(pr) = store.active_task_pr(&session.id).await? {
store
.mark_ci_incidents_blocked(
&pr.id,
time::OffsetDateTime::now_utc(),
&session.status_reason,
)
.await?;
}
}
Ok(())
}
async fn finish_failed(
store: &SharedStore,
session: &mut TaskSession,
lease: &ChildWriteLease,
harness: &mut dyn Harness,
error: &str,
capture: Option<&crate::trace::CaptureHandle>,
) -> Result<()> {
finish_capture(capture, "failed");
let _ = harness.stop().await;
set_and_record_status(store, session, lease, TaskSessionStatus::Failed, error).await?;
store
.append_task_event_for_lease(
&session.id,
lease,
&TaskEventKind::Failed {
error: error.to_string(),
resumable: true,
},
)
.await?;
if let Some(process) = &mut session.latest_process {
process.state = ChildLeaseState::Finished;
process.outcome = Some(ChildBodyOutcome::Failed {
reason: error.to_string(),
});
}
store.finish_task_process(session, lease).await?;
anyhow::bail!(error.to_string())
}
/// The Blocked reason for an infrastructure failure, naming the failing
/// capability and the safe next action. `pr_number` keeps the attached PR
/// visible so a resume after the capability recovers picks up the same PR.
fn infra_blocked_reason(capability: &str, detail: &str, pr_number: Option<u32>) -> String {
let pr_note = pr_number
.map(|n| format!(" Pull request #{n} stays attached."))
.unwrap_or_default();
format!("ci-fix blocked by {capability}: {detail}.{pr_note}")
}
/// Stop the body and transition the Task to Blocked for an infrastructure
/// failure (provider outage, GitHub observation failure), keeping the active PR
/// attached so a resume after the capability recovers picks up the same PR.
/// Returns `Ok(())` — a clean stop, not an error — so the runner does not also
/// record an unhandled failure.
async fn finish_infra_blocked(
store: &SharedStore,
session: &mut TaskSession,
lease: &ChildWriteLease,
harness: &mut dyn Harness,
capability: &str,
detail: &str,
) -> Result<()> {
let _ = harness.stop().await;
let pr_number = store
.active_task_pr(&session.id)
.await?
.and_then(|pr| pr.github().map(|g| g.number));
let reason = infra_blocked_reason(capability, detail, pr_number);
set_and_record_status(store, session, lease, TaskSessionStatus::Blocked, &reason).await?;
if let Some(process) = &mut session.latest_process {
process.state = ChildLeaseState::Finished;
process.outcome = Some(ChildBodyOutcome::Failed {
reason: reason.clone(),
});
}
store.finish_task_process(session, lease).await?;
Ok(())
}
/// Handle a body failure with disconnect-class recovery: classify the failure,
/// and if it's a disconnect/hollow-body with a configured backup agent, hand
/// the next generation to the backup instead of leaving the body failed for
/// the supervisor to respawn the same flaky provider.
#[allow(clippy::too_many_arguments)] // capture is a terminal-path output, not a knob
async fn handle_body_failure(
store: &SharedStore,
session: &mut TaskSession,
lease: &ChildWriteLease,
harness: &mut dyn Harness,
wave: &Wave,
reason: &str,
turn_had_durable_side_effect: bool,
capture: Option<&crate::trace::CaptureHandle>,
) -> Result<()> {
finish_capture(capture, "failed");
let wave_config = read_wave_config(Path::new(wave.repo()), wave.name());
let backup_agent = wave_config.as_ref().and_then(|c| c.backup_agent.as_deref());
let decision = classify_disconnect_recovery(
reason,
&session.agent,
turn_had_durable_side_effect,
backup_agent,
);
match decision {
RecoveryDecision::HandoffToBackup { agent, provider } => {
let _ = harness.stop().await;
set_and_record_status(store, session, lease, TaskSessionStatus::Failed, reason).await?;
store
.append_task_event_for_lease(
&session.id,
lease,
&TaskEventKind::Failed {
error: reason.to_string(),
resumable: true,
},
)
.await?;
if let Some(process) = &mut session.latest_process {
process.state = ChildLeaseState::Finished;
process.outcome = Some(ChildBodyOutcome::Failed {
reason: reason.to_string(),
});
}
store.finish_task_process(session, lease).await?;
let request = ChildBodyHandoffRequest {
agent: agent.clone(),
provider: provider.clone(),
reason: format!(
"disconnect-class failure; handing off from {} to {agent}",
session.agent
),
};
*session = store.handoff_task_body(&session.id, &request).await?;
Ok(())
}
RecoveryDecision::Stop => {
let non_convergence = format!(
"{reason}; not replay-safe (durable side effects this turn) and no backup agent configured"
);
finish_failed(store, session, lease, harness, &non_convergence, None).await
}
RecoveryDecision::AllowRetry => {
finish_failed(store, session, lease, harness, reason, None).await
}
RecoveryDecision::Normal => {
// Not a disconnect-class failure — a provider outage during a
// PR/ci-fix iteration is an infrastructure blocker: keep the PR
// attached and block actionably so a resume when the provider
// recovers picks up the same PR. Without a PR, fall back to the
// generic failed path.
if store.active_task_pr(&session.id).await?.is_some() {
return finish_infra_blocked(store, session, lease, harness, "provider", reason)
.await;
}
finish_failed(store, session, lease, harness, reason, None).await
}
}
}
async fn finish_abandoned(
store: &SharedStore,
session: &mut TaskSession,
lease: &ChildWriteLease,
harness: &mut dyn Harness,
reason: String,
) -> Result<()> {
let _ = harness.interrupt().await;
let _ = harness.stop().await;
set_and_record_status(
store,
session,
lease,
TaskSessionStatus::Abandoned,
format!("Task Session explicitly abandoned: {reason}"),
)
.await?;
if let Some(process) = &mut session.latest_process {
process.state = ChildLeaseState::Finished;
process.outcome = Some(ChildBodyOutcome::Interrupted { reason });
}
store.finish_task_process(session, lease).await?;
Ok(())
}
async fn finish_command_stop(
store: &SharedStore,
session: &mut TaskSession,
lease: &ChildWriteLease,
harness: &mut dyn Harness,
stop: CommandStop,
capture: Option<&crate::trace::CaptureHandle>,
) -> Result<()> {
finish_capture(
capture,
match stop {
CommandStop::Interrupted => "interrupted",
_ => "completed",
},
);
match stop {
CommandStop::Interrupted => {
let _ = harness.stop().await;
set_and_record_status(
store,
session,
lease,
TaskSessionStatus::Waiting,
"Task turn interrupted; waiting for resume or another instruction",
)
.await?;
if let Some(process) = &mut session.latest_process {
process.state = ChildLeaseState::Finished;
process.outcome = Some(ChildBodyOutcome::Interrupted {
reason: "Task turn interrupted".to_string(),
});
}
store.finish_task_process(session, lease).await?;
Ok(())
}
CommandStop::Abandoned(reason) => {
finish_abandoned(store, session, lease, harness, reason).await
}
}
}
/// The durable ci-fix wake this generation was born to service.
///
/// Held for the body's life. The command it names stays `Claimed` throughout,
/// which is exactly what lets a successor generation reclaim it and re-derive
/// this same handle after a crash.
#[derive(Debug, Clone)]
pub(crate) struct CiFixWake {
/// Which durable command this body is repairing.
///
/// The turn's whole span: claimed at the arm, settled by `settle_ci_fix_turn`
/// at the exit, and nothing in between transitions it. Naming the command —
/// rather than re-deriving the identity at the exit — is what makes the body
/// settle the wake it actually serviced, even if the PR moved on underneath.
pub command_id: ChildCommandId,
/// The incident this wake repairs. Carried forward from the `CiFix` command
/// so settlement can record the repaired head against the exact incident,
/// without re-deriving the identity from a PR that may have moved on.
pub incident_identity: String,
pub pr_number: u32,
pub head_sha: String,
pub failing_checks: Vec<CiCheck>,
}
/// Take the claimed ci-fix wake that names the PR's *current* failure, and arm
/// the body for it. Returns the wake plus the remaining commands for
/// `absorb_commands`.
///
/// Selection is by incident identity, never by position. Several wakes can be
/// claimable at once — a head that failed, was pushed to, and failed again mints
/// a fresh identity each time, and any of those commands can still be unsettled.
/// Taking the first one would seed the body with an obsolete head and failing
/// set, and then the *current* wake would be superseded as a stray, spending its
/// identity for good: `ensure_child_ci_fix_command` would find the spent command
/// and never relaunch, so the live failure would never be repaired. Every
/// non-matching wake is superseded here, where the reason is known to be
/// staleness rather than a live-body race.
///
/// The matching command is **left `Claimed`**. It is not delivered and not
/// accepted: there is no provider call here to be ambiguous about, so
/// `Delivering` would be a lie that `reconcile_stale_deliveries` later turns into
/// `Uncertain`, stranding an automatic wake on a human. `Claimed` is also what
/// makes the repair restartable — `claim_child_commands_in` reassigns
/// `persisted`/`claimed` rows, so a crash mid-turn hands this same command to the
/// next generation, which lands right back here and re-selects the ci-fix flow.
///
/// The identity of the failure this PR reads as *now*. `None` means no wake is
/// warranted: green, moved on, gone, or not `wake_legal`. One authority, shared
/// with [`holds_current_ci_fix_wake`], so the mint and the match cannot drift.
async fn current_ci_incident_identity(
store: &SharedStore,
session: &TaskSession,
) -> Result<Option<String>> {
Ok(store
.active_task_pr(&session.id)
.await?
.as_ref()
.and_then(crate::ops::task::current_ci_incident)
.map(|incident| incident.identity))
}
/// Whether any claimed wake names the PR's current failure.
///
/// The read-only half of [`arm_ci_fix_wake`]'s selection: the arm supersedes and
/// stamps, so it cannot answer this mid-turn without committing to a repair.
async fn holds_current_ci_fix_wake(
store: &SharedStore,
session: &TaskSession,
claimed: &[ChildCommand],
) -> Result<bool> {
let Some(current) = current_ci_incident_identity(store, session).await? else {
return Ok(false);
};
Ok(claimed.iter().any(|command| {
matches!(
&command.kind,
ChildCommandKind::CiFix {
incident_identity, ..
} if incident_identity == ¤t
)
}))
}
async fn arm_ci_fix_wake(
store: &SharedStore,
session: &TaskSession,
lease: &ChildWriteLease,
claimed: Vec<ChildCommand>,
) -> Result<(Option<CiFixWake>, Vec<ChildCommand>)> {
if !claimed
.iter()
.any(|command| matches!(command.kind, ChildCommandKind::CiFix { .. }))
{
return Ok((None, claimed));
}
let current = current_ci_incident_identity(store, session).await?;
let mut matched = None;
let mut remaining = Vec::with_capacity(claimed.len());
for command in claimed {
let ChildCommandKind::CiFix {
incident_identity, ..
} = &command.kind
else {
remaining.push(command);
continue;
};
if current.as_deref() == Some(incident_identity.as_str()) {
// At most one command can carry a given identity — `ensure_` is what
// guarantees it — so this matches once.
matched = Some(command);
continue;
}
ChildTarget::Task(&session.id, lease)
.supersede_command(
store,
command.id,
"the PR's current failure no longer matches this wake; the head or failing set moved on",
)
.await?;
}
let Some(command) = matched else {
return Ok((None, remaining));
};
let ChildCommandKind::CiFix {
incident_identity,
pr_number,
head_sha,
failing_checks,
} = command.kind.clone()
else {
unreachable!("matched a CiFix command");
};
// Absorb never sees this command, and absorb is what normally records a
// claim. Without this the wake would leave no trace in the event stream.
store
.append_task_event_for_lease(
&session.id,
lease,
&TaskEventKind::CommandChanged {
command_id: command.id.clone(),
state: ChildCommandState::Claimed,
effect: None,
error: None,
},
)
.await?;
// A body now exists for this failure. Body birth is the response milestone.
//
// A missed stamp fails the arm rather than running an unmeasurable repair.
// The wake was linked to this incident before any launch could happen, so a
// row that is gone now means the evidence was pruned underneath a live wake —
// and a repair whose response nothing records is precisely what the ledger
// exists to make impossible. The command stays `Claimed`, so a successor
// generation reclaims it once the ledger is coherent again.
if !store
.mark_ci_incident_responded(&incident_identity, time::OffsetDateTime::now_utc())
.await?
{
anyhow::bail!(
"ci-fix wake {} names incident {incident_identity}, which is no longer recorded; \
refusing to run a repair whose response nothing can measure",
command.id
);
}
Ok((
Some(CiFixWake {
command_id: command.id,
incident_identity,
pr_number,
head_sha,
failing_checks,
}),
remaining,
))
}
/// What a finished ci-fix turn does to the wake that started it. Named apart from
/// the writes so the verdict is decided before anything is durable.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum CiFixVerdict {
/// The turn left the PR in a good place.
Accept,
/// The turn ran and the PR is still blocked, or the reading was too degraded
/// to tell. Either way the failure is a human's now, not a retry loop's.
Fail,
/// The PR settled or went away; the repair stopped mattering.
Supersede,
/// Not an outcome — the turn never reached one, so recovery re-runs it.
LeaveClaimed,
}
/// End a bounded ci-fix turn: park the body, then settle the wake that started it.
///
/// The exit half of [`arm_ci_fix_wake`], and the only place a `CiFix` command goes
/// terminal by running. It settles `wake.command_id` — the command this body was
/// born for — rather than re-deriving the identity from a PR that may since name a
/// different failure. The verdict rides [`decide_open_pr_status`], which already
/// reads an open PR the way the lifecycle does: `Waiting` accepts, `Blocked` fails
/// with the same operator-facing reason, a settled or vanished PR supersedes.
///
/// Any terminal state spends the incident identity for good — `ensure_` mints no
/// second wake for an identity that already has one, however it settled. That is
/// the bound: one failing head, one automatic repair; a *new* failure is a new
/// identity and wakes again.
///
/// **The invariant.** The park and the settlement are two writes with no
/// transaction around them, so their order is the whole crash contract: park the
/// Session *before* terminalizing the wake. A death in that window then leaves the
/// wake `Claimed`, which is what [`relaunch_on_duplicate`] documents as recovery's
/// job. The reverse order would spend the identity while the Session still read
/// `Running`, and the successor — finding no claimable wake — would silently
/// resume the generic lifecycle and leave the PR red with nobody repairing it.
/// Interruption is the same case reached deliberately: no outcome, so no
/// settlement.
///
/// The body parks: no gate, no successor step, no PR rotation. A repair pushes to
/// an existing branch; advancing the Task on its behalf would credit a repair as
/// progress the Task never made.
///
/// **Where the caller must stand.** That last paragraph is only true because the
/// sole call site sits *above* the runner's parent-lifecycle loop, not at its
/// tail. Every path in that loop reads, validates, or rewrites a cursor a repair
/// body does not own: its playhead is `ci-fix` while the phase's flow is
/// `task-kickoff`, `task`, or `task-gate`. Three phases proved it three
/// different ways — Gate and Iterate *rejected* the playhead
/// (W2-280/W2-298, W2-303), and Kickoff silently *replaced* it, spending a
/// `task_clarify` turn and stranding the wake `Claimed` (W2-309). So a new path
/// that touches the Task's lifecycle belongs **below** that call; one placed
/// above re-opens the class, and Kickoff is the proof it re-opens quietly —
/// it validated nothing, so it survived #1054's cursor guard with a green suite.
///
/// This decides *when* a repair ends, never *what* the head deserves:
/// `decide_open_pr_status` owns the verdict and the wake names its own incident
/// from [`arm_ci_fix_wake`], so nothing here re-derives "is this failing head
/// actionable?".
#[allow(clippy::too_many_arguments)] // capture is a terminal-path output, not a knob
async fn settle_ci_fix_turn(
store: &SharedStore,
session: &mut TaskSession,
lease: &ChildWriteLease,
wake: &CiFixWake,
observed_pr: Option<&crate::task::TaskPr>,
head_before_turn: Option<&str>,
status: Lifecycle,
capture: Option<&crate::trace::CaptureHandle>,
) -> Result<()> {
// The authoritative post-turn head, when it moved past the incident head. Set
// only when the fresh reconcile proved advancement, so it names the head the
// repair body shipped for this incident.
let mut repaired_head: Option<String> = None;
let (verdict, settled_status, reason) = match observed_pr
.filter(|pr| pr.phase() == PrPhase::Open)
{
Some(pr) if status != Lifecycle::Interrupted => {
let head_advanced = match (head_before_turn, pr.head_sha()) {
// No baseline: this body never saw a head to move.
(None, _) => false,
(Some(start), Some(current)) => start != current,
(Some(_), None) => false,
};
if head_advanced {
repaired_head = pr.head_sha().map(str::to_string);
}
// Reconcile names a degraded read on the session; for a turn that just
// ran, that reading could not have verified a repair.
let degraded = match &session.observation {
Observation::Degraded { reason, .. } => Some(reason.as_str()),
_ => None,
};
let (settled_status, reason) =
crate::ops::task::decide_open_pr_status(pr, degraded, head_advanced);
let verdict = match settled_status {
TaskSessionStatus::Blocked => CiFixVerdict::Fail,
_ => CiFixVerdict::Accept,
};
(verdict, settled_status, reason)
}
Some(_) => (
CiFixVerdict::LeaveClaimed,
TaskSessionStatus::Waiting,
format!(
"ci-fix turn on pull request #{} was interrupted; the repair resumes on resume",
wake.pr_number
),
),
None => (
CiFixVerdict::Supersede,
TaskSessionStatus::Waiting,
format!(
"pull request #{} settled or is no longer attached; the ci-fix wake no longer applies",
wake.pr_number
),
),
};
set_and_record_status(store, session, lease, settled_status, &reason).await?;
// The head the body shipped is durable attribution on the incident, tied to
// its wake command and generation. First-write in the store, so a retry or a
// later push never rewrites which head settled it.
if let Some(head) = &repaired_head {
store
.mark_ci_incident_repaired(
&wake.incident_identity,
head,
time::OffsetDateTime::now_utc(),
)
.await?;
}
let target = ChildTarget::Task(&session.id, lease);
match verdict {
CiFixVerdict::Accept => {
target
.accept_command(store, wake.command_id.clone(), None)
.await?
}
CiFixVerdict::Fail => {
target
.fail_command(store, wake.command_id.clone(), None, &reason)
.await?
}
CiFixVerdict::Supersede => {
target
.supersede_command(store, wake.command_id.clone(), &reason)
.await?
}
CiFixVerdict::LeaveClaimed => {}
}
// The body's outcome describes this turn, not the wake's verdict: a repair
// that ran to the end Completed, whatever it found. Only a turn cut short was
// Interrupted — and that is the one that leaves the wake reclaimable.
let outcome = if status == Lifecycle::Interrupted {
ChildBodyOutcome::Interrupted { reason }
} else {
ChildBodyOutcome::Completed
};
finish_parked(store, session, lease, None, outcome, capture).await
}
/// The seed for a `ci-fix` turn: the PR the skill must repair plus the failing
/// required checks (names + log URLs) so it resolves the exact failure on the
/// current head without re-deriving it.
///
/// The failure comes from the wake command, not from `pr.fresh_ci()`. The
/// observation row is mutable and moves on; the command is the immutable record
/// of the failure this body was born for. Re-reading the observation here let a
/// body repair a different failure than the one that woke it — the seed and the
/// dedup key now describe the same thing.
fn ci_fix_seed(
session: &TaskSession,
pr: &crate::task::TaskPr,
wake: &CiFixWake,
wave_name: &str,
) -> String {
let url = pr.github().map(|github| github.url.as_str()).unwrap_or("");
let number = wake.pr_number;
let head = wake.head_sha.as_str();
let failing = wake
.failing_checks
.iter()
.map(|check| match &check.url {
Some(link) => format!("- {} ({link})", check.name),
None => format!("- {}", check.name),
})
.collect::<Vec<_>>()
.join("\n");
format!(
"Fix the failing required CI checks on Linear task {identifier}'s open pull request.\n\n\
Run the ci-fix skill: reproduce the latest failure on the current head, make the smallest correct fix, run targeted then proportional checks, and push the same branch. Report an infrastructure or credential blocker rather than weakening tests.\n\n\
PR #{number}: {url}\nBranch: {branch}\nHead commit: {head}\nFailing required checks:\n{failing}\n\n\
Wave: {wave}\nTask Session: {session_id}\nWorktree: {worktree}\n\n\
Push fixes to the same branch; do not open a new PR or rotate the serial branch. When the push lands, the Task returns to waiting on the new head.",
identifier = session.launch.issue.identifier,
number = number,
url = url,
branch = pr.branch,
head = head,
failing = if failing.is_empty() { "- (none reported)".to_string() } else { failing },
wave = wave_name,
session_id = session.id,
worktree = session.worktree.display(),
)
}
fn task_seed(
session: &TaskSession,
pr: &crate::task::TaskPr,
wave_name: &str,
directive: &ChildDirective,
) -> String {
let placement = pr
.parent_pr_id
.as_ref()
.map(|parent| format!("Stack parent PR: {parent} (land the parent first)"))
.unwrap_or_else(|| "Stack parent PR: none (rooted on main)".to_string());
let gate_proposal = session
.gate_proposal
.as_ref()
.map(|proposal| {
format!(
"Gate proposal: {} — {}",
proposal.status.as_str(),
proposal.reason
)
})
.unwrap_or_else(|| "Gate proposal: none".to_string());
format!(
"Advance Linear task {identifier}: {title}\n\n{description}\n\nLinear Project: {project} ({project_id})\n{project_context}\n\nCurrent directive v{directive_version} ({directive_kind}):\n{directive_text}\n\nAcknowledge this direction before continuing with `lf task acknowledge {identifier} --directive {directive_version} --summary \"<how the plan changed>\"`.\n\nPM snapshot synced at: {snapshot_synced_at}\nWave: {wave}\nTask Session: {session_id}\nLifecycle phase: {lifecycle_phase} (epoch {phase_epoch}, gate cycle {gate_cycle})\nInteraction policy: {interaction_policy}\n{gate_proposal}\nWorktree: {worktree}\nPR {pr_sequence}: {pr_branch}\nBase commit: {base_commit}\n{placement}\n\nThis PR owns one serial branch. Bare `lf pr land --next <slug>` ships it and keeps the Task open; `lf pr land -c` proposes completing the Task after merge. `lf pr abandon` discards only this PR. `lf task complete {identifier} --summary \"...\"` proposes completion for clean work that needs no PR. Gate approves settlement or returns the same Task to iteration. If this PR already merged out of band and follow-up work remains, `lf pr next [slug]` rotates to the next serial PR, carrying committed and uncommitted follow-up forward. The runner owns branch rotation between PRs.",
identifier = session.launch.issue.identifier,
title = session.launch.issue.title,
description = session.launch.issue.description,
project = session.launch.project.name,
project_id = session.launch.project.id.as_str(),
project_context = session.launch.project.prompt_context,
directive_version = directive.version,
directive_kind = directive.kind.as_str(),
directive_text = directive.text,
snapshot_synced_at = session.launch.pm_snapshot_synced_at,
wave = wave_name,
session_id = session.id,
lifecycle_phase = session.lifecycle_phase.as_str(),
phase_epoch = session.phase_epoch,
gate_cycle = session.gate_cycle,
interaction_policy = session.phase_plan().interaction_policy.as_str(),
gate_proposal = gate_proposal,
worktree = session.worktree.display(),
pr_sequence = pr.sequence,
pr_branch = pr.branch,
base_commit = pr.base_commit,
placement = placement,
)
}
fn progress_summary(text: &str) -> String {
const MAX_CHARS: usize = 2_000;
let text = text.trim();
if text.chars().count() <= MAX_CHARS {
return text.to_string();
}
let mut summary: String = text.chars().take(MAX_CHARS - 1).collect();
summary.push('…');
summary
}
/// The failed-PR ci-fix lifecycle driven end to end. In-crate because the
/// functions it proves — `arm_ci_fix_wake` among them — are private to this
/// module tree; see the module's own header.
#[cfg(test)]
mod ci_fix_lifecycle_tests;
#[cfg(test)]
mod tests {
use std::collections::VecDeque;
use std::path::PathBuf;
use std::sync::Arc;
use anyhow::Result;
use async_trait::async_trait;
use loopflow_test_support::TestRepo;
use time::OffsetDateTime;
use super::{
absorb_commands, apply_input, apply_next_pending, ci_fix_seed, handle_attachment,
handle_body_failure, human_interaction_review_protocol, infra_blocked_reason,
interaction_review_prompt, prepare_task_flow_step, progress_summary, resume_task_phase,
run_task_session_inner, start_prepared_task_step, task_seed, CommandStop, PreparedTaskStep,
};
use crate::chat::types::{ConversationEvent, Lifecycle, TurnUsage};
use crate::child_session::{
ChildBodyHandoffRequest, ChildCommand, ChildCommandEffect, ChildCommandKind,
ChildCommandSource, ChildCommandState, ChildDecisionId, ChildDirective, ChildRef,
ChildWriteLease,
};
use crate::engine::agent::AgentConfig;
use crate::harness::{Capabilities, Harness};
use crate::id::WaveId;
use crate::interaction_review::InteractionReviewId;
use crate::project_session::{ProjectSession, ProjectSessionId, ProjectSessionStatus};
use crate::session_context::{
LinearIssueId, LinearIssueSnapshot, LinearProjectId, LinearProjectSnapshot,
ProjectLaunchReceipt, TaskLaunchReceipt,
};
use crate::store::{open_store, SharedStore, StorageConfig};
use crate::task::{
PmWritebackState, TaskEventKind, TaskGateProposal, TaskLifecyclePhase, TaskLifecyclePlan,
TaskPr, TaskPrId, TaskSession, TaskSessionId, TaskSessionStatus,
};
use crate::wave::playhead::Playhead;
use crate::wave::Wave;
struct ScriptedHarness {
supports_steer: bool,
sent: Vec<String>,
interrupts: usize,
fail_send: bool,
fail_interrupt: bool,
}
struct RunnerUsageHarness {
events: tokio::sync::mpsc::UnboundedSender<ConversationEvent>,
inputs: usize,
}
#[async_trait]
impl Harness for RunnerUsageHarness {
async fn start(&mut self, _config: &AgentConfig) -> Result<()> {
Ok(())
}
async fn send_input(&mut self, _content: &str) -> Result<()> {
self.inputs += 1;
if self.inputs == 1 {
self.events
.send(ConversationEvent::TurnStarted {
turn_id: "spending-turn".to_string(),
})
.unwrap();
self.events
.send(ConversationEvent::TurnCompleted {
turn_id: "spending-turn".to_string(),
status: Lifecycle::Completed,
})
.unwrap();
self.events
.send(ConversationEvent::TurnUsage {
turn_id: "spending-turn".to_string(),
usage: TurnUsage {
input_tokens: 321,
output_tokens: 45,
..TurnUsage::default()
},
})
.unwrap();
} else {
self.events
.send(ConversationEvent::Error {
code: "script_complete".to_string(),
message: "end the runner fixture".to_string(),
})
.unwrap();
}
Ok(())
}
async fn interrupt(&mut self) -> Result<()> {
Ok(())
}
async fn stop(&mut self) -> Result<()> {
Ok(())
}
fn capabilities(&self) -> Capabilities {
Capabilities {
supports_steer: false,
}
}
fn provider_session_id(&self) -> Option<String> {
Some("runner-provider-session".to_string())
}
}
impl ScriptedHarness {
fn new(supports_steer: bool) -> Self {
Self {
supports_steer,
sent: Vec::new(),
interrupts: 0,
fail_send: false,
fail_interrupt: false,
}
}
}
#[async_trait]
impl Harness for ScriptedHarness {
async fn start(&mut self, _config: &AgentConfig) -> Result<()> {
Ok(())
}
async fn send_input(&mut self, content: &str) -> Result<()> {
if self.fail_send {
anyhow::bail!("scripted send failed");
}
self.sent.push(content.to_string());
Ok(())
}
async fn interrupt(&mut self) -> Result<()> {
self.interrupts += 1;
if self.fail_interrupt {
anyhow::bail!("scripted interrupt failed");
}
Ok(())
}
async fn stop(&mut self) -> Result<()> {
Ok(())
}
fn capabilities(&self) -> Capabilities {
Capabilities {
supports_steer: self.supports_steer,
}
}
fn provider_session_id(&self) -> Option<String> {
Some("provider-session".to_string())
}
}
async fn seed_conformance_session(
store: SharedStore,
provider: &str,
worktree: PathBuf,
directive_text: Option<&str>,
) -> (TaskSession, ChildWriteLease) {
let wave = Wave::new(
WaveId::new(),
format!("wave-{provider}"),
"/repo".to_string(),
);
store.create_wave(&wave).await.unwrap();
let now = OffsetDateTime::from_unix_timestamp(OffsetDateTime::now_utc().unix_timestamp())
.unwrap();
let project_snapshot = LinearProjectSnapshot {
id: LinearProjectId::new(format!("project-{provider}")).unwrap(),
slug: "control".to_string(),
name: "Control".to_string(),
prompt_context: "Provider-neutral control".to_string(),
};
let project = ProjectSession {
id: ProjectSessionId::new(),
launch: ProjectLaunchReceipt {
project: project_snapshot.clone(),
pm_snapshot_synced_at: now.unix_timestamp(),
},
wave_id: wave.id().clone(),
current_directive_version: 0,
incorporated_directive_version: 0,
status: ProjectSessionStatus::Created,
status_reason: "reserved".to_string(),
status_at: now,
iteration: 0,
observation_cursor: 0,
last_state_fingerprint: None,
agent: provider.to_string(),
provider: provider.to_string(),
provider_session_id: None,
latest_process: None,
abandon_intent: None,
created_at: now,
updated_at: now,
};
store.create_project_session(&project).await.unwrap();
let mut session = TaskSession {
id: TaskSessionId::new(),
launch: TaskLaunchReceipt {
issue: LinearIssueSnapshot {
id: LinearIssueId::new(format!("issue-{provider}")).unwrap(),
identifier: format!("{provider}-123"),
title: "Conformance".to_string(),
description: "Exercise provider-neutral control".to_string(),
},
project: project_snapshot,
pm_snapshot_synced_at: now.unix_timestamp(),
},
pm_writeback: PmWritebackState::Current,
wave_id: wave.id().clone(),
project_session_id: project.id,
current_directive_version: 0,
incorporated_directive_version: 0,
status: TaskSessionStatus::Waiting,
status_reason: "ready for provider".to_string(),
status_at: now,
worktree,
workspace_slug: format!("test-{provider}"),
lifecycle: crate::task::TaskLifecyclePlan::standard("task"),
lifecycle_phase: crate::task::TaskLifecyclePhase::Iterate,
phase_epoch: 1,
phase_cursor: 0,
phase_iteration: 0,
gate_cycle: 0,
gate_proposal: None,
agent: provider.to_string(),
provider: provider.to_string(),
provider_session_id: Some("provider-session".to_string()),
latest_process: None,
abandon_intent: None,
created_at: now,
updated_at: now,
observation: crate::task::Observation::NotRequired,
};
let pr = TaskPr {
id: TaskPrId::new(),
task_session_id: session.id.clone(),
sequence: 1,
slug: session.workspace_slug.clone(),
branch: format!("test/{provider}"),
base_commit: "deadbeef".to_string(),
parent_pr_id: None,
publication: None,
merge_commit: None,
abandoned_at: None,
created_at: now,
updated_at: now,
ci_observation: None,
github_observation: None,
linear_attachment_id: None,
linear_comment_id: None,
linear_link_error: None,
};
store.create_task_session(&session, &pr).await.unwrap();
if let Some(text) = directive_text {
let command = ChildCommand::new(
ChildRef::Task(session.id.clone()),
ChildCommandSource::Human,
ChildCommandKind::Steer {
text: text.to_string(),
},
);
let directive = ChildDirective::replacement(
ChildRef::Task(session.id.clone()),
1,
text.to_string(),
command.source.clone(),
command.id.clone(),
);
store
.create_child_command_with_directive(&command, &directive)
.await
.unwrap();
session.current_directive_version = 1;
store.update_task_session(&session).await.unwrap();
}
session.begin_generation(format!("task-{provider}"));
let lease = store
.reserve_task_process(&session, TaskSessionStatus::Waiting)
.await
.unwrap()
.unwrap();
(session, lease)
}
async fn conformance_session(provider: &str) -> (SharedStore, TaskSession, ChildWriteLease) {
let dir = tempfile::tempdir().unwrap();
let path = dir.keep().join("registry.db");
let store = Arc::new(open_store(&StorageConfig::sqlite(path)).await.unwrap());
let (mut session, lease) = seed_conformance_session(
store.clone(),
provider,
PathBuf::from(format!("/repo.{provider}")),
None,
)
.await;
if let Some(process) = &mut session.latest_process {
process.state = crate::child_session::ChildLeaseState::Active;
}
session.set_status(TaskSessionStatus::Running, "provider active");
store.activate_task_process(&session, &lease).await.unwrap();
(store, session, lease)
}
#[tokio::test]
async fn task_runner_records_reported_turn_usage() {
let ledger = crate::journal::TestLedgerGuard::new();
let repo = TestRepo::new();
let store = Arc::new(
open_store(&StorageConfig::sqlite(ledger.home().join("loopflow.db")))
.await
.unwrap(),
);
let (session, lease) = seed_conformance_session(
store,
"codex",
repo.path().to_path_buf(),
Some("record this provider turn"),
)
.await;
crate::journal::emit(
repo.path(),
crate::journal::LfNode::Run,
crate::journal::LfEventType::Started,
crate::journal::LfEventFields::default(),
);
run_task_session_inner(
session.id,
&lease,
Box::new(|name, _approval, events| {
assert_eq!(name, "codex");
Ok(Box::new(RunnerUsageHarness { events, inputs: 0 }))
}),
)
.await
.unwrap();
let trace_store = crate::journal::open_ledger().unwrap();
let launches = trace_store.agent_launches_since(0).unwrap();
assert_eq!(launches.len(), 1);
assert_eq!(launches[0].provider, "codex");
assert_eq!(launches[0].model, None);
let turns = trace_store
.agent_turns_for_launches(&[launches[0].id.clone()])
.unwrap();
assert_eq!(turns.len(), 1);
assert_eq!(turns[0].provider_input_tokens, Some(321));
assert_eq!(turns[0].provider_output_tokens, Some(45));
crate::journal::emit(
repo.path(),
crate::journal::LfNode::Run,
crate::journal::LfEventType::Completed,
crate::journal::LfEventFields::default(),
);
}
async fn prepared_gate_review(
lifecycle: TaskLifecyclePlan,
) -> (
tempfile::TempDir,
SharedStore,
TaskSession,
ChildWriteLease,
Playhead,
PreparedTaskStep,
) {
let repo = tempfile::tempdir().unwrap();
for args in [
["init", "-b", "main"].as_slice(),
["config", "user.email", "loopflow@example.com"].as_slice(),
["config", "user.name", "Loopflow Test"].as_slice(),
] {
assert!(std::process::Command::new("git")
.args(args)
.current_dir(repo.path())
.status()
.unwrap()
.success());
}
std::fs::write(repo.path().join("README.md"), "review evidence\n").unwrap();
assert!(std::process::Command::new("git")
.args(["add", "README.md"])
.current_dir(repo.path())
.status()
.unwrap()
.success());
assert!(std::process::Command::new("git")
.args(["commit", "-m", "evidence"])
.current_dir(repo.path())
.status()
.unwrap()
.success());
let (store, mut session, lease) = conformance_session("codex").await;
let command = ChildCommand::new(
ChildRef::Task(session.id.clone()),
ChildCommandSource::Human,
ChildCommandKind::Steer {
text: "Prepare the Task for review".to_string(),
},
);
let directive = ChildDirective::replacement(
ChildRef::Task(session.id.clone()),
1,
"Prepare the Task for review".to_string(),
command.source.clone(),
command.id.clone(),
);
store
.create_child_command_with_directive(&command, &directive)
.await
.unwrap();
session.current_directive_version = 1;
session.worktree = repo.path().to_path_buf();
session.lifecycle = lifecycle;
session.lifecycle_phase = TaskLifecyclePhase::Gate;
session.phase_epoch = 3;
session.gate_cycle = 1;
session.gate_proposal = Some(TaskGateProposal {
status: TaskSessionStatus::Waiting,
reason: "prove the delivered behavior".to_string(),
});
store
.update_task_session_for_lease(&session, &lease)
.await
.unwrap();
let flow = resume_task_phase(&session).unwrap();
let prepared =
prepare_task_flow_step(&store, &mut session, &lease, "test-wave", &flow, None)
.await
.unwrap();
(repo, store, session, lease, flow, prepared)
}
#[test]
fn progress_summary_bounds_wave_visible_text() {
let summary = progress_summary(&"x".repeat(2_500));
assert_eq!(summary.chars().count(), 2_000);
assert!(summary.ends_with('…'));
}
#[test]
fn infra_blocked_reason_names_capability_and_keeps_pr_attached() {
// Provider outage: the capability and safe next action are visible,
// and the attached PR is named so a resume recovers onto the same PR.
let reason = infra_blocked_reason(
"provider",
"provider turn failed; resume when the provider recovers",
Some(900),
);
assert!(reason.contains("blocked by provider"), "reason: {reason}");
assert!(
reason.contains("resume when the provider recovers"),
"reason: {reason}"
);
assert!(reason.contains("#900 stays attached"), "reason: {reason}");
// GitHub observation failure uses the same shape with its capability.
let gh = infra_blocked_reason("github-observation", "gh pr checks: HTTP 502", Some(7));
assert!(gh.contains("blocked by github-observation"), "reason: {gh}");
assert!(gh.contains("#7 stays attached"), "reason: {gh}");
// No PR attached (e.g. a no-PR task that still hit an infra failure):
// the reason names the capability without a PR note.
let no_pr = infra_blocked_reason("provider", "turn failed", None);
assert!(no_pr.contains("blocked by provider"), "reason: {no_pr}");
assert!(!no_pr.contains("stays attached"), "reason: {no_pr}");
}
#[test]
fn deferred_review_prompt_assigns_the_skill_and_two_way_protocol() {
let review_id = InteractionReviewId::new();
let prompt = interaction_review_prompt(&review_id, "demo", "Prove each Done When.");
assert!(prompt.contains("interactive `demo` exercise"));
assert!(prompt.contains(&format!("lf project review message {review_id}")));
assert!(prompt.contains(&format!("lf project review complete {review_id}")));
assert!(prompt.contains("Prove each Done When."));
}
#[test]
fn human_review_prompt_keeps_the_decision_with_the_human() {
let review_id = InteractionReviewId::new();
let prompt = human_interaction_review_protocol(&review_id, "demo");
assert!(prompt.contains("existing Task provider session"));
assert!(prompt.contains("FIFO follow-up messages"));
assert!(prompt.contains(&format!("lf task review reply {review_id}")));
assert!(prompt.contains(&format!("lf task review complete {review_id}")));
assert!(prompt.contains("requested changes return"));
}
#[tokio::test]
async fn headless_interactive_step_opens_parent_review_with_current_evidence() {
let (repo, store, session, _lease, _flow, prepared) =
prepared_gate_review(TaskLifecyclePlan::headless("task")).await;
let review = prepared.review.expect("demo is deferred to the Project");
assert_eq!(review.phase, TaskLifecyclePhase::Gate);
assert_eq!(review.phase_epoch, 3);
assert_eq!(review.step, "demo");
assert_eq!(
review.reviewer,
crate::interaction_review::InteractionReviewer::Project(
session.project_session_id.clone()
)
);
assert_eq!(review.evidence.worktree, repo.path());
assert_eq!(
review.evidence.head_commit,
crate::engine::git::rev_parse(repo.path(), "HEAD").unwrap()
);
assert!(review.prompt.contains("lf project review message"));
assert!(review.prompt.contains("lf project review complete"));
assert!(session.status_reason.contains(review.id.as_str()));
assert!(session.status_reason.contains("Project review"));
assert_eq!(
store
.interaction_review_at(&session.id, 3, 0, 0)
.await
.unwrap()
.unwrap()
.id,
review.id
);
}
#[tokio::test]
async fn standard_interactive_step_starts_human_review_in_existing_provider_session() {
let (_repo, store, mut session, lease, mut flow, prepared) =
prepared_gate_review(TaskLifecyclePlan::standard("task")).await;
let review = prepared.review.clone().expect("demo requires human review");
let replayed =
prepare_task_flow_step(&store, &mut session, &lease, "test-wave", &flow, None)
.await
.unwrap();
assert_eq!(
replayed.review.as_ref().map(|review| &review.id),
Some(&review.id)
);
assert!(replayed.turn.input.contains(review.id.as_str()));
let mut harness = ScriptedHarness::new(true);
let started = start_prepared_task_step(
&store,
&mut session,
&lease,
&mut harness,
&mut flow,
replayed,
)
.await
.unwrap();
assert_eq!(
review.reviewer,
crate::interaction_review::InteractionReviewer::Human
);
assert_eq!(review.policy, crate::engine::InteractionPolicy::Require);
assert_eq!(started.review, Some(review.id.clone()));
assert!(started.provider_turn_active);
assert_eq!(harness.sent.len(), 1);
assert!(harness.sent[0].contains(review.id.as_str()));
assert!(harness.sent[0].contains("lf task review complete"));
assert!(flow.active.is_some());
assert_eq!(
store
.get_interaction_review(&review.id)
.await
.unwrap()
.unwrap()
.status,
crate::interaction_review::InteractionReviewStatus::Active
);
assert!(session.status_reason.contains("Human review"));
}
#[tokio::test]
async fn attached_human_review_input_is_fifo_followup_not_steer() {
let (_repo, store, session, lease, _flow, prepared) =
prepared_gate_review(TaskLifecyclePlan::standard("task")).await;
let review = prepared.review.expect("demo requires human review");
handle_attachment(
&store,
&session,
&lease,
"Show the login path from the product.".to_string(),
)
.await
.unwrap();
let commands = store
.list_child_commands(&ChildRef::Task(session.id.clone()))
.await
.unwrap();
let message = commands.last().expect("attached review message is durable");
assert_eq!(message.source, ChildCommandSource::Attachment);
assert!(matches!(
&message.kind,
ChildCommandKind::FollowUp { text }
if text.contains(review.id.as_str()) && text.contains("Show the login path")
));
assert_eq!(
store
.get_task_session(&session.id)
.await
.unwrap()
.unwrap()
.current_directive_version,
1
);
}
#[tokio::test]
async fn ci_fix_seed_carries_the_pr_and_the_failing_checks() {
let (_store, session, _lease) = conformance_session("codex").await;
let now = time::OffsetDateTime::now_utc();
let pr = TaskPr {
id: TaskPrId::new(),
task_session_id: session.id.clone(),
sequence: 1,
slug: "ship".to_string(),
branch: "jack/ship".to_string(),
base_commit: "base".to_string(),
parent_pr_id: None,
publication: Some(crate::task::PrPublication {
requested_at: now,
after_merge: crate::task::AfterMerge::Review,
next_slug: None,
github: Some(crate::task::GithubPr {
number: 916,
url: "https://github.com/loopflow/loopflow/pull/916".to_string(),
head_sha: Some("headsha".to_string()),
}),
}),
merge_commit: None,
abandoned_at: None,
// The observation has already moved on to a different head and a
// different failure — exactly the drift that used to reach the seed.
ci_observation: Some(crate::task::CiObservation {
head_sha: "movedhead".to_string(),
state: crate::task::CiState::Failing,
failing_checks: vec![crate::task::CiCheck {
name: "some-other-check".to_string(),
url: Some("https://ci/other".to_string()),
}],
observed_at: now,
}),
github_observation: None,
linear_attachment_id: None,
linear_comment_id: None,
linear_link_error: None,
created_at: now,
updated_at: now,
};
// The runner loads the `ci-fix` flow by name; it must be a registered builtin.
assert!(
crate::engine::builtins::get_builtin_flow("ci-fix").is_some(),
"the ci-fix builtin flow must resolve"
);
let wake = super::CiFixWake {
command_id: crate::child_session::ChildCommandId::new(),
incident_identity: "github:ci:test/repo:916:headsha:deadbeef".to_string(),
pr_number: 916,
head_sha: "headsha".to_string(),
failing_checks: vec![crate::task::CiCheck {
name: "rust-test".to_string(),
url: Some("https://ci/rust".to_string()),
}],
};
let seed = ci_fix_seed(&session, &pr, &wake, "product");
// The skill resolves the exact failure from the injected metadata.
assert!(seed.contains("#916"), "seed names the PR");
assert!(seed.contains("jack/ship"), "seed names the branch");
assert!(seed.contains("headsha"), "seed names the head commit");
assert!(
seed.contains("rust-test"),
"seed names the failing leaf check"
);
assert!(seed.contains("https://ci/rust"), "seed carries the log URL");
assert!(
seed.contains("ci-fix skill"),
"seed points at the ci-fix skill"
);
// The seed follows the wake command, not the observation row. The command
// is the immutable record of the failure this body was born for; the row
// is mutable and moves on. When they disagree the command wins, or a body
// repairs a failure other than the one that woke it.
assert!(
!seed.contains("movedhead"),
"seed must not carry the observation's newer head"
);
assert!(
!seed.contains("some-other-check"),
"seed must not carry the observation's newer failure"
);
}
#[tokio::test]
async fn failed_claude_task_hands_off_to_codex_with_directive_and_active_pr() {
let (store, mut failed, _lease) = conformance_session("claude").await;
let session_id = failed.id.clone();
let worktree = failed.worktree.clone();
failed.set_status(
TaskSessionStatus::Failed,
"Claude quota exhausted after preserving durable state",
);
store.update_task_session(&failed).await.unwrap();
let command = ChildCommand::new(
ChildRef::Task(session_id.clone()),
ChildCommandSource::Human,
ChildCommandKind::Steer {
text: "Keep the existing directive and continue PR2".to_string(),
},
);
let directive = ChildDirective::replacement(
ChildRef::Task(session_id.clone()),
1,
"Keep the existing directive and continue PR2".to_string(),
command.source.clone(),
command.id.clone(),
);
store
.create_child_command_with_directive(&command, &directive)
.await
.unwrap();
let active_pr_before = store.active_task_pr(&session_id).await.unwrap().unwrap();
let request = ChildBodyHandoffRequest {
agent: "codex".to_string(),
provider: "codex".to_string(),
reason: "Claude quota exhausted".to_string(),
};
let mut resumed = store
.handoff_task_body(&session_id, &request)
.await
.unwrap();
let active_pr_after = store.active_task_pr(&session_id).await.unwrap().unwrap();
let current_directive = store
.child_directives(&ChildRef::Task(session_id.clone()))
.await
.unwrap()
.into_iter()
.find(|candidate| candidate.version == resumed.current_directive_version)
.expect("current directive survives provider death");
let seed = task_seed(
&resumed,
&active_pr_after,
"wave-claude",
¤t_directive,
);
assert_eq!(resumed.id, session_id);
assert_eq!(resumed.worktree, worktree);
assert_eq!(resumed.agent, "codex");
assert_eq!(resumed.provider, "codex");
assert_eq!(resumed.provider_session_id, None);
assert_eq!(active_pr_after.id, active_pr_before.id);
assert_eq!(active_pr_after.branch, active_pr_before.branch);
assert!(seed.contains("Keep the existing directive and continue PR2"));
assert!(seed.contains(&active_pr_before.branch));
assert!(seed.contains(session_id.as_str()));
assert_eq!(resumed.begin_generation("lf-task-codex".to_string()), 2);
assert_eq!(resumed.latest_process.unwrap().generation, 2);
let events = store.task_events_after(&session_id, 0).await.unwrap();
assert!(matches!(
events.last().map(|event| &event.kind),
Some(TaskEventKind::BodyHandedOff { handoff })
if handoff.from_agent == "claude"
&& handoff.to_agent == "codex"
&& handoff.from_provider == "claude"
&& handoff.to_provider == "codex"
&& handoff.reason == "Claude quota exhausted"
));
}
#[tokio::test]
async fn attached_task_direction_is_versioned_before_provider_input() {
let (store, session, lease) = conformance_session("codex").await;
handle_attachment(&store, &session, &lease, "fix the parser first".to_string())
.await
.unwrap();
let current = store.get_task_session(&session.id).await.unwrap().unwrap();
let directives = store
.child_directives(&ChildRef::Task(session.id.clone()))
.await
.unwrap();
assert_eq!(current.current_directive_version, 1);
assert_eq!(directives.len(), 1);
assert_eq!(directives[0].text, "fix the parser first");
assert_eq!(directives[0].source, ChildCommandSource::Attachment);
assert!(directives[0].command_id.is_some());
}
#[tokio::test]
async fn provider_control_conformance_reports_honest_steer_effects() {
for (provider, supports_steer, expected_effect) in [
("codex", true, ChildCommandEffect::LiveSteer),
("claude", false, ChildCommandEffect::Replacement),
("opencode", false, ChildCommandEffect::Replacement),
] {
let (store, session, lease) = conformance_session(provider).await;
let command = ChildCommand::new(
ChildRef::Task(session.id.clone()),
ChildCommandSource::Human,
ChildCommandKind::Steer {
text: "change direction".to_string(),
},
);
store.create_child_command(&command).await.unwrap();
let commands = store
.claim_child_commands(&ChildRef::Task(session.id.clone()), 1)
.await
.unwrap();
let mut harness = ScriptedHarness::new(supports_steer);
let mut pending = VecDeque::new();
absorb_commands(
&store,
&session,
&lease,
commands,
&mut harness,
true,
&mut pending,
)
.await
.unwrap();
if let Some(input) = pending.pop_front() {
apply_input(
&store,
&session,
&lease,
&mut harness,
&input.text,
input.command_id.map(|id| (id, input.effect)),
input.decision,
)
.await
.unwrap();
}
let receipt = store.get_child_command(&command.id).await.unwrap().unwrap();
assert_eq!(receipt.state, ChildCommandState::Accepted, "{provider}");
assert_eq!(receipt.effect, Some(expected_effect), "{provider}");
assert_eq!(harness.sent, vec!["change direction"], "{provider}");
assert_eq!(
harness.interrupts,
usize::from(!supports_steer),
"{provider}"
);
}
}
#[tokio::test]
async fn task_follow_up_is_fifo_and_never_interrupts() {
for provider in ["codex", "claude", "opencode"] {
let (store, session, lease) = conformance_session(provider).await;
let first = ChildCommand::new(
ChildRef::Task(session.id.clone()),
ChildCommandSource::Human,
ChildCommandKind::FollowUp {
text: "first".to_string(),
},
);
let second = ChildCommand::new(
ChildRef::Task(session.id.clone()),
ChildCommandSource::Human,
ChildCommandKind::FollowUp {
text: "second".to_string(),
},
);
store.create_child_command(&first).await.unwrap();
store.create_child_command(&second).await.unwrap();
let commands = store
.claim_child_commands(&ChildRef::Task(session.id.clone()), 1)
.await
.unwrap();
let mut harness = ScriptedHarness::new(provider == "codex");
let mut pending = VecDeque::new();
absorb_commands(
&store,
&session,
&lease,
commands,
&mut harness,
true,
&mut pending,
)
.await
.unwrap();
assert_eq!(harness.interrupts, 0, "{provider}");
assert!(harness.sent.is_empty(), "{provider}");
for expected in ["first", "second"] {
assert!(
apply_next_pending(&store, &session, &lease, &mut harness, &mut pending,)
.await
.unwrap()
);
assert_eq!(harness.sent.last().map(String::as_str), Some(expected));
}
assert_eq!(
store
.get_child_command(&first.id)
.await
.unwrap()
.unwrap()
.effect,
Some(ChildCommandEffect::NextTurn),
"{provider}"
);
}
}
#[tokio::test]
async fn task_replacement_supersedes_queued_input() {
let (store, session, lease) = conformance_session("claude").await;
let first = ChildCommand::new(
ChildRef::Task(session.id.clone()),
ChildCommandSource::Human,
ChildCommandKind::FollowUp {
text: "A".to_string(),
},
);
let second = ChildCommand::new(
ChildRef::Task(session.id.clone()),
ChildCommandSource::Human,
ChildCommandKind::FollowUp {
text: "B".to_string(),
},
);
store.create_child_command(&first).await.unwrap();
store.create_child_command(&second).await.unwrap();
let mut harness = ScriptedHarness::new(false);
let mut pending = VecDeque::new();
let commands = store
.claim_child_commands(&ChildRef::Task(session.id.clone()), 1)
.await
.unwrap();
absorb_commands(
&store,
&session,
&lease,
commands,
&mut harness,
true,
&mut pending,
)
.await
.unwrap();
let replacement = ChildCommand::new(
ChildRef::Task(session.id.clone()),
ChildCommandSource::Human,
ChildCommandKind::Interrupt {
replacement: Some("C".to_string()),
},
);
store
.supersede_and_create_child_command(&replacement)
.await
.unwrap();
let commands = store
.claim_child_commands(&ChildRef::Task(session.id.clone()), 1)
.await
.unwrap();
absorb_commands(
&store,
&session,
&lease,
commands,
&mut harness,
true,
&mut pending,
)
.await
.unwrap();
let input = pending.pop_front().expect("replacement input");
assert_eq!(input.command_id.as_ref(), Some(&replacement.id));
assert_eq!(input.text, "C");
assert!(pending.is_empty());
assert_eq!(harness.interrupts, 1);
for superseded in [&first, &second] {
assert_eq!(
store
.get_child_command(&superseded.id)
.await
.unwrap()
.unwrap()
.state,
ChildCommandState::Superseded
);
}
}
#[tokio::test]
async fn bare_task_interrupt_stops_one_turn_without_abandoning_the_session() {
let (store, session, lease) = conformance_session("codex").await;
let command = ChildCommand::new(
ChildRef::Task(session.id.clone()),
ChildCommandSource::Human,
ChildCommandKind::Interrupt { replacement: None },
);
store.create_child_command(&command).await.unwrap();
let commands = store
.claim_child_commands(&ChildRef::Task(session.id.clone()), 1)
.await
.unwrap();
let mut harness = ScriptedHarness::new(true);
let mut pending = VecDeque::new();
let stop = absorb_commands(
&store,
&session,
&lease,
commands,
&mut harness,
true,
&mut pending,
)
.await
.unwrap();
assert_eq!(stop, Some(CommandStop::Interrupted));
assert_eq!(harness.interrupts, 1);
assert!(pending.is_empty());
assert_eq!(
store
.get_child_command(&command.id)
.await
.unwrap()
.unwrap()
.state,
ChildCommandState::Accepted
);
assert!(!session.status.is_terminal());
}
#[tokio::test]
async fn task_decisions_resume_every_provider_without_losing_lineage() {
for (provider, supports_steer) in [("codex", true), ("claude", false), ("opencode", false)]
{
let (store, session, lease) = conformance_session(provider).await;
let decision_id = ChildDecisionId::new();
let command = ChildCommand::new(
ChildRef::Task(session.id.clone()),
ChildCommandSource::Human,
ChildCommandKind::Decide {
decision_id: decision_id.clone(),
choice: "revise".to_string(),
message: Some("cover the boundary".to_string()),
},
);
store.create_child_command(&command).await.unwrap();
let commands = store
.claim_child_commands(&ChildRef::Task(session.id.clone()), 1)
.await
.unwrap();
let mut harness = ScriptedHarness::new(supports_steer);
let mut pending = VecDeque::new();
absorb_commands(
&store,
&session,
&lease,
commands,
&mut harness,
true,
&mut pending,
)
.await
.unwrap();
if let Some(input) = pending.pop_front() {
apply_input(
&store,
&session,
&lease,
&mut harness,
&input.text,
input.command_id.map(|id| (id, input.effect)),
input.decision,
)
.await
.unwrap();
}
assert_eq!(
harness.interrupts,
usize::from(!supports_steer),
"{provider}"
);
assert_eq!(
store
.get_child_command(&command.id)
.await
.unwrap()
.unwrap()
.effect,
Some(ChildCommandEffect::Decision),
"{provider}"
);
assert!(
store
.task_events_after(&session.id, 0)
.await
.unwrap()
.iter()
.any(|event| matches!(
&event.kind,
TaskEventKind::DecisionResolved {
decision_id: resolved,
choice,
message: Some(message),
} if resolved == &decision_id
&& choice == "revise"
&& message == "cover the boundary"
)),
"{provider}"
);
}
}
#[tokio::test]
async fn task_provider_control_failures_settle_the_receipt() {
let (store, session, lease) = conformance_session("claude").await;
let command = ChildCommand::new(
ChildRef::Task(session.id.clone()),
ChildCommandSource::Human,
ChildCommandKind::Steer {
text: "change direction".to_string(),
},
);
store.create_child_command(&command).await.unwrap();
let commands = store
.claim_child_commands(&ChildRef::Task(session.id.clone()), 1)
.await
.unwrap();
let mut harness = ScriptedHarness::new(false);
harness.fail_interrupt = true;
let error = absorb_commands(
&store,
&session,
&lease,
commands,
&mut harness,
true,
&mut VecDeque::new(),
)
.await
.expect_err("interrupt failure should fail control");
assert!(error.to_string().contains("scripted interrupt failed"));
let receipt = store.get_child_command(&command.id).await.unwrap().unwrap();
assert_eq!(receipt.state, ChildCommandState::Failed);
assert_eq!(receipt.effect, Some(ChildCommandEffect::Replacement));
assert!(receipt
.error
.as_deref()
.is_some_and(|error| error.contains("scripted interrupt failed")));
}
fn task_handoff_request(
session: &TaskSession,
lease: &ChildWriteLease,
) -> crate::interactive_handoff::OpenInteractiveHandoff {
crate::interactive_handoff::OpenInteractiveHandoff {
parent: crate::interactive_handoff::InteractiveHandoffParent::Task(session.id.clone()),
home: crate::engine::wave_home::WaveHome::parse("jack@local").unwrap(),
cwd: session.worktree.clone(),
provider: session.provider.clone(),
provider_session_id: session.provider_session_id.clone(),
body_generation: lease.generation,
reason: "Needs an interactive login".to_string(),
environment: std::collections::BTreeMap::new(),
attach_argv: vec!["tmux".to_string(), "attach".to_string()],
}
}
#[tokio::test]
async fn parked_on_interactive_handoff_tracks_the_open_row() {
let (store, session, lease) = conformance_session("codex").await;
assert!(!super::parked_on_interactive_handoff(&store, &session)
.await
.unwrap());
let (handoff, created) = store
.open_interactive_handoff(task_handoff_request(&session, &lease))
.await
.unwrap();
assert!(created);
assert!(super::parked_on_interactive_handoff(&store, &session)
.await
.unwrap());
// A terminal-but-unclaimed handoff still parks the body: the next birth
// reconcile must claim the wake exactly once, not this turn's advance.
store
.finish_interactive_handoff(
&handoff.id,
&crate::interactive_handoff::InteractiveHandoffOutcome::Completed {
summary: "human finished the login".to_string(),
},
)
.await
.unwrap();
assert!(super::parked_on_interactive_handoff(&store, &session)
.await
.unwrap());
// Once a generation claims the wake, the rendezvous is fully resolved and
// the parent runs normally.
assert!(store
.claim_interactive_handoff_wake(&handoff.id, lease.generation)
.await
.unwrap());
assert!(!super::parked_on_interactive_handoff(&store, &session)
.await
.unwrap());
}
#[tokio::test]
async fn handed_back_interactive_work_resumes_the_same_flow_step() {
let (store, mut session, lease) = conformance_session("codex").await;
let mut flow = resume_task_phase(&session).unwrap();
let (handoff, _) = store
.open_interactive_handoff(task_handoff_request(&session, &lease))
.await
.unwrap();
store
.finish_interactive_handoff(
&handoff.id,
&crate::interactive_handoff::InteractiveHandoffOutcome::HandedBack {
summary: "Finish the remaining review fixes".to_string(),
},
)
.await
.unwrap();
let parked = super::reconcile_interactive_rendezvous_at_birth(
&store,
&mut session,
&lease,
&mut flow,
)
.await
.unwrap();
assert!(!parked);
assert_eq!(session.phase_cursor, 0);
assert_eq!(session.phase_iteration, 0);
assert!(store
.get_interactive_handoff(&handoff.id)
.await
.unwrap()
.unwrap()
.wake_claimed_at
.is_some());
}
#[tokio::test]
async fn completed_handoff_cannot_advance_past_interaction_review() {
let (_repo, store, mut session, lease, mut flow, prepared) =
prepared_gate_review(TaskLifecyclePlan::standard("task")).await;
let review = prepared.review.expect("demo requires human review");
let (handoff, _) = store
.open_interactive_handoff(task_handoff_request(&session, &lease))
.await
.unwrap();
store
.finish_interactive_handoff(
&handoff.id,
&crate::interactive_handoff::InteractiveHandoffOutcome::Completed {
summary: "human finished the login".to_string(),
},
)
.await
.unwrap();
let parked = super::reconcile_interactive_rendezvous_at_birth(
&store,
&mut session,
&lease,
&mut flow,
)
.await
.unwrap();
assert!(!parked);
assert_eq!(session.phase_cursor, review.step_index);
assert_eq!(session.phase_iteration, review.phase_iteration);
assert_eq!(
store
.interaction_review_at(
&session.id,
session.phase_epoch,
session.phase_iteration,
session.phase_cursor,
)
.await
.unwrap()
.map(|current| current.id),
Some(review.id)
);
assert!(store
.get_interactive_handoff(&handoff.id)
.await
.unwrap()
.unwrap()
.wake_claimed_at
.is_some());
}
#[tokio::test]
async fn finish_parked_settles_the_body_without_a_terminal_status() {
let (store, mut session, lease) = conformance_session("codex").await;
session.set_status(
TaskSessionStatus::Waiting,
"interactive handoff open; waiting for a human",
);
store
.update_task_session_for_lease(&session, &lease)
.await
.unwrap();
let outcome = crate::child_session::ChildBodyOutcome::Interrupted {
reason: session.status_reason.clone(),
};
super::finish_parked(&store, &mut session, &lease, None, outcome, None)
.await
.unwrap();
assert_eq!(session.status, TaskSessionStatus::Waiting);
assert!(!session.status.is_terminal());
let process = session.latest_process.as_ref().unwrap();
assert_eq!(
process.state,
crate::child_session::ChildLeaseState::Finished
);
// The parked body leaves the Session durably non-terminal, so a later
// resume can reconcile the handoff outcome.
let persisted = store.get_task_session(&session.id).await.unwrap().unwrap();
assert_eq!(persisted.status, TaskSessionStatus::Waiting);
}
/// A disconnect-class failure with `backup_agent` configured hands the
/// next generation to the backup, records `BodyHandedOff`, retains the
/// failed opencode generation, and fences out a late write from the dead
/// body — all in one test so the recovery contract is visible end-to-end.
#[tokio::test]
async fn disconnect_failure_with_backup_hands_off_and_fences_old_writer() {
let repo = tempfile::tempdir().unwrap();
let wave_name = "wave-opencode";
let wave_dir = repo.path().join("wave").join(wave_name);
std::fs::create_dir_all(&wave_dir).unwrap();
std::fs::write(
wave_dir.join("GOAL.md"),
"---\nbackup_agent: claude:opus\n---\n\n# Goal\n",
)
.unwrap();
let store_path = repo.path().join("registry.db");
let store: SharedStore = Arc::new(
open_store(&StorageConfig::sqlite(store_path))
.await
.unwrap(),
);
let wave = Wave::new(
WaveId::new(),
wave_name.to_string(),
repo.path().display().to_string(),
);
store.create_wave(&wave).await.unwrap();
let now = OffsetDateTime::from_unix_timestamp(OffsetDateTime::now_utc().unix_timestamp())
.unwrap();
let project_snapshot = LinearProjectSnapshot {
id: LinearProjectId::new("project-opencode").unwrap(),
slug: "control".to_string(),
name: "Control".to_string(),
prompt_context: "test".to_string(),
};
let project = ProjectSession {
id: ProjectSessionId::new(),
launch: ProjectLaunchReceipt {
project: project_snapshot.clone(),
pm_snapshot_synced_at: now.unix_timestamp(),
},
wave_id: wave.id().clone(),
current_directive_version: 0,
incorporated_directive_version: 0,
status: ProjectSessionStatus::Created,
status_reason: "reserved".to_string(),
status_at: now,
iteration: 0,
observation_cursor: 0,
last_state_fingerprint: None,
agent: "opencode".to_string(),
provider: "opencode".to_string(),
provider_session_id: None,
latest_process: None,
abandon_intent: None,
created_at: now,
updated_at: now,
};
store.create_project_session(&project).await.unwrap();
let mut session = TaskSession {
id: TaskSessionId::new(),
launch: TaskLaunchReceipt {
issue: LinearIssueSnapshot {
id: LinearIssueId::new("issue-opencode").unwrap(),
identifier: "OP-123".to_string(),
title: "Conformance".to_string(),
description: "test".to_string(),
},
project: project_snapshot,
pm_snapshot_synced_at: now.unix_timestamp(),
},
pm_writeback: PmWritebackState::Current,
wave_id: wave.id().clone(),
project_session_id: project.id,
current_directive_version: 0,
incorporated_directive_version: 0,
status: TaskSessionStatus::Waiting,
status_reason: "ready".to_string(),
status_at: now,
worktree: PathBuf::from(repo.path().join("worktree").display().to_string()),
workspace_slug: "test-opencode".to_string(),
lifecycle: TaskLifecyclePlan::standard("task"),
lifecycle_phase: TaskLifecyclePhase::Iterate,
phase_epoch: 1,
phase_cursor: 0,
phase_iteration: 0,
gate_cycle: 0,
gate_proposal: None,
agent: "opencode:glm-5.2".to_string(),
provider: "opencode".to_string(),
provider_session_id: Some("provider-session".to_string()),
latest_process: None,
abandon_intent: None,
created_at: now,
updated_at: now,
observation: crate::task::Observation::NotRequired,
};
let pr = crate::task::TaskPr {
id: crate::task::TaskPrId::new(),
task_session_id: session.id.clone(),
sequence: 1,
slug: session.workspace_slug.clone(),
branch: "test/opencode".to_string(),
base_commit: "deadbeef".to_string(),
parent_pr_id: None,
publication: None,
merge_commit: None,
abandoned_at: None,
created_at: now,
updated_at: now,
ci_observation: None,
github_observation: None,
linear_attachment_id: None,
linear_comment_id: None,
linear_link_error: None,
};
store.create_task_session(&session, &pr).await.unwrap();
session.begin_generation("lf-task-opencode".to_string());
let lease = store
.reserve_task_process(&session, TaskSessionStatus::Waiting)
.await
.unwrap()
.unwrap();
if let Some(process) = &mut session.latest_process {
process.state = crate::child_session::ChildLeaseState::Active;
}
session.set_status(TaskSessionStatus::Running, "provider active");
store.activate_task_process(&session, &lease).await.unwrap();
// Drive a disconnect-class failure with the backup configured.
let mut harness = ScriptedHarness::new(false);
let result = handle_body_failure(
&store,
&mut session,
&lease,
&mut harness,
&wave,
"opencode_disconnected: stream died mid-turn",
true, // durable side effect → backup is the preferred path
None,
)
.await;
// The body handed off, not failed — Ok(()) so the supervisor
// relaunches with the backup agent.
assert!(result.is_ok(), "handoff should return Ok");
// The session now carries the backup agent.
assert_eq!(session.agent, "claude:opus");
assert_eq!(session.provider, "claude");
assert_eq!(
session.provider_session_id, None,
"provider session cleared on cross-provider handoff"
);
// The failed opencode generation is retained as evidence.
let process = session.latest_process.as_ref().expect("process retained");
assert_eq!(
process.state,
crate::child_session::ChildLeaseState::Finished
);
assert!(matches!(
&process.outcome,
Some(crate::child_session::ChildBodyOutcome::Failed { reason })
if reason.contains("opencode_disconnected")
));
// The BodyHandedOff event is in the ledger.
let events = store.task_events_after(&session.id, 0).await.unwrap();
assert!(
events.iter().any(|event| matches!(
&event.kind,
TaskEventKind::BodyHandedOff { handoff }
if handoff.from_agent == "opencode:glm-5.2"
&& handoff.to_agent == "claude:opus"
&& handoff.reason.contains("disconnect-class failure")
)),
"BodyHandedOff event must be recorded; events: {events:?}"
);
// Fencing: a late write from the dead opencode body is rejected.
// The process is Finished, so the old lease can no longer update.
let mut late_session = store.get_task_session(&session.id).await.unwrap().unwrap();
late_session.status_reason = "late write from dead body".to_string();
let write_result = store
.update_task_session_for_lease(&late_session, &lease)
.await;
assert!(
write_result.is_err(),
"a late write from the dead generation must be fenced out"
);
}
}