benchmarks: extend telemetry workload and codec measurements

Checkpoint telemetry benchmark updates for payload encoding, compression, batching, fanout, storage, and concurrency alongside the migrated transport dependencies. Keep performance instrumentation separate from runtime correctness and paid qualification claims.

The checkpoint review exercised telemetry transport behavior, not the complete benchmark suite; no fresh benchmark performance result is claimed.
This commit is contained in:
Zachery Aaron Shores-Chmielewski 2026-09-14 12:21:09 +03:00
parent 6b8b1de897
commit a8d0a13ff1
3 changed files with 674 additions and 6 deletions

View file

@ -8,7 +8,14 @@ autobenches = false
[dependencies] [dependencies]
serde = { version = "1", features = ["derive"] } 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 = { path = "../.." }
swactor-engine = { path = "../../crates/engine" }
swactor-transport = { path = "../../crates/transport" } swactor-transport = { path = "../../crates/transport" }
telemetry = { path = "../../crates/telemetry" } telemetry = { path = "../../crates/telemetry" }

View file

@ -33,7 +33,7 @@ mod tests {
.map(|workload| workload.name()) .map(|workload| workload.name())
.collect(); .collect();
assert_eq!(codec_names.len(), 53); 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/direct/fixed/encode/8b".to_owned()));
assert!(codec_names.contains(&"codec/registry/json-nested/receive/65536b".to_owned())); assert!(codec_names.contains(&"codec/registry/json-nested/receive/65536b".to_owned()));
assert!( assert!(
@ -42,6 +42,21 @@ mod tests {
assert!( assert!(
telemetry_names.contains(&"telemetry/mux/concurrent/producers-8/total-4096".to_owned()) 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; let mut names = codec_names;
names.extend(telemetry_names); names.extend(telemetry_names);
@ -52,6 +67,58 @@ mod tests {
assert!(names.iter().all(|name| !name.contains(char::is_whitespace))); 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] #[test]
fn benchmark_commands_name_the_package() { fn benchmark_commands_name_the_package() {
assert!(BENCHMARK_COMMAND.contains("-p swactor-benchmarks")); assert!(BENCHMARK_COMMAND.contains("-p swactor-benchmarks"));

View file

@ -1,14 +1,23 @@
use std::collections::HashSet; use std::collections::HashSet;
use std::sync::{Arc, Barrier}; use std::sync::{Arc, Barrier, mpsc};
use std::thread::{self, JoinHandle}; 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::ingest::Consumer;
use telemetry::transport::Delivery; use telemetry::transport::Delivery;
use telemetry::wire::{decode_delivery, encode_delivery}; use telemetry::wire::{decode_delivery, encode_delivery};
use telemetry::{ use telemetry::{
ChannelContent, ChannelId, Lifetime, Mux, Position, StreamId, TelemetryEndpoint, ChannelContent, ChannelId, Lifetime, Mux, NodeId, Position, StreamDescriptor, StreamId,
TelemetrySubscription, StreamOrigin, TelemetryEndpoint, TelemetrySubscription,
}; };
use crate::{SetupPolicy, WorkUnits, Workload}; use crate::{SetupPolicy, WorkUnits, Workload};
@ -18,6 +27,12 @@ const BATCH_SIZES: [usize; 3] = [1, 64, 1024];
const STORE_BATCH: usize = 1024; const STORE_BATCH: usize = 1024;
const CONCURRENT_TOTAL: usize = 4096; const CONCURRENT_TOTAL: usize = 4096;
const CONCURRENT_PAYLOAD_SIZE: usize = 32; 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)] #[derive(Clone, Copy, Debug)]
pub enum StoreKind { 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)] #[derive(Clone, Copy, Debug)]
pub enum TelemetryWorkload { pub enum TelemetryWorkload {
Mux { Mux {
@ -71,9 +135,16 @@ pub enum TelemetryWorkload {
operation: WireOperation, operation: WireOperation,
payload_size: usize, payload_size: usize,
}, },
WireBatch,
Concurrent { Concurrent {
producers: usize, producers: usize,
}, },
EngineConcurrent,
PayloadCodec {
codec: PayloadCodec,
operation: PayloadCodecOperation,
},
Compression(CompressionCodec),
} }
pub struct MuxState { pub struct MuxState {
@ -105,6 +176,12 @@ pub struct WireState {
encoded: Vec<u8>, encoded: Vec<u8>,
} }
pub struct WireBatchState {
descriptor: StreamDescriptor,
events: Vec<TelemetryEvent>,
payload_bytes: usize,
}
pub struct ConcurrentState { pub struct ConcurrentState {
endpoint: TelemetryEndpoint, endpoint: TelemetryEndpoint,
start: Arc<Barrier>, start: Arc<Barrier>,
@ -112,12 +189,52 @@ pub struct ConcurrentState {
total: usize, total: usize,
} }
pub struct EngineConcurrentState {
_engine: Engine,
handle: EngineHandle,
endpoint: Arc<TelemetryEndpoint>,
subscriptions: Vec<TelemetrySubscription>,
submissions: Vec<Vec<(ChannelId, Vec<u8>)>>,
}
#[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<MyelinBenchmarkRecord>,
encoded: Vec<Vec<u8>>,
}
pub struct CompressionState {
codec: CompressionCodec,
chunks: Vec<Vec<u8>>,
raw_bytes: usize,
}
pub enum TelemetryState { pub enum TelemetryState {
Mux(MuxState), Mux(MuxState),
Fanout(FanoutState), Fanout(FanoutState),
Store(StoreState), Store(StoreState),
Wire(WireState), Wire(WireState),
WireBatch(WireBatchState),
Concurrent(ConcurrentState), Concurrent(ConcurrentState),
EngineConcurrent(EngineConcurrentState),
PayloadCodec(PayloadCodecState),
Compression(CompressionState),
} }
pub enum TelemetryOutput { pub enum TelemetryOutput {
@ -143,11 +260,22 @@ pub enum TelemetryOutput {
}, },
WireEncoded(Vec<u8>), WireEncoded(Vec<u8>),
WireDecoded(StreamId, Frame), WireDecoded(StreamId, Frame),
WireBatch(Vec<u8>),
Concurrent { Concurrent {
accepted: usize, accepted: usize,
frames: Vec<Frame>, frames: Vec<Frame>,
dropped: u64, dropped: u64,
}, },
EngineConcurrent {
accepted: usize,
drained: usize,
delivered: usize,
events: Vec<Vec<TelemetryEvent>>,
dropped: u64,
},
PayloadEncoded(Vec<Vec<u8>>),
PayloadDecoded(Vec<MyelinBenchmarkRecord>),
Compressed(Vec<Vec<u8>>),
} }
fn deterministic_payload(size: usize) -> Vec<u8> { fn deterministic_payload(size: usize) -> Vec<u8> {
@ -163,6 +291,54 @@ fn marked_payload(size: usize, marker: u64) -> Vec<u8> {
payload 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<u8> {
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<u8> {
encode_payload(codec, &myelin_record(sequence, channel))
}
fn marker(payload: &[u8]) -> u64 { fn marker(payload: &[u8]) -> u64 {
u64::from_le_bytes(payload[..8].try_into().unwrap()) 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::<Vec<_>>();
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::<Vec<_>>();
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::<Vec<_>>();
let encoded = records
.iter()
.map(|record| encode_payload(codec, record))
.collect();
PayloadCodecState {
codec,
records,
encoded,
}
}
fn execute_mux(state: &mut MuxState) -> TelemetryOutput { fn execute_mux(state: &mut MuxState) -> TelemetryOutput {
if state.reject_full { if state.reject_full {
let rejected = (0..state.batch) 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::<usize>(),
)
})
}
impl Workload for TelemetryWorkload { impl Workload for TelemetryWorkload {
type State = TelemetryState; type State = TelemetryState;
type Output = TelemetryOutput; type Output = TelemetryOutput;
@ -458,9 +925,24 @@ impl Workload for TelemetryWorkload {
operation, operation,
payload_size, payload_size,
} => format!("telemetry/wire/{}/{}b", operation.label(), 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 } => { Self::Concurrent { producers } => {
format!("telemetry/mux/concurrent/producers-{producers}/total-{CONCURRENT_TOTAL}") 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::Store(kind) => TelemetryState::Store(setup_store(*kind)),
Self::Wire { payload_size, .. } => TelemetryState::Wire(setup_wire(*payload_size)), Self::Wire { payload_size, .. } => TelemetryState::Wire(setup_wire(*payload_size)),
Self::WireBatch => TelemetryState::WireBatch(setup_wire_batch()),
Self::Concurrent { producers } => { Self::Concurrent { producers } => {
TelemetryState::Concurrent(setup_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)) => { (Self::Wire { operation, .. }, TelemetryState::Wire(state)) => {
execute_wire(state, *operation) execute_wire(state, *operation)
} }
(Self::WireBatch, TelemetryState::WireBatch(state)) => execute_wire_batch(state),
(Self::Concurrent { .. }, TelemetryState::Concurrent(state)) => { (Self::Concurrent { .. }, TelemetryState::Concurrent(state)) => {
execute_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"), _ => panic!("telemetry workload and state mismatch"),
} }
} }
@ -649,6 +1147,81 @@ impl Workload for TelemetryWorkload {
state.total * CONCURRENT_PAYLOAD_SIZE 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::<Vec<_>>();
let positions = frames
.iter()
.map(|frame| frame.position.0)
.collect::<Vec<_>>();
let payloads = frames
.iter()
.map(|frame| frame.payload.as_slice())
.collect::<HashSet<_>>();
let channel_ids = frames
.iter()
.map(|frame| frame.channel.channel)
.collect::<HashSet<_>>();
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::<Vec<_>>();
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"), _ => panic!("telemetry workload, state, and output mismatch"),
} }
} }
@ -685,12 +1258,18 @@ impl Workload for TelemetryWorkload {
STORE_BATCH as u64 STORE_BATCH as u64
}), }),
Self::Wire { payload_size, .. } => WorkUnits::Bytes(*payload_size 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::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 { fn setup_policy(&self) -> SetupPolicy {
if matches!(self, Self::Concurrent { .. }) { if matches!(self, Self::Concurrent { .. } | Self::EngineConcurrent) {
SetupPolicy::PerExecution SetupPolicy::PerExecution
} else { } else {
SetupPolicy::Batched SetupPolicy::Batched
@ -747,9 +1326,24 @@ pub fn workloads() -> Vec<TelemetryWorkload> {
payload_size, 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] { for producers in [1, 2, 4, 8] {
workloads.push(TelemetryWorkload::Concurrent { producers }); workloads.push(TelemetryWorkload::Concurrent { producers });
} }
workloads.push(TelemetryWorkload::EngineConcurrent);
workloads workloads
} }