Skip to main content

appcore_control_plane/
memory.rs

1// =============================================================================
2//        #######
3//     ###       ###     F: memory.rs
4//    ##   ## ##   ##    P: AppCore-Runtime
5//         ## ##
6//                       C: 2026/07/22 15:41:18 by dnettoRaw
7//    ##   ## ##   ##    U: 2026/08/02 13:24:05 by dnettoRaw
8//      ###########      S: 1.0.1-rc.8
9// =============================================================================
10
11use super::*;
12
13/// Deterministic in-memory control plane for embedded hosts and tests.
14#[derive(Debug, Clone, Default)]
15pub struct InMemoryControlPlane {
16    state: std::sync::Arc<Mutex<InMemoryState>>,
17}
18
19#[derive(Debug, Clone, Default, serde::Serialize, serde::Deserialize)]
20#[serde(deny_unknown_fields)]
21pub(crate) struct InMemoryState {
22    registrations: BTreeMap<String, CoreRegistration>,
23    service_leases: BTreeMap<String, ServiceLeaseSlot>,
24}
25
26#[derive(Debug, Clone, Default, serde::Serialize, serde::Deserialize)]
27struct ServiceLeaseSlot {
28    lease: Option<ServiceLeaderLease>,
29    last_epoch: u64,
30}
31
32impl InMemoryControlPlane {
33    /// Returns the number of registered runtime instances.
34    pub fn registrations_len(&self) -> ControlPlaneResult<usize> {
35        Ok(lock_state(&self.state)?.registrations.len())
36    }
37
38    pub(crate) fn from_state(state: InMemoryState) -> Self {
39        Self {
40            state: std::sync::Arc::new(Mutex::new(state)),
41        }
42    }
43
44    pub(crate) fn snapshot(&self) -> ControlPlaneResult<InMemoryState> {
45        Ok(lock_state(&self.state)?.clone())
46    }
47
48    pub(crate) fn prune_registrations(&self, cutoff_ms: u64) -> ControlPlaneResult<usize> {
49        let mut state = lock_state(&self.state)?;
50        let before = state.registrations.len();
51        state
52            .registrations
53            .retain(|_, registration| registration.registered_at_ms >= cutoff_ms);
54        Ok(before.saturating_sub(state.registrations.len()))
55    }
56}
57
58impl ControlPlaneProvider for InMemoryControlPlane {
59    fn register<'a>(
60        &'a self,
61        registration: CoreRegistration,
62    ) -> ControlPlaneFuture<'a, CorePresence> {
63        Box::pin(async move {
64            let presence = CorePresence {
65                identity: registration.manifest.identity.clone(),
66                operation_mode: registration.operation_mode,
67                healthy: is_routable(registration.operation_mode),
68                last_seen_ms: registration.registered_at_ms,
69            };
70            let key = instance_key(&presence.identity);
71            lock_state(&self.state)?
72                .registrations
73                .insert(key, registration);
74            Ok(presence)
75        })
76    }
77
78    fn heartbeat<'a>(
79        &'a self,
80        request: HeartbeatRequest,
81    ) -> ControlPlaneFuture<'a, HeartbeatResponse> {
82        Box::pin(async move {
83            let mut state = lock_state(&self.state)?;
84            let registration_key = instance_key(&request.identity);
85            if let Some(registration) = state.registrations.get_mut(&registration_key) {
86                registration.registered_at_ms = request.sent_at_ms;
87                registration.operation_mode = request.operation_mode;
88            }
89            Ok(HeartbeatResponse {
90                accepted: true,
91                server_time_ms: request.sent_at_ms,
92                operation_mode: request.operation_mode,
93            })
94        })
95    }
96
97    fn discover_peers<'a>(
98        &'a self,
99        identity: &'a CoreIdentity,
100    ) -> ControlPlaneFuture<'a, PeerDirectory> {
101        Box::pin(async move {
102            let state = lock_state(&self.state)?;
103            let peers = state
104                .registrations
105                .values()
106                .filter(|registration| {
107                    registration.manifest.identity.tenant_id == identity.tenant_id
108                        && registration.manifest.identity.cluster_id == identity.cluster_id
109                        && registration.manifest.identity.instance_id != identity.instance_id
110                })
111                .map(|registration| PeerRecord {
112                    identity: registration.manifest.identity.clone(),
113                    endpoints: registration.manifest.endpoints.clone(),
114                    capabilities: registration.manifest.capabilities.clone(),
115                    healthy: is_routable(registration.operation_mode),
116                    last_seen_ms: registration.registered_at_ms,
117                    metadata: registration.manifest.metadata.clone(),
118                })
119                .collect::<Vec<_>>();
120            Ok(PeerDirectory {
121                tenant_id: identity.tenant_id.clone(),
122                cluster_id: Some(identity.cluster_id.clone()),
123                peers,
124                refreshed_at_ms: 0,
125            })
126        })
127    }
128
129    fn acquire_or_renew_service_lease<'a>(
130        &'a self,
131        identity: &'a CoreIdentity,
132        service_id: &'a ServiceId,
133        ttl_ms: u64,
134        now_ms: u64,
135    ) -> ControlPlaneFuture<'a, ServiceLeaderLease> {
136        Box::pin(async move {
137            let expires_at_ms = lease_expiration(now_ms, ttl_ms)?;
138            let mut state = lock_state(&self.state)?;
139            let key = service_lease_key(&identity.tenant_id, &identity.cluster_id, service_id);
140            let slot = state.service_leases.entry(key).or_default();
141            let lease = match slot.lease.as_ref() {
142                Some(current)
143                    if current.expires_at_ms > now_ms
144                        && current.holder_core_id == identity.core_id =>
145                {
146                    ServiceLeaderLease {
147                        expires_at_ms,
148                        ..current.clone()
149                    }
150                }
151                Some(current) if current.expires_at_ms > now_ms => {
152                    return Err(ControlPlaneError::LeaseUnavailable);
153                }
154                _ => {
155                    let epoch = slot.last_epoch.checked_add(1).ok_or_else(|| {
156                        ControlPlaneError::Conflict(
157                            "service lease fencing epoch exhausted".to_string(),
158                        )
159                    })?;
160                    ServiceLeaderLease {
161                        service_id: service_id.clone(),
162                        tenant_id: identity.tenant_id.clone(),
163                        cluster_id: identity.cluster_id.clone(),
164                        holder_core_id: identity.core_id.clone(),
165                        epoch,
166                        acquired_at_ms: now_ms,
167                        expires_at_ms,
168                    }
169                }
170            };
171            slot.last_epoch = slot.last_epoch.max(lease.epoch);
172            slot.lease = Some(lease.clone());
173            Ok(lease)
174        })
175    }
176
177    fn release_service_lease<'a>(
178        &'a self,
179        lease: ServiceLeaderLease,
180    ) -> ControlPlaneFuture<'a, ()> {
181        Box::pin(async move {
182            let mut state = lock_state(&self.state)?;
183            let key = service_lease_key(&lease.tenant_id, &lease.cluster_id, &lease.service_id);
184            let Some(slot) = state.service_leases.get_mut(&key) else {
185                return Ok(());
186            };
187            let matches_current = slot.lease.as_ref().is_some_and(|current| {
188                current.holder_core_id == lease.holder_core_id && current.epoch == lease.epoch
189            });
190            if !matches_current {
191                return Err(ControlPlaneError::Conflict(
192                    "service lease release does not match current epoch and holder".to_string(),
193                ));
194            }
195            slot.lease = None;
196            Ok(())
197        })
198    }
199}
200
201fn lease_expiration(now_ms: u64, ttl_ms: u64) -> ControlPlaneResult<u64> {
202    if ttl_ms == 0 {
203        return Err(ControlPlaneError::Rejected(
204            "lease ttl must be greater than zero".to_string(),
205        ));
206    }
207    now_ms.checked_add(ttl_ms).ok_or_else(|| {
208        ControlPlaneError::Rejected("lease expiration exceeds the clock range".to_string())
209    })
210}
211
212fn instance_key(identity: &CoreIdentity) -> String {
213    format!(
214        "{}:{}:{}",
215        identity.tenant_id.as_str(),
216        identity.cluster_id.as_str(),
217        identity.instance_id.as_str()
218    )
219}
220
221fn service_lease_key(
222    tenant_id: &TenantId,
223    cluster_id: &ClusterId,
224    service_id: &ServiceId,
225) -> String {
226    format!(
227        "{}:{}:{}",
228        tenant_id.as_str(),
229        cluster_id.as_str(),
230        service_id.as_str()
231    )
232}
233
234fn lock_state(
235    state: &Mutex<InMemoryState>,
236) -> ControlPlaneResult<std::sync::MutexGuard<'_, InMemoryState>> {
237    state
238        .lock()
239        .map_err(|_| ControlPlaneError::Transport("control plane state poisoned".to_string()))
240}
241
242fn is_routable(mode: RuntimeOperationalMode) -> bool {
243    matches!(
244        mode,
245        RuntimeOperationalMode::ReadWrite
246            | RuntimeOperationalMode::ReadOnly
247            | RuntimeOperationalMode::Syncing
248    )
249}