Skip to main content

ferrum_interfaces/vnext/execution/
resolution.rs

1use 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    /// Resolves one node through the typed planning registry. Provider
77    /// selection is core-owned; `preferred_provider` is only a compatibility
78    /// preference and cannot inject a provider or resource estimate.
79    #[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}