1use crate::policy::enforce_requirements;
12use crate::{
13 CapabilityError, CapabilityProvider, CapabilityRegistry, CapabilityRequest, CapabilityResponse,
14 CapabilityResult, CapabilitySelectionPolicy, DefaultCapabilitySelectionPolicy,
15 RemoteCapabilityInvoker,
16};
17use appcore_contracts::ServiceId;
18use appcore_core::{
19 CapabilityDescriptor, CapabilityRequirements, CoreCompatibilityPolicy, CoreIdentity,
20};
21use appcore_distributed_contracts::{PeerRecord, ServiceLeadershipGuard};
22use std::sync::Arc;
23
24pub struct CapabilityResolver {
26 registry: CapabilityRegistry,
27 peers: Vec<PeerRecord>,
28 selector: Arc<dyn CapabilitySelectionPolicy>,
29}
30
31impl CapabilityResolver {
32 pub fn new(registry: CapabilityRegistry) -> Self {
34 Self {
35 registry,
36 peers: Vec::new(),
37 selector: Arc::new(DefaultCapabilitySelectionPolicy::default()),
38 }
39 }
40
41 pub fn with_peers(mut self, peers: Vec<PeerRecord>) -> Self {
43 self.peers = peers;
44 self
45 }
46
47 pub fn with_selector(mut self, selector: Arc<dyn CapabilitySelectionPolicy>) -> Self {
49 self.selector = selector;
50 self
51 }
52
53 pub fn resolve(
55 &self,
56 identity: &CoreIdentity,
57 service_id: &ServiceId,
58 request: &CapabilityRequest,
59 leadership: Option<&dyn ServiceLeadershipGuard>,
60 now_ms: u64,
61 ) -> CapabilityResult<CapabilityProvider> {
62 let candidates = self.candidates(identity, request);
63 let Some(provider) = self.selector.select(&candidates) else {
64 return Err(CapabilityError::ProviderUnavailable(
65 request.capability.clone(),
66 ));
67 };
68 enforce_requirements(
69 identity,
70 service_id,
71 request,
72 provider.descriptor(),
73 provider.core_id(),
74 leadership,
75 true,
76 now_ms,
77 )?;
78 Ok(provider)
79 }
80
81 pub fn handle_local(
83 &self,
84 identity: &CoreIdentity,
85 service_id: &ServiceId,
86 request: &CapabilityRequest,
87 leadership: Option<&dyn ServiceLeadershipGuard>,
88 now_ms: u64,
89 ) -> CapabilityResult<CapabilityResponse> {
90 let provider = self.resolve(identity, service_id, request, leadership, now_ms)?;
91 match provider {
92 CapabilityProvider::Local { .. } => {
93 let Some(local) = self.registry.get(&request.capability) else {
94 return Err(CapabilityError::HandlerNotFound(request.capability.clone()));
95 };
96 local.handle(request)
97 }
98 CapabilityProvider::Remote { peer, .. } => Ok(CapabilityResponse::accepted(
99 Vec::new(),
100 Some(peer.identity.core_id.clone()),
101 )),
102 }
103 }
104
105 pub fn handle(
107 &self,
108 identity: &CoreIdentity,
109 service_id: &ServiceId,
110 request: &CapabilityRequest,
111 leadership: Option<&dyn ServiceLeadershipGuard>,
112 remote_invoker: Option<&dyn RemoteCapabilityInvoker>,
113 now_ms: u64,
114 ) -> CapabilityResult<CapabilityResponse> {
115 let provider = self.resolve(identity, service_id, request, leadership, now_ms)?;
116 match provider {
117 CapabilityProvider::Local { .. } => {
118 let Some(local) = self.registry.get(&request.capability) else {
119 return Err(CapabilityError::HandlerNotFound(request.capability.clone()));
120 };
121 local.handle(request)
122 }
123 CapabilityProvider::Remote { peer, .. } => {
124 let Some(invoker) = remote_invoker else {
125 return Err(CapabilityError::RemoteEndpointUnavailable(
126 request.capability.clone(),
127 ));
128 };
129 invoker.invoke_remote(&peer, request)
130 }
131 }
132 }
133
134 fn candidates(
135 &self,
136 identity: &CoreIdentity,
137 request: &CapabilityRequest,
138 ) -> Vec<CapabilityProvider> {
139 let mut candidates = Vec::new();
140 if let Some(local) = self.registry.get(&request.capability) {
141 if local.is_healthy() && local.descriptor().mode == request.mode {
142 candidates.push(CapabilityProvider::Local {
143 core_id: identity.core_id.clone(),
144 descriptor: local.descriptor().clone(),
145 });
146 }
147 }
148
149 for peer in &self.peers {
150 if !peer.healthy {
151 continue;
152 }
153 if let Some(descriptor) = peer.capabilities.iter().find(|descriptor| {
154 descriptor.name == request.capability && descriptor.mode == request.mode
155 }) {
156 if !remote_descriptor_is_compatible(identity, peer, descriptor) {
157 continue;
158 }
159 let preferred = peer
160 .metadata
161 .get("preferred")
162 .map(|value| value == "true")
163 .unwrap_or(false);
164 candidates.push(CapabilityProvider::Remote {
165 peer: Box::new(peer.clone()),
166 descriptor: descriptor.clone(),
167 preferred,
168 });
169 }
170 }
171 candidates
172 }
173}
174
175fn remote_descriptor_is_compatible(
176 identity: &CoreIdentity,
177 peer: &PeerRecord,
178 descriptor: &CapabilityDescriptor,
179) -> bool {
180 let require_same_cluster = match descriptor.visibility {
181 appcore_core::CapabilityVisibility::Local => return false,
182 appcore_core::CapabilityVisibility::Cluster => true,
183 appcore_core::CapabilityVisibility::Tenant => false,
184 };
185 let policy = CoreCompatibilityPolicy {
186 require_same_cluster,
187 required_capability: Some(descriptor.name.clone()),
188 };
189 identity
190 .ensure_compatible(
191 &peer.identity,
192 &policy,
193 &peer
194 .capabilities
195 .iter()
196 .map(|capability| capability.name.clone())
197 .collect::<Vec<_>>(),
198 )
199 .is_ok()
200}
201
202pub fn requirements_for_read_only() -> CapabilityRequirements {
204 CapabilityRequirements {
205 requires_leader: false,
206 read_only: true,
207 idempotency_required: false,
208 }
209}