test with new analysis layout
All checks were successful
Build and Deploy MITM Webserver / traffic_target (push) Successful in 1s
Build and Deploy MITM Webserver / build (push) Successful in 12s

This commit is contained in:
2026-03-31 22:15:01 +02:00
parent cedc52eb8d
commit 72bab93aa3
6 changed files with 2632 additions and 145 deletions

View File

@@ -4,6 +4,8 @@ import asyncio
import base64
import json
import logging
import math
import statistics
from datetime import datetime, timezone
from typing import Any, Dict, List, Optional
@@ -63,6 +65,46 @@ def _analysis_ethernet_protocol_name(eth_type_raw: Any) -> Optional[str]:
return None
def _analysis_host_key(ip_address: Any, mac_address: Any) -> str:
return f"{ip_address or 'no-ip'}|{mac_address or 'no-mac'}"
def _classify_discovery_activity(
protocol_name: Optional[str],
ethernet_protocol_name: Optional[str],
ip_protocol_name: Optional[str],
src_port: Optional[int],
dst_port: Optional[int],
) -> Optional[str]:
protocol_upper = str(protocol_name or "").upper()
ethernet_upper = str(ethernet_protocol_name or "").upper()
ip_upper = str(ip_protocol_name or "").upper()
ports = {int(port) for port in (src_port, dst_port) if port is not None}
if protocol_upper == "ARP" or ethernet_upper == "ARP":
return "ARP"
if protocol_upper in {"DHCP", "DHCPV6"} or ports & {67, 68, 546, 547}:
return "DHCP"
if protocol_upper == "MDNS" or 5353 in ports:
return "mDNS"
if protocol_upper == "SSDP" or 1900 in ports:
return "SSDP"
if protocol_upper == "LLMNR" or 5355 in ports:
return "LLMNR"
if protocol_upper == "NBNS" or ports & {137, 138}:
return "NBNS"
if protocol_upper == "ICMPV6" or ip_upper == "ICMPV6":
return "ICMPv6 Discovery"
return None
def _safe_ratio(numerator: float, denominator: float) -> float:
if denominator <= 0:
return 0.0
return numerator / denominator
def _serialize_row_for_broadcast(row: Dict[str, Any]) -> Dict[str, Any]:
serialized = dict(row)
_normalize_json_fields(serialized)
@@ -1208,6 +1250,834 @@ class DatabasePool:
),
)
async def analyze_conversations(
self,
*,
since: Optional[datetime] = None,
limit: int = 300,
) -> List[Dict[str, Any]]:
"""Aggregate directional conversations between endpoints."""
if self._pool is None:
await self.init_pool()
try:
async with self._pool.acquire() as conn:
rows = await conn.fetch(
"""
WITH aggregated AS (
SELECT
ingress_if,
egress_if,
src_ip::text AS src_ip_address,
src_mac::text AS src_mac_address,
dst_ip::text AS dst_ip_address,
dst_mac::text AS dst_mac_address,
src_port,
dst_port,
NULLIF(app_protocol::text, '') AS app_protocol_name,
ip_proto_raw,
eth_type_raw,
ARRAY_REMOVE(ARRAY_AGG(DISTINCT NULLIF(app_hostname::text, '')), NULL) AS hostnames,
COUNT(*) AS packet_count,
COALESCE(SUM(length), 0) AS byte_count,
MIN(timestamp) AS first_seen,
MAX(timestamp) AS last_seen,
SUM(CASE WHEN COALESCE(NULLIF(verdict::text, ''), 'unknown') = 'accept' THEN 1 ELSE 0 END) AS accept_count,
SUM(CASE WHEN COALESCE(NULLIF(verdict::text, ''), 'unknown') = 'drop' THEN 1 ELSE 0 END) AS drop_count,
SUM(CASE WHEN COALESCE(NULLIF(verdict::text, ''), 'unknown') = 'reject' THEN 1 ELSE 0 END) AS reject_count,
SUM(CASE WHEN COALESCE(NULLIF(verdict::text, ''), 'unknown') NOT IN ('accept', 'drop', 'reject') THEN 1 ELSE 0 END) AS unknown_count
FROM packets
WHERE ($1::timestamptz IS NULL OR timestamp >= $1)
AND (src_ip IS NOT NULL OR src_mac IS NOT NULL)
AND (dst_ip IS NOT NULL OR dst_mac IS NOT NULL)
GROUP BY
ingress_if,
egress_if,
src_ip::text,
src_mac::text,
dst_ip::text,
dst_mac::text,
src_port,
dst_port,
app_protocol_name,
ip_proto_raw,
eth_type_raw
)
SELECT *
FROM aggregated
ORDER BY packet_count DESC, byte_count DESC, last_seen DESC
LIMIT $2
""",
since,
limit,
)
except Exception:
logger.exception("DB conversation analysis failed")
raise
result: List[Dict[str, Any]] = []
for row in rows:
record = dict(row)
protocol_name = _analysis_protocol_name(
record.get("app_protocol_name"),
record.get("ip_proto_raw"),
record.get("eth_type_raw"),
)
first_seen_raw = record.get("first_seen")
last_seen_raw = record.get("last_seen")
result.append(
{
"ingress_interface": record.get("ingress_if"),
"egress_interface": record.get("egress_if"),
"src_ip_address": record.get("src_ip_address"),
"src_mac_address": record.get("src_mac_address"),
"dst_ip_address": record.get("dst_ip_address"),
"dst_mac_address": record.get("dst_mac_address"),
"src_port": record.get("src_port"),
"dst_port": record.get("dst_port"),
"protocol": str(protocol_name),
"ethernet_protocol": _analysis_ethernet_protocol_name(record.get("eth_type_raw")),
"ip_protocol": _analysis_ip_protocol_name(record.get("ip_proto_raw")),
"hostnames": list(record.get("hostnames") or []),
"packet_count": int(record.get("packet_count") or 0),
"byte_count": int(record.get("byte_count") or 0),
"first_seen": first_seen_raw.isoformat() if hasattr(first_seen_raw, "isoformat") else first_seen_raw,
"last_seen": last_seen_raw.isoformat() if hasattr(last_seen_raw, "isoformat") else last_seen_raw,
"accept_count": int(record.get("accept_count") or 0),
"drop_count": int(record.get("drop_count") or 0),
"reject_count": int(record.get("reject_count") or 0),
"unknown_count": int(record.get("unknown_count") or 0),
}
)
return result
async def analyze_host_intelligence(
self,
*,
since: Optional[datetime] = None,
limit_hosts: int = 40,
) -> List[Dict[str, Any]]:
"""Build host-centric intelligence with peers, services, and hostname enrichment."""
if self._pool is None:
await self.init_pool()
try:
async with self._pool.acquire() as conn:
host_rows = await conn.fetch(
"""
WITH observations AS (
SELECT
src_ip::text AS ip_address,
src_mac::text AS mac_address,
ingress_if AS iface,
COALESCE(length, 0) AS packet_length,
timestamp,
'source' AS role
FROM packets
WHERE ($1::timestamptz IS NULL OR timestamp >= $1)
AND ingress_if IS NOT NULL
AND (src_ip IS NOT NULL OR src_mac IS NOT NULL)
UNION ALL
SELECT
dst_ip::text AS ip_address,
dst_mac::text AS mac_address,
egress_if AS iface,
COALESCE(length, 0) AS packet_length,
timestamp,
'destination' AS role
FROM packets
WHERE ($1::timestamptz IS NULL OR timestamp >= $1)
AND egress_if IS NOT NULL
AND (dst_ip IS NOT NULL OR dst_mac IS NOT NULL)
)
SELECT
ip_address,
mac_address,
COUNT(*) AS packet_count,
COALESCE(SUM(packet_length), 0) AS byte_count,
MIN(timestamp) AS first_seen,
MAX(timestamp) AS last_seen,
ARRAY_REMOVE(ARRAY_AGG(DISTINCT iface), NULL) AS interfaces,
SUM(CASE WHEN role = 'source' THEN 1 ELSE 0 END) AS source_count,
SUM(CASE WHEN role = 'destination' THEN 1 ELSE 0 END) AS destination_count
FROM observations
WHERE COALESCE(mac_address, '') <> 'ff:ff:ff:ff:ff:ff'
GROUP BY ip_address, mac_address
ORDER BY packet_count DESC, byte_count DESC, last_seen DESC
LIMIT $2
""",
since,
limit_hosts,
)
host_keys = [
_analysis_host_key(record.get("ip_address"), record.get("mac_address"))
for record in (dict(row) for row in host_rows)
]
if not host_keys:
return []
detail_rows = await conn.fetch(
"""
SELECT
src_ip::text AS src_ip_address,
src_mac::text AS src_mac_address,
dst_ip::text AS dst_ip_address,
dst_mac::text AS dst_mac_address,
ingress_if,
egress_if,
src_port,
dst_port,
NULLIF(app_protocol::text, '') AS app_protocol_name,
ip_proto_raw,
eth_type_raw,
NULLIF(app_hostname::text, '') AS app_hostname,
COALESCE(length, 0) AS packet_length,
timestamp,
COALESCE(NULLIF(verdict::text, ''), 'unknown') AS verdict_name
FROM packets
WHERE ($1::timestamptz IS NULL OR timestamp >= $1)
AND (
(COALESCE(src_ip::text, '') || '|' || COALESCE(src_mac::text, '')) = ANY($2::text[])
OR (COALESCE(dst_ip::text, '') || '|' || COALESCE(dst_mac::text, '')) = ANY($2::text[])
)
ORDER BY timestamp DESC
""",
since,
host_keys,
)
except Exception:
logger.exception("DB host intelligence analysis failed")
raise
host_index: Dict[str, Dict[str, Any]] = {}
for row in host_rows:
record = dict(row)
host_key = _analysis_host_key(record.get("ip_address"), record.get("mac_address"))
first_seen_raw = record.get("first_seen")
last_seen_raw = record.get("last_seen")
host_index[host_key] = {
"ip_address": record.get("ip_address"),
"mac_address": record.get("mac_address"),
"packet_count": int(record.get("packet_count") or 0),
"byte_count": int(record.get("byte_count") or 0),
"first_seen": first_seen_raw.isoformat() if hasattr(first_seen_raw, "isoformat") else first_seen_raw,
"last_seen": last_seen_raw.isoformat() if hasattr(last_seen_raw, "isoformat") else last_seen_raw,
"interfaces": list(record.get("interfaces") or []),
"source_count": int(record.get("source_count") or 0),
"destination_count": int(record.get("destination_count") or 0),
"hostnames": set(),
"_protocol_index": {},
"_peer_index": {},
"_service_index": {},
}
for row in detail_rows:
record = dict(row)
protocol_name = str(
_analysis_protocol_name(
record.get("app_protocol_name"),
record.get("ip_proto_raw"),
record.get("eth_type_raw"),
)
)
app_hostname = record.get("app_hostname")
packet_length = int(record.get("packet_length") or 0)
timestamp_raw = record.get("timestamp")
timestamp = timestamp_raw.isoformat() if hasattr(timestamp_raw, "isoformat") else timestamp_raw
src_key = _analysis_host_key(record.get("src_ip_address"), record.get("src_mac_address"))
dst_key = _analysis_host_key(record.get("dst_ip_address"), record.get("dst_mac_address"))
for role, host_key, peer_ip, peer_mac, service_port in (
("source", src_key, record.get("dst_ip_address"), record.get("dst_mac_address"), record.get("src_port")),
("destination", dst_key, record.get("src_ip_address"), record.get("src_mac_address"), record.get("dst_port")),
):
host_record = host_index.get(host_key)
if host_record is None:
continue
if app_hostname:
host_record["hostnames"].add(str(app_hostname))
protocol_record = host_record["_protocol_index"].setdefault(
protocol_name,
{"label": protocol_name, "packet_count": 0},
)
protocol_record["packet_count"] += 1
peer_key = _analysis_host_key(peer_ip, peer_mac)
peer_record = host_record["_peer_index"].get(peer_key)
if peer_record is None:
peer_record = {
"ip_address": peer_ip,
"mac_address": peer_mac,
"packet_count": 0,
"byte_count": 0,
"last_seen": timestamp,
"protocols": set(),
}
host_record["_peer_index"][peer_key] = peer_record
peer_record["packet_count"] += 1
peer_record["byte_count"] += packet_length
if timestamp and (
peer_record.get("last_seen") in (None, "")
or str(timestamp) > str(peer_record.get("last_seen"))
):
peer_record["last_seen"] = timestamp
peer_record["protocols"].add(protocol_name)
if role == "destination" and service_port is not None:
service_key = f"{service_port}|{protocol_name}"
service_record = host_record["_service_index"].get(service_key)
if service_record is None:
service_record = {
"port": int(service_port),
"protocol": protocol_name,
"packet_count": 0,
"byte_count": 0,
"last_seen": timestamp,
"hostnames": set(),
}
host_record["_service_index"][service_key] = service_record
service_record["packet_count"] += 1
service_record["byte_count"] += packet_length
if app_hostname:
service_record["hostnames"].add(str(app_hostname))
if timestamp and (
service_record.get("last_seen") in (None, "")
or str(timestamp) > str(service_record.get("last_seen"))
):
service_record["last_seen"] = timestamp
hosts: List[Dict[str, Any]] = []
for host in host_index.values():
protocols = sorted(
host["_protocol_index"].values(),
key=lambda item: (-int(item.get("packet_count") or 0), str(item.get("label") or "")),
)[:5]
peers = sorted(
host["_peer_index"].values(),
key=lambda item: (-int(item.get("packet_count") or 0), str(item.get("last_seen") or "")),
)[:6]
services = sorted(
host["_service_index"].values(),
key=lambda item: (-int(item.get("packet_count") or 0), str(item.get("last_seen") or ""), int(item.get("port") or 0)),
)[:6]
for peer in peers:
peer["protocols"] = sorted(peer["protocols"])
for service in services:
service["hostnames"] = sorted(service["hostnames"])
hosts.append(
{
"ip_address": host.get("ip_address"),
"mac_address": host.get("mac_address"),
"packet_count": int(host.get("packet_count") or 0),
"byte_count": int(host.get("byte_count") or 0),
"first_seen": host.get("first_seen"),
"last_seen": host.get("last_seen"),
"interfaces": sorted(host.get("interfaces") or []),
"source_count": int(host.get("source_count") or 0),
"destination_count": int(host.get("destination_count") or 0),
"hostnames": sorted(host["hostnames"]),
"top_protocols": protocols,
"peers": peers,
"services": services,
}
)
return sorted(
hosts,
key=lambda item: (
-int(item.get("packet_count") or 0),
-int(item.get("byte_count") or 0),
str(item.get("last_seen") or ""),
str(item.get("ip_address") or ""),
),
)
async def analyze_discovery_activity(
self,
*,
since: Optional[datetime] = None,
limit: int = 300,
) -> List[Dict[str, Any]]:
"""Aggregate discovery and local service advertisement traffic."""
if self._pool is None:
await self.init_pool()
try:
async with self._pool.acquire() as conn:
rows = await conn.fetch(
"""
SELECT
ingress_if,
egress_if,
src_ip::text AS src_ip_address,
src_mac::text AS src_mac_address,
dst_ip::text AS dst_ip_address,
dst_mac::text AS dst_mac_address,
src_port,
dst_port,
NULLIF(app_protocol::text, '') AS app_protocol_name,
ip_proto_raw,
eth_type_raw,
NULLIF(app_hostname::text, '') AS app_hostname,
COUNT(*) AS packet_count,
COALESCE(SUM(length), 0) AS byte_count,
MIN(timestamp) AS first_seen,
MAX(timestamp) AS last_seen
FROM packets
WHERE ($1::timestamptz IS NULL OR timestamp >= $1)
AND (
eth_type_raw = 2054
OR NULLIF(app_protocol::text, '') IN ('MDNS', 'SSDP', 'LLMNR', 'NBNS', 'DHCP', 'DHCPV6')
OR src_port IN (67, 68, 137, 138, 5353, 5355, 1900, 546, 547)
OR dst_port IN (67, 68, 137, 138, 5353, 5355, 1900, 546, 547)
OR ip_proto_raw = 58
)
GROUP BY
ingress_if,
egress_if,
src_ip::text,
src_mac::text,
dst_ip::text,
dst_mac::text,
src_port,
dst_port,
app_protocol_name,
ip_proto_raw,
eth_type_raw,
app_hostname
ORDER BY packet_count DESC, last_seen DESC
LIMIT $2
""",
since,
limit,
)
except Exception:
logger.exception("DB discovery activity analysis failed")
raise
grouped: Dict[str, Dict[str, Any]] = {}
for row in rows:
record = dict(row)
protocol_name = _analysis_protocol_name(
record.get("app_protocol_name"),
record.get("ip_proto_raw"),
record.get("eth_type_raw"),
)
ip_protocol_name = _analysis_ip_protocol_name(record.get("ip_proto_raw"))
ethernet_protocol_name = _analysis_ethernet_protocol_name(record.get("eth_type_raw"))
category = _classify_discovery_activity(
str(protocol_name),
ethernet_protocol_name,
ip_protocol_name,
record.get("src_port"),
record.get("dst_port"),
)
if category is None:
continue
key = "|".join(
[
category,
str(record.get("ingress_if") or ""),
str(record.get("egress_if") or ""),
str(record.get("src_ip_address") or ""),
str(record.get("src_mac_address") or ""),
str(record.get("dst_ip_address") or ""),
str(record.get("dst_mac_address") or ""),
str(record.get("src_port") or ""),
str(record.get("dst_port") or ""),
]
)
first_seen_raw = record.get("first_seen")
last_seen_raw = record.get("last_seen")
activity = grouped.get(key)
if activity is None:
activity = {
"category": category,
"protocol": str(protocol_name),
"ingress_interface": record.get("ingress_if"),
"egress_interface": record.get("egress_if"),
"src_ip_address": record.get("src_ip_address"),
"src_mac_address": record.get("src_mac_address"),
"dst_ip_address": record.get("dst_ip_address"),
"dst_mac_address": record.get("dst_mac_address"),
"src_port": record.get("src_port"),
"dst_port": record.get("dst_port"),
"hostnames": set(),
"packet_count": 0,
"byte_count": 0,
"first_seen": first_seen_raw.isoformat() if hasattr(first_seen_raw, "isoformat") else first_seen_raw,
"last_seen": last_seen_raw.isoformat() if hasattr(last_seen_raw, "isoformat") else last_seen_raw,
}
grouped[key] = activity
activity["packet_count"] += int(record.get("packet_count") or 0)
activity["byte_count"] += int(record.get("byte_count") or 0)
if record.get("app_hostname"):
activity["hostnames"].add(str(record.get("app_hostname")))
last_seen = last_seen_raw.isoformat() if hasattr(last_seen_raw, "isoformat") else last_seen_raw
if last_seen and (
activity.get("last_seen") in (None, "")
or str(last_seen) > str(activity.get("last_seen"))
):
activity["last_seen"] = last_seen
activities = []
for activity in grouped.values():
activity["hostnames"] = sorted(activity["hostnames"])
activities.append(activity)
return sorted(
activities,
key=lambda item: (
-int(item.get("packet_count") or 0),
str(item.get("category") or ""),
str(item.get("last_seen") or ""),
),
)
async def analyze_anomalies(
self,
*,
since: Optional[datetime] = None,
limit: int = 50,
) -> Dict[str, List[Dict[str, Any]]]:
"""Compute lightweight anomaly candidates from observed traffic."""
if self._pool is None:
await self.init_pool()
try:
async with self._pool.acquire() as conn:
scan_rows = await conn.fetch(
"""
SELECT
src_ip::text AS src_ip_address,
src_mac::text AS src_mac_address,
COUNT(*) AS packet_count,
COUNT(DISTINCT COALESCE(dst_ip::text, '') || '|' || COALESCE(dst_mac::text, '')) AS target_host_count,
COUNT(DISTINCT COALESCE(dst_port, -1)) FILTER (WHERE dst_port IS NOT NULL) AS target_port_count,
MIN(timestamp) AS first_seen,
MAX(timestamp) AS last_seen
FROM packets
WHERE ($1::timestamptz IS NULL OR timestamp >= $1)
AND (src_ip IS NOT NULL OR src_mac IS NOT NULL)
AND (dst_ip IS NOT NULL OR dst_mac IS NOT NULL)
GROUP BY src_ip::text, src_mac::text
HAVING COUNT(DISTINCT COALESCE(dst_ip::text, '') || '|' || COALESCE(dst_mac::text, '')) >= 5
OR COUNT(DISTINCT COALESCE(dst_port, -1)) FILTER (WHERE dst_port IS NOT NULL) >= 8
ORDER BY target_host_count DESC, target_port_count DESC, packet_count DESC, last_seen DESC
LIMIT $2
""",
since,
limit,
)
rare_service_rows = await conn.fetch(
"""
WITH service_counts AS (
SELECT
dst_ip::text AS dst_ip_address,
dst_mac::text AS dst_mac_address,
dst_port,
NULLIF(app_protocol::text, '') AS app_protocol_name,
ip_proto_raw,
eth_type_raw,
COUNT(*) AS packet_count,
COUNT(DISTINCT COALESCE(src_ip::text, '') || '|' || COALESCE(src_mac::text, '')) AS client_count,
ARRAY_REMOVE(ARRAY_AGG(DISTINCT NULLIF(app_hostname::text, '')), NULL) AS hostnames,
MAX(timestamp) AS last_seen
FROM packets
WHERE ($1::timestamptz IS NULL OR timestamp >= $1)
AND dst_port IS NOT NULL
AND (dst_ip IS NOT NULL OR dst_mac IS NOT NULL)
GROUP BY
dst_ip::text,
dst_mac::text,
dst_port,
app_protocol_name,
ip_proto_raw,
eth_type_raw
)
SELECT *
FROM service_counts
WHERE client_count <= 2
AND packet_count <= 20
ORDER BY packet_count ASC, client_count ASC, last_seen DESC
LIMIT $2
""",
since,
limit,
)
reset_rows = await conn.fetch(
"""
SELECT
src_ip::text AS src_ip_address,
src_mac::text AS src_mac_address,
dst_ip::text AS dst_ip_address,
dst_mac::text AS dst_mac_address,
dst_port,
COUNT(*) AS total_packets,
SUM(
CASE
WHEN COALESCE(dpi_metadata -> 'tcp' ->> 'packet_type', '') IN ('RST', 'RST-ACK')
OR COALESCE(dpi_metadata -> 'tcp' -> 'flag_names', '[]'::jsonb) ? 'RST'
THEN 1
ELSE 0
END
) AS reset_count,
MAX(timestamp) AS last_seen
FROM packets
WHERE ($1::timestamptz IS NULL OR timestamp >= $1)
AND ip_proto_raw = 6
AND (src_ip IS NOT NULL OR src_mac IS NOT NULL)
AND (dst_ip IS NOT NULL OR dst_mac IS NOT NULL)
GROUP BY
src_ip::text,
src_mac::text,
dst_ip::text,
dst_mac::text,
dst_port
HAVING SUM(
CASE
WHEN COALESCE(dpi_metadata -> 'tcp' ->> 'packet_type', '') IN ('RST', 'RST-ACK')
OR COALESCE(dpi_metadata -> 'tcp' -> 'flag_names', '[]'::jsonb) ? 'RST'
THEN 1
ELSE 0
END
) >= 2
ORDER BY reset_count DESC, total_packets DESC, last_seen DESC
LIMIT $2
""",
since,
limit,
)
drop_rows = await conn.fetch(
"""
SELECT
src_ip::text AS src_ip_address,
src_mac::text AS src_mac_address,
dst_ip::text AS dst_ip_address,
dst_mac::text AS dst_mac_address,
NULLIF(app_protocol::text, '') AS app_protocol_name,
ip_proto_raw,
eth_type_raw,
COUNT(*) AS total_packets,
SUM(CASE WHEN COALESCE(NULLIF(verdict::text, ''), 'unknown') = 'drop' THEN 1 ELSE 0 END) AS drop_count,
SUM(CASE WHEN COALESCE(NULLIF(verdict::text, ''), 'unknown') = 'reject' THEN 1 ELSE 0 END) AS reject_count,
MAX(timestamp) AS last_seen
FROM packets
WHERE ($1::timestamptz IS NULL OR timestamp >= $1)
AND (src_ip IS NOT NULL OR src_mac IS NOT NULL)
AND (dst_ip IS NOT NULL OR dst_mac IS NOT NULL)
GROUP BY
src_ip::text,
src_mac::text,
dst_ip::text,
dst_mac::text,
app_protocol_name,
ip_proto_raw,
eth_type_raw
HAVING SUM(CASE WHEN COALESCE(NULLIF(verdict::text, ''), 'unknown') IN ('drop', 'reject') THEN 1 ELSE 0 END) >= 2
ORDER BY (SUM(CASE WHEN COALESCE(NULLIF(verdict::text, ''), 'unknown') IN ('drop', 'reject') THEN 1 ELSE 0 END)) DESC, last_seen DESC
LIMIT $2
""",
since,
limit,
)
beacon_rows = await conn.fetch(
"""
SELECT
src_ip::text AS src_ip_address,
src_mac::text AS src_mac_address,
dst_ip::text AS dst_ip_address,
dst_mac::text AS dst_mac_address,
dst_port,
NULLIF(app_protocol::text, '') AS app_protocol_name,
ip_proto_raw,
eth_type_raw,
ARRAY_AGG(EXTRACT(EPOCH FROM timestamp) ORDER BY timestamp) AS observed_seconds,
COUNT(*) AS packet_count,
MIN(timestamp) AS first_seen,
MAX(timestamp) AS last_seen
FROM packets
WHERE ($1::timestamptz IS NULL OR timestamp >= $1)
AND (src_ip IS NOT NULL OR src_mac IS NOT NULL)
AND (dst_ip IS NOT NULL OR dst_mac IS NOT NULL)
GROUP BY
src_ip::text,
src_mac::text,
dst_ip::text,
dst_mac::text,
dst_port,
app_protocol_name,
ip_proto_raw,
eth_type_raw
HAVING COUNT(*) >= 4
ORDER BY packet_count DESC, last_seen DESC
LIMIT 1000
""",
since,
)
except Exception:
logger.exception("DB anomaly analysis failed")
raise
scan_candidates: List[Dict[str, Any]] = []
for row in scan_rows:
record = dict(row)
first_seen_raw = record.get("first_seen")
last_seen_raw = record.get("last_seen")
scan_candidates.append(
{
"src_ip_address": record.get("src_ip_address"),
"src_mac_address": record.get("src_mac_address"),
"packet_count": int(record.get("packet_count") or 0),
"target_host_count": int(record.get("target_host_count") or 0),
"target_port_count": int(record.get("target_port_count") or 0),
"first_seen": first_seen_raw.isoformat() if hasattr(first_seen_raw, "isoformat") else first_seen_raw,
"last_seen": last_seen_raw.isoformat() if hasattr(last_seen_raw, "isoformat") else last_seen_raw,
}
)
rare_services: List[Dict[str, Any]] = []
for row in rare_service_rows:
record = dict(row)
protocol_name = _analysis_protocol_name(
record.get("app_protocol_name"),
record.get("ip_proto_raw"),
record.get("eth_type_raw"),
)
last_seen_raw = record.get("last_seen")
rare_services.append(
{
"dst_ip_address": record.get("dst_ip_address"),
"dst_mac_address": record.get("dst_mac_address"),
"dst_port": record.get("dst_port"),
"protocol": str(protocol_name),
"packet_count": int(record.get("packet_count") or 0),
"client_count": int(record.get("client_count") or 0),
"hostnames": list(record.get("hostnames") or []),
"last_seen": last_seen_raw.isoformat() if hasattr(last_seen_raw, "isoformat") else last_seen_raw,
}
)
reset_heavy_paths: List[Dict[str, Any]] = []
for row in reset_rows:
record = dict(row)
last_seen_raw = record.get("last_seen")
total_packets = int(record.get("total_packets") or 0)
reset_count = int(record.get("reset_count") or 0)
reset_heavy_paths.append(
{
"src_ip_address": record.get("src_ip_address"),
"src_mac_address": record.get("src_mac_address"),
"dst_ip_address": record.get("dst_ip_address"),
"dst_mac_address": record.get("dst_mac_address"),
"dst_port": record.get("dst_port"),
"total_packets": total_packets,
"reset_count": reset_count,
"reset_ratio": round(_safe_ratio(reset_count, total_packets), 3),
"last_seen": last_seen_raw.isoformat() if hasattr(last_seen_raw, "isoformat") else last_seen_raw,
}
)
drop_heavy_paths: List[Dict[str, Any]] = []
for row in drop_rows:
record = dict(row)
protocol_name = _analysis_protocol_name(
record.get("app_protocol_name"),
record.get("ip_proto_raw"),
record.get("eth_type_raw"),
)
last_seen_raw = record.get("last_seen")
total_packets = int(record.get("total_packets") or 0)
drop_count = int(record.get("drop_count") or 0)
reject_count = int(record.get("reject_count") or 0)
failure_count = drop_count + reject_count
drop_heavy_paths.append(
{
"src_ip_address": record.get("src_ip_address"),
"src_mac_address": record.get("src_mac_address"),
"dst_ip_address": record.get("dst_ip_address"),
"dst_mac_address": record.get("dst_mac_address"),
"protocol": str(protocol_name),
"total_packets": total_packets,
"drop_count": drop_count,
"reject_count": reject_count,
"failure_ratio": round(_safe_ratio(failure_count, total_packets), 3),
"last_seen": last_seen_raw.isoformat() if hasattr(last_seen_raw, "isoformat") else last_seen_raw,
}
)
beacon_candidates: List[Dict[str, Any]] = []
for row in beacon_rows:
record = dict(row)
observed_seconds = [float(value) for value in (record.get("observed_seconds") or []) if value is not None]
if len(observed_seconds) < 4:
continue
intervals = [
observed_seconds[index] - observed_seconds[index - 1]
for index in range(1, len(observed_seconds))
if observed_seconds[index] - observed_seconds[index - 1] > 0
]
if len(intervals) < 3:
continue
avg_interval = sum(intervals) / len(intervals)
if avg_interval < 1 or avg_interval > 3600:
continue
if len(intervals) == 1:
jitter_ratio = 0.0
else:
jitter_ratio = _safe_ratio(statistics.pstdev(intervals), avg_interval)
if math.isnan(jitter_ratio) or jitter_ratio > 0.25:
continue
protocol_name = _analysis_protocol_name(
record.get("app_protocol_name"),
record.get("ip_proto_raw"),
record.get("eth_type_raw"),
)
first_seen_raw = record.get("first_seen")
last_seen_raw = record.get("last_seen")
beacon_candidates.append(
{
"src_ip_address": record.get("src_ip_address"),
"src_mac_address": record.get("src_mac_address"),
"dst_ip_address": record.get("dst_ip_address"),
"dst_mac_address": record.get("dst_mac_address"),
"dst_port": record.get("dst_port"),
"protocol": str(protocol_name),
"packet_count": int(record.get("packet_count") or 0),
"avg_interval_seconds": round(avg_interval, 2),
"jitter_ratio": round(jitter_ratio, 3),
"first_seen": first_seen_raw.isoformat() if hasattr(first_seen_raw, "isoformat") else first_seen_raw,
"last_seen": last_seen_raw.isoformat() if hasattr(last_seen_raw, "isoformat") else last_seen_raw,
}
)
beacon_candidates = sorted(
beacon_candidates,
key=lambda item: (
item.get("jitter_ratio", 1.0),
-int(item.get("packet_count") or 0),
str(item.get("last_seen") or ""),
),
)[:limit]
return {
"scan_candidates": scan_candidates,
"beacon_candidates": beacon_candidates,
"rare_services": rare_services,
"reset_heavy_paths": reset_heavy_paths,
"drop_heavy_paths": drop_heavy_paths,
}
async def clear_all_packets(self, reset_identity: bool = True) -> bool:
"""Truncate the packet table and optionally reset identity counters."""
if self._pool is None: