1use std::collections::HashMap;
21
22use crate::ast::{CompiledSlotOp, CompiledU64Op, PortType, ScratchBuf, ScratchElem};
23
24#[derive(Default)]
28pub(crate) struct P2Extras {
29 pub(crate) output_types: HashMap<String, PortType>,
30 pub(crate) externs: crate::compile::externs::Externs,
33 pub(crate) input_dependents: Vec<Vec<usize>>,
36 pub(crate) attribution: std::sync::Arc<crate::compile::Attribution>,
38}
39
40pub(crate) enum StepOp {
45 U64(CompiledU64Op),
46 Slot(CompiledSlotOp),
47 Copy,
50}
51pub(crate) struct P2Step {
53 pub(crate) name: String,
55 pub(crate) op: StepOp,
56 pub(crate) input_slots: Vec<usize>,
57 pub(crate) output_slots: Vec<usize>,
58 pub(crate) scratch: Vec<ScratchElem>,
60 pub(crate) ref_output_starts: Vec<usize>,
65 pub(crate) accepts_none: bool,
68 pub(crate) volatile: bool,
71 pub(crate) constant: bool,
74 pub(crate) side: bool,
78}
79
80struct CompiledStep {
81 op: StepOp,
82 input_slots: Vec<usize>,
83 output_slots: Vec<usize>,
84 scratch_range: (usize, usize),
85 accepts_none: bool,
87 volatile: bool,
89 constant: bool,
91 side: bool,
94}
95
96type ResolvedOutput = (
99 usize,
100 crate::ast::PortType,
101 Option<std::sync::Arc<[usize]>>,
102 bool,
103);
104
105struct KernelCore {
111 engine: crate::compile::select::Engine,
116 buffer: Vec<u64>,
117 coord_count: usize,
118 steps: std::sync::Arc<[CompiledStep]>,
119 output_map: HashMap<String, usize>,
120 gather_buf: Vec<u64>,
121 scatter_buf: Vec<u64>,
122 scratch: Vec<ScratchBuf>,
126 ref_slots: Vec<bool>,
129 ref_scratch: Vec<(usize, usize)>,
132 output_types: HashMap<String, PortType>,
134 externs: crate::compile::externs::Externs,
136 traversals: std::sync::Arc<[crate::dsl::traversal::Traversal]>,
139 resolved_outputs: Vec<Option<ResolvedOutput>>,
142 drive: crate::compile::Drive,
146 none: Vec<bool>,
149 ran: Vec<u64>,
152 epoch: u64,
158 all_ran: bool,
160 clean: Vec<bool>,
164 use_clean: bool,
168 plan: std::sync::Arc<crate::compile::Invalidation>,
171 slot_step: std::sync::Arc<[Option<usize>]>,
173 sites: std::sync::Arc<crate::compile::Attribution>,
175 cur_step: usize,
177 all: std::sync::Arc<[usize]>,
179 dirty: std::sync::Arc<[Vec<usize>]>,
185 volatile_steps: std::sync::Arc<[usize]>,
187 any_none: bool,
191}
192
193impl Clone for KernelCore {
194 fn clone(&self) -> Self {
195 let mut core = KernelCore {
196 engine: self.engine,
197 buffer: self.buffer.clone(),
198 coord_count: self.coord_count,
199 steps: self.steps.clone(),
200 output_map: self.output_map.clone(),
201 gather_buf: self.gather_buf.clone(),
202 scatter_buf: self.scatter_buf.clone(),
203 scratch: self.scratch.clone(),
204 ref_slots: self.ref_slots.clone(),
205 ref_scratch: self.ref_scratch.clone(),
206 output_types: self.output_types.clone(),
207 externs: self.externs.clone(),
208 traversals: self.traversals.clone(),
209 resolved_outputs: self.resolved_outputs.clone(),
210 drive: self.drive.clone(),
211 none: self.none.clone(),
212 ran: self.ran.clone(),
213 epoch: self.epoch,
214 all_ran: self.all_ran,
215 clean: self.clean.clone(),
216 use_clean: self.use_clean,
217 plan: self.plan.clone(),
218 slot_step: self.slot_step.clone(),
219 sites: self.sites.clone(),
220 cur_step: self.cur_step,
221 all: self.all.clone(),
222 dirty: self.dirty.clone(),
223 volatile_steps: self.volatile_steps.clone(),
224 any_none: self.any_none,
225 };
226 core.republish_refs();
227 core
228 }
229}
230
231impl KernelCore {
232 crate::compile::shared_core_methods!();
233 #[inline]
236 fn step_can_fail(&self, _i: usize) -> bool {
237 true
238 }
239
240 fn program_identity(&self) -> usize {
243 std::sync::Arc::as_ptr(&self.steps) as *const () as usize
244 }
245
246 #[inline]
251 fn failing_node(&self) -> usize {
252 self.cur_step
253 }
254
255 #[inline]
257 fn run_order(&mut self, order: &[usize]) {
258 let steps = &self.steps;
259 let none_free = !self.any_none;
260 for &i in order {
261 if self.all_ran || self.ran[i] == self.epoch {
262 continue;
263 }
264 let step = &steps[i];
265 if (self.use_clean || step.side) && self.clean[i] && !step.volatile {
268 self.ran[i] = self.epoch;
269 continue;
270 }
271 self.cur_step = i;
272 if none_free {
273 run_step_fast(
274 step,
275 &mut self.buffer,
276 &mut self.gather_buf,
277 &mut self.scatter_buf,
278 &mut self.scratch,
279 );
280 } else {
281 run_step(
282 step,
283 &mut self.buffer,
284 &mut self.none,
285 &mut self.gather_buf,
286 &mut self.scatter_buf,
287 &mut self.scratch,
288 );
289 }
290 self.ran[i] = self.epoch;
291 self.clean[i] = !step.volatile;
292 }
293 }
294
295 #[inline]
301 fn run_fresh(&mut self) {
302 let steps = &self.steps;
303 for (i, step) in steps.iter().enumerate() {
304 if step.side {
305 if self.clean[i] && !step.volatile {
306 continue;
307 }
308 self.clean[i] = !step.volatile;
309 }
310 self.cur_step = i;
311 run_step_fast(
312 step,
313 &mut self.buffer,
314 &mut self.gather_buf,
315 &mut self.scatter_buf,
316 &mut self.scratch,
317 );
318 }
319 self.all_ran = true;
320 }
321
322 fn plan(&self) -> crate::EnginePlan {
324 crate::EnginePlan {
325 closure_steps: self.steps.len(),
326 ..Default::default()
327 }
328 }
329}
330
331#[allow(clippy::too_many_arguments)]
334fn build_core(
335 coord_count: usize,
336 total_slots: usize,
337 steps: Vec<P2Step>,
338 output_map: HashMap<String, usize>,
339 ref_slots: Vec<bool>,
340 extras: P2Extras,
341 use_clean: bool,
342 engine: crate::compile::select::Engine,
343) -> Result<KernelCore, crate::KernelError> {
344 let P2Extras {
345 output_types,
346 externs,
347 input_dependents,
348 attribution,
349 } = extras;
350 let max_inputs = steps.iter().map(|s| s.input_slots.len()).max().unwrap_or(0);
351 let max_outputs = steps
352 .iter()
353 .map(|s| s.output_slots.len())
354 .max()
355 .unwrap_or(0);
356 let mut scratch: Vec<ScratchBuf> = Vec::new();
357 let mut ref_scratch: Vec<(usize, usize)> = Vec::new();
358 let compiled_steps: Vec<CompiledStep> = steps
359 .into_iter()
360 .map(|step| {
361 let start = scratch.len();
362 scratch.extend(step.scratch.iter().map(|e| ScratchBuf::new(*e)));
363 ref_scratch.extend(crate::compile::assembly::scratch_pairs(
364 &step.name,
365 &step.ref_output_starts,
366 &step.scratch,
367 start,
368 ));
369 CompiledStep {
370 op: step.op,
371 input_slots: step.input_slots,
372 output_slots: step.output_slots,
373 scratch_range: (start, scratch.len()),
374 accepts_none: step.accepts_none,
375 volatile: step.volatile,
376 constant: step.constant,
377 side: step.side,
378 }
379 })
380 .collect();
381 let mut slot_step: Vec<Option<usize>> = vec![None; total_slots];
382 for (i, step) in compiled_steps.iter().enumerate() {
383 for &s in &step.output_slots {
384 slot_step[s] = Some(i);
385 }
386 }
387 let step_inputs: Vec<&[usize]> = compiled_steps
388 .iter()
389 .map(|s| s.input_slots.as_slice())
390 .collect();
391 let step_outputs: Vec<&[usize]> = compiled_steps
392 .iter()
393 .map(|s| s.output_slots.as_slice())
394 .collect();
395 let plan = crate::compile::Invalidation::from_provenance(
396 input_dependents,
397 &step_inputs,
398 &step_outputs,
399 &output_map,
400 total_slots,
401 );
402 let dirty: Vec<Vec<usize>> = plan
403 .input_dependents
404 .iter()
405 .map(|deps| {
406 if use_clean {
407 deps.clone()
408 } else {
409 deps.iter()
410 .copied()
411 .filter(|&i| compiled_steps[i].side)
412 .collect()
413 }
414 })
415 .collect();
416 let volatile_steps: Vec<usize> = (0..compiled_steps.len())
417 .filter(|&i| compiled_steps[i].volatile)
418 .collect();
419 let mut buffer = vec![0u64; total_slots];
420 let mut none = vec![false; total_slots];
421 let any_none = externs.seed(&mut buffer, Some(&mut none));
422 let step_count = compiled_steps.len();
423 let constants: Vec<usize> = compiled_steps
424 .iter()
425 .enumerate()
426 .filter(|(_, s)| s.constant)
427 .map(|(i, _)| i)
428 .collect();
429 let mut core = KernelCore {
430 engine,
431 buffer,
432 coord_count,
433 steps: compiled_steps.into(),
434 output_map,
435 gather_buf: vec![0u64; max_inputs],
436 scatter_buf: vec![0u64; max_outputs],
437 scratch,
438 ref_slots,
439 ref_scratch,
440 output_types,
441 externs,
442 traversals: Vec::new().into(),
443 resolved_outputs: Vec::new(),
444 drive: crate::compile::Drive {
445 coords: Vec::new(),
446 stale: true,
447 },
448 none,
449 ran: vec![0; step_count],
450 epoch: 0,
451 all_ran: false,
452 clean: vec![false; step_count],
453 use_clean,
454 plan: std::sync::Arc::new(plan),
455 slot_step: slot_step.into(),
456 sites: attribution,
457 cur_step: 0,
458 all: (0..step_count).collect::<Vec<usize>>().into(),
459 dirty: dirty.into(),
460 volatile_steps: volatile_steps.into(),
461 any_none,
462 };
463 core.begin_epoch();
468 core.fold_steps(&constants)?;
469 core.drive.stale = true;
470 Ok(core)
471}
472
473fn compute_slot_provenance(
476 coord_count: usize,
477 total_slots: usize,
478 input_dependents: &[Vec<usize>],
479 steps: &[CompiledStep],
480) -> Vec<crate::kernel::ProvMask> {
481 let outs: Vec<&[usize]> = steps.iter().map(|s| s.output_slots.as_slice()).collect();
482 crate::compile::slot_provenance(coord_count, total_slots, &outs, input_dependents)
483}
484
485macro_rules! closure_writes {
494 () => {
495 pub fn set_input(
500 &mut self,
501 name: &str,
502 value: crate::ast::Value,
503 ) -> Result<(), crate::kernel::WriteError> {
504 let slot = self.core.set_extern(name, value)?;
505 self.mark_input_changed(slot);
506 Ok(())
507 }
508
509 pub fn set_input_at(
511 &mut self,
512 index: usize,
513 value: crate::ast::Value,
514 ) -> Result<(), crate::kernel::WriteError> {
515 let slot = self.core.set_extern_at(index, value)?;
516 self.mark_input_changed(slot);
517 Ok(())
518 }
519
520 fn mark_all_dirty(&mut self) {
524 for i in 0..self.core.coord_count {
525 self.mark_input_changed(i);
526 }
527 }
528 };
529}
530
531#[derive(Clone)]
536pub struct CompiledKernelRaw {
538 core: KernelCore,
539}
540
541impl CompiledKernelRaw {
542 pub(crate) fn new(
543 coord_count: usize,
544 total_slots: usize,
545 steps: Vec<P2Step>,
546 output_map: HashMap<String, usize>,
547 ref_slots: Vec<bool>,
548 extras: P2Extras,
549 ) -> Result<Self, crate::KernelError> {
550 Ok(Self {
551 core: build_core(
552 coord_count,
553 total_slots,
554 steps,
555 output_map,
556 ref_slots,
557 extras,
558 false,
559 Engine::Closures(Provenance::Raw),
560 )?,
561 })
562 }
563
564 fn mark_input_changed(&mut self, slot: usize) {
567 self.core.dirty_input(slot);
568 }
569
570 #[inline]
573 fn set_coords(&mut self, coords: &[u64]) {
574 for (i, &c) in coords
575 .iter()
576 .enumerate()
577 .take(self.core.externs.coordinate_slots())
578 {
579 if self.core.buffer[i] != c {
580 self.core.buffer[i] = c;
581 self.core.dirty_input(i);
582 }
583 }
584 }
585
586 #[inline]
588 pub fn eval(&mut self, coords: &[u64]) {
589 self.set_coords(coords);
590 self.core.drive.stale = true;
591 self.core.eval_all();
592 }
593
594 #[inline]
596 pub fn eval_for_slot(&mut self, coords: &[u64], slot: usize) -> u64 {
597 self.core.guard_ref_slot(slot);
598 self.eval(coords);
599 self.core.buffer[slot]
600 }
601
602 crate::compile::kernel_accessors!(set_coords);
603 closure_writes!();
604}
605
606#[derive(Clone)]
612pub struct CompiledKernelPush {
615 core: KernelCore,
616}
617
618impl CompiledKernelPush {
619 pub(crate) fn new(
620 coord_count: usize,
621 total_slots: usize,
622 steps: Vec<P2Step>,
623 output_map: HashMap<String, usize>,
624 input_dependents: Vec<Vec<usize>>,
625 ref_slots: Vec<bool>,
626 extras: P2Extras,
627 ) -> Result<Self, crate::KernelError> {
628 let _ = input_dependents;
630 Ok(Self {
631 core: build_core(
632 coord_count,
633 total_slots,
634 steps,
635 output_map,
636 ref_slots,
637 extras,
638 true,
639 Engine::Closures(Provenance::Push),
640 )?,
641 })
642 }
643
644 #[inline]
645 fn set_coords(&mut self, coords: &[u64]) {
646 for (i, &c) in coords
647 .iter()
648 .enumerate()
649 .take(self.core.externs.coordinate_slots())
650 {
651 if self.core.buffer[i] != c {
652 self.core.buffer[i] = c;
653 self.core.dirty_input(i);
654 }
655 }
656 }
657
658 fn mark_input_changed(&mut self, slot: usize) {
660 self.core.dirty_input(slot);
661 }
662
663 #[inline]
665 pub fn eval(&mut self, coords: &[u64]) {
666 self.set_coords(coords);
667 self.core.drive.stale = true;
668 self.core.eval_all();
669 }
670
671 #[inline]
673 pub fn eval_for_slot(&mut self, coords: &[u64], slot: usize) -> u64 {
674 self.core.guard_ref_slot(slot);
675 self.eval(coords);
676 self.core.buffer[slot]
677 }
678
679 crate::compile::kernel_accessors!(set_coords);
680 closure_writes!();
681}
682
683#[derive(Clone)]
690pub struct CompiledKernelPull {
693 core: KernelCore,
694 slot_provenance: Vec<crate::kernel::ProvMask>,
695 changed_mask: crate::kernel::ProvMask,
696 force_run: bool,
699}
700
701impl CompiledKernelPull {
702 pub(crate) fn new(
703 coord_count: usize,
704 total_slots: usize,
705 steps: Vec<P2Step>,
706 output_map: HashMap<String, usize>,
707 input_dependents: &[Vec<usize>],
708 ref_slots: Vec<bool>,
709 extras: P2Extras,
710 ) -> Result<Self, crate::KernelError> {
711 let core = build_core(
712 coord_count,
713 total_slots,
714 steps,
715 output_map,
716 ref_slots,
717 extras,
718 false,
719 Engine::Closures(Provenance::Pull),
720 )?;
721 let slot_provenance =
722 compute_slot_provenance(coord_count, total_slots, input_dependents, &core.steps);
723 Ok(Self {
724 core,
725 slot_provenance,
726 changed_mask: crate::kernel::ProvMask::all_below(coord_count), force_run: false,
728 })
729 }
730
731 #[inline]
734 fn set_coords(&mut self, coords: &[u64]) {
735 self.changed_mask.clear();
736 for (i, &c) in coords
737 .iter()
738 .enumerate()
739 .take(self.core.externs.coordinate_slots())
740 {
741 if self.core.buffer[i] != c {
742 self.core.buffer[i] = c;
743 self.changed_mask.set(i);
744 self.core.dirty_input(i);
745 }
746 }
747 }
748
749 fn mark_input_changed(&mut self, slot: usize) {
752 self.core.dirty_input(slot);
753 self.force_run = true;
754 }
755
756 #[inline]
758 pub fn eval(&mut self, coords: &[u64]) {
759 self.set_coords(coords);
760 self.force_run = false;
761 self.core.drive.stale = true;
762 self.core.eval_all();
763 }
764
765 #[inline]
768 pub fn eval_for_slot(&mut self, coords: &[u64], slot: usize) -> u64 {
769 self.core.guard_ref_slot(slot);
770 self.set_coords(coords);
771 if !self.force_run
772 && slot < self.slot_provenance.len()
773 && !self.slot_provenance[slot].intersects(&self.changed_mask)
774 {
775 return self.core.buffer[slot];
776 }
777 self.force_run = false;
778 self.core.drive.stale = true;
779 self.core.eval_all();
780 self.core.buffer[slot]
781 }
782
783 crate::compile::kernel_accessors!(set_coords);
784 closure_writes!();
785}
786
787#[derive(Clone)]
793pub struct CompiledKernelPushPull {
795 core: KernelCore,
796 slot_provenance: Vec<crate::kernel::ProvMask>,
797 changed_mask: crate::kernel::ProvMask,
798 force_run: bool,
801}
802
803impl CompiledKernelPushPull {
804 pub(crate) fn new(
805 coord_count: usize,
806 total_slots: usize,
807 steps: Vec<P2Step>,
808 output_map: HashMap<String, usize>,
809 input_dependents: Vec<Vec<usize>>,
810 ref_slots: Vec<bool>,
811 extras: P2Extras,
812 ) -> Result<Self, crate::KernelError> {
813 let core = build_core(
814 coord_count,
815 total_slots,
816 steps,
817 output_map,
818 ref_slots,
819 extras,
820 true,
821 Engine::Closures(Provenance::PushPull),
822 )?;
823 let slot_provenance =
824 compute_slot_provenance(coord_count, total_slots, &input_dependents, &core.steps);
825 Ok(Self {
826 core,
827 slot_provenance,
828 changed_mask: crate::kernel::ProvMask::all_below(coord_count),
829 force_run: false,
830 })
831 }
832
833 #[inline]
834 fn set_coords(&mut self, coords: &[u64]) {
835 self.changed_mask.clear();
836 for (i, &c) in coords
837 .iter()
838 .enumerate()
839 .take(self.core.externs.coordinate_slots())
840 {
841 if self.core.buffer[i] != c {
842 self.core.buffer[i] = c;
843 self.changed_mask.set(i);
844 self.core.dirty_input(i);
845 }
846 }
847 }
848
849 fn mark_input_changed(&mut self, slot: usize) {
852 self.core.dirty_input(slot);
853 self.force_run = true;
854 }
855
856 #[inline]
858 pub fn eval(&mut self, coords: &[u64]) {
859 self.set_coords(coords);
860 self.force_run = false;
861 self.core.drive.stale = true;
862 self.core.eval_all();
863 }
864
865 #[inline]
867 pub fn eval_for_slot(&mut self, coords: &[u64], slot: usize) -> u64 {
868 self.core.guard_ref_slot(slot);
869 self.set_coords(coords);
870 if !self.force_run
871 && slot < self.slot_provenance.len()
872 && !self.slot_provenance[slot].intersects(&self.changed_mask)
873 {
874 return self.core.buffer[slot];
875 }
876 self.force_run = false;
877 self.core.drive.stale = true;
878 self.core.eval_all();
879 self.core.buffer[slot]
880 }
881
882 crate::compile::kernel_accessors!(set_coords);
883 closure_writes!();
884}
885
886use crate::compile::select::{Engine, Provenance};
889
890crate::compile::impl_kernel_trait!(CompiledKernelRaw);
891crate::compile::impl_kernel_trait!(CompiledKernelPush);
892crate::compile::impl_kernel_trait!(CompiledKernelPull);
893crate::compile::impl_kernel_trait!(CompiledKernelPushPull);
894crate::compile::impl_slot_kernel!(CompiledKernelRaw);
895crate::compile::impl_slot_kernel!(CompiledKernelPush);
896crate::compile::impl_slot_kernel!(CompiledKernelPull);
897crate::compile::impl_slot_kernel!(CompiledKernelPushPull);
898
899#[inline(always)]
903fn run_step(
904 step: &CompiledStep,
905 buffer: &mut [u64],
906 none: &mut [bool],
907 gather: &mut [u64],
908 scatter: &mut [u64],
909 scratch: &mut [ScratchBuf],
910) {
911 let mut any_none = false;
912 for (i, &s) in step.input_slots.iter().enumerate() {
913 gather[i] = buffer[s];
914 any_none |= none[s];
915 }
916 if any_none && !step.accepts_none {
917 for &s in &step.output_slots {
918 none[s] = true;
919 }
920 return;
921 }
922 if matches!(step.op, StepOp::Copy) {
923 for (&i, &o) in step.input_slots.iter().zip(&step.output_slots) {
924 buffer[o] = buffer[i];
925 none[o] = false;
926 }
927 return;
928 }
929 let (n_in, n_out) = (step.input_slots.len(), step.output_slots.len());
930 match &step.op {
931 StepOp::Copy => unreachable!(),
932 StepOp::U64(op) => op(&gather[..n_in], &mut scatter[..n_out]),
933 StepOp::Slot(op) => op(
934 &gather[..n_in],
935 &mut scatter[..n_out],
936 &mut scratch[step.scratch_range.0..step.scratch_range.1],
937 ),
938 }
939 for (i, &s) in step.output_slots.iter().enumerate() {
940 buffer[s] = scatter[i];
941 none[s] = false;
942 }
943}
944
945#[inline(always)]
947fn run_step_fast(
948 step: &CompiledStep,
949 buffer: &mut [u64],
950 gather: &mut [u64],
951 scatter: &mut [u64],
952 scratch: &mut [ScratchBuf],
953) {
954 if matches!(step.op, StepOp::Copy) {
955 for (&i, &o) in step.input_slots.iter().zip(&step.output_slots) {
956 buffer[o] = buffer[i];
957 }
958 return;
959 }
960 for (i, &s) in step.input_slots.iter().enumerate() {
961 gather[i] = buffer[s];
962 }
963 let (n_in, n_out) = (step.input_slots.len(), step.output_slots.len());
964 match &step.op {
965 StepOp::Copy => unreachable!(),
966 StepOp::U64(op) => op(&gather[..n_in], &mut scatter[..n_out]),
967 StepOp::Slot(op) => op(
968 &gather[..n_in],
969 &mut scatter[..n_out],
970 &mut scratch[step.scratch_range.0..step.scratch_range.1],
971 ),
972 }
973 for (i, &s) in step.output_slots.iter().enumerate() {
974 buffer[s] = scatter[i];
975 }
976}