#![cfg(not(test))]
mod acceptance;
mod agent;
mod ai_command_runner;
mod analyzer;
mod archive_layout;
mod bounded_git;
mod embedded_skills;
mod install_skills;
mod cli;
mod client;
mod command_queue;
mod completion;
mod config;
mod dependency_targets;
mod error;
mod error_history;
mod events;
mod execution;
mod history;
mod hooks;
mod ids;
mod lifecycle_integration;
mod log_viewer;
mod openspec;
mod openspec_cmd;
mod orchestration;
mod orchestrator;
mod parallel;
mod parallel_run_service;
mod permission;
mod process_manager;
#[allow(dead_code)]
mod repo_lock;
#[allow(dead_code, unused_imports)]
mod runtime;
mod shell_command;
mod spec_delta;
#[cfg(test)]
mod spec_test_annotations;
mod stall;
#[allow(dead_code)]
mod startup_preflight;
mod stream_json_textifier;
mod task_file;
mod task_parser;
mod templates;
mod tui;
#[allow(dead_code, unused_imports)]
mod upstream;
mod vcs;
#[cfg(feature = "web-monitoring")]
mod web;
mod worktree_ops;
#[cfg(test)]
mod test_support;
use clap::{CommandFactory, Parser};
use cli::{
install_skills_legacy_error, Cli, Commands, InstallSkillsTarget, InternalCompleteCommands,
LogsArgs, TuiArgs, VERSION_WITH_BUILD,
};
use config::OrchestratorConfig;
use error::Result;
use install_skills::{run_install_skills, InstallSkillsOptions};
use lifecycle_integration::{
LifecycleContext, LifecycleEvent, LifecycleExecutionMode, LifecycleIntegration, LifecycleState,
};
use orchestrator::Orchestrator;
use parallel::PostArchiveAction;
use std::path::Path;
#[cfg(feature = "web-monitoring")]
use std::path::PathBuf;
use tracing::{error, info, Level};
use tracing_subscriber::{filter::LevelFilter, prelude::*};
fn tui_post_archive_action(args: &TuiArgs) -> PostArchiveAction {
args.push
.clone()
.map(|remote| PostArchiveAction::PushToRemote { remote })
.unwrap_or_default()
}
fn lifecycle_process_context() -> LifecycleContext {
match std::env::current_dir() {
Ok(dir) => LifecycleContext::workspace(dir.display().to_string()),
Err(_) => LifecycleContext::default(),
}
}
async fn resolve_tui_upstream_runtime(
args: &TuiArgs,
) -> std::result::Result<Option<upstream::UpstreamRuntime>, String> {
let upstream_config = args.upstream_integration().map_err(|err| err.to_string())?;
let git_dir_exists = true;
let repo_root = match std::env::current_dir() {
Ok(root) => root,
Err(err) => {
if upstream_config.is_some() {
return Err(format!(
"upstream integration requires a readable workspace: {err}"
));
}
return Ok(None);
}
};
match upstream_config {
Some(upstream_config) => upstream::prepare_upstream_integration(
upstream_config,
&repo_root,
args.push.clone(),
git_dir_exists,
false,
)
.await
.map(Some)
.map_err(|err| err.to_string()),
None => {
{
if let Err(err) = upstream::ensure_no_unpushed_upstream_recovery(&repo_root).await {
if matches!(err, upstream::UpstreamStartupError::Invalid(_)) {
return Err(err.to_string());
}
tracing::debug!("Upstream recovery scan unavailable, continuing: {}", err);
}
}
Ok(None)
}
}
}
async fn launch_tui(args: TuiArgs) -> Result<()> {
let post_archive_action = tui_post_archive_action(&args);
let config = OrchestratorConfig::load(args.config.as_deref())?;
ensure_state_root_ready(&config);
init_logging(false, config.get_state_base_dir())?;
log_startup("tui");
tui::log_deduplicator::configure_logging(config.get_logging());
let changes = openspec::list_changes_native()?;
#[cfg(feature = "web-monitoring")]
let started = start_local_api(LocalApiOptions::from(&args), &changes).await;
let startup: std::result::Result<Option<upstream::UpstreamRuntime>, String> =
resolve_tui_upstream_runtime(&args).await;
let upstream_runtime = match startup {
Ok(runtime) => runtime,
Err(err) => {
eprintln!("Error: {err}");
#[cfg(feature = "web-monitoring")]
if let Some((handle, _)) = started {
handle.shutdown().await;
}
std::process::exit(1);
}
};
#[cfg(feature = "web-monitoring")]
let (web_url, web_state_opt) = match &started {
Some((handle, state)) => (handle.tcp_url().map(str::to_string), Some(state.clone())),
None => (None, None),
};
#[cfg(feature = "web-monitoring")]
if let Some(state) = &web_state_opt {
publish_execution_contract(
state,
args.push.as_deref(),
upstream_runtime
.as_ref()
.map(|runtime| runtime.config.remote.as_str()),
)
.await;
}
#[cfg(not(feature = "web-monitoring"))]
let web_url: Option<String> = {
if args.web {
eprintln!(
"Warning: Web monitoring is not enabled. Compile with --features web-monitoring"
);
}
None
};
let lifecycle = LifecycleIntegration::start(
config.get_lifecycle_integration(),
LifecycleExecutionMode::Tui,
);
lifecycle.handle().publish(LifecycleEvent::ProcessStarted {
context: lifecycle_process_context(),
});
let result = tui::run_tui(
changes,
config,
web_url,
#[cfg(feature = "web-monitoring")]
web_state_opt,
post_archive_action,
upstream_runtime,
lifecycle.handle(),
)
.await;
lifecycle.shutdown().await;
#[cfg(feature = "web-monitoring")]
if let Some((handle, _)) = started {
handle.shutdown().await;
}
result
}
#[cfg(feature = "web-monitoring")]
async fn publish_execution_contract(
state: &web::WebState,
push_remote: Option<&str>,
upstream_remote: Option<&str>,
) {
let Ok(repo_root) = std::env::current_dir() else {
return;
};
let Ok(Some(base_branch)) = vcs::git::commands::get_current_branch(&repo_root).await else {
return;
};
state.set_execution_contract(
web::remote_control_api::dto::OwnerExecutionContract::resolve(
base_branch,
push_remote,
upstream_remote,
),
);
}
#[cfg(feature = "web-monitoring")]
fn resolve_listener_plan(
tcp: bool,
unix_socket: Option<&Path>,
no_unix_socket: bool,
) -> std::result::Result<web::ListenerPlan, String> {
#[cfg(unix)]
{
let workspace = std::env::current_dir()
.map_err(|e| format!("failed to resolve the current directory: {e}"))?;
let common_dir = repo_lock::discover_common_dir(&workspace);
let selection = web::unix_socket::resolve_unix_socket(
unix_socket,
no_unix_socket,
common_dir.as_deref(),
)?;
Ok(web::ListenerPlan {
unix_socket: selection.path().map(Path::to_path_buf),
tcp,
})
}
#[cfg(not(unix))]
{
let _ = (unix_socket, no_unix_socket);
Ok(web::ListenerPlan { tcp })
}
}
#[cfg(feature = "web-monitoring")]
struct LocalApiOptions {
tcp: bool,
port: u16,
bind: String,
auth_token: Option<String>,
auth_token_env: Option<String>,
allowed_origins: Vec<String>,
unix_socket: Option<PathBuf>,
no_unix_socket: bool,
}
#[cfg(feature = "web-monitoring")]
impl From<&TuiArgs> for LocalApiOptions {
fn from(args: &TuiArgs) -> Self {
Self {
tcp: args.web,
port: args.web_port,
bind: args.web_bind.clone(),
auth_token: args.web_auth_token.clone(),
auth_token_env: args.web_auth_token_env.clone(),
allowed_origins: args.web_allowed_origins.clone(),
unix_socket: args.web_unix_socket.clone(),
no_unix_socket: args.no_web_unix_socket,
}
}
}
#[cfg(feature = "web-monitoring")]
impl From<&cli::RunArgs> for LocalApiOptions {
fn from(args: &cli::RunArgs) -> Self {
Self {
tcp: args.web,
port: args.web_port,
bind: args.web_bind.clone(),
auth_token: args.web_auth_token.clone(),
auth_token_env: args.web_auth_token_env.clone(),
allowed_origins: args.web_allowed_origins.clone(),
unix_socket: args.web_unix_socket.clone(),
no_unix_socket: args.no_web_unix_socket,
}
}
}
#[cfg(feature = "web-monitoring")]
async fn start_local_api(
options: LocalApiOptions,
changes: &[openspec::Change],
) -> Option<(web::ServerHandle, std::sync::Arc<web::WebState>)> {
let plan = match resolve_listener_plan(
options.tcp,
options.unix_socket.as_deref(),
options.no_unix_socket,
) {
Ok(plan) => plan,
Err(error) => {
eprintln!("Error: {error}");
std::process::exit(1);
}
};
if plan.is_empty() {
return None;
}
let config = web::WebConfig::enabled(options.port, options.bind)
.with_tcp_enabled(options.tcp)
.with_auth(
options.auth_token,
options.auth_token_env,
options.allowed_origins,
);
let state = std::sync::Arc::new(web::WebState::new(changes));
match web::start_listeners(config, plan, state.clone()).await {
Ok(handle) => {
for endpoint in handle.endpoints() {
info!("Local API available at: {}", endpoint);
}
Some((handle, state))
}
Err(error) => {
eprintln!("Error: {error}");
std::process::exit(1);
}
}
}
fn ensure_state_root_ready(config: &OrchestratorConfig) {
if let Err(err) = config::defaults::ensure_state_root_usable(config.get_state_base_dir()) {
eprintln!("Error: {err}");
std::process::exit(1);
}
}
fn init_logging(enable_stdout: bool, state_base_dir: Option<&str>) -> Result<()> {
use config::defaults::{cleanup_old_logs, get_log_file_path};
use std::fs::{create_dir_all, File};
use tracing_subscriber::fmt::writer::MakeWriterExt;
let repo_root = std::env::current_dir().ok();
let log_path = get_log_file_path(state_base_dir, repo_root.as_deref())
.map_err(|e| error::OrchestratorError::Io(std::io::Error::other(e.to_string())))?;
if let Some(parent) = log_path.parent() {
create_dir_all(parent).map_err(|e| {
error::OrchestratorError::Io(std::io::Error::other(format!(
"Failed to create log directory '{}': {}",
parent.display(),
e
)))
})?;
}
if let Err(e) = cleanup_old_logs(state_base_dir, repo_root.as_deref(), 7) {
tracing::warn!("Failed to clean up old logs: {}", e);
}
let file = File::options()
.create(true)
.append(true)
.open(&log_path)
.map_err(|e| {
error::OrchestratorError::Io(std::io::Error::other(format!(
"Failed to open log file '{}': {}",
log_path.display(),
e
)))
})?;
let file_layer = tracing_subscriber::fmt::layer()
.with_writer(file.with_max_level(Level::DEBUG))
.with_ansi(false)
.with_target(true)
.with_thread_ids(false)
.with_file(true)
.with_line_number(true);
let registry = tracing_subscriber::registry().with(file_layer);
if enable_stdout {
let stdout_layer = tracing_subscriber::fmt::layer()
.with_writer(std::io::stdout)
.with_ansi(true)
.with_target(false)
.with_thread_ids(false)
.with_file(false)
.with_line_number(false)
.with_filter(LevelFilter::INFO);
registry.with(stdout_layer).init();
} else {
registry.init();
}
Ok(())
}
fn log_startup(mode: &str) {
info!("Starting cflx {} mode={}.", VERSION_WITH_BUILD, mode);
}
fn repository_lock_invocation(cli: &Cli) -> repo_lock::InvocationKind {
match &cli.command {
None => repo_lock::InvocationKind::DefaultTui,
Some(Commands::Tui(_)) => repo_lock::InvocationKind::Tui,
Some(Commands::Run(_)) => repo_lock::InvocationKind::Run,
_ => repo_lock::InvocationKind::Other,
}
}
fn enforce_owner_startup_preflight(cli: &Cli) {
let kind = repository_lock_invocation(cli);
if let Some(message) = startup_preflight::owner_startup_error(kind) {
eprintln!("Error: {message}");
std::process::exit(1);
}
}
fn acquire_repository_lock(cli: &Cli) {
let kind = repository_lock_invocation(cli);
let mode = match repo_lock::classify_invocation(kind) {
repo_lock::LockDecision::Bypass => return,
repo_lock::LockDecision::Acquire(mode) => mode,
};
let Ok(workspace) = std::env::current_dir() else {
return;
};
match repo_lock::acquire(&workspace, mode) {
Ok(None) => {}
Ok(Some(lock)) => repo_lock::install(lock),
Err(err @ repo_lock::LockError::Conflict(_)) => {
eprintln!("{err}");
std::process::exit(1);
}
Err(err) => {
eprintln!("Warning: repository lock unavailable, continuing: {err}");
}
}
}
fn run_completion_subcommand(args: cli::CompletionArgs) {
let shell = clap_complete::Shell::from(args.shell);
let mut command = Cli::command();
let mut stdout = std::io::stdout();
clap_complete::generate(shell, &mut command, "cflx", &mut stdout);
print_dynamic_completion_hooks(args.shell);
}
fn print_dynamic_completion_hooks(shell: cli::CompletionShell) {
match shell {
cli::CompletionShell::Bash => print!("{}", BASH_DYNAMIC_COMPLETION_HOOK),
cli::CompletionShell::Zsh => print!("{}", ZSH_DYNAMIC_COMPLETION_HOOK),
cli::CompletionShell::Fish => print!("{}", FISH_DYNAMIC_COMPLETION_HOOK),
cli::CompletionShell::PowerShell => print!("{}", POWERSHELL_DYNAMIC_COMPLETION_HOOK),
}
}
fn run_internal_complete_subcommand(args: cli::InternalCompleteArgs) {
match args.command {
InternalCompleteCommands::ChangeIds(change_args) => {
let scope = completion::ChangeIdCandidateScope::from_flags(
change_args.active,
change_args.archived,
);
let cwd = match std::env::current_dir() {
Ok(cwd) => cwd,
Err(_) => return,
};
for candidate in completion::discover_change_id_candidates(
&cwd,
scope,
change_args.prefix.as_deref(),
) {
println!("{candidate}");
}
}
}
}
const BASH_DYNAMIC_COMPLETION_HOOK: &str = r#"
# cflx dynamic OpenSpec change-id completion hook
_cflx_static_completion() {
_cflx "$@"
}
_cflx_dynamic_change_ids() {
local scope="$1"
local prefix="$2"
local -a cmd=(cflx __complete change-ids --prefix "$prefix")
case "$scope" in
active) cmd+=(--active) ;;
all) cmd+=(--active --archived) ;;
esac
mapfile -t COMPREPLY < <("${cmd[@]}" 2>/dev/null)
}
_cflx_dynamic_run_change_ids() {
local current="${COMP_WORDS[COMP_CWORD]}"
local prefix="${current##*,}"
local before="${current%,*}"
_cflx_dynamic_change_ids active "$prefix"
if [[ "$before" != "$current" ]]; then
local i
for i in "${!COMPREPLY[@]}"; do COMPREPLY[$i]="$before,${COMPREPLY[$i]}"; done
fi
}
_cflx_dynamic_completion() {
local cur="${COMP_WORDS[COMP_CWORD]}"
local prev="${COMP_WORDS[COMP_CWORD-1]}"
if [[ "$prev" == "--change" ]]; then
_cflx_dynamic_run_change_ids
return
fi
if [[ ${COMP_CWORD} -ge 3 && "${COMP_WORDS[1]}" == "openspec" ]]; then
case "${COMP_WORDS[2]}" in
show)
if [[ "$cur" != -* ]]; then
_cflx_dynamic_change_ids all "$cur"
return
fi
;;
validate|archive)
if [[ "$cur" != -* ]]; then
_cflx_dynamic_change_ids active "$cur"
return
fi
;;
esac
fi
_cflx_static_completion "$@"
}
complete -F _cflx_dynamic_completion -o bashdefault -o default cflx
# Surfaces: cflx run --change -> _cflx_dynamic_run_change_ids; cflx openspec show -> active+archived;
# cflx openspec validate/archive -> active. Candidate command: cflx __complete change-ids
"#;
const ZSH_DYNAMIC_COMPLETION_HOOK: &str = r#"
# cflx dynamic OpenSpec change-id completion hook
_cflx_static_completion() {
_cflx "$@"
}
_cflx_dynamic_change_ids() {
local scope="$1"
local prefix="$2"
local -a cmd=(cflx __complete change-ids --prefix "$prefix")
case "$scope" in
active) cmd+=(--active) ;;
all) cmd+=(--active --archived) ;;
esac
compadd -- "${(@f)$(${cmd[@]} 2>/dev/null)}"
}
_cflx_dynamic_run_change_ids() {
local current="${words[CURRENT]}"
local prefix="${current##*,}"
local before="${current%,*}"
local -a candidates
candidates=("${(@f)$(cflx __complete change-ids --active --prefix "$prefix" 2>/dev/null)}")
if [[ "$before" != "$current" ]]; then
candidates=("${(@)^candidates/#/$before,}")
fi
compadd -- "${candidates[@]}"
}
_cflx_dynamic_completion() {
if [[ "${words[CURRENT-1]}" == "--change" ]]; then
_cflx_dynamic_run_change_ids
return
fi
if [[ ${CURRENT} -ge 4 && "${words[2]}" == "openspec" ]]; then
case "${words[3]}" in
show)
[[ "${words[CURRENT]}" == -* ]] || { _cflx_dynamic_change_ids all "${words[CURRENT]}"; return; }
;;
validate|archive)
[[ "${words[CURRENT]}" == -* ]] || { _cflx_dynamic_change_ids active "${words[CURRENT]}"; return; }
;;
esac
fi
_cflx_static_completion "$@"
}
compdef _cflx_dynamic_completion cflx
# Surfaces: cflx run --change -> _cflx_dynamic_run_change_ids; cflx openspec show -> active+archived;
# cflx openspec validate/archive -> active. Candidate command: cflx __complete change-ids
"#;
const FISH_DYNAMIC_COMPLETION_HOOK: &str = r#"
# cflx dynamic OpenSpec change-id completion hook
function __cflx_dynamic_change_ids
set -l scope $argv[1]
set -l prefix $argv[2]
set -l cmd cflx __complete change-ids --prefix "$prefix"
switch $scope
case active
set cmd $cmd --active
case all
set cmd $cmd --active --archived
end
$cmd 2>/dev/null
end
complete -c cflx -n '__fish_seen_subcommand_from run; and __fish_seen_argument --change' -a '(__cflx_dynamic_change_ids active (string split -r -m1 , (commandline -ct))[-1])'
complete -c cflx -n '__fish_seen_subcommand_from openspec; and __fish_seen_subcommand_from show' -a '(__cflx_dynamic_change_ids all (commandline -ct))'
complete -c cflx -n '__fish_seen_subcommand_from openspec; and __fish_seen_subcommand_from validate archive' -a '(__cflx_dynamic_change_ids active (commandline -ct))'
"#;
const POWERSHELL_DYNAMIC_COMPLETION_HOOK: &str = r#"
# cflx dynamic OpenSpec change-id completion hook
function __CflxDynamicChangeIds($Scope, $Prefix) {
$args = @('__complete', 'change-ids', '--prefix', $Prefix)
if ($Scope -eq 'active') { $args += '--active' }
if ($Scope -eq 'all') { $args += @('--active', '--archived') }
& cflx @args 2>$null
}
function __CflxDynamicRunChangeIds($WordToComplete) {
$prefix = $WordToComplete -replace '^.*,', ''
$before = $WordToComplete -replace ',?[^,]*$', ''
foreach ($candidate in (__CflxDynamicChangeIds active $prefix)) {
if ($before) { "$before,$candidate" } else { $candidate }
}
}
Register-ArgumentCompleter -Native -CommandName 'cflx' -ScriptBlock {
param($wordToComplete, $commandAst, $cursorPosition)
$elements = @($commandAst.CommandElements | ForEach-Object { $_.Extent.Text })
$commandText = $elements -join ' '
if ($commandText -match '^cflx\s+run\b' -and ($elements -contains '--change')) {
__CflxDynamicRunChangeIds $wordToComplete | ForEach-Object {
[CompletionResult]::new($_, $_, [CompletionResultType]::ParameterValue, 'OpenSpec active change ID')
}
return
}
if ($commandText -match '^cflx\s+openspec\s+show\b' -and $wordToComplete -notlike '-*') {
__CflxDynamicChangeIds all $wordToComplete | ForEach-Object {
[CompletionResult]::new($_, $_, [CompletionResultType]::ParameterValue, 'OpenSpec change ID')
}
return
}
if ($commandText -match '^cflx\s+openspec\s+(validate|archive)\b' -and $wordToComplete -notlike '-*') {
__CflxDynamicChangeIds active $wordToComplete | ForEach-Object {
[CompletionResult]::new($_, $_, [CompletionResultType]::ParameterValue, 'OpenSpec active change ID')
}
return
}
}
# Surfaces: cflx run --change -> active comma-token candidates; cflx openspec show -> active+archived;
# cflx openspec validate/archive -> active. Candidate command: cflx __complete change-ids
"#;
fn run_openapi_subcommand() {
#[cfg(feature = "web-monitoring")]
{
use std::io::Write;
let document = web::openapi::document_yaml();
let mut stdout = std::io::stdout();
if let Err(e) = stdout
.write_all(document.as_bytes())
.and_then(|()| stdout.flush())
{
eprintln!("Error: failed to write the OpenAPI document to stdout: {e}");
std::process::exit(1);
}
}
#[cfg(not(feature = "web-monitoring"))]
{
eprintln!(
"Error: OpenAPI support is unavailable in this build. \
Rebuild with `--features web-monitoring` to export the /api/v2 schema."
);
std::process::exit(1);
}
}
fn run_logs_subcommand(args: LogsArgs, custom_config_path: Option<&Path>) {
let state_base_dir = match OrchestratorConfig::load_storage_settings(custom_config_path) {
Ok(config) => config.get_state_base_dir().map(str::to_string),
Err(e) => {
eprintln!("Error: {e}");
std::process::exit(1);
}
};
let options = log_viewer::LogViewerOptions {
print_path: args.path,
last: args.last,
follow: args.follow,
today: args.today,
project: args.project,
repo_root: std::env::current_dir().ok(),
state_base_dir,
};
if let Err(e) = log_viewer::run_logs_command(&options, &mut std::io::stdout()) {
eprintln!("Error: {e}");
std::process::exit(1);
}
}
#[tokio::main]
async fn main() -> Result<()> {
let cli = match Cli::try_parse() {
Ok(cli) => cli,
Err(error) => client::exit_on_parse_error(error),
};
if let Err(error) = cli.validate_upstream_option_placement() {
eprintln!("Error: {error}");
std::process::exit(2);
}
enforce_owner_startup_preflight(&cli);
acquire_repository_lock(&cli);
match cli.command {
Some(Commands::Completion(args)) => {
run_completion_subcommand(args);
}
Some(Commands::Complete(args)) => {
run_internal_complete_subcommand(args);
}
Some(Commands::Openapi) => {
run_openapi_subcommand();
}
Some(Commands::Client(args)) => {
std::process::exit(client::run(args).await);
}
None => {
launch_tui(TuiArgs {
config: cli.config,
web: cli.web,
web_port: cli.web_port,
web_bind: cli.web_bind,
web_auth_token: cli.web_auth_token,
web_auth_token_env: cli.web_auth_token_env,
web_allowed_origins: cli.web_allowed_origins,
web_unix_socket: cli.web_unix_socket,
no_web_unix_socket: cli.no_web_unix_socket,
push: cli.push,
integrate_upstream: cli.integrate_upstream,
integrate_upstream_default_remote: cli.integrate_upstream_default_remote,
upstream_verify_command: cli.upstream_verify_command,
})
.await?;
}
Some(Commands::Tui(tui_args)) => launch_tui(tui_args).await?,
Some(Commands::Run(args)) => {
let config = OrchestratorConfig::load(args.config.as_deref())?;
ensure_state_root_ready(&config);
init_logging(true, config.get_state_base_dir())?;
log_startup("run");
#[cfg(feature = "web-monitoring")]
let started_api = {
let initial_changes = openspec::list_changes_native()?;
start_local_api(LocalApiOptions::from(&args), &initial_changes).await
};
#[cfg(feature = "web-monitoring")]
let web_state_arc = started_api.as_ref().map(|(_, state)| state.clone());
#[cfg(feature = "web-monitoring")]
if let Some(web_state) = web_state_arc.as_ref() {
web_state
.set_execution_facts(std::sync::Arc::new(
orchestration::execution_facts::ExecutionFactsStore::new(),
))
.await;
publish_execution_contract(
web_state,
args.push.as_deref(),
args.upstream_integration()
.ok()
.flatten()
.as_ref()
.map(|config| config.remote.as_str()),
)
.await;
}
#[cfg(not(feature = "web-monitoring"))]
if args.web {
eprintln!(
"Warning: Web monitoring is not enabled. Compile with --features web-monitoring"
);
}
let vcs_override = match args.vcs.parse::<vcs::VcsBackend>() {
Ok(backend) => Some(backend),
Err(err) => {
eprintln!("Error: {}", err);
std::process::exit(1);
}
};
let upstream_integration = match args.upstream_integration() {
Ok(config) => config,
Err(err) => {
eprintln!("Error: {}", err);
std::process::exit(1);
}
};
let upstream_runtime = match upstream_integration {
Some(config) => {
let repo_root = std::env::current_dir()?;
match upstream::prepare_upstream_integration(
config,
&repo_root,
args.push.clone(),
true,
args.dry_run,
)
.await
{
Ok(runtime) => Some(runtime),
Err(err) => {
eprintln!("Error: {}", err);
std::process::exit(1);
}
}
}
None => {
if !args.dry_run {
let repo_root = std::env::current_dir()?;
if let Err(err) =
upstream::ensure_no_unpushed_upstream_recovery(&repo_root).await
{
if matches!(err, upstream::UpstreamStartupError::Invalid(_)) {
eprintln!("Error: {}", err);
std::process::exit(1);
}
tracing::debug!(
"Upstream recovery scan unavailable, continuing: {}",
err
);
}
}
None
}
};
let lifecycle = LifecycleIntegration::start(
config.get_lifecycle_integration(),
LifecycleExecutionMode::Run,
);
let lifecycle_handle = lifecycle.handle();
let workspace_context = lifecycle_process_context().workspace;
lifecycle_handle.publish(LifecycleEvent::ProcessStarted {
context: lifecycle_process_context(),
});
use std::sync::atomic::{AtomicBool, AtomicU8, Ordering};
use std::sync::Arc;
let run_state = Arc::new(AtomicU8::new(1)); let graceful_stop_flag = Arc::new(AtomicBool::new(false));
let force_stop_flag = Arc::new(AtomicBool::new(false));
let restart_requested = Arc::new(AtomicBool::new(false));
let signal_stop = Arc::new(AtomicBool::new(false));
#[cfg(unix)]
{
let signal_stop_sigterm = signal_stop.clone();
tokio::spawn(async move {
let mut sigterm =
tokio::signal::unix::signal(tokio::signal::unix::SignalKind::terminate())
.expect("Failed to install SIGTERM handler");
sigterm.recv().await;
info!("Received SIGTERM, shutting down gracefully...");
signal_stop_sigterm.store(true, Ordering::SeqCst);
});
}
{
let signal_stop_sigint = signal_stop.clone();
tokio::spawn(async move {
let _ = tokio::signal::ctrl_c().await;
info!("Received SIGINT (Ctrl+C), shutting down gracefully...");
signal_stop_sigint.store(true, Ordering::SeqCst);
});
}
let change_ids = args.normalized_target_changes();
let config_path = args.config.clone();
let max_iterations = args.max_iterations;
let max_concurrent = args.max_concurrent;
let dry_run = args.dry_run;
let no_resume = args.no_resume;
let post_archive_action = args
.push
.clone()
.map(|remote| parallel::PostArchiveAction::PushToRemote { remote })
.unwrap_or_default();
loop {
if signal_stop.load(Ordering::SeqCst) {
info!("Signal stop detected, exiting");
break;
}
info!("Starting orchestrator");
lifecycle_handle
.publish_state(LifecycleState::Working, lifecycle_process_context());
let mut orchestrator = Orchestrator::new(
change_ids.clone(),
config_path.clone(),
max_iterations,
max_concurrent,
dry_run,
vcs_override,
no_resume,
post_archive_action.clone(),
)?;
orchestrator
.set_lifecycle_handle(lifecycle_handle.clone(), workspace_context.clone());
if let Some(runtime) = upstream_runtime.clone() {
orchestrator.set_upstream_integration(runtime);
}
#[cfg(feature = "web-monitoring")]
if let Some(ref web_state) = web_state_arc {
orchestrator.set_web_state(web_state.clone()).await;
}
let cancel_token = tokio_util::sync::CancellationToken::new();
let monitor_token = cancel_token.clone();
let monitor_force = force_stop_flag.clone();
let monitor_signal = signal_stop.clone();
let monitor_handle = tokio::spawn(async move {
loop {
tokio::time::sleep(tokio::time::Duration::from_millis(100)).await;
if monitor_signal.load(Ordering::SeqCst)
|| monitor_force.load(Ordering::SeqCst)
{
if monitor_force.load(Ordering::SeqCst) {
info!("Force stop detected, cancelling orchestrator");
} else {
info!("Signal received, cancelling orchestrator");
}
monitor_token.cancel();
break;
}
}
});
let result = orchestrator
.run(cancel_token, Some(graceful_stop_flag.clone()))
.await;
monitor_handle.abort();
run_state.store(0, Ordering::SeqCst);
match result {
Err(e) => {
error!("Orchestrator error: {}", e);
lifecycle_handle
.publish_state(LifecycleState::Blocked, lifecycle_process_context());
loop {
if restart_requested.load(Ordering::SeqCst) {
info!("Retry requested after error, will restart orchestrator");
break;
}
if force_stop_flag.load(Ordering::SeqCst)
|| signal_stop.load(Ordering::SeqCst)
{
info!("Stop requested in error state, exiting");
lifecycle.shutdown().await;
return Err(e);
}
tokio::time::sleep(tokio::time::Duration::from_millis(100)).await;
}
info!("Continuing after error due to retry request");
}
Ok(()) => {
info!("Orchestrator completed successfully");
}
}
if restart_requested.swap(false, Ordering::SeqCst) {
info!("Restarting orchestrator due to web control request");
run_state.store(1, Ordering::SeqCst); graceful_stop_flag.store(false, Ordering::SeqCst);
force_stop_flag.store(false, Ordering::SeqCst);
continue; }
break;
}
lifecycle.shutdown().await;
#[cfg(feature = "web-monitoring")]
if let Some((handle, _)) = started_api {
handle.shutdown().await;
}
}
Some(Commands::Logs(args)) => {
run_logs_subcommand(args, cli.config.as_deref());
}
Some(Commands::Init(args)) => {
let config_path = Path::new(".cflx.jsonc");
if config_path.exists() && !args.force {
eprintln!(
"Error: Configuration file '{}' already exists.",
config_path.display()
);
eprintln!("Use --force to overwrite the existing file.");
std::process::exit(1);
}
let content = templates::get_template_content(args.template);
std::fs::write(config_path, content)?;
println!(
"Created configuration file '{}' with {:?} template.",
config_path.display(),
args.template
);
}
Some(Commands::InstallSkills(args)) => {
if let Some(src) = &args.legacy_source {
eprintln!("{}", install_skills_legacy_error(src));
std::process::exit(1);
}
let target = match args.target() {
InstallSkillsTarget::Agents => install_skills::InstallTarget::Agents,
InstallSkillsTarget::Claude => install_skills::InstallTarget::Claude,
};
let opts = InstallSkillsOptions {
global: args.global,
target,
project_root: None, };
if let Err(e) = run_install_skills(opts) {
eprintln!("Error: {e}");
std::process::exit(1);
}
}
Some(Commands::Openspec(args)) => {
use cli::{EvidenceMode, OpenspecCommands};
match args.command {
OpenspecCommands::List(list_args) => {
if let Err(e) = crate::openspec_cmd::cmd_list(list_args.specs) {
eprintln!("Error: {}", e);
std::process::exit(1);
}
}
OpenspecCommands::Show(show_args) => {
if let Err(e) = crate::openspec_cmd::cmd_show(
&show_args.change_id,
show_args.json,
show_args.deltas_only,
) {
eprintln!("Error: {}", e);
std::process::exit(1);
}
}
OpenspecCommands::Validate(val_args) => {
let strict = val_args.strict || val_args.archive_gate;
let evidence = if val_args.archive_gate {
"error"
} else {
match val_args.evidence {
EvidenceMode::Off => "off",
EvidenceMode::Warn => "warn",
EvidenceMode::Error => "error",
}
};
let (is_valid, exit_code) = crate::openspec_cmd::cmd_validate(
val_args.change_id.as_deref(),
strict,
evidence,
);
if !is_valid {
std::process::exit(exit_code);
}
}
OpenspecCommands::Archive(arc_args) => {
if !arc_args.yes {
eprintln!("Error: --yes flag is required (non-interactive only)");
std::process::exit(1);
}
if let Err(e) =
crate::openspec_cmd::cmd_archive(&arc_args.change_id, arc_args.skip_specs)
{
eprintln!("Error: {}", e);
std::process::exit(1);
}
}
OpenspecCommands::Verify(verify_args) => {
match crate::openspec_cmd::cmd_verify(
&verify_args.change_id,
verify_args.verification_id.as_deref(),
verify_args.plan,
verify_args.json,
)
.await
{
Ok(0) => {}
Ok(code) => std::process::exit(code),
Err(e) => {
eprintln!("Error: {}", e);
std::process::exit(1);
}
}
}
}
}
Some(Commands::CheckConflicts(args)) => {
let changes = openspec::list_changes_native()?;
let mut all_deltas = Vec::new();
for change in &changes {
match spec_delta::parse_change_deltas(&change.id) {
Ok(deltas) => all_deltas.extend(deltas),
Err(e) => {
eprintln!("Error parsing deltas for change '{}': {}", change.id, e);
std::process::exit(1);
}
}
}
let conflicts = spec_delta::detect_conflicts(&all_deltas);
if args.json {
match spec_delta::format_conflicts_json(&conflicts) {
Ok(json) => {
println!("{}", json);
}
Err(e) => {
eprintln!("Error formatting JSON output: {}", e);
std::process::exit(1);
}
}
} else {
let output = spec_delta::format_conflicts_human(&conflicts);
println!("{}", output);
}
if !conflicts.is_empty() {
std::process::exit(2);
}
}
}
Ok(())
}