use std::collections::{BTreeMap, BTreeSet};
use std::path::{Component, Path, PathBuf};
use std::process::Command;
use std::sync::Arc;
use std::time::{Duration, Instant};
use crate::child_session::{
body_progress_age, count_recovery_attempts, observe, plan_body_recovery,
plan_stranded_recovery, task_write_lease_from_env, BodyEvidence, BodyRecoveryPlan,
ChildBodyOutcome, ChildLeaseState, ChildProcessGeneration, ChildRef, ChildWriteLease,
StrandedPlan, DEFAULT_STALL_AFTER, MAX_RECOVERY_ATTEMPTS,
};
use crate::durable::{AttentionRoute, AuthenticatedRequest, ControlCtx, Review};
use crate::engine::config::{load_config_or_default, parse_agent};
use crate::engine::git::{
checkout, checkout_new_branch_from, cherry_pick_range, current_branch, delete_local_branch,
fetch, get_default_branch, is_ancestor, is_clean, merge_base, push_with_upstream, ref_exists,
rev_parse, stash_including_untracked, stash_pop,
};
use crate::engine::naming::sanitize_for_branch;
use crate::engine::process::{
start_lf_session_with_env, tmux_installed, tmux_live_sessions, tmux_session_exists,
tmux_session_slug,
};
use crate::engine::worktrees::{
create_from_placement_plan, plan_placement, PlacementStrategy, WorktreeSegment,
};
use crate::engine::{expand_flow, load_flow, ConcreteStep};
use crate::ops::error::{OpsError, OpsResult};
use crate::session_context::{
LinearIssueId, LinearIssueSnapshot, LinearProjectId, LinearProjectSnapshot, TaskLaunchReceipt,
};
use crate::store::{
open_existing_store, open_registry_for_authority, RegistryUnavailable, SharedStore, StoreError,
};
use crate::task::actions::{
derive_task_actions, ReviewGateState, TaskActionEvidence, TaskActionModel,
};
use crate::task::{
AfterMerge, CiCheck, CiIncident, CiObservation, CiState, GithubObservation,
GithubObservationResult, GithubPr, Observation, PmWritebackOperation, PmWritebackState,
PrPhase, PrPublication, TaskEventKind, TaskPr, TaskPrId, TaskSession, TaskSessionId,
TaskSessionStatus, TaskSessionSuccession,
};
use crate::wave::Wave;
use sha2::{Digest, Sha256};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum TaskWaitUntil {
Open,
Terminal,
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct TaskLaunchOptions {
pub name: Option<String>,
pub flow: Option<String>,
pub stack_on: Option<String>,
pub directive: Option<String>,
pub headless: bool,
}
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)]
pub struct TaskControlResult {
pub issue_id: String,
pub session_id: String,
pub receipt: super::child::WorkControlReceipt,
pub observation: Observation,
}
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)]
pub struct TaskSessionSnapshot {
pub issue_id: String,
pub issue_identifier: String,
pub session_id: String,
pub project_id: String,
pub project: String,
pub pm_snapshot_synced_at: i64,
pub pm_writeback: crate::task::PmWritebackState,
pub wave: String,
pub project_session_id: String,
pub routing_project_session_id: Option<String>,
pub project_route_succeeded: bool,
pub predecessor_session_id: Option<String>,
pub successor_session_id: Option<String>,
pub status: TaskSessionStatus,
pub status_reason: String,
pub status_at: time::OffsetDateTime,
pub worktree: String,
pub workspace_slug: String,
pub lifecycle: crate::task::TaskLifecyclePlan,
pub lifecycle_phase: crate::task::TaskLifecyclePhase,
pub phase_epoch: u32,
pub phase_cursor: u32,
pub phase_iteration: u32,
pub gate_cycle: u32,
pub gate_proposal: Option<crate::task::TaskGateProposal>,
pub prs: Vec<TaskPr>,
pub active_pr: Option<TaskPrId>,
pub agent: String,
pub provider: String,
pub provider_session_id: Option<String>,
pub process_alive: bool,
pub latest_process: Option<ChildProcessGeneration>,
pub latest_event: Option<crate::task::TaskEvent>,
pub created_at: time::OffsetDateTime,
pub updated_at: time::OffsetDateTime,
pub observation: Observation,
pub actions: TaskActionModel,
}
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)]
pub struct TaskChangedFile {
pub path: String,
pub committed: bool,
pub staged: bool,
pub unstaged: bool,
pub untracked: bool,
}
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)]
pub struct TaskChangesSnapshot {
pub issue_identifier: String,
pub session_id: String,
pub base_commit: String,
pub head_commit: String,
pub files: Vec<TaskChangedFile>,
}
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)]
pub struct TaskDiffSnapshot {
pub issue_identifier: String,
pub session_id: String,
pub path: Option<String>,
pub patch: String,
pub binary: bool,
pub truncated: bool,
}
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)]
pub struct TaskFileSnapshot {
pub issue_identifier: String,
pub session_id: String,
pub path: String,
pub content: Option<String>,
pub binary: bool,
pub size_bytes: u64,
pub truncated: bool,
}
#[derive(Debug, Clone, Copy)]
struct TaskWorkspace<'a> {
issue_identifier: &'a str,
session_id: &'a crate::task::TaskSessionId,
worktree: &'a Path,
base_commit: &'a str,
}
impl<'a> TaskWorkspace<'a> {
fn new(session: &'a TaskSession, pr: &'a TaskPr) -> Self {
Self {
issue_identifier: &session.launch.issue.identifier,
session_id: &session.id,
worktree: &session.worktree,
base_commit: &pr.base_commit,
}
}
}
fn active_pr(session: &TaskSession) -> OpsResult<TaskPr> {
let session_id = session.id.clone();
block_on_task(async move {
task_store()
.await?
.active_task_pr(&session_id)
.await
.map_err(|error| task_error(format!("failed to read active PR: {error}")))?
.ok_or_else(|| task_error("Task Session has no active PR"))
})
}
fn task_error(message: impl Into<String>) -> OpsError {
OpsError::Message(message.into())
}
fn block_on_task<T>(future: impl std::future::Future<Output = OpsResult<T>>) -> OpsResult<T> {
tokio::runtime::Runtime::new()
.map_err(|error| task_error(format!("failed to build task runtime: {error}")))?
.block_on(future)
}
async fn task_store() -> OpsResult<SharedStore> {
open_existing_store().await.map(Arc::new).ok_or_else(|| {
task_error("no Loopflow registry on this machine; start the owning Wave first")
})
}
#[derive(Debug, Clone)]
pub struct StackedRebase {
pub fork_base: String,
pub child: TaskPr,
pub parent_branch: Option<String>,
}
pub fn task_stack(worktree: &Path) -> OpsResult<Option<StackedRebase>> {
block_on_task(async move {
let TaskAuthority::Authority { store, session, .. } =
resolve_task_authority(worktree).await?
else {
return Ok(None);
};
let Some(active) = store
.active_task_pr(&session.id)
.await
.map_err(|error| task_error(error.to_string()))?
else {
return Ok(None);
};
let Some(parent_id) = active.parent_pr_id.clone() else {
return Ok(None);
};
let parent = store
.get_task_pr(&parent_id)
.await
.map_err(|error| task_error(error.to_string()))?
.ok_or_else(|| task_error(format!("stack parent {parent_id} is missing")))?;
let mut parent_session = store
.get_task_session(&parent.task_session_id)
.await
.map_err(|error| task_error(error.to_string()))?
.ok_or_else(|| task_error("stack parent Task Session is missing"))?;
reconcile_task_pr(&store, &mut parent_session).await?;
let parent = store
.get_task_pr(&parent_id)
.await
.map_err(|error| task_error(error.to_string()))?
.ok_or_else(|| task_error(format!("stack parent {parent_id} disappeared")))?;
let merged = parent.merge_commit.is_some();
let closed = parent.abandoned_at.is_some();
if closed && !merged {
return Err(task_error(format!(
"stack parent {} closed without merging; re-place the child deliberately",
parent.branch
)));
}
Ok(Some(StackedRebase {
fork_base: active.base_commit.clone(),
child: active,
parent_branch: (!merged).then_some(parent.branch),
}))
})
}
pub fn stacked_collapse(worktree: &Path) -> OpsResult<Option<StackedRebase>> {
let stacked = task_stack(worktree)?;
if let Some(stacked) = &stacked {
if let Some(parent) = &stacked.parent_branch {
return Err(task_error(format!(
"Task PR is stacked on {parent}, which has not merged; land the parent first"
)));
}
}
Ok(stacked)
}
pub fn record_stack_rebase(
stacked: &StackedRebase,
new_base: &str,
clear_parent: bool,
) -> OpsResult<()> {
let pr_id = stacked.child.id.clone();
let new_base = new_base.to_string();
block_on_task(async move {
let store = Arc::new(
open_registry_for_authority()
.await
.map_err(registry_authority_error)?,
);
store
.rebase_task_pr(
&pr_id,
&new_base,
clear_parent,
time::OffsetDateTime::now_utc(),
)
.await
.map_err(|error| task_error(error.to_string()))?;
Ok(())
})
}
async fn owning_wave(store: &SharedStore, session: &TaskSession) -> OpsResult<Wave> {
store
.get_wave(&session.wave_id)
.await
.map_err(|error| task_error(format!("failed to read owning Wave: {error}")))?
.ok_or_else(|| task_error(format!("owning Wave {} is not registered", session.wave_id)))
}
fn _defer_task_interactions(session: &mut TaskSession) -> OpsResult<bool> {
if session.lifecycle.all_interactions_deferred() {
return Ok(false);
}
if session.status.is_terminal() {
return Err(task_error(format!(
"Task {} is {}; terminal Tasks cannot change interaction policy",
session.launch.issue.identifier,
session.status.as_str()
)));
}
if session.status.is_process_active() {
return Err(task_error(format!(
"Task {} has an active body; interrupt or wait for it before marking the Task headless",
session.launch.issue.identifier
)));
}
session.lifecycle.defer_all_interactions();
session.updated_at = time::OffsetDateTime::now_utc();
Ok(true)
}
fn _refuse_current_human_review(session: &TaskSession, review: Option<&Review>) -> OpsResult<()> {
if review.is_some_and(|review| review.attention == AttentionRoute::User) {
return Err(task_error(format!(
"Task {} routes its current interactive step to the User; close it before changing interaction policy",
session.launch.issue.identifier
)));
}
Ok(())
}
pub fn task_run(repo: &Path, issue: &str, options: TaskLaunchOptions) -> OpsResult<TaskSession> {
let TaskLaunchOptions {
name,
flow,
stack_on,
directive,
headless,
} = options;
let directive = directive
.map(|directive| {
let directive = directive.trim().to_string();
if directive.is_empty() {
Err(task_error("directive cannot be empty"))
} else {
Ok(directive)
}
})
.transpose()?;
let (existing, terminal_predecessor_id) = block_on_task(async {
let store = task_store().await?;
let mut existing = store
.get_task_session_by_issue(issue)
.await
.map_err(|error| task_error(format!("failed to read task registry: {error}")))?;
if let Some(session) = &mut existing {
if session.status.is_terminal() {
let predecessor_id = session.id.clone();
return Ok((None, Some(predecessor_id)));
}
if let Some(requested) = name.as_deref() {
let requested = parse_workspace_slug(requested)?;
if requested.as_str() != session.workspace_slug {
return Err(task_error(format!(
"Task {} already uses workspace name {:?}",
session.launch.issue.identifier, session.workspace_slug
)));
}
}
if let Some(requested) = flow.as_deref() {
let requested = resolve_task_flow(&session.worktree, Some(requested))?;
if requested != session.lifecycle.iterate.flow {
return Err(task_error(format!(
"Task {} already uses flow {:?}",
session.launch.issue.identifier, session.lifecycle.iterate.flow
)));
}
}
if let Some(requested) = stack_on.as_deref() {
let active = store
.active_task_pr(&session.id)
.await
.map_err(|error| task_error(error.to_string()))?
.ok_or_else(|| task_error("existing Task has no active PR"))?;
let parent_id = active.parent_pr_id.as_ref().ok_or_else(|| {
task_error(format!(
"Task {} is rooted on main, not stacked on {requested}",
session.launch.issue.identifier
))
})?;
let parent = store
.get_task_pr(parent_id)
.await
.map_err(|error| task_error(error.to_string()))?
.ok_or_else(|| task_error(format!("stack parent {parent_id} is missing")))?;
let parent_session = store
.get_task_session(&parent.task_session_id)
.await
.map_err(|error| task_error(error.to_string()))?
.ok_or_else(|| task_error("stack parent Task Session is missing"))?;
if requested != parent_session.launch.issue.identifier
&& requested != parent_session.launch.issue.id.as_str()
{
return Err(task_error(format!(
"Task {} is stacked on {}, not {requested}",
session.launch.issue.identifier, parent_session.launch.issue.identifier
)));
}
}
if directive.is_some() {
return Err(task_error(format!(
"Task {} already exists; use `lf task steer {} <new-direction>`",
session.launch.issue.identifier, session.launch.issue.identifier,
)));
}
if headless && !session.lifecycle.all_interactions_deferred() {
let work = store
.work_for_child(&ChildRef::Task(session.id.clone()))
.await
.map_err(|error| task_error(error.to_string()))?;
let review = store
.review(&work)
.await
.map_err(|error| task_error(error.to_string()))?;
_refuse_current_human_review(session, review.as_ref())?;
}
if headless && _defer_task_interactions(session)? {
store
.update_task_session(session)
.await
.map_err(|error| task_error(error.to_string()))?;
}
}
Ok((existing, None))
})?;
if let Some(existing) = existing {
return task_status(existing.launch.issue.id.as_str());
}
let main_repo = crate::ops::project::ensure_clean_main(repo, "Task start")
.map_err(|error| task_error(error.to_string()))?;
let resolved_flow = resolve_task_flow(&main_repo, flow.as_deref())?;
let resolved =
crate::ops::task_pm::resolve_task(&main_repo, issue, crate::ops::pm::PmRefresh::Auto)?;
let segment = match name.as_deref() {
Some(name) => parse_workspace_slug(name)?,
None => match &terminal_predecessor_id {
Some(predecessor_id) => succession_workspace_slug(&resolved.item.name, predecessor_id)?,
None => derive_workspace_slug(&resolved.item.name)?,
},
};
let workspace_slug = segment.as_str().to_string();
let mut plan = plan_placement(&main_repo, segment)
.map_err(|error| task_error(format!("failed to plan task worktree: {error}")))?;
if plan.strategy != PlacementStrategy::Create {
return Err(task_error(format!(
"task worktree or branch already exists without a Task Session: {} ({})",
plan.worktree_path.display(),
plan.branch
)));
}
let default_branch =
get_default_branch(&main_repo).map_err(|error| task_error(error.to_string()))?;
let stack_parent = stack_on
.as_deref()
.map(|parent_issue| {
block_on_task(async {
let store = task_store().await?;
let parent_session = store
.get_task_session_by_issue(parent_issue)
.await
.map_err(|error| task_error(format!("failed to read parent Task: {error}")))?
.ok_or_else(|| {
task_error(format!(
"stack parent {parent_issue:?} has no Task Session; run it first"
))
})?;
if parent_session.launch.issue.id.as_str() == resolved.item.id {
return Err(task_error("a Task cannot stack on itself"));
}
let parent = store
.active_task_pr(&parent_session.id)
.await
.map_err(|error| task_error(format!("failed to read parent PR: {error}")))?
.ok_or_else(|| task_error("stack parent has no active PR"))?;
if parent.github().is_none() {
return Err(task_error(format!(
"open the parent PR from {} before stacking work on it",
parent_session.worktree.display()
)));
}
Ok(parent)
})
})
.transpose()?;
let (base_ref, base_commit) = match &stack_parent {
Some(parent) => {
fetch(&main_repo, "origin", &parent.branch).map_err(|error| {
task_error(format!(
"failed to fetch parent branch {}: {error}",
parent.branch
))
})?;
let base_ref = format!("origin/{}", parent.branch);
let base_commit = rev_parse(&main_repo, &base_ref).map_err(|error| {
task_error(format!("failed to resolve task base {base_ref}: {error}"))
})?;
(base_ref, base_commit)
}
None => {
let (base_ref, base_commit) = resolve_upstream_base(&main_repo, &default_branch)?;
if base_ref.starts_with("origin/") {
refuse_if_canonical_ahead(&main_repo, &default_branch)?;
}
(base_ref, base_commit)
}
};
plan.base_ref = base_ref.clone();
let project_session = crate::ops::project::ensure_project_session_for_task(
&main_repo,
crate::ops::task_pm::ResolvedProject {
snapshot: resolved.snapshot.clone(),
project: resolved.project.clone(),
},
)?;
let project_session_id = task_project_session_id(&project_session)?;
let wave_id = project_session.wave_id.clone();
let config = load_config_or_default(Some(&main_repo));
let agent = config.agent.as_deref().unwrap_or("claude:opus");
let (provider, _) = parse_agent(agent);
let agent = agent.to_string();
let directive = directive.unwrap_or_else(|| {
format!(
"Complete {}: {}\n\n{}",
resolved.item.identifier, resolved.item.name, resolved.item.description
)
});
block_on_task(async move {
let store = task_store().await?;
let predecessor = match store
.get_task_session_by_issue(&resolved.item.id)
.await
.map_err(|error| task_error(format!("failed to read task registry: {error}")))?
{
Some(existing) if !existing.status.is_terminal() => return Ok(existing),
Some(terminal) => Some(terminal),
None => None,
};
let now = time::OffsetDateTime::now_utc();
let mut session = TaskSession {
id: crate::task::TaskSessionId::new(),
launch: TaskLaunchReceipt {
issue: LinearIssueSnapshot {
id: LinearIssueId::new(resolved.item.id.clone())
.map_err(|error| task_error(error.to_string()))?,
identifier: resolved.item.identifier.clone(),
title: resolved.item.name.clone(),
description: resolved.item.description.clone(),
},
project: LinearProjectSnapshot {
id: LinearProjectId::new(resolved.project.id.clone())
.map_err(|error| task_error(error.to_string()))?,
slug: resolved.project.slug.clone(),
name: resolved.project.name.clone(),
prompt_context: project_context(&resolved.project),
},
pm_snapshot_synced_at: resolved.snapshot.synced_at,
},
wave_id,
project_session_id,
pm_writeback: PmWritebackState::Current,
status: TaskSessionStatus::Created,
status_reason: "Linear task reserved before placement".to_string(),
status_at: now,
worktree: plan.worktree_path.clone(),
workspace_slug: workspace_slug.clone(),
lifecycle: if headless {
crate::task::TaskLifecyclePlan::headless(resolved_flow)
} else {
crate::task::TaskLifecyclePlan::standard(resolved_flow)
},
lifecycle_phase: crate::task::TaskLifecyclePhase::Kickoff,
phase_epoch: 1,
phase_cursor: 0,
phase_iteration: 0,
gate_cycle: 0,
gate_proposal: None,
agent,
provider,
provider_session_id: None,
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: workspace_slug,
branch: plan.branch.clone(),
base_commit,
parent_pr_id: stack_parent.as_ref().map(|parent| parent.id.clone()),
publication: None,
merge_commit: None,
abandoned_at: None,
ci_observation: None,
github_observation: None,
linear_attachment_id: None,
linear_comment_id: None,
linear_link_error: None,
created_at: now,
updated_at: now,
};
let succession = if let Some(predecessor) = &predecessor {
store
.reserve_task_session_successor(
predecessor,
&session,
&pr,
crate::durable::Author::User,
&directive,
)
.await
.map_err(|error| task_error(format!("failed to reserve task successor: {error}")))?
} else {
match store
.create_task_session_with_steer(
&session,
&pr,
crate::durable::Author::User,
&directive,
)
.await
{
Ok(()) => TaskSessionSuccession {
session: session.clone(),
created: true,
},
Err(StoreError::Sqlite(_)) => {
if let Some(existing) = store
.get_task_session_by_issue(&resolved.item.id)
.await
.map_err(|error| {
task_error(format!("failed to recover task reservation: {error}"))
})?
{
if !existing.status.is_terminal() {
return Ok(existing);
}
}
return Err(task_error(
"task reservation collided with another task placement",
));
}
Err(error) => return Err(task_error(format!("failed to reserve task: {error}"))),
}
};
if !succession.created {
return Ok(succession.session);
}
store
.append_task_event(
&session.id,
&TaskEventKind::PrStarted {
pr_id: pr.id,
sequence: pr.sequence,
branch: pr.branch,
base_commit: pr.base_commit,
},
)
.await
.map_err(|error| task_error(error.to_string()))?;
if let Err(error) = create_from_placement_plan(&main_repo, &plan) {
record_task_failure(
&store,
&mut session,
format!("worktree creation failed: {error}"),
error.to_string(),
)
.await?;
return Err(task_error(format!(
"failed to create task worktree: {error}"
)));
}
launch_task_process(&store, &mut session, None).await?;
wait_until_running(&store, &session.id).await
})
}
fn task_project_session_id(
project: &crate::project_session::ProjectSession,
) -> OpsResult<crate::project_session::ProjectSessionId> {
match std::env::var("LF_PROJECT_SESSION_ID") {
Ok(value) => {
let session_id =
crate::project_session::ProjectSessionId::parse(&value).map_err(|error| {
task_error(format!("invalid ambient Project Session id: {error}"))
})?;
if session_id != project.id {
return Err(task_error(format!(
"Project Session {session_id} cannot supervise a Task under Project Session {}",
project.id
)));
}
Ok(session_id)
}
Err(std::env::VarError::NotPresent) => Ok(project.id.clone()),
Err(std::env::VarError::NotUnicode(_)) => {
Err(task_error("ambient Project Session id is not valid UTF-8"))
}
}
}
pub(crate) fn project_context(project: &crate::pm::PmProject) -> String {
let mut context = format!("Definition:\n{}", project.definition.trim());
if !project.krs.is_empty() {
context.push_str("\n\nKRs:");
for kr in &project.krs {
let mark = if kr.holds { "x" } else { " " };
context.push_str(&format!("\n- [{mark}] {}", kr.text));
}
}
context
}
pub fn task_start(
repo: &Path,
title: String,
project_id: &str,
options: TaskLaunchOptions,
) -> OpsResult<TaskSession> {
let main = crate::ops::project::ensure_clean_main(repo, "Task start")
.map_err(|error| task_error(error.to_string()))?;
let project =
crate::ops::task_pm::resolve_project(&main, project_id, crate::ops::pm::PmRefresh::Auto)?;
crate::ops::project::require_registered_wave(&project.snapshot.wave)
.map_err(|error| task_error(error.to_string()))?;
let marker = format!(
"<!-- loopflow-task-start:{} -->",
hex::encode(Sha256::digest(
format!("{}\0{}", project.project.id, title).as_bytes()
))
);
let created = crate::ops::task_pm::create_and_load_task(
&main,
&project.snapshot.wave,
&project.project.slug,
&title,
&marker,
)?;
task_run(&main, &created.item.id, options)
}
fn resolve_task_flow(repo: &Path, requested: Option<&str>) -> OpsResult<String> {
let requested = requested.unwrap_or("task");
let definition = load_flow(requested, repo)
.map_err(|error| task_error(format!("failed to load Task flow {requested:?}: {error}")))?;
let steps = expand_flow(&definition, repo).map_err(|error| {
task_error(format!("failed to expand Task flow {requested:?}: {error}"))
})?;
if steps.is_empty() {
return Err(task_error(format!("Task flow {requested:?} has no steps")));
}
if let Some(step) = steps
.iter()
.find(|step| !matches!(step, ConcreteStep::Skill(_)))
{
return Err(task_error(format!(
"Task flow {requested:?} contains {step:?}; durable Task flows currently require skills"
)));
}
Ok(definition.name)
}
fn parse_workspace_slug(value: &str) -> OpsResult<WorktreeSegment> {
let value = value.trim();
let words = value.split('-').filter(|word| !word.is_empty()).count();
if sanitize_for_branch(value) != value
|| value.contains(['.', '_', '/'])
|| !(2..=5).contains(&words)
{
return Err(task_error(
"workspace name must be 2-5 lowercase kebab-case words",
));
}
WorktreeSegment::parse(value).map_err(|error| task_error(error.to_string()))
}
fn derive_workspace_slug(title: &str) -> OpsResult<WorktreeSegment> {
derive_workspace_slug_with_cap(title, 5)
}
fn derive_workspace_slug_with_cap(title: &str, max_words: usize) -> OpsResult<WorktreeSegment> {
let sanitized = sanitize_for_branch(title);
let mut words = sanitized
.split('-')
.filter(|word| !word.is_empty())
.take(max_words)
.collect::<Vec<_>>();
if words.len() == 1 {
words.push("task");
}
parse_workspace_slug(&words.join("-"))
}
fn succession_workspace_slug(
title: &str,
predecessor_id: &TaskSessionId,
) -> OpsResult<WorktreeSegment> {
let base = derive_workspace_slug_with_cap(title, 4)?;
let id = predecessor_id.as_str();
let tail = &id[id.len().saturating_sub(8)..];
parse_workspace_slug(&format!("{}-s{}", base.as_str(), tail))
}
fn parse_pr_slug(value: &str) -> OpsResult<String> {
let value = value.trim();
let words = value.split('-').filter(|word| !word.is_empty()).count();
if sanitize_for_branch(value) != value
|| value.contains(['.', '_', '/'])
|| !(1..=5).contains(&words)
{
return Err(task_error(
"next PR name must be 1-5 lowercase kebab-case words",
));
}
Ok(value.to_string())
}
async fn update_task_pr_with_authority(
store: &SharedStore,
pr: &TaskPr,
lease: Option<&ChildWriteLease>,
) -> Result<(), StoreError> {
match lease {
Some(lease) => store.update_task_pr_for_lease(pr, lease).await,
None => store.update_task_pr(pr).await,
}
}
async fn settle_task_pr_with_authority(
store: &SharedStore,
settled: &TaskPr,
next: Option<&TaskPr>,
lease: Option<&ChildWriteLease>,
) -> Result<(), StoreError> {
match lease {
Some(lease) => store.settle_task_pr_for_lease(settled, next, lease).await,
None => store.settle_task_pr(settled, next).await,
}
}
async fn append_task_event_with_authority(
store: &SharedStore,
session_id: &crate::task::TaskSessionId,
event: &TaskEventKind,
lease: Option<&ChildWriteLease>,
) -> Result<(), StoreError> {
match lease {
Some(lease) => {
store
.append_task_event_for_lease(session_id, lease, event)
.await?;
}
None => {
store.append_task_event(session_id, event).await?;
}
}
Ok(())
}
async fn update_task_session_with_authority(
store: &SharedStore,
session: &TaskSession,
lease: Option<&ChildWriteLease>,
) -> Result<(), StoreError> {
match lease {
Some(lease) => store.update_task_session_for_lease(session, lease).await,
None => store.update_task_session(session).await,
}
}
async fn complete_task_session_after_pr_with_authority(
store: &SharedStore,
session: &TaskSession,
pr: &TaskPr,
lease: Option<&ChildWriteLease>,
) -> Result<(), StoreError> {
match lease {
Some(lease) => {
store
.complete_task_session_after_pr_for_lease(session, pr, lease)
.await
}
None => store.complete_task_session_after_pr(session, pr).await,
}
}
async fn complete_task_session_with_authority(
store: &SharedStore,
session: &TaskSession,
skipped_pr: Option<&TaskPr>,
lease: Option<&ChildWriteLease>,
) -> Result<(), StoreError> {
match lease {
Some(lease) => {
store
.complete_task_session_for_lease(session, skipped_pr, lease)
.await
}
None => store.complete_task_session(session, skipped_pr).await,
}
}
fn ambient_task_write_lease(session: &TaskSession) -> OpsResult<Option<ChildWriteLease>> {
let Some(value) = std::env::var_os("LF_TASK_SESSION_ID") else {
return Ok(None);
};
let value = value
.into_string()
.map_err(|_| task_error("ambient Task Session id is not valid UTF-8"))?;
let id = crate::task::TaskSessionId::parse(&value)
.map_err(|error| task_error(format!("invalid ambient Task Session id: {error}")))?;
if id != session.id {
return Err(task_error(format!(
"ambient Task Session {id} cannot mutate {}",
session.id
)));
}
task_write_lease_from_env()
.map(Some)
.map_err(|error| task_error(format!("ambient Task Session has no authority: {error}")))
}
async fn task_for_worktree(
store: &SharedStore,
repo: &Path,
) -> OpsResult<Option<(TaskSession, Option<ChildWriteLease>)>> {
let checkout = repo.canonicalize().unwrap_or_else(|_| repo.to_path_buf());
if let Some(value) = std::env::var_os("LF_TASK_SESSION_ID") {
let value = value
.into_string()
.map_err(|_| task_error("ambient Task Session id is not valid UTF-8"))?;
let id = crate::task::TaskSessionId::parse(&value)
.map_err(|error| task_error(format!("invalid ambient Task Session id: {error}")))?;
let session = store
.get_task_session(&id)
.await
.map_err(|error| task_error(format!("failed to read ambient Task Session: {error}")))?
.ok_or_else(|| task_error(format!("ambient Task Session {id} is not registered")))?;
let worktree = session
.worktree
.canonicalize()
.unwrap_or_else(|_| session.worktree.clone());
if checkout != worktree {
return Err(task_error(format!(
"ambient Task Session {id} owns {}, not {}",
session.worktree.display(),
repo.display()
)));
}
let lease = task_write_lease_from_env().map_err(|error| {
task_error(format!("ambient Task Session has no authority: {error}"))
})?;
return Ok(Some((session, Some(lease))));
}
let worktree_keys: BTreeSet<String> = store
.list_task_sessions(None)
.await
.map_err(|error| task_error(format!("failed to inspect Task Sessions: {error}")))?
.into_iter()
.filter(|session| {
session
.worktree
.canonicalize()
.unwrap_or_else(|_| session.worktree.clone())
== checkout
})
.map(|session| session.worktree.display().to_string())
.collect();
let mut current = Vec::new();
for worktree in worktree_keys {
if let Some(session) = store
.get_task_session_by_worktree(&worktree)
.await
.map_err(|error| task_error(format!("failed to resolve Task worktree: {error}")))?
{
current.push(session);
}
}
if current.len() > 1 {
return Err(task_error(format!(
"multiple Task Sessions claim worktree {}",
repo.display()
)));
}
Ok(current.pop().map(|session| (session, None)))
}
#[derive(Debug)]
enum TaskAuthority {
NotATaskWorktree,
Authority {
store: SharedStore,
session: Box<TaskSession>,
lease: Option<ChildWriteLease>,
},
}
fn registry_authority_error(err: RegistryUnavailable) -> OpsError {
task_error(match err {
RegistryUnavailable::MissingFile { path } => format!(
"Task PR authority refused: the shared Loopflow registry {} is missing. \
Start the owning Wave (it creates the registry) or run `lf doctor`.",
path.display()
),
RegistryUnavailable::Unresolved { error } => format!(
"Task PR authority refused: the shared Loopflow registry path is not usable: {error}. \
Fix LF_DB_PATH/LF_HOME or run `lf doctor`."
),
RegistryUnavailable::Incompatible { path, error } => format!(
"Task PR authority refused: the shared Loopflow registry {} is present but \
inaccessible or schema-incompatible: {error}. Run `lf doctor`.",
path.display()
),
})
}
async fn resolve_task_authority(repo: &Path) -> OpsResult<TaskAuthority> {
let ambient = std::env::var_os("LF_TASK_SESSION_ID").is_some();
let store = match open_registry_for_authority().await {
Ok(store) => Arc::new(store),
Err(RegistryUnavailable::MissingFile { .. }) if !ambient => {
return Ok(TaskAuthority::NotATaskWorktree);
}
Err(err) => return Err(registry_authority_error(err)),
};
match task_for_worktree(&store, repo).await? {
Some((session, lease)) => Ok(TaskAuthority::Authority {
store,
session: Box::new(session),
lease,
}),
None => Ok(TaskAuthority::NotATaskWorktree),
}
}
pub(crate) fn request_task_pr_publication(
repo: &Path,
after_merge: AfterMerge,
next_slug: Option<&str>,
) -> OpsResult<bool> {
let next_slug = next_slug.map(parse_pr_slug).transpose()?;
if after_merge == AfterMerge::CompleteTask && next_slug.is_some() {
return Err(task_error("--complete and --next cannot be used together"));
}
block_on_task(async move {
let TaskAuthority::Authority {
store,
session,
lease,
} = resolve_task_authority(repo).await?
else {
return Ok(false);
};
let mut pr = store
.active_task_pr(&session.id)
.await
.map_err(|error| task_error(format!("failed to read active PR: {error}")))?
.ok_or_else(|| {
task_error(format!(
"Task {} has no active PR",
session.launch.issue.identifier
))
})?;
let branch = crate::engine::git::current_branch(repo)?
.ok_or_else(|| task_error("Task worktree is not on a branch"))?;
if pr.branch != branch {
return Err(task_error(format!(
"Task {} active PR expects branch {:?}, but the worktree is on another branch",
session.launch.issue.identifier, pr.branch
)));
}
let now = time::OffsetDateTime::now_utc();
pr.publication = Some(PrPublication {
requested_at: pr
.publication
.as_ref()
.map_or(now, |publication| publication.requested_at),
after_merge,
next_slug,
github: pr.github().cloned(),
});
pr.updated_at = now;
match lease.as_ref() {
Some(lease) => store.update_task_pr_for_lease(&pr, lease).await,
None => store.update_task_pr(&pr).await,
}
.map_err(|error| task_error(format!("failed to request PR publication: {error}")))?;
Ok(true)
})
}
fn has_remote(repo: &Path) -> OpsResult<bool> {
Ok(!git_output(repo, &["remote"])?.trim().is_empty())
}
fn resolve_upstream_base(repo: &Path, default_branch: &str) -> OpsResult<(String, String)> {
let base_ref = if has_remote(repo)? {
fetch(repo, "origin", default_branch)
.map_err(|error| task_error(format!("failed to fetch task base: {error}")))?;
format!("origin/{default_branch}")
} else {
format!("refs/heads/{default_branch}")
};
let base_commit = rev_parse(repo, &base_ref)
.map_err(|error| task_error(format!("failed to resolve task base {base_ref}: {error}")))?;
Ok((base_ref, base_commit))
}
fn refuse_if_canonical_ahead(repo: &Path, default_branch: &str) -> OpsResult<()> {
if rev_parse(repo, &format!("refs/heads/{default_branch}")).is_err() {
return Ok(());
}
let range = format!("origin/{default_branch}..{default_branch}");
let ahead = git_output(repo, &["log", "--oneline", "--no-decorate", &range])?;
let ahead = ahead.trim();
if !ahead.is_empty() {
return Err(task_error(format!(
"canonical {default_branch} is ahead of origin/{default_branch}; new Task worktrees \
would inherit these unpushed commit(s):\n{ahead}\nThis is a control-plane violation. \
Push or reset {default_branch} to origin/{default_branch} before placing Task worktrees."
)));
}
Ok(())
}
pub(crate) fn verify_task_pr_range(repo: &Path) -> OpsResult<()> {
let repo = repo.to_path_buf();
block_on_task(async move {
let TaskAuthority::Authority {
store,
session,
lease,
} = resolve_task_authority(&repo).await?
else {
return Ok(());
};
verify_task_pr_range_with_authority(&store, &session, lease.as_ref(), &repo).await
})
}
pub(crate) fn require_task_pr_range_nonempty(repo: &Path) -> OpsResult<()> {
let repo = repo.to_path_buf();
block_on_task(async move {
let TaskAuthority::Authority {
store,
session,
lease,
} = resolve_task_authority(&repo).await?
else {
return Ok(());
};
require_task_pr_range_nonempty_with_authority(&store, &session, lease.as_ref(), &repo).await
})
}
async fn resolve_verifier_upstream(
store: &SharedStore,
pr: &TaskPr,
repo: &Path,
default_branch: &str,
) -> OpsResult<(String, String)> {
if let Some(parent_id) = pr.parent_pr_id.as_ref() {
let parent = store
.get_task_pr(parent_id)
.await
.map_err(|error| task_error(format!("failed to read stack parent: {error}")))?
.ok_or_else(|| task_error(format!("stack parent {parent_id} is missing")))?;
let parent_live = parent.merge_commit.is_none() && parent.abandoned_at.is_none();
if parent_live {
let base_ref = if has_remote(repo)? {
fetch(repo, "origin", &parent.branch).map_err(|error| {
task_error(format!("failed to fetch parent branch: {error}"))
})?;
format!("origin/{}", parent.branch)
} else {
format!("refs/heads/{}", parent.branch)
};
let tip = rev_parse(repo, &base_ref).map_err(|error| {
task_error(format!(
"failed to resolve parent branch {base_ref}: {error}"
))
})?;
return Ok((base_ref, tip));
}
}
resolve_upstream_base(repo, default_branch)
}
async fn verify_task_pr_range_with_authority(
store: &SharedStore,
session: &TaskSession,
lease: Option<&ChildWriteLease>,
repo: &Path,
) -> OpsResult<()> {
let mut pr = store
.active_task_pr(&session.id)
.await
.map_err(|error| task_error(format!("failed to read active PR: {error}")))?
.ok_or_else(|| {
task_error(format!(
"Task {} has no active PR",
session.launch.issue.identifier
))
})?;
let branch =
current_branch(repo)?.ok_or_else(|| task_error("Task worktree is not on a branch"))?;
if pr.branch != branch {
return Err(task_error(format!(
"Task {} active PR expects branch {:?}, but the worktree is on {:?}",
session.launch.issue.identifier, pr.branch, branch
)));
}
let default_branch = get_default_branch(repo)?;
let (base_ref, upstream) = resolve_verifier_upstream(store, &pr, repo, &default_branch).await?;
let head = rev_parse(repo, "HEAD")
.map_err(|error| task_error(format!("failed to resolve Task HEAD: {error}")))?;
let base = pr.base_commit.clone();
let identifier = &session.launch.issue.identifier;
let short = |sha: &str| sha.chars().take(12).collect::<String>();
let merge_base = crate::engine::git::merge_base(repo, &upstream, &head).map_err(|_| {
task_error(format!(
"Task {identifier} branch {branch:?} shares no history with {base_ref}; \
re-cut the branch from {base_ref} before publishing"
))
})?;
if merge_base == base {
return Ok(());
}
if crate::engine::git::is_ancestor(repo, &merge_base, &base)? {
let range = format!("{merge_base}..{base}");
let commits = git_output(repo, &["log", "--oneline", "--no-decorate", &range])?;
let files = git_output(repo, &["diff", "--name-only", &range])?;
let commits = commits.trim();
let files = files.trim();
return Err(task_error(format!(
"Task {identifier} PR range is contaminated: recorded base {} carries commit(s) \
not on {base_ref}, which would leak into the PR:\n{commits}\naffecting files:\n{files}\n\
Refused before any push. Recover with:\n git rebase --onto {base_ref} {} {branch}",
short(&base),
short(&base),
)));
}
if crate::engine::git::is_ancestor(repo, &base, &merge_base)? {
pr.base_commit = merge_base.clone();
pr.updated_at = time::OffsetDateTime::now_utc();
match lease {
Some(lease) => store.heal_task_pr_base_for_lease(&pr, lease).await,
None => store.heal_task_pr_base(&pr).await,
}
.map_err(|error| task_error(format!("failed to heal Task PR base: {error}")))?;
return Ok(());
}
let base_side = format!("{merge_base}..{base}");
let upstream_side = format!("{base}..{merge_base}");
let base_commits = git_output(repo, &["log", "--oneline", "--no-decorate", &base_side])?;
let base_files = git_output(repo, &["diff", "--name-only", &base_side])?;
let upstream_commits =
git_output(repo, &["log", "--oneline", "--no-decorate", &upstream_side])?;
let upstream_files = git_output(repo, &["diff", "--name-only", &upstream_side])?;
Err(task_error(format!(
"Task {identifier} PR base {} and {base_ref} have diverged with no common lineage at \
the recorded base. Refused before any push.\n\
Commits on the recorded base not on {base_ref}:\n{base_commits}\
affecting files:\n{base_files}\n\
Commits on {base_ref} not reachable from the recorded base:\n{upstream_commits}\
affecting files:\n{upstream_files}\n\
Recover with:\n git rebase --onto {base_ref} {} {branch}",
short(&base),
short(&base),
)))
}
async fn require_task_pr_range_nonempty_with_authority(
store: &SharedStore,
session: &TaskSession,
lease: Option<&ChildWriteLease>,
repo: &Path,
) -> OpsResult<()> {
verify_task_pr_range_with_authority(store, session, lease, repo).await?;
let pr = store
.active_task_pr(&session.id)
.await
.map_err(|error| task_error(format!("failed to read active PR: {error}")))?
.ok_or_else(|| {
task_error(format!(
"Task {} has no active PR",
session.launch.issue.identifier
))
})?;
let base = &pr.base_commit;
let identifier = &session.launch.issue.identifier;
let short = base.chars().take(12).collect::<String>();
let head = rev_parse(repo, "HEAD")
.map_err(|error| task_error(format!("failed to resolve Task HEAD: {error}")))?;
if head == *base {
return Err(task_error(format!(
"Task {identifier} PR range is empty: HEAD is the recorded base {short}, so the PR has \
no commits to publish. Commit the Task's work, or complete the Task directly if the \
work is done. Refused before any GitHub side effect."
)));
}
let range = format!("{base}..HEAD");
let status = Command::new("git")
.args(["diff", "--quiet", &range])
.current_dir(repo)
.status()?;
if status.success() {
return Err(task_error(format!(
"Task {identifier} PR range is empty: the tree at HEAD matches the recorded base \
{short}, so the PR has no changes to publish. Commit the Task's work, or complete the \
Task directly if the work is done. Refused before any GitHub side effect."
)));
}
Ok(())
}
pub(crate) fn attach_task_github_pr(
repo: &Path,
github_pr: Option<&crate::ops::pr::PrInfo>,
) -> OpsResult<bool> {
block_on_task(async move {
let TaskAuthority::Authority {
store,
session,
lease,
} = resolve_task_authority(repo).await?
else {
return Ok(false);
};
let mut pr = store
.active_task_pr(&session.id)
.await
.map_err(|error| task_error(format!("failed to read active PR: {error}")))?
.ok_or_else(|| {
task_error(format!(
"Task {} has no active PR",
session.launch.issue.identifier
))
})?;
let github_pr = github_pr.ok_or_else(|| {
task_error(format!(
"GitHub PR for Task {} could not be read after creation or update",
session.launch.issue.identifier
))
})?;
if github_pr.branch != pr.branch {
return Err(task_error(format!(
"Task {} active PR expects branch {:?}, but GitHub reported {:?}",
session.launch.issue.identifier, pr.branch, github_pr.branch
)));
}
let number = u32::try_from(github_pr.number).map_err(|_| {
task_error(format!(
"pull request #{} exceeds supported range",
github_pr.number
))
})?;
let url = github_pr.url.clone();
let opened = pr
.github()
.is_none_or(|github| github.number != number || github.url != url);
let publication = pr.publication.as_mut().ok_or_else(|| {
task_error(format!(
"Task {} has no durable PR publication request",
session.launch.issue.identifier
))
})?;
publication.github = Some(GithubPr {
number,
url: url.clone(),
head_sha: github_pr.head_sha.clone(),
});
link_pr_to_linear(&store, &session, &mut pr).await;
pr.updated_at = time::OffsetDateTime::now_utc();
match lease.as_ref() {
Some(lease) => store.update_task_pr_for_lease(&pr, lease).await,
None => store.update_task_pr(&pr).await,
}
.map_err(|error| task_error(format!("failed to attach GitHub PR: {error}")))?;
if opened {
let event = TaskEventKind::PrOpened {
pr_id: pr.id,
sequence: pr.sequence,
number,
url,
};
match lease.as_ref() {
Some(lease) => {
store
.append_task_event_for_lease(&session.id, lease, &event)
.await
}
None => store.append_task_event(&session.id, &event).await,
}
.map_err(|error| task_error(error.to_string()))?;
}
Ok(true)
})
}
pub(crate) fn abandon_task_pr(
repo: &Path,
force: bool,
progress: &impl crate::ops::progress::Progress,
) -> OpsResult<bool> {
block_on_task(async move {
let TaskAuthority::Authority {
store,
session,
lease,
} = resolve_task_authority(repo).await?
else {
return Ok(false);
};
let mut session = session;
let mut pr = store
.active_task_pr(&session.id)
.await
.map_err(|error| task_error(format!("failed to read active PR: {error}")))?
.ok_or_else(|| {
task_error(format!(
"Task {} has no active PR to abandon",
session.launch.issue.identifier
))
})?;
let branch =
current_branch(repo)?.ok_or_else(|| task_error("Task worktree is not on a branch"))?;
if branch != pr.branch {
return Err(task_error(format!(
"Task {} active PR expects branch {:?}, but the worktree is on {:?}",
session.launch.issue.identifier, pr.branch, branch
)));
}
let dirty = !is_clean(repo)?;
if dirty && !force {
return Err(task_error("uncommitted changes; use --force"));
}
if let Some(lease) = lease.as_ref() {
store
.validate_child_write_lease(&ChildRef::Task(session.id.clone()), lease)
.await
.map_err(|error| task_error(format!("Task body lost write authority: {error}")))?;
}
if dirty {
progress.status("Discarding uncommitted Task PR changes...");
for args in [
["reset", "--hard", "HEAD"].as_slice(),
["clean", "-fd"].as_slice(),
] {
let output = Command::new("git").args(args).current_dir(repo).output()?;
if !output.status.success() {
return Err(task_error(format!(
"failed to discard Task PR changes with `git {}`: {}",
args.join(" "),
String::from_utf8_lossy(&output.stderr).trim()
)));
}
}
}
progress.status("Closing Task PR...");
let _ = Command::new("gh")
.args(["pr", "close", &branch])
.current_dir(repo)
.status();
let now = time::OffsetDateTime::now_utc();
pr.abandoned_at = Some(now);
pr.updated_at = now;
match lease.as_ref() {
Some(lease) => store.settle_task_pr_for_lease(&pr, None, lease).await,
None => store.settle_task_pr(&pr, None).await,
}
.map_err(|error| task_error(format!("failed to settle Task PR: {error}")))?;
if !session.status.is_process_active() {
let from = session.status;
session.set_status(
TaskSessionStatus::Waiting,
format!("PR branch {branch:?} was abandoned; another PR may follow"),
);
store
.update_task_session(&session)
.await
.map_err(|error| task_error(error.to_string()))?;
if from != session.status {
store
.append_task_event(
&session.id,
&TaskEventKind::StatusChanged {
from,
to: session.status,
reason: session.status_reason.clone(),
},
)
.await
.map_err(|error| task_error(error.to_string()))?;
}
}
Ok(true)
})
}
async fn record_task_failure(
store: &SharedStore,
session: &mut TaskSession,
reason: impl Into<String>,
error: String,
) -> OpsResult<()> {
let from = session.status;
session.set_status(TaskSessionStatus::Failed, reason);
store
.update_task_session(session)
.await
.map_err(|store_error| task_error(store_error.to_string()))?;
store
.append_task_event(
&session.id,
&TaskEventKind::StatusChanged {
from,
to: TaskSessionStatus::Failed,
reason: session.status_reason.clone(),
},
)
.await
.map_err(|store_error| task_error(store_error.to_string()))?;
store
.append_task_event(
&session.id,
&TaskEventKind::Failed {
error,
resumable: true,
},
)
.await
.map_err(|store_error| task_error(store_error.to_string()))?;
Ok(())
}
pub(crate) async fn relaunch_inactive_process(
store: &SharedStore,
session: &mut TaskSession,
) -> OpsResult<()> {
let Some(_) = ensure_working_pr(store, session).await? else {
return Err(task_error(format!(
"Task {} is {}; terminal Task Sessions cannot start a process",
session.launch.issue.identifier,
session.status.as_str()
)));
};
launch_task_process(store, session, None).await
}
async fn relaunch_for_ci_incident(
store: &SharedStore,
session: &mut TaskSession,
incident_id: String,
) -> OpsResult<()> {
let Some(_) = ensure_working_pr(store, session).await? else {
return Err(task_error(format!(
"Task {} is terminal and cannot repair CI",
session.launch.issue.identifier
)));
};
launch_task_process(
store,
session,
Some(crate::durable::RunTrigger::CiIncident { incident_id }),
)
.await
}
async fn release_dead_revoked_task_lease(
store: &SharedStore,
session: &mut TaskSession,
) -> OpsResult<()> {
let Some(revoked) = session
.latest_process
.as_ref()
.filter(|process| process.state == ChildLeaseState::Revoked)
.cloned()
else {
return Ok(());
};
if let Some(finished) = super::child::release_dead_revoked_child_body(
store,
&ChildRef::Task(session.id.clone()),
&revoked,
)
.await?
{
session.latest_process = Some(finished);
}
Ok(())
}
async fn launch_task_process(
store: &SharedStore,
session: &mut TaskSession,
trigger: Option<crate::durable::RunTrigger>,
) -> OpsResult<()> {
let execution = crate::engine::process::current_home_execution_context()
.map_err(|error| task_error(format!("cannot resolve current lf binary: {error}")))?;
let tmux_name = format!(
"lf-task-{}-{}",
tmux_session_slug(&session.launch.issue.identifier),
&session.id.as_str()[3..11]
);
release_dead_revoked_task_lease(store, session).await?;
let from = session.status;
let mut launch = session.clone();
let generation = launch.begin_generation(tmux_name.clone());
let reservation = match trigger {
Some(trigger) => {
store
.reserve_task_process_for_trigger(&launch, from, trigger)
.await
}
None => store.reserve_task_process(&launch, from).await,
};
let Some(lease) = reservation
.map_err(|error| task_error(format!("failed to reserve task process: {error}")))?
else {
let current = store
.get_task_session(&session.id)
.await
.map_err(|error| task_error(format!("failed to reread task process: {error}")))?
.ok_or_else(|| task_error("Task Session disappeared during process reservation"))?;
if current.status.is_process_active() {
*session = current;
return Ok(());
}
if current.status.is_terminal() {
return Err(task_error(format!(
"task {} became {}; terminal Task Sessions cannot start a process",
current.launch.issue.identifier,
current.status.as_str()
)));
}
if let Some(process) = current.latest_process.as_ref() {
if process.state != ChildLeaseState::Finished {
return Err(task_error(format!(
"task {} holds a `{}` lease on body generation {}; a new generation \
cannot be reserved until that lease is released",
current.launch.issue.identifier,
process.state.as_str(),
process.generation
)));
}
}
return Err(task_error(format!(
"task {} changed from {} to {} during process reservation; retry the command",
current.launch.issue.identifier,
from.as_str(),
current.status.as_str()
)));
};
*session = launch;
let argv = vec![
execution.lf_bin.to_string_lossy().to_string(),
"__task".to_string(),
session.id.to_string(),
"--generation".to_string(),
generation.to_string(),
];
let generation_text = generation.to_string();
let run_lease = lease.run_token.env_value().to_string();
let control_bin = execution.lf_bin.to_string_lossy().to_string();
let db_path = execution.db_path.to_string_lossy().to_string();
let lf_home = execution.lf_home.to_string_lossy().to_string();
let wave_home = match owning_wave(store, session).await {
Ok(wave) => crate::engine::wave_config::read_wave_home(Path::new(wave.repo()), wave.name())
.to_string(),
Err(_) => crate::engine::wave_config::default_local_home(&session.worktree).to_string(),
};
let environment = [
(
crate::engine::wave_context::WAVE_ID_ENV,
session.wave_id.as_str(),
),
("LF_TASK_SESSION_ID", session.id.as_str()),
("LF_TASK_GENERATION", generation_text.as_str()),
(
crate::child_session::TASK_LEASE_TOKEN_ENV,
lease.token.as_str(),
),
(crate::durable::RUN_CONTEXT_ENV, "agent"),
(crate::durable::RUN_LEASE_ENV, run_lease.as_str()),
(crate::store::CONTROL_BIN_ENV, control_bin.as_str()),
(crate::store::CONTROL_DB_PATH_ENV, db_path.as_str()),
(crate::store::CONTROL_HOME_ENV, lf_home.as_str()),
(crate::engine::wave_home::WAVE_HOME_ENV, wave_home.as_str()),
];
if let Err(error) =
start_lf_session_with_env(&tmux_name, &session.worktree, &argv, &environment).await
{
session.latest_process = Some(
super::child::revoke_and_reap_child_body(
store,
&ChildRef::Task(session.id.clone()),
crate::child_session::ChildBodyOutcome::Lost {
reason: format!("task process launch failed: {error}"),
},
)
.await?,
);
record_task_failure(
store,
session,
format!("task process launch failed: {error}"),
error.to_string(),
)
.await?;
return Err(task_error(format!(
"failed to launch task process: {error}"
)));
}
Ok(())
}
async fn wait_until_running(
store: &SharedStore,
session_id: &crate::task::TaskSessionId,
) -> OpsResult<TaskSession> {
let deadline = tokio::time::Instant::now() + super::child::CHILD_STARTUP_GRACE;
loop {
let session = store
.get_task_session(session_id)
.await
.map_err(|error| task_error(format!("failed to observe task startup: {error}")))?
.ok_or_else(|| task_error("task session disappeared during startup"))?;
if session.status != TaskSessionStatus::Starting {
return if session.status == TaskSessionStatus::Running {
Ok(session)
} else {
Err(task_error(format!(
"task {} did not start: {}",
session.launch.issue.identifier, session.status_reason
)))
};
}
if tokio::time::Instant::now() >= deadline {
return Err(task_error(format!(
"task {} process did not report running within 10 seconds",
session.launch.issue.identifier
)));
}
tokio::time::sleep(Duration::from_millis(100)).await;
}
}
pub(crate) async fn reconcile_process_liveness(
store: &SharedStore,
session: &mut TaskSession,
) -> OpsResult<()> {
if session
.latest_process
.as_ref()
.is_some_and(super::child::child_body_reservation_is_fresh)
{
return Ok(());
}
let alive = match session.latest_process.as_ref() {
Some(process) => tmux_session_exists(&process.tmux_name)
.await
.map_err(|error| task_error(error.to_string()))?,
None => false,
};
if alive {
return Ok(());
}
let lost_reason = "task process disappeared before recording a terminal outcome";
if session.latest_process.as_ref().is_some_and(|process| {
matches!(
process.state,
crate::child_session::ChildLeaseState::Legacy
| crate::child_session::ChildLeaseState::Reserved
| crate::child_session::ChildLeaseState::Active
)
}) {
let outcome = super::child::lost_child_body_outcome(
session
.latest_process
.as_ref()
.expect("matched child process must still be present"),
lost_reason,
);
session.latest_process = Some(
super::child::revoke_and_reap_child_body(
store,
&ChildRef::Task(session.id.clone()),
outcome,
)
.await?,
);
}
if !session.status.is_process_active() {
return Ok(());
}
let active = store
.active_task_pr(&session.id)
.await
.map_err(|error| task_error(format!("failed to read active PR: {error}")))?;
if active.as_ref().is_none_or(|pr| pr.phase() == PrPhase::Open) {
let from = session.status;
session.set_status(
TaskSessionStatus::Waiting,
match active {
Some(pr) => format!("PR {} is open; waiting for review", pr.sequence),
None => "the previous PR settled; another PR may follow".to_string(),
},
);
store
.update_task_session(session)
.await
.map_err(|error| task_error(error.to_string()))?;
store
.append_task_event(
&session.id,
&TaskEventKind::StatusChanged {
from,
to: session.status,
reason: session.status_reason.clone(),
},
)
.await
.map_err(|error| task_error(error.to_string()))?;
return Ok(());
}
let reason = "task process is missing; Loopflow will recover this Task Session";
record_task_failure(store, session, reason, reason.to_string()).await
}
const RECOVERY_ATTEMPT_WINDOW: u32 = 64;
pub(crate) async fn reconcile_project_tasks(
store: &SharedStore,
project: &crate::project_session::ProjectSession,
) -> OpsResult<Vec<TaskSession>> {
let project_tasks = |tasks: Vec<TaskSession>| {
tasks
.into_iter()
.filter(|task| task.launch.project.id == project.launch.project.id)
.collect::<Vec<_>>()
};
let mut tasks = project_tasks(
store
.list_task_sessions(Some(&project.wave_id))
.await
.map_err(|error| task_error(format!("failed to list supervised Tasks: {error}")))?,
);
for task in &mut tasks {
if task.status.is_terminal() {
continue;
}
if let Err(error) = task_recovery_adoption(store, task).await {
tracing::warn!(
task = %task.launch.issue.identifier,
%error,
"supervisor skipped Task recovery: unsafe worktree/branch state"
);
continue;
}
let observed = reconcile_task_pr(store, task).await?;
refuse_dirty_between_prs(store, task).await?;
reconcile_process_liveness(store, task).await?;
reconcile_task_completion(store, task, None).await?;
if task.status.is_terminal() {
continue;
}
let no_active_pr = if observed.is_none() {
store
.active_task_pr(&task.id)
.await
.map_err(|error| task_error(format!("failed to read active PR: {error}")))?
.is_none()
} else {
false
};
let settled = observed.as_ref().is_some_and(TaskPr::is_settled) || no_active_pr;
let completing = observed.as_ref().is_some_and(|pr| {
pr.is_settled()
&& pr
.publication
.as_ref()
.is_some_and(|publication| publication.after_merge == AfterMerge::CompleteTask)
});
if settled && !completing {
ensure_working_pr(store, task).await?;
if !task.status.is_process_active() {
relaunch_inactive_process(store, task).await?;
}
} else {
let Some(pr) = observed.as_ref() else {
continue;
};
if pr.review_ready()
&& task.lifecycle_phase == crate::task::TaskLifecyclePhase::Gate
&& task.phase_cursor == 0
&& task.phase_iteration == 0
{
if !task.status.is_process_active() {
relaunch_inactive_process(store, task).await?;
}
} else {
route_ci_incident(store, task, pr).await?;
}
}
}
let refreshed = store
.list_task_sessions(Some(&project.wave_id))
.await
.map_err(|error| task_error(format!("failed to reread supervised Tasks: {error}")))?;
Ok(project_tasks(refreshed))
}
pub(crate) async fn supervise_project_task_bodies(
store: &SharedStore,
project: &crate::project_session::ProjectSession,
) -> OpsResult<usize> {
let tasks = reconcile_project_tasks(store, project).await?;
if !tmux_installed() {
return Ok(0);
}
let live_sessions = tmux_live_sessions()
.await
.map_err(|error| task_error(format!("failed to observe Task bodies: {error}")))?;
let now = time::OffsetDateTime::now_utc();
let mut recovered = 0;
let mine = tasks
.into_iter()
.filter(|task| task.project_session_id == project.id)
.collect::<Vec<_>>();
for mut task in mine
.iter()
.filter(|task| {
task.latest_process.as_ref().is_some_and(|process| {
!live_sessions.contains(&process.tmux_name)
&& !super::child::child_body_reservation_is_fresh(process)
})
})
.cloned()
{
match recover_stranded_task_body(store, &mut task).await {
Ok(true) => recovered += 1,
Ok(false) => {}
Err(error) => {
tracing::warn!(
project_session = %project.id,
task = %task.launch.issue.identifier,
error = %error,
"stranded Task recovery failed"
);
}
}
}
for task in mine.into_iter().filter(|task| {
task.status.is_process_active()
&& task.latest_process.as_ref().is_some_and(|process| {
process.state == ChildLeaseState::Active
&& live_sessions.contains(&process.tmux_name)
})
}) {
let latest_event = store
.latest_task_event(&task.id)
.await
.map_err(|error| task_error(format!("failed to read Task progress: {error}")))?;
let observation = observe(
&BodyEvidence {
intent: task.status.body_intent(),
observable: true,
process_alive: true,
progress_age: body_progress_age(
latest_event.as_ref().map(|event| event.created_at),
task.status_at,
now,
),
step: Some(task.lifecycle_phase.as_str().to_string()),
reason: task.status_reason.clone(),
},
DEFAULT_STALL_AFTER,
);
match recover_stalled_task_body(
store,
task,
&observation,
latest_event.as_ref().map(|event| event.id),
)
.await
{
Ok(true) => recovered += 1,
Ok(false) => {}
Err(error) => {
tracing::warn!(
project_session = %project.id,
error = %error,
"Task body recovery failed"
);
}
}
}
Ok(recovered)
}
async fn recover_stranded_task_body(
store: &SharedStore,
task: &mut TaskSession,
) -> OpsResult<bool> {
let events = store
.recent_task_events(&task.id, RECOVERY_ATTEMPT_WINDOW)
.await
.map_err(|error| task_error(format!("failed to read Task recovery history: {error}")))?;
let attempts = count_recovery_attempts(&events);
release_dead_revoked_task_lease(store, task).await?;
let plan = plan_stranded_recovery(
task.status.body_intent(),
true,
false,
task.latest_process.as_ref(),
attempts,
);
let (attempt, reason) = match plan {
StrandedPlan::LeaveAlone => return Ok(false),
StrandedPlan::Surface { reason } => {
if task.status == TaskSessionStatus::Failed && task.status_reason == reason {
return Ok(false);
}
record_task_failure(store, task, reason.clone(), reason).await?;
return Ok(false);
}
StrandedPlan::Redispatch { attempt } => {
let generation = task
.latest_process
.as_ref()
.map_or(0, |process| process.generation);
(
attempt,
format!(
"body generation {generation} died without recording an outcome; \
recovering the same Task Session (attempt {attempt}/{MAX_RECOVERY_ATTEMPTS})"
),
)
}
};
if let Err(error) = task_recovery_adoption(store, task).await {
tracing::info!(
task = %task.launch.issue.identifier,
"not recovering stranded Task: {error}"
);
return Ok(false);
}
if let Some(process) = task.latest_process.as_ref() {
if matches!(
process.state,
ChildLeaseState::Legacy | ChildLeaseState::Reserved | ChildLeaseState::Active
) {
let outcome = super::child::lost_child_body_outcome(process, &reason);
task.latest_process = Some(
super::child::revoke_and_reap_child_body(
store,
&ChildRef::Task(task.id.clone()),
outcome,
)
.await?,
);
}
}
let generation = task
.latest_process
.as_ref()
.map_or(0, |process| process.generation);
store
.append_task_event(
&task.id,
&TaskEventKind::BodyRecoveryAttempted {
generation,
attempt,
reason: reason.clone(),
},
)
.await
.map_err(|error| task_error(error.to_string()))?;
super::child::redispatch_task_body(store, task).await?;
tracing::info!(
task = %task.launch.issue.identifier,
attempt,
"recovered a stranded Task body"
);
Ok(true)
}
async fn recover_stalled_task_body(
store: &SharedStore,
task: TaskSession,
observation: &crate::child_session::BodyObservation,
latest_event_id: Option<i64>,
) -> OpsResult<bool> {
let generation = task
.latest_process
.as_ref()
.map(|process| process.generation)
.ok_or_else(|| task_error("stalled Task has no process generation"))?;
let plan = plan_body_recovery(observation);
if plan == BodyRecoveryPlan::LeaveAlone {
return Ok(false);
}
let active_pr = store
.active_task_pr(&task.id)
.await
.map_err(|error| task_error(format!("failed to inspect Task PR: {error}")))?;
if let Some(reason) = task.supervisor_restart_bar(active_pr.as_ref()) {
tracing::info!(task = %task.launch.issue.identifier, "not recovering Task body: {reason}");
return Ok(false);
}
let progress_age = observation.progress_age_secs.unwrap_or_default();
if let Err(error) = task_recovery_adoption(store, &task).await {
tracing::info!(
task = %task.launch.issue.identifier,
"not recovering Task body: {error}"
);
return Ok(false);
}
let reason = format!(
"body generation {generation} stalled after {progress_age}s without durable progress; recovering from current Work input"
);
let outcome = ChildBodyOutcome::Superseded {
reason: reason.clone(),
};
let Some(revoked) = store
.revoke_task_process_if_unchanged(
&task.id,
generation,
task.status_at,
latest_event_id,
&outcome,
)
.await
.map_err(|error| task_error(format!("failed to claim stalled Task body: {error}")))?
else {
return Ok(false);
};
if let Err(error) =
super::child::reap_revoked_child_body(store, &ChildRef::Task(task.id.clone()), revoked)
.await
{
let mut current = store
.get_task_session(&task.id)
.await
.map_err(|store_error| task_error(store_error.to_string()))?
.ok_or_else(|| task_error("Task Session disappeared during recovery"))?;
let failure = format!(
"body generation {generation} lease was revoked after a stall but its body could not \
be reaped or proven gone: {error}; the lease stays blocked and releases itself once \
the body is verifiably absent"
);
record_task_failure(store, &mut current, failure.clone(), failure).await?;
return Err(error);
}
let mut current = store
.get_task_session(&task.id)
.await
.map_err(|error| task_error(error.to_string()))?
.ok_or_else(|| task_error("Task Session disappeared during recovery"))?;
let from = current.status;
current.set_status(TaskSessionStatus::Waiting, reason);
store
.update_task_session(¤t)
.await
.map_err(|error| task_error(error.to_string()))?;
store
.append_task_event(
¤t.id,
&TaskEventKind::StatusChanged {
from,
to: TaskSessionStatus::Waiting,
reason: current.status_reason.clone(),
},
)
.await
.map_err(|error| task_error(error.to_string()))?;
if let Err(error) = relaunch_inactive_process(store, &mut current).await {
let mut persisted = store
.get_task_session(&task.id)
.await
.map_err(|store_error| task_error(store_error.to_string()))?
.ok_or_else(|| task_error("Task Session disappeared during relaunch"))?;
if persisted.status == TaskSessionStatus::Waiting {
let failure = format!(
"body generation {generation} was reaped after a stall but its successor could not start: {error}"
);
record_task_failure(store, &mut persisted, failure.clone(), failure).await?;
}
return Err(error);
}
Ok(true)
}
pub(crate) async fn reconcile_task_pr(
store: &SharedStore,
session: &mut TaskSession,
) -> OpsResult<Option<TaskPr>> {
reconcile_task_pr_with_authority(
store,
session,
None,
crate::ops::pr::PrReadFreshness::Cached,
)
.await
}
pub(crate) fn decide_open_pr_status(
pr: &TaskPr,
github_degraded: Option<&str>,
head_advanced: bool,
) -> (TaskSessionStatus, String) {
let number = pr
.github()
.expect("open Task PR requires a GitHub receipt")
.number;
if let Some(reason) = github_degraded {
return (
TaskSessionStatus::Blocked,
format!(
"ci-fix blocked by github-observation: {reason}. Resume when GitHub recovers; pull request #{number} stays attached."
),
);
}
let failing = pr
.ci_observation
.as_ref()
.is_some_and(|observation| observation.state == CiState::Failing);
if failing && !head_advanced {
return (
TaskSessionStatus::Blocked,
format!(
"CI failing on pull request #{number}; the Task body did not repair the head. Needs a new directive or human review; pull request #{number} stays attached."
),
);
}
(
TaskSessionStatus::Waiting,
format!("pull request #{number} is open for review"),
)
}
pub(crate) fn current_ci_incident(pr: &TaskPr) -> Option<CiIncident> {
let observation = pr.fresh_ci().filter(|reading| reading.wake_legal())?;
ci_incident(pr, observation)
}
pub(crate) async fn route_ci_incident(
store: &SharedStore,
session: &TaskSession,
pr: &TaskPr,
) -> OpsResult<()> {
let Some(incident) = current_ci_incident(pr) else {
return Ok(());
};
if !session.status.is_process_active() {
let mut session = session.clone();
relaunch_for_ci_incident(store, &mut session, incident.identity.clone()).await?;
}
let work = store
.work_for_child(&ChildRef::Task(session.id.clone()))
.await
.map_err(|error| task_error(error.to_string()))?;
let run = store
.current_run(&work)
.await
.map_err(|error| task_error(error.to_string()))?
.ok_or_else(|| task_error("actionable CI has no active Run"))?;
store
.claim_ci_incident(&incident.identity, &run.id, time::OffsetDateTime::now_utc())
.await
.map_err(|error| task_error(error.to_string()))?;
Ok(())
}
pub(crate) async fn reconcile_task_pr_for_lease(
store: &SharedStore,
session: &mut TaskSession,
lease: &ChildWriteLease,
) -> OpsResult<Option<TaskPr>> {
reconcile_task_pr_with_authority(
store,
session,
Some(lease),
crate::ops::pr::PrReadFreshness::Cached,
)
.await
}
pub(crate) async fn reconcile_task_pr_fresh_for_lease(
store: &SharedStore,
session: &mut TaskSession,
lease: &ChildWriteLease,
) -> OpsResult<Option<TaskPr>> {
reconcile_task_pr_with_authority(
store,
session,
Some(lease),
crate::ops::pr::PrReadFreshness::Fresh,
)
.await
}
fn observe_required_checks(
worktree: &Path,
branch: &str,
head_sha: Option<&str>,
now: time::OffsetDateTime,
) -> Option<CiObservation> {
let head_sha = head_sha?.to_string();
let checks = crate::ops::pr::merge_gate_state(worktree, branch)?;
let state = if checks.failing {
CiState::Failing
} else if checks.pending {
CiState::Pending
} else {
CiState::Passing
};
Some(CiObservation {
head_sha,
state,
failing_checks: checks
.failing_leaves
.into_iter()
.map(|check| CiCheck {
name: check.name,
url: check.url,
})
.collect(),
observed_at: now,
})
}
fn ci_incident(pr: &TaskPr, observation: &CiObservation) -> Option<CiIncident> {
if observation.state != CiState::Failing {
return None;
}
let identity = pr.pr_identity()?;
let failure_set = observation.failure_set();
let mut digest = Sha256::new();
for check in &failure_set {
digest.update(check.as_bytes());
digest.update([0]);
}
Some(CiIncident {
identity: format!(
"github:ci:{}:{}:{}:{}",
identity.repo,
identity.number,
observation.head_sha,
hex::encode(digest.finalize())
),
task_session_id: pr.task_session_id.clone(),
pr_id: pr.id.clone(),
repo: identity.repo,
pr_number: identity.number,
failed_head_sha: observation.head_sha.clone(),
repaired_head_sha: None,
failure_set,
provider_completed_at: None,
poll_observed_at: Some(observation.observed_at),
webhook_received_at: None,
claimed_run_id: None,
responded_at: None,
green_at: None,
merged_at: None,
blocked_at: None,
blocked_reason: None,
created_at: observation.observed_at,
updated_at: observation.observed_at,
})
}
const PR_OBSERVATION_TTL: time::Duration = time::Duration::seconds(60);
const PR_OBSERVATION_DEGRADED_BACKOFF: time::Duration = time::Duration::minutes(5);
fn cached_github_observation(pr: &TaskPr, now: time::OffsetDateTime) -> Option<Observation> {
let observation = pr.github_observation.as_ref()?;
let retry_at = observation.checked_at
+ match observation.result {
GithubObservationResult::Fresh => PR_OBSERVATION_TTL,
GithubObservationResult::Degraded { .. } => PR_OBSERVATION_DEGRADED_BACKOFF,
};
if retry_at <= now {
return None;
}
Some(match &observation.result {
GithubObservationResult::Fresh => Observation::Cached {
observed_at: observation.checked_at,
},
GithubObservationResult::Degraded { reason } => Observation::Degraded {
reason: reason.clone(),
cached_as_of: pr.updated_at,
retry_at,
},
})
}
async fn reconcile_subject(
store: &SharedStore,
session: &TaskSession,
) -> OpsResult<Option<TaskPr>> {
if let Some(active) = store
.active_task_pr(&session.id)
.await
.map_err(|error| task_error(format!("failed to read active PR: {error}")))?
{
return Ok(Some(active));
}
let prs = store
.task_prs(&session.id)
.await
.map_err(|error| task_error(format!("failed to read Task PRs: {error}")))?;
Ok(prs
.into_iter()
.next_back()
.filter(|pr| pr.phase() == PrPhase::Abandoned && pr.github().is_some()))
}
async fn reconcile_task_pr_with_authority(
store: &SharedStore,
session: &mut TaskSession,
lease: Option<&ChildWriteLease>,
freshness: crate::ops::pr::PrReadFreshness,
) -> OpsResult<Option<TaskPr>> {
let Some(mut pr) = reconcile_subject(store, session).await? else {
return Ok(None);
};
let Some(number) = pr.github().map(|github| github.number) else {
session.observation = Observation::NotRequired;
return Ok(Some(pr));
};
let now = time::OffsetDateTime::now_utc();
if matches!(freshness, crate::ops::pr::PrReadFreshness::Cached) {
if let Some(observation) = cached_github_observation(&pr, now) {
session.observation = observation;
return Ok(Some(pr));
}
}
let previous = pr.clone();
let github_pr = match crate::ops::pr::observe_pr_by_number(
&session.worktree,
number,
&pr.branch,
freshness,
) {
crate::ops::pr::PrObservation::Fresh(info) => {
pr.github_observation = Some(GithubObservation {
checked_at: now,
result: GithubObservationResult::Fresh,
});
session.observation = Observation::Fresh { observed_at: now };
info
}
crate::ops::pr::PrObservation::NotFound => {
pr.github_observation = Some(GithubObservation {
checked_at: now,
result: GithubObservationResult::Fresh,
});
pr.updated_at = now;
update_task_pr_with_authority(store, &pr, lease)
.await
.map_err(|error| task_error(error.to_string()))?;
session.observation = Observation::Fresh { observed_at: now };
return Ok(Some(pr));
}
crate::ops::pr::PrObservation::Degraded { reason } => {
let retry_at = now + PR_OBSERVATION_DEGRADED_BACKOFF;
pr.github_observation = Some(GithubObservation {
checked_at: now,
result: GithubObservationResult::Degraded {
reason: reason.clone(),
},
});
update_task_pr_with_authority(store, &pr, lease)
.await
.map_err(|error| task_error(error.to_string()))?;
session.observation = Observation::Degraded {
reason,
cached_as_of: pr.updated_at,
retry_at,
};
return Ok(Some(pr));
}
};
let number = u32::try_from(github_pr.number).map_err(|_| {
task_error(format!(
"pull request #{} exceeds supported range",
github_pr.number
))
})?;
let url = github_pr.url.clone();
let previous_phase = previous.phase();
let previous_github = previous.github().cloned();
let previous_session_status = session.status;
let previous_status_reason = session.status_reason.clone();
let previous_pm_writeback = session.pm_writeback.clone();
let publication = pr.publication.get_or_insert(PrPublication {
requested_at: now,
after_merge: AfterMerge::Review,
next_slug: None,
github: None,
});
publication.github = Some(GithubPr {
number,
url: url.clone(),
head_sha: github_pr.head_sha.clone(),
});
let mut observed_incident = None;
let mut green_at = None;
let mut merged_at = None;
let pr_event = match github_pr.state.as_str() {
"merged" => {
let merge_commit = github_pr.merge_commit.clone().ok_or_else(|| {
task_error(format!(
"GitHub reports pull request #{} merged without a merge commit",
github_pr.number
))
})?;
pr.merge_commit = Some(merge_commit.clone());
pr.ci_observation = None;
merged_at = Some(now);
let completes = pr
.publication
.as_ref()
.is_some_and(|publication| publication.after_merge == AfterMerge::CompleteTask)
&& matches!(
committed_follow_up_range(&session.worktree, &pr)?,
CommittedFollowUp::ProvenEmpty
);
if completes {
let gate = review_gate(store, session).await?;
if gate.satisfied {
session.set_status(
TaskSessionStatus::Completed,
format!(
"pull request #{} merged and completed the Task",
github_pr.number
),
);
reconcile_pm_writeback(store, session, Some(&url)).await;
} else if !session.status.is_process_active() {
session.set_status(
TaskSessionStatus::Waiting,
format!(
"pull request #{} merged; awaiting gate before completion: {}",
github_pr.number,
gate.reason()
),
);
}
} else if !session.status.is_process_active() {
session.set_status(
TaskSessionStatus::Waiting,
format!(
"pull request #{} merged; another PR may follow",
github_pr.number
),
);
}
Some(TaskEventKind::PrMerged {
pr_id: pr.id.clone(),
sequence: pr.sequence,
number,
url: url.clone(),
merge_commit,
})
}
"closed" => {
pr.abandoned_at = pr.abandoned_at.or(Some(now));
pr.ci_observation = None;
if !session.status.is_process_active() {
session.set_status(
TaskSessionStatus::Waiting,
format!("pull request #{} closed without merge", github_pr.number),
);
}
None
}
_ => {
pr.abandoned_at = None;
if !session.status.is_process_active() {
session.set_status(
TaskSessionStatus::Waiting,
format!("pull request #{} is open for review", github_pr.number),
);
}
if let Some(ci_observation) = observe_required_checks(
&session.worktree,
&pr.branch,
github_pr.head_sha.as_deref(),
now,
) {
if ci_observation.state == CiState::Passing {
green_at = Some(ci_observation.observed_at);
} else {
observed_incident = ci_incident(&pr, &ci_observation);
}
pr.ci_observation = Some(ci_observation);
}
Some(TaskEventKind::PrOpened {
pr_id: pr.id.clone(),
sequence: pr.sequence,
number,
url: url.clone(),
})
}
};
let pr_changed = pr != previous;
let completes_task = pr.phase() == PrPhase::Merged
&& session.status == TaskSessionStatus::Completed
&& pr
.publication
.as_ref()
.is_some_and(|publication| publication.after_merge == AfterMerge::CompleteTask);
let mut session_saved_with_pr = false;
if pr_changed {
pr.updated_at = now;
if completes_task {
complete_task_session_after_pr_with_authority(store, session, &pr, lease)
.await
.map_err(|error| task_error(error.to_string()))?;
session_saved_with_pr = true;
} else if pr.is_settled() {
settle_task_pr_with_authority(store, &pr, None, lease)
.await
.map_err(|error| task_error(error.to_string()))?;
} else {
update_task_pr_with_authority(store, &pr, lease)
.await
.map_err(|error| task_error(error.to_string()))?;
}
}
if let Some(incident) = observed_incident {
store
.observe_ci_incident(&incident)
.await
.map_err(|error| task_error(error.to_string()))?;
}
if let Some(green_at) = green_at {
store
.mark_ci_incidents_green(&pr.id, green_at)
.await
.map_err(|error| task_error(error.to_string()))?;
}
if let Some(merged_at) = merged_at {
store
.mark_ci_incidents_merged(&pr.id, merged_at)
.await
.map_err(|error| task_error(error.to_string()))?;
}
if !session_saved_with_pr
&& (session.status != previous_session_status
|| session.status_reason != previous_status_reason
|| session.pm_writeback != previous_pm_writeback)
{
update_task_session_with_authority(store, session, lease)
.await
.map_err(|error| task_error(error.to_string()))?;
}
if session.status != previous_session_status {
append_task_event_with_authority(
store,
&session.id,
&TaskEventKind::StatusChanged {
from: previous_session_status,
to: session.status,
reason: session.status_reason.clone(),
},
lease,
)
.await
.map_err(|error| task_error(error.to_string()))?;
}
if pr_changed {
if let Some(event) = pr_event {
let should_append = match &event {
TaskEventKind::PrOpened { .. } => previous_github.as_ref() != pr.github(),
TaskEventKind::PrMerged { .. } => previous_phase != PrPhase::Merged,
_ => true,
};
if should_append {
append_task_event_with_authority(store, &session.id, &event, lease)
.await
.map_err(|error| task_error(error.to_string()))?;
}
}
}
if previous_session_status != TaskSessionStatus::Completed
&& session.status == TaskSessionStatus::Completed
{
append_task_event_with_authority(
store,
&session.id,
&TaskEventKind::Completed {
summary: "pull request merge completed the Task".to_string(),
},
lease,
)
.await
.map_err(|error| task_error(error.to_string()))?;
}
Ok(Some(pr))
}
fn next_pr_slug(settled: &TaskPr, slug_override: Option<&str>) -> String {
slug_override
.map(str::to_string)
.or_else(|| {
settled
.publication
.as_ref()
.and_then(|publication| publication.next_slug.clone())
})
.unwrap_or_else(|| (settled.sequence + 1).to_string())
}
fn deterministic_next_branch(
session: &TaskSession,
settled: &TaskPr,
slug_override: Option<&str>,
) -> OpsResult<String> {
let slug = next_pr_slug(settled, slug_override);
let author = settled
.branch
.split_once('/')
.map(|(author, _)| author)
.ok_or_else(|| {
task_error(format!(
"Task PR branch {:?} has no author prefix",
settled.branch
))
})?;
Ok(format!("{author}/{}-{slug}", session.workspace_slug))
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) enum TaskRecoveryAdoption {
Active { branch: String },
BetweenPrs { settled: String, next: String },
}
pub(crate) async fn task_recovery_adoption(
store: &SharedStore,
session: &TaskSession,
) -> OpsResult<TaskRecoveryAdoption> {
let worktree = &session.worktree;
let identifier = &session.launch.issue.identifier;
if !worktree.exists() {
return Err(task_error(format!(
"Task {identifier} worktree {} is missing; recovery refused before moving any ownership",
worktree.display()
)));
}
if let Some(state) = crate::engine::git::intervention_state(worktree)
.map_err(|error| task_error(format!("failed to inspect Task worktree state: {error}")))?
{
return Err(task_error(format!(
"Task {identifier} worktree {} is mid-{state}; resolve or abort it before resuming, \
recovery refused before moving any ownership",
worktree.display()
)));
}
let current = current_branch(worktree)
.map_err(|error| task_error(format!("failed to inspect Task branch: {error}")))?
.ok_or_else(|| {
task_error(format!(
"Task {identifier} worktree {} is detached; recovery needs a branch",
worktree.display()
))
})?;
if !ref_exists(worktree, &format!("refs/heads/{current}"))
.map_err(|error| task_error(format!("failed to inspect Task branch: {error}")))?
{
return Err(task_error(format!(
"Task {identifier} worktree {} is on branch {current:?} which no longer exists; \
re-create it or recover the worktree before resuming",
worktree.display()
)));
}
if let Some(active) = store
.active_task_pr(&session.id)
.await
.map_err(|error| task_error(format!("failed to read active PR: {error}")))?
{
if current != active.branch {
return Err(task_error(format!(
"Task {identifier} active PR expects branch {:?}, but the worktree is on \
{current:?}; recovery refused before moving any ownership",
active.branch
)));
}
return Ok(TaskRecoveryAdoption::Active {
branch: active.branch,
});
}
let prs = store
.task_prs(&session.id)
.await
.map_err(|error| task_error(format!("failed to read Task PRs: {error}")))?;
let settled = prs
.last()
.cloned()
.ok_or_else(|| task_error("Task Session has no PR history"))?;
if !settled.is_settled() {
return Err(task_error(format!(
"Task PR {} is neither active nor settled",
settled.id
)));
}
let next = deterministic_next_branch(session, &settled, None)?;
if current != settled.branch && current != next {
return Err(task_error(format!(
"Task {identifier} between-PR recovery expected settled branch {:?} or next branch \
{next:?}, but the worktree is on {current:?}; recovery refused before moving any \
ownership",
settled.branch
)));
}
if !is_clean(worktree)
.map_err(|error| task_error(format!("failed to inspect Task worktree: {error}")))?
{
return Err(task_error(format!(
"Task {identifier} cannot recover between PRs while {} has uncommitted changes; \
carry them forward with `lf pr next` or commit before resuming, recovery refused \
before moving any ownership",
worktree.display()
)));
}
Ok(TaskRecoveryAdoption::BetweenPrs {
settled: settled.branch,
next,
})
}
pub(crate) async fn refuse_dirty_between_prs(
store: &SharedStore,
session: &TaskSession,
) -> OpsResult<()> {
if store
.active_task_pr(&session.id)
.await
.map_err(|error| task_error(format!("failed to read active PR: {error}")))?
.is_some()
{
return Ok(());
}
if is_clean(&session.worktree)
.map_err(|error| task_error(format!("failed to inspect Task worktree: {error}")))?
{
return Ok(());
}
Err(task_error(format!(
"Task {} cannot recover between PRs while {} has uncommitted changes; carry them \
forward with `lf pr next` or commit before resuming",
session.launch.issue.identifier,
session.worktree.display()
)))
}
pub(crate) async fn ensure_working_pr(
store: &SharedStore,
session: &mut TaskSession,
) -> OpsResult<Option<TaskPr>> {
ensure_working_pr_with_authority(store, session, None, RotateOptions::runner()).await
}
pub(crate) async fn ensure_working_pr_for_lease(
store: &SharedStore,
session: &mut TaskSession,
lease: &ChildWriteLease,
) -> OpsResult<Option<TaskPr>> {
ensure_working_pr_with_authority(store, session, Some(lease), RotateOptions::runner()).await
}
#[derive(Debug, Clone, Default)]
pub(crate) struct RotateOptions {
carry_dirty: bool,
slug_override: Option<String>,
}
impl RotateOptions {
fn runner() -> Self {
Self::default()
}
}
enum CommittedFollowUp {
ProvenEmpty,
Range { from: String, to: String },
Unprovable { reason: &'static str },
}
fn commits_past(
worktree: &Path,
branch: &str,
cut: &str,
not_ancestor: &'static str,
) -> OpsResult<CommittedFollowUp> {
let tip = rev_parse(worktree, branch)
.map_err(|error| task_error(format!("failed to resolve settled branch tip: {error}")))?;
if tip == cut {
return Ok(CommittedFollowUp::ProvenEmpty);
}
let ancestor = is_ancestor(worktree, cut, branch)
.map_err(|error| task_error(format!("failed to check follow-up ancestry: {error}")))?;
if !ancestor {
return Ok(CommittedFollowUp::Unprovable {
reason: not_ancestor,
});
}
Ok(CommittedFollowUp::Range {
from: cut.to_string(),
to: branch.to_string(),
})
}
fn committed_follow_up_range(worktree: &Path, settled: &TaskPr) -> OpsResult<CommittedFollowUp> {
let Some(head_sha) = settled.github().and_then(|github| github.head_sha.clone()) else {
return Ok(CommittedFollowUp::Unprovable {
reason: "the published pull request head is missing",
});
};
commits_past(
worktree,
&settled.branch,
&head_sha,
"the published pull request head is not an ancestor of the settled branch",
)
}
fn unpublished_work(worktree: &Path, pr: &TaskPr) -> OpsResult<CommittedFollowUp> {
commits_past(
worktree,
&pr.branch,
&pr.base_commit,
"the recorded base is not an ancestor of the unpublished branch",
)
}
fn fork_point(worktree: &Path, base_ref: &str, branch: &str) -> OpsResult<String> {
merge_base(worktree, base_ref, branch).map_err(|error| {
task_error(format!(
"{branch:?} shares no history with {base_ref}: {error}"
))
})
}
async fn heal_incoherent_base(
store: &SharedStore,
session: &TaskSession,
pr: TaskPr,
lease: Option<&ChildWriteLease>,
) -> OpsResult<TaskPr> {
if pr.phase() != PrPhase::Working || !session.worktree.exists() {
return Ok(pr);
}
if is_ancestor(&session.worktree, &pr.base_commit, &pr.branch).unwrap_or(false) {
return Ok(pr);
}
let Ok(default_branch) = get_default_branch(&session.worktree) else {
return Ok(pr);
};
let Ok((base_ref, _)) = resolve_upstream_base(&session.worktree, &default_branch) else {
return Ok(pr);
};
let Ok(fork) = fork_point(&session.worktree, &base_ref, &pr.branch) else {
tracing::warn!(
task = %session.launch.issue.identifier,
branch = %pr.branch,
base = %pr.base_commit,
"Task PR base is incoherent and shares no history with the upstream; \
leaving the row for the completion gate to refuse"
);
return Ok(pr);
};
let ancestry = |commit: &str, descendant: &str| {
is_ancestor(&session.worktree, commit, descendant).unwrap_or(false)
};
if !(ancestry(&fork, &pr.base_commit) && ancestry(&pr.base_commit, &base_ref)) {
tracing::warn!(
task = %session.launch.issue.identifier,
branch = %pr.branch,
base = %pr.base_commit,
"Task PR base is incoherent but is not on the upstream line, so no past mint \
wrote it; leaving the row for the completion gate to refuse"
);
return Ok(pr);
}
let mut healed = pr;
tracing::info!(
task = %session.launch.issue.identifier,
branch = %healed.branch,
from = %healed.base_commit,
to = %fork,
"healing a Task PR base that is not an ancestor of its branch"
);
healed.base_commit = fork;
healed.updated_at = time::OffsetDateTime::now_utc();
match lease {
Some(lease) => store.heal_task_pr_base_for_lease(&healed, lease).await,
None => store.heal_task_pr_base(&healed).await,
}
.map_err(|error| task_error(format!("failed to heal Task PR base: {error}")))?;
Ok(healed)
}
pub(crate) fn no_active_pr_resume_refusal(
identifier: &str,
active: Option<&TaskPr>,
latest: Option<&TaskPr>,
) -> Option<String> {
if active.is_some() {
return None;
}
let suffix = match latest {
Some(pr) => {
let which = pr
.github()
.map(|github| format!("pull request #{}", github.number))
.unwrap_or_else(|| format!("PR sequence {}", pr.sequence));
format!("{which} {}", pr.phase().as_str())
}
None => "no PR history recorded".to_string(),
};
Some(format!(
"Task {identifier} has no active PR to resume; {suffix}"
))
}
fn roll_back_failed_rotation(
worktree: &Path,
settled_branch: &str,
recovery_branch: &str,
stashed: bool,
) -> OpsResult<()> {
checkout(worktree, settled_branch)
.map_err(|error| task_error(format!("failed to restore settled branch: {error}")))?;
delete_local_branch(worktree, recovery_branch)
.map_err(|error| task_error(format!("failed to remove recovery branch: {error}")))?;
if stashed {
stash_pop(worktree)
.map_err(|error| task_error(format!("failed to restore follow-up edits: {error}")))?;
}
Ok(())
}
async fn ensure_working_pr_with_authority(
store: &SharedStore,
session: &mut TaskSession,
lease: Option<&ChildWriteLease>,
rotate: RotateOptions,
) -> OpsResult<Option<TaskPr>> {
reconcile_task_pr_with_authority(
store,
session,
lease,
crate::ops::pr::PrReadFreshness::Cached,
)
.await?;
if session.status.is_terminal() {
return Ok(None);
}
if let Some(active) = store
.active_task_pr(&session.id)
.await
.map_err(|error| task_error(format!("failed to read active PR: {error}")))?
{
return Ok(Some(
heal_incoherent_base(store, session, active, lease).await?,
));
}
let prs = store
.task_prs(&session.id)
.await
.map_err(|error| task_error(format!("failed to read Task PRs: {error}")))?;
let settled = prs
.last()
.cloned()
.ok_or_else(|| task_error("Task Session has no PR history"))?;
if !settled.is_settled() {
return Err(task_error(format!(
"Task PR {} is neither active nor settled",
settled.id
)));
}
if let (PrPhase::Abandoned, Some(github)) = (settled.phase(), settled.github()) {
if let Observation::Degraded { reason, .. } = &session.observation {
return Err(task_error(format!(
"cannot confirm pull request #{} is closed before starting the next PR: {reason}. \
Retry once GitHub is readable; if the PR was reopened, it continues as-is.",
github.number
)));
}
}
let committed_carry = committed_follow_up_range(&session.worktree, &settled)?;
if settled
.publication
.as_ref()
.is_some_and(|publication| publication.after_merge == AfterMerge::CompleteTask)
&& !matches!(&committed_carry, CommittedFollowUp::Range { .. })
{
return Ok(None);
}
if let Some(lease) = lease {
store
.validate_child_write_lease(&ChildRef::Task(session.id.clone()), lease)
.await
.map_err(|error| task_error(format!("Task body lost write authority: {error}")))?;
}
let sequence = settled.sequence + 1;
let slug = next_pr_slug(&settled, rotate.slug_override.as_deref());
let branch = deterministic_next_branch(session, &settled, rotate.slug_override.as_deref())?;
let default_branch = get_default_branch(&session.worktree)
.map_err(|error| task_error(format!("failed to resolve default branch: {error}")))?;
let (base_ref, _) = resolve_upstream_base(&session.worktree, &default_branch)?;
if !rotate.carry_dirty
&& !is_clean(&session.worktree)
.map_err(|error| task_error(format!("failed to inspect Task worktree: {error}")))?
{
return Err(task_error(format!(
"Task {} cannot rotate PRs while {} has uncommitted changes",
session.launch.issue.identifier,
session.worktree.display()
)));
}
let current = current_branch(&session.worktree)
.map_err(|error| task_error(format!("failed to inspect Task branch: {error}")))?
.ok_or_else(|| task_error("Task worktree is detached"))?;
if current != branch {
if current != settled.branch {
return Err(task_error(format!(
"Task {} expected settled branch {:?} or recovery branch {:?}, but {} is on {:?}",
session.launch.issue.identifier,
settled.branch,
branch,
session.worktree.display(),
current
)));
}
let local_ref = format!("refs/heads/{branch}");
let remote_ref = format!("refs/remotes/origin/{branch}");
let collision = ref_exists(&session.worktree, &local_ref)
.map_err(|error| task_error(format!("failed to inspect branch collision: {error}")))?
|| ref_exists(&session.worktree, &remote_ref).map_err(|error| {
task_error(format!("failed to inspect branch collision: {error}"))
})?;
if collision {
return Err(task_error(format!(
"next PR branch {branch:?} already exists; retry the settling command with a clearer --next name"
)));
}
let stashed = stash_including_untracked(&session.worktree)
.map_err(|error| task_error(format!("failed to stash follow-up edits: {error}")))?;
if let Err(error) = checkout_new_branch_from(&session.worktree, &branch, &base_ref) {
let recovered = current_branch(&session.worktree)
.map_err(|read_error| {
task_error(format!("failed to inspect recovery branch: {read_error}"))
})?
.as_deref()
== Some(branch.as_str());
if !recovered {
if stashed {
stash_pop(&session.worktree).map_err(|recovery_error| {
task_error(format!(
"failed to rotate Task worktree: {error}; restoring follow-up edits \
also failed: {recovery_error}"
))
})?;
}
return Err(task_error(format!(
"failed to rotate Task worktree: {error}; follow-up edits were restored"
)));
}
}
if let CommittedFollowUp::Range { from, to } = &committed_carry {
if let Err(error) = cherry_pick_range(&session.worktree, from, to) {
roll_back_failed_rotation(&session.worktree, &settled.branch, &branch, stashed)
.map_err(|recovery_error| {
task_error(format!(
"failed to carry committed follow-up from {:?} onto {branch}: {error}; \
automatic recovery also failed: {recovery_error}",
settled.branch
))
})?;
return Err(task_error(format!(
"failed to carry committed follow-up from {:?} onto {branch}: {error}; \
restored {:?} with its follow-up edits so the rotation can be retried",
settled.branch, settled.branch
)));
}
}
if stashed {
stash_pop(&session.worktree).map_err(|error| {
task_error(format!(
"carried the committed follow-up but could not reapply dirty edits: {error}; \
the recovery branch and retained stash are in {} for conflict resolution",
session.worktree.display()
))
})?;
}
}
let base_commit = fork_point(&session.worktree, &base_ref, &branch)?;
push_with_upstream(&session.worktree, "origin", &branch)
.map_err(|error| task_error(format!("failed to push next PR branch: {error}")))?;
let now = time::OffsetDateTime::now_utc();
let next = TaskPr {
id: TaskPrId::new(),
task_session_id: session.id.clone(),
sequence,
slug,
branch,
base_commit,
parent_pr_id: None,
publication: None,
merge_commit: None,
abandoned_at: None,
ci_observation: None,
github_observation: None,
linear_attachment_id: None,
linear_comment_id: None,
linear_link_error: None,
created_at: now,
updated_at: now,
};
match settle_task_pr_with_authority(store, &settled, Some(&next), lease).await {
Ok(()) => {
append_task_event_with_authority(
store,
&session.id,
&TaskEventKind::PrStarted {
pr_id: next.id.clone(),
sequence: next.sequence,
branch: next.branch.clone(),
base_commit: next.base_commit.clone(),
},
lease,
)
.await
.map_err(|error| task_error(error.to_string()))?;
Ok(Some(next))
}
Err(error) => {
let recovered = store
.task_prs(&session.id)
.await
.map_err(|read_error| task_error(read_error.to_string()))?
.into_iter()
.find(|pr| pr.sequence == sequence);
match recovered {
Some(pr)
if pr.branch == next.branch
&& pr.base_commit == next.base_commit
&& pr.phase() == PrPhase::Working =>
{
Ok(Some(pr))
}
_ => Err(task_error(format!(
"failed to record next Task PR after branch rotation: {error}"
))),
}
}
}
}
pub fn pr_next(repo: &Path, slug: Option<&str>) -> OpsResult<TaskPr> {
let slug_override = slug.map(parse_pr_slug).transpose()?;
let repo = repo.to_path_buf();
block_on_task(async move {
let store = task_store().await?;
let (mut session, lease) = task_for_worktree(&store, &repo)
.await?
.ok_or_else(|| task_error("no Task Session owns this worktree"))?;
reconcile_task_pr_with_authority(
&store,
&mut session,
lease.as_ref(),
crate::ops::pr::PrReadFreshness::Cached,
)
.await?;
if let Some(active) = store
.active_task_pr(&session.id)
.await
.map_err(|error| task_error(format!("failed to read active PR: {error}")))?
{
let which = active
.github()
.map(|github| format!("#{}", github.number))
.unwrap_or_else(|| format!("sequence {}", active.sequence));
return Err(task_error(format!(
"current PR {which} is not merged yet; land it or wait for the merge before `lf pr next`"
)));
}
if session.status.is_terminal() {
return Err(task_error(format!(
"Task {} is already {}; nothing to rotate",
session.launch.issue.identifier,
session.status.as_str()
)));
}
let rotate = RotateOptions {
carry_dirty: true,
slug_override,
};
ensure_working_pr_with_authority(&store, &mut session, lease.as_ref(), rotate)
.await?
.ok_or_else(|| task_error("Task has no settled PR to rotate from"))
})
}
pub fn task_status(issue: &str) -> OpsResult<TaskSession> {
block_on_task(async move {
let store = task_store().await?;
let mut session = store
.get_task_session_by_issue(issue)
.await
.map_err(|error| task_error(format!("failed to read task status: {error}")))?
.ok_or_else(|| task_error(format!("no Task Session exists for {issue:?}")))?;
reconcile_task_pr(&store, &mut session).await?;
reconcile_process_liveness(&store, &mut session).await?;
reconcile_task_completion(&store, &mut session, None).await?;
Ok(session)
})
}
pub fn task_complete(issue: &str, summary: String) -> OpsResult<TaskSession> {
let summary = summary.trim().to_string();
if summary.is_empty() {
return Err(task_error("completion summary cannot be empty"));
}
block_on_task(async move {
let store = task_store().await?;
let mut session = store
.get_task_session_by_issue(issue)
.await
.map_err(|error| task_error(format!("failed to read Task Session: {error}")))?
.ok_or_else(|| task_error(format!("no Task Session exists for {issue:?}")))?;
reconcile_task_pr(&store, &mut session).await?;
let lease = ambient_task_write_lease(&session)?;
if let Some(lease) = lease.as_ref() {
store
.validate_child_write_lease(&ChildRef::Task(session.id.clone()), lease)
.await
.map_err(|error| task_error(format!("Task body lost write authority: {error}")))?;
}
if session.status == TaskSessionStatus::Completed {
return Ok(session);
}
if session.status == TaskSessionStatus::Abandoned {
return Err(task_error(format!(
"Task {} is abandoned and cannot be completed",
session.launch.issue.identifier
)));
}
if !is_clean(&session.worktree)
.map_err(|error| task_error(format!("failed to inspect Task worktree: {error}")))?
{
return Err(task_error(
"Task worktree has uncommitted changes; publish or explicitly abandon them first",
));
}
let gate = task_completion_gate(&store, &session).await?;
if let Some(refusal) = gate.refusal(&session.launch.issue.identifier) {
return Err(task_error(refusal));
}
let from = session.status;
session.set_status(TaskSessionStatus::Completed, summary.clone());
reconcile_pm_writeback(&store, &mut session, None).await;
complete_task_session_with_authority(
&store,
&session,
gate.discardable_successor.as_ref(),
lease.as_ref(),
)
.await
.map_err(|error| task_error(format!("failed to complete Task Session: {error}")))?;
append_task_event_with_authority(
&store,
&session.id,
&TaskEventKind::StatusChanged {
from,
to: TaskSessionStatus::Completed,
reason: session.status_reason.clone(),
},
lease.as_ref(),
)
.await
.map_err(|error| task_error(error.to_string()))?;
append_task_event_with_authority(
&store,
&session.id,
&TaskEventKind::Completed { summary },
lease.as_ref(),
)
.await
.map_err(|error| task_error(error.to_string()))?;
Ok(session)
})
}
fn pr_link_state_label(pr: &TaskPr) -> String {
match pr.phase() {
PrPhase::Merged => "Merged".to_string(),
PrPhase::Abandoned => "Abandoned".to_string(),
_ => {
let completes = pr
.publication
.as_ref()
.is_some_and(|publication| publication.after_merge == AfterMerge::CompleteTask);
if completes {
"Open · completes task on merge".to_string()
} else {
"Open · in review".to_string()
}
}
}
}
async fn link_pr_to_linear(store: &SharedStore, session: &TaskSession, pr: &mut TaskPr) {
let Some(github) = pr.github().cloned() else {
return;
};
let state = pr_link_state_label(pr);
let title = format!("GitHub PR #{}", github.number);
let body = format!("[GitHub PR #{}]({}) — {}", github.number, github.url, state);
let wave = match owning_wave(store, session).await {
Ok(wave) => wave,
Err(error) => {
pr.linear_link_error = Some(error.to_string());
return;
}
};
let prior = crate::ops::pm::PrLinkageIds {
attachment_id: pr.linear_attachment_id.clone(),
comment_id: pr.linear_comment_id.clone(),
};
let request = crate::ops::pm::PrLinkRequest {
issue_id: session.launch.issue.id.as_str().to_string(),
url: github.url.clone(),
title,
subtitle: state,
body,
};
let outcome =
crate::ops::pm::pm_link_pr_async(&session.worktree, wave.name(), &request, &prior).await;
if let Some(error) = &outcome.error {
tracing::warn!(
issue = session.launch.issue.identifier,
pr = github.number,
"Linear link degraded; the GitHub PR is published and the next publish retries: {error}"
);
}
pr.linear_attachment_id = outcome.ids.attachment_id;
pr.linear_comment_id = outcome.ids.comment_id;
pr.linear_link_error = outcome.error;
}
fn writeback_state(result: OpsResult<()>) -> PmWritebackState {
writeback_state_for(PmWritebackOperation::CompleteTask, result)
}
fn writeback_state_for(operation: PmWritebackOperation, result: OpsResult<()>) -> PmWritebackState {
match result {
Ok(()) => PmWritebackState::Current,
Err(error) => PmWritebackState::Pending {
operation,
error: error.to_string(),
},
}
}
pub(crate) async fn reconcile_pm_writeback(
store: &SharedStore,
session: &mut TaskSession,
pr_url: Option<&str>,
) {
let Ok(wave) = owning_wave(store, session).await else {
session.pm_writeback = PmWritebackState::Pending {
operation: PmWritebackOperation::CompleteTask,
error: format!("owning Wave {} is not registered", session.wave_id),
};
return;
};
session.pm_writeback = writeback_state(
crate::ops::task_pm::complete_task(
&session.worktree,
wave.name(),
session.launch.issue.id.as_str(),
pr_url,
)
.await,
);
}
async fn retry_pm_writeback(store: &SharedStore, session: &mut TaskSession) {
let Ok(prs) = store.task_prs(&session.id).await else {
return;
};
let pr_url = prs
.iter()
.rev()
.find_map(|pr| pr.github().map(|github| github.url.as_str()));
let Ok(wave) = owning_wave(store, session).await else {
session.pm_writeback = PmWritebackState::Pending {
operation: PmWritebackOperation::CompleteTask,
error: format!("owning Wave {} is not registered", session.wave_id),
};
return;
};
session.pm_writeback = {
let operation = match &session.pm_writeback {
PmWritebackState::Pending { operation, .. } => *operation,
PmWritebackState::Current => PmWritebackOperation::CompleteTask,
};
let result = match operation {
PmWritebackOperation::CompleteTask => {
crate::ops::task_pm::retry_complete_task(
&session.worktree,
wave.name(),
session.launch.issue.id.as_str(),
pr_url,
)
.await
}
PmWritebackOperation::ReopenTask => {
crate::ops::task_pm::retry_reopen_task(
&session.worktree,
wave.name(),
session.launch.issue.id.as_str(),
)
.await
}
};
writeback_state_for(operation, result)
};
session.updated_at = time::OffsetDateTime::now_utc();
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct CompletionGate {
pub satisfied: bool,
pub blockers: Vec<String>,
pub discardable_successor: Option<TaskPr>,
}
impl CompletionGate {
pub fn reason(&self) -> String {
if self.blockers.is_empty() {
String::new()
} else {
self.blockers.join("; ")
}
}
pub(crate) fn refusal(&self, identifier: &str) -> Option<String> {
(!self.satisfied).then(|| {
format!(
"Task {identifier} cannot complete until its gates close: {}",
self.reason()
)
})
}
}
impl CompletionGate {
fn unsatisfied(blockers: Vec<String>) -> Self {
Self {
satisfied: false,
blockers,
discardable_successor: None,
}
}
}
async fn review_gate(store: &SharedStore, session: &TaskSession) -> OpsResult<CompletionGate> {
let work = store
.work_for_child(&ChildRef::Task(session.id.clone()))
.await
.map_err(|error| task_error(error.to_string()))?;
let review = store
.review(&work)
.await
.map_err(|error| task_error(error.to_string()))?;
Ok(if review.is_none() {
CompletionGate {
satisfied: true,
blockers: Vec::new(),
discardable_successor: None,
}
} else {
CompletionGate::unsatisfied(vec![
"current interactive flow step has not been closed".to_string()
])
})
}
pub(crate) async fn task_completion_gate(
store: &SharedStore,
session: &TaskSession,
) -> OpsResult<CompletionGate> {
let mut gate = review_gate(store, session).await?;
let prs = store
.task_prs(&session.id)
.await
.map_err(|error| task_error(format!("failed to read Task PRs: {error}")))?;
let has_merged_predecessor = prs.iter().any(|pr| pr.phase() == PrPhase::Merged);
if let Some(newest) = prs.last() {
if newest.phase() == PrPhase::Merged
&& newest
.publication
.as_ref()
.is_some_and(|publication| publication.after_merge == AfterMerge::CompleteTask)
{
let number = newest
.github()
.map(|github| github.number)
.unwrap_or_default();
match committed_follow_up_range(&session.worktree, newest)? {
CommittedFollowUp::ProvenEmpty => {}
CommittedFollowUp::Range { .. } => gate.blockers.push(format!(
"follow-up work is committed past merged pull request #{number}"
)),
CommittedFollowUp::Unprovable { .. }
if session.status == TaskSessionStatus::Completed => {}
CommittedFollowUp::Unprovable { reason } => gate.blockers.push(format!(
"cannot prove merged pull request #{number} has no committed follow-up: {reason}"
)),
}
}
}
if let Some(pr) = store
.active_task_pr(&session.id)
.await
.map_err(|error| task_error(format!("failed to read active PR: {error}")))?
{
let which = pr
.github()
.map(|github| format!("#{}", github.number))
.unwrap_or_else(|| format!("sequence {}", pr.sequence));
match pr.phase() {
PrPhase::Open => gate.blockers.push(format!(
"pull request {which} is open for review; merge it or run `lf pr abandon`"
)),
PrPhase::Publishing => gate.blockers.push(format!(
"pull request {which} is still publishing; wait for it to land or run `lf pr abandon`"
)),
PrPhase::Working => match unpublished_work(&session.worktree, &pr)? {
CommittedFollowUp::ProvenEmpty if has_merged_predecessor => {
gate.discardable_successor = Some(pr.clone());
}
CommittedFollowUp::ProvenEmpty => gate.blockers.push(format!(
"pull request {which} is unpublished; publish and merge it or run `lf pr abandon`"
)),
CommittedFollowUp::Range { .. } => gate.blockers.push(format!(
"follow-up work is committed on unpublished pull request {which}; \
publish and merge it or run `lf pr abandon`"
)),
CommittedFollowUp::Unprovable { reason } => gate.blockers.push(format!(
"cannot prove unpublished pull request {which} is empty: {reason}"
)),
},
PrPhase::Merged | PrPhase::Abandoned => {}
}
}
gate.satisfied = gate.blockers.is_empty();
Ok(gate)
}
async fn merged_completing_pr(
store: &SharedStore,
session: &TaskSession,
) -> OpsResult<Option<TaskPr>> {
let prs = store
.task_prs(&session.id)
.await
.map_err(|error| task_error(format!("failed to read Task PRs: {error}")))?;
Ok(prs.into_iter().find(|pr| {
pr.phase() == PrPhase::Merged
&& pr
.publication
.as_ref()
.is_some_and(|publication| publication.after_merge == AfterMerge::CompleteTask)
}))
}
async fn advance_completion_after_gate(
store: &SharedStore,
session: &mut TaskSession,
lease: Option<&ChildWriteLease>,
) -> OpsResult<bool> {
if session.status.is_terminal() || session.status.is_process_active() {
return Ok(false);
}
let Some(pr) = merged_completing_pr(store, session).await? else {
return Ok(false);
};
let gate = task_completion_gate(store, session).await?;
if !gate.satisfied {
return Ok(false);
}
if gate.discardable_successor.is_some() {
return Ok(false);
}
let from = session.status;
let url = pr.github().map(|github| github.url.clone());
session.set_status(
TaskSessionStatus::Completed,
format!(
"pull request #{} merged and completed the Task",
pr.github().map(|github| github.number).unwrap_or_default()
),
);
reconcile_pm_writeback(store, session, url.as_deref()).await;
complete_task_session_after_pr_with_authority(store, session, &pr, lease)
.await
.map_err(|error| task_error(error.to_string()))?;
append_task_event_with_authority(
store,
&session.id,
&TaskEventKind::StatusChanged {
from,
to: TaskSessionStatus::Completed,
reason: session.status_reason.clone(),
},
lease,
)
.await
.map_err(|error| task_error(error.to_string()))?;
Ok(true)
}
async fn repair_premature_completion(
store: &SharedStore,
session: &mut TaskSession,
lease: Option<&ChildWriteLease>,
) -> OpsResult<bool> {
if session.status != TaskSessionStatus::Completed {
return Ok(false);
}
if matches!(session.pm_writeback, PmWritebackState::Pending { .. }) {
return Ok(false);
}
let gate = task_completion_gate(store, session).await?;
if gate.satisfied {
return Ok(false);
}
let from = session.status;
session.set_status(
TaskSessionStatus::Waiting,
format!("reopened: completion outran its gates ({})", gate.reason()),
);
session.pm_writeback = PmWritebackState::Pending {
operation: PmWritebackOperation::ReopenTask,
error: "premature completion pending reopen".to_string(),
};
if let Some(lease) = lease {
store
.update_task_session_for_lease(session, lease)
.await
.map_err(|error| task_error(error.to_string()))?;
} else {
store
.update_task_session(session)
.await
.map_err(|error| task_error(error.to_string()))?;
}
append_task_event_with_authority(
store,
&session.id,
&TaskEventKind::StatusChanged {
from,
to: TaskSessionStatus::Waiting,
reason: session.status_reason.clone(),
},
lease,
)
.await
.map_err(|error| task_error(error.to_string()))?;
Ok(true)
}
pub(crate) async fn reconcile_task_completion(
store: &SharedStore,
session: &mut TaskSession,
lease: Option<&ChildWriteLease>,
) -> OpsResult<()> {
if session.status == TaskSessionStatus::Completed
&& matches!(session.pm_writeback, PmWritebackState::Pending { .. })
{
retry_pm_writeback(store, session).await;
if let Some(lease) = lease {
store
.update_task_session_for_lease(session, lease)
.await
.map_err(|error| task_error(error.to_string()))?;
} else {
store
.update_task_session(session)
.await
.map_err(|error| task_error(error.to_string()))?;
}
return Ok(());
}
if repair_premature_completion(store, session, lease).await? {
return Ok(());
}
advance_completion_after_gate(store, session, lease).await?;
Ok(())
}
pub fn task_snapshot(session: &TaskSession) -> OpsResult<TaskSessionSnapshot> {
let session = session.clone();
block_on_task(async move {
let store = task_store().await?;
let wave = owning_wave(&store, &session).await?;
let process_alive = if session.status.is_process_active() {
match session.latest_process.as_ref() {
Some(process) => tmux_session_exists(&process.tmux_name)
.await
.map_err(|error| task_error(error.to_string()))?,
None => false,
}
} else {
false
};
let latest_event = store
.task_events_after(&session.id, 0)
.await
.map_err(|error| task_error(format!("failed to read task events: {error}")))?
.into_iter()
.last();
let prs = store
.task_prs(&session.id)
.await
.map_err(|error| task_error(format!("failed to read Task PRs: {error}")))?;
let latest = prs.last();
let active = prs.iter().find(|pr| pr.is_active());
let active_pr = active.map(|pr| pr.id.clone());
let predecessor_phase = match active.and_then(|pr| pr.parent_pr_id.as_ref()) {
Some(parent_id) => store
.get_task_pr(parent_id)
.await
.map_err(|error| task_error(format!("failed to read parent PR: {error}")))?
.map(|pr| pr.phase()),
None => None,
};
let work = store
.work_for_child(&ChildRef::Task(session.id.clone()))
.await
.map_err(|error| task_error(format!("failed to resolve Task Work: {error}")))?;
let review_gate = store
.review(&work)
.await
.map_err(|error| task_error(format!("failed to read review gate: {error}")))?
.map(|_| ReviewGateState::Active);
let completion_gate = task_completion_gate(&store, &session).await?;
let completion_refusal = completion_gate.refusal(&session.launch.issue.identifier);
let resume_refusal =
no_active_pr_resume_refusal(&session.launch.issue.identifier, active, latest);
let action_evidence = TaskActionEvidence {
status: session.status,
latest_pr_phase: latest.map(|pr| pr.phase()),
latest_pr_after_merge: latest
.and_then(|pr| pr.publication.as_ref())
.map(|p| p.after_merge),
latest_pr_next_slug: latest
.and_then(|pr| pr.publication.as_ref())
.and_then(|p| p.next_slug.as_deref()),
completion_refusal: completion_refusal.as_deref(),
resume_refusal: resume_refusal.as_deref(),
pending_directive: false,
ci: active.and_then(|pr| pr.fresh_ci()),
process_alive: if session.status.is_process_active() {
Some(process_alive)
} else {
None
},
predecessor_phase,
review_gate,
abandon_intent: session.abandon_intent.is_some(),
local_progress_unsettled: None,
};
let actions = derive_task_actions(&action_evidence);
let (predecessor_session_id, successor_session_id) = store
.task_session_chain_neighbors(&session.id)
.await
.map_err(|error| task_error(format!("failed to read Task chain: {error}")))?;
let (routing_project_session_id, project_route_succeeded) =
match crate::ops::project::resolve_task_project_route(store.as_ref(), &session).await {
Ok(route) => (Some(route.current.to_string()), route.succeeded),
Err(_) => (None, false),
};
Ok(TaskSessionSnapshot {
issue_id: session.launch.issue.id.as_str().to_string(),
issue_identifier: session.launch.issue.identifier,
session_id: session.id.to_string(),
project_id: session.launch.project.id.as_str().to_string(),
project: session.launch.project.slug,
pm_snapshot_synced_at: session.launch.pm_snapshot_synced_at,
pm_writeback: session.pm_writeback,
wave: wave.name().to_string(),
project_session_id: session.project_session_id.to_string(),
routing_project_session_id,
project_route_succeeded,
predecessor_session_id,
successor_session_id,
status: session.status,
status_reason: session.status_reason,
status_at: session.status_at,
worktree: session.worktree.display().to_string(),
workspace_slug: session.workspace_slug,
lifecycle: session.lifecycle,
lifecycle_phase: session.lifecycle_phase,
phase_epoch: session.phase_epoch,
phase_cursor: session.phase_cursor,
phase_iteration: session.phase_iteration,
gate_cycle: session.gate_cycle,
gate_proposal: session.gate_proposal,
prs,
active_pr,
agent: session.agent,
provider: session.provider,
provider_session_id: session.provider_session_id,
process_alive,
latest_process: session.latest_process,
latest_event,
created_at: session.created_at,
updated_at: session.updated_at,
observation: session.observation,
actions,
})
})
}
pub fn task_changes(issue: &str) -> OpsResult<TaskChangesSnapshot> {
let session = task_status(issue)?;
let pr = active_pr(&session)?;
changes_snapshot(TaskWorkspace::new(&session, &pr))
}
fn changes_snapshot(workspace: TaskWorkspace<'_>) -> OpsResult<TaskChangesSnapshot> {
let mut files = BTreeMap::<String, TaskChangedFile>::new();
record_changed_paths(
workspace.worktree,
&[
"diff",
"--name-only",
"-z",
&format!("{}..HEAD", workspace.base_commit),
],
&mut files,
|file| file.committed = true,
)?;
record_changed_paths(
workspace.worktree,
&["diff", "--cached", "--name-only", "-z"],
&mut files,
|file| file.staged = true,
)?;
record_changed_paths(
workspace.worktree,
&["diff", "--name-only", "-z"],
&mut files,
|file| file.unstaged = true,
)?;
record_changed_paths(
workspace.worktree,
&["ls-files", "--others", "--exclude-standard", "-z"],
&mut files,
|file| file.untracked = true,
)?;
let head_commit = git_output(workspace.worktree, &["rev-parse", "HEAD"])?
.trim()
.to_string();
Ok(TaskChangesSnapshot {
issue_identifier: workspace.issue_identifier.to_string(),
session_id: workspace.session_id.to_string(),
base_commit: workspace.base_commit.to_string(),
head_commit,
files: files.into_values().collect(),
})
}
pub fn task_diff(issue: &str, path: Option<&str>) -> OpsResult<TaskDiffSnapshot> {
let session = task_status(issue)?;
let pr = active_pr(&session)?;
diff_snapshot(TaskWorkspace::new(&session, &pr), path)
}
fn diff_snapshot(workspace: TaskWorkspace<'_>, path: Option<&str>) -> OpsResult<TaskDiffSnapshot> {
const MAX_PATCH_BYTES: usize = 1_000_000;
let relative = path.map(validate_task_relative_path).transpose()?;
let mut args = vec![
"diff".to_string(),
"--no-ext-diff".to_string(),
"--no-color".to_string(),
workspace.base_commit.to_string(),
"--".to_string(),
];
if let Some(path) = &relative {
args.push(path.clone());
}
let mut patch = git_output_owned(workspace.worktree, &args)?;
let untracked = untracked_paths(workspace.worktree)?;
let include_untracked = untracked
.into_iter()
.filter(|candidate| relative.as_ref().is_none_or(|path| path == candidate));
for path in include_untracked {
let output = Command::new("git")
.current_dir(workspace.worktree)
.args(["diff", "--no-index", "--no-color", "--", "/dev/null", &path])
.output()
.map_err(|error| {
task_error(format!("failed to diff untracked file {path}: {error}"))
})?;
if !output.status.success() && output.status.code() != Some(1) {
return Err(task_error(format!(
"failed to diff untracked file {path}: {}",
String::from_utf8_lossy(&output.stderr).trim()
)));
}
patch.extend_from_slice(&output.stdout);
}
let truncated = patch.len() > MAX_PATCH_BYTES;
if truncated {
patch.truncate(MAX_PATCH_BYTES);
}
let patch = String::from_utf8_lossy(&patch).into_owned();
let binary = patch.contains("Binary files ") || patch.contains("GIT binary patch");
Ok(TaskDiffSnapshot {
issue_identifier: workspace.issue_identifier.to_string(),
session_id: workspace.session_id.to_string(),
path: relative,
patch,
binary,
truncated,
})
}
pub fn task_file(issue: &str, path: &str) -> OpsResult<TaskFileSnapshot> {
let session = task_status(issue)?;
let pr = active_pr(&session)?;
file_snapshot(TaskWorkspace::new(&session, &pr), path)
}
fn file_snapshot(workspace: TaskWorkspace<'_>, path: &str) -> OpsResult<TaskFileSnapshot> {
const MAX_FILE_BYTES: usize = 1_000_000;
let relative = validate_task_relative_path(path)?;
let root = workspace
.worktree
.canonicalize()
.map_err(|error| task_error(format!("cannot resolve Task worktree: {error}")))?;
let absolute = root
.join(&relative)
.canonicalize()
.map_err(|error| task_error(format!("cannot open Task file {relative:?}: {error}")))?;
if !absolute.starts_with(&root) || !absolute.is_file() {
return Err(task_error(format!(
"Task file {relative:?} does not resolve to a file inside the Task worktree"
)));
}
let bytes = std::fs::read(&absolute)
.map_err(|error| task_error(format!("cannot read Task file {relative:?}: {error}")))?;
let size_bytes = bytes.len() as u64;
let binary = bytes.iter().take(8_192).any(|byte| *byte == 0);
let truncated = bytes.len() > MAX_FILE_BYTES;
let visible = &bytes[..bytes.len().min(MAX_FILE_BYTES)];
let content = (!binary).then(|| String::from_utf8_lossy(visible).into_owned());
Ok(TaskFileSnapshot {
issue_identifier: workspace.issue_identifier.to_string(),
session_id: workspace.session_id.to_string(),
path: relative,
content,
binary,
size_bytes,
truncated,
})
}
fn validate_task_relative_path(path: &str) -> OpsResult<String> {
let path = Path::new(path);
if path.as_os_str().is_empty() || path.is_absolute() {
return Err(task_error(
"Task paths must stay relative to the Task worktree",
));
}
let mut normalized = PathBuf::new();
for component in path.components() {
match component {
Component::Normal(value) => normalized.push(value),
Component::CurDir => {}
_ => {
return Err(task_error(
"Task paths must stay relative to the Task worktree",
))
}
}
}
if normalized.as_os_str().is_empty() {
return Err(task_error("Task paths must name a file"));
}
Ok(normalized.to_string_lossy().to_string())
}
fn record_changed_paths(
worktree: &Path,
args: &[&str],
files: &mut BTreeMap<String, TaskChangedFile>,
mark: impl Fn(&mut TaskChangedFile),
) -> OpsResult<()> {
for path in nul_paths(&git_output_bytes(worktree, args)?) {
let file = files.entry(path.clone()).or_insert(TaskChangedFile {
path,
committed: false,
staged: false,
unstaged: false,
untracked: false,
});
mark(file);
}
Ok(())
}
fn untracked_paths(worktree: &Path) -> OpsResult<Vec<String>> {
Ok(nul_paths(&git_output_bytes(
worktree,
&["ls-files", "--others", "--exclude-standard", "-z"],
)?))
}
fn nul_paths(output: &[u8]) -> Vec<String> {
output
.split(|byte| *byte == 0)
.filter(|path| !path.is_empty())
.map(|path| String::from_utf8_lossy(path).into_owned())
.collect()
}
fn git_output(worktree: &Path, args: &[&str]) -> OpsResult<String> {
Ok(String::from_utf8_lossy(&git_output_bytes(worktree, args)?).into_owned())
}
fn git_output_bytes(worktree: &Path, args: &[&str]) -> OpsResult<Vec<u8>> {
let output = Command::new("git")
.current_dir(worktree)
.args(args)
.output()
.map_err(|error| task_error(format!("failed to run git {}: {error}", args.join(" "))))?;
if !output.status.success() {
return Err(task_error(format!(
"git {} failed: {}",
args.join(" "),
String::from_utf8_lossy(&output.stderr).trim()
)));
}
Ok(output.stdout)
}
fn git_output_owned(worktree: &Path, args: &[String]) -> OpsResult<Vec<u8>> {
let output = Command::new("git")
.current_dir(worktree)
.args(args)
.output()
.map_err(|error| task_error(format!("failed to run git diff: {error}")))?;
if !output.status.success() {
return Err(task_error(format!(
"git diff failed: {}",
String::from_utf8_lossy(&output.stderr).trim()
)));
}
Ok(output.stdout)
}
fn queue_task_steer(issue: &str, message: String) -> OpsResult<TaskControlResult> {
block_on_task(async move {
let store = task_store().await?;
let mut session = store
.get_task_session_by_issue(issue)
.await
.map_err(|error| task_error(format!("failed to resolve task: {error}")))?
.ok_or_else(|| task_error(format!("no Task Session exists for {issue:?}")))?;
reconcile_task_pr(&store, &mut session).await?;
reconcile_process_liveness(&store, &mut session).await?;
let receipt =
super::child::append_steer(&store, ChildRef::Task(session.id.clone()), &message)
.await?;
if !session.status.is_process_active() {
relaunch_inactive_process(&store, &mut session).await?;
}
Ok(TaskControlResult {
issue_id: session.launch.issue.identifier.clone(),
session_id: session.id.to_string(),
receipt: super::child::WorkControlReceipt::Steer { receipt },
observation: session.observation.clone(),
})
})
}
pub fn task_steer(issue: &str, message: String) -> OpsResult<TaskControlResult> {
queue_task_steer(issue, message)
}
pub fn task_interrupt(issue: &str) -> OpsResult<TaskControlResult> {
block_on_task(async move {
let store = task_store().await?;
let mut session = store
.get_task_session_by_issue(issue)
.await
.map_err(|error| task_error(format!("failed to resolve task: {error}")))?
.ok_or_else(|| task_error(format!("no Task Session exists for {issue:?}")))?;
reconcile_task_pr(&store, &mut session).await?;
reconcile_process_liveness(&store, &mut session).await?;
let work = store
.work_for_child(&ChildRef::Task(session.id.clone()))
.await
.map_err(|error| task_error(error.to_string()))?;
let run = store
.current_run(&work)
.await
.map_err(|error| task_error(error.to_string()))?
.ok_or_else(|| task_error("Task has no active Run to interrupt"))?;
let request = AuthenticatedRequest::cli();
let receipt = store
.interrupt(&ControlCtx::User(&request), &work, &run.id)
.await
.map_err(|error| task_error(error.to_string()))?;
Ok(TaskControlResult {
issue_id: session.launch.issue.identifier.clone(),
session_id: session.id.to_string(),
receipt: super::child::WorkControlReceipt::Interrupt { receipt },
observation: session.observation,
})
})
}
pub fn task_recover(issue: &str, reason: Option<String>) -> OpsResult<TaskSession> {
let issue = issue.to_string();
block_on_task(async move {
let store = task_store().await?;
_recover_abandoned_task(&store, &issue, reason).await
})
}
async fn _recover_abandoned_task(
store: &SharedStore,
issue: &str,
reason: Option<String>,
) -> OpsResult<TaskSession> {
let reason = reason
.map(|value| value.trim().to_string())
.map(|value| {
if value.is_empty() {
Err(task_error("recovery reason cannot be empty"))
} else {
Ok(value)
}
})
.transpose()?;
let predecessor = store
.get_task_session_by_issue(issue)
.await
.map_err(|error| task_error(format!("failed to resolve task: {error}")))?
.ok_or_else(|| task_error(format!("no Task Session exists for {issue:?}")))?;
if predecessor.status != TaskSessionStatus::Abandoned {
if predecessor.status == TaskSessionStatus::Completed {
return Err(task_error(format!(
"Task {} is completed; start a new Task rather than recovering it",
predecessor.launch.issue.identifier
)));
}
let (prior_id, _) = store
.task_session_chain_neighbors(&predecessor.id)
.await
.map_err(|error| task_error(format!("failed to read Task chain: {error}")))?;
if let Some(prior_id) = prior_id {
let prior_id = TaskSessionId::parse(&prior_id)
.map_err(|error| task_error(format!("invalid Task predecessor: {error}")))?;
if let Some(prior) = store
.get_task_session(&prior_id)
.await
.map_err(|error| task_error(format!("failed to read Task predecessor: {error}")))?
{
if prior.status == TaskSessionStatus::Abandoned
&& prior.worktree == predecessor.worktree
{
return Ok(predecessor);
}
}
}
return Err(task_error(format!(
"Task {} is {}; resume it with `lf task resume {}`. Recover is only for abandoned Tasks.",
predecessor.launch.issue.identifier,
predecessor.status.as_str(),
predecessor.launch.issue.identifier
)));
}
task_recovery_adoption(store, &predecessor).await?;
let mut carried = store
.boundary_seed_for_child(&ChildRef::Task(predecessor.id.clone()))
.await
.map_err(|error| task_error(error.to_string()))?
.render();
if carried.is_empty() {
carried = format!(
"Continue {}: {}",
predecessor.launch.issue.identifier, predecessor.launch.issue.title
);
}
let now = time::OffsetDateTime::now_utc();
let status_reason = reason.map_or_else(
|| format!("recovered from {}; resume to continue", predecessor.id),
|reason| {
format!(
"recovered from {}: {reason}; resume to continue",
predecessor.id
)
},
);
let successor = TaskSession {
id: TaskSessionId::new(),
launch: predecessor.launch.clone(),
pm_writeback: PmWritebackState::Current,
wave_id: predecessor.wave_id.clone(),
project_session_id: predecessor.project_session_id.clone(),
status: TaskSessionStatus::Waiting,
status_reason,
status_at: now,
worktree: predecessor.worktree.clone(),
workspace_slug: predecessor.workspace_slug.clone(),
lifecycle: predecessor.lifecycle.clone(),
lifecycle_phase: crate::task::TaskLifecyclePhase::Kickoff,
phase_epoch: 1,
phase_cursor: 0,
phase_iteration: 0,
gate_cycle: 0,
gate_proposal: None,
agent: predecessor.agent.clone(),
provider: predecessor.provider.clone(),
provider_session_id: None,
latest_process: None,
abandon_intent: None,
created_at: now,
updated_at: now,
observation: Observation::NotRequired,
};
let succession = store
.recover_task_session_successor(
&predecessor,
&successor,
crate::durable::Author::User,
&carried,
)
.await
.map_err(|error| task_error(format!("failed to recover Task: {error}")))?;
Ok(succession.session)
}
pub fn task_resume(
issue: &str,
model: Option<String>,
reason: Option<String>,
) -> OpsResult<TaskControlResult> {
let issue = issue.to_string();
block_on_task(async move { resume_task_async(&issue, model, reason).await })
}
pub(crate) async fn resume_task_async(
issue: &str,
model: Option<String>,
reason: Option<String>,
) -> OpsResult<TaskControlResult> {
let store = task_store().await?;
let mut session = store
.get_task_session_by_issue(issue)
.await
.map_err(|error| task_error(format!("failed to resolve task: {error}")))?
.ok_or_else(|| task_error(format!("no Task Session exists for {issue:?}")))?;
task_recovery_adoption(&store, &session).await?;
reconcile_task_pr(&store, &mut session).await?;
let prs = store
.task_prs(&session.id)
.await
.map_err(|error| task_error(format!("failed to read Task PRs: {error}")))?;
let latest = prs.last();
let active = prs.iter().find(|pr| pr.is_active());
if let Some(refusal) =
no_active_pr_resume_refusal(&session.launch.issue.identifier, active, latest)
{
return Err(task_error(refusal));
}
refuse_dirty_between_prs(&store, &session).await?;
reconcile_process_liveness(&store, &mut session).await?;
let issue_id = session.launch.issue.identifier.clone();
let observation = session.observation.clone();
let session_id = session.id.to_string();
let run = super::child::resume_session(
&store,
super::child::ChildSession::Task(Box::new(session)),
model,
reason,
)
.await?;
Ok(TaskControlResult {
issue_id,
session_id,
receipt: super::child::WorkControlReceipt::Resume { run },
observation,
})
}
pub fn task_abandon(issue: &str, reason: String) -> OpsResult<TaskControlResult> {
let reason = reason.trim();
if reason.is_empty() {
return Err(task_error("`lf task abandon --reason` cannot be empty"));
}
let reason = reason.to_string();
block_on_task(async move {
let store = task_store().await?;
let mut session = store
.get_task_session_by_issue(issue)
.await
.map_err(|error| task_error(format!("failed to resolve task: {error}")))?
.ok_or_else(|| task_error(format!("no Task Session exists for {issue:?}")))?;
reconcile_task_pr(&store, &mut session).await?;
reconcile_process_liveness(&store, &mut session).await?;
let work = store
.work_for_child(&ChildRef::Task(session.id.clone()))
.await
.map_err(|error| task_error(error.to_string()))?;
let basis = store
.current_epoch(&work)
.await
.map_err(|error| task_error(error.to_string()))?
.current_basis;
let receipt = store
.abandon(&work, &reason, &basis)
.await
.map_err(|error| task_error(error.to_string()))?;
if !session.status.is_process_active() {
session.set_status(
TaskSessionStatus::Abandoned,
format!("Task explicitly abandoned: {}", reason.trim()),
);
store
.update_task_session(&session)
.await
.map_err(|error| task_error(error.to_string()))?;
}
Ok(TaskControlResult {
issue_id: session.launch.issue.identifier.clone(),
session_id: session.id.to_string(),
receipt: super::child::WorkControlReceipt::Abandon { receipt },
observation: session.observation,
})
})
}
pub fn task_wait(
issue: &str,
until: TaskWaitUntil,
timeout: Option<Duration>,
) -> OpsResult<TaskSession> {
let started = Instant::now();
loop {
let session = task_status(issue)?;
let reached = match until {
TaskWaitUntil::Open => {
session.status.is_terminal()
|| active_pr(&session).is_ok_and(|pr| pr.phase() == PrPhase::Open)
}
TaskWaitUntil::Terminal => session.status.is_terminal(),
};
if reached || timeout.is_some_and(|limit| started.elapsed() >= limit) {
return Ok(session);
}
std::thread::sleep(Duration::from_secs(1));
}
}
pub fn task_attach(issue: &str) -> OpsResult<()> {
let session = task_status(issue)?;
if !session.status.is_process_active() {
return Err(task_error(format!(
"task {} is {}; resume it before attaching",
session.launch.issue.identifier,
session.status.as_str()
)));
}
let tmux_name = session
.latest_process
.as_ref()
.map(|process| process.tmux_name.as_str())
.filter(|name| !name.is_empty())
.ok_or_else(|| task_error("task has no attachable process; resume it first"))?;
let status = std::process::Command::new("tmux")
.args(["attach-session", "-t", tmux_name])
.status()
.map_err(|error| task_error(format!("failed to attach to task: {error}")))?;
if !status.success() {
return Err(task_error(format!("tmux attach failed for {tmux_name}")));
}
Ok(())
}
#[cfg(test)]
mod tests {
use std::ffi::OsString;
use std::os::unix::fs::PermissionsExt;
use std::os::unix::process::CommandExt;
use std::path::Path;
use std::process::Command;
use std::sync::Arc;
use std::time::Duration;
use super::{
_defer_task_interactions, _recover_abandoned_task, cached_github_observation,
changes_snapshot, count_recovery_attempts, decide_open_pr_status, derive_workspace_slug,
diff_snapshot, ensure_working_pr, ensure_working_pr_with_authority, file_snapshot,
next_pr_slug, parse_pr_slug, parse_workspace_slug, project_context,
reconcile_process_liveness, reconcile_task_pr, recover_stalled_task_body,
recover_stranded_task_body, refuse_dirty_between_prs, refuse_if_canonical_ahead,
require_task_pr_range_nonempty_with_authority, resolve_task_flow, resolve_upstream_base,
resume_task_async, succession_workspace_slug, supervise_project_task_bodies,
task_recovery_adoption, task_snapshot, unpublished_work,
verify_task_pr_range_with_authority, CommittedFollowUp, RotateOptions,
TaskRecoveryAdoption, TaskWorkspace,
};
use crate::child_session::{
observe, BodyEvidence, BodyIntent, ChildBodyOutcome, ChildLeaseState,
ChildProcessGeneration, ChildRef, MAX_RECOVERY_ATTEMPTS,
};
use crate::engine::git::is_ancestor;
use crate::id::WaveId;
use crate::pm::{PmKr, PmProject};
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::actions::TaskAction;
use crate::task::{
AfterMerge, CiCheck, CiObservation, CiState, GithubObservation, GithubObservationResult,
GithubPr, Observation, PmWritebackState, PrPhase, PrPublication, TaskEventKind, TaskPr,
TaskPrId, TaskSession, TaskSessionId, TaskSessionStatus,
};
use crate::wave::Wave;
use loopflow_test_support::TestRepo;
use time::OffsetDateTime;
fn git(repo: &Path, args: &[&str]) -> String {
let output = Command::new("git")
.current_dir(repo)
.args(args)
.output()
.expect("run git");
assert!(
output.status.success(),
"git {} failed: {}",
args.join(" "),
String::from_utf8_lossy(&output.stderr)
);
String::from_utf8_lossy(&output.stdout).trim().to_string()
}
struct TaskLaunchEnv {
previous: Vec<(&'static str, Option<OsString>)>,
_bin: tempfile::TempDir,
}
impl TaskLaunchEnv {
fn install(home: &Path) -> Self {
let bin = tempfile::tempdir().expect("fake tmux bin");
let tmux = bin.path().join("tmux");
std::fs::write(
&tmux,
"#!/bin/sh\nif [ \"$1\" = \"list-sessions\" ]; then printf 'lf-live-idle-control\\n'; fi\nexit 0\n",
)
.expect("write fake tmux");
let mut permissions = std::fs::metadata(&tmux)
.expect("stat fake tmux")
.permissions();
permissions.set_mode(0o755);
std::fs::set_permissions(&tmux, permissions).expect("make fake tmux executable");
let keys = [
"PATH",
"LF_BIN",
"LF_HOME",
"LF_DB_PATH",
"LF_CONTROL_BIN",
"LF_CONTROL_HOME",
"LF_CONTROL_DB_PATH",
crate::provider_account::lease::ACCOUNT_LEASE_ENV,
];
let previous = keys
.into_iter()
.map(|key| (key, std::env::var_os(key)))
.collect::<Vec<_>>();
let path = std::env::var_os("PATH").map_or_else(
|| bin.path().as_os_str().to_os_string(),
|path| {
let mut paths = vec![bin.path().to_path_buf()];
paths.extend(std::env::split_paths(&path));
std::env::join_paths(paths).expect("join fake tmux PATH")
},
);
let db = home.join("loopflow.db");
std::env::set_var("PATH", path);
std::env::set_var("LF_BIN", "/usr/bin/true");
std::env::set_var("LF_HOME", home);
std::env::set_var("LF_DB_PATH", &db);
std::env::set_var("LF_CONTROL_BIN", "/usr/bin/true");
std::env::set_var("LF_CONTROL_HOME", home);
std::env::set_var("LF_CONTROL_DB_PATH", db);
std::env::remove_var(crate::provider_account::lease::ACCOUNT_LEASE_ENV);
Self {
previous,
_bin: bin,
}
}
}
impl Drop for TaskLaunchEnv {
fn drop(&mut self) {
for (key, value) in self.previous.drain(..) {
match value {
Some(value) => std::env::set_var(key, value),
None => std::env::remove_var(key),
}
}
}
}
struct StoreEnvGuard {
_lock: std::sync::MutexGuard<'static, ()>,
previous: Vec<(&'static str, Option<OsString>)>,
}
impl StoreEnvGuard {
fn new(home: &Path) -> Self {
let lock = crate::journal::test_env_lock();
let names = [
"LF_HOME",
"LF_DB_PATH",
crate::store::CONTROL_HOME_ENV,
crate::store::CONTROL_DB_PATH_ENV,
];
let previous = names
.into_iter()
.map(|name| (name, std::env::var_os(name)))
.collect();
std::env::set_var("LF_HOME", home);
std::env::set_var("LF_DB_PATH", home.join("loopflow.db"));
std::env::remove_var(crate::store::CONTROL_HOME_ENV);
std::env::remove_var(crate::store::CONTROL_DB_PATH_ENV);
Self {
_lock: lock,
previous,
}
}
}
impl Drop for StoreEnvGuard {
fn drop(&mut self) {
for (name, value) in &self.previous {
match value {
Some(value) => std::env::set_var(name, value),
None => std::env::remove_var(name),
}
}
}
}
#[test]
fn task_flow_selection_accepts_skill_flows_and_rejects_ops() {
let repo = tempfile::tempdir().unwrap();
assert_eq!(
resolve_task_flow(repo.path(), Some("code")).unwrap(),
"code"
);
let error = resolve_task_flow(repo.path(), Some("deploy")).unwrap_err();
assert!(error
.to_string()
.contains("durable Task flows currently require skills"));
}
#[tokio::test]
async fn launch_task_process_ignores_control_bin_and_resolves_current_home() {
let home = tempfile::tempdir().unwrap();
let store: SharedStore = Arc::new(
open_store(&StorageConfig::sqlite(home.path().join("loopflow.db")))
.await
.unwrap(),
);
let now = OffsetDateTime::now_utc();
let mut session = TaskSession {
id: TaskSessionId::new(),
launch: TaskLaunchReceipt {
issue: LinearIssueSnapshot {
id: LinearIssueId::new("issue-no-pin").unwrap(),
identifier: "INF-NO-PIN".to_string(),
title: "Resolve through current lf".to_string(),
description: "Never read the pinned binary".to_string(),
},
project: LinearProjectSnapshot {
id: LinearProjectId::new("project-no-pin").unwrap(),
slug: "no-pin".to_string(),
name: "No pin".to_string(),
prompt_context: String::new(),
},
pm_snapshot_synced_at: now.unix_timestamp(),
},
pm_writeback: PmWritebackState::Current,
wave_id: WaveId::new(),
project_session_id: ProjectSessionId::new(),
status: TaskSessionStatus::Waiting,
status_reason: "ready".to_string(),
status_at: now,
worktree: home.path().join("worktree"),
workspace_slug: "task-no-pin".to_string(),
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: "codex".to_string(),
provider: "codex".to_string(),
provider_session_id: None,
latest_process: None,
abandon_intent: None,
created_at: now,
updated_at: now,
observation: crate::task::Observation::NotRequired,
};
let previous_control_bin = std::env::var_os("LF_CONTROL_BIN");
let previous_lf_bin = std::env::var_os("LF_BIN");
std::env::set_var("LF_CONTROL_BIN", "/bin/sh");
std::env::set_var("LF_BIN", "/loopflow-test/does-not-exist/lf");
let result = super::launch_task_process(&store, &mut session, None).await;
match previous_lf_bin {
Some(value) => std::env::set_var("LF_BIN", value),
None => std::env::remove_var("LF_BIN"),
}
match previous_control_bin {
Some(value) => std::env::set_var("LF_CONTROL_BIN", value),
None => std::env::remove_var("LF_CONTROL_BIN"),
}
let error = result.expect_err("launch must fail when the current Home lf is missing");
assert!(
error
.to_string()
.contains("cannot resolve current lf binary"),
"launch must resolve the current Home lf, not the LF_CONTROL_BIN pin: {error}"
);
assert!(session.latest_process.is_none());
}
fn changed_workspace() -> (tempfile::TempDir, String, TaskSessionId) {
let repo = tempfile::tempdir().expect("create temp repo");
git(repo.path(), &["init", "-b", "main"]);
git(repo.path(), &["config", "user.name", "Loopflow Test"]);
git(
repo.path(),
&["config", "user.email", "loopflow@example.com"],
);
std::fs::write(repo.path().join("tracked.txt"), "base\n").expect("write base");
git(repo.path(), &["add", "tracked.txt"]);
git(repo.path(), &["commit", "-m", "base"]);
let base = git(repo.path(), &["rev-parse", "HEAD"]);
std::fs::write(repo.path().join("committed.txt"), "committed\n")
.expect("write committed file");
git(repo.path(), &["add", "committed.txt"]);
git(repo.path(), &["commit", "-m", "task commit"]);
std::fs::write(repo.path().join("tracked.txt"), "staged\n").expect("write staged");
git(repo.path(), &["add", "tracked.txt"]);
std::fs::write(repo.path().join("tracked.txt"), "unstaged\n").expect("write unstaged");
std::fs::write(repo.path().join("untracked.txt"), "untracked\n").expect("write untracked");
(repo, base, TaskSessionId::new())
}
async fn rotation_task(
repo: &TestRepo,
branch: &str,
base_commit: &str,
) -> (tempfile::TempDir, SharedStore, TaskSession, TaskPr) {
rotation_task_with_lease(repo, branch, base_commit, None).await
}
async fn rotation_task_with_lease(
repo: &TestRepo,
branch: &str,
base_commit: &str,
lease: Option<(
crate::child_session::ChildLeaseState,
TaskSessionStatus,
Option<ChildBodyOutcome>,
)>,
) -> (tempfile::TempDir, SharedStore, TaskSession, TaskPr) {
let home = tempfile::tempdir().expect("task home");
let store = Arc::new(
open_store(&StorageConfig::sqlite(home.path().join("loopflow.db")))
.await
.expect("open store"),
);
let now = OffsetDateTime::now_utc();
let lease_seed = lease.map(|(state, status, outcome)| {
(
status,
ChildProcessGeneration {
generation: 1,
pid: None,
process_group_id: None,
tmux_name: format!("dead-lease-{}", WaveId::new()),
agent: "codex".to_string(),
provider: "codex".to_string(),
provider_session_id: None,
started_at: now - time::Duration::hours(1),
state,
outcome,
provenance: None,
},
)
});
let wave = Wave::new(
WaveId::new(),
"task-pr-rotation".to_string(),
repo.path().display().to_string(),
);
let project = ProjectSession {
id: ProjectSessionId::new(),
launch: ProjectLaunchReceipt {
project: LinearProjectSnapshot {
id: LinearProjectId::new(format!("project-{}", WaveId::new()))
.expect("project id"),
slug: "task-pr-rotation".to_string(),
name: "Task PR rotation".to_string(),
prompt_context: "Keep one stable worktree.".to_string(),
},
pm_snapshot_synced_at: now.unix_timestamp(),
},
wave_id: wave.id().clone(),
status: ProjectSessionStatus::Running,
status_reason: "test project is running".to_string(),
status_at: now,
iteration: 1,
observation_cursor: 0,
last_state_fingerprint: None,
agent: "codex".to_string(),
provider: "codex".to_string(),
provider_session_id: Some("task-pr-rotation".to_string()),
latest_process: Some(ChildProcessGeneration {
generation: 1,
pid: None,
process_group_id: None,
tmux_name: "task-pr-rotation".to_string(),
agent: "codex".to_string(),
provider: "codex".to_string(),
provider_session_id: Some("task-pr-rotation".to_string()),
started_at: now,
state: crate::child_session::ChildLeaseState::Active,
outcome: None,
provenance: None,
}),
abandon_intent: None,
created_at: now,
updated_at: now,
};
let session = TaskSession {
id: TaskSessionId::new(),
launch: TaskLaunchReceipt {
issue: LinearIssueSnapshot {
id: LinearIssueId::new(format!("issue-{}", WaveId::new())).expect("issue id"),
identifier: "INF-ROTATE".to_string(),
title: "Rotate Task PRs".to_string(),
description: "Keep the worktree stable.".to_string(),
},
project: project.launch.project.clone(),
pm_snapshot_synced_at: now.unix_timestamp(),
},
pm_writeback: PmWritebackState::Current,
wave_id: wave.id().clone(),
project_session_id: project.id.clone(),
status: lease_seed
.as_ref()
.map(|(status, _)| *status)
.unwrap_or(TaskSessionStatus::Waiting),
status_reason: if lease_seed.is_some() {
"recovered from a vanished body".to_string()
} else {
"first PR settled".to_string()
},
status_at: now,
worktree: repo.path().to_path_buf(),
workspace_slug: "task-pr-proof".to_string(),
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: "codex".to_string(),
provider: "codex".to_string(),
provider_session_id: None,
latest_process: lease_seed.map(|(_, process)| process),
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: branch.to_string(),
base_commit: base_commit.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_wave(&wave).await.expect("create wave");
store
.create_project_session(&project)
.await
.expect("create project");
store
.create_task_session(&session, &pr)
.await
.expect("create Task");
(home, store, session, pr)
}
async fn parked_gate_task(
repo: &TestRepo,
branch: &str,
ci_state: Option<CiState>,
phase_cursor: u32,
) -> (
tempfile::TempDir,
SharedStore,
ProjectSession,
TaskSession,
TaskPr,
) {
let base = repo.head_sha();
repo.create_branch(branch);
let (home, store, session, mut pr) =
gate_task_fixture(repo, branch, &base, false, phase_cursor).await;
let project = store
.get_project_session(&session.project_session_id)
.await
.expect("read Project")
.expect("Project exists");
let now = OffsetDateTime::now_utc();
let head = repo.head_sha();
pr.publication = Some(PrPublication {
requested_at: now,
after_merge: AfterMerge::Review,
next_slug: None,
github: Some(GithubPr {
number: 1054,
url: "https://github.com/loopflowstudio/loopflow/pull/1054".to_string(),
head_sha: Some(head.clone()),
}),
});
pr.github_observation = Some(GithubObservation {
checked_at: now,
result: GithubObservationResult::Fresh,
});
pr.ci_observation = ci_state.map(|state| CiObservation {
head_sha: head,
state,
failing_checks: if state == CiState::Failing {
vec![CiCheck {
name: "rust-test".to_string(),
url: Some("https://ci.example/rust-test".to_string()),
}]
} else {
Vec::new()
},
observed_at: now,
});
pr.updated_at = now;
store.update_task_pr(&pr).await.expect("publish fixture PR");
if let Some(incident) = super::current_ci_incident(&pr) {
store
.observe_ci_incident(&incident)
.await
.expect("record failing CI incident");
}
(home, store, project, session, pr)
}
async fn settle_pr(store: &SharedStore, mut pr: TaskPr, merge: &str, next_slug: Option<&str>) {
pr.publication = Some(PrPublication {
requested_at: OffsetDateTime::now_utc(),
after_merge: AfterMerge::Review,
next_slug: next_slug.map(str::to_string),
github: Some(GithubPr {
number: 900,
url: "https://example.com/pr/900".to_string(),
head_sha: None,
}),
});
pr.merge_commit = Some(merge.to_string());
pr.updated_at = OffsetDateTime::now_utc();
store.settle_task_pr(&pr, None).await.expect("settle PR");
}
async fn settle_completing_pr(store: &SharedStore, mut pr: TaskPr, merge: &str) {
pr.publication = Some(PrPublication {
requested_at: OffsetDateTime::now_utc(),
after_merge: AfterMerge::CompleteTask,
next_slug: None,
github: Some(GithubPr {
number: 1037,
url: "https://example.com/pr/1037".to_string(),
head_sha: None,
}),
});
pr.merge_commit = Some(merge.to_string());
pr.updated_at = OffsetDateTime::now_utc();
store.settle_task_pr(&pr, None).await.expect("settle PR");
}
#[tokio::test]
async fn merged_task_resume_refuses_before_reserving_a_generation() {
let repo = TestRepo::new();
let base = repo.head_sha();
let branch = "jack/no-doomed-resume";
repo.create_branch(branch);
let (home, store, session, pr) = rotation_task(&repo, branch, &base).await;
settle_completing_pr(&store, pr, "merge-1037").await;
let _env = StoreEnvGuard::new(home.path());
let before = store
.get_task_session(&session.id)
.await
.expect("read before")
.expect("Task exists")
.latest_process;
let error = resume_task_async(
&session.launch.issue.identifier,
None,
Some("status recommended Resume".to_string()),
)
.await
.expect_err("merged Task has no active PR to resume");
let after = store
.get_task_session(&session.id)
.await
.expect("read after")
.expect("Task exists")
.latest_process;
assert_eq!(
error.to_string(),
"Task INF-ROTATE has no active PR to resume; pull request #1037 merged"
);
assert_eq!(after, before, "a refusal must not reserve a generation");
}
async fn stacked_rotation_task(
repo: &TestRepo,
parent_branch: &str,
child_branch: &str,
parent_base: &str,
) -> (
tempfile::TempDir,
SharedStore,
TaskSession,
TaskPr,
TaskSession,
TaskPr,
) {
repo.create_branch(parent_branch);
repo.create_file("parent.txt", "parent work\n");
repo.stage_all();
repo.commit("parent commit");
repo.push_new_branch(parent_branch);
let parent_tip = repo.head_sha();
let (home, store, mut parent_session, mut parent_pr) =
rotation_task(repo, parent_branch, parent_base).await;
parent_pr.publication = Some(PrPublication {
requested_at: OffsetDateTime::now_utc(),
after_merge: AfterMerge::Review,
next_slug: None,
github: Some(GithubPr {
number: 900,
url: "https://example.com/pr/900".to_string(),
head_sha: None,
}),
});
parent_pr.updated_at = OffsetDateTime::now_utc();
store
.update_task_pr(&parent_pr)
.await
.expect("publish parent PR");
parent_session.worktree = std::path::PathBuf::from("/dummy/parent-worktree");
store
.update_task_session(&parent_session)
.await
.expect("reparent parent worktree");
repo.checkout("main");
git(repo.path(), &["branch", child_branch, &parent_tip]);
repo.checkout(child_branch);
let now = OffsetDateTime::now_utc();
let child_session = TaskSession {
id: TaskSessionId::new(),
launch: TaskLaunchReceipt {
issue: LinearIssueSnapshot {
id: LinearIssueId::new(format!("issue-{}", WaveId::new())).expect("issue id"),
identifier: "INF-STACK".to_string(),
title: "Stacked child".to_string(),
description: "Stacked on the parent.".to_string(),
},
project: parent_session.launch.project.clone(),
pm_snapshot_synced_at: now.unix_timestamp(),
},
pm_writeback: PmWritebackState::Current,
wave_id: parent_session.wave_id.clone(),
project_session_id: parent_session.project_session_id.clone(),
status: TaskSessionStatus::Waiting,
status_reason: "stacked child".to_string(),
status_at: now,
worktree: repo.path().to_path_buf(),
workspace_slug: "stacked-child".to_string(),
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: "codex".to_string(),
provider: "codex".to_string(),
provider_session_id: None,
latest_process: None,
abandon_intent: None,
created_at: now,
updated_at: now,
observation: crate::task::Observation::NotRequired,
};
let child_pr = TaskPr {
id: TaskPrId::new(),
task_session_id: child_session.id.clone(),
sequence: 1,
slug: "stacked-child".to_string(),
branch: child_branch.to_string(),
base_commit: parent_tip,
parent_pr_id: Some(parent_pr.id.clone()),
publication: None,
merge_commit: None,
abandoned_at: None,
ci_observation: None,
github_observation: None,
linear_attachment_id: None,
linear_comment_id: None,
linear_link_error: None,
created_at: now,
updated_at: now,
};
store
.create_task_session(&child_session, &child_pr)
.await
.expect("create child Task");
(
home,
store,
parent_session,
parent_pr,
child_session,
child_pr,
)
}
#[tokio::test]
async fn nonempty_refuses_when_head_is_the_recorded_base() {
let repo = TestRepo::new();
let base = repo.head_sha();
let branch = "jack/empty-head";
repo.create_branch(branch);
let (_home, store, session, _pr) = rotation_task(&repo, branch, &base).await;
let err =
require_task_pr_range_nonempty_with_authority(&store, &session, None, repo.path())
.await
.expect_err("HEAD == base must refuse as empty");
assert!(
err.to_string().contains("empty"),
"expected empty-range refusal, got: {err}"
);
}
#[tokio::test]
async fn nonempty_refuses_a_range_with_no_tree_change() {
let repo = TestRepo::new();
let base = repo.head_sha();
let branch = "jack/no-tree-change";
repo.create_branch(branch);
repo.create_file("ephemeral.txt", "gone\n");
repo.stage_all();
repo.commit("add ephemeral");
git(repo.path(), &["rm", "ephemeral.txt"]);
repo.commit("remove ephemeral");
let (_home, store, session, _pr) = rotation_task(&repo, branch, &base).await;
let err =
require_task_pr_range_nonempty_with_authority(&store, &session, None, repo.path())
.await
.expect_err("zero net tree change must refuse as empty");
assert!(
err.to_string().contains("empty"),
"expected empty-range refusal for zero tree change, got: {err}"
);
}
#[tokio::test]
async fn nonempty_refuses_even_when_the_pr_already_has_a_github_number() {
let repo = TestRepo::new();
let base = repo.head_sha();
let branch = "jack/empty-existing";
repo.create_branch(branch);
let (_home, store, session, mut pr) = rotation_task(&repo, branch, &base).await;
pr.publication = Some(PrPublication {
requested_at: OffsetDateTime::now_utc(),
after_merge: AfterMerge::Review,
next_slug: None,
github: Some(GithubPr {
number: 925,
url: "https://example.com/pr/925".to_string(),
head_sha: None,
}),
});
pr.updated_at = OffsetDateTime::now_utc();
store.update_task_pr(&pr).await.expect("set github number");
let err =
require_task_pr_range_nonempty_with_authority(&store, &session, None, repo.path())
.await
.expect_err("an existing PR with an empty range must refuse");
assert!(
err.to_string().contains("empty"),
"expected empty-range refusal despite github number, got: {err}"
);
}
#[tokio::test]
async fn nonempty_passes_for_a_real_range() {
let repo = TestRepo::new();
let base = repo.head_sha();
let branch = "jack/real-range";
repo.create_branch(branch);
repo.create_file("task.txt", "real work\n");
repo.stage_all();
repo.commit("real task commit");
let (_home, store, session, _pr) = rotation_task(&repo, branch, &base).await;
require_task_pr_range_nonempty_with_authority(&store, &session, None, repo.path())
.await
.expect("a real range must pass the non-empty check");
}
#[tokio::test]
async fn stacked_child_measures_from_live_parent_tip() {
let repo = TestRepo::new();
let origin_tip = repo.head_sha();
let (_home, store, _, _, child_session, _) =
stacked_rotation_task(&repo, "jack/stack-parent", "jack/stack-child", &origin_tip)
.await;
repo.create_file("child.txt", "child work\n");
repo.stage_all();
repo.commit("child commit");
require_task_pr_range_nonempty_with_authority(&store, &child_session, None, repo.path())
.await
.expect("stacked child with own work passes against the parent tip");
}
#[tokio::test]
async fn stacked_child_refuses_when_empty_against_live_parent() {
let repo = TestRepo::new();
let origin_tip = repo.head_sha();
let (_home, store, _, _, child_session, _) = stacked_rotation_task(
&repo,
"jack/stack-parent-empty",
"jack/stack-child-empty",
&origin_tip,
)
.await;
let err = require_task_pr_range_nonempty_with_authority(
&store,
&child_session,
None,
repo.path(),
)
.await
.expect_err("empty stacked child must refuse against the parent tip");
assert!(
err.to_string().contains("empty"),
"expected empty-range refusal for stacked child, got: {err}"
);
}
#[tokio::test]
async fn stacked_child_measures_from_origin_after_parent_collapsed() {
let repo = TestRepo::new();
let origin_tip = repo.head_sha();
let (_home, store, _, mut parent_pr, child_session, mut child_pr) = stacked_rotation_task(
&repo,
"jack/collapse-parent",
"jack/collapse-child",
&origin_tip,
)
.await;
repo.create_file("child.txt", "child work\n");
repo.stage_all();
repo.commit("child commit");
repo.checkout("main");
repo.create_file("parent.txt", "parent work\n");
repo.stage_all();
repo.commit("merge parent into main");
repo.push();
let main_tip = repo.head_sha();
repo.checkout("jack/collapse-child");
git(
repo.path(),
&[
"rebase",
"--onto",
"origin/main",
&parent_pr.base_commit,
"jack/collapse-child",
],
);
parent_pr.merge_commit = Some("merge-sha".to_string());
parent_pr.updated_at = OffsetDateTime::now_utc();
store
.update_task_pr(&parent_pr)
.await
.expect("mark parent merged");
child_pr.base_commit = main_tip;
child_pr.updated_at = OffsetDateTime::now_utc();
store
.heal_task_pr_base(&child_pr)
.await
.expect("heal child base to main");
require_task_pr_range_nonempty_with_authority(&store, &child_session, None, repo.path())
.await
.expect("collapsed child with own work passes against origin/main");
}
#[tokio::test]
async fn stacked_child_refuses_when_empty_after_parent_collapsed() {
let repo = TestRepo::new();
let origin_tip = repo.head_sha();
let (_home, store, _, mut parent_pr, child_session, mut child_pr) = stacked_rotation_task(
&repo,
"jack/collapse-parent-empty",
"jack/collapse-child-empty",
&origin_tip,
)
.await;
repo.checkout("main");
repo.create_file("parent.txt", "parent work\n");
repo.stage_all();
repo.commit("merge parent into main");
repo.push();
let main_tip = repo.head_sha();
repo.checkout("jack/collapse-child-empty");
git(
repo.path(),
&[
"rebase",
"--onto",
"origin/main",
&parent_pr.base_commit,
"jack/collapse-child-empty",
],
);
parent_pr.merge_commit = Some("merge-sha".to_string());
parent_pr.updated_at = OffsetDateTime::now_utc();
store
.update_task_pr(&parent_pr)
.await
.expect("mark parent merged");
child_pr.base_commit = main_tip;
child_pr.updated_at = OffsetDateTime::now_utc();
store
.heal_task_pr_base(&child_pr)
.await
.expect("heal child base to main");
let err = require_task_pr_range_nonempty_with_authority(
&store,
&child_session,
None,
repo.path(),
)
.await
.expect_err("empty collapsed child must refuse against origin/main");
assert!(
err.to_string().contains("empty"),
"expected empty-range refusal for collapsed child, got: {err}"
);
}
#[test]
fn task_context_captures_project_definition_and_kr_state() {
let project = PmProject {
id: "project-1".to_string(),
slug: "pr".to_string(),
name: "PR".to_string(),
summary: "Ship reliably".to_string(),
definition: "Every task has one durable session.".to_string(),
krs: vec![
PmKr {
text: "Review resumes the same session".to_string(),
holds: true,
},
PmKr {
text: "Merge wakes the Wave".to_string(),
holds: false,
},
],
initiative_ids: vec!["initiative-1".to_string()],
team_ids: None,
};
assert_eq!(
project_context(&project),
"Definition:\nEvery task has one durable session.\n\nKRs:\n- [x] Review resumes the same session\n- [ ] Merge wakes the Wave"
);
}
#[tokio::test]
async fn idle_task_can_defer_its_remaining_interactive_steps() {
let repo = TestRepo::new();
let base = repo.head_sha();
let (_home, _store, mut session, _pr) =
rotation_task(&repo, "jack/task-headless", &base).await;
assert!(_defer_task_interactions(&mut session).unwrap());
assert!(session.lifecycle.all_interactions_deferred());
assert!(!_defer_task_interactions(&mut session).unwrap());
session.lifecycle = crate::task::TaskLifecyclePlan::standard("task");
session.status = TaskSessionStatus::Completed;
assert!(_defer_task_interactions(&mut session)
.unwrap_err()
.to_string()
.contains("terminal Tasks cannot change interaction policy"));
}
async fn dead_lease_task(
repo: &TestRepo,
branch: &str,
base: &str,
status: TaskSessionStatus,
lease_state: crate::child_session::ChildLeaseState,
outcome: Option<ChildBodyOutcome>,
identity: (Option<u32>, Option<u32>),
) -> (tempfile::TempDir, SharedStore, TaskSession) {
let (pid, process_group_id) = identity;
let (home, store, base_session, _pr) = rotation_task(repo, branch, base).await;
let now = OffsetDateTime::now_utc();
let mut session = base_session.clone();
session.id = TaskSessionId::new();
session.workspace_slug = "dead-lease-proof".to_string();
session.worktree = repo.path().join(format!("dead-{}", session.id));
session.launch.issue.id =
LinearIssueId::new(format!("issue-{}", session.id)).expect("issue id");
session.launch.issue.identifier = format!("INF-DEAD-{}", session.id);
session.set_status(status, "recovered from a vanished body");
session.latest_process = Some(ChildProcessGeneration {
generation: 1,
pid,
process_group_id,
tmux_name: format!("dead-lease-{}", session.id),
agent: session.agent.clone(),
provider: session.provider.clone(),
provider_session_id: None,
started_at: now - time::Duration::hours(1),
state: lease_state,
outcome,
provenance: None,
});
let pr = TaskPr {
id: TaskPrId::new(),
task_session_id: session.id.clone(),
sequence: 1,
slug: session.workspace_slug.clone(),
branch: format!("{branch}-dead"),
base_commit: base.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
.expect("create dead-lease Task");
(home, store, session)
}
fn pin_unlaunchable_body() -> Option<std::ffi::OsString> {
let previous = std::env::var_os("LF_BIN");
std::env::set_var("LF_BIN", "/loopflow-test/does-not-exist/lf");
previous
}
fn restore_lf_bin(previous: Option<std::ffi::OsString>) {
match previous {
Some(value) => std::env::set_var("LF_BIN", value),
None => std::env::remove_var("LF_BIN"),
}
}
async fn stranded_task(
repo: &TestRepo,
branch: &str,
) -> (tempfile::TempDir, SharedStore, TaskSession) {
let base = repo.head_sha();
let (home, store, session, _pr) = rotation_task_with_lease(
repo,
branch,
&base,
Some((ChildLeaseState::Active, TaskSessionStatus::Running, None)),
)
.await;
git(repo.path(), &["checkout", "-b", branch]);
(home, store, session)
}
#[allow(clippy::await_holding_lock)] #[tokio::test]
async fn a_killed_body_under_a_live_task_is_redispatched_with_no_human_action() {
let _env_lock = crate::journal::test_env_lock();
let repo = TestRepo::new();
let (_home, store, mut session) = stranded_task(&repo, "jack/w2-267-recovers").await;
let pinned = pin_unlaunchable_body();
let outcome = recover_stranded_task_body(&store, &mut session).await;
restore_lf_bin(pinned);
assert!(outcome.is_err(), "the pinned-missing binary must fail");
let events = store.recent_task_events(&session.id, 64).await.unwrap();
assert_eq!(
count_recovery_attempts(&events),
1,
"recovery must durably record the attempt it made"
);
assert!(
events.iter().any(|event| matches!(
event.kind,
TaskEventKind::BodyRecoveryAttempted { attempt: 1, .. }
)),
"recovery must be observable, not silent"
);
}
#[allow(clippy::await_holding_lock)] #[tokio::test]
async fn an_unlaunchable_strand_exhausts_instead_of_minting_dead_generations() {
let _env_lock = crate::journal::test_env_lock();
let repo = TestRepo::new();
let (_home, store, mut session) = stranded_task(&repo, "jack/w2-267-exhausts").await;
for expected in 1..=MAX_RECOVERY_ATTEMPTS {
let mut current = store.get_task_session(&session.id).await.unwrap().unwrap();
let pinned = pin_unlaunchable_body();
let _ = recover_stranded_task_body(&store, &mut current).await;
restore_lf_bin(pinned);
let events = store.recent_task_events(&session.id, 64).await.unwrap();
assert_eq!(
count_recovery_attempts(&events),
expected,
"attempt {expected} must be recorded"
);
}
let mut spent = store.get_task_session(&session.id).await.unwrap().unwrap();
let pinned = pin_unlaunchable_body();
let recovered = recover_stranded_task_body(&store, &mut spent).await;
restore_lf_bin(pinned);
let recovered = recovered.expect("a spent budget surfaces rather than erroring");
assert!(!recovered, "an exhausted strand must not redispatch");
let events = store.recent_task_events(&session.id, 64).await.unwrap();
assert_eq!(
count_recovery_attempts(&events),
MAX_RECOVERY_ATTEMPTS,
"exhaustion must not append a further attempt"
);
let persisted = store.get_task_session(&session.id).await.unwrap().unwrap();
assert_eq!(persisted.status, TaskSessionStatus::Failed);
assert!(
persisted.status_reason.contains("did not survive"),
"an exhausted strand must say why: {}",
persisted.status_reason
);
session = persisted;
let _ = &session;
}
#[tokio::test]
async fn a_completed_task_whose_body_was_reaped_triggers_nothing() {
let repo = TestRepo::new();
let base = repo.head_sha();
let (_home, store, mut session, _pr) = rotation_task_with_lease(
&repo,
"jack/w2-267-completed",
&base,
Some((
ChildLeaseState::Finished,
TaskSessionStatus::Completed,
Some(ChildBodyOutcome::Lost {
reason: "task process disappeared before recording a terminal outcome"
.to_string(),
}),
)),
)
.await;
git(repo.path(), &["checkout", "-b", "jack/w2-267-completed"]);
let before = store.get_task_session(&session.id).await.unwrap().unwrap();
assert!(
matches!(
before
.latest_process
.as_ref()
.and_then(|p| p.outcome.as_ref()),
Some(ChildBodyOutcome::Lost { .. })
),
"the trap only exists when a completed Task carries a lost body"
);
let recovered = recover_stranded_task_body(&store, &mut session)
.await
.expect("classify a completed Task");
assert!(!recovered, "a completed Task must never be recovered");
let after = store.get_task_session(&session.id).await.unwrap().unwrap();
assert_eq!(after.status, TaskSessionStatus::Completed);
assert_eq!(
after.latest_process.map(|process| process.generation),
before.latest_process.map(|process| process.generation),
"a completed Task's body must not be replaced"
);
let events = store.recent_task_events(&session.id, 64).await.unwrap();
assert_eq!(
count_recovery_attempts(&events),
0,
"a completed Task must record no recovery attempt"
);
}
fn exited_process_group() -> u32 {
let mut child = std::process::Command::new("/bin/sh")
.arg("-c")
.arg("exit 0")
.process_group(0)
.spawn()
.expect("spawn a short-lived process group");
let group = child.id();
child.wait().expect("reap the short-lived group");
group
}
#[tokio::test]
async fn a_revoked_lease_over_a_dead_group_releases_and_only_then_reserves() {
let repo = TestRepo::new();
let base = repo.head_sha();
let mut stranger = std::process::Command::new("/bin/sh")
.arg("-c")
.arg("sleep 60")
.spawn()
.expect("spawn an unrelated live process");
let recycled_pid = stranger.id();
let (_home, store, session) = dead_lease_task(
&repo,
"jack/eng-4-released",
&base,
TaskSessionStatus::Waiting,
crate::child_session::ChildLeaseState::Revoked,
Some(ChildBodyOutcome::Superseded {
reason: "body generation 1 stalled; recovering the same Task Session".to_string(),
}),
(Some(recycled_pid), Some(exited_process_group())),
)
.await;
let target = ChildRef::Task(session.id.clone());
let revoked = session
.latest_process
.as_ref()
.expect("seeded generation")
.clone();
let _tmux = crate::engine::process::FakeTmux::no_session();
let mut blocked = session.clone();
let generation = blocked.begin_generation("lf-task-eng4-blocked".to_string());
assert_eq!(generation, 2);
assert!(store
.reserve_task_process(&blocked, TaskSessionStatus::Waiting)
.await
.expect("reserve against a revoked lease")
.is_none());
let finished =
crate::ops::child::release_dead_revoked_child_body(&store, &target, &revoked)
.await
.expect("probe and release a lease over a dead group")
.expect("a provably absent body releases its lease");
stranger.kill().expect("kill the unrelated process");
stranger.wait().expect("reap the unrelated process");
assert_eq!(
finished.state,
crate::child_session::ChildLeaseState::Finished
);
assert_eq!(finished.generation, revoked.generation);
let persisted = store.get_task_session(&session.id).await.unwrap().unwrap();
let persisted_process = persisted.latest_process.as_ref().expect("generation");
assert_eq!(
persisted_process.state,
crate::child_session::ChildLeaseState::Finished
);
assert!(matches!(
persisted_process.outcome,
Some(ChildBodyOutcome::Superseded { .. })
));
assert_eq!(persisted.status, TaskSessionStatus::Waiting);
let mut launch = persisted.clone();
launch.begin_generation("lf-task-eng4-successor".to_string());
assert!(store
.reserve_task_process(&launch, TaskSessionStatus::Waiting)
.await
.expect("reserve after the release")
.is_some());
}
#[tokio::test]
async fn only_the_matching_revoked_generation_can_finish() {
let repo = TestRepo::new();
let base = repo.head_sha();
let (_home, store, session) = dead_lease_task(
&repo,
"jack/eng-4-cas",
&base,
TaskSessionStatus::Waiting,
crate::child_session::ChildLeaseState::Revoked,
Some(ChildBodyOutcome::Superseded {
reason: "stalled".to_string(),
}),
(None, Some(exited_process_group())),
)
.await;
let target = ChildRef::Task(session.id.clone());
let revoked = session.latest_process.as_ref().expect("generation").clone();
let _tmux = crate::engine::process::FakeTmux::no_session();
let mismatched = ChildProcessGeneration {
generation: revoked.generation + 1,
..revoked.clone()
};
assert!(
crate::ops::child::release_dead_revoked_child_body(&store, &target, &mismatched)
.await
.is_err(),
"a release must not settle a generation the store is not holding"
);
crate::ops::child::release_dead_revoked_child_body(&store, &target, &revoked)
.await
.expect("release the matching generation")
.expect("a provably absent body releases");
let settled = store.get_task_session(&session.id).await.unwrap().unwrap();
let finished = settled.latest_process.as_ref().expect("generation").clone();
assert!(
crate::ops::child::release_dead_revoked_child_body(&store, &target, &finished)
.await
.expect("a finished lease is a no-op, not an error")
.is_none()
);
}
#[tokio::test]
async fn a_live_group_keeps_its_lease_and_the_refusal_names_it() {
let repo = TestRepo::new();
let base = repo.head_sha();
let mut child = std::process::Command::new("/bin/sh")
.arg("-c")
.arg("sleep 60")
.process_group(0)
.spawn()
.expect("spawn a live process group");
let live_group = child.id();
let (_home, store, mut session) = dead_lease_task(
&repo,
"jack/eng-4-live",
&base,
TaskSessionStatus::Waiting,
crate::child_session::ChildLeaseState::Revoked,
Some(ChildBodyOutcome::Superseded {
reason: "stalled".to_string(),
}),
(Some(live_group), Some(live_group)),
)
.await;
let revoked = session.latest_process.as_ref().expect("generation").clone();
let _tmux = crate::engine::process::FakeTmux::no_session();
let released = crate::ops::child::release_dead_revoked_child_body(
&store,
&ChildRef::Task(session.id.clone()),
&revoked,
)
.await
.expect("probe a live group");
super::release_dead_revoked_task_lease(&store, &mut session)
.await
.expect("a live body is not an error, it is a reason to hold");
child.kill().expect("kill the live group");
child.wait().expect("reap the live group");
assert!(released.is_none(), "a live body must keep its lease");
assert_eq!(
session.latest_process.as_ref().map(|process| process.state),
Some(crate::child_session::ChildLeaseState::Revoked)
);
let persisted = store.get_task_session(&session.id).await.unwrap().unwrap();
assert_eq!(
persisted.latest_process.map(|process| process.state),
Some(crate::child_session::ChildLeaseState::Revoked),
"the lease must still be doing its job in the store"
);
}
#[tokio::test]
async fn an_unprovable_identity_does_not_release_the_lease() {
let repo = TestRepo::new();
let base = repo.head_sha();
let unaddressable = u32::try_from(i32::MAX).expect("i32::MAX fits u32") + 1;
let (_home, store, session) = dead_lease_task(
&repo,
"jack/eng-4-unprovable",
&base,
TaskSessionStatus::Waiting,
crate::child_session::ChildLeaseState::Revoked,
Some(ChildBodyOutcome::Superseded {
reason: "stalled".to_string(),
}),
(None, Some(unaddressable)),
)
.await;
let revoked = session.latest_process.as_ref().expect("generation").clone();
let _tmux = crate::engine::process::FakeTmux::no_session();
let released = crate::ops::child::release_dead_revoked_child_body(
&store,
&ChildRef::Task(session.id.clone()),
&revoked,
)
.await
.expect("an unprovable probe is not an error");
assert!(released.is_none());
let persisted = store.get_task_session(&session.id).await.unwrap().unwrap();
assert_eq!(
persisted.latest_process.map(|process| process.state),
Some(crate::child_session::ChildLeaseState::Revoked)
);
}
#[tokio::test]
#[allow(clippy::await_holding_lock)] async fn resume_revokes_a_dead_legacy_lease_on_a_waiting_task() {
let _env_lock = crate::journal::test_env_lock();
let repo = TestRepo::new();
let base = repo.head_sha();
let (_home, store, mut session) = dead_lease_task(
&repo,
"jack/w2-135",
&base,
TaskSessionStatus::Waiting,
crate::child_session::ChildLeaseState::Legacy,
None,
(None, None),
)
.await;
reconcile_process_liveness(&store, &mut session)
.await
.expect("reconcile a waiting task with a dead legacy lease");
assert_eq!(
session.latest_process.as_ref().map(|process| process.state),
Some(crate::child_session::ChildLeaseState::Finished)
);
assert_eq!(session.status, TaskSessionStatus::Waiting);
let persisted = store.get_task_session(&session.id).await.unwrap().unwrap();
assert_eq!(
persisted.latest_process.map(|process| process.state),
Some(crate::child_session::ChildLeaseState::Finished)
);
assert_eq!(persisted.status, TaskSessionStatus::Waiting);
}
#[tokio::test]
#[allow(clippy::await_holding_lock)] async fn resume_revokes_a_dead_active_lease_on_a_failed_task() {
let _env_lock = crate::journal::test_env_lock();
let repo = TestRepo::new();
let base = repo.head_sha();
let (_home, store, mut session) = dead_lease_task(
&repo,
"jack/w2-122",
&base,
TaskSessionStatus::Failed,
crate::child_session::ChildLeaseState::Active,
None,
(None, None),
)
.await;
reconcile_process_liveness(&store, &mut session)
.await
.expect("reconcile a failed task with a dead active lease");
assert_eq!(
session.latest_process.as_ref().map(|process| process.state),
Some(crate::child_session::ChildLeaseState::Finished)
);
assert_eq!(session.status, TaskSessionStatus::Failed);
let persisted = store.get_task_session(&session.id).await.unwrap().unwrap();
assert_eq!(
persisted.latest_process.map(|process| process.state),
Some(crate::child_session::ChildLeaseState::Finished)
);
assert_eq!(persisted.status, TaskSessionStatus::Failed);
}
#[tokio::test]
async fn progress_wins_the_race_against_stall_recovery() {
let repo = TestRepo::new();
let base = repo.head_sha();
let (_home, store, mut session, _pr) =
rotation_task(&repo, "jack/progress-race", &base).await;
session.begin_generation(format!("progress-race-{}", session.id));
let lease = store
.reserve_task_process(&session, TaskSessionStatus::Waiting)
.await
.unwrap()
.expect("reserve body");
session
.latest_process
.as_mut()
.expect("reserved process")
.state = ChildLeaseState::Active;
session.set_status(TaskSessionStatus::Running, "provider is alive");
store.activate_task_process(&session, &lease).await.unwrap();
session = store.get_task_session(&session.id).await.unwrap().unwrap();
let observed_event_id = store
.latest_task_event(&session.id)
.await
.unwrap()
.map(|event| event.id);
store
.append_task_event(
&session.id,
&TaskEventKind::Progress {
summary: "body advanced before revocation".to_string(),
},
)
.await
.unwrap();
let revoked = store
.revoke_task_process_if_unchanged(
&session.id,
1,
session.status_at,
observed_event_id,
&ChildBodyOutcome::Superseded {
reason: "stale observation".to_string(),
},
)
.await
.unwrap();
assert!(revoked.is_none());
let persisted = store.get_task_session(&session.id).await.unwrap().unwrap();
assert_eq!(
persisted.latest_process.map(|process| process.state),
Some(ChildLeaseState::Active),
);
assert_eq!(persisted.status, TaskSessionStatus::Running);
}
#[test]
fn readable_task_names_are_semantic_and_bounded() {
assert_eq!(
derive_workspace_slug("Release scoped migrations across every target")
.unwrap()
.as_str(),
"release-scoped-migrations-across-every"
);
assert_eq!(
derive_workspace_slug("Investigate").unwrap().as_str(),
"investigate-task"
);
assert!(parse_workspace_slug("one").is_err());
assert!(parse_workspace_slug("release_scoped-migrations").is_err());
assert_eq!(
parse_pr_slug("released-upgrade-proof").unwrap(),
"released-upgrade-proof"
);
assert!(parse_pr_slug("released/upgrade").is_err());
}
#[test]
fn succession_slug_is_distinct_capped_and_per_predecessor() {
let title = "Release scoped migrations across every target";
let base = derive_workspace_slug(title).unwrap();
assert_eq!(base.as_str(), "release-scoped-migrations-across-every");
let pred_a = TaskSessionId::new();
let pred_b = TaskSessionId::new();
let succ_a = succession_workspace_slug(title, &pred_a).unwrap();
let succ_b = succession_workspace_slug(title, &pred_b).unwrap();
assert_ne!(succ_a.as_str(), base.as_str());
assert_ne!(succ_a.as_str(), succ_b.as_str());
let words = succ_a.as_str().split('-').count();
assert!(
(2..=5).contains(&words),
"got {words} words in {}",
succ_a.as_str()
);
let capped =
succession_workspace_slug("one two three four five six seven", &pred_a).unwrap();
assert_eq!(
capped.as_str().split('-').count(),
5,
"four base words + one suffix word"
);
}
#[tokio::test]
async fn settled_pr_rotates_the_same_worktree_from_fetched_main() {
let repo = TestRepo::new();
let base = repo.head_sha();
let first_branch = "jack/task-pr-proof";
repo.create_branch(first_branch);
repo.create_file("first.txt", "first PR\n");
repo.stage_all();
repo.commit("first PR");
let (_home, store, mut session, mut first) =
rotation_task(&repo, first_branch, &base).await;
first.publication = Some(PrPublication {
requested_at: OffsetDateTime::now_utc(),
after_merge: AfterMerge::Review,
next_slug: Some("follow-up-proof".to_string()),
github: Some(GithubPr {
number: 911,
url: "https://example.com/pr/911".to_string(),
head_sha: None,
}),
});
first.merge_commit = Some("merge-911".to_string());
first.updated_at = OffsetDateTime::now_utc();
store
.settle_task_pr(&first, None)
.await
.expect("settle first PR");
let second = ensure_working_pr(&store, &mut session)
.await
.expect("rotate PR")
.expect("working PR");
assert_eq!(session.worktree, repo.path());
assert_eq!(second.sequence, 2);
assert_eq!(second.branch, "jack/task-pr-proof-follow-up-proof");
assert_eq!(second.base_commit, base);
assert_eq!(second.phase(), PrPhase::Working);
assert_eq!(
git(repo.path(), &["rev-parse", "--abbrev-ref", "HEAD"]),
second.branch
);
let prs = store.task_prs(&session.id).await.expect("read PR history");
assert_eq!(prs.iter().map(|pr| pr.sequence).collect::<Vec<_>>(), [1, 2]);
assert_eq!(prs[0].phase(), PrPhase::Merged);
assert_eq!(prs[1].id, second.id);
}
struct RemintFixture {
_home: tempfile::TempDir,
repo: TestRepo,
store: SharedStore,
session: TaskSession,
settled: TaskPr,
b1: String,
b2: String,
}
impl RemintFixture {
const SUCCESSOR: &str = "jack/task-pr-proof-2";
async fn new(carry: Option<&str>) -> Self {
let repo = TestRepo::new();
let b1 = repo.head_sha();
let settled_branch = "jack/task-pr-proof";
repo.create_branch(settled_branch);
repo.create_file("first.txt", "first PR\n");
repo.stage_all();
repo.commit("first PR");
let (home, store, session, mut settled) =
rotation_task(&repo, settled_branch, &b1).await;
settled.publication = Some(PrPublication {
requested_at: OffsetDateTime::now_utc(),
after_merge: AfterMerge::Review,
next_slug: None,
github: Some(GithubPr {
number: 1050,
url: "https://example.com/pr/1050".to_string(),
head_sha: None,
}),
});
settled.merge_commit = Some("merge-1050".to_string());
settled.updated_at = OffsetDateTime::now_utc();
git(repo.path(), &["checkout", "-b", Self::SUCCESSOR, &b1]);
if let Some(file) = carry {
repo.create_file(file, "work a partial rotation already carried\n");
repo.stage_all();
repo.commit("carried follow-up");
}
repo.checkout("main");
repo.create_file("elsewhere.txt", "another PR merged\n");
repo.stage_all();
repo.commit("main advances");
repo.push();
let b2 = repo.head_sha();
assert_ne!(b1, b2, "origin/main must move or the defect cannot show");
repo.checkout(Self::SUCCESSOR);
Self {
_home: home,
repo,
store,
session,
settled,
b1,
b2,
}
}
async fn settle(&self, next: Option<&TaskPr>) {
self.store
.settle_task_pr(&self.settled, next)
.await
.expect("settle merged PR");
}
async fn rotate(&mut self) -> TaskPr {
ensure_working_pr(&self.store, &mut self.session)
.await
.expect("rotate")
.expect("working PR")
}
fn cut(&self, pr: &TaskPr) -> CommittedFollowUp {
unpublished_work(&self.session.worktree, pr).expect("classify")
}
fn successor_row(&self, base: &str) -> TaskPr {
let mut pr = self.settled.clone();
pr.id = TaskPrId::new();
pr.sequence = 2;
pr.slug = "2".to_string();
pr.branch = Self::SUCCESSOR.to_string();
pr.base_commit = base.to_string();
pr.publication = None;
pr.merge_commit = None;
pr.created_at = OffsetDateTime::now_utc();
pr.updated_at = OffsetDateTime::now_utc();
pr
}
fn sibling_commit(&self) -> String {
git(
self.repo.path(),
&["checkout", "-b", "jack/sibling", &self.b1],
);
self.repo
.create_file("sibling.txt", "another branch's work\n");
self.repo.stage_all();
self.repo.commit("sibling work");
let tip = self.repo.head_sha();
self.repo.checkout(Self::SUCCESSOR);
tip
}
}
#[tokio::test]
async fn a_re_mint_after_main_advances_records_the_base_its_branch_forks_from() {
let mut fx = RemintFixture::new(None).await;
fx.settle(None).await;
let next = fx.rotate().await;
assert_eq!(next.sequence, 2);
assert_eq!(next.branch, RemintFixture::SUCCESSOR);
assert_eq!(
next.base_commit, fx.b1,
"the fork point of the reused branch, not the upstream tip"
);
assert_ne!(next.base_commit, fx.b2);
assert!(
is_ancestor(fx.repo.path(), &next.base_commit, &next.branch).expect("ancestry"),
"the recorded base must be an ancestor of its branch"
);
assert!(matches!(fx.cut(&next), CommittedFollowUp::ProvenEmpty));
}
#[tokio::test]
async fn a_re_mint_holding_committed_work_still_reads_range() {
let mut fx = RemintFixture::new(Some("carried.txt")).await;
fx.settle(None).await;
let carried_tip = fx.repo.head_sha();
let next = fx.rotate().await;
assert_eq!(
next.base_commit, fx.b1,
"the fork point, not the carried tip"
);
assert_ne!(next.base_commit, carried_tip);
assert!(matches!(fx.cut(&next), CommittedFollowUp::Range { .. }));
}
#[tokio::test]
async fn an_incoherent_recorded_base_is_healed_when_the_active_pr_is_adopted() {
let mut fx = RemintFixture::new(None).await;
let incoherent = fx.successor_row(&fx.b2);
fx.settle(Some(&incoherent)).await;
assert!(
matches!(fx.cut(&incoherent), CommittedFollowUp::Unprovable { .. }),
"fixture must start wedged, or it proves nothing"
);
let adopted = fx.rotate().await;
assert_eq!(
adopted.id, incoherent.id,
"the same row, healed — not a new one"
);
assert_eq!(adopted.base_commit, fx.b1);
assert!(matches!(fx.cut(&adopted), CommittedFollowUp::ProvenEmpty));
let stored = fx
.store
.task_prs(&fx.session.id)
.await
.expect("read PR history")
.into_iter()
.find(|pr| pr.sequence == 2)
.expect("successor row");
assert_eq!(stored.base_commit, fx.b1, "the heal must reach the store");
}
#[tokio::test]
async fn a_recorded_base_off_the_upstream_line_is_refused_not_healed() {
let mut fx = RemintFixture::new(None).await;
let sibling = fx.sibling_commit();
let foreign = fx.successor_row(&sibling);
fx.settle(Some(&foreign)).await;
assert!(
is_ancestor(fx.repo.path(), &fx.b1, &sibling).expect("ancestry"),
"M <= B must hold, or an M<=B-only heal would refuse this anyway and the test proves nothing"
);
assert!(
!is_ancestor(fx.repo.path(), &sibling, "origin/main").expect("ancestry"),
"B must be off the upstream line"
);
assert!(matches!(
fx.cut(&foreign),
CommittedFollowUp::Unprovable { .. }
));
let adopted = fx.rotate().await;
assert_eq!(
adopted.base_commit, sibling,
"a base no past mint could have written must not be healed"
);
assert!(
matches!(fx.cut(&adopted), CommittedFollowUp::Unprovable { .. }),
"contamination stays Unprovable — the gate keeps refusing"
);
let stored = fx
.store
.task_prs(&fx.session.id)
.await
.expect("read PR history")
.into_iter()
.find(|pr| pr.sequence == 2)
.expect("successor row");
assert_eq!(stored.base_commit, sibling, "the store row is untouched");
}
#[tokio::test]
async fn rotate_forward_carries_uncommitted_follow_up_edits() {
let repo = TestRepo::new();
let base = repo.head_sha();
let first_branch = "jack/task-pr-proof";
repo.create_branch(first_branch);
repo.create_file("first.txt", "first PR\n");
repo.stage_all();
repo.commit("first PR");
let (_home, store, mut session, mut first) =
rotation_task(&repo, first_branch, &base).await;
first.publication = Some(PrPublication {
requested_at: OffsetDateTime::now_utc(),
after_merge: AfterMerge::Review,
next_slug: None,
github: Some(GithubPr {
number: 907,
url: "https://example.com/pr/907".to_string(),
head_sha: None,
}),
});
first.merge_commit = Some("merge-907".to_string());
first.updated_at = OffsetDateTime::now_utc();
store
.settle_task_pr(&first, None)
.await
.expect("settle first PR");
repo.create_file("follow-up.txt", "second PR work\n");
assert!(ensure_working_pr(&store, &mut session.clone())
.await
.is_err());
let second = ensure_working_pr_with_authority(
&store,
&mut session,
None,
RotateOptions {
carry_dirty: true,
slug_override: Some("keep-going".to_string()),
},
)
.await
.expect("rotate forward")
.expect("working PR");
assert_eq!(second.sequence, 2);
assert_eq!(second.branch, "jack/task-pr-proof-keep-going");
assert_eq!(second.base_commit, base);
assert_eq!(second.phase(), PrPhase::Working);
assert_eq!(
git(repo.path(), &["rev-parse", "--abbrev-ref", "HEAD"]),
second.branch
);
assert_eq!(
std::fs::read_to_string(repo.path().join("follow-up.txt")).expect("follow-up survives"),
"second PR work\n"
);
let prs = store.task_prs(&session.id).await.expect("read PR history");
assert_eq!(prs.iter().map(|pr| pr.sequence).collect::<Vec<_>>(), [1, 2]);
assert_eq!(prs[0].phase(), PrPhase::Merged);
}
#[tokio::test]
async fn reconcile_degrades_and_preserves_cache_when_the_github_read_fails() {
let repo = TestRepo::new();
let base = repo.head_sha();
let branch = "jack/task-pr-proof";
repo.create_branch(branch);
repo.create_file("first.txt", "first PR\n");
repo.stage_all();
repo.commit("first PR");
let (_home, store, mut session, mut pr) = rotation_task(&repo, branch, &base).await;
pr.publication = Some(PrPublication {
requested_at: OffsetDateTime::now_utc(),
after_merge: AfterMerge::Review,
next_slug: None,
github: Some(GithubPr {
number: 914,
url: "https://example.com/pr/914".to_string(),
head_sha: None,
}),
});
pr.updated_at = OffsetDateTime::now_utc();
store.update_task_pr(&pr).await.expect("publish PR");
let observed = reconcile_task_pr(&store, &mut session)
.await
.expect("reconcile does not error on a failed GitHub read")
.expect("the cached PR is preserved");
assert_eq!(observed.phase(), PrPhase::Open);
assert_eq!(observed.merge_commit, None);
match &session.observation {
Observation::Degraded { reason, .. } => assert!(!reason.is_empty()),
other => panic!("a failed GitHub read must degrade freshness, got {other:?}"),
}
}
#[tokio::test]
async fn reconcile_skips_the_remote_read_for_an_unpublished_working_pr() {
let repo = TestRepo::new();
let base = repo.head_sha();
let branch = "jack/task-pr-proof";
repo.create_branch(branch);
let (_home, store, mut session, pr) = rotation_task(&repo, branch, &base).await;
assert!(pr.github().is_none(), "fixture PR is unpublished");
let observed = reconcile_task_pr(&store, &mut session)
.await
.expect("reconcile succeeds")
.expect("working PR preserved");
assert_eq!(observed.phase(), PrPhase::Working);
assert_eq!(session.observation, Observation::NotRequired);
}
#[test]
fn github_observation_cache_expires_fresh_reads_before_degraded_circuits() {
let now = OffsetDateTime::now_utc();
let mut pr = TaskPr {
id: TaskPrId::new(),
task_session_id: TaskSessionId::new(),
sequence: 1,
slug: "cache-proof".to_string(),
branch: "jack/cache-proof".to_string(),
base_commit: "base".to_string(),
parent_pr_id: None,
publication: None,
merge_commit: None,
abandoned_at: None,
ci_observation: None,
github_observation: Some(GithubObservation {
checked_at: now - time::Duration::seconds(59),
result: GithubObservationResult::Fresh,
}),
linear_attachment_id: None,
linear_comment_id: None,
linear_link_error: None,
created_at: now,
updated_at: now - time::Duration::hours(1),
};
assert!(matches!(
cached_github_observation(&pr, now),
Some(Observation::Cached { .. })
));
pr.github_observation.as_mut().unwrap().checked_at = now - time::Duration::seconds(60);
assert_eq!(cached_github_observation(&pr, now), None);
pr.github_observation = Some(GithubObservation {
checked_at: now - time::Duration::minutes(4),
result: GithubObservationResult::Degraded {
reason: "rate limit exhausted".to_string(),
},
});
assert!(matches!(
cached_github_observation(&pr, now),
Some(Observation::Degraded { .. })
));
pr.github_observation.as_mut().unwrap().checked_at = now - time::Duration::minutes(5);
assert_eq!(cached_github_observation(&pr, now), None);
}
#[tokio::test]
async fn rotate_carries_committed_follow_up_and_dirty_edits_after_an_out_of_band_merge() {
let repo = TestRepo::new();
let base = repo.head_sha();
let settled_branch = "jack/task-pr-proof";
repo.create_branch(settled_branch);
repo.create_file("merged.txt", "merged work\n");
repo.stage_all();
repo.commit("merged work");
let merged_tip = repo.head_sha();
git(repo.path(), &["push", "origin", "jack/task-pr-proof:main"]);
repo.create_file("follow1.txt", "follow-up one\n");
repo.stage_all();
repo.commit("follow-up one");
repo.create_file("follow2.txt", "follow-up two\n");
repo.stage_all();
repo.commit("follow-up two");
repo.create_file("wip.txt", "uncommitted work\n");
let (_home, store, mut session, mut settled) =
rotation_task(&repo, settled_branch, &base).await;
settled.publication = Some(PrPublication {
requested_at: OffsetDateTime::now_utc(),
after_merge: AfterMerge::Review,
next_slug: Some("keep-going".to_string()),
github: Some(GithubPr {
number: 907,
url: "https://example.com/pr/907".to_string(),
head_sha: Some(merged_tip.clone()),
}),
});
settled.merge_commit = Some("merge-907".to_string());
settled.updated_at = OffsetDateTime::now_utc();
store
.settle_task_pr(&settled, None)
.await
.expect("settle merged PR");
let next = ensure_working_pr_with_authority(
&store,
&mut session,
None,
RotateOptions {
carry_dirty: true,
slug_override: Some("keep-going".to_string()),
},
)
.await
.expect("rotate forward")
.expect("working PR");
assert_eq!(next.sequence, 2);
assert_eq!(
git(repo.path(), &["rev-parse", "--abbrev-ref", "HEAD"]),
next.branch
);
assert_eq!(
std::fs::read_to_string(repo.path().join("follow1.txt")).expect("follow1 carried"),
"follow-up one\n"
);
assert_eq!(
std::fs::read_to_string(repo.path().join("follow2.txt")).expect("follow2 carried"),
"follow-up two\n"
);
assert_eq!(
std::fs::read_to_string(repo.path().join("wip.txt")).expect("dirty edit carried"),
"uncommitted work\n"
);
let beyond = git(
repo.path(),
&["log", "origin/main..HEAD", "--oneline", "--format=%s"],
);
let subjects: Vec<&str> = beyond.lines().collect();
assert_eq!(subjects, vec!["follow-up two", "follow-up one"]);
assert!(repo.path().join("merged.txt").exists());
}
#[tokio::test]
async fn a_completing_pr_rotates_only_to_carry_committed_follow_up() {
let repo = TestRepo::new();
let base = repo.head_sha();
let settled_branch = "jack/task-pr-proof";
repo.create_branch(settled_branch);
repo.create_file("merged.txt", "merged work\n");
repo.stage_all();
repo.commit("merged work");
let merged_tip = repo.head_sha();
git(repo.path(), &["push", "origin", "jack/task-pr-proof:main"]);
let (_home, store, mut session, mut settled) =
rotation_task(&repo, settled_branch, &base).await;
settled.publication = Some(PrPublication {
requested_at: OffsetDateTime::now_utc(),
after_merge: AfterMerge::CompleteTask,
next_slug: None,
github: Some(GithubPr {
number: 1042,
url: "https://example.com/pr/1042".to_string(),
head_sha: Some(merged_tip.clone()),
}),
});
settled.merge_commit = Some("merge-1042".to_string());
settled.updated_at = OffsetDateTime::now_utc();
store
.settle_task_pr(&settled, None)
.await
.expect("settle merged completing PR");
let none =
ensure_working_pr_with_authority(&store, &mut session, None, RotateOptions::runner())
.await
.expect("completing PR with no follow-up");
assert!(
none.is_none(),
"a completing PR with no follow-up must not mint a successor"
);
repo.create_file("follow-up.txt", "acknowledged follow-up\n");
repo.stage_all();
repo.commit("follow-up the directive asked for");
let next =
ensure_working_pr_with_authority(&store, &mut session, None, RotateOptions::runner())
.await
.expect("rotate to carry follow-up")
.expect("successor PR");
assert_eq!(next.sequence, 2);
assert_eq!(
git(repo.path(), &["rev-parse", "--abbrev-ref", "HEAD"]),
next.branch
);
assert_eq!(
std::fs::read_to_string(repo.path().join("follow-up.txt")).expect("follow-up carried"),
"acknowledged follow-up\n"
);
let beyond = git(
repo.path(),
&["log", "origin/main..HEAD", "--oneline", "--format=%s"],
);
assert_eq!(
beyond.lines().collect::<Vec<_>>(),
vec!["follow-up the directive asked for"]
);
}
#[tokio::test]
async fn failed_committed_carry_restores_the_settled_branch_for_retry() {
let repo = TestRepo::new();
let base = repo.head_sha();
let settled_branch = "jack/task-pr-proof";
repo.create_branch(settled_branch);
repo.create_file("shared.txt", "merged\n");
repo.stage_all();
repo.commit("merged work");
let merged_tip = repo.head_sha();
git(repo.path(), &["checkout", "-b", "main-update", &merged_tip]);
repo.create_file("shared.txt", "main moved on\n");
repo.stage_all();
repo.commit("main update");
git(repo.path(), &["push", "origin", "main-update:main"]);
git(repo.path(), &["checkout", settled_branch]);
repo.create_file("shared.txt", "follow-up edit\n");
repo.stage_all();
repo.commit("post-merge follow-up");
repo.create_file("wip.txt", "dirty follow-up\n");
let (_home, store, mut session, mut settled) =
rotation_task(&repo, settled_branch, &base).await;
settled.publication = Some(PrPublication {
requested_at: OffsetDateTime::now_utc(),
after_merge: AfterMerge::Review,
next_slug: None,
github: Some(GithubPr {
number: 918,
url: "https://example.com/pr/918".to_string(),
head_sha: Some(merged_tip),
}),
});
settled.merge_commit = Some("merge-918".to_string());
settled.updated_at = OffsetDateTime::now_utc();
store
.settle_task_pr(&settled, None)
.await
.expect("settle merged PR");
let error = ensure_working_pr_with_authority(
&store,
&mut session,
None,
RotateOptions {
carry_dirty: true,
slug_override: Some("retry".to_string()),
},
)
.await
.expect_err("conflicting carry must fail");
assert!(error.to_string().contains("restored"));
assert_eq!(
git(repo.path(), &["rev-parse", "--abbrev-ref", "HEAD"]),
settled_branch
);
assert_eq!(
std::fs::read_to_string(repo.path().join("shared.txt")).expect("commit restored"),
"follow-up edit\n"
);
assert_eq!(
std::fs::read_to_string(repo.path().join("wip.txt")).expect("dirty edit restored"),
"dirty follow-up\n"
);
assert_eq!(
git(repo.path(), &["branch", "--list", "*/task-pr-proof-retry"]),
""
);
assert_eq!(git(repo.path(), &["stash", "list"]), "");
}
#[test]
fn task_workspace_reports_committed_staged_unstaged_and_untracked_files() {
let (repo, base, session_id) = changed_workspace();
let workspace = TaskWorkspace {
issue_identifier: "INF-123",
session_id: &session_id,
worktree: repo.path(),
base_commit: &base,
};
let snapshot = changes_snapshot(workspace).expect("inspect task changes");
let committed = snapshot
.files
.iter()
.find(|file| file.path == "committed.txt")
.expect("committed file");
assert!(committed.committed);
let tracked = snapshot
.files
.iter()
.find(|file| file.path == "tracked.txt")
.expect("tracked file");
assert!(tracked.staged && tracked.unstaged);
let untracked = snapshot
.files
.iter()
.find(|file| file.path == "untracked.txt")
.expect("untracked file");
assert!(untracked.untracked);
let patch = diff_snapshot(workspace, None).expect("inspect task diff");
assert!(patch.patch.contains("committed.txt"));
assert!(patch.patch.contains("tracked.txt"));
assert!(patch.patch.contains("untracked.txt"));
let untracked_patch =
diff_snapshot(workspace, Some("./untracked.txt")).expect("inspect one file");
assert_eq!(untracked_patch.path.as_deref(), Some("untracked.txt"));
assert!(untracked_patch.patch.contains("+untracked"));
}
#[test]
fn task_file_reads_only_files_inside_the_task_worktree() {
let (repo, base, session_id) = changed_workspace();
let workspace = TaskWorkspace {
issue_identifier: "INF-123",
session_id: &session_id,
worktree: repo.path(),
base_commit: &base,
};
let file = file_snapshot(workspace, "./tracked.txt").expect("read task file");
assert_eq!(file.path, "tracked.txt");
assert_eq!(file.content.as_deref(), Some("unstaged\n"));
assert!(file_snapshot(workspace, "../outside.txt").is_err());
#[cfg(unix)]
{
let outside = tempfile::NamedTempFile::new().expect("outside file");
std::os::unix::fs::symlink(outside.path(), repo.path().join("outside-link"))
.expect("create outside symlink");
}
}
#[tokio::test]
async fn verify_refuses_a_base_carrying_a_foreign_commit() {
let repo = TestRepo::new();
repo.create_file("foreign.txt", "not this task's work\n");
repo.stage_all();
repo.commit("foreign canonical-main commit");
let contaminated_base = repo.head_sha();
let branch = "jack/contaminated";
repo.create_branch(branch);
repo.create_file("task.txt", "task work\n");
repo.stage_all();
repo.commit("task commit");
let (_home, store, session, _pr) = rotation_task(&repo, branch, &contaminated_base).await;
let err = verify_task_pr_range_with_authority(&store, &session, None, repo.path())
.await
.expect_err("contaminated base must refuse");
let message = err.to_string();
assert!(
message.contains("contaminated"),
"expected contamination refusal, got: {message}"
);
assert!(
message.contains("foreign canonical-main commit"),
"refusal must name the foreign commit, got: {message}"
);
assert!(
message.contains("rebase --onto"),
"refusal must print the recovery action, got: {message}"
);
}
#[tokio::test]
async fn verify_heals_a_stale_base_after_origin_advances() {
let repo = TestRepo::new();
let stale_base = repo.head_sha();
repo.create_file("upstream.txt", "landed upstream\n");
repo.stage_all();
repo.commit("upstream advance");
repo.push();
let advanced = repo.head_sha();
let branch = "jack/stale-base";
repo.create_branch(branch);
repo.create_file("task.txt", "task work\n");
repo.stage_all();
repo.commit("task commit");
let (_home, store, session, _pr) = rotation_task(&repo, branch, &stale_base).await;
verify_task_pr_range_with_authority(&store, &session, None, repo.path())
.await
.expect("stale-but-compatible base verifies");
let healed = store
.active_task_pr(&session.id)
.await
.expect("read active PR")
.expect("active PR exists");
assert_eq!(
healed.base_commit, advanced,
"the stale base should heal forward to the current fork point"
);
}
#[tokio::test]
async fn verify_falls_back_to_local_main_without_a_remote() {
let repo = TestRepo::new();
let base = repo.head_sha();
git(repo.path(), &["remote", "remove", "origin"]);
let branch = "jack/no-remote";
repo.create_branch(branch);
repo.create_file("task.txt", "task work\n");
repo.stage_all();
repo.commit("task commit");
let (_home, store, session, _pr) = rotation_task(&repo, branch, &base).await;
verify_task_pr_range_with_authority(&store, &session, None, repo.path())
.await
.expect("no-remote repo verifies against local main");
}
#[tokio::test]
async fn verify_passes_for_a_rotated_continuation_pr() {
let repo = TestRepo::new();
let base = repo.head_sha();
let first_branch = "jack/task-pr-proof";
repo.create_branch(first_branch);
repo.create_file("first.txt", "first PR\n");
repo.stage_all();
repo.commit("first PR");
let (_home, store, mut session, mut first) =
rotation_task(&repo, first_branch, &base).await;
first.publication = Some(PrPublication {
requested_at: OffsetDateTime::now_utc(),
after_merge: AfterMerge::Review,
next_slug: Some("follow-up".to_string()),
github: Some(GithubPr {
number: 938,
url: "https://example.com/pr/938".to_string(),
head_sha: None,
}),
});
first.merge_commit = Some("merge-938".to_string());
first.updated_at = OffsetDateTime::now_utc();
store
.settle_task_pr(&first, None)
.await
.expect("settle first PR");
let second = ensure_working_pr(&store, &mut session)
.await
.expect("rotate PR")
.expect("working PR");
assert_eq!(second.sequence, 2);
repo.create_file("second.txt", "second PR work\n");
repo.stage_all();
repo.commit("second PR commit");
verify_task_pr_range_with_authority(&store, &session, None, repo.path())
.await
.expect("rotated continuation PR verifies");
}
#[tokio::test]
async fn verify_refuses_divergent_ancestry_naming_both_sides() {
let repo = TestRepo::new();
let origin_tip = repo.head_sha();
repo.create_file("foreign.txt", "not this task's work\n");
repo.stage_all();
repo.commit("foreign canonical-main commit");
let contaminated_base = repo.head_sha();
let branch = "jack/divergent";
repo.create_branch(branch);
repo.create_file("task.txt", "task work\n");
repo.stage_all();
repo.commit("task commit");
repo.checkout("main");
git(repo.path(), &["reset", "--hard", &origin_tip]);
repo.create_file("upstream.txt", "landed upstream\n");
repo.stage_all();
repo.commit("upstream advance");
repo.push();
repo.checkout(branch);
git(
repo.path(),
&[
"rebase",
"--onto",
"origin/main",
&contaminated_base,
branch,
],
);
let (_home, store, session, _pr) = rotation_task(&repo, branch, &contaminated_base).await;
let err = verify_task_pr_range_with_authority(&store, &session, None, repo.path())
.await
.expect_err("divergent ancestry must refuse");
let message = err.to_string();
assert!(
message.contains("diverged"),
"expected divergence refusal, got: {message}"
);
assert!(
message.contains("foreign canonical-main commit"),
"refusal must name the base-side foreign commit, got: {message}"
);
assert!(
message.contains("foreign.txt"),
"refusal must name the base-side file, got: {message}"
);
assert!(
message.contains("upstream advance"),
"refusal must name the upstream-side commit, got: {message}"
);
assert!(
message.contains("upstream.txt"),
"refusal must name the upstream-side file, got: {message}"
);
assert!(
message.contains("rebase --onto"),
"refusal must print the recovery action, got: {message}"
);
}
#[tokio::test]
async fn verify_refuses_contaminated_range_without_a_remote() {
let repo = TestRepo::new();
let base = repo.head_sha();
git(repo.path(), &["remote", "remove", "origin"]);
repo.create_file("foreign.txt", "not this task's work\n");
repo.stage_all();
repo.commit("foreign local-main commit");
let contaminated_base = repo.head_sha();
let branch = "jack/no-remote-contaminated";
repo.create_branch(branch);
repo.create_file("task.txt", "task work\n");
repo.stage_all();
repo.commit("task commit");
repo.checkout("main");
git(repo.path(), &["reset", "--hard", &base]);
repo.checkout(branch);
let (_home, store, session, _pr) = rotation_task(&repo, branch, &contaminated_base).await;
let err = verify_task_pr_range_with_authority(&store, &session, None, repo.path())
.await
.expect_err("no-remote contaminated range must refuse");
let message = err.to_string();
assert!(
message.contains("contaminated"),
"expected contamination refusal, got: {message}"
);
assert!(
message.contains("foreign local-main commit"),
"refusal must name the foreign commit, got: {message}"
);
}
#[tokio::test]
async fn verify_refuses_contaminated_range_after_squash_merged_parent() {
let repo = TestRepo::new();
let first_branch = "jack/pr-one";
repo.create_branch(first_branch);
repo.create_file("pr1-a.txt", "a\n");
repo.stage_all();
repo.commit("PR1 first commit");
repo.create_file("pr1-b.txt", "b\n");
repo.stage_all();
repo.commit("PR1 second commit");
let pr1_tip = repo.head_sha();
repo.checkout("main");
git(repo.path(), &["merge", "--squash", first_branch]);
repo.stage_all();
repo.commit("squash-merge PR1");
repo.push();
let second_branch = "jack/pr-two";
git(repo.path(), &["branch", second_branch, &pr1_tip]);
repo.checkout(second_branch);
repo.create_file("pr2.txt", "PR2 work\n");
repo.stage_all();
repo.commit("PR2 commit");
let (_home, store, session, _pr) = rotation_task(&repo, second_branch, &pr1_tip).await;
let err = verify_task_pr_range_with_authority(&store, &session, None, repo.path())
.await
.expect_err("squash-merged parent contamination must refuse");
let message = err.to_string();
assert!(
message.contains("contaminated"),
"expected contamination refusal after squash-merge, got: {message}"
);
assert!(
message.contains("PR1 first commit"),
"refusal must name the first pre-squash commit, got: {message}"
);
assert!(
message.contains("PR1 second commit"),
"refusal must name the second pre-squash commit, got: {message}"
);
assert!(
message.contains("rebase --onto"),
"refusal must print the recovery action, got: {message}"
);
}
#[test]
fn placement_base_anchors_on_origin_when_a_remote_exists() {
let repo = TestRepo::new();
let origin = repo.head_sha();
repo.create_file("local.txt", "unpushed\n");
repo.stage_all();
repo.commit("unpushed local commit");
let (base_ref, base_commit) =
resolve_upstream_base(repo.path(), "main").expect("resolve base");
assert_eq!(base_ref, "origin/main");
assert_eq!(
base_commit, origin,
"the base must anchor on fetched origin, not the ahead-of-origin local tip"
);
}
#[test]
fn placement_base_falls_back_to_local_main_without_a_remote() {
let repo = TestRepo::new();
git(repo.path(), &["remote", "remove", "origin"]);
repo.create_file("local.txt", "local only\n");
repo.stage_all();
repo.commit("local commit");
let local_tip = repo.head_sha();
let (base_ref, base_commit) =
resolve_upstream_base(repo.path(), "main").expect("resolve base");
assert_eq!(base_ref, "refs/heads/main");
assert_eq!(base_commit, local_tip);
}
#[test]
fn placement_refuses_when_canonical_main_is_ahead_of_origin() {
let repo = TestRepo::new();
repo.create_file("ahead.txt", "unpushed canonical work\n");
repo.stage_all();
repo.commit("unpushed canonical commit");
let err = refuse_if_canonical_ahead(repo.path(), "main")
.expect_err("ahead-of-origin canonical main must refuse placement");
let message = err.to_string();
assert!(
message.contains("ahead of origin/main"),
"expected control-plane refusal, got: {message}"
);
assert!(
message.contains("unpushed canonical commit"),
"refusal must name the unpushed commit, got: {message}"
);
}
#[tokio::test]
async fn recover_refuses_a_non_abandoned_task() {
let repo = TestRepo::new();
let base = repo.head_sha();
let branch = "jack/task-pr-proof";
repo.create_branch(branch);
let (_home, store, session, _pr) = rotation_task(&repo, branch, &base).await;
let error = _recover_abandoned_task(&store, &session.launch.issue.identifier, None)
.await
.expect_err("a waiting Task resumes instead of recovering");
assert!(error.to_string().contains("lf task resume"), "{error}");
assert_eq!(
store
.get_task_session_by_issue(&session.launch.issue.identifier)
.await
.unwrap()
.unwrap()
.id,
session.id
);
}
#[tokio::test]
async fn recover_abandoned_task_adopts_existing_worktree_pr_and_direction() {
let repo = TestRepo::new();
let base = repo.head_sha();
let branch = "jack/task-pr-proof";
repo.create_branch(branch);
repo.create_file("recovery.txt", "work survives abandonment\n");
repo.stage_all();
repo.commit("task work");
let (_home, store, session, pr) = rotation_task(&repo, branch, &base).await;
let work = store
.work_for_child(&ChildRef::Task(session.id.clone()))
.await
.unwrap();
store
.append_steer(
&work,
crate::durable::Author::User,
"finish the existing PR",
None,
)
.await
.expect("record current direction");
let mut abandoned = store.get_task_session(&session.id).await.unwrap().unwrap();
abandoned.set_status(TaskSessionStatus::Abandoned, "operator stopped the attempt");
store.update_task_session(&abandoned).await.unwrap();
let successor = _recover_abandoned_task(
&store,
&abandoned.launch.issue.identifier,
Some("the work is still valid".to_string()),
)
.await
.expect("recover abandoned Task");
assert_ne!(successor.id, abandoned.id);
assert_eq!(successor.status, TaskSessionStatus::Waiting);
assert_eq!(successor.worktree, abandoned.worktree);
assert_eq!(successor.workspace_slug, abandoned.workspace_slug);
assert!(successor.latest_process.is_none());
assert_eq!(
successor.lifecycle_phase,
crate::task::TaskLifecyclePhase::Kickoff
);
assert!(successor.status_reason.contains("the work is still valid"));
assert!(store.task_prs(&abandoned.id).await.unwrap().is_empty());
let adopted_prs = store.task_prs(&successor.id).await.unwrap();
assert_eq!(adopted_prs.len(), 1);
assert_eq!(adopted_prs[0].id, pr.id);
let successor_work = store
.work_for_child(&ChildRef::Task(successor.id.clone()))
.await
.unwrap();
let carried = store.boundary_seed(&successor_work).await.unwrap();
assert!(carried.render().contains("finish the existing PR"));
let repeated = _recover_abandoned_task(
&store,
&successor.launch.issue.identifier,
Some("the work is still valid".to_string()),
)
.await
.expect("recovery retry converges");
assert_eq!(repeated.id, successor.id);
let mut completed = repeated;
completed.set_status(TaskSessionStatus::Completed, "recovered work landed");
store.update_task_session(&completed).await.unwrap();
let error = _recover_abandoned_task(
&store,
&completed.launch.issue.identifier,
Some("try again".to_string()),
)
.await
.expect_err("a completed recovered Task cannot be recovered again");
assert!(error.to_string().contains("is completed"), "{error}");
}
fn checkout_branch(repo: &TestRepo, branch: &str) {
let status = Command::new("git")
.current_dir(repo.path())
.args(["checkout", branch])
.status()
.expect("checkout");
assert!(status.success(), "checkout {branch} failed");
}
#[tokio::test]
async fn recovery_refuses_an_unrelated_branch_before_moving_ownership() {
let repo = TestRepo::new();
let base = repo.head_sha();
let first_branch = "jack/task-pr-proof";
repo.create_branch(first_branch);
repo.create_file("first.txt", "first PR\n");
repo.stage_all();
repo.commit("first PR");
let (_home, store, session, first) = rotation_task_with_lease(
&repo,
first_branch,
&base,
Some((
crate::child_session::ChildLeaseState::Active,
TaskSessionStatus::Waiting,
None,
)),
)
.await;
settle_pr(&store, first, "merge-unrelated", None).await;
repo.create_branch("jack/unrelated");
let prs_before = store.task_prs(&session.id).await.expect("read PRs");
let lease_before = store
.get_task_session(&session.id)
.await
.unwrap()
.unwrap()
.latest_process
.clone()
.expect("dead lease seeded");
let err = task_recovery_adoption(&store, &session)
.await
.expect_err("unrelated branch must refuse");
let message = err.to_string();
assert!(
message.contains("between-PR recovery expected settled branch"),
"expected unrelated-branch refusal, got: {message}"
);
assert!(
message.contains("jack/unrelated"),
"refusal must name the current branch, got: {message}"
);
assert!(
message.contains("refused before moving any ownership"),
"refusal must name the contract, got: {message}"
);
assert_eq!(
store.task_prs(&session.id).await.expect("reread PRs"),
prs_before,
"PR sequence untouched"
);
let after = store.get_task_session(&session.id).await.unwrap().unwrap();
assert_eq!(
after.latest_process,
Some(lease_before),
"dead lease untouched — the gate is what prevents the reap"
);
assert_eq!(after.status, TaskSessionStatus::Waiting, "status untouched");
assert_eq!(
git(repo.path(), &["rev-parse", "--abbrev-ref", "HEAD"]),
"jack/unrelated",
"worktree branch untouched"
);
}
#[tokio::test]
async fn supervised_restart_refuses_an_unrelated_branch() {
let repo = TestRepo::new();
let base = repo.head_sha();
let first_branch = "jack/task-pr-proof";
repo.create_branch(first_branch);
repo.create_file("first.txt", "first PR\n");
repo.stage_all();
repo.commit("first PR");
let (_home, store, session, first) = rotation_task_with_lease(
&repo,
first_branch,
&base,
Some((ChildLeaseState::Active, TaskSessionStatus::Running, None)),
)
.await;
settle_pr(&store, first, "merge-supervised", None).await;
repo.create_branch("jack/unrelated");
let session = store.get_task_session(&session.id).await.unwrap().unwrap();
let lease_before = session.latest_process.clone().expect("active lease seeded");
let observation = observe(
&BodyEvidence {
intent: BodyIntent::Active,
observable: true,
process_alive: true,
progress_age: Duration::from_secs(31 * 60),
step: Some("task_pursue".to_string()),
reason: "body is alive but stalled".to_string(),
},
Duration::from_secs(30 * 60),
);
let latest_event_id = store
.latest_task_event(&session.id)
.await
.unwrap()
.map(|event| event.id);
assert!(
!recover_stalled_task_body(&store, session.clone(), &observation, latest_event_id)
.await
.expect("an unsafe worktree declines recovery; it does not fail the supervisor"),
"an unrelated branch must not be restarted"
);
let after = store.get_task_session(&session.id).await.unwrap().unwrap();
assert_eq!(
after.latest_process,
Some(lease_before),
"lease untouched — the gate is what prevents the reap"
);
assert_eq!(after.status, TaskSessionStatus::Running, "status untouched");
}
#[tokio::test]
async fn recovery_refuses_a_dirty_worktree_between_prs() {
let repo = TestRepo::new();
let base = repo.head_sha();
let first_branch = "jack/task-pr-proof";
repo.create_branch(first_branch);
repo.create_file("first.txt", "first PR\n");
repo.stage_all();
repo.commit("first PR");
let (_home, store, session, first) =
rotation_task_with_lease(&repo, first_branch, &base, None).await;
settle_pr(&store, first, "merge-dirty", None).await;
repo.create_file("follow-up.txt", "uncommitted\n");
let err = task_recovery_adoption(&store, &session)
.await
.expect_err("dirty between-PR worktree must refuse");
let message = err.to_string();
assert!(
message.contains("uncommitted changes"),
"expected dirty refusal, got: {message}"
);
assert!(
message.contains("lf pr next"),
"refusal must name the recovery action, got: {message}"
);
assert!(
message.contains("refused before moving any ownership"),
"refusal must name the contract, got: {message}"
);
}
#[tokio::test]
async fn recovery_refuses_a_missing_branch() {
let repo = TestRepo::new();
let base = repo.head_sha();
let first_branch = "jack/task-pr-proof";
repo.create_branch(first_branch);
repo.create_file("first.txt", "first PR\n");
repo.stage_all();
repo.commit("first PR");
let (_home, store, session, first) =
rotation_task_with_lease(&repo, first_branch, &base, None).await;
settle_pr(&store, first, "merge-missing", None).await;
Command::new("git")
.current_dir(repo.path())
.args(["update-ref", "-d", &format!("refs/heads/{first_branch}")])
.status()
.expect("delete branch ref");
assert_eq!(
git(repo.path(), &["symbolic-ref", "--short", "HEAD"]),
first_branch,
"HEAD still names the deleted branch"
);
let err = task_recovery_adoption(&store, &session)
.await
.expect_err("missing branch must refuse");
let message = err.to_string();
assert!(
message.contains("no longer exists"),
"expected missing-branch refusal, got: {message}"
);
assert!(
message.contains(first_branch),
"refusal must name the missing branch, got: {message}"
);
}
#[tokio::test]
async fn recovery_adopts_an_active_pr_branch_and_allows_ongoing_work() {
let repo = TestRepo::new();
let base = repo.head_sha();
let branch = "jack/task-pr-proof";
repo.create_branch(branch);
repo.create_file("first.txt", "first PR\n");
repo.stage_all();
repo.commit("first PR");
let (_home, store, session, _pr) =
rotation_task_with_lease(&repo, branch, &base, None).await;
repo.create_file("wip.txt", "ongoing\n");
let adoption = task_recovery_adoption(&store, &session)
.await
.expect("active PR on its branch is adopted, dirty work allowed");
assert_eq!(
adoption,
TaskRecoveryAdoption::Active {
branch: branch.to_string()
}
);
}
#[tokio::test]
async fn recovery_refuses_a_crash_boundary() {
let repo = TestRepo::new();
let base = repo.head_sha();
let branch = "jack/task-pr-proof";
repo.create_branch(branch);
repo.create_file("first.txt", "first PR\n");
repo.stage_all();
repo.commit("first PR");
let (_home, store, session, _pr) =
rotation_task_with_lease(&repo, branch, &base, None).await;
repo.checkout("main");
repo.create_file("first.txt", "main wins\n");
repo.stage_all();
repo.commit("main advance");
checkout_branch(&repo, branch);
let rebase = Command::new("git")
.current_dir(repo.path())
.args(["rebase", "main"])
.output()
.expect("run rebase");
assert!(
!rebase.status.success(),
"rebase must conflict to seed a crash boundary"
);
let err = task_recovery_adoption(&store, &session)
.await
.expect_err("crash boundary must refuse");
let message = err.to_string();
assert!(
message.contains("mid-rebase"),
"expected crash-boundary refusal, got: {message}"
);
assert!(
message.contains("refused before moving any ownership"),
"refusal must name the contract, got: {message}"
);
let _ = Command::new("git")
.current_dir(repo.path())
.args(["rebase", "--abort"])
.status();
}
#[tokio::test]
async fn between_prs_recovery_selects_the_deterministic_next_branch() {
let repo = TestRepo::new();
let base = repo.head_sha();
let first_branch = "jack/task-pr-proof";
repo.create_branch(first_branch);
repo.create_file("first.txt", "first PR\n");
repo.stage_all();
repo.commit("first PR");
let (_home, store, session, first) =
rotation_task_with_lease(&repo, first_branch, &base, None).await;
settle_pr(&store, first, "merge-942", Some("follow-up")).await;
let settled = store
.task_prs(&session.id)
.await
.expect("read settled PR")
.pop()
.expect("settled PR");
let next = format!("{first_branch}-follow-up");
assert_eq!(next_pr_slug(&settled, None), "follow-up");
Command::new("git")
.current_dir(repo.path())
.args(["checkout", "-b", &next])
.status()
.expect("cut next branch");
let adoption = task_recovery_adoption(&store, &session)
.await
.expect("partial rotation onto the next branch is adopted");
assert_eq!(
adoption,
TaskRecoveryAdoption::BetweenPrs {
settled: first_branch.to_string(),
next
}
);
}
#[tokio::test]
async fn refuse_dirty_between_prs_blocks_after_an_out_of_band_merge() {
let repo = TestRepo::new();
let base = repo.head_sha();
let first_branch = "jack/task-pr-proof";
repo.create_branch(first_branch);
repo.create_file("first.txt", "first PR\n");
repo.stage_all();
repo.commit("first PR");
let (_home, store, session, first) =
rotation_task_with_lease(&repo, first_branch, &base, None).await;
repo.create_file("wip.txt", "ongoing\n");
refuse_dirty_between_prs(&store, &session)
.await
.expect("active PR with dirty work is allowed");
settle_pr(&store, first, "merge-out-of-band", None).await;
let err = refuse_dirty_between_prs(&store, &session)
.await
.expect_err("dirty between-PR after merge must refuse");
assert!(
err.to_string().contains("lf pr next"),
"expected dirty between-PR refusal, got: {err}"
);
}
use super::{
complete_task_session_with_authority, reconcile_task_completion, task_completion_gate,
};
use crate::task::{TaskGateProposal, TaskLifecyclePhase, TaskLifecyclePlan};
async fn gate_task_fixture(
repo: &TestRepo,
branch: &str,
base_commit: &str,
settle: bool,
phase_cursor: u32,
) -> (tempfile::TempDir, SharedStore, TaskSession, TaskPr) {
let home = tempfile::tempdir().expect("task home");
let db_path = home.path().join("loopflow.db");
let store = Arc::new(
open_store(&StorageConfig::sqlite(db_path.clone()))
.await
.expect("open store"),
);
let now = OffsetDateTime::now_utc();
let wave = Wave::new(
WaveId::new(),
"completion-gate".to_string(),
repo.path().display().to_string(),
);
let project = ProjectSession {
id: ProjectSessionId::new(),
launch: ProjectLaunchReceipt {
project: LinearProjectSnapshot {
id: LinearProjectId::new(format!("project-{}", WaveId::new()))
.expect("project id"),
slug: "completion-gate".to_string(),
name: "Completion gate".to_string(),
prompt_context: "Keep the gate honest.".to_string(),
},
pm_snapshot_synced_at: now.unix_timestamp(),
},
wave_id: wave.id().clone(),
status: ProjectSessionStatus::Running,
status_reason: "test project is running".to_string(),
status_at: now,
iteration: 1,
observation_cursor: 0,
last_state_fingerprint: None,
agent: "codex".to_string(),
provider: "codex".to_string(),
provider_session_id: Some("completion-gate".to_string()),
latest_process: Some(ChildProcessGeneration {
generation: 1,
pid: None,
process_group_id: None,
tmux_name: "completion-gate".to_string(),
agent: "codex".to_string(),
provider: "codex".to_string(),
provider_session_id: Some("completion-gate".to_string()),
started_at: now,
state: crate::child_session::ChildLeaseState::Active,
outcome: None,
provenance: None,
}),
abandon_intent: None,
created_at: now,
updated_at: now,
};
let session = TaskSession {
id: TaskSessionId::new(),
launch: TaskLaunchReceipt {
issue: LinearIssueSnapshot {
id: LinearIssueId::new(format!("issue-{}", WaveId::new())).expect("issue id"),
identifier: "INF-GATE".to_string(),
title: "Prove the completion gate".to_string(),
description: "Completion waits on the review gate.".to_string(),
},
project: project.launch.project.clone(),
pm_snapshot_synced_at: now.unix_timestamp(),
},
pm_writeback: PmWritebackState::Current,
wave_id: wave.id().clone(),
project_session_id: project.id.clone(),
status: TaskSessionStatus::Waiting,
status_reason: "merged; awaiting gate".to_string(),
status_at: now,
worktree: repo.path().to_path_buf(),
workspace_slug: "completion-gate".to_string(),
lifecycle: TaskLifecyclePlan::standard("task"),
lifecycle_phase: TaskLifecyclePhase::Gate,
phase_epoch: 2,
phase_cursor,
phase_iteration: 0,
gate_cycle: 1,
gate_proposal: Some(TaskGateProposal {
status: TaskSessionStatus::Completed,
reason: "merged and reviewed".to_string(),
}),
agent: "codex".to_string(),
provider: "codex".to_string(),
provider_session_id: None,
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: branch.to_string(),
base_commit: base_commit.to_string(),
parent_pr_id: None,
publication: None,
merge_commit: None,
abandoned_at: None,
ci_observation: None,
github_observation: None,
linear_attachment_id: None,
linear_comment_id: None,
linear_link_error: None,
created_at: now,
updated_at: now,
};
store.create_wave(&wave).await.expect("create wave");
store
.create_project_session(&project)
.await
.expect("create project");
store
.create_task_session(&session, &pr)
.await
.expect("create Task");
if !settle {
return (home, store, session, pr);
}
let mut merged = pr.clone();
merged.publication = Some(PrPublication {
requested_at: now,
after_merge: AfterMerge::CompleteTask,
next_slug: None,
github: Some(GithubPr {
number: 912,
url: "https://example.com/pr/912".to_string(),
head_sha: Some(repo.head_sha()),
}),
});
merged.merge_commit = Some("merge-912".to_string());
merged.updated_at = OffsetDateTime::now_utc();
store
.update_task_pr(&merged)
.await
.expect("settle merged PR");
(home, store, session, merged)
}
async fn gate_task(
repo: &TestRepo,
branch: &str,
base_commit: &str,
) -> (tempfile::TempDir, SharedStore, TaskSession, TaskPr) {
gate_task_fixture(repo, branch, base_commit, true, 0).await
}
#[test]
fn merged_snapshot_with_a_clear_gate_recommends_completion() {
let repo = TestRepo::new();
let branch = "jack/merged-complete-proof";
repo.create_branch(branch);
let base = repo.head_sha();
let runtime = tokio::runtime::Runtime::new().expect("test runtime");
let (home, session) = runtime.block_on(async {
let (home, _store, session, _pr) = gate_task(&repo, branch, &base).await;
(home, session)
});
drop(runtime);
let _env = StoreEnvGuard::new(home.path());
let snapshot = task_snapshot(&session).expect("snapshot completable Task");
assert_eq!(snapshot.active_pr, None);
assert_eq!(snapshot.actions.recommended, Some(TaskAction::Complete));
assert!(
snapshot
.actions
.status(TaskAction::Complete)
.expect("Complete action")
.available
);
assert_eq!(
snapshot
.actions
.status(TaskAction::Resume)
.expect("Resume action")
.reason,
"Task INF-GATE has no active PR to resume; pull request #912 merged"
);
}
#[tokio::test]
async fn completion_gate_is_satisfied_with_no_prs_and_no_reviews() {
let repo = TestRepo::new();
let branch = "jack/gate-proof";
repo.create_branch(branch);
let (_home, store, mut session, _pr) = gate_task(&repo, branch, &repo.head_sha()).await;
session.status = TaskSessionStatus::Waiting;
let gate = task_completion_gate(&store, &session).await.expect("gate");
assert!(
gate.satisfied,
"gate should be satisfied: {:?}",
gate.blockers
);
}
const SUCCESSOR_BRANCH: &str = "jack/gate-proof-2";
async fn merged_task_with_successor(
repo: &TestRepo,
successor_base: Option<&str>,
) -> (tempfile::TempDir, SharedStore, TaskSession, TaskPr, TaskPr) {
let merged_branch = "jack/gate-proof";
let base = repo.head_sha();
repo.create_branch(merged_branch);
repo.create_file("merged.txt", "the real work\n");
repo.stage_all();
repo.commit("merged work");
let (home, store, mut session, merged) = gate_task(repo, merged_branch, &base).await;
session.status = TaskSessionStatus::Waiting;
let mut merged = merged;
merged
.publication
.as_mut()
.expect("published PR")
.after_merge = AfterMerge::Review;
merged.updated_at = OffsetDateTime::now_utc();
store
.update_task_pr(&merged)
.await
.expect("reland merged PR for review");
let merged = store
.get_task_pr(&merged.id)
.await
.expect("read merged PR")
.expect("merged row");
match successor_base {
Some(base) => git(repo.path(), &["checkout", "-b", SUCCESSOR_BRANCH, base]),
None => {
git(repo.path(), &["checkout", "--orphan", SUCCESSOR_BRANCH]);
git(
repo.path(),
&["commit", "--allow-empty", "-m", "rewritten successor"],
)
}
};
let now = OffsetDateTime::now_utc();
let successor = TaskPr {
id: TaskPrId::new(),
task_session_id: session.id.clone(),
sequence: merged.sequence + 1,
slug: "successor".to_string(),
branch: SUCCESSOR_BRANCH.to_string(),
base_commit: successor_base.unwrap_or(&base).to_string(),
parent_pr_id: None,
publication: None,
merge_commit: None,
abandoned_at: None,
ci_observation: None,
github_observation: None,
linear_attachment_id: None,
linear_comment_id: None,
linear_link_error: None,
created_at: now,
updated_at: now,
};
store
.settle_task_pr(&merged, Some(&successor))
.await
.expect("rotate successor");
(home, store, session, merged, successor)
}
async fn pr_phase(store: &SharedStore, id: &TaskPrId) -> PrPhase {
store
.get_task_pr(id)
.await
.expect("read PR")
.expect("PR row")
.phase()
}
#[tokio::test]
async fn a_merged_task_settles_over_a_proven_empty_successor() {
let repo = TestRepo::new();
let base = repo.head_sha();
let (_home, store, mut session, _merged, successor) =
merged_task_with_successor(&repo, Some(&base)).await;
let gate = task_completion_gate(&store, &session)
.await
.expect("gate over empty successor");
assert!(
gate.satisfied,
"the rotation's empty artifact must not block merged work: {:?}",
gate.blockers
);
assert_eq!(
gate.discardable_successor.as_ref().map(|pr| &pr.id),
Some(&successor.id)
);
assert_eq!(
pr_phase(&store, &successor.id).await,
PrPhase::Working,
"the gate must classify without mutating"
);
session.set_status(TaskSessionStatus::Completed, "merged work settled");
complete_task_session_with_authority(
&store,
&session,
gate.discardable_successor.as_ref(),
None,
)
.await
.expect("complete over the empty successor");
let prs = store.task_prs(&session.id).await.expect("read PRs");
assert_eq!(
prs.len(),
1,
"the artifact must be gone and no replacement minted: {prs:?}"
);
assert_eq!(prs[0].phase(), PrPhase::Merged);
assert!(store
.active_task_pr(&session.id)
.await
.expect("read active PR")
.is_none());
reconcile_task_completion(&store, &mut session, None)
.await
.expect("reconcile after settling");
assert_eq!(session.status, TaskSessionStatus::Completed);
assert_eq!(
store.task_prs(&session.id).await.expect("read PRs").len(),
1
);
}
#[tokio::test]
async fn a_successor_that_is_not_proven_empty_is_never_discardable() {
for (case, commits_work, expected) in [
("range", true, "follow-up work is committed on unpublished"),
(
"unprovable",
false,
"recorded base is not an ancestor of the unpublished branch",
),
] {
let repo = TestRepo::new();
let base = repo.head_sha();
let (_home, store, session, _merged, successor) =
merged_task_with_successor(&repo, commits_work.then_some(base.as_str())).await;
if commits_work {
repo.create_file("follow-up.txt", "work that must not be dropped\n");
repo.stage_all();
repo.commit("follow-up work");
}
let gate = task_completion_gate(&store, &session)
.await
.unwrap_or_else(|error| panic!("{case}: {error}"));
assert!(
gate.blockers
.iter()
.any(|blocker| blocker.contains(expected)),
"{case}: expected blocker containing {expected:?}, got {:?}",
gate.blockers
);
assert_eq!(
pr_phase(&store, &successor.id).await,
PrPhase::Working,
"{case} must leave the row untouched"
);
if commits_work {
assert_eq!(
super::rev_parse(repo.path(), SUCCESSOR_BRANCH).expect("resolve tip"),
repo.head_sha(),
"{case} must preserve the committed follow-up"
);
}
}
}
#[tokio::test]
async fn completion_gate_blocks_when_follow_up_range_is_unprovable() {
let repo = TestRepo::new();
let branch = "jack/gate-proof";
let base = repo.head_sha();
repo.create_branch(branch);
let (_home, store, session, pr) = gate_task(&repo, branch, &base).await;
let mut unprovable = pr.clone();
unprovable
.publication
.as_mut()
.and_then(|publication| publication.github.as_mut())
.expect("published github PR")
.head_sha = None;
unprovable.updated_at = OffsetDateTime::now_utc();
store
.update_task_pr(&unprovable)
.await
.expect("remove published head");
let gate = task_completion_gate(&store, &session)
.await
.expect("gate with missing head");
assert!(!gate.satisfied);
assert!(
gate.reason()
.contains("published pull request head is missing"),
"missing head must fail closed: {}",
gate.reason()
);
git(
repo.path(),
&["checkout", "--orphan", "jack/rewritten-head"],
);
git(
repo.path(),
&["commit", "--allow-empty", "-m", "rewritten published head"],
);
let rewritten_head = repo.head_sha();
git(repo.path(), &["checkout", branch]);
unprovable
.publication
.as_mut()
.and_then(|publication| publication.github.as_mut())
.expect("published github PR")
.head_sha = Some(rewritten_head);
unprovable.updated_at = OffsetDateTime::now_utc();
store
.update_task_pr(&unprovable)
.await
.expect("record rewritten published head");
let gate = task_completion_gate(&store, &session)
.await
.expect("gate with rewritten head");
assert!(!gate.satisfied);
assert!(
gate.reason()
.contains("published pull request head is not an ancestor"),
"rewritten head must fail closed: {}",
gate.reason()
);
}
#[tokio::test]
#[allow(clippy::await_holding_lock)] async fn completion_is_withheld_over_work_committed_past_the_merged_tip() {
let repo = TestRepo::new();
let branch = "jack/gate-proof";
let base = repo.head_sha();
repo.create_branch(branch);
repo.create_file("merged.txt", "merged work\n");
repo.stage_all();
repo.commit("merged work");
let merged_tip = repo.head_sha();
let (_home, store, mut session, pr) = gate_task(&repo, branch, &base).await;
let mut open = pr.clone();
open.merge_commit = None;
open.github_observation = None;
open.publication
.as_mut()
.and_then(|publication| publication.github.as_mut())
.expect("published github PR")
.head_sha = Some(merged_tip.clone());
store.update_task_pr(&open).await.expect("open fixture PR");
repo.create_file("follow-up.txt", "acknowledged follow-up\n");
repo.stage_all();
repo.commit("follow-up the directive asked for");
git(
repo.path(),
&[
"remote",
"set-url",
"origin",
"https://github.com/test/repo.git",
],
);
let bin = tempfile::tempdir().expect("fake gh bin");
let gh = bin.path().join("gh");
std::fs::write(
&gh,
format!(
"#!/bin/sh\nif [ \"$1\" = \"--version\" ]; then exit 0; fi\nprintf '%s\\n' '{{\"merged\":true,\"state\":\"closed\",\"merge_commit_sha\":\"merge-912\",\"number\":912,\"html_url\":\"https://example.com/pr/912\",\"head\":{{\"sha\":\"{merged_tip}\"}}}}'\n"
),
)
.expect("write fake gh");
let mut permissions = std::fs::metadata(&gh).expect("stat fake gh").permissions();
permissions.set_mode(0o755);
std::fs::set_permissions(&gh, permissions).expect("make fake gh executable");
let _env_lock = crate::journal::test_env_lock();
let previous_path = std::env::var_os("PATH");
let path = match &previous_path {
Some(previous) => format!("{}:{}", bin.path().display(), previous.to_string_lossy()),
None => bin.path().display().to_string(),
};
std::env::set_var("PATH", path);
let observed = reconcile_task_pr(&store, &mut session).await;
match previous_path {
Some(path) => std::env::set_var("PATH", path),
None => std::env::remove_var("PATH"),
}
let observed = observed
.expect("observe merged PR")
.expect("persist merged PR");
assert_eq!(observed.phase(), PrPhase::Merged);
assert_ne!(
session.status,
TaskSessionStatus::Completed,
"merge observation must not complete over committed follow-up"
);
reconcile_task_completion(&store, &mut session, None)
.await
.expect("reconcile with follow-up outstanding");
assert_ne!(
session.status,
TaskSessionStatus::Completed,
"reconcile must not complete over committed follow-up"
);
assert_eq!(session.pm_writeback, PmWritebackState::Current);
assert!(repo.path().join("follow-up.txt").exists());
}
#[tokio::test]
#[allow(clippy::await_holding_lock)] async fn project_supervision_keeps_pending_and_unknown_gate_tasks_parked() {
let _env_lock = crate::journal::test_env_lock();
for (label, state) in [("pending", Some(CiState::Pending)), ("unknown", None)] {
let repo = TestRepo::new();
let branch = format!("jack/gate-{label}");
let (home, store, project, session, _pr) =
parked_gate_task(&repo, &branch, state, 0).await;
let _launch_env = TaskLaunchEnv::install(home.path());
supervise_project_task_bodies(&store, &project)
.await
.expect("supervise parked Gate");
let persisted = store
.get_task_session(&session.id)
.await
.expect("read Task")
.expect("Task exists");
assert_eq!(persisted.status, TaskSessionStatus::Waiting, "{label}");
assert_eq!(
persisted.lifecycle_phase,
crate::task::TaskLifecyclePhase::Gate,
"{label}"
);
assert_eq!(persisted.phase_cursor, 0, "{label}");
}
}
#[tokio::test]
#[allow(clippy::await_holding_lock)] async fn project_supervision_reserves_one_typed_ci_run_for_red() {
let _env_lock = crate::journal::test_env_lock();
let repo = TestRepo::new();
let (home, store, project, session, _pr) =
parked_gate_task(&repo, "jack/gate-red", Some(CiState::Failing), 0).await;
let _launch_env = TaskLaunchEnv::install(home.path());
supervise_project_task_bodies(&store, &project)
.await
.expect("first red supervision");
supervise_project_task_bodies(&store, &project)
.await
.expect("duplicate red supervision");
let incidents = store
.ci_incidents_since(OffsetDateTime::UNIX_EPOCH, None, None)
.await
.expect("read CI incidents");
assert_eq!(incidents.len(), 1, "the incident identity deduplicates");
let work = store
.work_for_child(&ChildRef::Task(session.id.clone()))
.await
.expect("Task Work");
let run = store
.current_run(&work)
.await
.expect("read Run")
.expect("CI reserves a Run");
assert!(matches!(
&run.trigger,
crate::durable::RunTrigger::CiIncident { incident_id }
if incident_id == &incidents[0].incident.identity
));
assert_eq!(incidents[0].incident.claimed_run_id.as_ref(), Some(&run.id));
let persisted = store
.get_task_session(&session.id)
.await
.expect("read Task")
.expect("Task exists");
assert_eq!(
persisted
.latest_process
.as_ref()
.expect("ci-fix launch reserves a generation")
.generation,
1
);
}
#[tokio::test]
#[allow(clippy::await_holding_lock)] async fn project_supervision_claims_live_run_without_reserving_a_sibling() {
let _env_lock = crate::journal::test_env_lock();
let repo = TestRepo::new();
let (home, store, project, mut session, _pr) =
parked_gate_task(&repo, "jack/gate-live-idle", Some(CiState::Failing), 0).await;
session.begin_generation("lf-live-idle-control".to_string());
let lease = store
.reserve_task_process(&session, TaskSessionStatus::Waiting)
.await
.expect("reserve live control body")
.expect("waiting Task reserves one generation");
if let Some(process) = session.latest_process.as_mut() {
process.state = crate::child_session::ChildLeaseState::Active;
}
session.set_status(
TaskSessionStatus::Running,
"Project review owns the Gate step; provider turn is idle",
);
store
.activate_task_process(&session, &lease)
.await
.expect("activate live control body");
let work = store
.work_for_child(&ChildRef::Task(session.id.clone()))
.await
.expect("Task Work");
let run = store
.current_run(&work)
.await
.expect("read active Run")
.expect("active process has a Run");
let _launch_env = TaskLaunchEnv::install(home.path());
supervise_project_task_bodies(&store, &project)
.await
.expect("first live-idle red supervision");
supervise_project_task_bodies(&store, &project)
.await
.expect("duplicate live-idle red supervision");
let incidents = store
.ci_incidents_since(OffsetDateTime::UNIX_EPOCH, None, None)
.await
.expect("read CI incidents");
let incident = incidents
.iter()
.find(|row| row.incident.task_session_id == session.id)
.expect("live-idle incident exists");
assert_eq!(
incident.incident.claimed_run_id.as_ref(),
Some(&run.id),
"the current Run owns actionable CI"
);
assert!(incident.incident.responded_at.is_some());
let current = store
.current_run(&work)
.await
.expect("read current Run")
.expect("Run remains current");
assert_eq!(current.id, run.id, "CI must not reserve a sibling Run");
let persisted = store
.get_task_session(&session.id)
.await
.expect("read Task")
.expect("Task exists");
assert_eq!(
persisted.status,
TaskSessionStatus::Running,
"{}",
persisted.status_reason
);
assert_eq!(
persisted
.latest_process
.as_ref()
.expect("control generation remains active")
.generation,
lease.generation,
"supervision must not interrupt or replace the live control body"
);
}
#[tokio::test]
#[allow(clippy::await_holding_lock)] async fn project_supervision_relaunches_a_fresh_green_gate_exactly_once() {
let _env_lock = crate::journal::test_env_lock();
let repo = TestRepo::new();
let (home, store, project, session, _pr) =
parked_gate_task(&repo, "jack/gate-green", Some(CiState::Passing), 0).await;
let _launch_env = TaskLaunchEnv::install(home.path());
supervise_project_task_bodies(&store, &project)
.await
.expect("first green supervision");
supervise_project_task_bodies(&store, &project)
.await
.expect("duplicate green supervision");
let persisted = store
.get_task_session(&session.id)
.await
.expect("read Task")
.expect("Task exists");
assert_eq!(persisted.status, TaskSessionStatus::Starting);
assert_eq!(
persisted
.latest_process
.as_ref()
.expect("green Gate relaunch reserves a generation")
.generation,
1,
"the second observation must not reserve another generation"
);
}
#[tokio::test]
#[allow(clippy::await_holding_lock)] async fn project_supervision_does_not_relaunch_green_after_gate_advances() {
let _env_lock = crate::journal::test_env_lock();
let repo = TestRepo::new();
let (home, store, project, session, _pr) =
parked_gate_task(&repo, "jack/gate-advanced", Some(CiState::Passing), 1).await;
let _launch_env = TaskLaunchEnv::install(home.path());
supervise_project_task_bodies(&store, &project)
.await
.expect("supervise advanced Gate");
let persisted = store
.get_task_session(&session.id)
.await
.expect("read Task")
.expect("Task exists");
assert_eq!(persisted.phase_cursor, 1);
assert!(
persisted.latest_process.is_none(),
"green observation cannot reopen an advanced Gate"
);
}
fn open_pr_with_ci(number: u32, head_sha: &str, state: Option<CiState>) -> TaskPr {
let now = OffsetDateTime::now_utc();
TaskPr {
id: TaskPrId::new(),
task_session_id: TaskSessionId::new(),
sequence: 1,
slug: "w2-231".to_string(),
branch: "feature".to_string(),
base_commit: "base".to_string(),
parent_pr_id: None,
publication: Some(PrPublication {
requested_at: now,
after_merge: AfterMerge::Review,
next_slug: None,
github: Some(GithubPr {
number,
url: format!("https://github.com/loopflow/loopflow/pull/{number}"),
head_sha: Some(head_sha.to_string()),
}),
}),
merge_commit: None,
abandoned_at: None,
github_observation: None,
ci_observation: state.map(|state| CiObservation {
head_sha: head_sha.to_string(),
state,
failing_checks: if state == CiState::Failing {
vec![CiCheck {
name: "build".to_string(),
url: Some("https://ci/build".to_string()),
}]
} else {
Vec::new()
},
observed_at: now,
}),
linear_attachment_id: None,
linear_comment_id: None,
linear_link_error: None,
created_at: now,
updated_at: now,
}
}
#[test]
fn decide_open_pr_status_blocks_on_degraded_github_observation() {
let pr = open_pr_with_ci(42, "sha-1", Some(CiState::Failing));
let (status, reason) =
decide_open_pr_status(&pr, Some("GitHub API rate limit exceeded"), false);
assert_eq!(status, TaskSessionStatus::Blocked);
assert!(reason.contains("github-observation"), "reason: {reason}");
assert!(reason.contains("rate limit"), "reason: {reason}");
assert!(reason.contains("#42 stays attached"), "reason: {reason}");
}
#[test]
fn decide_open_pr_status_blocks_on_no_change_failing_ci() {
let pr = open_pr_with_ci(42, "sha-1", Some(CiState::Failing));
let (status, reason) = decide_open_pr_status(&pr, None, false);
assert_eq!(status, TaskSessionStatus::Blocked);
assert!(reason.contains("did not repair"), "reason: {reason}");
assert!(reason.contains("#42 stays attached"), "reason: {reason}");
}
#[test]
fn decide_open_pr_status_waits_when_head_advanced_even_if_failing() {
let pr = open_pr_with_ci(42, "sha-2", Some(CiState::Failing));
let (status, reason) = decide_open_pr_status(&pr, None, true);
assert_eq!(status, TaskSessionStatus::Waiting);
assert!(reason.contains("open for review"), "reason: {reason}");
}
#[test]
fn decide_open_pr_status_waits_on_healthy_ci() {
let pending = open_pr_with_ci(42, "sha-1", Some(CiState::Pending));
assert_eq!(
decide_open_pr_status(&pending, None, false).0,
TaskSessionStatus::Waiting
);
let passing = open_pr_with_ci(42, "sha-1", Some(CiState::Passing));
assert_eq!(
decide_open_pr_status(&passing, None, false).0,
TaskSessionStatus::Waiting
);
let unknown = open_pr_with_ci(42, "sha-1", None);
assert_eq!(
decide_open_pr_status(&unknown, None, false).0,
TaskSessionStatus::Waiting
);
}
#[test]
fn decide_open_pr_status_degraded_dominates_healthy_ci() {
let pr = open_pr_with_ci(42, "sha-1", Some(CiState::Passing));
let (status, reason) = decide_open_pr_status(&pr, Some("network unreachable"), true);
assert_eq!(status, TaskSessionStatus::Blocked);
assert!(reason.contains("github-observation"), "reason: {reason}");
}
}