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 bfdb1df..2b0ca10 100644 --- a/src/howl.rs +++ b/src/howl.rs @@ -2,13 +2,17 @@ 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, }; 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::{ + Vetoes, VetoesArgs, finish_vetoes_buffer, +}; + use isis_streaming_data_types::flatbuffers_generated::run_start_pl72::{ RunStart, RunStartArgs, SpectraDetectorMapping, SpectraDetectorMappingArgs, finish_run_start_buffer, @@ -25,6 +29,8 @@ use rdkafka::producer::{BaseRecord, DefaultProducerContext, ThreadedProducer}; use serde_json::json; use uuid::Uuid; +const VETO_COUNT: usize = 32; + fn generate_run_start<'a>( fbb: &'a mut FlatBufferBuilder<'_>, det_max: i32, @@ -114,6 +120,44 @@ fn generate_run_stop<'a>(fbb: &'a mut FlatBufferBuilder<'_>, job_id: &str) -> &' fbb.finished_data() } +fn get_veto_names_fbb<'a>( + veto_names: &[String], + fbb: &mut FlatBufferBuilder<'a>, + buf: &mut Vec>, +) { + buf.clear(); + + for i in 0..VETO_COUNT { + let name = veto_names + .get(i) + .cloned() + .unwrap_or_else(|| format!("saluki_veto_{i}")); + + buf.push(fbb.create_string(&name.to_string())); + } +} + +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]); + vetoes = (vetoes << 1) | active as u32; + } + + vetoes +} + +fn get_active_vetoes(conf: &HowlConfig) -> u32 { + let mut vetoes = 0; + + for i in 0..VETO_COUNT { + vetoes = (vetoes << 1) | conf.enabled_vetoes[i] as u32; + } + + vetoes +} + fn produce_messages( producer: &ThreadedProducer, fbb: &mut FlatBufferBuilder, @@ -121,6 +165,7 @@ fn produce_messages( frame: u32, conf: &HowlConfig, current_job_id: &mut String, + vetoes_mask: &u32, ) { // get current time let now_nanos = SystemTime::now() @@ -133,12 +178,7 @@ fn produce_messages( match producer.send( BaseRecord::to(conf.event_topic) .key("") - .payload(generate_fake_metadata( - rng, - fbb, - now_nanos, - conf.veto_probability, - )) + .payload(generate_fake_metadata(vetoes_mask, fbb, now_nanos)) .timestamp(now_nanos / 1_000_000), ) { Ok(_) => {} @@ -244,61 +284,61 @@ fn generate_fake_events<'a>( } fn generate_fake_metadata<'a>( - rng: &mut ThreadRng, + vetoes_mask: &u32, 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 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), // active proton_charge: Some(0.1), }; let pu00 = Pu00Message::create(fbb, &args); finish_pu_00_message_buffer(fbb, pu00); + fbb.finished_data() } -pub struct HowlConfig<'a> { - pub broker: &'a str, - pub event_topic: &'a str, - pub run_info_topic: &'a str, - 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 event_message_config: &'a EventMessageConfig, - pub fast: bool, - pub kafka_config: Option>, -} +fn generate_veto_config<'a>( + veto_names: &[String], + fbb: &'a mut FlatBufferBuilder<'_>, + 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); -pub fn howl(conf: &HowlConfig) { - // create producer - let mut fbb = FlatBufferBuilder::new(); - let mut rng = rand::rng(); + let args = VetoesArgs { + timestamp: timestamp_ns, + vetoes: *vetoes_mask, // enabled + veto_names: Some(fbb.create_vector(&veto_names_fbb)), + }; + let vc00 = Vetoes::create(fbb, &args); + finish_vetoes_buffer(fbb, vc00); - 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"); + fbb.finished_data() +} +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(&mut rng, &mut fbb, now_nanos, conf.veto_probability).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}"); @@ -311,39 +351,60 @@ 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"); +} +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, + fbb, + now_nanos, + 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:?}"); @@ -360,12 +421,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, + current_job_id, + vetoes_mask, ); let now = SystemTime::now() .duration_since(SystemTime::UNIX_EPOCH) @@ -386,3 +448,66 @@ 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 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) + .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, &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); + } + } + + // 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(&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, + ); +} diff --git a/src/main.rs b/src/main.rs index 0ac9c45..de0146a 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, num_args = 0..33, value_delimiter = ' ')] + veto_probability: Vec, + /// Enabled vetoes + #[arg(long, num_args = 0..33, value_delimiter = ' ')] + enabled_vetoes: Vec, + /// Veto names + #[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)] fast: bool, @@ -167,6 +173,8 @@ async fn main() { det_min, det_max, veto_probability, + enabled_vetoes, + veto_names, fast, kafka_config, } => howl(&HowlConfig { @@ -174,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, @@ -185,6 +194,8 @@ async fn main() { det_max, }, veto_probability, + enabled_vetoes, + veto_names, fast, }), Commands::Count {