try to deduplicate db rows
All checks were successful
Build and Deploy MITM Webserver / traffic_target (push) Successful in 0s
Build and Deploy MITM Webserver / build (push) Successful in 11s

This commit is contained in:
2026-03-10 21:33:41 +01:00
parent 32aa17e0cf
commit c2fb35155c
2 changed files with 39 additions and 121 deletions

View File

@@ -11,6 +11,8 @@ import asyncpg
from asyncpg.pool import Pool from asyncpg.pool import Pool
from pydantic import ValidationError from pydantic import ValidationError
from src.Models.etherType import EtherTypeEnum, ethertype_from_int
from src.Models.ip_protocol import protocol_from_number
from src.Models.packets import PacketDBModel from src.Models.packets import PacketDBModel
logger = logging.getLogger("af_packet_sniffer") logger = logging.getLogger("af_packet_sniffer")
@@ -93,6 +95,24 @@ def _attach_derived_fields(payload: Dict[str, Any]) -> None:
if flow_id is not None: if flow_id is not None:
payload["flow_id"] = flow_id payload["flow_id"] = flow_id
if payload.get("raw_present") is None:
payload["raw_present"] = payload.get("raw") is not None
if payload.get("eth_type") in (None, "") and payload.get("eth_type_raw") is not None:
try:
payload["eth_type"] = ethertype_from_int(int(payload["eth_type_raw"]))
except Exception:
payload["eth_type"] = EtherTypeEnum.UNKNOWN
if payload.get("ip_proto") in (None, "") and payload.get("ip_proto_raw") is not None:
try:
payload["ip_proto"] = protocol_from_number(int(payload["ip_proto_raw"]))
except Exception:
payload["ip_proto"] = int(payload["ip_proto_raw"])
if payload.get("app_master_protocol") in (None, "") and payload.get("app_protocol") not in (None, ""):
payload["app_master_protocol"] = payload.get("app_protocol")
class DatabasePool: class DatabasePool:
"""Asyncpg connection pool wrapper used by the packet APIs.""" """Asyncpg connection pool wrapper used by the packet APIs."""
@@ -178,19 +198,15 @@ class DatabasePool:
src_mac, src_mac,
dst_mac, dst_mac,
eth_type_raw, eth_type_raw,
eth_type,
vlan_id, vlan_id,
src_ip, src_ip,
dst_ip, dst_ip,
ip_proto_raw, ip_proto_raw,
ip_proto,
src_port, src_port,
dst_port, dst_port,
length, length,
raw_present,
capture_sources, capture_sources,
app_protocol, app_protocol,
app_master_protocol,
app_category, app_category,
app_confidence, app_confidence,
app_hostname, app_hostname,
@@ -202,8 +218,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,$35,$36,$37,$38::jsonb, $21,$22,$23,$24,$25,$26,$27,$28,$29,$30,$31,$32,$33::jsonb,$34::jsonb,$35::jsonb,$36
$39::jsonb,$40::jsonb,$41
) )
ON CONFLICT (correlation_key) DO UPDATE SET ON CONFLICT (correlation_key) DO UPDATE SET
updated_at = NOW(), updated_at = NOW(),
@@ -229,16 +244,13 @@ class DatabasePool:
src_mac = COALESCE(EXCLUDED.src_mac, packets.src_mac), src_mac = COALESCE(EXCLUDED.src_mac, packets.src_mac),
dst_mac = COALESCE(EXCLUDED.dst_mac, packets.dst_mac), dst_mac = COALESCE(EXCLUDED.dst_mac, packets.dst_mac),
eth_type_raw = COALESCE(EXCLUDED.eth_type_raw, packets.eth_type_raw), eth_type_raw = COALESCE(EXCLUDED.eth_type_raw, packets.eth_type_raw),
eth_type = COALESCE(EXCLUDED.eth_type, packets.eth_type),
vlan_id = COALESCE(EXCLUDED.vlan_id, packets.vlan_id), vlan_id = COALESCE(EXCLUDED.vlan_id, packets.vlan_id),
src_ip = COALESCE(EXCLUDED.src_ip, packets.src_ip), src_ip = COALESCE(EXCLUDED.src_ip, packets.src_ip),
dst_ip = COALESCE(EXCLUDED.dst_ip, packets.dst_ip), dst_ip = COALESCE(EXCLUDED.dst_ip, packets.dst_ip),
ip_proto_raw = COALESCE(EXCLUDED.ip_proto_raw, packets.ip_proto_raw), ip_proto_raw = COALESCE(EXCLUDED.ip_proto_raw, packets.ip_proto_raw),
ip_proto = COALESCE(EXCLUDED.ip_proto, packets.ip_proto),
src_port = COALESCE(EXCLUDED.src_port, packets.src_port), src_port = COALESCE(EXCLUDED.src_port, packets.src_port),
dst_port = COALESCE(EXCLUDED.dst_port, packets.dst_port), dst_port = COALESCE(EXCLUDED.dst_port, packets.dst_port),
length = COALESCE(EXCLUDED.length, packets.length), length = COALESCE(EXCLUDED.length, packets.length),
raw_present = COALESCE(EXCLUDED.raw_present, FALSE) OR COALESCE(packets.raw_present, FALSE),
capture_sources = ( capture_sources = (
SELECT ARRAY( SELECT ARRAY(
SELECT DISTINCT source SELECT DISTINCT source
@@ -249,7 +261,6 @@ class DatabasePool:
) )
), ),
app_protocol = COALESCE(EXCLUDED.app_protocol, packets.app_protocol), app_protocol = COALESCE(EXCLUDED.app_protocol, packets.app_protocol),
app_master_protocol = COALESCE(EXCLUDED.app_master_protocol, packets.app_master_protocol),
app_category = COALESCE(EXCLUDED.app_category, packets.app_category), app_category = COALESCE(EXCLUDED.app_category, packets.app_category),
app_confidence = COALESCE(EXCLUDED.app_confidence, packets.app_confidence), app_confidence = COALESCE(EXCLUDED.app_confidence, packets.app_confidence),
app_hostname = COALESCE(EXCLUDED.app_hostname, packets.app_hostname), app_hostname = COALESCE(EXCLUDED.app_hostname, packets.app_hostname),
@@ -280,19 +291,15 @@ class DatabasePool:
pkt_info.get("src_mac"), pkt_info.get("src_mac"),
pkt_info.get("dst_mac"), pkt_info.get("dst_mac"),
pkt_info.get("eth_type_raw"), pkt_info.get("eth_type_raw"),
_db_text(pkt_info.get("eth_type")),
pkt_info.get("vlan_id"), pkt_info.get("vlan_id"),
pkt_info.get("src_ip"), pkt_info.get("src_ip"),
pkt_info.get("dst_ip"), pkt_info.get("dst_ip"),
pkt_info.get("protocol_raw"), pkt_info.get("protocol_raw"),
_db_text(pkt_info.get("protocol_name") or pkt_info.get("protocol")),
pkt_info.get("src_port"), pkt_info.get("src_port"),
pkt_info.get("dst_port"), pkt_info.get("dst_port"),
pkt_info.get("length"), pkt_info.get("length"),
pkt_info.get("raw_present"),
pkt_info.get("capture_sources"), pkt_info.get("capture_sources"),
pkt_info.get("app_protocol"), pkt_info.get("app_protocol"),
pkt_info.get("app_master_protocol"),
pkt_info.get("app_category"), pkt_info.get("app_category"),
pkt_info.get("app_confidence"), pkt_info.get("app_confidence"),
pkt_info.get("app_hostname"), pkt_info.get("app_hostname"),
@@ -349,22 +356,21 @@ class DatabasePool:
SET SET
updated_at = NOW(), updated_at = NOW(),
app_protocol = COALESCE(packets.app_protocol, $11), app_protocol = COALESCE(packets.app_protocol, $11),
app_master_protocol = COALESCE(packets.app_master_protocol, $12), app_category = COALESCE(packets.app_category, $12),
app_category = COALESCE(packets.app_category, $13), app_confidence = COALESCE(packets.app_confidence, $13),
app_confidence = COALESCE(packets.app_confidence, $14), app_hostname = COALESCE(packets.app_hostname, $14),
app_hostname = COALESCE(packets.app_hostname, $15), app_is_encrypted = COALESCE(packets.app_is_encrypted, $15),
app_is_encrypted = COALESCE(packets.app_is_encrypted, $16),
dpi_metadata = CASE dpi_metadata = CASE
WHEN $17::jsonb IS NULL THEN packets.dpi_metadata WHEN $16::jsonb IS NULL THEN packets.dpi_metadata
WHEN packets.dpi_metadata IS NULL THEN $17::jsonb WHEN packets.dpi_metadata IS NULL THEN $16::jsonb
ELSE packets.dpi_metadata || $17::jsonb ELSE packets.dpi_metadata || $16::jsonb
END, END,
capture_sources = ( capture_sources = (
SELECT ARRAY( SELECT ARRAY(
SELECT DISTINCT source SELECT DISTINCT source
FROM unnest( FROM unnest(
COALESCE(packets.capture_sources, ARRAY[]::text[]) || COALESCE(packets.capture_sources, ARRAY[]::text[]) ||
COALESCE($18::text[], ARRAY[]::text[]) COALESCE($17::text[], ARRAY[]::text[])
) AS source ) AS source
) )
) )
@@ -383,13 +389,12 @@ class DatabasePool:
AND timestamp BETWEEN $9 AND $10 AND timestamp BETWEEN $9 AND $10
AND ( AND (
packets.app_protocol IS NULL packets.app_protocol IS NULL
OR packets.app_master_protocol IS NULL
OR packets.app_category IS NULL OR packets.app_category IS NULL
OR packets.app_confidence IS NULL OR packets.app_confidence IS NULL
OR packets.app_hostname IS NULL OR packets.app_hostname IS NULL
OR packets.app_is_encrypted IS NULL OR packets.app_is_encrypted IS NULL
OR ($17::jsonb IS NOT NULL) OR ($16::jsonb IS NOT NULL)
OR (COALESCE(array_length($18::text[], 1), 0) > 0) OR (COALESCE(array_length($17::text[], 1), 0) > 0)
) )
RETURNING * RETURNING *
""", """,
@@ -404,7 +409,6 @@ class DatabasePool:
lower_bound, lower_bound,
upper_bound, upper_bound,
enrichment.get("app_protocol"), enrichment.get("app_protocol"),
enrichment.get("app_master_protocol"),
enrichment.get("app_category"), enrichment.get("app_category"),
enrichment.get("app_confidence"), enrichment.get("app_confidence"),
enrichment.get("app_hostname"), enrichment.get("app_hostname"),
@@ -465,29 +469,24 @@ class DatabasePool:
THEN COALESCE($7, packets.app_protocol) THEN COALESCE($7, packets.app_protocol)
ELSE packets.app_protocol ELSE packets.app_protocol
END, END,
app_master_protocol = CASE
WHEN packets.app_master_protocol IS NULL OR packets.app_master_protocol IN ('TCP', 'UDP', 'IP', 'IPv6', 'ETH')
THEN COALESCE($8, packets.app_master_protocol)
ELSE packets.app_master_protocol
END,
app_category = CASE app_category = CASE
WHEN packets.app_category IS NULL OR packets.app_category IN ('Transport', 'Network', 'Protocol') WHEN packets.app_category IS NULL OR packets.app_category IN ('Transport', 'Network', 'Protocol')
THEN COALESCE($9, packets.app_category) THEN COALESCE($8, packets.app_category)
ELSE packets.app_category ELSE packets.app_category
END, END,
app_confidence = CASE app_confidence = CASE
WHEN packets.app_protocol IS NULL OR packets.app_protocol IN ('TCP', 'UDP', 'IP', 'IPv6', 'ETH') WHEN packets.app_protocol IS NULL OR packets.app_protocol IN ('TCP', 'UDP', 'IP', 'IPv6', 'ETH')
THEN COALESCE($10, packets.app_confidence) THEN COALESCE($9, packets.app_confidence)
ELSE packets.app_confidence ELSE packets.app_confidence
END, END,
app_hostname = COALESCE(packets.app_hostname, $11), app_hostname = COALESCE(packets.app_hostname, $10),
app_is_encrypted = COALESCE(packets.app_is_encrypted, $12), app_is_encrypted = COALESCE(packets.app_is_encrypted, $11),
capture_sources = ( capture_sources = (
SELECT ARRAY( SELECT ARRAY(
SELECT DISTINCT source SELECT DISTINCT source
FROM unnest( FROM unnest(
COALESCE(packets.capture_sources, ARRAY[]::text[]) || COALESCE(packets.capture_sources, ARRAY[]::text[]) ||
COALESCE($13::text[], ARRAY[]::text[]) COALESCE($12::text[], ARRAY[]::text[])
) AS source ) AS source
) )
) )
@@ -505,14 +504,12 @@ class DatabasePool:
AND ( AND (
packets.app_protocol IS NULL packets.app_protocol IS NULL
OR packets.app_protocol IN ('TCP', 'UDP', 'IP', 'IPv6', 'ETH') OR packets.app_protocol IN ('TCP', 'UDP', 'IP', 'IPv6', 'ETH')
OR packets.app_master_protocol IS NULL
OR packets.app_master_protocol IN ('TCP', 'UDP', 'IP', 'IPv6', 'ETH')
OR packets.app_category IS NULL OR packets.app_category IS NULL
OR packets.app_category IN ('Transport', 'Network', 'Protocol') OR packets.app_category IN ('Transport', 'Network', 'Protocol')
OR packets.app_confidence IS NULL OR packets.app_confidence IS NULL
OR packets.app_hostname IS NULL OR packets.app_hostname IS NULL
OR packets.app_is_encrypted IS NULL OR packets.app_is_encrypted IS NULL
OR (COALESCE(array_length($13::text[], 1), 0) > 0) OR (COALESCE(array_length($12::text[], 1), 0) > 0)
) )
RETURNING * RETURNING *
""", """,
@@ -523,7 +520,6 @@ class DatabasePool:
lower_bound, lower_bound,
upper_bound, upper_bound,
enrichment.get("app_protocol"), enrichment.get("app_protocol"),
enrichment.get("app_master_protocol"),
enrichment.get("app_category"), enrichment.get("app_category"),
enrichment.get("app_confidence"), enrichment.get("app_confidence"),
enrichment.get("app_hostname"), enrichment.get("app_hostname"),

View File

@@ -48,6 +48,7 @@ CREATE TABLE IF NOT EXISTS packets (
correlation_key TEXT UNIQUE, correlation_key TEXT UNIQUE,
packet_id TEXT, packet_id TEXT,
packet_uid TEXT, packet_uid TEXT,
flow_id TEXT,
correlation_source VARCHAR(32), correlation_source VARCHAR(32),
skb_mark BIGINT, skb_mark BIGINT,
@@ -66,7 +67,6 @@ CREATE TABLE IF NOT EXISTS packets (
src_mac MACADDR, src_mac MACADDR,
dst_mac MACADDR, dst_mac MACADDR,
eth_type_raw INTEGER, eth_type_raw INTEGER,
eth_type VARCHAR(64),
-- VLAN -- VLAN
vlan_id INTEGER, vlan_id INTEGER,
@@ -75,7 +75,6 @@ CREATE TABLE IF NOT EXISTS packets (
src_ip INET, src_ip INET,
dst_ip INET, dst_ip INET,
ip_proto_raw INTEGER, ip_proto_raw INTEGER,
ip_proto VARCHAR(64),
-- Transport layer -- Transport layer
src_port INTEGER, src_port INTEGER,
@@ -83,12 +82,10 @@ CREATE TABLE IF NOT EXISTS packets (
-- Packet metadata -- Packet metadata
length INTEGER, length INTEGER,
raw_present BOOLEAN DEFAULT FALSE,
capture_sources TEXT[], capture_sources TEXT[],
-- DPI / flow enrichment metadata -- DPI / flow enrichment metadata
app_protocol VARCHAR(128), app_protocol VARCHAR(128),
app_master_protocol VARCHAR(128),
app_category VARCHAR(128), app_category VARCHAR(128),
app_confidence VARCHAR(64), app_confidence VARCHAR(64),
app_hostname VARCHAR(255), app_hostname VARCHAR(255),
@@ -101,84 +98,9 @@ CREATE TABLE IF NOT EXISTS packets (
-- Full packet dump -- Full packet dump
raw BYTEA raw BYTEA
); );
ALTER TABLE packets ADD COLUMN IF NOT EXISTS updated_at TIMESTAMPTZ DEFAULT NOW();
ALTER TABLE packets ADD COLUMN IF NOT EXISTS correlation_key TEXT;
ALTER TABLE packets ADD COLUMN IF NOT EXISTS packet_id TEXT;
ALTER TABLE packets ADD COLUMN IF NOT EXISTS packet_uid TEXT;
ALTER TABLE packets ADD COLUMN IF NOT EXISTS correlation_source VARCHAR(32);
ALTER TABLE packets ADD COLUMN IF NOT EXISTS skb_mark BIGINT;
ALTER TABLE packets ADD COLUMN IF NOT EXISTS capture_iface VARCHAR(64);
ALTER TABLE packets DROP COLUMN IF EXISTS iface;
ALTER TABLE packets ADD COLUMN IF NOT EXISTS app_protocol VARCHAR(128);
ALTER TABLE packets ADD COLUMN IF NOT EXISTS app_master_protocol VARCHAR(128);
ALTER TABLE packets ADD COLUMN IF NOT EXISTS app_category VARCHAR(128);
ALTER TABLE packets ADD COLUMN IF NOT EXISTS app_confidence VARCHAR(64);
ALTER TABLE packets ADD COLUMN IF NOT EXISTS app_hostname VARCHAR(255);
ALTER TABLE packets ADD COLUMN IF NOT EXISTS app_is_encrypted BOOLEAN;
ALTER TABLE packets ADD COLUMN IF NOT EXISTS app_risk_score INTEGER;
ALTER TABLE packets ADD COLUMN IF NOT EXISTS dpi_metadata JSONB;
ALTER TABLE packets ALTER COLUMN dpi_metadata TYPE JSONB USING
CASE
WHEN dpi_metadata IS NULL THEN NULL
WHEN pg_typeof(dpi_metadata)::text = 'jsonb' THEN dpi_metadata
ELSE dpi_metadata::jsonb
END;
ALTER TABLE packets ADD COLUMN IF NOT EXISTS capture_metadata JSONB;
ALTER TABLE packets ALTER COLUMN capture_metadata TYPE JSONB USING
CASE
WHEN capture_metadata IS NULL THEN NULL
WHEN pg_typeof(capture_metadata)::text = 'jsonb' THEN capture_metadata
ELSE capture_metadata::jsonb
END;
ALTER TABLE packets ADD COLUMN IF NOT EXISTS eth_type_raw INTEGER;
ALTER TABLE packets ADD COLUMN IF NOT EXISTS ip_proto_raw INTEGER;
ALTER TABLE packets ADD COLUMN IF NOT EXISTS raw_present BOOLEAN DEFAULT FALSE;
ALTER TABLE packets ADD COLUMN IF NOT EXISTS capture_sources TEXT[];
ALTER TABLE packets ADD COLUMN IF NOT EXISTS ingress_if VARCHAR(64);
ALTER TABLE packets ADD COLUMN IF NOT EXISTS egress_if VARCHAR(64);
ALTER TABLE packets DROP COLUMN IF EXISTS observed_ifaces;
ALTER TABLE packets ADD COLUMN IF NOT EXISTS verdict VARCHAR(32);
ALTER TABLE packets ADD COLUMN IF NOT EXISTS verdict_reason VARCHAR(128);
ALTER TABLE packets ADD COLUMN IF NOT EXISTS verdict_confidence VARCHAR(32);
ALTER TABLE packets ADD COLUMN IF NOT EXISTS ingress_seen_at TIMESTAMPTZ;
ALTER TABLE packets ADD COLUMN IF NOT EXISTS egress_seen_at TIMESTAMPTZ;
ALTER TABLE packets ADD COLUMN IF NOT EXISTS verdict_seen_at TIMESTAMPTZ;
ALTER TABLE packets ADD COLUMN IF NOT EXISTS telemetry_metadata JSONB;
ALTER TABLE packets ALTER COLUMN telemetry_metadata TYPE JSONB USING
CASE
WHEN telemetry_metadata IS NULL THEN NULL
WHEN pg_typeof(telemetry_metadata)::text = 'jsonb' THEN telemetry_metadata
ELSE telemetry_metadata::jsonb
END;
ALTER TABLE packets DROP COLUMN IF EXISTS direction;
UPDATE packets
SET raw_present = (raw IS NOT NULL)
WHERE raw_present IS DISTINCT FROM (raw IS NOT NULL);
UPDATE packets
SET correlation_key = CASE
WHEN packet_id IS NOT NULL AND packet_id <> '' THEN 'pid:' || packet_id
WHEN packet_uid IS NOT NULL AND packet_uid <> '' THEN 'uid:' || packet_uid
ELSE 'row:' || id::text
END
WHERE correlation_key IS NULL OR correlation_key = '';
UPDATE packets
SET correlation_source = CASE
WHEN packet_id IS NOT NULL AND packet_id <> '' THEN 'kernel_mark'
ELSE 'legacy_hash'
END
WHERE correlation_source IS NULL OR correlation_source = '';
UPDATE packets
SET capture_sources = ARRAY_REMOVE(ARRAY[
CASE WHEN raw IS NOT NULL THEN 'af_packet' END,
CASE WHEN telemetry_metadata IS NOT NULL THEN 'telemetry' END
], NULL)
WHERE capture_sources IS NULL OR array_length(capture_sources, 1) IS NULL;
UPDATE packets
SET updated_at = COALESCE(verdict_seen_at, egress_seen_at, ingress_seen_at, timestamp, NOW())
WHERE updated_at IS NULL;
CREATE UNIQUE INDEX IF NOT EXISTS idx_packets_correlation_key ON packets(correlation_key); CREATE UNIQUE INDEX IF NOT EXISTS idx_packets_correlation_key ON packets(correlation_key);
CREATE UNIQUE INDEX IF NOT EXISTS idx_packets_packet_uid ON packets(packet_uid); CREATE UNIQUE INDEX IF NOT EXISTS idx_packets_packet_uid ON packets(packet_uid);
CREATE INDEX IF NOT EXISTS idx_packets_flow_id ON packets(flow_id);
EOF EOF