fix eth type and db
All checks were successful
Build and Deploy MITM Webserver / build (push) Successful in 1m40s
All checks were successful
Build and Deploy MITM Webserver / build (push) Successful in 1m40s
This commit is contained in:
@@ -12,13 +12,13 @@ class PacketDBModel(BaseModel):
|
||||
id: Union[int, str]
|
||||
timestamp: datetime = Field(..., description="Packet timestamp in ISO format.")
|
||||
packet_uid: str = Field(..., description="Stable packet identity used for upserts/correlation.")
|
||||
iface: Optional[str] = None
|
||||
ingress_if: Optional[str] = None
|
||||
egress_if: Optional[str] = None
|
||||
observed_ifaces: Optional[list[str]] = None
|
||||
src_mac: Optional[str] = None
|
||||
dst_mac: Optional[str] = None
|
||||
eth_type_raw: Optional[int] = Field(None, description="Numeric Ethernet type from the frame header.")
|
||||
eth_type: Optional[Union[int, str]] = None
|
||||
ip_proto_raw: Optional[int] = Field(None, description="Numeric IP protocol / next-header value.")
|
||||
ip_proto: Optional[Union[int, str]] = None
|
||||
src_ip: Optional[IPvAnyAddress] = None
|
||||
dst_ip: Optional[IPvAnyAddress] = None
|
||||
@@ -36,7 +36,6 @@ class PacketDBModel(BaseModel):
|
||||
app_risk_score: Optional[int] = Field(None, description="Count/score of detected nDPI risks.")
|
||||
dpi_metadata: Optional[dict] = Field(None, description="Raw DPI metadata from nDPI.")
|
||||
telemetry_metadata: Optional[dict] = Field(None, description="Kernel telemetry details from eBPF collector.")
|
||||
direction: Optional[str] = None
|
||||
verdict: Optional[str] = None
|
||||
verdict_reason: Optional[str] = None
|
||||
verdict_confidence: Optional[str] = None
|
||||
@@ -51,13 +50,13 @@ class PacketDBModel(BaseModel):
|
||||
"id": 123,
|
||||
"timestamp": "2026-03-05T12:34:56.789Z",
|
||||
"packet_uid": "9f6d3af0d3c81cb20ee8e7d32df7c56414460542",
|
||||
"iface": "eth0",
|
||||
"ingress_if": "eth0",
|
||||
"egress_if": "eth1",
|
||||
"observed_ifaces": ["eth0", "eth1"],
|
||||
"src_mac": "aa:bb:cc:dd:ee:ff",
|
||||
"dst_mac": "11:22:33:44:55:66",
|
||||
"eth_type_raw": 2048,
|
||||
"eth_type": "IPv4",
|
||||
"ip_proto_raw": 6,
|
||||
"ip_proto": "TCP",
|
||||
"src_ip": "192.168.1.10",
|
||||
"dst_ip": "192.168.1.1",
|
||||
@@ -75,7 +74,6 @@ class PacketDBModel(BaseModel):
|
||||
"app_risk_score": 0,
|
||||
"dpi_metadata": {"method": "GET"},
|
||||
"telemetry_metadata": {"event_type": "egress", "iface": "eth1"},
|
||||
"direction": "forwarded",
|
||||
"verdict": "accept",
|
||||
"verdict_reason": "egress-observed",
|
||||
"verdict_confidence": "high",
|
||||
|
||||
@@ -91,11 +91,8 @@ class DatabasePool:
|
||||
"""
|
||||
INSERT INTO packets (
|
||||
packet_uid,
|
||||
iface,
|
||||
ingress_if,
|
||||
egress_if,
|
||||
observed_ifaces,
|
||||
direction,
|
||||
verdict,
|
||||
verdict_reason,
|
||||
verdict_confidence,
|
||||
@@ -104,10 +101,12 @@ class DatabasePool:
|
||||
verdict_seen_at,
|
||||
src_mac,
|
||||
dst_mac,
|
||||
eth_type_raw,
|
||||
eth_type,
|
||||
vlan_id,
|
||||
src_ip,
|
||||
dst_ip,
|
||||
ip_proto_raw,
|
||||
ip_proto,
|
||||
src_port,
|
||||
dst_port,
|
||||
@@ -125,14 +124,11 @@ class DatabasePool:
|
||||
) VALUES(
|
||||
$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::jsonb,$31::jsonb,$32
|
||||
$24,$25,$26,$27,$28,$29::jsonb,$30::jsonb,$31
|
||||
)
|
||||
ON CONFLICT (packet_uid) DO UPDATE SET
|
||||
iface = COALESCE(EXCLUDED.iface, packets.iface),
|
||||
ingress_if = COALESCE(EXCLUDED.ingress_if, packets.ingress_if),
|
||||
egress_if = COALESCE(EXCLUDED.egress_if, packets.egress_if),
|
||||
observed_ifaces = COALESCE(EXCLUDED.observed_ifaces, packets.observed_ifaces),
|
||||
direction = COALESCE(EXCLUDED.direction, packets.direction),
|
||||
verdict = COALESCE(EXCLUDED.verdict, packets.verdict),
|
||||
verdict_reason = COALESCE(EXCLUDED.verdict_reason, packets.verdict_reason),
|
||||
verdict_confidence = COALESCE(EXCLUDED.verdict_confidence, packets.verdict_confidence),
|
||||
@@ -141,10 +137,12 @@ class DatabasePool:
|
||||
verdict_seen_at = COALESCE(EXCLUDED.verdict_seen_at, packets.verdict_seen_at),
|
||||
src_mac = COALESCE(EXCLUDED.src_mac, packets.src_mac),
|
||||
dst_mac = COALESCE(EXCLUDED.dst_mac, packets.dst_mac),
|
||||
eth_type_raw = COALESCE(EXCLUDED.eth_type_raw, packets.eth_type_raw),
|
||||
eth_type = COALESCE(EXCLUDED.eth_type, packets.eth_type),
|
||||
vlan_id = COALESCE(EXCLUDED.vlan_id, packets.vlan_id),
|
||||
src_ip = COALESCE(EXCLUDED.src_ip, packets.src_ip),
|
||||
dst_ip = COALESCE(EXCLUDED.dst_ip, packets.dst_ip),
|
||||
ip_proto_raw = COALESCE(EXCLUDED.ip_proto_raw, packets.ip_proto_raw),
|
||||
ip_proto = COALESCE(EXCLUDED.ip_proto, packets.ip_proto),
|
||||
src_port = COALESCE(EXCLUDED.src_port, packets.src_port),
|
||||
dst_port = COALESCE(EXCLUDED.dst_port, packets.dst_port),
|
||||
@@ -162,11 +160,8 @@ class DatabasePool:
|
||||
RETURNING id, timestamp
|
||||
""",
|
||||
pkt_info["packet_uid"],
|
||||
pkt_info.get("iface"),
|
||||
pkt_info.get("ingress_if"),
|
||||
pkt_info.get("egress_if"),
|
||||
pkt_info.get("observed_ifaces"),
|
||||
pkt_info.get("direction"),
|
||||
pkt_info.get("verdict"),
|
||||
pkt_info.get("verdict_reason"),
|
||||
pkt_info.get("verdict_confidence"),
|
||||
@@ -175,10 +170,12 @@ class DatabasePool:
|
||||
pkt_info.get("verdict_seen_at"),
|
||||
pkt_info.get("src_mac"),
|
||||
pkt_info.get("dst_mac"),
|
||||
pkt_info.get("eth_type_raw"),
|
||||
_db_text(pkt_info.get("eth_type")),
|
||||
pkt_info.get("vlan_id"),
|
||||
pkt_info.get("src_ip"),
|
||||
pkt_info.get("dst_ip"),
|
||||
pkt_info.get("protocol_raw"),
|
||||
_db_text(pkt_info.get("protocol_name") or pkt_info.get("protocol")),
|
||||
pkt_info.get("src_port"),
|
||||
pkt_info.get("dst_port"),
|
||||
|
||||
@@ -10,6 +10,8 @@ from datetime import datetime, timezone
|
||||
from typing import Any, Dict, List, Optional
|
||||
|
||||
import src.shared_objects as shared_objects
|
||||
from src.Models.etherType import EtherTypeEnum, ethertype_from_int
|
||||
from src.Models.ip_protocol import protocol_from_number
|
||||
from src.utilities.packet_identity import build_packet_uid
|
||||
|
||||
logger = logging.getLogger("packet_tracker")
|
||||
@@ -84,22 +86,27 @@ class PacketTracker:
|
||||
continue
|
||||
if payload.get(key) is None:
|
||||
payload[key] = value
|
||||
|
||||
if payload.get("eth_type") is None and payload.get("eth_type_raw") is not None:
|
||||
try:
|
||||
payload["eth_type"] = ethertype_from_int(int(payload["eth_type_raw"]))
|
||||
except Exception:
|
||||
payload["eth_type"] = EtherTypeEnum.UNKNOWN
|
||||
|
||||
if payload.get("protocol") is None and payload.get("protocol_raw") is not None:
|
||||
try:
|
||||
payload["protocol"] = protocol_from_number(int(payload["protocol_raw"]))
|
||||
except Exception:
|
||||
payload["protocol"] = int(payload["protocol_raw"])
|
||||
event_type = event.get("event_type")
|
||||
iface = event.get("iface")
|
||||
if iface:
|
||||
observed_ifaces = payload.setdefault("observed_ifaces", [])
|
||||
if iface not in observed_ifaces:
|
||||
observed_ifaces.append(iface)
|
||||
|
||||
if event_type == "ingress":
|
||||
payload["ingress_if"] = iface
|
||||
payload["iface"] = iface
|
||||
payload["ingress_seen_at"] = _utcnow()
|
||||
payload["direction"] = "ingress"
|
||||
elif event_type == "egress":
|
||||
payload["egress_if"] = iface
|
||||
payload["egress_seen_at"] = _utcnow()
|
||||
payload["direction"] = "forwarded"
|
||||
payload["verdict"] = "accept"
|
||||
payload["verdict_reason"] = "egress-observed"
|
||||
payload["verdict_confidence"] = "high"
|
||||
@@ -109,13 +116,11 @@ class PacketTracker:
|
||||
payload["verdict_reason"] = event.get("reason") or "kfree_skb"
|
||||
payload["verdict_confidence"] = "high"
|
||||
payload["verdict_seen_at"] = _utcnow()
|
||||
payload["direction"] = "dropped"
|
||||
elif event_type == "reject":
|
||||
payload["verdict"] = "reject"
|
||||
payload["verdict_reason"] = event.get("reason") or "netfilter-reject"
|
||||
payload["verdict_confidence"] = event.get("verdict_confidence") or "medium"
|
||||
payload["verdict_seen_at"] = _utcnow()
|
||||
payload["direction"] = "rejected"
|
||||
|
||||
entry["last_observed_at"] = now_ts
|
||||
entry["dirty"] = True
|
||||
@@ -127,7 +132,6 @@ class PacketTracker:
|
||||
"packet_uid": packet_uid,
|
||||
"payload": {
|
||||
"packet_uid": packet_uid,
|
||||
"observed_ifaces": [],
|
||||
"verdict": "pending",
|
||||
"verdict_reason": None,
|
||||
"verdict_confidence": None,
|
||||
@@ -145,7 +149,7 @@ class PacketTracker:
|
||||
payload = entry["payload"]
|
||||
changed = False
|
||||
for key, value in pkt_info.items():
|
||||
if key == "observed_ifaces":
|
||||
if key == "iface":
|
||||
continue
|
||||
if value is None:
|
||||
continue
|
||||
@@ -157,16 +161,10 @@ class PacketTracker:
|
||||
changed = True
|
||||
|
||||
iface = pkt_info.get("iface")
|
||||
if iface:
|
||||
observed_ifaces = payload.setdefault("observed_ifaces", [])
|
||||
if iface not in observed_ifaces:
|
||||
observed_ifaces.append(iface)
|
||||
changed = True
|
||||
if not payload.get("ingress_if"):
|
||||
payload["ingress_if"] = iface
|
||||
payload["iface"] = iface
|
||||
payload["ingress_seen_at"] = _utcnow()
|
||||
changed = True
|
||||
if iface and not payload.get("ingress_if"):
|
||||
payload["ingress_if"] = iface
|
||||
payload["ingress_seen_at"] = _utcnow()
|
||||
changed = True
|
||||
|
||||
payload["last_observed_at"] = now_ts
|
||||
entry["last_observed_at"] = now_ts
|
||||
@@ -208,7 +206,6 @@ class PacketTracker:
|
||||
payload["verdict_reason"] = reject_reason
|
||||
payload["verdict_confidence"] = "medium"
|
||||
payload["verdict_seen_at"] = _utcnow()
|
||||
payload["direction"] = "rejected"
|
||||
match["last_observed_at"] = now_ts
|
||||
match["dirty"] = True
|
||||
match["finalized"] = True
|
||||
@@ -262,7 +259,6 @@ class PacketTracker:
|
||||
entry["payload"]["verdict_reason"] = "timeout"
|
||||
entry["payload"]["verdict_confidence"] = "low"
|
||||
entry["payload"]["verdict_seen_at"] = _utcnow()
|
||||
entry["payload"]["direction"] = "observed"
|
||||
entry["dirty"] = True
|
||||
|
||||
should_flush = entry["dirty"] and (
|
||||
@@ -288,8 +284,6 @@ class PacketTracker:
|
||||
|
||||
def _persist(self, entry: Dict[str, Any]) -> None:
|
||||
payload = dict(entry["payload"])
|
||||
if payload.get("observed_ifaces") == []:
|
||||
payload["observed_ifaces"] = None
|
||||
|
||||
web_loop = getattr(shared_objects, "web_loop", None)
|
||||
web_db = getattr(shared_objects, "db", None)
|
||||
|
||||
Reference in New Issue
Block a user