1use super::*;
12use crate::worker::ControlPlaneWorker;
13
14#[derive(Debug, Clone, Copy, PartialEq, Eq)]
16pub struct RetryPolicy {
17 pub max_attempts: usize,
19 pub initial_backoff_ms: u64,
21 pub max_backoff_ms: u64,
23}
24
25impl Default for RetryPolicy {
26 fn default() -> Self {
27 Self {
28 max_attempts: 3,
29 initial_backoff_ms: 50,
30 max_backoff_ms: 500,
31 }
32 }
33}
34
35#[derive(Debug, Clone, PartialEq, Eq)]
37pub struct ControlPlaneHttpConfig {
38 pub base_url: String,
40 pub timeout_ms: u64,
42 pub retry_policy: RetryPolicy,
44}
45
46#[derive(Debug, Clone, PartialEq, Eq)]
48pub struct HttpControlPlaneRequest {
49 pub method: String,
51 pub path: String,
53 pub body: Vec<u8>,
55 pub timeout_ms: u64,
57}
58
59#[derive(Debug, Clone, PartialEq, Eq)]
61pub struct HttpControlPlaneResponse {
62 pub status_code: u16,
64 pub body: Vec<u8>,
66}
67
68pub trait HttpTransport: Send + Sync {
70 fn send_json(
72 &self,
73 base_url: &str,
74 request: HttpControlPlaneRequest,
75 ) -> ControlPlaneResult<HttpControlPlaneResponse>;
76
77 fn send_json_traced(
79 &self,
80 base_url: &str,
81 request: HttpControlPlaneRequest,
82 _trace: Option<&TraceContext>,
83 ) -> ControlPlaneResult<HttpControlPlaneResponse> {
84 self.send_json(base_url, request)
85 }
86
87 fn send_json_traced_cancellable(
89 &self,
90 base_url: &str,
91 request: HttpControlPlaneRequest,
92 trace: Option<&TraceContext>,
93 cancellation: &CancellationToken,
94 ) -> ControlPlaneResult<HttpControlPlaneResponse> {
95 if cancellation.is_cancelled() {
96 return Err(ControlPlaneError::Transport(
97 "control-plane request cancelled".to_string(),
98 ));
99 }
100 self.send_json_traced(base_url, request, trace)
101 }
102}
103
104#[derive(Debug)]
106pub struct HttpControlPlaneClient<T> {
107 config: ControlPlaneHttpConfig,
108 transport: Arc<T>,
109 worker: ControlPlaneWorker,
110 cancellation: CancellationToken,
111}
112
113impl<T> Clone for HttpControlPlaneClient<T> {
114 fn clone(&self) -> Self {
115 Self {
116 config: self.config.clone(),
117 transport: Arc::clone(&self.transport),
118 worker: self.worker.clone(),
119 cancellation: self.cancellation.clone(),
120 }
121 }
122}
123
124impl<T> HttpControlPlaneClient<T>
125where
126 T: HttpTransport,
127{
128 pub fn new(config: ControlPlaneHttpConfig, transport: T) -> Self {
130 Self {
131 config,
132 transport: Arc::new(transport),
133 worker: ControlPlaneWorker::new(),
134 cancellation: CancellationToken::new(),
135 }
136 }
137
138 pub fn with_cancellation_token(mut self, cancellation: CancellationToken) -> Self {
140 self.cancellation = cancellation;
141 self
142 }
143
144 pub fn cancel(&self) {
146 self.cancellation.cancel();
147 }
148
149 pub fn is_cancelled(&self) -> bool {
151 self.cancellation.is_cancelled()
152 }
153
154 fn post<Req, Resp>(&self, path: &str, value: &Req) -> ControlPlaneResult<Resp>
155 where
156 Req: Serialize,
157 Resp: for<'de> Deserialize<'de>,
158 {
159 self.post_traced(path, value, None)
160 }
161
162 fn post_traced<Req, Resp>(
163 &self,
164 path: &str,
165 value: &Req,
166 trace: Option<&TraceContext>,
167 ) -> ControlPlaneResult<Resp>
168 where
169 Req: Serialize,
170 Resp: for<'de> Deserialize<'de>,
171 {
172 let body = serde_json::to_vec(value)
173 .map_err(|error| ControlPlaneError::Transport(error.to_string()))?;
174 let request = HttpControlPlaneRequest {
175 method: "POST".to_string(),
176 path: path.to_string(),
177 body,
178 timeout_ms: self.config.timeout_ms,
179 };
180 self.send_with_retry::<Resp>(request, trace)
181 }
182
183 fn send_with_retry<Resp>(
184 &self,
185 request: HttpControlPlaneRequest,
186 trace: Option<&TraceContext>,
187 ) -> ControlPlaneResult<Resp>
188 where
189 Resp: for<'de> Deserialize<'de>,
190 {
191 let attempts = self.config.retry_policy.max_attempts.max(1);
192 let mut backoff_ms = self.config.retry_policy.initial_backoff_ms;
193 let mut last_error = ControlPlaneError::Offline;
194 for attempt in 0..attempts {
195 match self.transport.send_json_traced_cancellable(
196 &self.config.base_url,
197 request.clone(),
198 trace,
199 &self.cancellation,
200 ) {
201 Ok(response) if (200..300).contains(&response.status_code) => {
202 return serde_json::from_slice(&response.body)
203 .map_err(|error| ControlPlaneError::InvalidResponse(error.to_string()));
204 }
205 Ok(response) => {
206 last_error = ControlPlaneError::Rejected(format!(
207 "http_status={}",
208 response.status_code
209 ));
210 }
211 Err(error) => last_error = error,
212 }
213
214 if attempt + 1 < attempts {
215 if self
216 .cancellation
217 .wait_timeout(Duration::from_millis(backoff_ms))
218 {
219 return Err(ControlPlaneError::Transport(
220 "control-plane request cancelled".to_string(),
221 ));
222 }
223 backoff_ms =
224 (backoff_ms.saturating_mul(2)).min(self.config.retry_policy.max_backoff_ms);
225 }
226 }
227 Err(last_error)
228 }
229}
230
231impl<T> ControlPlaneProvider for HttpControlPlaneClient<T>
232where
233 T: HttpTransport + 'static,
234{
235 fn register<'a>(
236 &'a self,
237 registration: CoreRegistration,
238 ) -> ControlPlaneFuture<'a, CorePresence> {
239 let client = self.clone();
240 self.worker
241 .enqueue(move || client.post(CONTROL_REGISTER_PATH, ®istration))
242 }
243
244 fn heartbeat<'a>(
245 &'a self,
246 request: HeartbeatRequest,
247 ) -> ControlPlaneFuture<'a, HeartbeatResponse> {
248 let client = self.clone();
249 self.worker
250 .enqueue(move || client.post(CONTROL_HEARTBEAT_PATH, &request))
251 }
252
253 fn discover_peers<'a>(
254 &'a self,
255 identity: &'a CoreIdentity,
256 ) -> ControlPlaneFuture<'a, PeerDirectory> {
257 let client = self.clone();
258 let identity = identity.clone();
259 self.worker
260 .enqueue(move || client.post(CONTROL_PEERS_PATH, &identity))
261 }
262
263 fn acquire_or_renew_service_lease<'a>(
264 &'a self,
265 identity: &'a CoreIdentity,
266 service_id: &'a ServiceId,
267 ttl_ms: u64,
268 now_ms: u64,
269 ) -> ControlPlaneFuture<'a, ServiceLeaderLease> {
270 let client = self.clone();
271 let request = ServiceLeaseRequest {
272 identity: identity.clone(),
273 service_id: service_id.clone(),
274 ttl_ms,
275 now_ms,
276 };
277 self.worker
278 .enqueue(move || client.post(CONTROL_SERVICE_LEASE_PATH, &request))
279 }
280
281 fn release_service_lease<'a>(
282 &'a self,
283 lease: ServiceLeaderLease,
284 ) -> ControlPlaneFuture<'a, ()> {
285 let client = self.clone();
286 self.worker.enqueue(move || {
287 let _: EmptyResponse = client.post(CONTROL_SERVICE_LEASE_RELEASE_PATH, &lease)?;
288 Ok(())
289 })
290 }
291
292 fn register_traced<'a>(
293 &'a self,
294 registration: CoreRegistration,
295 trace: Option<&'a TraceContext>,
296 ) -> ControlPlaneFuture<'a, CorePresence> {
297 let client = self.clone();
298 let trace = trace.cloned();
299 self.worker.enqueue(move || {
300 client.post_traced(CONTROL_REGISTER_PATH, ®istration, trace.as_ref())
301 })
302 }
303
304 fn heartbeat_traced<'a>(
305 &'a self,
306 request: HeartbeatRequest,
307 trace: Option<&'a TraceContext>,
308 ) -> ControlPlaneFuture<'a, HeartbeatResponse> {
309 let client = self.clone();
310 let trace = trace.cloned();
311 self.worker
312 .enqueue(move || client.post_traced(CONTROL_HEARTBEAT_PATH, &request, trace.as_ref()))
313 }
314
315 fn discover_peers_traced<'a>(
316 &'a self,
317 identity: &'a CoreIdentity,
318 trace: Option<&'a TraceContext>,
319 ) -> ControlPlaneFuture<'a, PeerDirectory> {
320 let client = self.clone();
321 let identity = identity.clone();
322 let trace = trace.cloned();
323 self.worker
324 .enqueue(move || client.post_traced(CONTROL_PEERS_PATH, &identity, trace.as_ref()))
325 }
326
327 fn acquire_or_renew_service_lease_traced<'a>(
328 &'a self,
329 identity: &'a CoreIdentity,
330 service_id: &'a ServiceId,
331 ttl_ms: u64,
332 now_ms: u64,
333 trace: Option<&'a TraceContext>,
334 ) -> ControlPlaneFuture<'a, ServiceLeaderLease> {
335 let client = self.clone();
336 let trace = trace.cloned();
337 let request = ServiceLeaseRequest {
338 identity: identity.clone(),
339 service_id: service_id.clone(),
340 ttl_ms,
341 now_ms,
342 };
343 self.worker.enqueue(move || {
344 client.post_traced(CONTROL_SERVICE_LEASE_PATH, &request, trace.as_ref())
345 })
346 }
347
348 fn release_service_lease_traced<'a>(
349 &'a self,
350 lease: ServiceLeaderLease,
351 trace: Option<&'a TraceContext>,
352 ) -> ControlPlaneFuture<'a, ()> {
353 let client = self.clone();
354 let trace = trace.cloned();
355 self.worker.enqueue(move || {
356 let _: EmptyResponse =
357 client.post_traced(CONTROL_SERVICE_LEASE_RELEASE_PATH, &lease, trace.as_ref())?;
358 Ok(())
359 })
360 }
361}