try overload fix timestamp null error
This commit is contained in:
@@ -20,6 +20,10 @@ from src.Models.packets import PacketDBModel
|
|||||||
logger = logging.getLogger("packet_capture")
|
logger = logging.getLogger("packet_capture")
|
||||||
|
|
||||||
|
|
||||||
|
def _utcnow() -> datetime:
|
||||||
|
return datetime.now(timezone.utc)
|
||||||
|
|
||||||
|
|
||||||
def _db_text(value: Any) -> Any:
|
def _db_text(value: Any) -> Any:
|
||||||
if value is None:
|
if value is None:
|
||||||
return None
|
return None
|
||||||
@@ -168,6 +172,9 @@ def _derive_flow_id(payload: Dict[str, Any]) -> Optional[str]:
|
|||||||
|
|
||||||
|
|
||||||
def _attach_derived_fields(payload: Dict[str, Any]) -> None:
|
def _attach_derived_fields(payload: Dict[str, Any]) -> None:
|
||||||
|
if payload.get("timestamp") in (None, ""):
|
||||||
|
payload["timestamp"] = _utcnow()
|
||||||
|
|
||||||
flow_id = _derive_flow_id(payload)
|
flow_id = _derive_flow_id(payload)
|
||||||
if flow_id is not None:
|
if flow_id is not None:
|
||||||
current_flow_id = payload.get("flow_id")
|
current_flow_id = payload.get("flow_id")
|
||||||
@@ -228,6 +235,20 @@ class DatabasePool:
|
|||||||
max_size=self._max_size,
|
max_size=self._max_size,
|
||||||
)
|
)
|
||||||
async with self._pool.acquire() as conn:
|
async with self._pool.acquire() as conn:
|
||||||
|
await conn.execute(
|
||||||
|
"""
|
||||||
|
DO $$
|
||||||
|
BEGIN
|
||||||
|
IF to_regclass('packets') IS NOT NULL THEN
|
||||||
|
UPDATE packets
|
||||||
|
SET timestamp = COALESCE(updated_at, NOW())
|
||||||
|
WHERE timestamp IS NULL;
|
||||||
|
END IF;
|
||||||
|
END $$;
|
||||||
|
"""
|
||||||
|
)
|
||||||
|
await conn.execute("ALTER TABLE IF EXISTS packets ALTER COLUMN timestamp SET DEFAULT NOW()")
|
||||||
|
await conn.execute("ALTER TABLE IF EXISTS packets ALTER COLUMN timestamp SET NOT NULL")
|
||||||
await conn.execute(
|
await conn.execute(
|
||||||
"""
|
"""
|
||||||
ALTER TABLE IF EXISTS packets
|
ALTER TABLE IF EXISTS packets
|
||||||
|
|||||||
@@ -13,6 +13,7 @@ import signal
|
|||||||
import socket
|
import socket
|
||||||
import subprocess
|
import subprocess
|
||||||
import sys
|
import sys
|
||||||
|
from datetime import datetime, timezone
|
||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
from typing import Iterable
|
from typing import Iterable
|
||||||
|
|
||||||
@@ -504,6 +505,8 @@ def _build_payload(event: Event) -> dict[str, object] | None:
|
|||||||
|
|
||||||
payload: dict[str, object] = {
|
payload: dict[str, object] = {
|
||||||
"event_type": _event_name(int(event.event_type)),
|
"event_type": _event_name(int(event.event_type)),
|
||||||
|
"timestamp": datetime.now(timezone.utc).isoformat(),
|
||||||
|
"kernel_ts_ns": int(event.ts_ns),
|
||||||
"iface": iface,
|
"iface": iface,
|
||||||
"skb_mark": int(event.skb_mark) or None,
|
"skb_mark": int(event.skb_mark) or None,
|
||||||
"length": int(event.length),
|
"length": int(event.length),
|
||||||
|
|||||||
@@ -53,6 +53,18 @@ def _parse_observation_timestamp(value: Any) -> datetime:
|
|||||||
return datetime.max.replace(tzinfo=timezone.utc)
|
return datetime.max.replace(tzinfo=timezone.utc)
|
||||||
|
|
||||||
|
|
||||||
|
def _coerce_payload_timestamp(value: Any) -> datetime:
|
||||||
|
if isinstance(value, datetime):
|
||||||
|
return value if value.tzinfo is not None else value.replace(tzinfo=timezone.utc)
|
||||||
|
if value not in (None, ""):
|
||||||
|
try:
|
||||||
|
parsed = datetime.fromisoformat(str(value).replace("Z", "+00:00"))
|
||||||
|
return parsed if parsed.tzinfo is not None else parsed.replace(tzinfo=timezone.utc)
|
||||||
|
except Exception:
|
||||||
|
pass
|
||||||
|
return _utcnow()
|
||||||
|
|
||||||
|
|
||||||
def _bridge_af_packet_observation_groups(payload: Dict[str, Any]) -> Dict[str, set[str]]:
|
def _bridge_af_packet_observation_groups(payload: Dict[str, Any]) -> Dict[str, set[str]]:
|
||||||
groups: Dict[str, set[str]] = {}
|
groups: Dict[str, set[str]] = {}
|
||||||
observations = payload.get("capture_observations") or []
|
observations = payload.get("capture_observations") or []
|
||||||
@@ -269,6 +281,10 @@ class PacketTracker:
|
|||||||
payload["skb_mark"] = event.get("skb_mark") or payload.get("skb_mark")
|
payload["skb_mark"] = event.get("skb_mark") or payload.get("skb_mark")
|
||||||
payload["telemetry_metadata"] = event
|
payload["telemetry_metadata"] = event
|
||||||
payload["last_observed_at"] = now_ts
|
payload["last_observed_at"] = now_ts
|
||||||
|
event_timestamp = _coerce_payload_timestamp(event.get("timestamp"))
|
||||||
|
current_timestamp = payload.get("timestamp")
|
||||||
|
if current_timestamp in (None, "") or event_timestamp < _coerce_payload_timestamp(current_timestamp):
|
||||||
|
payload["timestamp"] = event_timestamp
|
||||||
self._add_capture_source(payload, "telemetry")
|
self._add_capture_source(payload, "telemetry")
|
||||||
self._add_capture_observation(
|
self._add_capture_observation(
|
||||||
payload,
|
payload,
|
||||||
@@ -336,7 +352,7 @@ class PacketTracker:
|
|||||||
return {
|
return {
|
||||||
"correlation_key": correlation_key,
|
"correlation_key": correlation_key,
|
||||||
"payload": {
|
"payload": {
|
||||||
"timestamp": None,
|
"timestamp": _utcnow(),
|
||||||
"correlation_key": correlation_key,
|
"correlation_key": correlation_key,
|
||||||
"correlation_source": None,
|
"correlation_source": None,
|
||||||
"packet_id": None,
|
"packet_id": None,
|
||||||
|
|||||||
Reference in New Issue
Block a user