Skip to main content

huginn_net_http/process/
flow.rs

1use crate::error::HuginnNetHttpError;
2use crate::http::common::HttpProcessor;
3use crate::http::observable::{ObservableHttpRequest, ObservableHttpResponse};
4use crate::http1::process as http1_process;
5use crate::http2::process as http2_process;
6use pnet::packet::ip::IpNextHeaderProtocols;
7use pnet::packet::ipv4::Ipv4Packet;
8use pnet::packet::ipv6::Ipv6Packet;
9use pnet::packet::tcp::TcpPacket;
10use pnet::packet::Packet;
11use std::net::IpAddr;
12#[cfg(any(feature = "p0f-request", feature = "p0f-response"))]
13use std::time::Duration;
14#[cfg(any(feature = "p0f-request", feature = "p0f-response"))]
15use tracing::debug;
16use ttl_cache::TtlCache;
17
18/// FlowKey: (Client IP, Server IP, Client Port, Server Port)
19pub type FlowKey = (IpAddr, IpAddr, u16, u16);
20
21use crate::http::common::HttpParser;
22
23/// Valid first bytes for an HTTP request payload.
24///
25/// Union of first letters of:
26/// - HTTP/1.x methods (RFC 7231): `GET, POST, PUT, DELETE, HEAD, OPTIONS, TRACE, CONNECT, PATCH`
27/// - WebDAV (RFC 4918): `PROPFIND, PROPPATCH, MKCOL, COPY, MOVE, LOCK, UNLOCK`
28/// - WebDAV Versioning (RFC 3253): `REPORT`
29/// - HTTP/2 connection preface (RFC 7540): `"PRI * HTTP/2.0..."` (starts with `P`)
30///
31/// Same fast-reject strategy used by nDPI's `http_fs = "CDGHLMOPRU"` table in
32/// `src/lib/protocols/http.c` (we additionally include `T` for `TRACE`).
33#[cfg(feature = "p0f-request")]
34const HTTP_REQUEST_FIRST_BYTES: &[u8] = b"CDGHLMOPRTU";
35
36/// Maximum legal frame length for an HTTP/2 frame using the default `SETTINGS_MAX_FRAME_SIZE`
37/// (RFC 7540 ยง6.5.2). Used to discriminate HTTP/2 frames from random binary noise.
38#[cfg(feature = "p0f-response")]
39const HTTP2_MAX_FRAME_LEN: u32 = 16384;
40
41/// Minimum payload size we attempt to parse. Below this we treat the data as
42/// incomplete and wait for more TCP segments; avoids paying the trait-dispatch
43/// cost on fragments that obviously can't contain an HTTP request or response.
44#[cfg(feature = "p0f-request")]
45const MIN_HTTP_PAYLOAD_LEN: usize = 4;
46
47/// Cheap pre-filter: does this payload plausibly start an HTTP request (or HTTP/2
48/// connection preface)? Rejects binary protocols (TLS, SSH, etc.) on flows tracked
49/// by `huginn-net-http` before reaching the trait-dispatched parser path.
50///
51/// Also rejects fragments smaller than `MIN_HTTP_PAYLOAD_LEN`, so callers don't
52/// need to do their own length guard; the function is self-validating.
53#[cfg(feature = "p0f-request")]
54fn looks_like_http_request(data: &[u8]) -> bool {
55    if data.len() < MIN_HTTP_PAYLOAD_LEN {
56        return false;
57    }
58    HTTP_REQUEST_FIRST_BYTES.contains(&data[0])
59}
60
61/// Cheap pre-filter: does this payload plausibly start an HTTP response?
62/// - HTTP/1.x responses always start with the literal `"HTTP"` (then `/1.0` or `/1.1`).
63/// - HTTP/2 responses are sequences of frames with no magic prefix, so we accept any
64///   payload whose first 9 bytes look like a valid frame header (reasonable length and
65///   a known frame type byte). Same structural test used by
66///   `crate::http2::process::looks_like_http2_response`.
67#[cfg(feature = "p0f-response")]
68fn looks_like_http_response(data: &[u8]) -> bool {
69    if data.starts_with(b"HTTP") {
70        return true;
71    }
72    if data.len() < 9 {
73        return false;
74    }
75    let frame_length = u32::from_be_bytes([0, data[0], data[1], data[2]]);
76    if frame_length > HTTP2_MAX_FRAME_LEN {
77        return false;
78    }
79    matches!(data[3], 0..=10)
80}
81
82/// HTTP parser that automatically detects and processes different HTTP versions
83pub struct HttpProcessors {
84    parsers: Vec<Box<dyn HttpParser>>,
85}
86
87impl HttpProcessors {
88    pub fn new() -> Self {
89        Self {
90            parsers: vec![Box::new(Http1ParserAdapter::new()), Box::new(Http2ParserAdapter::new())],
91        }
92    }
93
94    /// Parse HTTP request data using the appropriate parser
95    #[cfg(feature = "p0f-request")]
96    #[inline]
97    pub fn parse_request(&self, data: &[u8]) -> Option<ObservableHttpRequest> {
98        // Fast-reject: skip the trait-dispatched parser loop for payloads that
99        // obviously can't start an HTTP request or HTTP/2 preface.
100        if !looks_like_http_request(data) {
101            return None;
102        }
103        for parser in &self.parsers {
104            if parser.can_parse(data) {
105                if let Some(result) = parser.parse_request(data) {
106                    return Some(result);
107                }
108            }
109        }
110        None
111    }
112
113    /// Parse HTTP response data using the appropriate parser
114    #[cfg(feature = "p0f-response")]
115    #[inline]
116    pub fn parse_response(&self, data: &[u8]) -> Option<ObservableHttpResponse> {
117        // Fast-reject: skip the parser loop for payloads that don't look like an
118        // HTTP/1.x response line or a valid HTTP/2 frame header.
119        if !looks_like_http_response(data) {
120            return None;
121        }
122        for parser in &self.parsers {
123            if parser.can_parse(data) {
124                if let Some(result) = parser.parse_response(data) {
125                    return Some(result);
126                }
127            }
128        }
129        None
130    }
131
132    /// Get all supported HTTP versions
133    pub fn supported_versions(&self) -> Vec<crate::http::Version> {
134        self.parsers.iter().map(|p| p.supported_version()).collect()
135    }
136}
137
138/// Adapter that bridges HTTP/1.x processor to the unified HttpParser interface
139struct Http1ParserAdapter {
140    processor: http1_process::Http1Processor,
141}
142
143impl Http1ParserAdapter {
144    fn new() -> Self {
145        Self { processor: http1_process::Http1Processor::new() }
146    }
147}
148
149impl HttpParser for Http1ParserAdapter {
150    fn supported_version(&self) -> crate::http::Version {
151        crate::http::Version::V11
152    }
153
154    fn can_parse(&self, data: &[u8]) -> bool {
155        self.processor.can_process_request(data) || self.processor.can_process_response(data)
156    }
157
158    fn name(&self) -> &'static str {
159        "HTTP/1.x"
160    }
161
162    fn parse_request(&self, data: &[u8]) -> Option<ObservableHttpRequest> {
163        self.processor.process_request(data).ok().flatten()
164    }
165
166    fn parse_response(&self, data: &[u8]) -> Option<ObservableHttpResponse> {
167        self.processor.process_response(data).ok().flatten()
168    }
169}
170
171/// Adapter that bridges HTTP/2 processor to the unified HttpParser interface
172struct Http2ParserAdapter {
173    processor: http2_process::Http2Processor,
174}
175
176impl Http2ParserAdapter {
177    fn new() -> Self {
178        Self { processor: http2_process::Http2Processor::new() }
179    }
180}
181
182impl HttpParser for Http2ParserAdapter {
183    fn supported_version(&self) -> crate::http::Version {
184        crate::http::Version::V20
185    }
186
187    fn can_parse(&self, data: &[u8]) -> bool {
188        self.processor.can_process_request(data) || self.processor.can_process_response(data)
189    }
190
191    fn name(&self) -> &'static str {
192        "HTTP/2"
193    }
194
195    fn parse_request(&self, data: &[u8]) -> Option<ObservableHttpRequest> {
196        self.processor.process_request(data).ok().flatten()
197    }
198
199    fn parse_response(&self, data: &[u8]) -> Option<ObservableHttpResponse> {
200        self.processor.process_response(data).ok().flatten()
201    }
202}
203
204impl Default for HttpProcessors {
205    fn default() -> Self {
206        Self::new()
207    }
208}
209
210pub struct ObservableHttpPackage {
211    #[cfg(feature = "p0f-request")]
212    pub http_request: Option<ObservableHttpRequest>,
213    #[cfg(feature = "p0f-response")]
214    pub http_response: Option<ObservableHttpResponse>,
215}
216
217impl ObservableHttpPackage {
218    #[must_use]
219    pub const fn empty() -> Self {
220        Self {
221            #[cfg(feature = "p0f-request")]
222            http_request: None,
223            #[cfg(feature = "p0f-response")]
224            http_response: None,
225        }
226    }
227}
228
229#[cfg(any(feature = "p0f-request", feature = "p0f-response"))]
230#[derive(Clone)]
231struct TcpData {
232    sequence: u32,
233    data: Vec<u8>,
234}
235
236#[cfg_attr(
237    not(any(feature = "p0f-request", feature = "p0f-response")),
238    allow(dead_code)
239)]
240pub struct TcpFlow {
241    client_ip: IpAddr,
242    server_ip: IpAddr,
243    client_port: u16,
244    server_port: u16,
245    #[cfg(feature = "p0f-request")]
246    client_data: Vec<TcpData>,
247    #[cfg(feature = "p0f-response")]
248    server_data: Vec<TcpData>,
249    #[cfg(feature = "p0f-request")]
250    client_http_parsed: bool,
251    #[cfg(feature = "p0f-response")]
252    server_http_parsed: bool,
253}
254
255#[cfg(any(feature = "p0f-request", feature = "p0f-response"))]
256impl TcpFlow {
257    #[cfg_attr(not(feature = "p0f-request"), allow(unused_variables))]
258    fn init(
259        src_ip: IpAddr,
260        src_port: u16,
261        dst_ip: IpAddr,
262        dst_port: u16,
263        tcp_data: TcpData,
264    ) -> TcpFlow {
265        TcpFlow {
266            client_ip: src_ip,
267            server_ip: dst_ip,
268            client_port: src_port,
269            server_port: dst_port,
270            #[cfg(feature = "p0f-request")]
271            client_data: vec![tcp_data],
272            #[cfg(feature = "p0f-response")]
273            server_data: Vec::new(),
274            #[cfg(feature = "p0f-request")]
275            client_http_parsed: false,
276            #[cfg(feature = "p0f-response")]
277            server_http_parsed: false,
278        }
279    }
280
281    fn is_fully_parsed(&self) -> bool {
282        #[cfg(all(feature = "p0f-request", feature = "p0f-response"))]
283        {
284            self.client_http_parsed && self.server_http_parsed
285        }
286        #[cfg(all(feature = "p0f-request", not(feature = "p0f-response")))]
287        {
288            self.client_http_parsed
289        }
290        #[cfg(all(not(feature = "p0f-request"), feature = "p0f-response"))]
291        {
292            self.server_http_parsed
293        }
294    }
295}
296
297#[cfg(any(feature = "p0f-request", feature = "p0f-response"))]
298fn ordered_payload(data: &[TcpData]) -> Vec<u8> {
299    let mut sorted_data = data.to_vec();
300    sorted_data.sort_by_key(|tcp_data| tcp_data.sequence);
301    let mut full_data = Vec::new();
302    for tcp_data in sorted_data {
303        full_data.extend_from_slice(&tcp_data.data);
304    }
305    full_data
306}
307
308#[inline]
309pub fn process_http_ipv4(
310    packet: &Ipv4Packet,
311    http_flows: &mut TtlCache<FlowKey, TcpFlow>,
312    processors: &HttpProcessors,
313) -> Result<ObservableHttpPackage, HuginnNetHttpError> {
314    if packet.get_next_level_protocol() != IpNextHeaderProtocols::Tcp {
315        return Err(HuginnNetHttpError::UnsupportedProtocol("IPv4".to_string()));
316    }
317    if let Some(tcp) = TcpPacket::new(packet.payload()) {
318        process_tcp_packet(
319            http_flows,
320            tcp,
321            IpAddr::V4(packet.get_source()),
322            IpAddr::V4(packet.get_destination()),
323            processors,
324        )
325    } else {
326        Ok(ObservableHttpPackage::empty())
327    }
328}
329
330#[inline]
331pub fn process_http_ipv6(
332    packet: &Ipv6Packet,
333    http_flows: &mut TtlCache<FlowKey, TcpFlow>,
334    processors: &HttpProcessors,
335) -> Result<ObservableHttpPackage, HuginnNetHttpError> {
336    if packet.get_next_header() != IpNextHeaderProtocols::Tcp {
337        return Err(HuginnNetHttpError::UnsupportedProtocol("IPv6".to_string()));
338    }
339    if let Some(tcp) = TcpPacket::new(packet.payload()) {
340        process_tcp_packet(
341            http_flows,
342            tcp,
343            IpAddr::V6(packet.get_source()),
344            IpAddr::V6(packet.get_destination()),
345            processors,
346        )
347    } else {
348        Ok(ObservableHttpPackage::empty())
349    }
350}
351
352#[cfg_attr(
353    not(any(feature = "p0f-request", feature = "p0f-response")),
354    allow(unused_variables)
355)]
356fn process_tcp_packet(
357    http_flows: &mut TtlCache<FlowKey, TcpFlow>,
358    tcp: TcpPacket,
359    src_ip: IpAddr,
360    dst_ip: IpAddr,
361    processors: &HttpProcessors,
362) -> Result<ObservableHttpPackage, HuginnNetHttpError> {
363    #[cfg(not(any(feature = "p0f-request", feature = "p0f-response")))]
364    {
365        Ok(ObservableHttpPackage::empty())
366    }
367
368    #[cfg(any(feature = "p0f-request", feature = "p0f-response"))]
369    {
370        let src_port: u16 = tcp.get_source();
371        let dst_port: u16 = tcp.get_destination();
372        let mut observable_http_package = ObservableHttpPackage::empty();
373
374        let flow_key: FlowKey = (src_ip, dst_ip, src_port, dst_port);
375        let (tcp_flow, is_client) = {
376            if let Some(flow) = http_flows.get_mut(&flow_key) {
377                (Some(flow), true)
378            } else {
379                let reversed_key: FlowKey = (dst_ip, src_ip, dst_port, src_port);
380                if let Some(flow) = http_flows.get_mut(&reversed_key) {
381                    (Some(flow), false)
382                } else {
383                    (None, false)
384                }
385            }
386        };
387
388        if let Some(flow) = tcp_flow {
389            if !tcp.payload().is_empty() {
390                let tcp_data =
391                    TcpData { sequence: tcp.get_sequence(), data: Vec::from(tcp.payload()) };
392
393                if is_client && src_ip == flow.client_ip && src_port == flow.client_port {
394                    #[cfg(feature = "p0f-request")]
395                    {
396                        if !flow.client_http_parsed {
397                            flow.client_data.push(tcp_data);
398                            let full_data = ordered_payload(&flow.client_data);
399                            match parse_http_request(&full_data, processors) {
400                                Ok(Some(http_request_parsed)) => {
401                                    observable_http_package.http_request =
402                                        Some(http_request_parsed);
403                                    flow.client_http_parsed = true;
404                                }
405                                Ok(None) => {}
406                                Err(_e) => {}
407                            }
408                        } else {
409                            debug!("CLIENT: HTTP already parsed, discarding additional data");
410                        }
411                    }
412                    #[cfg(not(feature = "p0f-request"))]
413                    let _ = (&tcp_data, processors);
414                } else if src_ip == flow.server_ip && src_port == flow.server_port {
415                    #[cfg(feature = "p0f-response")]
416                    {
417                        // Only add data and parse if not already parsed.
418                        if !flow.server_http_parsed {
419                            flow.server_data.push(tcp_data);
420                            let full_data = ordered_payload(&flow.server_data);
421                            match parse_http_response(&full_data, processors) {
422                                Ok(Some(http_response_parsed)) => {
423                                    observable_http_package.http_response =
424                                        Some(http_response_parsed);
425                                    flow.server_http_parsed = true;
426                                }
427                                Ok(None) => {
428                                    debug!("SERVER: Data not complete yet, waiting for more");
429                                }
430                                Err(_e) => {}
431                            }
432                        } else {
433                            debug!("SERVER: HTTP already parsed, discarding additional data");
434                        }
435                    }
436                    #[cfg(not(feature = "p0f-response"))]
437                    let _ = (&tcp_data, processors);
438                }
439
440                if flow.is_fully_parsed() {
441                    debug!("All enabled HTTP sides parsed, removing flow from http_flows early");
442                    http_flows.remove(&flow_key);
443                    return Ok(observable_http_package);
444                }
445
446                // Clean up on connection close
447                if tcp.get_flags()
448                    & (pnet::packet::tcp::TcpFlags::FIN | pnet::packet::tcp::TcpFlags::RST)
449                    != 0
450                {
451                    debug!("Connection closed or reset");
452                    http_flows.remove(&flow_key);
453                }
454            }
455        } else if tcp.get_flags() & pnet::packet::tcp::TcpFlags::SYN != 0 {
456            let tcp_data: TcpData =
457                TcpData { sequence: tcp.get_sequence(), data: Vec::from(tcp.payload()) };
458            let flow: TcpFlow = TcpFlow::init(src_ip, src_port, dst_ip, dst_port, tcp_data);
459            http_flows.insert(flow_key, flow, Duration::new(60, 0));
460        }
461
462        Ok(observable_http_package)
463    }
464}
465
466#[cfg(feature = "p0f-request")]
467fn parse_http_request(
468    data: &[u8],
469    processors: &HttpProcessors,
470) -> Result<Option<ObservableHttpRequest>, HuginnNetHttpError> {
471    match processors.parse_request(data) {
472        Some(request) => {
473            debug!("Successfully parsed HTTP request using polymorphic parser");
474            Ok(Some(request))
475        }
476        None => {
477            debug!("No HTTP parser could handle request data");
478            Ok(None)
479        }
480    }
481}
482
483#[cfg(feature = "p0f-response")]
484fn parse_http_response(
485    data: &[u8],
486    processors: &HttpProcessors,
487) -> Result<Option<ObservableHttpResponse>, HuginnNetHttpError> {
488    match processors.parse_response(data) {
489        Some(response) => {
490            debug!("Successfully parsed HTTP response using polymorphic parser");
491            Ok(Some(response))
492        }
493        None => {
494            debug!("No HTTP parser could handle response data");
495            Ok(None)
496        }
497    }
498}