test new flow stream analysis
This commit is contained in:
@@ -7,6 +7,7 @@ from fastapi import APIRouter, HTTPException, Query
|
||||
from pydantic import BaseModel, Field
|
||||
|
||||
import src.shared_objects as shared
|
||||
from src.Models.packets import PacketDBModel
|
||||
|
||||
router = APIRouter()
|
||||
|
||||
@@ -128,8 +129,11 @@ class ConversationEvidence(BaseModel):
|
||||
ethernet_protocol: Optional[str] = None
|
||||
ip_protocol: Optional[str] = None
|
||||
hostnames: List[str] = Field(default_factory=list)
|
||||
flow_ids: List[str] = Field(default_factory=list)
|
||||
flow_count: int = 0
|
||||
packet_count: int
|
||||
byte_count: int
|
||||
duration_ms: int = 0
|
||||
first_seen: datetime
|
||||
last_seen: datetime
|
||||
accept_count: int = 0
|
||||
@@ -150,6 +154,43 @@ class ConversationAnalysisResponse(BaseModel):
|
||||
)
|
||||
|
||||
|
||||
class ConversationFlowSummaryEvidence(BaseModel):
|
||||
flow_id: str
|
||||
protocol: str
|
||||
packet_count: int
|
||||
byte_count: int
|
||||
first_seen: datetime
|
||||
last_seen: datetime
|
||||
client_label: str
|
||||
server_label: str
|
||||
request_count: int = 0
|
||||
response_count: int = 0
|
||||
|
||||
|
||||
class ConversationFlowEventEvidence(BaseModel):
|
||||
flow_id: str
|
||||
timestamp: datetime
|
||||
kind: str
|
||||
label: str
|
||||
src_label: str
|
||||
dst_label: str
|
||||
packet_id: Optional[str] = None
|
||||
|
||||
|
||||
class ConversationFlowDetailResponse(BaseModel):
|
||||
since: Optional[datetime] = None
|
||||
packets: List[PacketDBModel] = Field(default_factory=list)
|
||||
flows: List[ConversationFlowSummaryEvidence] = Field(default_factory=list)
|
||||
events: List[ConversationFlowEventEvidence] = Field(default_factory=list)
|
||||
notes: List[str] = Field(
|
||||
default_factory=lambda: [
|
||||
"This detail view is reconstructed from ordered captured packets for one directional conversation or flow.",
|
||||
"Events are inferred from HTTP metadata and TCP packet types, so lower-layer traffic may have fewer high-level annotations.",
|
||||
"If multiple flow ids exist for the same directional tuple, the drawer shows all matching packets in timestamp order.",
|
||||
]
|
||||
)
|
||||
|
||||
|
||||
class LabelCountEvidence(BaseModel):
|
||||
label: str
|
||||
packet_count: int
|
||||
@@ -453,6 +494,69 @@ async def analysis_conversations(
|
||||
return ConversationAnalysisResponse(since=since, conversations=conversations)
|
||||
|
||||
|
||||
@router.get("/conversation-flow-detail", response_model=ConversationFlowDetailResponse)
|
||||
async def analysis_conversation_flow_detail(
|
||||
flow_id: Optional[str] = Query(
|
||||
None,
|
||||
description="Specific flow_id to inspect. If omitted, the directional conversation tuple is used.",
|
||||
),
|
||||
src_ip_address: Optional[str] = Query(None),
|
||||
src_mac_address: Optional[str] = Query(None),
|
||||
dst_ip_address: Optional[str] = Query(None),
|
||||
dst_mac_address: Optional[str] = Query(None),
|
||||
src_port: Optional[int] = Query(None, ge=0, le=65535),
|
||||
dst_port: Optional[int] = Query(None, ge=0, le=65535),
|
||||
protocol: Optional[str] = Query(None),
|
||||
since_minutes: Optional[int] = Query(
|
||||
None,
|
||||
ge=1,
|
||||
le=60 * 24 * 30,
|
||||
description="Analyze only packets seen within the last N minutes. Omit to cover all captured history.",
|
||||
),
|
||||
limit_packets: int = Query(
|
||||
1500,
|
||||
ge=1,
|
||||
le=10000,
|
||||
description="Maximum number of packets returned for this conversation detail view.",
|
||||
),
|
||||
) -> ConversationFlowDetailResponse:
|
||||
"""Return ordered packet detail, subflows, and derived request/response events for one conversation."""
|
||||
db = shared.db
|
||||
if db is None:
|
||||
raise HTTPException(status_code=503, detail="Database not available")
|
||||
|
||||
if flow_id is None and all(
|
||||
value is None
|
||||
for value in [src_ip_address, src_mac_address, dst_ip_address, dst_mac_address, src_port, dst_port]
|
||||
):
|
||||
raise HTTPException(
|
||||
status_code=400,
|
||||
detail="Provide either flow_id or enough directional conversation fields to identify the conversation.",
|
||||
)
|
||||
|
||||
since: Optional[datetime] = None
|
||||
if since_minutes is not None:
|
||||
since = datetime.now(timezone.utc) - timedelta(minutes=since_minutes)
|
||||
|
||||
try:
|
||||
result = await db.fetch_conversation_flow_detail(
|
||||
flow_id=flow_id,
|
||||
src_ip_address=src_ip_address,
|
||||
src_mac_address=src_mac_address,
|
||||
dst_ip_address=dst_ip_address,
|
||||
dst_mac_address=dst_mac_address,
|
||||
src_port=src_port,
|
||||
dst_port=dst_port,
|
||||
protocol=protocol,
|
||||
since=since,
|
||||
limit_packets=limit_packets,
|
||||
)
|
||||
except Exception as exc:
|
||||
raise HTTPException(status_code=500, detail=f"Failed to analyze conversation flow detail: {exc}") from exc
|
||||
|
||||
return ConversationFlowDetailResponse(since=since, **result)
|
||||
|
||||
|
||||
@router.get("/host-intelligence", response_model=HostIntelligenceAnalysisResponse)
|
||||
async def analysis_host_intelligence(
|
||||
since_minutes: Optional[int] = Query(
|
||||
|
||||
@@ -1278,6 +1278,8 @@ class DatabasePool:
|
||||
ip_proto_raw,
|
||||
eth_type_raw,
|
||||
ARRAY_REMOVE(ARRAY_AGG(DISTINCT NULLIF(app_hostname::text, '')), NULL) AS hostnames,
|
||||
ARRAY_REMOVE(ARRAY_AGG(DISTINCT NULLIF(flow_id::text, '')), NULL) AS flow_ids,
|
||||
COUNT(DISTINCT NULLIF(flow_id::text, '')) AS flow_count,
|
||||
COUNT(*) AS packet_count,
|
||||
COALESCE(SUM(length), 0) AS byte_count,
|
||||
MIN(timestamp) AS first_seen,
|
||||
@@ -1339,8 +1341,21 @@ class DatabasePool:
|
||||
"ethernet_protocol": _analysis_ethernet_protocol_name(record.get("eth_type_raw")),
|
||||
"ip_protocol": _analysis_ip_protocol_name(record.get("ip_proto_raw")),
|
||||
"hostnames": list(record.get("hostnames") or []),
|
||||
"flow_ids": list(record.get("flow_ids") or []),
|
||||
"flow_count": int(record.get("flow_count") or 0),
|
||||
"packet_count": int(record.get("packet_count") or 0),
|
||||
"byte_count": int(record.get("byte_count") or 0),
|
||||
"duration_ms": max(
|
||||
0,
|
||||
int(
|
||||
(
|
||||
last_seen_raw - first_seen_raw
|
||||
).total_seconds()
|
||||
* 1000
|
||||
)
|
||||
if hasattr(last_seen_raw, "timestamp") and hasattr(first_seen_raw, "timestamp")
|
||||
else 0,
|
||||
),
|
||||
"first_seen": first_seen_raw.isoformat() if hasattr(first_seen_raw, "isoformat") else first_seen_raw,
|
||||
"last_seen": last_seen_raw.isoformat() if hasattr(last_seen_raw, "isoformat") else last_seen_raw,
|
||||
"accept_count": int(record.get("accept_count") or 0),
|
||||
@@ -1351,6 +1366,161 @@ class DatabasePool:
|
||||
)
|
||||
return result
|
||||
|
||||
async def fetch_conversation_flow_detail(
|
||||
self,
|
||||
*,
|
||||
flow_id: Optional[str] = None,
|
||||
src_ip_address: Optional[str] = None,
|
||||
src_mac_address: Optional[str] = None,
|
||||
dst_ip_address: Optional[str] = None,
|
||||
dst_mac_address: Optional[str] = None,
|
||||
src_port: Optional[int] = None,
|
||||
dst_port: Optional[int] = None,
|
||||
protocol: Optional[str] = None,
|
||||
since: Optional[datetime] = None,
|
||||
limit_packets: int = 1500,
|
||||
) -> Dict[str, Any]:
|
||||
"""Fetch packet-level detail for one conversation or flow."""
|
||||
if self._pool is None:
|
||||
await self.init_pool()
|
||||
|
||||
query = """
|
||||
SELECT *
|
||||
FROM packets
|
||||
WHERE ($1::timestamptz IS NULL OR timestamp >= $1)
|
||||
AND (
|
||||
($2::text IS NOT NULL AND flow_id = $2)
|
||||
OR (
|
||||
$2::text IS NULL
|
||||
AND (
|
||||
(
|
||||
src_ip::text IS NOT DISTINCT FROM $3
|
||||
AND src_mac::text IS NOT DISTINCT FROM $4
|
||||
AND dst_ip::text IS NOT DISTINCT FROM $5
|
||||
AND dst_mac::text IS NOT DISTINCT FROM $6
|
||||
AND src_port IS NOT DISTINCT FROM $7
|
||||
AND dst_port IS NOT DISTINCT FROM $8
|
||||
)
|
||||
OR (
|
||||
src_ip::text IS NOT DISTINCT FROM $5
|
||||
AND src_mac::text IS NOT DISTINCT FROM $6
|
||||
AND dst_ip::text IS NOT DISTINCT FROM $3
|
||||
AND dst_mac::text IS NOT DISTINCT FROM $4
|
||||
AND src_port IS NOT DISTINCT FROM $8
|
||||
AND dst_port IS NOT DISTINCT FROM $7
|
||||
)
|
||||
)
|
||||
AND (
|
||||
$9::text IS NULL
|
||||
OR COALESCE(NULLIF(app_protocol::text, ''), '') = $9
|
||||
OR COALESCE(NULLIF(app_protocol::text, ''), '') = ''
|
||||
)
|
||||
)
|
||||
)
|
||||
ORDER BY timestamp ASC, id ASC
|
||||
LIMIT $10
|
||||
"""
|
||||
|
||||
async with self._pool.acquire() as conn:
|
||||
rows = await conn.fetch(
|
||||
query,
|
||||
since,
|
||||
flow_id,
|
||||
src_ip_address,
|
||||
src_mac_address,
|
||||
dst_ip_address,
|
||||
dst_mac_address,
|
||||
src_port,
|
||||
dst_port,
|
||||
protocol,
|
||||
limit_packets,
|
||||
)
|
||||
|
||||
packets: List[PacketDBModel] = []
|
||||
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)):
|
||||
data["raw_b64"] = base64.b64encode(raw_val).decode("ascii")
|
||||
data.pop("raw", None)
|
||||
|
||||
try:
|
||||
packets.append(PacketDBModel(**data))
|
||||
except ValidationError as exc:
|
||||
logger.warning("Skipping flow-detail packet row validation failure (id=%s): %s", data.get("id"), exc)
|
||||
|
||||
flow_summaries: Dict[str, Dict[str, Any]] = {}
|
||||
request_response_events: List[Dict[str, Any]] = []
|
||||
|
||||
for packet in packets:
|
||||
subflow_id = packet.flow_id or f"tuple:{packet.src_ip}:{packet.src_port}->{packet.dst_ip}:{packet.dst_port}"
|
||||
summary = flow_summaries.get(subflow_id)
|
||||
if summary is None:
|
||||
summary = {
|
||||
"flow_id": subflow_id,
|
||||
"protocol": packet.app_protocol or packet.ip_proto or packet.eth_type or "UNKNOWN",
|
||||
"packet_count": 0,
|
||||
"byte_count": 0,
|
||||
"first_seen": packet.timestamp.isoformat(),
|
||||
"last_seen": packet.timestamp.isoformat(),
|
||||
"client_label": f"{packet.src_ip or packet.src_mac or 'unknown'}:{packet.src_port or '-'}",
|
||||
"server_label": f"{packet.dst_ip or packet.dst_mac or 'unknown'}:{packet.dst_port or '-'}",
|
||||
"request_count": 0,
|
||||
"response_count": 0,
|
||||
}
|
||||
flow_summaries[subflow_id] = summary
|
||||
|
||||
summary["packet_count"] += 1
|
||||
summary["byte_count"] += int(packet.length or 0)
|
||||
packet_ts = packet.timestamp.isoformat()
|
||||
if packet_ts < str(summary["first_seen"]):
|
||||
summary["first_seen"] = packet_ts
|
||||
if packet_ts > str(summary["last_seen"]):
|
||||
summary["last_seen"] = packet_ts
|
||||
|
||||
dpi_metadata = packet.dpi_metadata if isinstance(packet.dpi_metadata, dict) else {}
|
||||
http_meta = dpi_metadata.get("http") if isinstance(dpi_metadata.get("http"), dict) else {}
|
||||
tcp_meta = dpi_metadata.get("tcp") if isinstance(dpi_metadata.get("tcp"), dict) else {}
|
||||
|
||||
event_label = None
|
||||
event_kind = None
|
||||
if http_meta.get("method"):
|
||||
event_kind = "request"
|
||||
event_label = f"{http_meta.get('method')} {http_meta.get('uri') or http_meta.get('path') or ''}".strip()
|
||||
summary["request_count"] += 1
|
||||
elif http_meta.get("response_code") is not None:
|
||||
event_kind = "response"
|
||||
event_label = f"{http_meta.get('response_code')} {http_meta.get('response_phrase') or ''}".strip()
|
||||
summary["response_count"] += 1
|
||||
elif tcp_meta.get("packet_type"):
|
||||
event_kind = "tcp"
|
||||
event_label = str(tcp_meta.get("packet_type"))
|
||||
|
||||
if event_label:
|
||||
request_response_events.append(
|
||||
{
|
||||
"flow_id": subflow_id,
|
||||
"timestamp": packet.timestamp.isoformat(),
|
||||
"kind": event_kind,
|
||||
"label": event_label,
|
||||
"src_label": f"{packet.src_ip or packet.src_mac or 'unknown'}:{packet.src_port or '-'}",
|
||||
"dst_label": f"{packet.dst_ip or packet.dst_mac or 'unknown'}:{packet.dst_port or '-'}",
|
||||
"packet_id": packet.id,
|
||||
}
|
||||
)
|
||||
|
||||
return {
|
||||
"packets": packets,
|
||||
"flows": sorted(
|
||||
flow_summaries.values(),
|
||||
key=lambda item: (str(item.get("first_seen") or ""), str(item.get("flow_id") or "")),
|
||||
),
|
||||
"events": request_response_events,
|
||||
}
|
||||
|
||||
async def analyze_host_intelligence(
|
||||
self,
|
||||
*,
|
||||
|
||||
@@ -3,6 +3,7 @@ import axios from 'axios';
|
||||
import {
|
||||
AnomalyAnalysisResponse,
|
||||
ConversationAnalysisResponse,
|
||||
ConversationFlowDetailResponse,
|
||||
DiscoveryAnalysisResponse,
|
||||
HostIntelligenceAnalysisResponse,
|
||||
InterfaceHostAnalysisResponse,
|
||||
@@ -235,6 +236,46 @@ export const fetchConversationAnalysis = async (
|
||||
return res.data;
|
||||
};
|
||||
|
||||
export const fetchConversationFlowDetail = async ({
|
||||
flowId,
|
||||
srcIpAddress,
|
||||
srcMacAddress,
|
||||
dstIpAddress,
|
||||
dstMacAddress,
|
||||
srcPort,
|
||||
dstPort,
|
||||
protocol,
|
||||
sinceMinutes = null,
|
||||
limitPackets = 1500,
|
||||
}: {
|
||||
flowId?: string | null;
|
||||
srcIpAddress?: string | null;
|
||||
srcMacAddress?: string | null;
|
||||
dstIpAddress?: string | null;
|
||||
dstMacAddress?: string | null;
|
||||
srcPort?: number | null;
|
||||
dstPort?: number | null;
|
||||
protocol?: string | null;
|
||||
sinceMinutes?: number | null;
|
||||
limitPackets?: number;
|
||||
}): Promise<ConversationFlowDetailResponse> => {
|
||||
const res = await api.get<ConversationFlowDetailResponse>('/analysis/conversation-flow-detail', {
|
||||
params: {
|
||||
flow_id: flowId ?? undefined,
|
||||
src_ip_address: srcIpAddress ?? undefined,
|
||||
src_mac_address: srcMacAddress ?? undefined,
|
||||
dst_ip_address: dstIpAddress ?? undefined,
|
||||
dst_mac_address: dstMacAddress ?? undefined,
|
||||
src_port: srcPort ?? undefined,
|
||||
dst_port: dstPort ?? undefined,
|
||||
protocol: protocol ?? undefined,
|
||||
since_minutes: sinceMinutes ?? undefined,
|
||||
limit_packets: limitPackets,
|
||||
},
|
||||
});
|
||||
return res.data;
|
||||
};
|
||||
|
||||
export const fetchHostIntelligenceAnalysis = async (
|
||||
sinceMinutes: number | null = null,
|
||||
limitHosts = 40,
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
@@ -1,3 +1,5 @@
|
||||
import type { PacketRow } from './packets';
|
||||
|
||||
export interface InterfaceHostEvidence {
|
||||
ip_address?: string | null;
|
||||
mac_address?: string | null;
|
||||
@@ -94,8 +96,11 @@ export interface ConversationEvidence {
|
||||
ethernet_protocol?: string | null;
|
||||
ip_protocol?: string | null;
|
||||
hostnames: string[];
|
||||
flow_ids: string[];
|
||||
flow_count: number;
|
||||
packet_count: number;
|
||||
byte_count: number;
|
||||
duration_ms: number;
|
||||
first_seen: string;
|
||||
last_seen: string;
|
||||
accept_count: number;
|
||||
@@ -110,6 +115,37 @@ export interface ConversationAnalysisResponse {
|
||||
notes: string[];
|
||||
}
|
||||
|
||||
export interface ConversationFlowSummaryEvidence {
|
||||
flow_id: string;
|
||||
protocol: string;
|
||||
packet_count: number;
|
||||
byte_count: number;
|
||||
first_seen: string;
|
||||
last_seen: string;
|
||||
client_label: string;
|
||||
server_label: string;
|
||||
request_count: number;
|
||||
response_count: number;
|
||||
}
|
||||
|
||||
export interface ConversationFlowEventEvidence {
|
||||
flow_id: string;
|
||||
timestamp: string;
|
||||
kind: string;
|
||||
label: string;
|
||||
src_label: string;
|
||||
dst_label: string;
|
||||
packet_id?: string | null;
|
||||
}
|
||||
|
||||
export interface ConversationFlowDetailResponse {
|
||||
since?: string | null;
|
||||
packets: PacketRow[];
|
||||
flows: ConversationFlowSummaryEvidence[];
|
||||
events: ConversationFlowEventEvidence[];
|
||||
notes: string[];
|
||||
}
|
||||
|
||||
export interface LabelCountEvidence {
|
||||
label: string;
|
||||
packet_count: number;
|
||||
|
||||
Reference in New Issue
Block a user