diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml
index 6bbd951..a998529 100644
--- a/.github/workflows/ci.yml
+++ b/.github/workflows/ci.yml
@@ -63,3 +63,4 @@ jobs:
run: |
python tests/test_e2e.py
python tests/test_modbus_e2e.py
+ python tests/test_modbus_stream_e2e.py
diff --git a/Cargo.toml b/Cargo.toml
index 1d77555..8da8603 100644
--- a/Cargo.toml
+++ b/Cargo.toml
@@ -3,12 +3,12 @@ name = "flowprep"
version = "0.3.0"
edition = "2024"
rust-version = "1.85"
-description = "Convert network telemetry (pcap, flow CSVs, vendor exports) into ML-ready canonical NetFlow parquet"
+description = "Convert network telemetry (pcap, flow CSVs, OT protocols) into ML-ready flow and protocol observations"
license = "Apache-2.0"
repository = "https://github.com/DeepTempo/flowprep"
homepage = "https://deeptempo.ai"
readme = "README.md"
-keywords = ["netflow", "pcap", "parquet", "arrow", "security"]
+keywords = ["netflow", "pcap", "parquet", "modbus", "security"]
categories = ["command-line-utilities", "network-programming"]
exclude = [".github/"]
diff --git a/README.md b/README.md
index cb2a69b..4ca0304 100644
--- a/README.md
+++ b/README.md
@@ -5,7 +5,7 @@
flowprep
- Network telemetry → ML-ready canonical NetFlow parquet.
+ Network telemetry → ML-ready flow and protocol observations.
@@ -17,9 +17,8 @@
---
flowprep converts the network telemetry you actually have — packet captures,
-flow CSVs, vendor exports — into a single, clean, typed, unit-normalized
-parquet table that you can hand directly to a model, a notebook, or a data
-pipeline.
+flow CSVs, vendor exports — into clean, typed, versioned observations that you
+can hand directly to a model, a notebook, or a data pipeline.
It is built and maintained by [DeepTempo](https://deeptempo.ai), where it is
used in production as the ingestion front door for our **LogLM**: flow
@@ -66,7 +65,7 @@ cargo build --release
./demo.sh
```
-Six subcommands:
+Seven subcommands:
```bash
# 1. Raw packet captures -> bidirectional flow records
@@ -77,19 +76,23 @@ flowprep modbus capture.pcap modbus.parquet
# Non-standard server ports are explicit, never guessed
flowprep modbus capture.pcap modbus.parquet --server-port 1502
-# 3. Any aliased flow table (CSV, parquet, Zeek TSV log, Argus .binetflow)
+# 3. Continuously decode a packet-capture stream -> immediate NDJSON events
+tcpdump -U -i -w - 'tcp port 502' \
+ | flowprep modbus-stream --sensor-id plant-a --output modbus.ndjson
+
+# 4. Any aliased flow table (CSV, parquet, Zeek TSV log, Argus .binetflow)
# -> the canonical schema
flowprep canonicalize cic_export.csv flows.parquet
flowprep canonicalize conn.log.labeled flows.parquet
flowprep canonicalize capture.binetflow flows.parquet
-# 4. OCSF Network Activity events (JSON/NDJSON) -> the canonical schema
+# 5. OCSF Network Activity events (JSON/NDJSON) -> the canonical schema
flowprep ocsf network_activity.ndjson flows.parquet
-# 5. nfdump/nfcapd binary flow files -> the canonical schema
+# 6. nfdump/nfcapd binary flow files -> the canonical schema
flowprep nfcapd nfcapd.202401011200 flows.parquet
-# 6. Inspect any parquet file from the terminal, no Python required
+# 7. Inspect any parquet file from the terminal, no Python required
flowprep peek flows.parquet -n 20
```
@@ -198,6 +201,39 @@ Direction is based only on the configured server port (502 by default). Set
`--server-port` for a known non-standard deployment; the decoder does not guess
roles from payload content.
+### Lightweight continuous Modbus sensor
+
+`modbus-stream` is the long-running counterpart to the offline `modbus`
+command. It consumes a continuous PCAP/PCAPNG byte stream from stdin, a FIFO,
+or a file and emits `modbus_stream_event/v1` NDJSON. A `request_observed` event
+is flushed as soon as a complete Modbus request ADU is decoded; a later
+`transaction_observed` event carries the paired response, exception, orphan, or
+timeout result. Consumers do not need to wait for a PLC response before seeing
+the requested operation. Both event types carry the same requested
+`coil_values` or `register_values` as the offline observation contract.
+
+The default hot path is one synchronous decode/output loop in one long-lived
+flowprep process. It does not invoke Arrow or Parquet, create temporary files,
+poll devices, start an async runtime or worker pool, or spawn child processes.
+`--flush-every 1` is the latency-first default. A larger value can improve bulk
+throughput, but delays sparse events and should only be selected deliberately.
+
+Output is synchronously backpressured instead of being placed on an unbounded
+in-process queue. That bounds sensor memory, but a deployment must give the
+capture producer a bounded ring buffer or explicit drop policy if its downstream
+consumer can stall. `sensor_run_id` changes on restart, while `event_sequence`
+is monotonic within a run, so append-mode files can be consumed without treating
+restarted sequence numbers as duplicates.
+
+This MVP accepts capture bytes; it does not open the capture interface itself.
+It can consume an existing sensor's PCAP stream without another DeepTempo
+process, or be paired with a packet-buffered producer such as `tcpdump -U`.
+It does not yet invoke LogLM or include end-to-end model inference in its latency
+measurement. Request timeouts advance on capture timestamps when traffic arrives
+(and all pending requests are finalized at EOF), so an otherwise idle blocking
+stream can delay the terminal timeout event; the immediate request event is not
+delayed.
+
### Zeek logs and research exports
`canonicalize` also reads **Zeek TSV logs** (`conn.log`, including labeled
@@ -233,6 +269,8 @@ Plus any passthrough label columns present in the source.
The Modbus protocol-observation contract is separate from canonical NetFlow and
lives at [`schemas/modbus/v1/schema.json`](schemas/modbus/v1/schema.json).
+The continuous NDJSON envelope is versioned separately at
+[`schemas/modbus_stream/v1/schema.json`](schemas/modbus_stream/v1/schema.json).
## Example: a real research dataset
@@ -316,9 +354,12 @@ cargo build --release
python3 -m venv .venv && .venv/bin/pip install dpkt pyarrow
.venv/bin/python tests/test_e2e.py
.venv/bin/python tests/test_modbus_e2e.py
+.venv/bin/python tests/test_modbus_stream_e2e.py
-# throughput benchmark
+# throughput benchmarks
.venv/bin/python tests/bench_pcap.py
+.venv/bin/python tests/bench_modbus_stream.py
+.venv/bin/python tests/bench_modbus_stream_latency.py
```
## Release
diff --git a/schemas/modbus_stream/v1/schema.json b/schemas/modbus_stream/v1/schema.json
new file mode 100644
index 0000000..5d8a59d
--- /dev/null
+++ b/schemas/modbus_stream/v1/schema.json
@@ -0,0 +1,69 @@
+{
+ "schema_version": "modbus_stream_event/v1",
+ "description": "Low-latency NDJSON envelope for passive Modbus/TCP sensor events. The observation fields are defined by modbus_observation/v1.",
+ "transport": "newline-delimited JSON",
+ "observation_schema": "../../modbus/v1/schema.json",
+ "event_types": {
+ "request_observed": "Emitted as soon as a complete request ADU is decoded. response_status is pending.",
+ "transaction_observed": "Emitted when a response is paired, an orphan response is observed, a request times out, a transaction ID is reused, or the input stream ends."
+ },
+ "envelope_fields": [
+ {
+ "name": "schema_version",
+ "type": "utf8",
+ "nullable": false
+ },
+ {
+ "name": "observation_schema_version",
+ "type": "utf8",
+ "nullable": false
+ },
+ {
+ "name": "event_type",
+ "type": "utf8",
+ "nullable": false
+ },
+ {
+ "name": "sensor_id",
+ "type": "utf8",
+ "nullable": false
+ },
+ {
+ "name": "sensor_run_id",
+ "type": "utf8",
+ "nullable": false,
+ "description": "Changes on every process start; combine with event_sequence for run-local uniqueness."
+ },
+ {
+ "name": "event_sequence",
+ "type": "uint64",
+ "nullable": false
+ },
+ {
+ "name": "emitted_at_usec",
+ "type": "int64",
+ "nullable": false
+ },
+ {
+ "name": "sensor_processing_usec",
+ "type": "uint64",
+ "nullable": false,
+ "description": "Elapsed time from availability of the containing PCAP packet block through parsing and event construction, measured before NDJSON serialization."
+ }
+ ],
+ "hot_path": {
+ "process_model": "single long-lived process",
+ "queueing": "none; output applies synchronous bounded backpressure",
+ "default_flush_every": 1,
+ "request_emission": "before response pairing",
+ "excluded": [
+ "Arrow",
+ "Parquet",
+ "temporary files",
+ "polling",
+ "async runtime",
+ "worker pool",
+ "child process spawning"
+ ]
+ }
+}
diff --git a/src/main.rs b/src/main.rs
index e162e2b..723146d 100644
--- a/src/main.rs
+++ b/src/main.rs
@@ -12,7 +12,7 @@ use clap::{Parser, Subcommand};
#[derive(Parser)]
#[command(
name = "flowprep",
- about = "Convert network telemetry into ML-ready canonical NetFlow parquet"
+ about = "Convert network telemetry into ML-ready flow and protocol observations"
)]
struct Cli {
#[command(subcommand)]
@@ -31,6 +31,27 @@ enum Command {
#[arg(long, default_value_t = 502)]
server_port: u16,
},
+ /// continuously decode a PCAP stream -> low-latency NDJSON events
+ ModbusStream {
+ /// PCAP/PCAPNG input path, or '-' for stdin
+ #[arg(long, default_value = "-")]
+ input: String,
+ /// NDJSON output path (append mode), or '-' for stdout
+ #[arg(long, default_value = "-")]
+ output: String,
+ /// TCP port used by Modbus servers in this capture
+ #[arg(long, default_value_t = 502)]
+ server_port: u16,
+ /// Stable deployment identity attached to every event
+ #[arg(long, default_value = "flowprep-local")]
+ sensor_id: String,
+ /// Time before an unmatched request becomes a terminal event
+ #[arg(long, default_value_t = 5000)]
+ request_timeout_ms: u64,
+ /// Flush NDJSON after this many events; 1 minimizes latency
+ #[arg(long, default_value_t = 1)]
+ flush_every: usize,
+ },
/// aliased parquet/CSV flow table -> canonical parquet
Canonicalize { input: String, output: String },
/// OCSF Network Activity JSON/NDJSON -> canonical parquet
@@ -68,6 +89,22 @@ fn main() {
server_port,
} => modbus::modbus_to_parquet(input, output, *server_port)
.map(|summary| println!("Wrote Modbus {summary} to {output}")),
+ Command::ModbusStream {
+ input,
+ output,
+ server_port,
+ sensor_id,
+ request_timeout_ms,
+ flush_every,
+ } => modbus::modbus_stream_to_ndjson(
+ input,
+ output,
+ *server_port,
+ sensor_id,
+ *request_timeout_ms,
+ *flush_every,
+ )
+ .map(|summary| eprintln!("Modbus stream finished: {summary}")),
Command::Canonicalize { input, output } => canonicalize::canonicalize_file(input, output)
.map(|n| println!("Wrote {n} flows to {output}")),
Command::Ocsf { input, output } => {
diff --git a/src/modbus.rs b/src/modbus.rs
index 8ca7869..022f9c0 100644
--- a/src/modbus.rs
+++ b/src/modbus.rs
@@ -1,4 +1,4 @@
-//! Passive Modbus/TCP decoding from offline PCAP/PCAPNG captures.
+//! Passive Modbus/TCP decoding from offline files or a continuous PCAP stream.
//!
//! This module deliberately writes a protocol-observation schema rather than
//! adding application fields to canonical NetFlow. Direction is inferred only
@@ -6,9 +6,11 @@
//! to, polls, or otherwise interacts with an OT device.
use std::collections::{BTreeMap, HashMap};
-use std::fs::File;
+use std::fs::{File, OpenOptions};
+use std::io::{BufWriter, Cursor, Read, Write};
use std::net::{Ipv4Addr, Ipv6Addr};
use std::sync::Arc;
+use std::time::{Instant, SystemTime, UNIX_EPOCH};
use arrow::array::{
ArrayRef, BooleanArray, BooleanBuilder, Int32Array, Int32Builder, Int64Array, ListBuilder,
@@ -18,7 +20,9 @@ use arrow::datatypes::{DataType, Field, Schema};
use arrow::error::ArrowError;
use arrow::record_batch::RecordBatch;
use etherparse::{NetSlice, SlicedPacket, TransportSlice};
-use pcap_parser::{Block, PcapBlockOwned, PcapError, create_reader};
+use pcap_parser::traits::PcapReaderIterator;
+use pcap_parser::{Block, LegacyPcapReader, PcapBlockOwned, PcapError, PcapNGReader};
+use serde::Serialize;
use serde_json::Value;
use crate::writer::write_parquet;
@@ -29,13 +33,21 @@ const LINKTYPE_ETHERNET: u16 = 1;
const MAX_MODBUS_LENGTH: usize = 254; // unit identifier + PDU
const MAX_PENDING_STREAM_BYTES: usize = 1024 * 1024;
const SCHEMA_VERSION: &str = "modbus_observation/v1";
+const STREAM_SCHEMA_VERSION: &str = "modbus_stream_event/v1";
const DIRECTION_BASIS: &str = "configured_server_port";
+const STREAM_SWEEP_INTERVAL_USEC: i64 = 1_000_000;
+const STREAM_IDLE_TIMEOUT_USEC: i64 = 120 * 1_000_000;
const SCHEMA_JSON: &str = include_str!(concat!(
env!("CARGO_MANIFEST_DIR"),
"/schemas/modbus/v1/schema.json"
));
+const STREAM_SCHEMA_JSON: &str = include_str!(concat!(
+ env!("CARGO_MANIFEST_DIR"),
+ "/schemas/modbus_stream/v1/schema.json"
+));
+
#[derive(Debug, Default, Clone)]
pub struct DecodeSummary {
pub observations: usize,
@@ -53,6 +65,31 @@ pub struct DecodeSummary {
pub incomplete_streams: usize,
}
+#[derive(Debug, Default, Clone)]
+pub struct StreamSummary {
+ pub decode: DecodeSummary,
+ pub stream_events: usize,
+ pub request_events: usize,
+ pub terminal_events: usize,
+ pub average_output_usec: u64,
+ pub maximum_output_usec: u64,
+}
+
+impl std::fmt::Display for StreamSummary {
+ fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
+ write!(
+ f,
+ "{} stream events ({} immediate requests, {} terminal); average output {}us, max {}us; {}",
+ self.stream_events,
+ self.request_events,
+ self.terminal_events,
+ self.average_output_usec,
+ self.maximum_output_usec,
+ self.decode,
+ )
+ }
+}
+
impl std::fmt::Display for DecodeSummary {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(
@@ -144,6 +181,7 @@ struct TcpStreamState {
pending_segments: BTreeMap>,
decoded_bytes: Vec,
pending_warning: Option,
+ last_seen_timestamp: i64,
}
impl TcpStreamState {
@@ -752,11 +790,43 @@ impl Observation {
}
}
+#[derive(Clone, Copy, Debug, PartialEq, Eq)]
+enum DecoderMode {
+ Batch,
+ Stream,
+}
+
+#[derive(Clone, Copy, Debug, PartialEq, Eq)]
+enum StreamEventKind {
+ RequestObserved,
+ TransactionObserved,
+}
+
+impl StreamEventKind {
+ fn as_str(self) -> &'static str {
+ match self {
+ Self::RequestObserved => "request_observed",
+ Self::TransactionObserved => "transaction_observed",
+ }
+ }
+}
+
+#[derive(Clone, Debug)]
+struct StreamEvent {
+ kind: StreamEventKind,
+ observation: Observation,
+}
+
struct Decoder {
server_port: u16,
+ mode: DecoderMode,
streams: HashMap,
pending: HashMap,
observations: Vec,
+ stream_events: Vec,
+ request_timeout_usec: Option,
+ maximum_capture_timestamp: Option,
+ last_sweep_timestamp: Option,
summary: DecodeSummary,
}
@@ -764,14 +834,30 @@ impl Decoder {
fn new(server_port: u16) -> Self {
Self {
server_port,
+ mode: DecoderMode::Batch,
streams: HashMap::new(),
pending: HashMap::new(),
observations: Vec::new(),
+ stream_events: Vec::new(),
+ request_timeout_usec: None,
+ maximum_capture_timestamp: None,
+ last_sweep_timestamp: None,
summary: DecodeSummary::default(),
}
}
+ fn new_streaming(server_port: u16, request_timeout_usec: i64) -> Self {
+ let mut decoder = Self::new(server_port);
+ decoder.mode = DecoderMode::Stream;
+ decoder.request_timeout_usec = Some(request_timeout_usec);
+ decoder
+ }
+
fn ingest(&mut self, packet: TcpPayloadPacket) {
+ self.maximum_capture_timestamp = Some(
+ self.maximum_capture_timestamp
+ .map_or(packet.timestamp, |current| current.max(packet.timestamp)),
+ );
let direction = match (
packet.src_port == self.server_port,
packet.dest_port == self.server_port,
@@ -812,10 +898,14 @@ impl Decoder {
let payload_sequence = packet.sequence_number.wrapping_add(u32::from(packet.syn));
let extract = if packet.payload.is_empty() {
+ if let Some(stream) = self.streams.get_mut(&stream_key) {
+ stream.last_seen_timestamp = packet.timestamp;
+ }
ExtractReport::default()
} else {
self.summary.tcp_payload_packets += 1;
let stream = self.streams.entry(stream_key.clone()).or_default();
+ stream.last_seen_timestamp = packet.timestamp;
let feed = stream.feed(payload_sequence, &packet.payload);
self.summary.retransmitted_segments += feed.retransmitted_segments;
self.summary.out_of_order_segments += feed.out_of_order_segments;
@@ -854,6 +944,7 @@ impl Decoder {
.is_some_and(|stream| stream.has_incomplete_data());
self.summary.incomplete_streams += usize::from(incomplete);
}
+ self.sweep_stream_state();
}
fn handle_request(
@@ -871,12 +962,18 @@ impl Decoder {
};
let observation =
Observation::from_request(conversation, &adu, decoded, timestamp, packet_number);
+ if self.mode == DecoderMode::Stream {
+ self.stream_events.push(StreamEvent {
+ kind: StreamEventKind::RequestObserved,
+ observation: observation.clone(),
+ });
+ }
if let Some(mut replaced) = self.pending.insert(key, observation) {
add_warning(
&mut replaced.parser_warning,
"transaction_id_reused_before_response",
);
- self.observations.push(replaced);
+ self.publish_terminal(replaced);
}
}
@@ -893,7 +990,7 @@ impl Decoder {
unit_id: adu.unit_id,
};
let Some(mut observation) = self.pending.remove(&key) else {
- self.observations.push(Observation::from_orphan_response(
+ self.publish_terminal(Observation::from_orphan_response(
conversation,
&adu,
timestamp,
@@ -957,16 +1054,93 @@ impl Decoder {
add_warning(&mut observation.parser_warning, &warning);
}
}
- self.observations.push(observation);
+ self.publish_terminal(observation);
}
- fn finish(mut self) -> (Vec, DecodeSummary) {
+ fn publish_terminal(&mut self, observation: Observation) {
+ self.summary.observations += 1;
+ if observation.request_seen && observation.response_seen {
+ self.summary.complete += 1;
+ } else if observation.request_seen {
+ self.summary.request_only += 1;
+ } else if observation.response_seen {
+ self.summary.response_only += 1;
+ }
+ if observation.response_status == "exception" {
+ self.summary.exceptions += 1;
+ }
+
+ match self.mode {
+ DecoderMode::Batch => self.observations.push(observation),
+ DecoderMode::Stream => self.stream_events.push(StreamEvent {
+ kind: StreamEventKind::TransactionObserved,
+ observation,
+ }),
+ }
+ }
+
+ fn sweep_stream_state(&mut self) {
+ let (Some(timeout), Some(watermark)) =
+ (self.request_timeout_usec, self.maximum_capture_timestamp)
+ else {
+ return;
+ };
+ if self
+ .last_sweep_timestamp
+ .is_some_and(|previous| watermark.saturating_sub(previous) < STREAM_SWEEP_INTERVAL_USEC)
+ {
+ return;
+ }
+ self.last_sweep_timestamp = Some(watermark);
+
+ let request_cutoff = watermark.saturating_sub(timeout);
+ let expired_keys: Vec = self
+ .pending
+ .iter()
+ .filter(|(_, observation)| {
+ observation
+ .request_timestamp
+ .is_some_and(|timestamp| timestamp <= request_cutoff)
+ })
+ .map(|(key, _)| key.clone())
+ .collect();
+ for key in expired_keys {
+ if let Some(mut observation) = self.pending.remove(&key) {
+ add_warning(&mut observation.parser_warning, "request_timeout");
+ self.publish_terminal(observation);
+ }
+ }
+
+ let stream_cutoff = watermark.saturating_sub(STREAM_IDLE_TIMEOUT_USEC);
+ let mut incomplete_streams = 0;
+ self.streams.retain(|_, stream| {
+ let keep = stream.last_seen_timestamp > stream_cutoff;
+ if !keep && stream.has_incomplete_data() {
+ incomplete_streams += 1;
+ }
+ keep
+ });
+ self.summary.incomplete_streams += incomplete_streams;
+ }
+
+ fn take_stream_events(&mut self) -> Vec {
+ std::mem::take(&mut self.stream_events)
+ }
+
+ fn finalize(&mut self) {
self.summary.incomplete_streams += self
.streams
.values()
.filter(|stream| stream.has_incomplete_data())
.count();
- self.observations.extend(self.pending.into_values());
+ let pending = std::mem::take(&mut self.pending);
+ for observation in pending.into_values() {
+ self.publish_terminal(observation);
+ }
+ }
+
+ fn finish(mut self) -> (Vec, DecodeSummary) {
+ self.finalize();
self.observations.sort_by_key(|observation| {
(
observation.timestamp,
@@ -975,30 +1149,13 @@ impl Decoder {
observation.unit_id,
)
});
-
- self.summary.observations = self.observations.len();
- self.summary.complete = self
- .observations
- .iter()
- .filter(|observation| observation.request_seen && observation.response_seen)
- .count();
- self.summary.request_only = self
- .observations
- .iter()
- .filter(|observation| observation.request_seen && !observation.response_seen)
- .count();
- self.summary.response_only = self
- .observations
- .iter()
- .filter(|observation| !observation.request_seen && observation.response_seen)
- .count();
- self.summary.exceptions = self
- .observations
- .iter()
- .filter(|observation| observation.response_status == "exception")
- .count();
(self.observations, self.summary)
}
+
+ fn finish_stream(mut self) -> (Vec, DecodeSummary) {
+ self.finalize();
+ (self.stream_events, self.summary)
+ }
}
pub fn modbus_to_parquet(input: &str, output: &str, server_port: u16) -> Result {
@@ -1007,8 +1164,28 @@ pub fn modbus_to_parquet(input: &str, output: &str, server_port: u16) -> Result<
}
let file = File::open(input)?;
- let mut reader = create_reader(1 << 20, file)?;
let mut decoder = Decoder::new(server_port);
+ decode_capture(file, &mut decoder, |_, _| Ok(()))?;
+ let (observations, summary) = decoder.finish();
+ if observations.is_empty() {
+ return Err(format!(
+ "no decodable Modbus/TCP observations found on server port {server_port} \
+ ({} payload packets, {} malformed bytes, {} incomplete streams)",
+ summary.tcp_payload_packets, summary.malformed_bytes, summary.incomplete_streams,
+ )
+ .into());
+ }
+ let batch = observations_to_batch(&observations)?;
+ write_parquet(&batch, output)?;
+ Ok(summary)
+}
+
+fn decode_capture(input: R, decoder: &mut Decoder, mut after_packet: F) -> Result<()>
+where
+ R: Read + Send,
+ F: FnMut(&mut Decoder, Instant) -> Result<()>,
+{
+ let mut reader = create_incremental_reader(input)?;
let mut linktype = LINKTYPE_ETHERNET;
let mut legacy_nanos = false;
let mut packet_number = 0_i64;
@@ -1023,6 +1200,7 @@ pub fn modbus_to_parquet(input: &str, output: &str, server_port: u16) -> Result<
}
PcapBlockOwned::Legacy(packet) => {
packet_number += 1;
+ let processing_started = Instant::now();
let fractional_usec = if legacy_nanos {
(packet.ts_usec / 1000) as i64
} else {
@@ -1033,6 +1211,7 @@ pub fn modbus_to_parquet(input: &str, output: &str, server_port: u16) -> Result<
parse_tcp_payload(packet.data, linktype, timestamp, packet_number)
{
decoder.ingest(packet);
+ after_packet(decoder, processing_started)?;
}
}
PcapBlockOwned::NG(Block::InterfaceDescription(description)) => {
@@ -1040,6 +1219,7 @@ pub fn modbus_to_parquet(input: &str, output: &str, server_port: u16) -> Result<
}
PcapBlockOwned::NG(Block::EnhancedPacket(packet)) => {
packet_number += 1;
+ let processing_started = Instant::now();
// pcapng's default if_tsresol is microseconds, matching
// the existing flow reader. Interface-specific options
// are intentionally left for a later capture-layer pass.
@@ -1048,6 +1228,7 @@ pub fn modbus_to_parquet(input: &str, output: &str, server_port: u16) -> Result<
parse_tcp_payload(packet.data, linktype, timestamp, packet_number)
{
decoder.ingest(packet);
+ after_packet(decoder, processing_started)?;
}
}
_ => {}
@@ -1063,19 +1244,312 @@ pub fn modbus_to_parquet(input: &str, output: &str, server_port: u16) -> Result<
Err(error) => return Err(format!("pcap parse error: {error:?}").into()),
}
}
+ Ok(())
+}
- let (observations, summary) = decoder.finish();
- if observations.is_empty() {
- return Err(format!(
- "no decodable Modbus/TCP observations found on server port {server_port} \
- ({} payload packets, {} malformed bytes, {} incomplete streams)",
- summary.tcp_payload_packets, summary.malformed_bytes, summary.incomplete_streams,
- )
- .into());
+fn create_incremental_reader<'a, R>(mut input: R) -> Result>
+where
+ R: Read + Send + 'a,
+{
+ const PCAPNG_MAGIC: [u8; 4] = [0x0a, 0x0d, 0x0d, 0x0a];
+ const PCAP_MAGICS: [[u8; 4]; 4] = [
+ [0xd4, 0xc3, 0xb2, 0xa1],
+ [0xa1, 0xb2, 0xc3, 0xd4],
+ [0x4d, 0x3c, 0xb2, 0xa1],
+ [0xa1, 0xb2, 0x3c, 0x4d],
+ ];
+
+ let mut first = [0_u8; 12];
+ input.read_exact(&mut first[..4])?;
+ if PCAP_MAGICS.contains(&first[..4].try_into().unwrap()) {
+ let mut header = vec![0_u8; 24];
+ header[..4].copy_from_slice(&first[..4]);
+ input.read_exact(&mut header[4..])?;
+ let replay = Cursor::new(header).chain(input);
+ return Ok(Box::new(LegacyPcapReader::new(1 << 20, replay)?));
}
- let batch = observations_to_batch(&observations)?;
- write_parquet(&batch, output)?;
- Ok(summary)
+ if first[..4] != PCAPNG_MAGIC {
+ return Err("capture header is neither PCAP nor PCAPNG".into());
+ }
+
+ input.read_exact(&mut first[4..12])?;
+ let byte_order_magic = &first[8..12];
+ let total_length = if byte_order_magic == [0x4d, 0x3c, 0x2b, 0x1a] {
+ u32::from_le_bytes(first[4..8].try_into().unwrap()) as usize
+ } else if byte_order_magic == [0x1a, 0x2b, 0x3c, 0x4d] {
+ u32::from_be_bytes(first[4..8].try_into().unwrap()) as usize
+ } else {
+ return Err("pcapng section header has an invalid byte-order magic".into());
+ };
+ if !(28..=(1 << 20)).contains(&total_length) {
+ return Err(format!("pcapng section header length {total_length} is invalid").into());
+ }
+ let mut header = vec![0_u8; total_length];
+ header[..12].copy_from_slice(&first);
+ input.read_exact(&mut header[12..])?;
+ let replay = Cursor::new(header).chain(input);
+ Ok(Box::new(PcapNGReader::new(1 << 20, replay)?))
+}
+
+#[derive(Serialize)]
+struct StreamRecord<'a> {
+ schema_version: &'static str,
+ observation_schema_version: &'static str,
+ event_type: &'static str,
+ sensor_id: &'a str,
+ sensor_run_id: &'a str,
+ event_sequence: u64,
+ emitted_at_usec: i64,
+ sensor_processing_usec: u64,
+ timestamp: i64,
+ client_ip: &'a str,
+ client_port: u16,
+ server_ip: &'a str,
+ server_port: u16,
+ transaction_id: u16,
+ unit_id: u8,
+ function_code: u8,
+ function_name: &'a str,
+ operation: &'a str,
+ address: Option,
+ quantity: Option,
+ write_address: Option,
+ write_quantity: Option,
+ coil_values: Option<&'a [bool]>,
+ register_values: Option<&'a [u16]>,
+ diagnostic_subfunction: Option,
+ device_id_code: Option,
+ device_id_object: Option,
+ request_seen: bool,
+ response_seen: bool,
+ request_timestamp: Option,
+ response_timestamp: Option,
+ latency_usec: Option,
+ response_status: &'a str,
+ exception_code: Option,
+ exception_name: Option<&'a str>,
+ vendor_name: Option<&'a str>,
+ product_code: Option<&'a str>,
+ revision: Option<&'a str>,
+ vendor_url: Option<&'a str>,
+ product_name: Option<&'a str>,
+ model_name: Option<&'a str>,
+ user_application_name: Option<&'a str>,
+ request_packet: Option,
+ response_packet: Option,
+ direction_basis: &'static str,
+ parser_warning: Option<&'a str>,
+}
+
+struct NdjsonSink {
+ writer: BufWriter,
+ sensor_id: String,
+ sensor_run_id: String,
+ event_sequence: u64,
+ flush_every: usize,
+ since_flush: usize,
+ request_events: usize,
+ terminal_events: usize,
+ total_output_usec: u128,
+ maximum_output_usec: u64,
+}
+
+impl NdjsonSink {
+ fn new(output: W, sensor_id: &str, flush_every: usize) -> Result {
+ if sensor_id.trim().is_empty() {
+ return Err("sensor ID must not be empty".into());
+ }
+ if flush_every == 0 {
+ return Err("flush-every must be at least 1".into());
+ }
+ let spec: Value =
+ serde_json::from_str(STREAM_SCHEMA_JSON).expect("embedded stream schema is valid");
+ assert_eq!(spec["schema_version"].as_str(), Some(STREAM_SCHEMA_VERSION));
+ let started_at = epoch_microseconds()?;
+ Ok(Self {
+ writer: BufWriter::with_capacity(16 * 1024, output),
+ sensor_id: sensor_id.to_string(),
+ sensor_run_id: format!("{sensor_id}-{}-{started_at}", std::process::id()),
+ event_sequence: 0,
+ flush_every,
+ since_flush: 0,
+ request_events: 0,
+ terminal_events: 0,
+ total_output_usec: 0,
+ maximum_output_usec: 0,
+ })
+ }
+
+ fn write_event(&mut self, event: &StreamEvent, processing_started: Instant) -> Result<()> {
+ let output_started = Instant::now();
+ self.event_sequence += 1;
+ match event.kind {
+ StreamEventKind::RequestObserved => self.request_events += 1,
+ StreamEventKind::TransactionObserved => self.terminal_events += 1,
+ }
+
+ let observation = &event.observation;
+ let response_status = if event.kind == StreamEventKind::RequestObserved {
+ "pending"
+ } else {
+ observation.response_status.as_str()
+ };
+ let record = StreamRecord {
+ schema_version: STREAM_SCHEMA_VERSION,
+ observation_schema_version: SCHEMA_VERSION,
+ event_type: event.kind.as_str(),
+ sensor_id: &self.sensor_id,
+ sensor_run_id: &self.sensor_run_id,
+ event_sequence: self.event_sequence,
+ emitted_at_usec: epoch_microseconds()?,
+ sensor_processing_usec: elapsed_microseconds(processing_started),
+ timestamp: observation.timestamp,
+ client_ip: &observation.conversation.client_ip,
+ client_port: observation.conversation.client_port,
+ server_ip: &observation.conversation.server_ip,
+ server_port: observation.conversation.server_port,
+ transaction_id: observation.transaction_id,
+ unit_id: observation.unit_id,
+ function_code: observation.function_code,
+ function_name: &observation.function_name,
+ operation: &observation.operation,
+ address: observation.address,
+ quantity: observation.quantity,
+ write_address: observation.write_address,
+ write_quantity: observation.write_quantity,
+ coil_values: observation.coil_values.as_deref(),
+ register_values: observation.register_values.as_deref(),
+ diagnostic_subfunction: observation.diagnostic_subfunction,
+ device_id_code: observation.device_id_code,
+ device_id_object: observation.device_id_object,
+ request_seen: observation.request_seen,
+ response_seen: observation.response_seen,
+ request_timestamp: observation.request_timestamp,
+ response_timestamp: observation.response_timestamp,
+ latency_usec: observation.latency_usec,
+ response_status,
+ exception_code: observation.exception_code,
+ exception_name: observation.exception_name.as_deref(),
+ vendor_name: observation.identity.vendor_name.as_deref(),
+ product_code: observation.identity.product_code.as_deref(),
+ revision: observation.identity.revision.as_deref(),
+ vendor_url: observation.identity.vendor_url.as_deref(),
+ product_name: observation.identity.product_name.as_deref(),
+ model_name: observation.identity.model_name.as_deref(),
+ user_application_name: observation.identity.user_application_name.as_deref(),
+ request_packet: observation.request_packet,
+ response_packet: observation.response_packet,
+ direction_basis: DIRECTION_BASIS,
+ parser_warning: observation.parser_warning.as_deref(),
+ };
+ serde_json::to_writer(&mut self.writer, &record)?;
+ self.writer.write_all(b"\n")?;
+ self.since_flush += 1;
+ if self.since_flush >= self.flush_every {
+ self.writer.flush()?;
+ self.since_flush = 0;
+ }
+
+ let output_usec = elapsed_microseconds(output_started);
+ self.total_output_usec += u128::from(output_usec);
+ self.maximum_output_usec = self.maximum_output_usec.max(output_usec);
+ Ok(())
+ }
+
+ fn finish(mut self, decode: DecodeSummary) -> Result {
+ self.writer.flush()?;
+ let stream_events = self.request_events + self.terminal_events;
+ let average_output_usec = if stream_events == 0 {
+ 0
+ } else {
+ (self.total_output_usec / stream_events as u128) as u64
+ };
+ Ok(StreamSummary {
+ decode,
+ stream_events,
+ request_events: self.request_events,
+ terminal_events: self.terminal_events,
+ average_output_usec,
+ maximum_output_usec: self.maximum_output_usec,
+ })
+ }
+}
+
+fn elapsed_microseconds(started: Instant) -> u64 {
+ started.elapsed().as_micros().min(u128::from(u64::MAX)) as u64
+}
+
+fn epoch_microseconds() -> Result {
+ let micros = SystemTime::now().duration_since(UNIX_EPOCH)?.as_micros();
+ i64::try_from(micros).map_err(|_| "system timestamp exceeds int64 microseconds".into())
+}
+
+pub fn modbus_stream_to_ndjson(
+ input: &str,
+ output: &str,
+ server_port: u16,
+ sensor_id: &str,
+ request_timeout_ms: u64,
+ flush_every: usize,
+) -> Result {
+ if server_port == 0 {
+ return Err("Modbus server port must be between 1 and 65535".into());
+ }
+ if request_timeout_ms == 0 {
+ return Err("request timeout must be at least 1 millisecond".into());
+ }
+ let timeout_usec = request_timeout_ms
+ .checked_mul(1000)
+ .and_then(|value| i64::try_from(value).ok())
+ .ok_or("request timeout is too large")?;
+
+ let input: Box = if input == "-" {
+ Box::new(std::io::stdin())
+ } else {
+ Box::new(File::open(input)?)
+ };
+ let output: Box = if output == "-" {
+ Box::new(std::io::stdout())
+ } else {
+ Box::new(OpenOptions::new().create(true).append(true).open(output)?)
+ };
+ stream_reader_to_writer(
+ input,
+ output,
+ server_port,
+ sensor_id,
+ timeout_usec,
+ flush_every,
+ )
+}
+
+fn stream_reader_to_writer(
+ input: R,
+ output: W,
+ server_port: u16,
+ sensor_id: &str,
+ request_timeout_usec: i64,
+ flush_every: usize,
+) -> Result
+where
+ R: Read + Send,
+ W: Write,
+{
+ let mut decoder = Decoder::new_streaming(server_port, request_timeout_usec);
+ let mut sink = NdjsonSink::new(output, sensor_id, flush_every)?;
+ decode_capture(input, &mut decoder, |decoder, processing_started| {
+ for event in decoder.take_stream_events() {
+ sink.write_event(&event, processing_started)?;
+ }
+ Ok(())
+ })?;
+
+ let (events, decode) = decoder.finish_stream();
+ let finalization_started = Instant::now();
+ for event in events {
+ sink.write_event(&event, finalization_started)?;
+ }
+ sink.finish(decode)
}
fn parse_tcp_payload(
@@ -1467,6 +1941,44 @@ mod tests {
bytes
}
+ fn tcp_payload_packet(
+ timestamp: i64,
+ packet_number: i64,
+ request_direction: bool,
+ client_port: u16,
+ sequence_number: u32,
+ payload: Vec,
+ ) -> TcpPayloadPacket {
+ let (src_ip, dest_ip, src_port, dest_port) = if request_direction {
+ (
+ "10.0.0.1".to_string(),
+ "10.0.0.2".to_string(),
+ client_port,
+ 502,
+ )
+ } else {
+ (
+ "10.0.0.2".to_string(),
+ "10.0.0.1".to_string(),
+ 502,
+ client_port,
+ )
+ };
+ TcpPayloadPacket {
+ timestamp,
+ packet_number,
+ src_ip,
+ dest_ip,
+ src_port,
+ dest_port,
+ sequence_number,
+ syn: false,
+ fin: false,
+ rst: false,
+ payload,
+ }
+ }
+
#[test]
fn reassembles_fragmented_and_coalesced_adus() {
let first = adu(1, 7, &[3, 0, 10, 0, 2]);
@@ -1712,4 +2224,123 @@ mod tests {
Some("truncated_function_payload")
);
}
+
+ #[test]
+ fn streaming_emits_request_before_the_response() {
+ let mut decoder = Decoder::new_streaming(502, 5_000_000);
+ decoder.ingest(tcp_payload_packet(
+ 1_000_000,
+ 1,
+ true,
+ 40000,
+ 1000,
+ adu(7, 1, &[16, 0, 100, 0, 2, 4, 0, 1, 0, 2]),
+ ));
+ let request_events = decoder.take_stream_events();
+ assert_eq!(request_events.len(), 1);
+ assert_eq!(request_events[0].kind, StreamEventKind::RequestObserved);
+ assert_eq!(request_events[0].observation.operation, "write");
+ assert!(!request_events[0].observation.response_seen);
+
+ decoder.ingest(tcp_payload_packet(
+ 1_001_000,
+ 2,
+ false,
+ 40000,
+ 5000,
+ adu(7, 1, &[16, 0, 100, 0, 2]),
+ ));
+ let terminal_events = decoder.take_stream_events();
+ assert_eq!(terminal_events.len(), 1);
+ assert_eq!(
+ terminal_events[0].kind,
+ StreamEventKind::TransactionObserved
+ );
+ assert_eq!(terminal_events[0].observation.response_status, "ok");
+ assert_eq!(terminal_events[0].observation.latency_usec, Some(1000));
+ }
+
+ #[test]
+ fn streaming_timeout_bounds_unmatched_request_state() {
+ let mut decoder = Decoder::new_streaming(502, 5_000_000);
+ decoder.ingest(tcp_payload_packet(
+ 0,
+ 1,
+ true,
+ 40000,
+ 1000,
+ adu(1, 1, &[3, 0, 0, 0, 1]),
+ ));
+ decoder.take_stream_events();
+
+ decoder.ingest(tcp_payload_packet(
+ 6_000_000,
+ 2,
+ true,
+ 40001,
+ 2000,
+ adu(2, 1, &[3, 0, 1, 0, 1]),
+ ));
+ let events = decoder.take_stream_events();
+ let timeout = events
+ .iter()
+ .find(|event| event.observation.transaction_id == 1)
+ .expect("the old request is emitted as terminal");
+ assert_eq!(timeout.kind, StreamEventKind::TransactionObserved);
+ assert_eq!(timeout.observation.response_status, "missing_response");
+ assert_eq!(
+ timeout.observation.parser_warning.as_deref(),
+ Some("request_timeout")
+ );
+ assert_eq!(decoder.pending.len(), 1);
+ }
+
+ #[test]
+ fn stream_json_has_provenance_and_pending_semantics() {
+ let conversation = ConversationKey {
+ client_ip: "10.0.0.1".to_string(),
+ client_port: 40000,
+ server_ip: "10.0.0.2".to_string(),
+ server_port: 502,
+ };
+ let raw = RawAdu {
+ transaction_id: 1,
+ unit_id: 1,
+ pdu: vec![6, 0x04, 0x01, 0, 10],
+ warning: None,
+ };
+ let event = StreamEvent {
+ kind: StreamEventKind::RequestObserved,
+ observation: Observation::from_request(
+ conversation,
+ &raw,
+ decode_request(&raw.pdu),
+ 1,
+ 2,
+ ),
+ };
+ let mut output = Vec::new();
+ let summary = {
+ let mut sink = NdjsonSink::new(&mut output, "unit-test-sensor", 1).unwrap();
+ sink.write_event(&event, Instant::now()).unwrap();
+ sink.finish(DecodeSummary::default()).unwrap()
+ };
+ assert_eq!(summary.request_events, 1);
+ let record: Value = serde_json::from_slice(&output).unwrap();
+ assert_eq!(record["schema_version"], STREAM_SCHEMA_VERSION);
+ assert_eq!(record["sensor_id"], "unit-test-sensor");
+ assert_eq!(record["event_type"], "request_observed");
+ assert_eq!(record["response_status"], "pending");
+ assert_eq!(record["event_sequence"], 1);
+ assert_eq!(record["address"], 1025);
+ assert_eq!(record["register_values"], serde_json::json!([10]));
+ assert!(record["coil_values"].is_null());
+ }
+
+ #[test]
+ fn embedded_stream_schema_declares_the_runtime_version() {
+ let spec: Value = serde_json::from_str(STREAM_SCHEMA_JSON).unwrap();
+ assert_eq!(spec["schema_version"], STREAM_SCHEMA_VERSION);
+ assert_eq!(spec["hot_path"]["default_flush_every"], 1);
+ }
}
diff --git a/tests/bench_modbus_stream.py b/tests/bench_modbus_stream.py
new file mode 100644
index 0000000..990bb5f
--- /dev/null
+++ b/tests/bench_modbus_stream.py
@@ -0,0 +1,122 @@
+"""Local throughput comparison for immediate-flush and small-batch stream modes."""
+
+import io
+import os
+import subprocess
+import sys
+import time
+
+import dpkt
+
+from test_modbus_e2e import modbus_adu, tcp_packet
+
+
+_REPO = os.path.join(os.path.dirname(os.path.abspath(__file__)), "..")
+_DEFAULT_BIN = os.path.join(_REPO, "target", "release", "flowprep")
+FLOWPREP_BIN = os.environ.get("FLOWPREP_BIN", _DEFAULT_BIN)
+TRANSACTIONS = int(os.environ.get("MODBUS_BENCH_TRANSACTIONS", "5000"))
+
+
+def build_capture(transaction_count):
+ capture = io.BytesIO()
+ writer = dpkt.pcap.Writer(capture)
+ client_sequence = 1000
+ server_sequence = 5000
+ base = 1_750_000_000.0
+ for index in range(transaction_count):
+ transaction_id = index % 65536
+ address = index % 65536
+ request = modbus_adu(
+ transaction_id,
+ 1,
+ bytes(
+ [
+ 16,
+ (address >> 8) & 0xFF,
+ address & 0xFF,
+ 0,
+ 2,
+ 4,
+ 0,
+ 10,
+ 0,
+ 20,
+ ]
+ ),
+ )
+ response = modbus_adu(
+ transaction_id,
+ 1,
+ bytes([16, (address >> 8) & 0xFF, address & 0xFF, 0, 2]),
+ )
+ writer.writepkt(
+ tcp_packet(
+ "10.20.0.10",
+ "10.20.0.20",
+ 40000,
+ 502,
+ client_sequence,
+ request,
+ ),
+ ts=base + index / 1_000_000,
+ )
+ client_sequence += len(request)
+ writer.writepkt(
+ tcp_packet(
+ "10.20.0.20",
+ "10.20.0.10",
+ 502,
+ 40000,
+ server_sequence,
+ response,
+ ),
+ ts=base + index / 1_000_000 + 0.000001,
+ )
+ server_sequence += len(response)
+ return capture.getvalue()
+
+
+def run(raw, flush_every):
+ started = time.perf_counter()
+ result = subprocess.run(
+ [
+ FLOWPREP_BIN,
+ "modbus-stream",
+ "--sensor-id",
+ "benchmark",
+ "--flush-every",
+ str(flush_every),
+ ],
+ input=raw,
+ stdout=subprocess.DEVNULL,
+ stderr=subprocess.PIPE,
+ )
+ elapsed = time.perf_counter() - started
+ assert result.returncode == 0, result.stderr.decode()
+ events = TRANSACTIONS * 2
+ return {
+ "flush_every": flush_every,
+ "elapsed_seconds": elapsed,
+ "transactions_per_second": TRANSACTIONS / elapsed,
+ "events_per_second": events / elapsed,
+ "wall_usec_per_event": elapsed * 1_000_000 / events,
+ "summary": result.stderr.decode().strip(),
+ }
+
+
+def main():
+ raw = build_capture(TRANSACTIONS)
+ print(f"capture transactions={TRANSACTIONS} bytes={len(raw)}")
+ for flush_every in (1, 64):
+ result = run(raw, flush_every)
+ print(
+ f"flush_every={flush_every} "
+ f"transactions/s={result['transactions_per_second']:.0f} "
+ f"events/s={result['events_per_second']:.0f} "
+ f"wall_us/event={result['wall_usec_per_event']:.2f}"
+ )
+ print(result["summary"])
+
+
+if __name__ == "__main__":
+ sys.exit(main())
diff --git a/tests/bench_modbus_stream_latency.py b/tests/bench_modbus_stream_latency.py
new file mode 100644
index 0000000..36e803c
--- /dev/null
+++ b/tests/bench_modbus_stream_latency.py
@@ -0,0 +1,96 @@
+"""Measure incremental packet-to-NDJSON latency for the long-lived sensor.
+
+This intentionally feeds one complete PCAP record at a time and waits for its
+event. The external measurement includes pipe scheduling, JSON serialization,
+and the default per-event flush; ``sensor_processing_usec`` is also summarized
+to separate the in-process parser path from transport/test-harness overhead.
+"""
+
+import json
+import math
+import os
+import subprocess
+import sys
+import time
+
+from bench_modbus_stream import FLOWPREP_BIN, build_capture
+from test_modbus_stream_e2e import split_legacy_pcap
+
+
+TRANSACTIONS = int(os.environ.get("MODBUS_LATENCY_TRANSACTIONS", "1000"))
+WARMUP_TRANSACTIONS = int(os.environ.get("MODBUS_LATENCY_WARMUP", "10"))
+BUDGET_USEC = int(os.environ.get("MODBUS_LATENCY_BUDGET_USEC", "50000"))
+
+
+def percentile(values, percent):
+ ordered = sorted(values)
+ index = max(0, math.ceil(len(ordered) * percent / 100) - 1)
+ return ordered[index]
+
+
+def summarize(label, values):
+ print(
+ f"{label}: p50={percentile(values, 50)}us "
+ f"p95={percentile(values, 95)}us "
+ f"p99={percentile(values, 99)}us max={max(values)}us"
+ )
+
+
+def main():
+ total_transactions = WARMUP_TRANSACTIONS + TRANSACTIONS
+ header, records = split_legacy_pcap(build_capture(total_transactions))
+ assert len(records) == total_transactions * 2
+
+ process = subprocess.Popen(
+ [FLOWPREP_BIN, "modbus-stream", "--sensor-id", "latency-benchmark"],
+ stdin=subprocess.PIPE,
+ stdout=subprocess.PIPE,
+ stderr=subprocess.PIPE,
+ bufsize=0,
+ )
+ process.stdin.write(header)
+ process.stdin.flush()
+
+ external_usec = []
+ internal_usec = []
+ first_event_usec = None
+ warmup_records = WARMUP_TRANSACTIONS * 2
+ for index, record in enumerate(records):
+ started = time.perf_counter_ns()
+ process.stdin.write(record)
+ process.stdin.flush()
+ line = process.stdout.readline()
+ elapsed_usec = (time.perf_counter_ns() - started) // 1000
+ assert line, "sensor stdout closed before the expected event"
+ event = json.loads(line)
+ assert event["register_values"] == [10, 20]
+ if index == 0:
+ first_event_usec = elapsed_usec
+ if index >= warmup_records:
+ external_usec.append(elapsed_usec)
+ internal_usec.append(event["sensor_processing_usec"])
+
+ process.stdin.close()
+ return_code = process.wait(timeout=5)
+ stderr = process.stderr.read().decode().strip()
+ assert return_code == 0, stderr
+ assert len(external_usec) == TRANSACTIONS * 2
+
+ print(
+ f"transactions={TRANSACTIONS} events={len(external_usec)} "
+ f"budget={BUDGET_USEC}us first_event_with_startup={first_event_usec}us"
+ )
+ summarize("external packet-to-flushed-NDJSON", external_usec)
+ summarize("internal packet parse-to-record", internal_usec)
+ print(stderr)
+
+ p99 = percentile(external_usec, 99)
+ if p99 >= BUDGET_USEC:
+ print(f"LATENCY BUDGET FAILED: p99 {p99}us >= {BUDGET_USEC}us", file=sys.stderr)
+ return 1
+ print(f"LATENCY BUDGET PASSED: p99 {p99}us < {BUDGET_USEC}us")
+ return 0
+
+
+if __name__ == "__main__":
+ sys.exit(main())
diff --git a/tests/test_modbus_stream_e2e.py b/tests/test_modbus_stream_e2e.py
new file mode 100644
index 0000000..be84ec5
--- /dev/null
+++ b/tests/test_modbus_stream_e2e.py
@@ -0,0 +1,285 @@
+"""End-to-end proof for the low-latency, long-lived Modbus sensor path."""
+
+import io
+import json
+import os
+import select
+import struct
+import subprocess
+import sys
+import tempfile
+import time
+
+import dpkt
+
+from test_modbus_e2e import modbus_adu, tcp_packet
+
+
+_REPO = os.path.join(os.path.dirname(os.path.abspath(__file__)), "..")
+_DEFAULT_BIN = os.path.join(_REPO, "target", "release", "flowprep")
+FLOWPREP_BIN = os.environ.get("FLOWPREP_BIN", _DEFAULT_BIN)
+
+
+def build_two_packet_capture():
+ request = modbus_adu(42, 1, bytes([16, 0, 100, 0, 2, 4, 0, 1, 0, 2]))
+ response = modbus_adu(42, 1, bytes([16, 0, 100, 0, 2]))
+ capture = io.BytesIO()
+ writer = dpkt.pcap.Writer(capture)
+ writer.writepkt(
+ tcp_packet("10.20.0.10", "10.20.0.20", 40000, 502, 1000, request),
+ ts=1_750_000_000.0,
+ )
+ writer.writepkt(
+ tcp_packet("10.20.0.20", "10.20.0.10", 502, 40000, 5000, response),
+ ts=1_750_000_000.001,
+ )
+ return capture.getvalue()
+
+
+def build_incremental_capture():
+ """Build one warm-up transaction followed by the measured transaction."""
+ warm_request = modbus_adu(41, 1, bytes([16, 0, 90, 0, 2, 4, 0, 1, 0, 2]))
+ warm_response = modbus_adu(41, 1, bytes([16, 0, 90, 0, 2]))
+ request = modbus_adu(42, 1, bytes([16, 0, 100, 0, 2, 4, 0, 1, 0, 2]))
+ response = modbus_adu(42, 1, bytes([16, 0, 100, 0, 2]))
+ capture = io.BytesIO()
+ writer = dpkt.pcap.Writer(capture)
+ writer.writepkt(
+ tcp_packet("10.20.0.10", "10.20.0.20", 40000, 502, 1000, warm_request),
+ ts=1_750_000_000.0,
+ )
+ writer.writepkt(
+ tcp_packet("10.20.0.20", "10.20.0.10", 502, 40000, 5000, warm_response),
+ ts=1_750_000_000.001,
+ )
+ writer.writepkt(
+ tcp_packet(
+ "10.20.0.10",
+ "10.20.0.20",
+ 40000,
+ 502,
+ 1000 + len(warm_request),
+ request,
+ ),
+ ts=1_750_000_000.002,
+ )
+ writer.writepkt(
+ tcp_packet(
+ "10.20.0.20",
+ "10.20.0.10",
+ 502,
+ 40000,
+ 5000 + len(warm_response),
+ response,
+ ),
+ ts=1_750_000_000.003,
+ )
+ return capture.getvalue()
+
+
+def split_legacy_pcap(raw):
+ magic = raw[:4]
+ if magic in (b"\xd4\xc3\xb2\xa1", b"\x4d\x3c\xb2\xa1"):
+ byte_order = "<"
+ elif magic in (b"\xa1\xb2\xc3\xd4", b"\xa1\xb2\x3c\x4d"):
+ byte_order = ">"
+ else:
+ raise AssertionError(f"unexpected pcap magic: {magic.hex()}")
+
+ header = raw[:24]
+ records = []
+ offset = 24
+ while offset < len(raw):
+ record_header = raw[offset : offset + 16]
+ assert len(record_header) == 16
+ _, _, captured_length, _ = struct.unpack(f"{byte_order}IIII", record_header)
+ end = offset + 16 + captured_length
+ records.append(raw[offset:end])
+ offset = end
+ return header, records
+
+
+def read_event_with_deadline(process, timeout_seconds=1.0):
+ readable, _, _ = select.select([process.stdout], [], [], timeout_seconds)
+ assert readable, "sensor did not emit an event before the deadline"
+ line = process.stdout.readline()
+ assert line, "sensor stdout closed before an event was emitted"
+ return json.loads(line)
+
+
+def prove_incremental_emission():
+ header, records = split_legacy_pcap(build_incremental_capture())
+ assert len(records) == 4
+ process = subprocess.Popen(
+ [
+ FLOWPREP_BIN,
+ "modbus-stream",
+ "--sensor-id",
+ "stream-e2e",
+ "--request-timeout-ms",
+ "1000",
+ ],
+ stdin=subprocess.PIPE,
+ stdout=subprocess.PIPE,
+ stderr=subprocess.PIPE,
+ bufsize=0,
+ )
+ startup_started = time.perf_counter_ns()
+ process.stdin.write(header)
+ process.stdin.write(records[0])
+ process.stdin.flush()
+ warm_request = read_event_with_deadline(process, timeout_seconds=2.0)
+ startup_usec = (time.perf_counter_ns() - startup_started) // 1000
+ assert warm_request["event_type"] == "request_observed"
+ assert warm_request["event_sequence"] == 1
+
+ # Complete one transaction as an explicit readiness handshake. The budget
+ # below therefore measures packet-to-event latency in the already-running
+ # background sensor, not one-time executable startup and dynamic linking.
+ process.stdin.write(records[1])
+ process.stdin.flush()
+ warm_transaction = read_event_with_deadline(process)
+ assert warm_transaction["event_type"] == "transaction_observed"
+ assert warm_transaction["event_sequence"] == 2
+
+ started = time.perf_counter_ns()
+ process.stdin.write(records[2])
+ process.stdin.flush()
+ request = read_event_with_deadline(process)
+ request_latency_usec = (time.perf_counter_ns() - started) // 1000
+
+ assert request["schema_version"] == "modbus_stream_event/v1"
+ assert request["observation_schema_version"] == "modbus_observation/v1"
+ assert request["event_type"] == "request_observed"
+ assert request["response_status"] == "pending"
+ assert request["operation"] == "write"
+ assert request["address"] == 100 and request["quantity"] == 2
+ assert request["register_values"] == [1, 2]
+ assert request["coil_values"] is None
+ assert request["sensor_id"] == "stream-e2e"
+ assert request["event_sequence"] == 3
+ assert request["sensor_processing_usec"] < 50_000
+ assert request_latency_usec < 50_000, (
+ f"request event took {request_latency_usec}us; sensor budget is <50000us"
+ )
+
+ started = time.perf_counter_ns()
+ process.stdin.write(records[3])
+ process.stdin.flush()
+ transaction = read_event_with_deadline(process)
+ response_latency_usec = (time.perf_counter_ns() - started) // 1000
+
+ assert transaction["event_type"] == "transaction_observed"
+ assert transaction["response_status"] == "ok"
+ assert transaction["request_seen"] is True
+ assert transaction["response_seen"] is True
+ assert transaction["register_values"] == [1, 2]
+ assert transaction["latency_usec"] == 1000
+ assert transaction["event_sequence"] == 4
+ assert transaction["sensor_run_id"] == request["sensor_run_id"]
+ assert transaction["sensor_processing_usec"] < 50_000
+ assert response_latency_usec < 50_000, (
+ f"transaction event took {response_latency_usec}us; sensor budget is <50000us"
+ )
+
+ process.stdin.close()
+ return_code = process.wait(timeout=5)
+ stderr = process.stderr.read().decode()
+ assert return_code == 0, stderr
+ assert "4 stream events" in stderr
+ assert "2 immediate requests" in stderr
+ return startup_usec, request_latency_usec, response_latency_usec
+
+
+def prove_restart_safe_append():
+ raw = build_two_packet_capture()
+ with tempfile.TemporaryDirectory(prefix="flowprep_modbus_stream_") as tempdir:
+ capture = os.path.join(tempdir, "input.pcap")
+ output = os.path.join(tempdir, "events.ndjson")
+ with open(capture, "wb") as handle:
+ handle.write(raw)
+
+ command = [
+ FLOWPREP_BIN,
+ "modbus-stream",
+ "--input",
+ capture,
+ "--output",
+ output,
+ "--sensor-id",
+ "restart-e2e",
+ ]
+ first = subprocess.run(command, capture_output=True, text=True)
+ second = subprocess.run(command, capture_output=True, text=True)
+ assert first.returncode == 0, first.stderr
+ assert second.returncode == 0, second.stderr
+
+ with open(output) as handle:
+ events = [json.loads(line) for line in handle]
+ assert len(events) == 4
+ assert [event["event_sequence"] for event in events] == [1, 2, 1, 2]
+ assert all(event["register_values"] == [1, 2] for event in events)
+ first_run = {event["sensor_run_id"] for event in events[:2]}
+ second_run = {event["sensor_run_id"] for event in events[2:]}
+ assert len(first_run) == 1 and len(second_run) == 1
+ assert first_run != second_run
+
+
+def prove_pcapng_input():
+ request = modbus_adu(77, 3, bytes([15, 0, 10, 0, 10, 2, 0x55, 0x03]))
+ response = modbus_adu(77, 3, bytes([15, 0, 10, 0, 10]))
+ with tempfile.TemporaryDirectory(prefix="flowprep_modbus_stream_pcapng_") as tempdir:
+ capture = os.path.join(tempdir, "input.pcapng")
+ with open(capture, "wb") as handle:
+ writer = dpkt.pcapng.Writer(handle)
+ writer.writepkt(
+ tcp_packet("10.20.0.10", "10.20.0.20", 40000, 502, 1000, request),
+ ts=1_750_000_000.0,
+ )
+ writer.writepkt(
+ tcp_packet("10.20.0.20", "10.20.0.10", 502, 40000, 5000, response),
+ ts=1_750_000_000.001,
+ )
+
+ result = subprocess.run(
+ [FLOWPREP_BIN, "modbus-stream", "--input", capture],
+ capture_output=True,
+ text=True,
+ )
+ assert result.returncode == 0, result.stderr
+ events = [json.loads(line) for line in result.stdout.splitlines()]
+ assert [event["event_type"] for event in events] == [
+ "request_observed",
+ "transaction_observed",
+ ]
+ assert events[1]["transaction_id"] == 77
+ assert events[1]["response_status"] == "ok"
+ assert events[0]["coil_values"] == [
+ True,
+ False,
+ True,
+ False,
+ True,
+ False,
+ True,
+ False,
+ True,
+ True,
+ ]
+ assert events[1]["coil_values"] == events[0]["coil_values"]
+ assert events[1]["register_values"] is None
+
+
+def main():
+ startup_usec, request_usec, response_usec = prove_incremental_emission()
+ prove_restart_safe_append()
+ prove_pcapng_input()
+ print(
+ "MODBUS STREAM E2E PASSED "
+ f"(one-time startup {startup_usec}us; "
+ f"request emission {request_usec}us, transaction emission {response_usec}us)"
+ )
+
+
+if __name__ == "__main__":
+ sys.exit(main())