From 0be5f5139ebdd66952001a4f42baf081eb813c72 Mon Sep 17 00:00:00 2001 From: mikemiles-dev Date: Sat, 22 Nov 2025 10:15:11 -0600 Subject: [PATCH] fix: Several memory and performance optimizations. --- Cargo.toml | 2 +- RELEASES.md | 3 + src/lib.rs | 10 +- src/netflow_common.rs | 173 ++++++++++++++++----------- src/static_versions/v5.rs | 4 +- src/static_versions/v7.rs | 4 +- src/variable_versions/data_number.rs | 13 +- src/variable_versions/ipfix.rs | 69 +++++++---- src/variable_versions/v9.rs | 41 +++---- 9 files changed, 184 insertions(+), 135 deletions(-) diff --git a/Cargo.toml b/Cargo.toml index 64967f64..2f9ff0f5 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,7 +1,7 @@ [package] name = "netflow_parser" description = "Parser for Netflow Cisco V5, V7, V9, IPFIX" -version = "0.6.4" +version = "0.6.5" edition = "2024" authors = ["michael.mileusnich@gmail.com"] license = "MIT OR Apache-2.0" diff --git a/RELEASES.md b/RELEASES.md index 26011768..5ce5afed 100644 --- a/RELEASES.md +++ b/RELEASES.md @@ -1,3 +1,6 @@ +# 0.6.5 +* Several memory and performance optimizations. + # 0.6.4 * Removed uneeded DataNumber Parsing for Durations. * Renamed methods DurationMicros and DurationNanos into DurationMicrosNTP and DurationNanosNTP. diff --git a/src/lib.rs b/src/lib.rs index f1254b2f..ce9affab 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -266,10 +266,10 @@ pub struct NetflowParser { } #[derive(Debug, Clone)] -pub enum ParsedNetflow { +pub enum ParsedNetflow<'a> { Success { packet: NetflowPacket, - remaining: Vec, + remaining: &'a [u8], }, Error { error: NetflowParseError, @@ -334,10 +334,10 @@ impl NetflowParser { } let mut packets = Vec::new(); - let mut remaining = packet.to_vec(); + let mut remaining = packet; while !remaining.is_empty() { - match self.parse_packet_by_version(&remaining) { + match self.parse_packet_by_version(remaining) { ParsedNetflow::Success { packet, remaining: new_remaining, @@ -362,7 +362,7 @@ impl NetflowParser { } #[inline] - fn parse_packet_by_version(&mut self, packet: &[u8]) -> ParsedNetflow { + fn parse_packet_by_version<'a>(&mut self, packet: &'a [u8]) -> ParsedNetflow<'a> { match GenericNetflowHeader::parse(packet) { Ok((packet, header)) if self.allowed_versions.contains(&header.version) => { match header.version { diff --git a/src/netflow_common.rs b/src/netflow_common.rs index 3c3b5ba1..507dcbae 100644 --- a/src/netflow_common.rs +++ b/src/netflow_common.rs @@ -1,4 +1,3 @@ -use std::collections::HashMap; use std::net::IpAddr; use crate::NetflowPacket; @@ -116,6 +115,14 @@ impl From<&V7> for NetflowCommon { } } +/// Helper function to find a field value in a V9 flow record by field type +fn find_v9_field<'a>( + fields: &'a [(V9Field, FieldValue)], + field: V9Field, +) -> Option<&'a FieldValue> { + fields.iter().find(|(f, _)| *f == field).map(|(_, v)| v) +} + impl From<&V9> for NetflowCommon { fn from(value: &V9) -> Self { // Convert V9 to NetflowCommon @@ -124,42 +131,33 @@ impl From<&V9> for NetflowCommon { for flowset in &value.flowsets { if let V9FlowSetBody::Data(data) = &flowset.body { for data_field in &data.fields { - let value_map: HashMap = - data_field.clone().into_iter().collect(); flowsets.push(NetflowCommonFlowSet { - src_addr: value_map - .get(&V9Field::Ipv4SrcAddr) - .or_else(|| value_map.get(&V9Field::Ipv6SrcAddr)) + src_addr: find_v9_field(data_field, V9Field::Ipv4SrcAddr) + .or_else(|| find_v9_field(data_field, V9Field::Ipv6SrcAddr)) .and_then(|v| v.try_into().ok()), - dst_addr: value_map - .get(&V9Field::Ipv4DstAddr) - .or_else(|| value_map.get(&V9Field::Ipv6DstAddr)) + dst_addr: find_v9_field(data_field, V9Field::Ipv4DstAddr) + .or_else(|| find_v9_field(data_field, V9Field::Ipv6DstAddr)) .and_then(|v| v.try_into().ok()), - src_port: value_map - .get(&V9Field::L4SrcPort) + src_port: find_v9_field(data_field, V9Field::L4SrcPort) .and_then(|v| v.try_into().ok()), - dst_port: value_map - .get(&V9Field::L4DstPort) + dst_port: find_v9_field(data_field, V9Field::L4DstPort) .and_then(|v| v.try_into().ok()), - protocol_number: value_map - .get(&V9Field::Protocol) + protocol_number: find_v9_field(data_field, V9Field::Protocol) .and_then(|v| v.try_into().ok()), - protocol_type: value_map.get(&V9Field::Protocol).and_then(|v| { - v.try_into() - .ok() - .map(|proto: u8| ProtocolTypes::from(proto)) - }), - first_seen: value_map - .get(&V9Field::FirstSwitched) + protocol_type: find_v9_field(data_field, V9Field::Protocol).and_then( + |v| { + v.try_into() + .ok() + .map(|proto: u8| ProtocolTypes::from(proto)) + }, + ), + first_seen: find_v9_field(data_field, V9Field::FirstSwitched) .and_then(|v| v.try_into().ok()), - last_seen: value_map - .get(&V9Field::LastSwitched) + last_seen: find_v9_field(data_field, V9Field::LastSwitched) .and_then(|v| v.try_into().ok()), - src_mac: value_map - .get(&V9Field::InSrcMac) + src_mac: find_v9_field(data_field, V9Field::InSrcMac) .and_then(|v| v.try_into().ok()), - dst_mac: value_map - .get(&V9Field::InDstMac) + dst_mac: find_v9_field(data_field, V9Field::InDstMac) .and_then(|v| v.try_into().ok()), }); } @@ -174,6 +172,14 @@ impl From<&V9> for NetflowCommon { } } +/// Helper function to find a field value in an IPFix flow record by field type +fn find_ipfix_field<'a>( + fields: &'a [(IPFixField, FieldValue)], + field: IPFixField, +) -> Option<&'a FieldValue> { + fields.iter().find(|(f, _)| *f == field).map(|(_, v)| v) +} + impl From<&IPFix> for NetflowCommon { fn from(value: &IPFix) -> Self { // Convert IPFix to NetflowCommon @@ -183,52 +189,73 @@ impl From<&IPFix> for NetflowCommon { for flowset in &value.flowsets { if let IPFixFlowSetBody::Data(data) = &flowset.body { for data_field in &data.fields { - let value_map: HashMap = - data_field.clone().into_iter().collect(); flowsets.push(NetflowCommonFlowSet { - src_addr: value_map - .get(&IPFixField::IANA(IANAIPFixField::SourceIpv4address)) - .or_else(|| { - value_map - .get(&IPFixField::IANA(IANAIPFixField::SourceIpv6address)) - }) - .and_then(|v| v.try_into().ok()), - dst_addr: value_map - .get(&IPFixField::IANA(IANAIPFixField::DestinationIpv4address)) - .or_else(|| { - value_map.get(&IPFixField::IANA( - IANAIPFixField::DestinationIpv6address, - )) - }) - .and_then(|v| v.try_into().ok()), - src_port: value_map - .get(&IPFixField::IANA(IANAIPFixField::SourceTransportPort)) - .and_then(|v| v.try_into().ok()), - dst_port: value_map - .get(&IPFixField::IANA(IANAIPFixField::DestinationTransportPort)) - .and_then(|v| v.try_into().ok()), - protocol_number: value_map - .get(&IPFixField::IANA(IANAIPFixField::ProtocolIdentifier)) - .and_then(|v| v.try_into().ok()), - protocol_type: value_map - .get(&IPFixField::IANA(IANAIPFixField::ProtocolIdentifier)) - .and_then(|v| { - v.try_into() - .ok() - .map(|proto: u8| ProtocolTypes::from(proto)) - }), - first_seen: value_map - .get(&IPFixField::IANA(IANAIPFixField::FlowStartSysUpTime)) - .and_then(|v| v.try_into().ok()), - last_seen: value_map - .get(&IPFixField::IANA(IANAIPFixField::FlowEndSysUpTime)) - .and_then(|v| v.try_into().ok()), - src_mac: value_map - .get(&IPFixField::IANA(IANAIPFixField::SourceMacaddress)) - .and_then(|v| v.try_into().ok()), - dst_mac: value_map - .get(&IPFixField::IANA(IANAIPFixField::DestinationMacaddress)) - .and_then(|v| v.try_into().ok()), + src_addr: find_ipfix_field( + data_field, + IPFixField::IANA(IANAIPFixField::SourceIpv4address), + ) + .or_else(|| { + find_ipfix_field( + data_field, + IPFixField::IANA(IANAIPFixField::SourceIpv6address), + ) + }) + .and_then(|v| v.try_into().ok()), + dst_addr: find_ipfix_field( + data_field, + IPFixField::IANA(IANAIPFixField::DestinationIpv4address), + ) + .or_else(|| { + find_ipfix_field( + data_field, + IPFixField::IANA(IANAIPFixField::DestinationIpv6address), + ) + }) + .and_then(|v| v.try_into().ok()), + src_port: find_ipfix_field( + data_field, + IPFixField::IANA(IANAIPFixField::SourceTransportPort), + ) + .and_then(|v| v.try_into().ok()), + dst_port: find_ipfix_field( + data_field, + IPFixField::IANA(IANAIPFixField::DestinationTransportPort), + ) + .and_then(|v| v.try_into().ok()), + protocol_number: find_ipfix_field( + data_field, + IPFixField::IANA(IANAIPFixField::ProtocolIdentifier), + ) + .and_then(|v| v.try_into().ok()), + protocol_type: find_ipfix_field( + data_field, + IPFixField::IANA(IANAIPFixField::ProtocolIdentifier), + ) + .and_then(|v| { + v.try_into() + .ok() + .map(|proto: u8| ProtocolTypes::from(proto)) + }), + first_seen: find_ipfix_field( + data_field, + IPFixField::IANA(IANAIPFixField::FlowStartSysUpTime), + ) + .and_then(|v| v.try_into().ok()), + last_seen: find_ipfix_field( + data_field, + IPFixField::IANA(IANAIPFixField::FlowEndSysUpTime), + ) + .and_then(|v| v.try_into().ok()), + src_mac: find_ipfix_field( + data_field, + IPFixField::IANA(IANAIPFixField::SourceMacaddress), + ) + .and_then(|v| v.try_into().ok()), + dst_mac: find_ipfix_field( + data_field, + IPFixField::IANA(IANAIPFixField::DestinationMacaddress), + ) + .and_then(|v| v.try_into().ok()), }); } } diff --git a/src/static_versions/v5.rs b/src/static_versions/v5.rs index f5d73b0f..0bbdf2c8 100644 --- a/src/static_versions/v5.rs +++ b/src/static_versions/v5.rs @@ -15,11 +15,11 @@ use std::net::Ipv4Addr; pub struct V5Parser; impl V5Parser { - pub fn parse(packet: &[u8]) -> ParsedNetflow { + pub fn parse(packet: &[u8]) -> ParsedNetflow<'_> { match V5::parse(packet) { Ok((remaining, v5)) => ParsedNetflow::Success { packet: NetflowPacket::V5(v5), - remaining: remaining.to_vec(), + remaining, }, Err(e) => ParsedNetflow::Error { error: NetflowParseError::Partial(PartialParse { diff --git a/src/static_versions/v7.rs b/src/static_versions/v7.rs index 857a10f5..55be985d 100644 --- a/src/static_versions/v7.rs +++ b/src/static_versions/v7.rs @@ -16,11 +16,11 @@ use std::net::Ipv4Addr; pub struct V7Parser; impl V7Parser { - pub fn parse(packet: &[u8]) -> ParsedNetflow { + pub fn parse(packet: &[u8]) -> ParsedNetflow<'_> { match V7::parse(packet) { Ok((remaining, v7)) => ParsedNetflow::Success { packet: NetflowPacket::V7(v7), - remaining: remaining.to_vec(), + remaining, }, Err(e) => ParsedNetflow::Error { error: NetflowParseError::Partial(PartialParse { diff --git a/src/variable_versions/data_number.rs b/src/variable_versions/data_number.rs index ed6ddec7..81497a46 100644 --- a/src/variable_versions/data_number.rs +++ b/src/variable_versions/data_number.rs @@ -254,13 +254,12 @@ impl FieldValue { } FieldDataType::String => { let (i, taken) = take(field_length)(remaining)?; - let s = String::from_utf8_lossy(taken).to_string(); - let s: String = s.chars().filter(|&c| !c.is_control()).collect(); - let s = if s.starts_with("P4") { - s.trim_start_matches("P4").to_string() - } else { - s - }; + // Single-pass string construction: filter control chars and strip "P4" prefix + let s: String = String::from_utf8_lossy(taken) + .chars() + .filter(|&c| !c.is_control()) + .collect(); + let s = s.strip_prefix("P4").unwrap_or(&s).to_string(); (i, FieldValue::String(s)) } FieldDataType::Ip4Addr => { diff --git a/src/variable_versions/ipfix.rs b/src/variable_versions/ipfix.rs index db51e09d..dd53e605 100644 --- a/src/variable_versions/ipfix.rs +++ b/src/variable_versions/ipfix.rs @@ -43,11 +43,11 @@ pub struct IPFixParser { } impl IPFixParser { - pub fn parse(&mut self, packet: &[u8]) -> ParsedNetflow { + pub fn parse<'a>(&mut self, packet: &'a [u8]) -> ParsedNetflow<'a> { match IPFix::parse(packet, self) { Ok((remaining, ipfix)) => ParsedNetflow::Success { packet: NetflowPacket::IPFix(ipfix), - remaining: remaining.to_vec(), + remaining, }, Err(e) => ParsedNetflow::Error { error: NetflowParseError::Partial(PartialParse { @@ -59,21 +59,29 @@ impl IPFixParser { } } - /// Add a template (Template, OptionsTemplate, V9Template, or V9OptionsTemplate) to the parser generically. - pub fn add_templates(&mut self, templates: Vec) { - let any = &templates as &dyn std::any::Any; - if let Some(ts) = any.downcast_ref::>() { - self.templates - .extend(ts.iter().map(|t| (t.template_id, t.clone()))); - } else if let Some(ts) = any.downcast_ref::>() { + /// Add templates to the parser from a slice. + fn add_ipfix_templates(&mut self, templates: &[Template]) { + for t in templates { + self.templates.insert(t.template_id, t.clone()); + } + } + + fn add_ipfix_options_templates(&mut self, templates: &[OptionsTemplate]) { + for t in templates { self.ipfix_options_templates - .extend(ts.iter().map(|t| (t.template_id, t.clone()))); - } else if let Some(ts) = any.downcast_ref::>() { - self.v9_templates - .extend(ts.iter().map(|t| (t.template_id, t.clone()))); - } else if let Some(ts) = any.downcast_ref::>() { - self.v9_options_templates - .extend(ts.iter().map(|t| (t.template_id, t.clone()))); + .insert(t.template_id, t.clone()); + } + } + + fn add_v9_templates(&mut self, templates: &[V9Template]) { + for t in templates { + self.v9_templates.insert(t.template_id, t.clone()); + } + } + + fn add_v9_options_templates(&mut self, templates: &[V9OptionsTemplate]) { + for t in templates { + self.v9_options_templates.insert(t.template_id, t.clone()); } } } @@ -155,7 +163,7 @@ impl FlowSetBody { single_variant: fn(T) -> FlowSetBody, multi_variant: fn(Vec) -> FlowSetBody, validate: fn(&T) -> bool, - add_templates: fn(&mut IPFixParser, Vec), + add_templates: fn(&mut IPFixParser, &[T]), ) -> IResult<&'a [u8], FlowSetBody> where T: Clone, @@ -168,10 +176,10 @@ impl FlowSetBody { nom::error::ErrorKind::Verify, ))); } - add_templates(parser, templates.clone()); + add_templates(parser, &templates); match templates.len() { 1 => { - if let Some(template) = templates.first().cloned() { + if let Some(template) = templates.into_iter().next() { Ok((i, single_variant(template))) } else { Err(nom::Err::Error(nom::error::Error::new( @@ -197,7 +205,7 @@ impl FlowSetBody { FlowSetBody::Template, FlowSetBody::Templates, |t: &Template| t.is_valid(), - |parser, templates| parser.add_templates(templates), + |parser, templates| parser.add_ipfix_templates(templates), ), DATA_TEMPLATE_V9_ID => Self::parse_templates( i, @@ -206,7 +214,7 @@ impl FlowSetBody { FlowSetBody::V9Template, FlowSetBody::V9Templates, |_t: &V9Template| true, - |parser, templates| parser.add_templates(templates), + |parser, templates| parser.add_v9_templates(templates), ), OPTIONS_TEMPLATE_V9_ID => Self::parse_templates( i, @@ -215,7 +223,7 @@ impl FlowSetBody { FlowSetBody::V9OptionsTemplate, FlowSetBody::V9OptionsTemplates, |_t: &V9OptionsTemplate| true, - |parser, templates| parser.add_templates(templates), + |parser, templates| parser.add_v9_options_templates(templates), ), OPTIONS_TEMPLATE_IPFIX_ID => Self::parse_templates( i, @@ -224,7 +232,7 @@ impl FlowSetBody { FlowSetBody::OptionsTemplate, FlowSetBody::OptionsTemplates, |t: &OptionsTemplate| t.is_valid(), - |parser, templates| parser.add_templates(templates), + |parser, templates| parser.add_ipfix_options_templates(templates), ), // Parse Data _ => { @@ -406,11 +414,22 @@ impl<'a> FieldParser { if template_fields.is_empty() { return Ok((i, Vec::new())); } - let mut res = Vec::new(); + + // Estimate capacity based on input size and template field count + let template_size: usize = template_fields + .iter() + .map(|f| usize::from(f.field_length)) + .sum(); + let estimated_records = if template_size > 0 { + i.len() / template_size + } else { + 0 + }; + let mut res = Vec::with_capacity(estimated_records); // Try to parse as much as we can, but if it fails, just return what we have so far. while !i.is_empty() { - let mut vec = Vec::new(); + let mut vec = Vec::with_capacity(template_fields.len()); for field in template_fields.iter() { let field_res = field.parse_as_field_value(i); if field_res.is_err() { diff --git a/src/variable_versions/v9.rs b/src/variable_versions/v9.rs index ce75a669..fa2bcd19 100644 --- a/src/variable_versions/v9.rs +++ b/src/variable_versions/v9.rs @@ -27,11 +27,11 @@ pub type V9FieldPair = (V9Field, FieldValue); pub type V9FlowRecord = Vec; impl V9Parser { - pub fn parse(&mut self, packet: &[u8]) -> ParsedNetflow { + pub fn parse<'a>(&mut self, packet: &'a [u8]) -> ParsedNetflow<'a> { match V9::parse(packet, self) { Ok((remaining, v9)) => ParsedNetflow::Success { packet: NetflowPacket::V9(v9), - remaining: remaining.to_vec(), + remaining, }, Err(e) => ParsedNetflow::Error { error: NetflowParseError::Partial(PartialParse { @@ -183,22 +183,20 @@ impl FlowSetBody { match id { DATA_TEMPLATE_V9_ID => { let (i, templates) = Templates::parse(i)?; - parser.templates.extend( - templates + for template in &templates.templates { + parser .templates - .iter() - .map(|template| (template.template_id, template.clone())), - ); + .insert(template.template_id, template.clone()); + } Ok((i, FlowSetBody::Template(templates))) } OPTIONS_TEMPLATE_V9_ID => { let (i, options_templates) = OptionsTemplates::parse(i)?; - parser.options_templates.extend( - options_templates - .templates - .iter() - .map(|template| (template.template_id, template.clone())), - ); + for template in &options_templates.templates { + parser + .options_templates + .insert(template.template_id, template.clone()); + } Ok((i, FlowSetBody::OptionsTemplate(options_templates))) } _ => { @@ -301,7 +299,7 @@ impl<'a> ScopeParser { input: &'a [u8], template: &OptionsTemplate, ) -> IResult<&'a [u8], Vec> { - let mut result = Vec::new(); + let mut result = Vec::with_capacity(template.scope_fields.len()); let mut remaining = input; for template_field in template.scope_fields.iter() { let (i, scope_field) = ScopeDataField::parse(remaining, template_field)?; @@ -319,7 +317,7 @@ impl<'a> OptionsFieldParser { input: &'a [u8], template: &OptionsTemplate, ) -> IResult<&'a [u8], Vec>> { - let mut result = Vec::new(); + let mut result = Vec::with_capacity(template.option_fields.len()); let mut remaining = input; for template_field in template.option_fields.iter() { let (i, field_value) = template_field.parse_as_field_value(remaining)?; @@ -492,13 +490,16 @@ impl<'a> FieldParser { mut input: &'a [u8], template: &Template, ) -> IResult<&'a [u8], Vec>> { - let tempalte_total_size = usize::from(template.get_total_size()); - if tempalte_total_size == 0 { + let template_total_size = usize::from(template.get_total_size()); + if template_total_size == 0 { return Err(nom::Err::Error(NomError::new(input, ErrorKind::Verify))); } - let mut res = Vec::new(); - for _ in 0..tempalte_total_size { + // Calculate how many complete records we can parse based on input length + let record_count = input.len() / template_total_size; + let mut res = Vec::with_capacity(record_count); + + for _ in 0..record_count { match Self::parse_data_fields(input, template) { Ok((remaining, record)) => { input = remaining; @@ -534,7 +535,7 @@ impl<'a> FieldParser { mut input: &'a [u8], template: &Template, ) -> IResult<&'a [u8], V9FlowRecord> { - let mut res = Vec::new(); + let mut res = Vec::with_capacity(template.fields.len()); for template_field in template.fields.iter() { let (new_input, field_value) = template_field.parse_as_field_value(input)?;