use axum::extract::{Path, Query, State};
use axum::response::{IntoResponse, Response};
use axum::{Extension, Json};
use super::auth::CorrelationId;
use super::dto::{
ApiError, ErrorCode, ExecutionSinkSpec, ProposalSubscriptionParams,
ProposalSubscriptionRequest, ProposalSubscriptionResponse,
};
use super::{ApiTransport, RemoteControlState};
use crate::web::completion_sink::{ProposalSubscriptionView, SinkRefusal};
#[utoipa::path(
get,
path = "/api/v2/proposals/{change_id}/subscription",
tag = "remote-control",
params(
("change_id" = String, Path, description = "Proposal the subscription is keyed by"),
("instance_id" = String, Query, description = "Owner incarnation the caller believes it is addressing. Required: inspection asserts the same complete binding a registration does"),
),
responses(
(status = 200, description = "Current subscription for this proposal. The registered argv is returned only to a request that arrived on the owner Unix socket; `subscribed` reports presence on either transport", body = ProposalSubscriptionResponse),
(status = 409, description = "The presented instance binding is not this owner incarnation", body = ApiError),
(status = 422, description = "The request omitted part of the binding", body = ApiError),
)
)]
pub async fn get_subscription(
State(state): State<RemoteControlState>,
Extension(CorrelationId(correlation)): Extension<CorrelationId>,
Path(change_id): Path<String>,
Query(params): Query<ProposalSubscriptionParams>,
transport: Option<Extension<ApiTransport>>,
) -> Response {
let Some(registry) = state.completion_sinks.get() else {
return unsupported(&correlation);
};
let Some(instance_id) = params.instance_id else {
return incomplete_binding("reading", &correlation);
};
match registry.view_proposal_subscription(&change_id, &instance_id) {
Ok(view) => body(&state, view, discloses_argv(&transport)).into_response(),
Err(refusal) => refusal_response(refusal, &change_id, &correlation),
}
}
#[utoipa::path(
put,
path = "/api/v2/proposals/{change_id}/subscription",
tag = "remote-control",
params(("change_id" = String, Path, description = "Proposal the subscription is keyed by")),
request_body = ProposalSubscriptionRequest,
responses(
(status = 200, description = "Subscription registered or replaced", body = ProposalSubscriptionResponse),
(status = 403, description = "This transport may not register an executable argv", body = ApiError),
(status = 409, description = "The presented instance binding is not this owner incarnation", body = ApiError),
(status = 422, description = "The argv is not an acceptable bounded command", body = ApiError),
)
)]
pub async fn put_subscription(
State(state): State<RemoteControlState>,
Extension(CorrelationId(correlation)): Extension<CorrelationId>,
Path(change_id): Path<String>,
transport: Option<Extension<ApiTransport>>,
Json(request): Json<ProposalSubscriptionRequest>,
) -> Response {
let Some(registry) = state.completion_sinks.get() else {
return unsupported(&correlation);
};
if let Some(refusal) = require_owner_socket(transport, &correlation) {
return refusal;
}
match registry.set_proposal_subscription(
&change_id,
&request.instance_id,
ExecutionSinkSpec {
command: request.command,
notify_blocked: request.notify_blocked,
},
) {
Ok(view) => body(&state, view, true).into_response(),
Err(refusal) => refusal_response(refusal, &change_id, &correlation),
}
}
#[utoipa::path(
delete,
path = "/api/v2/proposals/{change_id}/subscription",
tag = "remote-control",
params(
("change_id" = String, Path, description = "Proposal the subscription is keyed by"),
("instance_id" = String, Query, description = "Owner incarnation the caller believes it is addressing. Required"),
),
responses(
(status = 200, description = "Subscription cleared", body = ProposalSubscriptionResponse),
(status = 403, description = "This transport may not mutate a subscription", body = ApiError),
(status = 409, description = "The presented instance binding is not this owner incarnation", body = ApiError),
(status = 422, description = "The request omitted part of the binding", body = ApiError),
)
)]
pub async fn delete_subscription(
State(state): State<RemoteControlState>,
Extension(CorrelationId(correlation)): Extension<CorrelationId>,
Path(change_id): Path<String>,
Query(params): Query<ProposalSubscriptionParams>,
transport: Option<Extension<ApiTransport>>,
) -> Response {
let Some(registry) = state.completion_sinks.get() else {
return unsupported(&correlation);
};
if let Some(refusal) = require_owner_socket(transport, &correlation) {
return refusal;
}
let Some(instance_id) = params.instance_id else {
return incomplete_binding("clearing", &correlation);
};
match registry.clear_proposal_subscription(&change_id, &instance_id) {
Ok(view) => body(&state, view, true).into_response(),
Err(refusal) => refusal_response(refusal, &change_id, &correlation),
}
}
fn incomplete_binding(operation: &str, correlation: &str) -> Response {
ApiError::new(
ErrorCode::ValidationFailed,
format!("{operation} a proposal subscription requires the instance_id it is bound to"),
correlation,
)
.into_response()
}
fn discloses_argv(transport: &Option<Extension<ApiTransport>>) -> bool {
matches!(transport, Some(Extension(ApiTransport::Unix)))
}
fn require_owner_socket(
transport: Option<Extension<ApiTransport>>,
correlation: &str,
) -> Option<Response> {
match transport.map(|Extension(transport)| transport) {
Some(ApiTransport::Unix) => None,
_ => Some(
ApiError::new(
ErrorCode::TransportNotPermitted,
"a proposal subscription stores an argv this owner will execute, so it may be \
set or cleared only over the owner's Unix socket",
correlation,
)
.into_response(),
),
}
}
fn body(
state: &RemoteControlState,
view: ProposalSubscriptionView,
disclose_argv: bool,
) -> Json<ProposalSubscriptionResponse> {
let execution_state = state
.completion_sinks
.get()
.map(|registry| registry.execution_state(&view.change_id))
.unwrap_or(super::dto::ChangeExecutionState::Unknown);
let subscribed = view.sink.is_some();
Json(ProposalSubscriptionResponse {
instance_id: state.projection.instance_id().to_string(),
change_id: view.change_id,
sink: match disclose_argv {
true => view.sink,
false => None,
},
subscribed,
execution_id: view.execution_id,
execution_state,
terminal_dispatched: view.terminal_dispatched,
delivered_events: view.delivered_events,
})
}
fn refusal_response(refusal: SinkRefusal, change_id: &str, correlation: &str) -> Response {
match refusal {
SinkRefusal::InstanceMismatch => ApiError::new(
ErrorCode::ExecutionBindingMismatch,
"the presented instance_id is not this owner incarnation, so a subscription for \
'{change_id}' would belong to a process this caller never observed"
.replace("{change_id}", change_id),
correlation,
)
.into_response(),
SinkRefusal::InvalidCommand(detail) => {
ApiError::new(ErrorCode::ValidationFailed, detail, correlation).into_response()
}
SinkRefusal::UnknownExecution | SinkRefusal::BindingMismatch { .. } => ApiError::new(
ErrorCode::ExecutionBindingMismatch,
format!("the subscription binding presented for '{change_id}' is not this owner's"),
correlation,
)
.into_response(),
}
}
fn unsupported(correlation: &str) -> Response {
ApiError::new(
ErrorCode::CommandExecutorUnbound,
"this process has no completion-sink registry bound, so no proposal subscription can be \
held",
correlation,
)
.into_response()
}