use std::time::Duration;
use axum::body::Bytes;
use axum::extract::{Path, State};
use axum::http::header;
use axum::response::{IntoResponse, Response};
use axum::{Extension, Json};
use chrono::Utc;
use super::auth::CorrelationId;
use super::dto::{
is_valid_correlation_id, ApiError, CommandRecord, CommandRequest, CommandState, ErrorCode,
};
use super::executor::{Applied, CommandFailure, ExecutionSummary, PendingCommand};
use super::projection::{Admission, Projection};
use super::RemoteControlState;
const SYNCHRONOUS_GRACE: Duration = Duration::from_millis(250);
fn no_store(status: axum::http::StatusCode, body: impl serde::Serialize) -> Response {
(status, [(header::CACHE_CONTROL, "no-store")], Json(body)).into_response()
}
fn record_response(record: &CommandRecord) -> Response {
no_store(record.http_status(), record)
}
#[utoipa::path(
post,
path = "/api/v2/commands",
tag = "remote-control",
request_body = CommandRequest,
responses(
(status = 200, description = "Completed synchronously, replayed, or an explicit no-op", body = CommandRecord),
(status = 202, description = "Accepted and still running", body = CommandRecord),
(status = 409, description = "Stale revision, lifecycle, eligibility, or idempotency conflict", body = ApiError),
(status = 422, description = "Schema or correlation-ID validation failed", body = ApiError),
(status = 503, description = "No record slot could be reserved", body = ApiError)
)
)]
pub async fn submit_command(
State(state): State<RemoteControlState>,
Extension(correlation): Extension<CorrelationId>,
body: Bytes,
) -> Response {
let correlation_id = correlation.0;
let request: CommandRequest = match serde_json::from_slice(&body) {
Ok(request) => request,
Err(error) => {
return ApiError::new(
ErrorCode::ValidationFailed,
format!("command envelope failed typed validation: {error}"),
&correlation_id,
)
.into_response()
}
};
if request.idempotency_key.is_empty() || request.idempotency_key.len() > 200 {
return ApiError::new(
ErrorCode::ValidationFailed,
"idempotency_key must be 1-200 characters",
&correlation_id,
)
.into_response();
}
if let Some(supplied) = request.correlation_id.as_deref() {
if !is_valid_correlation_id(supplied) {
return ApiError::new(
ErrorCode::ValidationFailed,
"correlation_id must be 1-64 characters matching [A-Za-z0-9._:-]",
&correlation_id,
)
.into_response();
}
}
let correlation_id = request.correlation_id.clone().unwrap_or(correlation_id);
match state.projection.resolve_replay(&request, Utc::now()) {
Some(Admission::Replay(record)) => return record_response(&record),
Some(Admission::IdempotencyMismatch) => {
return ApiError::new(
ErrorCode::IdempotencyMismatch,
"idempotency_key is already bound to a different command identity",
&correlation_id,
)
.with_revision(state.projection.revision())
.into_response()
}
_ => {}
}
let gate = state.gate.hold().await;
let record = match state
.projection
.admit(&request, &correlation_id, Utc::now())
{
Admission::Replay(record) => return record_response(&record),
Admission::IdempotencyMismatch => {
return ApiError::new(
ErrorCode::IdempotencyMismatch,
"idempotency_key is already bound to a different command identity",
&correlation_id,
)
.with_revision(state.projection.revision())
.into_response()
}
Admission::Stale(current) => {
return ApiError::new(
ErrorCode::StaleRevision,
format!("expected_revision {} is stale", request.expected_revision),
&correlation_id,
)
.with_revision(current)
.into_response()
}
Admission::Capacity => {
return ApiError::new(
ErrorCode::RegistryCapacity,
"no command slot could be reserved without evicting in-progress work",
&correlation_id,
)
.with_revision(state.projection.revision())
.into_response()
}
Admission::Admitted(record) => *record,
};
let command_id = record.command_id.clone();
let projection = state.projection.clone();
let settlement: Option<PendingCommand> =
match state.executor.begin(&request.command, gate).await {
Applied::Ordinary(gate) => {
let executor = state.executor.clone();
let spec = request.command.clone();
Some(Box::pin(
async move { executor.execute_held(&spec, gate).await },
))
}
Applied::Pending(pending) => Some(pending),
Applied::Settled(result) => {
settle(&projection, &command_id, result);
None
}
};
if let Some(pending) = settlement {
let settle_id = command_id.clone();
let (done_tx, done_rx) = tokio::sync::oneshot::channel();
tokio::spawn(async move {
let result = pending.await;
settle(&projection, &settle_id, result);
let _ = done_tx.send(());
});
let _ = tokio::time::timeout(SYNCHRONOUS_GRACE, done_rx).await;
}
match state.projection.command(&command_id) {
Some(settled) => record_response(&settled),
None => record_response(&record),
}
}
fn settle(
projection: &Projection,
command_id: &str,
result: Result<ExecutionSummary, CommandFailure>,
) {
let (command_state, detail, error_code, revision, typed_result) = match result {
Ok(summary) if summary.changed => (
CommandState::Succeeded,
summary.detail,
None,
summary.result_revision,
summary.result,
),
Ok(summary) => (
CommandState::NoOp,
summary.detail,
None,
summary.result_revision,
summary.result,
),
Err(failure) => (
CommandState::Failed,
Some(failure.message),
Some(failure.error_code),
failure.result_revision,
None,
),
};
projection.complete_command(
command_id,
command_state,
detail,
error_code,
revision,
typed_result,
);
}
#[utoipa::path(
get,
path = "/api/v2/commands/{command_id}",
tag = "remote-control",
params(("command_id" = String, Path, description = "Command ID")),
responses(
(status = 200, description = "Command record", body = CommandRecord),
(status = 404, description = "No such command in this incarnation", body = ApiError)
)
)]
pub async fn get_command(
State(state): State<RemoteControlState>,
Extension(correlation): Extension<CorrelationId>,
Path(command_id): Path<String>,
) -> Response {
match state.projection.command(&command_id) {
Some(record) => no_store(axum::http::StatusCode::OK, record),
None => ApiError::new(
ErrorCode::NotFound,
format!("command '{command_id}' is not known to this instance"),
&correlation.0,
)
.into_response(),
}
}