diff --git a/Cargo.lock b/Cargo.lock index 637f86209..83b93ccb4 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -457,6 +457,10 @@ name = "e2e" version = "0.1.0-dev" dependencies = [ "anyhow", + "huntsman-nn-core", + "rand 0.9.5", + "rmp-serde", + "serde", "spider-client", "spider-core", "tokio", diff --git a/taskfiles/test.yaml b/taskfiles/test.yaml index b239f8960..99021b244 100644 --- a/taskfiles/test.yaml +++ b/taskfiles/test.yaml @@ -248,7 +248,8 @@ tasks: "{{.G_TDL_PACKAGES_DIR}}/complex/libcomplex.so" cp "{{.G_RUST_RELEASE_DIR}}/libintegration_test_tasks.so" \ "{{.G_TDL_PACKAGES_DIR}}/integration_test_tasks/libintegration_test_tasks.so" - cargo nextest run --all --all-features --run-ignored all --release + cargo nextest run --all --all-features --run-ignored all --release \ + -E 'not (package(e2e) & kind(test))' - |- for f in ${SPIDER_TEST_INSTRUMENT_OUTPUT_DIR}/*; do if [ -f "$f" ]; then diff --git a/tests/huntsman/e2e/Cargo.toml b/tests/huntsman/e2e/Cargo.toml index 5ca6c0285..b7aa0846c 100644 --- a/tests/huntsman/e2e/Cargo.toml +++ b/tests/huntsman/e2e/Cargo.toml @@ -6,6 +6,10 @@ publish = false [dependencies] anyhow = { workspace = true } +huntsman-nn-core = { path = "../../../examples/huntsman/nn/core" } +rand = { workspace = true } +rmp-serde = { workspace = true } +serde = { workspace = true } spider-client = { workspace = true } spider-core = { workspace = true } tokio = { workspace = true, features = ["macros", "rt-multi-thread", "sync", "time"] } diff --git a/tests/huntsman/e2e/src/lib.rs b/tests/huntsman/e2e/src/lib.rs index 9fb2353d1..03e35b335 100644 --- a/tests/huntsman/e2e/src/lib.rs +++ b/tests/huntsman/e2e/src/lib.rs @@ -1,7 +1,10 @@ //! End-to-end integration-test harness for the huntsman suites. +pub mod nn; +pub mod payload_serde; pub mod test_driver; mod types; +pub use payload_serde::*; pub use test_driver::SpiderTestDriver; pub use types::*; diff --git a/tests/huntsman/e2e/src/nn/mod.rs b/tests/huntsman/e2e/src/nn/mod.rs new file mode 100644 index 000000000..3d45023b4 --- /dev/null +++ b/tests/huntsman/e2e/src/nn/mod.rs @@ -0,0 +1,11 @@ +//! Self-contained neural-network model for the end-to-end test. +//! +//! [`NeuralNetwork`] builds a layered `neuron::dense_*` task graph and reproduces it in-process via +//! [`NeuralNetwork::simulate`]. + +mod network; +mod neuron; +mod wiring; + +pub use network::NeuralNetwork; +pub use neuron::Neuron; diff --git a/tests/huntsman/e2e/src/nn/network.rs b/tests/huntsman/e2e/src/nn/network.rs new file mode 100644 index 000000000..c811d4cf2 --- /dev/null +++ b/tests/huntsman/e2e/src/nn/network.rs @@ -0,0 +1,179 @@ +//! The neural-network model: a layered topology of `neuron::dense_*` neurons whose Spider +//! [`TaskGraph`] and in-process simulation describe the same DAG. + +use huntsman_nn_core::NUM_INPUTS; +use rand::SeedableRng; +use rand::rngs::StdRng; +use spider_core::task::DataTypeDescriptor; +use spider_core::task::TaskDescriptor; +use spider_core::task::TaskGraph; +use spider_core::task::TaskIndex; +use spider_core::task::TaskInputOutputIndex; +use spider_core::task::TdlContext; +use spider_core::task::ValueTypeDescriptor; + +use crate::nn::Neuron; +use crate::nn::wiring; + +/// A randomly-wired, layered neural network of `neuron::dense_*` neurons. +pub struct NeuralNetwork { + /// The layers in layer order. + layers: Vec, +} + +impl NeuralNetwork { + /// Factory function. + /// + /// Validates `layer_specs` via [`wiring::validate`] and generates the inner-layer fan-in + /// wiring deterministically from `seed`. + /// + /// # Returns + /// + /// The newly created [`NeuralNetwork`] on success. + /// + /// # Errors + /// + /// Returns an error if: + /// + /// * Forwards [`wiring::validate`]'s return values on failure. + pub fn new(layer_specs: Vec<(usize, Neuron)>, seed: u64) -> anyhow::Result { + let sizes: Vec = layer_specs.iter().map(|(size, _)| *size).collect(); + wiring::validate(&sizes)?; + let mut rng = StdRng::seed_from_u64(seed); + let fan_ins = wiring::build_wiring(&sizes, &mut rng); + let layers = layer_specs + .into_iter() + .zip(fan_ins) + .map(|((neuron_count, activation), fan_in)| Layer { + neuron_count, + activation, + fan_in, + }) + .collect(); + Ok(Self { layers }) + } + + /// # Returns + /// + /// The number of graph inputs. + #[must_use] + pub fn num_graph_inputs(&self) -> usize { + self.layers[0].neuron_count * NUM_INPUTS + } + + /// Builds the Spider [`TaskGraph`] for this network. + /// + /// # Returns + /// + /// The [`TaskGraph`] for this network on success. + /// + /// # Errors + /// + /// Returns an error if: + /// + /// * Forwards [`TaskGraph::new`]'s return values on failure. + /// * Forwards [`TaskGraph::insert_task`]'s return values on failure. + pub fn to_task_graph(&self) -> anyhow::Result { + let float64 = DataTypeDescriptor::Value(ValueTypeDescriptor::float64()); + let mut graph = TaskGraph::new(None, None)?; + let first_layer = &self.layers[0]; + let mut prev_layer: Vec = Vec::with_capacity(first_layer.neuron_count); + + for _ in 0..first_layer.neuron_count { + let task_idx = graph.insert_task(TaskDescriptor { + tdl_context: TdlContext { + package: PACKAGE.to_owned(), + task_func: first_layer.activation.task_name().to_owned(), + }, + execution_policy: None, + inputs: vec![float64.clone(); NUM_INPUTS], + outputs: vec![float64.clone()], + input_sources: None, + })?; + prev_layer.push(task_idx); + } + + for layer in self.layers.iter().skip(1) { + let mut curr_layer = Vec::with_capacity(layer.neuron_count); + for j in 0..layer.neuron_count { + let input_sources: Vec = layer.fan_in[j] + .iter() + .map(|&src| TaskInputOutputIndex { + task_idx: prev_layer[src], + position: 0, + }) + .collect(); + let task_idx = graph.insert_task(TaskDescriptor { + tdl_context: TdlContext { + package: PACKAGE.to_owned(), + task_func: layer.activation.task_name().to_owned(), + }, + execution_policy: None, + inputs: vec![float64.clone(); NUM_INPUTS], + outputs: vec![float64.clone()], + input_sources: Some(input_sources), + })?; + curr_layer.push(task_idx); + } + prev_layer = curr_layer; + } + + Ok(graph) + } + + /// Computes the network's outputs from graph inputs. + /// + /// # Returns + /// + /// The network's outputs on success. + /// + /// # Errors + /// + /// Returns an error if: + /// + /// * [`anyhow::Error`] if `inputs` length is not [`Self::num_graph_inputs`]. + pub fn simulate(&self, inputs: &[f64]) -> anyhow::Result> { + let expected = self.num_graph_inputs(); + anyhow::ensure!( + inputs.len() == expected, + "expected {expected} graph inputs, got {}", + inputs.len(), + ); + + let first_layer = &self.layers[0]; + let mut layer_outputs: Vec = (0..first_layer.neuron_count) + .map(|i| { + let start = i * NUM_INPUTS; + let mut neuron_inputs = [0.0_f64; NUM_INPUTS]; + neuron_inputs.copy_from_slice(&inputs[start..start + NUM_INPUTS]); + first_layer.activation.evaluate_func()(&neuron_inputs) + }) + .collect(); + + for layer in self.layers.iter().skip(1) { + layer_outputs = (0..layer.neuron_count) + .map(|i| { + let neuron_inputs: [f64; NUM_INPUTS] = + std::array::from_fn(|j| layer_outputs[layer.fan_in[i][j]]); + layer.activation.evaluate_func()(&neuron_inputs) + }) + .collect(); + } + + Ok(layer_outputs) + } +} + +/// Name of the TDL package supplying the `neuron::dense_*` tasks. +const PACKAGE: &str = "nn"; + +/// One layer of the network: its neuron count, activation, and per-neuron fan-in. +struct Layer { + /// Number of neurons in this layer. + neuron_count: usize, + /// Activation applied by every neuron in this layer. + activation: Neuron, + /// Per-neuron fan-in, listing previous-layer output indices feeding each neuron. + /// Empty for layer 0, which reads graph inputs directly. + fan_in: Vec>, +} diff --git a/tests/huntsman/e2e/src/nn/neuron.rs b/tests/huntsman/e2e/src/nn/neuron.rs new file mode 100644 index 000000000..f6338a1ee --- /dev/null +++ b/tests/huntsman/e2e/src/nn/neuron.rs @@ -0,0 +1,46 @@ +//! Activation functions for the end-to-end neural-network test workload. +//! +//! Each [`Neuron`] pairs a Spider `neuron::dense_*` task with the in-process +//! `huntsman_nn_core::dense_*` evaluation function so the task graph and [`super::NeuralNetwork`]'s +//! simulation share one source of truth. + +use huntsman_nn_core::NUM_INPUTS; + +/// A dense-layer neuron activation. +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum Neuron { + /// Rectified-linear activation. + Relu, + + /// Logistic-sigmoid activation. + Sigmoid, + + /// Identity (no-op) activation. + Identity, +} + +impl Neuron { + /// # Returns + /// + /// The `neuron::dense_*` task function name that evaluates this activation. + #[must_use] + pub const fn task_name(self) -> &'static str { + match self { + Self::Relu => "neuron::dense_relu", + Self::Sigmoid => "neuron::dense_sigmoid", + Self::Identity => "neuron::dense_identity", + } + } + + /// # Returns + /// + /// The `huntsman_nn_core::dense_*` function that evaluates this activation. + #[must_use] + pub fn evaluate_func(self) -> fn(&[f64; NUM_INPUTS]) -> f64 { + match self { + Self::Relu => huntsman_nn_core::dense_relu, + Self::Sigmoid => huntsman_nn_core::dense_sigmoid, + Self::Identity => huntsman_nn_core::dense_identity, + } + } +} diff --git a/tests/huntsman/e2e/src/nn/wiring.rs b/tests/huntsman/e2e/src/nn/wiring.rs new file mode 100644 index 000000000..bb2521fe5 --- /dev/null +++ b/tests/huntsman/e2e/src/nn/wiring.rs @@ -0,0 +1,147 @@ +//! Topology wiring for the end-to-end neural-network test workload. +//! +//! Layer 0 neurons read graph inputs directly; inner-layer neurons each draw a fixed fan-in +//! from the previous layer's outputs. + +use huntsman_nn_core::NUM_INPUTS; +use rand::rngs::StdRng; +use rand::seq::SliceRandom; + +/// Validates the layer sizes against the following invariants: +/// +/// * The layer list is non-empty. +/// * Each layer's size is at least the neuron fan-in [`NUM_INPUTS`]. +/// * Each consecutive pair of layers fully covers the previous layer's outputs (`next_size * +/// NUM_INPUTS >= prev_size`). +/// +/// # Errors +/// +/// Returns an error if: +/// +/// * [`anyhow::Error`] if an invariant is violated. +pub fn validate(sizes: &[usize]) -> anyhow::Result<()> { + anyhow::ensure!(!sizes.is_empty(), "at least one layer is required"); + for (i, &size) in sizes.iter().enumerate() { + anyhow::ensure!( + size >= NUM_INPUTS, + "layer {i} size {size} is smaller than the neuron fan-in {NUM_INPUTS}", + ); + } + for (i, window) in sizes.windows(2).enumerate() { + let layer_index = i + 1; + let prev_size = window[0]; + let next_size = window[1]; + let coverage = next_size.checked_mul(NUM_INPUTS).ok_or_else(|| { + anyhow::anyhow!( + "layer {layer_index}'s {next_size} * fan-in {NUM_INPUTS} overflows usize", + ) + })?; + anyhow::ensure!( + coverage >= prev_size, + "layer {layer_index}'s {next_size} * fan-in {NUM_INPUTS} cannot cover previous size \ + {prev_size}", + ); + } + Ok(()) +} + +/// Builds the per-neuron fan-in wiring for every layer. +/// +/// # Returns +/// +/// The per-neuron fan-in wiring, indexed `[layer][neuron][fan_in]`. `wiring[0]` is empty since +/// layer 0 reads graph inputs directly. +pub fn build_wiring(sizes: &[usize], rng: &mut StdRng) -> Vec>> { + let mut wiring = Vec::with_capacity(sizes.len()); + wiring.push(Vec::new()); + for window in sizes.windows(2) { + let prev_size = window[0]; + let next_size = window[1]; + wiring.push(generate_layer_wiring(rng, prev_size, next_size)); + } + wiring +} + +/// Deals numbers from a shuffled deck, reshuffling a fresh permutation whenever the current deck +/// runs out. +struct Dealer { + range_size: usize, + deck: Vec, +} + +impl Dealer { + /// Factory function. + /// + /// # Returns + /// + /// The created [`Dealer`] with a freshly shuffled deck of 0..`range_size`. + const fn new(range_size: usize) -> Self { + Self { + range_size, + deck: Vec::new(), + } + } + + /// Draws the next number, reshuffling when the current deck is exhausted. + /// + /// # Returns + /// + /// The drawn number. + /// + /// # Panics + /// + /// Panics if a reshuffled deck is empty. + fn draw(&mut self, rng: &mut StdRng) -> usize { + if self.deck.is_empty() { + self.deck = shuffled_range(rng, self.range_size); + } + self.deck.pop().expect("reshuffled deck must be non-empty") + } +} + +/// Generates the fan-in for each neuron of one inner layer. +/// +/// # Returns +/// +/// One fan-in vector per neuron in the layer, each containing previous-layer output indices. +/// +/// # Panics +/// +/// Panics if the layer invariants do not hold, i.e. `next_size * NUM_INPUTS < prev_size` or +/// `prev_size < NUM_INPUTS`. +fn generate_layer_wiring(rng: &mut StdRng, prev_size: usize, next_size: usize) -> Vec> { + assert!( + prev_size >= NUM_INPUTS + && next_size + .checked_mul(NUM_INPUTS) + .is_some_and(|coverage| coverage >= prev_size), + "layer invariants do not hold", + ); + let mut slots: Vec> = vec![Vec::new(); next_size]; + let mut dealer = Dealer::new(prev_size); + + for neuron_slots in &mut slots { + for _ in 0..NUM_INPUTS { + let output = loop { + let candidate = dealer.draw(rng); + if !neuron_slots.contains(&candidate) { + break candidate; + } + }; + neuron_slots.push(output); + } + } + + slots +} + +/// Produces a random permutation of the numbers in 0..`range_size`. +/// +/// # Returns +/// +/// The random permutation of 0..`range_size`. +fn shuffled_range(rng: &mut StdRng, range_size: usize) -> Vec { + let mut deck: Vec = (0..range_size).collect(); + deck.shuffle(rng); + deck +} diff --git a/tests/huntsman/e2e/src/payload_serde.rs b/tests/huntsman/e2e/src/payload_serde.rs new file mode 100644 index 000000000..5132b7032 --- /dev/null +++ b/tests/huntsman/e2e/src/payload_serde.rs @@ -0,0 +1,152 @@ +//! Spider task input/output wire-format codec. +//! +//! Converts between a Rust value and the `MessagePack`-encoded +//! [`TaskInput::ValuePayload`] / [`TaskOutput`] payload Spider exchanges over a single job +//! input/output boundary. + +use serde::Serialize; +use serde::de::DeserializeOwned; +use spider_core::types::io::TaskInput; +use spider_core::types::io::TaskOutput; + +/// Encodes `value` as a `MessagePack` [`TaskInput::ValuePayload`]. +/// +/// # Type Parameters +/// +/// * `T` - A serializable input value type. +/// +/// # Returns +/// +/// The msgpack-encoded [`TaskInput::ValuePayload`] on success. +/// +/// # Errors +/// +/// Returns an error if: +/// +/// * Forwards [`rmp_serde::to_vec`]'s return values on failure. +pub fn encode_input(value: &T) -> anyhow::Result { + Ok(TaskInput::ValuePayload(rmp_serde::to_vec(value)?)) +} + +/// Decodes a `MessagePack` [`TaskOutput`] payload into `T`. +/// +/// # Type Parameters +/// +/// * `T` - The deserialized output value type. +/// +/// # Returns +/// +/// The decoded `T` on success. +/// +/// # Errors +/// +/// Returns an error if: +/// +/// * Forwards [`rmp_serde::from_slice`]'s return values on failure. +pub fn decode_output(output: &TaskOutput) -> anyhow::Result +where + T: DeserializeOwned, { + Ok(rmp_serde::from_slice(output)?) +} + +#[cfg(test)] +mod tests { + use serde::Deserialize; + use serde::Serialize; + use serde::de::DeserializeOwned; + use spider_core::types::io::TaskInput; + use spider_core::types::io::TaskOutput; + + use super::decode_output; + use super::encode_input; + + #[derive(Debug, PartialEq, Serialize, Deserialize)] + struct Sample { + flag: bool, + count: i64, + label: String, + } + + /// Round-trips `value` through [`encode_input`] then [`decode_output`]. + fn round_trip(value: &T) -> T + where + T: Serialize + DeserializeOwned, { + let encoded = encode_input(value).expect("encode_input should succeed"); + let TaskInput::ValuePayload(bytes) = encoded; + decode_output(&bytes).expect("decode_output should succeed") + } + + #[test] + fn encode_input_wraps_msgpack_bytes_in_value_payload() { + let value = 42.0_f64; + let encoded = encode_input(&value).expect("encode_input should succeed"); + let expected = rmp_serde::to_vec(&value).expect("rmp_serde::to_vec should succeed"); + assert_eq!(encoded, TaskInput::ValuePayload(expected)); + } + + #[test] + fn round_trip_preserves_floats() { + let values = [ + 0.0_f64, + -0.0, + 1.5, + -2.25, + f64::INFINITY, + f64::NEG_INFINITY, + f64::NAN, + f64::MIN, + f64::MAX, + f64::MIN_POSITIVE, + ]; + for &value in &values { + let got = round_trip(&value); + assert_eq!( + got.to_bits(), + value.to_bits(), + "float {value:?} not preserved" + ); + } + } + + #[test] + fn round_trip_preserves_struct() { + let value = Sample { + flag: true, + count: -7, + label: "hello".to_owned(), + }; + assert_eq!(round_trip(&value), value); + } + + #[test] + fn round_trip_preserves_multiple_values_end_to_end() { + let values: Vec = vec![1.5, -2.25, 3.0, 0.0, f64::INFINITY]; + let inputs: Vec = values + .iter() + .map(encode_input) + .collect::>>() + .expect("encoding all values should succeed"); + let outputs: Vec = inputs + .into_iter() + .map(|TaskInput::ValuePayload(bytes)| bytes) + .collect(); + let decoded: Vec = outputs + .iter() + .map(decode_output) + .collect::>>() + .expect("decoding all values should succeed"); + assert_eq!( + decoded.iter().map(|v| v.to_bits()).collect::>(), + values.iter().map(|v| v.to_bits()).collect::>(), + ); + } + + #[test] + fn decode_output_errors_on_empty_payload() { + let result = decode_output::(&TaskOutput::new()); + assert!( + result.is_err(), + "decoding an empty payload should fail, not panic" + ); + } +} diff --git a/tests/huntsman/e2e/tests/nn.rs b/tests/huntsman/e2e/tests/nn.rs new file mode 100644 index 000000000..ffd0b062c --- /dev/null +++ b/tests/huntsman/e2e/tests/nn.rs @@ -0,0 +1,136 @@ +//! End-to-end test: a layered `neuron::dense_*` task graph run through Spider must match the +//! in-process simulation. + +use std::time::Duration; + +use anyhow::Context; +use anyhow::bail; +use e2e::JobSubmission; +use e2e::SpiderTestDriver; +use e2e::TerminationResult; +use e2e::decode_output; +use e2e::encode_input; +use e2e::nn::NeuralNetwork; +use e2e::nn::Neuron; +use rand::Rng; +use rand::SeedableRng; +use rand::rngs::StdRng; +use tokio::task::JoinSet; + +/// Relative-tolerance float comparison. +const REL_TOL: f64 = 1.0e-12; + +/// Number of layers in the test network. +const NUM_LAYERS: usize = 10; + +/// Neurons per layer in the test network. +const LAYER_SIZE: usize = 1000; + +/// Maximum duration of one neural-network job. +const JOB_TIMEOUT: Duration = Duration::from_secs(300); + +/// Number of neural-network job batches. +const NUM_BATCHES: usize = 3; + +/// Number of concurrent neural-network jobs in each batch. +const NUM_JOBS_PER_BATCH: usize = 8; + +#[tokio::test] +async fn test_nn() -> anyhow::Result<()> { + for batch_index in 0..NUM_BATCHES { + let mut jobs = JoinSet::new(); + for job_index in 0..NUM_JOBS_PER_BATCH { + let seed = u64::try_from(batch_index * NUM_JOBS_PER_BATCH + job_index) + .expect("neural-network job index does not fit in u64"); + jobs.spawn(async move { + run_neural_network_job(seed).await.with_context(|| { + format!( + "neural-network job {job_index} in batch {batch_index} with seed {seed} \ + failed" + ) + }) + }); + } + while let Some(result) = jobs.join_next().await { + result.context("neural-network job task panicked")??; + } + } + + Ok(()) +} + +/// Runs a neural-network job and validates its outputs against the in-process simulation. +/// +/// # Errors +/// +/// Returns an error if: +/// +/// * Forwards [`NeuralNetwork::new`]'s return values on failure. +/// * Forwards [`NeuralNetwork::simulate`]'s return values on failure. +/// * Forwards [`NeuralNetwork::to_task_graph`]'s return values on failure. +/// * Forwards [`encode_input`]'s return values on failure. +/// * Forwards [`SpiderTestDriver::run`]'s return values on failure. +async fn run_neural_network_job(seed: u64) -> anyhow::Result<()> { + let layer_specs = (0..NUM_LAYERS) + .map(|i| { + ( + LAYER_SIZE, + match i % 3 { + 0 => Neuron::Relu, + 1 => Neuron::Sigmoid, + _ => Neuron::Identity, + }, + ) + }) + .collect::>(); + let nn = NeuralNetwork::new(layer_specs, seed)?; + let inputs = random_f64s(nn.num_graph_inputs(), seed); + let expected = nn.simulate(&inputs)?; + let task_graph = nn.to_task_graph()?; + let job = JobSubmission { + resource_group_id: "e2e-nn".to_owned(), + task_graph, + inputs: inputs + .iter() + .map(encode_input) + .collect::>>()?, + }; + + SpiderTestDriver::run(job, JOB_TIMEOUT, async move |_job_id, result| { + let outputs = match result { + TerminationResult::Success(outputs) => outputs, + TerminationResult::Failure(message) => bail!("job failed: {message}"), + TerminationResult::Cancelled => bail!("job cancelled"), + }; + let actual: Vec = outputs + .iter() + .map(decode_output) + .collect::>>()?; + anyhow::ensure!( + actual.len() == expected.len(), + "expected {} outputs, got {}", + expected.len(), + actual.len(), + ); + for (&got, &exp) in actual.iter().zip(expected.iter()) { + let diff = (got - exp).abs(); + let tol = REL_TOL * (1.0 + exp.abs()); + assert!( + got.is_finite() && exp.is_finite() && diff <= tol, + "output mismatch: got={got}, expected={exp}, diff={diff}, tol={tol}", + ); + } + Ok(()) + }) + .await?; + + Ok(()) +} + +/// # Returns +/// +/// `count` number of deterministic random `f64` values seeded by `seed`. +fn random_f64s(count: usize, seed: u64) -> Vec { + let mut rng = StdRng::seed_from_u64(seed); + (0..count).map(|_| rng.random::()).collect() +}