use netflow_parser::{
InMemoryTemplateStore, NetflowPacket, NetflowParser, TemplateEvent, TemplateProtocol,
};
use std::sync::Arc;
use std::sync::Mutex;
fn main() {
println!("=== Horizontal scale-out demo: two parsers sharing a TemplateStore ===\n");
let store = Arc::new(InMemoryTemplateStore::new());
println!("[replica A] starting up, will learn one template");
let mut replica_a = NetflowParser::builder()
.with_template_store(Arc::clone(&store) as _)
.build()
.expect("build replica A");
let template_packet = build_v9_template_packet(256, &[(8, 4), (12, 4), (1, 8)]);
let result = replica_a.parse_bytes(&template_packet);
if let Some(err) = result.error {
panic!("template parse failed: {err}");
}
println!(
"[replica A] learned template 256, store now has {} entr(ies)\n",
store.len()
);
drop(replica_a);
println!("[replica B] starting cold (no in-process template cache)");
let restored_log: Arc<Mutex<Vec<(TemplateProtocol, u16)>>> =
Arc::new(Mutex::new(Vec::new()));
let restored_log_for_hook = Arc::clone(&restored_log);
let mut replica_b = NetflowParser::builder()
.with_template_store(Arc::clone(&store) as _)
.on_template_event(move |event| {
if let TemplateEvent::Restored {
template_id: Some(id),
protocol,
} = event
{
restored_log_for_hook
.lock()
.expect("poisoned")
.push((*protocol, *id));
}
Ok(())
})
.build()
.expect("build replica B");
let data_payload = [
0x0A, 0x00, 0x00, 0x01, 0x0A, 0x00, 0x00, 0x02, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x10, 0x00,
];
let data_packet = build_v9_data_packet(256, &data_payload);
let result = replica_b.parse_bytes(&data_packet);
if let Some(err) = result.error {
panic!("data parse on replica B failed: {err}");
}
let v9 = result
.packets
.into_iter()
.find_map(|p| match p {
NetflowPacket::V9(v) => Some(v),
_ => None,
})
.expect("expected a V9 packet");
let flowset_count = v9.flowsets.len();
println!(
"[replica B] decoded data packet against restored template ({} flowset(s))",
flowset_count
);
let metrics = replica_b.v9_cache_info().metrics;
println!("\n[replica B] cache metrics after read-through:");
println!(" hits = {}", metrics.hits);
println!(" misses = {}", metrics.misses);
println!(
" template_store_restored = {}",
metrics.template_store_restored
);
println!(
" template_store_codec_err = {}",
metrics.template_store_codec_errors
);
println!(
" template_store_backend_err = {}",
metrics.template_store_backend_errors
);
let restored = restored_log.lock().expect("poisoned");
println!("\n[replica B] TemplateEvent::Restored events:");
for (protocol, id) in restored.iter() {
println!(" {:?} template_id={}", protocol, id);
}
println!("\nDone. The same protocol works for IPFIX and IPFIX-options templates.");
println!(
"Hot-path overhead when no store is configured is a single Option::is_none branch."
);
}
fn build_v9_template_packet(template_id: u16, fields: &[(u16, u16)]) -> Vec<u8> {
let template_record_len = 4 + fields.len() * 4; let flowset_len = 4 + template_record_len; let mut pkt = Vec::new();
pkt.extend_from_slice(&9u16.to_be_bytes()); pkt.extend_from_slice(&1u16.to_be_bytes()); pkt.extend_from_slice(&0u32.to_be_bytes()); pkt.extend_from_slice(&0u32.to_be_bytes()); pkt.extend_from_slice(&0u32.to_be_bytes()); pkt.extend_from_slice(&0u32.to_be_bytes()); pkt.extend_from_slice(&0u16.to_be_bytes()); pkt.extend_from_slice(&(flowset_len as u16).to_be_bytes());
pkt.extend_from_slice(&template_id.to_be_bytes());
pkt.extend_from_slice(&(fields.len() as u16).to_be_bytes());
for &(ft, fl) in fields {
pkt.extend_from_slice(&ft.to_be_bytes());
pkt.extend_from_slice(&fl.to_be_bytes());
}
pkt
}
fn build_v9_data_packet(template_id: u16, payload: &[u8]) -> Vec<u8> {
let flowset_len = 4 + payload.len();
let mut pkt = Vec::new();
pkt.extend_from_slice(&9u16.to_be_bytes());
pkt.extend_from_slice(&1u16.to_be_bytes());
pkt.extend_from_slice(&0u32.to_be_bytes());
pkt.extend_from_slice(&0u32.to_be_bytes());
pkt.extend_from_slice(&0u32.to_be_bytes());
pkt.extend_from_slice(&0u32.to_be_bytes());
pkt.extend_from_slice(&template_id.to_be_bytes());
pkt.extend_from_slice(&(flowset_len as u16).to_be_bytes());
pkt.extend_from_slice(payload);
pkt
}