1use appcore_contracts::ServiceId;
14use appcore_types::{
15 CapabilityDescriptor, ClusterId, CoreId, CoreIdentity, DistributedCoreManifest, PeerEndpoint,
16 RuntimeOperationalMode, TenantId, TraceContext,
17};
18use serde::{Deserialize, Serialize};
19use std::collections::BTreeMap;
20use std::future::Future;
21use std::pin::Pin;
22
23pub const CONTROL_PLANE_PROTOCOL_VERSION: u16 = 1;
25pub const CONTROL_REGISTER_PATH: &str = "/v1/control/register";
27pub const CONTROL_HEARTBEAT_PATH: &str = "/v1/control/heartbeat";
29pub const CONTROL_PEERS_PATH: &str = "/v1/control/peers";
31pub const CONTROL_SERVICE_LEASE_PATH: &str = "/v1/control/service-lease";
33pub const CONTROL_SERVICE_LEASE_RELEASE_PATH: &str = "/v1/control/service-lease/release";
35
36pub type ControlPlaneResult<T> = Result<T, ControlPlaneError>;
38pub type ControlPlaneFuture<'a, T> =
40 Pin<Box<dyn Future<Output = ControlPlaneResult<T>> + Send + 'a>>;
41
42#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
44pub enum ControlPlaneError {
45 #[error("control plane is offline")]
47 Offline,
48 #[error("control plane request timed out")]
50 Timeout,
51 #[error("control plane rejected operation: {0}")]
53 Rejected(String),
54 #[error("control plane state conflict: {0}")]
56 Conflict(String),
57 #[error("invalid control plane response: {0}")]
59 InvalidResponse(String),
60 #[error("control plane transport failed: {0}")]
62 Transport(String),
63 #[error("control plane lease is unavailable")]
65 LeaseUnavailable,
66}
67
68#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
70pub struct CoreRegistration {
71 pub manifest: DistributedCoreManifest,
73 pub registered_at_ms: u64,
75 pub operation_mode: RuntimeOperationalMode,
77}
78
79#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
81pub struct CorePresence {
82 pub identity: CoreIdentity,
84 pub operation_mode: RuntimeOperationalMode,
86 pub healthy: bool,
88 pub last_seen_ms: u64,
90}
91
92#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
94pub struct HeartbeatRequest {
95 pub identity: CoreIdentity,
97 pub operation_mode: RuntimeOperationalMode,
99 pub sent_at_ms: u64,
101}
102
103#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
105pub struct HeartbeatResponse {
106 pub accepted: bool,
108 pub server_time_ms: u64,
110 pub operation_mode: RuntimeOperationalMode,
112}
113
114#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
116pub struct PeerDirectory {
117 pub tenant_id: TenantId,
119 pub cluster_id: Option<ClusterId>,
121 pub peers: Vec<PeerRecord>,
123 pub refreshed_at_ms: u64,
125}
126
127#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
129pub struct PeerRecord {
130 pub identity: CoreIdentity,
132 pub endpoints: Vec<PeerEndpoint>,
134 pub capabilities: Vec<CapabilityDescriptor>,
136 pub healthy: bool,
138 pub last_seen_ms: u64,
140 pub metadata: BTreeMap<String, String>,
142}
143
144#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
146pub struct ServiceLeaderLease {
147 pub service_id: ServiceId,
149 pub tenant_id: TenantId,
151 pub cluster_id: ClusterId,
153 pub holder_core_id: CoreId,
155 pub epoch: u64,
157 pub acquired_at_ms: u64,
159 pub expires_at_ms: u64,
161}
162
163#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
165pub struct ServiceLeaseRequest {
166 pub identity: CoreIdentity,
168 pub service_id: ServiceId,
170 pub ttl_ms: u64,
172 pub now_ms: u64,
174}
175
176#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
178pub struct EmptyResponse {}
179
180pub trait ControlPlaneProvider: Send + Sync {
182 fn register<'a>(
184 &'a self,
185 registration: CoreRegistration,
186 ) -> ControlPlaneFuture<'a, CorePresence>;
187
188 fn heartbeat<'a>(
190 &'a self,
191 request: HeartbeatRequest,
192 ) -> ControlPlaneFuture<'a, HeartbeatResponse>;
193
194 fn discover_peers<'a>(
196 &'a self,
197 identity: &'a CoreIdentity,
198 ) -> ControlPlaneFuture<'a, PeerDirectory>;
199
200 fn acquire_or_renew_service_lease<'a>(
202 &'a self,
203 identity: &'a CoreIdentity,
204 service_id: &'a ServiceId,
205 ttl_ms: u64,
206 now_ms: u64,
207 ) -> ControlPlaneFuture<'a, ServiceLeaderLease>;
208
209 fn release_service_lease<'a>(&'a self, lease: ServiceLeaderLease)
211 -> ControlPlaneFuture<'a, ()>;
212
213 fn register_traced<'a>(
215 &'a self,
216 registration: CoreRegistration,
217 _trace: Option<&'a TraceContext>,
218 ) -> ControlPlaneFuture<'a, CorePresence> {
219 self.register(registration)
220 }
221
222 fn heartbeat_traced<'a>(
224 &'a self,
225 request: HeartbeatRequest,
226 _trace: Option<&'a TraceContext>,
227 ) -> ControlPlaneFuture<'a, HeartbeatResponse> {
228 self.heartbeat(request)
229 }
230
231 fn discover_peers_traced<'a>(
233 &'a self,
234 identity: &'a CoreIdentity,
235 _trace: Option<&'a TraceContext>,
236 ) -> ControlPlaneFuture<'a, PeerDirectory> {
237 self.discover_peers(identity)
238 }
239
240 fn acquire_or_renew_service_lease_traced<'a>(
242 &'a self,
243 identity: &'a CoreIdentity,
244 service_id: &'a ServiceId,
245 ttl_ms: u64,
246 now_ms: u64,
247 _trace: Option<&'a TraceContext>,
248 ) -> ControlPlaneFuture<'a, ServiceLeaderLease> {
249 self.acquire_or_renew_service_lease(identity, service_id, ttl_ms, now_ms)
250 }
251
252 fn release_service_lease_traced<'a>(
254 &'a self,
255 lease: ServiceLeaderLease,
256 _trace: Option<&'a TraceContext>,
257 ) -> ControlPlaneFuture<'a, ()> {
258 self.release_service_lease(lease)
259 }
260}
261
262pub trait DiscoveryProvider: Send + Sync {
264 fn discover<'a>(&'a self, identity: &'a CoreIdentity) -> ControlPlaneFuture<'a, PeerDirectory>;
266
267 fn discover_traced<'a>(
269 &'a self,
270 identity: &'a CoreIdentity,
271 trace: Option<&'a TraceContext>,
272 ) -> ControlPlaneFuture<'a, PeerDirectory>;
273}
274
275impl<T> DiscoveryProvider for T
276where
277 T: ControlPlaneProvider + ?Sized,
278{
279 fn discover<'a>(&'a self, identity: &'a CoreIdentity) -> ControlPlaneFuture<'a, PeerDirectory> {
280 self.discover_peers(identity)
281 }
282
283 fn discover_traced<'a>(
284 &'a self,
285 identity: &'a CoreIdentity,
286 trace: Option<&'a TraceContext>,
287 ) -> ControlPlaneFuture<'a, PeerDirectory> {
288 self.discover_peers_traced(identity, trace)
289 }
290}
291
292pub trait ServiceLeadershipGuard: Send + Sync {
294 fn current_service_lease(&self, service_id: &ServiceId) -> Option<ServiceLeaderLease>;
296
297 fn check_service_write_permission(
299 &self,
300 service_id: &ServiceId,
301 tenant_id: &TenantId,
302 cluster_id: &ClusterId,
303 core_id: &CoreId,
304 min_epoch: Option<u64>,
305 now_ms: u64,
306 ) -> LeadershipDecision;
307}
308
309#[derive(Debug, Clone, Copy, PartialEq, Eq)]
311pub enum LeadershipDecision {
312 Allowed,
314 NoLease,
316 Expired,
318 StaleEpoch,
320 WrongHolder,
322}
323
324#[cfg(test)]
325#[path = "tests.rs"]
326mod tests;