test fix timestampt to capture time not db insert time
This commit is contained in:
@@ -7,6 +7,7 @@ import socket
|
|||||||
import selectors
|
import selectors
|
||||||
import errno
|
import errno
|
||||||
import struct
|
import struct
|
||||||
|
from datetime import datetime, timezone
|
||||||
from typing import Dict, List, Optional, Any, TypedDict, Union
|
from typing import Dict, List, Optional, Any, TypedDict, Union
|
||||||
from uuid import uuid4
|
from uuid import uuid4
|
||||||
|
|
||||||
@@ -61,6 +62,7 @@ sessions: Dict[str, Dict[str, Any]] = {}
|
|||||||
class PacketInfo(TypedDict, total=False):
|
class PacketInfo(TypedDict, total=False):
|
||||||
"""TypedDict for parsed packet data used by persistence and telemetry."""
|
"""TypedDict for parsed packet data used by persistence and telemetry."""
|
||||||
|
|
||||||
|
timestamp: datetime
|
||||||
correlation_key: str
|
correlation_key: str
|
||||||
correlation_source: str
|
correlation_source: str
|
||||||
packet_id: Optional[str]
|
packet_id: Optional[str]
|
||||||
@@ -162,6 +164,16 @@ def _safe_get_attr(layer, attr: str):
|
|||||||
return None
|
return None
|
||||||
|
|
||||||
|
|
||||||
|
def _packet_timestamp(pkt: Any) -> datetime:
|
||||||
|
try:
|
||||||
|
packet_time = float(getattr(pkt, "time", 0.0) or 0.0)
|
||||||
|
if packet_time > 0:
|
||||||
|
return datetime.fromtimestamp(packet_time, tz=timezone.utc)
|
||||||
|
except Exception:
|
||||||
|
pass
|
||||||
|
return datetime.now(timezone.utc)
|
||||||
|
|
||||||
|
|
||||||
def _merge_enrichment(pkt_info: PacketInfo, enrichment: Dict[str, Any]) -> None:
|
def _merge_enrichment(pkt_info: PacketInfo, enrichment: Dict[str, Any]) -> None:
|
||||||
"""Populate enrichment fields without discarding existing metadata."""
|
"""Populate enrichment fields without discarding existing metadata."""
|
||||||
for key, value in enrichment.items():
|
for key, value in enrichment.items():
|
||||||
@@ -199,6 +211,7 @@ def parse_packet(pkt, bridge_label: str, capture_metadata: Optional[Dict[str, An
|
|||||||
logger.debug("Packet captured on %s (bridge_label %s)", pkt_iface, bridge_label)
|
logger.debug("Packet captured on %s (bridge_label %s)", pkt_iface, bridge_label)
|
||||||
|
|
||||||
pkt_info: PacketInfo = {
|
pkt_info: PacketInfo = {
|
||||||
|
"timestamp": _packet_timestamp(pkt),
|
||||||
"iface": pkt_iface,
|
"iface": pkt_iface,
|
||||||
"capture_iface": None,
|
"capture_iface": None,
|
||||||
"length": len(pkt),
|
"length": len(pkt),
|
||||||
|
|||||||
@@ -156,6 +156,7 @@ class DatabasePool:
|
|||||||
row = await conn.fetchrow(
|
row = await conn.fetchrow(
|
||||||
"""
|
"""
|
||||||
INSERT INTO packets (
|
INSERT INTO packets (
|
||||||
|
timestamp,
|
||||||
correlation_key,
|
correlation_key,
|
||||||
packet_id,
|
packet_id,
|
||||||
packet_uid,
|
packet_uid,
|
||||||
@@ -197,11 +198,16 @@ class DatabasePool:
|
|||||||
raw
|
raw
|
||||||
) VALUES(
|
) VALUES(
|
||||||
$1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11,$12,$13,$14,$15,$16,$17,$18,$19,$20,
|
$1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11,$12,$13,$14,$15,$16,$17,$18,$19,$20,
|
||||||
$21,$22,$23,$24,$25,$26,$27,$28,$29,$30,$31,$32,$33,$34,$35,$36::jsonb,
|
$21,$22,$23,$24,$25,$26,$27,$28,$29,$30,$31,$32,$33,$34,$35,$36,$37::jsonb,
|
||||||
$37::jsonb,$38::jsonb,$39
|
$38::jsonb,$39::jsonb,$40
|
||||||
)
|
)
|
||||||
ON CONFLICT (correlation_key) DO UPDATE SET
|
ON CONFLICT (correlation_key) DO UPDATE SET
|
||||||
updated_at = NOW(),
|
updated_at = NOW(),
|
||||||
|
timestamp = CASE
|
||||||
|
WHEN packets.timestamp IS NULL THEN EXCLUDED.timestamp
|
||||||
|
WHEN EXCLUDED.timestamp IS NULL THEN packets.timestamp
|
||||||
|
ELSE LEAST(packets.timestamp, EXCLUDED.timestamp)
|
||||||
|
END,
|
||||||
packet_id = COALESCE(EXCLUDED.packet_id, packets.packet_id),
|
packet_id = COALESCE(EXCLUDED.packet_id, packets.packet_id),
|
||||||
packet_uid = COALESCE(EXCLUDED.packet_uid, packets.packet_uid),
|
packet_uid = COALESCE(EXCLUDED.packet_uid, packets.packet_uid),
|
||||||
correlation_source = COALESCE(EXCLUDED.correlation_source, packets.correlation_source),
|
correlation_source = COALESCE(EXCLUDED.correlation_source, packets.correlation_source),
|
||||||
@@ -250,6 +256,7 @@ class DatabasePool:
|
|||||||
raw = COALESCE(EXCLUDED.raw, packets.raw)
|
raw = COALESCE(EXCLUDED.raw, packets.raw)
|
||||||
RETURNING *
|
RETURNING *
|
||||||
""",
|
""",
|
||||||
|
pkt_info.get("timestamp"),
|
||||||
pkt_info["correlation_key"],
|
pkt_info["correlation_key"],
|
||||||
pkt_info.get("packet_id"),
|
pkt_info.get("packet_id"),
|
||||||
pkt_info.get("packet_uid"),
|
pkt_info.get("packet_uid"),
|
||||||
|
|||||||
@@ -175,6 +175,7 @@ class PacketTracker:
|
|||||||
return {
|
return {
|
||||||
"correlation_key": correlation_key,
|
"correlation_key": correlation_key,
|
||||||
"payload": {
|
"payload": {
|
||||||
|
"timestamp": None,
|
||||||
"correlation_key": correlation_key,
|
"correlation_key": correlation_key,
|
||||||
"correlation_source": None,
|
"correlation_source": None,
|
||||||
"packet_id": None,
|
"packet_id": None,
|
||||||
@@ -239,6 +240,12 @@ class PacketTracker:
|
|||||||
continue
|
continue
|
||||||
if value is None:
|
if value is None:
|
||||||
continue
|
continue
|
||||||
|
if key == "timestamp":
|
||||||
|
current_ts = payload.get("timestamp")
|
||||||
|
if current_ts is None or value < current_ts:
|
||||||
|
payload["timestamp"] = value
|
||||||
|
changed = True
|
||||||
|
continue
|
||||||
if key == "raw" and payload.get("raw") is not None:
|
if key == "raw" and payload.get("raw") is not None:
|
||||||
continue
|
continue
|
||||||
if payload.get(key) == value:
|
if payload.get(key) == value:
|
||||||
|
|||||||
File diff suppressed because it is too large
Load Diff
Reference in New Issue
Block a user