ferrum_interfaces/vnext/execution/
resolution.rs1use super::{
2 invalid_plan, node_weight_requirements, provider_resource_estimator_input_fingerprint,
3 validate_active_sequence_ceiling, validate_program_bindings, validate_scheduled_token_ceiling,
4 validate_semantic_binding, workspace_base_id, BTreeMap, BTreeSet, BufferUsage,
5 CapabilityCatalog, CapabilityId, DynamicStorageRequirement, ExecutionWeightPlan, NodeId,
6 OperationPlanningHandle, OperationPlanningRegistry, OperationRegistryAuthority,
7 OperationResourceEstimateRequest, PlanNodeResolution, PlanProviderRejectReason,
8 PreparedModelFamily, ProviderCompatibilityRequest, ProviderId, ProviderResourcePlan,
9 ResolvedValueBinding, RuntimePolicy, VNextError, WeightSchema,
10};
11
12impl PlanNodeResolution {
13 pub(super) fn from_provider_resolution(
14 operation_registry_authority: OperationRegistryAuthority,
15 node_id: NodeId,
16 values: Vec<ResolvedValueBinding>,
17 required_capabilities: BTreeSet<CapabilityId>,
18 preferred_provider: Option<ProviderId>,
19 mut provider_resource_candidates: Vec<ProviderResourcePlan>,
20 provider_resolution_rejections: BTreeMap<ProviderId, PlanProviderRejectReason>,
21 ) -> Result<Self, VNextError> {
22 if values.is_empty() {
23 return Err(invalid_plan(format!(
24 "node `{node_id}` has no physical value resolution"
25 )));
26 }
27 provider_resource_candidates
28 .sort_by(|left, right| left.provider_id.cmp(&right.provider_id));
29 if provider_resource_candidates.is_empty()
30 || provider_resource_candidates
31 .windows(2)
32 .any(|pair| pair[0].provider_id == pair[1].provider_id)
33 {
34 return Err(invalid_plan(format!(
35 "node `{node_id}` has empty or duplicate provider resource candidates"
36 )));
37 }
38 for candidate in &provider_resource_candidates {
39 candidate.validate_shape()?;
40 if provider_resolution_rejections.contains_key(candidate.provider_id()) {
41 return Err(invalid_plan(format!(
42 "provider `{}` is both a resource candidate and rejected resolution",
43 candidate.provider_id()
44 )));
45 }
46 }
47 for reason in provider_resolution_rejections.values() {
48 match reason {
49 PlanProviderRejectReason::StorageIncompatible { resource_ids }
50 if !resource_ids.is_empty()
51 && !resource_ids.windows(2).any(|pair| pair[0] >= pair[1]) => {}
52 PlanProviderRejectReason::StorageIncompatible { .. } => {
53 return Err(invalid_plan(
54 "provider storage rejection resources are empty or non-canonical",
55 ))
56 }
57 PlanProviderRejectReason::NotRegistered
58 | PlanProviderRejectReason::Incompatible(_) => {
59 return Err(invalid_plan(
60 "resolution-local rejection must describe storage incompatibility",
61 ))
62 }
63 }
64 }
65 Ok(Self {
66 operation_registry_authority,
67 node_id,
68 values,
69 required_capabilities,
70 preferred_provider,
71 provider_resource_candidates,
72 provider_resolution_rejections,
73 })
74 }
75
76 #[allow(clippy::too_many_arguments)]
80 pub fn resolve<P: RuntimePolicy>(
81 family: &PreparedModelFamily,
82 catalog: &CapabilityCatalog,
83 policy: &P,
84 registry: &OperationPlanningHandle<'_>,
85 node_id: NodeId,
86 values: Vec<ResolvedValueBinding>,
87 required_capabilities: BTreeSet<CapabilityId>,
88 preferred_provider: Option<ProviderId>,
89 ) -> Result<Self, VNextError> {
90 let prepared_family_fingerprint = family.fingerprint()?;
91 let execution_weights = ExecutionWeightPlan::identity(family)?;
92 Self::resolve_with_family_fingerprint(
93 family,
94 execution_weights.schema(),
95 &prepared_family_fingerprint,
96 catalog,
97 policy,
98 registry,
99 node_id,
100 values,
101 required_capabilities,
102 preferred_provider,
103 )
104 }
105
106 #[allow(clippy::too_many_arguments)]
107 pub(super) fn resolve_with_family_fingerprint<P: RuntimePolicy>(
108 family: &PreparedModelFamily,
109 execution_weight_schema: &WeightSchema,
110 prepared_family_fingerprint: &str,
111 catalog: &CapabilityCatalog,
112 policy: &P,
113 registry: &OperationPlanningHandle<'_>,
114 node_id: NodeId,
115 values: Vec<ResolvedValueBinding>,
116 required_capabilities: BTreeSet<CapabilityId>,
117 preferred_provider: Option<ProviderId>,
118 ) -> Result<Self, VNextError> {
119 policy.validate()?;
120 validate_active_sequence_ceiling(policy.maximum_active_sequences())?;
121 validate_scheduled_token_ceiling(policy.maximum_scheduled_tokens())?;
122 if values.is_empty() {
123 return Err(invalid_plan(format!(
124 "node `{node_id}` has no physical value resolution"
125 )));
126 }
127 let program_node = family
128 .program()
129 .blocks()
130 .iter()
131 .flat_map(|block| &block.nodes)
132 .find(|node| node.id == node_id)
133 .ok_or_else(|| invalid_plan(format!("program has no node `{node_id}`")))?;
134 let operation = catalog.operation_for_node(&program_node.id, &program_node.operation_id)?;
135
136 let contracts = registry.contracts_for(&program_node.operation_id);
137 if contracts.len() != 1 {
138 return Err(invalid_plan(format!(
139 "operation `{}` requires exactly one typed contract registration, found {}",
140 program_node.operation_id,
141 contracts.len()
142 )));
143 }
144 let contract = contracts[0];
145 if contract.descriptor() != operation {
146 return Err(invalid_plan(format!(
147 "typed contract for operation `{}` differs from the capability catalog",
148 program_node.operation_id
149 )));
150 }
151 contract.validate_signature(&operation.inputs, &operation.outputs)?;
152 operation.validate_attributes(&program_node.attributes)?;
153 operation.validate_resolved_bindings(&values)?;
154 validate_program_bindings(program_node, &values)?;
155 for binding in &values {
156 validate_semantic_binding(family, execution_weight_schema, binding)?;
157 }
158
159 let (required_weight_formats, required_quantization_formats) =
160 node_weight_requirements(family, &values)?;
161 let compatibility_request = ProviderCompatibilityRequest::new(
162 program_node.operation_id.clone(),
163 program_node.required_version,
164 operation
165 .provider
166 .required_capabilities
167 .union(&required_capabilities)
168 .cloned()
169 .collect(),
170 required_weight_formats,
171 required_quantization_formats,
172 policy.execution_determinism_requirement(),
173 )?;
174 let report = catalog.provider_compatibility(compatibility_request)?;
175 report.require_compatible_for_node(&catalog.device().id, &program_node.id)?;
176 let profile_available = |requirement: &DynamicStorageRequirement| {
177 policy
178 .dynamic_storage_profile_order()
179 .iter()
180 .any(|profile| {
181 catalog.device().dynamic_storage_profiles.contains(profile)
182 && requirement.accepts(*profile)
183 })
184 };
185 let mut provider_resource_candidates = Vec::new();
186 let mut provider_resolution_rejections = BTreeMap::new();
187 for provider_id in report.compatible_provider_ids() {
188 let descriptor = catalog
189 .providers_for_node(&program_node.id, &program_node.operation_id)?
190 .iter()
191 .find(|provider| provider.provider_id() == provider_id)
192 .ok_or_else(|| invalid_plan("compatible provider is absent from the catalog"))?;
193 let value_storage_conflicts = values
194 .iter()
195 .filter(|binding| binding.usage() != BufferUsage::Weights)
196 .filter(|binding| {
197 descriptor
198 .dynamic_storage_for(binding.role(), binding.ordinal())
199 .is_none_or(|requirement| !profile_available(requirement))
200 })
201 .flat_map(|binding| binding.storage().components())
202 .map(|component| component.resource_id().clone())
203 .collect::<BTreeSet<_>>();
204 if !value_storage_conflicts.is_empty() {
205 provider_resolution_rejections.insert(
206 provider_id.clone(),
207 PlanProviderRejectReason::StorageIncompatible {
208 resource_ids: value_storage_conflicts.into_iter().collect(),
209 },
210 );
211 continue;
212 }
213 let estimators = registry.estimators_for(provider_id);
214 if estimators.len() != 1 || estimators[0].descriptor() != descriptor {
215 return Err(invalid_plan(format!(
216 "compatible provider `{provider_id}` lacks one exact resource estimator registration"
217 )));
218 }
219 let estimator_input_fingerprint = provider_resource_estimator_input_fingerprint(
220 family,
221 prepared_family_fingerprint,
222 operation,
223 program_node,
224 provider_id,
225 &values,
226 &required_capabilities,
227 )?;
228 let estimate_request = OperationResourceEstimateRequest::new(
229 &program_node.id,
230 operation,
231 &values,
232 &program_node.attributes,
233 &estimator_input_fingerprint,
234 )?;
235 let estimate = estimators[0].estimate_resources(estimate_request)?;
236 let provider_resources = ProviderResourcePlan::from_provider_output(
237 descriptor,
238 &estimator_input_fingerprint,
239 estimate,
240 )?;
241 let minimum_alignment = operation.resources.minimum_value_alignment_bytes;
242 if provider_resources.value_alignment_bytes() < minimum_alignment
243 || provider_resources.value_alignment_bytes() % minimum_alignment != 0
244 || !operation
245 .resources
246 .scratch
247 .accepts(provider_resources.scratch().is_some())
248 || !operation
249 .resources
250 .binding
251 .accepts(provider_resources.binding().is_some())
252 || !operation
253 .resources
254 .persistent
255 .accepts(provider_resources.persistent().is_some())
256 {
257 return Err(invalid_plan(format!(
258 "compatible provider `{provider_id}` returned resources outside its operation contract"
259 )));
260 }
261 let mut workspace_storage_conflicts = BTreeSet::new();
262 for (kind, workspace) in [
263 ("scratch", provider_resources.scratch()),
264 ("binding", provider_resources.binding()),
265 ("persistent", provider_resources.persistent()),
266 ] {
267 if workspace.is_some_and(|workspace| !profile_available(workspace.storage())) {
268 workspace_storage_conflicts.insert(workspace_base_id(
269 &program_node.id,
270 kind,
271 provider_resources.estimate_fingerprint(),
272 )?);
273 }
274 }
275 if !workspace_storage_conflicts.is_empty() {
276 provider_resolution_rejections.insert(
277 provider_id.clone(),
278 PlanProviderRejectReason::StorageIncompatible {
279 resource_ids: workspace_storage_conflicts.into_iter().collect(),
280 },
281 );
282 continue;
283 }
284 provider_resource_candidates.push(provider_resources);
285 }
286 if provider_resource_candidates.is_empty() {
287 return Err(invalid_plan(format!(
288 "node `{node_id}` has no provider whose binding and workspace storage requirements intersect runtime offers and policy"
289 )));
290 }
291 Self::from_provider_resolution(
292 registry.authority().clone(),
293 node_id,
294 values,
295 required_capabilities,
296 preferred_provider,
297 provider_resource_candidates,
298 provider_resolution_rejections,
299 )
300 }
301
302 pub fn node_id(&self) -> &NodeId {
303 &self.node_id
304 }
305
306 pub fn values(&self) -> &[ResolvedValueBinding] {
307 &self.values
308 }
309
310 pub fn required_capabilities(&self) -> &BTreeSet<CapabilityId> {
311 &self.required_capabilities
312 }
313
314 pub fn preferred_provider(&self) -> Option<&ProviderId> {
315 self.preferred_provider.as_ref()
316 }
317
318 pub fn provider_resource_candidates(&self) -> &[ProviderResourcePlan] {
319 &self.provider_resource_candidates
320 }
321
322 pub fn provider_resolution_rejections(
323 &self,
324 ) -> &BTreeMap<ProviderId, PlanProviderRejectReason> {
325 &self.provider_resolution_rejections
326 }
327}