diff --git a/backend/src/utilities/database.py b/backend/src/utilities/database.py index b7b81c3..be23f44 100644 --- a/backend/src/utilities/database.py +++ b/backend/src/utilities/database.py @@ -32,6 +32,7 @@ def _serialize_row_for_broadcast(row: Dict[str, Any]) -> Dict[str, Any]: serialized = dict(row) _normalize_json_fields(serialized) _attach_derived_fields(serialized) + serialized.pop("capture_session_id", None) raw_val = serialized.get("raw") if isinstance(raw_val, (bytes, bytearray)): @@ -184,6 +185,7 @@ class DatabasePool: packet_id, packet_uid, flow_id, + capture_session_id, correlation_source, skb_mark, capture_iface, @@ -218,7 +220,7 @@ class DatabasePool: raw ) VALUES( $1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11,$12,$13,$14,$15,$16,$17,$18,$19,$20, - $21,$22,$23,$24,$25,$26,$27,$28,$29,$30,$31,$32,$33,$34::jsonb,$35::jsonb,$36::jsonb,$37 + $21,$22,$23,$24,$25,$26,$27,$28,$29,$30,$31,$32,$33,$34,$35::jsonb,$36::jsonb,$37::jsonb,$38 ) ON CONFLICT (correlation_key) DO UPDATE SET updated_at = NOW(), @@ -230,6 +232,7 @@ class DatabasePool: packet_id = COALESCE(EXCLUDED.packet_id, packets.packet_id), packet_uid = COALESCE(EXCLUDED.packet_uid, packets.packet_uid), flow_id = COALESCE(EXCLUDED.flow_id, packets.flow_id), + capture_session_id = COALESCE(EXCLUDED.capture_session_id, packets.capture_session_id), correlation_source = COALESCE(EXCLUDED.correlation_source, packets.correlation_source), skb_mark = COALESCE(EXCLUDED.skb_mark, packets.skb_mark), capture_iface = COALESCE(EXCLUDED.capture_iface, packets.capture_iface), @@ -277,6 +280,7 @@ class DatabasePool: pkt_info.get("packet_id"), pkt_info.get("packet_uid"), pkt_info.get("flow_id"), + pkt_info.get("capture_session_id"), pkt_info.get("correlation_source"), pkt_info.get("skb_mark"), pkt_info.get("capture_iface"), @@ -365,6 +369,20 @@ class DatabasePool: WHEN packets.dpi_metadata IS NULL THEN $16::jsonb ELSE packets.dpi_metadata || $16::jsonb END, + flow_id = COALESCE( + packets.flow_id, + CASE + WHEN COALESCE($16::jsonb -> 'tcp' ->> 'stream', $16::jsonb -> 'tshark' ->> 'tcp_stream') IS NOT NULL + AND COALESCE($16::jsonb -> 'tcp' ->> 'stream', $16::jsonb -> 'tshark' ->> 'tcp_stream') <> '' THEN + COALESCE(NULLIF(packets.capture_session_id, '') || ':', '') || + 'tcp:' || COALESCE($16::jsonb -> 'tcp' ->> 'stream', $16::jsonb -> 'tshark' ->> 'tcp_stream') + WHEN COALESCE($16::jsonb -> 'udp' ->> 'stream', $16::jsonb -> 'tshark' ->> 'udp_stream') IS NOT NULL + AND COALESCE($16::jsonb -> 'udp' ->> 'stream', $16::jsonb -> 'tshark' ->> 'udp_stream') <> '' THEN + COALESCE(NULLIF(packets.capture_session_id, '') || ':', '') || + 'udp:' || COALESCE($16::jsonb -> 'udp' ->> 'stream', $16::jsonb -> 'tshark' ->> 'udp_stream') + ELSE NULL + END + ), capture_sources = ( SELECT ARRAY( SELECT DISTINCT source @@ -481,6 +499,10 @@ class DatabasePool: END, app_hostname = COALESCE(packets.app_hostname, $10), app_is_encrypted = COALESCE(packets.app_is_encrypted, $11), + flow_id = COALESCE( + packets.flow_id, + COALESCE(NULLIF(packets.capture_session_id, '') || ':', '') || $3 || ':' || $4 + ), capture_sources = ( SELECT ARRAY( SELECT DISTINCT source diff --git a/setup_database.sh b/setup_database.sh index 0b8065c..94a04b2 100755 --- a/setup_database.sh +++ b/setup_database.sh @@ -49,6 +49,7 @@ CREATE TABLE IF NOT EXISTS packets ( packet_id TEXT, packet_uid TEXT, flow_id TEXT, + capture_session_id TEXT, correlation_source VARCHAR(32), skb_mark BIGINT,