add analysis base api
This commit is contained in:
@@ -632,6 +632,107 @@ class DatabasePool:
|
||||
|
||||
return result
|
||||
|
||||
async def infer_interface_hosts(
|
||||
self,
|
||||
*,
|
||||
since: Optional[datetime] = None,
|
||||
limit_per_interface: int = 100,
|
||||
) -> List[Dict[str, Any]]:
|
||||
"""Infer which IP/MAC endpoints are attached to each observed interface."""
|
||||
if self._pool is None:
|
||||
await self.init_pool()
|
||||
|
||||
async with self._pool.acquire() as conn:
|
||||
rows = await conn.fetch(
|
||||
"""
|
||||
WITH observations AS (
|
||||
SELECT
|
||||
ingress_if AS iface,
|
||||
src_ip::text AS ip_address,
|
||||
src_mac::text AS mac_address,
|
||||
timestamp,
|
||||
'source_on_ingress' AS evidence
|
||||
FROM packets
|
||||
WHERE ingress_if IS NOT NULL
|
||||
AND (src_ip IS NOT NULL OR src_mac IS NOT NULL)
|
||||
AND ($1::timestamptz IS NULL OR timestamp >= $1)
|
||||
|
||||
UNION ALL
|
||||
|
||||
SELECT
|
||||
egress_if AS iface,
|
||||
dst_ip::text AS ip_address,
|
||||
dst_mac::text AS mac_address,
|
||||
timestamp,
|
||||
'destination_on_egress' AS evidence
|
||||
FROM packets
|
||||
WHERE egress_if IS NOT NULL
|
||||
AND (dst_ip IS NOT NULL OR dst_mac IS NOT NULL)
|
||||
AND ($1::timestamptz IS NULL OR timestamp >= $1)
|
||||
),
|
||||
filtered AS (
|
||||
SELECT *
|
||||
FROM observations
|
||||
WHERE iface IS NOT NULL
|
||||
AND COALESCE(mac_address, '') <> 'ff:ff:ff:ff:ff:ff'
|
||||
AND (
|
||||
COALESCE(ip_address, '') <> ''
|
||||
OR COALESCE(mac_address, '') <> ''
|
||||
)
|
||||
),
|
||||
aggregated AS (
|
||||
SELECT
|
||||
iface,
|
||||
ip_address,
|
||||
mac_address,
|
||||
COUNT(*) AS packet_count,
|
||||
MAX(timestamp) AS last_seen,
|
||||
SUM(CASE WHEN evidence = 'source_on_ingress' THEN 1 ELSE 0 END) AS source_on_ingress_count,
|
||||
SUM(CASE WHEN evidence = 'destination_on_egress' THEN 1 ELSE 0 END) AS destination_on_egress_count
|
||||
FROM filtered
|
||||
GROUP BY iface, ip_address, mac_address
|
||||
),
|
||||
ranked AS (
|
||||
SELECT
|
||||
*,
|
||||
ROW_NUMBER() OVER (
|
||||
PARTITION BY iface
|
||||
ORDER BY packet_count DESC, last_seen DESC, ip_address, mac_address
|
||||
) AS row_num
|
||||
FROM aggregated
|
||||
)
|
||||
SELECT
|
||||
iface,
|
||||
ip_address,
|
||||
mac_address,
|
||||
packet_count,
|
||||
last_seen,
|
||||
source_on_ingress_count,
|
||||
destination_on_egress_count
|
||||
FROM ranked
|
||||
WHERE row_num <= $2
|
||||
ORDER BY iface, packet_count DESC, last_seen DESC, ip_address, mac_address
|
||||
""",
|
||||
since,
|
||||
limit_per_interface,
|
||||
)
|
||||
|
||||
grouped: Dict[str, List[Dict[str, Any]]] = {}
|
||||
for row in rows:
|
||||
record = dict(row)
|
||||
iface = str(record.pop("iface"))
|
||||
if hasattr(record.get("last_seen"), "isoformat"):
|
||||
record["last_seen"] = record["last_seen"].isoformat()
|
||||
grouped.setdefault(iface, []).append(record)
|
||||
|
||||
return [
|
||||
{
|
||||
"interface": iface,
|
||||
"hosts": hosts,
|
||||
}
|
||||
for iface, hosts in sorted(grouped.items())
|
||||
]
|
||||
|
||||
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:
|
||||
|
||||
Reference in New Issue
Block a user