1use 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 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 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 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
906pub 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(¤t).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 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(®istry_a, ExecutionRealm::HostStream)
3752 .compile(&graph_a)
3753 .unwrap();
3754 let b = Compiler::new(®istry_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(®istry_a, ExecutionRealm::HostStream)
3782 .compile(&graph)
3783 .unwrap();
3784 let b = Compiler::new(®istry_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(®istry, 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(®istry, 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(®istry, 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(®istry, 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(®istry, 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(®istry, 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(®istry, 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(®istry, 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(®istry, 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(®istry, 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(®istry, 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(®istry, 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 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 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(®istry, 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(®istry, ExecutionRealm::HostStream)
4089 .compile(&graph)
4090 .unwrap();
4091 let plain = Compiler::new(®istry, 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(®istry, 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(®istry, 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(®istry, 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(®istry, 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(®istry, 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(®istry, 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(®istry, 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(®istry, 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(®istry, ExecutionRealm::HostStream)
4338 .compile(&graph)
4339 .unwrap();
4340 let mcu = Compiler::new(®istry, ExecutionRealm::McuAot)
4341 .compile(&graph)
4342 .unwrap();
4343 let durable = Compiler::new(®istry, 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(®istry, 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(®istry, 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(®istry, 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(®istry, 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(®istry, 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(®istry, 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(®istry, 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(®istry, 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(®istry, 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(®istry, 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(®istry, 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(®istry, 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(®istry, 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(®istry, 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(®istry, 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(®istry, 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(®istry, 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(®istry, 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(®istry, 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(®istry, 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(®istry, 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}