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;
use tracing_subscriber::EnvFilter;
use loopflow::journal::{self, with_runtime, LfEventFields, LfEventType, LfNode};
use loopflow::lf::{Cli, Commands, InstallCommand, TaskCommand, WaveCommand};
use loopflow::ops::chapter::{
chapter_history, empty_plan, new_chapter, update_plan, NewChapterRequest,
};
use loopflow::ops::task_execution::TaskExecutionState;
use loopflow::work::chapter::{ChapterId, ChapterPhase};
#[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 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 in_directory_runtime<T>(
command: &[String],
run: impl FnOnce(&Path) -> anyhow::Result<T>,
) -> anyhow::Result<T> {
let directory = loopflow::repo::working_directory()?;
with_runtime(&directory, command, || run(&directory))
}
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 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 run_target(
name: &str,
kind: Option<TargetKind>,
message: Option<&str>,
cli: &Cli,
command: &[String],
binding: Option<&loopflow::ops::WorkBinding>,
) -> anyhow::Result<()> {
use loopflow::lf::discovery::{discover_skill, discover_target, Target};
let repo_root = loopflow::repo::working_directory()?;
let target = match kind {
Some(TargetKind::Skill) => Target::Skill(discover_skill(&repo_root, name)?),
Some(TargetKind::Flow) => Target::Flow(loopflow::engine::load_flow(name, &repo_root)?),
None => discover_target(&repo_root, name)?,
};
with_runtime(&repo_root, command, || match target {
Target::Skill(_) => with_skill_runtime(&repo_root, name, || {
match binding {
Some(binding) => {
loopflow::lf::commands::run::run_bound(Some(name), message, cli, binding)?
}
None => loopflow::lf::commands::run::run(Some(name), message, cli)?,
}
if binding.is_none() && std::env::var_os(loopflow::durable::RUN_ID_ENV).is_none() {
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(())
}),
Target::Flow(flow) => {
loopflow::lf::commands::flow::run(&flow, message, cli, &repo_root, binding)
}
})
}
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>,
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, 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> or wave:<selector>")
})?;
if !matches!(kind, "task" | "wave") {
anyhow::bail!("invalid --as kind {kind:?}; expected task or wave");
}
if value.trim().is_empty() {
anyhow::bail!("{kind} selector cannot be empty");
}
Ok(())
}
fn require_bound_invocation(command: &Option<Commands>) -> anyhow::Result<()> {
match command {
Some(Commands::Skill { .. } | Commands::Flow { .. } | Commands::External(_) | Commands::Inline { .. }) => Ok(()),
_ => anyhow::bail!("Work selectors run a skill, flow or inline prompt; use `lf --task LOO-123 design`, `lf --task LOO-123 code` or `lf --wave product : \"question\"`"),
}
}
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)?;
print_task_snapshot(&snapshot, json)
}
fn print_task_snapshot(
snapshot: &loopflow::ops::task::TaskSnapshot,
json: bool,
) -> anyhow::Result<()> {
if json {
println!("{}", serde_json::to_string_pretty(&snapshot)?);
} else {
let pm_writeback = match &snapshot.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 body = format!(
"agent {}, provider {}",
snapshot.agent.as_deref().unwrap_or("default"),
snapshot.provider
);
let status = match snapshot.execution.state {
TaskExecutionState::Blocked => "blocked".to_string(),
TaskExecutionState::Stalled => "stalled".to_string(),
_ => snapshot.status.to_string(),
};
println!(
"{} {}\n task: {}\n body: {}\n worktree: {}\n branch: {}\n PM writeback: {}",
snapshot.issue_identifier,
status,
snapshot.task_id,
body,
snapshot.worktree,
branch,
pm_writeback,
);
println!(" execution: {}", snapshot.execution.reason);
if let Some(run) = &snapshot.execution.run_id {
println!(" worker Run: {run}");
}
for run in &snapshot.runs {
println!(
" Run: {} {} {} {}",
run.id,
run.label(),
run.surface,
run.status()
);
}
if snapshot.runs_truncated {
println!(" Run history truncated; inspect exact Run IDs for older evidence");
}
println!(" project: {}", snapshot.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 run_wave_command(repo: &Path, command: &WaveCommand) -> anyhow::Result<()> {
let (receipt, json, dry_run) = match command {
WaveCommand::List { .. } | WaveCommand::Status { .. } => {
unreachable!("read commands dispatch without runtime capture")
}
WaveCommand::Probe { wave, json } => {
return loopflow::lf::commands::home::probe_cmd(wave.as_deref(), *json, Some(repo));
}
WaveCommand::Connect {
wave,
wave_flag,
all,
team_key,
team_name,
} => {
return loopflow::lf::commands::ops::connect_wave(
repo,
wave.as_deref().or(wave_flag.as_deref()),
*all,
team_key.as_deref(),
team_name.as_deref(),
)
}
WaveCommand::Sync {
wave,
wave_flag,
all,
} => {
return loopflow::lf::commands::ops::sync_planning(
repo,
wave.as_deref().or(wave_flag.as_deref()),
*all,
false,
false,
)
}
WaveCommand::Rename { wave, title } => {
return loopflow::lf::commands::ops::rename_wave(repo, wave, title)
}
WaveCommand::Forget { .. }
| WaveCommand::Place { .. }
| WaveCommand::Relocate { .. }
| WaveCommand::Retire { .. } => {
return loopflow::lf::commands::placement::wave(repo, command)
}
WaveCommand::Recover {
name,
cancel,
reason,
} => {
let repo = loopflow::engine::worktrees::main_repo_root(repo)?;
let wave = loopflow::ops::normalize_wave_name(name)
.ok_or_else(|| anyhow::anyhow!("invalid wave name: '{name}'"))?;
let report = match cancel {
Some(seq) => loopflow::controller::wave::recovery::cancel(
&repo,
&wave,
*seq,
reason
.as_deref()
.expect("Clap requires a cancellation reason"),
)?,
None => loopflow::controller::wave::recovery::inspect(&repo, &wave)?,
};
println!("{}", serde_json::to_string_pretty(&report)?);
return Ok(());
}
WaveCommand::NewChapter {
wave,
chapter,
plan,
dry_run,
json,
} => {
let content = match plan {
Some(path) => serde_json::from_slice(&std::fs::read(path)?)?,
None => empty_plan(),
};
let request = NewChapterRequest {
wave: wave.clone(),
chapter: ChapterId::parse(chapter).map_err(anyhow::Error::msg)?,
content,
};
(new_chapter(repo, &request, *dry_run)?, *json, *dry_run)
}
WaveCommand::History { wave, json } => {
let history = chapter_history(repo, wave.as_deref())?;
if *json {
println!("{}", serde_json::to_string_pretty(&history)?);
} else {
for chapter in history {
println!("{} {:?}", chapter.id.as_str(), chapter.phase);
}
}
return Ok(());
}
WaveCommand::UpdatePlan { wave, plan } => {
let content = serde_json::from_slice(&std::fs::read(plan)?)?;
update_plan(repo, wave.as_deref(), &content)?;
return Ok(());
}
};
if json {
println!("{}", serde_json::to_string_pretty(&receipt)?);
} else {
println!(
"{} · chapter {} · {:?}",
receipt.wave,
receipt.id.as_str(),
receipt.phase
);
for task in &receipt.tasks {
println!(
" {} {:?} {}",
task.task.identifier, task.disposition, task.reason
);
}
if let Some(error) = &receipt.error {
println!(" pending: {error}");
}
}
if !dry_run && receipt.phase != ChapterPhase::Complete {
anyhow::bail!("chapter transition is incomplete; retry the same chapter id");
}
Ok(())
}
fn read_task_draft() -> anyhow::Result<String> {
let limit = loopflow::ops::task::MAX_FILE_BYTES;
let mut bytes = Vec::new();
std::io::stdin()
.take(limit as u64 + 1)
.read_to_end(&mut bytes)?;
anyhow::ensure!(bytes.len() <= limit, "Draft exceeds 1 MB");
Ok(String::from_utf8(bytes)?)
}
fn run_task_command(repo: &Path, command: &TaskCommand, agent: Option<&str>) -> anyhow::Result<()> {
match command {
TaskCommand::Worker { .. } => unreachable!("Task worker dispatches at process entry"),
TaskCommand::Checkout {
issue,
name,
stack_on,
directive,
json,
} => {
let task = loopflow::ops::task::task_checkout(
repo,
issue,
loopflow::ops::task::TaskCheckoutOptions {
name: name.clone(),
stack_on: stack_on.clone(),
directive: directive.clone(),
},
)?;
print_task(&task, *json)
}
TaskCommand::Run {
issue,
name,
flow,
stack_on,
directive,
reason,
json,
} => {
let task = loopflow::ops::task::task_run(
repo,
issue,
loopflow::ops::task::TaskLaunchOptions {
reason: reason.clone(),
agent: agent.map(str::to_string),
name: name.clone(),
flow: flow.clone(),
stack_on: stack_on.clone(),
directive: directive.clone(),
},
)?;
print_task_snapshot(&task, *json)
}
TaskCommand::Create {
wave,
title,
notes,
run,
name,
flow,
stack_on,
json,
} => {
let report = match notes {
Some(notes) => Some(notes.clone()),
None => piped_task_report()?,
};
let result = loopflow::ops::task::task_create(
repo,
wave.as_deref(),
title.clone(),
report,
run.then(|| loopflow::ops::task::TaskLaunchOptions {
reason: None,
agent: agent.map(str::to_string),
name: name.clone(),
flow: flow.clone(),
stack_on: stack_on.clone(),
directive: None,
}),
)?;
match result {
loopflow::ops::task::TaskCreateResult::Started(task) => {
print_task_snapshot(&task, *json)
}
loopflow::ops::task::TaskCreateResult::Created(issue) => {
if *json {
println!("{}", serde_json::to_string_pretty(&issue)?);
} else {
println!("{} · {}", issue.identifier, issue.name);
}
Ok(())
}
}
}
TaskCommand::Status { issue, json } => {
let task = loopflow::ops::task::task_status(issue.as_deref())?;
print_task(&task, *json)
}
TaskCommand::Changes { issue, base, json } => {
let snapshot = loopflow::ops::task::task_changes(issue, base)?;
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,
base,
draft,
json,
} => {
let content = draft.then(read_task_draft).transpose()?;
let snapshot =
loopflow::ops::task::task_diff(issue, path.as_deref(), base, content.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,
recoveries,
json,
} => {
let snapshot = loopflow::ops::task::task_file(issue, path, *recoveries)?;
if *json {
println!("{}", serde_json::to_string_pretty(&snapshot)?);
} else if snapshot.state != loopflow::ops::task::TaskFileState::Text {
anyhow::bail!("{}: {:?}", snapshot.path, snapshot.state);
} else {
print!("{}", snapshot.content.as_deref().unwrap_or_default());
}
Ok(())
}
TaskCommand::Save {
issue,
path,
revision,
json,
} => {
let content = read_task_draft()?;
let saved = loopflow::ops::task::task_save(issue, path, revision, &content)?;
if *json {
println!("{}", serde_json::to_string_pretty(&saved)?);
} else {
println!("{}", saved.message);
}
Ok(())
}
TaskCommand::Complete {
issue,
summary,
json,
} => match loopflow::ops::task::task_complete(repo, issue, summary.clone())? {
Some(task) => print_task(&task, *json),
None => {
let resolved = loopflow::ops::pm::pm_resolve_task(repo, issue)?;
if *json {
println!("{}", serde_json::to_string_pretty(&resolved.item)?);
} else {
println!(
"{}: completed {}",
resolved.item.identifier, resolved.item.name
);
}
Ok(())
}
},
TaskCommand::Delete { issue } => {
let identifier = loopflow::ops::task::task_delete(repo, issue)?;
println!("{identifier}: deleted");
Ok(())
}
TaskCommand::Edit {
issue,
title,
notes,
wave,
} => {
let result = loopflow::ops::task::task_edit(
repo,
issue,
wave.as_deref(),
title.clone(),
notes.clone(),
)?;
println!("{}: updated task {}", result.wave, result.id);
Ok(())
}
TaskCommand::Comment {
issue,
message,
wave,
json,
} => {
let result = loopflow::ops::task::task_comment(
repo,
issue,
wave.as_deref(),
message.as_deref(),
)?;
if *json {
println!("{}", serde_json::to_string(&result)?);
} else if result.comments.is_empty() {
println!("{}: no comments", result.identifier);
} else {
for comment in &result.comments {
let author = match &comment.author {
loopflow::ops::pm::TaskCommentAuthor::Person { name } => {
name.as_deref().unwrap_or("unnamed person")
}
loopflow::ops::pm::TaskCommentAuthor::Integration => "integration",
};
let date = comment.created_at.as_deref().unwrap_or("date unavailable");
println!("── {author} · {date}\n{}\n", comment.body.trim_end());
}
}
Ok(())
}
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::Restart {
issue,
advice,
flow,
json,
} => {
let task = loopflow::ops::task::task_restart(
issue,
advice.clone(),
flow.clone(),
agent.map(str::to_string),
)?;
print_task_snapshot(&task, *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 { .. }));
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(),
)?;
if !matches!(
&cli.command,
Some(Commands::Screenshot { .. } | Commands::ScreenshotSupervisor { .. })
) {
loopflow::store::isolate_branch_data()?;
}
}
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.as_ref() {
None => loopflow::lf::commands::install::latest(),
Some(InstallCommand::Schedule { frequency }) => {
loopflow::lf::commands::install::schedule(*frequency)
}
Some(InstallCommand::RecoverSwitch { switch }) => {
loopflow::lf::commands::install::recover_switch(switch)
}
Some(InstallCommand::Preflight { json }) => {
loopflow::lf::commands::install::preflight(*json)
}
Some(InstallCommand::LocalPreflight { store, json }) => {
loopflow::lf::commands::install::local_preflight(store, *json)
}
Some(InstallCommand::AdvanceSwitch { switch }) => {
loopflow::lf::commands::install::advance_switch(switch)
}
Some(InstallCommand::Promote {
from_build,
coordinated_build,
fresh,
reuse_home,
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,
reuse_home.as_deref(),
),
Some(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.wave.is_some()
&& matches!(
&cli.command,
Some(Commands::Skill { .. })
| Some(Commands::Flow { .. })
| 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.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)?;
direct_binding = Some(binding);
_bound_cwd = Some(cwd);
}
let result = {
match &cli.command {
Some(Commands::Inline { prompt }) => {
let text = prompt.join(" ");
in_directory_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::User {
cmd: loopflow::lf::UserCommand::Name { json },
}) => {
let name = loopflow::engine::config::load_user_name()?;
if *json {
println!("{}", serde_json::to_string(&name)?);
} else if let Some(name) = name {
println!("{name}");
}
Ok(())
}
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,
}) => {
let repo =
loopflow::repo::require_repo_root(&std::env::current_dir()?, "lf rebase")?;
with_runtime(&repo, &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 }) => loopflow::lf::commands::auth::run(cmd),
Some(Commands::Release { cmd }) => {
in_repo_runtime(&args, |_| loopflow::lf::commands::ops::run_release(cmd))
}
Some(Commands::Repo { cmd }) => {
in_directory_runtime(&args, |_| loopflow::lf::commands::ops::run_repo(cmd))
}
Some(Commands::Home { cmd }) => loopflow::lf::commands::home::run(cmd),
Some(Commands::SyncSkills { yes, no_prune }) => {
loopflow::lf::commands::ops::run_sync_skills(*yes, *no_prune)
}
Some(Commands::Cron {
cmd:
cmd @ (loopflow::lf::CronCommand::List { .. }
| loopflow::lf::CronCommand::Remove { .. }),
}) => loopflow::lf::commands::ops::cron_cmd(cmd),
Some(Commands::Cron { cmd }) => {
in_repo_runtime(&args, |_| loopflow::lf::commands::ops::cron_cmd(cmd))
}
Some(Commands::Wave {
cmd: WaveCommand::List { json, all, current },
}) => loopflow::lf::commands::waves::ls(*json, *all, *current),
Some(Commands::Wave {
cmd:
WaveCommand::Status {
wave,
chapter,
json,
sync,
no_sync: _,
},
}) => {
if let Some(chapter) = chapter {
let id = loopflow::work::chapter::ChapterId::parse(chapter)
.map_err(anyhow::Error::msg)?;
let repo = loopflow::repo::find_repo_root()?;
let snapshot = loopflow::ops::chapter::chapter_snapshot(
&repo,
wave.as_deref(),
Some(&id),
)?;
if *json {
println!("{}", serde_json::to_string_pretty(&snapshot)?);
} else {
println!(
"{} · chapter {} · observed {}",
snapshot.wave,
snapshot.id.as_str(),
snapshot.observed_at
);
for kr in &snapshot.content.krs {
println!("[{}] {}", if kr.holds { "x" } else { " " }, kr.text);
}
println!("Metrics evaluated at {}", snapshot.metrics_evaluated_at);
print!(
"{}",
loopflow::lf::commands::waves::metric_portfolio_text(&snapshot.metrics)
);
for task in snapshot.tasks {
println!(
"{} {:?} {}",
task.task.identifier, task.disposition, task.reason
);
}
}
Ok(())
} else {
let refreshed = if *sync {
Some(loopflow::lf::commands::ops::refresh_status(
wave.as_deref(),
)?)
} else {
None
};
loopflow::lf::commands::waves::status(
refreshed.as_deref().or(wave.as_deref()),
*json,
)
}
}
Some(Commands::Wave {
cmd:
cmd @ (WaveCommand::Connect { .. }
| WaveCommand::Sync { .. }
| WaveCommand::Rename { .. }
| WaveCommand::Forget { .. }
| WaveCommand::Place { .. }
| WaveCommand::Relocate { .. }
| WaveCommand::Retire { .. }),
}) => in_directory_runtime(&args, |repo| run_wave_command(repo, cmd)),
Some(Commands::Wave { cmd }) => {
in_repo_runtime(&args, |repo| run_wave_command(repo, cmd))
}
Some(Commands::ChatConnect {
waves,
wave_ids,
json,
}) => in_repo_runtime(&args, |repo| {
loopflow::lf::commands::home::connect_chat(waves, wave_ids, *json, repo)
}),
Some(Commands::Resident { name }) => {
in_repo_runtime(&args, |_| loopflow::controller::wave::resident::run(name))
}
Some(Commands::Task {
cmd: TaskCommand::Worker { task_id },
}) => in_repo_runtime(&args, |_| {
tokio::runtime::Runtime::new()?
.block_on(loopflow::controller::task::run_worker(task_id.clone()))
}),
Some(Commands::Task {
cmd:
cmd @ (TaskCommand::Changes { .. }
| TaskCommand::Diff { .. }
| TaskCommand::File { .. }
| TaskCommand::Save { .. }),
}) => run_task_command(&std::env::current_dir()?, cmd, cli.model.as_deref()),
Some(Commands::Task { cmd }) => in_repo_runtime(&args, |repo| {
run_task_command(repo, cmd, cli.model.as_deref())
}),
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, planning }) => {
if *planning {
let repo = loopflow::repo::working_directory()?;
loopflow::lf::commands::ops::sync_planning(&repo, None, true, true, *json)
} else {
loopflow::lf::commands::doctor::run(*json)
}
}
Some(Commands::Catalog) => loopflow::lf::commands::list::show_all(),
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 {
active,
watch,
run,
parent,
events,
final_answer,
resume,
task,
wave,
project,
json,
}) => match run {
None if *active => {
loopflow::lf::commands::runs::list_active(*json, *watch, task.as_deref())
}
Some(run) if *resume => loopflow::lf::commands::runs::resume_run(run),
Some(run) => {
loopflow::lf::commands::runs::inspect(run, *events, *final_answer, *json)
}
None => loopflow::lf::commands::runs::list(
*json,
wave.as_deref(),
project.as_deref(),
task.as_deref(),
parent.as_deref(),
),
},
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,
history,
json,
limit,
epoch,
target,
}) => loopflow::lf::commands::chat::run(
text,
loopflow::lf::commands::chat::ChatOptions {
follow: *follow,
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,
json,
}) => {
if matches!(name.as_str(), "decide" | "route" | "blocked" | "resume") {
let directory = loopflow::repo::working_directory()?;
return loopflow::lf::commands::flow::control(name, rest, &cli, &directory);
}
if name == "list" {
if !rest.is_empty() {
return Err(anyhow::anyhow!("usage: lf flow list [--json]"));
}
let directory = loopflow::repo::working_directory()?;
return loopflow::lf::commands::flow::list(&directory, *json);
}
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"));
}
let directory = loopflow::repo::working_directory()?;
return match name.as_str() {
"show" => loopflow::lf::commands::flow::show(target, &directory),
"validate" => loopflow::lf::commands::flow::validate(target, &directory),
_ => unreachable!(),
};
}
let message = join_args(rest);
run_target(
name,
Some(TargetKind::Flow),
message.as_deref(),
&cli,
&args,
direct_binding.as_ref(),
)
}
Some(Commands::Skill { name, args: rest }) => {
let message = join_args(rest);
run_target(
name,
Some(TargetKind::Skill),
message.as_deref(),
&cli,
&args,
direct_binding.as_ref(),
)
}
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);
run_target(
&name,
None,
message.as_deref(),
&cli,
&args,
direct_binding.as_ref(),
)
}
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, PrCommand, TaskCommand, WaveCommand};
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",
"repo",
"task",
"flow",
"skill",
"chat",
"usage",
"top",
"catalog",
"wave",
"runs",
"help",
] {
assert!(tables.commands.contains_key(command), "command {command}");
}
for flag in [
"--docs",
"-m",
"-M",
"--model",
"--max-turns",
"-w",
"-W",
"--wave",
"--task",
] {
assert!(tables.top_level.value.contains(flag), "value flag {flag}");
}
for flag in [
"-c",
"-C",
"--clipboard",
"--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(),
"--wave".to_string(),
"context".to_string(),
];
assert_eq!(
reorder_args(args),
vec!["lf", "--task", "LOO-123", "--wave", "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 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", "__chat-connect", "product", "intelligence", "--json"])
.unwrap();
assert!(matches!(
start.command,
Some(Commands::ChatConnect { waves, wave_ids, json })
if waves == ["product", "intelligence"] && wave_ids.is_empty() && json
));
let remote_start = Cli::try_parse_from([
"lf",
"__chat-connect",
"product",
"--wave-id",
"product=wave_00000000000000000000000000000001",
"--json",
])
.unwrap();
assert!(matches!(
remote_start.command,
Some(Commands::ChatConnect { waves, wave_ids, json })
if waves == ["product"]
&& wave_ids == ["product=wave_00000000000000000000000000000001"]
&& json
));
let ssh = Cli::try_parse_from([
"lf",
"ssh",
"home_00000000000000000000000000000001",
"__chat-connect",
"product",
])
.unwrap();
assert!(matches!(
ssh.command,
Some(Commands::Ssh { lf_args, .. }) if lf_args == ["__chat-connect", "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"
);
}
#[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_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", "wave", "--wave", "systems", "sync"]
.map(String::from)
.to_vec();
let reordered = reorder_args(args);
assert_eq!(reordered, vec!["lf", "wave", "sync", "--wave", "systems"]);
assert!(matches!(
Cli::try_parse_from(reordered).unwrap().command,
Some(Commands::Wave {
cmd: WaveCommand::Sync { .. }
})
));
let args: Vec<String> = [
"lf", "task", "--wave", "systems", "create", "--title", "file it",
]
.map(String::from)
.to_vec();
let reordered = reorder_args(args);
assert_eq!(
reordered,
vec!["lf", "task", "create", "--wave", "systems", "--title", "file it"]
);
assert!(matches!(
Cli::try_parse_from(reordered).unwrap().command,
Some(Commands::Task {
cmd: TaskCommand::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", "wave", "status", "systems"]
.map(String::from)
.to_vec();
assert_eq!(reorder_args(args), vec!["lf", "wave", "status", "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"]);
}
}