use crate::ast::{PortType, Value};
use crate::kernel::{SharedCell, SharedCellEntry};
#[derive(Debug, Clone, PartialEq)]
pub enum WriteError {
UnknownWire {
key: String,
known: Vec<String>,
},
TypeMismatch {
slot: String,
expected: PortType,
got: PortType,
},
CoordinateSlot {
slot: String,
},
ConstSlot {
slot: String,
},
FromParent {
slot: String,
expected: PortType,
got: PortType,
},
}
impl std::fmt::Display for WriteError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
WriteError::UnknownWire { key, known } => {
write!(f, "unknown wire '{key}': no input slot by this name")?;
if !known.is_empty() {
write!(f, "; this kernel's are {known:?}")?;
}
Ok(())
}
WriteError::CoordinateSlot { slot } => {
write!(
f,
"'{slot}' is a coordinate: advance it with set_inputs, not by name"
)
}
WriteError::ConstSlot { slot } => {
write!(
f,
"'{slot}' holds a const, which only initialization writes: write the \
inputs it reads and call init()"
)
}
WriteError::FromParent {
slot,
expected,
got,
} => {
write!(
f,
"the parent's '{slot}' is {got:?}, but the child declares its input \
'{slot}' as {expected:?}: declare the child's input with the parent's \
type, or convert the value in the parent"
)
}
WriteError::TypeMismatch {
slot,
expected,
got,
} => {
write!(
f,
"type mismatch writing to slot '{slot}': expected {expected:?}, got {got:?} (no auto-adapter available)"
)?;
if matches!(got, PortType::VecF32 | PortType::VecI32)
&& !matches!(
expected,
PortType::VecF32
| PortType::VecI32
| PortType::Str
| PortType::Bytes
| PortType::Json
)
{
write!(
f,
" — collection → scalar requires an explicit \
reduction node in the program (the library \
provides none; `vec_dot` and `vec_norm` are \
the vector reductions that exist)"
)?;
}
Ok(())
}
}
}
}
impl std::error::Error for WriteError {}
pub trait WireKey: sealed::Sealed {
fn resolve<M: Metadata + ?Sized>(self, metadata: &M) -> Option<usize>;
}
mod sealed {
pub trait Sealed {}
impl Sealed for usize {}
impl Sealed for &str {}
impl Sealed for String {}
impl Sealed for &String {}
}
impl WireKey for usize {
#[inline]
fn resolve<M: Metadata + ?Sized>(self, _: &M) -> Option<usize> {
Some(self)
}
}
impl WireKey for &str {
#[inline]
fn resolve<M: Metadata + ?Sized>(self, metadata: &M) -> Option<usize> {
metadata.find_input(self)
}
}
impl WireKey for String {
#[inline]
fn resolve<M: Metadata + ?Sized>(self, metadata: &M) -> Option<usize> {
metadata.find_input(&self)
}
}
impl WireKey for &String {
#[inline]
fn resolve<M: Metadata + ?Sized>(self, metadata: &M) -> Option<usize> {
metadata.find_input(self)
}
}
pub trait Metadata {
fn find_input(&self, name: &str) -> Option<usize>;
fn input_names(&self) -> Vec<String>;
fn output_names(&self) -> Vec<String>;
fn coord_count(&self) -> usize;
fn input_port_type(&self, name: &str) -> Option<PortType>;
fn input_port_type_by_idx(&self, idx: usize) -> Option<PortType>;
fn output_port_type(&self, name: &str) -> Option<PortType>;
}
pub trait Dataflow: Metadata {
fn get_wire_idx(&self, idx: usize) -> Value;
#[inline]
fn get_wire<W: WireKey>(&self, key: W) -> Option<Value> {
key.resolve(self).map(|idx| self.get_wire_idx(idx))
}
}
pub trait Construction: Sized {
type Error;
fn root(matter: super::subcontext::PolydatMatter<'_>) -> Result<Self, Self::Error>;
fn subscope(
&self,
matter: super::subcontext::PolydatMatter<'_>,
) -> Result<Box<dyn Kernel>, Self::Error>;
}
pub trait Kernel: Send + Sync + internals::KernelInternals {
fn engine(&self) -> crate::compile::select::Engine;
fn set_inputs(&mut self, coords: &[u64]);
fn set_input(&mut self, name: &str, value: Value) -> Result<(), WriteError>;
fn set_cursor(
&mut self,
name: &str,
partition: &crate::iteration::cursor_partition::Partition,
) -> Result<(), crate::kernel::WriteError>;
fn eval(&mut self);
fn pull(&mut self, name: &str) -> Value;
fn input_names(&self) -> Vec<String>;
fn output_names(&self) -> Vec<String>;
fn output_type(&self, name: &str) -> Option<PortType>;
fn externs(&self) -> Vec<(String, PortType)>;
fn cursor_schemas(&self) -> &[crate::iteration::source::SourceSchema];
fn plan(&self) -> crate::EnginePlan;
fn input_value(&self, name: &str) -> Option<Value>;
fn input_index(&self, name: &str) -> Option<usize> {
self.input_names().iter().position(|n| n == name)
}
fn set_input_at(&mut self, index: usize, value: Value) -> Result<(), WriteError> {
let name =
self.input_names()
.get(index)
.cloned()
.ok_or_else(|| WriteError::UnknownWire {
key: format!("wire[{index}]"),
known: self.input_names(),
})?;
self.set_input(&name, value)
}
fn output_index(&self, name: &str) -> Option<usize> {
self.output_names().iter().position(|n| n == name)
}
fn pull_at(&mut self, index: usize) -> Value {
let name = self
.output_names()
.get(index)
.cloned()
.unwrap_or_else(|| panic!("no output at index {index}"));
self.pull(&name)
}
fn const_inits(&self) -> &[crate::kernel::ConstInit];
fn init_input_at(&mut self, index: usize, value: Value) -> Result<(), WriteError>;
fn init(&mut self) -> Result<(), crate::KernelError> {
for i in 0..self.const_inits().len() {
let (source, slot, fallback, register) = {
let c = &self.const_inits()[i];
(c.source_index, c.slot_index, c.fallback_index, c.register)
};
if register && register_written(self, i) {
continue;
}
let own =
std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| self.pull_at(source)))
.map_err(|payload| crate::KernelError::ConstInit {
name: self.const_inits()[i].name.clone(),
reason: crate::kernel::panic_message(&payload),
})?;
let value = match own {
Value::None => fallback
.and_then(|f| self.input_value_at(f))
.unwrap_or(Value::None),
v => v,
};
if let (Value::None, engine @ crate::Engine::PureNative(_)) = (&value, self.engine()) {
return Err(crate::KernelError::Refused {
engine,
reason: format!(
"the const '{}' has no value, and pure native code cannot carry a \
`None`; give it a value or run this program on `native`",
self.const_inits()[i].name
),
});
}
self.init_input_at(slot, value)
.map_err(crate::KernelError::Write)?;
}
Ok(())
}
fn traversals(&self) -> &[crate::dsl::traversal::Traversal];
fn traverse(&mut self, index: usize) -> Result<crate::kernel::TraversalStream, String>;
fn traverse_all(&mut self) -> Result<Vec<crate::kernel::TraversalStream>, String> {
(0..self.traversals().len())
.map(|i| self.traverse(i))
.collect()
}
fn invalidate_all(&mut self);
fn shared_cells(&self) -> Vec<SharedCellEntry>;
fn output_cell(&self, _name: &str) -> Option<SharedCell> {
None
}
fn output_modifier(&self, _name: &str) -> crate::dsl::ast::BindingModifier {
crate::dsl::ast::BindingModifier::NONE
}
fn cells_in_scope(&self) -> Vec<SharedCellEntry> {
self.shared_cells()
}
fn set_transit_cells(&mut self, _cells: Vec<SharedCellEntry>) {}
fn scope_coordinates(&self) -> &[super::ScopeCoord] {
&[]
}
fn extend_scope_coordinates(&mut self, _outer: &[super::ScopeCoord]) {}
fn input_port_type(&self, _name: &str) -> Option<PortType> {
None
}
fn bind_input_cell(&mut self, _name: &str, _cell: SharedCell) -> bool {
false
}
fn attach_shared_cell(&mut self, name: &str, cell: SharedCell) -> Result<(), String>;
fn into_program(self: Box<Self>) -> std::sync::Arc<dyn KernelProgram>;
fn ledger(&self) -> &std::sync::Arc<crate::kernel::CompileLedger>;
fn resources(&self) -> &crate::resource::ResourceScope;
fn canonical_hash(&self) -> [u8; 32];
fn instance_hash(&self, ancestors: &[&dyn Kernel]) -> [u8; 32] {
let chain: Vec<[u8; 32]> = ancestors.iter().map(|a| a.canonical_hash()).collect();
crate::kernel::instance_hash_of(self.canonical_hash(), &chain)
}
fn is_equivalent_to(&self, other: &dyn Kernel) -> bool {
self.canonical_hash() == other.canonical_hash()
}
fn is_subset_of(&self, parent: &dyn Kernel) -> bool {
if self.is_equivalent_to(parent) {
return true;
}
let inputs = Kernel::input_names(self);
if Kernel::output_names(self)
.iter()
.any(|name| !inputs.contains(name))
{
return false;
}
let parent_inputs = Kernel::input_names(parent);
inputs.iter().all(|name| parent_inputs.contains(name))
}
fn coord_count(&self) -> usize;
fn input_value_at(&self, index: usize) -> Option<Value>;
fn input_default_at(&self, index: usize) -> Option<Value>;
fn input_is_cell_bound(&self, index: usize) -> bool;
fn reset_inputs(&mut self);
fn fork(&self) -> Box<dyn Kernel>;
fn publish_broadcasts(&mut self);
fn commit_write_throughs(&mut self) -> Result<(), String>;
fn input_type_origin(&self, name: &str) -> Option<crate::kernel::TypeOrigin>;
fn program_id(&self) -> ProgramId;
fn as_interpreter(&self) -> Option<&crate::kernel::PolydatKernel> {
None
}
fn as_interpreter_mut(&mut self) -> Option<&mut crate::kernel::PolydatKernel> {
None
}
}
fn register_written<K: Kernel + ?Sized>(kernel: &K, index: usize) -> bool {
let name = &kernel.const_inits()[index].slot;
kernel.shared_cells().iter().any(|entry| {
&entry.name == name
&& entry
.cell
.revision
.load(std::sync::atomic::Ordering::Acquire)
> 0
})
}
#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)]
pub struct ProgramId(pub(crate) usize);
pub(crate) mod internals {
use crate::ast::{PortType, Value};
pub trait KernelInternals {
fn set_traversals(
&mut self,
traversals: Vec<crate::dsl::traversal::Traversal>,
producers: Vec<crate::dsl::traversal::Producer>,
);
fn slot_value(&self, _slot: usize, _ty: PortType) -> Value {
Value::None
}
fn folded_value(&self, name: &str) -> Option<Value>;
fn set_cursor_extent(&mut self, index: usize, extent: u64);
fn reset_to_program(&mut self) {}
fn set_inherited_outputs(&mut self, names: Vec<String>);
fn set_write_throughs(&mut self, pairs: Vec<(String, String)>);
}
}
pub trait KernelProgram: Send + Sync {
fn engine(&self) -> crate::compile::select::Engine;
fn create_kernel(self: std::sync::Arc<Self>) -> Box<dyn Kernel> {
let mut kernel = self.create_uninitialized();
if let Err(e) = kernel.init() {
panic!("{e}");
}
kernel
}
fn create_uninitialized(self: std::sync::Arc<Self>) -> Box<dyn Kernel>;
fn as_interpreter(
self: std::sync::Arc<Self>,
) -> Option<std::sync::Arc<crate::kernel::PolydatProgram>> {
None
}
fn ledger(&self) -> &std::sync::Arc<crate::kernel::CompileLedger>;
fn resources(&self) -> &crate::resource::ResourceScope;
fn canonical_hash(&self) -> [u8; 32];
fn program_id(&self) -> ProgramId;
}
pub(crate) struct SharedKernel<K>(pub(crate) K);
impl<K: Kernel + Clone + Send + Sync + 'static> KernelProgram for SharedKernel<K> {
fn engine(&self) -> crate::compile::select::Engine {
self.0.engine()
}
fn create_uninitialized(self: std::sync::Arc<Self>) -> Box<dyn Kernel> {
let mut kernel = self.0.clone();
kernel.reset_to_program();
Box::new(kernel)
}
fn ledger(&self) -> &std::sync::Arc<crate::kernel::CompileLedger> {
self.0.ledger()
}
fn resources(&self) -> &crate::resource::ResourceScope {
self.0.resources()
}
fn canonical_hash(&self) -> [u8; 32] {
self.0.canonical_hash()
}
fn program_id(&self) -> ProgramId {
self.0.program_id()
}
}