use crate::client::envelope::{Operation, Outcome, ResultEnvelope};
use crate::client::session::{
describe_api_error, observe_bounded, Connection, MAX_RECONCILE_ATTEMPTS,
};
use crate::client::transport::{encode_query_value, HttpResponse, TransportError};
use crate::web::remote_control_api::dto::{
ApiError, ErrorCode, ProposalSubscriptionRequest, ProposalSubscriptionResponse,
};
#[derive(Debug, Clone)]
pub enum Intent {
Set {
command: Vec<String>,
notify_blocked: bool,
},
Get,
Clear,
}
impl Intent {
pub fn operation(&self) -> Operation {
match self {
Self::Set { .. } => Operation::SubscribeSet,
Self::Get => Operation::SubscribeGet,
Self::Clear => Operation::SubscribeClear,
}
}
pub fn action(&self) -> &'static str {
match self {
Self::Set { .. } => "set",
Self::Get => "get",
Self::Clear => "clear",
}
}
fn success(&self) -> Outcome {
match self {
Self::Clear => Outcome::Cleared,
_ => Outcome::Subscribed,
}
}
}
pub fn validate_request(change_ids: &[String], intent: &Intent) -> Result<(), String> {
crate::client::control::validate_targets(change_ids)?;
match intent {
Intent::Set { command, .. } => {
crate::web::completion_sink::validate_command(command).map_err(
|refusal| match refusal {
crate::web::completion_sink::SinkRefusal::InvalidCommand(detail) => detail,
other => format!("the callback argv is not acceptable: {other:?}"),
},
)
}
Intent::Get | Intent::Clear => Ok(()),
}
}
pub async fn run(
connection: &Connection,
change_ids: &[String],
expected_instance: Option<&str>,
intent: Intent,
) -> ResultEnvelope {
let operation = intent.operation();
if let Err(message) = validate_request(change_ids, &intent) {
return ResultEnvelope::new(operation, Outcome::UsageError).with_message(message);
}
let observation = match observe_bounded(connection, None, MAX_RECONCILE_ATTEMPTS).await {
Ok(observation) => observation,
Err(error) => return error.into_envelope(operation),
};
let instance_id = observation.instance_id.clone();
if let Some(expected) = expected_instance {
if expected != instance_id {
return ResultEnvelope::new(operation, Outcome::OwnerRestarted)
.with_instance(Some(instance_id.clone()))
.with_message(
"the socket is serving a different owner incarnation, so subscriptions \
registered with the one you named no longer exist. It is not completion: \
nothing about any proposal was observed",
)
.with_detail(serde_json::json!({
"expected_instance_id": expected,
"observed_instance_id": instance_id,
}));
}
}
if !observation.proposal_subscriptions_available() {
return ResultEnvelope::new(operation, Outcome::UnsupportedOwner)
.with_instance(Some(instance_id))
.with_message(
"this owner does not serve proposal-scoped subscriptions, so no callback can be \
registered with it. Observe completion with `cflx client wait` instead",
);
}
let mut settled: Vec<serde_json::Value> = Vec::new();
let mut episode: Option<String> = None;
for change_id in change_ids {
let response = request(connection, change_id, &instance_id, &intent).await;
let response = match response {
Ok(response) => response,
Err(error) => {
let outcome = match error {
TransportError::NotListening { .. } => Outcome::OwnerNotRunning,
_ => Outcome::TransportError,
};
return stop(
operation,
&instance_id,
change_id,
&settled,
outcome,
error.to_string(),
);
}
};
match classify(response, &instance_id) {
Ok(body) => {
episode = body.execution_id.clone();
settled.push(serde_json::to_value(&body).unwrap_or_default())
}
Err(refusal) => {
return stop(
operation,
&instance_id,
change_id,
&settled,
refusal.outcome,
refusal.message,
)
}
}
}
let envelope = ResultEnvelope::new(operation, intent.success())
.with_instance(Some(instance_id))
.with_message(match &intent {
Intent::Set { .. } => format!(
"{} proposal subscriptions are registered; delivery is notification only and \
resumes no agent",
settled.len()
),
Intent::Get => format!("{} proposal subscriptions were read", settled.len()),
Intent::Clear => format!(
"{} proposal subscriptions were removed; a callback already running keeps its \
own bounds",
settled.len()
),
})
.with_detail(serde_json::json!({
"action": intent.action(),
"subscriptions": settled,
}));
match change_ids {
[only] => envelope.with_change(only.clone()).with_execution(episode),
_ => envelope,
}
}
fn stop(
operation: Operation,
instance_id: &str,
change_id: &str,
settled: &[serde_json::Value],
outcome: Outcome,
message: String,
) -> ResultEnvelope {
if settled.is_empty() {
return ResultEnvelope::new(operation, outcome)
.with_instance(Some(instance_id.to_string()))
.with_change(change_id)
.with_message(message)
.with_detail(serde_json::json!({ "subscriptions": settled }));
}
ResultEnvelope::new(operation, Outcome::PartialIntent)
.with_instance(Some(instance_id.to_string()))
.with_change(change_id)
.with_message(format!(
"{} proposals settled before '{change_id}' stopped the request: {message}. Nothing \
was rolled back",
settled.len()
))
.with_detail(serde_json::json!({
"subscriptions": settled,
"stopped_at": change_id,
"rolled_back": false,
}))
}
async fn request(
connection: &Connection,
change_id: &str,
instance_id: &str,
intent: &Intent,
) -> Result<HttpResponse, TransportError> {
let path = format!(
"/api/v2/proposals/{}/subscription",
encode_query_value(change_id)
);
let binding = format!("instance_id={}", encode_query_value(instance_id));
match intent {
Intent::Set {
command,
notify_blocked,
} => {
let body = serde_json::to_string(&ProposalSubscriptionRequest {
instance_id: instance_id.to_string(),
command: command.clone(),
notify_blocked: *notify_blocked,
})
.expect("a closed struct of owned values is encodable");
connection.client().put_json(&path, &body).await
}
Intent::Get => connection.client().get(&format!("{path}?{binding}")).await,
Intent::Clear => {
connection
.client()
.delete(&format!("{path}?{binding}"))
.await
}
}
}
#[derive(Debug)]
struct Refusal {
outcome: Outcome,
message: String,
}
fn classify(
response: HttpResponse,
instance_id: &str,
) -> Result<ProposalSubscriptionResponse, Refusal> {
if response.status == 200 {
return match response.json::<ProposalSubscriptionResponse>() {
Ok(body) => {
if body.instance_id != instance_id {
return Err(Refusal {
outcome: Outcome::OwnerRestarted,
message: "the socket began serving a different owner incarnation while \
the subscription was being registered"
.to_string(),
});
}
Ok(body)
}
Err(error) => Err(Refusal {
outcome: Outcome::IncompatibleOwner,
message: format!(
"the owner's subscription response did not match this build's contract: \
{error}"
),
}),
};
}
let code = serde_json::from_slice::<ApiError>(&response.body).map(|error| error.error_code);
let message = describe_api_error(&response.body, "the owner refused the subscription request");
let outcome = match (response.status, code) {
(401 | 403, Ok(ErrorCode::TransportNotPermitted)) => Outcome::TransportNotPermitted,
(401 | 403, _) => Outcome::AuthenticationFailed,
(404, _) => Outcome::UnsupportedOwner,
(409, Ok(ErrorCode::ExecutionBindingMismatch)) => Outcome::OwnerRestarted,
(409, _) => Outcome::TargetIneligible,
(422, _) => Outcome::UsageError,
(503, _) => Outcome::UnsupportedOwner,
_ => Outcome::TransportError,
};
Err(Refusal { outcome, message })
}
#[cfg(test)]
mod tests {
use super::*;
fn response(status: u16, body: serde_json::Value) -> HttpResponse {
HttpResponse {
status,
body: serde_json::to_vec(&body).unwrap(),
}
}
fn api_error(code: &str) -> serde_json::Value {
serde_json::json!({
"error_code": code,
"message": "refused",
"correlation_id": "abc",
})
}
fn set(command: &[&str]) -> Intent {
Intent::Set {
command: command.iter().map(|part| part.to_string()).collect(),
notify_blocked: false,
}
}
#[test]
fn each_intent_reports_its_own_operation_and_success_token() {
assert_eq!(Intent::Get.operation(), Operation::SubscribeGet);
assert_eq!(Intent::Get.success(), Outcome::Subscribed);
assert_eq!(set(&["/bin/true"]).operation(), Operation::SubscribeSet);
assert_eq!(set(&["/bin/true"]).success(), Outcome::Subscribed);
assert_eq!(Intent::Clear.operation(), Operation::SubscribeClear);
assert_eq!(Intent::Clear.success(), Outcome::Cleared);
assert_eq!(Intent::Clear.action(), "clear");
}
#[test]
fn a_request_is_validated_completely_before_any_owner_is_contacted() {
let targets = ["alpha".to_string(), "beta".to_string()];
assert!(validate_request(&targets, &set(&["/absolute/callback", "--flag"])).is_ok());
assert!(validate_request(&[], &Intent::Get).is_err());
let duplicated = ["alpha", "alpha"].map(str::to_string).to_vec();
assert!(validate_request(&duplicated, &Intent::Clear).is_err());
let too_many: Vec<String> = (0..65).map(|n| format!("change-{n}")).collect();
assert!(validate_request(&too_many, &Intent::Get).is_err());
let empty = validate_request(&targets, &set(&[])).expect_err("an empty argv is refused");
assert!(empty.contains("program"), "{empty}");
assert!(validate_request(&targets, &Intent::Get).is_ok());
assert!(validate_request(&targets, &Intent::Clear).is_ok());
}
#[test]
fn a_typed_refusal_maps_onto_its_own_stable_outcome() {
let cases = [
(
403,
"transport_not_permitted",
Outcome::TransportNotPermitted,
),
(403, "forbidden", Outcome::AuthenticationFailed),
(401, "unauthorized", Outcome::AuthenticationFailed),
(404, "not_found", Outcome::UnsupportedOwner),
(409, "execution_binding_mismatch", Outcome::OwnerRestarted),
(422, "validation_failed", Outcome::UsageError),
(503, "command_executor_unbound", Outcome::UnsupportedOwner),
];
for (status, code, expected) in cases {
let refusal =
classify(response(status, api_error(code)), "i-1").expect_err("{status} {code}");
assert_eq!(refusal.outcome, expected, "{status} {code}");
}
}
#[test]
fn a_successful_registration_carries_the_owner_binding_and_the_episode() {
let body = serde_json::json!({
"instance_id": "i-1",
"change_id": "alpha",
"sink": {"command": ["/bin/true"], "notify_blocked": false},
"subscribed": true,
"execution_id": "e-1",
"execution_state": "active",
"terminal_dispatched": false,
"delivered_events": [],
});
let parsed = classify(response(200, body), "i-1").expect("a registered subscription");
assert_eq!(parsed.change_id, "alpha");
assert!(parsed.subscribed);
assert_eq!(parsed.execution_id.as_deref(), Some("e-1"));
assert_eq!(parsed.sink.unwrap().command, vec!["/bin/true".to_string()]);
}
#[test]
fn a_subscription_before_admission_reports_no_episode_and_still_succeeds() {
let body = serde_json::json!({
"instance_id": "i-1",
"change_id": "alpha",
"sink": {"command": ["/bin/true"], "notify_blocked": true},
"subscribed": true,
"execution_state": "unknown",
"terminal_dispatched": false,
"delivered_events": [],
});
let parsed = classify(response(200, body), "i-1").expect("a pre-admission subscription");
assert!(parsed.subscribed);
assert_eq!(parsed.execution_id, None);
assert!(!parsed.terminal_dispatched);
}
#[test]
fn an_answer_from_another_incarnation_is_reported_as_a_restart() {
let body = serde_json::json!({
"instance_id": "i-2",
"change_id": "alpha",
"sink": serde_json::Value::Null,
"subscribed": false,
"execution_state": "unknown",
"terminal_dispatched": false,
"delivered_events": [],
});
let refusal = classify(response(200, body), "i-1").expect_err("a replaced owner");
assert_eq!(refusal.outcome, Outcome::OwnerRestarted);
}
#[test]
fn a_redacted_read_still_reports_presence() {
let body = serde_json::json!({
"instance_id": "i-1",
"change_id": "alpha",
"sink": serde_json::Value::Null,
"subscribed": true,
"execution_state": "queued",
"terminal_dispatched": false,
"delivered_events": [],
});
let parsed = classify(response(200, body), "i-1").expect("a redacted read");
assert!(parsed.subscribed, "presence is not the secret; the argv is");
assert!(parsed.sink.is_none());
}
#[test]
fn a_refusal_before_the_first_mutation_is_not_reported_as_partial() {
let envelope = stop(
Operation::SubscribeSet,
"i-1",
"alpha",
&[],
Outcome::TransportNotPermitted,
"TCP may not register an argv".to_string(),
);
assert_eq!(envelope.outcome, Outcome::TransportNotPermitted);
assert_eq!(envelope.change_id.as_deref(), Some("alpha"));
assert!(envelope.detail["subscriptions"]
.as_array()
.unwrap()
.is_empty());
}
#[test]
fn a_refusal_after_a_settled_registration_is_partial_and_claims_no_rollback() {
let settled = vec![serde_json::json!({"change_id": "alpha", "subscribed": true})];
let envelope = stop(
Operation::SubscribeSet,
"i-1",
"beta",
&settled,
Outcome::TransportError,
"the socket closed".to_string(),
);
assert_eq!(envelope.outcome, Outcome::PartialIntent);
assert_eq!(envelope.detail["rolled_back"], false);
assert_eq!(envelope.detail["stopped_at"], "beta");
assert_eq!(envelope.detail["subscriptions"][0]["change_id"], "alpha");
}
}