pub mod envelope;
use std::ffi::OsString;
use std::path::{Path, PathBuf};
use clap::CommandFactory;
#[cfg(feature = "web-monitoring")]
pub mod completion;
#[cfg(feature = "web-monitoring")]
pub mod control;
#[cfg(feature = "web-monitoring")]
pub mod mcp;
#[cfg(feature = "web-monitoring")]
pub mod repo;
#[cfg(feature = "web-monitoring")]
mod session;
#[cfg(feature = "web-monitoring")]
pub mod subscribe;
#[cfg(feature = "web-monitoring")]
mod transport;
#[cfg(feature = "web-monitoring")]
mod wait;
use envelope::{Operation, Outcome, ResultEnvelope};
use crate::cli::{ClientArgs, ClientCommands, ClientSubscribeCommands};
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub enum RouteSelector {
#[default]
Default,
Project(PathBuf),
Socket(PathBuf),
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct RouteError {
pub message: String,
}
impl RouteError {
fn new(message: impl Into<String>) -> Self {
Self {
message: message.into(),
}
}
}
impl std::fmt::Display for RouteError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str(&self.message)
}
}
impl RouteSelector {
pub fn from_inputs(
project_dir: Option<&Path>,
unix_socket: Option<&Path>,
) -> Result<Self, RouteError> {
match (project_dir, unix_socket) {
(Some(_), Some(_)) => Err(RouteError::new(
"project_dir and unix_socket select two different routes, so supplying both in \
one call is ambiguous. Name the project directory, or the socket, not both",
)),
(Some(project), None) => {
if project.as_os_str().is_empty() {
return Err(RouteError::new("project_dir is empty"));
}
if !project.is_absolute() {
return Err(RouteError::new(format!(
"project_dir '{}' is relative. It must be absolute: this process's \
working directory has nothing to do with the project the work belongs to",
project.display()
)));
}
Ok(Self::Project(project.to_path_buf()))
}
(None, Some(socket)) => {
if socket.as_os_str().is_empty() {
return Err(RouteError::new("unix_socket is empty"));
}
Ok(Self::Socket(socket.to_path_buf()))
}
(None, None) => Ok(Self::Default),
}
}
pub fn or_default(self, fallback: &RouteSelector) -> Self {
match self {
Self::Default => fallback.clone(),
explicit => explicit,
}
}
}
#[cfg(feature = "web-monitoring")]
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ProjectRoute {
pub socket: PathBuf,
pub repo_root: PathBuf,
}
#[cfg(feature = "web-monitoring")]
pub fn resolve_project(project_dir: &Path) -> Result<ProjectRoute, RouteError> {
if !project_dir.is_absolute() {
return Err(RouteError::new(format!(
"project_dir '{}' is relative. It must be absolute",
project_dir.display()
)));
}
let metadata = std::fs::metadata(project_dir).map_err(|error| {
RouteError::new(format!(
"project_dir '{}' could not be read: {error}",
project_dir.display()
))
})?;
if !metadata.is_dir() {
return Err(RouteError::new(format!(
"project_dir '{}' is not a directory",
project_dir.display()
)));
}
let canonical =
std::fs::canonicalize(project_dir).unwrap_or_else(|_| project_dir.to_path_buf());
let repo_root = session::discover_repo_root(&canonical).ok_or_else(|| {
RouteError::new(format!(
"project_dir '{}' is not inside a usable Git working tree, so no owner socket can \
be derived from it",
project_dir.display()
))
})?;
let common_dir = crate::repo_lock::discover_common_dir(&canonical).ok_or_else(|| {
RouteError::new(format!(
"the Git common directory of project_dir '{}' could not be resolved",
project_dir.display()
))
})?;
Ok(ProjectRoute {
socket: crate::web::unix_socket::default_socket_path(&common_dir),
repo_root,
})
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum OutputMode {
Human,
Json,
}
impl OutputMode {
pub fn from_json_flag(json: bool) -> Self {
if json {
Self::Json
} else {
Self::Human
}
}
}
pub fn emit(envelope: &ResultEnvelope, mode: OutputMode) -> i32 {
use std::io::Write;
let line = match mode {
OutputMode::Json => envelope.to_json_line(),
OutputMode::Human => envelope.to_human_line(),
};
let mut stdout = std::io::stdout();
if writeln!(stdout, "{line}")
.and_then(|()| stdout.flush())
.is_err()
{
return Outcome::TransportError.exit_code();
}
if !envelope.ok {
if let Some(message) = &envelope.message {
eprintln!("cflx client: {}: {message}", envelope.outcome.as_str());
} else {
eprintln!("cflx client: {}", envelope.outcome.as_str());
}
}
envelope.exit_code()
}
const JSON_FLAG: &str = "--json";
const END_OF_OPTIONS: &str = "--";
fn operation_of(subcommand: &str) -> Option<Operation> {
match subcommand {
"status" => Some(Operation::Status),
"mark" => Some(Operation::ControlMark),
"unmark" => Some(Operation::ControlUnmark),
"start" => Some(Operation::ControlStart),
"stop" => Some(Operation::ControlStop),
"force-stop" => Some(Operation::ControlForceStop),
"force-stop-change" => Some(Operation::ControlForceStopChange),
"wait" => Some(Operation::Wait),
_ => None,
}
}
fn subscribe_operation_of(subcommand: Option<&str>) -> Operation {
match subcommand {
Some("set") => Operation::SubscribeSet,
Some("clear") => Operation::SubscribeClear,
_ => Operation::SubscribeGet,
}
}
pub fn json_usage_operation(argv: &[OsString]) -> Option<Operation> {
let selects_json = argv
.iter()
.skip(1)
.take_while(|arg| *arg != END_OF_OPTIONS)
.any(|arg| arg == JSON_FLAG);
if !selects_json {
return None;
}
let matches = crate::cli::Cli::command()
.ignore_errors(true)
.try_get_matches_from(argv)
.ok()?;
let client = matches.subcommand_matches("client")?;
if let Some(subscribe) = client.subcommand_matches("subscribe") {
return Some(subscribe_operation_of(subscribe.subcommand_name()));
}
Some(
client
.subcommand_name()
.and_then(operation_of)
.unwrap_or(Operation::Status),
)
}
fn summarize_parse_error(error: &clap::Error) -> String {
error
.render()
.to_string()
.lines()
.map(str::trim)
.find(|line| !line.is_empty())
.map(|line| line.strip_prefix("error: ").unwrap_or(line).to_string())
.unwrap_or_else(|| "the invocation could not be parsed".to_string())
}
pub fn usage_error_envelope(error: &clap::Error, operation: Operation) -> ResultEnvelope {
ResultEnvelope::new(operation, Outcome::UsageError).with_message(summarize_parse_error(error))
}
fn is_usage_failure(kind: clap::error::ErrorKind) -> bool {
use clap::error::ErrorKind;
!matches!(
kind,
ErrorKind::DisplayHelp
| ErrorKind::DisplayVersion
| ErrorKind::DisplayHelpOnMissingArgumentOrSubcommand
)
}
pub fn exit_on_parse_error(error: clap::Error) -> ! {
if is_usage_failure(error.kind()) {
let argv: Vec<OsString> = std::env::args_os().collect();
if let Some(operation) = json_usage_operation(&argv) {
let envelope = usage_error_envelope(&error, operation);
std::process::exit(emit(&envelope, OutputMode::Json));
}
}
error.exit()
}
pub async fn run(args: ClientArgs) -> i32 {
if matches!(args.command, ClientCommands::Mcp(_)) {
return serve_mcp(args).await;
}
let (operation, mode, change_id) = match &args.command {
ClientCommands::Status(status) => (
Operation::Status,
OutputMode::from_json_flag(status.json),
None,
),
ClientCommands::Mark(mark) => (
Operation::ControlMark,
OutputMode::from_json_flag(mark.json),
single(&mark.change_ids),
),
ClientCommands::Unmark(unmark) => (
Operation::ControlUnmark,
OutputMode::from_json_flag(unmark.json),
single(&unmark.change_ids),
),
ClientCommands::Start(start) => (
Operation::ControlStart,
OutputMode::from_json_flag(start.json),
None,
),
ClientCommands::Stop(stop) => (
Operation::ControlStop,
OutputMode::from_json_flag(stop.json),
None,
),
ClientCommands::ForceStop(force) => (
Operation::ControlForceStop,
OutputMode::from_json_flag(force.json),
None,
),
ClientCommands::ForceStopChange(force) => (
Operation::ControlForceStopChange,
OutputMode::from_json_flag(force.json),
Some(force.change_id.clone()),
),
ClientCommands::Wait(wait) => (
Operation::Wait,
OutputMode::from_json_flag(wait.json),
Some(wait.change_id.clone()),
),
ClientCommands::Subscribe(subscribe) => match &subscribe.command {
ClientSubscribeCommands::Set(set) => (
Operation::SubscribeSet,
OutputMode::from_json_flag(set.json),
single(&set.change_ids),
),
ClientSubscribeCommands::Get(get) => (
Operation::SubscribeGet,
OutputMode::from_json_flag(get.json),
single(&get.change_ids),
),
ClientSubscribeCommands::Clear(clear) => (
Operation::SubscribeClear,
OutputMode::from_json_flag(clear.json),
single(&clear.change_ids),
),
},
ClientCommands::Mcp(_) => unreachable!("the MCP session is served before this point"),
};
let envelope = execute(args, operation).await;
let envelope = match change_id {
Some(change_id) if envelope.change_id.is_none() => envelope.with_change(change_id),
_ => envelope,
};
emit(&envelope, mode)
}
#[cfg(not(feature = "web-monitoring"))]
async fn execute(_args: ClientArgs, operation: Operation) -> ResultEnvelope {
ResultEnvelope::new(operation, Outcome::FeatureUnavailable).with_message(
"this build has no local /api/v2 support, so it cannot reach an existing owner. \
Rebuild with `--features web-monitoring`",
)
}
#[cfg(feature = "web-monitoring")]
async fn serve_mcp(args: ClientArgs) -> i32 {
let default_route = match RouteSelector::from_inputs(
args.project_dir.as_deref(),
args.unix_socket.as_deref(),
) {
Ok(selector) => selector,
Err(error) => {
eprintln!("cflx client mcp: {error}");
return Outcome::UsageError.exit_code();
}
};
mcp::run(default_route, args.auth_token_env).await
}
#[cfg(not(feature = "web-monitoring"))]
async fn serve_mcp(_args: ClientArgs) -> i32 {
eprintln!(
"cflx client mcp: this build has no local /api/v2 support, so it cannot reach an \
existing owner. Rebuild with `--features web-monitoring`"
);
Outcome::FeatureUnavailable.exit_code()
}
#[cfg(feature = "web-monitoring")]
async fn execute(args: ClientArgs, operation: Operation) -> ResultEnvelope {
let selector = match RouteSelector::from_inputs(
args.project_dir.as_deref(),
args.unix_socket.as_deref(),
) {
Ok(selector) => selector,
Err(error) => {
return ResultEnvelope::new(operation, Outcome::UsageError).with_message(error.message)
}
};
let connection =
match session::Connection::resolve_route(&selector, args.auth_token_env.as_deref()) {
Ok(connection) => connection,
Err(refusal) => return refusal.into_envelope(operation),
};
match args.command {
ClientCommands::Status(_) => session::status(&connection).await,
ClientCommands::Mark(mark) => {
control::run(&connection, control::Action::Mark, &mark.change_ids).await
}
ClientCommands::Unmark(unmark) => {
control::run(&connection, control::Action::Unmark, &unmark.change_ids).await
}
ClientCommands::Start(_) => control::run(&connection, control::Action::Start, &[]).await,
ClientCommands::Stop(_) => control::run(&connection, control::Action::Stop, &[]).await,
ClientCommands::ForceStop(_) => {
control::run(&connection, control::Action::ForceStop, &[]).await
}
ClientCommands::ForceStopChange(force) => {
control::run(
&connection,
control::Action::ForceStopChange,
std::slice::from_ref(&force.change_id),
)
.await
}
ClientCommands::Wait(wait) => {
wait::run(&connection, &wait.change_id, wait.timeout.deadline()).await
}
ClientCommands::Subscribe(args) => match args.command {
ClientSubscribeCommands::Set(set) => {
subscribe::run(
&connection,
&set.change_ids,
Some(&set.instance_id),
subscribe::Intent::Set {
command: set.command,
notify_blocked: set.blocked,
},
)
.await
}
ClientSubscribeCommands::Get(get) => {
subscribe::run(
&connection,
&get.change_ids,
Some(&get.instance_id),
subscribe::Intent::Get,
)
.await
}
ClientSubscribeCommands::Clear(clear) => {
subscribe::run(
&connection,
&clear.change_ids,
Some(&clear.instance_id),
subscribe::Intent::Clear,
)
.await
}
},
ClientCommands::Mcp(_) => unreachable!("the MCP session is served before this point"),
}
}
#[cfg(feature = "web-monitoring")]
fn single(change_ids: &[String]) -> Option<String> {
match change_ids {
[only] => Some(only.clone()),
_ => None,
}
}
#[cfg(test)]
mod tests {
use super::*;
fn argv(args: &[&str]) -> Vec<OsString> {
std::iter::once("cflx")
.chain(args.iter().copied())
.map(OsString::from)
.collect()
}
#[test]
fn a_client_json_invocation_is_recognized_with_the_operation_it_named() {
assert_eq!(
json_usage_operation(&argv(&["client", "status", "--json"])),
Some(Operation::Status)
);
assert_eq!(
json_usage_operation(&argv(&["client", "mark", "../escape", "--json"])),
Some(Operation::ControlMark)
);
assert_eq!(
json_usage_operation(&argv(&["client", "unmark", "../escape", "--json"])),
Some(Operation::ControlUnmark)
);
for (verb, operation) in [
("start", Operation::ControlStart),
("stop", Operation::ControlStop),
("force-stop", Operation::ControlForceStop),
] {
assert_eq!(
json_usage_operation(&argv(&["client", verb, "--nope", "--json"])),
Some(operation),
"{verb}"
);
}
assert_eq!(
json_usage_operation(&argv(&[
"client",
"wait",
"alpha",
"--timeout",
"abc",
"--json"
])),
Some(Operation::Wait)
);
assert_eq!(
json_usage_operation(&argv(&["client", "mark", "--json"])),
Some(Operation::ControlMark)
);
assert_eq!(
json_usage_operation(&argv(&["client", "subscribe", "set", "--json"])),
Some(Operation::SubscribeSet)
);
assert_eq!(
json_usage_operation(&argv(&["client", "subscribe", "clear", "--json"])),
Some(Operation::SubscribeClear)
);
assert_eq!(
json_usage_operation(&argv(&["client", "subscribe", "--json"])),
Some(Operation::SubscribeGet)
);
assert_eq!(
json_usage_operation(&argv(&["client", "--json"])),
Some(Operation::Status)
);
}
#[test]
fn human_and_non_client_invocations_keep_clap_to_themselves() {
assert_eq!(json_usage_operation(&argv(&["client", "status"])), None);
assert_eq!(
json_usage_operation(&argv(&["client", "mark", "../escape"])),
None
);
assert_eq!(
json_usage_operation(&argv(&["openspec", "show", "alpha", "--json"])),
None
);
assert_eq!(json_usage_operation(&argv(&["--json"])), None);
}
#[test]
fn json_intent_is_a_whole_argument_never_a_substring_of_a_value() {
for value in ["--jsonish", "not--json", "x--json", "--json=true"] {
assert_eq!(
json_usage_operation(&argv(&["client", "mark", value])),
None,
"{value} must not select JSON mode"
);
}
}
#[test]
fn help_and_version_are_answers_rather_than_usage_failures() {
use clap::error::ErrorKind;
assert!(!is_usage_failure(ErrorKind::DisplayHelp));
assert!(!is_usage_failure(ErrorKind::DisplayVersion));
assert!(!is_usage_failure(
ErrorKind::DisplayHelpOnMissingArgumentOrSubcommand
));
for kind in [
ErrorKind::InvalidValue,
ErrorKind::UnknownArgument,
ErrorKind::MissingRequiredArgument,
ErrorKind::ValueValidation,
ErrorKind::InvalidSubcommand,
] {
assert!(is_usage_failure(kind), "{kind:?}");
}
}
#[test]
fn a_usage_envelope_is_one_line_and_carries_only_the_problem() {
let error = <crate::cli::Cli as clap::Parser>::try_parse_from(argv(&[
"client",
"mark",
"../escape",
"--json",
]))
.expect_err("an escaping change ID must be rejected");
let envelope = usage_error_envelope(&error, Operation::ControlMark);
assert_eq!(envelope.outcome, Outcome::UsageError);
assert!(!envelope.ok);
assert_eq!(envelope.exit_code(), 2);
let message = envelope.message.clone().unwrap();
assert!(!message.contains('\n'), "{message}");
assert!(!message.starts_with("error: "), "{message}");
assert!(!message.contains("Usage:"), "{message}");
assert!(!envelope.to_json_line().contains('\n'));
}
#[test]
fn a_call_naming_no_selector_takes_the_default_route() {
assert_eq!(
RouteSelector::from_inputs(None, None).unwrap(),
RouteSelector::Default
);
}
#[test]
fn a_project_directory_is_the_normal_selector_and_a_socket_the_override() {
assert_eq!(
RouteSelector::from_inputs(Some(Path::new("/srv/project-b")), None).unwrap(),
RouteSelector::Project(PathBuf::from("/srv/project-b"))
);
assert_eq!(
RouteSelector::from_inputs(None, Some(Path::new("/tmp/owner.sock"))).unwrap(),
RouteSelector::Socket(PathBuf::from("/tmp/owner.sock"))
);
}
#[test]
fn two_selectors_in_one_call_are_refused_rather_than_silently_ranked() {
let error = RouteSelector::from_inputs(
Some(Path::new("/srv/project-b")),
Some(Path::new("/tmp/project-a.sock")),
)
.expect_err("two routes in one call are ambiguous");
assert!(error.message.contains("project_dir"), "{error}");
assert!(error.message.contains("unix_socket"), "{error}");
}
#[test]
fn a_relative_or_empty_project_directory_is_refused() {
for value in ["relative/project", "./project", "..", ""] {
let error = RouteSelector::from_inputs(Some(Path::new(value)), None)
.expect_err("{value} must be refused");
assert!(
error.message.contains("empty") || error.message.contains("relative"),
"{value}: {error}"
);
}
assert!(RouteSelector::from_inputs(None, Some(Path::new(""))).is_err());
}
#[test]
fn a_call_scoped_selector_overrides_the_namespace_default_without_mutating_it() {
let namespace = RouteSelector::Socket(PathBuf::from("/tmp/project-a.sock"));
let call = RouteSelector::from_inputs(Some(Path::new("/srv/project-b")), None)
.unwrap()
.or_default(&namespace);
assert_eq!(
call,
RouteSelector::Project(PathBuf::from("/srv/project-b"))
);
let plain = RouteSelector::from_inputs(None, None)
.unwrap()
.or_default(&namespace);
assert_eq!(plain, namespace);
}
#[test]
fn the_cli_refuses_both_route_options_through_the_existing_usage_contract() {
let error = <crate::cli::Cli as clap::Parser>::try_parse_from(argv(&[
"client",
"--project-dir",
"/srv/project-b",
"--unix-socket",
"/tmp/project-a.sock",
"status",
]))
.expect_err("two routes on one namespace are ambiguous");
assert_eq!(error.kind(), clap::error::ErrorKind::ArgumentConflict);
let envelope = usage_error_envelope(&error, Operation::Status);
assert_eq!(envelope.outcome, Outcome::UsageError);
assert_eq!(envelope.exit_code(), 2);
for route in [
["--project-dir", "/srv/project-b"],
["--unix-socket", "/tmp/project-a.sock"],
] {
<crate::cli::Cli as clap::Parser>::try_parse_from(argv(&[
"client", route[0], route[1], "status",
]))
.unwrap_or_else(|error| panic!("{route:?} must parse: {error}"));
}
}
#[test]
fn the_json_flag_selects_the_machine_contract() {
assert_eq!(OutputMode::from_json_flag(true), OutputMode::Json);
assert_eq!(OutputMode::from_json_flag(false), OutputMode::Human);
}
#[test]
fn an_unsuccessful_envelope_reports_its_own_exit_status() {
let envelope = ResultEnvelope::new(Operation::Wait, Outcome::Timeout);
assert_eq!(envelope.exit_code(), Outcome::Timeout.exit_code());
assert!(!envelope.ok);
}
}