test path visu
This commit is contained in:
@@ -85,6 +85,36 @@ class InterfaceHostProtocolAnalysisResponse(BaseModel):
|
||||
)
|
||||
|
||||
|
||||
class InterfaceProtocolPathEvidence(BaseModel):
|
||||
ingress_interface: Optional[str] = Field(None, description="Observed ingress interface for the packet path.")
|
||||
egress_interface: Optional[str] = Field(None, description="Observed egress interface for the packet path.")
|
||||
src_ip_address: Optional[str] = Field(None, description="Observed source IP address.")
|
||||
src_mac_address: Optional[str] = Field(None, description="Observed source MAC address.")
|
||||
dst_ip_address: Optional[str] = Field(None, description="Observed destination IP address.")
|
||||
dst_mac_address: Optional[str] = Field(None, description="Observed destination MAC address.")
|
||||
protocol: str = Field(..., description="Detected application or fallback protocol for the packet path.")
|
||||
ethernet_protocol: Optional[str] = Field(None, description="Dominant Ethernet protocol associated with this path.")
|
||||
ip_protocol: Optional[str] = Field(None, description="Dominant IP protocol associated with this path.")
|
||||
packet_count: int = Field(..., description="Packet observations supporting this end-to-end path.")
|
||||
last_seen: datetime = Field(..., description="Most recent packet timestamp supporting this path.")
|
||||
accept_count: int = Field(0, description="Packets with verdict=accept for this path.")
|
||||
drop_count: int = Field(0, description="Packets with verdict=drop for this path.")
|
||||
reject_count: int = Field(0, description="Packets with verdict=reject for this path.")
|
||||
unknown_count: int = Field(0, description="Packets with verdict pending/unknown or without a verdict.")
|
||||
|
||||
|
||||
class InterfaceProtocolPathAnalysisResponse(BaseModel):
|
||||
since: Optional[datetime] = Field(None, description="Only packets at or after this timestamp were analyzed.")
|
||||
paths: List[InterfaceProtocolPathEvidence] = Field(default_factory=list)
|
||||
notes: List[str] = Field(
|
||||
default_factory=lambda: [
|
||||
"This Sankey view is built from packet paths, not from inferred interface-host attachment.",
|
||||
"Each row represents a grouped ingress -> source endpoint -> protocol -> destination endpoint -> egress path.",
|
||||
"Protocols prefer app_protocol and fall back to lower-layer protocol names.",
|
||||
]
|
||||
)
|
||||
|
||||
|
||||
@router.get("/interface-hosts", response_model=InterfaceHostAnalysisResponse)
|
||||
async def analysis_interface_hosts(
|
||||
since_minutes: Optional[int] = Query(
|
||||
@@ -159,3 +189,39 @@ async def analysis_interface_host_protocols(
|
||||
|
||||
interfaces = [InterfaceProtocolAttachment(**row) for row in rows]
|
||||
return InterfaceHostProtocolAnalysisResponse(since=since, interfaces=interfaces)
|
||||
|
||||
|
||||
@router.get("/interface-protocol-paths", response_model=InterfaceProtocolPathAnalysisResponse)
|
||||
async def analysis_interface_protocol_paths(
|
||||
since_minutes: Optional[int] = Query(
|
||||
None,
|
||||
ge=1,
|
||||
le=60 * 24 * 30,
|
||||
description="Analyze only packets seen within the last N minutes. Omit to cover all captured history.",
|
||||
),
|
||||
limit_paths: int = Query(
|
||||
500,
|
||||
ge=1,
|
||||
le=5000,
|
||||
description="Maximum number of grouped packet paths returned for the Sankey view.",
|
||||
),
|
||||
) -> InterfaceProtocolPathAnalysisResponse:
|
||||
"""Aggregate directional packet paths for the Sankey diagram."""
|
||||
db = shared.db
|
||||
if db is None:
|
||||
raise HTTPException(status_code=503, detail="Database not available")
|
||||
|
||||
since: Optional[datetime] = None
|
||||
if since_minutes is not None:
|
||||
since = datetime.now(timezone.utc) - timedelta(minutes=since_minutes)
|
||||
|
||||
try:
|
||||
rows = await db.infer_interface_protocol_paths(
|
||||
since=since,
|
||||
limit_paths=limit_paths,
|
||||
)
|
||||
except Exception as exc:
|
||||
raise HTTPException(status_code=500, detail=f"Failed to infer interface protocol paths: {exc}") from exc
|
||||
|
||||
paths = [InterfaceProtocolPathEvidence(**row) for row in rows]
|
||||
return InterfaceProtocolPathAnalysisResponse(since=since, paths=paths)
|
||||
|
||||
@@ -1075,6 +1075,139 @@ class DatabasePool:
|
||||
|
||||
return result
|
||||
|
||||
async def infer_interface_protocol_paths(
|
||||
self,
|
||||
*,
|
||||
since: Optional[datetime] = None,
|
||||
limit_paths: int = 500,
|
||||
) -> List[Dict[str, Any]]:
|
||||
"""Aggregate directional packet paths for Sankey rendering."""
|
||||
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,
|
||||
NULLIF(app_protocol::text, '') AS app_protocol_name,
|
||||
ip_proto_raw,
|
||||
eth_type_raw,
|
||||
COUNT(*) AS packet_count,
|
||||
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 (ingress_if IS NOT NULL OR egress_if IS NOT NULL)
|
||||
AND (src_ip IS NOT NULL OR src_mac IS NOT NULL)
|
||||
AND (dst_ip IS NOT NULL OR dst_mac IS NOT NULL)
|
||||
AND COALESCE(src_mac::text, '') <> 'ff:ff:ff:ff:ff:ff'
|
||||
AND COALESCE(dst_mac::text, '') <> 'ff:ff:ff:ff:ff:ff'
|
||||
GROUP BY
|
||||
ingress_if,
|
||||
egress_if,
|
||||
src_ip::text,
|
||||
src_mac::text,
|
||||
dst_ip::text,
|
||||
dst_mac::text,
|
||||
app_protocol_name,
|
||||
ip_proto_raw,
|
||||
eth_type_raw
|
||||
)
|
||||
SELECT *
|
||||
FROM aggregated
|
||||
ORDER BY packet_count DESC, last_seen DESC
|
||||
LIMIT $2
|
||||
""",
|
||||
since,
|
||||
limit_paths,
|
||||
)
|
||||
except Exception:
|
||||
logger.exception("DB interface-protocol-path 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"))
|
||||
last_seen_raw = record.get("last_seen")
|
||||
last_seen = last_seen_raw.isoformat() if hasattr(last_seen_raw, "isoformat") else last_seen_raw
|
||||
path_key = "|".join(
|
||||
[
|
||||
str(record.get("ingress_if") or ""),
|
||||
str(record.get("src_ip_address") or ""),
|
||||
str(record.get("src_mac_address") or ""),
|
||||
str(protocol_name or ""),
|
||||
str(record.get("dst_ip_address") or ""),
|
||||
str(record.get("dst_mac_address") or ""),
|
||||
str(record.get("egress_if") or ""),
|
||||
]
|
||||
)
|
||||
|
||||
path_record = grouped.get(path_key)
|
||||
if path_record is None:
|
||||
path_record = {
|
||||
"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"),
|
||||
"protocol": str(protocol_name),
|
||||
"ethernet_protocol": ethernet_protocol_name,
|
||||
"ip_protocol": ip_protocol_name,
|
||||
"packet_count": 0,
|
||||
"last_seen": last_seen,
|
||||
"accept_count": 0,
|
||||
"drop_count": 0,
|
||||
"reject_count": 0,
|
||||
"unknown_count": 0,
|
||||
}
|
||||
grouped[path_key] = path_record
|
||||
|
||||
path_record["packet_count"] += int(record.get("packet_count") or 0)
|
||||
path_record["accept_count"] += int(record.get("accept_count") or 0)
|
||||
path_record["drop_count"] += int(record.get("drop_count") or 0)
|
||||
path_record["reject_count"] += int(record.get("reject_count") or 0)
|
||||
path_record["unknown_count"] += int(record.get("unknown_count") or 0)
|
||||
if path_record.get("ethernet_protocol") is None and ethernet_protocol_name is not None:
|
||||
path_record["ethernet_protocol"] = ethernet_protocol_name
|
||||
if path_record.get("ip_protocol") is None and ip_protocol_name is not None:
|
||||
path_record["ip_protocol"] = ip_protocol_name
|
||||
if last_seen and (
|
||||
path_record.get("last_seen") in (None, "")
|
||||
or str(last_seen) > str(path_record.get("last_seen"))
|
||||
):
|
||||
path_record["last_seen"] = last_seen
|
||||
|
||||
return sorted(
|
||||
grouped.values(),
|
||||
key=lambda item: (
|
||||
-int(item.get("packet_count") or 0),
|
||||
str(item.get("last_seen") or ""),
|
||||
str(item.get("ingress_interface") or ""),
|
||||
str(item.get("src_ip_address") or ""),
|
||||
str(item.get("dst_ip_address") or ""),
|
||||
str(item.get("egress_interface") or ""),
|
||||
),
|
||||
)
|
||||
|
||||
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