Skip to content

Commit c881e2c

Browse files
meirdevmeir
andauthored
fix the time_flow_start and end. add fallback to sampling rate (#6)
Co-authored-by: meir <meire@flow-sec.com>
1 parent afcfb01 commit c881e2c

8 files changed

Lines changed: 80 additions & 18 deletions

File tree

‎Cargo.lock‎

Lines changed: 4 additions & 4 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

‎crates/rustflow/Cargo.toml‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
11
[package]
22
name = "rustflow"
3-
version = "0.1.0"
3+
version = "0.3.1"
44
edition = "2024"
55
description = "High-performance flow collector library for NetFlow, IPFIX, and sFlow"
66
license = "BSD-3-Clause"

‎crates/rustflow_collector/Cargo.toml‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
11
[package]
22
name = "rustflow_collector"
3-
version = "0.3.0"
3+
version = "0.3.1"
44
edition = "2024"
55

66
[dependencies]

‎crates/rustflow_core/Cargo.toml‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
11
[package]
22
name = "rustflow_core"
3-
version = "0.3.0"
3+
version = "0.3.1"
44
edition = "2024"
55

66
[dependencies]

‎crates/rustflow_core/src/common/common_flow.rs‎

Lines changed: 62 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -296,6 +296,17 @@ pub struct NetFlowV9Context<'a> {
296296
}
297297

298298
impl NetFlowV9Context<'_> {
299+
/// Convert uptime-based timestamp to absolute nanoseconds since epoch.
300+
///
301+
/// In NetFlow v9, `FlowStartSysUpTime` and `FlowEndSysUpTime` are system
302+
/// uptime values in milliseconds. We convert them to absolute time using:
303+
/// `absolute_time = unix_seconds - (system_uptime - uptime_value)`
304+
fn uptime_to_absolute_ns(&self, uptime_ms: u32) -> Option<i64> {
305+
let unix_time_ns = self.header.unix_seconds.timestamp_nanos_opt()?;
306+
let offset_ms = self.header.system_uptime as i64 - uptime_ms as i64;
307+
Some(unix_time_ns - (offset_ms * 1_000_000))
308+
}
309+
299310
pub fn convert(&self, record: &V9DataRecord, template_id: u16) -> CommonFlow {
300311
use InformationElement::*;
301312

@@ -346,10 +357,22 @@ impl NetFlowV9Context<'_> {
346357
flow.bgp_next_hop = Some(IpAddr::V4(*addr));
347358
}
348359
}
349-
FlowStartSysUpTime | FlowStartSeconds | FlowStartMilliseconds => {
360+
FlowStartSysUpTime => {
361+
flow.time_flow_start_ns =
362+
extract_u32(value).and_then(|v| self.uptime_to_absolute_ns(v));
363+
}
364+
FlowEndSysUpTime => {
365+
flow.time_flow_end_ns =
366+
extract_u32(value).and_then(|v| self.uptime_to_absolute_ns(v));
367+
}
368+
FlowStartSeconds
369+
| FlowStartMilliseconds
370+
| FlowStartMicroseconds
371+
| FlowStartNanoseconds => {
350372
flow.time_flow_start_ns = extract_datetime_ns(value);
351373
}
352-
FlowEndSysUpTime | FlowEndSeconds | FlowEndMilliseconds => {
374+
FlowEndSeconds | FlowEndMilliseconds | FlowEndMicroseconds
375+
| FlowEndNanoseconds => {
353376
flow.time_flow_end_ns = extract_datetime_ns(value);
354377
}
355378
SourceIpv6Address => {
@@ -374,7 +397,7 @@ impl NetFlowV9Context<'_> {
374397
}
375398
IcmpTypeIpv4 | IcmpTypeIpv6 => flow.icmp_type = extract_u8(value),
376399
IcmpCodeIpv4 | IcmpCodeIpv6 => flow.icmp_code = extract_u8(value),
377-
SamplingInterval => {
400+
SamplingInterval | SamplingPacketInterval | SamplerRandomInterval => {
378401
if self.sampling_rate.is_none() {
379402
flow.sampling_rate = extract_u32(value);
380403
}
@@ -414,8 +437,13 @@ impl NetFlowV9Context<'_> {
414437

415438
pub fn extract_v9_sampling_rate(record: &V9DataRecord) -> Option<u32> {
416439
let sampling_interval_id: u16 = InformationElement::SamplingInterval.into();
440+
let sampling_packet_interval_id: u16 = InformationElement::SamplingPacketInterval.into();
441+
let sampler_random_interval_id: u16 = InformationElement::SamplerRandomInterval.into();
417442
for (field_type, _, value) in &record.0 {
418-
if *field_type == sampling_interval_id {
443+
if *field_type == sampling_interval_id
444+
|| *field_type == sampling_packet_interval_id
445+
|| *field_type == sampler_random_interval_id
446+
{
419447
return extract_u32(value);
420448
}
421449
}
@@ -483,6 +511,15 @@ pub struct IpfixContext<'a> {
483511
}
484512

485513
impl IpfixContext<'_> {
514+
/// Convert delta microseconds to absolute nanoseconds since epoch.
515+
///
516+
/// Delta fields represent time backwards from export_time:
517+
/// `absolute_time = export_time - delta_microseconds`
518+
fn delta_to_absolute_ns(&self, delta_us: u32) -> Option<i64> {
519+
let export_time_ns = self.header.export_time.timestamp_nanos_opt()?;
520+
Some(export_time_ns - (delta_us as i64 * 1_000))
521+
}
522+
486523
pub fn convert(&self, record: &IpfixDataRecord, template_id: u16) -> CommonFlow {
487524
use InformationElement::*;
488525

@@ -533,12 +570,24 @@ impl IpfixContext<'_> {
533570
flow.bgp_next_hop = Some(IpAddr::V4(*addr));
534571
}
535572
}
536-
FlowStartSysUpTime | FlowStartSeconds | FlowStartMilliseconds => {
573+
FlowStartSeconds
574+
| FlowStartMilliseconds
575+
| FlowStartMicroseconds
576+
| FlowStartNanoseconds => {
537577
flow.time_flow_start_ns = ipfix_extract_datetime_ns(value);
538578
}
539-
FlowEndSysUpTime | FlowEndSeconds | FlowEndMilliseconds => {
579+
FlowEndSeconds | FlowEndMilliseconds | FlowEndMicroseconds
580+
| FlowEndNanoseconds => {
540581
flow.time_flow_end_ns = ipfix_extract_datetime_ns(value);
541582
}
583+
FlowStartDeltaMicroseconds => {
584+
flow.time_flow_start_ns =
585+
ipfix_extract_u32(value).and_then(|v| self.delta_to_absolute_ns(v));
586+
}
587+
FlowEndDeltaMicroseconds => {
588+
flow.time_flow_end_ns =
589+
ipfix_extract_u32(value).and_then(|v| self.delta_to_absolute_ns(v));
590+
}
542591
SourceIpv6Address => {
543592
if let IpfixFieldValue::Ipv6Address(addr) = value {
544593
flow.src_addr = Some(IpAddr::V6(*addr));
@@ -561,7 +610,7 @@ impl IpfixContext<'_> {
561610
}
562611
IcmpTypeIpv4 | IcmpTypeIpv6 => flow.icmp_type = ipfix_extract_u8(value),
563612
IcmpCodeIpv4 | IcmpCodeIpv6 => flow.icmp_code = ipfix_extract_u8(value),
564-
SamplingInterval => {
613+
SamplingInterval | SamplingPacketInterval | SamplerRandomInterval => {
565614
if self.sampling_rate.is_none() {
566615
flow.sampling_rate = ipfix_extract_u32(value);
567616
}
@@ -651,8 +700,13 @@ fn ipfix_extract_datetime_ns(value: &IpfixFieldValue) -> Option<i64> {
651700

652701
pub fn extract_ipfix_sampling_rate(record: &IpfixDataRecord) -> Option<u32> {
653702
let sampling_interval_id: u16 = InformationElement::SamplingInterval.into();
703+
let sampling_packet_interval_id: u16 = InformationElement::SamplingPacketInterval.into();
704+
let sampler_random_interval_id: u16 = InformationElement::SamplerRandomInterval.into();
654705
for (_, field_type, _, value) in &record.0 {
655-
if *field_type == sampling_interval_id {
706+
if *field_type == sampling_interval_id
707+
|| *field_type == sampling_packet_interval_id
708+
|| *field_type == sampler_random_interval_id
709+
{
656710
return ipfix_extract_u32(value);
657711
}
658712
}

‎crates/rustflow_core/src/common/information_element.rs‎

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -32,6 +32,7 @@ pub enum InformationElement {
3232
FlowLabelIpv6 = 31,
3333
IcmpTypeCodeIpv4 = 32,
3434
SamplingInterval = 34,
35+
SamplerRandomInterval = 50,
3536
MinimumTtl = 52,
3637
MaximumTtl = 53,
3738
FragmentIdentification = 54,
@@ -49,6 +50,12 @@ pub enum InformationElement {
4950
FlowEndSeconds = 151,
5051
FlowStartMilliseconds = 152,
5152
FlowEndMilliseconds = 153,
53+
FlowStartMicroseconds = 154,
54+
FlowEndMicroseconds = 155,
55+
FlowStartNanoseconds = 156,
56+
FlowEndNanoseconds = 157,
57+
FlowStartDeltaMicroseconds = 158,
58+
FlowEndDeltaMicroseconds = 159,
5259
IcmpTypeIpv4 = 176,
5360
IcmpCodeIpv4 = 177,
5461
IcmpTypeIpv6 = 178,

‎crates/rustflow_core/src/netflow_v9/parser.rs‎

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -237,7 +237,8 @@ pub struct NetFlowV9Packet {
237237
pub struct Header {
238238
pub version: u16,
239239
pub count: u16,
240-
pub system_uptime: DateTime<Utc>,
240+
/// Milliseconds since device boot
241+
pub system_uptime: u32,
241242
pub unix_seconds: DateTime<Utc>,
242243
pub sequence_number: u32,
243244
pub source_id: u32,
@@ -246,7 +247,7 @@ pub struct Header {
246247
fn parse_header(input: &[u8]) -> IResult<&[u8], Header> {
247248
let (input, version) = verify_version(input, NETFLOW_V9_VERSION)?;
248249
let (input, count) = be_u16(input)?;
249-
let (input, system_uptime) = timestamp_secs(input)?;
250+
let (input, system_uptime) = be_u32(input)?;
250251
let (input, unix_seconds) = timestamp_secs(input)?;
251252
let (input, sequence_number) = be_u32(input)?;
252253
let (input, source_id) = be_u32(input)?;

‎crates/rustflow_exporter/Cargo.toml‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
11
[package]
22
name = "rustflow_exporter"
3-
version = "0.3.0"
3+
version = "0.3.1"
44
edition = "2024"
55

66
[dependencies]

0 commit comments

Comments
 (0)