use std::sync::Arc;
use std::time::Duration;
use agent_deploy_contract::{
ActivateReleaseRequest, Application, ApplicationPage, ApprovalRequest, Artifact,
CancelDeploymentRequest, CreateApplicationRequest, CreateDeploymentRequest,
CreateEnvironmentRequest, DeleteApplicationRequest, DeleteResult, Deployment, DeploymentPage,
Environment, GetOperationRequest, LogPage, Operation, OperationState, PromoteReleaseRequest,
Promotion, Release, ReleasePage, ResumeDeploymentRequest, RollbackDeploymentRequest,
RouteBinding, UpdateApplicationRequest, paths,
};
use async_trait::async_trait;
use reqwest::Client;
use crate::operation::{OperationHandle, OperationObservation, OperationPoller, OperationProgress};
use crate::transport::{
CallOptions, ClientOptions, HttpTransport, InfraClientError, ServiceEndpoint,
};
#[derive(Clone, Debug)]
pub struct DeployClient {
transport: HttpTransport,
}
impl DeployClient {
pub async fn create_application(
&self,
request: &CreateApplicationRequest,
idempotency_key: &str,
) -> Result<Application, InfraClientError> {
self.transport
.post_json_with_options(
paths::APPS,
request,
CallOptions::default().idempotency_key(idempotency_key),
)
.await
}
pub async fn application(&self, app_id: &str) -> Result<Application, InfraClientError> {
self.transport
.get_json(&expand(paths::APP, &[app_id]))
.await
}
pub async fn applications(
&self,
cursor: Option<&str>,
limit: Option<u32>,
) -> Result<ApplicationPage, InfraClientError> {
self.transport
.get_json(&page_path(paths::APPS, cursor, limit))
.await
}
pub async fn update_application(
&self,
app_id: &str,
request: &UpdateApplicationRequest,
) -> Result<Application, InfraClientError> {
self.transport
.patch_json_with_options(
&expand(paths::APP, &[app_id]),
request,
CallOptions::default(),
)
.await
}
pub async fn delete_application(
&self,
app_id: &str,
expected_version: i64,
) -> Result<DeleteResult, InfraClientError> {
self.transport
.delete_json_with_options(
&expand(paths::APP, &[app_id]),
&DeleteApplicationRequest { expected_version },
CallOptions::default(),
)
.await
}
pub async fn deployments(
&self,
app_id: &str,
cursor: Option<&str>,
limit: Option<u32>,
) -> Result<DeploymentPage, InfraClientError> {
self.transport
.get_json(&page_path(
&expand(paths::APP_DEPLOYMENTS, &[app_id]),
cursor,
limit,
))
.await
}
pub(crate) fn new_with_endpoint(
http: Client,
endpoint: ServiceEndpoint,
options: ClientOptions,
) -> Self {
let endpoint = endpoint.with_default_credential_audience("deploy");
Self {
transport: HttpTransport::new_with_options(http, "deploy", endpoint, options),
}
}
pub async fn deploy(
&self,
app_id: &str,
request: &CreateDeploymentRequest,
idempotency_key: &str,
) -> Result<OperationHandle<Operation>, InfraClientError> {
let path = expand(paths::APP_DEPLOYMENTS, &[app_id]);
let operation: Operation = self
.transport
.post_json_with_options(
&path,
request,
CallOptions::default().idempotency_key(idempotency_key),
)
.await?;
Ok(self.handle(operation))
}
pub async fn operation(
&self,
operation_id: &str,
) -> Result<OperationHandle<Operation>, InfraClientError> {
let operation = DeployPoller {
transport: self.transport.clone(),
}
.poll(operation_id)
.await?;
Ok(self.handle(operation))
}
pub async fn deployment(&self, deployment_id: &str) -> Result<Deployment, InfraClientError> {
self.transport
.get_json(&expand(paths::DEPLOYMENT, &[deployment_id]))
.await
}
pub async fn cancel_deployment(
&self,
deployment_id: &str,
request: &CancelDeploymentRequest,
idempotency_key: &str,
) -> Result<OperationHandle<Operation>, InfraClientError> {
let operation = self
.transport
.post_json_with_options(
&expand(paths::DEPLOYMENT_CANCEL, &[deployment_id]),
request,
CallOptions::default().idempotency_key(idempotency_key),
)
.await?;
Ok(self.handle(operation))
}
pub async fn rollback(
&self,
deployment_id: &str,
request: &RollbackDeploymentRequest,
idempotency_key: &str,
) -> Result<OperationHandle<Operation>, InfraClientError> {
let operation = self
.transport
.post_json_with_options(
&expand(paths::DEPLOYMENT_ROLLBACK, &[deployment_id]),
request,
CallOptions::default().idempotency_key(idempotency_key),
)
.await?;
Ok(self.handle(operation))
}
pub async fn resume_deployment(
&self,
deployment_id: &str,
expected_version: i64,
idempotency_key: &str,
) -> Result<OperationHandle<Operation>, InfraClientError> {
let operation = self
.transport
.post_json_with_options(
&expand(paths::DEPLOYMENT_RESUME, &[deployment_id]),
&ResumeDeploymentRequest { expected_version },
CallOptions::default().idempotency_key(idempotency_key),
)
.await?;
Ok(self.handle(operation))
}
pub async fn create_environment(
&self,
app_id: &str,
request: &CreateEnvironmentRequest,
idempotency_key: &str,
) -> Result<Environment, InfraClientError> {
self.transport
.post_json_with_options(
&expand(paths::APP_ENVIRONMENTS, &[app_id]),
request,
CallOptions::default().idempotency_key(idempotency_key),
)
.await
}
pub async fn environment(&self, environment_id: &str) -> Result<Environment, InfraClientError> {
self.transport
.get_json(&expand(paths::ENVIRONMENT, &[environment_id]))
.await
}
pub async fn artifact(&self, digest: &str) -> Result<Artifact, InfraClientError> {
self.transport
.get_json(&expand(paths::ARTIFACT, &[digest]))
.await
}
pub async fn release(
&self,
environment_id: &str,
release_id: &str,
) -> Result<Release, InfraClientError> {
self.transport
.get_json(&expand(paths::RELEASE, &[environment_id, release_id]))
.await
}
pub async fn releases(
&self,
environment_id: &str,
cursor: Option<&str>,
limit: Option<u32>,
) -> Result<ReleasePage, InfraClientError> {
self.transport
.get_json(&page_path(
&expand(paths::RELEASES, &[environment_id]),
cursor,
limit,
))
.await
}
pub async fn activate_release(
&self,
environment_id: &str,
release_id: &str,
request: &ActivateReleaseRequest,
idempotency_key: &str,
) -> Result<RouteBinding, InfraClientError> {
self.transport
.post_json_with_options(
&expand(paths::RELEASE_ACTIVATE, &[environment_id, release_id]),
request,
CallOptions::default().idempotency_key(idempotency_key),
)
.await
}
pub async fn route_binding(
&self,
app_id: &str,
environment: &str,
) -> Result<RouteBinding, InfraClientError> {
self.transport
.get_json(&expand(paths::ROUTE, &[app_id, environment]))
.await
}
pub async fn promote_release(
&self,
request: &PromoteReleaseRequest,
idempotency_key: &str,
) -> Result<Promotion, InfraClientError> {
self.transport
.post_json_with_options(
paths::PROMOTIONS,
request,
CallOptions::default().idempotency_key(idempotency_key),
)
.await
}
pub async fn approve_promotion(
&self,
promotion_id: &str,
request: &ApprovalRequest,
) -> Result<Promotion, InfraClientError> {
self.transport
.post_json_with_options(
&expand(paths::PROMOTION_APPROVE, &[promotion_id]),
request,
CallOptions::default().idempotency_key(format!(
"approve:{promotion_id}:{}",
request.expected_version
)),
)
.await
}
pub async fn logs(
&self,
deployment_id: &str,
after: Option<i64>,
limit: Option<u32>,
) -> Result<LogPage, InfraClientError> {
let mut path = expand(paths::DEPLOYMENT_LOGS, &[deployment_id]);
let mut query = Vec::new();
if let Some(after) = after {
query.push(format!("after={after}"));
}
if let Some(limit) = limit {
query.push(format!("limit={}", limit.min(1_000)));
}
if !query.is_empty() {
path.push('?');
path.push_str(&query.join("&"));
}
self.transport.get_json(&path).await
}
pub async fn next_logs(
&self,
deployment_id: &str,
current: &LogPage,
limit: Option<u32>,
) -> Result<Option<LogPage>, InfraClientError> {
let Some(cursor) = current.next_cursor.as_deref() else {
return Ok(None);
};
let after = cursor
.parse::<i64>()
.map_err(|_| InfraClientError::Protocol {
service: "deploy",
message: "deploy returned a malformed log cursor".into(),
})?;
self.logs(deployment_id, Some(after), limit).await.map(Some)
}
fn handle(&self, operation: Operation) -> OperationHandle<Operation> {
OperationHandle::new(
operation.id.clone(),
Some(operation),
Arc::new(DeployPoller {
transport: self.transport.clone(),
}),
)
}
}
#[derive(Clone, Debug)]
struct DeployPoller {
transport: HttpTransport,
}
#[async_trait]
impl OperationPoller<Operation> for DeployPoller {
async fn poll(&self, operation_id: &str) -> Result<Operation, InfraClientError> {
let path = expand(paths::OPERATION, &[operation_id]);
match self.transport.get_json(&path).await {
Ok(operation) => Ok(operation),
Err(InfraClientError::HttpStatus { status: 405, .. }) => {
self.transport
.post_json_idempotent(
paths::OPERATIONS,
&GetOperationRequest {
operation_id: operation_id.to_string(),
},
)
.await
}
Err(error) => Err(error),
}
}
async fn cancel(&self, operation_id: &str) -> Result<(), InfraClientError> {
let operation = self.poll(operation_id).await?;
let _: Operation = self
.transport
.post_json_with_options(
&expand(paths::DEPLOYMENT_CANCEL, &[&operation.deployment_id]),
&CancelDeploymentRequest {
reason: Some("SDK cancellation".into()),
version: Some(operation.version),
},
CallOptions::default().idempotency_key(format!("cancel:{operation_id}")),
)
.await?;
Ok(())
}
fn observe(&self, operation: &Operation) -> OperationObservation {
let progress = match operation.state {
OperationState::Succeeded => OperationProgress::Succeeded,
OperationState::Failed => OperationProgress::Failed,
OperationState::Canceled => OperationProgress::Canceled,
OperationState::Queued | OperationState::Running | OperationState::Canceling => {
OperationProgress::Pending
}
};
OperationObservation {
progress,
next_poll_after: operation.next_poll_after_ms.map(Duration::from_millis),
error_code: operation.error_code.clone(),
}
}
}
fn expand(template: &str, values: &[&str]) -> String {
let mut result = template.to_string();
for value in values {
let Some(start) = result.find('{') else {
break;
};
let Some(relative_end) = result[start..].find('}') else {
break;
};
result.replace_range(start..=start + relative_end, &encode_path_segment(value));
}
result
}
fn encode_path_segment(value: &str) -> String {
let mut encoded = String::with_capacity(value.len());
for byte in value.bytes() {
if byte.is_ascii_alphanumeric() || matches!(byte, b'-' | b'.' | b'_' | b'~') {
encoded.push(char::from(byte));
} else {
use std::fmt::Write as _;
write!(&mut encoded, "%{byte:02X}").expect("writing to String cannot fail");
}
}
encoded
}
fn page_path(path: &str, cursor: Option<&str>, limit: Option<u32>) -> String {
let mut result = path.to_string();
let mut query = Vec::new();
if let Some(cursor) = cursor {
query.push(format!("cursor={}", encode_path_segment(cursor)));
}
if let Some(limit) = limit {
query.push(format!("limit={}", limit.min(100)));
}
if !query.is_empty() {
result.push('?');
result.push_str(&query.join("&"));
}
result
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn expands_contract_templates_without_duplicating_paths() {
assert_eq!(
expand(paths::APP_DEPLOYMENTS, &["app-1"]),
"/internal/v1/deploy/apps/app-1/deployments"
);
assert_eq!(
expand(paths::APP_DEPLOYMENTS, &["tenant/app 1"]),
"/internal/v1/deploy/apps/tenant%2Fapp%201/deployments"
);
}
}