diff --git a/backend/src/utilities/database.py b/backend/src/utilities/database.py index be23f44..53a740a 100644 --- a/backend/src/utilities/database.py +++ b/backend/src/utilities/database.py @@ -356,6 +356,34 @@ class DatabasePool: async with self._pool.acquire() as conn: rows = await conn.fetch( """ + WITH candidate_rows AS ( + SELECT id + FROM packets + WHERE + ( + ($1::int IS NOT NULL AND ip_proto_raw = $1) + OR + ($1::int IS NULL AND $2::int IS NOT NULL AND eth_type_raw = $2) + ) + AND (capture_iface = $3 OR ingress_if = $3 OR egress_if = $3) + AND src_ip = $4::inet + AND dst_ip = $5::inet + AND COALESCE(src_port, 0) = $6 + AND COALESCE(dst_port, 0) = $7 + AND length = $8 + AND timestamp BETWEEN $9 AND $10 + AND ( + app_protocol IS NULL + OR app_category IS NULL + OR app_confidence IS NULL + OR app_hostname IS NULL + OR app_is_encrypted IS NULL + OR ($16::jsonb IS NOT NULL) + OR (COALESCE(array_length($17::text[], 1), 0) > 0) + ) + ORDER BY id + FOR UPDATE SKIP LOCKED + ) UPDATE packets SET updated_at = NOW(), @@ -392,29 +420,9 @@ class DatabasePool: ) AS source ) ) - WHERE - ( - ($1::int IS NOT NULL AND ip_proto_raw = $1) - OR - ($1::int IS NULL AND $2::int IS NOT NULL AND eth_type_raw = $2) - ) - AND (capture_iface = $3 OR ingress_if = $3 OR egress_if = $3) - AND src_ip = $4::inet - AND dst_ip = $5::inet - AND COALESCE(src_port, 0) = $6 - AND COALESCE(dst_port, 0) = $7 - AND length = $8 - AND timestamp BETWEEN $9 AND $10 - AND ( - packets.app_protocol IS NULL - OR packets.app_category IS NULL - OR packets.app_confidence IS NULL - OR packets.app_hostname IS NULL - OR packets.app_is_encrypted IS NULL - OR ($16::jsonb IS NOT NULL) - OR (COALESCE(array_length($17::text[], 1), 0) > 0) - ) - RETURNING * + FROM candidate_rows + WHERE packets.id = candidate_rows.id + RETURNING packets.* """, protocol, eth_type_raw, @@ -479,6 +487,33 @@ class DatabasePool: async with self._pool.acquire() as conn: rows = await conn.fetch( """ + WITH candidate_rows AS ( + SELECT id + FROM packets + WHERE + ip_proto_raw = $1 + AND (capture_iface = $2 OR ingress_if = $2 OR egress_if = $2) + AND timestamp BETWEEN $5 AND $6 + AND ( + CASE + WHEN $3 = 'tcp' THEN COALESCE(dpi_metadata -> 'tcp' ->> 'stream', '') + WHEN $3 = 'udp' THEN COALESCE(dpi_metadata -> 'udp' ->> 'stream', '') + ELSE '' + END + ) = $4 + AND ( + app_protocol IS NULL + OR app_protocol IN ('TCP', 'UDP', 'IP', 'IPv6', 'ETH') + OR app_category IS NULL + OR app_category IN ('Transport', 'Network', 'Protocol') + OR app_confidence IS NULL + OR app_hostname IS NULL + OR app_is_encrypted IS NULL + OR (COALESCE(array_length($12::text[], 1), 0) > 0) + ) + ORDER BY id + FOR UPDATE SKIP LOCKED + ) UPDATE packets SET updated_at = NOW(), @@ -512,28 +547,9 @@ class DatabasePool: ) AS source ) ) - WHERE - ip_proto_raw = $1 - AND (capture_iface = $2 OR ingress_if = $2 OR egress_if = $2) - AND timestamp BETWEEN $5 AND $6 - AND ( - CASE - WHEN $3 = 'tcp' THEN COALESCE(packets.dpi_metadata -> 'tcp' ->> 'stream', '') - WHEN $3 = 'udp' THEN COALESCE(packets.dpi_metadata -> 'udp' ->> 'stream', '') - ELSE '' - END - ) = $4 - AND ( - packets.app_protocol IS NULL - OR packets.app_protocol IN ('TCP', 'UDP', 'IP', 'IPv6', 'ETH') - OR packets.app_category IS NULL - OR packets.app_category IN ('Transport', 'Network', 'Protocol') - OR packets.app_confidence IS NULL - OR packets.app_hostname IS NULL - OR packets.app_is_encrypted IS NULL - OR (COALESCE(array_length($12::text[], 1), 0) > 0) - ) - RETURNING * + FROM candidate_rows + WHERE packets.id = candidate_rows.id + RETURNING packets.* """, protocol, iface,