use std::collections::{HashMap, HashSet};
use std::io::{IsTerminal, Read};
use std::path::Path;
use std::sync::{Arc, OnceLock};
use clap::{CommandFactory, Parser};
use tracing::{debug, warn};
use tracing_subscriber::EnvFilter;
use loopflow::journal::{self, LfEventFields, LfEventType, LfNode};
use loopflow::lf::{Cli, Commands, InstallCommand, ProjectCommand, TaskCommand};
#[derive(Clone, Default)]
struct FlagTables {
value: HashSet<String>,
boolean: HashSet<String>,
}
impl FlagTables {
fn insert(&mut self, arg: &clap::Arg) {
let flags = if arg.get_action().takes_values() {
&mut self.value
} else {
&mut self.boolean
};
if let Some(short) = arg.get_short() {
flags.insert(format!("-{short}"));
}
for alias in arg.get_all_short_aliases().into_iter().flatten() {
flags.insert(format!("-{alias}"));
}
if let Some(long) = arg.get_long() {
flags.insert(format!("--{long}"));
}
for alias in arg.get_all_aliases().into_iter().flatten() {
flags.insert(format!("--{alias}"));
}
}
fn contains(&self, arg: &str) -> bool {
self.boolean.contains(flag_name(arg)) || self.takes_value(arg)
}
fn takes_value(&self, arg: &str) -> bool {
self.value.contains(flag_name(arg))
}
fn extend(&mut self, other: &Self) {
self.value.extend(other.value.iter().cloned());
self.boolean.extend(other.boolean.iter().cloned());
}
}
#[derive(Clone, Default)]
struct CommandArgTables {
direct: FlagTables,
recursive: FlagTables,
subcommands: HashMap<String, CommandArgTables>,
}
struct ArgTables {
commands: HashMap<String, CommandArgTables>,
top_level: FlagTables,
}
fn command_arg_tables(command: &clap::Command) -> CommandArgTables {
let mut direct = FlagTables::default();
for arg in command.get_arguments() {
direct.insert(arg);
}
let mut recursive = direct.clone();
let mut subcommands = HashMap::new();
for subcommand in command.get_subcommands() {
let child = command_arg_tables(subcommand);
recursive.extend(&child.recursive);
let names = std::iter::once(subcommand.get_name()).chain(subcommand.get_all_aliases());
for name in names {
subcommands.insert(name.to_string(), child.clone());
}
}
CommandArgTables {
direct,
recursive,
subcommands,
}
}
fn arg_tables() -> &'static ArgTables {
static TABLES: OnceLock<ArgTables> = OnceLock::new();
TABLES.get_or_init(|| {
let mut cli = Cli::command();
cli.build();
let mut commands = HashMap::new();
for sub in cli.get_subcommands() {
let table = command_arg_tables(sub);
let names = std::iter::once(sub.get_name()).chain(sub.get_all_aliases());
for name in names {
commands.insert(name.to_string(), table.clone());
}
}
let mut top_level = FlagTables::default();
for arg in cli.get_arguments() {
top_level.insert(arg);
}
ArgTables {
commands,
top_level,
}
})
}
fn flag_name(arg: &str) -> &str {
arg.split_once('=').map_or(arg, |(name, _)| name)
}
fn has_inline_value(arg: &str) -> bool {
arg.starts_with('-') && arg.contains('=')
}
fn is_value_flag(arg: &str) -> bool {
arg_tables().top_level.takes_value(arg)
}
fn is_known_flag(arg: &str) -> bool {
arg_tables().top_level.contains(arg)
}
fn push_flag(args: &[String], output: &mut Vec<String>, index: &mut usize, takes_value: bool) {
output.push(args[*index].clone());
if takes_value && !has_inline_value(&args[*index]) && *index + 1 < args.len() {
*index += 1;
output.push(args[*index].clone());
}
}
fn first_target_index(args: &[String]) -> Option<usize> {
let mut index = 0;
while index < args.len() {
let arg = &args[index];
if arg == "--" {
return None;
}
if !arg.starts_with('-') {
return Some(index);
}
if is_value_flag(arg) && !has_inline_value(arg) {
index += 1;
}
index += 1;
}
None
}
fn normalize_ssh_args(mut args: Vec<String>) -> Vec<String> {
if args.len() <= 1 {
return args;
}
let rest = &args[1..];
let Some(command_index) = first_target_index(rest) else {
return args;
};
if rest[command_index] != "ssh" {
return args;
}
let ssh_args = &arg_tables()
.commands
.get("ssh")
.expect("ssh command has derived argument metadata")
.direct;
let mut index = command_index + 2;
while index < args.len() {
let arg = &args[index];
if arg == "--" {
return args;
}
if arg.starts_with('-') {
let takes_value = ssh_args.takes_value(arg) || is_value_flag(arg);
if takes_value && !has_inline_value(arg) {
index += 1;
}
index += 1;
continue;
}
args.insert(index + 1, "--".to_string());
return args;
}
args
}
#[derive(Clone, Copy)]
struct SelectedCommand<'a> {
index: usize,
args: &'a CommandArgTables,
}
fn selected_command_path<'a>(
rest: &[String],
command_index: usize,
command: &'a CommandArgTables,
) -> Vec<SelectedCommand<'a>> {
let mut path = vec![SelectedCommand {
index: command_index,
args: command,
}];
let mut index = command_index + 1;
while index < rest.len() {
let arg = &rest[index];
if arg == "--" {
break;
}
let current = path.last().expect("command path is never empty").args;
if arg.starts_with('-') {
let recognized = current.recursive.contains(arg) || is_known_flag(arg);
let takes_value = current.recursive.takes_value(arg) || is_value_flag(arg);
if recognized && takes_value && !has_inline_value(arg) && index + 1 < rest.len() {
index += 1;
}
index += 1;
continue;
}
let Some(child) = current.subcommands.get(arg) else {
break;
};
path.push(SelectedCommand { index, args: child });
index += 1;
}
path
}
fn deepest_command_before(path: &[SelectedCommand<'_>], index: usize) -> usize {
path.iter()
.rposition(|command| command.index < index)
.expect("the top-level command precedes its arguments")
}
fn local_flag_owner(path: &[SelectedCommand<'_>], arg: &str, current: usize) -> Option<usize> {
if path[current].args.direct.contains(arg) {
return Some(current);
}
path.iter()
.rposition(|command| command.args.direct.contains(arg))
}
fn reorder_command_args(
program: String,
rest: &[String],
command_index: usize,
command: &CommandArgTables,
) -> Vec<String> {
let path = selected_command_path(rest, command_index, command);
let mut moved_globals = Vec::new();
let mut moved_locals: HashMap<usize, Vec<String>> = HashMap::new();
let mut retained = vec![true; rest.len()];
let mut index = command_index + 1;
while index < rest.len() {
let arg = &rest[index];
if arg == "--" {
break;
}
let current = deepest_command_before(&path, index);
let local_owner = local_flag_owner(&path, arg, current);
let destination = if let Some(owner) = local_owner {
(owner != current).then_some(Some(owner))
} else if is_known_flag(arg) {
Some(None)
} else {
None
};
let Some(destination) = destination else {
index += 1;
continue;
};
let takes_value = destination.map_or_else(
|| is_value_flag(arg),
|owner| path[owner].args.direct.takes_value(arg),
);
let mut moved = Vec::new();
push_flag(rest, &mut moved, &mut index, takes_value);
let moved_start = index + 1 - moved.len();
retained[moved_start..=index].fill(false);
if let Some(owner) = destination {
moved_locals.entry(owner).or_default().extend(moved);
} else {
moved_globals.extend(moved);
}
index += 1;
}
let mut result = vec![program];
result.extend_from_slice(&rest[..command_index]);
result.extend(moved_globals);
for (index, arg) in rest.iter().enumerate().skip(command_index) {
if retained[index] {
result.push(arg.clone());
}
if let Some(owner) = path.iter().position(|command| command.index == index) {
if let Some(flags) = moved_locals.remove(&owner) {
result.extend(flags);
}
}
}
result
}
fn reorder_args(args: Vec<String>) -> Vec<String> {
if args.len() <= 1 {
return args;
}
let program = args[0].clone();
let rest = &args[1..];
let Some(target_index) = first_target_index(rest) else {
return args;
};
if let Some(command) = arg_tables().commands.get(rest[target_index].as_str()) {
if rest[target_index] == "ssh" {
return args;
}
return reorder_command_args(program, rest, target_index, command);
}
let mut flags_before: Vec<String> = Vec::new();
let mut skill_and_args: Vec<String> = Vec::new();
let mut flags_after: Vec<String> = Vec::new();
let mut i = 0;
let mut found_skill = false;
while i < rest.len() {
let arg = &rest[i];
if arg == "--" {
skill_and_args.extend_from_slice(&rest[i..]);
break;
}
if !found_skill {
if arg.starts_with('-') {
flags_before.push(arg.clone());
if is_value_flag(arg) && !has_inline_value(arg) && i + 1 < rest.len() {
i += 1;
flags_before.push(rest[i].clone());
}
} else {
found_skill = true;
skill_and_args.push(arg.clone());
}
} else {
if arg.starts_with('-') {
if is_known_flag(arg) {
flags_after.push(arg.clone());
if is_value_flag(arg) && !has_inline_value(arg) && i + 1 < rest.len() {
i += 1;
flags_after.push(rest[i].clone());
}
} else {
skill_and_args.push(arg.clone());
}
} else {
skill_and_args.push(arg.clone());
}
}
i += 1;
}
let mut result = vec![program];
result.extend(flags_before);
result.extend(flags_after);
result.extend(skill_and_args);
result
}
fn join_args(args: &[String]) -> Option<String> {
if args.is_empty() {
None
} else {
Some(args.join(" "))
}
}
fn with_runtime<T>(
repo_root: &std::path::Path,
command: &[String],
run: impl FnOnce() -> anyhow::Result<T>,
) -> anyhow::Result<T> {
let attribution = loopflow::work::wave::context::run_attribution(Some(repo_root));
if let Some(failure) = attribution.failure.as_deref() {
warn!(
error = failure,
"ambient wave identity failed validation; run attributed to no wave \
— pass --wave <name> to recover"
);
}
journal::emit(
repo_root,
LfNode::Run,
LfEventType::Started,
LfEventFields {
wave_name: attribution.wave,
error: attribution.failure,
worktree: Some(repo_root.display().to_string()),
command: Some(command.to_vec()),
..LfEventFields::default()
},
);
let result = run();
match &result {
Ok(_) => journal::emit(
repo_root,
LfNode::Run,
LfEventType::Completed,
LfEventFields::default(),
),
Err(err) => journal::emit(
repo_root,
LfNode::Run,
LfEventType::Errored,
LfEventFields {
error: Some(err.to_string()),
..LfEventFields::default()
},
),
}
result
}
fn in_repo_runtime<T>(
command: &[String],
run: impl FnOnce(&std::path::Path) -> anyhow::Result<T>,
) -> anyhow::Result<T> {
let repo_root = loopflow::lf::commands::util::find_repo_root()?;
with_runtime(&repo_root, command, || run(&repo_root))
}
fn run_default_agent(cli: &Cli, command: &[String]) -> anyhow::Result<()> {
let repo_root = loopflow::lf::commands::util::find_repo_root()?;
let moved = loopflow::engine::worktrees::move_default_agent_to_worktree(&repo_root)?;
match moved {
Some(worktree) => {
eprintln!("moved to `{}`", worktree.path.display());
let _cwd = CwdGuard::enter(&worktree.path)?;
with_runtime(&worktree.path, command, || {
loopflow::lf::commands::run::run(Some("loopflow"), None, cli)
})
}
None => with_runtime(&repo_root, command, || {
loopflow::lf::commands::run::run(Some("loopflow"), None, cli)
}),
}
}
fn run_work_runner_entrypoint<T>(
command: &[String],
kind: &str,
work_id: &str,
run: impl FnOnce(loopflow::durable::WorkRef) -> anyhow::Result<T>,
) -> anyhow::Result<T> {
in_repo_runtime(command, |_| {
let work = match kind {
"project" => loopflow::durable::WorkRef::Project(work_id.parse()?),
"task" => loopflow::durable::WorkRef::Task(work_id.parse()?),
_ => anyhow::bail!("unsupported hidden Work kind: {kind}"),
};
run(work)
})
}
fn with_skill_runtime<T>(
repo_root: &std::path::Path,
skill_name: &str,
run: impl FnOnce() -> anyhow::Result<T>,
) -> anyhow::Result<T> {
journal::emit(
repo_root,
LfNode::Skill,
LfEventType::Started,
LfEventFields {
skill: Some(skill_name.to_string()),
index: Some(0),
..LfEventFields::default()
},
);
let result = run();
match &result {
Ok(_) => journal::emit(
repo_root,
LfNode::Skill,
LfEventType::Completed,
LfEventFields {
skill: Some(skill_name.to_string()),
index: Some(0),
..LfEventFields::default()
},
),
Err(err) => journal::emit(
repo_root,
LfNode::Skill,
LfEventType::Errored,
LfEventFields {
skill: Some(skill_name.to_string()),
index: Some(0),
error: Some(err.to_string()),
..LfEventFields::default()
},
),
}
result
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum TargetKind {
Flow,
Skill,
}
fn require_target_kind(name: &str, kind: TargetKind) -> anyhow::Result<()> {
let repo_root = loopflow::lf::commands::util::find_repo_root()?;
match (
loopflow::lf::discovery::discover_target(&repo_root, name)?,
kind,
) {
(loopflow::lf::discovery::Target::Flow(_), TargetKind::Flow) => Ok(()),
(loopflow::lf::discovery::Target::Skill(_), TargetKind::Skill) => Ok(()),
(loopflow::lf::discovery::Target::Flow(_), TargetKind::Skill) => {
Err(anyhow::anyhow!("'{name}' is a flow — run `lf flow {name}`"))
}
(loopflow::lf::discovery::Target::Skill(_), TargetKind::Flow) => Err(anyhow::anyhow!(
"'{name}' is a skill — run `lf skill {name}`"
)),
}
}
fn run_target(
name: &str,
message: Option<&str>,
cli: &Cli,
command: &[String],
) -> anyhow::Result<()> {
let repo_root = loopflow::lf::commands::util::find_repo_root()?;
run_target_in_repo(&repo_root, name, message, cli, command)
}
fn run_target_in_repo(
repo_root: &Path,
name: &str,
message: Option<&str>,
cli: &Cli,
command: &[String],
) -> anyhow::Result<()> {
match loopflow::lf::discovery::discover_target(repo_root, name)? {
loopflow::lf::discovery::Target::Skill(_) => with_runtime(repo_root, command, || {
with_skill_runtime(repo_root, name, || {
loopflow::lf::commands::run::run(Some(name), message, cli)?;
if cli.work_subject_selector().is_some() {
return Ok(());
}
let options = loopflow::ops::CommitOptions {
add: true,
message: Some(format!("lf commit: {name}")),
..loopflow::ops::CommitOptions::for_task(name)
};
loopflow::ops::commit_workflow(repo_root, &options, &loopflow::ops::NullProgress)?;
Ok(())
})
}),
loopflow::lf::discovery::Target::Flow(flow) => with_runtime(repo_root, command, || {
loopflow::lf::commands::flow::run(&flow, message, cli, repo_root)
}),
}
}
fn run_bound_target_in_repo(
repo_root: &Path,
name: &str,
message: Option<&str>,
cli: &Cli,
command: &[String],
binding: &loopflow::ops::WorkBinding,
) -> anyhow::Result<()> {
match loopflow::lf::discovery::discover_target(repo_root, name)? {
loopflow::lf::discovery::Target::Skill(_) => with_runtime(repo_root, command, || {
with_skill_runtime(repo_root, name, || {
loopflow::lf::commands::run::run_bound(Some(name), message, cli, binding)?;
let options = loopflow::ops::CommitOptions {
add: true,
message: Some(format!("lf commit: {name}")),
..loopflow::ops::CommitOptions::for_task(name)
};
loopflow::ops::commit_workflow(repo_root, &options, &loopflow::ops::NullProgress)?;
Ok(())
})
}),
loopflow::lf::discovery::Target::Flow(_) => anyhow::bail!(
"Work selectors support one named skill invocation, not multi-step flow {name:?}"
),
}
}
fn prepare_work_binding(selector: &str, repo: &Path) -> anyhow::Result<loopflow::ops::WorkBinding> {
let runtime = tokio::runtime::Runtime::new()
.map_err(|error| anyhow::anyhow!("cannot resolve --as {selector}: {error}"))?;
runtime.block_on(async {
let store = loopflow::store::open_existing_store()
.await
.ok_or_else(|| {
anyhow::anyhow!("cannot resolve --as {selector}: planning registry unavailable")
})?;
let store = Arc::new(store);
loopflow::ops::resolve_work_binding(&store, repo, selector)
.await
.map_err(anyhow::Error::from)
})
}
fn prepare_hierarchical_work_binding(
task: Option<&str>,
project: Option<&str>,
wave: Option<&str>,
repo: &Path,
) -> anyhow::Result<loopflow::ops::WorkBinding> {
let runtime = tokio::runtime::Runtime::new()
.map_err(|error| anyhow::anyhow!("cannot resolve Work selection: {error}"))?;
runtime.block_on(async {
let store = loopflow::store::open_existing_store()
.await
.ok_or_else(|| {
anyhow::anyhow!("cannot resolve Work selection: planning registry unavailable")
})?;
let store = Arc::new(store);
loopflow::ops::resolve_work_selection(
&store,
repo,
loopflow::ops::WorkSelection {
task,
project,
wave,
},
)
.await
.map_err(anyhow::Error::from)
})
}
fn validate_work_selector(selector: &str) -> anyhow::Result<()> {
let (kind, value) = selector.split_once(':').ok_or_else(|| {
anyhow::anyhow!(
"invalid --as {selector:?}; use task:<selector>, project:<selector>, or wave:<selector>"
)
})?;
if !matches!(kind, "task" | "project" | "wave") {
anyhow::bail!("invalid --as kind {kind:?}; expected task, project, or wave");
}
if value.trim().is_empty() {
anyhow::bail!("{kind} selector cannot be empty");
}
Ok(())
}
fn require_bound_invocation(command: &Option<Commands>, repo: &Path) -> anyhow::Result<()> {
let name = match command {
Some(Commands::Skill { name, .. }) => name.clone(),
Some(Commands::External(args)) => loopflow::lf::commands::run::split_skill_args(args)?.0,
Some(Commands::Inline { .. }) => return Ok(()),
_ => {
anyhow::bail!("`lf --as` starts one skill or inline prompt; use `lf --as task:LOO-123 implement` or `lf --as project:api : \"question\"`")
}
};
match loopflow::lf::discovery::discover_target(repo, &name)? {
loopflow::lf::discovery::Target::Skill(_) => Ok(()),
loopflow::lf::discovery::Target::Flow(_) => anyhow::bail!(
"Work selectors support one named skill invocation, not multi-step flow {name:?}"
),
}
}
struct CwdGuard(std::path::PathBuf);
impl CwdGuard {
fn enter(path: &Path) -> anyhow::Result<Self> {
if !path.is_absolute() {
anyhow::bail!("bound Work cwd must be absolute: {}", path.display());
}
let previous = std::env::current_dir()?;
std::env::set_current_dir(path)
.map_err(|error| anyhow::anyhow!("enter bound Work cwd {}: {error}", path.display()))?;
Ok(Self(previous))
}
}
impl Drop for CwdGuard {
fn drop(&mut self) {
let _ = std::env::set_current_dir(&self.0);
}
}
struct EnvGuard {
key: &'static str,
previous: Option<std::ffi::OsString>,
}
impl EnvGuard {
fn set(key: &'static str, value: impl Into<String>) -> Self {
let previous = std::env::var_os(key);
std::env::set_var(key, value.into());
Self { key, previous }
}
}
impl Drop for EnvGuard {
fn drop(&mut self) {
if let Some(value) = &self.previous {
std::env::set_var(self.key, value);
} else {
std::env::remove_var(self.key);
}
}
}
fn build_local_account_lease(
selection: &loopflow::provider_account::lease::AccountSelection,
) -> anyhow::Result<
Option<(
loopflow::provider_account::lease::AccountLeaseBroker,
EnvGuard,
)>,
> {
use loopflow::provider_account::lease;
let runtime = tokio::runtime::Runtime::new()?;
let Some(broker) = runtime.block_on(lease::AccountLeaseBroker::start_root(selection))? else {
return Ok(None);
};
let guard = EnvGuard::set(lease::ACCOUNT_LEASE_ENV, broker.local_env_value()?);
Ok(Some((broker, guard)))
}
fn parse_duration(value: &str) -> anyhow::Result<std::time::Duration> {
let value = value.trim();
let (number, multiplier) = if let Some(number) = value.strip_suffix('s') {
(number, 1)
} else if let Some(number) = value.strip_suffix('m') {
(number, 60)
} else if let Some(number) = value.strip_suffix('h') {
(number, 60 * 60)
} else {
(value, 1)
};
let amount: u64 = number
.parse()
.map_err(|_| anyhow::anyhow!("invalid duration {value:?}; use seconds, 10s, 5m, or 1h"))?;
Ok(std::time::Duration::from_secs(
amount.saturating_mul(multiplier),
))
}
fn format_task_pr_line(pr: &loopflow::work::task::TaskPr) -> String {
let provider = pr
.github()
.map(|github| format!("GitHub #{}", github.number))
.unwrap_or_else(|| "not opened on GitHub".to_string());
let placement = pr
.parent_pr_id
.as_ref()
.map(|parent| format!(" stacked on {parent}"))
.unwrap_or_default();
let linkage = pr
.linear_link_error
.as_ref()
.map(|error| format!(" Linear link degraded: {error}"))
.unwrap_or_default();
format!(
" PR {}: {} {} {}{}{}",
pr.sequence,
pr.phase().as_str(),
provider,
pr.branch,
placement,
linkage,
)
}
fn print_task(task: &loopflow::work::task::Task, json: bool) -> anyhow::Result<()> {
let snapshot = loopflow::ops::task::task_snapshot(task)?;
if json {
println!("{}", serde_json::to_string_pretty(&snapshot)?);
} else {
let pm_writeback = match &task.pm_writeback {
loopflow::work::task::PmWritebackState::Current => "current".to_string(),
loopflow::work::task::PmWritebackState::Pending { error, .. } => {
format!("pending: {error}")
}
};
let branch = snapshot
.active_pr
.as_ref()
.and_then(|active| snapshot.prs.iter().find(|pr| &pr.id == active))
.map(|pr| pr.branch.as_str())
.unwrap_or("none");
let (phase, cycle, flow, iteration, step, body) = match snapshot.controller.as_ref() {
Some(controller) => (
controller.lifecycle_phase.as_str().to_string(),
match controller.lifecycle_phase {
loopflow::controller::task::TaskLifecyclePhase::First => "0".to_string(),
loopflow::controller::task::TaskLifecyclePhase::Loop => {
(controller.gate_cycle + 1).to_string()
}
loopflow::controller::task::TaskLifecyclePhase::Finally => {
controller.gate_cycle.to_string()
}
_ => "unknown".to_string(),
},
controller
.lifecycle
.phase(controller.lifecycle_phase)
.flow
.clone(),
(controller.phase_iteration + 1).to_string(),
(controller.phase_cursor + 1).to_string(),
format!(
"agent {}, provider {}",
controller.agent, controller.provider
),
),
None => (
"none".to_string(),
"none".to_string(),
"none".to_string(),
"none".to_string(),
"none".to_string(),
"no end-to-end controller".to_string(),
),
};
println!(
"{} {}\n task: {}\n phase: {} cycle {}\n flow: {} (iteration {}, step {})\n body: {}\n worktree: {}\n branch: {}\n PM writeback: {}\n reason: {}",
task.plan.identifier,
snapshot.status,
task.id,
phase,
cycle,
flow,
iteration,
step,
body,
task.worktree.display(),
branch,
pm_writeback,
snapshot.status,
);
println!(" project: {}", task.project_id);
for pr in &snapshot.prs {
println!("{}", format_task_pr_line(pr));
}
match &snapshot.observation {
loopflow::work::task::Observation::Cached { observed_at } => {
println!(" PR observation: cached from {observed_at}");
}
loopflow::work::task::Observation::Degraded {
reason,
cached_as_of,
retry_at,
} => {
println!(
" PR observation: degraded — {reason} (cached from {cached_as_of}; retry after {retry_at})"
);
}
loopflow::work::task::Observation::NotRequired
| loopflow::work::task::Observation::Fresh { .. } => {}
}
let actions = &snapshot.actions;
if let Some(recommended) = actions.recommended {
println!(" action: {} ({})", recommended.as_str(), actions.reason);
}
}
Ok(())
}
fn print_task_control(
result: &loopflow::ops::task::TaskControlResult,
json: bool,
) -> anyhow::Result<()> {
if json {
println!("{}", serde_json::to_string_pretty(result)?);
} else {
println!(
"{} → {} ({})",
result.receipt.label(),
result.issue_id,
result.receipt.action(),
);
match &result.observation {
loopflow::work::task::Observation::Cached { observed_at } => {
println!("PR observation: cached from {observed_at}");
}
loopflow::work::task::Observation::Degraded {
reason,
cached_as_of,
retry_at,
} => {
println!(
"PR observation: degraded — {reason} (cached from {cached_as_of}; retry after {retry_at})"
);
}
loopflow::work::task::Observation::NotRequired
| loopflow::work::task::Observation::Fresh { .. } => {}
}
}
Ok(())
}
fn print_project(project: &loopflow::work::project::Project, json: bool) -> anyhow::Result<()> {
let snapshot = loopflow::ops::project::project_snapshot(project)?;
if json {
println!("{}", serde_json::to_string_pretty(&snapshot)?);
} else {
let body = match (&snapshot.agent, &snapshot.provider) {
(Some(agent), Some(provider)) => format!("agent {agent}, provider {provider}"),
_ => "no end-to-end controller".to_string(),
};
println!(
"{} {}\n project: {}\n body: {}\n iteration: {}\n reason: {}",
project.plan.slug,
snapshot.status,
project.id,
body,
snapshot
.iteration
.map_or_else(|| "none".to_string(), |iteration| iteration.to_string()),
snapshot.reason,
);
if let Some(failure) = &snapshot.last_failure {
println!(
" last failure at {}: {}",
failure.occurred_at, failure.message
);
}
}
Ok(())
}
fn print_project_control(
result: &loopflow::ops::project::ProjectControlResult,
json: bool,
) -> anyhow::Result<()> {
if json {
println!("{}", serde_json::to_string_pretty(result)?);
} else {
println!(
"{} → {} ({})",
result.receipt_id, result.external_project_id, result.action,
);
}
Ok(())
}
fn run_project_command(repo: &Path, command: &ProjectCommand) -> anyhow::Result<()> {
match command {
ProjectCommand::Prepare {
project_id,
directive,
json,
} => {
let project =
loopflow::ops::project::project_prepare(repo, project_id, directive.clone())?;
print_project(&project, *json)
}
ProjectCommand::Run {
project_id,
directive,
json,
} => {
let project = loopflow::ops::project::project_run(repo, project_id, directive.clone())?;
print_project(&project, *json)
}
ProjectCommand::Start {
title,
wave,
directive,
json,
} => {
let project = loopflow::ops::project::project_start(
repo,
title,
wave.as_deref(),
directive.clone(),
)?;
print_project(&project, *json)
}
ProjectCommand::Status { project_id, json } => {
let project = loopflow::ops::project::project_status(project_id)?;
print_project(&project, *json)
}
ProjectCommand::Steer {
project_id,
message,
json,
} => {
let result = loopflow::ops::project::project_steer(project_id, message.clone())?;
print_project_control(&result, *json)
}
ProjectCommand::Interrupt { project_id, json } => {
let result = loopflow::ops::project::project_interrupt(project_id)?;
print_project_control(&result, *json)
}
ProjectCommand::Wait {
project_id,
until,
timeout,
json,
} => {
let until = if until == "waiting" {
loopflow::ops::project::ProjectWaitUntil::Waiting
} else {
loopflow::ops::project::ProjectWaitUntil::Terminal
};
let timeout = timeout.as_deref().map(parse_duration).transpose()?;
let project = loopflow::ops::project::project_wait(project_id, until, timeout)?;
print_project(&project, *json)
}
ProjectCommand::Resume {
project_id,
model,
reason,
json,
} => {
let result =
loopflow::ops::project::project_resume(project_id, model.clone(), reason.clone())?;
print_project_control(&result, *json)
}
ProjectCommand::Attach { project_id } => {
loopflow::ops::project::project_attach(project_id).map_err(Into::into)
}
ProjectCommand::Abandon {
project_id,
reason,
json,
} => {
let result = loopflow::ops::project::project_abandon(project_id, reason.clone())?;
print_project_control(&result, *json)
}
ProjectCommand::Promote { .. } => {
anyhow::bail!("project promote is handled by the authored promotion flow")
}
}
}
fn task_cycle(fix: bool, feature: bool) -> Option<loopflow::ops::task::TaskCycle> {
if fix {
Some(loopflow::ops::task::TaskCycle::Fix)
} else if feature {
Some(loopflow::ops::task::TaskCycle::Feature)
} else {
None
}
}
fn run_task_command(repo: &Path, command: &TaskCommand) -> anyhow::Result<()> {
match command {
TaskCommand::Prepare {
issue,
name,
stack_on,
directive,
json,
} => {
let task = loopflow::ops::task::task_prepare(
repo,
issue,
loopflow::ops::task::TaskPrepareOptions {
name: name.clone(),
stack_on: stack_on.clone(),
directive: directive.clone(),
},
)?;
print_task(&task, *json)
}
TaskCommand::Run {
issue,
name,
fix,
feature,
first,
loop_,
finally,
stack_on,
directive,
json,
} => {
let task = loopflow::ops::task::task_run(
repo,
issue,
loopflow::ops::task::TaskLaunchOptions {
name: name.clone(),
flows: loopflow::ops::task::TaskFlowOverrides::for_cycle(
task_cycle(*fix, *feature),
first.clone(),
loop_.clone(),
finally.clone(),
),
stack_on: stack_on.clone(),
directive: directive.clone(),
},
)?;
print_task(&task, *json)
}
TaskCommand::Start {
project_id,
title,
name,
fix,
feature,
first,
loop_,
finally,
stack_on,
directive,
json,
} => {
let task = loopflow::ops::task::task_start(
repo,
project_id,
title.clone(),
piped_task_report()?,
loopflow::ops::task::TaskLaunchOptions {
name: name.clone(),
flows: loopflow::ops::task::TaskFlowOverrides::for_cycle(
task_cycle(*fix, *feature),
first.clone(),
loop_.clone(),
finally.clone(),
),
stack_on: stack_on.clone(),
directive: directive.clone(),
},
)?;
print_task(&task, *json)
}
TaskCommand::Status { issue, json } => {
let task = loopflow::ops::task::task_status(issue)?;
print_task(&task, *json)
}
TaskCommand::Changes { issue, json } => {
let snapshot = loopflow::ops::task::task_changes(issue)?;
if *json {
println!("{}", serde_json::to_string_pretty(&snapshot)?);
} else if snapshot.files.is_empty() {
println!("{} has no changes from its recorded base", issue);
} else {
for file in snapshot.files {
let mut states = Vec::new();
if file.committed {
states.push("committed");
}
if file.staged {
states.push("staged");
}
if file.unstaged {
states.push("unstaged");
}
if file.untracked {
states.push("untracked");
}
println!("{}\t{}", states.join(","), file.path);
}
}
Ok(())
}
TaskCommand::Diff { issue, path, json } => {
let snapshot = loopflow::ops::task::task_diff(issue, path.as_deref())?;
if *json {
println!("{}", serde_json::to_string_pretty(&snapshot)?);
} else {
print!("{}", snapshot.patch);
if snapshot.truncated {
eprintln!("\n[diff truncated at 1 MB]");
}
}
Ok(())
}
TaskCommand::File { issue, path, json } => {
let snapshot = loopflow::ops::task::task_file(issue, path)?;
if *json {
println!("{}", serde_json::to_string_pretty(&snapshot)?);
} else if snapshot.binary {
anyhow::bail!("{} is binary", snapshot.path);
} else {
print!("{}", snapshot.content.as_deref().unwrap_or_default());
if snapshot.truncated {
eprintln!("\n[file truncated at 1 MB]");
}
}
Ok(())
}
TaskCommand::Complete {
issue,
summary,
json,
} => {
let task = loopflow::ops::task::task_complete(issue, summary.clone())?;
print_task(&task, *json)
}
TaskCommand::Steer {
issue,
message,
json,
} => {
let result = loopflow::ops::task::task_steer(issue, message.clone())?;
print_task_control(&result, *json)
}
TaskCommand::Interrupt { issue, json } => {
let result = loopflow::ops::task::task_interrupt(issue)?;
print_task_control(&result, *json)
}
TaskCommand::Wait {
issue,
until,
timeout,
json,
} => {
let until = if until == "submitted" {
loopflow::ops::task::TaskWaitUntil::Open
} else {
loopflow::ops::task::TaskWaitUntil::Terminal
};
let timeout = timeout.as_deref().map(parse_duration).transpose()?;
let task = loopflow::ops::task::task_wait(issue, until, timeout)?;
print_task(&task, *json)
}
TaskCommand::Resume {
issue,
model,
reason,
json,
} => {
let result = loopflow::ops::task::task_resume(issue, model.clone(), reason.clone())?;
print_task_control(&result, *json)
}
TaskCommand::Restart {
issue,
advice,
json,
} => {
let task = loopflow::ops::task::task_restart(issue, advice.clone())?;
print_task(&task, *json)
}
TaskCommand::Recover {
issue,
reason,
json,
} => {
let task = loopflow::ops::task::task_recover(issue, reason.clone())?;
print_task(&task, *json)
}
TaskCommand::Abandon {
issue,
reason,
json,
} => {
let result = loopflow::ops::task::task_abandon(issue, reason.clone())?;
print_task_control(&result, *json)
}
}
}
fn piped_task_report() -> anyhow::Result<Option<String>> {
if std::io::stdin().is_terminal() {
return Ok(None);
}
let mut report = String::new();
std::io::stdin().read_to_string(&mut report)?;
Ok((!report.trim().is_empty()).then_some(report))
}
fn main() -> anyhow::Result<()> {
loopflow::machine_install::dispatch_entry_gate(&loopflow::machine_install::ArtifactRole::Cli)?;
let filter = EnvFilter::try_from_default_env()
.unwrap_or_else(|_| EnvFilter::new("lf=info,loopflow=info"));
tracing_subscriber::fmt()
.with_env_filter(filter)
.with_writer(std::io::stderr)
.without_time()
.init();
let args = reorder_args(normalize_ssh_args(std::env::args().collect()));
let mut cli = Cli::parse_from(args.clone());
let bypasses_machine_startup_gate = matches!(
&cli.command,
Some(Commands::Install {
cmd: InstallCommand::RecoverSwitch { .. }
| InstallCommand::AdvanceSwitch { .. }
| InstallCommand::LocalPreflight { .. }
| InstallCommand::Promote { .. }
})
);
if !bypasses_machine_startup_gate {
let switch_id = std::env::var(loopflow::machine_install::INSTALL_SWITCH_ENV)
.ok()
.filter(|value| !value.is_empty());
loopflow::machine_install::authorize_current_for_switch(
&loopflow::machine_install::ArtifactRole::Cli,
switch_id.as_deref(),
)?;
}
ctrlc::set_handler(|| {
loopflow::engine::agent::run_interrupt_cleanups();
loopflow::engine::agent::kill_child_if_running();
std::process::exit(130);
})
.expect("failed to set Ctrl+C handler");
match &cli.command {
Some(Commands::Screenshot { screenshot }) => {
return loopflow::lf::commands::screenshot::run(screenshot);
}
Some(Commands::ScreenshotSupervisor { screenshot }) => {
return loopflow::lf::commands::screenshot::run_supervisor(screenshot);
}
_ => {}
}
let explicit_wave = cli
.wave
.as_deref()
.map(loopflow::work::wave::context::resolve_explicit_wave)
.transpose()?;
if let Some(wave) = &explicit_wave {
cli.wave = Some(wave.name().to_string());
}
let _explicit_wave_env = explicit_wave.as_ref().map(|wave| {
EnvGuard::set(
loopflow::work::wave::context::WAVE_ID_ENV,
wave.id().to_string(),
)
});
let mut preferred_accounts = cli.account.clone();
let mut restricted_accounts = cli.only_account.clone();
if let Some(Commands::Ssh {
origin_account,
origin_only_account,
..
}) = &cli.command
{
preferred_accounts.extend(origin_account.iter().cloned());
restricted_accounts.extend(origin_only_account.iter().cloned());
}
let account_selection = loopflow::provider_account::lease::AccountSelection::from_flags(
&preferred_accounts,
&restricted_accounts,
)?;
let inherited_account_lease = loopflow::provider_account::lease::account_lease_active();
let _forwarded_account_selection = if inherited_account_lease && !account_selection.is_default()
{
Some(EnvGuard::set(
loopflow::provider_account::lease::ACCOUNT_SELECTION_ENV,
account_selection.env_value()?,
))
} else {
None
};
if cli.account_lease_probe {
return loopflow::provider_account::lease::probe_forwarded_authority()
.map_err(anyhow::Error::from);
}
debug!(?cli, "parsed CLI arguments");
if let Some(Commands::Install { cmd }) = &cli.command {
return match cmd {
InstallCommand::RecoverSwitch { switch } => {
loopflow::lf::commands::install::recover_switch(switch)
}
InstallCommand::Preflight { json } => loopflow::lf::commands::install::preflight(*json),
InstallCommand::LocalPreflight { store, json } => {
loopflow::lf::commands::install::local_preflight(store, *json)
}
InstallCommand::AdvanceSwitch { switch } => {
loopflow::lf::commands::install::advance_switch(switch)
}
InstallCommand::Promote {
from_build,
coordinated_build,
fresh,
cli_target,
daemon_source,
daemon_target,
app_source,
app_target,
legacy_app_target,
sync_skills,
preview,
} => loopflow::lf::commands::install::promote(
loopflow::lf::commands::install::PromotionArtifacts {
cli_target,
daemon_source,
daemon_target,
app_source: app_source.as_deref(),
app_target: app_target.as_deref(),
legacy_app_target: legacy_app_target.as_deref(),
},
*sync_skills,
*preview,
from_build.as_deref(),
coordinated_build.as_deref(),
*fresh,
),
InstallCommand::Rollback {
cli_target,
candidate,
daemon_target,
daemon_candidate,
} => loopflow::lf::commands::install::rollback(
cli_target,
candidate,
daemon_target,
daemon_candidate,
),
};
}
loopflow::lf::commands::home::validate_expected_home_process()?;
let is_ssh = matches!(cli.command, Some(loopflow::lf::Commands::Ssh { .. }));
let _local_account_lease =
if is_ssh || inherited_account_lease || account_selection.is_default() {
None
} else {
build_local_account_lease(&account_selection)?
};
let mut direct_binding = None;
let mut _bound_cwd = None;
let selects_direct_work = cli.as_work.is_some()
|| cli.task.is_some()
|| cli.project.is_some()
|| (cli.wave.is_some()
&& matches!(
&cli.command,
Some(Commands::Skill { .. }) | Some(Commands::External(_))
));
if selects_direct_work {
let repo = loopflow::lf::commands::util::find_repo_root()?;
let mut binding = if let Some(selector) = cli.as_work.as_deref() {
validate_work_selector(selector)?;
prepare_work_binding(selector, &repo)?
} else {
prepare_hierarchical_work_binding(
cli.task.as_deref(),
cli.project.as_deref(),
cli.wave.as_deref(),
&repo,
)?
};
if let Some(cwd) = cli.bound_cwd.clone() {
binding.cwd = cwd;
}
if let Some(wave) = &explicit_wave {
if wave.id() != &binding.wave_id {
anyhow::bail!(
"--wave {} does not own --as {}:{}",
wave.name(),
binding.work.kind(),
binding.work.id(),
);
}
}
cli.wave = Some(binding.wave_name.clone());
if cli.model.is_none() {
cli.model = binding.agent.clone();
}
let cwd = CwdGuard::enter(&binding.cwd)?;
require_bound_invocation(&cli.command, &binding.cwd)?;
direct_binding = Some(binding);
_bound_cwd = Some(cwd);
}
let result = if cli.list {
in_repo_runtime(&args, |_| loopflow::lf::commands::list::show_all())
} else {
match &cli.command {
Some(Commands::Inline { prompt }) => {
let text = prompt.join(" ");
in_repo_runtime(&args, |_| match direct_binding.as_ref() {
Some(binding) => {
loopflow::lf::commands::run::run_bound(None, Some(&text), &cli, binding)
}
None => loopflow::lf::commands::run::run(None, Some(&text), &cli),
})
}
Some(Commands::Desktop) => loopflow::lf::commands::desktop::run(),
Some(Commands::ProviderSession) => {
loopflow::lf::commands::runs::observe_provider_session()
}
Some(Commands::Ask { ask }) => loopflow::lf::commands::ask::run(ask),
Some(Commands::Session { cmd }) => loopflow::lf::commands::session::run(cmd),
Some(Commands::Pr { cmd }) => in_repo_runtime(&args, |_| {
loopflow::lf::commands::ops::run_pr(cmd.as_ref(), cli.model.as_deref())
}),
Some(Commands::Wt { cmd }) => {
in_repo_runtime(&args, |_| loopflow::lf::commands::ops::run_wt(cmd))
}
Some(Commands::Rebase {
plan,
manual,
continue_rebase,
abort,
adopt,
onto,
}) => in_repo_runtime(&args, |_| {
loopflow::lf::commands::ops::run_rebase(
onto.as_deref(),
*plan,
*manual,
*continue_rebase,
*abort,
*adopt,
)
}),
Some(Commands::Commit {
message,
push,
no_add,
}) => in_repo_runtime(&args, |_| {
loopflow::lf::commands::ops::run_commit(
message.as_deref(),
*push,
*no_add,
cli.model.as_deref(),
)
}),
Some(Commands::Auth { cmd }) => {
in_repo_runtime(&args, |_| loopflow::lf::commands::auth::run(cmd))
}
Some(Commands::Profile { cmd }) => in_repo_runtime(&args, |repo| {
loopflow::lf::commands::profile::run(cmd, repo)
}),
Some(Commands::Route { cmd }) => in_repo_runtime(&args, |repo| {
loopflow::lf::commands::profile::run_route(cmd, repo)
}),
Some(Commands::Release { cmd }) => {
in_repo_runtime(&args, |_| loopflow::lf::commands::ops::run_release(cmd))
}
Some(Commands::Pm { cmd }) => {
in_repo_runtime(&args, |_| loopflow::lf::commands::ops::run_pm(cmd))
}
Some(Commands::Home { cmd }) => {
in_repo_runtime(&args, |repo| loopflow::lf::commands::home::run(cmd, repo))
}
Some(Commands::SyncSkills { yes, no_prune }) => in_repo_runtime(&args, |_| {
loopflow::lf::commands::ops::run_sync_skills(*yes, *no_prune)
}),
Some(Commands::Cron { cmd }) => {
in_repo_runtime(&args, |_| loopflow::lf::commands::ops::cron_cmd(cmd))
}
Some(Commands::Wave { name, force }) => {
in_repo_runtime(&args, |_| loopflow::controller::wave::run(name, *force))
}
Some(Commands::Start {
waves,
wave_ids,
json,
}) => in_repo_runtime(&args, |repo| {
loopflow::lf::commands::home::start(waves, wave_ids, *json, repo)
}),
Some(Commands::Stop { name }) => {
in_repo_runtime(&args, |repo| loopflow::lf::commands::home::stop(name, repo))
}
Some(Commands::Pause { name, json }) => in_repo_runtime(&args, |repo| {
loopflow::lf::commands::wave_intent::run(name, true, *json, repo)
}),
Some(Commands::Resume { name, json }) => in_repo_runtime(&args, |repo| {
loopflow::lf::commands::wave_intent::run(name, false, *json, repo)
}),
Some(Commands::Resident { name }) => {
in_repo_runtime(&args, |_| loopflow::controller::wave::resident::run(name))
}
Some(Commands::FlowStep { flow, index, seed }) => in_repo_runtime(&args, |repo| {
loopflow::lf::commands::flow::run_step(flow, *index, seed, &cli, repo)
}),
Some(Commands::Project {
cmd: ProjectCommand::Promote { slug, wave },
}) => in_repo_runtime(&args, |repo| {
let parent = loopflow::work::wave::context::resolve_managed_wave_sync(
Some(repo),
wave.as_deref(),
)
.map(|wave| wave.name().to_string())
.map_err(|err| match err {
loopflow::work::wave::context::WaveResolveError::NoContext => {
anyhow::anyhow!("cannot determine parent wave; pass --wave <name>")
}
other => anyhow::Error::from(other),
})?;
let message = format!(
"Promote project '{slug}' from parent wave '{parent}'. Complete the authored migration, PM move, parent link, and residency checks."
);
loopflow::ops::project::prepare_promotion(repo, &parent, slug)
.map_err(anyhow::Error::from)?;
run_target("project-promote", Some(&message), &cli, &args)?;
let residency = loopflow::ops::project::complete_promotion(repo, &parent, slug)
.map_err(anyhow::Error::from)?;
println!("promoted {slug} from {parent}; residency: {residency}");
Ok(())
}),
Some(Commands::Project { cmd }) => {
in_repo_runtime(&args, |repo| run_project_command(repo, cmd))
}
Some(Commands::Task { cmd }) => {
in_repo_runtime(&args, |repo| run_task_command(repo, cmd))
}
Some(Commands::Work { cmd }) => {
in_repo_runtime(&args, |repo| loopflow::lf::commands::work::run(cmd, repo))
}
Some(Commands::WorkRunner { kind, work_id }) => {
run_work_runner_entrypoint(&args, kind, work_id, |work| {
tokio::runtime::Runtime::new()?.block_on(loopflow::controller::run_work(work))
})
}
Some(Commands::Tokens { json, days }) => {
loopflow::lf::commands::tokens::run(*json, *days)
}
Some(Commands::Usage {
json,
days,
wave,
project,
task,
}) => loopflow::lf::commands::usage::run(
*json,
*days,
wave.as_deref(),
project.as_deref(),
task.as_deref(),
),
Some(Commands::TelemetryScorecard { json }) => in_repo_runtime(&args, |repo| {
let item = loopflow::engine::flow::Op {
command: "__telemetry-scorecard".to_string(),
args: if *json {
vec!["--json".to_string()]
} else {
Vec::new()
},
};
loopflow::ops::execute_flow_ops(repo, &item, &loopflow::ops::NullProgress)
.map_err(Into::into)
}),
Some(Commands::Ci {
since,
wave,
repo,
json,
}) => loopflow::lf::commands::ci::run(since, wave.as_deref(), repo.as_deref(), *json),
Some(Commands::Ps { json }) => loopflow::lf::commands::top::run_ps(*json),
Some(Commands::Top { json }) => loopflow::lf::commands::top::run_top(*json),
Some(Commands::Prune { dry_run, json }) => {
loopflow::lf::commands::top::run_prune(*json, *dry_run)
}
Some(Commands::Doctor { json }) => loopflow::lf::commands::doctor::run(*json),
Some(Commands::List) => {
in_repo_runtime(&args, |_| loopflow::lf::commands::list::show_all())
}
Some(Commands::Ls { json, all }) => loopflow::lf::commands::waves::ls(*json, *all),
Some(Commands::Status { wave, json }) => {
loopflow::lf::commands::waves::status(wave.as_deref(), *json)
}
Some(Commands::Roadmap { wave, json, all }) => {
loopflow::lf::commands::waves::roadmap(wave.as_deref(), *json, *all)
}
Some(Commands::Activity {
since,
limit,
wave,
project,
task,
json,
}) => loopflow::lf::commands::activity::run(
since,
*limit,
wave.as_deref(),
project.as_deref(),
task.as_deref(),
*json,
),
Some(Commands::Runs {
run,
events,
resume,
task,
project,
wave,
json,
}) => loopflow::lf::commands::runs::list(
*json,
wave.as_deref(),
project.as_deref(),
task.as_deref(),
run.as_deref(),
*events,
*resume,
),
Some(Commands::Replay { run }) => loopflow::lf::commands::replay::run(run),
Some(Commands::Reply {
wave,
text,
agent,
max_turns,
}) => loopflow::lf::commands::reply::run(wave, text, agent.clone(), *max_turns),
Some(Commands::Chat {
text,
follow,
steer,
history,
json,
limit,
epoch,
target,
}) => loopflow::lf::commands::chat::run(
text,
loopflow::lf::commands::chat::ChatOptions {
follow: *follow,
steer: *steer,
history: *history,
json: *json,
limit: *limit,
epoch: epoch.as_deref(),
},
target,
),
Some(Commands::Install { .. }) => {
unreachable!("install dispatches before home routing")
}
Some(Commands::Screenshot { .. } | Commands::ScreenshotSupervisor { .. }) => {
unreachable!("screenshot dispatches before home routing")
}
Some(Commands::RetiredOp { .. }) => unreachable!("retired op cannot parse"),
Some(Commands::Ssh {
target,
repo,
secret,
forward_agent,
origin_account: _,
origin_only_account: _,
lf_args,
}) => loopflow::lf::commands::ssh::run(
target,
repo.as_deref(),
secret,
*forward_agent,
&account_selection,
lf_args,
),
Some(Commands::Flow { name, args: rest }) => {
if matches!(name.as_str(), "show" | "validate") {
let target = rest
.first()
.ok_or_else(|| anyhow::anyhow!("usage: lf flow {name} FLOW"))?;
if rest.len() != 1 {
return Err(anyhow::anyhow!("usage: lf flow {name} FLOW"));
}
return in_repo_runtime(&args, |repo| match name.as_str() {
"show" => loopflow::lf::commands::flow::show(target, repo),
"validate" => loopflow::lf::commands::flow::validate(target, repo),
_ => unreachable!(),
});
}
require_target_kind(name, TargetKind::Flow)?;
let message = join_args(rest);
run_target(name, message.as_deref(), &cli, &args)
}
Some(Commands::Skill { name, args: rest }) => {
require_target_kind(name, TargetKind::Skill)?;
let message = join_args(rest);
if let Some(binding) = direct_binding.as_ref() {
run_bound_target_in_repo(
&binding.cwd,
name,
message.as_deref(),
&cli,
&args,
binding,
)
} else {
run_target(name, message.as_deref(), &cli, &args)
}
}
Some(Commands::External(external_args)) => {
match loopflow::lf::commands::run::split_skill_args(external_args) {
Ok((name, skill_args)) => {
let message = join_args(&skill_args);
if let Some(binding) = direct_binding.as_ref() {
run_bound_target_in_repo(
&binding.cwd,
&name,
message.as_deref(),
&cli,
&args,
binding,
)
} else {
run_target(&name, message.as_deref(), &cli, &args)
}
}
Err(err) => Err(err),
}
}
None => run_default_agent(&cli, &args),
}
};
result
}
#[cfg(test)]
mod tests {
use super::{
arg_tables, format_task_pr_line, normalize_ssh_args, reorder_args, validate_work_selector,
CwdGuard, EnvGuard,
};
use clap::Parser;
use loopflow::lf::{Cli, Commands, PmCommand, PmTaskCommand, PrCommand};
use loopflow::work::task::{GithubPr, PrPublication, TaskId, TaskPr, TaskPrId};
#[test]
fn bound_work_selector_requires_a_readable_registry() {
let _lock = PROCESS_STATE_LOCK
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let directory = tempfile::tempdir().unwrap();
let registry_blocker = directory.path().join("unreadable-registry");
std::fs::create_dir(®istry_blocker).unwrap();
let _db = EnvGuard::set("LF_DB_PATH", registry_blocker.display().to_string());
validate_work_selector("task:LOO-265").unwrap();
let error = super::prepare_work_binding("task:LOO-265", directory.path())
.expect_err("--as must not degrade to raw attribution");
assert!(error.to_string().contains("planning registry unavailable"));
}
static PROCESS_STATE_LOCK: std::sync::Mutex<()> = std::sync::Mutex::new(());
fn published_pr() -> TaskPr {
let now = time::OffsetDateTime::now_utc();
TaskPr {
id: TaskPrId::new(),
task_id: TaskId::new(),
sequence: 1,
slug: "linear-pr-linkage".to_string(),
branch: "jack/linear-pr-linkage".to_string(),
base_commit: "abc".to_string(),
parent_pr_id: None,
publication: Some(PrPublication {
requested_at: now,
presentation: None,
github: Some(GithubPr {
number: 931,
url: "https://github.com/loopflowstudio/loopflow/pull/931".to_string(),
head_sha: None,
}),
merge: 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,
}
}
#[test]
fn only_bare_lf_selects_the_default_agent_surface() {
let bare = Cli::try_parse_from(["lf"]).unwrap();
assert!(bare.command.is_none());
let code = Cli::try_parse_from(["lf", "code"]).unwrap();
assert!(matches!(
code.command,
Some(Commands::External(parts)) if parts == ["code"]
));
}
#[test]
fn task_pr_line_is_quiet_when_the_linear_link_is_healthy() {
let line = format_task_pr_line(&published_pr());
assert!(line.contains("GitHub #931"), "{line}");
assert!(!line.contains("Linear"), "{line}");
}
#[test]
fn task_pr_line_names_a_degraded_linear_link() {
let mut pr = published_pr();
pr.linear_link_error = Some("linear token expired".to_string());
let line = format_task_pr_line(&pr);
assert!(line.contains("Linear link degraded"), "{line}");
assert!(line.contains("linear token expired"), "{line}");
assert!(line.contains("GitHub #931"), "{line}");
}
#[test]
fn derived_tables_cover_commands_flags_and_aliases() {
let tables = arg_tables();
for command in [
":",
"desktop",
"screenshot",
"pr",
"wt",
"rebase",
"commit",
"auth",
"release",
"pm",
"task",
"project",
"flow",
"skill",
"chat",
"usage",
"top",
"list",
"ls",
"status",
"runs",
"help",
] {
assert!(tables.commands.contains_key(command), "command {command}");
}
for flag in [
"-d",
"-D",
"--direction",
"--docs",
"-m",
"-M",
"--model",
"--max-turns",
"-w",
"-W",
"--wave",
"--project",
"--task",
] {
assert!(tables.top_level.value.contains(flag), "value flag {flag}");
}
for flag in [
"-l",
"--list",
"-c",
"-C",
"--clipboard",
"--no-direction",
"--yolo",
"-i",
"-I",
"-b",
"-B",
"--tui",
"--ide",
"--chrome",
"--no-chrome",
"--diff-files",
"--no-diff-files",
"--diff",
"--no-diff",
"--no-loopflow",
"-h",
"--help",
"-V",
"--version",
] {
assert!(tables.top_level.boolean.contains(flag), "bool flag {flag}");
}
}
#[test]
fn reorder_args_uppercase_value_alias_after_skill() {
let args = vec![
"lf".to_string(),
"debug".to_string(),
"-M".to_string(),
"codex".to_string(),
];
assert_eq!(reorder_args(args), vec!["lf", "-M", "codex", "debug"]);
}
#[test]
fn reorder_args_moves_work_selector_after_skill_to_global_position() {
let args = vec![
"lf".to_string(),
"implement".to_string(),
"--task".to_string(),
"LOO-123".to_string(),
"--project".to_string(),
"context".to_string(),
];
assert_eq!(
reorder_args(args),
vec![
"lf",
"--task",
"LOO-123",
"--project",
"context",
"implement"
]
);
}
#[test]
fn bound_cwd_is_entered_for_the_invocation_and_restored_afterward() {
let _lock = PROCESS_STATE_LOCK.lock().unwrap();
let previous = std::env::current_dir().unwrap();
let directory = tempfile::tempdir().unwrap();
let guard = CwdGuard::enter(directory.path()).unwrap();
assert_eq!(
std::env::current_dir().unwrap(),
directory.path().canonicalize().unwrap()
);
drop(guard);
assert_eq!(std::env::current_dir().unwrap(), previous);
}
#[test]
fn reorder_args_uppercase_bool_alias_after_skill() {
let args = vec!["lf".to_string(), "debug".to_string(), "-C".to_string()];
assert_eq!(reorder_args(args), vec!["lf", "-C", "debug"]);
}
#[test]
fn desktop_remains_an_explicit_app_command() {
let cli = Cli::try_parse_from(["lf", "desktop"]).unwrap();
assert!(matches!(cli.command, Some(Commands::Desktop)));
}
#[test]
fn bare_lf_has_a_terminal_control_skill() {
let cli = Cli::try_parse_from(["lf"]).unwrap();
assert!(cli.command.is_none());
let skill = loopflow::engine::builtins::get_builtin_skill("loopflow")
.expect("builtin terminal control skill");
assert!(skill.contains("lf session list --json"));
assert!(skill.contains("lf session open <session-id> --json"));
assert!(skill.contains("Keep this conversation open"));
}
#[test]
fn reorder_args_flag_after_skill() {
let args = vec!["lf".to_string(), "debug".to_string(), "-c".to_string()];
let result = reorder_args(args);
assert_eq!(result, vec!["lf", "-c", "debug"]);
}
#[test]
fn wave_and_resident_are_distinct_entrypoints() {
let served = Cli::try_parse_from(["lf", "wave", "goals"]).unwrap();
assert!(matches!(
served.command,
Some(Commands::Wave { name, force: false }) if name == "goals"
));
let forced = Cli::try_parse_from(["lf", "wave", "goals", "--force"]).unwrap();
assert!(matches!(
forced.command,
Some(Commands::Wave { force: true, .. })
));
let stopped = Cli::try_parse_from(["lf", "stop", "goals"]).unwrap();
assert!(matches!(
stopped.command,
Some(Commands::Stop { name }) if name == "goals"
));
let paused = Cli::try_parse_from(["lf", "pause", "goals", "--json"]).unwrap();
assert!(matches!(
paused.command,
Some(Commands::Pause { name, json: true }) if name == "goals"
));
let resumed = Cli::try_parse_from(["lf", "resume", "goals", "--json"]).unwrap();
assert!(matches!(
resumed.command,
Some(Commands::Resume { name, json: true }) if name == "goals"
));
let body = Cli::try_parse_from(["lf", "__resident", "goals"]).unwrap();
assert!(matches!(
body.command,
Some(Commands::Resident { name }) if name == "goals"
));
}
#[test]
fn start_accepts_local_and_home_bound_waves() {
let start =
Cli::try_parse_from(["lf", "start", "product", "intelligence", "--json"]).unwrap();
assert!(matches!(
start.command,
Some(Commands::Start { waves, wave_ids, json })
if waves == ["product", "intelligence"] && wave_ids.is_empty() && json
));
let remote_start = Cli::try_parse_from([
"lf",
"start",
"product",
"--wave-id",
"product=wave_00000000000000000000000000000001",
"--json",
])
.unwrap();
assert!(matches!(
remote_start.command,
Some(Commands::Start { waves, wave_ids, json })
if waves == ["product"]
&& wave_ids == ["product=wave_00000000000000000000000000000001"]
&& json
));
let ssh = Cli::try_parse_from([
"lf",
"ssh",
"home_00000000000000000000000000000001",
"start",
"product",
])
.unwrap();
assert!(matches!(
ssh.command,
Some(Commands::Ssh { lf_args, .. }) if lf_args == ["start", "product"]
));
}
#[test]
fn ssh_help_prefers_home_identity() {
let help = Cli::try_parse_from(["lf", "ssh", "--help"])
.expect_err("help exits through clap")
.to_string();
assert!(help.contains("<TARGET>"));
assert!(help.contains("HomeId (preferred), SSH alias, or user@host"));
}
#[test]
fn old_serve_surface_is_no_longer_a_builtin_command() {
let cli = Cli::try_parse_from(["lf", "serve", "goals"]).expect("falls through to external");
assert!(
matches!(cli.command, Some(Commands::External(parts)) if parts[0] == "serve"),
"`serve` survives only as an external verb, not a built-in"
);
assert!(matches!(
Cli::try_parse_from(["lf", "wave", "goals"])
.expect("the replacement")
.command,
Some(Commands::Wave { .. })
));
}
#[test]
fn retired_op_namespace_names_its_replacement() {
let removed = Cli::try_parse_from(["lf", "op", "next"])
.expect_err("`lf op next` cannot parse")
.to_string();
assert!(
removed.contains("no replacement") && removed.contains("lf task run"),
"`lf op next` should state the removal and how work is dispatched now: {removed}"
);
let landed = Cli::try_parse_from(["lf", "op", "land"])
.expect_err("`lf op land` cannot parse")
.to_string();
assert!(
landed.contains("`lf pr land`"),
"`lf op land` should name `lf pr land`: {landed}"
);
let bare = Cli::try_parse_from(["lf", "op"])
.expect_err("bare `lf op` cannot parse")
.to_string();
assert!(
bare.contains("top-level"),
"bare `lf op` should say the operations are top-level: {bare}"
);
}
#[test]
fn removed_dispatch_flag_is_rejected() {
assert!(Cli::try_parse_from(["lf", "--dispatch", "implement", "ship it"]).is_err());
}
#[test]
fn reorder_args_flag_before_skill() {
let args = vec!["lf".to_string(), "-c".to_string(), "debug".to_string()];
let result = reorder_args(args);
assert_eq!(result, vec!["lf", "-c", "debug"]);
}
#[test]
fn reorder_args_value_flag_before_skill() {
let args = vec![
"lf".to_string(),
"-m".to_string(),
"codex".to_string(),
"implement".to_string(),
];
let result = reorder_args(args);
assert_eq!(result, vec!["lf", "-m", "codex", "implement"]);
}
#[test]
fn reorder_args_value_flag_after_skill() {
let args = vec![
"lf".to_string(),
"debug".to_string(),
"-m".to_string(),
"codex".to_string(),
];
let result = reorder_args(args);
assert_eq!(result, vec!["lf", "-m", "codex", "debug"]);
}
#[test]
fn reorder_args_repeatable_account_flags_after_skill() {
let args = [
"lf",
"implement",
"--account",
"claude=personal",
"--account",
"codex=reserve",
]
.map(String::from)
.to_vec();
assert_eq!(
reorder_args(args),
vec![
"lf",
"--account",
"claude=personal",
"--account",
"codex=reserve",
"implement"
]
);
}
#[test]
fn reorder_args_mixed_flags() {
let args = vec![
"lf".to_string(),
"-i".to_string(),
"implement".to_string(),
"-c".to_string(),
"-m".to_string(),
"claude".to_string(),
];
let result = reorder_args(args);
assert_eq!(result, vec!["lf", "-i", "-c", "-m", "claude", "implement"]);
}
#[test]
fn reorder_args_no_direction_flag_after_skill() {
let args = vec![
"lf".to_string(),
"implement".to_string(),
"--no-direction".to_string(),
];
let result = reorder_args(args);
assert_eq!(result, vec!["lf", "--no-direction", "implement"]);
}
#[test]
fn reorder_args_no_loopflow_flag_after_skill() {
let args = vec![
"lf".to_string(),
"gate".to_string(),
"--no-loopflow".to_string(),
];
let result = reorder_args(args);
assert_eq!(result, vec!["lf", "--no-loopflow", "gate"]);
}
#[test]
fn reorder_args_skill_with_args() {
let args = vec![
"lf".to_string(),
"implement:".to_string(),
"add".to_string(),
"logout".to_string(),
"-c".to_string(),
];
let result = reorder_args(args);
assert_eq!(result, vec!["lf", "-c", "implement:", "add", "logout"]);
}
#[test]
fn reorder_args_known_command_unchanged() {
let args = vec![
"lf".to_string(),
"commit".to_string(),
"-m".to_string(),
"msg".to_string(),
];
let result = reorder_args(args);
assert_eq!(result, vec!["lf", "commit", "-m", "msg"]);
}
#[test]
fn reorder_args_preserves_the_ssh_target_boundary() {
let args = vec![
"lf".to_string(),
"ssh".to_string(),
"build-vm".to_string(),
"--account".to_string(),
"remote@company".to_string(),
"task".to_string(),
"pursue".to_string(),
];
assert_eq!(reorder_args(args.clone()), args);
}
#[test]
fn normalize_ssh_args_makes_the_target_a_hard_boundary() {
let args = [
"lf",
"ssh",
"--account",
"origin@example.com",
"build-vm",
"--account",
"remote@example.com",
"task",
"pursue",
]
.map(str::to_string)
.to_vec();
let normalized = normalize_ssh_args(args);
assert_eq!(
normalized,
[
"lf",
"ssh",
"--account",
"origin@example.com",
"build-vm",
"--",
"--account",
"remote@example.com",
"task",
"pursue",
]
);
let cli = Cli::try_parse_from(normalized).expect("parse normalized SSH command");
assert!(matches!(
cli.command,
Some(Commands::Ssh {
origin_account,
lf_args,
..
}) if origin_account == ["origin@example.com"]
&& lf_args == ["--account", "remote@example.com", "task", "pursue"]
));
}
#[test]
fn reorder_args_preserves_local_collision_after_leading_global() {
let args: Vec<String> = ["lf", "--wave", "goals", "commit", "-m", "ship it"]
.map(String::from)
.to_vec();
assert_eq!(
reorder_args(args),
vec!["lf", "--wave", "goals", "commit", "-m", "ship it"]
);
}
#[test]
fn reorder_args_stops_at_double_dash() {
let args: Vec<String> = ["lf", "skill", "debug", "--", "--wave", "literal"]
.map(String::from)
.to_vec();
assert_eq!(
reorder_args(args),
vec!["lf", "skill", "debug", "--", "--wave", "literal"]
);
}
#[test]
fn reorder_args_moves_flags_to_nested_owners() {
let args: Vec<String> = ["lf", "pm", "--wave", "systems", "show"]
.map(String::from)
.to_vec();
let reordered = reorder_args(args);
assert_eq!(reordered, vec!["lf", "pm", "show", "--wave", "systems"]);
assert!(matches!(
Cli::try_parse_from(reordered).unwrap().command,
Some(Commands::Pm {
cmd: PmCommand::Show { .. }
})
));
let args: Vec<String> = [
"lf",
"pm",
"task",
"--wave",
"systems",
"create",
"--project",
"wave-chat",
"--title",
"file it",
]
.map(String::from)
.to_vec();
let reordered = reorder_args(args);
assert_eq!(
reordered,
vec![
"lf",
"pm",
"task",
"create",
"--wave",
"systems",
"--project",
"wave-chat",
"--title",
"file it"
]
);
assert!(matches!(
Cli::try_parse_from(reordered).unwrap().command,
Some(Commands::Pm {
cmd: PmCommand::Task {
cmd: PmTaskCommand::Create { .. }
}
})
));
let args: Vec<String> = ["lf", "pr", "-m", "codex", "open"]
.map(String::from)
.to_vec();
let reordered = reorder_args(args);
assert_eq!(reordered, vec!["lf", "pr", "open", "-m", "codex"]);
assert!(matches!(
Cli::try_parse_from(reordered).unwrap().command,
Some(Commands::Pr {
cmd: Some(PrCommand::Open { .. })
})
));
let args: Vec<String> = ["lf", "pr", "--strict", "submit"]
.map(String::from)
.to_vec();
let reordered = reorder_args(args);
assert_eq!(reordered, vec!["lf", "pr", "submit", "--strict"]);
assert!(matches!(
Cli::try_parse_from(reordered).unwrap().command,
Some(Commands::Pr {
cmd: Some(PrCommand::Submit { strict: true, .. })
})
));
let args: Vec<String> = ["lf", "wt", "--force", "rm", "old-tree"]
.map(String::from)
.to_vec();
assert_eq!(
reorder_args(args),
vec!["lf", "wt", "rm", "--force", "old-tree"]
);
}
#[test]
fn reorder_args_leaves_explicit_targeting_alone() {
let args: Vec<String> = ["lf", "chat", "--wave", "systems", "shipped it"]
.map(String::from)
.to_vec();
assert_eq!(
reorder_args(args),
vec!["lf", "chat", "--wave", "systems", "shipped it"]
);
let args: Vec<String> = ["lf", "pm", "show", "--wave", "systems"]
.map(String::from)
.to_vec();
assert_eq!(
reorder_args(args),
vec!["lf", "pm", "show", "--wave", "systems"]
);
}
#[test]
fn reorder_args_no_skill() {
let args = vec!["lf".to_string(), "-l".to_string()];
let result = reorder_args(args);
assert_eq!(result, vec!["lf", "-l"]);
}
}