From be527359f386b49a962e1f37c7e11816fd345e26 Mon Sep 17 00:00:00 2001 From: nxq64494 Date: Tue, 15 Sep 2026 14:37:11 +0100 Subject: [PATCH 01/18] Args for veto_probability, enabled_vetoes and veto_names arrays. Does not send vc00 yet. --- src/howl.rs | 51 ++++++++++++++++++++++++++++++++++++++++++++++++--- src/main.rs | 16 +++++++++++++--- 2 files changed, 61 insertions(+), 6 deletions(-) diff --git a/src/howl.rs b/src/howl.rs index bfdb1df..5e38eb4 100644 --- a/src/howl.rs +++ b/src/howl.rs @@ -9,6 +9,7 @@ use isis_streaming_data_types::flatbuffers_generated::events_ev44::{ use isis_streaming_data_types::flatbuffers_generated::pulse_metadata_pu00::{ Pu00Message, Pu00MessageArgs, finish_pu_00_message_buffer, }; +// use isis_streaming_data_types::flatbuffers_generated::veto_configuration_vc00::{}; use isis_streaming_data_types::flatbuffers_generated::run_start_pl72::{ RunStart, RunStartArgs, SpectraDetectorMapping, SpectraDetectorMappingArgs, finish_run_start_buffer, @@ -114,6 +115,31 @@ fn generate_run_stop<'a>(fbb: &'a mut FlatBufferBuilder<'_>, job_id: &str) -> &' fbb.finished_data() } +fn get_veto_probability(conf: &HowlConfig, frame: u32) -> f64 { + conf.veto_probability + .get(frame as usize) + .map(|&prob| match conf.enabled_vetoes.get(frame as usize) { + Some(&enabled) => { + if enabled { + prob + } else { + 0.0 + } + } + None => prob, + }) + .unwrap_or(0.0) +} + +fn get_veto_name(conf: &HowlConfig, frame: u32, buf: &mut String) { + buf.clear(); + buf.push_str( + conf.veto_names + .get(frame as usize) + .unwrap_or(&("saluki_veto_".to_string() + &frame.to_string())), + ) +} + fn produce_messages( producer: &ThreadedProducer, fbb: &mut FlatBufferBuilder, @@ -121,6 +147,7 @@ fn produce_messages( frame: u32, conf: &HowlConfig, current_job_id: &mut String, + veto_name: &mut String, ) { // get current time let now_nanos = SystemTime::now() @@ -130,6 +157,10 @@ fn produce_messages( .try_into() .expect("This will fail after April 11th, 2262"); + let veto_probability = get_veto_probability(conf, frame); + + get_veto_name(conf, 0, veto_name); + match producer.send( BaseRecord::to(conf.event_topic) .key("") @@ -137,7 +168,7 @@ fn produce_messages( rng, fbb, now_nanos, - conf.veto_probability, + veto_probability, )) .timestamp(now_nanos / 1_000_000), ) { @@ -200,6 +231,7 @@ fn produce_messages( error!("Failed to send run start: {}", err.0); } } + // match producer send new vc00 } } @@ -250,6 +282,7 @@ fn generate_fake_metadata<'a>( veto_probability: f64, ) -> &'a [u8] { fbb.reset(); + let is_vetoed = rng.random_range(0.0..1.0) < veto_probability; let args = Pu00MessageArgs { reference_time: timestamp_ns, @@ -271,7 +304,9 @@ pub struct HowlConfig<'a> { pub messages_per_frame: u32, pub frames_per_second: u32, pub frames_per_run: u32, - pub veto_probability: f64, // 1 = always vetoed, 0 = never vetoed + pub veto_probability: Vec, + pub enabled_vetoes: Vec, + pub veto_names: Vec, pub event_message_config: &'a EventMessageConfig, pub fast: bool, pub kafka_config: Option>, @@ -294,8 +329,13 @@ pub fn howl(conf: &HowlConfig) { as u32; debug!("ev44 size is {ev44_size} bytes"); + let veto_probability = get_veto_probability(conf, 0); + + let mut veto_name = String::new(); + get_veto_name(conf, 0, &mut veto_name); + let pu00_size = - generate_fake_metadata(&mut rng, &mut fbb, now_nanos, conf.veto_probability).len() as u32; + generate_fake_metadata(&mut rng, &mut fbb, now_nanos, veto_probability).len() as u32; debug!("pu00 size is {pu00_size} bytes"); // calculate overall rate (with both ev44 and pu00) @@ -344,6 +384,8 @@ pub fn howl(conf: &HowlConfig) { ) .expect("Failed to enqueue run start message"); + // match producer send new vc00 + let target_frame_time = Duration::from_secs_f64(1.0 / conf.frames_per_second as f64); debug!("Target frame time: {target_frame_time:?}"); @@ -354,6 +396,8 @@ pub fn howl(conf: &HowlConfig) { .expect("Failed to get system time"); debug!("Target time: {target_time:?}"); + let mut veto_name = String::new(); + loop { target_time += target_frame_time; debug!("New target: {target_time:?}"); @@ -366,6 +410,7 @@ pub fn howl(conf: &HowlConfig) { frames, conf, &mut current_job_id, + &mut veto_name, ); let now = SystemTime::now() .duration_since(SystemTime::UNIX_EPOCH) diff --git a/src/main.rs b/src/main.rs index 0ac9c45..1c8f960 100644 --- a/src/main.rs +++ b/src/main.rs @@ -95,9 +95,15 @@ enum Commands { /// Maximum detector ID #[arg(long, default_value = "1000")] det_max: i32, - /// Veto probability (0 = never vetoed; 1 = always vetoed) - #[arg(long, default_value = "0.0")] - veto_probability: f64, + /// Veto probabilities + #[arg(long, default_value = "[0.0]")] + veto_probability: Vec, + /// Enabled vetoes + #[arg(long, default_value = "false")] + enabled_vetoes: Vec, + /// Veto names + #[arg(long, default_value = "")] + veto_names: Vec, /// Enable howl fast mode (Disables randomised ev44 blob generation) #[arg(long, action=clap::ArgAction::SetTrue)] fast: bool, @@ -167,6 +173,8 @@ async fn main() { det_min, det_max, veto_probability, + enabled_vetoes, + veto_names, fast, kafka_config, } => howl(&HowlConfig { @@ -185,6 +193,8 @@ async fn main() { det_max, }, veto_probability, + enabled_vetoes, + veto_names, fast, }), Commands::Count { From 844fef8b4344bb850089e99ea2a6f693fc17d31f Mon Sep 17 00:00:00 2001 From: nxq64494 Date: Wed, 16 Sep 2026 16:32:10 +0100 Subject: [PATCH 02/18] educated guess as to how vetoes/masking works --- src/howl.rs | 92 +++++++++++++++++++++++++++++++---------------------- 1 file changed, 54 insertions(+), 38 deletions(-) diff --git a/src/howl.rs b/src/howl.rs index 5e38eb4..91ee365 100644 --- a/src/howl.rs +++ b/src/howl.rs @@ -9,7 +9,10 @@ use isis_streaming_data_types::flatbuffers_generated::events_ev44::{ use isis_streaming_data_types::flatbuffers_generated::pulse_metadata_pu00::{ Pu00Message, Pu00MessageArgs, finish_pu_00_message_buffer, }; -// use isis_streaming_data_types::flatbuffers_generated::veto_configuration_vc00::{}; +use isis_streaming_data_types::flatbuffers_generated::veto_configuration_vc00::{ + Vetoes, VetoesArgs, finish_vetoes_buffer, +}; + use isis_streaming_data_types::flatbuffers_generated::run_start_pl72::{ RunStart, RunStartArgs, SpectraDetectorMapping, SpectraDetectorMappingArgs, finish_run_start_buffer, @@ -26,6 +29,8 @@ use rdkafka::producer::{BaseRecord, DefaultProducerContext, ThreadedProducer}; use serde_json::json; use uuid::Uuid; +const VETO_COUNT: i32 = 32; + fn generate_run_start<'a>( fbb: &'a mut FlatBufferBuilder<'_>, det_max: i32, @@ -115,7 +120,7 @@ fn generate_run_stop<'a>(fbb: &'a mut FlatBufferBuilder<'_>, job_id: &str) -> &' fbb.finished_data() } -fn get_veto_probability(conf: &HowlConfig, frame: u32) -> f64 { +fn get_veto_probability(conf: &HowlConfig, frame: i32) -> f64 { conf.veto_probability .get(frame as usize) .map(|&prob| match conf.enabled_vetoes.get(frame as usize) { @@ -131,23 +136,38 @@ fn get_veto_probability(conf: &HowlConfig, frame: u32) -> f64 { .unwrap_or(0.0) } -fn get_veto_name(conf: &HowlConfig, frame: u32, buf: &mut String) { +fn get_veto_names(conf: &HowlConfig, buf: &mut Vec) { buf.clear(); - buf.push_str( - conf.veto_names - .get(frame as usize) - .unwrap_or(&("saluki_veto_".to_string() + &frame.to_string())), - ) + + for i in 0..VETO_COUNT { + let name = conf + .veto_names + .get(i as usize) + .cloned() + .unwrap_or_else(|| format!("saluki_veto_{i}")); + + buf.push(name.to_string()); + } +} + +fn get_vetoes(conf: &HowlConfig, rng: &mut ThreadRng, vetoes: &mut Vec) { + vetoes.clear(); + for i in 0..VETO_COUNT { + vetoes.push(rng.random_range(0.0..1.0) < get_veto_probability(conf, i)); + } +} + +fn get_vetoes_mask(vetoes: &Vec) -> bool { + vetoes.iter().all(|&b| b == vetoes[0]) } fn produce_messages( producer: &ThreadedProducer, fbb: &mut FlatBufferBuilder, rng: &mut ThreadRng, - frame: u32, + frame: i64, conf: &HowlConfig, current_job_id: &mut String, - veto_name: &mut String, ) { // get current time let now_nanos = SystemTime::now() @@ -157,19 +177,10 @@ fn produce_messages( .try_into() .expect("This will fail after April 11th, 2262"); - let veto_probability = get_veto_probability(conf, frame); - - get_veto_name(conf, 0, veto_name); - match producer.send( BaseRecord::to(conf.event_topic) .key("") - .payload(generate_fake_metadata( - rng, - fbb, - now_nanos, - veto_probability, - )) + .payload(generate_fake_metadata(conf, rng, fbb, now_nanos)) .timestamp(now_nanos / 1_000_000), ) { Ok(_) => {} @@ -198,7 +209,7 @@ fn produce_messages( } } - if conf.frames_per_run > 0 && frame.is_multiple_of(conf.frames_per_run) { + if conf.frames_per_run > 0 && (frame as u32).is_multiple_of(conf.frames_per_run) { info!( "Starting new run after {} simulated frames", conf.frames_per_run @@ -231,7 +242,6 @@ fn produce_messages( error!("Failed to send run start: {}", err.0); } } - // match producer send new vc00 } } @@ -246,7 +256,7 @@ pub struct EventMessageConfig { fn generate_fake_events<'a>( fbb: &'a mut FlatBufferBuilder<'_>, rng: &mut ThreadRng, - msg_id: u32, + msg_id: i64, conf: &EventMessageConfig, timestamp_ns: i64, ) -> &'a [u8] { @@ -264,7 +274,7 @@ fn generate_fake_events<'a>( let args = Event44MessageArgs { source_name: Some(fbb.create_string("saluki")), - message_id: msg_id as i64, + message_id: msg_id, reference_time: Some(fbb.create_vector(&[timestamp_ns])), reference_time_index: Some(fbb.create_vector(&[0])), time_of_flight: Some(fbb.create_vector(&tofs)), @@ -276,24 +286,39 @@ fn generate_fake_events<'a>( } fn generate_fake_metadata<'a>( + conf: &HowlConfig, rng: &mut ThreadRng, fbb: &'a mut FlatBufferBuilder<'_>, timestamp_ns: i64, - veto_probability: f64, ) -> &'a [u8] { fbb.reset(); - let is_vetoed = rng.random_range(0.0..1.0) < veto_probability; + let mut vetoes = Vec::new(); + get_vetoes(&conf, rng, &mut vetoes); + let vetoes_mask = if get_vetoes_mask(&vetoes) { 1 } else { 0 }; + let args = Pu00MessageArgs { reference_time: timestamp_ns, message_id: 0, source_name: Some(fbb.create_string("saluki")), period_number: Some(0), - vetos: Some(if is_vetoed { 1 } else { 0 }), + vetos: Some(vetoes_mask), proton_charge: Some(0.1), }; let pu00 = Pu00Message::create(fbb, &args); finish_pu_00_message_buffer(fbb, pu00); + + let mut veto_names: Vec = Vec::new(); + get_veto_names(&conf, &mut veto_names); + + let args = VetoesArgs { + timestamp: timestamp_ns, + vetoes: vetoes_mask, + // names: veto_names, + }; + let vc00 = Vetoes::create(fbb, &args); + finish_vetoes_buffer(fbb, vc00); + fbb.finished_data() } @@ -329,13 +354,7 @@ pub fn howl(conf: &HowlConfig) { as u32; debug!("ev44 size is {ev44_size} bytes"); - let veto_probability = get_veto_probability(conf, 0); - - let mut veto_name = String::new(); - get_veto_name(conf, 0, &mut veto_name); - - let pu00_size = - generate_fake_metadata(&mut rng, &mut fbb, now_nanos, veto_probability).len() as u32; + let pu00_size = generate_fake_metadata(&conf, &mut rng, &mut fbb, now_nanos).len() as u32; debug!("pu00 size is {pu00_size} bytes"); // calculate overall rate (with both ev44 and pu00) @@ -389,15 +408,13 @@ pub fn howl(conf: &HowlConfig) { let target_frame_time = Duration::from_secs_f64(1.0 / conf.frames_per_second as f64); debug!("Target frame time: {target_frame_time:?}"); - let mut frames: u32 = 0; + let mut frames: i64 = 0; let mut target_time = SystemTime::now() .duration_since(SystemTime::UNIX_EPOCH) .expect("Failed to get system time"); debug!("Target time: {target_time:?}"); - let mut veto_name = String::new(); - loop { target_time += target_frame_time; debug!("New target: {target_time:?}"); @@ -410,7 +427,6 @@ pub fn howl(conf: &HowlConfig) { frames, conf, &mut current_job_id, - &mut veto_name, ); let now = SystemTime::now() .duration_since(SystemTime::UNIX_EPOCH) From 37f960c8161aad138d13e4bb8061265d4a311d89 Mon Sep 17 00:00:00 2001 From: nxq64494 Date: Wed, 16 Sep 2026 16:36:10 +0100 Subject: [PATCH 03/18] make clippy happy --- src/howl.rs | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/src/howl.rs b/src/howl.rs index 91ee365..d058812 100644 --- a/src/howl.rs +++ b/src/howl.rs @@ -157,7 +157,7 @@ fn get_vetoes(conf: &HowlConfig, rng: &mut ThreadRng, vetoes: &mut Vec) { } } -fn get_vetoes_mask(vetoes: &Vec) -> bool { +fn get_vetoes_mask(vetoes: &[bool]) -> bool { vetoes.iter().all(|&b| b == vetoes[0]) } @@ -294,7 +294,7 @@ fn generate_fake_metadata<'a>( fbb.reset(); let mut vetoes = Vec::new(); - get_vetoes(&conf, rng, &mut vetoes); + get_vetoes(conf, rng, &mut vetoes); let vetoes_mask = if get_vetoes_mask(&vetoes) { 1 } else { 0 }; let args = Pu00MessageArgs { @@ -309,7 +309,7 @@ fn generate_fake_metadata<'a>( finish_pu_00_message_buffer(fbb, pu00); let mut veto_names: Vec = Vec::new(); - get_veto_names(&conf, &mut veto_names); + get_veto_names(conf, &mut veto_names); let args = VetoesArgs { timestamp: timestamp_ns, @@ -354,7 +354,7 @@ pub fn howl(conf: &HowlConfig) { as u32; debug!("ev44 size is {ev44_size} bytes"); - let pu00_size = generate_fake_metadata(&conf, &mut rng, &mut fbb, now_nanos).len() as u32; + let pu00_size = generate_fake_metadata(conf, &mut rng, &mut fbb, now_nanos).len() as u32; debug!("pu00 size is {pu00_size} bytes"); // calculate overall rate (with both ev44 and pu00) From decb820be4c9ad2d47a562e51893cdb006986268 Mon Sep 17 00:00:00 2001 From: nxq64494 Date: Thu, 17 Sep 2026 10:22:48 +0100 Subject: [PATCH 04/18] sends vc00 (untested) --- Cargo.lock | 4 ++-- Cargo.toml | 2 +- src/howl.rs | 66 +++++++++++++++++++++++++++++++++++++---------------- src/main.rs | 1 + 4 files changed, 50 insertions(+), 23 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 102295d..95375d5 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -441,9 +441,9 @@ checksum = "a6cb138bb79a146c1bd460005623e142ef0181e3d0219cb493e02f7d08a35695" [[package]] name = "isis_streaming_data_types" -version = "0.1.2" +version = "0.1.3" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "4a40f55610cfa45cc165129d63849bd6a805a918c1e1f64c7de58cba09d181f9" +checksum = "7189171301072e5e8c6541db5f7f086a45045a20c190cf946a208f8cc4b8f142" dependencies = [ "flatbuffers", ] diff --git a/Cargo.toml b/Cargo.toml index 7c5cdfc..3ddd67b 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -10,7 +10,7 @@ log = "0.4.29" uuid = { version = "1.22.0", features = ["v4"] } env_logger = "0.11.9" clap-verbosity-flag = "3.0.4" -isis_streaming_data_types = "0.1.1" +isis_streaming_data_types = "0.1.3" rand = { version = "0.10.1", features = ["thread_rng"]} rand_distr = "0.6.0" flatbuffers = "25.12.19" diff --git a/src/howl.rs b/src/howl.rs index d058812..5562489 100644 --- a/src/howl.rs +++ b/src/howl.rs @@ -2,7 +2,7 @@ use crate::KafkaOption; use std::thread; use std::time::{Duration, SystemTime}; -use flatbuffers::FlatBufferBuilder; +use flatbuffers::{FlatBufferBuilder, WIPOffset}; use isis_streaming_data_types::flatbuffers_generated::events_ev44::{ Event44Message, Event44MessageArgs, finish_event_44_message_buffer, }; @@ -136,17 +136,20 @@ fn get_veto_probability(conf: &HowlConfig, frame: i32) -> f64 { .unwrap_or(0.0) } -fn get_veto_names(conf: &HowlConfig, buf: &mut Vec) { +fn get_veto_names_fbb<'a>( + veto_names: &Vec, + fbb: &mut FlatBufferBuilder<'a>, + buf: &mut Vec>, +) { buf.clear(); for i in 0..VETO_COUNT { - let name = conf - .veto_names + let name = veto_names .get(i as usize) .cloned() .unwrap_or_else(|| format!("saluki_veto_{i}")); - buf.push(name.to_string()); + buf.push(fbb.create_string(&name.to_string())); } } @@ -168,6 +171,7 @@ fn produce_messages( frame: i64, conf: &HowlConfig, current_job_id: &mut String, + vetoes_mask: &u32, ) { // get current time let now_nanos = SystemTime::now() @@ -180,7 +184,7 @@ fn produce_messages( match producer.send( BaseRecord::to(conf.event_topic) .key("") - .payload(generate_fake_metadata(conf, rng, fbb, now_nanos)) + .payload(generate_fake_metadata(vetoes_mask, fbb, now_nanos)) .timestamp(now_nanos / 1_000_000), ) { Ok(_) => {} @@ -286,35 +290,39 @@ fn generate_fake_events<'a>( } fn generate_fake_metadata<'a>( - conf: &HowlConfig, - rng: &mut ThreadRng, + vetoes_mask: &u32, fbb: &'a mut FlatBufferBuilder<'_>, timestamp_ns: i64, ) -> &'a [u8] { fbb.reset(); - let mut vetoes = Vec::new(); - get_vetoes(conf, rng, &mut vetoes); - let vetoes_mask = if get_vetoes_mask(&vetoes) { 1 } else { 0 }; - let args = Pu00MessageArgs { reference_time: timestamp_ns, message_id: 0, source_name: Some(fbb.create_string("saluki")), period_number: Some(0), - vetos: Some(vetoes_mask), + vetos: Some(*vetoes_mask), proton_charge: Some(0.1), }; let pu00 = Pu00Message::create(fbb, &args); finish_pu_00_message_buffer(fbb, pu00); - let mut veto_names: Vec = Vec::new(); - get_veto_names(conf, &mut veto_names); + fbb.finished_data() +} + +fn generate_veto_config<'a>( + veto_names: &Vec, + fbb: &'a mut FlatBufferBuilder<'_>, + timestamp_ns: i64, + vetoes_mask: &u32, +) -> &'a [u8] { + let mut veto_names_fbb: Vec> = Vec::new(); + get_veto_names_fbb(veto_names, fbb, &mut veto_names_fbb); let args = VetoesArgs { timestamp: timestamp_ns, - vetoes: vetoes_mask, - // names: veto_names, + vetoes: *vetoes_mask, + veto_names: Some(fbb.create_vector(&veto_names_fbb)), }; let vc00 = Vetoes::create(fbb, &args); finish_vetoes_buffer(fbb, vc00); @@ -326,6 +334,7 @@ pub struct HowlConfig<'a> { pub broker: &'a str, pub event_topic: &'a str, pub run_info_topic: &'a str, + pub veto_config_topic: &'a str, pub messages_per_frame: u32, pub frames_per_second: u32, pub frames_per_run: u32, @@ -342,6 +351,10 @@ pub fn howl(conf: &HowlConfig) { let mut fbb = FlatBufferBuilder::new(); let mut rng = rand::rng(); + let mut vetoes = Vec::new(); + get_vetoes(conf, &mut rng, &mut vetoes); + let vetoes_mask = if get_vetoes_mask(&vetoes) { 1 } else { 0 }; + let now_nanos = SystemTime::now() .duration_since(SystemTime::UNIX_EPOCH) .expect("Failed to get system time") @@ -354,7 +367,7 @@ pub fn howl(conf: &HowlConfig) { as u32; debug!("ev44 size is {ev44_size} bytes"); - let pu00_size = generate_fake_metadata(conf, &mut rng, &mut fbb, now_nanos).len() as u32; + let pu00_size = generate_fake_metadata(&vetoes_mask, &mut fbb, now_nanos).len() as u32; debug!("pu00 size is {pu00_size} bytes"); // calculate overall rate (with both ev44 and pu00) @@ -403,8 +416,6 @@ pub fn howl(conf: &HowlConfig) { ) .expect("Failed to enqueue run start message"); - // match producer send new vc00 - let target_frame_time = Duration::from_secs_f64(1.0 / conf.frames_per_second as f64); debug!("Target frame time: {target_frame_time:?}"); @@ -415,6 +426,20 @@ pub fn howl(conf: &HowlConfig) { .expect("Failed to get system time"); debug!("Target time: {target_time:?}"); + producer + .send( + BaseRecord::to(conf.veto_config_topic) + .key("") + .payload(generate_veto_config( + &conf.veto_names, + &mut fbb, + now_nanos, + &vetoes_mask, + )) + .timestamp(now_nanos / 1_000_000), + ) + .expect("Failed to enqueue run veto configuration message"); + loop { target_time += target_frame_time; debug!("New target: {target_time:?}"); @@ -427,6 +452,7 @@ pub fn howl(conf: &HowlConfig) { frames, conf, &mut current_job_id, + &vetoes_mask, ); let now = SystemTime::now() .duration_since(SystemTime::UNIX_EPOCH) diff --git a/src/main.rs b/src/main.rs index 1c8f960..2d2f4ec 100644 --- a/src/main.rs +++ b/src/main.rs @@ -182,6 +182,7 @@ async fn main() { broker: &broker, event_topic: &format!("{topic_prefix}_rawEvents"), run_info_topic: &format!("{topic_prefix}_runInfo"), + veto_config_topic: &format!("{topic_prefix}_vetoConfig"), messages_per_frame, frames_per_second, frames_per_run, From ae2ed22ae560502916afbecadf972330e024d828 Mon Sep 17 00:00:00 2001 From: nxq64494 Date: Thu, 17 Sep 2026 10:24:28 +0100 Subject: [PATCH 05/18] make clippy happy --- src/howl.rs | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/src/howl.rs b/src/howl.rs index 5562489..f9da42a 100644 --- a/src/howl.rs +++ b/src/howl.rs @@ -137,7 +137,7 @@ fn get_veto_probability(conf: &HowlConfig, frame: i32) -> f64 { } fn get_veto_names_fbb<'a>( - veto_names: &Vec, + veto_names: &[String], fbb: &mut FlatBufferBuilder<'a>, buf: &mut Vec>, ) { @@ -311,7 +311,7 @@ fn generate_fake_metadata<'a>( } fn generate_veto_config<'a>( - veto_names: &Vec, + veto_names: &[String], fbb: &'a mut FlatBufferBuilder<'_>, timestamp_ns: i64, vetoes_mask: &u32, From 66825643b03b9fa0846df0f7f51ba5882b1a78d6 Mon Sep 17 00:00:00 2001 From: nxq64494 Date: Thu, 17 Sep 2026 10:56:23 +0100 Subject: [PATCH 06/18] Shorten howl function --- src/howl.rs | 195 ++++++++++++++++++++++++++++++++-------------------- 1 file changed, 122 insertions(+), 73 deletions(-) diff --git a/src/howl.rs b/src/howl.rs index f9da42a..65ae6d4 100644 --- a/src/howl.rs +++ b/src/howl.rs @@ -330,47 +330,20 @@ fn generate_veto_config<'a>( fbb.finished_data() } -pub struct HowlConfig<'a> { - pub broker: &'a str, - pub event_topic: &'a str, - pub run_info_topic: &'a str, - pub veto_config_topic: &'a str, - pub messages_per_frame: u32, - pub frames_per_second: u32, - pub frames_per_run: u32, - pub veto_probability: Vec, - pub enabled_vetoes: Vec, - pub veto_names: Vec, - pub event_message_config: &'a EventMessageConfig, - pub fast: bool, - pub kafka_config: Option>, -} - -pub fn howl(conf: &HowlConfig) { - // create producer - let mut fbb = FlatBufferBuilder::new(); - let mut rng = rand::rng(); - - let mut vetoes = Vec::new(); - get_vetoes(conf, &mut rng, &mut vetoes); - let vetoes_mask = if get_vetoes_mask(&vetoes) { 1 } else { 0 }; - - let now_nanos = SystemTime::now() - .duration_since(SystemTime::UNIX_EPOCH) - .expect("Failed to get system time") - .as_nanos() - .try_into() - .expect("This will fail after April 11th, 2262"); - +fn calculate_data_rate( + fbb: &mut FlatBufferBuilder<'_>, + rng: &mut ThreadRng, + conf: &HowlConfig, + timestamp_ns: i64, + vetoes_mask: &u32, +) { let ev44_size = - generate_fake_events(&mut fbb, &mut rng, 0, conf.event_message_config, now_nanos).len() - as u32; + generate_fake_events(fbb, rng, 0, conf.event_message_config, timestamp_ns).len() as u32; debug!("ev44 size is {ev44_size} bytes"); - let pu00_size = generate_fake_metadata(&vetoes_mask, &mut fbb, now_nanos).len() as u32; + let pu00_size = generate_fake_metadata(vetoes_mask, fbb, timestamp_ns).len() as u32; debug!("pu00 size is {pu00_size} bytes"); - // calculate overall rate (with both ev44 and pu00) let rate_bytes_per_sec = ev44_size * conf.messages_per_frame * conf.frames_per_second + pu00_size * conf.frames_per_second; debug!("bytes per second: {rate_bytes_per_sec}"); @@ -383,62 +356,69 @@ pub fn howl(conf: &HowlConfig) { ); println!("Each pu00 is {pu00_size} bytes"); println!("Each ev44 is {ev44_size} bytes"); +} - let mut config: ClientConfig = ClientConfig::new(); - config.set("bootstrap.servers", conf.broker); - - if let Some(kafka_options) = &conf.kafka_config { - for option in kafka_options { - println!( - "Setting Kafka config option {}={}", - option.key, option.value - ); - config.set(&option.key, &option.value); - } - } - - let producer: ThreadedProducer = - config.create().expect("Producer creation error"); - - let mut current_job_id = Uuid::new_v4().to_string(); - +fn send_run_start( + producer: &mut ThreadedProducer, + fbb: &mut FlatBufferBuilder<'_>, + conf: &HowlConfig, + current_job_id: &str, + now_nanos: i64, +) { producer .send( BaseRecord::to(conf.run_info_topic) .key("") .payload(generate_run_start( - &mut fbb, + fbb, conf.event_message_config.det_max, conf.event_topic, - ¤t_job_id, + current_job_id, )) .timestamp(now_nanos / 1_000_000), ) .expect("Failed to enqueue run start message"); +} - let target_frame_time = Duration::from_secs_f64(1.0 / conf.frames_per_second as f64); - debug!("Target frame time: {target_frame_time:?}"); - - let mut frames: i64 = 0; - - let mut target_time = SystemTime::now() - .duration_since(SystemTime::UNIX_EPOCH) - .expect("Failed to get system time"); - debug!("Target time: {target_time:?}"); - +fn send_veto_config( + producer: &mut ThreadedProducer, + fbb: &mut FlatBufferBuilder<'_>, + conf: &HowlConfig, + vetoes_mask: &u32, + now_nanos: i64, +) { producer .send( BaseRecord::to(conf.veto_config_topic) .key("") .payload(generate_veto_config( &conf.veto_names, - &mut fbb, + fbb, now_nanos, - &vetoes_mask, + vetoes_mask, )) .timestamp(now_nanos / 1_000_000), ) .expect("Failed to enqueue run veto configuration message"); +} + +fn howl_begin( + producer: &mut ThreadedProducer, + fbb: &mut FlatBufferBuilder<'_>, + rng: &mut ThreadRng, + conf: &HowlConfig, + current_job_id: &mut String, + vetoes_mask: &u32, +) { + let target_frame_time = Duration::from_secs_f64(1.0 / conf.frames_per_second as f64); + debug!("Target frame time: {target_frame_time:?}"); + + let mut frames: i64 = 0; + + let mut target_time = SystemTime::now() + .duration_since(SystemTime::UNIX_EPOCH) + .expect("Failed to get system time"); + debug!("Target time: {target_time:?}"); loop { target_time += target_frame_time; @@ -446,13 +426,13 @@ pub fn howl(conf: &HowlConfig) { frames += 1; debug!("current job id: {current_job_id}"); produce_messages( - &producer, - &mut fbb, - &mut rng, + producer, + fbb, + rng, frames, conf, - &mut current_job_id, - &vetoes_mask, + current_job_id, + vetoes_mask, ); let now = SystemTime::now() .duration_since(SystemTime::UNIX_EPOCH) @@ -473,3 +453,72 @@ pub fn howl(conf: &HowlConfig) { } } } + +pub struct HowlConfig<'a> { + pub broker: &'a str, + pub event_topic: &'a str, + pub run_info_topic: &'a str, + pub veto_config_topic: &'a str, + pub messages_per_frame: u32, + pub frames_per_second: u32, + pub frames_per_run: u32, + pub veto_probability: Vec, + pub enabled_vetoes: Vec, + pub veto_names: Vec, + pub event_message_config: &'a EventMessageConfig, + pub fast: bool, + pub kafka_config: Option>, +} + +pub fn howl(conf: &HowlConfig) { + let mut fbb = FlatBufferBuilder::new(); + let mut rng = rand::rng(); + + let mut vetoes = Vec::new(); + get_vetoes(conf, &mut rng, &mut vetoes); + let vetoes_mask = if get_vetoes_mask(&vetoes) { 1 } else { 0 }; + + let now_nanos = SystemTime::now() + .duration_since(SystemTime::UNIX_EPOCH) + .expect("Failed to get system time") + .as_nanos() + .try_into() + .expect("This will fail after April 11th, 2262"); + + calculate_data_rate(&mut fbb, &mut rng, conf, now_nanos, &vetoes_mask); + + let mut config: ClientConfig = ClientConfig::new(); + config.set("bootstrap.servers", conf.broker); + + if let Some(kafka_options) = &conf.kafka_config { + for option in kafka_options { + println!( + "Setting Kafka config option {}={}", + option.key, option.value + ); + config.set(&option.key, &option.value); + } + } + + // create producer + let mut producer: ThreadedProducer = + config.create().expect("Producer creation error"); + + let mut current_job_id = Uuid::new_v4().to_string(); + + // send run start + send_run_start(&mut producer, &mut fbb, conf, ¤t_job_id, now_nanos); + + // send veto config + send_veto_config(&mut producer, &mut fbb, conf, &vetoes_mask, now_nanos); + + // start howling + howl_begin( + &mut producer, + &mut fbb, + &mut rng, + conf, + &mut current_job_id, + &vetoes_mask, + ); +} From b8716ad5c5bab2fb9a062ed149fb26dffecd8fcb Mon Sep 17 00:00:00 2001 From: nxq64494 Date: Thu, 17 Sep 2026 11:26:02 +0100 Subject: [PATCH 07/18] Change frames back to u32 --- src/howl.rs | 10 +++++----- 1 file changed, 5 insertions(+), 5 deletions(-) diff --git a/src/howl.rs b/src/howl.rs index 65ae6d4..cb490c0 100644 --- a/src/howl.rs +++ b/src/howl.rs @@ -168,7 +168,7 @@ fn produce_messages( producer: &ThreadedProducer, fbb: &mut FlatBufferBuilder, rng: &mut ThreadRng, - frame: i64, + frame: u32, conf: &HowlConfig, current_job_id: &mut String, vetoes_mask: &u32, @@ -213,7 +213,7 @@ fn produce_messages( } } - if conf.frames_per_run > 0 && (frame as u32).is_multiple_of(conf.frames_per_run) { + if conf.frames_per_run > 0 && frame.is_multiple_of(conf.frames_per_run) { info!( "Starting new run after {} simulated frames", conf.frames_per_run @@ -260,7 +260,7 @@ pub struct EventMessageConfig { fn generate_fake_events<'a>( fbb: &'a mut FlatBufferBuilder<'_>, rng: &mut ThreadRng, - msg_id: i64, + msg_id: u32, conf: &EventMessageConfig, timestamp_ns: i64, ) -> &'a [u8] { @@ -278,7 +278,7 @@ fn generate_fake_events<'a>( let args = Event44MessageArgs { source_name: Some(fbb.create_string("saluki")), - message_id: msg_id, + message_id: msg_id as i64, reference_time: Some(fbb.create_vector(&[timestamp_ns])), reference_time_index: Some(fbb.create_vector(&[0])), time_of_flight: Some(fbb.create_vector(&tofs)), @@ -413,7 +413,7 @@ fn howl_begin( let target_frame_time = Duration::from_secs_f64(1.0 / conf.frames_per_second as f64); debug!("Target frame time: {target_frame_time:?}"); - let mut frames: i64 = 0; + let mut frames: u32 = 0; let mut target_time = SystemTime::now() .duration_since(SystemTime::UNIX_EPOCH) From 795b8f6f2a6276435d8a1cc76b1b046eb235b984 Mon Sep 17 00:00:00 2001 From: nxq64494 Date: Thu, 17 Sep 2026 13:36:05 +0100 Subject: [PATCH 08/18] Fixed arguments + enabling/disabling vetoes --- src/howl.rs | 33 ++++++++++++++++++--------------- src/main.rs | 6 +++--- 2 files changed, 21 insertions(+), 18 deletions(-) diff --git a/src/howl.rs b/src/howl.rs index cb490c0..9587b7c 100644 --- a/src/howl.rs +++ b/src/howl.rs @@ -121,19 +121,17 @@ fn generate_run_stop<'a>(fbb: &'a mut FlatBufferBuilder<'_>, job_id: &str) -> &' } fn get_veto_probability(conf: &HowlConfig, frame: i32) -> f64 { - conf.veto_probability - .get(frame as usize) - .map(|&prob| match conf.enabled_vetoes.get(frame as usize) { - Some(&enabled) => { - if enabled { - prob - } else { - 0.0 - } - } - None => prob, - }) - .unwrap_or(0.0) + let idx = frame as usize; + + let enabled = conf.enabled_vetoes.get(idx).copied().unwrap_or(false); // Assume no veto if not found + let prob = conf.veto_probability.get(idx).copied().unwrap_or(0.0); // Assume 0% probability of veto if not found + + if enabled { + // If explicitly found to be enabled, 100% chance of veto + return 1.0; + } + + prob } fn get_veto_names_fbb<'a>( @@ -155,13 +153,17 @@ fn get_veto_names_fbb<'a>( fn get_vetoes(conf: &HowlConfig, rng: &mut ThreadRng, vetoes: &mut Vec) { vetoes.clear(); + + let mut veto: bool; for i in 0..VETO_COUNT { - vetoes.push(rng.random_range(0.0..1.0) < get_veto_probability(conf, i)); + veto = rng.random_range(0.0..1.0) < get_veto_probability(conf, i); + vetoes.push(veto); + println!("{}", veto); } } fn get_vetoes_mask(vetoes: &[bool]) -> bool { - vetoes.iter().all(|&b| b == vetoes[0]) + vetoes.iter().all(|&b| b == vetoes[0]) // are all entries equal to the first } fn produce_messages( @@ -316,6 +318,7 @@ fn generate_veto_config<'a>( timestamp_ns: i64, vetoes_mask: &u32, ) -> &'a [u8] { + fbb.reset(); let mut veto_names_fbb: Vec> = Vec::new(); get_veto_names_fbb(veto_names, fbb, &mut veto_names_fbb); diff --git a/src/main.rs b/src/main.rs index 2d2f4ec..de0146a 100644 --- a/src/main.rs +++ b/src/main.rs @@ -96,13 +96,13 @@ enum Commands { #[arg(long, default_value = "1000")] det_max: i32, /// Veto probabilities - #[arg(long, default_value = "[0.0]")] + #[arg(long, num_args = 0..33, value_delimiter = ' ')] veto_probability: Vec, /// Enabled vetoes - #[arg(long, default_value = "false")] + #[arg(long, num_args = 0..33, value_delimiter = ' ')] enabled_vetoes: Vec, /// Veto names - #[arg(long, default_value = "")] + #[arg(long, num_args = 0..33, value_delimiter = ' ')] veto_names: Vec, /// Enable howl fast mode (Disables randomised ev44 blob generation) #[arg(long, action=clap::ArgAction::SetTrue)] From 454aa19cc173c9cfddef459c33143fffb9e2c33b Mon Sep 17 00:00:00 2001 From: nxq64494 Date: Thu, 17 Sep 2026 14:49:40 +0100 Subject: [PATCH 09/18] Vetos behave correctly --- src/howl.rs | 52 ++++++++++++++++++++++------------------------------ 1 file changed, 22 insertions(+), 30 deletions(-) diff --git a/src/howl.rs b/src/howl.rs index 9587b7c..94905f9 100644 --- a/src/howl.rs +++ b/src/howl.rs @@ -120,20 +120,6 @@ fn generate_run_stop<'a>(fbb: &'a mut FlatBufferBuilder<'_>, job_id: &str) -> &' fbb.finished_data() } -fn get_veto_probability(conf: &HowlConfig, frame: i32) -> f64 { - let idx = frame as usize; - - let enabled = conf.enabled_vetoes.get(idx).copied().unwrap_or(false); // Assume no veto if not found - let prob = conf.veto_probability.get(idx).copied().unwrap_or(0.0); // Assume 0% probability of veto if not found - - if enabled { - // If explicitly found to be enabled, 100% chance of veto - return 1.0; - } - - prob -} - fn get_veto_names_fbb<'a>( veto_names: &[String], fbb: &mut FlatBufferBuilder<'a>, @@ -151,19 +137,23 @@ fn get_veto_names_fbb<'a>( } } -fn get_vetoes(conf: &HowlConfig, rng: &mut ThreadRng, vetoes: &mut Vec) { - vetoes.clear(); +fn get_enabled_vetoes(conf: &HowlConfig, rng: &mut ThreadRng, vetoes: &mut u32) { + let mut vtemp = *vetoes; - let mut veto: bool; for i in 0..VETO_COUNT { - veto = rng.random_range(0.0..1.0) < get_veto_probability(conf, i); - vetoes.push(veto); - println!("{}", veto); + let active = rng.random_bool(conf.veto_probability[i as usize]); + vtemp = (vtemp << 1) | active as u32; } + + *vetoes = vtemp; } -fn get_vetoes_mask(vetoes: &[bool]) -> bool { - vetoes.iter().all(|&b| b == vetoes[0]) // are all entries equal to the first +fn get_active_vetoes(conf: &HowlConfig, vetoes: &mut u32) { + let mut vtemp = *vetoes; + + for i in 0..VETO_COUNT { + vtemp = (vtemp << 1) | conf.enabled_vetoes[i as usize] as u32; // as 1 or 0 for each 32 bit + } } fn produce_messages( @@ -303,7 +293,7 @@ fn generate_fake_metadata<'a>( message_id: 0, source_name: Some(fbb.create_string("saluki")), period_number: Some(0), - vetos: Some(*vetoes_mask), + vetos: Some(*vetoes_mask), // active proton_charge: Some(0.1), }; let pu00 = Pu00Message::create(fbb, &args); @@ -324,7 +314,7 @@ fn generate_veto_config<'a>( let args = VetoesArgs { timestamp: timestamp_ns, - vetoes: *vetoes_mask, + vetoes: *vetoes_mask, // enable veto_names: Some(fbb.create_vector(&veto_names_fbb)), }; let vc00 = Vetoes::create(fbb, &args); @@ -477,9 +467,11 @@ pub fn howl(conf: &HowlConfig) { let mut fbb = FlatBufferBuilder::new(); let mut rng = rand::rng(); - let mut vetoes = Vec::new(); - get_vetoes(conf, &mut rng, &mut vetoes); - let vetoes_mask = if get_vetoes_mask(&vetoes) { 1 } else { 0 }; + let mut active_vetoes: u32 = 0; + let mut enabled_vetoes: u32 = 0; + + get_active_vetoes(conf, &mut active_vetoes); + get_enabled_vetoes(conf, &mut rng, &mut enabled_vetoes); let now_nanos = SystemTime::now() .duration_since(SystemTime::UNIX_EPOCH) @@ -488,7 +480,7 @@ pub fn howl(conf: &HowlConfig) { .try_into() .expect("This will fail after April 11th, 2262"); - calculate_data_rate(&mut fbb, &mut rng, conf, now_nanos, &vetoes_mask); + calculate_data_rate(&mut fbb, &mut rng, conf, now_nanos, &active_vetoes); let mut config: ClientConfig = ClientConfig::new(); config.set("bootstrap.servers", conf.broker); @@ -513,7 +505,7 @@ pub fn howl(conf: &HowlConfig) { send_run_start(&mut producer, &mut fbb, conf, ¤t_job_id, now_nanos); // send veto config - send_veto_config(&mut producer, &mut fbb, conf, &vetoes_mask, now_nanos); + send_veto_config(&mut producer, &mut fbb, conf, &enabled_vetoes, now_nanos); // start howling howl_begin( @@ -522,6 +514,6 @@ pub fn howl(conf: &HowlConfig) { &mut rng, conf, &mut current_job_id, - &vetoes_mask, + &active_vetoes, ); } From 25ef6133574b7dfab921c4216d3e4a7bcd7e7a23 Mon Sep 17 00:00:00 2001 From: nxq64494 Date: Thu, 17 Sep 2026 14:55:36 +0100 Subject: [PATCH 10/18] Change comments --- src/howl.rs | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/src/howl.rs b/src/howl.rs index 94905f9..e0b6bf0 100644 --- a/src/howl.rs +++ b/src/howl.rs @@ -152,7 +152,7 @@ fn get_active_vetoes(conf: &HowlConfig, vetoes: &mut u32) { let mut vtemp = *vetoes; for i in 0..VETO_COUNT { - vtemp = (vtemp << 1) | conf.enabled_vetoes[i as usize] as u32; // as 1 or 0 for each 32 bit + vtemp = (vtemp << 1) | conf.enabled_vetoes[i as usize] as u32; } } @@ -314,7 +314,7 @@ fn generate_veto_config<'a>( let args = VetoesArgs { timestamp: timestamp_ns, - vetoes: *vetoes_mask, // enable + vetoes: *vetoes_mask, // enabled veto_names: Some(fbb.create_vector(&veto_names_fbb)), }; let vc00 = Vetoes::create(fbb, &args); From cb1a304ba1002e77af3a5df81d464ffa016e0ee3 Mon Sep 17 00:00:00 2001 From: nxq64494 Date: Thu, 17 Sep 2026 15:31:12 +0100 Subject: [PATCH 11/18] Remove redundant comments --- src/howl.rs | 5 ----- 1 file changed, 5 deletions(-) diff --git a/src/howl.rs b/src/howl.rs index e0b6bf0..a92d633 100644 --- a/src/howl.rs +++ b/src/howl.rs @@ -501,13 +501,8 @@ pub fn howl(conf: &HowlConfig) { let mut current_job_id = Uuid::new_v4().to_string(); - // send run start send_run_start(&mut producer, &mut fbb, conf, ¤t_job_id, now_nanos); - - // send veto config send_veto_config(&mut producer, &mut fbb, conf, &enabled_vetoes, now_nanos); - - // start howling howl_begin( &mut producer, &mut fbb, From 43b6b67402d0301973ea23c63f30a163f852b22e Mon Sep 17 00:00:00 2001 From: nxq64494 Date: Fri, 18 Sep 2026 12:53:49 +0100 Subject: [PATCH 12/18] requested changes --- src/howl.rs | 29 ++++++++++++++--------------- 1 file changed, 14 insertions(+), 15 deletions(-) diff --git a/src/howl.rs b/src/howl.rs index a92d633..2b0ca10 100644 --- a/src/howl.rs +++ b/src/howl.rs @@ -29,7 +29,7 @@ use rdkafka::producer::{BaseRecord, DefaultProducerContext, ThreadedProducer}; use serde_json::json; use uuid::Uuid; -const VETO_COUNT: i32 = 32; +const VETO_COUNT: usize = 32; fn generate_run_start<'a>( fbb: &'a mut FlatBufferBuilder<'_>, @@ -129,7 +129,7 @@ fn get_veto_names_fbb<'a>( for i in 0..VETO_COUNT { let name = veto_names - .get(i as usize) + .get(i) .cloned() .unwrap_or_else(|| format!("saluki_veto_{i}")); @@ -137,23 +137,25 @@ fn get_veto_names_fbb<'a>( } } -fn get_enabled_vetoes(conf: &HowlConfig, rng: &mut ThreadRng, vetoes: &mut u32) { - let mut vtemp = *vetoes; +fn get_enabled_vetoes(conf: &HowlConfig, rng: &mut ThreadRng) -> u32 { + let mut vetoes = 0; for i in 0..VETO_COUNT { - let active = rng.random_bool(conf.veto_probability[i as usize]); - vtemp = (vtemp << 1) | active as u32; + let active = rng.random_bool(conf.veto_probability[i]); + vetoes = (vetoes << 1) | active as u32; } - *vetoes = vtemp; + vetoes } -fn get_active_vetoes(conf: &HowlConfig, vetoes: &mut u32) { - let mut vtemp = *vetoes; +fn get_active_vetoes(conf: &HowlConfig) -> u32 { + let mut vetoes = 0; for i in 0..VETO_COUNT { - vtemp = (vtemp << 1) | conf.enabled_vetoes[i as usize] as u32; + vetoes = (vetoes << 1) | conf.enabled_vetoes[i] as u32; } + + vetoes } fn produce_messages( @@ -467,11 +469,8 @@ pub fn howl(conf: &HowlConfig) { let mut fbb = FlatBufferBuilder::new(); let mut rng = rand::rng(); - let mut active_vetoes: u32 = 0; - let mut enabled_vetoes: u32 = 0; - - get_active_vetoes(conf, &mut active_vetoes); - get_enabled_vetoes(conf, &mut rng, &mut enabled_vetoes); + let active_vetoes = get_active_vetoes(conf); + let enabled_vetoes = get_enabled_vetoes(conf, &mut rng); let now_nanos = SystemTime::now() .duration_since(SystemTime::UNIX_EPOCH) From 2e2297c5e3ca21b9c352c93c82f60c5e64cc0477 Mon Sep 17 00:00:00 2001 From: nxq64494 Date: Tue, 22 Sep 2026 10:23:16 +0100 Subject: [PATCH 13/18] Index out of bounds + reuse veto_count glob --- src/howl.rs | 13 ++++++++----- src/main.rs | 8 ++++---- 2 files changed, 12 insertions(+), 9 deletions(-) diff --git a/src/howl.rs b/src/howl.rs index 2b0ca10..f2a0af6 100644 --- a/src/howl.rs +++ b/src/howl.rs @@ -29,7 +29,7 @@ use rdkafka::producer::{BaseRecord, DefaultProducerContext, ThreadedProducer}; use serde_json::json; use uuid::Uuid; -const VETO_COUNT: usize = 32; +pub const VETO_COUNT: usize = 32; fn generate_run_start<'a>( fbb: &'a mut FlatBufferBuilder<'_>, @@ -131,7 +131,7 @@ fn get_veto_names_fbb<'a>( let name = veto_names .get(i) .cloned() - .unwrap_or_else(|| format!("saluki_veto_{i}")); + .unwrap_or(format!("saluki_veto_{i}")); buf.push(fbb.create_string(&name.to_string())); } @@ -139,10 +139,11 @@ fn get_veto_names_fbb<'a>( fn get_enabled_vetoes(conf: &HowlConfig, rng: &mut ThreadRng) -> u32 { let mut vetoes = 0; + let mut prob; for i in 0..VETO_COUNT { - let active = rng.random_bool(conf.veto_probability[i]); - vetoes = (vetoes << 1) | active as u32; + prob = conf.veto_probability.get(i).cloned().unwrap_or(0.0); + vetoes = (vetoes << 1) | rng.random_bool(prob) as u32; } vetoes @@ -150,9 +151,11 @@ fn get_enabled_vetoes(conf: &HowlConfig, rng: &mut ThreadRng) -> u32 { fn get_active_vetoes(conf: &HowlConfig) -> u32 { let mut vetoes = 0; + let mut active; for i in 0..VETO_COUNT { - vetoes = (vetoes << 1) | conf.enabled_vetoes[i] as u32; + active = conf.enabled_vetoes.get(i).cloned().unwrap_or(false); + vetoes = (vetoes << 1) | active as u32; } vetoes diff --git a/src/main.rs b/src/main.rs index de0146a..80a2db6 100644 --- a/src/main.rs +++ b/src/main.rs @@ -8,7 +8,7 @@ use crate::cli_utils::BrokerAndOptionalTopic; use crate::cli_utils::KafkaOption; use crate::consume::ConsumeConfig; use crate::count::count; -use crate::howl::{EventMessageConfig, HowlConfig, howl}; +use crate::howl::{EventMessageConfig, HowlConfig, VETO_COUNT, howl}; use crate::sniff::sniff; use clap::{Parser, Subcommand}; use cli_utils::{BrokerAndTopic, parse_broker_spec, parse_broker_spec_optional_topic}; @@ -96,13 +96,13 @@ enum Commands { #[arg(long, default_value = "1000")] det_max: i32, /// Veto probabilities - #[arg(long, num_args = 0..33, value_delimiter = ' ')] + #[arg(long, num_args = 0..VETO_COUNT+1, value_delimiter = ' ')] veto_probability: Vec, /// Enabled vetoes - #[arg(long, num_args = 0..33, value_delimiter = ' ')] + #[arg(long, num_args = 0..VETO_COUNT+1, value_delimiter = ' ')] enabled_vetoes: Vec, /// Veto names - #[arg(long, num_args = 0..33, value_delimiter = ' ')] + #[arg(long, num_args = 0..VETO_COUNT+1, value_delimiter = ' ')] veto_names: Vec, /// Enable howl fast mode (Disables randomised ev44 blob generation) #[arg(long, action=clap::ArgAction::SetTrue)] From dbcaee36c22548baa7820c93f1f2726f79dc58f7 Mon Sep 17 00:00:00 2001 From: nxq64494 Date: Wed, 23 Sep 2026 13:18:55 +0100 Subject: [PATCH 14/18] Make separate func for setting kafka options --- src/cli_utils.rs | 16 ++++++++++++++++ src/consume.rs | 12 ++---------- src/howl.rs | 18 +++++------------- 3 files changed, 23 insertions(+), 23 deletions(-) diff --git a/src/cli_utils.rs b/src/cli_utils.rs index b19b553..4e5e986 100644 --- a/src/cli_utils.rs +++ b/src/cli_utils.rs @@ -1,4 +1,5 @@ use anyhow::{Context, Result, bail}; +use rdkafka::ClientConfig; use std::str::FromStr; pub(crate) fn parse_broker_spec(s: &str) -> Result { @@ -68,6 +69,21 @@ pub struct KafkaOption { pub value: String, } +pub fn set_kafka_options( + client_config: &mut ClientConfig, + kafka_config: &Option>, +) { + if let Some(kafka_options) = kafka_config { + for option in kafka_options { + println!( + "Setting Kafka config option {}={}", + option.key, option.value + ); + client_config.set(&option.key, &option.value); + } + } +} + impl FromStr for KafkaOption { type Err = String; diff --git a/src/consume.rs b/src/consume.rs index 3cf9993..7172b1f 100644 --- a/src/consume.rs +++ b/src/consume.rs @@ -1,5 +1,5 @@ use crate::KafkaOption; -use crate::cli_utils::BrokerAndTopic; +use crate::cli_utils::{BrokerAndTopic, set_kafka_options}; use isis_streaming_data_types::{deserialize_message, get_schema_id}; use log::{debug, error, info}; @@ -36,15 +36,7 @@ pub fn consume(config: &ConsumeConfig) { client_config.set("group.id", Uuid::new_v4().to_string()); client_config.set("bootstrap.servers", config.topic.broker()); - if let Some(kafka_options) = &config.kafka_config { - for option in kafka_options { - println!( - "Setting Kafka config option {}={}", - option.key, option.value - ); - client_config.set(&option.key, &option.value); - } - } + set_kafka_options(&mut client_config, &config.kafka_config); let consumer: BaseConsumer = client_config.create().expect("Base creation failed"); diff --git a/src/howl.rs b/src/howl.rs index f2a0af6..ac9f36d 100644 --- a/src/howl.rs +++ b/src/howl.rs @@ -13,6 +13,7 @@ use isis_streaming_data_types::flatbuffers_generated::veto_configuration_vc00::{ Vetoes, VetoesArgs, finish_vetoes_buffer, }; +use crate::cli_utils::set_kafka_options; use isis_streaming_data_types::flatbuffers_generated::run_start_pl72::{ RunStart, RunStartArgs, SpectraDetectorMapping, SpectraDetectorMappingArgs, finish_run_start_buffer, @@ -484,22 +485,13 @@ pub fn howl(conf: &HowlConfig) { calculate_data_rate(&mut fbb, &mut rng, conf, now_nanos, &active_vetoes); - let mut config: ClientConfig = ClientConfig::new(); - config.set("bootstrap.servers", conf.broker); - - if let Some(kafka_options) = &conf.kafka_config { - for option in kafka_options { - println!( - "Setting Kafka config option {}={}", - option.key, option.value - ); - config.set(&option.key, &option.value); - } - } + let mut client_config: ClientConfig = ClientConfig::new(); + client_config.set("bootstrap.servers", conf.broker); + set_kafka_options(&mut client_config, &conf.kafka_config); // create producer let mut producer: ThreadedProducer = - config.create().expect("Producer creation error"); + client_config.create().expect("Producer creation error"); let mut current_job_id = Uuid::new_v4().to_string(); From 287ab53e002b7449e4e0f3a1ce0560128c493d56 Mon Sep 17 00:00:00 2001 From: nxq64494 Date: Wed, 23 Sep 2026 13:20:23 +0100 Subject: [PATCH 15/18] Deduplicate --- src/consume.rs | 1 - src/count.rs | 13 ++----------- src/sniff.rs | 13 ++----------- 3 files changed, 4 insertions(+), 23 deletions(-) diff --git a/src/consume.rs b/src/consume.rs index 7172b1f..84ea6ea 100644 --- a/src/consume.rs +++ b/src/consume.rs @@ -35,7 +35,6 @@ pub fn consume(config: &ConsumeConfig) { let mut client_config = ClientConfig::new(); client_config.set("group.id", Uuid::new_v4().to_string()); client_config.set("bootstrap.servers", config.topic.broker()); - set_kafka_options(&mut client_config, &config.kafka_config); let consumer: BaseConsumer = client_config.create().expect("Base creation failed"); diff --git a/src/count.rs b/src/count.rs index d707bea..7b04278 100644 --- a/src/count.rs +++ b/src/count.rs @@ -1,5 +1,5 @@ use crate::KafkaOption; -use crate::cli_utils::BrokerAndTopic; +use crate::cli_utils::{set_kafka_options, BrokerAndTopic}; use futures::stream::StreamExt; use log::error; use rdkafka::consumer::{Consumer, DefaultConsumerContext, StreamConsumer}; @@ -15,16 +15,7 @@ pub async fn count( let mut config = ClientConfig::new(); config.set("group.id", Uuid::new_v4().to_string()); config.set("bootstrap.servers", topic.broker()); - - if let Some(kafka_options) = kafka_config { - for option in kafka_options { - println!( - "Setting Kafka config option {}={}", - option.key, option.value - ); - config.set(&option.key, &option.value); - } - } + set_kafka_options(&mut config, &kafka_config); let consumer: StreamConsumer = config.create().expect("Consumer creation failed"); diff --git a/src/sniff.rs b/src/sniff.rs index 2a42914..4aad732 100644 --- a/src/sniff.rs +++ b/src/sniff.rs @@ -1,5 +1,5 @@ use crate::KafkaOption; -use crate::cli_utils::BrokerAndOptionalTopic; +use crate::cli_utils::{set_kafka_options, BrokerAndOptionalTopic}; use rdkafka::ClientConfig; use rdkafka::consumer::{BaseConsumer, Consumer}; use std::time::Duration; @@ -7,16 +7,7 @@ use std::time::Duration; pub fn sniff(broker: &BrokerAndOptionalTopic, kafka_config: Option>) { let mut config = ClientConfig::new(); config.set("bootstrap.servers", broker.broker()); - - if let Some(kafka_options) = kafka_config { - for option in kafka_options { - println!( - "Setting Kafka config option {}={}", - option.key, option.value - ); - config.set(&option.key, &option.value); - } - } + set_kafka_options(&mut config, &kafka_config); let consumer: BaseConsumer = config.create().expect("Consumer creation failed"); From bebc771242347f3225bbb443b34ce75441d0ebe2 Mon Sep 17 00:00:00 2001 From: nxq64494 Date: Wed, 23 Sep 2026 14:10:17 +0100 Subject: [PATCH 16/18] Test --- src/count.rs | 2 +- src/howl.rs | 159 ++++++++++++++++++++++++++++----------------------- src/sniff.rs | 2 +- 3 files changed, 89 insertions(+), 74 deletions(-) diff --git a/src/count.rs b/src/count.rs index 7b04278..7941bd3 100644 --- a/src/count.rs +++ b/src/count.rs @@ -1,5 +1,5 @@ use crate::KafkaOption; -use crate::cli_utils::{set_kafka_options, BrokerAndTopic}; +use crate::cli_utils::{BrokerAndTopic, set_kafka_options}; use futures::stream::StreamExt; use log::error; use rdkafka::consumer::{Consumer, DefaultConsumerContext, StreamConsumer}; diff --git a/src/howl.rs b/src/howl.rs index ac9f36d..8760e5d 100644 --- a/src/howl.rs +++ b/src/howl.rs @@ -168,8 +168,8 @@ fn produce_messages( rng: &mut ThreadRng, frame: u32, conf: &HowlConfig, - current_job_id: &mut String, vetoes_mask: &u32, + enabled_vetoes: &u32, ) { // get current time let now_nanos = SystemTime::now() @@ -179,72 +179,24 @@ fn produce_messages( .try_into() .expect("This will fail after April 11th, 2262"); - match producer.send( - BaseRecord::to(conf.event_topic) - .key("") - .payload(generate_fake_metadata(vetoes_mask, fbb, now_nanos)) - .timestamp(now_nanos / 1_000_000), - ) { - Ok(_) => {} - Err(err) => { - error!("Failed to send messages: {}", err.0); - } - } - - let ev44 = generate_fake_events(fbb, rng, frame, conf.event_message_config, now_nanos).to_vec(); - - for _ in 0..conf.messages_per_frame { - match producer.send( - BaseRecord::to(conf.event_topic) - .key("") - .payload(if conf.fast { - ev44.as_slice() - } else { - generate_fake_events(fbb, rng, frame, conf.event_message_config, now_nanos) - }) - .timestamp(now_nanos / 1_000_000), - ) { - Ok(_) => {} - Err(err) => { - error!("Failed to send messages: {}", err.0); - } - } - } - if conf.frames_per_run > 0 && frame.is_multiple_of(conf.frames_per_run) { + let current_job_id = Uuid::new_v4().to_string(); + info!( "Starting new run after {} simulated frames", conf.frames_per_run ); - match producer.send( - BaseRecord::to(conf.run_info_topic) - .key("") - .payload(generate_run_stop(fbb, current_job_id)) - .timestamp(now_nanos / 1_000_000), - ) { - Ok(_) => {} - Err(err) => { - error!("Failed to send run stop: {}", err.0); - } - } - *current_job_id = Uuid::new_v4().to_string(); - match producer.send( - BaseRecord::to(conf.run_info_topic) - .key("") - .payload(generate_run_start( - fbb, - conf.event_message_config.det_max, - conf.event_topic, - current_job_id, - )) - .timestamp(now_nanos / 1_000_000), - ) { - Ok(_) => {} - Err(err) => { - error!("Failed to send run start: {}", err.0); - } + + if frame != 0 { + send_run_stop(producer, fbb, conf, ¤t_job_id, now_nanos); } + + send_run_start(producer, fbb, conf, ¤t_job_id, now_nanos); + send_veto_config(producer, fbb, conf, enabled_vetoes, now_nanos); } + + send_run_metadata(producer, fbb, conf, vetoes_mask, now_nanos); + send_run_data(producer, fbb, conf, rng, frame, now_nanos); } pub struct EventMessageConfig { @@ -357,8 +309,56 @@ fn calculate_data_rate( println!("Each ev44 is {ev44_size} bytes"); } +fn send_run_metadata( + producer: &ThreadedProducer, + fbb: &mut FlatBufferBuilder<'_>, + conf: &HowlConfig, + vetoes_mask: &u32, + now_nanos: i64, +) { + match producer.send( + BaseRecord::to(conf.event_topic) + .key("") + .payload(generate_fake_metadata(vetoes_mask, fbb, now_nanos)) + .timestamp(now_nanos / 1_000_000), + ) { + Ok(_) => {} + Err(err) => { + error!("Failed to send messages: {}", err.0); + } + } +} + +fn send_run_data( + producer: &ThreadedProducer, + fbb: &mut FlatBufferBuilder<'_>, + conf: &HowlConfig, + rng: &mut ThreadRng, + frame: u32, + now_nanos: i64, +) { + let ev44 = generate_fake_events(fbb, rng, frame, conf.event_message_config, now_nanos).to_vec(); + + for _ in 0..conf.messages_per_frame { + match producer.send( + BaseRecord::to(conf.event_topic) + .key("") + .payload(if conf.fast { + ev44.as_slice() + } else { + generate_fake_events(fbb, rng, frame, conf.event_message_config, now_nanos) + }) + .timestamp(now_nanos / 1_000_000), + ) { + Ok(_) => {} + Err(err) => { + error!("Failed to send messages: {}", err.0); + } + } + } +} fn send_run_start( - producer: &mut ThreadedProducer, + producer: &ThreadedProducer, fbb: &mut FlatBufferBuilder<'_>, conf: &HowlConfig, current_job_id: &str, @@ -379,8 +379,28 @@ fn send_run_start( .expect("Failed to enqueue run start message"); } +fn send_run_stop( + producer: &ThreadedProducer, + fbb: &mut FlatBufferBuilder<'_>, + conf: &HowlConfig, + current_job_id: &str, + now_nanos: i64, +) { + match producer.send( + BaseRecord::to(conf.run_info_topic) + .key("") + .payload(generate_run_stop(fbb, current_job_id)) + .timestamp(now_nanos / 1_000_000), + ) { + Ok(_) => {} + Err(err) => { + error!("Failed to send run stop: {}", err.0); + } + } +} + fn send_veto_config( - producer: &mut ThreadedProducer, + producer: &ThreadedProducer, fbb: &mut FlatBufferBuilder<'_>, conf: &HowlConfig, vetoes_mask: &u32, @@ -406,8 +426,8 @@ fn howl_begin( fbb: &mut FlatBufferBuilder<'_>, rng: &mut ThreadRng, conf: &HowlConfig, - current_job_id: &mut String, - vetoes_mask: &u32, + active_vetoes: &u32, + enabled_vetoes: &u32, ) { let target_frame_time = Duration::from_secs_f64(1.0 / conf.frames_per_second as f64); debug!("Target frame time: {target_frame_time:?}"); @@ -423,15 +443,14 @@ fn howl_begin( target_time += target_frame_time; debug!("New target: {target_time:?}"); frames += 1; - debug!("current job id: {current_job_id}"); produce_messages( producer, fbb, rng, frames, conf, - current_job_id, - vetoes_mask, + active_vetoes, + enabled_vetoes, ); let now = SystemTime::now() .duration_since(SystemTime::UNIX_EPOCH) @@ -493,16 +512,12 @@ pub fn howl(conf: &HowlConfig) { let mut producer: ThreadedProducer = client_config.create().expect("Producer creation error"); - let mut current_job_id = Uuid::new_v4().to_string(); - - send_run_start(&mut producer, &mut fbb, conf, ¤t_job_id, now_nanos); - send_veto_config(&mut producer, &mut fbb, conf, &enabled_vetoes, now_nanos); howl_begin( &mut producer, &mut fbb, &mut rng, conf, - &mut current_job_id, &active_vetoes, + &enabled_vetoes, ); } diff --git a/src/sniff.rs b/src/sniff.rs index 4aad732..95a80ab 100644 --- a/src/sniff.rs +++ b/src/sniff.rs @@ -1,5 +1,5 @@ use crate::KafkaOption; -use crate::cli_utils::{set_kafka_options, BrokerAndOptionalTopic}; +use crate::cli_utils::{BrokerAndOptionalTopic, set_kafka_options}; use rdkafka::ClientConfig; use rdkafka::consumer::{BaseConsumer, Consumer}; use std::time::Duration; From 15b8794e070d60584d7c11cdb3f0c6a0fca49bd3 Mon Sep 17 00:00:00 2001 From: nxq64494 Date: Wed, 23 Sep 2026 15:24:12 +0100 Subject: [PATCH 17/18] Test: Ensure job id is same for start/end --- src/howl.rs | 23 ++++++++++++++--------- 1 file changed, 14 insertions(+), 9 deletions(-) diff --git a/src/howl.rs b/src/howl.rs index 8760e5d..27bd310 100644 --- a/src/howl.rs +++ b/src/howl.rs @@ -162,14 +162,16 @@ fn get_active_vetoes(conf: &HowlConfig) -> u32 { vetoes } +#[allow(clippy::too_many_arguments)] fn produce_messages( producer: &ThreadedProducer, fbb: &mut FlatBufferBuilder, rng: &mut ThreadRng, frame: u32, conf: &HowlConfig, - vetoes_mask: &u32, + active_vetoes: &u32, enabled_vetoes: &u32, + current_job_id: &mut String, ) { // get current time let now_nanos = SystemTime::now() @@ -180,23 +182,23 @@ fn produce_messages( .expect("This will fail after April 11th, 2262"); if conf.frames_per_run > 0 && frame.is_multiple_of(conf.frames_per_run) { - let current_job_id = Uuid::new_v4().to_string(); - info!( "Starting new run after {} simulated frames", conf.frames_per_run ); if frame != 0 { - send_run_stop(producer, fbb, conf, ¤t_job_id, now_nanos); + send_run_stop(producer, fbb, conf, current_job_id, now_nanos); } - send_run_start(producer, fbb, conf, ¤t_job_id, now_nanos); + *current_job_id = Uuid::new_v4().to_string(); + + send_run_start(producer, fbb, conf, current_job_id, now_nanos); send_veto_config(producer, fbb, conf, enabled_vetoes, now_nanos); } - send_run_metadata(producer, fbb, conf, vetoes_mask, now_nanos); - send_run_data(producer, fbb, conf, rng, frame, now_nanos); + send_frame_metadata(producer, fbb, conf, active_vetoes, now_nanos); + send_frame_data(producer, fbb, conf, rng, frame, now_nanos); } pub struct EventMessageConfig { @@ -309,7 +311,7 @@ fn calculate_data_rate( println!("Each ev44 is {ev44_size} bytes"); } -fn send_run_metadata( +fn send_frame_metadata( producer: &ThreadedProducer, fbb: &mut FlatBufferBuilder<'_>, conf: &HowlConfig, @@ -329,7 +331,7 @@ fn send_run_metadata( } } -fn send_run_data( +fn send_frame_data( producer: &ThreadedProducer, fbb: &mut FlatBufferBuilder<'_>, conf: &HowlConfig, @@ -439,6 +441,8 @@ fn howl_begin( .expect("Failed to get system time"); debug!("Target time: {target_time:?}"); + let mut current_job_id = Uuid::new_v4().to_string(); + loop { target_time += target_frame_time; debug!("New target: {target_time:?}"); @@ -451,6 +455,7 @@ fn howl_begin( conf, active_vetoes, enabled_vetoes, + &mut current_job_id, ); let now = SystemTime::now() .duration_since(SystemTime::UNIX_EPOCH) From 37a69df4946147cd1b9e00828a638731241ad555 Mon Sep 17 00:00:00 2001 From: nxq64494 Date: Wed, 23 Sep 2026 15:27:56 +0100 Subject: [PATCH 18/18] Test: small change --- src/howl.rs | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/src/howl.rs b/src/howl.rs index 27bd310..40da051 100644 --- a/src/howl.rs +++ b/src/howl.rs @@ -189,10 +189,9 @@ fn produce_messages( if frame != 0 { send_run_stop(producer, fbb, conf, current_job_id, now_nanos); + *current_job_id = Uuid::new_v4().to_string(); } - *current_job_id = Uuid::new_v4().to_string(); - send_run_start(producer, fbb, conf, current_job_id, now_nanos); send_veto_config(producer, fbb, conf, enabled_vetoes, now_nanos); }