use super::*;
use crate::worker::ControlPlaneWorker;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct RetryPolicy {
pub max_attempts: usize,
pub initial_backoff_ms: u64,
pub max_backoff_ms: u64,
}
impl Default for RetryPolicy {
fn default() -> Self {
Self {
max_attempts: 3,
initial_backoff_ms: 50,
max_backoff_ms: 500,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ControlPlaneHttpConfig {
pub base_url: String,
pub timeout_ms: u64,
pub retry_policy: RetryPolicy,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct HttpControlPlaneRequest {
pub method: String,
pub path: String,
pub body: Vec<u8>,
pub timeout_ms: u64,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct HttpControlPlaneResponse {
pub status_code: u16,
pub body: Vec<u8>,
}
pub trait HttpTransport: Send + Sync {
fn send_json(
&self,
base_url: &str,
request: HttpControlPlaneRequest,
) -> ControlPlaneResult<HttpControlPlaneResponse>;
fn send_json_traced(
&self,
base_url: &str,
request: HttpControlPlaneRequest,
_trace: Option<&TraceContext>,
) -> ControlPlaneResult<HttpControlPlaneResponse> {
self.send_json(base_url, request)
}
fn send_json_traced_cancellable(
&self,
base_url: &str,
request: HttpControlPlaneRequest,
trace: Option<&TraceContext>,
cancellation: &CancellationToken,
) -> ControlPlaneResult<HttpControlPlaneResponse> {
if cancellation.is_cancelled() {
return Err(ControlPlaneError::Transport(
"control-plane request cancelled".to_string(),
));
}
self.send_json_traced(base_url, request, trace)
}
}
#[derive(Debug)]
pub struct HttpControlPlaneClient<T> {
config: ControlPlaneHttpConfig,
transport: Arc<T>,
worker: ControlPlaneWorker,
cancellation: CancellationToken,
}
impl<T> Clone for HttpControlPlaneClient<T> {
fn clone(&self) -> Self {
Self {
config: self.config.clone(),
transport: Arc::clone(&self.transport),
worker: self.worker.clone(),
cancellation: self.cancellation.clone(),
}
}
}
impl<T> HttpControlPlaneClient<T>
where
T: HttpTransport,
{
pub fn new(config: ControlPlaneHttpConfig, transport: T) -> Self {
Self {
config,
transport: Arc::new(transport),
worker: ControlPlaneWorker::new(),
cancellation: CancellationToken::new(),
}
}
pub fn with_cancellation_token(mut self, cancellation: CancellationToken) -> Self {
self.cancellation = cancellation;
self
}
pub fn cancel(&self) {
self.cancellation.cancel();
}
pub fn is_cancelled(&self) -> bool {
self.cancellation.is_cancelled()
}
fn post<Req, Resp>(&self, path: &str, value: &Req) -> ControlPlaneResult<Resp>
where
Req: Serialize,
Resp: for<'de> Deserialize<'de>,
{
self.post_traced(path, value, None)
}
fn post_traced<Req, Resp>(
&self,
path: &str,
value: &Req,
trace: Option<&TraceContext>,
) -> ControlPlaneResult<Resp>
where
Req: Serialize,
Resp: for<'de> Deserialize<'de>,
{
let body = serde_json::to_vec(value)
.map_err(|error| ControlPlaneError::Transport(error.to_string()))?;
let request = HttpControlPlaneRequest {
method: "POST".to_string(),
path: path.to_string(),
body,
timeout_ms: self.config.timeout_ms,
};
self.send_with_retry::<Resp>(request, trace)
}
fn send_with_retry<Resp>(
&self,
request: HttpControlPlaneRequest,
trace: Option<&TraceContext>,
) -> ControlPlaneResult<Resp>
where
Resp: for<'de> Deserialize<'de>,
{
let attempts = self.config.retry_policy.max_attempts.max(1);
let mut backoff_ms = self.config.retry_policy.initial_backoff_ms;
let mut last_error = ControlPlaneError::Offline;
for attempt in 0..attempts {
match self.transport.send_json_traced_cancellable(
&self.config.base_url,
request.clone(),
trace,
&self.cancellation,
) {
Ok(response) if (200..300).contains(&response.status_code) => {
return serde_json::from_slice(&response.body)
.map_err(|error| ControlPlaneError::InvalidResponse(error.to_string()));
}
Ok(response) => {
last_error = ControlPlaneError::Rejected(format!(
"http_status={}",
response.status_code
));
}
Err(error) => last_error = error,
}
if attempt + 1 < attempts {
if self
.cancellation
.wait_timeout(Duration::from_millis(backoff_ms))
{
return Err(ControlPlaneError::Transport(
"control-plane request cancelled".to_string(),
));
}
backoff_ms =
(backoff_ms.saturating_mul(2)).min(self.config.retry_policy.max_backoff_ms);
}
}
Err(last_error)
}
}
impl<T> ControlPlaneProvider for HttpControlPlaneClient<T>
where
T: HttpTransport + 'static,
{
fn register<'a>(
&'a self,
registration: CoreRegistration,
) -> ControlPlaneFuture<'a, CorePresence> {
let client = self.clone();
self.worker
.enqueue(move || client.post(CONTROL_REGISTER_PATH, ®istration))
}
fn heartbeat<'a>(
&'a self,
request: HeartbeatRequest,
) -> ControlPlaneFuture<'a, HeartbeatResponse> {
let client = self.clone();
self.worker
.enqueue(move || client.post(CONTROL_HEARTBEAT_PATH, &request))
}
fn discover_peers<'a>(
&'a self,
identity: &'a CoreIdentity,
) -> ControlPlaneFuture<'a, PeerDirectory> {
let client = self.clone();
let identity = identity.clone();
self.worker
.enqueue(move || client.post(CONTROL_PEERS_PATH, &identity))
}
fn acquire_or_renew_service_lease<'a>(
&'a self,
identity: &'a CoreIdentity,
service_id: &'a ServiceId,
ttl_ms: u64,
now_ms: u64,
) -> ControlPlaneFuture<'a, ServiceLeaderLease> {
let client = self.clone();
let request = ServiceLeaseRequest {
identity: identity.clone(),
service_id: service_id.clone(),
ttl_ms,
now_ms,
};
self.worker
.enqueue(move || client.post(CONTROL_SERVICE_LEASE_PATH, &request))
}
fn release_service_lease<'a>(
&'a self,
lease: ServiceLeaderLease,
) -> ControlPlaneFuture<'a, ()> {
let client = self.clone();
self.worker.enqueue(move || {
let _: EmptyResponse = client.post(CONTROL_SERVICE_LEASE_RELEASE_PATH, &lease)?;
Ok(())
})
}
fn register_traced<'a>(
&'a self,
registration: CoreRegistration,
trace: Option<&'a TraceContext>,
) -> ControlPlaneFuture<'a, CorePresence> {
let client = self.clone();
let trace = trace.cloned();
self.worker.enqueue(move || {
client.post_traced(CONTROL_REGISTER_PATH, ®istration, trace.as_ref())
})
}
fn heartbeat_traced<'a>(
&'a self,
request: HeartbeatRequest,
trace: Option<&'a TraceContext>,
) -> ControlPlaneFuture<'a, HeartbeatResponse> {
let client = self.clone();
let trace = trace.cloned();
self.worker
.enqueue(move || client.post_traced(CONTROL_HEARTBEAT_PATH, &request, trace.as_ref()))
}
fn discover_peers_traced<'a>(
&'a self,
identity: &'a CoreIdentity,
trace: Option<&'a TraceContext>,
) -> ControlPlaneFuture<'a, PeerDirectory> {
let client = self.clone();
let identity = identity.clone();
let trace = trace.cloned();
self.worker
.enqueue(move || client.post_traced(CONTROL_PEERS_PATH, &identity, trace.as_ref()))
}
fn acquire_or_renew_service_lease_traced<'a>(
&'a self,
identity: &'a CoreIdentity,
service_id: &'a ServiceId,
ttl_ms: u64,
now_ms: u64,
trace: Option<&'a TraceContext>,
) -> ControlPlaneFuture<'a, ServiceLeaderLease> {
let client = self.clone();
let trace = trace.cloned();
let request = ServiceLeaseRequest {
identity: identity.clone(),
service_id: service_id.clone(),
ttl_ms,
now_ms,
};
self.worker.enqueue(move || {
client.post_traced(CONTROL_SERVICE_LEASE_PATH, &request, trace.as_ref())
})
}
fn release_service_lease_traced<'a>(
&'a self,
lease: ServiceLeaderLease,
trace: Option<&'a TraceContext>,
) -> ControlPlaneFuture<'a, ()> {
let client = self.clone();
let trace = trace.cloned();
self.worker.enqueue(move || {
let _: EmptyResponse =
client.post_traced(CONTROL_SERVICE_LEASE_RELEASE_PATH, &lease, trace.as_ref())?;
Ok(())
})
}
}