Skip to main content

blut_graph_core/
compile.rs

1// SPDX-License-Identifier: AGPL-3.0-or-later
2
3use alloc::collections::{BTreeMap, BTreeSet};
4use alloc::string::{String, ToString};
5use alloc::vec::Vec;
6use core::fmt;
7
8use crate::model::{
9    AuthorizedPlan, BufferId, BufferPlan, CompiledNode, CompiledPlan, CompiledPortContract, Edge,
10    ExecutionRealm, Graph, GraphId, KernelDescriptor, KernelId, Layout, NodeDescriptor, NodeId,
11    NodeTypeRef, OutputBinding, PlanId, PortDescriptor, StateContract, StateScope, StepId, Target,
12};
13
14#[derive(Clone, Debug, PartialEq, Eq)]
15pub enum CompileError {
16    UnsupportedGraphVersion(u32),
17    DuplicateNode(NodeId),
18    UnknownNode(NodeId),
19    UnknownDescriptor(String, u32),
20    DuplicateDescriptor(String, u32),
21    InvalidDescriptor(String, u32),
22    InvalidConfig(NodeId, crate::ConfigError),
23    DuplicateKernel(KernelId),
24    InvalidKernelContract(KernelId),
25    UnknownPort(NodeId, String),
26    InvalidPortSize(NodeId, String),
27    TypeMismatch(String, String),
28    PortCapacityMismatch(NodeId, String, NodeId, String),
29    LayoutUnavailable(NodeId, String),
30    MissingInput(NodeId, String),
31    DuplicateInput(NodeId, String),
32    DuplicateInvocation(NodeId, String),
33    Cycle,
34    CapabilityMissing(String),
35    CapabilityUnsupported(String),
36    TargetUnsupported(NodeId, Target),
37    KernelUnavailable(NodeId, Target),
38    ProofMissing(NodeId, String),
39    PolicyMissing(NodeId, String),
40    FidelityInsufficient(NodeId),
41    UnsafeRetry(NodeId),
42    SearchLimitExceeded,
43    CompileLimitExceeded,
44    ResourceOverflow,
45    EmptyGraph,
46    InvalidGraphContract,
47    InvalidState(NodeId),
48    InvalidSession,
49    InvalidFeedback(NodeId, String),
50    PortContractMismatch(NodeId, String, NodeId, String),
51    UnknownSubgraph(crate::SubgraphId),
52    InvalidSubgraph(crate::SubgraphId),
53    SubgraphDepthExceeded,
54    SubgraphEntryLimitExceeded,
55}
56
57impl fmt::Display for CompileError {
58    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
59        write!(f, "{self:?}")
60    }
61}
62
63#[cfg(feature = "std")]
64impl std::error::Error for CompileError {}
65
66#[derive(Clone, Debug, Default)]
67pub struct KernelRegistry {
68    descriptors: BTreeMap<(String, u32), NodeDescriptor>,
69    kernels: BTreeMap<KernelId, KernelDescriptor>,
70    subgraphs: BTreeMap<crate::SubgraphId, crate::SubgraphSchema>,
71}
72
73impl KernelRegistry {
74    /// Return exact normalized semantic contract trusted by this registry.
75    ///
76    /// Product runtimes use this after compilation to ensure live policy grants
77    /// are required by the semantic node they authorize.
78    pub fn descriptor(&self, node_type: &NodeTypeRef) -> Option<&NodeDescriptor> {
79        self.descriptors
80            .get(&(node_type.type_name.clone(), node_type.version))
81    }
82
83    pub fn register_descriptor(
84        &mut self,
85        mut descriptor: NodeDescriptor,
86    ) -> Result<(), CompileError> {
87        for port in descriptor
88            .inputs
89            .iter_mut()
90            .chain(descriptor.outputs.iter_mut())
91        {
92            port.layouts.sort_unstable();
93            port.layouts.dedup();
94            normalize_contract(&mut port.proof, &mut port.policy);
95        }
96        descriptor.capabilities.sort_unstable();
97        descriptor.capabilities.dedup();
98        descriptor.targets.sort_unstable();
99        descriptor.targets.dedup();
100        descriptor.proof.requires.sort_unstable();
101        descriptor.proof.requires.dedup();
102        descriptor.proof.provides.sort_unstable();
103        descriptor.proof.provides.dedup();
104        descriptor.proof.invalidates.sort_unstable();
105        descriptor.proof.invalidates.dedup();
106        descriptor.policy.requires.sort_unstable();
107        descriptor.policy.requires.dedup();
108        descriptor.policy.adds.sort_unstable();
109        descriptor.policy.adds.dedup();
110        descriptor.failure.domains.sort_unstable();
111        descriptor.failure.domains.dedup();
112        if let Some(lowering) = &mut descriptor.subgraph {
113            lowering.input_map.sort_unstable();
114            lowering.input_map.dedup();
115            lowering.output_map.sort_unstable();
116            lowering.output_map.dedup();
117            lowering.config_map.sort_unstable();
118            lowering.config_map.dedup();
119        }
120        descriptor.config.normalize().map_err(|_| {
121            CompileError::InvalidDescriptor(descriptor.type_name.clone(), descriptor.version)
122        })?;
123        let key = (descriptor.type_name.clone(), descriptor.version);
124        if self.descriptors.contains_key(&key) {
125            return Err(CompileError::DuplicateDescriptor(key.0, key.1));
126        }
127        fn invalid_ports(ports: &[PortDescriptor]) -> bool {
128            let mut names = BTreeSet::new();
129            ports.iter().any(|port| {
130                port.name.is_empty()
131                    || port.semantic_type.is_empty()
132                    || port.max_bytes == 0
133                    || port.layouts.is_empty()
134                    || !valid_port_contract(port)
135                    || !names.insert(port.name.as_str())
136            })
137        }
138        if descriptor.type_name.is_empty()
139            || descriptor.version == 0
140            || descriptor.resources.threads == 0
141            || !valid_contract_names(&descriptor.proof, &descriptor.policy)
142            || invalid_ports(&descriptor.inputs)
143            || invalid_ports(&descriptor.outputs)
144            || (descriptor.effect == crate::model::Effect::AtMostOnce && descriptor.retry_limit > 0)
145            || descriptor
146                .failure
147                .domains
148                .iter()
149                .any(|domain| domain.is_empty())
150            || (descriptor.partiality == crate::model::Partiality::ExplicitGaps
151                && descriptor.failure.domains.is_empty())
152            || !valid_state_contract(&descriptor.state)
153        {
154            return Err(CompileError::InvalidDescriptor(key.0, key.1));
155        }
156        if descriptor.subgraph.as_ref().is_some_and(|lowering| {
157            validate_subgraph_lowering(&descriptor, lowering, &self.subgraphs, &self.descriptors)
158                .is_err()
159        }) {
160            return Err(CompileError::InvalidDescriptor(key.0, key.1));
161        }
162        self.descriptors.insert(key, descriptor);
163        Ok(())
164    }
165
166    pub fn register_subgraph(
167        &mut self,
168        mut schema: crate::SubgraphSchema,
169    ) -> Result<(), CompileError> {
170        schema.nodes.sort_by_key(|node| node.id);
171        schema
172            .edges
173            .sort_by_key(|edge| (edge.from.clone(), edge.to.clone()));
174        schema.inputs.sort_unstable();
175        schema.outputs.sort_unstable();
176        let duplicate_nodes = schema.nodes.windows(2).any(|pair| pair[0].id == pair[1].id);
177        let duplicate_edges = schema.edges.windows(2).any(|pair| pair[0] == pair[1]);
178        let duplicate_interface = |ports: &[crate::SubgraphInterfacePort]| {
179            ports.windows(2).any(|pair| pair[0].name == pair[1].name)
180                || ports
181                    .iter()
182                    .map(|port| &port.inner)
183                    .collect::<BTreeSet<_>>()
184                    .len()
185                    != ports.len()
186        };
187        if schema.version == 0
188            || schema.nodes.is_empty()
189            || duplicate_nodes
190            || duplicate_edges
191            || duplicate_interface(&schema.inputs)
192            || duplicate_interface(&schema.outputs)
193            || schema.id != subgraph_identity(&schema)
194            || self.subgraphs.contains_key(&schema.id)
195        {
196            return Err(CompileError::InvalidSubgraph(schema.id));
197        }
198        let mut descriptors = BTreeMap::new();
199        for node in &schema.nodes {
200            let descriptor = self
201                .descriptors
202                .get(&(node.node_type.type_name.clone(), node.node_type.version))
203                .ok_or(CompileError::InvalidSubgraph(schema.id))?;
204            let invalid_config = match descriptor.config.canonicalize(&node.config) {
205                Ok(canonical) => canonical != node.config,
206                Err(_) => true,
207            };
208            let invalid_child = match node.child {
209                Some(child) if child == schema.id => true,
210                Some(child) => self.subgraphs.get(&child).is_none_or(|child| {
211                    !subgraph_implements_descriptor(child, descriptor, &self.descriptors)
212                }),
213                None => false,
214            };
215            if invalid_config || invalid_child {
216                return Err(CompileError::InvalidSubgraph(schema.id));
217            }
218            descriptors.insert(node.id, descriptor);
219        }
220        let mut bound_inputs = BTreeSet::new();
221        for edge in &schema.edges {
222            let from = descriptors
223                .get(&edge.from.node)
224                .ok_or(CompileError::InvalidSubgraph(schema.id))?;
225            let to = descriptors
226                .get(&edge.to.node)
227                .ok_or(CompileError::InvalidSubgraph(schema.id))?;
228            let output = from
229                .outputs
230                .iter()
231                .find(|port| port.name == edge.from.port)
232                .ok_or(CompileError::InvalidSubgraph(schema.id))?;
233            let input = to
234                .inputs
235                .iter()
236                .find(|port| port.name == edge.to.port)
237                .ok_or(CompileError::InvalidSubgraph(schema.id))?;
238            if !port_contract_satisfies(output, input) || !bound_inputs.insert(edge.to.clone()) {
239                return Err(CompileError::InvalidSubgraph(schema.id));
240            }
241        }
242        for port in &schema.inputs {
243            let descriptor = descriptors
244                .get(&port.inner.node)
245                .ok_or(CompileError::InvalidSubgraph(schema.id))?;
246            if port.name.is_empty()
247                || !descriptor
248                    .inputs
249                    .iter()
250                    .any(|input| input.name == port.inner.port)
251                || !bound_inputs.insert(port.inner.clone())
252            {
253                return Err(CompileError::InvalidSubgraph(schema.id));
254            }
255        }
256        for port in &schema.outputs {
257            let descriptor = descriptors
258                .get(&port.inner.node)
259                .ok_or(CompileError::InvalidSubgraph(schema.id))?;
260            if port.name.is_empty()
261                || !descriptor
262                    .outputs
263                    .iter()
264                    .any(|output| output.name == port.inner.port)
265            {
266                return Err(CompileError::InvalidSubgraph(schema.id));
267            }
268        }
269        if descriptors.iter().any(|(node, descriptor)| {
270            descriptor.inputs.iter().any(|input| {
271                !input.optional
272                    && !bound_inputs.contains(&crate::PortRef {
273                        node: *node,
274                        port: input.name.clone(),
275                    })
276            })
277        }) || !subgraph_is_acyclic(&schema)
278        {
279            return Err(CompileError::InvalidSubgraph(schema.id));
280        }
281        self.subgraphs.insert(schema.id, schema);
282        Ok(())
283    }
284
285    pub fn register_kernel(&mut self, mut kernel: KernelDescriptor) -> Result<(), CompileError> {
286        kernel.input_layouts.sort_unstable();
287        kernel.input_layouts.dedup();
288        kernel.output_layouts.sort_unstable();
289        kernel.output_layouts.dedup();
290        let id = kernel.id;
291        if self.kernels.contains_key(&id) {
292            return Err(CompileError::DuplicateKernel(id));
293        }
294        let conversion_role = kernel.conversion.is_some();
295        if conversion_role != kernel.implements.is_empty()
296            || kernel.resources.threads == 0
297            || kernel.conversion.as_ref().is_some_and(|conversion| {
298                conversion.max_input_bytes == 0
299                    || conversion.max_output_bytes == 0
300                    || conversion.semantic_type.is_empty()
301                    || conversion.from == conversion.to
302                    || !kernel.input_layouts.contains(&conversion.from)
303                    || !kernel.output_layouts.contains(&conversion.to)
304                    || kernel.determinism != crate::model::Determinism::BitExact
305            })
306        {
307            return Err(CompileError::InvalidKernelContract(id));
308        }
309        self.kernels.insert(id, kernel);
310        Ok(())
311    }
312
313    /// Apply one outer node instance's canonical configuration to its declared
314    /// inner DAG. This is an explicit reference graph, not physical inlining:
315    /// callers compile it normally with fusion enabled or disabled.
316    pub fn materialize_subgraph(
317        &self,
318        instance: &crate::NodeInstance,
319    ) -> Result<crate::MaterializedSubgraph, CompileError> {
320        let descriptor = self
321            .descriptors
322            .get(&(instance.descriptor.clone(), instance.descriptor_version))
323            .ok_or_else(|| {
324                CompileError::UnknownDescriptor(
325                    instance.descriptor.clone(),
326                    instance.descriptor_version,
327                )
328            })?;
329        let lowering = descriptor
330            .subgraph
331            .as_ref()
332            .ok_or(CompileError::InvalidSubgraph(crate::SubgraphId([0; 32])))?;
333        validate_subgraph_lowering(descriptor, lowering, &self.subgraphs, &self.descriptors)?;
334        let schema = self
335            .subgraphs
336            .get(&lowering.subgraph)
337            .ok_or(CompileError::UnknownSubgraph(lowering.subgraph))?;
338        if schema.nodes.iter().any(|node| node.child.is_some()) {
339            return Err(CompileError::InvalidSubgraph(schema.id));
340        }
341        let outer_config = descriptor
342            .config
343            .canonicalize(&instance.config)
344            .map_err(|error| CompileError::InvalidConfig(instance.id, error))?;
345        let mut nodes = schema
346            .nodes
347            .iter()
348            .map(|node| crate::NodeInstance {
349                id: node.id,
350                descriptor: node.node_type.type_name.clone(),
351                descriptor_version: node.node_type.version,
352                config: node.config.clone(),
353            })
354            .collect::<Vec<_>>();
355        for binding in &lowering.config_map {
356            let value = outer_config
357                .get(&binding.outer)
358                .ok_or(CompileError::InvalidSubgraph(schema.id))?
359                .clone();
360            let node = nodes
361                .iter_mut()
362                .find(|node| node.id == binding.node)
363                .ok_or(CompileError::InvalidSubgraph(schema.id))?;
364            node.config.insert(binding.inner.clone(), value);
365        }
366        for node in &mut nodes {
367            let inner = self
368                .descriptors
369                .get(&(node.descriptor.clone(), node.descriptor_version))
370                .ok_or(CompileError::InvalidSubgraph(schema.id))?;
371            node.config = inner
372                .config
373                .canonicalize(&node.config)
374                .map_err(|_| CompileError::InvalidSubgraph(schema.id))?;
375        }
376
377        let map_interfaces =
378            |maps: &[crate::PortMap], interfaces: &[crate::SubgraphInterfacePort]| {
379                maps.iter()
380                    .map(|map| {
381                        interfaces
382                            .iter()
383                            .find(|interface| interface.name == map.inner)
384                            .map(|interface| crate::SubgraphInterfacePort {
385                                name: map.outer.clone(),
386                                inner: interface.inner.clone(),
387                            })
388                            .ok_or(CompileError::InvalidSubgraph(schema.id))
389                    })
390                    .collect::<Result<Vec<_>, _>>()
391            };
392        let mut inputs = map_interfaces(&lowering.input_map, &schema.inputs)?;
393        let mut outputs = map_interfaces(&lowering.output_map, &schema.outputs)?;
394        inputs.sort_unstable();
395        outputs.sort_unstable();
396        let mut required_capabilities = schema
397            .nodes
398            .iter()
399            .flat_map(|node| {
400                self.descriptors[&(node.node_type.type_name.clone(), node.node_type.version)]
401                    .capabilities
402                    .clone()
403            })
404            .collect::<Vec<_>>();
405        required_capabilities.sort_unstable();
406        required_capabilities.dedup();
407        Ok(crate::MaterializedSubgraph {
408            graph: crate::Graph {
409                version: 3,
410                nodes,
411                edges: schema.edges.clone(),
412                feedback: Vec::new(),
413                invocation_inputs: inputs
414                    .iter()
415                    .map(|interface| interface.inner.clone())
416                    .collect(),
417                required_capabilities,
418                required_proofs: descriptor.proof.requires.clone(),
419                policy: descriptor.policy.requires.clone(),
420                minimum_fidelity: descriptor.fidelity.minimum_input,
421                session: None,
422            },
423            inputs,
424            outputs,
425        })
426    }
427
428    /// Decode an untrusted physical plan, bind it to a trusted realm/PlanId,
429    /// and verify every executable step against this registry before use.
430    pub fn decode_authorized_plan(
431        &self,
432        bytes: &[u8],
433        limits: crate::PlanLimits,
434        authorization: crate::PlanAuthorization,
435    ) -> Result<AuthorizedPlan, crate::PlanDecodeError> {
436        let plan = CompiledPlan::from_authorized_aot_bytes(bytes, limits, authorization)?;
437        let target = plan.realm.target();
438        for node in &plan.nodes {
439            let kernel = self
440                .kernels
441                .get(&node.kernel)
442                .ok_or(crate::PlanDecodeError::UnauthorizedPlan)?;
443            if kernel.target != target
444                || kernel.implements != node.semantic_types
445                || kernel.implementation_id != node.implementation_id
446                || kernel.conversion != node.conversion
447                || kernel.resources != node.resources
448                || kernel.determinism != node.determinism
449                || kernel.lowering != node.lowering
450            {
451                return Err(crate::PlanDecodeError::UnauthorizedPlan);
452            }
453            if let Some(conversion) = &node.conversion {
454                if !kernel.implements.is_empty()
455                    || node.input_ports.as_slice() != ["input"]
456                    || node.output_ports.as_slice() != ["output"]
457                    || node.input_contracts.len() != 1
458                    || node.output_contracts.len() != 1
459                    || !conversion_contracts_match(
460                        &node.input_contracts[0],
461                        &node.output_contracts[0],
462                        conversion,
463                    )
464                    || node.partiality != crate::model::Partiality::Atomic
465                    || !node.failure.domains.is_empty()
466                    || node.state != StateContract::stateless()
467                    || !node.subgraph_path.is_empty()
468                {
469                    return Err(crate::PlanDecodeError::UnauthorizedPlan);
470                }
471                continue;
472            }
473            let mut semantic_descriptors = Vec::with_capacity(node.semantic_types.len());
474            for semantic_type in &node.semantic_types {
475                semantic_descriptors.push(
476                    self.descriptors
477                        .get(&(semantic_type.type_name.clone(), semantic_type.version))
478                        .ok_or(crate::PlanDecodeError::UnauthorizedPlan)?,
479                );
480            }
481            let first = semantic_descriptors
482                .first()
483                .ok_or(crate::PlanDecodeError::UnauthorizedPlan)?;
484            let last = semantic_descriptors
485                .last()
486                .ok_or(crate::PlanDecodeError::UnauthorizedPlan)?;
487            if semantic_descriptors.iter().zip(&node.semantic_configs).any(
488                |(descriptor, config)| match descriptor.config.canonicalize(config) {
489                    Ok(canonical) => canonical != *config,
490                    Err(_) => true,
491                },
492            ) {
493                return Err(crate::PlanDecodeError::UnauthorizedPlan);
494            }
495            if semantic_descriptors
496                .iter()
497                .any(|descriptor| kernel.determinism > descriptor.determinism)
498                || !fused_layouts_compatible(first, last, kernel)
499                || node.input_bindings.len() != first.inputs.len()
500                || node.output_bindings.len() != last.outputs.len()
501                || node.input_ports
502                    != first
503                        .inputs
504                        .iter()
505                        .map(|port| port.name.clone())
506                        .collect::<Vec<_>>()
507                || node.input_contracts.len() != first.inputs.len()
508                || node.output_contracts.len() != last.outputs.len()
509                || node
510                    .input_contracts
511                    .iter()
512                    .any(|contract| !kernel.input_layouts.contains(&contract.layout))
513                || node
514                    .output_contracts
515                    .iter()
516                    .any(|contract| !kernel.output_layouts.contains(&contract.layout))
517                || node
518                    .input_contracts
519                    .iter()
520                    .zip(&first.inputs)
521                    .any(|(compiled, port)| !compiled_port_matches(compiled, port))
522                || node
523                    .output_contracts
524                    .iter()
525                    .zip(&last.outputs)
526                    .any(|(compiled, port)| !compiled_port_matches(compiled, port))
527                || node.output_ports
528                    != last
529                        .outputs
530                        .iter()
531                        .map(|port| port.name.clone())
532                        .collect::<Vec<_>>()
533                || first
534                    .inputs
535                    .iter()
536                    .zip(&node.input_bindings)
537                    .enumerate()
538                    .any(|(input_index, (port, binding))| {
539                        (!port.optional && matches!(binding, crate::model::InputBinding::Absent))
540                            || match binding {
541                                crate::model::InputBinding::Invocation(invocation) => plan
542                                    .invocation_ports
543                                    .get(*invocation as usize)
544                                    .is_none_or(|invocation_port| {
545                                        invocation_port.node != node.semantic_nodes[0]
546                                            || invocation_port.port
547                                                != first.inputs[input_index].name
548                                    }),
549                                _ => false,
550                            }
551                    })
552            {
553                return Err(crate::PlanDecodeError::UnauthorizedPlan);
554            }
555            if semantic_descriptors.len() == 1 {
556                if node.effect != first.effect
557                    || node.retry_limit != first.retry_limit
558                    || node.state != first.state
559                    || node.subgraph_path
560                        != first
561                            .subgraph
562                            .iter()
563                            .map(|lowering| lowering.subgraph)
564                            .collect::<Vec<_>>()
565                    || node.partiality != first.partiality
566                    || node.failure != first.failure
567                {
568                    return Err(crate::PlanDecodeError::UnauthorizedPlan);
569                }
570                if let Some(lowering) = &first.subgraph {
571                    validate_subgraph_path(
572                        &self.subgraphs,
573                        lowering.subgraph,
574                        limits.max_subgraph_depth,
575                        limits.max_contract_entries,
576                    )
577                    .map_err(|_| crate::PlanDecodeError::UnauthorizedPlan)?;
578                    validate_subgraph_lowering(first, lowering, &self.subgraphs, &self.descriptors)
579                        .map_err(|_| crate::PlanDecodeError::UnauthorizedPlan)?;
580                }
581            } else {
582                let expected_failures = failure_domain_union(semantic_descriptors.iter().copied());
583                if node.effect != crate::model::Effect::Pure
584                    || node.retry_limit != 0
585                    || node.state != StateContract::stateless()
586                    || !node.subgraph_path.is_empty()
587                    || node.partiality != crate::model::Partiality::Atomic
588                    || node.failure.domains != expected_failures
589                    || semantic_descriptors.iter().any(|descriptor| {
590                        descriptor.effect != crate::model::Effect::Pure
591                            || descriptor.state.scope != StateScope::Stateless
592                            || descriptor.retry_limit != 0
593                            || descriptor.state.checkpointable()
594                            || descriptor.partiality != crate::model::Partiality::Atomic
595                            || descriptor.subgraph.is_some()
596                    })
597                {
598                    return Err(crate::PlanDecodeError::UnauthorizedPlan);
599                }
600            }
601        }
602        Ok(AuthorizedPlan::new(plan))
603    }
604}
605
606fn normalize_contract(proof: &mut crate::ProofContract, policy: &mut crate::PolicyContract) {
607    proof.requires.sort_unstable();
608    proof.requires.dedup();
609    proof.provides.sort_unstable();
610    proof.provides.dedup();
611    proof.invalidates.sort_unstable();
612    proof.invalidates.dedup();
613    policy.requires.sort_unstable();
614    policy.requires.dedup();
615    policy.adds.sort_unstable();
616    policy.adds.dedup();
617}
618
619fn valid_contract_names(proof: &crate::ProofContract, policy: &crate::PolicyContract) -> bool {
620    proof
621        .requires
622        .iter()
623        .chain(&proof.provides)
624        .chain(&proof.invalidates)
625        .chain(&policy.requires)
626        .chain(&policy.adds)
627        .all(|name| !name.is_empty())
628}
629
630fn subgraph_implements_descriptor(
631    schema: &crate::SubgraphSchema,
632    descriptor: &NodeDescriptor,
633    descriptors: &BTreeMap<(String, u32), NodeDescriptor>,
634) -> bool {
635    fn interface_matches(
636        schema: &crate::SubgraphSchema,
637        interface: &crate::SubgraphInterfacePort,
638        expected: &PortDescriptor,
639        descriptors: &BTreeMap<(String, u32), NodeDescriptor>,
640        input: bool,
641    ) -> bool {
642        let Some(node) = schema
643            .nodes
644            .iter()
645            .find(|node| node.id == interface.inner.node)
646        else {
647            return false;
648        };
649        let Some(descriptor) =
650            descriptors.get(&(node.node_type.type_name.clone(), node.node_type.version))
651        else {
652            return false;
653        };
654        let ports = if input {
655            &descriptor.inputs
656        } else {
657            &descriptor.outputs
658        };
659        let Some(inner) = ports.iter().find(|port| port.name == interface.inner.port) else {
660            return false;
661        };
662        let mut inner = inner.clone();
663        inner.name = interface.name.clone();
664        &inner == expected
665    }
666
667    schema.inputs.len() == descriptor.inputs.len()
668        && schema.outputs.len() == descriptor.outputs.len()
669        && descriptor.inputs.iter().all(|expected| {
670            schema
671                .inputs
672                .iter()
673                .find(|interface| interface.name == expected.name)
674                .is_some_and(|interface| {
675                    interface_matches(schema, interface, expected, descriptors, true)
676                })
677        })
678        && descriptor.outputs.iter().all(|expected| {
679            schema
680                .outputs
681                .iter()
682                .find(|interface| interface.name == expected.name)
683                .is_some_and(|interface| {
684                    interface_matches(schema, interface, expected, descriptors, false)
685                })
686        })
687}
688
689pub(crate) fn valid_port_contract(port: &PortDescriptor) -> bool {
690    let extent = &port.extent;
691    if extent.maximum_shape.len() != usize::from(extent.rank)
692        || extent.max_elements == 0
693        || extent.maximum_shape.contains(&0)
694    {
695        return false;
696    }
697    let shape_product = extent
698        .maximum_shape
699        .iter()
700        .try_fold(1u64, |product, size| product.checked_mul(*size));
701    if shape_product.is_none_or(|product| product > extent.max_elements) {
702        return false;
703    }
704    if port.lease.access == crate::LeaseAccess::ExclusiveWrite
705        && port.lease.lifetime == crate::LeaseLifetime::Session
706    {
707        return false;
708    }
709    valid_contract_names(&port.proof, &port.policy)
710        && !port.domain.root.is_empty()
711        && !port.domain.view.is_empty()
712}
713
714pub(crate) fn valid_state_contract(state: &StateContract) -> bool {
715    match state.scope {
716        StateScope::Stateless => {
717            state.max_bytes == 0
718                && state.checkpoint.mode == crate::CheckpointMode::Disabled
719                && state.checkpoint.max_snapshot_bytes == 0
720                && state.checkpoint.max_interval_invocations == 0
721        }
722        StateScope::Invocation => {
723            state.max_bytes > 0
724                && state.checkpoint.mode == crate::CheckpointMode::Disabled
725                && state.checkpoint.max_snapshot_bytes == 0
726                && state.checkpoint.max_interval_invocations == 0
727        }
728        StateScope::Session => {
729            state.max_bytes > 0
730                && match state.checkpoint.mode {
731                    crate::CheckpointMode::Disabled => {
732                        state.checkpoint.max_snapshot_bytes == 0
733                            && state.checkpoint.max_interval_invocations == 0
734                    }
735                    crate::CheckpointMode::Optional | crate::CheckpointMode::Required => {
736                        state.checkpoint.max_snapshot_bytes > 0
737                            && state.checkpoint.max_snapshot_bytes <= state.max_bytes
738                            && state.checkpoint.max_interval_invocations > 0
739                    }
740                }
741        }
742        StateScope::Durable => {
743            state.max_bytes > 0
744                && state.checkpoint.mode == crate::CheckpointMode::Required
745                && state.checkpoint.max_snapshot_bytes > 0
746                && state.checkpoint.max_snapshot_bytes <= state.max_bytes
747                && state.checkpoint.max_interval_invocations > 0
748        }
749    }
750}
751
752fn port_contract_satisfies(output: &PortDescriptor, input: &PortDescriptor) -> bool {
753    output.semantic_type == input.semantic_type
754        && output.domain == input.domain
755        && (!output.optional || input.optional)
756        && output.max_bytes <= input.max_bytes
757        && output.extent.rank == input.extent.rank
758        && output.extent.max_elements <= input.extent.max_elements
759        && output
760            .extent
761            .maximum_shape
762            .iter()
763            .zip(&input.extent.maximum_shape)
764            .all(|(actual, maximum)| actual <= maximum)
765        && (!output.extent.ragged || input.extent.ragged)
766        && (!output.extent.sparse || input.extent.sparse)
767        && input
768            .proof
769            .requires
770            .iter()
771            .all(|required| output.proof.provides.contains(required))
772        && input
773            .policy
774            .requires
775            .iter()
776            .all(|required| output.policy.adds.contains(required))
777        && output.fidelity.maximum_loss <= input.fidelity.maximum_loss
778        && output.fidelity.minimum_input >= input.fidelity.minimum_input
779        && output.lease == input.lease
780}
781
782pub(crate) fn compiled_port_contract_satisfies(
783    output: &CompiledPortContract,
784    input: &CompiledPortContract,
785) -> bool {
786    output.layout == input.layout
787        && output.semantic_type == input.semantic_type
788        && output.domain == input.domain
789        && (!output.optional || input.optional)
790        && output.max_bytes <= input.max_bytes
791        && output.extent.rank == input.extent.rank
792        && output.extent.max_elements <= input.extent.max_elements
793        && output
794            .extent
795            .maximum_shape
796            .iter()
797            .zip(&input.extent.maximum_shape)
798            .all(|(actual, maximum)| actual <= maximum)
799        && (!output.extent.ragged || input.extent.ragged)
800        && (!output.extent.sparse || input.extent.sparse)
801        && input
802            .proof
803            .requires
804            .iter()
805            .all(|required| output.proof.provides.contains(required))
806        && input
807            .policy
808            .requires
809            .iter()
810            .all(|required| output.policy.adds.contains(required))
811        && output.fidelity.maximum_loss <= input.fidelity.maximum_loss
812        && output.fidelity.minimum_input >= input.fidelity.minimum_input
813        && output.lease == input.lease
814}
815
816fn select_layout(port: &PortDescriptor, kernel_layouts: &[Layout]) -> Layout {
817    port.layouts
818        .iter()
819        .filter(|layout| kernel_layouts.contains(layout))
820        .copied()
821        .min()
822        .expect("kernel compatibility was checked")
823}
824
825fn compiled_port_contract(port: &PortDescriptor, layout: Layout) -> CompiledPortContract {
826    CompiledPortContract {
827        name: port.name.clone(),
828        semantic_type: port.semantic_type.clone(),
829        optional: port.optional,
830        layout,
831        max_bytes: port.max_bytes,
832        domain: port.domain.clone(),
833        proof: port.proof.clone(),
834        policy: port.policy.clone(),
835        fidelity: port.fidelity.clone(),
836        extent: port.extent.clone(),
837        lease: port.lease.clone(),
838    }
839}
840
841fn conversion_port_contract(
842    port: &PortDescriptor,
843    layout: Layout,
844    name: &str,
845    max_bytes: u64,
846) -> CompiledPortContract {
847    let mut contract = compiled_port_contract(port, layout);
848    contract.name = name.to_string();
849    contract.optional = false;
850    contract.max_bytes = max_bytes;
851    contract
852}
853
854fn compiled_port_matches(compiled: &CompiledPortContract, descriptor: &PortDescriptor) -> bool {
855    compiled.name == descriptor.name
856        && compiled.semantic_type == descriptor.semantic_type
857        && compiled.optional == descriptor.optional
858        && descriptor.layouts.contains(&compiled.layout)
859        && compiled.max_bytes == descriptor.max_bytes
860        && compiled.domain == descriptor.domain
861        && compiled.proof == descriptor.proof
862        && compiled.policy == descriptor.policy
863        && compiled.fidelity == descriptor.fidelity
864        && compiled.extent == descriptor.extent
865        && compiled.lease == descriptor.lease
866}
867
868fn conversion_contracts_match(
869    input: &CompiledPortContract,
870    output: &CompiledPortContract,
871    conversion: &crate::LayoutConversion,
872) -> bool {
873    input.name == "input"
874        && output.name == "output"
875        && !input.optional
876        && !output.optional
877        && input.semantic_type == conversion.semantic_type
878        && output.semantic_type == conversion.semantic_type
879        && input.layout == conversion.from
880        && output.layout == conversion.to
881        && input.max_bytes <= conversion.max_input_bytes
882        && output.max_bytes == conversion.max_output_bytes
883        && input.domain == output.domain
884        && input.proof == output.proof
885        && input.policy == output.policy
886        && input.fidelity == output.fidelity
887        && input.extent == output.extent
888        && input.lease == output.lease
889}
890
891fn synchronize_port_layouts(nodes: &mut [CompiledNode], buffers: &[BufferPlan]) {
892    for node in nodes {
893        for (contract, binding) in node.input_contracts.iter_mut().zip(&node.input_bindings) {
894            if let crate::InputBinding::Buffer(buffer) = binding {
895                contract.layout = buffers[buffer.0 as usize].layout;
896            }
897        }
898        for (contract, binding) in node.output_contracts.iter_mut().zip(&node.output_bindings) {
899            if let OutputBinding::Buffer(buffer) = binding {
900                contract.layout = buffers[buffer.0 as usize].layout;
901            }
902        }
903    }
904}
905
906/// Domain-separated semantic identity for a normalized hierarchical schema.
907pub fn subgraph_identity(schema: &crate::SubgraphSchema) -> crate::SubgraphId {
908    let mut hasher = blake3::Hasher::new_derive_key("blut.subgraph.v1");
909    put_u32(&mut hasher, schema.version);
910    let mut nodes = schema.nodes.clone();
911    nodes.sort_by_key(|node| node.id);
912    put_u32(&mut hasher, nodes.len() as u32);
913    for node in &nodes {
914        put_u32(&mut hasher, node.id.0);
915        put_str(&mut hasher, &node.node_type.type_name);
916        put_u32(&mut hasher, node.node_type.version);
917        put_u32(&mut hasher, node.config.len() as u32);
918        for (key, value) in &node.config {
919            put_str(&mut hasher, key);
920            hash_config_value(&mut hasher, value);
921        }
922        match node.child {
923            Some(child) => {
924                hasher.update(&[1]);
925                hasher.update(&child.0);
926            }
927            None => {
928                hasher.update(&[0]);
929            }
930        }
931    }
932    let mut edges = schema.edges.clone();
933    edges.sort_by_key(|edge| (edge.from.clone(), edge.to.clone()));
934    put_u32(&mut hasher, edges.len() as u32);
935    for edge in &edges {
936        put_port_ref(&mut hasher, &edge.from);
937        put_port_ref(&mut hasher, &edge.to);
938    }
939    for ports in [&schema.inputs, &schema.outputs] {
940        let mut ports = ports.clone();
941        ports.sort_unstable();
942        put_u32(&mut hasher, ports.len() as u32);
943        for port in &ports {
944            put_str(&mut hasher, &port.name);
945            put_port_ref(&mut hasher, &port.inner);
946        }
947    }
948    crate::SubgraphId(*hasher.finalize().as_bytes())
949}
950
951fn subgraph_is_acyclic(schema: &crate::SubgraphSchema) -> bool {
952    let mut indegree: BTreeMap<_, usize> = schema.nodes.iter().map(|node| (node.id, 0)).collect();
953    let mut outgoing: BTreeMap<_, Vec<_>> = BTreeMap::new();
954    for edge in &schema.edges {
955        let Some(degree) = indegree.get_mut(&edge.to.node) else {
956            return false;
957        };
958        *degree += 1;
959        outgoing
960            .entry(edge.from.node)
961            .or_default()
962            .push(edge.to.node);
963    }
964    let mut ready: Vec<_> = indegree
965        .iter()
966        .filter_map(|(node, degree)| (*degree == 0).then_some(*node))
967        .collect();
968    let mut visited = 0usize;
969    while let Some(node) = ready.pop() {
970        visited += 1;
971        for target in outgoing.get(&node).into_iter().flatten() {
972            let Some(degree) = indegree.get_mut(target) else {
973                return false;
974            };
975            *degree -= 1;
976            if *degree == 0 {
977                ready.push(*target);
978            }
979        }
980    }
981    visited == schema.nodes.len()
982}
983
984fn validate_subgraph_path(
985    schemas: &BTreeMap<crate::SubgraphId, crate::SubgraphSchema>,
986    root: crate::SubgraphId,
987    max_depth: usize,
988    max_entries: usize,
989) -> Result<(), CompileError> {
990    let mut pending = alloc::vec![(root, 1usize, BTreeSet::new())];
991    let mut searched = 0usize;
992    while let Some((id, depth, mut ancestors)) = pending.pop() {
993        let schema = schemas.get(&id).ok_or(CompileError::UnknownSubgraph(id))?;
994        searched = searched
995            .checked_add(schema.nodes.len())
996            .and_then(|count| count.checked_add(schema.edges.len()))
997            .and_then(|count| count.checked_add(schema.inputs.len()))
998            .and_then(|count| count.checked_add(schema.outputs.len()))
999            .ok_or(CompileError::SubgraphEntryLimitExceeded)?;
1000        if searched > max_entries {
1001            return Err(CompileError::SubgraphEntryLimitExceeded);
1002        }
1003        if depth > max_depth {
1004            return Err(CompileError::SubgraphDepthExceeded);
1005        }
1006        if !ancestors.insert(id) {
1007            return Err(CompileError::InvalidSubgraph(id));
1008        }
1009        for child in schema.nodes.iter().filter_map(|node| node.child) {
1010            pending.push((child, depth + 1, ancestors.clone()));
1011        }
1012    }
1013    Ok(())
1014}
1015
1016fn config_type_satisfies(outer: &crate::ConfigType, inner: &crate::ConfigType) -> bool {
1017    match (outer, inner) {
1018        (crate::ConfigType::Bool, crate::ConfigType::Bool) => true,
1019        (
1020            crate::ConfigType::I64 {
1021                minimum: outer_min,
1022                maximum: outer_max,
1023            },
1024            crate::ConfigType::I64 {
1025                minimum: inner_min,
1026                maximum: inner_max,
1027            },
1028        ) => outer_min >= inner_min && outer_max <= inner_max,
1029        (
1030            crate::ConfigType::U64 {
1031                minimum: outer_min,
1032                maximum: outer_max,
1033            },
1034            crate::ConfigType::U64 {
1035                minimum: inner_min,
1036                maximum: inner_max,
1037            },
1038        ) => outer_min >= inner_min && outer_max <= inner_max,
1039        (
1040            crate::ConfigType::Text {
1041                max_bytes: outer_max,
1042            },
1043            crate::ConfigType::Text {
1044                max_bytes: inner_max,
1045            },
1046        )
1047        | (
1048            crate::ConfigType::Bytes {
1049                max_bytes: outer_max,
1050            },
1051            crate::ConfigType::Bytes {
1052                max_bytes: inner_max,
1053            },
1054        ) => outer_max <= inner_max,
1055        (
1056            crate::ConfigType::Choice {
1057                values: outer_values,
1058            },
1059            crate::ConfigType::Choice {
1060                values: inner_values,
1061            },
1062        ) => outer_values
1063            .iter()
1064            .all(|value| inner_values.contains(value)),
1065        (
1066            crate::ConfigType::Choice {
1067                values: outer_values,
1068            },
1069            crate::ConfigType::Text { max_bytes },
1070        ) => outer_values
1071            .iter()
1072            .all(|value| value.len() <= *max_bytes as usize),
1073        _ => false,
1074    }
1075}
1076
1077fn validate_subgraph_lowering(
1078    descriptor: &NodeDescriptor,
1079    lowering: &crate::SubgraphLowering,
1080    schemas: &BTreeMap<crate::SubgraphId, crate::SubgraphSchema>,
1081    descriptors: &BTreeMap<(String, u32), NodeDescriptor>,
1082) -> Result<(), CompileError> {
1083    let schema = schemas
1084        .get(&lowering.subgraph)
1085        .ok_or(CompileError::UnknownSubgraph(lowering.subgraph))?;
1086    let input_names: BTreeSet<_> = descriptor
1087        .inputs
1088        .iter()
1089        .map(|port| port.name.as_str())
1090        .collect();
1091    let output_names: BTreeSet<_> = descriptor
1092        .outputs
1093        .iter()
1094        .map(|port| port.name.as_str())
1095        .collect();
1096    let mapped_inputs: BTreeSet<_> = lowering
1097        .input_map
1098        .iter()
1099        .map(|map| map.outer.as_str())
1100        .collect();
1101    let mapped_outputs: BTreeSet<_> = lowering
1102        .output_map
1103        .iter()
1104        .map(|map| map.outer.as_str())
1105        .collect();
1106    let inner_inputs: BTreeSet<_> = schema
1107        .inputs
1108        .iter()
1109        .map(|port| port.name.as_str())
1110        .collect();
1111    let inner_outputs: BTreeSet<_> = schema
1112        .outputs
1113        .iter()
1114        .map(|port| port.name.as_str())
1115        .collect();
1116    let outer_fields: BTreeMap<_, _> = descriptor
1117        .config
1118        .fields
1119        .iter()
1120        .map(|field| (field.name.as_str(), field))
1121        .collect();
1122    let mapped_outer_fields: BTreeSet<_> = lowering
1123        .config_map
1124        .iter()
1125        .map(|map| map.outer.as_str())
1126        .collect();
1127    let inner_nodes: BTreeMap<_, _> = schema.nodes.iter().map(|node| (node.id, node)).collect();
1128    let mut config_targets = BTreeSet::new();
1129    let invalid_config_map = lowering.config_map.iter().any(|map| {
1130        let Some(outer) = outer_fields.get(map.outer.as_str()) else {
1131            return true;
1132        };
1133        let Some(node) = inner_nodes.get(&map.node) else {
1134            return true;
1135        };
1136        let Some(inner_descriptor) =
1137            descriptors.get(&(node.node_type.type_name.clone(), node.node_type.version))
1138        else {
1139            return true;
1140        };
1141        let Some(inner) = inner_descriptor
1142            .config
1143            .fields
1144            .iter()
1145            .find(|field| field.name == map.inner)
1146        else {
1147            return true;
1148        };
1149        !config_targets.insert((map.node, map.inner.as_str()))
1150            || !config_type_satisfies(&outer.value_type, &inner.value_type)
1151    });
1152    if input_names != mapped_inputs
1153        || output_names != mapped_outputs
1154        || lowering.input_map.len() != input_names.len()
1155        || lowering.output_map.len() != output_names.len()
1156        || lowering
1157            .input_map
1158            .iter()
1159            .any(|map| !inner_inputs.contains(map.inner.as_str()))
1160        || lowering
1161            .output_map
1162            .iter()
1163            .any(|map| !inner_outputs.contains(map.inner.as_str()))
1164        || lowering
1165            .input_map
1166            .iter()
1167            .map(|map| map.inner.as_str())
1168            .collect::<BTreeSet<_>>()
1169            .len()
1170            != lowering.input_map.len()
1171        || lowering
1172            .output_map
1173            .iter()
1174            .map(|map| map.inner.as_str())
1175            .collect::<BTreeSet<_>>()
1176            .len()
1177            != lowering.output_map.len()
1178        || invalid_config_map
1179        || mapped_outer_fields != outer_fields.keys().copied().collect()
1180    {
1181        return Err(CompileError::InvalidSubgraph(lowering.subgraph));
1182    }
1183    Ok(())
1184}
1185
1186pub struct Compiler<'a> {
1187    registry: &'a KernelRegistry,
1188    realm: ExecutionRealm,
1189    max_peak_bytes: u64,
1190    fuse: bool,
1191    limits: CompileLimits,
1192}
1193
1194#[derive(Clone, Copy, Debug, PartialEq, Eq)]
1195pub struct CompileLimits {
1196    pub max_search_states: usize,
1197    pub max_conversion_states: usize,
1198    pub max_semantic_nodes: usize,
1199    pub max_steps: usize,
1200    pub max_buffers: usize,
1201    pub max_subgraph_depth: usize,
1202    pub max_subgraph_entries: usize,
1203    pub max_feedback_edges: usize,
1204    pub max_persistent_state_bytes: u64,
1205}
1206
1207type LoweredCandidate = (
1208    u64,
1209    usize,
1210    Vec<KernelId>,
1211    Vec<CompiledNode>,
1212    Vec<BufferPlan>,
1213);
1214
1215impl Default for CompileLimits {
1216    fn default() -> Self {
1217        Self {
1218            max_search_states: 65_536,
1219            max_conversion_states: 65_536,
1220            max_semantic_nodes: 65_536,
1221            max_steps: 65_536,
1222            max_buffers: 262_144,
1223            max_subgraph_depth: 16,
1224            max_subgraph_entries: 65_536,
1225            max_feedback_edges: 65_536,
1226            max_persistent_state_bytes: 64 * 1024 * 1024,
1227        }
1228    }
1229}
1230
1231impl<'a> Compiler<'a> {
1232    pub const fn new(registry: &'a KernelRegistry, realm: ExecutionRealm) -> Self {
1233        Self {
1234            registry,
1235            realm,
1236            max_peak_bytes: u64::MAX,
1237            fuse: true,
1238            limits: CompileLimits {
1239                max_search_states: 65_536,
1240                max_conversion_states: 65_536,
1241                max_semantic_nodes: 65_536,
1242                max_steps: 65_536,
1243                max_buffers: 262_144,
1244                max_subgraph_depth: 16,
1245                max_subgraph_entries: 65_536,
1246                max_feedback_edges: 65_536,
1247                max_persistent_state_bytes: 64 * 1024 * 1024,
1248            },
1249        }
1250    }
1251
1252    pub const fn with_memory_limit(mut self, bytes: u64) -> Self {
1253        self.max_peak_bytes = bytes;
1254        self
1255    }
1256
1257    pub const fn with_fusion(mut self, enabled: bool) -> Self {
1258        self.fuse = enabled;
1259        self
1260    }
1261
1262    pub const fn with_limits(mut self, limits: CompileLimits) -> Self {
1263        self.limits = limits;
1264        self
1265    }
1266
1267    pub fn compile(&self, graph: &Graph) -> Result<AuthorizedPlan, CompileError> {
1268        if graph.version != 3 {
1269            return Err(CompileError::UnsupportedGraphVersion(graph.version));
1270        }
1271        if graph.nodes.is_empty() {
1272            return Err(CompileError::EmptyGraph);
1273        }
1274        if graph
1275            .required_proofs
1276            .iter()
1277            .chain(&graph.policy)
1278            .any(|name| name.is_empty())
1279            || graph
1280                .required_capabilities
1281                .iter()
1282                .any(|capability| capability.0.is_empty())
1283        {
1284            return Err(CompileError::InvalidGraphContract);
1285        }
1286        if graph.nodes.len() > self.limits.max_semantic_nodes {
1287            return Err(CompileError::CompileLimitExceeded);
1288        }
1289
1290        let mut normalized_graph = graph.clone();
1291        let mut seen_nodes = BTreeSet::new();
1292        for node in &normalized_graph.nodes {
1293            if !seen_nodes.insert(node.id) {
1294                return Err(CompileError::DuplicateNode(node.id));
1295            }
1296        }
1297
1298        let mut descriptors = BTreeMap::new();
1299        let target = self.realm.target();
1300        let required_caps: BTreeSet<_> = graph.required_capabilities.iter().collect();
1301        for node in &normalized_graph.nodes {
1302            let key = (node.descriptor.clone(), node.descriptor_version);
1303            let descriptor = self
1304                .registry
1305                .descriptors
1306                .get(&key)
1307                .ok_or_else(|| CompileError::UnknownDescriptor(key.0.clone(), key.1))?;
1308            if !descriptor.targets.contains(&target) {
1309                return Err(CompileError::TargetUnsupported(node.id, target));
1310            }
1311            if descriptor.retry_limit > 0
1312                && matches!(descriptor.effect, crate::model::Effect::AtMostOnce)
1313            {
1314                return Err(CompileError::UnsafeRetry(node.id));
1315            }
1316            if matches!(
1317                descriptor.state.scope,
1318                StateScope::Session | StateScope::Durable
1319            ) && normalized_graph.session.is_none()
1320            {
1321                return Err(CompileError::InvalidState(node.id));
1322            }
1323            if let Some(lowering) = &descriptor.subgraph {
1324                validate_subgraph_path(
1325                    &self.registry.subgraphs,
1326                    lowering.subgraph,
1327                    self.limits.max_subgraph_depth,
1328                    self.limits.max_subgraph_entries,
1329                )?;
1330                validate_subgraph_lowering(
1331                    descriptor,
1332                    lowering,
1333                    &self.registry.subgraphs,
1334                    &self.registry.descriptors,
1335                )?;
1336            }
1337            for capability in &descriptor.capabilities {
1338                if !required_caps.contains(capability) {
1339                    return Err(CompileError::CapabilityMissing(capability.0.clone()));
1340                }
1341            }
1342            descriptors.insert(node.id, descriptor);
1343        }
1344        for node in &mut normalized_graph.nodes {
1345            node.config = descriptors[&node.id]
1346                .config
1347                .canonicalize(&node.config)
1348                .map_err(|error| CompileError::InvalidConfig(node.id, error))?;
1349        }
1350        if normalized_graph.session.as_ref().is_some_and(|session| {
1351            session.namespace.is_empty()
1352                || session.max_concurrent_sessions == 0
1353                || session.max_idle_millis == 0
1354        }) {
1355            return Err(CompileError::InvalidSession);
1356        }
1357        if !normalized_graph.feedback.is_empty() && normalized_graph.session.is_none() {
1358            return Err(CompileError::InvalidSession);
1359        }
1360        let mut nodes = BTreeMap::new();
1361        for node in &normalized_graph.nodes {
1362            nodes.insert(node.id, node);
1363        }
1364        let supplied_caps: BTreeSet<_> = descriptors
1365            .values()
1366            .flat_map(|descriptor| descriptor.capabilities.iter())
1367            .collect();
1368        if let Some(extra) = required_caps.difference(&supplied_caps).next() {
1369            return Err(CompileError::CapabilityUnsupported(extra.0.clone()));
1370        }
1371
1372        let invocation_ports = self.verify_edges(&normalized_graph, &descriptors)?;
1373        let order = topological_order(&normalized_graph)?;
1374        let (proofs, policy, fidelity) =
1375            propagate_contracts(&normalized_graph, &order, &descriptors)?;
1376        let (mut compiled_nodes, buffers, peak_bytes) =
1377            self.select_and_lower(&normalized_graph, &order, &nodes, &descriptors, target)?;
1378        let (feedback, feedback_bytes) = lower_feedback(
1379            &normalized_graph,
1380            &mut compiled_nodes,
1381            self.limits.max_feedback_edges,
1382        )?;
1383        let node_state_bytes = compiled_nodes
1384            .iter()
1385            .filter(|node| matches!(node.state.scope, StateScope::Session | StateScope::Durable))
1386            .try_fold(0u64, |total, node| {
1387                total
1388                    .checked_add(node.state.max_bytes)
1389                    .ok_or(CompileError::ResourceOverflow)
1390            })?;
1391        let persistent_state_bytes = node_state_bytes
1392            .checked_add(feedback_bytes)
1393            .ok_or(CompileError::ResourceOverflow)?;
1394        if persistent_state_bytes > self.limits.max_persistent_state_bytes {
1395            return Err(CompileError::ResourceOverflow);
1396        }
1397        let graph_id = GraphId(hash_graph(&normalized_graph, &descriptors));
1398        let mut plan = CompiledPlan {
1399            schema_version: 3,
1400            graph_id,
1401            plan_id: PlanId([0; 32]),
1402            realm: self.realm,
1403            order,
1404            nodes: compiled_nodes,
1405            buffers,
1406            feedback,
1407            invocation_ports,
1408            propagated_proofs: proofs,
1409            propagated_policy: policy,
1410            resulting_fidelity: fidelity,
1411            peak_bytes,
1412            persistent_state_bytes,
1413            session: normalized_graph.session.clone(),
1414        };
1415        plan.plan_id = PlanId(hash_plan(&plan));
1416        Ok(AuthorizedPlan::new(plan))
1417    }
1418
1419    fn verify_edges(
1420        &self,
1421        graph: &Graph,
1422        descriptors: &BTreeMap<NodeId, &NodeDescriptor>,
1423    ) -> Result<Vec<crate::model::PortRef>, CompileError> {
1424        let mut bound = BTreeSet::new();
1425        for edge in &graph.edges {
1426            let from = descriptors
1427                .get(&edge.from.node)
1428                .ok_or_else(|| CompileError::UnknownPort(edge.from.node, edge.from.port.clone()))?;
1429            let to = descriptors
1430                .get(&edge.to.node)
1431                .ok_or_else(|| CompileError::UnknownPort(edge.to.node, edge.to.port.clone()))?;
1432            let output = find_port(&from.outputs, edge.from.node, &edge.from.port)?;
1433            let input = find_port(&to.inputs, edge.to.node, &edge.to.port)?;
1434            if output.max_bytes == 0 {
1435                return Err(CompileError::InvalidPortSize(
1436                    edge.from.node,
1437                    edge.from.port.clone(),
1438                ));
1439            }
1440            if output.semantic_type != input.semantic_type {
1441                return Err(CompileError::TypeMismatch(
1442                    output.semantic_type.clone(),
1443                    input.semantic_type.clone(),
1444                ));
1445            }
1446            if !port_contract_satisfies(output, input) {
1447                return Err(CompileError::PortContractMismatch(
1448                    edge.from.node,
1449                    edge.from.port.clone(),
1450                    edge.to.node,
1451                    edge.to.port.clone(),
1452                ));
1453            }
1454            if !bound.insert(edge.to.clone()) {
1455                return Err(CompileError::DuplicateInput(
1456                    edge.to.node,
1457                    edge.to.port.clone(),
1458                ));
1459            }
1460        }
1461        for feedback in &graph.feedback {
1462            let from = descriptors.get(&feedback.from.node).ok_or_else(|| {
1463                CompileError::UnknownPort(feedback.from.node, feedback.from.port.clone())
1464            })?;
1465            let to = descriptors.get(&feedback.to.node).ok_or_else(|| {
1466                CompileError::UnknownPort(feedback.to.node, feedback.to.port.clone())
1467            })?;
1468            let output = find_port(&from.outputs, feedback.from.node, &feedback.from.port)?;
1469            let input = find_port(&to.inputs, feedback.to.node, &feedback.to.port)?;
1470            if feedback.delay.invocations == 0
1471                || (matches!(&feedback.delay.initial, crate::DelayInitial::Absent)
1472                    && !input.optional)
1473                || !port_contract_satisfies(output, input)
1474                || output.max_bytes > input.max_bytes
1475                || !bound.insert(feedback.to.clone())
1476            {
1477                return Err(CompileError::InvalidFeedback(
1478                    feedback.to.node,
1479                    feedback.to.port.clone(),
1480                ));
1481            }
1482        }
1483        let mut invocation_ports = graph.invocation_inputs.clone();
1484        invocation_ports.sort_unstable();
1485        for pair in invocation_ports.windows(2) {
1486            if pair[0] == pair[1] {
1487                return Err(CompileError::DuplicateInvocation(
1488                    pair[0].node,
1489                    pair[0].port.clone(),
1490                ));
1491            }
1492        }
1493        for invocation in &invocation_ports {
1494            let descriptor = descriptors.get(&invocation.node).ok_or_else(|| {
1495                CompileError::UnknownPort(invocation.node, invocation.port.clone())
1496            })?;
1497            find_port(&descriptor.inputs, invocation.node, &invocation.port)?;
1498            if bound.contains(invocation) {
1499                return Err(CompileError::DuplicateInput(
1500                    invocation.node,
1501                    invocation.port.clone(),
1502                ));
1503            }
1504        }
1505        let invocation_set: BTreeSet<_> = invocation_ports.iter().cloned().collect();
1506        for (node_id, descriptor) in descriptors {
1507            for input in descriptor.inputs.iter().filter(|port| !port.optional) {
1508                let port_ref = crate::model::PortRef {
1509                    node: *node_id,
1510                    port: input.name.clone(),
1511                };
1512                if !bound.contains(&port_ref) && !invocation_set.contains(&port_ref) {
1513                    return Err(CompileError::MissingInput(*node_id, input.name.clone()));
1514                }
1515            }
1516        }
1517        Ok(invocation_ports)
1518    }
1519
1520    fn kernel_candidates(
1521        &self,
1522        order: &[NodeId],
1523        nodes: &BTreeMap<NodeId, &crate::model::NodeInstance>,
1524        target: Target,
1525    ) -> Result<BTreeMap<NodeId, Vec<&KernelDescriptor>>, CompileError> {
1526        let mut candidates = BTreeMap::new();
1527        for node_id in order {
1528            let node = nodes[node_id];
1529            let mut node_candidates: Vec<_> = self
1530                .registry
1531                .kernels
1532                .values()
1533                .filter(|kernel| {
1534                    kernel.implements.as_slice()
1535                        == [NodeTypeRef {
1536                            type_name: node.descriptor.clone(),
1537                            version: node.descriptor_version,
1538                        }]
1539                        && kernel.target == target
1540                        && kernel.determinism
1541                            <= self.registry.descriptors
1542                                [&(node.descriptor.clone(), node.descriptor_version)]
1543                                .determinism
1544                        && descriptor_layouts_compatible(
1545                            self.registry
1546                                .descriptors
1547                                .get(&(node.descriptor.clone(), node.descriptor_version))
1548                                .expect("descriptor was verified"),
1549                            kernel,
1550                        )
1551                })
1552                .collect();
1553            node_candidates.sort_by_key(|kernel| {
1554                (
1555                    kernel.resources.peak_bytes,
1556                    kernel.resources.scratch_bytes,
1557                    kernel.implementation_id,
1558                    kernel.id,
1559                )
1560            });
1561            if node_candidates.is_empty() {
1562                return Err(CompileError::KernelUnavailable(*node_id, target));
1563            }
1564            candidates.insert(*node_id, node_candidates);
1565        }
1566        Ok(candidates)
1567    }
1568
1569    fn select_and_lower(
1570        &self,
1571        graph: &Graph,
1572        order: &[NodeId],
1573        instances: &BTreeMap<NodeId, &crate::model::NodeInstance>,
1574        descriptors: &BTreeMap<NodeId, &NodeDescriptor>,
1575        target: Target,
1576    ) -> Result<(Vec<CompiledNode>, Vec<BufferPlan>, u64), CompileError> {
1577        let candidates = self.kernel_candidates(order, instances, target)?;
1578        let assignment_count = order.iter().try_fold(1usize, |count, node| {
1579            let count = count
1580                .checked_mul(candidates[node].len())
1581                .ok_or(CompileError::SearchLimitExceeded)?;
1582            if count > self.limits.max_search_states {
1583                return Err(CompileError::SearchLimitExceeded);
1584            }
1585            Ok(count)
1586        })?;
1587        let region_limit = self
1588            .limits
1589            .max_search_states
1590            .checked_div(assignment_count)
1591            .filter(|limit| *limit > 0)
1592            .ok_or(CompileError::SearchLimitExceeded)?;
1593
1594        let mut best: Option<LoweredCandidate> = None;
1595        let mut saw_layout_failure = false;
1596        let mut saw_feedback_failure = false;
1597        let mut saw_resource_failure = false;
1598        for ordinal in 0..assignment_count {
1599            let mut remainder = ordinal;
1600            let mut selected = alloc::vec![0usize; order.len()];
1601            for (index, node) in order.iter().enumerate().rev() {
1602                selected[index] = remainder % candidates[node].len();
1603                remainder /= candidates[node].len();
1604            }
1605            let assignment: BTreeMap<_, _> = order
1606                .iter()
1607                .enumerate()
1608                .map(|(index, node)| (*node, candidates[node][selected[index]]))
1609                .collect();
1610            match lower_physical_plan(
1611                PhysicalLowering {
1612                    graph,
1613                    order,
1614                    descriptors,
1615                    kernels: &assignment,
1616                    registry: self.registry,
1617                    target,
1618                },
1619                self.fuse,
1620                PhysicalSearchLimits {
1621                    region_candidates: region_limit,
1622                    conversion_states: self.limits.max_conversion_states,
1623                    max_peak_bytes: self.max_peak_bytes,
1624                    max_steps: self.limits.max_steps,
1625                    max_buffers: self.limits.max_buffers,
1626                },
1627            ) {
1628                Ok((mut nodes, buffers, peak))
1629                    if peak <= self.max_peak_bytes
1630                        && nodes.len() <= self.limits.max_steps
1631                        && buffers.len() <= self.limits.max_buffers =>
1632                {
1633                    if !align_feedback_layouts(graph, descriptors, self.registry, &mut nodes)? {
1634                        saw_feedback_failure = true;
1635                        continue;
1636                    }
1637                    if !feedback_physical_compatible(graph, &nodes)? {
1638                        saw_feedback_failure = true;
1639                        continue;
1640                    }
1641                    let implementation_order = nodes.iter().map(|node| node.kernel).collect();
1642                    let score = (peak, nodes.len(), implementation_order);
1643                    if best
1644                        .as_ref()
1645                        .is_none_or(|current| score < (current.0, current.1, current.2.clone()))
1646                    {
1647                        best = Some((score.0, score.1, score.2, nodes, buffers));
1648                    }
1649                }
1650                Ok((nodes, buffers, _))
1651                    if nodes.len() > self.limits.max_steps
1652                        || buffers.len() > self.limits.max_buffers =>
1653                {
1654                    return Err(CompileError::CompileLimitExceeded);
1655                }
1656                Ok(_) | Err(CompileError::ResourceOverflow) => saw_resource_failure = true,
1657                Err(CompileError::LayoutUnavailable(..)) => saw_layout_failure = true,
1658                Err(error) => return Err(error),
1659            }
1660        }
1661        match best {
1662            Some((peak, _, _, nodes, buffers)) => Ok((nodes, buffers, peak)),
1663            None if saw_feedback_failure => {
1664                let feedback = graph.feedback.first().expect("failure requires feedback");
1665                Err(CompileError::InvalidFeedback(
1666                    feedback.to.node,
1667                    feedback.to.port.clone(),
1668                ))
1669            }
1670            None if saw_resource_failure => Err(CompileError::ResourceOverflow),
1671            None if saw_layout_failure => Err(CompileError::LayoutUnavailable(
1672                *order.first().expect("non-empty graph"),
1673                "physical-lowering".to_string(),
1674            )),
1675            None => Err(CompileError::KernelUnavailable(
1676                *order.first().expect("non-empty graph"),
1677                target,
1678            )),
1679        }
1680    }
1681}
1682
1683fn descriptor_layouts_compatible(descriptor: &NodeDescriptor, kernel: &KernelDescriptor) -> bool {
1684    descriptor.inputs.iter().all(|port| {
1685        port.layouts
1686            .iter()
1687            .any(|layout| kernel.input_layouts.contains(layout))
1688    }) && descriptor.outputs.iter().all(|port| {
1689        port.layouts
1690            .iter()
1691            .any(|layout| kernel.output_layouts.contains(layout))
1692    })
1693}
1694
1695fn align_feedback_layouts(
1696    graph: &Graph,
1697    descriptors: &BTreeMap<NodeId, &NodeDescriptor>,
1698    registry: &KernelRegistry,
1699    nodes: &mut [CompiledNode],
1700) -> Result<bool, CompileError> {
1701    let mut groups: BTreeMap<_, Vec<_>> = BTreeMap::new();
1702    for edge in &graph.feedback {
1703        groups.entry(edge.from.clone()).or_default().push(edge);
1704    }
1705    for (source, edges) in groups {
1706        let from_index = nodes
1707            .iter()
1708            .position(|node| node.semantic_nodes.contains(&source.node))
1709            .ok_or(CompileError::UnknownNode(source.node))?;
1710        let from_port = nodes[from_index]
1711            .output_ports
1712            .iter()
1713            .position(|port| port == &source.port)
1714            .ok_or_else(|| CompileError::UnknownPort(source.node, source.port.clone()))?;
1715        let from_descriptor = descriptors[&source.node]
1716            .outputs
1717            .iter()
1718            .find(|port| port.name == source.port)
1719            .ok_or_else(|| CompileError::UnknownPort(source.node, source.port.clone()))?;
1720        let from_kernel = registry
1721            .kernels
1722            .get(&nodes[from_index].kernel)
1723            .ok_or_else(|| CompileError::InvalidFeedback(source.node, source.port.clone()))?;
1724        let producer_fixed = matches!(
1725            nodes[from_index].output_bindings[from_port],
1726            OutputBinding::Buffer(_)
1727        );
1728        let mut layouts: Vec<_> = from_descriptor
1729            .layouts
1730            .iter()
1731            .filter(|layout| from_kernel.output_layouts.contains(layout))
1732            .copied()
1733            .collect();
1734        let mut consumers = Vec::with_capacity(edges.len());
1735        for edge in edges {
1736            let to_index = nodes
1737                .iter()
1738                .position(|node| node.semantic_nodes.contains(&edge.to.node))
1739                .ok_or(CompileError::UnknownNode(edge.to.node))?;
1740            let to_port = nodes[to_index]
1741                .input_ports
1742                .iter()
1743                .position(|port| port == &edge.to.port)
1744                .ok_or_else(|| CompileError::UnknownPort(edge.to.node, edge.to.port.clone()))?;
1745            let to_descriptor = descriptors[&edge.to.node]
1746                .inputs
1747                .iter()
1748                .find(|port| port.name == edge.to.port)
1749                .ok_or_else(|| CompileError::UnknownPort(edge.to.node, edge.to.port.clone()))?;
1750            let to_kernel = registry
1751                .kernels
1752                .get(&nodes[to_index].kernel)
1753                .ok_or_else(|| CompileError::InvalidFeedback(edge.to.node, edge.to.port.clone()))?;
1754            layouts.retain(|layout| {
1755                to_descriptor.layouts.contains(layout) && to_kernel.input_layouts.contains(layout)
1756            });
1757            consumers.push((to_index, to_port));
1758        }
1759        layouts.sort_unstable();
1760        layouts.dedup();
1761        let selected = if producer_fixed {
1762            let current = nodes[from_index].output_contracts[from_port].layout;
1763            layouts.contains(&current).then_some(current)
1764        } else {
1765            layouts.first().copied()
1766        };
1767        let Some(selected) = selected else {
1768            return Ok(false);
1769        };
1770        nodes[from_index].output_contracts[from_port].layout = selected;
1771        for (to_index, to_port) in consumers {
1772            nodes[to_index].input_contracts[to_port].layout = selected;
1773        }
1774    }
1775    Ok(true)
1776}
1777
1778fn feedback_physical_compatible(
1779    graph: &Graph,
1780    nodes: &[CompiledNode],
1781) -> Result<bool, CompileError> {
1782    for edge in &graph.feedback {
1783        let from = nodes
1784            .iter()
1785            .find(|node| node.semantic_nodes.contains(&edge.from.node))
1786            .ok_or(CompileError::UnknownNode(edge.from.node))?;
1787        let from_port = from
1788            .output_ports
1789            .iter()
1790            .position(|port| port == &edge.from.port)
1791            .ok_or_else(|| CompileError::UnknownPort(edge.from.node, edge.from.port.clone()))?;
1792        let to = nodes
1793            .iter()
1794            .find(|node| node.semantic_nodes.contains(&edge.to.node))
1795            .ok_or(CompileError::UnknownNode(edge.to.node))?;
1796        let to_port = to
1797            .input_ports
1798            .iter()
1799            .position(|port| port == &edge.to.port)
1800            .ok_or_else(|| CompileError::UnknownPort(edge.to.node, edge.to.port.clone()))?;
1801        if !compiled_port_contract_satisfies(
1802            &from.output_contracts[from_port],
1803            &to.input_contracts[to_port],
1804        ) {
1805            return Ok(false);
1806        }
1807    }
1808    Ok(true)
1809}
1810
1811fn find_port<'a>(
1812    ports: &'a [PortDescriptor],
1813    node: NodeId,
1814    name: &str,
1815) -> Result<&'a PortDescriptor, CompileError> {
1816    ports
1817        .iter()
1818        .find(|port| port.name == name)
1819        .ok_or_else(|| CompileError::UnknownPort(node, name.to_string()))
1820}
1821
1822fn topological_order(graph: &Graph) -> Result<Vec<NodeId>, CompileError> {
1823    let mut indegree: BTreeMap<NodeId, usize> =
1824        graph.nodes.iter().map(|node| (node.id, 0)).collect();
1825    let mut outgoing: BTreeMap<NodeId, Vec<NodeId>> = BTreeMap::new();
1826    for Edge { from, to } in &graph.edges {
1827        if !indegree.contains_key(&from.node) || !indegree.contains_key(&to.node) {
1828            return Err(CompileError::UnknownNode(
1829                if !indegree.contains_key(&from.node) {
1830                    from.node
1831                } else {
1832                    to.node
1833                },
1834            ));
1835        }
1836        *indegree.get_mut(&to.node).expect("checked") += 1;
1837        outgoing.entry(from.node).or_default().push(to.node);
1838    }
1839    for values in outgoing.values_mut() {
1840        values.sort();
1841    }
1842    let mut ready: BTreeSet<NodeId> = indegree
1843        .iter()
1844        .filter_map(|(id, count)| (*count == 0).then_some(*id))
1845        .collect();
1846    let mut result = Vec::with_capacity(indegree.len());
1847    while let Some(id) = ready.pop_first() {
1848        result.push(id);
1849        if let Some(next) = outgoing.get(&id) {
1850            for target in next {
1851                let degree = indegree.get_mut(target).expect("edge target checked");
1852                *degree -= 1;
1853                if *degree == 0 {
1854                    ready.insert(*target);
1855                }
1856            }
1857        }
1858    }
1859    if result.len() != indegree.len() {
1860        return Err(CompileError::Cycle);
1861    }
1862    Ok(result)
1863}
1864
1865fn propagate_contracts(
1866    graph: &Graph,
1867    order: &[NodeId],
1868    descriptors: &BTreeMap<NodeId, &NodeDescriptor>,
1869) -> Result<(Vec<String>, Vec<String>, u16), CompileError> {
1870    #[derive(Clone)]
1871    struct State {
1872        proofs: BTreeSet<String>,
1873        policy: BTreeSet<String>,
1874        fidelity: u16,
1875    }
1876
1877    let initial = State {
1878        proofs: graph.required_proofs.iter().cloned().collect(),
1879        policy: graph.policy.iter().cloned().collect(),
1880        fidelity: u16::MAX,
1881    };
1882    let mut predecessors: BTreeMap<NodeId, Vec<NodeId>> = BTreeMap::new();
1883    let mut has_successor = BTreeSet::new();
1884    for edge in &graph.edges {
1885        predecessors
1886            .entry(edge.to.node)
1887            .or_default()
1888            .push(edge.from.node);
1889        has_successor.insert(edge.from.node);
1890    }
1891    for values in predecessors.values_mut() {
1892        values.sort_unstable();
1893        values.dedup();
1894    }
1895    let mut states: BTreeMap<NodeId, State> = BTreeMap::new();
1896    for id in order {
1897        let descriptor = descriptors[id];
1898        let mut state = match predecessors.get(id).map(Vec::as_slice).unwrap_or_default() {
1899            [] => initial.clone(),
1900            [first, rest @ ..] => {
1901                let mut joined = states[first].clone();
1902                for predecessor in rest {
1903                    joined.proofs = joined
1904                        .proofs
1905                        .intersection(&states[predecessor].proofs)
1906                        .cloned()
1907                        .collect();
1908                    joined
1909                        .policy
1910                        .extend(states[predecessor].policy.iter().cloned());
1911                    joined.fidelity = joined.fidelity.min(states[predecessor].fidelity);
1912                }
1913                joined
1914            }
1915        };
1916        for required in &descriptor.proof.requires {
1917            if !state.proofs.contains(required) {
1918                return Err(CompileError::ProofMissing(*id, required.clone()));
1919            }
1920        }
1921        for required in &descriptor.policy.requires {
1922            if !state.policy.contains(required) {
1923                return Err(CompileError::PolicyMissing(*id, required.clone()));
1924            }
1925        }
1926        if state.fidelity < descriptor.fidelity.minimum_input {
1927            return Err(CompileError::FidelityInsufficient(*id));
1928        }
1929        state.fidelity = state
1930            .fidelity
1931            .saturating_sub(descriptor.fidelity.maximum_loss);
1932        for invalidated in &descriptor.proof.invalidates {
1933            state.proofs.remove(invalidated);
1934        }
1935        state
1936            .proofs
1937            .extend(descriptor.proof.provides.iter().cloned());
1938        state.policy.extend(descriptor.policy.adds.iter().cloned());
1939        states.insert(*id, state);
1940    }
1941    let mut terminals = order.iter().filter(|id| !has_successor.contains(id));
1942    let first = *terminals.next().expect("non-empty graph has a terminal");
1943    let mut result = states[&first].clone();
1944    for terminal in terminals {
1945        result.proofs = result
1946            .proofs
1947            .intersection(&states[terminal].proofs)
1948            .cloned()
1949            .collect();
1950        result
1951            .policy
1952            .extend(states[terminal].policy.iter().cloned());
1953        result.fidelity = result.fidelity.min(states[terminal].fidelity);
1954    }
1955    if result.fidelity < graph.minimum_fidelity {
1956        return Err(CompileError::FidelityInsufficient(
1957            *order.last().expect("non-empty graph"),
1958        ));
1959    }
1960    Ok((
1961        result.proofs.into_iter().collect(),
1962        result.policy.into_iter().collect(),
1963        result.fidelity,
1964    ))
1965}
1966
1967struct PhysicalLowering<'a> {
1968    graph: &'a Graph,
1969    order: &'a [NodeId],
1970    descriptors: &'a BTreeMap<NodeId, &'a NodeDescriptor>,
1971    kernels: &'a BTreeMap<NodeId, &'a KernelDescriptor>,
1972    registry: &'a KernelRegistry,
1973    target: Target,
1974}
1975
1976#[derive(Clone, Copy)]
1977struct PhysicalSearchLimits {
1978    region_candidates: usize,
1979    conversion_states: usize,
1980    max_peak_bytes: u64,
1981    max_steps: usize,
1982    max_buffers: usize,
1983}
1984
1985fn build_semantic_region_candidates(
1986    context: &PhysicalLowering<'_>,
1987    fuse: bool,
1988    max_candidates: usize,
1989) -> Result<Vec<Vec<CompiledNode>>, CompileError> {
1990    let graph = context.graph;
1991    let order = context.order;
1992    let descriptors = context.descriptors;
1993    let kernels = context.kernels;
1994    let registry = context.registry;
1995    let target = context.target;
1996    let instances: BTreeMap<_, _> = graph.nodes.iter().map(|node| (node.id, node)).collect();
1997    let mut pending = alloc::vec![(0usize, Vec::<(usize, &KernelDescriptor)>::new())];
1998    let mut complete = Vec::new();
1999    while let Some((position, fused_choices)) = pending.pop() {
2000        if position == order.len() {
2001            let mut regions = Vec::with_capacity(order.len());
2002            let mut cursor = 0usize;
2003            let mut choice = 0usize;
2004            while cursor < order.len() {
2005                if fused_choices
2006                    .get(choice)
2007                    .is_some_and(|(start, _)| *start == cursor)
2008                {
2009                    let (_, kernel) = fused_choices[choice];
2010                    let end = cursor + kernel.implements.len();
2011                    regions.push(semantic_region(
2012                        &order[cursor..end],
2013                        kernel,
2014                        descriptors,
2015                        &instances,
2016                    ));
2017                    cursor = end;
2018                    choice += 1;
2019                } else {
2020                    let id = order[cursor];
2021                    regions.push(semantic_region(
2022                        &order[cursor..cursor + 1],
2023                        kernels[&id],
2024                        descriptors,
2025                        &instances,
2026                    ));
2027                    cursor += 1;
2028                }
2029            }
2030            complete.push(regions);
2031            if complete.len() > max_candidates {
2032                return Err(CompileError::SearchLimitExceeded);
2033            }
2034            continue;
2035        }
2036
2037        let mut alternatives = Vec::new();
2038        if fuse {
2039            for kernel in registry.kernels.values().filter(|kernel| {
2040                kernel.target == target
2041                    && kernel.implements.len() >= 2
2042                    && position + kernel.implements.len() <= order.len()
2043            }) {
2044                let ids = &order[position..position + kernel.implements.len()];
2045                if kernel
2046                    .implements
2047                    .iter()
2048                    .zip(ids)
2049                    .all(|(implemented, node)| {
2050                        implemented.type_name == instances[node].descriptor
2051                            && implemented.version == instances[node].descriptor_version
2052                    })
2053                    && ids
2054                        .iter()
2055                        .all(|node| kernel.determinism <= descriptors[node].determinism)
2056                    && fused_layouts_compatible(
2057                        descriptors[&ids[0]],
2058                        descriptors[ids.last().expect("non-empty")],
2059                        kernel,
2060                    )
2061                    && linear_fusion_is_safe(graph, ids, descriptors)
2062                {
2063                    alternatives.push(kernel);
2064                }
2065            }
2066        }
2067        alternatives.sort_by_key(|kernel| {
2068            (
2069                core::cmp::Reverse(kernel.implements.len()),
2070                kernel.resources.peak_bytes,
2071                kernel.resources.scratch_bytes,
2072                kernel.implementation_id,
2073                kernel.id,
2074            )
2075        });
2076        for kernel in alternatives.into_iter().rev() {
2077            let mut next = fused_choices.clone();
2078            next.push((position, kernel));
2079            pending.push((position + kernel.implements.len(), next));
2080            if pending.len().saturating_add(complete.len()) > max_candidates {
2081                return Err(CompileError::SearchLimitExceeded);
2082            }
2083        }
2084        pending.push((position + 1, fused_choices));
2085        if pending.len().saturating_add(complete.len()) > max_candidates {
2086            return Err(CompileError::SearchLimitExceeded);
2087        }
2088    }
2089    Ok(complete)
2090}
2091
2092fn semantic_region(
2093    ids: &[NodeId],
2094    kernel: &KernelDescriptor,
2095    descriptors: &BTreeMap<NodeId, &NodeDescriptor>,
2096    instances: &BTreeMap<NodeId, &crate::model::NodeInstance>,
2097) -> CompiledNode {
2098    let first = descriptors[&ids[0]];
2099    let last = descriptors[ids.last().expect("semantic region is non-empty")];
2100    let fused = ids.len() > 1;
2101    CompiledNode {
2102        id: StepId(0),
2103        semantic_nodes: ids.to_vec(),
2104        semantic_types: ids
2105            .iter()
2106            .map(|id| NodeTypeRef {
2107                type_name: instances[id].descriptor.clone(),
2108                version: instances[id].descriptor_version,
2109            })
2110            .collect(),
2111        semantic_configs: ids.iter().map(|id| instances[id].config.clone()).collect(),
2112        kernel: kernel.id,
2113        implementation_id: kernel.implementation_id,
2114        resources: kernel.resources.clone(),
2115        determinism: kernel.determinism,
2116        lowering: kernel.lowering.clone(),
2117        conversion: None,
2118        input_ports: first.inputs.iter().map(|port| port.name.clone()).collect(),
2119        output_ports: last.outputs.iter().map(|port| port.name.clone()).collect(),
2120        input_contracts: first
2121            .inputs
2122            .iter()
2123            .map(|port| compiled_port_contract(port, select_layout(port, &kernel.input_layouts)))
2124            .collect(),
2125        output_contracts: last
2126            .outputs
2127            .iter()
2128            .map(|port| compiled_port_contract(port, select_layout(port, &kernel.output_layouts)))
2129            .collect(),
2130        input_bindings: alloc::vec![crate::model::InputBinding::Absent; first.inputs.len()],
2131        output_bindings: alloc::vec![OutputBinding::Terminal; last.outputs.len()],
2132        partiality: if fused {
2133            crate::model::Partiality::Atomic
2134        } else {
2135            first.partiality
2136        },
2137        failure: if fused {
2138            crate::model::FailureContract {
2139                domains: failure_domain_union(ids.iter().map(|id| descriptors[id])),
2140            }
2141        } else {
2142            first.failure.clone()
2143        },
2144        effect: if fused {
2145            crate::model::Effect::Pure
2146        } else {
2147            first.effect
2148        },
2149        retry_limit: if fused { 0 } else { first.retry_limit },
2150        state: if fused {
2151            StateContract::stateless()
2152        } else {
2153            first.state.clone()
2154        },
2155        subgraph_path: ids
2156            .iter()
2157            .filter_map(|id| descriptors[id].subgraph.as_ref().map(|item| item.subgraph))
2158            .collect(),
2159    }
2160}
2161
2162fn linear_fusion_is_safe(
2163    graph: &Graph,
2164    ids: &[NodeId],
2165    descriptors: &BTreeMap<NodeId, &NodeDescriptor>,
2166) -> bool {
2167    if graph
2168        .feedback
2169        .iter()
2170        .any(|edge| ids.contains(&edge.from.node) || ids.contains(&edge.to.node))
2171    {
2172        return false;
2173    }
2174    if ids.iter().any(|id| {
2175        let descriptor = descriptors[id];
2176        descriptor.effect != crate::model::Effect::Pure
2177            || descriptor.partiality != crate::model::Partiality::Atomic
2178            || descriptor.state.scope != StateScope::Stateless
2179            || descriptor.retry_limit != 0
2180            || descriptor.state.checkpointable()
2181            || descriptor.subgraph.is_some()
2182    }) {
2183        return false;
2184    }
2185    ids.windows(2).all(|pair| {
2186        let from = pair[0];
2187        let to = pair[1];
2188        let outgoing: Vec<_> = graph
2189            .edges
2190            .iter()
2191            .filter(|edge| edge.from.node == from)
2192            .collect();
2193        let incoming: Vec<_> = graph
2194            .edges
2195            .iter()
2196            .filter(|edge| edge.to.node == to)
2197            .collect();
2198        descriptors[&from].outputs.len() == 1
2199            && descriptors[&to].inputs.len() == 1
2200            && outgoing.len() == 1
2201            && incoming.len() == 1
2202            && outgoing[0] == incoming[0]
2203            && outgoing[0].from.port == descriptors[&from].outputs[0].name
2204            && incoming[0].to.port == descriptors[&to].inputs[0].name
2205    })
2206}
2207
2208fn failure_domain_union<'a>(
2209    descriptors: impl IntoIterator<Item = &'a NodeDescriptor>,
2210) -> Vec<String> {
2211    let mut domains = descriptors
2212        .into_iter()
2213        .flat_map(|descriptor| descriptor.failure.domains.iter().cloned())
2214        .collect::<Vec<_>>();
2215    domains.sort_unstable();
2216    domains.dedup();
2217    domains
2218}
2219
2220fn fused_layouts_compatible(
2221    first: &NodeDescriptor,
2222    last: &NodeDescriptor,
2223    kernel: &KernelDescriptor,
2224) -> bool {
2225    first.inputs.iter().all(|port| {
2226        port.layouts
2227            .iter()
2228            .any(|layout| kernel.input_layouts.contains(layout))
2229    }) && last.outputs.iter().all(|port| {
2230        port.layouts
2231            .iter()
2232            .any(|layout| kernel.output_layouts.contains(layout))
2233    })
2234}
2235
2236#[derive(Clone)]
2237struct Route<'a> {
2238    consumer_region: usize,
2239    input_index: usize,
2240    path: Vec<&'a KernelDescriptor>,
2241}
2242
2243struct PortGroup<'a> {
2244    producer_region: usize,
2245    output_index: usize,
2246    layout: Layout,
2247    capacity_bytes: u64,
2248    contract: PortDescriptor,
2249    routes: Vec<Route<'a>>,
2250}
2251
2252fn lower_physical_plan(
2253    context: PhysicalLowering<'_>,
2254    fuse: bool,
2255    limits: PhysicalSearchLimits,
2256) -> Result<(Vec<CompiledNode>, Vec<BufferPlan>, u64), CompileError> {
2257    let region_candidates =
2258        build_semantic_region_candidates(&context, fuse, limits.region_candidates)?;
2259    let mut best: Option<LoweredCandidate> = None;
2260    let mut first_error = None;
2261    let mut saw_resource_limit = false;
2262    let mut saw_compile_limit = false;
2263    for semantic_regions in region_candidates {
2264        match lower_semantic_regions(&context, semantic_regions, limits.conversion_states) {
2265            Ok((nodes, buffers, _))
2266                if nodes.len() > limits.max_steps || buffers.len() > limits.max_buffers =>
2267            {
2268                saw_compile_limit = true;
2269            }
2270            Ok((_, _, peak)) if peak > limits.max_peak_bytes => {
2271                saw_resource_limit = true;
2272            }
2273            Ok((nodes, buffers, peak)) => {
2274                let kernels = nodes.iter().map(|node| node.kernel).collect::<Vec<_>>();
2275                let score = (peak, nodes.len(), kernels);
2276                if best
2277                    .as_ref()
2278                    .is_none_or(|current| score < (current.0, current.1, current.2.clone()))
2279                {
2280                    best = Some((score.0, score.1, score.2, nodes, buffers));
2281                }
2282            }
2283            Err(error) => {
2284                first_error.get_or_insert(error);
2285            }
2286        }
2287    }
2288    best.map(|(peak, _, _, nodes, buffers)| (nodes, buffers, peak))
2289        .ok_or_else(|| {
2290            if saw_resource_limit {
2291                CompileError::ResourceOverflow
2292            } else if saw_compile_limit {
2293                CompileError::CompileLimitExceeded
2294            } else {
2295                first_error.unwrap_or(CompileError::ResourceOverflow)
2296            }
2297        })
2298}
2299
2300fn lower_semantic_regions(
2301    context: &PhysicalLowering<'_>,
2302    semantic_regions: Vec<CompiledNode>,
2303    max_conversion_states: usize,
2304) -> Result<(Vec<CompiledNode>, Vec<BufferPlan>, u64), CompileError> {
2305    let graph = context.graph;
2306    let descriptors = context.descriptors;
2307    let registry = context.registry;
2308    let target = context.target;
2309    let owner: BTreeMap<_, _> = semantic_regions
2310        .iter()
2311        .enumerate()
2312        .flat_map(|(index, node)| node.semantic_nodes.iter().map(move |id| (*id, index)))
2313        .collect();
2314    let mut grouped: BTreeMap<crate::model::PortRef, Vec<&Edge>> = BTreeMap::new();
2315    for edge in &graph.edges {
2316        if owner[&edge.from.node] != owner[&edge.to.node] {
2317            grouped.entry(edge.from.clone()).or_default().push(edge);
2318        }
2319    }
2320    let mut groups = Vec::new();
2321    let mut conversion_states = 0usize;
2322    for (source, mut edges) in grouped {
2323        edges.sort_by_key(|edge| edge.to.clone());
2324        let producer_descriptor = descriptors[&source.node];
2325        let output_index = producer_descriptor
2326            .outputs
2327            .iter()
2328            .position(|port| port.name == source.port)
2329            .ok_or_else(|| CompileError::UnknownPort(source.node, source.port.clone()))?;
2330        let output = &producer_descriptor.outputs[output_index];
2331        if output.max_bytes == 0 {
2332            return Err(CompileError::InvalidPortSize(source.node, source.port));
2333        }
2334        let producer_region = owner[&source.node];
2335        let producer_kernel = registry
2336            .kernels
2337            .get(&semantic_regions[producer_region].kernel)
2338            .expect("compiled kernel is registered");
2339        let mut source_layouts: Vec<_> = output
2340            .layouts
2341            .iter()
2342            .filter(|layout| producer_kernel.output_layouts.contains(layout))
2343            .copied()
2344            .collect();
2345        source_layouts.sort_unstable();
2346        source_layouts.dedup();
2347
2348        let mut best: Option<(usize, u64, Layout, Vec<Route<'_>>)> = None;
2349        for source_layout in source_layouts {
2350            let mut routes = Vec::new();
2351            let mut conversion_count = 0usize;
2352            let mut conversion_workspace = 0u64;
2353            let mut valid = true;
2354            for edge in &edges {
2355                let consumer_descriptor = descriptors[&edge.to.node];
2356                let input_index = consumer_descriptor
2357                    .inputs
2358                    .iter()
2359                    .position(|port| port.name == edge.to.port)
2360                    .ok_or_else(|| CompileError::UnknownPort(edge.to.node, edge.to.port.clone()))?;
2361                let input = &consumer_descriptor.inputs[input_index];
2362                let consumer_region = owner[&edge.to.node];
2363                let consumer_kernel = registry
2364                    .kernels
2365                    .get(&semantic_regions[consumer_region].kernel)
2366                    .expect("compiled kernel is registered");
2367                let mut targets: Vec<_> = input
2368                    .layouts
2369                    .iter()
2370                    .filter(|layout| consumer_kernel.input_layouts.contains(layout))
2371                    .copied()
2372                    .collect();
2373                targets.sort_unstable();
2374                targets.dedup();
2375                let path =
2376                    if targets.contains(&source_layout) && output.max_bytes <= input.max_bytes {
2377                        Some(Vec::new())
2378                    } else {
2379                        conversion_path(
2380                            registry,
2381                            ConversionRequest {
2382                                target,
2383                                semantic_type: &output.semantic_type,
2384                                max_bytes: output.max_bytes,
2385                                from: source_layout,
2386                                targets: &targets,
2387                                target_max_bytes: input.max_bytes,
2388                            },
2389                            &mut ConversionBudget {
2390                                searched: &mut conversion_states,
2391                                max_states: max_conversion_states,
2392                            },
2393                        )?
2394                    };
2395                let Some(path) = path else {
2396                    valid = false;
2397                    break;
2398                };
2399                conversion_count = conversion_count
2400                    .checked_add(path.len())
2401                    .ok_or(CompileError::SearchLimitExceeded)?;
2402                for kernel in &path {
2403                    conversion_workspace = conversion_workspace.max(
2404                        kernel
2405                            .resources
2406                            .peak_bytes
2407                            .checked_add(kernel.resources.scratch_bytes)
2408                            .ok_or(CompileError::ResourceOverflow)?,
2409                    );
2410                }
2411                routes.push(Route {
2412                    consumer_region,
2413                    input_index,
2414                    path,
2415                });
2416            }
2417            if valid {
2418                let score = (conversion_count, conversion_workspace, source_layout);
2419                if best
2420                    .as_ref()
2421                    .is_none_or(|current| score < (current.0, current.1, current.2))
2422                {
2423                    best = Some((score.0, score.1, score.2, routes));
2424                }
2425            }
2426        }
2427        let Some((_, _, layout, routes)) = best else {
2428            return Err(CompileError::LayoutUnavailable(source.node, source.port));
2429        };
2430        groups.push(PortGroup {
2431            producer_region,
2432            output_index,
2433            layout,
2434            capacity_bytes: output.max_bytes,
2435            contract: output.clone(),
2436            routes,
2437        });
2438    }
2439
2440    let mut nodes = Vec::new();
2441    let mut semantic_steps = alloc::vec![StepId(0); semantic_regions.len()];
2442    let mut conversion_steps: BTreeMap<(usize, usize), Vec<StepId>> = BTreeMap::new();
2443    for region_index in 0..semantic_regions.len() {
2444        for (group_index, group) in groups.iter().enumerate() {
2445            for (route_index, route) in group.routes.iter().enumerate() {
2446                if route.consumer_region != region_index {
2447                    continue;
2448                }
2449                let mut steps = Vec::new();
2450                for (path_index, kernel) in route.path.iter().enumerate() {
2451                    let conversion = kernel
2452                        .conversion
2453                        .clone()
2454                        .expect("conversion path contains conversion kernels");
2455                    let id = StepId(nodes.len() as u32);
2456                    steps.push(id);
2457                    nodes.push(CompiledNode {
2458                        id,
2459                        semantic_nodes: Vec::new(),
2460                        semantic_types: Vec::new(),
2461                        semantic_configs: Vec::new(),
2462                        kernel: kernel.id,
2463                        implementation_id: kernel.implementation_id,
2464                        resources: kernel.resources.clone(),
2465                        determinism: kernel.determinism,
2466                        lowering: kernel.lowering.clone(),
2467                        conversion: Some(conversion.clone()),
2468                        input_ports: alloc::vec!["input".to_string()],
2469                        output_ports: alloc::vec!["output".to_string()],
2470                        input_contracts: alloc::vec![conversion_port_contract(
2471                            &group.contract,
2472                            conversion.from,
2473                            "input",
2474                            if path_index == 0 {
2475                                group.capacity_bytes
2476                            } else {
2477                                route.path[path_index - 1]
2478                                    .conversion
2479                                    .as_ref()
2480                                    .expect("conversion path")
2481                                    .max_output_bytes
2482                            },
2483                        )],
2484                        output_contracts: alloc::vec![conversion_port_contract(
2485                            &group.contract,
2486                            conversion.to,
2487                            "output",
2488                            conversion.max_output_bytes,
2489                        )],
2490                        input_bindings: alloc::vec![crate::model::InputBinding::Absent],
2491                        output_bindings: alloc::vec![OutputBinding::Terminal],
2492                        partiality: crate::model::Partiality::Atomic,
2493                        failure: crate::model::FailureContract {
2494                            domains: Vec::new(),
2495                        },
2496                        effect: crate::model::Effect::Pure,
2497                        retry_limit: 0,
2498                        state: StateContract::stateless(),
2499                        subgraph_path: Vec::new(),
2500                    });
2501                }
2502                conversion_steps.insert((group_index, route_index), steps);
2503            }
2504        }
2505        let mut semantic = semantic_regions[region_index].clone();
2506        semantic.id = StepId(nodes.len() as u32);
2507        semantic_steps[region_index] = semantic.id;
2508        nodes.push(semantic);
2509    }
2510
2511    let mut invocation_ports = graph.invocation_inputs.clone();
2512    invocation_ports.sort_unstable();
2513    invocation_ports.dedup();
2514    for (invocation_index, invocation) in invocation_ports.iter().enumerate() {
2515        let region = owner[&invocation.node];
2516        let semantic = &semantic_regions[region];
2517        if semantic.semantic_nodes.first() != Some(&invocation.node) {
2518            return Err(CompileError::DuplicateInput(
2519                invocation.node,
2520                invocation.port.clone(),
2521            ));
2522        }
2523        let input_index = descriptors[&invocation.node]
2524            .inputs
2525            .iter()
2526            .position(|port| port.name == invocation.port)
2527            .ok_or_else(|| CompileError::UnknownPort(invocation.node, invocation.port.clone()))?;
2528        set_invocation(
2529            &mut nodes,
2530            semantic_steps[region],
2531            input_index,
2532            invocation_index as u32,
2533        )?;
2534    }
2535
2536    let mut buffers = Vec::new();
2537    for (group_index, group) in groups.iter().enumerate() {
2538        let producer = semantic_steps[group.producer_region];
2539        let mut consumers = Vec::new();
2540        for (route_index, route) in group.routes.iter().enumerate() {
2541            let path = &conversion_steps[&(group_index, route_index)];
2542            consumers.push(
2543                path.first()
2544                    .copied()
2545                    .unwrap_or(semantic_steps[route.consumer_region]),
2546            );
2547        }
2548        consumers.sort_unstable();
2549        consumers.dedup();
2550        let source_buffer = push_buffer(
2551            &mut buffers,
2552            group.layout,
2553            group.capacity_bytes,
2554            producer,
2555            consumers,
2556        );
2557        set_output(&mut nodes, producer, group.output_index, source_buffer)?;
2558        for (route_index, route) in group.routes.iter().enumerate() {
2559            let path = &conversion_steps[&(group_index, route_index)];
2560            if path.is_empty() {
2561                set_input(
2562                    &mut nodes,
2563                    semantic_steps[route.consumer_region],
2564                    route.input_index,
2565                    source_buffer,
2566                )?;
2567                continue;
2568            }
2569            set_input(&mut nodes, path[0], 0, source_buffer)?;
2570            for (index, step) in path.iter().copied().enumerate() {
2571                let next = path
2572                    .get(index + 1)
2573                    .copied()
2574                    .unwrap_or(semantic_steps[route.consumer_region]);
2575                let conversion = nodes[step.0 as usize]
2576                    .conversion
2577                    .as_ref()
2578                    .expect("conversion step");
2579                let output = push_buffer(
2580                    &mut buffers,
2581                    conversion.to,
2582                    conversion.max_output_bytes,
2583                    step,
2584                    alloc::vec![next],
2585                );
2586                set_output(&mut nodes, step, 0, output)?;
2587                if index + 1 < path.len() {
2588                    set_input(&mut nodes, next, 0, output)?;
2589                } else {
2590                    set_input(&mut nodes, next, route.input_index, output)?;
2591                }
2592            }
2593        }
2594    }
2595    canonicalize_buffers(&mut nodes, &mut buffers);
2596    synchronize_port_layouts(&mut nodes, &buffers);
2597    assign_aliases(&mut buffers);
2598    let arena_bytes = buffers
2599        .iter()
2600        .filter(|buffer| buffer.aliases.is_none())
2601        .try_fold(0u64, |sum, buffer| {
2602            sum.checked_add(buffer.capacity_bytes)
2603                .ok_or(CompileError::ResourceOverflow)
2604        })?;
2605    let workspace = nodes.iter().try_fold(0u64, |peak, node| {
2606        let mut bytes = node
2607            .resources
2608            .peak_bytes
2609            .checked_add(node.resources.scratch_bytes)
2610            .ok_or(CompileError::ResourceOverflow)?;
2611        if node.state.scope == StateScope::Invocation {
2612            bytes = bytes
2613                .checked_add(node.state.max_bytes)
2614                .ok_or(CompileError::ResourceOverflow)?;
2615        }
2616        Ok::<_, CompileError>(peak.max(bytes))
2617    })?;
2618    let peak = arena_bytes
2619        .checked_add(workspace)
2620        .ok_or(CompileError::ResourceOverflow)?;
2621    Ok((nodes, buffers, peak))
2622}
2623
2624fn lower_feedback(
2625    graph: &Graph,
2626    nodes: &mut [CompiledNode],
2627    max_feedback_edges: usize,
2628) -> Result<(Vec<crate::FeedbackPlan>, u64), CompileError> {
2629    if graph.feedback.len() > max_feedback_edges {
2630        return Err(CompileError::CompileLimitExceeded);
2631    }
2632    let mut feedback_edges = graph.feedback.clone();
2633    feedback_edges.sort_by_key(|edge| (edge.from.clone(), edge.to.clone()));
2634    if feedback_edges
2635        .windows(2)
2636        .any(|pair| pair[0].from == pair[1].from && pair[0].to == pair[1].to)
2637    {
2638        return Err(CompileError::InvalidFeedback(
2639            feedback_edges[0].to.node,
2640            feedback_edges[0].to.port.clone(),
2641        ));
2642    }
2643    let mut plans = Vec::with_capacity(feedback_edges.len());
2644    let mut state_bytes = 0u64;
2645    for edge in feedback_edges {
2646        let from_step_index = nodes
2647            .iter()
2648            .position(|node| node.semantic_nodes.contains(&edge.from.node))
2649            .ok_or(CompileError::UnknownNode(edge.from.node))?;
2650        let from_port = nodes[from_step_index]
2651            .output_ports
2652            .iter()
2653            .position(|port| port == &edge.from.port)
2654            .ok_or_else(|| CompileError::UnknownPort(edge.from.node, edge.from.port.clone()))?;
2655        let to_step_index = nodes
2656            .iter()
2657            .position(|node| node.semantic_nodes.contains(&edge.to.node))
2658            .ok_or(CompileError::UnknownNode(edge.to.node))?;
2659        let to_port = nodes[to_step_index]
2660            .input_ports
2661            .iter()
2662            .position(|port| port == &edge.to.port)
2663            .ok_or_else(|| CompileError::UnknownPort(edge.to.node, edge.to.port.clone()))?;
2664        let output_contract = &nodes[from_step_index].output_contracts[from_port];
2665        let input_contract = &nodes[to_step_index].input_contracts[to_port];
2666        if !compiled_port_contract_satisfies(output_contract, input_contract) {
2667            return Err(CompileError::InvalidFeedback(
2668                edge.to.node,
2669                edge.to.port.clone(),
2670            ));
2671        }
2672        let value_bytes = output_contract.max_bytes;
2673        let bytes = value_bytes
2674            .checked_mul(u64::from(edge.delay.invocations))
2675            .ok_or(CompileError::ResourceOverflow)?;
2676        let id = feedback_id(plans.len())?;
2677        nodes[to_step_index].input_bindings[to_port] = crate::InputBinding::Feedback(id);
2678        plans.push(crate::FeedbackPlan {
2679            id,
2680            from_step: nodes[from_step_index].id,
2681            from_port: from_port as u32,
2682            to_step: nodes[to_step_index].id,
2683            to_port: to_port as u32,
2684            delay: edge.delay,
2685            state_bytes: bytes,
2686        });
2687        state_bytes = state_bytes
2688            .checked_add(bytes)
2689            .ok_or(CompileError::ResourceOverflow)?;
2690    }
2691    Ok((plans, state_bytes))
2692}
2693
2694fn feedback_id(index: usize) -> Result<crate::FeedbackId, CompileError> {
2695    u32::try_from(index)
2696        .map(crate::FeedbackId)
2697        .map_err(|_| CompileError::CompileLimitExceeded)
2698}
2699
2700struct ConversionRequest<'a> {
2701    target: Target,
2702    semantic_type: &'a str,
2703    max_bytes: u64,
2704    from: Layout,
2705    targets: &'a [Layout],
2706    target_max_bytes: u64,
2707}
2708
2709struct ConversionBudget<'a> {
2710    searched: &'a mut usize,
2711    max_states: usize,
2712}
2713
2714fn conversion_path<'a>(
2715    registry: &'a KernelRegistry,
2716    request: ConversionRequest<'_>,
2717    budget: &mut ConversionBudget<'_>,
2718) -> Result<Option<Vec<&'a KernelDescriptor>>, CompileError> {
2719    let mut kernels: Vec<_> = registry
2720        .kernels
2721        .values()
2722        .filter(|kernel| {
2723            kernel.target == request.target
2724                && kernel.implements.is_empty()
2725                && kernel.determinism == crate::model::Determinism::BitExact
2726                && kernel.conversion.as_ref().is_some_and(|conversion| {
2727                    conversion.semantic_type == request.semantic_type
2728                        && kernel.input_layouts.contains(&conversion.from)
2729                        && kernel.output_layouts.contains(&conversion.to)
2730                })
2731        })
2732        .collect();
2733    kernels.sort_by_key(|kernel| (kernel.implementation_id, kernel.id));
2734    let mut frontier = alloc::vec![(
2735        request.from,
2736        request.max_bytes,
2737        Vec::new(),
2738        alloc::vec![request.from],
2739    )];
2740    let mut solutions = Vec::new();
2741    while let Some((layout, capacity_bytes, path, visited)) = frontier.pop() {
2742        *budget.searched = budget
2743            .searched
2744            .checked_add(1)
2745            .ok_or(CompileError::SearchLimitExceeded)?;
2746        if *budget.searched > budget.max_states {
2747            return Err(CompileError::SearchLimitExceeded);
2748        }
2749        if request.targets.contains(&layout)
2750            && capacity_bytes <= request.target_max_bytes
2751            && !path.is_empty()
2752        {
2753            solutions.push(path);
2754            continue;
2755        }
2756        if visited.len() >= 5 {
2757            continue;
2758        }
2759        for kernel in kernels.iter().rev() {
2760            let conversion = kernel.conversion.as_ref().expect("filtered");
2761            if conversion.from == layout
2762                && conversion.max_input_bytes >= capacity_bytes
2763                && !visited.contains(&conversion.to)
2764            {
2765                let mut next_path = path.clone();
2766                next_path.push(*kernel);
2767                let mut next_visited = visited.clone();
2768                next_visited.push(conversion.to);
2769                frontier.push((
2770                    conversion.to,
2771                    conversion.max_output_bytes,
2772                    next_path,
2773                    next_visited,
2774                ));
2775            }
2776        }
2777    }
2778    Ok(solutions.into_iter().min_by_key(|path| {
2779        let workspace = path
2780            .iter()
2781            .map(|kernel| {
2782                kernel
2783                    .resources
2784                    .peak_bytes
2785                    .saturating_add(kernel.resources.scratch_bytes)
2786            })
2787            .max()
2788            .unwrap_or_default();
2789        let implementations: Vec<_> = path.iter().map(|kernel| kernel.implementation_id).collect();
2790        let ids: Vec<_> = path.iter().map(|kernel| kernel.id).collect();
2791        (path.len(), workspace, implementations, ids)
2792    }))
2793}
2794
2795fn push_buffer(
2796    buffers: &mut Vec<BufferPlan>,
2797    layout: Layout,
2798    capacity_bytes: u64,
2799    producer: StepId,
2800    mut consumers: Vec<StepId>,
2801) -> BufferId {
2802    consumers.sort_unstable();
2803    consumers.dedup();
2804    let last_consumer = *consumers.last().expect("physical buffer has a consumer");
2805    let id = BufferId(buffers.len() as u32);
2806    buffers.push(BufferPlan {
2807        id,
2808        layout,
2809        capacity_bytes,
2810        producer,
2811        consumers,
2812        last_consumer,
2813        aliases: None,
2814    });
2815    id
2816}
2817
2818fn set_input(
2819    nodes: &mut [CompiledNode],
2820    step: StepId,
2821    index: usize,
2822    buffer: BufferId,
2823) -> Result<(), CompileError> {
2824    let binding = nodes
2825        .get_mut(step.0 as usize)
2826        .and_then(|node| node.input_bindings.get_mut(index))
2827        .ok_or(CompileError::ResourceOverflow)?;
2828    *binding = crate::model::InputBinding::Buffer(buffer);
2829    Ok(())
2830}
2831
2832fn set_invocation(
2833    nodes: &mut [CompiledNode],
2834    step: StepId,
2835    index: usize,
2836    invocation: u32,
2837) -> Result<(), CompileError> {
2838    let binding = nodes
2839        .get_mut(step.0 as usize)
2840        .and_then(|node| node.input_bindings.get_mut(index))
2841        .ok_or(CompileError::ResourceOverflow)?;
2842    if !matches!(binding, crate::model::InputBinding::Absent) {
2843        return Err(CompileError::ResourceOverflow);
2844    }
2845    *binding = crate::model::InputBinding::Invocation(invocation);
2846    Ok(())
2847}
2848
2849fn set_output(
2850    nodes: &mut [CompiledNode],
2851    step: StepId,
2852    index: usize,
2853    buffer: BufferId,
2854) -> Result<(), CompileError> {
2855    let binding = nodes
2856        .get_mut(step.0 as usize)
2857        .and_then(|node| node.output_bindings.get_mut(index))
2858        .ok_or(CompileError::ResourceOverflow)?;
2859    *binding = OutputBinding::Buffer(buffer);
2860    Ok(())
2861}
2862
2863fn canonicalize_buffers(nodes: &mut [CompiledNode], buffers: &mut [BufferPlan]) {
2864    let mut order: Vec<_> = (0..buffers.len()).collect();
2865    order.sort_by_key(|index| {
2866        let buffer = &buffers[*index];
2867        (buffer.producer, buffer.last_consumer, buffer.id)
2868    });
2869    let mut remap = BTreeMap::new();
2870    for (new, old) in order.iter().copied().enumerate() {
2871        remap.insert(buffers[old].id, BufferId(new as u32));
2872    }
2873    let old = buffers.to_vec();
2874    for (new, old_index) in order.into_iter().enumerate() {
2875        buffers[new] = BufferPlan {
2876            id: BufferId(new as u32),
2877            aliases: None,
2878            ..old[old_index].clone()
2879        };
2880    }
2881    for node in nodes {
2882        for binding in &mut node.input_bindings {
2883            if let crate::model::InputBinding::Buffer(buffer) = binding {
2884                *buffer = remap[buffer];
2885            }
2886        }
2887        for binding in &mut node.output_bindings {
2888            if let OutputBinding::Buffer(buffer) = binding {
2889                *buffer = remap[buffer];
2890            }
2891        }
2892    }
2893}
2894
2895fn assign_aliases(buffers: &mut [BufferPlan]) {
2896    struct Slot {
2897        root: BufferId,
2898        layout: Layout,
2899        capacity_bytes: u64,
2900        available_after: StepId,
2901    }
2902    let mut slots: Vec<Slot> = Vec::new();
2903    for buffer in buffers {
2904        let candidate = slots
2905            .iter_mut()
2906            .filter(|slot| slot.available_after < buffer.producer)
2907            .filter(|slot| {
2908                slot.layout == buffer.layout && slot.capacity_bytes >= buffer.capacity_bytes
2909            })
2910            .min_by_key(|slot| (slot.capacity_bytes, slot.root));
2911        if let Some(slot) = candidate {
2912            buffer.aliases = Some(slot.root);
2913            slot.available_after = buffer.last_consumer;
2914        } else {
2915            slots.push(Slot {
2916                root: buffer.id,
2917                layout: buffer.layout,
2918                capacity_bytes: buffer.capacity_bytes,
2919                available_after: buffer.last_consumer,
2920            });
2921        }
2922    }
2923}
2924
2925fn hash_graph(graph: &Graph, descriptors: &BTreeMap<NodeId, &NodeDescriptor>) -> [u8; 32] {
2926    let mut hasher = blake3::Hasher::new_derive_key("blut.graph.v3");
2927    put_u32(&mut hasher, graph.version);
2928    let mut nodes = graph.nodes.clone();
2929    nodes.sort_by_key(|node| node.id);
2930    put_u32(&mut hasher, nodes.len() as u32);
2931    for node in nodes {
2932        put_u32(&mut hasher, node.id.0);
2933        put_str(&mut hasher, &node.descriptor);
2934        put_u32(&mut hasher, node.descriptor_version);
2935        hash_descriptor(&mut hasher, descriptors[&node.id]);
2936        put_u32(&mut hasher, node.config.len() as u32);
2937        for (key, value) in node.config {
2938            put_str(&mut hasher, &key);
2939            hash_config_value(&mut hasher, &value);
2940        }
2941    }
2942    let mut edges = graph.edges.clone();
2943    edges.sort_by_key(|edge| (edge.from.clone(), edge.to.clone()));
2944    put_u32(&mut hasher, edges.len() as u32);
2945    for edge in edges {
2946        put_u32(&mut hasher, edge.from.node.0);
2947        put_str(&mut hasher, &edge.from.port);
2948        put_u32(&mut hasher, edge.to.node.0);
2949        put_str(&mut hasher, &edge.to.port);
2950    }
2951    let mut feedback = graph.feedback.clone();
2952    feedback.sort_by_key(|edge| (edge.from.clone(), edge.to.clone()));
2953    put_u32(&mut hasher, feedback.len() as u32);
2954    for edge in feedback {
2955        put_port_ref(&mut hasher, &edge.from);
2956        put_port_ref(&mut hasher, &edge.to);
2957        put_u32(&mut hasher, edge.delay.invocations);
2958        hash_delay_initial(&mut hasher, &edge.delay.initial);
2959    }
2960    let mut invocation_ports = graph.invocation_inputs.clone();
2961    invocation_ports.sort_unstable();
2962    invocation_ports.dedup();
2963    put_u32(&mut hasher, invocation_ports.len() as u32);
2964    for port in invocation_ports {
2965        put_u32(&mut hasher, port.node.0);
2966        put_str(&mut hasher, &port.port);
2967    }
2968    let mut capabilities: Vec<_> = graph
2969        .required_capabilities
2970        .iter()
2971        .map(|capability| capability.0.as_str())
2972        .collect();
2973    capabilities.sort_unstable();
2974    capabilities.dedup();
2975    put_str_set(&mut hasher, &capabilities);
2976    let mut proofs: Vec<_> = graph.required_proofs.iter().map(String::as_str).collect();
2977    proofs.sort_unstable();
2978    proofs.dedup();
2979    put_str_set(&mut hasher, &proofs);
2980    let mut policy: Vec<_> = graph.policy.iter().map(String::as_str).collect();
2981    policy.sort_unstable();
2982    policy.dedup();
2983    put_str_set(&mut hasher, &policy);
2984    put_u32(&mut hasher, u32::from(graph.minimum_fidelity));
2985    hash_session(&mut hasher, graph.session.as_ref());
2986    *hasher.finalize().as_bytes()
2987}
2988
2989fn hash_descriptor(hasher: &mut blake3::Hasher, descriptor: &NodeDescriptor) {
2990    fn hash_ports(hasher: &mut blake3::Hasher, ports: &[PortDescriptor]) {
2991        let mut ports = ports.to_vec();
2992        ports.sort_by(|a, b| a.name.cmp(&b.name));
2993        put_u32(hasher, ports.len() as u32);
2994        for port in ports {
2995            put_str(hasher, &port.name);
2996            put_str(hasher, &port.semantic_type);
2997            put_u32(hasher, u32::from(port.optional));
2998            let mut layouts = port.layouts;
2999            layouts.sort_unstable();
3000            layouts.dedup();
3001            put_u32(hasher, layouts.len() as u32);
3002            for layout in layouts {
3003                put_u32(hasher, layout.token());
3004            }
3005            hasher.update(&port.max_bytes.to_le_bytes());
3006            hash_domain_type(hasher, &port.domain);
3007            hash_proof(hasher, &port.proof);
3008            hash_policy(hasher, &port.policy);
3009            hash_fidelity(hasher, &port.fidelity);
3010            hash_extent(hasher, &port.extent);
3011            hash_lease(hasher, &port.lease);
3012        }
3013    }
3014    hash_ports(hasher, &descriptor.inputs);
3015    hash_ports(hasher, &descriptor.outputs);
3016    let mut capabilities: Vec<_> = descriptor
3017        .capabilities
3018        .iter()
3019        .map(|capability| capability.0.as_str())
3020        .collect();
3021    capabilities.sort_unstable();
3022    capabilities.dedup();
3023    put_str_set(hasher, &capabilities);
3024    let mut targets = descriptor.targets.clone();
3025    targets.sort_unstable();
3026    targets.dedup();
3027    put_u32(hasher, targets.len() as u32);
3028    for target in targets {
3029        put_u32(hasher, target.token());
3030    }
3031    hasher.update(&descriptor.resources.peak_bytes.to_le_bytes());
3032    hasher.update(&descriptor.resources.scratch_bytes.to_le_bytes());
3033    put_u32(hasher, u32::from(descriptor.resources.threads));
3034    match &descriptor.resources.device {
3035        Some(device) => {
3036            hasher.update(&[1]);
3037            put_str(hasher, device);
3038        }
3039        None => {
3040            hasher.update(&[0]);
3041        }
3042    }
3043    put_u32(hasher, descriptor.determinism as u32);
3044    hash_config_schema(hasher, &descriptor.config);
3045    hash_state(hasher, &descriptor.state);
3046    match &descriptor.subgraph {
3047        Some(lowering) => {
3048            hasher.update(&[1]);
3049            hasher.update(&lowering.subgraph.0);
3050            hash_port_maps(hasher, &lowering.input_map);
3051            hash_port_maps(hasher, &lowering.output_map);
3052            put_u32(hasher, lowering.config_map.len() as u32);
3053            for map in &lowering.config_map {
3054                put_str(hasher, &map.outer);
3055                put_u32(hasher, map.node.0);
3056                put_str(hasher, &map.inner);
3057            }
3058        }
3059        None => {
3060            hasher.update(&[0]);
3061        }
3062    }
3063    let mut proof_requires: Vec<_> = descriptor
3064        .proof
3065        .requires
3066        .iter()
3067        .map(String::as_str)
3068        .collect();
3069    proof_requires.sort_unstable();
3070    proof_requires.dedup();
3071    put_str_set(hasher, &proof_requires);
3072    let mut proof_provides: Vec<_> = descriptor
3073        .proof
3074        .provides
3075        .iter()
3076        .map(String::as_str)
3077        .collect();
3078    proof_provides.sort_unstable();
3079    proof_provides.dedup();
3080    put_str_set(hasher, &proof_provides);
3081    let mut proof_invalidates: Vec<_> = descriptor
3082        .proof
3083        .invalidates
3084        .iter()
3085        .map(String::as_str)
3086        .collect();
3087    proof_invalidates.sort_unstable();
3088    proof_invalidates.dedup();
3089    put_str_set(hasher, &proof_invalidates);
3090    let mut policy_requires: Vec<_> = descriptor
3091        .policy
3092        .requires
3093        .iter()
3094        .map(String::as_str)
3095        .collect();
3096    policy_requires.sort_unstable();
3097    policy_requires.dedup();
3098    put_str_set(hasher, &policy_requires);
3099    let mut policy_adds: Vec<_> = descriptor.policy.adds.iter().map(String::as_str).collect();
3100    policy_adds.sort_unstable();
3101    policy_adds.dedup();
3102    put_str_set(hasher, &policy_adds);
3103    put_u32(hasher, u32::from(descriptor.fidelity.minimum_input));
3104    put_u32(hasher, u32::from(descriptor.fidelity.maximum_loss));
3105    put_u32(hasher, descriptor.partiality as u32);
3106    let mut failure_domains: Vec<_> = descriptor
3107        .failure
3108        .domains
3109        .iter()
3110        .map(String::as_str)
3111        .collect();
3112    failure_domains.sort_unstable();
3113    failure_domains.dedup();
3114    put_str_set(hasher, &failure_domains);
3115    put_u32(hasher, descriptor.effect as u32);
3116    put_u32(hasher, u32::from(descriptor.retry_limit));
3117}
3118
3119fn put_str_set(hasher: &mut blake3::Hasher, values: &[&str]) {
3120    put_u32(hasher, values.len() as u32);
3121    for value in values {
3122        put_str(hasher, value);
3123    }
3124}
3125
3126fn hash_config_value(hasher: &mut blake3::Hasher, value: &crate::ConfigValue) {
3127    match value {
3128        crate::ConfigValue::Bool(value) => {
3129            hasher.update(&[0, u8::from(*value)]);
3130        }
3131        crate::ConfigValue::I64(value) => {
3132            hasher.update(&[1]);
3133            hasher.update(&value.to_le_bytes());
3134        }
3135        crate::ConfigValue::U64(value) => {
3136            hasher.update(&[2]);
3137            hasher.update(&value.to_le_bytes());
3138        }
3139        crate::ConfigValue::Text(value) => {
3140            hasher.update(&[3]);
3141            put_str(hasher, value);
3142        }
3143        crate::ConfigValue::Bytes(value) => {
3144            hasher.update(&[4]);
3145            put_u32(hasher, value.len() as u32);
3146            hasher.update(value);
3147        }
3148    }
3149}
3150
3151fn hash_config_schema(hasher: &mut blake3::Hasher, schema: &crate::ConfigSchema) {
3152    put_u32(hasher, schema.fields.len() as u32);
3153    for field in &schema.fields {
3154        put_str(hasher, &field.name);
3155        match &field.value_type {
3156            crate::ConfigType::Bool => {
3157                hasher.update(&[0]);
3158            }
3159            crate::ConfigType::I64 { minimum, maximum } => {
3160                hasher.update(&[1]);
3161                hasher.update(&minimum.to_le_bytes());
3162                hasher.update(&maximum.to_le_bytes());
3163            }
3164            crate::ConfigType::U64 { minimum, maximum } => {
3165                hasher.update(&[2]);
3166                hasher.update(&minimum.to_le_bytes());
3167                hasher.update(&maximum.to_le_bytes());
3168            }
3169            crate::ConfigType::Text { max_bytes } => {
3170                hasher.update(&[3]);
3171                put_u32(hasher, *max_bytes);
3172            }
3173            crate::ConfigType::Choice { values } => {
3174                hasher.update(&[4]);
3175                put_u32(hasher, values.len() as u32);
3176                for value in values {
3177                    put_str(hasher, value);
3178                }
3179            }
3180            crate::ConfigType::Bytes { max_bytes } => {
3181                hasher.update(&[5]);
3182                put_u32(hasher, *max_bytes);
3183            }
3184        };
3185        hasher.update(&[u8::from(field.required)]);
3186        match &field.default {
3187            Some(value) => {
3188                hasher.update(&[1]);
3189                hash_config_value(hasher, value);
3190            }
3191            None => {
3192                hasher.update(&[0]);
3193            }
3194        }
3195    }
3196}
3197
3198fn hash_domain_type(hasher: &mut blake3::Hasher, domain: &crate::DomainType) {
3199    // Tokens are opaque, so the plan hash folds their BYTES. The pre-migration shape
3200    // hashed a numeric discriminant per enum variant plus the string only for the
3201    // `Unknown` escape; hashing the string uniformly is what makes the compiler
3202    // independent of any domain's variant list. Absolute plan ids therefore move
3203    // once, at this commit -- nothing pins one (every assertion in this crate is
3204    // relational: `assert_eq!(a.plan_id, b.plan_id)`), and the alpha has no
3205    // persisted plan cache to invalidate.
3206    put_str(hasher, domain.root.as_str());
3207    put_str(hasher, domain.view.as_str());
3208}
3209
3210fn hash_proof(hasher: &mut blake3::Hasher, proof: &crate::ProofContract) {
3211    let requires = proof
3212        .requires
3213        .iter()
3214        .map(String::as_str)
3215        .collect::<Vec<_>>();
3216    let provides = proof
3217        .provides
3218        .iter()
3219        .map(String::as_str)
3220        .collect::<Vec<_>>();
3221    let invalidates = proof
3222        .invalidates
3223        .iter()
3224        .map(String::as_str)
3225        .collect::<Vec<_>>();
3226    put_str_set(hasher, &requires);
3227    put_str_set(hasher, &provides);
3228    put_str_set(hasher, &invalidates);
3229}
3230
3231fn hash_policy(hasher: &mut blake3::Hasher, policy: &crate::PolicyContract) {
3232    let requires = policy
3233        .requires
3234        .iter()
3235        .map(String::as_str)
3236        .collect::<Vec<_>>();
3237    let adds = policy.adds.iter().map(String::as_str).collect::<Vec<_>>();
3238    put_str_set(hasher, &requires);
3239    put_str_set(hasher, &adds);
3240}
3241
3242fn hash_fidelity(hasher: &mut blake3::Hasher, fidelity: &crate::FidelityContract) {
3243    put_u32(hasher, u32::from(fidelity.minimum_input));
3244    put_u32(hasher, u32::from(fidelity.maximum_loss));
3245}
3246
3247fn hash_extent(hasher: &mut blake3::Hasher, extent: &crate::ExtentContract) {
3248    hasher.update(&[extent.rank]);
3249    put_u32(hasher, extent.maximum_shape.len() as u32);
3250    for size in &extent.maximum_shape {
3251        hasher.update(&size.to_le_bytes());
3252    }
3253    hasher.update(&extent.max_elements.to_le_bytes());
3254    hasher.update(&[u8::from(extent.ragged), u8::from(extent.sparse)]);
3255}
3256
3257fn hash_lease(hasher: &mut blake3::Hasher, lease: &crate::LeaseContract) {
3258    put_u32(hasher, lease.access as u32);
3259    put_u32(hasher, lease.lifetime as u32);
3260    hasher.update(&[
3261        u8::from(lease.zero_copy_permitted),
3262        u8::from(lease.contiguous_required),
3263    ]);
3264}
3265
3266fn hash_state(hasher: &mut blake3::Hasher, state: &StateContract) {
3267    put_u32(hasher, state.scope as u32);
3268    hasher.update(&state.max_bytes.to_le_bytes());
3269    put_u32(hasher, state.checkpoint.mode as u32);
3270    hasher.update(&state.checkpoint.max_snapshot_bytes.to_le_bytes());
3271    put_u32(hasher, state.checkpoint.max_interval_invocations);
3272}
3273
3274fn hash_port_maps(hasher: &mut blake3::Hasher, maps: &[crate::PortMap]) {
3275    put_u32(hasher, maps.len() as u32);
3276    for map in maps {
3277        put_str(hasher, &map.outer);
3278        put_str(hasher, &map.inner);
3279    }
3280}
3281
3282fn hash_delay_initial(hasher: &mut blake3::Hasher, initial: &crate::DelayInitial) {
3283    match initial {
3284        crate::DelayInitial::Absent => {
3285            hasher.update(&[0]);
3286        }
3287        crate::DelayInitial::Zeroed => {
3288            hasher.update(&[1]);
3289        }
3290        crate::DelayInitial::ContentId(id) => {
3291            hasher.update(&[2]);
3292            hasher.update(id);
3293        }
3294    };
3295}
3296
3297fn hash_session(hasher: &mut blake3::Hasher, session: Option<&crate::SessionContract>) {
3298    match session {
3299        Some(session) => {
3300            hasher.update(&[1]);
3301            put_str(hasher, &session.namespace);
3302            put_u32(hasher, session.max_concurrent_sessions);
3303            hasher.update(&session.max_idle_millis.to_le_bytes());
3304            hasher.update(&[u8::from(session.reset_on_plan_change)]);
3305        }
3306        None => {
3307            hasher.update(&[0]);
3308        }
3309    }
3310}
3311
3312fn hash_compiled_ports(hasher: &mut blake3::Hasher, ports: &[CompiledPortContract]) {
3313    put_u32(hasher, ports.len() as u32);
3314    for port in ports {
3315        put_str(hasher, &port.name);
3316        put_str(hasher, &port.semantic_type);
3317        hasher.update(&[u8::from(port.optional)]);
3318        put_u32(hasher, port.layout.token());
3319        hasher.update(&port.max_bytes.to_le_bytes());
3320        hash_domain_type(hasher, &port.domain);
3321        hash_proof(hasher, &port.proof);
3322        hash_policy(hasher, &port.policy);
3323        hash_fidelity(hasher, &port.fidelity);
3324        hash_extent(hasher, &port.extent);
3325        hash_lease(hasher, &port.lease);
3326    }
3327}
3328
3329fn put_port_ref(hasher: &mut blake3::Hasher, port: &crate::PortRef) {
3330    put_u32(hasher, port.node.0);
3331    put_str(hasher, &port.port);
3332}
3333
3334pub(crate) fn hash_plan(plan: &CompiledPlan) -> [u8; 32] {
3335    let mut hasher = blake3::Hasher::new_derive_key("blut.compiled-plan.v3");
3336    put_u32(&mut hasher, plan.schema_version);
3337    hasher.update(&plan.graph_id.0);
3338    put_u32(&mut hasher, plan.realm.token());
3339    put_u32(&mut hasher, plan.order.len() as u32);
3340    for node in &plan.order {
3341        put_u32(&mut hasher, node.0);
3342    }
3343    put_u32(&mut hasher, plan.nodes.len() as u32);
3344    for node in &plan.nodes {
3345        put_u32(&mut hasher, node.id.0);
3346        put_u32(&mut hasher, node.semantic_nodes.len() as u32);
3347        for semantic in &node.semantic_nodes {
3348            put_u32(&mut hasher, semantic.0);
3349        }
3350        put_u32(&mut hasher, node.semantic_types.len() as u32);
3351        for semantic_type in &node.semantic_types {
3352            put_str(&mut hasher, &semantic_type.type_name);
3353            put_u32(&mut hasher, semantic_type.version);
3354        }
3355        put_u32(&mut hasher, node.semantic_configs.len() as u32);
3356        for config in &node.semantic_configs {
3357            put_u32(&mut hasher, config.len() as u32);
3358            for (key, value) in config {
3359                put_str(&mut hasher, key);
3360                hash_config_value(&mut hasher, value);
3361            }
3362        }
3363        put_u32(&mut hasher, node.kernel.0);
3364        hasher.update(&node.implementation_id.0);
3365        hasher.update(&node.resources.peak_bytes.to_le_bytes());
3366        hasher.update(&node.resources.scratch_bytes.to_le_bytes());
3367        put_u32(&mut hasher, u32::from(node.resources.threads));
3368        match &node.resources.device {
3369            Some(device) => {
3370                hasher.update(&[1]);
3371                put_str(&mut hasher, device);
3372            }
3373            None => {
3374                hasher.update(&[0]);
3375            }
3376        }
3377        put_u32(&mut hasher, node.determinism as u32);
3378        put_str(&mut hasher, &node.lowering);
3379        match &node.conversion {
3380            Some(conversion) => {
3381                hasher.update(&[1]);
3382                put_str(&mut hasher, &conversion.semantic_type);
3383                put_u32(&mut hasher, conversion.from.token());
3384                put_u32(&mut hasher, conversion.to.token());
3385                hasher.update(&conversion.max_input_bytes.to_le_bytes());
3386                hasher.update(&conversion.max_output_bytes.to_le_bytes());
3387            }
3388            None => {
3389                hasher.update(&[0]);
3390            }
3391        }
3392        put_u32(&mut hasher, node.input_ports.len() as u32);
3393        for port in &node.input_ports {
3394            put_str(&mut hasher, port);
3395        }
3396        put_u32(&mut hasher, node.output_ports.len() as u32);
3397        for port in &node.output_ports {
3398            put_str(&mut hasher, port);
3399        }
3400        hash_compiled_ports(&mut hasher, &node.input_contracts);
3401        hash_compiled_ports(&mut hasher, &node.output_contracts);
3402        put_u32(&mut hasher, node.input_bindings.len() as u32);
3403        for binding in &node.input_bindings {
3404            match binding {
3405                crate::model::InputBinding::Buffer(buffer) => {
3406                    hasher.update(&[1]);
3407                    put_u32(&mut hasher, buffer.0);
3408                }
3409                crate::model::InputBinding::Invocation(invocation) => {
3410                    hasher.update(&[2]);
3411                    put_u32(&mut hasher, *invocation);
3412                }
3413                crate::model::InputBinding::Feedback(feedback) => {
3414                    hasher.update(&[3]);
3415                    put_u32(&mut hasher, feedback.0);
3416                }
3417                crate::model::InputBinding::Absent => {
3418                    hasher.update(&[0]);
3419                }
3420            }
3421        }
3422        put_u32(&mut hasher, node.output_bindings.len() as u32);
3423        for binding in &node.output_bindings {
3424            match binding {
3425                OutputBinding::Buffer(buffer) => {
3426                    hasher.update(&[1]);
3427                    put_u32(&mut hasher, buffer.0);
3428                }
3429                OutputBinding::Terminal => {
3430                    hasher.update(&[0]);
3431                }
3432            }
3433        }
3434        put_u32(&mut hasher, node.partiality as u32);
3435        let mut failure_domains = node.failure.domains.clone();
3436        failure_domains.sort_unstable();
3437        failure_domains.dedup();
3438        put_u32(&mut hasher, failure_domains.len() as u32);
3439        for domain in failure_domains {
3440            put_str(&mut hasher, &domain);
3441        }
3442        put_u32(&mut hasher, node.effect as u32);
3443        put_u32(&mut hasher, u32::from(node.retry_limit));
3444        hash_state(&mut hasher, &node.state);
3445        put_u32(&mut hasher, node.subgraph_path.len() as u32);
3446        for subgraph in &node.subgraph_path {
3447            hasher.update(&subgraph.0);
3448        }
3449    }
3450    put_u32(&mut hasher, plan.invocation_ports.len() as u32);
3451    for port in &plan.invocation_ports {
3452        put_u32(&mut hasher, port.node.0);
3453        put_str(&mut hasher, &port.port);
3454    }
3455    put_u32(&mut hasher, plan.buffers.len() as u32);
3456    for buffer in &plan.buffers {
3457        put_u32(&mut hasher, buffer.id.0);
3458        put_u32(&mut hasher, buffer.layout.token());
3459        hasher.update(&buffer.capacity_bytes.to_le_bytes());
3460        put_u32(&mut hasher, buffer.producer.0);
3461        put_u32(&mut hasher, buffer.consumers.len() as u32);
3462        for consumer in &buffer.consumers {
3463            put_u32(&mut hasher, consumer.0);
3464        }
3465        put_u32(&mut hasher, buffer.last_consumer.0);
3466        match buffer.aliases {
3467            Some(alias) => {
3468                hasher.update(&[1]);
3469                put_u32(&mut hasher, alias.0);
3470            }
3471            None => {
3472                hasher.update(&[0]);
3473            }
3474        };
3475    }
3476    put_u32(&mut hasher, plan.feedback.len() as u32);
3477    for feedback in &plan.feedback {
3478        put_u32(&mut hasher, feedback.id.0);
3479        put_u32(&mut hasher, feedback.from_step.0);
3480        put_u32(&mut hasher, feedback.from_port);
3481        put_u32(&mut hasher, feedback.to_step.0);
3482        put_u32(&mut hasher, feedback.to_port);
3483        put_u32(&mut hasher, feedback.delay.invocations);
3484        hash_delay_initial(&mut hasher, &feedback.delay.initial);
3485        hasher.update(&feedback.state_bytes.to_le_bytes());
3486    }
3487    put_u32(&mut hasher, plan.propagated_proofs.len() as u32);
3488    for proof in &plan.propagated_proofs {
3489        put_str(&mut hasher, proof);
3490    }
3491    put_u32(&mut hasher, plan.propagated_policy.len() as u32);
3492    for policy in &plan.propagated_policy {
3493        put_str(&mut hasher, policy);
3494    }
3495    put_u32(&mut hasher, u32::from(plan.resulting_fidelity));
3496    hasher.update(&plan.peak_bytes.to_le_bytes());
3497    hasher.update(&plan.persistent_state_bytes.to_le_bytes());
3498    hash_session(&mut hasher, plan.session.as_ref());
3499    *hasher.finalize().as_bytes()
3500}
3501
3502fn put_str(hasher: &mut blake3::Hasher, value: &str) {
3503    put_u32(hasher, value.len() as u32);
3504    hasher.update(value.as_bytes());
3505}
3506
3507fn put_u32(hasher: &mut blake3::Hasher, value: u32) {
3508    hasher.update(&value.to_le_bytes());
3509}
3510
3511#[cfg(test)]
3512mod tests {
3513    use alloc::collections::BTreeMap;
3514    use alloc::string::{String, ToString};
3515    use alloc::vec;
3516
3517    use super::*;
3518    use crate::model::{
3519        Capability, Determinism, DomainToken, DomainType, Effect, FidelityContract,
3520        ImplementationId, KernelDescriptor, Layout, NodeInstance, NodeTypeRef, PolicyContract,
3521        PortRef, ProofContract, ResourceEnvelope,
3522    };
3523
3524    struct SemanticKernels;
3525
3526    impl crate::KernelExecutor for SemanticKernels {
3527        type Value = u32;
3528
3529        fn execute(
3530            &mut self,
3531            node: &CompiledNode,
3532            inputs: &[Option<&Self::Value>],
3533        ) -> Result<Vec<Self::Value>, crate::ExecutionError> {
3534            let input = inputs
3535                .first()
3536                .and_then(|value| *value)
3537                .copied()
3538                .unwrap_or_default();
3539            let value = match node.implementation_id.0[0] {
3540                1 => input + 10,
3541                4 => input * 2,
3542                7 => input.saturating_sub(3),
3543                100 => (input + 10) * 2,
3544                other => panic!("unexpected test implementation {other}"),
3545            };
3546            Ok(vec![value; node.output_bindings.len()])
3547        }
3548    }
3549
3550    struct NoTransactions;
3551
3552    impl crate::TransactionalSink for NoTransactions {
3553        fn prepare(&mut self, _idempotency_key: &str) -> Result<(), crate::ExecutionError> {
3554            Ok(())
3555        }
3556
3557        fn commit(&mut self, _idempotency_key: &str) -> Result<String, crate::ExecutionError> {
3558            unreachable!("fixture nodes are pure")
3559        }
3560
3561        fn abort(&mut self, _idempotency_key: &str) {}
3562    }
3563
3564    fn descriptor(name: &str, input: bool) -> NodeDescriptor {
3565        NodeDescriptor {
3566            type_name: name.to_string(),
3567            version: 1,
3568            inputs: if input {
3569                vec![PortDescriptor {
3570                    name: "in".to_string(),
3571                    semantic_type: "sample.block".to_string(),
3572                    optional: false,
3573                    layouts: vec![Layout::Canonical],
3574                    max_bytes: 64,
3575                    ..PortDescriptor::opaque("in", "sample.block", 64)
3576                }]
3577            } else {
3578                vec![]
3579            },
3580            outputs: vec![PortDescriptor {
3581                name: "out".to_string(),
3582                semantic_type: "sample.block".to_string(),
3583                optional: false,
3584                layouts: vec![Layout::Canonical],
3585                max_bytes: 64,
3586                ..PortDescriptor::opaque("out", "sample.block", 64)
3587            }],
3588            capabilities: vec![Capability("sample".to_string())],
3589            targets: vec![Target::Host, Target::McuAot, Target::BlutDurable],
3590            resources: ResourceEnvelope::bounded(64, 0, 1),
3591            determinism: Determinism::BitExact,
3592            config: crate::ConfigSchema {
3593                fields: vec![crate::ConfigField {
3594                    name: "gain".into(),
3595                    value_type: crate::ConfigType::Text { max_bytes: 16 },
3596                    required: false,
3597                    default: None,
3598                }],
3599            },
3600            state: StateContract::stateless(),
3601            subgraph: None,
3602            proof: ProofContract {
3603                requires: vec![],
3604                provides: vec![format!("{name}.verified")],
3605                invalidates: vec![],
3606            },
3607            policy: PolicyContract {
3608                requires: vec![],
3609                adds: vec![],
3610            },
3611            fidelity: FidelityContract {
3612                minimum_input: 0,
3613                maximum_loss: 0,
3614            },
3615            partiality: crate::Partiality::Atomic,
3616            failure: crate::FailureContract { domains: vec![] },
3617            effect: Effect::Pure,
3618            retry_limit: 0,
3619        }
3620    }
3621
3622    fn fixture(reverse: bool) -> (KernelRegistry, Graph) {
3623        let mut registry = KernelRegistry::default();
3624        for (index, name) in ["source", "process", "sink"].into_iter().enumerate() {
3625            registry
3626                .register_descriptor(descriptor(name, index != 0))
3627                .unwrap();
3628            for (target_index, target) in [Target::Host, Target::McuAot, Target::BlutDurable]
3629                .into_iter()
3630                .enumerate()
3631            {
3632                registry
3633                    .register_kernel(KernelDescriptor {
3634                        id: KernelId((index * 3 + target_index) as u32),
3635                        implements: vec![NodeTypeRef {
3636                            type_name: name.to_string(),
3637                            version: 1,
3638                        }],
3639                        implementation_id: ImplementationId(
3640                            [(index * 3 + target_index + 1) as u8; 32],
3641                        ),
3642                        conversion: None,
3643                        target,
3644                        input_layouts: vec![Layout::Canonical],
3645                        output_layouts: vec![Layout::Canonical],
3646                        resources: ResourceEnvelope::bounded(64, 0, 1),
3647                        determinism: Determinism::BitExact,
3648                        lowering: "test".to_string(),
3649                    })
3650                    .unwrap();
3651            }
3652        }
3653        for (target_index, target) in [Target::Host, Target::McuAot, Target::BlutDurable]
3654            .into_iter()
3655            .enumerate()
3656        {
3657            registry
3658                .register_kernel(KernelDescriptor {
3659                    id: KernelId(100 + target_index as u32),
3660                    implements: vec![
3661                        NodeTypeRef {
3662                            type_name: "source".to_string(),
3663                            version: 1,
3664                        },
3665                        NodeTypeRef {
3666                            type_name: "process".to_string(),
3667                            version: 1,
3668                        },
3669                    ],
3670                    implementation_id: ImplementationId([100 + target_index as u8; 32]),
3671                    conversion: None,
3672                    target,
3673                    input_layouts: vec![Layout::Canonical],
3674                    output_layouts: vec![Layout::Canonical],
3675                    resources: ResourceEnvelope::bounded(64, 0, 1),
3676                    determinism: Determinism::BitExact,
3677                    lowering: "test.fused.source-process".to_string(),
3678                })
3679                .unwrap();
3680        }
3681        let mut nodes = vec![
3682            NodeInstance {
3683                id: NodeId(0),
3684                descriptor: "source".to_string(),
3685                descriptor_version: 1,
3686                config: BTreeMap::new(),
3687            },
3688            NodeInstance {
3689                id: NodeId(1),
3690                descriptor: "process".to_string(),
3691                descriptor_version: 1,
3692                config: BTreeMap::new(),
3693            },
3694            NodeInstance {
3695                id: NodeId(2),
3696                descriptor: "sink".to_string(),
3697                descriptor_version: 1,
3698                config: BTreeMap::new(),
3699            },
3700        ];
3701        if reverse {
3702            nodes.reverse();
3703        }
3704        let graph = Graph {
3705            version: 3,
3706            nodes,
3707            edges: vec![
3708                Edge {
3709                    from: PortRef {
3710                        node: NodeId(0),
3711                        port: "out".to_string(),
3712                    },
3713                    to: PortRef {
3714                        node: NodeId(1),
3715                        port: "in".to_string(),
3716                    },
3717                },
3718                Edge {
3719                    from: PortRef {
3720                        node: NodeId(1),
3721                        port: "out".to_string(),
3722                    },
3723                    to: PortRef {
3724                        node: NodeId(2),
3725                        port: "in".to_string(),
3726                    },
3727                },
3728            ],
3729            feedback: vec![],
3730            invocation_inputs: vec![],
3731            required_capabilities: vec![Capability("sample".to_string())],
3732            required_proofs: vec![],
3733            policy: vec![],
3734            minimum_fidelity: u16::MAX,
3735            session: None,
3736        };
3737        (registry, graph)
3738    }
3739
3740    #[test]
3741    fn deterministic_across_insertion_order() {
3742        let (registry_a, graph_a) = fixture(false);
3743        let (registry_template, graph_b) = fixture(true);
3744        let mut registry_b = KernelRegistry::default();
3745        for descriptor in registry_template.descriptors.values().rev() {
3746            registry_b.register_descriptor(descriptor.clone()).unwrap();
3747        }
3748        for kernel in registry_template.kernels.values().rev() {
3749            registry_b.register_kernel(kernel.clone()).unwrap();
3750        }
3751        let a = Compiler::new(&registry_a, ExecutionRealm::HostStream)
3752            .compile(&graph_a)
3753            .unwrap();
3754        let b = Compiler::new(&registry_b, ExecutionRealm::HostStream)
3755            .compile(&graph_b)
3756            .unwrap();
3757        assert_eq!(a.graph_id, b.graph_id);
3758        assert_eq!(a.plan_id, b.plan_id);
3759        assert_eq!(a.order, vec![NodeId(0), NodeId(1), NodeId(2)]);
3760        assert_eq!(a.nodes.len(), 2, "source and process must fuse");
3761    }
3762
3763    #[test]
3764    fn physical_lowering_is_independent_of_layout_declaration_order() {
3765        let (mut registry_a, graph) = fixture(false);
3766        for descriptor in registry_a.descriptors.values_mut() {
3767            for port in descriptor.inputs.iter_mut().chain(&mut descriptor.outputs) {
3768                port.layouts.push(Layout::TimeMajor);
3769            }
3770        }
3771        for kernel in registry_a.kernels.values_mut() {
3772            kernel.input_layouts.push(Layout::TimeMajor);
3773            kernel.output_layouts.push(Layout::TimeMajor);
3774        }
3775        let mut registry_b = registry_a.clone();
3776        for kernel in registry_b.kernels.values_mut() {
3777            kernel.input_layouts.reverse();
3778            kernel.output_layouts.reverse();
3779        }
3780
3781        let a = Compiler::new(&registry_a, ExecutionRealm::HostStream)
3782            .compile(&graph)
3783            .unwrap();
3784        let b = Compiler::new(&registry_b, ExecutionRealm::HostStream)
3785            .compile(&graph)
3786            .unwrap();
3787        assert_eq!(a.plan_id, b.plan_id);
3788        assert!(
3789            a.buffers
3790                .iter()
3791                .all(|buffer| buffer.layout == Layout::Canonical)
3792        );
3793    }
3794
3795    #[test]
3796    fn semantic_graph_identity_covers_graph_contracts() {
3797        let (registry, graph) = fixture(false);
3798        let baseline = Compiler::new(&registry, ExecutionRealm::HostStream)
3799            .compile(&graph)
3800            .unwrap()
3801            .graph_id;
3802
3803        let mut variants = Vec::new();
3804        let mut policy = graph.clone();
3805        policy.policy.push("export-controlled".to_string());
3806        variants.push(policy);
3807        let mut proofs = graph.clone();
3808        proofs.required_proofs.push("calibrated".to_string());
3809        variants.push(proofs);
3810        let mut fidelity = graph.clone();
3811        fidelity.minimum_fidelity -= 1;
3812        variants.push(fidelity);
3813
3814        for variant in variants {
3815            let identity = Compiler::new(&registry, ExecutionRealm::HostStream)
3816                .compile(&variant)
3817                .unwrap()
3818                .graph_id;
3819            assert_ne!(identity, baseline);
3820        }
3821
3822        let mut unsupported = graph.clone();
3823        unsupported
3824            .required_capabilities
3825            .push(Capability("accelerator".to_string()));
3826        assert_eq!(
3827            Compiler::new(&registry, ExecutionRealm::HostStream).compile(&unsupported),
3828            Err(CompileError::CapabilityUnsupported(
3829                "accelerator".to_string()
3830            ))
3831        );
3832    }
3833
3834    #[test]
3835    fn semantic_graph_rejects_empty_proof_and_policy_names() {
3836        let (registry, mut graph) = fixture(false);
3837        graph.required_proofs.push(String::new());
3838        assert_eq!(
3839            Compiler::new(&registry, ExecutionRealm::HostStream).compile(&graph),
3840            Err(CompileError::InvalidGraphContract)
3841        );
3842
3843        graph.required_proofs.clear();
3844        graph.policy.push(String::new());
3845        assert_eq!(
3846            Compiler::new(&registry, ExecutionRealm::HostStream).compile(&graph),
3847            Err(CompileError::InvalidGraphContract)
3848        );
3849
3850        graph.policy.clear();
3851        graph.required_capabilities.push(Capability(String::new()));
3852        assert_eq!(
3853            Compiler::new(&registry, ExecutionRealm::HostStream).compile(&graph),
3854            Err(CompileError::InvalidGraphContract)
3855        );
3856    }
3857
3858    #[test]
3859    fn instance_configuration_reaches_the_selected_physical_step() {
3860        let (registry, mut graph) = fixture(false);
3861        graph.nodes[1].config.insert(
3862            "gain".to_string(),
3863            crate::ConfigValue::Text("2".to_string()),
3864        );
3865        let plan = Compiler::new(&registry, ExecutionRealm::HostStream)
3866            .compile(&graph)
3867            .unwrap();
3868        assert_eq!(plan.nodes[0].semantic_nodes, vec![NodeId(0), NodeId(1)]);
3869        assert_eq!(
3870            plan.nodes[0].semantic_configs[1]["gain"],
3871            crate::ConfigValue::Text("2".into())
3872        );
3873    }
3874
3875    #[test]
3876    fn semantic_set_order_and_duplicates_do_not_change_identity() {
3877        let (registry, mut graph) = fixture(false);
3878        graph
3879            .policy
3880            .extend(["alpha".to_string(), "beta".to_string()]);
3881        graph
3882            .required_proofs
3883            .extend(["proof-a".to_string(), "proof-b".to_string()]);
3884        graph
3885            .required_capabilities
3886            .push(Capability("sample".to_string()));
3887        let canonical = Compiler::new(&registry, ExecutionRealm::HostStream)
3888            .compile(&graph)
3889            .unwrap();
3890
3891        let mut reordered = graph.clone();
3892        reordered.policy.reverse();
3893        reordered.policy.push("alpha".to_string());
3894        reordered.required_proofs.reverse();
3895        reordered.required_proofs.push("proof-a".to_string());
3896        reordered.required_capabilities.reverse();
3897        reordered
3898            .required_capabilities
3899            .push(Capability("sample".to_string()));
3900        let duplicate = Compiler::new(&registry, ExecutionRealm::HostStream)
3901            .compile(&reordered)
3902            .unwrap();
3903        assert_eq!(canonical.graph_id, duplicate.graph_id);
3904        assert_eq!(canonical.plan_id, duplicate.plan_id);
3905    }
3906
3907    #[test]
3908    fn descriptor_and_implementation_contracts_are_identity_bound() {
3909        let (registry, graph) = fixture(false);
3910        let baseline = Compiler::new(&registry, ExecutionRealm::HostStream)
3911            .compile(&graph)
3912            .unwrap();
3913
3914        let mut descriptor_changed = registry.clone();
3915        descriptor_changed
3916            .descriptors
3917            .get_mut(&("source".to_string(), 1))
3918            .unwrap()
3919            .policy
3920            .adds
3921            .push("regulated".to_string());
3922        let semantic_change = Compiler::new(&descriptor_changed, ExecutionRealm::HostStream)
3923            .compile(&graph)
3924            .unwrap();
3925        assert_ne!(baseline.graph_id, semantic_change.graph_id);
3926
3927        let mut implementation_changed = registry.clone();
3928        implementation_changed
3929            .kernels
3930            .get_mut(&KernelId(100))
3931            .unwrap()
3932            .implementation_id = ImplementationId([222; 32]);
3933        let physical_change = Compiler::new(&implementation_changed, ExecutionRealm::HostStream)
3934            .compile(&graph)
3935            .unwrap();
3936        assert_eq!(baseline.graph_id, physical_change.graph_id);
3937        assert_ne!(baseline.plan_id, physical_change.plan_id);
3938    }
3939
3940    #[test]
3941    fn descriptor_lookup_returns_normalized_registered_contract_only() {
3942        let mut registry = KernelRegistry::default();
3943        let mut contract = descriptor("lookup", false);
3944        contract.policy.requires = vec!["z".into(), "a".into(), "z".into()];
3945        let node_type = NodeTypeRef {
3946            type_name: contract.type_name.clone(),
3947            version: contract.version,
3948        };
3949        registry.register_descriptor(contract).unwrap();
3950
3951        assert_eq!(
3952            registry.descriptor(&node_type).unwrap().policy.requires,
3953            ["a", "z"]
3954        );
3955        assert!(
3956            registry
3957                .descriptor(&NodeTypeRef {
3958                    type_name: "missing".into(),
3959                    version: 1,
3960                })
3961                .is_none()
3962        );
3963    }
3964
3965    #[test]
3966    fn weaker_kernel_determinism_cannot_satisfy_a_bit_exact_node() {
3967        let (mut registry, graph) = fixture(false);
3968        registry.kernels.get_mut(&KernelId(0)).unwrap().determinism =
3969            Determinism::NumericallyEquivalent;
3970        assert_eq!(
3971            Compiler::new(&registry, ExecutionRealm::HostStream).compile(&graph),
3972            Err(CompileError::KernelUnavailable(NodeId(0), Target::Host))
3973        );
3974    }
3975
3976    #[test]
3977    fn duplicate_registry_identities_do_not_replace_the_first_definition() {
3978        let (mut registry, graph) = fixture(false);
3979        let descriptor_key = ("source".to_string(), 1);
3980        let original_descriptor = registry.descriptors[&descriptor_key].clone();
3981        let mut replacement_descriptor = original_descriptor.clone();
3982        replacement_descriptor.targets.clear();
3983        assert_eq!(
3984            registry.register_descriptor(replacement_descriptor),
3985            Err(CompileError::DuplicateDescriptor("source".to_string(), 1))
3986        );
3987        assert_eq!(registry.descriptors[&descriptor_key], original_descriptor);
3988
3989        let original_kernel = registry.kernels[&KernelId(0)].clone();
3990        let mut replacement_kernel = original_kernel.clone();
3991        replacement_kernel.implements[0].type_name = "replacement".to_string();
3992        assert_eq!(
3993            registry.register_kernel(replacement_kernel),
3994            Err(CompileError::DuplicateKernel(KernelId(0)))
3995        );
3996        assert_eq!(registry.kernels[&KernelId(0)], original_kernel);
3997        Compiler::new(&registry, ExecutionRealm::HostStream)
3998            .compile(&graph)
3999            .unwrap();
4000    }
4001
4002    #[test]
4003    fn descriptor_registration_rejects_empty_contract_names() {
4004        let mut descriptor_contract = descriptor("empty-descriptor-contract", false);
4005        descriptor_contract.policy.adds.push(String::new());
4006        let mut registry = KernelRegistry::default();
4007        assert_eq!(
4008            registry.register_descriptor(descriptor_contract),
4009            Err(CompileError::InvalidDescriptor(
4010                "empty-descriptor-contract".to_string(),
4011                1
4012            ))
4013        );
4014
4015        let mut port_contract = descriptor("empty-port-contract", false);
4016        port_contract.outputs[0].proof.provides.push(String::new());
4017        assert_eq!(
4018            registry.register_descriptor(port_contract),
4019            Err(CompileError::InvalidDescriptor(
4020                "empty-port-contract".to_string(),
4021                1
4022            ))
4023        );
4024    }
4025
4026    #[test]
4027    fn descriptor_registration_rejects_empty_domain_tokens() {
4028        // The compiler treats a domain token as opaque, so "is it a
4029        // legal token" reduces to "is it non-empty". That is the ONLY structural
4030        // claim graph-core still makes about domain vocabulary, and it replaces
4031        // the pre-migration rejection of `AbirRootType::Unknown("")`. Neither the old
4032        // check nor this one had a test until now; a guard nothing exercises is
4033        // indistinguishable from a guard that cannot fire.
4034        let mut registry = KernelRegistry::default();
4035
4036        let mut empty_root = descriptor("empty-domain-root", false);
4037        empty_root.outputs[0].domain.root = DomainToken::default();
4038        assert_eq!(
4039            registry.register_descriptor(empty_root),
4040            Err(CompileError::InvalidDescriptor(
4041                "empty-domain-root".to_string(),
4042                1
4043            ))
4044        );
4045
4046        let mut empty_view = descriptor("empty-domain-view", false);
4047        empty_view.outputs[0].domain.view = DomainToken::new("");
4048        assert_eq!(
4049            registry.register_descriptor(empty_view),
4050            Err(CompileError::InvalidDescriptor(
4051                "empty-domain-view".to_string(),
4052                1
4053            ))
4054        );
4055
4056        // A non-empty token graph-core has never heard of is ACCEPTED -- that is
4057        // the whole point of moving the vocabulary out of the compiler.
4058        let mut foreign = descriptor("foreign-domain-vocabulary", false);
4059        foreign.outputs[0].domain = DomainType::new("dicom-series", "frame");
4060        assert!(registry.register_descriptor(foreign).is_ok());
4061    }
4062
4063    #[test]
4064    fn rejects_cycle_before_lowering() {
4065        let (registry, mut graph) = fixture(false);
4066        graph.edges.push(Edge {
4067            from: PortRef {
4068                node: NodeId(2),
4069                port: "out".to_string(),
4070            },
4071            to: PortRef {
4072                node: NodeId(1),
4073                port: "in".to_string(),
4074            },
4075        });
4076        let error = Compiler::new(&registry, ExecutionRealm::HostStream)
4077            .compile(&graph)
4078            .unwrap_err();
4079        assert!(matches!(
4080            error,
4081            CompileError::DuplicateInput(..) | CompileError::Cycle
4082        ));
4083    }
4084
4085    #[test]
4086    fn fusion_preserves_semantic_identity() {
4087        let (registry, graph) = fixture(false);
4088        let fused = Compiler::new(&registry, ExecutionRealm::HostStream)
4089            .compile(&graph)
4090            .unwrap();
4091        let plain = Compiler::new(&registry, ExecutionRealm::HostStream)
4092            .with_fusion(false)
4093            .compile(&graph)
4094            .unwrap();
4095        assert_eq!(fused.graph_id, plain.graph_id);
4096        assert_ne!(fused.plan_id, plain.plan_id);
4097        assert_eq!(fused.order, plain.order);
4098        assert_eq!(fused.nodes[0].kernel, KernelId(100));
4099
4100        let mut fused_kernels = SemanticKernels;
4101        let mut fused_transactions = NoTransactions;
4102        let roots = BTreeMap::new();
4103        let fused_result = crate::PlanExecutor::new(&mut fused_kernels, &mut fused_transactions)
4104            .execute(&fused, [1; 32], roots.clone())
4105            .unwrap();
4106        let mut plain_kernels = SemanticKernels;
4107        let mut plain_transactions = NoTransactions;
4108        let plain_result = crate::PlanExecutor::new(&mut plain_kernels, &mut plain_transactions)
4109            .execute(&plain, [1; 32], roots)
4110            .unwrap();
4111        assert_eq!(fused_result.terminal_values, plain_result.terminal_values);
4112        assert_eq!(
4113            fused_result.receipt.completed_nodes,
4114            plain_result.receipt.completed_nodes
4115        );
4116    }
4117
4118    #[test]
4119    fn arbitrary_length_registered_fusion_is_selected_globally() {
4120        let (mut registry, graph) = fixture(false);
4121        registry
4122            .register_kernel(KernelDescriptor {
4123                id: KernelId(150),
4124                implements: ["source", "process", "sink"]
4125                    .into_iter()
4126                    .map(|type_name| NodeTypeRef {
4127                        type_name: type_name.to_string(),
4128                        version: 1,
4129                    })
4130                    .collect(),
4131                implementation_id: ImplementationId([150; 32]),
4132                conversion: None,
4133                target: Target::Host,
4134                input_layouts: vec![Layout::Canonical],
4135                output_layouts: vec![Layout::Canonical],
4136                resources: ResourceEnvelope::bounded(1, 0, 1),
4137                determinism: Determinism::BitExact,
4138                lowering: "test.fused.all".to_string(),
4139            })
4140            .unwrap();
4141        let plan = Compiler::new(&registry, ExecutionRealm::HostStream)
4142            .compile(&graph)
4143            .unwrap();
4144        assert_eq!(plan.nodes.len(), 1);
4145        assert_eq!(
4146            plan.nodes[0].semantic_nodes,
4147            vec![NodeId(0), NodeId(1), NodeId(2)]
4148        );
4149        assert_eq!(plan.nodes[0].kernel, KernelId(150));
4150    }
4151
4152    #[test]
4153    fn infeasible_cheapest_fused_kernel_does_not_mask_feasible_alternative() {
4154        let (mut registry, _graph) = fixture(false);
4155        registry
4156            .descriptors
4157            .get_mut(&("process".to_string(), 1))
4158            .unwrap()
4159            .outputs[0]
4160            .layouts = vec![Layout::Canonical, Layout::TimeMajor];
4161        registry
4162            .kernels
4163            .get_mut(&KernelId(100))
4164            .unwrap()
4165            .output_layouts = vec![Layout::TimeMajor];
4166        registry.kernels.get_mut(&KernelId(100)).unwrap().resources =
4167            ResourceEnvelope::bounded(1, 0, 1);
4168        registry
4169            .register_kernel(KernelDescriptor {
4170                id: KernelId(110),
4171                implements: vec![
4172                    NodeTypeRef {
4173                        type_name: "source".to_string(),
4174                        version: 1,
4175                    },
4176                    NodeTypeRef {
4177                        type_name: "process".to_string(),
4178                        version: 1,
4179                    },
4180                ],
4181                implementation_id: ImplementationId([110; 32]),
4182                conversion: None,
4183                target: Target::Host,
4184                input_layouts: vec![Layout::Canonical],
4185                output_layouts: vec![Layout::Canonical],
4186                resources: ResourceEnvelope::bounded(2, 0, 1),
4187                determinism: Determinism::BitExact,
4188                lowering: "test.fused.feasible".to_string(),
4189            })
4190            .unwrap();
4191        let graph = fixture(false).1;
4192        let plan = Compiler::new(&registry, ExecutionRealm::HostStream)
4193            .compile(&graph)
4194            .unwrap();
4195        assert_eq!(plan.nodes[0].kernel, KernelId(110));
4196    }
4197
4198    #[test]
4199    fn invocation_values_bind_only_declared_typed_ports() {
4200        let (mut registry, mut graph) = fixture(false);
4201        registry
4202            .descriptors
4203            .get_mut(&("source".to_string(), 1))
4204            .unwrap()
4205            .inputs = vec![PortDescriptor {
4206            name: "seed".to_string(),
4207            semantic_type: "sample.block".to_string(),
4208            optional: false,
4209            layouts: vec![Layout::Canonical],
4210            max_bytes: 64,
4211            ..PortDescriptor::opaque("seed", "sample.block", 64)
4212        }];
4213        let seed = PortRef {
4214            node: NodeId(0),
4215            port: "seed".to_string(),
4216        };
4217        graph.invocation_inputs = vec![seed.clone()];
4218        let plan = Compiler::new(&registry, ExecutionRealm::HostStream)
4219            .compile(&graph)
4220            .unwrap();
4221        assert_eq!(plan.invocation_ports, vec![seed.clone()]);
4222        assert_eq!(plan.nodes[0].input_ports, vec!["seed"]);
4223        assert_eq!(
4224            plan.nodes[0].input_bindings,
4225            vec![crate::InputBinding::Invocation(0)]
4226        );
4227
4228        let mut kernels = SemanticKernels;
4229        let mut transactions = NoTransactions;
4230        let mut values = BTreeMap::new();
4231        values.insert(seed, 5);
4232        let result = crate::PlanExecutor::new(&mut kernels, &mut transactions)
4233            .execute(&plan, [4; 32], values)
4234            .unwrap();
4235        assert_eq!(result.terminal_values[&NodeId(2)], vec![27]);
4236    }
4237
4238    #[test]
4239    fn infeasible_fused_implementation_falls_back_to_unfused_steps() {
4240        let (mut registry, graph) = fixture(false);
4241        registry.kernels.get_mut(&KernelId(100)).unwrap().resources =
4242            ResourceEnvelope::bounded(1024, 0, 1);
4243        let plan = Compiler::new(&registry, ExecutionRealm::HostStream)
4244            .with_memory_limit(192)
4245            .compile(&graph)
4246            .unwrap();
4247        assert_eq!(plan.nodes.len(), 3);
4248        assert!(plan.nodes.iter().all(|node| node.kernel != KernelId(100)));
4249        assert_eq!(plan.peak_bytes, 192);
4250    }
4251
4252    #[test]
4253    fn hard_step_limit_is_applied_during_global_plan_search() {
4254        let (mut registry, graph) = fixture(false);
4255        registry.kernels.get_mut(&KernelId(100)).unwrap().resources =
4256            ResourceEnvelope::bounded(1024, 0, 1);
4257        let plan = Compiler::new(&registry, ExecutionRealm::HostStream)
4258            .with_limits(CompileLimits {
4259                max_steps: 2,
4260                ..CompileLimits::default()
4261            })
4262            .compile(&graph)
4263            .unwrap();
4264        assert_eq!(plan.nodes.len(), 2);
4265        assert_eq!(plan.nodes[0].kernel, KernelId(100));
4266    }
4267
4268    #[test]
4269    fn fusion_never_erases_an_unconnected_observable_output() {
4270        let (mut registry, graph) = fixture(false);
4271        registry
4272            .descriptors
4273            .get_mut(&("source".to_string(), 1))
4274            .unwrap()
4275            .outputs
4276            .push(PortDescriptor {
4277                name: "audit".to_string(),
4278                semantic_type: "sample.block".to_string(),
4279                optional: false,
4280                layouts: vec![Layout::Canonical],
4281                max_bytes: 64,
4282                ..PortDescriptor::opaque("audit", "sample.block", 64)
4283            });
4284        let plan = Compiler::new(&registry, ExecutionRealm::HostStream)
4285            .compile(&graph)
4286            .unwrap();
4287        assert_eq!(plan.nodes.len(), 3);
4288        assert_eq!(
4289            plan.nodes[0].output_bindings,
4290            vec![OutputBinding::Buffer(BufferId(0)), OutputBinding::Terminal]
4291        );
4292    }
4293
4294    #[test]
4295    fn fanout_shares_one_sized_buffer_and_tracks_every_consumer() {
4296        let (registry, mut graph) = fixture(false);
4297        graph.edges.push(Edge {
4298            from: PortRef {
4299                node: NodeId(0),
4300                port: "out".to_string(),
4301            },
4302            to: PortRef {
4303                node: NodeId(2),
4304                port: "in".to_string(),
4305            },
4306        });
4307        graph
4308            .edges
4309            .retain(|edge| !(edge.from.node == NodeId(1) && edge.to.node == NodeId(2)));
4310        let plan = Compiler::new(&registry, ExecutionRealm::HostStream)
4311            .compile(&graph)
4312            .unwrap();
4313        assert_eq!(plan.buffers.len(), 1);
4314        assert_eq!(plan.buffers[0].capacity_bytes, 64);
4315        assert_eq!(plan.buffers[0].consumers, vec![StepId(1), StepId(2)]);
4316        assert_eq!(plan.peak_bytes, 128);
4317    }
4318
4319    #[test]
4320    fn connected_zero_sized_output_is_rejected() {
4321        let (mut registry, graph) = fixture(false);
4322        registry
4323            .descriptors
4324            .get_mut(&("source".to_string(), 1))
4325            .unwrap()
4326            .outputs[0]
4327            .max_bytes = 0;
4328        assert!(matches!(
4329            Compiler::new(&registry, ExecutionRealm::HostStream).compile(&graph),
4330            Err(CompileError::InvalidPortSize(NodeId(0), _))
4331        ));
4332    }
4333
4334    #[test]
4335    fn all_realms_compile_the_same_semantic_graph() {
4336        let (registry, graph) = fixture(false);
4337        let host = Compiler::new(&registry, ExecutionRealm::HostStream)
4338            .compile(&graph)
4339            .unwrap();
4340        let mcu = Compiler::new(&registry, ExecutionRealm::McuAot)
4341            .compile(&graph)
4342            .unwrap();
4343        let durable = Compiler::new(&registry, ExecutionRealm::BlutDurable)
4344            .compile(&graph)
4345            .unwrap();
4346        assert_eq!(host.graph_id, mcu.graph_id);
4347        assert_eq!(host.graph_id, durable.graph_id);
4348        assert_eq!(host.order, mcu.order);
4349        assert_eq!(host.order, durable.order);
4350        assert_ne!(host.plan_id, mcu.plan_id);
4351        assert_ne!(host.plan_id, durable.plan_id);
4352        for plan in [&host, &mcu, &durable] {
4353            assert_eq!(
4354                crate::CompiledPlan::from_aot_bytes(
4355                    &plan.to_aot_bytes().unwrap(),
4356                    crate::PlanLimits::default()
4357                )
4358                .unwrap(),
4359                *plan.as_plan()
4360            );
4361        }
4362    }
4363
4364    #[test]
4365    fn kernel_resources_are_included_in_peak_memory_admission() {
4366        let (registry, graph) = fixture(false);
4367        let plan = Compiler::new(&registry, ExecutionRealm::HostStream)
4368            .with_memory_limit(128)
4369            .compile(&graph)
4370            .unwrap();
4371        assert_eq!(plan.peak_bytes, 128);
4372        assert_eq!(
4373            Compiler::new(&registry, ExecutionRealm::HostStream)
4374                .with_memory_limit(127)
4375                .compile(&graph),
4376            Err(CompileError::ResourceOverflow)
4377        );
4378    }
4379
4380    fn fanout_graph(graph: &mut Graph) {
4381        graph.edges.retain(|edge| edge.from.node == NodeId(0));
4382        graph.edges.push(Edge {
4383            from: PortRef {
4384                node: NodeId(0),
4385                port: "out".to_string(),
4386            },
4387            to: PortRef {
4388                node: NodeId(2),
4389                port: "in".to_string(),
4390            },
4391        });
4392    }
4393
4394    #[test]
4395    fn proofs_do_not_leak_between_parallel_branches() {
4396        let (mut registry, mut graph) = fixture(false);
4397        fanout_graph(&mut graph);
4398        registry
4399            .descriptors
4400            .get_mut(&("process".to_string(), 1))
4401            .unwrap()
4402            .proof
4403            .provides
4404            .push("branch-only".to_string());
4405        registry
4406            .descriptors
4407            .get_mut(&("sink".to_string(), 1))
4408            .unwrap()
4409            .proof
4410            .requires
4411            .push("branch-only".to_string());
4412        assert_eq!(
4413            Compiler::new(&registry, ExecutionRealm::HostStream).compile(&graph),
4414            Err(CompileError::ProofMissing(
4415                NodeId(2),
4416                "branch-only".to_string()
4417            ))
4418        );
4419    }
4420
4421    #[test]
4422    fn invalidated_proof_cannot_satisfy_a_downstream_requirement() {
4423        let (mut registry, graph) = fixture(false);
4424        registry
4425            .descriptors
4426            .get_mut(&("source".to_string(), 1))
4427            .unwrap()
4428            .proof
4429            .provides
4430            .push("calibrated".to_string());
4431        registry
4432            .descriptors
4433            .get_mut(&("process".to_string(), 1))
4434            .unwrap()
4435            .proof
4436            .invalidates
4437            .push("calibrated".to_string());
4438        registry
4439            .descriptors
4440            .get_mut(&("sink".to_string(), 1))
4441            .unwrap()
4442            .proof
4443            .requires
4444            .push("calibrated".to_string());
4445        assert_eq!(
4446            Compiler::new(&registry, ExecutionRealm::HostStream).compile(&graph),
4447            Err(CompileError::ProofMissing(
4448                NodeId(2),
4449                "calibrated".to_string()
4450            ))
4451        );
4452    }
4453
4454    #[test]
4455    fn parallel_fidelity_uses_the_worst_branch_not_the_sum() {
4456        let (mut registry, mut graph) = fixture(false);
4457        fanout_graph(&mut graph);
4458        for name in ["process", "sink"] {
4459            registry
4460                .descriptors
4461                .get_mut(&(name.to_string(), 1))
4462                .unwrap()
4463                .fidelity
4464                .maximum_loss = 100;
4465        }
4466        graph.minimum_fidelity = u16::MAX - 100;
4467        let plan = Compiler::new(&registry, ExecutionRealm::HostStream)
4468            .compile(&graph)
4469            .unwrap();
4470        assert_eq!(plan.resulting_fidelity, u16::MAX - 100);
4471    }
4472
4473    #[test]
4474    fn physical_search_chooses_a_feasible_non_greedy_kernel_assignment() {
4475        let (mut registry, graph) = fixture(false);
4476        let source_descriptor = registry
4477            .descriptors
4478            .get_mut(&("source".to_string(), 1))
4479            .unwrap();
4480        source_descriptor.outputs[0].layouts = vec![Layout::Canonical, Layout::TimeMajor];
4481        let process_descriptor = registry
4482            .descriptors
4483            .get_mut(&("process".to_string(), 1))
4484            .unwrap();
4485        process_descriptor.inputs[0].layouts = vec![Layout::Canonical, Layout::TimeMajor];
4486        registry
4487            .kernels
4488            .get_mut(&KernelId(0))
4489            .unwrap()
4490            .output_layouts = vec![Layout::Canonical];
4491        registry
4492            .kernels
4493            .get_mut(&KernelId(3))
4494            .unwrap()
4495            .input_layouts = vec![Layout::TimeMajor];
4496        registry
4497            .register_kernel(KernelDescriptor {
4498                id: KernelId(200),
4499                implements: vec![NodeTypeRef {
4500                    type_name: "source".to_string(),
4501                    version: 1,
4502                }],
4503                implementation_id: ImplementationId([200; 32]),
4504                conversion: None,
4505                target: Target::Host,
4506                input_layouts: vec![Layout::Canonical],
4507                output_layouts: vec![Layout::TimeMajor],
4508                resources: ResourceEnvelope::bounded(65, 0, 1),
4509                determinism: Determinism::BitExact,
4510                lowering: "feasible-source".to_string(),
4511            })
4512            .unwrap();
4513
4514        let plan = Compiler::new(&registry, ExecutionRealm::HostStream)
4515            .with_fusion(false)
4516            .compile(&graph)
4517            .unwrap();
4518        assert_eq!(plan.nodes[0].kernel, KernelId(200));
4519        assert_eq!(plan.buffers[0].layout, Layout::TimeMajor);
4520    }
4521
4522    #[test]
4523    fn physical_search_limit_fails_closed() {
4524        let (registry, graph) = fixture(false);
4525        assert_eq!(
4526            Compiler::new(&registry, ExecutionRealm::HostStream)
4527                .with_limits(CompileLimits {
4528                    max_search_states: 1,
4529                    ..CompileLimits::default()
4530                })
4531                .compile(&graph),
4532            Err(CompileError::SearchLimitExceeded)
4533        );
4534    }
4535
4536    #[test]
4537    fn feedback_identity_conversion_is_checked() {
4538        assert_eq!(
4539            feedback_id(u32::MAX as usize),
4540            Ok(crate::FeedbackId(u32::MAX))
4541        );
4542        #[cfg(target_pointer_width = "64")]
4543        assert_eq!(
4544            feedback_id(u32::MAX as usize + 1),
4545            Err(CompileError::CompileLimitExceeded)
4546        );
4547    }
4548
4549    fn conversion_kernel(id: u32, from: Layout, to: Layout) -> KernelDescriptor {
4550        KernelDescriptor {
4551            id: KernelId(id),
4552            implements: vec![],
4553            implementation_id: ImplementationId([id as u8; 32]),
4554            conversion: Some(crate::LayoutConversion {
4555                semantic_type: "sample.block".to_string(),
4556                from,
4557                to,
4558                max_input_bytes: 64,
4559                max_output_bytes: 64,
4560            }),
4561            target: Target::Host,
4562            input_layouts: vec![from],
4563            output_layouts: vec![to],
4564            resources: ResourceEnvelope::bounded(8, 0, 1),
4565            determinism: Determinism::BitExact,
4566            lowering: format!("convert-{from:?}-{to:?}"),
4567        }
4568    }
4569
4570    #[test]
4571    fn layout_conversion_is_explicit_and_preserves_semantic_receipts() {
4572        struct ConversionExecutor;
4573        impl crate::KernelExecutor for ConversionExecutor {
4574            type Value = u32;
4575
4576            fn execute(
4577                &mut self,
4578                node: &CompiledNode,
4579                inputs: &[Option<&Self::Value>],
4580            ) -> Result<Vec<Self::Value>, crate::ExecutionError> {
4581                let input = inputs
4582                    .iter()
4583                    .flatten()
4584                    .next()
4585                    .map(|value| **value)
4586                    .unwrap_or_default();
4587                let value = if node.conversion.is_some() {
4588                    input
4589                } else {
4590                    input + 1
4591                };
4592                Ok(vec![value; node.output_bindings.len()])
4593            }
4594        }
4595
4596        let (mut registry, mut graph) = fixture(false);
4597        graph
4598            .edges
4599            .retain(|edge| !(edge.from.node == NodeId(1) && edge.to.node == NodeId(2)));
4600        graph.edges.push(Edge {
4601            from: PortRef {
4602                node: NodeId(0),
4603                port: "out".to_string(),
4604            },
4605            to: PortRef {
4606                node: NodeId(2),
4607                port: "in".to_string(),
4608            },
4609        });
4610        registry
4611            .descriptors
4612            .get_mut(&("source".to_string(), 1))
4613            .unwrap()
4614            .outputs[0]
4615            .layouts = vec![Layout::ChannelMajor];
4616        registry
4617            .descriptors
4618            .get_mut(&("process".to_string(), 1))
4619            .unwrap()
4620            .inputs[0]
4621            .layouts = vec![Layout::TimeMajor];
4622        registry
4623            .descriptors
4624            .get_mut(&("sink".to_string(), 1))
4625            .unwrap()
4626            .inputs[0]
4627            .layouts = vec![Layout::ChannelMajor];
4628        registry
4629            .kernels
4630            .get_mut(&KernelId(0))
4631            .unwrap()
4632            .output_layouts = vec![Layout::ChannelMajor];
4633        registry
4634            .kernels
4635            .get_mut(&KernelId(3))
4636            .unwrap()
4637            .input_layouts = vec![Layout::TimeMajor];
4638        registry
4639            .kernels
4640            .get_mut(&KernelId(6))
4641            .unwrap()
4642            .input_layouts = vec![Layout::ChannelMajor];
4643
4644        assert!(matches!(
4645            Compiler::new(&registry, ExecutionRealm::HostStream)
4646                .with_fusion(false)
4647                .compile(&graph),
4648            Err(CompileError::LayoutUnavailable(..))
4649        ));
4650        let mut expanding = conversion_kernel(210, Layout::ChannelMajor, Layout::TimeMajor);
4651        expanding.conversion.as_mut().unwrap().max_output_bytes = 128;
4652        registry.register_kernel(expanding).unwrap();
4653        assert!(matches!(
4654            Compiler::new(&registry, ExecutionRealm::HostStream)
4655                .with_fusion(false)
4656                .compile(&graph),
4657            Err(CompileError::LayoutUnavailable(..))
4658        ));
4659        registry
4660            .descriptors
4661            .get_mut(&("process".to_string(), 1))
4662            .unwrap()
4663            .inputs[0]
4664            .max_bytes = 128;
4665        let plan = Compiler::new(&registry, ExecutionRealm::HostStream)
4666            .with_fusion(false)
4667            .compile(&graph)
4668            .unwrap();
4669        assert_eq!(plan.nodes.len(), 4);
4670        assert_eq!(plan.nodes[1].semantic_nodes, Vec::<NodeId>::new());
4671        assert!(plan.nodes[1].conversion.is_some());
4672        let conversion_output = plan.nodes[1].output_bindings[0];
4673        let OutputBinding::Buffer(conversion_output) = conversion_output else {
4674            panic!("conversion output must be buffered");
4675        };
4676        assert_eq!(
4677            plan.buffers[conversion_output.0 as usize].capacity_bytes,
4678            128
4679        );
4680
4681        let roots = BTreeMap::new();
4682        let mut executor = ConversionExecutor;
4683        let mut transactions = NoTransactions;
4684        let result = crate::PlanExecutor::new(&mut executor, &mut transactions)
4685            .execute(&plan, [7; 32], roots)
4686            .unwrap();
4687        assert_eq!(result.terminal_values[&NodeId(1)], vec![2]);
4688        assert_eq!(result.terminal_values[&NodeId(2)], vec![2]);
4689        assert_eq!(result.receipt.completed_nodes, plan.order);
4690        assert_eq!(result.receipt.attempts.len(), 4);
4691        assert!(result.receipt.attempts[1].semantic_nodes.is_empty());
4692    }
4693
4694    #[test]
4695    fn alias_slot_lifetime_advances_after_every_reuse() {
4696        let mut buffers = vec![
4697            BufferPlan {
4698                id: BufferId(0),
4699                layout: Layout::Canonical,
4700                capacity_bytes: 64,
4701                producer: StepId(0),
4702                consumers: vec![StepId(1)],
4703                last_consumer: StepId(1),
4704                aliases: None,
4705            },
4706            BufferPlan {
4707                id: BufferId(1),
4708                layout: Layout::Canonical,
4709                capacity_bytes: 64,
4710                producer: StepId(2),
4711                consumers: vec![StepId(4)],
4712                last_consumer: StepId(4),
4713                aliases: None,
4714            },
4715            BufferPlan {
4716                id: BufferId(2),
4717                layout: Layout::Canonical,
4718                capacity_bytes: 64,
4719                producer: StepId(3),
4720                consumers: vec![StepId(5)],
4721                last_consumer: StepId(5),
4722                aliases: None,
4723            },
4724        ];
4725        assign_aliases(&mut buffers);
4726        assert_eq!(buffers[1].aliases, Some(BufferId(0)));
4727        assert_eq!(
4728            buffers[2].aliases, None,
4729            "overlapping lifetime needs a new slot"
4730        );
4731    }
4732
4733    #[test]
4734    fn registry_authorization_rejects_self_consistent_forged_resources() {
4735        let (registry, graph) = fixture(false);
4736        let plan = Compiler::new(&registry, ExecutionRealm::HostStream)
4737            .compile(&graph)
4738            .unwrap();
4739        let authorization = crate::PlanAuthorization {
4740            expected_realm: plan.realm,
4741            expected_plan_id: plan.plan_id,
4742        };
4743        assert_eq!(
4744            registry
4745                .decode_authorized_plan(
4746                    &plan.to_aot_bytes().unwrap(),
4747                    crate::PlanLimits::default(),
4748                    authorization,
4749                )
4750                .unwrap(),
4751            plan
4752        );
4753
4754        let mut forged = plan.clone().into_plan();
4755        forged.nodes[0].resources.peak_bytes += 1;
4756        forged.peak_bytes += 1;
4757        forged.plan_id = PlanId(hash_plan(&forged));
4758        let bytes = forged.to_aot_bytes().unwrap();
4759        let forged_authorization = crate::PlanAuthorization {
4760            expected_realm: forged.realm,
4761            expected_plan_id: forged.plan_id,
4762        };
4763        assert!(
4764            CompiledPlan::from_authorized_aot_bytes(
4765                &bytes,
4766                crate::PlanLimits::default(),
4767                forged_authorization,
4768            )
4769            .is_ok()
4770        );
4771        assert_eq!(
4772            registry.decode_authorized_plan(
4773                &bytes,
4774                crate::PlanLimits::default(),
4775                forged_authorization,
4776            ),
4777            Err(crate::PlanDecodeError::UnauthorizedPlan)
4778        );
4779    }
4780
4781    #[test]
4782    fn typed_configuration_defaults_are_identity_canonical_and_invalid_values_fail() {
4783        let (mut registry, graph) = fixture(false);
4784        registry
4785            .descriptors
4786            .get_mut(&("process".into(), 1))
4787            .unwrap()
4788            .config
4789            .fields[0]
4790            .default = Some(crate::ConfigValue::Text("1".into()));
4791
4792        let implicit = Compiler::new(&registry, ExecutionRealm::HostStream)
4793            .compile(&graph)
4794            .unwrap();
4795        let mut explicit_graph = graph.clone();
4796        explicit_graph.nodes[1]
4797            .config
4798            .insert("gain".into(), crate::ConfigValue::Text("1".into()));
4799        let explicit = Compiler::new(&registry, ExecutionRealm::HostStream)
4800            .compile(&explicit_graph)
4801            .unwrap();
4802        assert_eq!(implicit.graph_id, explicit.graph_id);
4803        assert_eq!(implicit.plan_id, explicit.plan_id);
4804
4805        explicit_graph.nodes[1]
4806            .config
4807            .insert("gain".into(), crate::ConfigValue::Text("x".repeat(17)));
4808        assert!(matches!(
4809            Compiler::new(&registry, ExecutionRealm::HostStream).compile(&explicit_graph),
4810            Err(CompileError::InvalidConfig(NodeId(1), _))
4811        ));
4812    }
4813
4814    #[test]
4815    fn per_port_domain_contract_mismatch_fails_before_kernel_selection() {
4816        let (mut registry, graph) = fixture(false);
4817        registry
4818            .descriptors
4819            .get_mut(&("process".into(), 1))
4820            .unwrap()
4821            .inputs[0]
4822            .domain
4823            .root = crate::DomainToken::new("recording");
4824        assert!(matches!(
4825            Compiler::new(&registry, ExecutionRealm::HostStream).compile(&graph),
4826            Err(CompileError::PortContractMismatch(
4827                NodeId(0),
4828                _,
4829                NodeId(1),
4830                _
4831            ))
4832        ));
4833    }
4834
4835    #[test]
4836    fn session_feedback_has_explicit_binding_and_exact_bounded_state() {
4837        let (mut registry, mut graph) = fixture(false);
4838        let source = registry.descriptors.get_mut(&("source".into(), 1)).unwrap();
4839        let mut history = PortDescriptor::opaque("history", "sample.block", 64);
4840        history.optional = true;
4841        source.inputs.push(history);
4842        registry
4843            .descriptors
4844            .get_mut(&("process".into(), 1))
4845            .unwrap()
4846            .state = StateContract {
4847            scope: StateScope::Session,
4848            max_bytes: 128,
4849            checkpoint: crate::CheckpointContract {
4850                mode: crate::CheckpointMode::Required,
4851                max_snapshot_bytes: 64,
4852                max_interval_invocations: 8,
4853            },
4854        };
4855        graph.feedback.push(crate::FeedbackEdge {
4856            from: PortRef {
4857                node: NodeId(1),
4858                port: "out".into(),
4859            },
4860            to: PortRef {
4861                node: NodeId(0),
4862                port: "history".into(),
4863            },
4864            delay: crate::DelayContract {
4865                invocations: 1,
4866                initial: crate::DelayInitial::Absent,
4867            },
4868        });
4869        graph.session = Some(crate::SessionContract {
4870            namespace: "patient-session".into(),
4871            max_concurrent_sessions: 16,
4872            max_idle_millis: 60_000,
4873            reset_on_plan_change: true,
4874        });
4875
4876        let plan = Compiler::new(&registry, ExecutionRealm::HostStream)
4877            .compile(&graph)
4878            .unwrap();
4879        assert_eq!(plan.feedback.len(), 1);
4880        assert_eq!(plan.persistent_state_bytes, 192);
4881        assert!(plan.nodes.iter().any(|step| {
4882            step.input_bindings
4883                .contains(&crate::InputBinding::Feedback(crate::FeedbackId(0)))
4884        }));
4885        let bytes = plan.to_aot_bytes().unwrap();
4886        assert!(CompiledPlan::from_aot_bytes(&bytes, crate::PlanLimits::default()).is_ok());
4887
4888        graph.feedback[0].delay.invocations = 3;
4889        let longer_delay = Compiler::new(&registry, ExecutionRealm::HostStream)
4890            .compile(&graph)
4891            .unwrap();
4892        assert_eq!(longer_delay.feedback[0].state_bytes, 192);
4893        assert_eq!(longer_delay.persistent_state_bytes, 320);
4894
4895        let mut overlapping_layouts = registry.clone();
4896        let mut overlapping_graph = graph.clone();
4897        overlapping_graph.feedback[0].from.node = NodeId(2);
4898        overlapping_graph.feedback.push(crate::FeedbackEdge {
4899            from: PortRef {
4900                node: NodeId(2),
4901                port: "out".into(),
4902            },
4903            to: PortRef {
4904                node: NodeId(1),
4905                port: "history-2".into(),
4906            },
4907            delay: crate::DelayContract {
4908                invocations: 3,
4909                initial: crate::DelayInitial::Absent,
4910            },
4911        });
4912        overlapping_layouts
4913            .descriptors
4914            .get_mut(&("source".into(), 1))
4915            .unwrap()
4916            .inputs[0]
4917            .layouts = vec![Layout::ChannelMajor, Layout::TimeMajor];
4918        overlapping_layouts
4919            .descriptors
4920            .get_mut(&("sink".into(), 1))
4921            .unwrap()
4922            .outputs[0]
4923            .layouts = vec![Layout::Canonical, Layout::TimeMajor];
4924        let mut history_2 = PortDescriptor::opaque("history-2", "sample.block", 64);
4925        history_2.optional = true;
4926        history_2.layouts = vec![Layout::ChannelMajor, Layout::TimeMajor];
4927        overlapping_layouts
4928            .descriptors
4929            .get_mut(&("process".into(), 1))
4930            .unwrap()
4931            .inputs
4932            .push(history_2);
4933        for kernel in overlapping_layouts.kernels.values_mut() {
4934            if kernel.implements.as_slice()
4935                == [NodeTypeRef {
4936                    type_name: "source".into(),
4937                    version: 1,
4938                }]
4939            {
4940                kernel.input_layouts = vec![Layout::ChannelMajor, Layout::TimeMajor];
4941            }
4942            if kernel.implements.as_slice()
4943                == [NodeTypeRef {
4944                    type_name: "sink".into(),
4945                    version: 1,
4946                }]
4947            {
4948                kernel.output_layouts = vec![Layout::Canonical, Layout::TimeMajor];
4949            }
4950            if kernel.implements.as_slice()
4951                == [NodeTypeRef {
4952                    type_name: "process".into(),
4953                    version: 1,
4954                }]
4955            {
4956                kernel.input_layouts =
4957                    vec![Layout::Canonical, Layout::ChannelMajor, Layout::TimeMajor];
4958            }
4959        }
4960        let overlapping = Compiler::new(&overlapping_layouts, ExecutionRealm::HostStream)
4961            .compile(&overlapping_graph)
4962            .unwrap();
4963        for feedback in &overlapping.feedback {
4964            assert_eq!(
4965                overlapping.nodes[feedback.from_step.0 as usize].output_contracts
4966                    [feedback.from_port as usize]
4967                    .layout,
4968                Layout::TimeMajor
4969            );
4970            assert_eq!(
4971                overlapping.nodes[feedback.to_step.0 as usize].input_contracts
4972                    [feedback.to_port as usize]
4973                    .layout,
4974                Layout::TimeMajor
4975            );
4976        }
4977
4978        let mut incompatible_layouts = registry.clone();
4979        incompatible_layouts
4980            .descriptors
4981            .get_mut(&("source".into(), 1))
4982            .unwrap()
4983            .inputs[0]
4984            .layouts = vec![Layout::TimeMajor];
4985        incompatible_layouts
4986            .descriptors
4987            .get_mut(&("process".into(), 1))
4988            .unwrap()
4989            .outputs[0]
4990            .layouts = vec![Layout::ChannelMajor];
4991        incompatible_layouts
4992            .descriptors
4993            .get_mut(&("sink".into(), 1))
4994            .unwrap()
4995            .inputs[0]
4996            .layouts = vec![Layout::ChannelMajor];
4997        for kernel in incompatible_layouts.kernels.values_mut() {
4998            if kernel.implements.as_slice()
4999                == [NodeTypeRef {
5000                    type_name: "source".into(),
5001                    version: 1,
5002                }]
5003            {
5004                kernel.input_layouts = vec![Layout::TimeMajor];
5005            }
5006            if kernel.implements.as_slice()
5007                == [NodeTypeRef {
5008                    type_name: "process".into(),
5009                    version: 1,
5010                }]
5011            {
5012                kernel.output_layouts = vec![Layout::ChannelMajor];
5013            }
5014            if kernel.implements.as_slice()
5015                == [NodeTypeRef {
5016                    type_name: "sink".into(),
5017                    version: 1,
5018                }]
5019            {
5020                kernel.input_layouts = vec![Layout::ChannelMajor];
5021            }
5022        }
5023        assert_eq!(
5024            Compiler::new(&incompatible_layouts, ExecutionRealm::HostStream)
5025                .compile(&graph)
5026                .unwrap_err(),
5027            CompileError::InvalidFeedback(NodeId(0), "history".into())
5028        );
5029
5030        registry
5031            .descriptors
5032            .get_mut(&("source".into(), 1))
5033            .unwrap()
5034            .inputs[0]
5035            .optional = false;
5036        assert!(matches!(
5037            Compiler::new(&registry, ExecutionRealm::HostStream).compile(&graph),
5038            Err(CompileError::InvalidFeedback(NodeId(0), _))
5039        ));
5040    }
5041
5042    #[test]
5043    fn hierarchical_lowering_identity_and_depth_are_bounded() {
5044        let (mut registry, graph) = fixture(false);
5045        let mut leaf = crate::SubgraphSchema {
5046            id: crate::SubgraphId([0; 32]),
5047            version: 1,
5048            nodes: vec![crate::SubgraphNode {
5049                id: NodeId(0),
5050                node_type: NodeTypeRef {
5051                    type_name: "process".into(),
5052                    version: 1,
5053                },
5054                config: BTreeMap::from([("gain".into(), crate::ConfigValue::Text("1".into()))]),
5055                child: None,
5056            }],
5057            edges: vec![],
5058            inputs: vec![crate::SubgraphInterfacePort {
5059                name: "in".into(),
5060                inner: PortRef {
5061                    node: NodeId(0),
5062                    port: "in".into(),
5063                },
5064            }],
5065            outputs: vec![crate::SubgraphInterfacePort {
5066                name: "out".into(),
5067                inner: PortRef {
5068                    node: NodeId(0),
5069                    port: "out".into(),
5070                },
5071            }],
5072        };
5073        leaf.id = subgraph_identity(&leaf);
5074        registry.register_subgraph(leaf.clone()).unwrap();
5075        let mut repeated = leaf.clone();
5076        repeated.nodes.push(crate::SubgraphNode {
5077            id: NodeId(1),
5078            node_type: NodeTypeRef {
5079                type_name: "process".into(),
5080                version: 1,
5081            },
5082            config: BTreeMap::from([("gain".into(), crate::ConfigValue::Text("1".into()))]),
5083            child: None,
5084        });
5085        repeated.id = subgraph_identity(&repeated);
5086        assert_ne!(
5087            leaf.id, repeated.id,
5088            "repeated instances are identity-bearing"
5089        );
5090        let mut parent = crate::SubgraphSchema {
5091            id: crate::SubgraphId([0; 32]),
5092            version: 1,
5093            nodes: vec![crate::SubgraphNode {
5094                id: NodeId(0),
5095                node_type: NodeTypeRef {
5096                    type_name: "process".into(),
5097                    version: 1,
5098                },
5099                config: BTreeMap::from([("gain".into(), crate::ConfigValue::Text("1".into()))]),
5100                child: Some(leaf.id),
5101            }],
5102            edges: vec![],
5103            inputs: vec![crate::SubgraphInterfacePort {
5104                name: "in".into(),
5105                inner: PortRef {
5106                    node: NodeId(0),
5107                    port: "in".into(),
5108                },
5109            }],
5110            outputs: vec![crate::SubgraphInterfacePort {
5111                name: "out".into(),
5112                inner: PortRef {
5113                    node: NodeId(0),
5114                    port: "out".into(),
5115                },
5116            }],
5117        };
5118        parent.id = subgraph_identity(&parent);
5119        registry.register_subgraph(parent.clone()).unwrap();
5120        registry
5121            .descriptors
5122            .get_mut(&("process".into(), 1))
5123            .unwrap()
5124            .subgraph = Some(crate::SubgraphLowering {
5125            subgraph: parent.id,
5126            input_map: vec![crate::PortMap {
5127                outer: "in".into(),
5128                inner: "in".into(),
5129            }],
5130            output_map: vec![crate::PortMap {
5131                outer: "out".into(),
5132                inner: "out".into(),
5133            }],
5134            config_map: vec![crate::SubgraphConfigMap {
5135                outer: "gain".into(),
5136                node: NodeId(0),
5137                inner: "gain".into(),
5138            }],
5139        });
5140
5141        assert_eq!(
5142            Compiler::new(&registry, ExecutionRealm::HostStream)
5143                .with_limits(CompileLimits {
5144                    max_subgraph_entries: 5,
5145                    ..CompileLimits::default()
5146                })
5147                .compile(&graph),
5148            Err(CompileError::SubgraphEntryLimitExceeded)
5149        );
5150        assert!(matches!(
5151            Compiler::new(&registry, ExecutionRealm::HostStream)
5152                .with_limits(CompileLimits {
5153                    max_subgraph_depth: 1,
5154                    ..CompileLimits::default()
5155                })
5156                .compile(&graph),
5157            Err(CompileError::SubgraphDepthExceeded)
5158        ));
5159        let plan = Compiler::new(&registry, ExecutionRealm::HostStream)
5160            .with_limits(CompileLimits {
5161                max_subgraph_depth: 2,
5162                ..CompileLimits::default()
5163            })
5164            .compile(&graph)
5165            .unwrap();
5166        assert!(
5167            plan.nodes
5168                .iter()
5169                .any(|step| step.subgraph_path == [parent.id])
5170        );
5171    }
5172}