add session id to db
This commit is contained in:
@@ -32,6 +32,7 @@ def _serialize_row_for_broadcast(row: Dict[str, Any]) -> Dict[str, Any]:
|
|||||||
serialized = dict(row)
|
serialized = dict(row)
|
||||||
_normalize_json_fields(serialized)
|
_normalize_json_fields(serialized)
|
||||||
_attach_derived_fields(serialized)
|
_attach_derived_fields(serialized)
|
||||||
|
serialized.pop("capture_session_id", None)
|
||||||
|
|
||||||
raw_val = serialized.get("raw")
|
raw_val = serialized.get("raw")
|
||||||
if isinstance(raw_val, (bytes, bytearray)):
|
if isinstance(raw_val, (bytes, bytearray)):
|
||||||
@@ -184,6 +185,7 @@ class DatabasePool:
|
|||||||
packet_id,
|
packet_id,
|
||||||
packet_uid,
|
packet_uid,
|
||||||
flow_id,
|
flow_id,
|
||||||
|
capture_session_id,
|
||||||
correlation_source,
|
correlation_source,
|
||||||
skb_mark,
|
skb_mark,
|
||||||
capture_iface,
|
capture_iface,
|
||||||
@@ -218,7 +220,7 @@ class DatabasePool:
|
|||||||
raw
|
raw
|
||||||
) VALUES(
|
) VALUES(
|
||||||
$1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11,$12,$13,$14,$15,$16,$17,$18,$19,$20,
|
$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
|
ON CONFLICT (correlation_key) DO UPDATE SET
|
||||||
updated_at = NOW(),
|
updated_at = NOW(),
|
||||||
@@ -230,6 +232,7 @@ class DatabasePool:
|
|||||||
packet_id = COALESCE(EXCLUDED.packet_id, packets.packet_id),
|
packet_id = COALESCE(EXCLUDED.packet_id, packets.packet_id),
|
||||||
packet_uid = COALESCE(EXCLUDED.packet_uid, packets.packet_uid),
|
packet_uid = COALESCE(EXCLUDED.packet_uid, packets.packet_uid),
|
||||||
flow_id = COALESCE(EXCLUDED.flow_id, packets.flow_id),
|
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),
|
correlation_source = COALESCE(EXCLUDED.correlation_source, packets.correlation_source),
|
||||||
skb_mark = COALESCE(EXCLUDED.skb_mark, packets.skb_mark),
|
skb_mark = COALESCE(EXCLUDED.skb_mark, packets.skb_mark),
|
||||||
capture_iface = COALESCE(EXCLUDED.capture_iface, packets.capture_iface),
|
capture_iface = COALESCE(EXCLUDED.capture_iface, packets.capture_iface),
|
||||||
@@ -277,6 +280,7 @@ class DatabasePool:
|
|||||||
pkt_info.get("packet_id"),
|
pkt_info.get("packet_id"),
|
||||||
pkt_info.get("packet_uid"),
|
pkt_info.get("packet_uid"),
|
||||||
pkt_info.get("flow_id"),
|
pkt_info.get("flow_id"),
|
||||||
|
pkt_info.get("capture_session_id"),
|
||||||
pkt_info.get("correlation_source"),
|
pkt_info.get("correlation_source"),
|
||||||
pkt_info.get("skb_mark"),
|
pkt_info.get("skb_mark"),
|
||||||
pkt_info.get("capture_iface"),
|
pkt_info.get("capture_iface"),
|
||||||
@@ -365,6 +369,20 @@ class DatabasePool:
|
|||||||
WHEN packets.dpi_metadata IS NULL THEN $16::jsonb
|
WHEN packets.dpi_metadata IS NULL THEN $16::jsonb
|
||||||
ELSE packets.dpi_metadata || $16::jsonb
|
ELSE packets.dpi_metadata || $16::jsonb
|
||||||
END,
|
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 = (
|
capture_sources = (
|
||||||
SELECT ARRAY(
|
SELECT ARRAY(
|
||||||
SELECT DISTINCT source
|
SELECT DISTINCT source
|
||||||
@@ -481,6 +499,10 @@ class DatabasePool:
|
|||||||
END,
|
END,
|
||||||
app_hostname = COALESCE(packets.app_hostname, $10),
|
app_hostname = COALESCE(packets.app_hostname, $10),
|
||||||
app_is_encrypted = COALESCE(packets.app_is_encrypted, $11),
|
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 = (
|
capture_sources = (
|
||||||
SELECT ARRAY(
|
SELECT ARRAY(
|
||||||
SELECT DISTINCT source
|
SELECT DISTINCT source
|
||||||
|
|||||||
@@ -49,6 +49,7 @@ CREATE TABLE IF NOT EXISTS packets (
|
|||||||
packet_id TEXT,
|
packet_id TEXT,
|
||||||
packet_uid TEXT,
|
packet_uid TEXT,
|
||||||
flow_id TEXT,
|
flow_id TEXT,
|
||||||
|
capture_session_id TEXT,
|
||||||
correlation_source VARCHAR(32),
|
correlation_source VARCHAR(32),
|
||||||
skb_mark BIGINT,
|
skb_mark BIGINT,
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user