use std::collections::{BTreeMap, BTreeSet, HashSet};
use std::fs::{File, OpenOptions};
use std::path::{Component, Path, PathBuf};
use std::process::Command;
use std::sync::Arc;
use std::time::{Duration, Instant};
use crate::child::ChildRef;
use crate::durable::{
AgentInvocation, AuthenticatedRequest, Containment, ContainmentObservation, ControlCtx, Run,
RunLease, RunState, RunTrigger, WaitOn, WorkRef, WorkStatus,
};
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, is_materially_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::{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::planning::{LinearIssueId, TaskPlan};
use crate::store::{
open_existing_store, open_registry_for_authority, RegistryUnavailable, SharedStore, Store,
StoreError,
};
use crate::task::actions::{derive_task_actions, TaskActionEvidence, TaskActionModel};
use crate::task::{
AfterMerge, CiCheck, CiIncident, CiObservation, CiState, GithubObservation,
GithubObservationResult, GithubPr, Observation, PmWritebackOperation, PmWritebackState,
PrMergeMode, PrMergeRequest, PrPhase, PrPublication, Task, TaskEventKind, TaskId, TaskPr,
TaskPrId,
};
use crate::wave::Wave;
use fs2::FileExt;
use sha2::{Digest, Sha256};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum TaskWaitUntil {
Open,
Terminal,
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct TaskFlowOverrides {
pub first: Option<String>,
pub loop_: Option<String>,
pub finally: Option<String>,
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct TaskLaunchOptions {
pub name: Option<String>,
pub flows: TaskFlowOverrides,
pub stack_on: Option<String>,
pub directive: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct TaskStartInput {
pub title: String,
pub report: String,
}
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)]
pub struct TaskControlResult {
pub issue_id: String,
pub task_id: String,
pub receipt: super::child::WorkControlReceipt,
pub observation: Observation,
}
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)]
pub struct TaskSnapshot {
pub issue_id: String,
pub issue_identifier: String,
pub task_id: String,
pub external_project_id: String,
pub project: String,
pub pm_snapshot_synced_at: i64,
pub pm_writeback: crate::task::PmWritebackState,
pub wave: String,
pub project_id: String,
pub status: WorkStatus,
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 invocation: Option<AgentInvocation>,
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 task_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 task_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 task_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,
task_id: &'a crate::task::TaskId,
worktree: &'a Path,
base_commit: &'a str,
}
impl<'a> TaskWorkspace<'a> {
fn new(task: &'a Task, pr: &'a TaskPr) -> Self {
Self {
issue_identifier: &task.plan.identifier,
task_id: &task.id,
worktree: &task.worktree,
base_commit: &pr.base_commit,
}
}
}
fn active_pr(task: &Task) -> OpsResult<TaskPr> {
let task_id = task.id.clone();
block_on_task(async move {
task_store()
.await?
.active_task_pr(&task_id)
.await
.map_err(|error| task_error(format!("failed to read active PR: {error}")))?
.ok_or_else(|| task_error("Task 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, task, .. } = resolve_task_authority(worktree).await?
else {
return Ok(None);
};
let Some(active) = store
.active_task_pr(&task.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_task = store
.get_task(&parent.task_id)
.await
.map_err(|error| task_error(error.to_string()))?
.ok_or_else(|| task_error("stack parent Task is missing"))?;
reconcile_task_pr(&store, &mut parent_task).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, task: &Task) -> OpsResult<Wave> {
store
.get_wave(&task.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", task.wave_id)))
}
async fn task_work_status(store: &Store, task: &Task) -> OpsResult<WorkStatus> {
let work = store
.work_for_child(&ChildRef::Task(task.id.clone()))
.await
.map_err(|error| task_error(error.to_string()))?;
store
.work_status(&work)
.await
.map_err(|error| task_error(error.to_string()))
}
pub fn task_run(repo: &Path, issue: &str, options: TaskLaunchOptions) -> OpsResult<Task> {
let TaskLaunchOptions {
name,
flows,
stack_on,
directive,
} = 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_by_issue(issue)
.await
.map_err(|error| task_error(format!("failed to read task registry: {error}")))?;
if let Some(task) = &mut existing {
let status = task_work_status(&store, task).await?;
if matches!(status, WorkStatus::Done | WorkStatus::Abandoned) {
let predecessor_id = task.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() != task.workspace_slug {
return Err(task_error(format!(
"Task {} already uses workspace name {:?}",
task.plan.identifier, task.workspace_slug
)));
}
}
ensure_task_flow_override(
&task.worktree,
&task.plan.identifier,
"first",
flows.first.as_deref(),
&task.lifecycle.first.flow,
)?;
ensure_task_flow_override(
&task.worktree,
&task.plan.identifier,
"loop",
flows.loop_.as_deref(),
&task.lifecycle.loop_.flow,
)?;
ensure_task_flow_override(
&task.worktree,
&task.plan.identifier,
"finally",
flows.finally.as_deref(),
&task.lifecycle.finally.flow,
)?;
if let Some(requested) = stack_on.as_deref() {
let active = store
.active_task_pr(&task.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}",
task.plan.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_task = store
.get_task(&parent.task_id)
.await
.map_err(|error| task_error(error.to_string()))?
.ok_or_else(|| task_error("stack parent Task is missing"))?;
if requested != parent_task.plan.identifier
&& requested != parent_task.plan.id.as_str()
{
return Err(task_error(format!(
"Task {} is stacked on {}, not {requested}",
task.plan.identifier, parent_task.plan.identifier
)));
}
}
if directive.is_some() {
return Err(task_error(format!(
"Task {} already exists; use `lf task steer {} <new-direction>`",
task.plan.identifier, task.plan.identifier,
)));
}
}
Ok((existing, None))
})?;
if let Some(existing) = existing {
return task_status(existing.plan.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 =
crate::ops::task_pm::resolve_task(&main_repo, issue, crate::ops::pm::PmRefresh::Auto)?;
let project_flows = resolved
.project
.flows
.clone()
.unwrap_or_else(crate::pm::ProjectFlowPlan::empty);
let lifecycle = resolve_task_lifecycle(&main_repo, &project_flows, &flows)?;
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: {} ({})",
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_task = store
.get_task_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; run it first"
))
})?;
if parent_task.plan.id.as_str() == resolved.item.id {
return Err(task_error("a Task cannot stack on itself"));
}
let parent = store
.active_task_pr(&parent_task.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_task.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 = crate::ops::project::ensure_project_for_task(
&main_repo,
crate::ops::task_pm::ResolvedProject {
snapshot: resolved.snapshot.clone(),
project: resolved.project.clone(),
},
)?;
let project_id = project.id.clone();
let wave_id = project.wave_id.clone();
let config = load_config_or_default(Some(&main_repo));
let agent = config.agent();
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_by_issue(&resolved.item.id)
.await
.map_err(|error| task_error(format!("failed to read task registry: {error}")))?
{
Some(existing)
if !matches!(
task_work_status(&store, &existing).await?,
WorkStatus::Done | WorkStatus::Abandoned
) =>
{
return Ok(existing)
}
Some(terminal) => Some(terminal),
None => None,
};
let now = time::OffsetDateTime::now_utc();
let task_id = predecessor
.as_ref()
.map(|task| task.id.clone())
.unwrap_or_else(crate::task::TaskId::new);
let sequence = if let Some(predecessor) = &predecessor {
store
.task_prs(&predecessor.id)
.await
.map_err(|error| task_error(format!("failed to read Task PR history: {error}")))?
.last()
.map_or(1, |pr| pr.sequence + 1)
} else {
1
};
let mut task = Task {
id: task_id,
plan: TaskPlan {
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(),
pm_snapshot_synced_at: resolved.snapshot.synced_at,
},
wave_id,
project_id,
pm_writeback: PmWritebackState::Current,
worktree: plan.worktree_path.clone(),
workspace_slug: workspace_slug.clone(),
lifecycle,
lifecycle_phase: crate::task::TaskLifecyclePhase::First,
phase_epoch: 1,
phase_cursor: 0,
phase_iteration: 0,
gate_cycle: 0,
gate_proposal: None,
agent,
provider,
provider_session_id: None,
abandon_intent: None,
created_at: predecessor.as_ref().map_or(now, |task| task.created_at),
updated_at: now,
observation: crate::task::Observation::NotRequired,
};
let pr = TaskPr {
id: TaskPrId::new(),
task_id: task.id.clone(),
sequence,
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,
};
if predecessor.is_some() {
store
.reopen_task(&task, Some(&pr), crate::durable::Author::User, &directive)
.await
.map_err(|error| task_error(format!("failed to reopen Task: {error}")))?;
} else {
match store
.create_task_with_steer(&task, &pr, crate::durable::Author::User, &directive)
.await
{
Ok(()) => {}
Err(StoreError::Sqlite(_)) => {
if let Some(existing) = store
.get_task_by_issue(&resolved.item.id)
.await
.map_err(|error| {
task_error(format!("failed to recover task reservation: {error}"))
})?
{
if !matches!(
task_work_status(&store, &existing).await?,
WorkStatus::Done | WorkStatus::Abandoned
) {
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}"))),
}
}
store
.append_task_event(
&task.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 task,
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 task, None).await?;
wait_until_running(&store, &task.id).await
})
}
pub(crate) fn project_context(project: &crate::pm::PmProject) -> String {
let mut context = format!("Definition:\n{}", project.definition.trim());
if let Some(flows) = project
.flows
.as_ref()
.filter(|flows| **flows != crate::pm::ProjectFlowPlan::empty())
{
context.push_str("\n\nProject Task flows:");
if let Some(first) = &flows.first {
context.push_str(&format!("\n- first: {first}"));
}
if let Some(loop_flow) = &flows.loop_ {
context.push_str(&format!("\n- loop: {loop_flow}"));
}
if let Some(finally) = &flows.finally {
context.push_str(&format!("\n- finally: {finally}"));
}
}
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,
project_id: &str,
title: Option<String>,
report: Option<String>,
options: TaskLaunchOptions,
) -> OpsResult<Task> {
let input = resolve_task_start_input(title.as_deref(), report.as_deref())?;
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 project_flows = project
.project
.flows
.clone()
.unwrap_or_else(crate::pm::ProjectFlowPlan::empty);
resolve_task_lifecycle(&main, &project_flows, &options.flows)?;
let marker = format!(
"<!-- loopflow-task-start:{} -->",
hex::encode(Sha256::digest(
format!("{}\0{}\0{}", project.project.id, input.title, input.report).as_bytes()
))
);
let created = crate::ops::task_pm::create_and_load_task(
&main,
&project.snapshot.wave,
&project.project.slug,
&input.title,
&input.report,
&marker,
)?;
task_run(&main, &created.item.id, options)
}
pub fn resolve_task_start_input(
explicit_title: Option<&str>,
piped_report: Option<&str>,
) -> OpsResult<TaskStartInput> {
let report = piped_report
.map(str::trim)
.filter(|value| !value.is_empty());
let title = explicit_title
.map(str::trim)
.filter(|value| !value.is_empty());
let title = match (title, report) {
(Some(title), _) => title.to_string(),
(None, Some(report)) => {
let first_line = report
.lines()
.map(str::trim)
.find(|line| !line.is_empty())
.expect("non-empty report has a meaningful line");
truncate_task_title(first_line, 100)
}
(None, None) => {
return Err(task_error(
"Task title or piped report is required: `pbpaste | lf task start <project>`",
))
}
};
let report = report.unwrap_or(&title).to_string();
Ok(TaskStartInput { title, report })
}
fn truncate_task_title(value: &str, max_chars: usize) -> String {
if value.chars().count() <= max_chars {
return value.to_string();
}
let mut title = value.chars().take(max_chars - 1).collect::<String>();
title.push('…');
title
}
fn ensure_task_flow_override(
repo: &Path,
issue_identifier: &str,
phase: &str,
requested: Option<&str>,
pinned: &str,
) -> OpsResult<()> {
let Some(requested) = requested else {
return Ok(());
};
let requested = resolve_task_flow(repo, requested, phase == "finally")?;
if requested != pinned {
return Err(task_error(format!(
"Task {} already pins {phase} flow {:?}",
issue_identifier, pinned
)));
}
Ok(())
}
fn resolve_task_lifecycle(
repo: &Path,
project: &crate::pm::ProjectFlowPlan,
overrides: &TaskFlowOverrides,
) -> OpsResult<crate::task::TaskLifecyclePlan> {
let first = overrides
.first
.as_deref()
.or(project.first.as_deref())
.unwrap_or("task-design");
let loop_flow = overrides
.loop_
.as_deref()
.or(project.loop_.as_deref())
.unwrap_or("slice");
let finally = overrides
.finally
.as_deref()
.or(project.finally.as_deref())
.unwrap_or("ship");
let first = resolve_task_flow(repo, first, false)?;
let loop_flow = resolve_task_flow(repo, loop_flow, false)?;
let finally = resolve_task_flow(repo, finally, true)?;
Ok(crate::task::TaskLifecyclePlan::standard(
first, loop_flow, finally,
))
}
fn validate_task_lifecycle(task: &Task) -> OpsResult<()> {
for (phase, flow, allow_ops) in [
("first", &task.lifecycle.first.flow, false),
("loop", &task.lifecycle.loop_.flow, false),
("finally", &task.lifecycle.finally.flow, true),
] {
resolve_task_flow(&task.worktree, flow, allow_ops).map_err(|error| {
task_error(format!(
"Task {} cannot launch: pinned {phase} flow {flow:?} is invalid: {error}",
task.plan.identifier
))
})?;
}
Ok(())
}
fn resolve_task_flow(repo: &Path, requested: &str, allow_ops: bool) -> OpsResult<String> {
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 allow_ops {
let first_op = steps
.iter()
.position(|step| matches!(step, ConcreteStep::Op(_)));
if matches!(first_op, Some(0))
|| first_op.is_some_and(|index| {
steps[index..]
.iter()
.any(|step| !matches!(step, ConcreteStep::Op(_)))
})
{
return Err(task_error(format!(
"Task finally flow {requested:?} must run one or more skills followed by optional ops"
)));
}
}
if let Some(step) = steps.iter().find(|step| {
!(matches!(step, ConcreteStep::Skill(_))
|| allow_ops && matches!(step, ConcreteStep::Op(_)))
}) {
return Err(task_error(format!(
"Task flow {requested:?} contains {step:?}; first/loop require skills and finally permits skills or ops"
)));
}
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: &TaskId) -> 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<&RunLease>,
) -> Result<(), StoreError> {
match lease {
Some(lease) => store.update_task_pr_for_run(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<&RunLease>,
) -> Result<(), StoreError> {
match lease {
Some(lease) => store.settle_task_pr_for_run(settled, next, lease).await,
None => store.settle_task_pr(settled, next).await,
}
}
async fn append_task_event_with_authority(
store: &SharedStore,
task_id: &crate::task::TaskId,
event: &TaskEventKind,
lease: Option<&RunLease>,
) -> Result<(), StoreError> {
match lease {
Some(lease) => {
store
.append_task_event_for_run(task_id, lease, event)
.await?;
}
None => {
store.append_task_event(task_id, event).await?;
}
}
Ok(())
}
async fn update_task_with_authority(
store: &SharedStore,
task: &Task,
lease: Option<&RunLease>,
) -> Result<(), StoreError> {
match lease {
Some(lease) => store.update_task_for_run(task, lease).await,
None => store.update_task(task).await,
}
}
async fn complete_task_after_pr_with_authority(
store: &SharedStore,
task: &Task,
pr: &TaskPr,
lease: Option<&RunLease>,
) -> Result<(), StoreError> {
match lease {
Some(lease) => store.complete_task_after_pr_for_run(task, pr, lease).await,
None => store.complete_task_after_pr(task, pr).await,
}
}
async fn complete_task_with_authority(
store: &SharedStore,
task: &Task,
skipped_pr: Option<&TaskPr>,
lease: Option<&RunLease>,
) -> Result<(), StoreError> {
match lease {
Some(lease) => store.complete_task_for_run(task, skipped_pr, lease).await,
None => store.complete_task(task, skipped_pr).await,
}
}
async fn ambient_task_run_lease(store: &SharedStore, task: &Task) -> OpsResult<Option<RunLease>> {
let Some(lease) = crate::ops::ambient_run_lease(store).await? else {
return Ok(None);
};
let work = store
.work_for_child(&ChildRef::Task(task.id.clone()))
.await
.map_err(|error| task_error(format!("failed to resolve Task Work: {error}")))?;
if lease.work != work {
return Err(task_error(format!(
"ambient Run {} cannot mutate Task {}",
lease.run_id, task.id
)));
}
Ok(Some(lease))
}
async fn task_for_worktree(
store: &SharedStore,
repo: &Path,
) -> OpsResult<Option<(Task, Option<RunLease>)>> {
let checkout = repo.canonicalize().unwrap_or_else(|_| repo.to_path_buf());
if let Some(lease) = crate::ops::ambient_run_lease(store).await? {
let WorkRef::Task(id) = &lease.work else {
return Err(task_error(format!(
"ambient Run {} owns {}, not Task Work",
lease.run_id,
lease.work.kind()
)));
};
let task = store
.get_task(id)
.await
.map_err(|error| task_error(format!("failed to read ambient Task: {error}")))?
.ok_or_else(|| task_error(format!("ambient Task {id} is not registered")))?;
let worktree = task
.worktree
.canonicalize()
.unwrap_or_else(|_| task.worktree.clone());
if checkout != worktree {
return Err(task_error(format!(
"ambient Task {id} owns {}, not {}",
task.worktree.display(),
repo.display()
)));
}
return Ok(Some((task, Some(lease))));
}
let worktree_keys: BTreeSet<String> = store
.list_tasks(None)
.await
.map_err(|error| task_error(format!("failed to inspect Tasks: {error}")))?
.into_iter()
.filter(|task| {
task.worktree
.canonicalize()
.unwrap_or_else(|_| task.worktree.clone())
== checkout
})
.map(|task| task.worktree.display().to_string())
.collect();
let mut current = Vec::new();
for worktree in worktree_keys {
if let Some(task) = store
.get_task_by_worktree(&worktree)
.await
.map_err(|error| task_error(format!("failed to resolve Task worktree: {error}")))?
{
current.push(task);
}
}
if current.len() > 1 {
return Err(task_error(format!(
"multiple Tasks claim worktree {}",
repo.display()
)));
}
Ok(current.pop().map(|task| (task, None)))
}
#[derive(Debug)]
enum TaskAuthority {
NotATaskWorktree,
Authority {
store: SharedStore,
task: Box<Task>,
lease: Option<RunLease>,
},
}
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(crate::durable::RUN_CONTEXT_ENV).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((task, lease)) => Ok(TaskAuthority::Authority {
store,
task: Box::new(task),
lease,
}),
None => Ok(TaskAuthority::NotATaskWorktree),
}
}
pub(crate) fn request_task_pr_publication(repo: &Path) -> OpsResult<bool> {
block_on_task(async move {
let TaskAuthority::Authority { store, task, lease } = resolve_task_authority(repo).await?
else {
return Ok(false);
};
let mut pr = store
.active_task_pr(&task.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", task.plan.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",
task.plan.identifier, pr.branch
)));
}
let now = time::OffsetDateTime::now_utc();
let github = pr.github().cloned();
let merge = pr
.publication
.as_ref()
.and_then(|publication| publication.merge.as_ref())
.filter(|request| {
github
.as_ref()
.and_then(|github| github.head_sha.as_deref())
== Some(request.head_sha.as_str())
})
.cloned();
pr.publication = Some(PrPublication {
requested_at: pr
.publication
.as_ref()
.map_or(now, |publication| publication.requested_at),
github,
merge,
});
pr.updated_at = now;
match lease.as_ref() {
Some(lease) => store.update_task_pr_for_run(&pr, lease).await,
None => store.update_task_pr(&pr).await,
}
.map_err(|error| task_error(format!("failed to request PR publication: {error}")))?;
Ok(true)
})
}
pub(crate) fn matching_task_pr_merge_request(
repo: &Path,
mode: PrMergeMode,
after_merge: AfterMerge,
next_slug: Option<&str>,
) -> OpsResult<Option<(u32, String)>> {
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, task, .. } = resolve_task_authority(repo).await?
else {
return Ok(None);
};
if !is_clean(repo)? {
return Ok(None);
}
let branch = current_branch(repo)?;
let head = rev_parse(repo, "HEAD")?;
let pr = store
.active_task_pr(&task.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", task.plan.identifier)))?;
if branch.as_deref() != Some(pr.branch.as_str()) {
return Ok(None);
}
let Some(github) = pr.github() else {
return Ok(None);
};
let Some(request) = pr.merge_request() else {
return Ok(None);
};
if github.head_sha.as_deref() != Some(head.as_str())
|| request.mode != mode
|| request.after_merge != after_merge
|| request.next_slug != next_slug
{
return Ok(None);
}
Ok(Some((github.number, head)))
})
}
pub(crate) fn clear_task_pr_merge_before_head_mutation(
repo: &Path,
mutation_is_unconditional: bool,
) -> OpsResult<bool> {
block_on_task(async move {
let TaskAuthority::Authority { store, task, lease } = resolve_task_authority(repo).await?
else {
return Ok(false);
};
clear_task_pr_merge(
&store,
&task,
lease.as_ref(),
repo,
mutation_is_unconditional,
)
.await
})
}
#[derive(Debug)]
pub(crate) struct TaskPrMutationGuard {
_file: File,
}
pub(crate) fn lock_task_pr_mutation(repo: &Path) -> OpsResult<TaskPrMutationGuard> {
let path = crate::engine::git::absolute_git_dir(repo)?.join("lf-pr-mutation.lock");
let file = OpenOptions::new()
.create(true)
.read(true)
.write(true)
.truncate(false)
.open(path)?;
match FileExt::try_lock_exclusive(&file) {
Ok(()) => Ok(TaskPrMutationGuard { _file: file }),
Err(error) if error.kind() == std::io::ErrorKind::WouldBlock => Err(OpsError::Message(
"another PR or branch-head mutation is already running for this worktree".to_string(),
)),
Err(error) => Err(error.into()),
}
}
async fn clear_task_pr_merge(
store: &SharedStore,
task: &Task,
lease: Option<&RunLease>,
repo: &Path,
mutation_is_unconditional: bool,
) -> OpsResult<bool> {
let mut pr = store
.active_task_pr(&task.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", task.plan.identifier)))?;
let Some(request) = pr
.publication
.as_ref()
.and_then(|publication| publication.merge.as_ref())
.cloned()
else {
return Ok(false);
};
if !mutation_is_unconditional {
let head = rev_parse(repo, "HEAD")?;
if is_clean(repo)? && head == request.head_sha {
return Ok(false);
}
}
if request.mode == PrMergeMode::Auto {
let number = pr
.github()
.expect("merge request validation requires GitHub PR")
.number;
crate::ops::pr::disable_auto_merge(repo, number)?;
}
pr.publication
.as_mut()
.expect("merge request requires publication")
.merge = None;
pr.updated_at = time::OffsetDateTime::now_utc();
match lease {
Some(lease) => store.update_task_pr_for_run(&pr, lease).await,
None => store.update_task_pr(&pr).await,
}
.map_err(|error| task_error(format!("failed to clear stale PR merge request: {error}")))?;
Ok(true)
}
pub(crate) fn request_task_pr_merge(
repo: &Path,
mode: PrMergeMode,
head_sha: Option<&str>,
after_merge: AfterMerge,
next_slug: Option<&str>,
) -> OpsResult<bool> {
let head_sha = head_sha.map(str::to_string);
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, task, lease } = resolve_task_authority(repo).await?
else {
return Ok(false);
};
let head_sha = head_sha
.filter(|head| !head.trim().is_empty())
.ok_or_else(|| {
task_error(format!(
"GitHub did not report the current head for Task {}; refusing to request a merge without an exact commit",
task.plan.identifier
))
})?;
let mut pr = store
.active_task_pr(&task.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", task.plan.identifier)))?;
let publication = pr.publication.as_mut().ok_or_else(|| {
task_error(format!(
"Task {} has no durable PR publication request",
task.plan.identifier
))
})?;
let github_head = publication
.github
.as_ref()
.and_then(|github| github.head_sha.as_deref());
if github_head != Some(head_sha.as_str()) {
return Err(task_error(format!(
"Task {} stored GitHub head {:?}, not requested merge head {}; refusing an unpinned settlement",
task.plan.identifier, github_head, head_sha
)));
}
if publication
.merge
.as_ref()
.is_some_and(|request| request.mode == PrMergeMode::Auto)
&& mode == PrMergeMode::User
{
let number = publication
.github
.as_ref()
.expect("merge request validation requires GitHub PR")
.number;
crate::ops::pr::disable_auto_merge(repo, number)?;
}
let now = time::OffsetDateTime::now_utc();
let requested_at = publication
.merge
.as_ref()
.filter(|request| {
request.mode == mode
&& request.head_sha == head_sha
&& request.after_merge == after_merge
&& request.next_slug == next_slug
})
.map_or(now, |request| request.requested_at);
publication.merge = Some(PrMergeRequest {
mode,
requested_at,
head_sha,
after_merge,
next_slug,
});
pr.updated_at = now;
match lease.as_ref() {
Some(lease) => store.update_task_pr_for_run(&pr, lease).await,
None => store.update_task_pr(&pr).await,
}
.map_err(|error| task_error(format!("failed to request PR merge: {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, task, lease } = resolve_task_authority(&repo).await?
else {
return Ok(());
};
verify_task_pr_range_with_authority(&store, &task, lease.as_ref(), &repo).await
})
}
pub(crate) fn verify_task_pr_range_without_healing(repo: &Path) -> OpsResult<()> {
let repo = repo.to_path_buf();
block_on_task(async move {
let TaskAuthority::Authority { store, task, lease } = resolve_task_authority(&repo).await?
else {
return Ok(());
};
verify_task_pr_range_with_authority_mode(
&store,
&task,
lease.as_ref(),
&repo,
StaleBaseAction::Refuse,
None,
)
.await
})
}
pub(crate) fn validate_task_pr_range_for_integration(
repo: &Path,
target_ref: &str,
target_sha: &str,
) -> OpsResult<()> {
verify_task_pr_range_for_integration(repo, target_ref, target_sha, StaleBaseAction::Accept)
}
pub(crate) fn record_task_pr_range_after_integration(
repo: &Path,
target_ref: &str,
target_sha: &str,
) -> OpsResult<()> {
verify_task_pr_range_for_integration(repo, target_ref, target_sha, StaleBaseAction::Heal)
}
fn verify_task_pr_range_for_integration(
repo: &Path,
target_ref: &str,
target_sha: &str,
stale_base: StaleBaseAction,
) -> OpsResult<()> {
let repo = repo.to_path_buf();
let target_ref = target_ref.to_string();
let target_sha = target_sha.to_string();
block_on_task(async move {
let TaskAuthority::Authority { store, task, lease } = resolve_task_authority(&repo).await?
else {
return Ok(());
};
verify_task_pr_range_with_authority_mode(
&store,
&task,
lease.as_ref(),
&repo,
stale_base,
Some((target_ref, target_sha)),
)
.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, task, lease } = resolve_task_authority(&repo).await?
else {
return Ok(());
};
require_task_pr_range_nonempty_with_authority(&store, &task, lease.as_ref(), &repo).await
})
}
pub(crate) fn require_task_pr_range_nonempty_without_healing(repo: &Path) -> OpsResult<()> {
let repo = repo.to_path_buf();
block_on_task(async move {
let TaskAuthority::Authority { store, task, lease } = resolve_task_authority(&repo).await?
else {
return Ok(());
};
require_task_pr_range_nonempty_with_authority_mode(
&store,
&task,
lease.as_ref(),
&repo,
StaleBaseAction::Refuse,
)
.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)
}
pub(crate) async fn verify_task_pr_range_with_authority(
store: &SharedStore,
task: &Task,
lease: Option<&RunLease>,
repo: &Path,
) -> OpsResult<()> {
verify_task_pr_range_with_authority_mode(store, task, lease, repo, StaleBaseAction::Heal, None)
.await
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum StaleBaseAction {
Accept,
Refuse,
Heal,
}
async fn verify_task_pr_range_with_authority_mode(
store: &SharedStore,
task: &Task,
lease: Option<&RunLease>,
repo: &Path,
stale_base: StaleBaseAction,
upstream_override: Option<(String, String)>,
) -> OpsResult<()> {
let mut pr = store
.active_task_pr(&task.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", task.plan.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 {:?}",
task.plan.identifier, pr.branch, branch
)));
}
let (base_ref, upstream) = match upstream_override {
Some(target) => target,
None => {
let default_branch = get_default_branch(repo)?;
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 = &task.plan.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)? {
match stale_base {
StaleBaseAction::Accept => return Ok(()),
StaleBaseAction::Refuse => {
return Err(task_error(format!(
"Task {identifier} PR base {} is stale behind the branch fork {}. Publication does not update integration metadata; run `lf rebase` before publishing.",
short(&base),
short(&merge_base),
)));
}
StaleBaseAction::Heal => {}
}
pr.base_commit = merge_base.clone();
pr.updated_at = time::OffsetDateTime::now_utc();
match lease {
Some(lease) => store.heal_task_pr_base_for_run(&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,
task: &Task,
lease: Option<&RunLease>,
repo: &Path,
) -> OpsResult<()> {
require_task_pr_range_nonempty_with_authority_mode(
store,
task,
lease,
repo,
StaleBaseAction::Heal,
)
.await
}
async fn require_task_pr_range_nonempty_with_authority_mode(
store: &SharedStore,
task: &Task,
lease: Option<&RunLease>,
repo: &Path,
stale_base: StaleBaseAction,
) -> OpsResult<()> {
verify_task_pr_range_with_authority_mode(store, task, lease, repo, stale_base, None).await?;
let pr = store
.active_task_pr(&task.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", task.plan.identifier)))?;
let base = &pr.base_commit;
let identifier = &task.plan.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, task, lease } = resolve_task_authority(repo).await?
else {
return Ok(false);
};
let mut pr = store
.active_task_pr(&task.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", task.plan.identifier)))?;
let github_pr = github_pr.ok_or_else(|| {
task_error(format!(
"GitHub PR for Task {} could not be read after creation or update",
task.plan.identifier
))
})?;
if github_pr.branch != pr.branch {
return Err(task_error(format!(
"Task {} active PR expects branch {:?}, but GitHub reported {:?}",
task.plan.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",
task.plan.identifier
))
})?;
invalidate_stale_merge_request(repo, publication, github_pr)?;
publication.github = Some(GithubPr {
number,
url: url.clone(),
head_sha: github_pr.head_sha.clone(),
});
link_pr_to_linear(&store, &task, &mut pr).await;
pr.updated_at = time::OffsetDateTime::now_utc();
match lease.as_ref() {
Some(lease) => store.update_task_pr_for_run(&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_run(&task.id, lease, &event)
.await
}
None => store.append_task_event(&task.id, &event).await,
}
.map_err(|error| task_error(error.to_string()))?;
}
Ok(true)
})
}
fn invalidate_stale_merge_request(
repo: &Path,
publication: &mut PrPublication,
github_pr: &crate::ops::pr::PrInfo,
) -> OpsResult<()> {
let Some(request) = publication.merge.as_ref() else {
return Ok(());
};
let observed_head = github_pr.head_sha.as_deref().ok_or_else(|| {
task_error(format!(
"GitHub did not report the current head for pull request #{}; refusing to change its head-pinned merge request",
github_pr.number
))
})?;
if observed_head == request.head_sha {
return Ok(());
}
if request.mode == PrMergeMode::Auto && matches!(github_pr.state.as_str(), "open" | "draft") {
let number = u32::try_from(github_pr.number).map_err(|_| {
task_error(format!(
"pull request #{} exceeds supported range",
github_pr.number
))
})?;
crate::ops::pr::disable_auto_merge(repo, number)?;
}
publication.merge = None;
Ok(())
}
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, task, lease } = resolve_task_authority(repo).await?
else {
return Ok(false);
};
let _mutation = lock_task_pr_mutation(repo)?;
let mut pr = store
.active_task_pr(&task.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",
task.plan.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 {:?}",
task.plan.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_run_lease(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_run(&pr, None, lease).await,
None => store.settle_task_pr(&pr, None).await,
}
.map_err(|error| task_error(format!("failed to settle Task PR: {error}")))?;
Ok(true)
})
}
async fn record_task_failure(
store: &SharedStore,
task: &mut Task,
_reason: impl Into<String>,
error: String,
) -> OpsResult<()> {
store
.update_task(task)
.await
.map_err(|store_error| task_error(store_error.to_string()))?;
store
.append_task_event(
&task.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,
task: &mut Task,
) -> OpsResult<()> {
relaunch_inactive_process_with_trigger(store, task, None).await
}
pub(crate) async fn resume_inactive_process(store: &SharedStore, task: &mut Task) -> OpsResult<()> {
relaunch_inactive_process_with_trigger(store, task, Some(RunTrigger::User)).await
}
async fn relaunch_inactive_process_with_trigger(
store: &SharedStore,
task: &mut Task,
trigger: Option<RunTrigger>,
) -> OpsResult<()> {
let Some(_) = ensure_working_pr(store, task).await? else {
return Err(task_error(format!(
"Task {} is terminal and cannot start a Run",
task.plan.identifier
)));
};
launch_task_process(store, task, trigger).await
}
const MAX_AUTOMATIC_TASK_RECOVERY_RUNS: usize = 3;
#[derive(Debug, Clone, PartialEq, Eq)]
enum AutomaticTaskRelaunch {
Idle,
Launch(RunTrigger),
Exhausted {
basis: crate::durable::Basis,
recoveries: usize,
error: String,
},
}
async fn task_recovery_chain(store: &SharedStore, latest: Run) -> OpsResult<(usize, Run)> {
let mut seen = HashSet::new();
let mut recoveries = 0;
let mut run = latest;
loop {
if !seen.insert(run.id.clone()) {
return Err(task_error(format!(
"Task recovery chain contains a cycle at Run {}",
run.id
)));
}
let Some(prior_run_id) = run.retry_of.clone() else {
return Ok((recoveries, run));
};
recoveries += 1;
run = store
.run_by_id(&prior_run_id)
.await
.map_err(|error| task_error(format!("failed to read Task recovery chain: {error}")))?;
}
}
fn failed_run_made_durable_progress(events: &[crate::task::TaskEvent]) -> bool {
events
.iter()
.skip(1)
.take_while(|event| !matches!(event.kind, TaskEventKind::Started))
.any(|event| !matches!(event.kind, TaskEventKind::Failed { .. }))
}
async fn plan_automatic_task_relaunch(
store: &SharedStore,
task: &Task,
) -> OpsResult<AutomaticTaskRelaunch> {
let work = store
.work_for_child(&ChildRef::Task(task.id.clone()))
.await
.map_err(|error| task_error(error.to_string()))?;
if !matches!(
store
.work_status(&work)
.await
.map_err(|error| task_error(error.to_string()))?,
WorkStatus::Ready
) {
return Ok(AutomaticTaskRelaunch::Idle);
}
let basis = store
.current_epoch(&work)
.await
.map_err(|error| task_error(error.to_string()))?
.current_basis;
let Some(crate::task::TaskEvent {
kind: TaskEventKind::Failed { error, resumable },
..
}) = store
.latest_task_event(&task.id)
.await
.map_err(|error| task_error(error.to_string()))?
else {
return Ok(AutomaticTaskRelaunch::Launch(RunTrigger::Input { basis }));
};
let Some(latest_run) = store
.latest_run(&work)
.await
.map_err(|error| task_error(error.to_string()))?
else {
return Ok(AutomaticTaskRelaunch::Launch(RunTrigger::Input { basis }));
};
let (recoveries, root) = task_recovery_chain(store, latest_run.clone()).await?;
let recent_events = store
.recent_task_events(&task.id, 16)
.await
.map_err(|store_error| task_error(store_error.to_string()))?;
let durable_progress = failed_run_made_durable_progress(&recent_events);
let input_advanced = match &root.trigger {
RunTrigger::Input { basis: prior } => {
prior.epoch_id != basis.epoch_id || prior.revision < basis.revision
}
_ => false,
};
if durable_progress || input_advanced {
return Ok(AutomaticTaskRelaunch::Launch(RunTrigger::Input { basis }));
}
if resumable && recoveries < MAX_AUTOMATIC_TASK_RECOVERY_RUNS {
return Ok(AutomaticTaskRelaunch::Launch(RunTrigger::Recovery {
prior_run_id: latest_run.id,
}));
}
Ok(AutomaticTaskRelaunch::Exhausted {
basis,
recoveries,
error,
})
}
async fn prepare_automatic_task_relaunch(
store: &SharedStore,
task: &Task,
) -> OpsResult<Option<RunTrigger>> {
match plan_automatic_task_relaunch(store, task).await? {
AutomaticTaskRelaunch::Idle => Ok(None),
AutomaticTaskRelaunch::Launch(trigger) => Ok(Some(trigger)),
AutomaticTaskRelaunch::Exhausted {
basis,
recoveries,
error,
} => {
let work = store
.work_for_child(&ChildRef::Task(task.id.clone()))
.await
.map_err(|store_error| task_error(store_error.to_string()))?;
let (_, lease) = store
.reserve_run(
&work,
RunTrigger::Input {
basis: basis.clone(),
},
)
.await
.map_err(|store_error| {
task_error(format!(
"failed to reserve exhausted Task recovery Run: {store_error}"
))
})?;
store
.advance_run(
&lease,
crate::durable::RunAdvance::Wait {
on: WaitOn::Input {
after: basis.clone(),
},
},
)
.await
.map_err(|store_error| {
task_error(format!(
"failed to park exhausted Task recovery: {store_error}"
))
})?;
tracing::warn!(
task = %task.plan.identifier,
recoveries,
limit = MAX_AUTOMATIC_TASK_RECOVERY_RUNS,
%error,
"automatic Task recovery exhausted; waiting for durable input"
);
Ok(None)
}
}
}
async fn relaunch_inactive_process_automatically(
store: &SharedStore,
task: &mut Task,
) -> OpsResult<()> {
let Some(_) = ensure_working_pr(store, task).await? else {
return Err(task_error(format!(
"Task {} is terminal and cannot start a Run",
task.plan.identifier
)));
};
validate_task_lifecycle(task)?;
let Some(trigger) = prepare_automatic_task_relaunch(store, task).await? else {
return Ok(());
};
launch_task_process(store, task, Some(trigger)).await
}
async fn relaunch_for_ci_incident(
store: &SharedStore,
task: &mut Task,
incident_id: String,
) -> OpsResult<()> {
let Some(_) = ensure_working_pr(store, task).await? else {
return Err(task_error(format!(
"Task {} is terminal and cannot repair CI",
task.plan.identifier
)));
};
launch_task_process(store, task, Some(RunTrigger::CiIncident { incident_id })).await
}
async fn launch_task_process(
store: &SharedStore,
task: &mut Task,
trigger: Option<RunTrigger>,
) -> OpsResult<()> {
let work = store
.work_for_child(&ChildRef::Task(task.id.clone()))
.await
.map_err(|error| task_error(error.to_string()))?;
if store
.current_run(&work)
.await
.map_err(|error| task_error(error.to_string()))?
.is_some()
{
return Ok(());
}
validate_task_lifecycle(task)?;
let trigger = match trigger {
Some(trigger) => trigger,
None => RunTrigger::Input {
basis: store
.current_epoch(&work)
.await
.map_err(|error| task_error(error.to_string()))?
.current_basis,
},
};
let (run, lease) = store
.reserve_run(&work, trigger)
.await
.map_err(|error| task_error(format!("failed to reserve Task Run: {error}")))?;
let tmux_name = format!(
"lf-task-{}-{}-{}",
tmux_session_slug(&task.plan.identifier),
&task.id.as_str()[3..11],
&run.id.as_str()[4..12]
);
store
.update_task_for_run(task, &lease)
.await
.map_err(|error| task_error(error.to_string()))?;
crate::ops::launch_in_run(
store,
&lease,
crate::ops::RunLaunch {
work: WorkRef::Task(task.id.clone()),
wave_id: task.wave_id.clone(),
cwd: task.worktree.clone(),
tmux_name,
agent: task.agent.clone(),
account_id: None,
resume_token: task.provider_session_id.clone(),
},
)
.await
.map(|_| ())
.map_err(|error| task_error(error.to_string()))
}
async fn wait_until_running(store: &SharedStore, task_id: &crate::task::TaskId) -> OpsResult<Task> {
let deadline = tokio::time::Instant::now() + super::child::CHILD_STARTUP_GRACE;
loop {
let task = store
.get_task(task_id)
.await
.map_err(|error| task_error(format!("failed to observe task startup: {error}")))?
.ok_or_else(|| task_error("task task disappeared during startup"))?;
match task_work_status(store, &task).await? {
WorkStatus::Running { .. } => return Ok(task),
WorkStatus::Done | WorkStatus::Abandoned => {
return Err(task_error(format!(
"task {} ended during startup",
task.plan.identifier
)))
}
WorkStatus::Ready | WorkStatus::Waiting { .. } => {}
}
if tokio::time::Instant::now() >= deadline {
return Err(task_error(format!(
"task {} process did not report running within 10 seconds",
task.plan.identifier
)));
}
tokio::time::sleep(Duration::from_millis(100)).await;
}
}
pub(crate) async fn reconcile_process_liveness(
store: &SharedStore,
task: &mut Task,
) -> OpsResult<()> {
let work = store
.work_for_child(&ChildRef::Task(task.id.clone()))
.await
.map_err(|error| task_error(error.to_string()))?;
let Some(run) = store
.current_run(&work)
.await
.map_err(|error| task_error(error.to_string()))?
else {
return Ok(());
};
if run.state == RunState::Reserved {
let still_starting =
run.created_at + time::Duration::seconds(10) > time::OffsetDateTime::now_utc();
if still_starting {
return Ok(());
}
}
if let Some(containment) = &run.containment {
let alive = match containment {
Containment::Tmux { name } => tmux_session_exists(name)
.await
.map_err(|error| task_error(error.to_string()))?,
Containment::ProcessGroup { .. } => true,
};
if alive {
return Ok(());
}
if run.started_at.is_some_and(|started_at| {
started_at + time::Duration::seconds(10) > time::OffsetDateTime::now_utc()
}) {
return Ok(());
}
}
store
.recover_run(&run.id, ContainmentObservation::Absent)
.await
.map_err(|error| task_error(error.to_string()))?;
mark_task_body_lost(store, task).await
}
async fn mark_task_body_lost(store: &SharedStore, task: &mut Task) -> OpsResult<()> {
let active = store
.active_task_pr(&task.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 && pr.merge_request().is_some())
{
return Ok(());
}
let reason = "task process is missing; Loopflow will recover this Task";
record_task_failure(store, task, reason, reason.to_string()).await?;
if let Err(error) = task_recovery_adoption(store, task).await {
tracing::info!(
task = %task.plan.identifier,
"not recovering missing Task body: {error}"
);
return Ok(());
}
relaunch_inactive_process_automatically(store, task).await
}
pub(crate) async fn reconcile_project_tasks(
store: &SharedStore,
project: &crate::project::Project,
) -> OpsResult<Vec<Task>> {
let project_tasks = |tasks: Vec<Task>| {
tasks
.into_iter()
.filter(|task| task.project_id == project.id)
.collect::<Vec<_>>()
};
let mut tasks = project_tasks(
store
.list_tasks(Some(&project.wave_id))
.await
.map_err(|error| task_error(format!("failed to list supervised Tasks: {error}")))?,
);
for task in &mut tasks {
if matches!(
task_work_status(store, task).await?,
WorkStatus::Done | WorkStatus::Abandoned
) {
continue;
}
if let Err(error) = task_recovery_adoption(store, task).await {
tracing::warn!(
task = %task.plan.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 matches!(
task_work_status(store, task).await?,
WorkStatus::Done | WorkStatus::Abandoned
) {
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.after_merge() == AfterMerge::CompleteTask);
if settled && !completing {
ensure_working_pr(store, task).await?;
if !matches!(
task_work_status(store, task).await?,
WorkStatus::Running { .. }
) {
relaunch_inactive_process_automatically(store, task).await?;
}
} else {
let Some(pr) = observed.as_ref() else {
continue;
};
route_ci_incident(store, task, pr).await?;
if pr.merge_request().is_none()
&& !matches!(
task_work_status(store, task).await?,
WorkStatus::Running { .. }
)
{
relaunch_inactive_process_automatically(store, task).await?;
}
}
}
let refreshed = store
.list_tasks(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::Project,
) -> OpsResult<usize> {
reconcile_project_tasks(store, project).await?;
Ok(0)
}
pub(crate) async fn reconcile_task_pr(
store: &SharedStore,
task: &mut Task,
) -> OpsResult<Option<TaskPr>> {
reconcile_task_pr_with_authority(store, task, None, crate::ops::pr::PrReadFreshness::Cached)
.await
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum OpenPrDisposition {
ObservationDegraded,
NeedsDirection,
}
pub(crate) fn decide_open_pr_status(
pr: &TaskPr,
github_degraded: Option<&str>,
head_advanced: bool,
) -> (Option<OpenPrDisposition>, String) {
let number = pr
.github()
.expect("open Task PR requires a GitHub PR record")
.number;
if let Some(reason) = github_degraded {
return (
Some(OpenPrDisposition::ObservationDegraded),
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 (
Some(OpenPrDisposition::NeedsDirection),
format!(
"CI failing on pull request #{number}; the Task body did not repair the head. Needs a new directive; pull request #{number} stays attached."
),
);
}
let reason = match pr.merge_request() {
Some(request) if request.mode == PrMergeMode::User => {
let short = request.head_sha.chars().take(12).collect::<String>();
format!("pull request #{number} awaits the user's explicit merge of head {short}")
}
Some(request) => {
let short = request.head_sha.chars().take(12).collect::<String>();
format!("pull request #{number} awaits GitHub auto-merge of head {short}")
}
None => format!("pull request #{number} is published; no merge was requested"),
};
(None, reason)
}
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,
task: &Task,
pr: &TaskPr,
) -> OpsResult<()> {
let Some(incident) = current_ci_incident(pr) else {
return Ok(());
};
match task_work_status(store, task).await? {
WorkStatus::Ready => {
let mut task = task.clone();
relaunch_for_ci_incident(store, &mut task, incident.identity.clone()).await?;
}
WorkStatus::Running { .. } => {}
WorkStatus::Waiting { .. } | WorkStatus::Done | WorkStatus::Abandoned => return Ok(()),
}
let work = store
.work_for_child(&ChildRef::Task(task.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_run(
store: &SharedStore,
task: &mut Task,
lease: &RunLease,
) -> OpsResult<Option<TaskPr>> {
reconcile_task_pr_with_authority(
store,
task,
Some(lease),
crate::ops::pr::PrReadFreshness::Cached,
)
.await
}
pub(crate) async fn reconcile_task_pr_fresh_for_run(
store: &SharedStore,
task: &mut Task,
lease: &RunLease,
) -> OpsResult<Option<TaskPr>> {
reconcile_task_pr_with_authority(
store,
task,
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 github = pr.github()?;
let repo = github_repo_slug_from_pr_url(&github.url)?;
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:{}:{}:{}:{}",
repo,
github.number,
observation.head_sha,
hex::encode(digest.finalize())
),
task_id: pr.task_id.clone(),
pr_id: pr.id.clone(),
repo,
pr_number: github.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,
})
}
fn github_repo_slug_from_pr_url(url: &str) -> Option<String> {
let path = url.split("github.com/").nth(1)?;
let mut parts = path.split('/');
let owner = parts.next().filter(|part| !part.is_empty())?;
let repo = parts.next().filter(|part| !part.is_empty())?;
Some(format!("{owner}/{repo}"))
}
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, task: &Task) -> OpsResult<Option<TaskPr>> {
if let Some(active) = store
.active_task_pr(&task.id)
.await
.map_err(|error| task_error(format!("failed to read active PR: {error}")))?
{
return Ok(Some(active));
}
let prs = store
.task_prs(&task.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,
task: &mut Task,
lease: Option<&RunLease>,
freshness: crate::ops::pr::PrReadFreshness,
) -> OpsResult<Option<TaskPr>> {
let _mutation = lock_task_pr_mutation(&task.worktree)?;
let Some(mut pr) = reconcile_subject(store, task).await? else {
return Ok(None);
};
let Some(number) = pr.github().map(|github| github.number) else {
task.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) {
task.observation = observation;
return Ok(Some(pr));
}
}
let previous = pr.clone();
let github_pr =
match crate::ops::pr::observe_pr_by_number(&task.worktree, number, &pr.branch, freshness) {
crate::ops::pr::PrObservation::Fresh(info) => {
pr.github_observation = Some(GithubObservation {
checked_at: now,
result: GithubObservationResult::Fresh,
});
task.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()))?;
task.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()))?;
task.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_gate_proposal = task.gate_proposal.clone();
let previous_pm_writeback = task.pm_writeback.clone();
let publication = pr.publication.get_or_insert(PrPublication {
requested_at: now,
github: None,
merge: None,
});
invalidate_stale_merge_request(&task.worktree, publication, &github_pr)?;
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.after_merge() == AfterMerge::CompleteTask
&& matches!(
committed_follow_up_range(&task.worktree, &pr)?,
CommittedFollowUp::ProvenEmpty
);
if completes {
let proposal = crate::task::TaskGateProposal {
done: true,
reason: format!(
"pull request #{} merged and completed the Task",
github_pr.number
),
};
match task.lifecycle_phase {
crate::task::TaskLifecyclePhase::First => {
task.enter_loop()
.map_err(|error| task_error(error.to_string()))?;
task.enter_finally(proposal)
.map_err(|error| task_error(error.to_string()))?;
}
crate::task::TaskLifecyclePhase::Loop => {
task.enter_finally(proposal)
.map_err(|error| task_error(error.to_string()))?;
}
crate::task::TaskLifecyclePhase::Finally => {
task.gate_proposal = Some(proposal);
task.updated_at = now;
}
}
reconcile_pm_writeback(store, task, Some(&url)).await;
}
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;
None
}
_ => {
pr.abandoned_at = None;
if let Some(ci_observation) = observe_required_checks(
&task.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;
if pr_changed {
pr.updated_at = now;
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 task.gate_proposal != previous_gate_proposal || task.pm_writeback != previous_pm_writeback {
update_task_with_authority(store, task, 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, &task.id, &event, 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.next_slug().map(str::to_string))
.unwrap_or_else(|| (settled.sequence + 1).to_string())
}
fn deterministic_next_branch(
task: &Task,
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}", task.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,
task: &Task,
) -> OpsResult<TaskRecoveryAdoption> {
let worktree = &task.worktree;
let identifier = &task.plan.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(&task.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(&task.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 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(task, &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, task: &Task) -> OpsResult<()> {
if store
.active_task_pr(&task.id)
.await
.map_err(|error| task_error(format!("failed to read active PR: {error}")))?
.is_some()
{
return Ok(());
}
if is_clean(&task.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",
task.plan.identifier,
task.worktree.display()
)))
}
pub(crate) async fn ensure_working_pr(
store: &SharedStore,
task: &mut Task,
) -> OpsResult<Option<TaskPr>> {
ensure_working_pr_with_authority(store, task, None, RotateOptions::runner()).await
}
pub(crate) async fn ensure_working_pr_for_run(
store: &SharedStore,
task: &mut Task,
lease: &RunLease,
) -> OpsResult<Option<TaskPr>> {
ensure_working_pr_with_authority(store, task, 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,
task: &Task,
pr: TaskPr,
lease: Option<&RunLease>,
) -> OpsResult<TaskPr> {
if pr.phase() != PrPhase::Working || !task.worktree.exists() {
return Ok(pr);
}
if is_ancestor(&task.worktree, &pr.base_commit, &pr.branch).unwrap_or(false) {
return Ok(pr);
}
let Ok(default_branch) = get_default_branch(&task.worktree) else {
return Ok(pr);
};
let Ok((base_ref, _)) = resolve_upstream_base(&task.worktree, &default_branch) else {
return Ok(pr);
};
let Ok(fork) = fork_point(&task.worktree, &base_ref, &pr.branch) else {
tracing::warn!(
task = %task.plan.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(&task.worktree, commit, descendant).unwrap_or(false)
};
if !(ancestry(&fork, &pr.base_commit) && ancestry(&pr.base_commit, &base_ref)) {
tracing::warn!(
task = %task.plan.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 = %task.plan.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_run(&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,
task: &mut Task,
lease: Option<&RunLease>,
rotate: RotateOptions,
) -> OpsResult<Option<TaskPr>> {
reconcile_task_pr_with_authority(store, task, lease, crate::ops::pr::PrReadFreshness::Cached)
.await?;
if matches!(
task_work_status(store, task).await?,
WorkStatus::Done | WorkStatus::Abandoned
) {
return Ok(None);
}
if let Some(active) = store
.active_task_pr(&task.id)
.await
.map_err(|error| task_error(format!("failed to read active PR: {error}")))?
{
return Ok(Some(
heal_incoherent_base(store, task, active, lease).await?,
));
}
let prs = store
.task_prs(&task.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 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, .. } = &task.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(&task.worktree, &settled)?;
if settled.after_merge() == AfterMerge::CompleteTask
&& !matches!(&committed_carry, CommittedFollowUp::Range { .. })
{
return Ok(None);
}
if let Some(lease) = lease {
store
.validate_run_lease(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(task, &settled, rotate.slug_override.as_deref())?;
let default_branch = get_default_branch(&task.worktree)
.map_err(|error| task_error(format!("failed to resolve default branch: {error}")))?;
let (base_ref, _) = resolve_upstream_base(&task.worktree, &default_branch)?;
if !rotate.carry_dirty
&& !is_clean(&task.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",
task.plan.identifier,
task.worktree.display()
)));
}
let current = current_branch(&task.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 {:?}",
task.plan.identifier,
settled.branch,
branch,
task.worktree.display(),
current
)));
}
let local_ref = format!("refs/heads/{branch}");
let remote_ref = format!("refs/remotes/origin/{branch}");
let collision = ref_exists(&task.worktree, &local_ref)
.map_err(|error| task_error(format!("failed to inspect branch collision: {error}")))?
|| ref_exists(&task.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(&task.worktree)
.map_err(|error| task_error(format!("failed to stash follow-up edits: {error}")))?;
if let Err(error) = checkout_new_branch_from(&task.worktree, &branch, &base_ref) {
let recovered = current_branch(&task.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(&task.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(&task.worktree, from, to) {
roll_back_failed_rotation(&task.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(&task.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",
task.worktree.display()
))
})?;
}
}
let base_commit = fork_point(&task.worktree, &base_ref, &branch)?;
let _mutation = lock_task_pr_mutation(&task.worktree)?;
push_with_upstream(&task.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_id: task.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,
&task.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(&task.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 task, lease) = task_for_worktree(&store, &repo)
.await?
.ok_or_else(|| task_error("no Task owns this worktree"))?;
reconcile_task_pr_with_authority(
&store,
&mut task,
lease.as_ref(),
crate::ops::pr::PrReadFreshness::Cached,
)
.await?;
if let Some(active) = store
.active_task_pr(&task.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 matches!(
task_work_status(&store, &task).await?,
WorkStatus::Done | WorkStatus::Abandoned
) {
return Err(task_error(format!(
"Task {} is terminal; nothing to rotate",
task.plan.identifier
)));
}
let rotate = RotateOptions {
carry_dirty: true,
slug_override,
};
ensure_working_pr_with_authority(&store, &mut task, 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<Task> {
block_on_task(async move {
let store = task_store().await?;
let mut task = store
.get_task_by_issue(issue)
.await
.map_err(|error| task_error(format!("failed to read task status: {error}")))?
.ok_or_else(|| task_error(format!("no Task exists for {issue:?}")))?;
reconcile_task_pr(&store, &mut task).await?;
reconcile_process_liveness(&store, &mut task).await?;
reconcile_task_completion(&store, &mut task, None).await?;
Ok(task)
})
}
pub(crate) fn find_discardable_final_task(repo: &Path) -> OpsResult<Option<String>> {
let repo = repo.to_path_buf();
block_on_task(async move {
let TaskAuthority::Authority { store, task, .. } = resolve_task_authority(&repo).await?
else {
return Ok(None);
};
if task.lifecycle_phase != crate::task::TaskLifecyclePhase::Finally {
return Ok(None);
}
let materially_clean = is_materially_clean(&task.worktree)
.map_err(|error| task_error(format!("failed to inspect Task worktree: {error}")))?;
if !materially_clean {
return Ok(None);
}
let gate = task_completion_gate(&store, &task).await?;
Ok((gate.satisfied && gate.discardable_successor.is_some())
.then(|| task.plan.identifier.clone()))
})
}
pub fn task_complete(issue: &str, summary: String) -> OpsResult<Task> {
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 task = store
.get_task_by_issue(issue)
.await
.map_err(|error| task_error(format!("failed to read Task: {error}")))?
.ok_or_else(|| task_error(format!("no Task exists for {issue:?}")))?;
reconcile_task_pr(&store, &mut task).await?;
let lease = ambient_task_run_lease(&store, &task).await?;
if let Some(lease) = lease.as_ref() {
store
.validate_run_lease(lease)
.await
.map_err(|error| task_error(format!("Task body lost write authority: {error}")))?;
}
let work = store
.work_for_child(&ChildRef::Task(task.id.clone()))
.await
.map_err(|error| task_error(error.to_string()))?;
match store
.work_status(&work)
.await
.map_err(|error| task_error(error.to_string()))?
{
WorkStatus::Done => return Ok(task),
WorkStatus::Abandoned => {
return Err(task_error(format!(
"Task {} is abandoned and cannot be completed",
task.plan.identifier
)))
}
WorkStatus::Ready | WorkStatus::Running { .. } | WorkStatus::Waiting { .. } => {}
}
if !is_clean(&task.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, &task).await?;
if let Some(refusal) = gate.refusal(&task.plan.identifier) {
return Err(task_error(refusal));
}
propose_task_done(&mut task, summary.clone())?;
complete_task_with_authority(
&store,
&task,
gate.discardable_successor.as_ref(),
lease.as_ref(),
)
.await
.map_err(|error| task_error(format!("failed to complete Task: {error}")))?;
if lease.is_none() {
reconcile_pm_writeback(&store, &mut task, None).await;
store
.update_task(&task)
.await
.map_err(|error| task_error(error.to_string()))?;
}
Ok(task)
})
}
fn propose_task_done(task: &mut Task, summary: String) -> OpsResult<()> {
let proposal = crate::task::TaskGateProposal {
done: true,
reason: summary,
};
match task.lifecycle_phase {
crate::task::TaskLifecyclePhase::First => {
task.enter_loop()
.map_err(|error| task_error(error.to_string()))?;
task.enter_finally(proposal)
.map_err(|error| task_error(error.to_string()))?;
}
crate::task::TaskLifecyclePhase::Loop => {
task.enter_finally(proposal)
.map_err(|error| task_error(error.to_string()))?;
}
crate::task::TaskLifecyclePhase::Finally => {
task.gate_proposal = Some(proposal);
task.updated_at = time::OffsetDateTime::now_utc();
}
}
Ok(())
}
fn pr_link_state_label(pr: &TaskPr) -> String {
match pr.phase() {
PrPhase::Merged => "Merged".to_string(),
PrPhase::Abandoned => "Abandoned".to_string(),
_ => {
let completes = pr.after_merge() == AfterMerge::CompleteTask;
if completes {
"Open · completes task on merge".to_string()
} else if let Some(request) = pr.merge_request() {
match request.mode {
PrMergeMode::User => "Open · user merge requested".to_string(),
PrMergeMode::Auto => "Open · auto-merge requested".to_string(),
}
} else {
"Open · published".to_string()
}
}
}
}
async fn link_pr_to_linear(store: &SharedStore, task: &Task, 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, task).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: task.plan.id.as_str().to_string(),
url: github.url.clone(),
title,
subtitle: state,
body,
};
let outcome =
crate::ops::pm::pm_link_pr_async(&task.worktree, wave.name(), &request, &prior).await;
if let Some(error) = &outcome.error {
tracing::warn!(
issue = task.plan.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,
task: &mut Task,
pr_url: Option<&str>,
) {
let Ok(wave) = owning_wave(store, task).await else {
task.pm_writeback = PmWritebackState::Pending {
operation: PmWritebackOperation::CompleteTask,
error: format!("owning Wave {} is not registered", task.wave_id),
};
return;
};
task.pm_writeback = writeback_state(
crate::ops::task_pm::complete_task(
&task.worktree,
wave.name(),
task.plan.id.as_str(),
pr_url,
)
.await,
);
}
async fn retry_pm_writeback(store: &SharedStore, task: &mut Task) {
let Ok(prs) = store.task_prs(&task.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, task).await else {
task.pm_writeback = PmWritebackState::Pending {
operation: PmWritebackOperation::CompleteTask,
error: format!("owning Wave {} is not registered", task.wave_id),
};
return;
};
task.pm_writeback = {
let operation = match &task.pm_writeback {
PmWritebackState::Pending { operation, .. } => *operation,
PmWritebackState::Current => PmWritebackOperation::CompleteTask,
};
let result = crate::ops::task_pm::retry_complete_task(
&task.worktree,
wave.name(),
task.plan.id.as_str(),
pr_url,
)
.await;
writeback_state_for(operation, result)
};
task.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()
)
})
}
}
pub(crate) async fn task_completion_gate(
store: &SharedStore,
task: &Task,
) -> OpsResult<CompletionGate> {
let mut gate = CompletionGate {
satisfied: true,
blockers: Vec::new(),
discardable_successor: None,
};
let work_done = task_work_status(store, task).await? == WorkStatus::Done;
let prs = store
.task_prs(&task.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.after_merge() == AfterMerge::CompleteTask {
let number = newest
.github()
.map(|github| github.number)
.unwrap_or_default();
match committed_follow_up_range(&task.worktree, newest)? {
CommittedFollowUp::ProvenEmpty => {}
CommittedFollowUp::Range { .. } => gate.blockers.push(format!(
"follow-up work is committed past merged pull request #{number}"
)),
CommittedFollowUp::Unprovable { .. } if work_done => {}
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(&task.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; 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(&task.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, task: &Task) -> OpsResult<Option<TaskPr>> {
let prs = store
.task_prs(&task.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.after_merge() == AfterMerge::CompleteTask))
}
async fn advance_completion_after_gate(
store: &SharedStore,
task: &mut Task,
lease: Option<&RunLease>,
) -> OpsResult<bool> {
let work = store
.work_for_child(&ChildRef::Task(task.id.clone()))
.await
.map_err(|error| task_error(error.to_string()))?;
match store
.work_status(&work)
.await
.map_err(|error| task_error(error.to_string()))?
{
WorkStatus::Done | WorkStatus::Abandoned => return Ok(false),
WorkStatus::Running { .. } if lease.is_none() => return Ok(false),
WorkStatus::Ready | WorkStatus::Waiting { .. } | WorkStatus::Running { .. } => {}
}
let Some(pr) = merged_completing_pr(store, task).await? else {
return Ok(false);
};
let gate = task_completion_gate(store, task).await?;
if !gate.satisfied {
return Ok(false);
}
if gate.discardable_successor.is_some() {
return Ok(false);
}
let url = pr.github().map(|github| github.url.clone());
let summary = format!(
"pull request #{} merged and completed the Task",
pr.github().map(|github| github.number).unwrap_or_default()
);
propose_task_done(task, summary.clone())?;
complete_task_after_pr_with_authority(store, task, &pr, lease)
.await
.map_err(|error| task_error(error.to_string()))?;
if lease.is_none() {
reconcile_pm_writeback(store, task, url.as_deref()).await;
store
.update_task(task)
.await
.map_err(|error| task_error(error.to_string()))?;
}
Ok(true)
}
pub(crate) async fn reconcile_task_completion(
store: &SharedStore,
task: &mut Task,
lease: Option<&RunLease>,
) -> OpsResult<()> {
let status = task_work_status(store, task).await?;
if status == WorkStatus::Done && matches!(task.pm_writeback, PmWritebackState::Pending { .. }) {
retry_pm_writeback(store, task).await;
if let Some(lease) = lease {
store
.update_task_for_run(task, lease)
.await
.map_err(|error| task_error(error.to_string()))?;
} else {
store
.update_task(task)
.await
.map_err(|error| task_error(error.to_string()))?;
}
return Ok(());
}
if !matches!(status, WorkStatus::Done | WorkStatus::Abandoned) {
advance_completion_after_gate(store, task, lease).await?;
}
Ok(())
}
pub fn task_snapshot(task: &Task) -> OpsResult<TaskSnapshot> {
let task = task.clone();
block_on_task(async move {
let store = task_store().await?;
let wave = owning_wave(&store, &task).await?;
let project = store
.get_project(&task.project_id)
.await
.map_err(|error| task_error(format!("failed to read owning Project: {error}")))?
.ok_or_else(|| task_error(format!("owning Project {} is missing", task.project_id)))?;
let work = store
.work_for_child(&ChildRef::Task(task.id.clone()))
.await
.map_err(|error| task_error(format!("failed to resolve Task Work: {error}")))?;
let run = store
.current_run(&work)
.await
.map_err(|error| task_error(error.to_string()))?;
let invocation = match &run {
Some(run) => store
.open_invocation_for_run(&run.id)
.await
.map_err(|error| task_error(error.to_string()))?,
None => None,
};
let process_alive = match run.as_ref().and_then(|run| run.containment.as_ref()) {
Some(Containment::Tmux { name }) => tmux_session_exists(name)
.await
.map_err(|error| task_error(error.to_string()))?,
Some(Containment::ProcessGroup { .. }) => true,
None => false,
};
let latest_event = store
.task_events_after(&task.id, 0)
.await
.map_err(|error| task_error(format!("failed to read task events: {error}")))?
.into_iter()
.last();
let prs = store
.task_prs(&task.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 completion_gate = task_completion_gate(&store, &task).await?;
let completion_refusal = completion_gate.refusal(&task.plan.identifier);
let resume_refusal = no_active_pr_resume_refusal(&task.plan.identifier, active, latest);
let work_status = store
.work_status(&work)
.await
.map_err(|error| task_error(format!("failed to derive Task Work status: {error}")))?;
let action_evidence = TaskActionEvidence {
status: work_status.clone(),
latest_pr_phase: latest.map(|pr| pr.phase()),
latest_pr_after_merge: latest
.filter(|pr| pr.phase() == PrPhase::Merged)
.map(TaskPr::after_merge),
latest_pr_merge_request: latest.and_then(TaskPr::merge_request),
completion_refusal: completion_refusal.as_deref(),
resume_refusal: resume_refusal.as_deref(),
ci: active.and_then(|pr| pr.fresh_ci()),
process_alive: if matches!(work_status, WorkStatus::Running { .. }) {
Some(process_alive)
} else {
None
},
predecessor_phase,
abandon_intent: task.abandon_intent.is_some(),
local_progress_unsettled: None,
};
let actions = derive_task_actions(&action_evidence);
Ok(TaskSnapshot {
issue_id: task.plan.id.as_str().to_string(),
issue_identifier: task.plan.identifier,
task_id: task.id.to_string(),
external_project_id: project.plan.id.as_str().to_string(),
project: project.plan.slug,
pm_snapshot_synced_at: task.plan.pm_snapshot_synced_at,
pm_writeback: task.pm_writeback,
wave: wave.name().to_string(),
project_id: task.project_id.to_string(),
status: work_status,
worktree: task.worktree.display().to_string(),
workspace_slug: task.workspace_slug,
lifecycle: task.lifecycle,
lifecycle_phase: task.lifecycle_phase,
phase_epoch: task.phase_epoch,
phase_cursor: task.phase_cursor,
phase_iteration: task.phase_iteration,
gate_cycle: task.gate_cycle,
gate_proposal: task.gate_proposal,
prs,
active_pr,
agent: task.agent,
provider: task.provider,
provider_session_id: task.provider_session_id,
process_alive,
invocation,
latest_event,
created_at: task.created_at,
updated_at: task.updated_at,
observation: task.observation,
actions,
})
})
}
pub fn task_changes(issue: &str) -> OpsResult<TaskChangesSnapshot> {
let task = task_status(issue)?;
let pr = active_pr(&task)?;
changes_snapshot(TaskWorkspace::new(&task, &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(),
task_id: workspace.task_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 task = task_status(issue)?;
let pr = active_pr(&task)?;
diff_snapshot(TaskWorkspace::new(&task, &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(),
task_id: workspace.task_id.to_string(),
path: relative,
patch,
binary,
truncated,
})
}
pub fn task_file(issue: &str, path: &str) -> OpsResult<TaskFileSnapshot> {
let task = task_status(issue)?;
let pr = active_pr(&task)?;
file_snapshot(TaskWorkspace::new(&task, &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(),
task_id: workspace.task_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 task = store
.get_task_by_issue(issue)
.await
.map_err(|error| task_error(format!("failed to resolve task: {error}")))?
.ok_or_else(|| task_error(format!("no Task exists for {issue:?}")))?;
reconcile_task_pr(&store, &mut task).await?;
reconcile_process_liveness(&store, &mut task).await?;
let receipt =
super::child::append_steer(&store, ChildRef::Task(task.id.clone()), &message).await?;
if !matches!(
task_work_status(&store, &task).await?,
WorkStatus::Running { .. }
) {
relaunch_inactive_process(&store, &mut task).await?;
}
Ok(TaskControlResult {
issue_id: task.plan.identifier.clone(),
task_id: task.id.to_string(),
receipt: super::child::WorkControlReceipt::Steer { receipt },
observation: task.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 task = store
.get_task_by_issue(issue)
.await
.map_err(|error| task_error(format!("failed to resolve task: {error}")))?
.ok_or_else(|| task_error(format!("no Task exists for {issue:?}")))?;
reconcile_task_pr(&store, &mut task).await?;
reconcile_process_liveness(&store, &mut task).await?;
let work = store
.work_for_child(&ChildRef::Task(task.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: task.plan.identifier.clone(),
task_id: task.id.to_string(),
receipt: super::child::WorkControlReceipt::Interrupt { receipt },
observation: task.observation,
})
})
}
pub fn task_recover(issue: &str, reason: Option<String>) -> OpsResult<Task> {
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<Task> {
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_by_issue(issue)
.await
.map_err(|error| task_error(format!("failed to resolve task: {error}")))?
.ok_or_else(|| task_error(format!("no Task exists for {issue:?}")))?;
match task_work_status(store, &predecessor).await? {
WorkStatus::Done => {
return Err(task_error(format!(
"Task {} is completed; start a new Task rather than recovering it",
predecessor.plan.identifier
)));
}
WorkStatus::Abandoned => {}
WorkStatus::Ready | WorkStatus::Running { .. } | WorkStatus::Waiting { .. } => {
return Ok(predecessor)
}
}
task_recovery_adoption(store, &predecessor)
.await
.map_err(|error| task_error(format!("validate Task recovery: {error}")))?;
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.plan.identifier, predecessor.plan.title
);
}
let now = time::OffsetDateTime::now_utc();
let _reason = reason;
let mut task = predecessor;
task.lifecycle_phase = crate::task::TaskLifecyclePhase::First;
task.phase_epoch += 1;
task.phase_cursor = 0;
task.phase_iteration = 0;
task.gate_cycle = 0;
task.gate_proposal = None;
task.provider_session_id = None;
task.abandon_intent = None;
task.updated_at = now;
task.observation = Observation::NotRequired;
store
.reopen_task(&task, None, crate::durable::Author::User, &carried)
.await
.map_err(|error| task_error(format!("failed to recover Task: {error}")))?;
Ok(task)
}
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 task = store
.get_task_by_issue(issue)
.await
.map_err(|error| task_error(format!("failed to resolve task: {error}")))?
.ok_or_else(|| task_error(format!("no Task exists for {issue:?}")))?;
task_recovery_adoption(&store, &task).await?;
reconcile_task_pr(&store, &mut task).await?;
let prs = store
.task_prs(&task.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(&task.plan.identifier, active, latest) {
return Err(task_error(refusal));
}
{
let _mutation = lock_task_pr_mutation(&task.worktree)?;
clear_task_pr_merge(&store, &task, None, &task.worktree, true).await?;
}
refuse_dirty_between_prs(&store, &task).await?;
reconcile_process_liveness(&store, &mut task).await?;
let issue_id = task.plan.identifier.clone();
let observation = task.observation.clone();
let task_id = task.id.to_string();
let run = super::child::resume_child(
&store,
super::child::Child::Task(Box::new(task)),
model,
reason,
)
.await?;
Ok(TaskControlResult {
issue_id,
task_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 task = store
.get_task_by_issue(issue)
.await
.map_err(|error| task_error(format!("failed to resolve task: {error}")))?
.ok_or_else(|| task_error(format!("no Task exists for {issue:?}")))?;
reconcile_task_pr(&store, &mut task).await?;
reconcile_process_liveness(&store, &mut task).await?;
let work = store
.work_for_child(&ChildRef::Task(task.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()))?;
Ok(TaskControlResult {
issue_id: task.plan.identifier.clone(),
task_id: task.id.to_string(),
receipt: super::child::WorkControlReceipt::Abandon { receipt },
observation: task.observation,
})
})
}
pub fn task_wait(issue: &str, until: TaskWaitUntil, timeout: Option<Duration>) -> OpsResult<Task> {
let started = Instant::now();
loop {
let task = task_status(issue)?;
let status = block_on_task(async {
let store = task_store().await?;
task_work_status(&store, &task).await
})?;
let reached = match until {
TaskWaitUntil::Open => {
matches!(status, WorkStatus::Done | WorkStatus::Abandoned)
|| active_pr(&task).is_ok_and(|pr| pr.phase() == PrPhase::Open)
}
TaskWaitUntil::Terminal => matches!(status, WorkStatus::Done | WorkStatus::Abandoned),
};
if reached || timeout.is_some_and(|limit| started.elapsed() >= limit) {
return Ok(task);
}
std::thread::sleep(Duration::from_secs(1));
}
}
#[cfg(test)]
mod tests {
use std::ffi::OsString;
#[cfg(unix)]
use std::os::unix::fs::PermissionsExt;
use super::{
ensure_task_flow_override, launch_task_process, lock_task_pr_mutation,
prepare_automatic_task_relaunch, resolve_task_lifecycle, resolve_task_start_input,
TaskFlowOverrides, MAX_AUTOMATIC_TASK_RECOVERY_RUNS,
};
use crate::child::ChildRef;
use crate::durable::{RunTrigger, WorkRef, WorkStatus};
use crate::planning::{LinearIssueId, LinearProjectId, ProjectPlan, TaskPlan};
use crate::pm::ProjectFlowPlan;
use crate::project::{Project, ProjectId};
use crate::store::{open_store, SharedStore, StorageConfig};
use crate::task::{
Observation, PmWritebackState, Task, TaskEventKind, TaskId, TaskLifecyclePhase,
TaskLifecyclePlan, TaskPr, TaskPrId,
};
use crate::wave::Wave;
struct TaskFixture {
_database: tempfile::TempDir,
database_path: std::path::PathBuf,
store: SharedStore,
task: Task,
work: WorkRef,
}
struct EnvRestore(Vec<(&'static str, Option<OsString>)>);
impl EnvRestore {
fn capture(names: &[&'static str]) -> Self {
Self(
names
.iter()
.map(|name| (*name, std::env::var_os(name)))
.collect(),
)
}
}
impl Drop for EnvRestore {
fn drop(&mut self) {
for (name, value) in &self.0 {
match value {
Some(value) => std::env::set_var(name, value),
None => std::env::remove_var(name),
}
}
}
}
async fn task_fixture(identifier: &str, loop_flow: &str) -> TaskFixture {
let repository =
std::fs::canonicalize(std::path::Path::new(env!("CARGO_MANIFEST_DIR")).join("../.."))
.unwrap();
let database = tempfile::tempdir().unwrap();
let database_path = database.path().join("registry.db");
let store = std::sync::Arc::new(
open_store(&StorageConfig::sqlite(database_path.clone()))
.await
.unwrap(),
);
let now = time::OffsetDateTime::now_utc();
let wave = Wave::new(
crate::id::WaveId::new(),
"task-recovery".to_string(),
repository.display().to_string(),
);
let project = Project {
id: ProjectId::new(),
plan: ProjectPlan {
id: LinearProjectId::new("task-recovery-project").unwrap(),
slug: "task-recovery".to_string(),
name: "Task recovery".to_string(),
prompt_context: "Keep automatic Task recovery bounded.".to_string(),
pm_snapshot_synced_at: now.unix_timestamp(),
},
wave_id: wave.id().clone(),
iteration: 0,
observation_cursor: 0,
last_state_fingerprint: None,
agent: "codex".to_string(),
provider: "codex".to_string(),
provider_session_id: None,
abandon_intent: None,
created_at: now,
updated_at: now,
};
let task = Task {
id: TaskId::new(),
plan: TaskPlan {
id: LinearIssueId::new(format!("{identifier}-issue")).unwrap(),
identifier: identifier.to_string(),
title: "Task recovery fixture".to_string(),
description: String::new(),
pm_snapshot_synced_at: now.unix_timestamp(),
},
pm_writeback: PmWritebackState::Current,
wave_id: wave.id().clone(),
project_id: project.id.clone(),
worktree: repository,
workspace_slug: "task-recovery-fixture".to_string(),
lifecycle: TaskLifecyclePlan::standard("task-design", loop_flow, "ship"),
lifecycle_phase: TaskLifecyclePhase::Loop,
phase_epoch: 4,
phase_cursor: 0,
phase_iteration: 0,
gate_cycle: 1,
gate_proposal: None,
agent: "codex".to_string(),
provider: "codex".to_string(),
provider_session_id: None,
abandon_intent: None,
created_at: now,
updated_at: now,
observation: Observation::NotRequired,
};
let pr = TaskPr {
id: TaskPrId::new(),
task_id: task.id.clone(),
sequence: 1,
slug: task.workspace_slug.clone(),
branch: format!("test/{}", task.workspace_slug),
base_commit: "deadbeef".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.unwrap();
store.create_project(&project).await.unwrap();
store.create_task(&task, &pr).await.unwrap();
let work = store
.work_for_child(&ChildRef::Task(task.id.clone()))
.await
.unwrap();
TaskFixture {
_database: database,
database_path,
store,
task,
work,
}
}
#[test]
fn piped_report_supplies_title_and_preserves_full_description() {
let report = "\n lf status rejects stored timestamp \n\nstack trace\nmore evidence\n";
let input = resolve_task_start_input(None, Some(report)).expect("resolve piped report");
assert_eq!(input.title, "lf status rejects stored timestamp");
assert_eq!(
input.report,
"lf status rejects stored timestamp \n\nstack trace\nmore evidence"
);
}
#[test]
fn project_flows_resolve_once_with_per_task_overrides() {
let repo = tempfile::tempdir().expect("temp repo");
let project = ProjectFlowPlan {
first: Some("incident".to_string()),
loop_: Some("ship-5whys".to_string()),
finally: Some("ship".to_string()),
};
let overrides = TaskFlowOverrides {
loop_: Some("slice".to_string()),
..TaskFlowOverrides::default()
};
let plan =
resolve_task_lifecycle(repo.path(), &project, &overrides).expect("resolve lifecycle");
assert_eq!(plan.first.flow, "incident");
assert_eq!(plan.loop_.flow, "slice");
assert_eq!(plan.finally.flow, "ship");
}
#[test]
fn finally_flow_rejects_ops_before_agent_work() {
let repo = tempfile::tempdir().expect("temp repo");
let flows = repo.path().join(".lf/flows");
std::fs::create_dir_all(&flows).expect("create flow directory");
std::fs::write(
flows.join("unsafe-finally.yaml"),
"- op: pr land -c\n- gate\n",
)
.expect("write flow");
let project = ProjectFlowPlan {
first: None,
loop_: None,
finally: Some("unsafe-finally".to_string()),
};
let error = resolve_task_lifecycle(repo.path(), &project, &TaskFlowOverrides::default())
.expect_err("reject unsafe finally flow");
assert!(error
.to_string()
.contains("one or more skills followed by optional ops"));
}
#[test]
fn started_task_rejects_a_different_flow_override() {
let repo = tempfile::tempdir().expect("temp repo");
ensure_task_flow_override(repo.path(), "INF-123", "loop", Some("slice"), "slice")
.expect("same pinned flow is idempotent");
let error =
ensure_task_flow_override(repo.path(), "INF-123", "loop", Some("ship-5whys"), "slice")
.expect_err("different flow must be rejected");
assert!(error
.to_string()
.contains("Task INF-123 already pins loop flow \"slice\""));
}
#[tokio::test]
async fn unavailable_persisted_flow_is_rejected_before_every_run_reservation() {
let TaskFixture {
store,
mut task,
work,
..
} = task_fixture("TEST-STALE", "retired-task-flow").await;
for _ in 0..2 {
let error = launch_task_process(&store, &mut task, None)
.await
.expect_err("the unavailable flow must fail startup validation");
assert!(error.to_string().contains(
"Task TEST-STALE cannot launch: pinned loop flow \"retired-task-flow\" is invalid: failed to load Task flow \"retired-task-flow\": flow not found: retired-task-flow"
));
assert!(store.current_run(&work).await.unwrap().is_none());
}
assert!(store
.recent_task_events(&task.id, 10)
.await
.unwrap()
.is_empty());
}
#[cfg(unix)]
#[allow(clippy::await_holding_lock)]
#[tokio::test]
async fn repaired_persisted_task_flows_launch_through_the_generic_path() {
let _env_lock = crate::journal::test_env_lock();
let _restore = EnvRestore::capture(&["LF_BIN", "LF_HOME", "LF_DB_PATH", "PATH"]);
let environment = tempfile::tempdir().unwrap();
let bin = environment.path().join("bin");
std::fs::create_dir(&bin).unwrap();
let tmux = bin.join("tmux");
std::fs::write(&tmux, "#!/bin/sh\nexit 0\n").unwrap();
std::fs::set_permissions(&tmux, std::fs::Permissions::from_mode(0o755)).unwrap();
std::env::set_var("LF_BIN", std::env::current_exe().unwrap());
std::env::set_var("LF_HOME", environment.path());
std::env::remove_var("LF_DB_PATH");
std::env::set_var("PATH", &bin);
for identifier in ["LOO-167", "LOO-193", "LOO-195"] {
let TaskFixture {
_database,
database_path,
store,
task,
work,
..
} = task_fixture(identifier, "task").await;
rusqlite::Connection::open(database_path)
.unwrap()
.execute_batch(&crate::store::migrations::migration_sql_for_test(
"repair_legacy_task_flow",
))
.unwrap();
let mut repaired = store.get_task(&task.id).await.unwrap().unwrap();
assert_eq!(repaired.lifecycle.loop_.flow, "slice");
launch_task_process(&store, &mut repaired, None)
.await
.expect("a repaired Task reaches the generic launch boundary");
let run = store.current_run(&work).await.unwrap().unwrap();
assert!(matches!(run.trigger, RunTrigger::Input { .. }));
assert!(store
.open_invocation_for_run(&run.id)
.await
.unwrap()
.is_some());
}
}
#[tokio::test]
async fn automatic_task_relaunch_exhausts_until_durable_progress_or_user_input() {
let TaskFixture {
store, task, work, ..
} = task_fixture("TEST-RECOVERY", "slice").await;
let basis = store.current_epoch(&work).await.unwrap().current_basis;
let (initial, initial_lease) = store
.reserve_run(
&work,
RunTrigger::Input {
basis: basis.clone(),
},
)
.await
.unwrap();
store
.append_task_event_for_run(&task.id, &initial_lease, &TaskEventKind::Started)
.await
.unwrap();
store
.fail_task_run(&task.id, &initial_lease, "always fails")
.await
.unwrap();
let mut prior_run_id = initial.id;
for recovery in 1..=MAX_AUTOMATIC_TASK_RECOVERY_RUNS {
let trigger = prepare_automatic_task_relaunch(&store, &task)
.await
.unwrap()
.expect("automatic recovery remains within its budget");
assert_eq!(
trigger,
RunTrigger::Recovery {
prior_run_id: prior_run_id.clone()
}
);
let (run, lease) = store.reserve_run(&work, trigger).await.unwrap();
assert_eq!(run.retry_of, Some(prior_run_id));
store
.append_task_event_for_run(&task.id, &lease, &TaskEventKind::Started)
.await
.unwrap();
store
.fail_task_run(
&task.id,
&lease,
&format!("automatic recovery {recovery} failed"),
)
.await
.unwrap();
prior_run_id = run.id;
}
assert!(prepare_automatic_task_relaunch(&store, &task)
.await
.unwrap()
.is_none());
assert!(matches!(
store.work_status(&work).await.unwrap(),
WorkStatus::Waiting { .. }
));
let parked_run = store.latest_run(&work).await.unwrap().unwrap();
for _ in 0..10 {
assert!(prepare_automatic_task_relaunch(&store, &task)
.await
.unwrap()
.is_none());
assert_eq!(
store.latest_run(&work).await.unwrap().unwrap().id,
parked_run.id
);
}
let (user_run, user_lease) = store
.reserve_run(&work, RunTrigger::User)
.await
.expect("explicit User input clears the exhausted wait");
store
.append_task_event_for_run(&task.id, &user_lease, &TaskEventKind::Started)
.await
.unwrap();
store
.fail_task_run(&task.id, &user_lease, "User-started Run failed")
.await
.unwrap();
let trigger = prepare_automatic_task_relaunch(&store, &task)
.await
.unwrap()
.expect("User input starts a fresh recovery budget");
assert_eq!(
trigger,
RunTrigger::Recovery {
prior_run_id: user_run.id
}
);
let (_, progress_lease) = store.reserve_run(&work, trigger).await.unwrap();
store
.append_task_event_for_run(&task.id, &progress_lease, &TaskEventKind::Started)
.await
.unwrap();
store
.append_task_event_for_run(
&task.id,
&progress_lease,
&TaskEventKind::Progress {
summary: "durable progress".to_string(),
},
)
.await
.unwrap();
store
.fail_task_run(&task.id, &progress_lease, "failed after durable progress")
.await
.unwrap();
assert_eq!(
prepare_automatic_task_relaunch(&store, &task)
.await
.unwrap(),
Some(RunTrigger::Input { basis })
);
}
#[test]
fn pr_mutation_lock_refuses_a_concurrent_writer() {
let repo = tempfile::tempdir().expect("temporary repository");
let status = std::process::Command::new("git")
.args(["init", "--quiet"])
.current_dir(repo.path())
.status()
.expect("initialize repository");
assert!(status.success());
let _first = lock_task_pr_mutation(repo.path()).expect("first mutation lock");
let error = lock_task_pr_mutation(repo.path()).expect_err("second writer must be refused");
assert!(error.to_string().contains("already running"));
}
}