db deadlock fix
This commit is contained in:
@@ -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,
|
||||
|
||||
Reference in New Issue
Block a user