diff --git a/backend/src/utilities/database.py b/backend/src/utilities/database.py index 288482b..516cb77 100644 --- a/backend/src/utilities/database.py +++ b/backend/src/utilities/database.py @@ -27,6 +27,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) raw_val = serialized.get("raw") if isinstance(raw_val, (bytes, bytearray)): @@ -40,6 +41,18 @@ def _serialize_row_for_broadcast(row: Dict[str, Any]) -> Dict[str, Any]: return serialized +def _normalize_json_fields(payload: Dict[str, Any]) -> None: + for key in ("dpi_metadata", "telemetry_metadata"): + value = payload.get(key) + if isinstance(value, str): + try: + parsed = json.loads(value) + except json.JSONDecodeError: + continue + if isinstance(parsed, dict): + payload[key] = parsed + + class DatabasePool: """Asyncpg connection pool wrapper used by the packet APIs.""" @@ -239,6 +252,7 @@ class DatabasePool: result: List[PacketDBModel] = [] for row in rows: data = dict(row) + _normalize_json_fields(data) raw_val = data.get("raw") if isinstance(raw_val, (bytes, bytearray)): diff --git a/setup_database.sh b/setup_database.sh index 1faf97b..3fe93ba 100755 --- a/setup_database.sh +++ b/setup_database.sh @@ -103,6 +103,12 @@ ALTER TABLE packets ADD COLUMN IF NOT EXISTS app_hostname VARCHAR(255); ALTER TABLE packets ADD COLUMN IF NOT EXISTS app_is_encrypted BOOLEAN; ALTER TABLE packets ADD COLUMN IF NOT EXISTS app_risk_score INTEGER; ALTER TABLE packets ADD COLUMN IF NOT EXISTS dpi_metadata JSONB; +ALTER TABLE packets ALTER COLUMN dpi_metadata TYPE JSONB USING + CASE + WHEN dpi_metadata IS NULL THEN NULL + WHEN pg_typeof(dpi_metadata)::text = 'jsonb' THEN dpi_metadata + ELSE dpi_metadata::jsonb + END; ALTER TABLE packets ADD COLUMN IF NOT EXISTS eth_type_raw INTEGER; ALTER TABLE packets ADD COLUMN IF NOT EXISTS ip_proto_raw INTEGER; ALTER TABLE packets ADD COLUMN IF NOT EXISTS ingress_if VARCHAR(64); @@ -115,6 +121,12 @@ ALTER TABLE packets ADD COLUMN IF NOT EXISTS ingress_seen_at TIMESTAMPTZ; ALTER TABLE packets ADD COLUMN IF NOT EXISTS egress_seen_at TIMESTAMPTZ; ALTER TABLE packets ADD COLUMN IF NOT EXISTS verdict_seen_at TIMESTAMPTZ; ALTER TABLE packets ADD COLUMN IF NOT EXISTS telemetry_metadata JSONB; +ALTER TABLE packets ALTER COLUMN telemetry_metadata TYPE JSONB USING + CASE + WHEN telemetry_metadata IS NULL THEN NULL + WHEN pg_typeof(telemetry_metadata)::text = 'jsonb' THEN telemetry_metadata + ELSE telemetry_metadata::jsonb + END; ALTER TABLE packets DROP COLUMN IF EXISTS direction; CREATE UNIQUE INDEX IF NOT EXISTS idx_packets_packet_uid ON packets(packet_uid); EOF