use std::path::Path;
use std::process::Output;
fn run_cli(cwd: &Path, args: &[&str], env: &[(&str, &str)]) -> Output {
let mut command = std::process::Command::new(env!("CARGO_BIN_EXE_cflx"));
command.args(args).current_dir(cwd);
command.env_remove("CFLX_CLIENT_TEST_TOKEN");
for (name, value) in env {
command.env(name, value);
}
command
.output()
.expect("the compiled cflx binary must be runnable")
}
fn stdout_of(output: &Output) -> String {
String::from_utf8_lossy(&output.stdout).to_string()
}
fn stderr_of(output: &Output) -> String {
String::from_utf8_lossy(&output.stderr).to_string()
}
#[test]
fn feature_disabled_or_enabled_client_never_takes_the_repository_lock() {
let tmp = tempfile::tempdir().expect("temp dir");
let output = run_cli(tmp.path(), &["client", "status", "--json"], &[]);
for artifact in ["cflx-owner.json", "cflx-api.sock", ".cflx"] {
assert!(
!tmp.path().join(artifact).exists(),
"client must not create {artifact}"
);
}
assert!(!output.status.success());
}
#[test]
fn client_subscribe_help_and_usage_documents_an_argv_callback_and_proposal_scope() {
let tmp = tempfile::tempdir().expect("temp dir");
let namespace = run_cli(tmp.path(), &["client", "--help"], &[]);
assert!(namespace.status.success(), "{}", stderr_of(&namespace));
let namespace_help = stdout_of(&namespace);
assert!(
namespace_help.contains("subscribe"),
"the namespace must offer the group:\n{namespace_help}"
);
let output = run_cli(tmp.path(), &["client", "subscribe", "--help"], &[]);
assert!(output.status.success(), "{}", stderr_of(&output));
let help = stdout_of(&output);
for expected in [
"set",
"get",
"clear",
"never shell source",
"execution completion, not process completion",
"transport_not_permitted",
"CFLX_EVENT_PATH",
"--unix-socket",
"--auth-token-env",
] {
assert!(
help.contains(expected),
"subscribe help must mention {expected}:\n{help}"
);
}
assert!(
help.contains("never resumes") || help.contains("resumes no agent"),
"subscribe help must say delivery does not resume an agent:\n{help}"
);
for forbidden in [
"--expected-revision",
"--idempotency-key",
"--command-type",
"--execution-mark",
"--auth-token ",
"--shell",
] {
assert!(
!help.contains(forbidden),
"subscribe help must not expose {forbidden}:\n{help}"
);
}
let set = run_cli(tmp.path(), &["client", "subscribe", "set", "--help"], &[]);
assert!(set.status.success(), "{}", stderr_of(&set));
let set_help = stdout_of(&set);
assert!(
set_help.contains("-- <COMMAND>..."),
"the separator must be part of the documented usage:\n{set_help}"
);
for expected in ["--blocked", "--instance-id", "--json", "<CHANGE_IDS>"] {
assert!(
set_help.contains(expected),
"set help must mention {expected}:\n{set_help}"
);
}
assert!(
std::fs::read_dir(tmp.path())
.expect("temp dir readable")
.next()
.is_none(),
"printing help must write nothing"
);
}
#[test]
fn top_level_help_synopsis_names_current_client_verbs_and_retires_enqueue() {
let tmp = tempfile::tempdir().expect("temp dir");
let top_level = run_cli(tmp.path(), &["--help"], &[]);
assert!(top_level.status.success(), "{}", stderr_of(&top_level));
let help = stdout_of(&top_level);
let synopsis = help
.lines()
.find(|line| line.trim_start().starts_with("client "))
.unwrap_or_else(|| {
panic!("the top-level help must describe the client namespace:\n{help}")
});
for verb in [
"status",
"mark",
"start",
"stop",
"wait",
"subscribe",
"mcp",
] {
assert!(
synopsis.contains(verb),
"the top-level synopsis must name {verb}:\n{synopsis}"
);
}
for retired in ["enqueue", "notify"] {
assert!(
!help.contains(retired),
"{retired} is retired and must not appear in the top-level help:\n{help}"
);
}
}
#[test]
fn client_control_help_documents_operator_verbs_and_retires_enqueue() {
let tmp = tempfile::tempdir().expect("temp dir");
let namespace = run_cli(tmp.path(), &["client", "--help"], &[]);
assert!(namespace.status.success(), "{}", stderr_of(&namespace));
let help = stdout_of(&namespace);
for verb in [
"status",
"mark",
"unmark",
"start",
"stop",
"force-stop",
"force-stop-change",
"wait",
"subscribe",
"mcp",
] {
assert!(
help.contains(verb),
"the namespace must offer {verb}:\n{help}"
);
}
for retired in ["enqueue", "notify"] {
assert!(
!help.contains(retired),
"{retired} is retired and must not be offered:\n{help}"
);
}
let mark = run_cli(tmp.path(), &["client", "mark", "--help"], &[]);
assert!(mark.status.success(), "{}", stderr_of(&mark));
let mark_help = stdout_of(&mark);
assert!(
mark_help.contains("preserves every unrelated mark"),
"mark help must state the target scope:\n{mark_help}"
);
assert!(
mark_help.contains("does not construct queue intent"),
"mark help must state that admission is not claimed:\n{mark_help}"
);
let start = run_cli(tmp.path(), &["client", "start", "--help"], &[]);
assert!(start.status.success(), "{}", stderr_of(&start));
let start_help = stdout_of(&start);
assert!(
start_help.contains("F5"),
"start help must name its TUI equivalent:\n{start_help}"
);
assert!(
start_help.contains("no target list"),
"start help must say it takes no targets:\n{start_help}"
);
}
#[test]
fn client_subscribe_help_and_usage_rejects_a_malformed_invocation_before_any_request() {
let tmp = tempfile::tempdir().expect("temp dir");
let instance = "a".repeat(32);
let oversized_instance = "a".repeat(129);
let cases: Vec<(Vec<&str>, &str, &str)> = vec![
(
vec![
"subscribe",
"set",
"alpha",
"--instance-id",
&instance,
"--json",
"--",
],
"subscribe_set",
"a callback with no program at all",
),
(
vec![
"subscribe",
"set",
"alpha",
"--instance-id",
&instance,
"--json",
"/bin/true",
],
"subscribe_set",
"a callback that skipped the separator",
),
(
vec!["subscribe", "set", "alpha", "--json", "--", "/bin/true"],
"subscribe_set",
"a registration that named no owner incarnation",
),
(
vec![
"subscribe",
"set",
"../escape",
"--instance-id",
&instance,
"--json",
"--",
"/bin/true",
],
"subscribe_set",
"an escaping change ID",
),
(
vec!["subscribe", "get", "alpha", "--instance-id", "", "--json"],
"subscribe_get",
"an empty instance ID",
),
(
vec![
"subscribe",
"clear",
"alpha",
"--instance-id",
"has space",
"--json",
],
"subscribe_clear",
"an instance ID with a separator in it",
),
(
vec![
"subscribe",
"get",
"alpha",
"--instance-id",
".reserved",
"--json",
],
"subscribe_get",
"an instance ID starting with a reserved character",
),
(
vec![
"subscribe",
"get",
"alpha",
"--instance-id",
&oversized_instance,
"--json",
],
"subscribe_get",
"an oversized instance ID",
),
(
vec!["subscribe", "get", "--instance-id", &instance, "--json"],
"subscribe_get",
"a request that named no proposal",
),
(
vec!["subscribe", "--json"],
"subscribe_get",
"a group with no operation",
),
];
for (action, operation, what) in cases {
let mut args = vec!["client"];
args.extend(action);
let output = run_cli(tmp.path(), &args, &[]);
assert_eq!(
output.status.code(),
Some(2),
"{what} must be a usage error, stderr={}",
stderr_of(&output)
);
let stdout = stdout_of(&output);
let parsed: serde_json::Value = serde_json::from_str(stdout.trim())
.unwrap_or_else(|e| panic!("{what}: stdout must be one envelope, got {stdout:?}: {e}"));
assert_eq!(parsed["outcome"], "usage_error", "{what}");
assert_eq!(parsed["operation"], operation, "{what}");
assert!(
std::fs::read_dir(tmp.path())
.expect("temp dir readable")
.next()
.is_none(),
"{what}: a rejected invocation must write nothing"
);
}
}
#[test]
fn client_subscribe_help_and_usage_keeps_a_callback_flag_out_of_the_output_mode() {
let tmp = tempfile::tempdir().expect("temp dir");
let absent = tmp.path().join("absent.sock");
let socket = absent.display().to_string();
let instance = "a".repeat(32);
let human = run_cli(
tmp.path(),
&[
"client",
"--unix-socket",
&socket,
"subscribe",
"set",
"alpha",
"--instance-id",
&instance,
"--",
"/bin/true",
"--json",
],
&[],
);
let stdout = stdout_of(&human);
assert!(
!stdout.trim_start().starts_with('{'),
"a callback argument must not select machine output:\n{stdout}"
);
assert!(
stdout.starts_with("subscribe_set: "),
"the human line still names the operation:\n{stdout}"
);
let machine = run_cli(
tmp.path(),
&[
"client",
"--unix-socket",
&socket,
"subscribe",
"set",
"alpha",
"--instance-id",
&instance,
"--json",
"--",
"/bin/true",
"--json",
],
&[],
);
let stdout = stdout_of(&machine);
let parsed: serde_json::Value = serde_json::from_str(stdout.trim())
.unwrap_or_else(|e| panic!("stdout must be one envelope, got {stdout:?}: {e}"));
assert_eq!(parsed["operation"], "subscribe_set");
assert!(!parsed["ok"].as_bool().expect("ok is a boolean"));
}
#[cfg(not(feature = "web-monitoring"))]
mod feature_disabled {
use super::*;
#[test]
fn feature_disabled_mcp_refuses_on_stderr_without_any_side_effect() {
let tmp = tempfile::tempdir().expect("temp dir");
let output = run_cli(tmp.path(), &["client", "mcp"], &[]);
assert_eq!(output.status.code(), Some(20));
assert!(
stdout_of(&output).is_empty(),
"the protocol stream must carry nothing at all"
);
assert!(
stderr_of(&output).contains("web-monitoring"),
"the refusal must name the missing feature"
);
assert!(
std::fs::read_dir(tmp.path())
.expect("temp dir readable")
.next()
.is_none(),
"a refused client command must write nothing"
);
}
#[test]
fn feature_disabled_client_refuses_before_any_side_effect() {
let tmp = tempfile::tempdir().expect("temp dir");
for action in [
vec!["client", "status", "--json"],
vec!["client", "mark", "alpha", "--json"],
vec!["client", "unmark", "alpha", "--json"],
vec!["client", "start", "--json"],
vec!["client", "stop", "--json"],
vec!["client", "force-stop", "--json"],
vec!["client", "force-stop-change", "alpha", "--json"],
vec!["client", "wait", "alpha", "--json"],
vec![
"client",
"subscribe",
"set",
"alpha",
"--instance-id",
"owner",
"--json",
"--",
"/bin/true",
],
vec![
"client",
"subscribe",
"get",
"alpha",
"--instance-id",
"owner",
"--json",
],
vec![
"client",
"subscribe",
"clear",
"alpha",
"--instance-id",
"owner",
"--json",
],
] {
let output = run_cli(tmp.path(), &action, &[]);
let stdout = stdout_of(&output);
let parsed: serde_json::Value =
serde_json::from_str(stdout.trim()).unwrap_or_else(|e| {
panic!("stdout must be one JSON envelope, got {stdout:?}: {e}")
});
assert_eq!(parsed["ok"], false);
assert_eq!(parsed["outcome"], "feature_unavailable");
assert_eq!(output.status.code(), Some(20));
assert!(
stderr_of(&output).contains("web-monitoring"),
"the refusal must name the missing feature"
);
assert!(
std::fs::read_dir(tmp.path())
.expect("temp dir readable")
.next()
.is_none(),
"a refused client command must write nothing"
);
}
}
}
#[cfg(feature = "web-monitoring")]
mod enabled {
use super::*;
use std::path::PathBuf;
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use std::sync::{Arc, Mutex};
use std::time::Duration;
use async_trait::async_trait;
use conflux::orchestration::execution_facts::ExecutionFactsStore;
use conflux::orchestration::operator_command::RunBoundaryLiveness;
use conflux::web::remote_control_api::auth::RemoteControlAuth;
use conflux::web::remote_control_api::dto::{
CommandSpec, ErrorCode, OwnerExecutionContract, TerminalMode,
};
use conflux::web::remote_control_api::executor::{
CommandFailure, ExecutionSummary, RemoteControlExecutor,
};
use conflux::web::remote_control_api::projection::Projection;
use conflux::web::remote_control_api::{router, RemoteControlRuntime, RemoteControlState};
#[derive(Default)]
struct SpyExecutor {
calls: Mutex<Vec<CommandSpec>>,
script: Mutex<Vec<Result<ExecutionSummary, CommandFailure>>>,
}
impl SpyExecutor {
fn new() -> Arc<Self> {
Arc::new(Self::default())
}
fn script(&self, outcomes: Vec<Result<ExecutionSummary, CommandFailure>>) {
*self.script.lock().unwrap() = outcomes;
}
fn calls(&self) -> Vec<CommandSpec> {
self.calls.lock().unwrap().clone()
}
fn call_count(&self) -> usize {
self.calls.lock().unwrap().len()
}
}
#[async_trait]
impl RemoteControlExecutor for SpyExecutor {
async fn execute(&self, command: &CommandSpec) -> Result<ExecutionSummary, CommandFailure> {
self.calls.lock().unwrap().push(command.clone());
let mut script = self.script.lock().unwrap();
if script.len() > 1 {
script.remove(0)
} else {
script
.first()
.cloned()
.unwrap_or_else(|| Ok(ExecutionSummary::changed("applied")))
}
}
}
#[derive(Default)]
struct Boundary {
running: AtomicBool,
}
impl Boundary {
fn set_running(&self, running: bool) {
self.running.store(running, Ordering::SeqCst);
}
}
impl RunBoundaryLiveness for Boundary {
fn boundary_running(&self) -> bool {
self.running.load(Ordering::SeqCst)
}
}
#[derive(Debug, Clone)]
struct CommandExchange {
command_type: String,
expected_revision: u64,
idempotency_key: String,
status: u16,
record_id: Option<String>,
error_code: Option<String>,
}
#[derive(Default)]
struct ApiSpy {
requests: Mutex<Vec<String>>,
exchanges: Mutex<Vec<CommandExchange>>,
before_command: Mutex<Vec<Box<dyn Fn() + Send + Sync>>>,
before_state: Mutex<Vec<Box<dyn Fn() + Send + Sync>>>,
}
impl ApiSpy {
fn new() -> Arc<Self> {
Arc::new(Self::default())
}
fn inject_before_commands(&self, script: Vec<Box<dyn Fn() + Send + Sync>>) {
*self.before_command.lock().unwrap() = script;
}
fn inject_before_state_reads(&self, script: Vec<Box<dyn Fn() + Send + Sync>>) {
*self.before_state.lock().unwrap() = script;
}
fn requests(&self) -> Vec<String> {
self.requests.lock().unwrap().clone()
}
fn exchanges(&self) -> Vec<CommandExchange> {
self.exchanges.lock().unwrap().clone()
}
async fn intercept(
self: Arc<Self>,
request: axum::extract::Request,
next: axum::middleware::Next,
) -> axum::response::Response {
let method = request.method().to_string();
let path = request.uri().path().to_string();
self.requests
.lock()
.unwrap()
.push(format!("{method} {path}"));
if method == "GET" && path == "/api/v2/state" {
let advance = {
let mut script = self.before_state.lock().unwrap();
(!script.is_empty()).then(|| script.remove(0))
};
if let Some(advance) = advance {
advance();
}
}
if method != "POST" || path != "/api/v2/commands" {
return next.run(request).await;
}
let (parts, body) = request.into_parts();
let submitted_bytes = axum::body::to_bytes(body, usize::MAX)
.await
.expect("the command request body is readable");
let submitted: serde_json::Value =
serde_json::from_slice(&submitted_bytes).expect("the command request body is JSON");
let advance = {
let mut script = self.before_command.lock().unwrap();
(!script.is_empty()).then(|| script.remove(0))
};
if let Some(advance) = advance {
advance();
}
let response = next
.run(axum::extract::Request::from_parts(
parts,
axum::body::Body::from(submitted_bytes),
))
.await;
let status = response.status().as_u16();
let (parts, body) = response.into_parts();
let answered_bytes = axum::body::to_bytes(body, usize::MAX)
.await
.expect("the command response body is readable");
let answered: serde_json::Value =
serde_json::from_slice(&answered_bytes).unwrap_or(serde_json::Value::Null);
self.exchanges.lock().unwrap().push(CommandExchange {
command_type: submitted["type"].as_str().unwrap_or_default().to_string(),
expected_revision: submitted["expected_revision"].as_u64().unwrap_or_default(),
idempotency_key: submitted["idempotency_key"]
.as_str()
.unwrap_or_default()
.to_string(),
status,
record_id: answered["command_id"].as_str().map(str::to_string),
error_code: answered["error_code"].as_str().map(str::to_string),
});
axum::response::Response::from_parts(parts, axum::body::Body::from(answered_bytes))
}
}
#[derive(Default, Clone, Copy)]
struct OwnerExtras {
sinks: bool,
transport: Option<conflux::web::remote_control_api::ApiTransport>,
}
impl OwnerExtras {
fn sink_capable() -> Self {
Self {
sinks: true,
transport: Some(conflux::web::remote_control_api::ApiTransport::Unix),
}
}
}
struct Owner {
socket: PathBuf,
projection: Arc<Projection>,
runtime: Arc<RemoteControlRuntime>,
boundary: Arc<Boundary>,
facts: Arc<ExecutionFactsStore>,
sinks: Option<Arc<conflux::web::completion_sink::CompletionSinkRegistry>>,
dispatch: AtomicUsize,
shutdown: tokio_util::sync::CancellationToken,
task: tokio::task::JoinHandle<()>,
_dir: Option<tempfile::TempDir>,
}
impl Owner {
async fn start(
executor: Option<Arc<dyn RemoteControlExecutor>>,
token: Option<&str>,
) -> Self {
let dir = tempfile::tempdir().expect("temp dir");
let socket = dir.path().join("cflx-api.sock");
let mut owner = Self::start_on(socket, executor, token).await;
owner._dir = Some(dir);
owner
}
async fn start_with(
executor: Option<Arc<dyn RemoteControlExecutor>>,
token: Option<&str>,
extras: OwnerExtras,
) -> Self {
let dir = tempfile::tempdir().expect("temp dir");
let socket = dir.path().join("cflx-api.sock");
let mut owner = Self::start_layered(socket, executor, token, None, extras).await;
owner._dir = Some(dir);
owner
}
async fn start_intercepted(
executor: Option<Arc<dyn RemoteControlExecutor>>,
token: Option<&str>,
api: Arc<ApiSpy>,
) -> Self {
let dir = tempfile::tempdir().expect("temp dir");
let socket = dir.path().join("cflx-api.sock");
let mut owner =
Self::start_layered(socket, executor, token, Some(api), OwnerExtras::default())
.await;
owner._dir = Some(dir);
owner
}
async fn start_on(
socket: PathBuf,
executor: Option<Arc<dyn RemoteControlExecutor>>,
token: Option<&str>,
) -> Self {
Self::start_layered(socket, executor, token, None, OwnerExtras::default()).await
}
async fn start_layered(
socket: PathBuf,
executor: Option<Arc<dyn RemoteControlExecutor>>,
token: Option<&str>,
api: Option<Arc<ApiSpy>>,
extras: OwnerExtras,
) -> Self {
let runtime = Arc::new(RemoteControlRuntime::new());
if let Some(executor) = executor {
runtime.bind(executor).await;
}
let facts = Arc::new(ExecutionFactsStore::new());
runtime.bind_execution_facts(facts.clone());
let boundary = Arc::new(Boundary::default());
runtime.bind_run_boundary(boundary.clone());
let sinks = extras.sinks.then(|| {
let registry = conflux::web::completion_sink::CompletionSinkRegistry::start(
runtime.projection().instance_id().to_string(),
facts.clone(),
runtime.execution_contract(),
);
runtime.bind_completion_sinks(registry.clone());
registry
});
let auth = RemoteControlAuth::new(token.map(str::to_string), &[])
.expect("test auth policy is valid");
let app = router(
RemoteControlState::new(runtime.projection(), Arc::new(auth), runtime.clone())
.with_gate(runtime.gate())
.with_execution_facts(runtime.execution_facts())
.with_execution_contract(runtime.execution_contract())
.with_completion_sinks(runtime.completion_sinks()),
);
let app = match api {
Some(api) => app.layer(axum::middleware::from_fn(
move |request: axum::extract::Request, next: axum::middleware::Next| {
let api = api.clone();
async move { api.intercept(request, next).await }
},
)),
None => app,
};
let app = match extras.transport {
Some(transport) => app.layer(axum::Extension(transport)),
None => app,
};
let listener = tokio::net::UnixListener::bind(&socket).expect("binds the test socket");
let shutdown = tokio_util::sync::CancellationToken::new();
let token_for_task = shutdown.clone();
let task = tokio::spawn(async move {
let _ = axum::serve(listener, app)
.with_graceful_shutdown(async move { token_for_task.cancelled().await })
.await;
});
Self {
socket,
projection: runtime.projection(),
runtime,
boundary,
facts,
sinks,
dispatch: AtomicUsize::new(0),
shutdown,
task,
_dir: None,
}
}
fn instance_id(&self) -> String {
self.projection.instance_id().to_string()
}
fn admit(&self, change_id: &str) -> String {
use conflux::events::{ExecutionEvent, OperatorCommandEffect};
use conflux::orchestration::state::{OrchestratorState, ReducerCommand};
let mut reducer = OrchestratorState::new(vec![change_id.to_string()], 0);
reducer.apply_command(ReducerCommand::AddToQueue(change_id.to_string()));
let event = ExecutionEvent::OperatorCommandApplied {
effect: OperatorCommandEffect::QueueDelta {
change_id: change_id.to_string(),
queued: true,
},
};
reducer.apply_execution_event(&event);
let dispatch = self.dispatch.fetch_add(1, Ordering::SeqCst) as u64 + 1;
self.facts
.observe(dispatch, &event, Some(&reducer), chrono::Utc::now());
self.facts
.change(change_id)
.execution_id
.expect("admission opens an execution episode")
}
async fn admit_registered(&self, change_id: &str) -> String {
let execution_id = self.admit(change_id);
let registry = self
.sinks
.as_ref()
.expect("a sink-capable owner holds a registry");
let instance_id = self.instance_id();
for _ in 0..1_000 {
if registry
.view(&execution_id, &instance_id, change_id)
.is_ok()
{
return execution_id;
}
tokio::time::sleep(Duration::from_millis(5)).await;
}
panic!("the registry never opened an entry for execution {execution_id}");
}
fn socket(&self) -> String {
self.socket.display().to_string()
}
fn publish(&self, snapshot: conflux::web::remote_control_api::dto::InstanceSnapshot) {
self.projection
.apply_state("test_snapshot", None, serde_json::json!({}), snapshot);
}
fn contract(&self, contract: OwnerExecutionContract) {
self.runtime.bind_execution_contract(contract);
}
async fn stop(self) {
self.shutdown.cancel();
let _ = self.task.await;
}
}
use conflux::web::remote_control_api::dto::{
AttentionState, BlockerKind, ChangeBlocker, ChangeResource, ChangeTiming, InstanceSnapshot,
ParallelEligibility, ParallelRuntimeState, QueueIntent, SnapshotTotals,
};
fn change(id: &str, app_mode: &str, display_status: &str) -> ChangeResource {
ChangeResource {
id: id.to_string(),
display_status: display_status.to_string(),
progress_status: "pending".to_string(),
completed_tasks: 0,
total_tasks: 2,
progress_percent: 0.0,
dependencies: Vec::new(),
iteration_number: None,
execution_marked: false,
queue_intent: QueueIntent::NotQueued,
attention: AttentionState::None,
blocker: None,
error_detail: None,
actions: conflux::web::remote_control_api::projection::change_actions_for_test(
app_mode,
display_status,
None,
),
parallel: ParallelEligibility::default(),
timing: ChangeTiming::default(),
latest_activity: None,
worktree: None,
}
}
fn blocked_change(id: &str, app_mode: &str, kind: BlockerKind) -> ChangeResource {
let blocker = ChangeBlocker {
status: "blocked".to_string(),
kind,
category: match kind {
BlockerKind::External => Some("pending_verification".to_string()),
_ => None,
},
detail: match kind {
BlockerKind::External => Some("waiting on the signing certificate".to_string()),
BlockerKind::Dependency => Some("waiting on 'upstream'".to_string()),
BlockerKind::None => None,
},
unblock_condition: match kind {
BlockerKind::External => Some("the certificate is issued".to_string()),
_ => None,
},
prerequisite_owner: match kind {
BlockerKind::External => Some("release".to_string()),
_ => None,
},
origin: Some("apply".to_string()),
resumable: true,
dependencies: match kind {
BlockerKind::Dependency => vec!["upstream".to_string()],
_ => Vec::new(),
},
};
ChangeResource {
actions: conflux::web::remote_control_api::projection::change_actions_for_test(
app_mode,
"blocked",
Some(&blocker),
),
blocker: Some(blocker),
..change(id, app_mode, "blocked")
}
}
fn snapshot(app_mode: &str, changes: Vec<ChangeResource>) -> InstanceSnapshot {
let total = changes.len();
InstanceSnapshot {
app_mode: app_mode.to_string(),
persistent_scheduler_idle: false,
is_resolving: false,
process_error: None,
parallel: ParallelRuntimeState::default(),
changes,
totals: SnapshotTotals {
total,
completed: 0,
in_progress: 0,
pending: total,
},
}
}
fn merged_contract(base: &str) -> OwnerExecutionContract {
OwnerExecutionContract {
base_branch: base.to_string(),
terminal_mode: TerminalMode::Merged,
remote: None,
pushed_branch: None,
}
}
fn envelope(output: &Output) -> serde_json::Value {
let stdout = stdout_of(output);
let trimmed = stdout.trim();
assert!(
!trimmed.contains('\n'),
"JSON stdout must be exactly one object, got:\n{stdout}"
);
serde_json::from_str(trimmed)
.unwrap_or_else(|e| panic!("stdout must be one JSON envelope, got {stdout:?}: {e}"))
}
fn neutral_cwd() -> tempfile::TempDir {
tempfile::tempdir().expect("temp dir")
}
#[test]
fn cli_surface_exposes_only_the_operator_verbs_wait_subscribe_and_mcp() {
let tmp = neutral_cwd();
let output = run_cli(tmp.path(), &["client", "--help"], &[]);
assert!(output.status.success(), "{}", stderr_of(&output));
let help = stdout_of(&output);
for expected in [
"status",
"mark",
"unmark",
"start",
"stop",
"force-stop",
"wait",
"subscribe",
"mcp",
"--project-dir",
"--unix-socket",
"--auth-token-env",
] {
assert!(
help.contains(expected),
"help must mention {expected}:\n{help}"
);
}
for forbidden in [
"--expected-revision",
"--idempotency-key",
"--command-type",
"--queue-intent",
"--execution-mark",
"--auth-token ",
] {
assert!(
!help.contains(forbidden),
"help must not expose {forbidden}:\n{help}"
);
}
for absent in [
"enqueue", "notify", "retry", "resolve", "worktree", "archive",
] {
assert!(
!help.contains(&format!(" {absent} ")),
"client must not offer a {absent} subcommand:\n{help}"
);
}
}
#[test]
fn cli_surface_documents_the_mcp_server_as_a_bounded_client() {
let tmp = neutral_cwd();
let output = run_cli(tmp.path(), &["client", "mcp", "--help"], &[]);
assert!(output.status.success(), "{}", stderr_of(&output));
let help = stdout_of(&output);
for tool in ["cflx_status", "cflx_control", "cflx_subscribe"] {
assert!(help.contains(tool), "help must name {tool}:\n{help}");
}
for retired in ["cflx_enqueue", "cflx_notify"] {
assert!(
!help.contains(retired),
"{retired} is retired and must not be named:\n{help}"
);
}
assert!(
help.contains("cflx_wait is deliberately absent"),
"help must explain why there is no wait tool:\n{help}"
);
assert!(
help.contains("token value is never accepted"),
"help must state the credential rule:\n{help}"
);
for forbidden in ["--expected-revision", "--idempotency-key", "--auth-token "] {
assert!(
!help.contains(forbidden),
"help must not expose {forbidden}:\n{help}"
);
}
}
#[test]
fn cli_surface_rejects_a_malformed_change_id_as_usage() {
let tmp = neutral_cwd();
for bad in ["../escape", ".hidden", "-leading", "has space", ""] {
let output = run_cli(tmp.path(), &["client", "mark", bad, "--json"], &[]);
assert_eq!(
output.status.code(),
Some(2),
"'{bad}' must be a usage error, stderr={}",
stderr_of(&output)
);
let parsed = envelope(&output);
assert_eq!(parsed["outcome"], "usage_error");
assert_eq!(parsed["operation"], "control_mark");
}
}
#[test]
fn cli_surface_rejects_a_malformed_timeout_as_usage() {
let tmp = neutral_cwd();
for bad in ["50ms", "99ms", "abc", "5x", "-1", "99999999999999h", ""] {
let output = run_cli(
tmp.path(),
&["client", "wait", "alpha", "--timeout", bad, "--json"],
&[],
);
assert_eq!(
output.status.code(),
Some(2),
"'{bad}' must be a usage error, stderr={}",
stderr_of(&output)
);
let parsed = envelope(&output);
assert_eq!(parsed["outcome"], "usage_error");
assert_eq!(parsed["operation"], "wait");
}
}
#[test]
fn cli_surface_accepts_every_zero_spelling_as_the_unbounded_sentinel() {
let tmp = neutral_cwd();
let socket = tmp.path().join("absent.sock");
for zero in ["0", "0s", "0ms", "0m", "0h"] {
let output = run_cli(
tmp.path(),
&[
"client",
"--unix-socket",
&socket.display().to_string(),
"wait",
"alpha",
"--timeout",
zero,
"--json",
],
&[],
);
let parsed = envelope(&output);
assert_eq!(
parsed["outcome"],
"owner_not_running",
"'{zero}' must select the unbounded sentinel, stderr={}",
stderr_of(&output)
);
}
}
#[test]
fn cli_surface_defaults_the_timeout_to_the_unbounded_sentinel() {
let tmp = neutral_cwd();
let socket = tmp.path().join("absent.sock");
let path = socket.display().to_string();
let omitted = run_cli(
tmp.path(),
&["client", "--unix-socket", &path, "wait", "alpha", "--json"],
&[],
);
assert_eq!(envelope(&omitted)["outcome"], "owner_not_running");
let help = run_cli(tmp.path(), &["client", "wait", "--help"], &[]);
let text = stdout_of(&help);
assert!(
text.contains("waits indefinitely"),
"the default must be documented as unbounded: {text}"
);
assert!(
text.contains("[default: 0]"),
"the default value must be visible: {text}"
);
}
#[test]
fn cli_surface_accepts_the_documented_duration_spellings() {
let tmp = neutral_cwd();
let socket = tmp.path().join("absent.sock");
for good in ["30", "500ms", "30s", "45m", "2h"] {
let output = run_cli(
tmp.path(),
&[
"client",
"--unix-socket",
&socket.display().to_string(),
"wait",
"alpha",
"--timeout",
good,
"--json",
],
&[],
);
assert_eq!(
envelope(&output)["outcome"],
"owner_not_running",
"'{good}' must parse"
);
}
}
#[test]
fn cli_surface_outside_a_repository_names_the_socket_option() {
let tmp = neutral_cwd();
let output = run_cli(tmp.path(), &["client", "status", "--json"], &[]);
let parsed = envelope(&output);
if parsed["outcome"] == "not_in_repository" {
assert_eq!(output.status.code(), Some(3));
assert!(stderr_of(&output).contains("--unix-socket"));
}
}
#[test]
fn json_usage_errors_emit_one_envelope_for_every_rejected_client_invocation() {
let tmp = neutral_cwd();
for (args, operation, what) in [
(
vec!["client", "mark", "../escape", "--json"],
"control_mark",
"an invalid change ID",
),
(
vec!["client", "unmark", "", "--json"],
"control_unmark",
"an empty change ID",
),
(
vec!["client", "wait", "alpha", "--timeout", "abc", "--json"],
"wait",
"an unparseable timeout",
),
(
vec!["client", "wait", "alpha", "--timeout", "50ms", "--json"],
"wait",
"a below-minimum timeout",
),
(
vec!["client", "mark", "--json"],
"control_mark",
"a missing required argument",
),
(
vec!["client", "status", "--json", "--not-an-option"],
"status",
"an unknown client option",
),
(
vec!["client", "--json"],
"status",
"the namespace with no operation",
),
] {
let output = run_cli(tmp.path(), &args, &[]);
let parsed = envelope(&output);
assert_eq!(parsed["schema_version"], 1, "{what}");
assert_eq!(parsed["ok"], false, "{what}");
assert_eq!(parsed["outcome"], "usage_error", "{what}");
assert_eq!(parsed["operation"], operation, "{what}");
assert!(parsed["detail"].is_object(), "{what}");
assert_eq!(output.status.code(), Some(2), "{what}");
assert!(
stderr_of(&output).contains("usage_error"),
"{what}: stderr={}",
stderr_of(&output)
);
assert!(
std::fs::read_dir(tmp.path())
.expect("temp dir readable")
.next()
.is_none(),
"{what}: a rejected invocation must write nothing"
);
}
}
#[test]
fn json_usage_errors_leave_human_and_non_client_parse_failures_alone() {
let tmp = neutral_cwd();
for args in [
vec!["client", "mark", "../escape"],
vec!["client", "wait", "alpha", "--timeout", "abc"],
vec!["client", "status", "--not-an-option"],
] {
let output = run_cli(tmp.path(), &args, &[]);
assert!(
stdout_of(&output).trim().is_empty(),
"a human parse failure must not print an envelope: {}",
stdout_of(&output)
);
assert!(!output.status.success());
assert!(
stderr_of(&output).contains("error:"),
"Clap's human diagnostic must survive: {}",
stderr_of(&output)
);
}
let unrelated = run_cli(
tmp.path(),
&["openspec", "show", "alpha", "--not-an-option", "--json"],
&[],
);
assert!(
stdout_of(&unrelated).trim().is_empty(),
"an unrelated top-level failure must not be rewritten as a client envelope: {}",
stdout_of(&unrelated)
);
assert!(!unrelated.status.success());
let substring = run_cli(
tmp.path(),
&[
"client",
"--unix-socket",
"/tmp/holds--json-in-its-name.sock",
"mark",
"../escape",
],
&[],
);
assert!(
stdout_of(&substring).trim().is_empty(),
"'--json' inside a value must not select JSON mode: {}",
stdout_of(&substring)
);
assert_eq!(substring.status.code(), Some(2));
}
#[test]
fn json_usage_errors_never_rewrite_help_or_version() {
let tmp = neutral_cwd();
let help = run_cli(tmp.path(), &["client", "status", "--help", "--json"], &[]);
assert!(help.status.success(), "{}", stderr_of(&help));
assert!(
stdout_of(&help).contains("Usage:"),
"help must still be help: {}",
stdout_of(&help)
);
let version = run_cli(tmp.path(), &["--version"], &[]);
assert!(version.status.success());
assert!(!stdout_of(&version).contains("usage_error"));
}
#[tokio::test]
async fn transport_reaches_a_real_owner_over_an_explicit_unix_socket() {
let owner = Owner::start(Some(SpyExecutor::new()), None).await;
owner.publish(snapshot(
"select",
vec![change("alpha", "select", "not queued")],
));
let tmp = neutral_cwd();
let socket = owner.socket();
let output = tokio::task::spawn_blocking({
let cwd = tmp.path().to_path_buf();
move || {
run_cli(
&cwd,
&["client", "--unix-socket", &socket, "status", "--json"],
&[],
)
}
})
.await
.unwrap();
let parsed = envelope(&output);
assert_eq!(parsed["outcome"], "observed");
assert_eq!(parsed["ok"], true);
assert_eq!(output.status.code(), Some(0));
assert_eq!(parsed["detail"]["changes"][0]["id"], "alpha");
owner.stop().await;
}
#[test]
fn transport_reports_owner_not_running_for_an_absent_socket() {
let tmp = neutral_cwd();
let socket = tmp.path().join("nobody.sock");
let output = run_cli(
tmp.path(),
&[
"client",
"--unix-socket",
&socket.display().to_string(),
"status",
"--json",
],
&[],
);
let parsed = envelope(&output);
assert_eq!(parsed["outcome"], "owner_not_running");
assert_eq!(output.status.code(), Some(4));
assert!(!socket.exists(), "a client must never create the socket");
}
#[tokio::test]
async fn transport_presents_the_environment_token_and_never_prints_it() {
let owner = Owner::start(Some(SpyExecutor::new()), Some("s3cret-token")).await;
owner.publish(snapshot(
"select",
vec![change("alpha", "select", "not queued")],
));
let tmp = neutral_cwd();
let socket = owner.socket();
let cwd = tmp.path().to_path_buf();
let authorized = tokio::task::spawn_blocking({
let socket = socket.clone();
let cwd = cwd.clone();
move || {
run_cli(
&cwd,
&[
"client",
"--unix-socket",
&socket,
"--auth-token-env",
"CFLX_CLIENT_TEST_TOKEN",
"status",
"--json",
],
&[("CFLX_CLIENT_TEST_TOKEN", "s3cret-token")],
)
}
})
.await
.unwrap();
assert_eq!(envelope(&authorized)["outcome"], "observed");
assert!(
!stdout_of(&authorized).contains("s3cret-token")
&& !stderr_of(&authorized).contains("s3cret-token"),
"the token must never appear in either stream"
);
let anonymous = tokio::task::spawn_blocking({
let socket = socket.clone();
let cwd = cwd.clone();
move || {
run_cli(
&cwd,
&["client", "--unix-socket", &socket, "status", "--json"],
&[],
)
}
})
.await
.unwrap();
assert_eq!(envelope(&anonymous)["outcome"], "authentication_failed");
assert_eq!(anonymous.status.code(), Some(5));
let wrong = tokio::task::spawn_blocking({
let socket = socket.clone();
move || {
run_cli(
&cwd,
&[
"client",
"--unix-socket",
&socket,
"--auth-token-env",
"CFLX_CLIENT_TEST_TOKEN",
"status",
"--json",
],
&[("CFLX_CLIENT_TEST_TOKEN", "wrong-token")],
)
}
})
.await
.unwrap();
assert_eq!(envelope(&wrong)["outcome"], "authentication_failed");
assert!(!stderr_of(&wrong).contains("wrong-token"));
owner.stop().await;
}
#[test]
fn transport_refuses_an_unset_token_variable_without_connecting() {
let tmp = neutral_cwd();
let socket = tmp.path().join("nobody.sock");
let output = run_cli(
tmp.path(),
&[
"client",
"--unix-socket",
&socket.display().to_string(),
"--auth-token-env",
"CFLX_CLIENT_TEST_TOKEN",
"status",
"--json",
],
&[],
);
assert_eq!(envelope(&output)["outcome"], "authentication_failed");
}
#[tokio::test]
async fn transport_reports_an_incompatible_owner_for_a_non_v2_endpoint() {
let dir = tempfile::tempdir().unwrap();
let socket = dir.path().join("impostor.sock");
let listener = tokio::net::UnixListener::bind(&socket).unwrap();
let app = axum::Router::new().route(
"/api/v2/capabilities",
axum::routing::get(|| async { "not a capabilities document" }),
);
let shutdown = tokio_util::sync::CancellationToken::new();
let task = tokio::spawn({
let shutdown = shutdown.clone();
async move {
let _ = axum::serve(listener, app)
.with_graceful_shutdown(async move { shutdown.cancelled().await })
.await;
}
});
let tmp = neutral_cwd();
let path = socket.display().to_string();
let output = tokio::task::spawn_blocking({
let cwd = tmp.path().to_path_buf();
move || {
run_cli(
&cwd,
&["client", "--unix-socket", &path, "status", "--json"],
&[],
)
}
})
.await
.unwrap();
assert_eq!(envelope(&output)["outcome"], "incompatible_owner");
assert_eq!(output.status.code(), Some(6));
shutdown.cancel();
let _ = task.await;
}
#[tokio::test]
async fn auth_header_validation_refuses_a_malformed_token_before_connecting() {
let dir = tempfile::tempdir().expect("temp dir");
let socket = dir.path().join("counting.sock");
let listener = tokio::net::UnixListener::bind(&socket).expect("binds the counting socket");
let accepted = Arc::new(AtomicUsize::new(0));
let shutdown = tokio_util::sync::CancellationToken::new();
let task = tokio::spawn({
let accepted = accepted.clone();
let stop = shutdown.clone();
async move {
loop {
tokio::select! {
_ = stop.cancelled() => break,
incoming = listener.accept() => {
if incoming.is_ok() {
accepted.fetch_add(1, Ordering::SeqCst);
} else {
break;
}
}
}
}
}
});
let path = socket.display().to_string();
for (token, what) in [
("s3cret\r\nX-Injected: yes", "a CRLF header injection"),
("s3cret\n", "a trailing line feed"),
("s3cret\rmore", "a bare carriage return"),
("s3cret\u{7f}", "DEL"),
("s3\u{1}cret", "another C0 control"),
("s3\tcret", "a horizontal tab"),
] {
let output = tokio::task::spawn_blocking({
let path = path.clone();
let dir = dir.path().to_path_buf();
let token = token.to_string();
move || {
run_cli(
&dir,
&[
"client",
"--unix-socket",
&path,
"--auth-token-env",
"CFLX_CLIENT_TEST_TOKEN",
"status",
"--json",
],
&[("CFLX_CLIENT_TEST_TOKEN", &token)],
)
}
})
.await
.unwrap();
let parsed = envelope(&output);
assert_eq!(parsed["outcome"], "authentication_failed", "{what}");
assert_eq!(output.status.code(), Some(5), "{what}");
let stdout = stdout_of(&output);
let stderr = stderr_of(&output);
for stream in [&stdout, &stderr] {
assert!(
!stream.contains("s3"),
"{what}: no fragment of the token value may be shown: {stream}"
);
assert!(
!stream.contains("X-Injected"),
"{what}: an injection attempt must not be echoed: {stream}"
);
}
assert!(
stderr.contains("CFLX_CLIENT_TEST_TOKEN"),
"{what}: {stderr}"
);
}
assert_eq!(
accepted.load(Ordering::SeqCst),
0,
"a malformed token must be refused before any connection is opened"
);
shutdown.cancel();
let _ = task.await;
}
#[tokio::test]
async fn auth_header_validation_keeps_a_valid_token_authenticating() {
const TOKEN: &str = "tok.en-plus~/+=:_9";
let owner = Owner::start(Some(SpyExecutor::new()), Some(TOKEN)).await;
owner.publish(snapshot(
"select",
vec![change("alpha", "select", "not queued")],
));
let tmp = neutral_cwd();
let socket = owner.socket();
let output = tokio::task::spawn_blocking({
let cwd = tmp.path().to_path_buf();
move || {
run_cli(
&cwd,
&[
"client",
"--unix-socket",
&socket,
"--auth-token-env",
"CFLX_CLIENT_TEST_TOKEN",
"status",
"--json",
],
&[("CFLX_CLIENT_TEST_TOKEN", TOKEN)],
)
}
})
.await
.unwrap();
assert_eq!(envelope(&output)["outcome"], "observed");
assert_eq!(output.status.code(), Some(0));
assert!(
!stdout_of(&output).contains(TOKEN) && !stderr_of(&output).contains(TOKEN),
"a valid token still must never be printed"
);
owner.stop().await;
}
#[tokio::test]
async fn output_contract_keeps_json_stdout_clean_and_diagnostics_on_stderr() {
let owner = Owner::start(Some(SpyExecutor::new()), None).await;
owner.publish(snapshot(
"select",
vec![change("alpha", "select", "not queued")],
));
let tmp = neutral_cwd();
let socket = owner.socket();
let success = tokio::task::spawn_blocking({
let cwd = tmp.path().to_path_buf();
let socket = socket.clone();
move || {
run_cli(
&cwd,
&["client", "--unix-socket", &socket, "status", "--json"],
&[],
)
}
})
.await
.unwrap();
let parsed = envelope(&success);
assert_eq!(parsed["schema_version"], 1);
assert_eq!(parsed["operation"], "status");
assert!(parsed["detail"].is_object());
assert!(
stderr_of(&success).trim().is_empty(),
"a successful run must print no diagnostics: {}",
stderr_of(&success)
);
let failure = tokio::task::spawn_blocking({
let cwd = tmp.path().to_path_buf();
let socket = socket.clone();
move || {
run_cli(
&cwd,
&[
"client",
"--unix-socket",
&socket,
"mark",
"ghost",
"--json",
],
&[],
)
}
})
.await
.unwrap();
let parsed = envelope(&failure);
assert_eq!(parsed["ok"], false);
assert_eq!(parsed["outcome"], "change_not_found");
assert_eq!(parsed["change_id"], "ghost");
assert_eq!(failure.status.code(), Some(9));
assert!(
stderr_of(&failure).contains("change_not_found"),
"the human diagnostic belongs on stderr"
);
owner.stop().await;
}
#[tokio::test]
async fn output_contract_human_mode_is_one_concise_line() {
let owner = Owner::start(Some(SpyExecutor::new()), None).await;
owner.publish(snapshot(
"select",
vec![change("alpha", "select", "not queued")],
));
let tmp = neutral_cwd();
let socket = owner.socket();
let output = tokio::task::spawn_blocking({
let cwd = tmp.path().to_path_buf();
move || run_cli(&cwd, &["client", "--unix-socket", &socket, "status"], &[])
})
.await
.unwrap();
let stdout = stdout_of(&output);
assert_eq!(stdout.lines().count(), 1, "human output must be one line");
assert!(stdout.starts_with("status: observed"), "{stdout}");
assert!(!stdout.contains("schema_version"));
owner.stop().await;
}
#[tokio::test]
async fn output_contract_maps_every_reached_outcome_to_its_documented_exit_code() {
let tmp = neutral_cwd();
let cwd = tmp.path().to_path_buf();
let cases: Vec<(bool, Vec<String>, &str, i32)> = vec![
(
false,
vec!["mark".into(), "alpha".into()],
"owner_not_command_capable",
7,
),
(
true,
vec!["mark".into(), "ghost".into()],
"change_not_found",
9,
),
(
true,
vec![
"wait".into(),
"alpha".into(),
"--timeout".into(),
"300ms".into(),
],
"unsupported_terminal_mode",
16,
),
];
for (bound, action, outcome, code) in cases {
let executor: Option<Arc<dyn RemoteControlExecutor>> = if bound {
Some(SpyExecutor::new())
} else {
None
};
let owner = Owner::start(executor, None).await;
owner.publish(snapshot(
"select",
vec![change("alpha", "select", "not queued")],
));
let socket = owner.socket();
let cwd = cwd.clone();
let output = tokio::task::spawn_blocking(move || {
let mut args = vec!["client".to_string(), "--unix-socket".to_string(), socket];
args.extend(action);
args.push("--json".to_string());
let borrowed: Vec<&str> = args.iter().map(String::as_str).collect();
run_cli(&cwd, &borrowed, &[])
})
.await
.unwrap();
let parsed = envelope(&output);
assert_eq!(parsed["outcome"], outcome);
assert_eq!(output.status.code(), Some(code), "outcome={outcome}");
owner.stop().await;
}
}
#[tokio::test]
async fn owner_contract_publishes_each_terminal_mode_and_omits_inapplicable_fields() {
for (contract, expected_mode, expects_remote, expects_branch) in [
(
OwnerExecutionContract::resolve("main", None, None),
"merged",
false,
false,
),
(
OwnerExecutionContract::resolve("main", None, Some("upstream")),
"base_published",
true,
false,
),
(
OwnerExecutionContract::resolve("main", Some("origin"), None),
"branch_pushed",
true,
true,
),
] {
let owner = Owner::start(Some(SpyExecutor::new()), None).await;
owner.publish(snapshot(
"select",
vec![change("alpha", "select", "not queued")],
));
owner.contract(contract);
let tmp = neutral_cwd();
let socket = owner.socket();
let output = tokio::task::spawn_blocking({
let cwd = tmp.path().to_path_buf();
move || {
run_cli(
&cwd,
&["client", "--unix-socket", &socket, "status", "--json"],
&[],
)
}
})
.await
.unwrap();
let published = envelope(&output)["detail"]["execution_contract"].clone();
assert_eq!(published["terminal_mode"], expected_mode);
assert_eq!(published["base_branch"], "main");
assert_eq!(
published.get("remote").is_some(),
expects_remote,
"mode {expected_mode} remote presence"
);
assert!(
published.get("pushed_branch").is_none(),
"an unscoped read must not invent a branch"
);
let _ = expects_branch;
owner.stop().await;
}
}
#[tokio::test]
async fn owner_contract_joins_the_revision_and_incarnation_of_the_snapshot() {
let owner = Owner::start(Some(SpyExecutor::new()), None).await;
owner.publish(snapshot(
"select",
vec![change("alpha", "select", "not queued")],
));
owner.contract(merged_contract("main"));
let tmp = neutral_cwd();
let socket = owner.socket();
let output = tokio::task::spawn_blocking({
let cwd = tmp.path().to_path_buf();
move || {
run_cli(
&cwd,
&["client", "--unix-socket", &socket, "status", "--json"],
&[],
)
}
})
.await
.unwrap();
let parsed = envelope(&output);
assert_eq!(
parsed["instance_id"].as_str().unwrap(),
owner.projection.instance_id(),
"the envelope must name the incarnation it observed"
);
assert_eq!(
parsed["detail"]["state_revision"].as_u64().unwrap(),
owner.projection.revision(),
"the contract must be joined at the snapshot's revision"
);
owner.stop().await;
}
#[tokio::test]
async fn owner_contract_reports_command_capability_and_a_distinct_unbound_error() {
let unbound = Owner::start(None, None).await;
unbound.publish(snapshot(
"select",
vec![change("alpha", "select", "not queued")],
));
let tmp = neutral_cwd();
let socket = unbound.socket();
let status = tokio::task::spawn_blocking({
let cwd = tmp.path().to_path_buf();
let socket = socket.clone();
move || {
run_cli(
&cwd,
&["client", "--unix-socket", &socket, "status", "--json"],
&[],
)
}
})
.await
.unwrap();
assert_eq!(
envelope(&status)["detail"]["command_execution_available"],
false
);
let marked = tokio::task::spawn_blocking({
let cwd = tmp.path().to_path_buf();
move || {
run_cli(
&cwd,
&[
"client",
"--unix-socket",
&socket,
"mark",
"alpha",
"--json",
],
&[],
)
}
})
.await
.unwrap();
assert_eq!(envelope(&marked)["outcome"], "owner_not_command_capable");
assert_eq!(marked.status.code(), Some(7));
unbound.stop().await;
}
#[test]
fn owner_contract_error_code_is_its_own_wire_token() {
assert_eq!(
ErrorCode::CommandExecutorUnbound.as_str(),
"command_executor_unbound"
);
assert_ne!(
ErrorCode::CommandExecutorUnbound.as_str(),
ErrorCode::LifecycleConflict.as_str()
);
assert!(conflux::web::remote_control_api::dto::ALL_ERROR_CODES
.contains(&ErrorCode::CommandExecutorUnbound));
}
#[tokio::test]
async fn status_is_read_only_and_submits_no_command() {
let spy = SpyExecutor::new();
let owner = Owner::start(Some(spy.clone()), None).await;
owner.publish(snapshot(
"running",
vec![
change("alpha", "running", "applying"),
change("beta", "running", "not queued"),
],
));
owner.boundary.set_running(true);
owner.contract(merged_contract("main"));
let tmp = neutral_cwd();
let socket = owner.socket();
let output = tokio::task::spawn_blocking({
let cwd = tmp.path().to_path_buf();
move || {
run_cli(
&cwd,
&["client", "--unix-socket", &socket, "status", "--json"],
&[],
)
}
})
.await
.unwrap();
let parsed = envelope(&output);
assert_eq!(parsed["outcome"], "observed");
assert_eq!(parsed["detail"]["process"]["app_mode"], "running");
assert_eq!(parsed["detail"]["process"]["scheduler_running"], true);
assert_eq!(parsed["detail"]["changes"].as_array().unwrap().len(), 2);
assert_eq!(spy.call_count(), 0, "status must submit no command");
owner.stop().await;
}
#[tokio::test]
async fn status_reconciles_a_snapshot_that_advances_between_reads() {
let spy = SpyExecutor::new();
let owner = Owner::start(Some(spy.clone()), None).await;
owner.publish(snapshot(
"select",
vec![change("alpha", "select", "not queued")],
));
let tmp = neutral_cwd();
let socket = owner.socket();
let churn_projection = owner.projection.clone();
let stop_churn = Arc::new(AtomicBool::new(false));
let churn_flag = stop_churn.clone();
let churn = tokio::spawn(async move {
let mut n = 0usize;
while !churn_flag.load(Ordering::SeqCst) {
n += 1;
churn_projection.apply_state(
"test_snapshot",
None,
serde_json::json!({}),
snapshot(
"select",
vec![change(&format!("alpha{}", n % 2), "select", "not queued")],
),
);
tokio::time::sleep(std::time::Duration::from_millis(1)).await;
}
});
let output = tokio::task::spawn_blocking({
let cwd = tmp.path().to_path_buf();
move || {
run_cli(
&cwd,
&["client", "--unix-socket", &socket, "status", "--json"],
&[],
)
}
})
.await
.unwrap();
stop_churn.store(true, Ordering::SeqCst);
let _ = churn.await;
let parsed = envelope(&output);
let outcome = parsed["outcome"].as_str().unwrap();
assert!(
outcome == "observed" || outcome == "observation_conflict",
"a racing read must reconcile or report a typed conflict, got {outcome}"
);
if outcome == "observation_conflict" {
assert_eq!(output.status.code(), Some(13));
}
assert_eq!(
spy.call_count(),
0,
"even a conflicted status submits nothing"
);
owner.stop().await;
}
#[tokio::test]
async fn status_reports_a_missing_owner_contract_without_failing() {
let owner = Owner::start(Some(SpyExecutor::new()), None).await;
owner.publish(snapshot(
"select",
vec![change("alpha", "select", "not queued")],
));
let tmp = neutral_cwd();
let socket = owner.socket();
let output = tokio::task::spawn_blocking({
let cwd = tmp.path().to_path_buf();
move || {
run_cli(
&cwd,
&["client", "--unix-socket", &socket, "status", "--json"],
&[],
)
}
})
.await
.unwrap();
let parsed = envelope(&output);
assert_eq!(parsed["outcome"], "observed");
assert!(parsed["detail"]["execution_contract"].is_null());
owner.stop().await;
}
async fn switch_after_first_post(
front: PathBuf,
first: PathBuf,
second: PathBuf,
) -> tokio_util::sync::CancellationToken {
use tokio::io::AsyncReadExt;
use tokio::io::AsyncWriteExt;
let listener = tokio::net::UnixListener::bind(&front).expect("binds the relay socket");
let cancel = tokio_util::sync::CancellationToken::new();
let stop = cancel.clone();
tokio::spawn(async move {
let switched = Arc::new(AtomicBool::new(false));
loop {
let accepted = tokio::select! {
_ = stop.cancelled() => break,
accepted = listener.accept() => accepted,
};
let Ok((mut client, _)) = accepted else { break };
let switched = switched.clone();
let first = first.clone();
let second = second.clone();
tokio::spawn(async move {
let mut head = vec![0u8; 8 * 1024];
let Ok(read) = client.read(&mut head).await else {
return;
};
let head = &head[..read];
let target = if switched.load(Ordering::SeqCst) {
second
} else {
first
};
if head.starts_with(b"POST ") {
switched.store(true, Ordering::SeqCst);
}
let Ok(mut upstream) = tokio::net::UnixStream::connect(&target).await else {
return;
};
if upstream.write_all(head).await.is_err() {
return;
}
let _ = tokio::io::copy_bidirectional(&mut client, &mut upstream).await;
});
}
});
cancel
}
async fn run_client(socket: &str, cwd: &Path, tail: &[&str]) -> Output {
let socket = socket.to_string();
let cwd = cwd.to_path_buf();
let tail: Vec<String> = tail.iter().map(|value| (*value).to_string()).collect();
tokio::task::spawn_blocking(move || {
let mut args = vec!["client".to_string(), "--unix-socket".to_string(), socket];
args.extend(tail);
let borrowed: Vec<&str> = args.iter().map(String::as_str).collect();
run_cli(&cwd, &borrowed, &[])
})
.await
.expect("the client invocation must not panic")
}
async fn control(owner: &Owner, tail: &[&str]) -> Output {
let cwd = neutral_cwd();
let mut args = tail.to_vec();
args.push("--json");
let output = run_client(&owner.socket(), cwd.path(), &args).await;
drop(cwd);
output
}
fn marked(id: &str, app_mode: &str, display_status: &str) -> ChangeResource {
let mut change = change(id, app_mode, display_status);
change.execution_marked = true;
change
}
fn submitted_names(spy: &SpyExecutor) -> Vec<&'static str> {
spy.calls()
.iter()
.map(|command| match command {
CommandSpec::SetExecutionMark { .. } => "set_execution_mark",
CommandSpec::SetAllExecutionMarks {} => "set_all_execution_marks",
CommandSpec::Start => "start",
CommandSpec::Stop => "stop",
CommandSpec::ForceStop => "force_stop",
CommandSpec::SetQueueIntent { .. } => "set_queue_intent",
CommandSpec::RetryChange { .. } => "retry_change",
CommandSpec::RetryErrors { .. } => "retry_errors",
CommandSpec::ForceStopChange { .. } => "force_stop_change",
_ => "other",
})
.collect()
}
#[tokio::test]
async fn control_mark_preserves_unrelated_marks_and_submits_only_mark_commands() {
let spy = SpyExecutor::new();
let owner = Owner::start(Some(spy.clone()), None).await;
owner.publish(snapshot(
"select",
vec![
change("alpha", "select", "not queued"),
marked("beta", "select", "not queued"),
change("gamma", "select", "not queued"),
],
));
let output = control(&owner, &["mark", "alpha", "gamma"]).await;
let parsed = envelope(&output);
assert_eq!(parsed["outcome"], "marked", "{}", stderr_of(&output));
assert_eq!(output.status.code(), Some(0));
assert_eq!(
submitted_names(&spy),
vec!["set_execution_mark", "set_execution_mark"],
"only the two named targets are written, and only as marks"
);
assert_eq!(
spy.calls(),
vec![
CommandSpec::SetExecutionMark {
change_id: "alpha".to_string(),
marked: true
},
CommandSpec::SetExecutionMark {
change_id: "gamma".to_string(),
marked: true
},
],
"nothing addressed beta, and nothing addressed lifecycle or queue state"
);
let targets = parsed["detail"]["targets"].as_array().unwrap();
assert_eq!(targets.len(), 2);
assert_eq!(targets[0]["change_id"], "alpha");
assert_eq!(targets[0]["changed"], true);
assert_eq!(targets[1]["change_id"], "gamma");
owner.stop().await;
}
#[tokio::test]
async fn control_mark_settlement_claims_no_admission() {
let spy = SpyExecutor::new();
let owner = Owner::start(Some(spy.clone()), None).await;
owner.publish(snapshot(
"select",
vec![change("alpha", "select", "not queued")],
));
let output = control(&owner, &["mark", "alpha"]).await;
let parsed = envelope(&output);
assert_eq!(parsed["outcome"], "marked", "{}", stderr_of(&output));
assert_eq!(parsed["operation"], "control_mark");
assert_eq!(parsed["change_id"], "alpha");
assert!(
parsed.as_object().unwrap().get("execution_id").is_none(),
"a mark write creates no execution episode to name: {parsed}"
);
let rendered = parsed.to_string();
for forbidden in ["queue_intent", "admitted", "already_admitted"] {
assert!(!rendered.contains(forbidden), "{forbidden}: {rendered}");
}
assert_eq!(submitted_names(&spy), vec!["set_execution_mark"]);
owner.stop().await;
}
#[tokio::test]
async fn control_mark_of_an_already_marked_target_settles_unchanged_without_a_command() {
let spy = SpyExecutor::new();
let owner = Owner::start(Some(spy.clone()), None).await;
owner.publish(snapshot(
"select",
vec![marked("alpha", "select", "not queued")],
));
let output = control(&owner, &["mark", "alpha"]).await;
let parsed = envelope(&output);
assert_eq!(parsed["outcome"], "unchanged", "{}", stderr_of(&output));
assert_eq!(output.status.code(), Some(0), "unchanged is a success");
assert_eq!(parsed["detail"]["targets"][0]["changed"], false);
assert!(
spy.calls().is_empty(),
"a settled desired state needs no command: {:?}",
spy.calls()
);
owner.stop().await;
}
#[tokio::test]
async fn control_mark_of_a_terminal_target_is_a_reasoned_unchanged_no_op() {
let spy = SpyExecutor::new();
spy.script(vec![Ok(ExecutionSummary::no_op(
"'alpha' is terminal and carries no next-run intent",
))]);
let owner = Owner::start(Some(spy.clone()), None).await;
owner.publish(snapshot(
"select",
vec![change("alpha", "select", "merged")],
));
let output = control(&owner, &["mark", "alpha"]).await;
let parsed = envelope(&output);
assert_eq!(parsed["outcome"], "unchanged", "{}", stderr_of(&output));
assert_eq!(
submitted_names(&spy),
vec!["set_execution_mark"],
"the shared service decides, so the command still goes out"
);
let reason = parsed["detail"]["targets"][0]["reason"].as_str().unwrap();
assert!(reason.contains("next-run intent"), "{reason}");
owner.stop().await;
}
#[tokio::test]
async fn control_unmark_is_target_scoped_and_stops_nothing() {
let spy = SpyExecutor::new();
let owner = Owner::start(Some(spy.clone()), None).await;
owner.publish(snapshot(
"select",
vec![
marked("alpha", "select", "not queued"),
marked("beta", "select", "not queued"),
],
));
let output = control(&owner, &["unmark", "alpha"]).await;
let parsed = envelope(&output);
assert_eq!(parsed["outcome"], "unmarked", "{}", stderr_of(&output));
assert_eq!(parsed["operation"], "control_unmark");
assert_eq!(
spy.calls(),
vec![CommandSpec::SetExecutionMark {
change_id: "alpha".to_string(),
marked: false
}],
"beta's mark is untouched, and nothing stopped or dequeued anything"
);
owner.stop().await;
}
#[tokio::test]
async fn control_refuses_an_unknown_proposal_before_any_command() {
let spy = SpyExecutor::new();
let owner = Owner::start(Some(spy.clone()), None).await;
owner.publish(snapshot(
"select",
vec![change("alpha", "select", "not queued")],
));
let output = control(&owner, &["mark", "alpha", "missing"]).await;
let parsed = envelope(&output);
assert_eq!(parsed["outcome"], "change_not_found");
assert_eq!(output.status.code(), Some(9));
assert_eq!(parsed["change_id"], "missing");
assert!(
spy.calls().is_empty(),
"classification happens before submission: {:?}",
spy.calls()
);
owner.stop().await;
}
#[tokio::test]
async fn control_refuses_a_malformed_target_list_before_contact() {
let spy = SpyExecutor::new();
let owner = Owner::start(Some(spy.clone()), None).await;
owner.publish(snapshot(
"select",
vec![change("alpha", "select", "not queued")],
));
let duplicated = control(&owner, &["mark", "alpha", "alpha"]).await;
let parsed = envelope(&duplicated);
assert_eq!(parsed["outcome"], "usage_error");
assert_eq!(duplicated.status.code(), Some(2));
assert!(
parsed["message"]
.as_str()
.unwrap()
.contains("more than once"),
"{parsed}"
);
let many: Vec<String> = (0..65).map(|n| format!("change-{n}")).collect();
let mut args = vec!["mark".to_string()];
args.extend(many);
let borrowed: Vec<&str> = args.iter().map(String::as_str).collect();
let oversized = control(&owner, &borrowed).await;
assert_eq!(envelope(&oversized)["outcome"], "usage_error");
assert!(spy.calls().is_empty());
owner.stop().await;
}
#[tokio::test]
async fn control_start_submits_only_the_shared_start_intent() {
let spy = SpyExecutor::new();
let owner = Owner::start(Some(spy.clone()), None).await;
owner.publish(snapshot(
"select",
vec![
marked("alpha", "select", "not queued"),
marked("beta", "select", "not queued"),
],
));
let output = control(&owner, &["start"]).await;
let parsed = envelope(&output);
assert_eq!(parsed["outcome"], "accepted", "{}", stderr_of(&output));
assert_eq!(parsed["operation"], "control_start");
assert_eq!(output.status.code(), Some(0));
assert_eq!(
spy.calls(),
vec![CommandSpec::Start],
"exactly the shared intent, with no mark rewriting on the way"
);
owner.stop().await;
}
#[tokio::test]
async fn control_stop_and_force_stop_submit_their_own_shared_intents() {
for (verb, operation, expected) in [
("stop", "control_stop", CommandSpec::Stop),
("force-stop", "control_force_stop", CommandSpec::ForceStop),
] {
let spy = SpyExecutor::new();
let owner = Owner::start(Some(spy.clone()), None).await;
owner.publish(snapshot(
"running",
vec![change("alpha", "running", "applying")],
));
let output = control(&owner, &[verb]).await;
let parsed = envelope(&output);
assert_eq!(
parsed["outcome"],
"accepted",
"{verb}: {}",
stderr_of(&output)
);
assert_eq!(parsed["operation"], operation);
assert_eq!(spy.calls(), vec![expected], "{verb}");
owner.stop().await;
}
}
#[tokio::test]
async fn control_lifecycle_refuses_a_target_list_before_contact() {
let spy = SpyExecutor::new();
let owner = Owner::start(Some(spy.clone()), None).await;
owner.publish(snapshot(
"select",
vec![change("alpha", "select", "not queued")],
));
for verb in ["start", "stop", "force-stop"] {
let output = control(&owner, &[verb, "alpha"]).await;
assert_eq!(envelope(&output)["outcome"], "usage_error", "{verb}");
assert_eq!(output.status.code(), Some(2), "{verb}");
}
assert!(spy.calls().is_empty());
owner.stop().await;
}
#[tokio::test]
async fn control_never_constructs_queue_intent_retry_or_a_bulk_mark() {
let spy = SpyExecutor::new();
let owner = Owner::start(Some(spy.clone()), None).await;
owner.publish(snapshot(
"running",
vec![
change("alpha", "running", "not queued"),
change("beta", "running", "error"),
],
));
let alpha = control(&owner, &["mark", "alpha"]).await;
assert_eq!(
envelope(&alpha)["outcome"],
"marked",
"{}",
stderr_of(&alpha)
);
let beta = control(&owner, &["mark", "beta"]).await;
assert_eq!(envelope(&beta)["outcome"], "marked", "{}", stderr_of(&beta));
let started = control(&owner, &["start"]).await;
assert_eq!(envelope(&started)["outcome"], "accepted");
let names = submitted_names(&spy);
assert_eq!(
names,
vec!["set_execution_mark", "set_execution_mark", "start"],
"the only commands a client may submit are marks and shared lifecycle intents"
);
for forbidden in [
"set_queue_intent",
"retry_change",
"retry_errors",
"set_all_execution_marks",
] {
assert!(
!names.contains(&forbidden),
"{forbidden} must be unreachable"
);
}
owner.stop().await;
}
#[tokio::test]
async fn control_reports_partial_intent_with_an_exact_audit_and_no_rollback() {
let spy = SpyExecutor::new();
spy.script(vec![
Ok(ExecutionSummary::changed("marked alpha")),
Ok(ExecutionSummary::changed("marked beta")),
Err(CommandFailure::new(
ErrorCode::TargetIneligible,
"gamma cannot be marked right now",
)),
]);
let owner = Owner::start(Some(spy.clone()), None).await;
owner.publish(snapshot(
"select",
vec![
change("alpha", "select", "not queued"),
change("beta", "select", "not queued"),
change("gamma", "select", "not queued"),
],
));
let output = control(&owner, &["mark", "alpha", "beta", "gamma"]).await;
let parsed = envelope(&output);
assert_eq!(parsed["outcome"], "partial_intent");
assert_eq!(output.status.code(), Some(15));
assert_eq!(parsed["detail"]["rolled_back"], false);
let commands = parsed["detail"]["commands_submitted"].as_array().unwrap();
assert_eq!(commands.len(), 3, "{parsed}");
for (index, change_id) in ["alpha", "beta", "gamma"].iter().enumerate() {
assert_eq!(commands[index]["command"], "set_execution_mark");
assert_eq!(commands[index]["change_id"], *change_id);
}
let targets = parsed["detail"]["targets"].as_array().unwrap();
assert_eq!(targets.len(), 2);
assert!(parsed["message"].as_str().unwrap().contains("rolled back"));
assert_eq!(spy.call_count(), 3, "{:?}", spy.calls());
owner.stop().await;
}
#[tokio::test]
async fn control_recomputes_after_a_stale_revision_without_repeating_a_settled_effect() {
let spy = SpyExecutor::new();
let api = ApiSpy::new();
let owner = Owner::start_intercepted(Some(spy.clone()), None, api.clone()).await;
owner.publish(snapshot(
"select",
vec![
change("alpha", "select", "not queued"),
change("beta", "select", "not queued"),
],
));
let projection = owner.projection.clone();
api.inject_before_commands(vec![Box::new(move || {
projection.apply_state(
"test_snapshot",
None,
serde_json::json!({}),
snapshot(
"select",
vec![
change("alpha", "select", "not queued"),
change("beta", "select", "not queued"),
change("gamma", "select", "not queued"),
],
),
);
})]);
let output = control(&owner, &["mark", "alpha", "beta"]).await;
let parsed = envelope(&output);
assert_eq!(parsed["outcome"], "marked", "{}", stderr_of(&output));
let commands = parsed["detail"]["commands_submitted"].as_array().unwrap();
assert_eq!(
commands.len(),
spy.call_count(),
"the audit equals the records the endpoint actually made: {parsed}"
);
assert_eq!(
submitted_names(&spy),
vec!["set_execution_mark", "set_execution_mark"],
"a stale submission produced no record, and nothing was submitted twice"
);
let exchanges = api.exchanges();
assert_eq!(
exchanges.len(),
3,
"one refused submission, its recomputation, and beta's: {exchanges:#?}"
);
assert_eq!(exchanges[0].status, 409, "{exchanges:#?}");
assert_eq!(exchanges[0].error_code.as_deref(), Some("stale_revision"));
assert!(
exchanges[0].record_id.is_none(),
"a refused submission creates no record to audit"
);
let recorded: Vec<&str> = exchanges
.iter()
.filter_map(|exchange| exchange.record_id.as_deref())
.collect();
let audited: Vec<&str> = commands
.iter()
.map(|entry| entry["command_id"].as_str().unwrap())
.collect();
assert_eq!(
audited, recorded,
"the audit is exactly the records the owner made, in order"
);
for exchange in &exchanges {
assert_eq!(exchange.command_type, "set_execution_mark");
}
let keys: std::collections::BTreeSet<&str> = exchanges
.iter()
.map(|exchange| exchange.idempotency_key.as_str())
.collect();
assert_eq!(keys.len(), exchanges.len(), "{exchanges:#?}");
assert!(
exchanges[1].expected_revision > exchanges[0].expected_revision,
"the recomputation submits against the revision it just reread: {exchanges:#?}"
);
assert!(
api.requests()
.iter()
.any(|request| request == "GET /api/v2/state"),
"a recomputation rereads the owner's state rather than guessing"
);
owner.stop().await;
}
#[tokio::test]
async fn control_aborts_when_the_socket_starts_serving_a_different_incarnation() {
let first_spy = SpyExecutor::new();
let second_spy = SpyExecutor::new();
let first = Owner::start(Some(first_spy.clone()), None).await;
let second = Owner::start(Some(second_spy.clone()), None).await;
for owner in [&first, &second] {
owner.publish(snapshot(
"select",
vec![
change("alpha", "select", "not queued"),
change("beta", "select", "not queued"),
change("gamma", "select", "not queued"),
],
));
}
let relay_dir = tempfile::tempdir().expect("temp dir");
let relay = relay_dir.path().join("relay.sock");
let cancel = switch_after_first_post(
relay.clone(),
PathBuf::from(first.socket()),
PathBuf::from(second.socket()),
)
.await;
let cwd = neutral_cwd();
let output = run_client(
relay.to_str().unwrap(),
cwd.path(),
&["mark", "alpha", "beta", "gamma", "--json"],
)
.await;
let parsed = envelope(&output);
assert_eq!(
parsed["outcome"],
"owner_restarted",
"{}",
stderr_of(&output)
);
assert_eq!(output.status.code(), Some(8));
assert_ne!(
parsed["detail"]["expected_instance_id"],
parsed["detail"]["observed_instance_id"]
);
assert_eq!(
parsed["detail"]["stopped_at"], "beta",
"the request stops at the target whose record revealed the switch"
);
assert_eq!(
first_spy.call_count(),
1,
"only the first mark reached the incarnation the client observed: {:?}",
first_spy.calls()
);
assert_eq!(
second_spy.call_count(),
1,
"exactly the in-flight command reached the replacement, and gamma was never \
submitted at all: {:?}",
second_spy.calls()
);
cancel.cancel();
first.stop().await;
second.stop().await;
}
#[tokio::test]
async fn control_reports_owner_not_command_capable_without_touching_the_workspace() {
let owner = Owner::start(None, None).await;
owner.publish(snapshot(
"select",
vec![change("alpha", "select", "not queued")],
));
let workspace = neutral_cwd();
let before: Vec<_> = std::fs::read_dir(workspace.path())
.unwrap()
.map(|entry| entry.unwrap().file_name())
.collect();
for tail in [vec!["mark", "alpha"], vec!["start"]] {
let output = run_client(
&owner.socket(),
workspace.path(),
&[tail.clone(), vec!["--json"]].concat(),
)
.await;
assert_eq!(
envelope(&output)["outcome"],
"owner_not_command_capable",
"{tail:?}"
);
assert_eq!(output.status.code(), Some(7), "{tail:?}");
}
let after: Vec<_> = std::fs::read_dir(workspace.path())
.unwrap()
.map(|entry| entry.unwrap().file_name())
.collect();
assert_eq!(
before, after,
"a refused control command must not touch the workspace"
);
owner.stop().await;
}
struct Fixture {
dir: tempfile::TempDir,
}
impl Fixture {
fn new() -> Self {
let dir = tempfile::tempdir().expect("temp dir");
let fixture = Self { dir };
fixture.git(&["init", "--initial-branch=main"]);
fixture.git(&["config", "user.email", "test@example.com"]);
fixture.git(&["config", "user.name", "Test"]);
fixture.git(&["config", "commit.gpgsign", "false"]);
fixture.write("README.md", "fixture\n");
fixture.git(&["add", "-A"]);
fixture.git(&["commit", "-m", "init"]);
fixture
}
fn path(&self) -> &Path {
self.dir.path()
}
fn git(&self, args: &[&str]) -> String {
let output = std::process::Command::new("git")
.args(args)
.current_dir(self.dir.path())
.output()
.expect("git must be available");
assert!(
output.status.success(),
"git {args:?} failed: {}",
String::from_utf8_lossy(&output.stderr)
);
String::from_utf8_lossy(&output.stdout).trim().to_string()
}
fn write(&self, relative: &str, contents: &str) {
let path = self.dir.path().join(relative);
std::fs::create_dir_all(path.parent().unwrap()).unwrap();
std::fs::write(path, contents).unwrap();
}
fn stage_active(&self, change_id: &str) {
self.write(
&format!("openspec/changes/{change_id}/proposal.md"),
"# Proposal\n",
);
self.git(&["add", "-A"]);
self.git(&["commit", "-m", "add change"]);
}
fn archive(&self, change_id: &str) {
archive_in(self.dir.path(), change_id);
}
fn contradict(&self, change_id: &str) {
self.write(
&format!("openspec/changes/{change_id}/proposal.md"),
"# Proposal\n",
);
self.write(
&format!("openspec/changes/archive/2026-01-01-{change_id}/proposal.md"),
"# Proposal\n",
);
self.git(&["add", "-A"]);
self.git(&["commit", "-m", "contradictory"]);
}
}
fn archive_in(repo: &Path, change_id: &str) {
let _ = std::fs::remove_dir_all(repo.join(format!("openspec/changes/{change_id}")));
let archived = repo.join(format!(
"openspec/changes/archive/2026-01-01-{change_id}/proposal.md"
));
std::fs::create_dir_all(archived.parent().unwrap()).unwrap();
std::fs::write(archived, "# Proposal\n").unwrap();
for args in [
["add", "-A"].as_slice(),
["commit", "-m", "archive change"].as_slice(),
] {
let output = std::process::Command::new("git")
.args(args)
.current_dir(repo)
.output()
.expect("git must be available");
assert!(
output.status.success(),
"git {args:?} failed: {}",
String::from_utf8_lossy(&output.stderr)
);
}
}
async fn wait_for(owner: &Owner, repo: &Path, change_id: &str, timeout: &str) -> Output {
let socket = owner.socket();
let repo = repo.to_path_buf();
let change_id = change_id.to_string();
let timeout = timeout.to_string();
tokio::task::spawn_blocking(move || {
run_cli(
&repo,
&[
"client",
"--unix-socket",
&socket,
"wait",
&change_id,
"--timeout",
&timeout,
"--json",
],
&[],
)
})
.await
.unwrap()
}
fn spawn_wait(
owner: &Owner,
repo: &Path,
change_id: &str,
timeout: Option<&str>,
) -> tokio::task::JoinHandle<Output> {
let socket = owner.socket();
let repo = repo.to_path_buf();
let change_id = change_id.to_string();
let timeout = timeout.map(str::to_string);
tokio::task::spawn_blocking(move || {
let mut args = vec![
"client".to_string(),
"--unix-socket".to_string(),
socket,
"wait".to_string(),
change_id,
];
if let Some(timeout) = timeout {
args.push("--timeout".to_string());
args.push(timeout);
}
args.push("--json".to_string());
let argv: Vec<&str> = args.iter().map(String::as_str).collect();
run_cli(&repo, &argv, &[])
})
}
#[tokio::test]
async fn wait_proves_local_integration_from_repository_evidence() {
let repo = Fixture::new();
repo.stage_active("alpha");
repo.archive("alpha");
let spy = SpyExecutor::new();
let owner = Owner::start(Some(spy.clone()), None).await;
owner.publish(snapshot(
"select",
vec![change("alpha", "select", "merged")],
));
owner.contract(merged_contract("main"));
let output = wait_for(&owner, repo.path(), "alpha", "30s").await;
let parsed = envelope(&output);
assert_eq!(parsed["outcome"], "completed");
assert_eq!(output.status.code(), Some(0));
assert_eq!(parsed["detail"]["terminal_mode"], "merged");
assert_eq!(parsed["detail"]["commands_submitted"], 0);
assert_eq!(spy.call_count(), 0, "wait must submit no command");
owner.stop().await;
}
#[tokio::test]
async fn wait_proves_branch_publication_without_claiming_base_integration() {
let repo = Fixture::new();
repo.stage_active("alpha");
repo.git(&["checkout", "-b", "alpha"]);
repo.archive("alpha");
repo.git(&["checkout", "main"]);
let remote = tempfile::tempdir().unwrap();
let remote_path = remote.path().join("origin.git");
std::process::Command::new("git")
.args(["init", "--bare", remote_path.to_str().unwrap()])
.output()
.unwrap();
repo.git(&["remote", "add", "origin", remote_path.to_str().unwrap()]);
repo.git(&["push", "origin", "alpha"]);
let spy = SpyExecutor::new();
let owner = Owner::start(Some(spy.clone()), None).await;
owner.publish(snapshot(
"select",
vec![change("alpha", "select", "pushed")],
));
owner.contract(OwnerExecutionContract::resolve(
"main",
Some("origin"),
None,
));
let output = wait_for(&owner, repo.path(), "alpha", "30s").await;
let parsed = envelope(&output);
assert_eq!(parsed["outcome"], "completed");
assert_eq!(parsed["detail"]["terminal_mode"], "branch_pushed");
assert_eq!(parsed["detail"]["pushed_branch"], "alpha");
let evidence = parsed["detail"]["evidence"].as_str().unwrap();
assert!(
evidence.contains("not base integration"),
"publication must not be reported as base integration: {evidence}"
);
assert_eq!(spy.call_count(), 0);
owner.stop().await;
}
#[tokio::test]
async fn wait_proves_base_publication_only_when_the_remote_matches() {
let repo = Fixture::new();
repo.stage_active("alpha");
repo.archive("alpha");
let remote = tempfile::tempdir().unwrap();
let remote_path = remote.path().join("upstream.git");
std::process::Command::new("git")
.args(["init", "--bare", remote_path.to_str().unwrap()])
.output()
.unwrap();
repo.git(&["remote", "add", "upstream", remote_path.to_str().unwrap()]);
let owner = Owner::start(Some(SpyExecutor::new()), None).await;
owner.publish(snapshot(
"select",
vec![change("alpha", "select", "merged")],
));
owner.contract(OwnerExecutionContract::resolve(
"main",
None,
Some("upstream"),
));
let unpublished = wait_for(&owner, repo.path(), "alpha", "300ms").await;
let parsed = envelope(&unpublished);
assert_eq!(parsed["outcome"], "timeout");
assert_eq!(unpublished.status.code(), Some(19));
repo.git(&["push", "upstream", "main"]);
let published = wait_for(&owner, repo.path(), "alpha", "30s").await;
let parsed = envelope(&published);
assert_eq!(parsed["outcome"], "completed");
assert_eq!(parsed["detail"]["terminal_mode"], "base_published");
owner.stop().await;
}
#[tokio::test]
async fn wait_does_not_treat_disappearance_as_completion() {
let repo = Fixture::new();
repo.stage_active("alpha");
let api = ApiSpy::new();
let spy = SpyExecutor::new();
let owner = Owner::start_intercepted(Some(spy.clone()), None, api.clone()).await;
owner.publish(snapshot(
"select",
vec![change("alpha", "select", "applying")],
));
owner.contract(merged_contract("main"));
let projection = owner.projection.clone();
let repo_path = repo.path().to_path_buf();
api.inject_before_state_reads(vec![
Box::new(|| {}),
Box::new(move || {
projection.apply_state(
"test_snapshot",
None,
serde_json::json!({}),
snapshot("select", vec![change("beta", "select", "not queued")]),
);
}),
Box::new(move || archive_in(&repo_path, "alpha")),
]);
let output = spawn_wait(&owner, repo.path(), "alpha", None)
.await
.unwrap();
let parsed = envelope(&output);
assert_eq!(
parsed["outcome"],
"completed",
"only repository evidence may end this wait, stderr={}",
stderr_of(&output)
);
let state_reads = api
.requests()
.iter()
.filter(|request| *request == "GET /api/v2/state")
.count();
assert!(
state_reads >= 3,
"the wait must have observed the disappearance and kept going rather than \
concluding at it: {state_reads} state reads"
);
assert_eq!(spy.call_count(), 0);
owner.stop().await;
}
#[tokio::test]
async fn wait_refuses_unknown_change_rather_than_waiting_for_a_row_that_cannot_appear() {
let repo = Fixture::new();
repo.stage_active("alpha");
let head_before = repo.git(&["rev-parse", "HEAD"]);
let spy = SpyExecutor::new();
let owner = Owner::start(Some(spy.clone()), None).await;
owner.publish(snapshot(
"select",
vec![change("alpha", "select", "not queued")],
));
owner.contract(merged_contract("main"));
let output = spawn_wait(&owner, repo.path(), "aaaa", None).await.unwrap();
let parsed = envelope(&output);
assert_eq!(
parsed["outcome"],
"change_not_found",
"stderr={}",
stderr_of(&output)
);
assert!(!parsed["ok"].as_bool().unwrap());
assert_eq!(output.status.code(), Some(9));
assert_eq!(parsed["change_id"], "aaaa");
assert!(
parsed["instance_id"]
.as_str()
.is_some_and(|instance| !instance.is_empty()),
"the refusal must name the owner it observed: {parsed}"
);
assert_eq!(parsed["detail"]["commands_submitted"], 0);
assert_eq!(spy.call_count(), 0, "wait must submit no command");
assert_eq!(repo.git(&["rev-parse", "HEAD"]), head_before);
assert_eq!(
repo.git(&["status", "--porcelain"]),
"",
"wait must leave the worktree clean"
);
owner.stop().await;
}
#[tokio::test]
async fn wait_refuses_unknown_change_before_its_deadline_rather_than_at_it() {
let repo = Fixture::new();
repo.stage_active("alpha");
let spy = SpyExecutor::new();
let owner = Owner::start(Some(spy.clone()), None).await;
owner.publish(snapshot(
"select",
vec![change("alpha", "select", "not queued")],
));
owner.contract(merged_contract("main"));
let output = wait_for(&owner, repo.path(), "aaaa", "5s").await;
let parsed = envelope(&output);
assert_eq!(
parsed["outcome"],
"change_not_found",
"an unknown target must be refused, not waited out, stderr={}",
stderr_of(&output)
);
assert_eq!(output.status.code(), Some(9));
assert_eq!(parsed["detail"]["commands_submitted"], 0);
assert_eq!(spy.call_count(), 0);
owner.stop().await;
}
#[tokio::test]
async fn wait_refuses_unknown_change_only_when_the_repository_cannot_certify_it() {
let repo = Fixture::new();
repo.archive("alpha");
let spy = SpyExecutor::new();
let owner = Owner::start(Some(spy.clone()), None).await;
owner.publish(snapshot(
"select",
vec![change("beta", "select", "not queued")],
));
owner.contract(merged_contract("main"));
let output = spawn_wait(&owner, repo.path(), "alpha", None)
.await
.unwrap();
let parsed = envelope(&output);
assert_eq!(
parsed["outcome"],
"completed",
"certified completion must outrank an unknown-target refusal, stderr={}",
stderr_of(&output)
);
assert_eq!(output.status.code(), Some(0));
assert_eq!(parsed["detail"]["terminal_mode"], "merged");
assert_eq!(parsed["detail"]["commands_submitted"], 0);
assert_eq!(spy.call_count(), 0);
owner.stop().await;
}
#[tokio::test]
async fn wait_refuses_unknown_change_without_refusing_a_known_unqueued_row() {
let repo = Fixture::new();
repo.stage_active("alpha");
let spy = SpyExecutor::new();
let owner = Owner::start(Some(spy.clone()), None).await;
owner.publish(snapshot(
"select",
vec![change("alpha", "select", "not queued")],
));
owner.contract(merged_contract("main"));
let output = wait_for(&owner, repo.path(), "alpha", "300ms").await;
let parsed = envelope(&output);
assert_eq!(
parsed["outcome"],
"timeout",
"a tracked idle row must keep observing, stderr={}",
stderr_of(&output)
);
assert_eq!(output.status.code(), Some(19));
assert_eq!(spy.call_count(), 0);
owner.stop().await;
}
#[tokio::test]
async fn wait_reports_contradictory_repository_evidence_as_its_own_outcome() {
let repo = Fixture::new();
repo.contradict("alpha");
let owner = Owner::start(Some(SpyExecutor::new()), None).await;
owner.publish(snapshot(
"select",
vec![change("alpha", "select", "merged")],
));
owner.contract(merged_contract("main"));
let output = wait_for(&owner, repo.path(), "alpha", "30s").await;
let parsed = envelope(&output);
assert_eq!(parsed["outcome"], "evidence_error");
assert_eq!(output.status.code(), Some(22));
owner.stop().await;
}
#[tokio::test]
async fn wait_reports_a_missing_base_branch_as_an_evidence_error() {
let repo = Fixture::new();
repo.stage_active("alpha");
repo.archive("alpha");
let owner = Owner::start(Some(SpyExecutor::new()), None).await;
owner.publish(snapshot(
"select",
vec![change("alpha", "select", "merged")],
));
owner.contract(merged_contract("a-branch-that-does-not-exist"));
let output = wait_for(&owner, repo.path(), "alpha", "30s").await;
assert_eq!(envelope(&output)["outcome"], "evidence_error");
owner.stop().await;
}
#[tokio::test]
async fn wait_reports_rejection_without_repairing_it() {
let repo = Fixture::new();
repo.stage_active("alpha");
let spy = SpyExecutor::new();
let owner = Owner::start(Some(spy.clone()), None).await;
let mut rejected = change("alpha", "select", "rejected");
rejected.error_detail = Some("the proposal was rejected in review".to_string());
owner.publish(snapshot("select", vec![rejected]));
owner.contract(merged_contract("main"));
let output = wait_for(&owner, repo.path(), "alpha", "30s").await;
let parsed = envelope(&output);
assert_eq!(parsed["outcome"], "change_rejected");
assert_eq!(output.status.code(), Some(17));
assert!(parsed["message"]
.as_str()
.unwrap()
.contains("rejected in review"));
assert_eq!(spy.call_count(), 0, "wait must never retry a rejection");
owner.stop().await;
}
#[tokio::test]
async fn wait_releases_immediately_on_every_status_that_needs_an_operator() {
let repo = Fixture::new();
repo.stage_active("alpha");
let head_before = repo.git(&["rev-parse", "HEAD"]);
let spy = SpyExecutor::new();
let owner = Owner::start(Some(spy.clone()), None).await;
owner.contract(merged_contract("main"));
for (status, error_detail) in [
("error", Some("the apply phase exited non-zero")),
("merge wait", Some("conflicting hunk in src/lib.rs")),
("stopped", None),
("stalled", None),
] {
let mut row = change("alpha", "select", status);
row.error_detail = error_detail.map(str::to_string);
owner.publish(snapshot("select", vec![row]));
let output = spawn_wait(&owner, repo.path(), "alpha", None)
.await
.unwrap();
let parsed = envelope(&output);
assert_eq!(
parsed["outcome"],
"change_requires_action",
"'{status}' must release the caller, stderr={}",
stderr_of(&output)
);
assert_eq!(output.status.code(), Some(27), "status '{status}'");
assert!(!parsed["ok"].as_bool().unwrap(), "status '{status}'");
assert_eq!(parsed["detail"]["observed_status"], status);
assert_eq!(parsed["detail"]["commands_submitted"], 0);
match error_detail {
Some(detail) => assert_eq!(parsed["detail"]["error_detail"], detail),
None => assert!(
parsed["detail"].get("error_detail").is_none(),
"'{status}' published no detail, so none must be invented"
),
}
}
assert_eq!(spy.call_count(), 0, "wait must submit no command");
assert_eq!(
repo.git(&["rev-parse", "HEAD"]),
head_before,
"wait must not touch the repository on its way out"
);
owner.stop().await;
}
#[tokio::test]
async fn wait_releases_when_a_live_change_transitions_into_a_manual_action_row() {
let repo = Fixture::new();
repo.stage_active("alpha");
let spy = SpyExecutor::new();
let owner = Owner::start(Some(spy.clone()), None).await;
owner.publish(snapshot(
"select",
vec![change("alpha", "select", "applying")],
));
owner.contract(merged_contract("main"));
let waiting = spawn_wait(&owner, repo.path(), "alpha", None);
tokio::time::sleep(Duration::from_millis(300)).await;
assert!(
!waiting.is_finished(),
"an active phase must keep the wait observing"
);
let mut failed = change("alpha", "select", "error");
failed.error_detail = Some("acceptance never passed".to_string());
owner.publish(snapshot("select", vec![failed]));
let output = waiting.await.unwrap();
let parsed = envelope(&output);
assert_eq!(
parsed["outcome"],
"change_requires_action",
"stderr={}",
stderr_of(&output)
);
assert_eq!(output.status.code(), Some(27));
assert_eq!(parsed["detail"]["observed_status"], "error");
assert_eq!(parsed["detail"]["error_detail"], "acceptance never passed");
assert_eq!(spy.call_count(), 0);
owner.stop().await;
}
#[tokio::test]
async fn wait_releases_external_blocker_without_waiting_for_operator() {
let repo = Fixture::new();
repo.stage_active("alpha");
let head_before = repo.git(&["rev-parse", "HEAD"]);
let spy = SpyExecutor::new();
let owner = Owner::start(Some(spy.clone()), None).await;
owner.publish(snapshot(
"select",
vec![blocked_change("alpha", "select", BlockerKind::External)],
));
owner.contract(merged_contract("main"));
let output = spawn_wait(&owner, repo.path(), "alpha", None)
.await
.unwrap();
let parsed = envelope(&output);
assert_eq!(
parsed["outcome"],
"change_requires_action",
"an external blocker must release an unbounded waiter, stderr={}",
stderr_of(&output)
);
assert_eq!(output.status.code(), Some(27));
assert!(!parsed["ok"].as_bool().unwrap());
assert_eq!(parsed["detail"]["observed_status"], "blocked");
assert_eq!(parsed["detail"]["commands_submitted"], 0);
assert_eq!(parsed["detail"]["blocker"]["kind"], "external");
assert_eq!(
parsed["detail"]["blocker"]["unblock_condition"],
"the certificate is issued"
);
assert_eq!(parsed["detail"]["blocker"]["prerequisite_owner"], "release");
assert_eq!(
parsed["detail"]["error_detail"],
"waiting on the signing certificate"
);
assert_eq!(spy.call_count(), 0, "wait must submit no command");
assert_eq!(
repo.git(&["rev-parse", "HEAD"]),
head_before,
"wait must not touch the repository on its way out"
);
owner.stop().await;
}
#[tokio::test]
async fn wait_releases_when_live_work_becomes_externally_blocked() {
let repo = Fixture::new();
repo.stage_active("alpha");
let spy = SpyExecutor::new();
let owner = Owner::start(Some(spy.clone()), None).await;
owner.publish(snapshot(
"select",
vec![change("alpha", "select", "applying")],
));
owner.contract(merged_contract("main"));
let waiting = spawn_wait(&owner, repo.path(), "alpha", None);
tokio::time::sleep(Duration::from_millis(300)).await;
assert!(
!waiting.is_finished(),
"an active phase must keep the wait observing"
);
owner.publish(snapshot(
"select",
vec![blocked_change("alpha", "select", BlockerKind::External)],
));
let output = waiting.await.unwrap();
let parsed = envelope(&output);
assert_eq!(
parsed["outcome"],
"change_requires_action",
"stderr={}",
stderr_of(&output)
);
assert_eq!(output.status.code(), Some(27));
assert_eq!(parsed["detail"]["observed_status"], "blocked");
assert_eq!(parsed["detail"]["blocker"]["kind"], "external");
assert_eq!(parsed["detail"]["commands_submitted"], 0);
assert_eq!(spy.call_count(), 0);
owner.stop().await;
}
#[tokio::test]
async fn wait_keeps_observing_every_status_the_owner_can_still_advance() {
let repo = Fixture::new();
repo.stage_active("held-not-queued");
let spy = SpyExecutor::new();
let owner = Owner::start(Some(spy.clone()), None).await;
owner.publish(snapshot(
"select",
vec![
change("held-not-queued", "select", "not queued"),
change("held-blocked", "select", "blocked"),
blocked_change("held-dependency", "select", BlockerKind::Dependency),
change("held-applying", "select", "applying"),
],
));
owner.contract(merged_contract("main"));
let (not_queued, blocked, dependency, applying) = tokio::join!(
wait_for(&owner, repo.path(), "held-not-queued", "250ms"),
wait_for(&owner, repo.path(), "held-blocked", "250ms"),
wait_for(&owner, repo.path(), "held-dependency", "250ms"),
wait_for(&owner, repo.path(), "held-applying", "250ms"),
);
for (output, change_id) in [
(¬_queued, "held-not-queued"),
(&blocked, "held-blocked"),
(&dependency, "held-dependency"),
(&applying, "held-applying"),
] {
let parsed = envelope(output);
assert_eq!(
parsed["outcome"],
"timeout",
"'{change_id}' must keep observing until the caller's deadline, stderr={}",
stderr_of(output)
);
assert_eq!(output.status.code(), Some(19), "{change_id}");
}
assert_eq!(spy.call_count(), 0);
owner.stop().await;
}
#[tokio::test]
async fn wait_releases_a_settled_merged_row_whose_evidence_never_arrives() {
let repo = Fixture::new();
repo.stage_active("alpha");
let spy = SpyExecutor::new();
let owner = Owner::start(Some(spy.clone()), None).await;
owner.publish(snapshot(
"select",
vec![change("alpha", "select", "merged")],
));
owner.contract(merged_contract("main"));
let output = spawn_wait(&owner, repo.path(), "alpha", None)
.await
.unwrap();
let parsed = envelope(&output);
assert_eq!(
parsed["outcome"],
"change_requires_action",
"an unbounded wait must not hold on an uncertifiable settled row, stderr={}",
stderr_of(&output)
);
assert_eq!(output.status.code(), Some(27));
assert_eq!(parsed["detail"]["observed_status"], "merged");
assert_eq!(parsed["detail"]["commands_submitted"], 0);
assert!(
!parsed["detail"]["error_detail"]
.as_str()
.unwrap_or_default()
.is_empty(),
"the oracle's own reason must survive into the release: {parsed}"
);
assert_eq!(spy.call_count(), 0);
owner.stop().await;
}
#[tokio::test]
async fn wait_certifies_a_merged_row_whose_evidence_lands_between_two_observations() {
let repo = Fixture::new();
repo.stage_active("alpha");
let api = ApiSpy::new();
let spy = SpyExecutor::new();
let owner = Owner::start_intercepted(Some(spy.clone()), None, api.clone()).await;
owner.publish(snapshot(
"select",
vec![change("alpha", "select", "merged")],
));
owner.contract(merged_contract("main"));
let repo_path = repo.path().to_path_buf();
api.inject_before_state_reads(vec![
Box::new(|| {}),
Box::new(move || archive_in(&repo_path, "alpha")),
]);
let output = spawn_wait(&owner, repo.path(), "alpha", None)
.await
.unwrap();
let parsed = envelope(&output);
assert_eq!(
parsed["outcome"],
"completed",
"late evidence must still certify, stderr={}",
stderr_of(&output)
);
assert_eq!(output.status.code(), Some(0));
assert_eq!(parsed["detail"]["commands_submitted"], 0);
assert_eq!(spy.call_count(), 0);
owner.stop().await;
}
#[tokio::test]
async fn wait_reports_a_fatal_process_error() {
let repo = Fixture::new();
repo.stage_active("alpha");
let spy = SpyExecutor::new();
let owner = Owner::start(Some(spy.clone()), None).await;
let mut failed = snapshot("select", vec![change("alpha", "select", "applying")]);
failed.process_error = Some("the orchestration run died".to_string());
owner.publish(failed);
owner.contract(merged_contract("main"));
let output = wait_for(&owner, repo.path(), "alpha", "30s").await;
let parsed = envelope(&output);
assert_eq!(parsed["outcome"], "process_failed");
assert_eq!(output.status.code(), Some(18));
assert_eq!(spy.call_count(), 0);
owner.stop().await;
}
#[tokio::test]
async fn wait_refuses_to_wait_on_an_owner_that_published_no_terminal_mode() {
let repo = Fixture::new();
repo.stage_active("alpha");
let owner = Owner::start(Some(SpyExecutor::new()), None).await;
owner.publish(snapshot(
"select",
vec![change("alpha", "select", "applying")],
));
let output = wait_for(&owner, repo.path(), "alpha", "30s").await;
let parsed = envelope(&output);
assert_eq!(parsed["outcome"], "unsupported_terminal_mode");
assert_eq!(output.status.code(), Some(16));
let message = parsed["message"].as_str().unwrap();
assert!(message.contains("could only end in a timeout"), "{message}");
owner.stop().await;
}
#[tokio::test]
async fn wait_times_out_without_mutating_anything() {
let repo = Fixture::new();
repo.stage_active("alpha");
let head_before = repo.git(&["rev-parse", "HEAD"]);
let spy = SpyExecutor::new();
let owner = Owner::start(Some(spy.clone()), None).await;
owner.publish(snapshot(
"select",
vec![change("alpha", "select", "applying")],
));
owner.contract(merged_contract("main"));
let output = wait_for(&owner, repo.path(), "alpha", "300ms").await;
let parsed = envelope(&output);
assert_eq!(parsed["outcome"], "timeout");
assert_eq!(output.status.code(), Some(19));
assert_eq!(parsed["detail"]["commands_submitted"], 0);
assert_eq!(spy.call_count(), 0);
assert_eq!(repo.git(&["rev-parse", "HEAD"]), head_before);
assert_eq!(
repo.git(&["status", "--porcelain"]),
"",
"wait must leave the worktree clean"
);
owner.stop().await;
}
#[tokio::test]
async fn wait_timeout_reports_the_latest_target_observation_and_nothing_else() {
let repo = Fixture::new();
repo.stage_active("alpha");
let spy = SpyExecutor::new();
let owner = Owner::start(Some(spy.clone()), None).await;
let execution_id = owner.admit("alpha");
owner.publish(snapshot(
"select",
vec![
change("alpha", "select", "applying"),
change("beta", "select", "archiving"),
],
));
owner.contract(merged_contract("main"));
let output = wait_for(&owner, repo.path(), "alpha", "300ms").await;
let parsed = envelope(&output);
assert_eq!(parsed["outcome"], "timeout");
assert_eq!(output.status.code(), Some(19));
let detail = &parsed["detail"];
assert_eq!(detail["commands_submitted"], 0);
assert_eq!(detail["timeout_ms"], 300);
assert_eq!(detail["timeout_stage"], "observing_owner");
let elapsed = detail["wait_elapsed_ms"]
.as_u64()
.expect("the measured wait must be a number");
assert!(
elapsed >= 300,
"a wait that reached its deadline cannot have taken less than it: {elapsed}"
);
let last = &detail["last_observation"];
assert_eq!(last["change"]["id"], "alpha");
assert_eq!(last["change"]["display_status"], "applying");
assert_eq!(last["execution"]["id"], "alpha");
assert_eq!(last["execution"]["execution_id"], execution_id);
assert!(
last["state_revision"].is_u64() && last["event_sequence"].is_u64(),
"the observation must identify the revision it was reconciled at: {last}"
);
assert!(
!last["observed_at"].as_str().unwrap_or_default().is_empty(),
"the observation must carry the owner's own observation instant: {last}"
);
let rendered = last.to_string();
assert!(
!rendered.contains("beta"),
"a timeout about 'alpha' must not publish another change: {rendered}"
);
assert_eq!(spy.call_count(), 0, "a timeout must submit no command");
owner.stop().await;
}
#[tokio::test]
async fn wait_timeout_reports_the_newest_observation_rather_than_the_first() {
let repo = Fixture::new();
repo.stage_active("alpha");
let api = ApiSpy::new();
let spy = SpyExecutor::new();
let owner = Owner::start_intercepted(Some(spy.clone()), None, api.clone()).await;
owner.publish(snapshot(
"select",
vec![change("alpha", "select", "applying")],
));
owner.contract(merged_contract("main"));
let projection = owner.projection.clone();
api.inject_before_state_reads(vec![
Box::new(|| {}),
Box::new(move || {
projection.apply_state(
"test_snapshot",
None,
serde_json::json!({}),
snapshot("select", vec![change("alpha", "select", "accepting")]),
);
}),
]);
let output = wait_for(&owner, repo.path(), "alpha", "900ms").await;
let parsed = envelope(&output);
assert_eq!(parsed["outcome"], "timeout");
assert_eq!(
parsed["detail"]["last_observation"]["change"]["display_status"], "accepting",
"the timeout must report the last thing the wait saw, not the first"
);
assert_eq!(parsed["detail"]["timeout_stage"], "observing_owner");
assert_eq!(spy.call_count(), 0);
owner.stop().await;
}
#[tokio::test]
async fn wait_timeout_inside_local_certification_reports_the_repository_stage() {
let repo = Fixture::new();
repo.stage_active("alpha");
repo.archive("alpha");
let head_before = repo.git(&["rev-parse", "HEAD"]);
let real_git = String::from_utf8_lossy(
&std::process::Command::new("sh")
.args(["-c", "command -v git"])
.output()
.expect("git must be on PATH")
.stdout,
)
.trim()
.to_string();
let shim_dir = tempfile::tempdir().expect("temp dir");
let shim = shim_dir.path().join("git");
std::fs::write(
&shim,
format!(
"#!/bin/sh\nfor arg in \"$@\"; do\n if [ \"$arg\" = ls-tree ]; then\n exec \
sleep 20\n fi\ndone\nexec {real_git} \"$@\"\n"
),
)
.expect("the shim is writable");
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt;
std::fs::set_permissions(&shim, std::fs::Permissions::from_mode(0o755))
.expect("the shim is executable");
}
let path = format!(
"{}:{}",
shim_dir.path().display(),
std::env::var("PATH").unwrap_or_default()
);
let spy = SpyExecutor::new();
let owner = Owner::start(Some(spy.clone()), None).await;
owner.publish(snapshot(
"select",
vec![change("alpha", "select", "merged")],
));
owner.contract(merged_contract("main"));
let socket = owner.socket();
let cwd = repo.path().to_path_buf();
let output = tokio::time::timeout(
DEADLINE_TEST_GUARD,
tokio::task::spawn_blocking(move || {
run_cli(
&cwd,
&[
"client",
"--unix-socket",
&socket,
"wait",
"alpha",
"--timeout",
"700ms",
"--json",
],
&[("PATH", path.as_str())],
)
}),
)
.await
.expect("a bounded wait must not hang")
.unwrap();
let parsed = envelope(&output);
assert_eq!(parsed["outcome"], "timeout");
assert_eq!(output.status.code(), Some(19));
assert_eq!(
parsed["detail"]["timeout_stage"], "repository_certification",
"a stalled local classification is not an owner read"
);
assert_eq!(
parsed["detail"]["last_observation"]["change"]["display_status"],
"merged"
);
assert_eq!(parsed["detail"]["commands_submitted"], 0);
assert_eq!(spy.call_count(), 0);
assert_eq!(repo.git(&["rev-parse", "HEAD"]), head_before);
owner.stop().await;
}
#[tokio::test]
async fn wait_without_a_timeout_keeps_observing_until_repository_evidence_arrives() {
let repo = Fixture::new();
repo.stage_active("alpha");
let spy = SpyExecutor::new();
let owner = Owner::start(Some(spy.clone()), None).await;
owner.publish(snapshot(
"select",
vec![change("alpha", "select", "applying")],
));
owner.contract(merged_contract("main"));
let waiting = spawn_wait(&owner, repo.path(), "alpha", None);
tokio::time::sleep(Duration::from_millis(500)).await;
assert!(
!waiting.is_finished(),
"an omitted timeout must not end the wait on its own"
);
repo.archive("alpha");
owner.publish(snapshot(
"select",
vec![change("alpha", "select", "merged")],
));
let output = waiting.await.unwrap();
let parsed = envelope(&output);
assert_eq!(
parsed["outcome"],
"completed",
"an unbounded wait must still settle from repository evidence, stderr={}",
stderr_of(&output)
);
assert_eq!(output.status.code(), Some(0));
assert_eq!(parsed["detail"]["commands_submitted"], 0);
assert_eq!(spy.call_count(), 0, "wait must submit no command");
owner.stop().await;
}
#[tokio::test]
async fn wait_with_an_explicit_zero_timeout_behaves_exactly_like_omission() {
let repo = Fixture::new();
repo.stage_active("alpha");
let spy = SpyExecutor::new();
let owner = Owner::start(Some(spy.clone()), None).await;
owner.publish(snapshot(
"select",
vec![change("alpha", "select", "applying")],
));
owner.contract(merged_contract("main"));
let waiting = spawn_wait(&owner, repo.path(), "alpha", Some("0"));
tokio::time::sleep(Duration::from_millis(500)).await;
assert!(
!waiting.is_finished(),
"`--timeout 0` must not be an immediate timeout"
);
repo.archive("alpha");
owner.publish(snapshot(
"select",
vec![change("alpha", "select", "merged")],
));
let output = waiting.await.unwrap();
let parsed = envelope(&output);
assert_eq!(
parsed["outcome"],
"completed",
"zero must select the same unbounded wait omission does, stderr={}",
stderr_of(&output)
);
assert_eq!(output.status.code(), Some(0));
assert_eq!(spy.call_count(), 0);
owner.stop().await;
}
#[tokio::test]
async fn wait_without_a_timeout_still_returns_its_typed_non_success_outcomes() {
let repo = Fixture::new();
repo.stage_active("alpha");
let spy = SpyExecutor::new();
let owner = Owner::start(Some(spy.clone()), None).await;
let mut rejected = change("alpha", "select", "rejected");
rejected.error_detail = Some("the proposal was rejected in review".to_string());
owner.publish(snapshot("select", vec![rejected]));
owner.contract(merged_contract("main"));
let output = spawn_wait(&owner, repo.path(), "alpha", None)
.await
.unwrap();
let parsed = envelope(&output);
assert_eq!(parsed["outcome"], "change_rejected");
assert_eq!(output.status.code(), Some(17));
assert_eq!(spy.call_count(), 0);
owner.stop().await;
}
#[tokio::test]
#[cfg_attr(not(feature = "heavy-tests"), ignore)]
async fn wait_unbounded_terminates_a_stalled_remote_lookup_without_reporting_timeout() {
let repo = Fixture::new();
repo.stage_active("alpha");
repo.archive("alpha");
let head_before = repo.git(&["rev-parse", "HEAD"]);
let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
.await
.expect("binds a loopback port");
let port = listener.local_addr().unwrap().port();
let (connected_tx, connected_rx) = tokio::sync::oneshot::channel();
let shutdown = tokio_util::sync::CancellationToken::new();
let task = tokio::spawn({
let stop = shutdown.clone();
async move {
let mut connected_tx = Some(connected_tx);
let mut held = Vec::new();
loop {
tokio::select! {
_ = stop.cancelled() => break,
incoming = listener.accept() => {
let Ok((stream, _)) = incoming else { break };
if let Some(tx) = connected_tx.take() {
let _ = tx.send(());
}
held.push(stream);
}
}
}
held
}
});
repo.git(&[
"remote",
"add",
"upstream",
&format!("git://127.0.0.1:{port}/stalled.git"),
]);
let spy = SpyExecutor::new();
let owner = Owner::start(Some(spy.clone()), None).await;
owner.publish(snapshot(
"select",
vec![change("alpha", "select", "merged")],
));
owner.contract(OwnerExecutionContract::resolve(
"main",
None,
Some("upstream"),
));
let waiting = spawn_wait(&owner, repo.path(), "alpha", None);
tokio::time::timeout(UNBOUNDED_SUBPROCESS_GUARD, connected_rx)
.await
.expect("git must have reached the stalled remote")
.expect("the fixture must report the connection");
shutdown.cancel();
let mut held = task.await.unwrap();
let mut stream = held.pop().expect("one accepted connection");
tokio::time::timeout(UNBOUNDED_SUBPROCESS_GUARD, async {
let mut sink = [0u8; 64];
loop {
let read = tokio::io::AsyncReadExt::read(&mut stream, &mut sink)
.await
.expect("reading the fixture connection");
if read == 0 {
break;
}
}
})
.await
.expect("the stalled git child must be terminated and reaped at its own budget");
let output = waiting.await.unwrap();
let parsed = envelope(&output);
assert_ne!(
parsed["outcome"],
"timeout",
"a per-subprocess expiry must never become the operation's timeout, stderr={}",
stderr_of(&output)
);
assert_eq!(parsed["detail"]["commands_submitted"], 0);
assert_eq!(spy.call_count(), 0, "wait must submit no command");
assert_eq!(repo.git(&["rev-parse", "HEAD"]), head_before);
assert_eq!(
repo.git(&["status", "--porcelain"]),
"",
"wait must leave the worktree clean"
);
owner.stop().await;
}
const UNBOUNDED_SUBPROCESS_GUARD: Duration = Duration::from_secs(120);
const DEADLINE_TEST_GUARD: Duration = Duration::from_secs(60);
#[tokio::test]
async fn wait_deadline_bounds_a_stalled_owner_read_rather_than_the_transport_valve() {
let repo = Fixture::new();
repo.stage_active("alpha");
let head_before = repo.git(&["rev-parse", "HEAD"]);
let dir = tempfile::tempdir().expect("temp dir");
let socket = dir.path().join("stalled.sock");
let listener = tokio::net::UnixListener::bind(&socket).expect("binds the stalled socket");
let accepted = Arc::new(AtomicUsize::new(0));
let shutdown = tokio_util::sync::CancellationToken::new();
let task = tokio::spawn({
let accepted = accepted.clone();
let stop = shutdown.clone();
async move {
let mut held = Vec::new();
loop {
tokio::select! {
_ = stop.cancelled() => break,
incoming = listener.accept() => {
let Ok((stream, _)) = incoming else { break };
accepted.fetch_add(1, Ordering::SeqCst);
held.push(stream);
}
}
}
}
});
let path = socket.display().to_string();
let output = tokio::time::timeout(
DEADLINE_TEST_GUARD,
tokio::task::spawn_blocking({
let cwd = repo.path().to_path_buf();
move || {
run_cli(
&cwd,
&[
"client",
"--unix-socket",
&path,
"wait",
"alpha",
"--timeout",
"500ms",
"--json",
],
&[],
)
}
}),
)
.await
.expect("a bounded wait must not hang")
.unwrap();
let parsed = envelope(&output);
assert_eq!(
parsed["outcome"], "timeout",
"a stalled owner must expire the operation, not the transport"
);
assert_eq!(output.status.code(), Some(19));
assert_eq!(parsed["detail"]["commands_submitted"], 0);
assert_eq!(parsed["detail"]["timeout_stage"], "initial_observation");
assert!(
parsed["detail"]["last_observation"].is_null(),
"a wait that saw nothing must say so: {parsed}"
);
assert!(
parsed["instance_id"].is_null(),
"an unobserved owner must not be named: {parsed}"
);
assert_eq!(parsed["detail"]["timeout_ms"], 500);
assert!(
accepted.load(Ordering::SeqCst) >= 1,
"the client must have been inside a request, not failing to connect"
);
assert_eq!(repo.git(&["rev-parse", "HEAD"]), head_before);
assert_eq!(
repo.git(&["status", "--porcelain"]),
"",
"a timed-out wait must leave the worktree clean"
);
shutdown.cancel();
let _ = task.await;
}
#[tokio::test]
async fn wait_deadline_terminates_a_stalled_remote_lookup_and_reaps_it() {
let repo = Fixture::new();
repo.stage_active("alpha");
repo.archive("alpha");
let head_before = repo.git(&["rev-parse", "HEAD"]);
let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
.await
.expect("binds a loopback port");
let port = listener.local_addr().unwrap().port();
let (connected_tx, connected_rx) = tokio::sync::oneshot::channel();
let shutdown = tokio_util::sync::CancellationToken::new();
let task = tokio::spawn({
let stop = shutdown.clone();
async move {
let mut connected_tx = Some(connected_tx);
let mut held = Vec::new();
loop {
tokio::select! {
_ = stop.cancelled() => break,
incoming = listener.accept() => {
let Ok((stream, _)) = incoming else { break };
if let Some(tx) = connected_tx.take() {
let _ = tx.send(());
}
held.push(stream);
}
}
}
held
}
});
repo.git(&[
"remote",
"add",
"upstream",
&format!("git://127.0.0.1:{port}/stalled.git"),
]);
let spy = SpyExecutor::new();
let owner = Owner::start(Some(spy.clone()), None).await;
owner.publish(snapshot(
"select",
vec![change("alpha", "select", "merged")],
));
owner.contract(OwnerExecutionContract::resolve(
"main",
None,
Some("upstream"),
));
let output = tokio::time::timeout(
DEADLINE_TEST_GUARD,
wait_for(&owner, repo.path(), "alpha", "2s"),
)
.await
.expect("an unbounded git ls-remote would hang here");
let parsed = envelope(&output);
assert_eq!(
parsed["outcome"], "timeout",
"the deadline owns the outcome; no later evidence error may replace it"
);
assert_eq!(output.status.code(), Some(19));
assert_eq!(parsed["detail"]["commands_submitted"], 0);
assert_eq!(
parsed["detail"]["timeout_stage"], "remote_verification",
"a stalled remote lookup is distinguishable from a local read"
);
assert_eq!(
parsed["detail"]["last_observation"]["change"]["display_status"],
"merged"
);
assert_eq!(spy.call_count(), 0, "wait must submit no command");
tokio::time::timeout(DEADLINE_TEST_GUARD, connected_rx)
.await
.expect("git must have reached the stalled remote")
.expect("the fixture must report the connection");
let held = tokio::time::timeout(DEADLINE_TEST_GUARD, async {
shutdown.cancel();
task.await.unwrap()
})
.await
.expect("the fixture must stop");
let mut stream = held.into_iter().next().expect("one accepted connection");
let mut sink = [0u8; 64];
loop {
let read = tokio::time::timeout(
DEADLINE_TEST_GUARD,
tokio::io::AsyncReadExt::read(&mut stream, &mut sink),
)
.await
.expect("the git child must have been terminated and reaped")
.expect("reading the fixture connection");
if read == 0 {
break;
}
}
assert_eq!(repo.git(&["rev-parse", "HEAD"]), head_before);
assert_eq!(
repo.git(&["status", "--porcelain"]),
"",
"a timed-out wait must leave the worktree clean"
);
owner.stop().await;
}
#[tokio::test]
async fn wait_reports_owner_replacement_when_repository_evidence_is_absent() {
let repo = Fixture::new();
repo.stage_active("alpha");
let endpoint = tempfile::tempdir().expect("temp dir");
let socket_path = endpoint.path().join("cflx-api.sock");
let spy = SpyExecutor::new();
let first = Owner::start_on(socket_path.clone(), Some(spy.clone()), None).await;
first.publish(snapshot(
"select",
vec![change("alpha", "select", "applying")],
));
first.contract(merged_contract("main"));
let repo_path = repo.path().to_path_buf();
let socket = socket_path.display().to_string();
let waiting = tokio::task::spawn_blocking(move || {
run_cli(
&repo_path,
&[
"client",
"--unix-socket",
&socket,
"wait",
"alpha",
"--timeout",
"20s",
"--json",
],
&[],
)
});
tokio::time::sleep(std::time::Duration::from_millis(200)).await;
let first_instance = first.projection.instance_id().to_string();
first.stop().await;
std::fs::remove_file(&socket_path).ok();
let second = Owner::start_on(socket_path.clone(), Some(SpyExecutor::new()), None).await;
assert_ne!(
second.projection.instance_id(),
first_instance,
"the replacement must be a different incarnation"
);
second.publish(snapshot(
"select",
vec![change("alpha", "select", "applying")],
));
second.contract(merged_contract("main"));
let output = waiting.await.unwrap();
let parsed = envelope(&output);
let outcome = parsed["outcome"].as_str().unwrap();
assert!(
outcome == "owner_restarted" || outcome == "owner_not_running",
"a replaced owner must never read as completion, got {outcome}"
);
assert_eq!(spy.call_count(), 0);
second.stop().await;
}
#[tokio::test]
async fn wait_succeeds_across_owner_replacement_on_repository_evidence_alone() {
let repo = Fixture::new();
repo.stage_active("alpha");
let endpoint = tempfile::tempdir().expect("temp dir");
let socket_path = endpoint.path().join("cflx-api.sock");
let spy = SpyExecutor::new();
let first = Owner::start_on(socket_path.clone(), Some(spy.clone()), None).await;
first.publish(snapshot(
"select",
vec![change("alpha", "select", "applying")],
));
first.contract(merged_contract("main"));
let repo_path = repo.path().to_path_buf();
let socket = socket_path.display().to_string();
let waiting = tokio::task::spawn_blocking(move || {
run_cli(
&repo_path,
&[
"client",
"--unix-socket",
&socket,
"wait",
"alpha",
"--timeout",
"20s",
"--json",
],
&[],
)
});
tokio::time::sleep(std::time::Duration::from_millis(200)).await;
repo.archive("alpha");
first.stop().await;
std::fs::remove_file(&socket_path).ok();
let second = Owner::start_on(socket_path.clone(), Some(SpyExecutor::new()), None).await;
second.publish(snapshot("select", vec![]));
second.contract(merged_contract("main"));
let output = waiting.await.unwrap();
let parsed = envelope(&output);
assert_eq!(
parsed["outcome"],
"completed",
"repository evidence alone must certify the success, stderr={}",
stderr_of(&output)
);
assert_eq!(output.status.code(), Some(0));
assert_eq!(spy.call_count(), 0);
second.stop().await;
}
#[tokio::test]
async fn wait_recovers_from_an_event_gap_by_rehydrating_every_resource() {
let repo = Fixture::new();
repo.stage_active("alpha");
let spy = SpyExecutor::new();
let owner = Owner::start(Some(spy.clone()), None).await;
owner.publish(snapshot(
"select",
vec![change("alpha", "select", "applying")],
));
owner.contract(merged_contract("main"));
let repo_path = repo.path().to_path_buf();
let socket = owner.socket();
let waiting = tokio::task::spawn_blocking(move || {
run_cli(
&repo_path,
&[
"client",
"--unix-socket",
&socket,
"wait",
"alpha",
"--timeout",
"20s",
"--json",
],
&[],
)
});
tokio::time::sleep(std::time::Duration::from_millis(150)).await;
for iteration in 0..1200 {
owner.publish(snapshot(
"select",
vec![change(
"alpha",
"select",
if iteration % 2 == 0 {
"applying"
} else {
"accepting"
},
)],
));
}
repo.archive("alpha");
owner.publish(snapshot(
"select",
vec![change("alpha", "select", "merged")],
));
let output = waiting.await.unwrap();
let parsed = envelope(&output);
assert_eq!(
parsed["outcome"],
"completed",
"a replay gap must not cost the wait its answer, stderr={}",
stderr_of(&output)
);
assert_eq!(spy.call_count(), 0);
owner.stop().await;
}
async fn subscribe_owner(spy: Arc<SpyExecutor>) -> Owner {
let owner = Owner::start_with(Some(spy), None, OwnerExtras::sink_capable()).await;
owner.publish(snapshot(
"select",
vec![
change("alpha", "select", "not queued"),
change("beta", "select", "not queued"),
],
));
owner.contract(merged_contract("main"));
owner
}
#[tokio::test]
async fn client_subscribe_routes_through_existing_owner() {
let spy = SpyExecutor::new();
let owner = subscribe_owner(spy.clone()).await;
let instance = owner.instance_id();
let tmp = neutral_cwd();
let socket = owner.socket();
let set = run_client(
&socket,
tmp.path(),
&[
"subscribe",
"set",
"alpha",
"--instance-id",
&instance,
"--json",
"--",
"/bin/true",
],
)
.await;
let parsed = envelope(&set);
assert_eq!(parsed["operation"], "subscribe_set");
assert_eq!(parsed["outcome"], "subscribed", "{}", stderr_of(&set));
assert_eq!(parsed["ok"], true);
assert_eq!(set.status.code(), Some(0));
assert_eq!(parsed["change_id"], "alpha");
assert_eq!(parsed["instance_id"], instance.as_str());
assert_eq!(
parsed["detail"]["subscriptions"][0]["sink"]["command"][0],
"/bin/true"
);
assert_eq!(
parsed["detail"]["subscriptions"][0]["sink"]["notify_blocked"],
false
);
assert!(
parsed["detail"]["subscriptions"][0]
.as_object()
.unwrap()
.get("execution_id")
.is_none(),
"{parsed}"
);
let read = run_client(
&socket,
tmp.path(),
&[
"subscribe",
"get",
"alpha",
"--instance-id",
&instance,
"--json",
],
)
.await;
let parsed = envelope(&read);
assert_eq!(parsed["operation"], "subscribe_get");
assert_eq!(parsed["outcome"], "subscribed", "{}", stderr_of(&read));
assert_eq!(parsed["detail"]["subscriptions"][0]["subscribed"], true);
assert_eq!(
parsed["detail"]["subscriptions"][0]["sink"]["command"][0],
"/bin/true"
);
let cleared = run_client(
&socket,
tmp.path(),
&[
"subscribe",
"clear",
"alpha",
"--instance-id",
&instance,
"--json",
],
)
.await;
let parsed = envelope(&cleared);
assert_eq!(parsed["operation"], "subscribe_clear");
assert_eq!(parsed["outcome"], "cleared", "{}", stderr_of(&cleared));
assert_eq!(cleared.status.code(), Some(0));
let after = run_client(
&socket,
tmp.path(),
&[
"subscribe",
"get",
"alpha",
"--instance-id",
&instance,
"--json",
],
)
.await;
let parsed = envelope(&after);
assert_eq!(parsed["detail"]["subscriptions"][0]["subscribed"], false);
assert!(parsed["detail"]["subscriptions"][0]["sink"].is_null());
assert_eq!(
spy.call_count(),
0,
"a subscription is observability, not a workflow command"
);
owner.stop().await;
}
#[tokio::test]
async fn client_subscribe_addresses_several_proposals_and_clears_only_those() {
let spy = SpyExecutor::new();
let owner = subscribe_owner(spy.clone()).await;
let instance = owner.instance_id();
let tmp = neutral_cwd();
let socket = owner.socket();
let set = run_client(
&socket,
tmp.path(),
&[
"subscribe",
"set",
"alpha",
"beta",
"--instance-id",
&instance,
"--json",
"--",
"/bin/true",
],
)
.await;
let parsed = envelope(&set);
assert_eq!(parsed["outcome"], "subscribed", "{}", stderr_of(&set));
let subscriptions = parsed["detail"]["subscriptions"].as_array().unwrap();
assert_eq!(subscriptions.len(), 2);
for (index, change_id) in ["alpha", "beta"].iter().enumerate() {
assert_eq!(subscriptions[index]["change_id"], *change_id);
assert_eq!(subscriptions[index]["subscribed"], true);
}
assert!(parsed.as_object().unwrap().get("change_id").is_none());
let cleared = run_client(
&socket,
tmp.path(),
&[
"subscribe",
"clear",
"alpha",
"--instance-id",
&instance,
"--json",
],
)
.await;
assert_eq!(envelope(&cleared)["outcome"], "cleared");
let read = run_client(
&socket,
tmp.path(),
&[
"subscribe",
"get",
"beta",
"--instance-id",
&instance,
"--json",
],
)
.await;
assert_eq!(
envelope(&read)["detail"]["subscriptions"][0]["subscribed"],
true,
"clearing alpha must leave beta alone"
);
assert_eq!(spy.call_count(), 0);
owner.stop().await;
}
#[tokio::test]
async fn client_subscribe_binds_the_episode_the_owner_later_opens() {
let spy = SpyExecutor::new();
let owner = subscribe_owner(spy.clone()).await;
let instance = owner.instance_id();
let tmp = neutral_cwd();
let socket = owner.socket();
let set = run_client(
&socket,
tmp.path(),
&[
"subscribe",
"set",
"alpha",
"--instance-id",
&instance,
"--json",
"--",
"/bin/true",
],
)
.await;
let parsed = envelope(&set);
assert_eq!(parsed["outcome"], "subscribed", "{}", stderr_of(&set));
assert!(
parsed.as_object().unwrap().get("execution_id").is_none(),
"no episode exists yet, and none may be synthesized: {parsed}"
);
let execution_id = owner.admit_registered("alpha").await;
let read = run_client(
&socket,
tmp.path(),
&[
"subscribe",
"get",
"alpha",
"--instance-id",
&instance,
"--json",
],
)
.await;
let parsed = envelope(&read);
assert_eq!(parsed["outcome"], "subscribed", "{}", stderr_of(&read));
assert_eq!(
parsed["execution_id"],
execution_id.as_str(),
"the read names the episode the owner bound: {parsed}"
);
assert_eq!(
parsed["detail"]["subscriptions"][0]["execution_id"],
execution_id.as_str()
);
assert_eq!(parsed["detail"]["subscriptions"][0]["subscribed"], true);
assert_eq!(
parsed["detail"]["subscriptions"][0]["terminal_dispatched"], false,
"the episode is live, not finished"
);
assert_eq!(spy.call_count(), 0);
owner.stop().await;
}
#[tokio::test]
async fn client_subscribe_preserves_every_argv_boundary_without_a_shell() {
let spy = SpyExecutor::new();
let owner = subscribe_owner(spy.clone()).await;
let instance = owner.instance_id();
let tmp = neutral_cwd();
let socket = owner.socket();
let argv = [
"/absolute/callback",
"--flag",
"one argument",
"*",
"$HOME",
"; rm -rf /",
"--json",
];
let mut args = vec![
"subscribe",
"set",
"alpha",
"--instance-id",
&instance,
"--json",
"--",
];
args.extend_from_slice(&argv);
let output = run_client(&socket, tmp.path(), &args).await;
let parsed = envelope(&output);
assert_eq!(parsed["outcome"], "subscribed", "{}", stderr_of(&output));
let stored = parsed["detail"]["subscriptions"][0]["sink"]["command"]
.as_array()
.expect("the socket that registered the argv is told it back")
.iter()
.map(|value| value.as_str().expect("argv elements are strings"))
.collect::<Vec<_>>();
assert_eq!(
stored, argv,
"each argument after `--` must survive as exactly one argv element"
);
assert_eq!(stored.last(), Some(&"--json"));
assert_eq!(spy.call_count(), 0);
owner.stop().await;
}
#[tokio::test]
async fn client_subscribe_blocked_delivery_is_an_explicit_opt_in() {
let spy = SpyExecutor::new();
let owner = subscribe_owner(spy.clone()).await;
let instance = owner.instance_id();
let tmp = neutral_cwd();
let socket = owner.socket();
let without = run_client(
&socket,
tmp.path(),
&[
"subscribe",
"set",
"alpha",
"--instance-id",
&instance,
"--json",
"--",
"/bin/true",
],
)
.await;
assert_eq!(
envelope(&without)["detail"]["subscriptions"][0]["sink"]["notify_blocked"],
false
);
let with = run_client(
&socket,
tmp.path(),
&[
"subscribe",
"set",
"alpha",
"--instance-id",
&instance,
"--blocked",
"--json",
"--",
"/bin/true",
],
)
.await;
assert_eq!(
envelope(&with)["detail"]["subscriptions"][0]["sink"]["notify_blocked"],
true
);
owner.stop().await;
}
#[tokio::test]
async fn client_subscribe_human_output_is_one_concise_line() {
let spy = SpyExecutor::new();
let owner = subscribe_owner(spy.clone()).await;
let instance = owner.instance_id();
let tmp = neutral_cwd();
let output = run_client(
&owner.socket(),
tmp.path(),
&[
"subscribe",
"set",
"alpha",
"--instance-id",
&instance,
"--",
"/bin/true",
],
)
.await;
let stdout = stdout_of(&output);
assert_eq!(stdout.lines().count(), 1, "human output must be one line");
assert!(stdout.starts_with("subscribe_set: subscribed"), "{stdout}");
assert!(!stdout.contains("schema_version"), "{stdout}");
assert_eq!(output.status.code(), Some(0));
owner.stop().await;
}
#[tokio::test]
async fn client_subscribe_reports_a_replaced_owner_incarnation() {
let spy = SpyExecutor::new();
let owner = subscribe_owner(spy.clone()).await;
let instance = owner.instance_id();
let tmp = neutral_cwd();
let socket = owner.socket();
let stale_instance = "0".repeat(32);
for action in ["set", "get", "clear"] {
let mut args = vec![
"subscribe",
action,
"alpha",
"--instance-id",
&stale_instance,
"--json",
];
if action == "set" {
args.extend_from_slice(&["--", "/bin/true"]);
}
let output = run_client(&socket, tmp.path(), &args).await;
let parsed = envelope(&output);
assert_eq!(parsed["outcome"], "owner_restarted", "{action}");
assert_eq!(output.status.code(), Some(8), "{action}");
assert_eq!(parsed["detail"]["expected_instance_id"], stale_instance);
assert_eq!(parsed["detail"]["observed_instance_id"], instance.as_str());
}
let read = run_client(
&socket,
tmp.path(),
&[
"subscribe",
"get",
"alpha",
"--instance-id",
&instance,
"--json",
],
)
.await;
assert_eq!(
envelope(&read)["detail"]["subscriptions"][0]["subscribed"],
false
);
assert_eq!(spy.call_count(), 0);
owner.stop().await;
}
#[tokio::test]
async fn client_subscribe_reports_an_owner_without_the_surface_as_unsupported() {
let spy = SpyExecutor::new();
let owner = Owner::start(Some(spy.clone()), None).await;
owner.publish(snapshot(
"select",
vec![change("alpha", "select", "not queued")],
));
owner.contract(merged_contract("main"));
let instance = owner.instance_id();
let tmp = neutral_cwd();
let socket = owner.socket();
for action in ["set", "get", "clear"] {
let mut args = vec![
"subscribe",
action,
"alpha",
"--instance-id",
&instance,
"--json",
];
if action == "set" {
args.extend_from_slice(&["--", "/bin/true"]);
}
let output = run_client(&socket, tmp.path(), &args).await;
let parsed = envelope(&output);
assert_eq!(parsed["outcome"], "unsupported_owner", "{action}");
assert_eq!(output.status.code(), Some(26), "{action}");
assert!(
stderr_of(&output).contains("unsupported_owner"),
"the diagnostic belongs on stderr"
);
}
assert_eq!(spy.call_count(), 0);
owner.stop().await;
}
#[tokio::test]
async fn client_subscribe_mutation_is_refused_off_the_owner_socket() {
let spy = SpyExecutor::new();
let owner = Owner::start_with(
Some(spy.clone()),
None,
OwnerExtras {
sinks: true,
transport: Some(conflux::web::remote_control_api::ApiTransport::Tcp),
},
)
.await;
owner.publish(snapshot(
"select",
vec![change("alpha", "select", "not queued")],
));
owner.contract(merged_contract("main"));
let instance = owner.instance_id();
let tmp = neutral_cwd();
let socket = owner.socket();
for action in ["set", "clear"] {
let mut args = vec![
"subscribe",
action,
"alpha",
"--instance-id",
&instance,
"--json",
];
if action == "set" {
args.extend_from_slice(&["--", "/bin/true"]);
}
let output = run_client(&socket, tmp.path(), &args).await;
let parsed = envelope(&output);
assert_eq!(parsed["outcome"], "transport_not_permitted", "{action}");
assert_eq!(output.status.code(), Some(25), "{action}");
}
let read = run_client(
&socket,
tmp.path(),
&[
"subscribe",
"get",
"alpha",
"--instance-id",
&instance,
"--json",
],
)
.await;
let parsed = envelope(&read);
assert_eq!(parsed["outcome"], "subscribed", "{}", stderr_of(&read));
assert_eq!(parsed["detail"]["subscriptions"][0]["subscribed"], false);
assert!(parsed["detail"]["subscriptions"][0]["sink"].is_null());
assert_eq!(spy.call_count(), 0);
owner.stop().await;
}
#[derive(Default)]
struct SmokeScheduler {
launches: Arc<Mutex<Vec<Vec<String>>>>,
running: AtomicBool,
}
impl SmokeScheduler {
fn launches(&self) -> Vec<Vec<String>> {
self.launches.lock().unwrap().clone()
}
}
#[async_trait]
impl conflux::orchestration::run_control::RunSchedulerPort for SmokeScheduler {
fn is_running(&self) -> bool {
self.running.load(Ordering::SeqCst)
}
async fn prepare_run(
&self,
targets: Vec<String>,
_explicit_retry: bool,
) -> Result<conflux::orchestration::run_control::RunPermit, String> {
let recorded = self.launches.clone();
Ok(conflux::orchestration::run_control::RunPermit::new(
move || recorded.lock().unwrap().push(targets),
))
}
async fn notify_scheduler(&self) {}
async fn cancel_run(&self) {}
fn set_graceful_stop(&self, _requested: bool) {}
async fn stop_activity(&self) -> conflux::tui::stop_classification::StopActivitySnapshot {
use conflux::tui::stop_classification::{
ExecutionEvidence, ShutdownWorkEvidence, StopActivitySnapshot,
};
StopActivitySnapshot {
execution_handles: ExecutionEvidence::Known { registered: 0 },
reducer_agent_execution_active: false,
shutdown_work: ShutdownWorkEvidence::Known { pending: false },
}
}
}
#[derive(Default)]
struct SmokeQueue {
added: Mutex<Vec<String>>,
}
#[async_trait]
impl conflux::orchestration::operator_command::QueuePort for SmokeQueue {
async fn add(&self, change_id: &str) -> bool {
self.added.lock().unwrap().push(change_id.to_string());
true
}
async fn remove(&self, _change_id: &str) -> bool {
true
}
async fn request_cancellation(
&self,
_change_id: &str,
) -> Result<Option<conflux::orchestration::operator_command::TerminationWaiter>, String>
{
Ok(None)
}
async fn notify_scheduler(&self) {}
}
#[tokio::test]
async fn production_owner_smoke_admits_one_change_through_the_shared_coordinator() {
use conflux::orchestration::operator_command::{
ExecutionMarkStore, NoopQueueHooks, OperatorCommandService, ParallelRuntime,
};
use conflux::orchestration::operator_coordinator::CoreMode;
use conflux::orchestration::run_control::{ResolveReservations, RunControlService};
use conflux::orchestration::state::OrchestratorState;
use conflux::web::state::WebState;
let reducer = Arc::new(tokio::sync::RwLock::new(OrchestratorState::new(
vec!["alpha".to_string()],
4,
)));
let marks = Arc::new(ExecutionMarkStore::new());
let parallel = Arc::new(ParallelRuntime::new());
let queue = Arc::new(SmokeQueue::default());
let scheduler = Arc::new(SmokeScheduler::default());
let service = Arc::new(
OperatorCommandService::new(
reducer.clone(),
queue.clone(),
Arc::new(NoopQueueHooks),
marks.clone(),
)
.with_parallel(parallel.clone()),
);
let run_control = Arc::new(RunControlService::new(
reducer.clone(),
service,
scheduler.clone(),
Arc::new(ResolveReservations::new()),
parallel,
));
let listing = conflux::openspec::Change {
id: "alpha".to_string(),
completed_tasks: 0,
total_tasks: 2,
last_modified: "now".to_string(),
dependencies: Vec::new(),
metadata: Default::default(),
};
let web_state = Arc::new(WebState::new(std::slice::from_ref(&listing)));
web_state.set_shared_state(reducer.clone()).await;
web_state.set_execution_marks(marks.clone()).await;
web_state
.seed_workspace_observation_for_tests(&[listing], "select")
.await;
web_state.sync_remote_control_projection().await;
let core_mode = Arc::new(CoreMode::new());
let (executor, application) = conflux::web::remote_control_api::executor::wired_for_test(
reducer.clone(),
run_control,
web_state.clone(),
core_mode,
);
let runtime = web_state.remote_control();
runtime.bind(Arc::new(executor)).await;
runtime.bind_gate(application.gate()).await;
runtime.bind_run_boundary(scheduler.clone());
let dir = tempfile::tempdir().expect("temp dir");
let socket_path = dir.path().join("cflx-api.sock");
let auth = RemoteControlAuth::new(None, &[]).expect("auth policy is valid");
let app = router(
RemoteControlState::new(runtime.projection(), Arc::new(auth), runtime.clone())
.with_gate(runtime.gate())
.with_execution_facts(runtime.execution_facts())
.with_execution_contract(runtime.execution_contract()),
);
let listener = tokio::net::UnixListener::bind(&socket_path).expect("binds");
let shutdown = tokio_util::sync::CancellationToken::new();
let serve = tokio::spawn({
let shutdown = shutdown.clone();
async move {
let _ = axum::serve(listener, app)
.with_graceful_shutdown(async move { shutdown.cancelled().await })
.await;
}
});
let workspace = tempfile::tempdir().expect("temp dir");
let socket = socket_path.display().to_string();
let output = tokio::task::spawn_blocking({
let cwd = workspace.path().to_path_buf();
move || {
run_cli(
&cwd,
&[
"client",
"--unix-socket",
&socket,
"mark",
"alpha",
"--json",
],
&[],
)
}
})
.await
.unwrap();
let parsed = envelope(&output);
assert_eq!(parsed["outcome"], "marked", "stderr={}", stderr_of(&output));
assert_eq!(output.status.code(), Some(0));
assert!(
marks.is_marked("alpha"),
"the shared execution-mark store must hold the operator's selection"
);
assert!(
scheduler.launches().is_empty(),
"marking is selection: nothing may be launched by it"
);
let socket = socket_path.display().to_string();
let started = tokio::task::spawn_blocking({
let cwd = workspace.path().to_path_buf();
move || {
run_cli(
&cwd,
&["client", "--unix-socket", &socket, "start", "--json"],
&[],
)
}
})
.await
.unwrap();
let parsed = envelope(&started);
assert_eq!(
parsed["outcome"],
"accepted",
"stderr={}",
stderr_of(&started)
);
assert!(
marks.is_marked("alpha"),
"the shared execution-mark store must hold the started change"
);
assert_eq!(
scheduler.launches(),
vec![vec!["alpha".to_string()]],
"exactly one launch, for exactly the requested change"
);
shutdown.cancel();
let _ = serve.await;
}
#[test]
fn wait_documentation_states_which_statuses_hold_and_which_release() {
let repo_root = Path::new(env!("CARGO_MANIFEST_DIR"));
let agents = std::fs::read_to_string(repo_root.join("AGENTS.md")).expect("AGENTS.md");
let readme = std::fs::read_to_string(repo_root.join("README.md")).expect("README.md");
for (document, name) in [(&agents, "AGENTS.md"), (&readme, "README.md")] {
assert!(
document.contains("change_requires_action"),
"{name} must name the new outcome a script has to branch on"
);
assert!(
document.contains("`27`"),
"{name} must give the exit status a shell branches on"
);
for released in ["error", "merge wait", "stopped", "stalled"] {
assert!(
document.contains(released),
"{name} must name '{released}' as a status that releases"
);
}
for held in ["queued", "blocked", "applying", "archiving"] {
assert!(
document.contains(held),
"{name} must name '{held}' as a status that keeps observing"
);
}
assert!(
document.contains("observed_status"),
"{name} must name the detail field that says which status released"
);
assert!(
document.contains("not repairing") || document.contains("not a repair"),
"{name} must say releasing is still observation-only"
);
assert!(
document.contains("external"),
"{name} must name the external blocker that releases a blocked row"
);
assert!(
document.contains("dependency"),
"{name} must name the blocked hold that keeps observing"
);
assert!(
document.contains("detail.blocker"),
"{name} must name the detail field carrying the prerequisite facts"
);
}
let cwd = neutral_cwd();
let help = stdout_of(&run_cli(cwd.path(), &["client", "wait", "--help"], &[]));
assert!(
help.contains("change_requires_action"),
"`client wait --help` must document the new outcome, got:\n{help}"
);
assert!(
help.contains("merge wait") && help.contains("stalled"),
"`client wait --help` must say which rows release, got:\n{help}"
);
assert!(
help.contains("external") && help.contains("dependency"),
"`client wait --help` must say which blocked rows release, got:\n{help}"
);
}
#[test]
fn wait_documentation_describes_the_machine_readable_timeout_detail() {
let repo_root = Path::new(env!("CARGO_MANIFEST_DIR"));
let agents = std::fs::read_to_string(repo_root.join("AGENTS.md")).expect("AGENTS.md");
let readme = std::fs::read_to_string(repo_root.join("README.md")).expect("README.md");
let cwd = neutral_cwd();
let help = stdout_of(&run_cli(cwd.path(), &["client", "wait", "--help"], &[]));
for (document, name) in [
(&agents, "AGENTS.md"),
(&readme, "README.md"),
(&help, "client wait --help"),
] {
for field in ["timeout_ms", "wait_elapsed_ms", "timeout_stage"] {
assert!(
document.contains(field),
"{name} must name the '{field}' a caller branches on"
);
}
for stage in [
"initial_observation",
"observing_owner",
"repository_certification",
"remote_verification",
] {
assert!(
document.contains(stage),
"{name} must name the '{stage}' timeout stage"
);
}
assert!(
document.contains("last_observation"),
"{name} must name the retained observation field"
);
assert!(
document.contains("null"),
"{name} must say what an unobserved wait reports instead"
);
assert!(
document.contains("after expiry"),
"{name} must say no read happens after the deadline"
);
assert!(
document.contains("lifecycle status"),
"{name} must say a timeout leaves the proposal's status alone"
);
}
}
#[test]
fn documentation_recommends_the_client_for_existing_owner_delegation() {
let repo_root = Path::new(env!("CARGO_MANIFEST_DIR"));
let agents = std::fs::read_to_string(repo_root.join("AGENTS.md")).expect("AGENTS.md");
assert!(
agents.contains("cflx client"),
"AGENTS.md must document the client namespace"
);
for example in [
"cflx client status",
"cflx client mark",
"cflx client start",
"cflx client wait",
"cflx client subscribe set",
] {
assert!(agents.contains(example), "AGENTS.md must show `{example}`");
}
for retired in ["cflx client enqueue", "cflx client notify", "cflx_enqueue"] {
assert!(
!agents.contains(retired),
"AGENTS.md must not recommend the retired `{retired}`"
);
}
assert!(
agents.contains("/api/v2"),
"AGENTS.md must retain the API as the low-level contract"
);
assert!(
agents.contains("cflx run") && agents.contains("owner"),
"AGENTS.md must keep the run/client ownership distinction"
);
for forbidden in ["expected_revision", "idempotency_key"] {
let recommended = agents
.split("## Delegating to an existing owner")
.nth(1)
.unwrap_or("");
assert!(
!recommended.contains(forbidden),
"the delegation guidance must not tell a caller to construct {forbidden}"
);
}
}
#[test]
fn documentation_separates_marking_from_starting_from_admission() {
let repo_root = Path::new(env!("CARGO_MANIFEST_DIR"));
let agents = std::fs::read_to_string(repo_root.join("AGENTS.md")).expect("AGENTS.md");
let readme = std::fs::read_to_string(repo_root.join("README.md")).expect("README.md");
for (document, name) in [(&agents, "AGENTS.md"), (&readme, "README.md")] {
assert!(
document.contains("preserve") || document.contains("preserves"),
"{name} must say a mark preserves unrelated marks"
);
assert!(
document.contains("unrelated mark"),
"{name} must name what a mark leaves alone"
);
assert!(
document.contains("F5"),
"{name} must describe Start through its TUI equivalent"
);
assert!(
document.contains("authoritative mark set"),
"{name} must say Start consumes the owner's own marks"
);
assert!(
document.contains("settlement"),
"{name} must attribute admission to owner-side settlement"
);
assert!(
!document.contains("cflx client enqueue"),
"{name} must not recommend the retired admission verb"
);
}
}
#[test]
fn documentation_states_the_mcp_and_subscription_contract() {
let repo_root = Path::new(env!("CARGO_MANIFEST_DIR"));
let agents = std::fs::read_to_string(repo_root.join("AGENTS.md")).expect("AGENTS.md");
let readme = std::fs::read_to_string(repo_root.join("README.md")).expect("README.md");
for document in [&agents, &readme] {
assert!(
document.contains("cflx client mcp"),
"every document must show the MCP server"
);
for tool in ["cflx_status", "cflx_control", "cflx_subscribe"] {
assert!(document.contains(tool), "{tool} must be documented");
}
for retired in ["cflx_enqueue", "cflx_notify_set"] {
assert!(
!document.contains(retired),
"the retired tool {retired} must not be documented"
);
}
assert!(
document.contains("cflx_wait is deliberately absent")
|| document.contains("cflx_wait` is deliberately absent")
|| document.contains("deliberately absent from MCP"),
"the documents must explain that MCP has no wait tool"
);
assert!(
document.contains("transport_not_permitted"),
"the transport refusal is the rule an operator hits first"
);
for variable in [
"CFLX_EVENT_PATH",
"CFLX_EVENT_TYPE",
"CFLX_EXECUTION_ID",
"CFLX_CHANGE_ID",
"CFLX_INSTANCE_ID",
] {
assert!(document.contains(variable), "{variable} must be documented");
}
assert!(
document.contains("process exit was never"),
"a resident TUI's liveness must be distinguished from completion"
);
assert!(
document.contains("resumes no agent")
|| document.contains("does not resume")
|| document.contains("resumes nothing")
|| document.contains("it does not resume"),
"every document must say delivery does not resume an agent"
);
assert!(
document.contains("--auth-token-env"),
"the credential rule must be stated"
);
assert!(
document.contains("before"),
"the pre-admission registration must be documented"
);
assert!(
document.contains("episode"),
"the delivery episode is the unit a reader has to understand"
);
assert!(
document.contains("`0400`") && document.contains("`0700`"),
"the event artifact's default permissions must be documented"
);
assert!(
document.contains("never reads it back") || document.contains("never read back"),
"and why those permissions are not an integrity guarantee"
);
assert!(
document.contains("full pipe"),
"an operator writing a chatty callback needs to know it will not wedge"
);
assert!(
document.contains("truncation diagnostic"),
"and that overflow is a diagnostic rather than a delivery failure"
);
assert!(
document.contains("reaped"),
"shutdown must promise termination and reaping, not just a wait"
);
}
let skill = std::fs::read_to_string(repo_root.join("skills/cflx-run/SKILL.md"))
.expect("the embedded cflx-run skill");
for (document, name) in [
(&agents, "AGENTS.md"),
(&readme, "README.md"),
(&skill, "the cflx-run skill"),
] {
assert!(
document.contains("cflx client subscribe set"),
"{name} must show the direct callback command"
);
assert!(
document.contains("cflx client subscribe clear"),
"{name} must show how a registration is withdrawn"
);
assert!(
document.contains("cflx_subscribe"),
"{name} must keep MCP as the alternative for a host with no shell"
);
for retired in ["auto-resume", "cflx_enqueue", "cflx client enqueue"] {
assert!(
!document.contains(retired),
"{name} must not carry retired guidance: {retired}"
);
}
}
assert!(
skill.contains("no `sh -c`"),
"the skill must say the callback is argv rather than shell source"
);
assert!(
skill.contains("not process completion"),
"and that a resident TUI's exit is not the signal"
);
assert!(
skill.contains("Nothing subscribes you to completion automatically"),
"the skill must state that registration is explicit"
);
assert!(
agents.contains("lifecycle adapter") && agents.contains("one-shot"),
"AGENTS.md must explain why subscriptions are not the lifecycle adapter"
);
assert!(
agents.contains("owner_restarted"),
"AGENTS.md must name the typed answer for a lost owner"
);
}
fn killable(id: &str, app_mode: &str, display_status: &str) -> ChangeResource {
ChangeResource {
actions:
conflux::web::remote_control_api::projection::change_actions_with_live_process_for_test(
app_mode,
display_status,
true,
),
..change(id, app_mode, display_status)
}
}
fn assert_message_is_unbroken(parsed: &serde_json::Value, context: &str) {
let message = parsed["message"]
.as_str()
.unwrap_or_else(|| panic!("{context}: the envelope carries no message: {parsed}"));
assert!(
!message.contains(" "),
"{context}: the message carries collapsed literal indentation: {message:?}"
);
}
#[tokio::test]
async fn force_stop_change_submits_only_the_targeted_command() {
let spy = SpyExecutor::new();
let owner = Owner::start(Some(spy.clone()), None).await;
owner.publish(snapshot(
"running",
vec![
killable("alpha", "running", "applying"),
killable("beta", "running", "applying"),
],
));
let output = control(&owner, &["force-stop-change", "alpha"]).await;
let parsed = envelope(&output);
assert_eq!(parsed["outcome"], "stopped", "{}", stderr_of(&output));
assert_eq!(parsed["operation"], "control_force_stop_change");
assert_eq!(parsed["change_id"], "alpha");
assert_eq!(output.status.code(), Some(0));
assert_message_is_unbroken(&parsed, "the settled stopped envelope");
assert_eq!(
spy.calls(),
vec![CommandSpec::ForceStopChange {
change_id: "alpha".to_string()
}],
"the client submits the shared targeted command and nothing else"
);
owner.stop().await;
}
#[tokio::test]
async fn force_stop_change_requires_exactly_one_target() {
let spy = SpyExecutor::new();
let owner = Owner::start(Some(spy.clone()), None).await;
owner.publish(snapshot(
"running",
vec![
killable("alpha", "running", "applying"),
killable("beta", "running", "applying"),
],
));
for tail in [
vec!["force-stop-change"],
vec!["force-stop-change", "alpha", "beta"],
] {
let output = control(&owner, &tail).await;
assert_eq!(output.status.code(), Some(2), "{tail:?}");
}
assert!(
spy.calls().is_empty(),
"a shape refusal must create no command record"
);
owner.stop().await;
}
#[tokio::test]
async fn force_stop_change_refuses_an_ineligible_target_before_submitting() {
let spy = SpyExecutor::new();
let owner = Owner::start(Some(spy.clone()), None).await;
owner.publish(snapshot(
"running",
vec![
change("alpha", "running", "applying"),
change("gamma", "running", "merge wait"),
],
));
for (target, expected_status) in [("alpha", "applying"), ("gamma", "merge wait")] {
let output = control(&owner, &["force-stop-change", target]).await;
let parsed = envelope(&output);
assert_eq!(parsed["outcome"], "target_ineligible", "{target}");
assert_eq!(parsed["operation"], "control_force_stop_change", "{target}");
assert_eq!(parsed["change_id"], target);
assert_eq!(
parsed["detail"]["observed_status"], expected_status,
"{target}"
);
assert_eq!(output.status.code(), Some(10), "{target}");
assert_message_is_unbroken(&parsed, &format!("the {target} refusal"));
}
let unknown = control(&owner, &["force-stop-change", "never-tracked"]).await;
assert_eq!(envelope(&unknown)["outcome"], "change_not_found");
assert_eq!(unknown.status.code(), Some(9));
assert!(
spy.calls().is_empty(),
"a refused target must create no command record"
);
owner.stop().await;
}
#[tokio::test]
async fn force_stop_change_is_never_a_spelling_of_process_wide_force_stop() {
let spy = SpyExecutor::new();
let owner = Owner::start(Some(spy.clone()), None).await;
owner.publish(snapshot(
"running",
vec![killable("alpha", "running", "applying")],
));
control(&owner, &["force-stop-change", "alpha"]).await;
control(&owner, &["force-stop"]).await;
assert_eq!(
submitted_names(&spy),
vec!["force_stop_change", "force_stop"],
"each verb submits its own command; neither can produce the other"
);
let widened = control(&owner, &["force-stop", "alpha"]).await;
assert_eq!(envelope(&widened)["outcome"], "usage_error");
assert_eq!(spy.call_count(), 2);
owner.stop().await;
}
#[test]
fn force_stop_change_help_documents_the_one_target_contract() {
let tmp = tempfile::tempdir().expect("temp dir");
let help = run_cli(tmp.path(), &["client", "force-stop-change", "--help"], &[]);
assert!(help.status.success(), "{}", stderr_of(&help));
let text = stdout_of(&help);
for phrase in [
"SIGKILL",
"one named proposal",
"stopped",
"effects_rolled_back",
] {
assert!(
text.contains(phrase),
"the help must document `{phrase}`:\n{text}"
);
}
}
#[test]
fn force_stop_change_documentation_states_the_one_target_contract() {
let repo_root = Path::new(env!("CARGO_MANIFEST_DIR"));
let agents = std::fs::read_to_string(repo_root.join("AGENTS.md")).expect("AGENTS.md");
let readme = std::fs::read_to_string(repo_root.join("README.md")).expect("README.md");
let skill = std::fs::read_to_string(repo_root.join("skills/cflx-run/SKILL.md"))
.expect("the bundled client skill");
let reference =
std::fs::read_to_string(repo_root.join("skills/cflx-run/references/cflx-run.md"))
.expect("the bundled client skill reference");
for (document, name) in [
(&agents, "AGENTS.md"),
(&readme, "README.md"),
(&skill, "skills/cflx-run/SKILL.md"),
(&reference, "skills/cflx-run/references/cflx-run.md"),
] {
assert!(
document.contains("force-stop-change") || document.contains("force_stop_change"),
"{name} must document the targeted control"
);
assert!(
document.contains("SIGKILL"),
"{name} must say the graceful window is bypassed"
);
assert!(
document.contains("effects_rolled_back"),
"{name} must say completed worktree effects survive"
);
assert!(
document.contains("`stopped`") || document.contains("outcome is `stopped`"),
"{name} must name the distinct success token"
);
}
assert!(
agents.contains("control_force_stop_change"),
"AGENTS.md must name the envelope operation a script branches on"
);
for (document, name) in [(&agents, "AGENTS.md"), (&readme, "README.md")] {
assert!(
document.contains("actions.force_stop_change"),
"{name} must point at the published per-change eligibility"
);
assert!(
document.contains("exactly one"),
"{name} must state the one-target rule"
);
assert!(
document.contains("stop_and_dequeue"),
"{name} must say the graceful control is unchanged"
);
}
}
}