use std::path::{Path, PathBuf};
use crate::client::envelope::{Operation, Outcome, ResultEnvelope};
use crate::client::transport::{TransportError, UnixApiClient};
use crate::client::RouteSelector;
use crate::web::remote_control_api::dto::{
ApiError, CapabilitiesResponse, ChangeResource, ExecutionContractResponse,
ExecutionStatusResponse, InstanceResponse, StateResponse, API_VERSION,
};
pub const MAX_RECONCILE_ATTEMPTS: usize = 5;
#[derive(Debug)]
pub struct ConnectionRefusal {
outcome: Outcome,
message: String,
}
impl ConnectionRefusal {
pub fn into_envelope(self, operation: Operation) -> ResultEnvelope {
ResultEnvelope::new(operation, self.outcome).with_message(self.message)
}
}
#[derive(Debug)]
pub struct Connection {
client: UnixApiClient,
repo_root: Option<PathBuf>,
}
impl Connection {
pub fn resolve_route(
selector: &RouteSelector,
auth_token_env: Option<&str>,
) -> Result<Self, ConnectionRefusal> {
let (socket, repo_root) = match selector {
RouteSelector::Project(project_dir) => {
let route = crate::client::resolve_project(project_dir).map_err(|error| {
ConnectionRefusal {
outcome: Outcome::NotInRepository,
message: error.message,
}
})?;
(route.socket, Some(route.repo_root))
}
RouteSelector::Socket(path) => {
let workspace = current_workspace()?;
(path.to_path_buf(), discover_repo_root(&workspace))
}
RouteSelector::Default => {
let workspace = current_workspace()?;
let repo_root = discover_repo_root(&workspace);
let socket = match crate::repo_lock::discover_common_dir(&workspace) {
Some(common_dir) => crate::web::unix_socket::default_socket_path(&common_dir),
None => {
return Err(ConnectionRefusal {
outcome: Outcome::NotInRepository,
message: "the default owner socket needs a Git repository: run \
inside one, name the project with --project-dir PATH, or \
name the socket with --unix-socket PATH"
.to_string(),
})
}
};
(socket, repo_root)
}
};
let token = match auth_token_env {
Some(name) => match std::env::var(name) {
Ok(value) if !value.is_empty() => Some(value),
_ => {
return Err(ConnectionRefusal {
outcome: Outcome::AuthenticationFailed,
message: format!(
"environment variable '{name}' is unset or empty, so no bearer token \
could be presented"
),
})
}
},
None => None,
};
let client = UnixApiClient::new(socket, token).map_err(|rejection| ConnectionRefusal {
outcome: Outcome::AuthenticationFailed,
message: match auth_token_env {
Some(name) => format!(
"the bearer token in environment variable '{name}' contains {rejection}, \
which cannot be sent as an HTTP header value. The value is not shown"
),
None => format!(
"the configured bearer token contains {rejection}, which cannot be sent as \
an HTTP header value. The value is not shown"
),
},
})?;
Ok(Self { client, repo_root })
}
pub fn client(&self) -> &UnixApiClient {
&self.client
}
pub fn repo_root(&self) -> Option<&Path> {
self.repo_root.as_deref()
}
}
fn current_workspace() -> Result<PathBuf, ConnectionRefusal> {
std::env::current_dir().map_err(|error| ConnectionRefusal {
outcome: Outcome::NotInRepository,
message: format!("the current directory could not be resolved: {error}"),
})
}
pub(crate) fn discover_repo_root(workspace: &Path) -> Option<PathBuf> {
let output = std::process::Command::new("git")
.args(["rev-parse", "--show-toplevel"])
.current_dir(workspace)
.output()
.ok()?;
if !output.status.success() {
return None;
}
let text = String::from_utf8(output.stdout).ok()?;
let trimmed = text.trim();
if trimmed.is_empty() {
return None;
}
Some(PathBuf::from(trimmed))
}
pub struct Observation {
pub instance_id: String,
pub state_revision: u64,
pub event_sequence: u64,
pub capabilities: CapabilitiesResponse,
pub instance: InstanceResponse,
pub state: StateResponse,
pub execution: ExecutionStatusResponse,
pub contract: ExecutionContractResponse,
}
#[cfg(test)]
pub(crate) fn observation_for_test(changes: Vec<ChangeResource>) -> Observation {
use crate::web::remote_control_api::dto::{
CapabilityLimits, CommandExecutionCapability, InstanceSnapshot, ProcessExecutionStatus,
};
let instance_id = "i-test".to_string();
Observation {
instance_id: instance_id.clone(),
state_revision: 1,
event_sequence: 1,
capabilities: CapabilitiesResponse {
api_version: API_VERSION.to_string(),
instance_id: instance_id.clone(),
commands: Vec::new(),
transports: Vec::new(),
error_codes: Vec::new(),
limits: CapabilityLimits {
max_events: 0,
max_logs: 0,
max_commands: 0,
max_idempotency_records: 0,
command_record_ttl_secs: 0,
max_correlation_id_len: 0,
},
authentication_required: false,
command_execution: CommandExecutionCapability { available: true },
execution_sinks: Default::default(),
proposal_subscriptions: crate::web::completion_sink::proposal_capability(),
worktrees: Default::default(),
parallel: crate::web::remote_control_api::dto::ParallelCapabilities {
max_concurrent: 1,
vcs_backend: "git".to_string(),
blocked_reasons: Vec::new(),
},
},
instance: InstanceResponse {
instance_id: instance_id.clone(),
started_at: String::new(),
pid: 0,
version: String::new(),
api_version: API_VERSION.to_string(),
},
state: StateResponse {
instance_id: instance_id.clone(),
state_revision: 1,
event_sequence: 1,
snapshot: InstanceSnapshot {
changes,
..InstanceSnapshot::empty()
},
},
execution: ExecutionStatusResponse {
instance_id: instance_id.clone(),
state_revision: 1,
event_sequence: 1,
observed_at: String::new(),
process: ProcessExecutionStatus {
app_mode: "select".to_string(),
scheduler_running: false,
has_active_work: false,
active_activities: Vec::new(),
latest_log: None,
},
changes: Vec::new(),
},
contract: ExecutionContractResponse {
instance_id,
state_revision: 1,
contract: None,
},
}
}
impl Observation {
pub fn change(&self, change_id: &str) -> Option<&ChangeResource> {
self.state
.snapshot
.changes
.iter()
.find(|change| change.id == change_id)
}
pub fn command_capable(&self) -> bool {
self.capabilities.command_execution.available
}
#[allow(dead_code)] pub fn execution_sinks_available(&self) -> bool {
self.capabilities.execution_sinks.available
}
pub fn proposal_subscriptions_available(&self) -> bool {
self.capabilities.proposal_subscriptions.available
}
}
#[derive(Debug)]
pub enum ObserveError {
Incoherent(String),
Incompatible(String),
Unauthenticated(String),
NotRunning(String),
Transport(String),
}
impl ObserveError {
pub fn outcome(&self) -> Outcome {
match self {
Self::Incoherent(_) => Outcome::ObservationConflict,
Self::Incompatible(_) => Outcome::IncompatibleOwner,
Self::Unauthenticated(_) => Outcome::AuthenticationFailed,
Self::NotRunning(_) => Outcome::OwnerNotRunning,
Self::Transport(_) => Outcome::TransportError,
}
}
pub fn message(&self) -> &str {
match self {
Self::Incoherent(detail)
| Self::Incompatible(detail)
| Self::Unauthenticated(detail)
| Self::NotRunning(detail)
| Self::Transport(detail) => detail,
}
}
pub fn is_transient(&self) -> bool {
matches!(self, Self::Incoherent(_))
}
pub fn into_envelope(self, operation: Operation) -> ResultEnvelope {
let outcome = self.outcome();
let message = self.message().to_string();
ResultEnvelope::new(operation, outcome).with_message(message)
}
}
impl From<TransportError> for ObserveError {
fn from(error: TransportError) -> Self {
match &error {
TransportError::NotListening { .. } => Self::NotRunning(error.to_string()),
_ => Self::Transport(error.to_string()),
}
}
}
async fn read<T: serde::de::DeserializeOwned>(
client: &UnixApiClient,
path_and_query: &str,
) -> Result<T, ObserveError> {
let response = client.get(path_and_query).await?;
match response.status {
200 => response.json::<T>().map_err(|error| {
ObserveError::Incompatible(format!(
"'{path_and_query}' did not match this build's contract: {error}"
))
}),
401 => Err(ObserveError::Unauthenticated(describe_api_error(
&response.body,
"the owner requires a bearer token; supply one with --auth-token-env NAME",
))),
403 => Err(ObserveError::Unauthenticated(describe_api_error(
&response.body,
"the owner refused the presented credentials",
))),
404 => Err(ObserveError::Incompatible(format!(
"the owner does not serve '{path_and_query}', so it is not a compatible build"
))),
status => Err(ObserveError::Transport(format!(
"'{path_and_query}' returned unexpected status {status}"
))),
}
}
pub fn describe_api_error(body: &[u8], fallback: &str) -> String {
match serde_json::from_slice::<ApiError>(body) {
Ok(error) => format!("{} ({})", error.message, error.error_code.as_str()),
Err(_) => fallback.to_string(),
}
}
pub async fn observe(
connection: &Connection,
change_id: Option<&str>,
) -> Result<Observation, ObserveError> {
let client = connection.client();
let capabilities: CapabilitiesResponse = read(client, "/api/v2/capabilities").await?;
if capabilities.api_version != API_VERSION {
return Err(ObserveError::Incompatible(format!(
"the owner serves API version '{}', but this client speaks '{API_VERSION}'",
capabilities.api_version
)));
}
let instance: InstanceResponse = read(client, "/api/v2/instance").await?;
let contract_path = match change_id {
Some(change_id) => format!(
"/api/v2/execution-contract?change_id={}",
crate::client::transport::encode_query_value(change_id)
),
None => "/api/v2/execution-contract".to_string(),
};
let contract: ExecutionContractResponse = read(client, &contract_path).await?;
let execution: ExecutionStatusResponse = read(client, "/api/v2/execution-status").await?;
let state: StateResponse = read(client, "/api/v2/state").await?;
let instances = [
capabilities.instance_id.as_str(),
instance.instance_id.as_str(),
contract.instance_id.as_str(),
execution.instance_id.as_str(),
state.instance_id.as_str(),
];
if instances.iter().any(|id| *id != state.instance_id) {
return Err(ObserveError::Incoherent(
"the owner's resources reported different process incarnations, so they cannot be \
one observation"
.to_string(),
));
}
if execution.event_sequence > state.event_sequence {
return Err(ObserveError::Incoherent(
"the owner's event cursor moved backwards between reads".to_string(),
));
}
if contract.state_revision != state.state_revision
|| execution.state_revision != state.state_revision
{
let contract: ExecutionContractResponse = read(client, &contract_path).await?;
let execution: ExecutionStatusResponse = read(client, "/api/v2/execution-status").await?;
if contract.instance_id != state.instance_id || execution.instance_id != state.instance_id {
return Err(ObserveError::Incoherent(
"the owner restarted while the observation was being reconciled".to_string(),
));
}
if contract.state_revision != state.state_revision
|| execution.state_revision != state.state_revision
{
return Err(ObserveError::Incoherent(format!(
"the owner advanced past revision {} while the observation was being reconciled",
state.state_revision
)));
}
return Ok(Observation {
instance_id: state.instance_id.clone(),
state_revision: state.state_revision,
event_sequence: state.event_sequence,
capabilities,
instance,
state,
execution,
contract,
});
}
Ok(Observation {
instance_id: state.instance_id.clone(),
state_revision: state.state_revision,
event_sequence: state.event_sequence,
capabilities,
instance,
state,
execution,
contract,
})
}
pub async fn observe_bounded(
connection: &Connection,
change_id: Option<&str>,
attempts: usize,
) -> Result<Observation, ObserveError> {
let mut last = None;
for _ in 0..attempts.max(1) {
match observe(connection, change_id).await {
Ok(observation) => return Ok(observation),
Err(error) if error.is_transient() => last = Some(error),
Err(error) => return Err(error),
}
}
Err(last.unwrap_or_else(|| {
ObserveError::Incoherent("no coherent observation was produced".to_string())
}))
}
pub async fn status(connection: &Connection) -> ResultEnvelope {
let observation = match observe_bounded(connection, None, MAX_RECONCILE_ATTEMPTS).await {
Ok(observation) => observation,
Err(error) => return error.into_envelope(Operation::Status),
};
let changes: Vec<serde_json::Value> = observation
.state
.snapshot
.changes
.iter()
.map(|change| {
let execution = observation
.execution
.changes
.iter()
.find(|status| status.id == change.id);
serde_json::json!({
"id": change.id,
"display_status": change.display_status,
"queue_intent": change.queue_intent,
"execution_marked": change.execution_marked,
"parallel_eligible": change.parallel.eligible,
"actions": change.actions,
"blocker": change.blocker,
"execution_state": execution.map(|status| status.execution_state),
"current_phase": execution.map(|status| status.current_phase),
"last_completed_phase": execution.and_then(|status| status.last_completed_phase),
})
})
.collect();
let detail = serde_json::json!({
"socket": connection.client().socket().display().to_string(),
"state_revision": observation.state_revision,
"event_sequence": observation.event_sequence,
"owner": {
"pid": observation.instance.pid,
"version": observation.instance.version,
"started_at": observation.instance.started_at,
"api_version": observation.instance.api_version,
},
"process": {
"app_mode": observation.execution.process.app_mode,
"scheduler_running": observation.execution.process.scheduler_running,
"has_active_work": observation.execution.process.has_active_work,
"active_activities": observation.execution.process.active_activities,
"process_error": observation.state.snapshot.process_error,
"is_resolving": observation.state.snapshot.is_resolving,
"persistent_scheduler_idle": observation.state.snapshot.persistent_scheduler_idle,
},
"command_execution_available": observation.command_capable(),
"authentication_required": observation.capabilities.authentication_required,
"execution_contract": observation.contract.contract,
"totals": observation.state.snapshot.totals,
"changes": changes,
});
ResultEnvelope::new(Operation::Status, Outcome::Observed)
.with_instance(Some(observation.instance_id))
.with_detail(detail)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::web::remote_control_api::dto::ErrorCode;
#[test]
fn an_absent_socket_outside_a_repository_names_both_choices() {
let tmp = tempfile::tempdir().unwrap();
let previous = std::env::current_dir().unwrap();
std::env::set_current_dir(tmp.path()).unwrap();
let resolved = Connection::resolve_route(&RouteSelector::Default, None);
std::env::set_current_dir(previous).unwrap();
if let Err(refusal) = resolved {
assert_eq!(refusal.outcome, Outcome::NotInRepository);
assert!(refusal.message.contains("--unix-socket"), "{refusal:?}");
}
}
#[test]
fn an_unset_token_variable_fails_closed_instead_of_connecting_anonymously() {
let refusal = Connection::resolve_route(
&RouteSelector::Socket(PathBuf::from("/tmp/cflx-client-test.sock")),
Some("CFLX_CLIENT_TEST_TOKEN_THAT_IS_NOT_SET"),
)
.expect_err("an unset variable must refuse");
assert_eq!(refusal.outcome, Outcome::AuthenticationFailed);
assert!(refusal.message.contains("unset or empty"), "{refusal:?}");
}
#[test]
fn an_explicit_socket_needs_no_repository_identity() {
let connection = Connection::resolve_route(
&RouteSelector::Socket(PathBuf::from("/tmp/explicit.sock")),
None,
)
.expect("an explicit socket always resolves");
assert_eq!(
connection.client().socket(),
Path::new("/tmp/explicit.sock")
);
assert!(!connection.client().has_token());
}
#[test]
fn a_typed_api_error_is_described_by_its_own_message_and_code() {
let body = serde_json::to_vec(&ApiError::new(
ErrorCode::Unauthorized,
"missing bearer credentials",
"corr-1",
))
.unwrap();
let described = describe_api_error(&body, "fallback");
assert!(described.contains("missing bearer credentials"));
assert!(described.contains("unauthorized"));
}
#[test]
fn an_untyped_body_is_never_echoed_back_to_the_caller() {
let described = describe_api_error(b"<html>secret</html>", "fallback text");
assert_eq!(described, "fallback text");
}
#[test]
fn only_incoherence_is_worth_rereading() {
assert!(ObserveError::Incoherent(String::new()).is_transient());
for error in [
ObserveError::Incompatible(String::new()),
ObserveError::Unauthenticated(String::new()),
ObserveError::NotRunning(String::new()),
ObserveError::Transport(String::new()),
] {
assert!(!error.is_transient(), "{error:?}");
}
}
#[test]
fn observation_failures_map_to_stable_outcomes() {
assert_eq!(
ObserveError::Incoherent(String::new()).outcome(),
Outcome::ObservationConflict
);
assert_eq!(
ObserveError::Incompatible(String::new()).outcome(),
Outcome::IncompatibleOwner
);
assert_eq!(
ObserveError::Unauthenticated(String::new()).outcome(),
Outcome::AuthenticationFailed
);
assert_eq!(
ObserveError::NotRunning(String::new()).outcome(),
Outcome::OwnerNotRunning
);
assert_eq!(
ObserveError::Transport(String::new()).outcome(),
Outcome::TransportError
);
}
}