Add capture diagnostics and replay fixtures
Add privacy-safe capture diagnostics and sanitized replay fixtures without changing the existing export format, record order, or UID generation.
This commit is contained in:
parent
73bb100cd6
commit
2b33ea63af
11 changed files with 445 additions and 13 deletions
|
|
@ -119,4 +119,5 @@ def decode_mitmproxy_flows(path: str | Path, flow_index: int | None = None) -> d
|
|||
"arc_rows": arc_rows,
|
||||
"arc_warnings": arc_warnings,
|
||||
"user_uid": session.user_uid or user_uid,
|
||||
"capture_diagnostics": session.diagnostic_report(),
|
||||
}
|
||||
|
|
|
|||
|
|
@ -10,6 +10,7 @@ from nte_history_exporter.decoder.boundary import annotate_groups
|
|||
from nte_history_exporter.export.csv_export import write_csv
|
||||
from nte_history_exporter.export.json_export import build_export_json
|
||||
from nte_history_exporter.live_capture.libpcap import LibpcapUnavailable
|
||||
from nte_history_exporter.live_capture.diagnostics import new_diagnostics_path, write_capture_diagnostics
|
||||
from nte_history_exporter.live_capture.runner import export_paths, run_live_capture
|
||||
from nte_history_exporter.update_check import check_for_update
|
||||
|
||||
|
|
@ -32,7 +33,11 @@ def build_parser() -> argparse.ArgumentParser:
|
|||
),
|
||||
)
|
||||
parser.add_argument("--copy-clipboard", action="store_true", help="copy a single live export to clipboard")
|
||||
parser.add_argument("--debug", action="store_true", help="also write research CSVs next to the JSON exports")
|
||||
parser.add_argument(
|
||||
"--debug",
|
||||
action="store_true",
|
||||
help="also write a research CSV and privacy-safe capture diagnostics",
|
||||
)
|
||||
parser.add_argument("--user-uid", default=None, help="override the auto-detected NTE user UID in the JSON export")
|
||||
return parser
|
||||
|
||||
|
|
@ -79,8 +84,11 @@ def main(argv: list[str] | None = None) -> int:
|
|||
resolved_user_uid = console.prompt_user_uid()
|
||||
|
||||
out_path, json_path = export_paths(kind, resolved_user_uid)
|
||||
diagnostics_path = None
|
||||
if args.debug:
|
||||
write_csv(out_path, rows)
|
||||
diagnostics_path = new_diagnostics_path(out_path.parent)
|
||||
write_capture_diagnostics(diagnostics_path, decoded["capture_diagnostics"])
|
||||
export = build_export_json(
|
||||
rows,
|
||||
warnings,
|
||||
|
|
@ -105,6 +113,7 @@ def main(argv: list[str] | None = None) -> int:
|
|||
print()
|
||||
if args.debug:
|
||||
console.print_note(f"CSV written: {out_path}")
|
||||
console.print_note(f"Diagnostics written: {diagnostics_path}")
|
||||
console.print_note(f"Export written: {json_path}")
|
||||
return 0
|
||||
|
||||
|
|
|
|||
90
src/nte_history_exporter/live_capture/diagnostics.py
Normal file
90
src/nte_history_exporter/live_capture/diagnostics.py
Normal file
|
|
@ -0,0 +1,90 @@
|
|||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
from collections import Counter
|
||||
from datetime import datetime
|
||||
from pathlib import Path
|
||||
from typing import Any, Iterable
|
||||
|
||||
|
||||
DIAGNOSTIC_FORMAT = "nte-capture-diagnostics"
|
||||
DIAGNOSTIC_FORMAT_VERSION = 1
|
||||
MAX_EVENTS = 200
|
||||
|
||||
|
||||
class CaptureDiagnostics:
|
||||
"""Collect bounded, privacy-safe observations about a capture session."""
|
||||
|
||||
def __init__(self) -> None:
|
||||
self.counters: Counter[str] = Counter()
|
||||
self.event_counts: Counter[str] = Counter()
|
||||
self.reason_counts: Counter[str] = Counter()
|
||||
self.events: list[dict[str, Any]] = []
|
||||
|
||||
def observe_packet(self, protocol: str) -> None:
|
||||
self.counters["packets_seen"] += 1
|
||||
if protocol == "udp":
|
||||
self.counters["udp_packets_seen"] += 1
|
||||
else:
|
||||
self.counters["non_udp_packets_ignored"] += 1
|
||||
|
||||
def add_event(
|
||||
self,
|
||||
code: str,
|
||||
packet_index: int,
|
||||
*,
|
||||
reason: bool = False,
|
||||
**fields: Any,
|
||||
) -> None:
|
||||
self.event_counts[code] += 1
|
||||
if reason:
|
||||
self.reason_counts[code] += 1
|
||||
if len(self.events) >= MAX_EVENTS:
|
||||
self.counters["events_omitted"] += 1
|
||||
return
|
||||
event: dict[str, Any] = {"packet_index": packet_index, "code": code}
|
||||
event.update({key: value for key, value in fields.items() if value is not None})
|
||||
self.events.append(event)
|
||||
|
||||
def report(self, pending_requests: Iterable[Any]) -> dict[str, Any]:
|
||||
pending = [
|
||||
{
|
||||
"kind": request.kind,
|
||||
"page": request.page,
|
||||
"response_candidates": request.response_candidates,
|
||||
"response_candidate_lengths": list(request.response_candidate_lengths),
|
||||
}
|
||||
for request in pending_requests
|
||||
]
|
||||
return {
|
||||
"format": DIAGNOSTIC_FORMAT,
|
||||
"format_version": DIAGNOSTIC_FORMAT_VERSION,
|
||||
"privacy": {
|
||||
"contains_network_addresses": False,
|
||||
"contains_network_ports": False,
|
||||
"contains_packet_timestamps": False,
|
||||
"contains_payload_bytes": False,
|
||||
"contains_user_uid": False,
|
||||
},
|
||||
"counters": dict(sorted(self.counters.items())),
|
||||
"event_counts": dict(sorted(self.event_counts.items())),
|
||||
"reason_counts": dict(sorted(self.reason_counts.items())),
|
||||
"events": list(self.events),
|
||||
"pending_requests": pending,
|
||||
}
|
||||
|
||||
|
||||
def new_diagnostics_path(output_dir: str | Path = "exports") -> Path:
|
||||
directory = Path(output_dir)
|
||||
directory.mkdir(parents=True, exist_ok=True)
|
||||
stamp = datetime.now().strftime("%Y%m%d_%H%M%S")
|
||||
path = directory / f"Capture_{stamp}.diagnostics.json"
|
||||
counter = 2
|
||||
while path.exists():
|
||||
path = directory / f"Capture_{stamp}_{counter}.diagnostics.json"
|
||||
counter += 1
|
||||
return path
|
||||
|
||||
|
||||
def write_capture_diagnostics(path: str | Path, report: dict[str, Any]) -> None:
|
||||
Path(path).write_text(json.dumps(report, ensure_ascii=False, indent=2), encoding="utf-8")
|
||||
|
|
@ -13,6 +13,7 @@ from nte_history_exporter.decoder.boundary import annotate_groups, select_contin
|
|||
from nte_history_exporter.export.csv_export import write_csv
|
||||
from nte_history_exporter.export.json_export import build_export_json
|
||||
from nte_history_exporter.live_capture.backends import open_capture_backend
|
||||
from nte_history_exporter.live_capture.diagnostics import new_diagnostics_path, write_capture_diagnostics
|
||||
from nte_history_exporter.live_capture.session import LiveHistorySession, UdpPacket
|
||||
from nte_history_exporter.live_capture.stop_key import StopKeyMonitor
|
||||
from nte_history_exporter.live_capture.windows_raw import detect_local_ipv4
|
||||
|
|
@ -135,6 +136,10 @@ def run_live_capture(
|
|||
)
|
||||
|
||||
exports = []
|
||||
diagnostics_path = None
|
||||
if write_debug_csv:
|
||||
diagnostics_path = new_diagnostics_path()
|
||||
write_capture_diagnostics(diagnostics_path, session.diagnostic_report())
|
||||
resolved_user_uid = user_uid or session.user_uid
|
||||
if session.kinds_seen() and not resolved_user_uid:
|
||||
resolved_user_uid = console.prompt_user_uid()
|
||||
|
|
@ -164,6 +169,7 @@ def run_live_capture(
|
|||
{
|
||||
"kind": kind,
|
||||
"csv_path": csv_path if write_debug_csv else None,
|
||||
"diagnostics_path": diagnostics_path if write_debug_csv else None,
|
||||
"json_path": json_path,
|
||||
"export": export,
|
||||
"payload": payload,
|
||||
|
|
@ -176,7 +182,9 @@ def run_live_capture(
|
|||
console.print_note("Make sure the capture backend is running, then reopen the")
|
||||
console.print_note("history screen and scroll from page 1. If no page messages")
|
||||
console.print_note("appear, return to the main menu and re-enter the game.")
|
||||
return {"exports": []}
|
||||
if diagnostics_path is not None:
|
||||
console.print_note(f"Diagnostics written: {diagnostics_path}")
|
||||
return {"exports": [], "diagnostics_path": diagnostics_path}
|
||||
|
||||
for item in exports:
|
||||
scan = item["export"]["scan"]
|
||||
|
|
@ -194,6 +202,8 @@ def run_live_capture(
|
|||
if item["csv_path"] is not None:
|
||||
console.print_note(f"CSV written: {item['csv_path']}")
|
||||
console.print_note(f"Export written: {item['json_path']}")
|
||||
if diagnostics_path is not None:
|
||||
console.print_note(f"Diagnostics written: {diagnostics_path}")
|
||||
|
||||
if copy_clipboard and len(exports) == 1:
|
||||
if copy_to_clipboard(exports[0]["payload"]):
|
||||
|
|
@ -203,7 +213,7 @@ def run_live_capture(
|
|||
elif copy_clipboard and len(exports) > 1:
|
||||
console.print_note("Multiple banners captured; clipboard copy skipped so one export")
|
||||
console.print_note("does not overwrite another.")
|
||||
return {"exports": exports}
|
||||
return {"exports": exports, "diagnostics_path": diagnostics_path}
|
||||
|
||||
|
||||
def export_paths(kind: str, user_uid: str | None = None) -> tuple[Path, Path]:
|
||||
|
|
|
|||
|
|
@ -17,9 +17,12 @@ from nte_history_exporter.decoder.protocol import (
|
|||
history_request_kind,
|
||||
is_history_request,
|
||||
request_page,
|
||||
response_contains_history_marker,
|
||||
)
|
||||
from nte_history_exporter.decoder.run import build_rows_from_pairs
|
||||
from nte_history_exporter.decoder.structured_protocol import FORK_MARKER, MONOPOLY_MARKER
|
||||
from nte_history_exporter.decoder.user_uid import extract_user_uid_candidates
|
||||
from nte_history_exporter.live_capture.diagnostics import CaptureDiagnostics
|
||||
|
||||
|
||||
@dataclass
|
||||
|
|
@ -61,6 +64,7 @@ class LiveHistorySession:
|
|||
self.unanswered_pages: dict[str, dict[int, str]] = {}
|
||||
self.user_uid: str | None = None
|
||||
self.user_uid_candidates: Counter[str] = Counter()
|
||||
self.diagnostics = CaptureDiagnostics()
|
||||
|
||||
def _mark_unanswered(self, request: PendingRequest) -> None:
|
||||
if request.response_candidates:
|
||||
|
|
@ -89,6 +93,14 @@ class LiveHistorySession:
|
|||
)
|
||||
if same_stream and (request.page == 1 or pending.page == request.page):
|
||||
self._mark_unanswered(pending)
|
||||
self.diagnostics.add_event(
|
||||
"REQUEST_REPLACED",
|
||||
request.request_msg,
|
||||
reason=True,
|
||||
kind=pending.kind,
|
||||
page=pending.page,
|
||||
response_candidates=pending.response_candidates,
|
||||
)
|
||||
else:
|
||||
retained.append(pending)
|
||||
self.pending = retained
|
||||
|
|
@ -97,6 +109,7 @@ class LiveHistorySession:
|
|||
|
||||
def process_packet(self, packet: UdpPacket) -> bool:
|
||||
self.packet_count += 1
|
||||
self.diagnostics.observe_packet(packet.protocol)
|
||||
candidates = extract_user_uid_candidates(packet.payload)
|
||||
if candidates:
|
||||
self.user_uid_candidates.update(candidates)
|
||||
|
|
@ -118,6 +131,13 @@ class LiveHistorySession:
|
|||
dst_port=packet.dst_port,
|
||||
)
|
||||
self._queue_request(req)
|
||||
self.diagnostics.counters["history_requests_recognized"] += 1
|
||||
self.diagnostics.add_event(
|
||||
"HISTORY_REQUEST_RECOGNIZED",
|
||||
self.packet_count,
|
||||
kind=req.kind,
|
||||
page=req.page,
|
||||
)
|
||||
self.last_page_seen = req.page
|
||||
return False
|
||||
|
||||
|
|
@ -135,10 +155,17 @@ class LiveHistorySession:
|
|||
dst_port=packet.dst_port,
|
||||
)
|
||||
self._queue_request(req)
|
||||
self.diagnostics.counters["history_requests_recognized"] += 1
|
||||
self.diagnostics.add_event(
|
||||
"HISTORY_REQUEST_RECOGNIZED",
|
||||
self.packet_count,
|
||||
kind=req.kind,
|
||||
page=req.page,
|
||||
)
|
||||
self.last_page_seen = req.page
|
||||
return False
|
||||
|
||||
if packet.dst_ip != self.local_ip or len(packet.payload) < 100:
|
||||
if packet.dst_ip != self.local_ip:
|
||||
return False
|
||||
|
||||
connection_candidates = [
|
||||
|
|
@ -154,6 +181,14 @@ class LiveHistorySession:
|
|||
if not connection_candidates:
|
||||
return False
|
||||
|
||||
if len(packet.payload) < 100:
|
||||
self._record_rejected_candidate(
|
||||
connection_candidates,
|
||||
packet.payload,
|
||||
"RESPONSE_TOO_SHORT",
|
||||
)
|
||||
return False
|
||||
|
||||
monopoly_records = decode_response_records(packet.payload)
|
||||
arc_records = parse_arc_response(packet.payload) if not monopoly_records else []
|
||||
if monopoly_records:
|
||||
|
|
@ -166,15 +201,36 @@ class LiveHistorySession:
|
|||
candidates = connection_candidates
|
||||
records = []
|
||||
|
||||
if not records or not candidates:
|
||||
for req in candidates:
|
||||
req.response_candidates += 1
|
||||
req.response_candidate_lengths = (
|
||||
*req.response_candidate_lengths[-4:],
|
||||
len(packet.payload),
|
||||
)
|
||||
if records and not candidates:
|
||||
self._record_rejected_candidate(
|
||||
connection_candidates,
|
||||
packet.payload,
|
||||
"RESPONSE_KIND_MISMATCH",
|
||||
)
|
||||
return False
|
||||
|
||||
if not records:
|
||||
reason = (
|
||||
"HISTORY_MARKER_PARSE_FAILED"
|
||||
if response_contains_history_marker(packet.payload)
|
||||
or MONOPOLY_MARKER in packet.payload
|
||||
or FORK_MARKER in packet.payload
|
||||
else "NO_HISTORY_MARKER"
|
||||
)
|
||||
self._record_rejected_candidate(candidates, packet.payload, reason)
|
||||
return False
|
||||
|
||||
decoder_modes = sorted({record.get("decoder_mode", "heuristic") for record in records})
|
||||
self.diagnostics.counters["history_responses_decoded"] += 1
|
||||
self.diagnostics.add_event(
|
||||
"HISTORY_RESPONSE_DECODED",
|
||||
self.packet_count,
|
||||
kind=candidates[0].kind,
|
||||
payload_length=len(packet.payload),
|
||||
record_count=len(records),
|
||||
decoder_modes=decoder_modes,
|
||||
)
|
||||
|
||||
page_count = max(1, (len(records) + 4) // 5)
|
||||
if len(records) < 5:
|
||||
selected = [candidates[-1]]
|
||||
|
|
@ -208,9 +264,43 @@ class LiveHistorySession:
|
|||
)
|
||||
self.last_match_time = packet.timestamp
|
||||
self.last_page_seen = req.page
|
||||
self.diagnostics.counters["pages_matched"] += 1
|
||||
self.diagnostics.add_event(
|
||||
"PAGE_RESPONSE_MATCHED",
|
||||
self.packet_count,
|
||||
kind=req.kind,
|
||||
page=req.page,
|
||||
record_count=slice_count,
|
||||
)
|
||||
|
||||
return bool(selected)
|
||||
|
||||
def _record_rejected_candidate(
|
||||
self,
|
||||
requests: list[PendingRequest],
|
||||
payload: bytes,
|
||||
reason: str,
|
||||
) -> None:
|
||||
for req in requests:
|
||||
req.response_candidates += 1
|
||||
req.response_candidate_lengths = (
|
||||
*req.response_candidate_lengths[-4:],
|
||||
len(payload),
|
||||
)
|
||||
self.diagnostics.counters["response_candidates_rejected"] += 1
|
||||
self.diagnostics.add_event(
|
||||
reason,
|
||||
self.packet_count,
|
||||
reason=True,
|
||||
kind=requests[0].kind if requests else None,
|
||||
page=requests[0].page if requests else None,
|
||||
payload_length=len(payload),
|
||||
matching_requests=len(requests),
|
||||
)
|
||||
|
||||
def diagnostic_report(self) -> dict[str, Any]:
|
||||
return self.diagnostics.report(self.pending)
|
||||
|
||||
def kinds_seen(self) -> list[str]:
|
||||
seen = []
|
||||
for pair in self.pairs:
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue