appcore_control_plane/
memory.rs1use super::*;
12
13#[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 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(®istration_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}