1use crate::policy::enforce_requirements;
14use crate::{
15 CapabilityError, CapabilityProvider, CapabilityRegistry, CapabilityRequest, CapabilityResponse,
16 CapabilityResult, CapabilitySelectionPolicy, LocalCapabilityProvider, RemoteCapabilityInvoker,
17};
18use appcore_contracts::ServiceId;
19use appcore_core::{
20 CapabilityDescriptor, CapabilityRequirements, CoreCompatibilityPolicy, CoreIdentity,
21};
22use appcore_distributed_contracts::{PeerRecord, ServiceLeadershipGuard};
23use std::sync::Arc;
24
25pub struct CapabilityResolver {
27 registry: CapabilityRegistry,
28 peers: Vec<PeerRecord>,
29 selector: Option<Arc<dyn CapabilitySelectionPolicy>>,
30}
31
32enum ExecutionProvider<'a> {
33 Local(&'a LocalCapabilityProvider),
34 Remote {
35 peer: &'a PeerRecord,
36 descriptor: &'a CapabilityDescriptor,
37 },
38 Selected(CapabilityProvider),
39}
40
41impl ExecutionProvider<'_> {
42 fn descriptor(&self) -> &CapabilityDescriptor {
43 match self {
44 Self::Local(provider) => provider.descriptor(),
45 Self::Remote { descriptor, .. } => descriptor,
46 Self::Selected(provider) => provider.descriptor(),
47 }
48 }
49
50 fn core_id<'a>(&'a self, identity: &'a CoreIdentity) -> &'a appcore_core::CoreId {
51 match self {
52 Self::Local(_) => &identity.core_id,
53 Self::Remote { peer, .. } => &peer.identity.core_id,
54 Self::Selected(provider) => provider.core_id(),
55 }
56 }
57}
58
59impl CapabilityResolver {
60 pub fn new(registry: CapabilityRegistry) -> Self {
62 Self {
63 registry,
64 peers: Vec::new(),
65 selector: None,
66 }
67 }
68
69 pub fn with_peers(mut self, peers: Vec<PeerRecord>) -> Self {
71 self.peers = peers;
72 self
73 }
74
75 pub fn with_selector(mut self, selector: Arc<dyn CapabilitySelectionPolicy>) -> Self {
77 self.selector = Some(selector);
78 self
79 }
80
81 pub fn resolve(
83 &self,
84 identity: &CoreIdentity,
85 service_id: &ServiceId,
86 request: &CapabilityRequest,
87 leadership: Option<&dyn ServiceLeadershipGuard>,
88 now_ms: u64,
89 ) -> CapabilityResult<CapabilityProvider> {
90 let provider = match &self.selector {
91 Some(selector) => selector.select(&self.candidates(identity, request)),
92 None => self.select_default(identity, request),
93 };
94 let Some(provider) = provider else {
95 return Err(CapabilityError::ProviderUnavailable(
96 request.capability.clone(),
97 ));
98 };
99 enforce_requirements(
100 identity,
101 service_id,
102 request,
103 provider.descriptor(),
104 provider.core_id(),
105 leadership,
106 true,
107 now_ms,
108 )?;
109 Ok(provider)
110 }
111
112 pub fn handle_local(
114 &self,
115 identity: &CoreIdentity,
116 service_id: &ServiceId,
117 request: &CapabilityRequest,
118 leadership: Option<&dyn ServiceLeadershipGuard>,
119 now_ms: u64,
120 ) -> CapabilityResult<CapabilityResponse> {
121 let provider =
122 self.resolve_for_execution(identity, service_id, request, leadership, now_ms)?;
123 match provider {
124 ExecutionProvider::Local(local) => local.handle(request),
125 ExecutionProvider::Remote { peer, .. } => Ok(CapabilityResponse::accepted(
126 Vec::new(),
127 Some(peer.identity.core_id.clone()),
128 )),
129 ExecutionProvider::Selected(CapabilityProvider::Local { .. }) => {
130 let Some(local) = self.registry.get(&request.capability) else {
131 return Err(CapabilityError::HandlerNotFound(request.capability.clone()));
132 };
133 local.handle(request)
134 }
135 ExecutionProvider::Selected(CapabilityProvider::Remote { peer, .. }) => Ok(
136 CapabilityResponse::accepted(Vec::new(), Some(peer.identity.core_id.clone())),
137 ),
138 }
139 }
140
141 pub fn handle(
143 &self,
144 identity: &CoreIdentity,
145 service_id: &ServiceId,
146 request: &CapabilityRequest,
147 leadership: Option<&dyn ServiceLeadershipGuard>,
148 remote_invoker: Option<&dyn RemoteCapabilityInvoker>,
149 now_ms: u64,
150 ) -> CapabilityResult<CapabilityResponse> {
151 let provider =
152 self.resolve_for_execution(identity, service_id, request, leadership, now_ms)?;
153 match provider {
154 ExecutionProvider::Local(local) => local.handle(request),
155 ExecutionProvider::Remote { peer, .. } => {
156 let Some(invoker) = remote_invoker else {
157 return Err(CapabilityError::RemoteEndpointUnavailable(
158 request.capability.clone(),
159 ));
160 };
161 invoker.invoke_remote(peer, request)
162 }
163 ExecutionProvider::Selected(CapabilityProvider::Local { .. }) => {
164 let Some(local) = self.registry.get(&request.capability) else {
165 return Err(CapabilityError::HandlerNotFound(request.capability.clone()));
166 };
167 local.handle(request)
168 }
169 ExecutionProvider::Selected(CapabilityProvider::Remote { peer, .. }) => {
170 let Some(invoker) = remote_invoker else {
171 return Err(CapabilityError::RemoteEndpointUnavailable(
172 request.capability.clone(),
173 ));
174 };
175 invoker.invoke_remote(&peer, request)
176 }
177 }
178 }
179
180 pub fn handle_owned(
185 &self,
186 identity: &CoreIdentity,
187 service_id: &ServiceId,
188 request: CapabilityRequest,
189 leadership: Option<&dyn ServiceLeadershipGuard>,
190 remote_invoker: Option<&dyn RemoteCapabilityInvoker>,
191 now_ms: u64,
192 ) -> CapabilityResult<CapabilityResponse> {
193 let provider =
194 self.resolve_for_execution(identity, service_id, &request, leadership, now_ms)?;
195 match provider {
196 ExecutionProvider::Local(local) => local.handle(&request),
197 ExecutionProvider::Remote { peer, .. } => {
198 let Some(invoker) = remote_invoker else {
199 return Err(CapabilityError::RemoteEndpointUnavailable(
200 request.capability.clone(),
201 ));
202 };
203 invoker.invoke_remote_owned(peer, request)
204 }
205 ExecutionProvider::Selected(CapabilityProvider::Local { .. }) => {
206 let Some(local) = self.registry.get(&request.capability) else {
207 return Err(CapabilityError::HandlerNotFound(request.capability.clone()));
208 };
209 local.handle(&request)
210 }
211 ExecutionProvider::Selected(CapabilityProvider::Remote { peer, .. }) => {
212 let Some(invoker) = remote_invoker else {
213 return Err(CapabilityError::RemoteEndpointUnavailable(
214 request.capability.clone(),
215 ));
216 };
217 invoker.invoke_remote_owned(&peer, request)
218 }
219 }
220 }
221
222 fn resolve_for_execution<'a>(
223 &'a self,
224 identity: &CoreIdentity,
225 service_id: &ServiceId,
226 request: &CapabilityRequest,
227 leadership: Option<&dyn ServiceLeadershipGuard>,
228 now_ms: u64,
229 ) -> CapabilityResult<ExecutionProvider<'a>> {
230 let provider = match &self.selector {
231 Some(selector) => selector
232 .select(&self.candidates(identity, request))
233 .map(ExecutionProvider::Selected),
234 None => self.select_default_for_execution(identity, request),
235 };
236 let Some(provider) = provider else {
237 return Err(CapabilityError::ProviderUnavailable(
238 request.capability.clone(),
239 ));
240 };
241 enforce_requirements(
242 identity,
243 service_id,
244 request,
245 provider.descriptor(),
246 provider.core_id(identity),
247 leadership,
248 true,
249 now_ms,
250 )?;
251 Ok(provider)
252 }
253
254 fn select_default_for_execution<'a>(
255 &'a self,
256 identity: &CoreIdentity,
257 request: &CapabilityRequest,
258 ) -> Option<ExecutionProvider<'a>> {
259 if let Some(local) = self.registry.get(&request.capability) {
260 if local.is_healthy() && local.descriptor().mode == request.mode {
261 return Some(ExecutionProvider::Local(local));
262 }
263 }
264
265 let mut fallback = None;
266 for peer in &self.peers {
267 let Some((descriptor, preferred)) = compatible_remote(identity, request, peer) else {
268 continue;
269 };
270 if preferred {
271 return Some(ExecutionProvider::Remote { peer, descriptor });
272 }
273 fallback.get_or_insert((peer, descriptor));
274 }
275 fallback.map(|(peer, descriptor)| ExecutionProvider::Remote { peer, descriptor })
276 }
277
278 fn select_default(
279 &self,
280 identity: &CoreIdentity,
281 request: &CapabilityRequest,
282 ) -> Option<CapabilityProvider> {
283 if let Some(local) = self.registry.get(&request.capability) {
284 if local.is_healthy() && local.descriptor().mode == request.mode {
285 return Some(CapabilityProvider::Local {
286 core_id: identity.core_id.clone(),
287 descriptor: local.descriptor().clone(),
288 });
289 }
290 }
291
292 let mut fallback = None;
293 for peer in &self.peers {
294 let Some((descriptor, preferred)) = compatible_remote(identity, request, peer) else {
295 continue;
296 };
297 if preferred {
298 return Some(remote_provider(peer, descriptor, true));
299 }
300 fallback.get_or_insert((peer, descriptor));
301 }
302 fallback.map(|(peer, descriptor)| remote_provider(peer, descriptor, false))
303 }
304
305 fn candidates(
306 &self,
307 identity: &CoreIdentity,
308 request: &CapabilityRequest,
309 ) -> Vec<CapabilityProvider> {
310 let mut candidates = Vec::new();
311 if let Some(local) = self.registry.get(&request.capability) {
312 if local.is_healthy() && local.descriptor().mode == request.mode {
313 candidates.push(CapabilityProvider::Local {
314 core_id: identity.core_id.clone(),
315 descriptor: local.descriptor().clone(),
316 });
317 }
318 }
319
320 for peer in &self.peers {
321 if let Some((descriptor, preferred)) = compatible_remote(identity, request, peer) {
322 candidates.push(CapabilityProvider::Remote {
323 peer: Box::new(peer.clone()),
324 descriptor: descriptor.clone(),
325 preferred,
326 });
327 }
328 }
329 candidates
330 }
331}
332
333fn compatible_remote<'a>(
334 identity: &CoreIdentity,
335 request: &CapabilityRequest,
336 peer: &'a PeerRecord,
337) -> Option<(&'a CapabilityDescriptor, bool)> {
338 if !peer.healthy {
339 return None;
340 }
341 let descriptor = peer.capabilities.iter().find(|descriptor| {
342 descriptor.name == request.capability && descriptor.mode == request.mode
343 })?;
344 if !remote_descriptor_is_compatible(identity, peer, descriptor) {
345 return None;
346 }
347 let preferred = peer
348 .metadata
349 .get("preferred")
350 .is_some_and(|value| value == "true");
351 Some((descriptor, preferred))
352}
353
354fn remote_provider(
355 peer: &PeerRecord,
356 descriptor: &CapabilityDescriptor,
357 preferred: bool,
358) -> CapabilityProvider {
359 CapabilityProvider::Remote {
360 peer: Box::new(peer.clone()),
361 descriptor: descriptor.clone(),
362 preferred,
363 }
364}
365
366fn remote_descriptor_is_compatible(
367 identity: &CoreIdentity,
368 peer: &PeerRecord,
369 descriptor: &CapabilityDescriptor,
370) -> bool {
371 let require_same_cluster = match descriptor.visibility {
372 appcore_core::CapabilityVisibility::Local => return false,
373 appcore_core::CapabilityVisibility::Cluster => true,
374 appcore_core::CapabilityVisibility::Tenant => false,
375 };
376 let policy = CoreCompatibilityPolicy {
377 require_same_cluster,
378 required_capability: None,
379 };
380 identity
381 .ensure_compatible(&peer.identity, &policy, &[])
382 .is_ok()
383}
384
385pub fn requirements_for_read_only() -> CapabilityRequirements {
387 CapabilityRequirements {
388 requires_leader: false,
389 read_only: true,
390 idempotency_required: false,
391 }
392}