use crate::client::envelope::{Operation, Outcome, ResultEnvelope};
use crate::client::session::{describe_api_error, observe, Connection, Observation};
use crate::client::transport::TransportError;
use crate::web::remote_control_api::dto::{
new_hex_id, ActionBlockedReason, CommandRecord, CommandRequest, CommandSpec, CommandState,
ErrorCode,
};
pub const MAX_TARGETS: usize = 64;
pub const MAX_REVISION_ATTEMPTS: usize = 3;
const SETTLEMENT_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(60);
const SETTLEMENT_POLL: std::time::Duration = std::time::Duration::from_millis(25);
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Action {
Mark,
Unmark,
Start,
Stop,
ForceStop,
ForceStopChange,
}
impl Action {
pub fn operation(self) -> Operation {
match self {
Self::Mark => Operation::ControlMark,
Self::Unmark => Operation::ControlUnmark,
Self::Start => Operation::ControlStart,
Self::Stop => Operation::ControlStop,
Self::ForceStop => Operation::ControlForceStop,
Self::ForceStopChange => Operation::ControlForceStopChange,
}
}
pub fn as_str(self) -> &'static str {
match self {
Self::Mark => "mark",
Self::Unmark => "unmark",
Self::Start => "start",
Self::Stop => "stop",
Self::ForceStop => "force_stop",
Self::ForceStopChange => "force_stop_change",
}
}
pub fn parse(value: &str) -> Option<Self> {
Some(match value {
"mark" => Self::Mark,
"unmark" => Self::Unmark,
"start" => Self::Start,
"stop" => Self::Stop,
"force_stop" => Self::ForceStop,
"force_stop_change" => Self::ForceStopChange,
_ => return None,
})
}
pub fn is_mark(self) -> bool {
matches!(self, Self::Mark | Self::Unmark)
}
pub fn is_single_target_lifecycle(self) -> bool {
matches!(self, Self::ForceStopChange)
}
fn desired_mark(self) -> Option<bool> {
match self {
Self::Mark => Some(true),
Self::Unmark => Some(false),
_ => None,
}
}
fn lifecycle_command(self) -> Option<CommandSpec> {
match self {
Self::Start => Some(CommandSpec::Start),
Self::Stop => Some(CommandSpec::Stop),
Self::ForceStop => Some(CommandSpec::ForceStop),
_ => None,
}
}
}
pub fn validate_targets(change_ids: &[String]) -> Result<(), String> {
if change_ids.is_empty() {
return Err("name at least one proposal".to_string());
}
if change_ids.len() > MAX_TARGETS {
return Err(format!(
"one request may name at most {MAX_TARGETS} proposals; {} were given",
change_ids.len()
));
}
let mut seen = std::collections::BTreeSet::new();
for change_id in change_ids {
if !seen.insert(change_id.as_str()) {
return Err(format!("'{change_id}' is named more than once"));
}
}
Ok(())
}
pub fn validate_request(action: Action, change_ids: &[String]) -> Result<(), String> {
if action.is_mark() {
return validate_targets(change_ids);
}
if action.is_single_target_lifecycle() {
validate_targets(change_ids)?;
if change_ids.len() != 1 {
return Err(format!(
"'{}' addresses exactly one proposal; {} were given. It is not a process-wide \
stop, and it accepts no target list",
action.as_str(),
change_ids.len()
));
}
return Ok(());
}
if !change_ids.is_empty() {
return Err(format!(
"'{}' consumes the owner's authoritative mark set and accepts no target list; \
mark the proposals you want first",
action.as_str()
));
}
Ok(())
}
pub async fn run(connection: &Connection, action: Action, change_ids: &[String]) -> ResultEnvelope {
let operation = action.operation();
if let Err(message) = validate_request(action, change_ids) {
return ResultEnvelope::new(operation, Outcome::UsageError).with_message(message);
}
if action.is_mark() {
return marks(connection, action, change_ids).await;
}
if action.is_single_target_lifecycle() {
return single_target_lifecycle(connection, action, &change_ids[0]).await;
}
lifecycle(connection, action).await
}
async fn single_target_lifecycle(
connection: &Connection,
action: Action,
change_id: &str,
) -> ResultEnvelope {
let operation = action.operation();
let mut first_instance: Option<String> = None;
for attempt in 0..MAX_REVISION_ATTEMPTS {
let observation = match observe(connection, None).await {
Ok(observation) => observation,
Err(error) if error.is_transient() && attempt + 1 < MAX_REVISION_ATTEMPTS => continue,
Err(error) => return error.into_envelope(operation),
};
match &first_instance {
None => first_instance = Some(observation.instance_id.clone()),
Some(expected) if *expected != observation.instance_id => {
return restarted_envelope(operation, expected, &observation.instance_id)
}
Some(_) => {}
}
let instance = Some(observation.instance_id.clone());
if !observation.command_capable() {
return not_command_capable(operation, instance);
}
let Some(change) = observation.change(change_id) else {
return ResultEnvelope::new(operation, Outcome::ChangeNotFound)
.with_instance(instance)
.with_change(change_id.to_string())
.with_message(format!(
"the owner does not track a proposal named '{change_id}'"
));
};
let eligibility = change.actions.force_stop_change;
if !eligibility.allowed {
return ResultEnvelope::new(operation, Outcome::TargetIneligible)
.with_instance(instance)
.with_change(change_id.to_string())
.with_message(format!(
"this owner cannot force-stop '{change_id}' right now ({}); its display \
status is '{}'",
describe_block(eligibility.blocked_reason),
change.display_status
))
.with_detail(serde_json::json!({
"action": action.as_str(),
"commands_submitted": Vec::<serde_json::Value>::new(),
"observed_status": change.display_status,
"blocked_reason": eligibility.blocked_reason,
}));
}
let mut audit = Vec::new();
match submit_and_settle(
connection,
CommandSpec::ForceStopChange {
change_id: change_id.to_string(),
},
action.as_str(),
Some(change_id),
observation.state_revision,
&observation.instance_id,
&mut audit,
)
.await
{
Ok(record) => {
return ResultEnvelope::new(operation, Outcome::Stopped)
.with_instance(instance)
.with_change(change_id.to_string())
.with_message(format!(
"'{change_id}' was force-stopped and dequeued; completed worktree \
effects were not rolled back"
))
.with_detail(serde_json::json!({
"action": action.as_str(),
"commands_submitted": audit,
"command_state": record.state,
"result_revision": record.result_revision,
"result": record.result,
"detail": record.detail,
}))
}
Err(SubmitFailure::Stale { .. }) => continue,
Err(failure) => {
return failure
.into_envelope(operation, instance)
.with_change(change_id.to_string())
.with_detail(serde_json::json!({
"action": action.as_str(),
"commands_submitted": audit,
}))
}
}
}
revision_conflict(operation, first_instance)
}
async fn lifecycle(connection: &Connection, action: Action) -> ResultEnvelope {
let operation = action.operation();
let command = action
.lifecycle_command()
.expect("a lifecycle action always names a command");
let mut first_instance: Option<String> = None;
for attempt in 0..MAX_REVISION_ATTEMPTS {
let observation = match observe(connection, None).await {
Ok(observation) => observation,
Err(error) if error.is_transient() && attempt + 1 < MAX_REVISION_ATTEMPTS => continue,
Err(error) => return error.into_envelope(operation),
};
match &first_instance {
None => first_instance = Some(observation.instance_id.clone()),
Some(expected) if *expected != observation.instance_id => {
return restarted_envelope(operation, expected, &observation.instance_id)
}
Some(_) => {}
}
let instance = Some(observation.instance_id.clone());
if !observation.command_capable() {
return not_command_capable(operation, instance);
}
let mut audit = Vec::new();
match submit_and_settle(
connection,
command.clone(),
action.as_str(),
None,
observation.state_revision,
&observation.instance_id,
&mut audit,
)
.await
{
Ok(record) => {
return ResultEnvelope::new(operation, Outcome::Accepted)
.with_instance(instance)
.with_message(format!(
"the owner accepted the shared '{}' intent",
action.as_str()
))
.with_detail(serde_json::json!({
"action": action.as_str(),
"commands_submitted": audit,
"command_state": record.state,
"result_revision": record.result_revision,
"detail": record.detail,
}))
}
Err(SubmitFailure::Stale { .. }) => continue,
Err(failure) => {
return failure
.into_envelope(operation, instance)
.with_detail(serde_json::json!({
"action": action.as_str(),
"commands_submitted": audit,
}))
}
}
}
revision_conflict(operation, first_instance)
}
#[derive(Debug, Clone, PartialEq, Eq)]
enum Plan {
Submit,
Satisfied,
}
#[derive(Debug, Clone, PartialEq, Eq)]
struct Refusal {
outcome: Outcome,
change_id: String,
message: String,
}
#[derive(Debug, Clone, PartialEq, Eq)]
struct TargetResult {
change_id: String,
changed: bool,
reason: String,
}
impl TargetResult {
fn to_json(&self) -> serde_json::Value {
serde_json::json!({
"change_id": self.change_id,
"changed": self.changed,
"reason": self.reason,
})
}
}
fn classify(
observation: &Observation,
change_id: &str,
desired: bool,
) -> Result<Plan, Box<Refusal>> {
let Some(change) = observation.change(change_id) else {
return Err(Box::new(Refusal {
outcome: Outcome::ChangeNotFound,
change_id: change_id.to_string(),
message: format!("the owner does not track a proposal named '{change_id}'"),
}));
};
let mark = &change.actions.set_execution_mark;
if !mark.allowed && !matches!(mark.blocked_reason, Some(ActionBlockedReason::FinalStatus)) {
return Err(Box::new(Refusal {
outcome: Outcome::TargetIneligible,
change_id: change_id.to_string(),
message: format!(
"this owner refuses execution-mark mutation right now ({}), so no proposal in \
this request was marked",
describe_block(mark.blocked_reason)
),
}));
}
if change.execution_marked == desired {
return Ok(Plan::Satisfied);
}
Ok(Plan::Submit)
}
fn describe_block(reason: Option<ActionBlockedReason>) -> String {
match reason {
Some(reason) => serde_json::to_value(reason)
.ok()
.and_then(|value| value.as_str().map(str::to_string))
.unwrap_or_else(|| format!("{reason:?}")),
None => "no reason published".to_string(),
}
}
async fn marks(connection: &Connection, action: Action, change_ids: &[String]) -> ResultEnvelope {
let operation = action.operation();
let desired = action
.desired_mark()
.expect("a mark action always names a desired state");
let mut settled: Vec<TargetResult> = Vec::new();
let mut audit: Vec<serde_json::Value> = Vec::new();
let mut first_instance: Option<String> = None;
for attempt in 0..MAX_REVISION_ATTEMPTS {
let observation = match observe(connection, None).await {
Ok(observation) => observation,
Err(error) if error.is_transient() && attempt + 1 < MAX_REVISION_ATTEMPTS => continue,
Err(error) if settled.is_empty() => return error.into_envelope(operation),
Err(error) => {
return partial(operation, first_instance, &settled, &audit, error.message())
}
};
match &first_instance {
None => first_instance = Some(observation.instance_id.clone()),
Some(expected) if *expected != observation.instance_id => {
return restarted_envelope(operation, expected, &observation.instance_id)
}
Some(_) => {}
}
let instance = Some(observation.instance_id.clone());
if !observation.command_capable() {
return not_command_capable(operation, instance);
}
let remaining: Vec<&String> = change_ids
.iter()
.filter(|change_id| !settled.iter().any(|result| result.change_id == **change_id))
.collect();
let mut plans = Vec::with_capacity(remaining.len());
for change_id in &remaining {
match classify(&observation, change_id, desired) {
Ok(plan) => plans.push((*change_id, plan)),
Err(refusal) if settled.is_empty() => {
return ResultEnvelope::new(operation, refusal.outcome)
.with_instance(instance)
.with_change(refusal.change_id.clone())
.with_message(refusal.message.clone())
.with_detail(serde_json::json!({
"action": action.as_str(),
"commands_submitted": Vec::<serde_json::Value>::new(),
"targets": Vec::<serde_json::Value>::new(),
}))
}
Err(refusal) => {
return partial(
operation,
first_instance,
&settled,
&audit,
&refusal.message,
)
}
}
}
let mut revision = observation.state_revision;
let mut stale = false;
for (change_id, plan) in plans {
if plan == Plan::Satisfied {
settled.push(TargetResult {
change_id: change_id.clone(),
changed: false,
reason: "execution mark already had the requested value".to_string(),
});
continue;
}
match submit_and_settle(
connection,
CommandSpec::SetExecutionMark {
change_id: change_id.clone(),
marked: desired,
},
"set_execution_mark",
Some(change_id),
revision,
&observation.instance_id,
&mut audit,
)
.await
{
Ok(record) => {
revision = record.result_revision.unwrap_or(revision);
let changed = matches!(record.state, CommandState::Succeeded);
settled.push(TargetResult {
change_id: change_id.clone(),
changed,
reason: record.detail.clone().unwrap_or_else(|| match changed {
true => "the execution mark was updated".to_string(),
false => "the owner settled the request unchanged".to_string(),
}),
});
}
Err(SubmitFailure::Stale { .. }) => {
stale = true;
break;
}
Err(SubmitFailure::Restarted(observed)) => {
let mut envelope =
restarted_envelope(operation, &observation.instance_id, &observed);
if let Some(detail) = envelope.detail.as_object_mut() {
detail.insert("action".to_string(), serde_json::json!(action.as_str()));
detail.insert("commands_submitted".to_string(), serde_json::json!(audit));
detail.insert("stopped_at".to_string(), serde_json::json!(change_id));
}
return envelope;
}
Err(failure) if settled.is_empty() => {
return failure
.into_envelope(operation, instance)
.with_change(change_id.clone())
.with_detail(serde_json::json!({
"action": action.as_str(),
"commands_submitted": audit,
"targets": Vec::<serde_json::Value>::new(),
}))
}
Err(failure) => {
return partial(
operation,
first_instance,
&settled,
&audit,
failure.message(),
)
}
}
}
if stale {
continue;
}
return succeeded(action, first_instance, &settled, &audit);
}
if settled.is_empty() {
return revision_conflict(operation, first_instance);
}
partial(
operation,
first_instance,
&settled,
&audit,
format!(
"the owner's state advanced past every observation this client made; \
{MAX_REVISION_ATTEMPTS} bounded recomputations were exhausted"
),
)
}
fn succeeded(
action: Action,
instance: Option<String>,
settled: &[TargetResult],
audit: &[serde_json::Value],
) -> ResultEnvelope {
let changed = settled.iter().filter(|result| result.changed).count();
let outcome = match (changed > 0, action) {
(false, _) => Outcome::Unchanged,
(true, Action::Mark) => Outcome::Marked,
(true, Action::Unmark) => Outcome::Unmarked,
(true, _) => Outcome::Accepted,
};
let envelope = ResultEnvelope::new(action.operation(), outcome)
.with_instance(instance)
.with_message(format!(
"{changed} of {} named proposals changed; the owner's own settlement may later \
admit stable marked work, and nothing here waited for it",
settled.len()
))
.with_detail(serde_json::json!({
"action": action.as_str(),
"commands_submitted": audit,
"targets": settled.iter().map(TargetResult::to_json).collect::<Vec<_>>(),
}));
match settled {
[only] => envelope.with_change(only.change_id.clone()),
_ => envelope,
}
}
fn partial(
operation: Operation,
instance: Option<String>,
settled: &[TargetResult],
audit: &[serde_json::Value],
reason: impl Into<String>,
) -> ResultEnvelope {
ResultEnvelope::new(operation, Outcome::PartialIntent)
.with_instance(instance)
.with_message(format!(
"{} of the named proposals settled before the request stopped: {}. Nothing was \
rolled back and the settled marks stand",
settled.len(),
reason.into()
))
.with_detail(serde_json::json!({
"commands_submitted": audit,
"targets": settled.iter().map(TargetResult::to_json).collect::<Vec<_>>(),
"rolled_back": false,
}))
}
fn restarted_envelope(operation: Operation, expected: &str, observed: &str) -> ResultEnvelope {
ResultEnvelope::new(operation, Outcome::OwnerRestarted)
.with_instance(Some(observed.to_string()))
.with_message(
"the socket began serving a different owner incarnation, which cannot prove whether \
a command submitted to the previous one settled",
)
.with_detail(serde_json::json!({
"expected_instance_id": expected,
"observed_instance_id": observed,
}))
}
fn not_command_capable(operation: Operation, instance: Option<String>) -> ResultEnvelope {
ResultEnvelope::new(operation, Outcome::OwnerNotCommandCapable)
.with_instance(instance)
.with_message(
"this owner serves reads but has no command executor bound, so it can run no \
control command. A headless `cflx run` is read-only for control purposes",
)
}
fn revision_conflict(operation: Operation, instance: Option<String>) -> ResultEnvelope {
ResultEnvelope::new(operation, Outcome::RevisionConflict)
.with_instance(instance)
.with_message(format!(
"the owner's state advanced past every observation this client made; \
{MAX_REVISION_ATTEMPTS} bounded recomputations were exhausted without a settled \
command"
))
}
#[derive(Debug)]
enum SubmitFailure {
Stale {
current: Option<u64>,
},
ExecutorUnbound(String),
Ineligible(String),
Failed(String),
Restarted(String),
Unauthenticated(String),
Transport(String),
}
impl SubmitFailure {
fn message(&self) -> String {
match self {
Self::Stale { current } => match current {
Some(current) => format!(
"the observed revision was already consumed; the owner is now at revision \
{current}"
),
None => "the observed revision was already consumed".to_string(),
},
Self::Restarted(observed) => format!(
"the command record belongs to owner incarnation '{observed}', so this socket \
began serving a different process"
),
Self::ExecutorUnbound(detail)
| Self::Ineligible(detail)
| Self::Failed(detail)
| Self::Unauthenticated(detail)
| Self::Transport(detail) => detail.clone(),
}
}
fn outcome(&self) -> Outcome {
match self {
Self::Stale { .. } => Outcome::RevisionConflict,
Self::ExecutorUnbound(_) => Outcome::OwnerNotCommandCapable,
Self::Ineligible(_) => Outcome::TargetIneligible,
Self::Failed(_) => Outcome::CommandFailed,
Self::Restarted(_) => Outcome::OwnerRestarted,
Self::Unauthenticated(_) => Outcome::AuthenticationFailed,
Self::Transport(_) => Outcome::TransportError,
}
}
fn into_envelope(self, operation: Operation, instance: Option<String>) -> ResultEnvelope {
let outcome = self.outcome();
let message = self.message();
ResultEnvelope::new(operation, outcome)
.with_instance(instance)
.with_message(message)
}
}
impl From<TransportError> for SubmitFailure {
fn from(error: TransportError) -> Self {
Self::Transport(error.to_string())
}
}
#[allow(clippy::too_many_arguments)]
async fn submit_and_settle(
connection: &Connection,
command: CommandSpec,
name: &str,
change_id: Option<&str>,
expected_revision: u64,
expected_instance: &str,
audit: &mut Vec<serde_json::Value>,
) -> Result<CommandRecord, SubmitFailure> {
let request = CommandRequest {
command,
expected_revision,
idempotency_key: format!("cflx-client-{}", new_hex_id()),
correlation_id: Some(format!("cflx-client-{}", &new_hex_id()[..16])),
};
let body = serde_json::to_string(&request).map_err(|error| {
SubmitFailure::Failed(format!("command envelope is not encodable: {error}"))
})?;
let response = connection
.client()
.post_json("/api/v2/commands", &body)
.await?;
if matches!(response.status, 401 | 403) {
return Err(SubmitFailure::Unauthenticated(describe_api_error(
&response.body,
"the owner refused the presented credentials",
)));
}
if let Ok(record) = response.json::<CommandRecord>() {
let mut entry = serde_json::json!({
"command": name,
"command_id": record.command_id,
});
if let (Some(object), Some(change_id)) = (entry.as_object_mut(), change_id) {
object.insert("change_id".to_string(), serde_json::json!(change_id));
}
audit.push(entry);
return settle(connection, record, expected_instance).await;
}
let error: crate::web::remote_control_api::dto::ApiError = response.json().map_err(|_| {
SubmitFailure::Failed(describe_api_error(
&response.body,
"the owner refused the command without a typed error",
))
})?;
Err(match error.error_code {
ErrorCode::StaleRevision => SubmitFailure::Stale {
current: error.current_revision,
},
ErrorCode::CommandExecutorUnbound => SubmitFailure::ExecutorUnbound(error.message),
ErrorCode::LifecycleConflict
| ErrorCode::TargetIneligible
| ErrorCode::RootBusy
| ErrorCode::NotFound => SubmitFailure::Ineligible(error.message),
_ => SubmitFailure::Failed(format!("{} ({})", error.message, error.error_code.as_str())),
})
}
async fn settle(
connection: &Connection,
record: CommandRecord,
expected_instance: &str,
) -> Result<CommandRecord, SubmitFailure> {
let mut record = record;
let deadline = tokio::time::Instant::now() + SETTLEMENT_TIMEOUT;
loop {
if record.instance_id != expected_instance {
return Err(SubmitFailure::Restarted(record.instance_id.clone()));
}
match record.state {
CommandState::Succeeded | CommandState::NoOp => return Ok(record),
CommandState::Failed => {
let code = record
.error_code
.map(|code| code.as_str())
.unwrap_or("unspecified");
let detail = record.detail.clone().unwrap_or_default();
return Err(match record.error_code {
Some(ErrorCode::CommandExecutorUnbound) => {
SubmitFailure::ExecutorUnbound(detail)
}
Some(
ErrorCode::LifecycleConflict
| ErrorCode::TargetIneligible
| ErrorCode::RootBusy,
) => SubmitFailure::Ineligible(detail),
Some(ErrorCode::StaleRevision) => SubmitFailure::Stale { current: None },
_ => SubmitFailure::Failed(format!("{detail} ({code})")),
});
}
CommandState::Running => {}
}
if tokio::time::Instant::now() >= deadline {
return Err(SubmitFailure::Failed(format!(
"command '{}' was still running after {SETTLEMENT_TIMEOUT:?}; its effect is \
unknown and no second command was submitted",
record.command_id
)));
}
tokio::time::sleep(SETTLEMENT_POLL).await;
let path = format!("/api/v2/commands/{}", record.command_id);
let response = connection.client().get(&path).await?;
if response.status != 200 {
return Err(SubmitFailure::Failed(describe_api_error(
&response.body,
"the command record could not be read back",
)));
}
record = response.json().map_err(|error| {
SubmitFailure::Failed(format!("the command record was not usable: {error}"))
})?;
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::web::remote_control_api::dto::{
ActionEligibility, AttentionState, ChangeActions, ChangeResource, ChangeTiming,
ParallelEligibility, QueueIntent,
};
fn change(id: &str, display_status: &str, marked: bool) -> ChangeResource {
ChangeResource {
id: id.to_string(),
display_status: display_status.to_string(),
progress_status: "pending".to_string(),
completed_tasks: 0,
total_tasks: 3,
progress_percent: 0.0,
dependencies: Vec::new(),
iteration_number: None,
execution_marked: marked,
queue_intent: QueueIntent::NotQueued,
attention: AttentionState::None,
blocker: None,
error_detail: None,
actions: crate::web::remote_control_api::projection::change_actions_for_test(
"select",
display_status,
None,
),
parallel: ParallelEligibility::default(),
timing: ChangeTiming::default(),
latest_activity: None,
worktree: None,
}
}
fn observation(changes: Vec<ChangeResource>) -> Observation {
crate::client::session::observation_for_test(changes)
}
#[test]
fn every_action_round_trips_through_its_wire_name() {
for action in [
Action::Mark,
Action::Unmark,
Action::Start,
Action::Stop,
Action::ForceStop,
] {
assert_eq!(Action::parse(action.as_str()), Some(action));
}
for retired in ["enqueue", "queue", "retry", "admit", "set_queue_intent"] {
assert_eq!(
Action::parse(retired),
None,
"{retired} must not be an action"
);
}
}
#[test]
fn only_mark_actions_address_named_proposals() {
assert!(Action::Mark.is_mark());
assert!(Action::Unmark.is_mark());
for action in [Action::Start, Action::Stop, Action::ForceStop] {
assert!(!action.is_mark());
assert_eq!(action.desired_mark(), None);
assert!(action.lifecycle_command().is_some());
}
assert_eq!(Action::Mark.desired_mark(), Some(true));
assert_eq!(Action::Unmark.desired_mark(), Some(false));
assert!(Action::Mark.lifecycle_command().is_none());
}
#[test]
fn lifecycle_actions_map_onto_the_shared_run_control_commands() {
assert_eq!(Action::Start.lifecycle_command(), Some(CommandSpec::Start));
assert_eq!(Action::Stop.lifecycle_command(), Some(CommandSpec::Stop));
assert_eq!(
Action::ForceStop.lifecycle_command(),
Some(CommandSpec::ForceStop)
);
}
#[test]
fn a_target_list_is_bounded_non_empty_and_distinct() {
assert!(validate_targets(&["alpha".to_string()]).is_ok());
assert!(validate_targets(&[]).is_err());
let duplicated = ["alpha", "beta", "alpha"].map(str::to_string).to_vec();
let error = validate_targets(&duplicated).expect_err("a duplicate is refused");
assert!(error.contains("more than once"), "{error}");
let at_limit: Vec<String> = (0..MAX_TARGETS).map(|n| format!("change-{n}")).collect();
assert!(validate_targets(&at_limit).is_ok());
let over: Vec<String> = (0..MAX_TARGETS + 1)
.map(|n| format!("change-{n}"))
.collect();
assert!(validate_targets(&over).is_err());
}
#[test]
fn a_target_already_in_the_desired_state_needs_no_command() {
let observed = observation(vec![change("alpha", "not queued", true)]);
assert_eq!(classify(&observed, "alpha", true).unwrap(), Plan::Satisfied);
assert_eq!(classify(&observed, "alpha", false).unwrap(), Plan::Submit);
}
#[test]
fn an_unknown_proposal_refuses_the_whole_request() {
let observed = observation(vec![change("alpha", "not queued", false)]);
let refusal = classify(&observed, "missing", true).expect_err("unknown");
assert_eq!(refusal.outcome, Outcome::ChangeNotFound);
assert_eq!(refusal.change_id, "missing");
}
#[test]
fn a_terminal_target_is_submitted_and_left_to_the_shared_no_op() {
for status in ["archived", "merged", "pushed", "rejected"] {
let observed = observation(vec![change("alpha", status, false)]);
assert!(
!observed
.change("alpha")
.unwrap()
.actions
.set_execution_mark
.allowed
);
assert_eq!(
classify(&observed, "alpha", true).unwrap(),
Plan::Submit,
"{status} must reach the shared service"
);
}
}
#[test]
fn a_mode_level_mark_refusal_stops_the_request_before_submission() {
let mut blocked = change("alpha", "not queued", false);
blocked.actions = ChangeActions {
set_execution_mark: ActionEligibility::blocked(ActionBlockedReason::StopPending),
..blocked.actions
};
let observed = observation(vec![blocked]);
let refusal = classify(&observed, "alpha", true).expect_err("mode refuses marking");
assert_eq!(refusal.outcome, Outcome::TargetIneligible);
assert!(refusal.message.contains("stop_pending"), "{refusal:?}");
assert!(
refusal
.message
.contains("no proposal in this request was marked"),
"{refusal:?}"
);
}
#[test]
fn a_request_that_moved_nothing_is_an_unchanged_success() {
let settled = vec![TargetResult {
change_id: "alpha".to_string(),
changed: false,
reason: "execution mark already had the requested value".to_string(),
}];
let envelope = succeeded(Action::Mark, Some("i-1".to_string()), &settled, &[]);
assert_eq!(envelope.outcome, Outcome::Unchanged);
assert!(envelope.ok);
assert_eq!(envelope.exit_code(), 0);
assert_eq!(envelope.change_id.as_deref(), Some("alpha"));
assert!(envelope.execution_id.is_none());
}
#[test]
fn a_multi_target_success_reports_each_target_and_names_no_single_change() {
let settled = vec![
TargetResult {
change_id: "alpha".to_string(),
changed: true,
reason: "the execution mark was updated".to_string(),
},
TargetResult {
change_id: "gamma".to_string(),
changed: false,
reason: "execution mark already had the requested value".to_string(),
},
];
let envelope = succeeded(Action::Mark, None, &settled, &[]);
assert_eq!(envelope.outcome, Outcome::Marked);
assert!(envelope.change_id.is_none());
assert_eq!(envelope.detail["targets"][0]["change_id"], "alpha");
assert_eq!(envelope.detail["targets"][0]["changed"], true);
assert_eq!(envelope.detail["targets"][1]["changed"], false);
let rendered = envelope.to_json_line();
for forbidden in ["queue_intent", "admitted", "execution_id"] {
assert!(!rendered.contains(forbidden), "{forbidden}: {rendered}");
}
}
#[test]
fn unmark_reports_its_own_success_token() {
let settled = vec![TargetResult {
change_id: "alpha".to_string(),
changed: true,
reason: "the execution mark was updated".to_string(),
}];
assert_eq!(
succeeded(Action::Unmark, None, &settled, &[]).outcome,
Outcome::Unmarked
);
}
#[test]
fn partial_intent_lists_created_records_in_order_without_claiming_rollback() {
let settled = vec![
TargetResult {
change_id: "alpha".to_string(),
changed: true,
reason: "the execution mark was updated".to_string(),
},
TargetResult {
change_id: "beta".to_string(),
changed: true,
reason: "the execution mark was updated".to_string(),
},
];
let audit = vec![
serde_json::json!({"command": "set_execution_mark", "command_id": "c-1", "change_id": "alpha"}),
serde_json::json!({"command": "set_execution_mark", "command_id": "c-2", "change_id": "beta"}),
serde_json::json!({"command": "set_execution_mark", "command_id": "c-3", "change_id": "gamma"}),
];
let envelope = partial(
Operation::ControlMark,
Some("i-1".to_string()),
&settled,
&audit,
"gamma was refused",
);
assert_eq!(envelope.outcome, Outcome::PartialIntent);
assert!(!envelope.ok);
assert_eq!(envelope.detail["rolled_back"], false);
let commands = envelope.detail["commands_submitted"].as_array().unwrap();
assert_eq!(commands.len(), 3);
assert_eq!(commands[0]["change_id"], "alpha");
assert_eq!(commands[2]["change_id"], "gamma");
assert!(envelope.message.as_ref().unwrap().contains("rolled back"));
}
#[test]
fn submission_failures_map_to_distinct_stable_outcomes() {
assert_eq!(
SubmitFailure::ExecutorUnbound(String::new()).outcome(),
Outcome::OwnerNotCommandCapable
);
assert_eq!(
SubmitFailure::Ineligible(String::new()).outcome(),
Outcome::TargetIneligible
);
assert_eq!(
SubmitFailure::Failed(String::new()).outcome(),
Outcome::CommandFailed
);
assert_eq!(
SubmitFailure::Restarted(String::new()).outcome(),
Outcome::OwnerRestarted
);
assert_eq!(
SubmitFailure::Stale { current: Some(4) }.outcome(),
Outcome::RevisionConflict
);
}
#[test]
fn a_lifecycle_action_refuses_a_target_list_before_contact() {
for action in [Action::Start, Action::Stop, Action::ForceStop] {
assert!(validate_request(action, &[]).is_ok(), "{}", action.as_str());
let error = validate_request(action, &["alpha".to_string()])
.expect_err("a lifecycle target list is refused");
assert!(error.contains("authoritative mark set"), "{error}");
}
for action in [Action::Mark, Action::Unmark] {
assert!(validate_request(action, &["alpha".to_string()]).is_ok());
assert!(validate_request(action, &[]).is_err());
}
assert_eq!(
ResultEnvelope::new(Operation::ControlStart, Outcome::UsageError).exit_code(),
2
);
}
}