diff --git a/README.md b/README.md index 2ffd974..a30c397 100644 --- a/README.md +++ b/README.md @@ -144,7 +144,7 @@ Decodes a `mitmproxy .flows` capture instead of listening live — used for rese | Flag | Effect | | --------- | --------------------------------------------------- | | `--live` | Capture live UDP traffic instead of reading a file. | -| `--debug` | Also write the full research CSV next to each JSON. | +| `--debug` | Also write the research CSV and privacy-safe capture diagnostics. | | `--user-uid ` | Override the auto-detected NTE user UID in the JSON export. | | `--copy-clipboard` | Copy a single live export JSON to clipboard after saving. | @@ -156,7 +156,10 @@ Advanced live-capture selection: --capture-backend raw Require the Windows raw-socket backend ``` -The `--debug` CSV holds any extra information that might be needed for fixing bugs. It contains no dangerous personal account data — only the raw bytes of the captured history page. +The `--debug` CSV holds decoded research fields, including raw captured history +records. A separate versioned `*.diagnostics.json` sidecar provides shareable +reason codes and counts without payloads, network addresses, ports, packet +timestamps, or user UID values. See [Capture diagnostics](docs/capture-diagnostics.md). The exporter automatically includes the shareable NTE user UID when it appears in the capture. If a short capture does not include it, the console asks before saving; you can also pass it explicitly with `--user-uid`. diff --git a/docs/capture-diagnostics.md b/docs/capture-diagnostics.md new file mode 100644 index 0000000..9f92b0d --- /dev/null +++ b/docs/capture-diagnostics.md @@ -0,0 +1,27 @@ +# Capture diagnostics + +Running the exporter with `--debug` writes a UID-free `Capture_*.diagnostics.json` +file in the export directory. This sidecar explains what the capture +pipeline recognized, rejected, and paired without changing the public export +format or record order. It is also written when no history page could be +exported. + +The report contains bounded events, aggregate counters, reason codes, packet +positions, payload lengths, history kinds, page numbers, and record counts. It +does not contain packet payloads, IP addresses, ports, packet timestamps, or a +user UID. The report format is versioned independently as +`nte-capture-diagnostics` version 1. + +Useful rejection codes include: + +- `RESPONSE_TOO_SHORT`: a matching inbound packet could not contain a history + response; +- `NO_HISTORY_MARKER`: a matching response candidate had no known history + marker; +- `HISTORY_MARKER_PARSE_FAILED`: a known marker was present but neither decoder + produced records; +- `RESPONSE_KIND_MISMATCH`: decoded data did not match the pending history kind; +- `REQUEST_REPLACED`: a recovery request superseded an unanswered request. + +Events are capped at 200 per session. Aggregate counters continue after that +limit and `events_omitted` records how many event entries were left out. diff --git a/src/nte_history_exporter/adapters/mitmproxy_flows.py b/src/nte_history_exporter/adapters/mitmproxy_flows.py index 19c35f7..fc2b930 100644 --- a/src/nte_history_exporter/adapters/mitmproxy_flows.py +++ b/src/nte_history_exporter/adapters/mitmproxy_flows.py @@ -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(), } diff --git a/src/nte_history_exporter/cli.py b/src/nte_history_exporter/cli.py index 09b7368..5b0090e 100644 --- a/src/nte_history_exporter/cli.py +++ b/src/nte_history_exporter/cli.py @@ -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 diff --git a/src/nte_history_exporter/live_capture/diagnostics.py b/src/nte_history_exporter/live_capture/diagnostics.py new file mode 100644 index 0000000..6f9ac6c --- /dev/null +++ b/src/nte_history_exporter/live_capture/diagnostics.py @@ -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") diff --git a/src/nte_history_exporter/live_capture/runner.py b/src/nte_history_exporter/live_capture/runner.py index 3f41f50..6c16d79 100644 --- a/src/nte_history_exporter/live_capture/runner.py +++ b/src/nte_history_exporter/live_capture/runner.py @@ -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]: diff --git a/src/nte_history_exporter/live_capture/session.py b/src/nte_history_exporter/live_capture/session.py index e793294..171e782 100644 --- a/src/nte_history_exporter/live_capture/session.py +++ b/src/nte_history_exporter/live_capture/session.py @@ -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: diff --git a/tests/fixtures/README.md b/tests/fixtures/README.md index abdedc9..2bec598 100644 --- a/tests/fixtures/README.md +++ b/tests/fixtures/README.md @@ -22,6 +22,12 @@ Do not replace this file with a real `.pcap`, `.flows`, or exported account history. Add new cases by constructing the smallest relevant payload, replacing all timestamps and endpoints, and extending the privacy assertions. +`synthetic_capture_diagnostics.json` is a smaller replay transcript for failure +paths. It deliberately contains a short response and a marker-free response so +the debug sidecar's reason codes and privacy contract can be tested without a +real capture. Repeated-byte payloads use `payload_byte` plus `payload_length` to +keep the fixture readable. + ## Synthetic NTE_Assets fixture `nte_assets/` mirrors only the six table paths and English localization file diff --git a/tests/fixtures/synthetic_capture_diagnostics.json b/tests/fixtures/synthetic_capture_diagnostics.json new file mode 100644 index 0000000..3766e4d --- /dev/null +++ b/tests/fixtures/synthetic_capture_diagnostics.json @@ -0,0 +1,60 @@ +{ + "schema_version": 1, + "description": "Minimal synthetic replay for capture diagnostic reason codes.", + "privacy": { + "synthetic": true, + "contains_user_uid": false, + "contains_raw_account_session": false + }, + "local_ip": "192.0.2.10", + "expected": { + "packets_seen": 3, + "history_requests_recognized": 1, + "response_candidates_rejected": 2, + "event_counts": { + "HISTORY_REQUEST_RECOGNIZED": 1, + "NO_HISTORY_MARKER": 1, + "RESPONSE_TOO_SHORT": 1 + }, + "reason_counts": { + "NO_HISTORY_MARKER": 1, + "RESPONSE_TOO_SHORT": 1 + }, + "pending_response_candidates": 2, + "pending_candidate_lengths": [80, 220] + }, + "packets": [ + { + "label": "permanent-page-1-request", + "timestamp": 1893456000.0, + "src_ip": "192.0.2.10", + "dst_ip": "198.51.100.20", + "src_port": 50000, + "dst_port": 40000, + "protocol": "udp", + "payload_hex": "00000000000000000000000000000000000000000000000000000000000000040000007c100000000400000000" + }, + { + "label": "matching-response-too-short", + "timestamp": 1893456000.1, + "src_ip": "198.51.100.20", + "dst_ip": "192.0.2.10", + "src_port": 40000, + "dst_port": 50000, + "protocol": "udp", + "payload_byte": "00", + "payload_length": 80 + }, + { + "label": "matching-response-without-history-marker", + "timestamp": 1893456000.2, + "src_ip": "198.51.100.20", + "dst_ip": "192.0.2.10", + "src_port": 40000, + "dst_port": 50000, + "protocol": "udp", + "payload_byte": "00", + "payload_length": 220 + } + ] +} diff --git a/tests/support.py b/tests/support.py index d64467d..58a9a55 100644 --- a/tests/support.py +++ b/tests/support.py @@ -60,6 +60,34 @@ def load_network_fixture(): return json.load(f) +def load_capture_diagnostics_fixture(): + path = FIXTURES / "synthetic_capture_diagnostics.json" + with path.open(encoding="utf-8") as f: + return json.load(f) + + +def diagnostic_fixture_session(): + fixture = load_capture_diagnostics_fixture() + session = LiveHistorySession(fixture["local_ip"]) + for packet in fixture["packets"]: + if "payload_hex" in packet: + payload = bytes.fromhex(packet["payload_hex"]) + else: + payload = bytes.fromhex(packet["payload_byte"]) * packet["payload_length"] + session.process_packet( + UdpPacket( + timestamp=packet["timestamp"], + src_ip=packet["src_ip"], + dst_ip=packet["dst_ip"], + src_port=packet["src_port"], + dst_port=packet["dst_port"], + payload=payload, + protocol=packet["protocol"], + ) + ) + return session + + def fixture_packets(scenario=None): fixture = load_network_fixture() packets = fixture["packets"] diff --git a/tests/test_capture_diagnostics.py b/tests/test_capture_diagnostics.py new file mode 100644 index 0000000..9a2fbdd --- /dev/null +++ b/tests/test_capture_diagnostics.py @@ -0,0 +1,108 @@ +from tests.support import * # noqa: F401,F403 + +from nte_history_exporter.live_capture.diagnostics import ( + new_diagnostics_path, + write_capture_diagnostics, +) + + +class CaptureDiagnosticsTests(unittest.TestCase): + def test_successful_network_replay_reports_capture_pipeline_counts(self): + report = fixture_session().diagnostic_report() + + self.assertEqual( + report["counters"], + { + "history_requests_recognized": 10, + "history_responses_decoded": 10, + "packets_seen": 20, + "pages_matched": 10, + "udp_packets_seen": 20, + }, + ) + self.assertEqual(report["reason_counts"], {}) + self.assertEqual(report["pending_requests"], []) + + def test_synthetic_replay_reports_actionable_rejection_reasons(self): + fixture = load_capture_diagnostics_fixture() + report = diagnostic_fixture_session().diagnostic_report() + expected = fixture["expected"] + + self.assertEqual(report["format"], "nte-capture-diagnostics") + self.assertEqual(report["format_version"], 1) + self.assertEqual(report["counters"]["packets_seen"], expected["packets_seen"]) + self.assertEqual( + report["counters"]["history_requests_recognized"], + expected["history_requests_recognized"], + ) + self.assertEqual( + report["counters"]["response_candidates_rejected"], + expected["response_candidates_rejected"], + ) + self.assertEqual(report["event_counts"], expected["event_counts"]) + self.assertEqual(report["reason_counts"], expected["reason_counts"]) + self.assertEqual(report["pending_requests"][0]["response_candidates"], 2) + self.assertEqual( + report["pending_requests"][0]["response_candidate_lengths"], + expected["pending_candidate_lengths"], + ) + + def test_diagnostic_fixture_contains_only_synthetic_network_identity(self): + fixture = load_capture_diagnostics_fixture() + self.assertTrue(fixture["privacy"]["synthetic"]) + self.assertFalse(fixture["privacy"]["contains_user_uid"]) + self.assertFalse(fixture["privacy"]["contains_raw_account_session"]) + + allowed_ips = {"192.0.2.10", "198.51.100.20"} + for packet in fixture["packets"]: + with self.subTest(label=packet["label"]): + self.assertIn(packet["src_ip"], allowed_ips) + self.assertIn(packet["dst_ip"], allowed_ips) + if "payload_hex" in packet: + payload = bytes.fromhex(packet["payload_hex"]) + else: + payload = bytes.fromhex(packet["payload_byte"]) * packet["payload_length"] + self.assertIsNone(extract_user_uid(payload)) + + def test_diagnostic_report_excludes_capture_identity_and_payload_data(self): + report = diagnostic_fixture_session().diagnostic_report() + serialized = json.dumps(report) + forbidden_keys = { + "src_ip", + "dst_ip", + "src_port", + "dst_port", + "timestamp", + "payload", + "payload_hex", + "user_uid", + } + + def assert_safe(value): + if isinstance(value, dict): + self.assertTrue(forbidden_keys.isdisjoint(value)) + for child in value.values(): + assert_safe(child) + elif isinstance(value, list): + for child in value: + assert_safe(child) + + assert_safe(report) + self.assertNotIn("192.0.2.10", serialized) + self.assertNotIn("198.51.100.20", serialized) + self.assertNotIn("1893456000", serialized) + + def test_diagnostics_are_debug_sidecar_not_public_export_data(self): + session = diagnostic_fixture_session() + export = build_export_json([], []) + self.assertNotIn("diagnostics", export) + self.assertNotIn("capture_diagnostics", export) + + with TemporaryDirectory() as temp_dir: + diagnostics_path = new_diagnostics_path(temp_dir) + write_capture_diagnostics(diagnostics_path, session.diagnostic_report()) + written = json.loads(diagnostics_path.read_text(encoding="utf-8")) + + self.assertEqual(written, session.diagnostic_report()) + self.assertRegex(diagnostics_path.name, r"^Capture_\d{8}_\d{6}\.diagnostics\.json$") + self.assertNotIn("user", diagnostics_path.name.casefold())