Skip to main content

appcore_control_plane/
client.rs

1// =============================================================================
2//        #######
3//     ###       ###     F: client.rs
4//    ##   ## ##   ##    P: AppCore-Runtime
5//         ## ##
6//                       C: 2026/07/22 15:41:18 by dnettoRaw
7//    ##   ## ##   ##    U: 2026/07/24 16:07:49 by dnettoRaw
8//      ###########      S: 1.0.1-rc.8
9// =============================================================================
10
11use super::*;
12use crate::worker::ControlPlaneWorker;
13
14/// Retry limits and exponential backoff bounds for control-plane requests.
15#[derive(Debug, Clone, Copy, PartialEq, Eq)]
16pub struct RetryPolicy {
17    /// Maximum number of attempts, including the initial request.
18    pub max_attempts: usize,
19    /// Delay before the first retry.
20    pub initial_backoff_ms: u64,
21    /// Maximum delay between attempts.
22    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/// Configuration for the generic HTTP control-plane client.
36#[derive(Debug, Clone, PartialEq, Eq)]
37pub struct ControlPlaneHttpConfig {
38    /// Base URL that hosts the stable control-plane endpoints.
39    pub base_url: String,
40    /// Per-attempt network timeout.
41    pub timeout_ms: u64,
42    /// Retry behavior for transport and non-success responses.
43    pub retry_policy: RetryPolicy,
44}
45
46/// Transport request produced by [`HttpControlPlaneClient`].
47#[derive(Debug, Clone, PartialEq, Eq)]
48pub struct HttpControlPlaneRequest {
49    /// HTTP method.
50    pub method: String,
51    /// Stable endpoint path relative to the configured base URL.
52    pub path: String,
53    /// Serialized JSON body.
54    pub body: Vec<u8>,
55    /// Per-attempt timeout.
56    pub timeout_ms: u64,
57}
58
59/// Bounded response returned by an [`HttpTransport`].
60#[derive(Debug, Clone, PartialEq, Eq)]
61pub struct HttpControlPlaneResponse {
62    /// HTTP status code.
63    pub status_code: u16,
64    /// Raw response body.
65    pub body: Vec<u8>,
66}
67
68/// Adapter contract used by the control-plane HTTP client.
69pub trait HttpTransport: Send + Sync {
70    /// Sends one JSON request to the configured base URL.
71    fn send_json(
72        &self,
73        base_url: &str,
74        request: HttpControlPlaneRequest,
75    ) -> ControlPlaneResult<HttpControlPlaneResponse>;
76
77    /// Sends one JSON request and propagates optional trace headers.
78    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    /// Sends one traced request with cooperative cancellation.
88    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/// Control-plane client that maps stable contracts onto an HTTP transport.
105#[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    /// Creates an HTTP client with a dedicated bounded worker queue.
129    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    /// Replaces the shared cancellation token used by requests and retries.
139    pub fn with_cancellation_token(mut self, cancellation: CancellationToken) -> Self {
140        self.cancellation = cancellation;
141        self
142    }
143
144    /// Cancels queued requests, active official transport I/O, and retry waits.
145    pub fn cancel(&self) {
146        self.cancellation.cancel();
147    }
148
149    /// Reports whether this client has been cancelled.
150    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, &registration))
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, &registration, 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}