Skip to content

Commit 2833f7c

Browse files
authored
fix: bound cumulative decoded output per message (#305)
* Bound cumulative decoded output per message * Add decoded-output performance gates * Verify benchmark decoded record counts * Reject decoded-output accounting overflow * Resolve decoded-output test lint warnings * fix pending replay body validation * isolate pending replay validation changes * fix: harden cumulative decoded-output limits * fix: classify IPFIX replay trigger boundaries
1 parent fff2ab0 commit 2833f7c

19 files changed

Lines changed: 3476 additions & 224 deletions

‎README.md‎

Lines changed: 29 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -258,6 +258,35 @@ The parser also automatically validates:
258258
- Each occurrence is decoded as a distinct value in the record
259259
- Descriptor order is significant and preserved
260260

261+
### Cumulative Decoded Output Limits
262+
263+
A small NetFlow v9 or IPFIX message can expand into many decoded field values,
264+
especially when a template contains zero-width or very small fields. The parser
265+
therefore applies two cumulative limits to each message, including pending data
266+
replayed when a template arrives:
267+
268+
- 65,536 decoded field values
269+
- 4 MiB of decoded field content
270+
271+
The content-byte limit counts field contents, not IPFIX variable-length prefixes.
272+
Both limits must be greater than zero. A message that exceeds either limit is
273+
rejected with `NetflowError::DecodedOutputLimitExceeded`; no packet from that
274+
message is returned.
275+
276+
```rust
277+
use netflow_parser::NetflowParser;
278+
279+
let parser = NetflowParser::builder()
280+
.with_max_decoded_field_values_per_message(32_768)
281+
.with_max_decoded_field_payload_bytes_per_message(2 * 1024 * 1024)
282+
.build()
283+
.expect("valid limits");
284+
```
285+
286+
Use the `with_v9_*` and `with_ipfix_*` variants to configure the protocols
287+
independently. These limits complement the per-FlowSet record limit; they bound
288+
the combined decoded output of all Sets or FlowSets in one message.
289+
261290
### Template TTL (Time-to-Live)
262291

263292
> **Note:** Only time-based TTL is supported. See [RELEASES.md](RELEASES.md) for details.

‎SECURITY.md‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -124,6 +124,10 @@ let parser = NetflowParser::builder()
124124
// Limit fields per template (DoS protection)
125125
.with_max_field_count(5000)
126126

127+
// Bound cumulative decoded output from each message
128+
.with_max_decoded_field_values_per_message(65_536)
129+
.with_max_decoded_field_payload_bytes_per_message(4 * 1024 * 1024)
130+
127131
// Limit error sample size (prevents memory exhaustion)
128132
.with_max_error_sample_size(256)
129133

@@ -154,6 +158,7 @@ The parser includes several DoS mitigations:
154158

155159
- **Template Field Count Limit:** Default 10,000 fields per template
156160
- **Template Total Size Validation:** Maximum 65,535 bytes per template
161+
- **Cumulative Decoded Output:** Defaults to 65,536 field values and 4 MiB of field content per message
157162
- **Error Sample Size Limit:** Default 256 bytes to prevent memory exhaustion
158163
- **LRU Template Cache:** Prevents unbounded cache growth
159164

‎benches/hot_path_bench.rs‎

Lines changed: 170 additions & 44 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,8 @@
11
use criterion::{BenchmarkId, Criterion, Throughput, criterion_group, criterion_main};
2-
use netflow_parser::NetflowParser;
2+
use netflow_parser::scoped_parser::AutoScopedParser;
3+
use netflow_parser::{NetflowPacket, NetflowParser};
34
use std::hint::black_box;
5+
use std::net::SocketAddr;
46

57
fn v9_template_packet() -> Vec<u8> {
68
vec![
@@ -167,55 +169,179 @@ fn ipfix_data_packet(flow_count: u16) -> Vec<u8> {
167169
packet
168170
}
169171

170-
fn bench_warm_v9_data_hot_path(c: &mut Criterion) {
171-
let template = v9_template_packet();
172-
let mut group = c.benchmark_group("Hot Path V9 Data");
173-
174-
for flow_count in [100u16, 500, 1000] {
175-
let data = v9_data_packet(flow_count);
176-
group.throughput(Throughput::Bytes(data.len() as u64));
177-
group.bench_with_input(BenchmarkId::from_parameter(flow_count), &data, |b, pkt| {
178-
let mut parser = NetflowParser::default();
179-
let template_result = parser.parse_bytes(&template);
180-
assert!(template_result.error.is_none());
181-
assert_eq!(template_result.packets.len(), 1);
182-
183-
b.iter(|| {
184-
let result = parser.parse_bytes(black_box(pkt));
185-
black_box(result.packets.len());
186-
});
187-
});
172+
#[derive(Clone, Copy)]
173+
enum Protocol {
174+
V9,
175+
Ipfix,
176+
}
177+
178+
impl Protocol {
179+
fn name(self) -> &'static str {
180+
match self {
181+
Self::V9 => "v9",
182+
Self::Ipfix => "ipfix",
183+
}
184+
}
185+
186+
fn template_packet(self) -> Vec<u8> {
187+
match self {
188+
Self::V9 => v9_template_packet(),
189+
Self::Ipfix => ipfix_template_packet(),
190+
}
188191
}
189192

190-
group.finish();
193+
fn data_packet(self, flow_count: u16) -> Vec<u8> {
194+
match self {
195+
Self::V9 => v9_data_packet(flow_count),
196+
Self::Ipfix => ipfix_data_packet(flow_count),
197+
}
198+
}
199+
200+
fn assert_decoded_records(self, packets: &[NetflowPacket], expected: usize) {
201+
assert_eq!(packets.len(), 1, "fixture must decode one outer packet");
202+
let actual: usize = match (self, &packets[0]) {
203+
(Self::V9, NetflowPacket::V9(packet)) => packet
204+
.flowsets
205+
.iter()
206+
.map(|flowset| match &flowset.body {
207+
netflow_parser::variable_versions::v9::FlowSetBody::Data(data) => {
208+
data.fields.len()
209+
}
210+
_ => 0,
211+
})
212+
.sum(),
213+
(Self::Ipfix, NetflowPacket::IPFix(packet)) => packet
214+
.flowsets
215+
.iter()
216+
.map(|flowset| match &flowset.body {
217+
netflow_parser::variable_versions::ipfix::FlowSetBody::Data(data) => {
218+
data.fields.len()
219+
}
220+
_ => 0,
221+
})
222+
.sum(),
223+
_ => panic!("fixture decoded as the wrong protocol"),
224+
};
225+
assert_eq!(actual, expected, "fixture decoded-record count mismatch");
226+
}
191227
}
192228

193-
fn bench_warm_ipfix_data_hot_path(c: &mut Criterion) {
194-
let template = ipfix_template_packet();
195-
let mut group = c.benchmark_group("Hot Path IPFIX Data");
196-
197-
for flow_count in [100u16, 500, 1000] {
198-
let data = ipfix_data_packet(flow_count);
199-
group.throughput(Throughput::Bytes(data.len() as u64));
200-
group.bench_with_input(BenchmarkId::from_parameter(flow_count), &data, |b, pkt| {
201-
let mut parser = NetflowParser::default();
202-
let template_result = parser.parse_bytes(&template);
203-
assert!(template_result.error.is_none());
204-
assert_eq!(template_result.packets.len(), 1);
205-
206-
b.iter(|| {
207-
let result = parser.parse_bytes(black_box(pkt));
208-
black_box(result.packets.len());
209-
});
210-
});
229+
#[derive(Clone, Copy)]
230+
enum Scenario {
231+
DirectParse,
232+
DirectIterator,
233+
AutoParse,
234+
AutoIterator,
235+
}
236+
237+
impl Scenario {
238+
fn name(self) -> &'static str {
239+
match self {
240+
Self::DirectParse => "direct/parse",
241+
Self::DirectIterator => "direct/iterator",
242+
Self::AutoParse => "auto/parse",
243+
Self::AutoIterator => "auto/iterator",
244+
}
211245
}
246+
}
212247

213-
group.finish();
248+
fn bench_warmed_hot_paths(c: &mut Criterion) {
249+
let source = SocketAddr::from(([192, 0, 2, 1], 2055));
250+
251+
for protocol in [Protocol::V9, Protocol::Ipfix] {
252+
let template = protocol.template_packet();
253+
for scenario in [
254+
Scenario::DirectParse,
255+
Scenario::DirectIterator,
256+
Scenario::AutoParse,
257+
Scenario::AutoIterator,
258+
] {
259+
let mut group =
260+
c.benchmark_group(format!("Hot Path/{}/{}", protocol.name(), scenario.name()));
261+
262+
for flow_count in [1u16, 1000] {
263+
let data = protocol.data_packet(flow_count);
264+
group.throughput(Throughput::Elements(u64::from(flow_count)));
265+
group.bench_with_input(
266+
BenchmarkId::from_parameter(flow_count),
267+
&data,
268+
|b, packet| match scenario {
269+
Scenario::DirectParse => {
270+
let mut parser = NetflowParser::default();
271+
assert!(parser.parse_bytes(&template).is_ok());
272+
let result = parser.parse_bytes(packet);
273+
assert!(result.is_ok());
274+
protocol.assert_decoded_records(
275+
&result.packets,
276+
usize::from(flow_count),
277+
);
278+
b.iter(|| {
279+
drop(black_box(
280+
parser.parse_bytes(black_box(packet.as_slice())),
281+
));
282+
});
283+
}
284+
Scenario::DirectIterator => {
285+
let mut parser = NetflowParser::default();
286+
assert!(parser.parse_bytes(&template).is_ok());
287+
let packets = parser
288+
.iter_packets(packet)
289+
.map(Result::unwrap)
290+
.collect::<Vec<_>>();
291+
protocol.assert_decoded_records(&packets, usize::from(flow_count));
292+
b.iter(|| {
293+
for result in parser.iter_packets(black_box(packet.as_slice()))
294+
{
295+
black_box(result.unwrap());
296+
}
297+
});
298+
}
299+
Scenario::AutoParse => {
300+
let mut parser = AutoScopedParser::new();
301+
assert!(parser.parse_from_source(source, &template).is_ok());
302+
let result = parser.parse_from_source(source, packet);
303+
assert!(result.is_ok());
304+
protocol.assert_decoded_records(
305+
&result.packets,
306+
usize::from(flow_count),
307+
);
308+
b.iter(|| {
309+
drop(black_box(
310+
parser.parse_from_source(
311+
source,
312+
black_box(packet.as_slice()),
313+
),
314+
));
315+
});
316+
}
317+
Scenario::AutoIterator => {
318+
let mut parser = AutoScopedParser::new();
319+
assert!(parser.parse_from_source(source, &template).is_ok());
320+
let packets = parser
321+
.iter_packets_from_source(source, packet)
322+
.unwrap()
323+
.map(Result::unwrap)
324+
.collect::<Vec<_>>();
325+
protocol.assert_decoded_records(&packets, usize::from(flow_count));
326+
b.iter(|| {
327+
let iterator = parser
328+
.iter_packets_from_source(
329+
source,
330+
black_box(packet.as_slice()),
331+
)
332+
.unwrap();
333+
for result in iterator {
334+
black_box(result.unwrap());
335+
}
336+
});
337+
}
338+
},
339+
);
340+
}
341+
group.finish();
342+
}
343+
}
214344
}
215345

216-
criterion_group!(
217-
benches,
218-
bench_warm_v9_data_hot_path,
219-
bench_warm_ipfix_data_hot_path
220-
);
346+
criterion_group!(benches, bench_warmed_hot_paths);
221347
criterion_main!(benches);

0 commit comments

Comments
 (0)