flow showing
This commit is contained in:
@@ -80,6 +80,7 @@ class PacketDBModel(BaseModel):
|
||||
length: Optional[int] = None
|
||||
raw_present: Optional[bool] = Field(None, description="Whether raw packet bytes were captured for this row.")
|
||||
capture_sources: Optional[list[str]] = Field(None, description="Capture sources that contributed to this row.")
|
||||
flow_id: Optional[str] = Field(None, description="Derived flow identifier from tshark stream metadata, if available.")
|
||||
raw_b64: Optional[str] = Field(None, description="Base64-encoded packet bytes.")
|
||||
app_protocol: Optional[str] = Field(None, description="Detected application protocol.")
|
||||
app_master_protocol: Optional[str] = Field(None, description="Detected application master protocol.")
|
||||
|
||||
@@ -29,6 +29,7 @@ def _db_text(value: Any) -> Any:
|
||||
def _serialize_row_for_broadcast(row: Dict[str, Any]) -> Dict[str, Any]:
|
||||
serialized = dict(row)
|
||||
_normalize_json_fields(serialized)
|
||||
_attach_derived_fields(serialized)
|
||||
|
||||
raw_val = serialized.get("raw")
|
||||
if isinstance(raw_val, (bytes, bytearray)):
|
||||
@@ -54,6 +55,42 @@ def _normalize_json_fields(payload: Dict[str, Any]) -> None:
|
||||
payload[key] = parsed
|
||||
|
||||
|
||||
def _derive_flow_id(payload: Dict[str, Any]) -> Optional[str]:
|
||||
dpi_metadata = payload.get("dpi_metadata")
|
||||
if not isinstance(dpi_metadata, dict):
|
||||
return None
|
||||
|
||||
tcp_meta = dpi_metadata.get("tcp")
|
||||
if isinstance(tcp_meta, dict):
|
||||
stream = tcp_meta.get("stream")
|
||||
if stream not in (None, "", []):
|
||||
return f"tcp:{stream}"
|
||||
|
||||
udp_meta = dpi_metadata.get("udp")
|
||||
if isinstance(udp_meta, dict):
|
||||
stream = udp_meta.get("stream")
|
||||
if stream not in (None, "", []):
|
||||
return f"udp:{stream}"
|
||||
|
||||
tshark_meta = dpi_metadata.get("tshark")
|
||||
if isinstance(tshark_meta, dict):
|
||||
tcp_stream = tshark_meta.get("tcp_stream")
|
||||
if tcp_stream not in (None, "", []):
|
||||
return f"tcp:{tcp_stream}"
|
||||
udp_stream = tshark_meta.get("udp_stream")
|
||||
if udp_stream not in (None, "", []):
|
||||
return f"udp:{udp_stream}"
|
||||
|
||||
return None
|
||||
|
||||
|
||||
def _attach_derived_fields(payload: Dict[str, Any]) -> None:
|
||||
if payload.get("flow_id") in (None, ""):
|
||||
flow_id = _derive_flow_id(payload)
|
||||
if flow_id is not None:
|
||||
payload["flow_id"] = flow_id
|
||||
|
||||
|
||||
class DatabasePool:
|
||||
"""Asyncpg connection pool wrapper used by the packet APIs."""
|
||||
|
||||
@@ -509,7 +546,7 @@ class DatabasePool:
|
||||
"""
|
||||
SELECT *
|
||||
FROM packets
|
||||
ORDER BY updated_at DESC, id DESC
|
||||
ORDER BY timestamp DESC, id DESC
|
||||
LIMIT $1
|
||||
""",
|
||||
limit,
|
||||
@@ -519,6 +556,7 @@ class DatabasePool:
|
||||
for row in rows:
|
||||
data = dict(row)
|
||||
_normalize_json_fields(data)
|
||||
_attach_derived_fields(data)
|
||||
|
||||
raw_val = data.get("raw")
|
||||
if isinstance(raw_val, (bytes, bytearray)):
|
||||
|
||||
Reference in New Issue
Block a user