diff --git a/tools/benchmarks/Cargo.toml b/tools/benchmarks/Cargo.toml index 898114c..11e23c3 100644 --- a/tools/benchmarks/Cargo.toml +++ b/tools/benchmarks/Cargo.toml @@ -8,7 +8,14 @@ autobenches = false [dependencies] serde = { version = "1", features = ["derive"] } +serde_json = "1" +ciborium = "0.2" +rmp-serde = "1" +lz4_flex = { version = "0.11", default-features = false, features = ["std"] } +zstd = { version = "0.13", default-features = false } +iroh-driver = { path = "../../crates/iroh-driver" } swactor = { path = "../.." } +swactor-engine = { path = "../../crates/engine" } swactor-transport = { path = "../../crates/transport" } telemetry = { path = "../../crates/telemetry" } diff --git a/tools/benchmarks/src/lib.rs b/tools/benchmarks/src/lib.rs index d5b43c5..6741a09 100644 --- a/tools/benchmarks/src/lib.rs +++ b/tools/benchmarks/src/lib.rs @@ -33,7 +33,7 @@ mod tests { .map(|workload| workload.name()) .collect(); assert_eq!(codec_names.len(), 53); - assert_eq!(telemetry_names.len(), 31); + assert_eq!(telemetry_names.len(), 42); assert!(codec_names.contains(&"codec/direct/fixed/encode/8b".to_owned())); assert!(codec_names.contains(&"codec/registry/json-nested/receive/65536b".to_owned())); assert!( @@ -42,6 +42,21 @@ mod tests { assert!( telemetry_names.contains(&"telemetry/mux/concurrent/producers-8/total-4096".to_owned()) ); + assert!(telemetry_names.contains( + &"telemetry/engine/multithread/channels-256/producers-8/total-16384".to_owned() + )); + assert!( + telemetry_names.contains( + &"telemetry/wire/subscription-batch/channels-256/frames-16384".to_owned() + ) + ); + assert!( + telemetry_names + .contains(&"telemetry/payload-codec/messagepack/encode/records-4096".to_owned()) + ); + assert!( + telemetry_names.contains(&"telemetry/wire/compression/zstd/frames-16384".to_owned()) + ); let mut names = codec_names; names.extend(telemetry_names); @@ -52,6 +67,58 @@ mod tests { assert!(names.iter().all(|name| !name.contains(char::is_whitespace))); } + #[test] + fn representative_telemetry_bandwidth_is_bounded() { + let (payload_bytes, wire_bytes) = telemetry::representative_wire_sizes(); + println!("telemetry payload_bytes={payload_bytes} wire_bytes={wire_bytes}"); + assert!(wire_bytes < payload_bytes / 4); + } + + #[test] + fn binary_payload_codec_sizes_beat_json() { + let sizes = telemetry::representative_payload_codec_sizes(); + println!("telemetry payload codec sizes: {sizes:?}"); + let size = |name| { + sizes + .iter() + .find_map(|(codec, bytes)| (*codec == name).then_some(*bytes)) + .expect("codec size") + }; + assert!(size("cbor") < size("json")); + assert!(size("messagepack") < size("json")); + assert!(size("messagepack") < size("cbor")); + } + + #[test] + fn binary_record_codecs_reduce_compressed_wire_bytes() { + let sizes = telemetry::representative_wire_codec_sizes(); + println!("telemetry wire codec sizes: {sizes:?}"); + let wire_size = |name| { + sizes + .iter() + .find_map(|(codec, _, wire_bytes)| (*codec == name).then_some(*wire_bytes)) + .expect("wire codec size") + }; + assert!(wire_size("cbor") < wire_size("json")); + assert!(wire_size("messagepack") < wire_size("json")); + assert!(wire_size("messagepack") < wire_size("cbor")); + } + + #[test] + fn compression_candidates_preserve_bytes_and_report_size() { + let sizes = telemetry::representative_compression_sizes(); + println!("telemetry compression sizes: {sizes:?}"); + assert!(sizes.iter().all(|(_, raw, compressed)| compressed < raw)); + let compressed_size = |name| { + sizes + .iter() + .find_map(|(codec, _, bytes)| (*codec == name).then_some(*bytes)) + .expect("compression size") + }; + assert!(compressed_size("zstd") < compressed_size("zstd-fast")); + assert!(compressed_size("zstd-fast") < compressed_size("lz4")); + } + #[test] fn benchmark_commands_name_the_package() { assert!(BENCHMARK_COMMAND.contains("-p swactor-benchmarks")); diff --git a/tools/benchmarks/src/telemetry.rs b/tools/benchmarks/src/telemetry.rs index db7e229..dc58235 100644 --- a/tools/benchmarks/src/telemetry.rs +++ b/tools/benchmarks/src/telemetry.rs @@ -1,14 +1,23 @@ use std::collections::HashSet; -use std::sync::{Arc, Barrier}; +use std::sync::{Arc, Barrier, mpsc}; use std::thread::{self, JoinHandle}; -use telemetry::frame::{Frame, TelemetryEvent}; +use serde::{Deserialize, Serialize}; + +use iroh_driver::telemetry_transport::{ + decode_event_records, encode_event_batch, encode_event_record, +}; +use swactor::config::RuntimeConfig; +use swactor::runtime::RuntimeParts; +use swactor_engine::{Engine, EngineHandle, TokioBackend, TokioConfig}; + +use telemetry::frame::{ChannelRef, Frame, FrameDelivery, TelemetryEvent}; use telemetry::ingest::Consumer; use telemetry::transport::Delivery; use telemetry::wire::{decode_delivery, encode_delivery}; use telemetry::{ - ChannelContent, ChannelId, Lifetime, Mux, Position, StreamId, TelemetryEndpoint, - TelemetrySubscription, + ChannelContent, ChannelId, Lifetime, Mux, NodeId, Position, StreamDescriptor, StreamId, + StreamOrigin, TelemetryEndpoint, TelemetrySubscription, }; use crate::{SetupPolicy, WorkUnits, Workload}; @@ -18,6 +27,12 @@ const BATCH_SIZES: [usize; 3] = [1, 64, 1024]; const STORE_BATCH: usize = 1024; const CONCURRENT_TOTAL: usize = 4096; const CONCURRENT_PAYLOAD_SIZE: usize = 32; +const ENGINE_WORKERS: usize = 4; +const ENGINE_PRODUCERS: usize = 8; +const ENGINE_CHANNELS: usize = 256; +const ENGINE_TOTAL: usize = 16_384; +const PAYLOAD_CODEC_TOTAL: usize = 4096; +const COMPRESSION_CHUNK_BYTES: usize = 256 * 1024; #[derive(Clone, Copy, Debug)] pub enum StoreKind { @@ -53,6 +68,55 @@ impl WireOperation { } } +#[derive(Clone, Copy, Debug)] +pub enum PayloadCodec { + Json, + Cbor, + MessagePack, +} + +impl PayloadCodec { + fn label(self) -> &'static str { + match self { + Self::Json => "json", + Self::Cbor => "cbor", + Self::MessagePack => "messagepack", + } + } +} + +#[derive(Clone, Copy, Debug)] +pub enum PayloadCodecOperation { + Encode, + Decode, +} + +impl PayloadCodecOperation { + fn label(self) -> &'static str { + match self { + Self::Encode => "encode", + Self::Decode => "decode", + } + } +} + +#[derive(Clone, Copy, Debug)] +pub enum CompressionCodec { + Lz4, + ZstdFast, + Zstd, +} + +impl CompressionCodec { + fn label(self) -> &'static str { + match self { + Self::Lz4 => "lz4", + Self::ZstdFast => "zstd-fast", + Self::Zstd => "zstd", + } + } +} + #[derive(Clone, Copy, Debug)] pub enum TelemetryWorkload { Mux { @@ -71,9 +135,16 @@ pub enum TelemetryWorkload { operation: WireOperation, payload_size: usize, }, + WireBatch, Concurrent { producers: usize, }, + EngineConcurrent, + PayloadCodec { + codec: PayloadCodec, + operation: PayloadCodecOperation, + }, + Compression(CompressionCodec), } pub struct MuxState { @@ -105,6 +176,12 @@ pub struct WireState { encoded: Vec, } +pub struct WireBatchState { + descriptor: StreamDescriptor, + events: Vec, + payload_bytes: usize, +} + pub struct ConcurrentState { endpoint: TelemetryEndpoint, start: Arc, @@ -112,12 +189,52 @@ pub struct ConcurrentState { total: usize, } +pub struct EngineConcurrentState { + _engine: Engine, + handle: EngineHandle, + endpoint: Arc, + subscriptions: Vec, + submissions: Vec)>>, +} + +#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] +pub struct MyelinBenchmarkDetail { + channel: usize, + message: String, +} + +#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] +pub struct MyelinBenchmarkRecord { + schema_version: u32, + event_type: String, + producer_sequence: u64, + node_id: u64, + status: String, + detail: MyelinBenchmarkDetail, +} + +pub struct PayloadCodecState { + codec: PayloadCodec, + records: Vec, + encoded: Vec>, +} + +pub struct CompressionState { + codec: CompressionCodec, + chunks: Vec>, + raw_bytes: usize, +} + pub enum TelemetryState { Mux(MuxState), Fanout(FanoutState), Store(StoreState), Wire(WireState), + WireBatch(WireBatchState), Concurrent(ConcurrentState), + EngineConcurrent(EngineConcurrentState), + PayloadCodec(PayloadCodecState), + Compression(CompressionState), } pub enum TelemetryOutput { @@ -143,11 +260,22 @@ pub enum TelemetryOutput { }, WireEncoded(Vec), WireDecoded(StreamId, Frame), + WireBatch(Vec), Concurrent { accepted: usize, frames: Vec, dropped: u64, }, + EngineConcurrent { + accepted: usize, + drained: usize, + delivered: usize, + events: Vec>, + dropped: u64, + }, + PayloadEncoded(Vec>), + PayloadDecoded(Vec), + Compressed(Vec>), } fn deterministic_payload(size: usize) -> Vec { @@ -163,6 +291,54 @@ fn marked_payload(size: usize, marker: u64) -> Vec { payload } +fn myelin_record(sequence: u64, channel: usize) -> MyelinBenchmarkRecord { + let event_type = match channel % 4 { + 0 => "HostCpuSample", + 1 => "RuntimeActorActivity", + 2 => "WorkerStep", + _ => "ProvisionLogLine", + }; + MyelinBenchmarkRecord { + schema_version: 1, + event_type: event_type.to_owned(), + producer_sequence: sequence, + node_id: sequence % 64, + status: "running".to_owned(), + detail: MyelinBenchmarkDetail { + channel, + message: "telemetry benchmark workload".to_owned(), + }, + } +} + +fn encode_payload(codec: PayloadCodec, record: &MyelinBenchmarkRecord) -> Vec { + match codec { + PayloadCodec::Json => serde_json::to_vec(record).expect("encode benchmark JSON"), + PayloadCodec::Cbor => { + let mut bytes = Vec::new(); + ciborium::ser::into_writer(record, &mut bytes).expect("encode benchmark CBOR"); + bytes + } + PayloadCodec::MessagePack => { + rmp_serde::to_vec_named(record).expect("encode benchmark MessagePack") + } + } +} + +fn decode_payload(codec: PayloadCodec, payload: &[u8]) -> MyelinBenchmarkRecord { + match codec { + PayloadCodec::Json => serde_json::from_slice(payload).expect("decode benchmark JSON"), + PayloadCodec::Cbor => ciborium::de::from_reader(payload).expect("decode benchmark CBOR"), + PayloadCodec::MessagePack => { + rmp_serde::from_slice(payload).expect("decode benchmark MessagePack") + } + } +} + +fn myelin_payload(codec: PayloadCodec, sequence: u64, channel: usize) -> Vec { + encode_payload(codec, &myelin_record(sequence, channel)) +} + fn marker(payload: &[u8]) -> u64 { u64::from_le_bytes(payload[..8].try_into().unwrap()) } @@ -311,6 +487,155 @@ fn setup_concurrent(producers: usize) -> ConcurrentState { } } +fn setup_engine_concurrent() -> EngineConcurrentState { + let endpoint = Arc::new(TelemetryEndpoint::with_capacity( + StreamId::new("bench-engine", Lifetime(1)), + ENGINE_TOTAL, + ENGINE_TOTAL, + )); + let channel_families = [ + "host.cpu", + "host.gpu", + "host.memory", + "host.net", + "host.storage", + "runtime.actors", + "myelin.worker.step", + "myelin.provisioning.logs", + ]; + let channels = (0..ENGINE_CHANNELS) + .map(|index| { + let family = channel_families[index % channel_families.len()]; + endpoint.register_channel( + format!("{family}.{}", index / channel_families.len()), + ChannelContent::MessagePackRecord { + schema: Some(family.to_owned()), + }, + ) + }) + .collect::>(); + let subscriptions = ["dashboard", "archive"] + .into_iter() + .map(|name| endpoint.subscribe_all_with_capacity(name, ENGINE_TOTAL)) + .collect(); + let per_producer = ENGINE_TOTAL / ENGINE_PRODUCERS; + let submissions = (0..ENGINE_PRODUCERS) + .map(|producer_index| { + (0..per_producer) + .map(|index| { + let marker = (producer_index * per_producer + index) as u64; + ( + channels[marker as usize % channels.len()], + myelin_payload( + PayloadCodec::MessagePack, + marker, + marker as usize % channels.len(), + ), + ) + }) + .collect() + }) + .collect(); + let parts = RuntimeParts::new(RuntimeConfig { + worker_count: ENGINE_WORKERS, + ..RuntimeConfig::default() + }); + let engine = Engine::new( + parts, + TokioBackend::new(TokioConfig { + worker_threads: ENGINE_WORKERS, + ..TokioConfig::default() + }) + .expect("benchmark Tokio backend"), + ) + .expect("benchmark engine"); + let handle = engine.handle(); + EngineConcurrentState { + _engine: engine, + handle, + endpoint, + subscriptions, + submissions, + } +} + +fn setup_wire_batch_with_codec(codec: PayloadCodec) -> WireBatchState { + let descriptor = StreamDescriptor { + stream: StreamId::new(NodeId::new("myelin-benchmark-node"), Lifetime(7)), + label: Some("myelin worker node".to_owned()), + origin: StreamOrigin::RemoteNode, + }; + let events = (0..ENGINE_TOTAL) + .map(|sequence| { + let channel = sequence % ENGINE_CHANNELS; + TelemetryEvent::Frame(FrameDelivery { + channel: ChannelRef { + stream: descriptor.stream.clone(), + channel: ChannelId((channel + 1) as u32), + }, + position: Position(sequence as u64), + payload: myelin_payload(codec, sequence as u64, channel), + }) + }) + .collect::>(); + let payload_bytes = events + .iter() + .map(|event| match event { + TelemetryEvent::Frame(frame) => frame.payload.len(), + _ => 0, + }) + .sum(); + WireBatchState { + descriptor, + events, + payload_bytes, + } +} + +fn setup_wire_batch() -> WireBatchState { + setup_wire_batch_with_codec(PayloadCodec::MessagePack) +} + +fn setup_compression(codec: CompressionCodec) -> CompressionState { + let state = setup_wire_batch(); + let mut chunks = vec![Vec::with_capacity(COMPRESSION_CHUNK_BYTES)]; + let mut record = Vec::new(); + for event in &state.events { + record.clear(); + encode_event_record(event, &mut record).expect("encode raw telemetry event"); + if !chunks.last().expect("compression chunk").is_empty() + && chunks.last().expect("compression chunk").len() + record.len() + > COMPRESSION_CHUNK_BYTES + { + chunks.push(Vec::with_capacity(COMPRESSION_CHUNK_BYTES)); + } + chunks + .last_mut() + .expect("compression chunk") + .extend_from_slice(&record); + } + let raw_bytes = chunks.iter().map(Vec::len).sum(); + CompressionState { + codec, + chunks, + raw_bytes, + } +} + +fn setup_payload_codec(codec: PayloadCodec) -> PayloadCodecState { + let records = (0..PAYLOAD_CODEC_TOTAL) + .map(|sequence| myelin_record(sequence as u64, sequence % ENGINE_CHANNELS)) + .collect::>(); + let encoded = records + .iter() + .map(|record| encode_payload(codec, record)) + .collect(); + PayloadCodecState { + codec, + records, + encoded, + } +} fn execute_mux(state: &mut MuxState) -> TelemetryOutput { if state.reject_full { let rejected = (0..state.batch) @@ -422,6 +747,148 @@ fn execute_concurrent(state: &mut ConcurrentState) -> TelemetryOutput { } } +fn execute_engine_concurrent(state: &mut EngineConcurrentState) -> TelemetryOutput { + let (done_tx, done_rx) = mpsc::channel(); + for submissions in std::mem::take(&mut state.submissions) { + let done_tx = done_tx.clone(); + let producer = state.endpoint.producer(); + state.handle.spawn(async move { + let accepted: usize = submissions + .into_iter() + .map(|(channel, payload)| usize::from(producer.submit_bytes(channel, payload))) + .sum(); + let _ = done_tx.send(accepted); + }); + } + drop(done_tx); + let accepted = done_rx.iter().sum(); + let tick = state.endpoint.tick(); + let events = state + .subscriptions + .iter() + .map(TelemetrySubscription::drain_available) + .collect(); + TelemetryOutput::EngineConcurrent { + accepted, + drained: tick.drained, + delivered: tick.delivered, + events, + dropped: state.endpoint.mux_dropped(), + } +} + +fn execute_wire_batch(state: &WireBatchState) -> TelemetryOutput { + let mut encoded = Vec::with_capacity(state.payload_bytes); + encode_event_batch(&state.events, &mut encoded).expect("encode telemetry subscription events"); + TelemetryOutput::WireBatch(encoded) +} + +fn execute_payload_codec( + state: &PayloadCodecState, + codec: PayloadCodec, + operation: PayloadCodecOperation, +) -> TelemetryOutput { + match operation { + PayloadCodecOperation::Encode => TelemetryOutput::PayloadEncoded( + state + .records + .iter() + .map(|record| encode_payload(codec, record)) + .collect(), + ), + PayloadCodecOperation::Decode => TelemetryOutput::PayloadDecoded( + state + .encoded + .iter() + .map(|payload| decode_payload(codec, payload)) + .collect(), + ), + } +} + +fn execute_compression(state: &CompressionState) -> TelemetryOutput { + let compressed = match state.codec { + CompressionCodec::Lz4 => state + .chunks + .iter() + .map(|raw| lz4_flex::block::compress(raw)) + .collect(), + CompressionCodec::ZstdFast | CompressionCodec::Zstd => { + let level = if matches!(state.codec, CompressionCodec::ZstdFast) { + -5 + } else { + 1 + }; + let mut compressor = + zstd::bulk::Compressor::new(level).expect("create benchmark zstd compressor"); + state + .chunks + .iter() + .map(|raw| { + compressor + .compress(raw) + .expect("compress telemetry batch with zstd") + }) + .collect() + } + }; + TelemetryOutput::Compressed(compressed) +} + +pub fn representative_wire_sizes() -> (usize, usize) { + let state = setup_wire_batch(); + let TelemetryOutput::WireBatch(encoded) = execute_wire_batch(&state) else { + unreachable!("wire batch execution returns wire bytes"); + }; + (state.payload_bytes, encoded.len()) +} + +pub fn representative_wire_codec_sizes() -> [(&'static str, usize, usize); 3] { + [ + PayloadCodec::Json, + PayloadCodec::Cbor, + PayloadCodec::MessagePack, + ] + .map(|codec| { + let state = setup_wire_batch_with_codec(codec); + let TelemetryOutput::WireBatch(encoded) = execute_wire_batch(&state) else { + unreachable!("wire batch execution returns wire bytes"); + }; + (codec.label(), state.payload_bytes, encoded.len()) + }) +} + +pub fn representative_compression_sizes() -> [(&'static str, usize, usize); 3] { + [ + CompressionCodec::Lz4, + CompressionCodec::ZstdFast, + CompressionCodec::Zstd, + ] + .map(|codec| { + let state = setup_compression(codec); + let TelemetryOutput::Compressed(compressed) = execute_compression(&state) else { + unreachable!("compression execution returns bytes"); + }; + let compressed_bytes = compressed.iter().map(Vec::len).sum(); + (codec.label(), state.raw_bytes, compressed_bytes) + }) +} + +pub fn representative_payload_codec_sizes() -> [(&'static str, usize); 3] { + [ + PayloadCodec::Json, + PayloadCodec::Cbor, + PayloadCodec::MessagePack, + ] + .map(|codec| { + let state = setup_payload_codec(codec); + ( + codec.label(), + state.encoded.iter().map(Vec::len).sum::(), + ) + }) +} + impl Workload for TelemetryWorkload { type State = TelemetryState; type Output = TelemetryOutput; @@ -458,9 +925,24 @@ impl Workload for TelemetryWorkload { operation, payload_size, } => format!("telemetry/wire/{}/{}b", operation.label(), payload_size), + Self::WireBatch => format!( + "telemetry/wire/subscription-batch/channels-{ENGINE_CHANNELS}/frames-{ENGINE_TOTAL}" + ), Self::Concurrent { producers } => { format!("telemetry/mux/concurrent/producers-{producers}/total-{CONCURRENT_TOTAL}") } + Self::EngineConcurrent => format!( + "telemetry/engine/multithread/channels-{ENGINE_CHANNELS}/producers-{ENGINE_PRODUCERS}/total-{ENGINE_TOTAL}" + ), + Self::PayloadCodec { codec, operation } => format!( + "telemetry/payload-codec/{}/{}/records-{PAYLOAD_CODEC_TOTAL}", + codec.label(), + operation.label(), + ), + Self::Compression(codec) => format!( + "telemetry/wire/compression/{}/frames-{ENGINE_TOTAL}", + codec.label() + ), } } @@ -484,9 +966,15 @@ impl Workload for TelemetryWorkload { )), Self::Store(kind) => TelemetryState::Store(setup_store(*kind)), Self::Wire { payload_size, .. } => TelemetryState::Wire(setup_wire(*payload_size)), + Self::WireBatch => TelemetryState::WireBatch(setup_wire_batch()), Self::Concurrent { producers } => { TelemetryState::Concurrent(setup_concurrent(*producers)) } + Self::EngineConcurrent => TelemetryState::EngineConcurrent(setup_engine_concurrent()), + Self::PayloadCodec { codec, .. } => { + TelemetryState::PayloadCodec(setup_payload_codec(*codec)) + } + Self::Compression(codec) => TelemetryState::Compression(setup_compression(*codec)), } } @@ -498,9 +986,19 @@ impl Workload for TelemetryWorkload { (Self::Wire { operation, .. }, TelemetryState::Wire(state)) => { execute_wire(state, *operation) } + (Self::WireBatch, TelemetryState::WireBatch(state)) => execute_wire_batch(state), (Self::Concurrent { .. }, TelemetryState::Concurrent(state)) => { execute_concurrent(state) } + (Self::EngineConcurrent, TelemetryState::EngineConcurrent(state)) => { + execute_engine_concurrent(state) + } + (Self::PayloadCodec { codec, operation }, TelemetryState::PayloadCodec(state)) => { + execute_payload_codec(state, *codec, *operation) + } + (Self::Compression(_), TelemetryState::Compression(state)) => { + execute_compression(state) + } _ => panic!("telemetry workload and state mismatch"), } } @@ -649,6 +1147,81 @@ impl Workload for TelemetryWorkload { state.total * CONCURRENT_PAYLOAD_SIZE ); } + (TelemetryState::WireBatch(state), TelemetryOutput::WireBatch(encoded)) => { + let decoded = decode_event_records(encoded, &state.descriptor) + .expect("decode telemetry subscription events"); + assert_eq!(decoded, state.events); + assert!(encoded.len() < state.payload_bytes / 4); + } + ( + TelemetryState::EngineConcurrent(_), + TelemetryOutput::EngineConcurrent { + accepted, + drained, + delivered, + events, + dropped, + }, + ) => { + assert_eq!(*accepted, ENGINE_TOTAL); + assert_eq!(*drained, ENGINE_TOTAL); + assert_eq!(*delivered, ENGINE_TOTAL * 2); + assert_eq!(*dropped, 0); + assert_eq!(events.len(), 2); + assert_eq!(events[0], events[1]); + assert_eq!(events[0].len(), ENGINE_TOTAL); + let frames = events[0] + .iter() + .map(|event| match event { + TelemetryEvent::Frame(frame) => frame, + other => panic!("unexpected engine event: {other:?}"), + }) + .collect::>(); + let positions = frames + .iter() + .map(|frame| frame.position.0) + .collect::>(); + let payloads = frames + .iter() + .map(|frame| frame.payload.as_slice()) + .collect::>(); + let channel_ids = frames + .iter() + .map(|frame| frame.channel.channel) + .collect::>(); + assert!(contiguous(&positions)); + assert_eq!(payloads.len(), ENGINE_TOTAL); + assert_eq!(channel_ids.len(), ENGINE_CHANNELS); + assert!(frames.iter().all(|frame| { + decode_payload(PayloadCodec::MessagePack, &frame.payload).producer_sequence + < ENGINE_TOTAL as u64 + })); + } + (TelemetryState::PayloadCodec(state), TelemetryOutput::PayloadEncoded(encoded)) => { + assert_eq!(encoded.len(), state.records.len()); + let decoded = encoded + .iter() + .map(|payload| decode_payload(state.codec, payload)) + .collect::>(); + assert_eq!(decoded, state.records); + } + (TelemetryState::PayloadCodec(state), TelemetryOutput::PayloadDecoded(decoded)) => { + assert_eq!(decoded, &state.records) + } + (TelemetryState::Compression(state), TelemetryOutput::Compressed(compressed)) => { + assert_eq!(compressed.len(), state.chunks.len()); + for (compressed, raw) in compressed.iter().zip(&state.chunks) { + let decoded = match state.codec { + CompressionCodec::Lz4 => lz4_flex::block::decompress(compressed, raw.len()) + .expect("decompress benchmark LZ4"), + CompressionCodec::ZstdFast | CompressionCodec::Zstd => { + zstd::bulk::decompress(compressed, raw.len()) + .expect("decompress benchmark zstd") + } + }; + assert_eq!(&decoded, raw); + } + } _ => panic!("telemetry workload, state, and output mismatch"), } } @@ -685,12 +1258,18 @@ impl Workload for TelemetryWorkload { STORE_BATCH as u64 }), Self::Wire { payload_size, .. } => WorkUnits::Bytes(*payload_size as u64), + Self::WireBatch => WorkUnits::Frames(ENGINE_TOTAL as u64), Self::Concurrent { .. } => WorkUnits::Operations(CONCURRENT_TOTAL as u64), + Self::EngineConcurrent => WorkUnits::Frames(ENGINE_TOTAL as u64), + Self::PayloadCodec { .. } => WorkUnits::Frames(PAYLOAD_CODEC_TOTAL as u64), + Self::Compression(_) => { + WorkUnits::Bytes(setup_compression(CompressionCodec::Lz4).raw_bytes as u64) + } } } fn setup_policy(&self) -> SetupPolicy { - if matches!(self, Self::Concurrent { .. }) { + if matches!(self, Self::Concurrent { .. } | Self::EngineConcurrent) { SetupPolicy::PerExecution } else { SetupPolicy::Batched @@ -747,9 +1326,24 @@ pub fn workloads() -> Vec { payload_size, }); } + workloads.push(TelemetryWorkload::WireBatch); + + for codec in [ + PayloadCodec::Json, + PayloadCodec::Cbor, + PayloadCodec::MessagePack, + ] { + for operation in [PayloadCodecOperation::Encode, PayloadCodecOperation::Decode] { + workloads.push(TelemetryWorkload::PayloadCodec { codec, operation }); + } + } + workloads.push(TelemetryWorkload::Compression(CompressionCodec::Lz4)); + workloads.push(TelemetryWorkload::Compression(CompressionCodec::ZstdFast)); + workloads.push(TelemetryWorkload::Compression(CompressionCodec::Zstd)); for producers in [1, 2, 4, 8] { workloads.push(TelemetryWorkload::Concurrent { producers }); } + workloads.push(TelemetryWorkload::EngineConcurrent); workloads }