0.1.1 - Add libpcap/Npcap backend and robust live capture

Introduce a cross-platform libpcap/Npcap capture backend with a Windows raw-socket fallback and plumbing to select backends via --capture-backend. Add ctypes-based libpcap wrapper (open_libpcap_capture, CaptureStats, LibpcapUnavailable) and a backend factory; refactor live capture runner to use the new backend, StopKeyMonitor, and improved clipboard handling. Make session and protocol decoders resilient to pipelined/multi-page responses and non-byte-aligned streams (alignment iterator, alignment-aware decoding, embedded-record trimming, slice support). Update mitmproxy flows adapter to reuse LiveHistorySession. Expand console messages and docs/README to explain cross-platform requirements, usage, and capture diagnostics. Minor: update package description in pyproject and adjust ARC/response handling and row bookkeeping to record row indices.
This commit is contained in:
Golumpa 2026-06-12 12:58:15 +01:00
parent 1144e33a57
commit 0e78906e3a
20 changed files with 1234 additions and 147 deletions

View file

@ -0,0 +1,74 @@
from __future__ import annotations
import socket
import sys
from dataclasses import dataclass
from typing import Iterator, Protocol
from nte_history_exporter.live_capture.libpcap import (
CaptureStats,
LibpcapUnavailable,
open_libpcap_capture,
)
from nte_history_exporter.live_capture.windows_raw import (
ParsedIpUdpPacket,
open_raw_udp_socket,
read_packets,
)
class CaptureBackend(Protocol):
name: str
detail: str
fallback_reason: str
def packets(self) -> Iterator[ParsedIpUdpPacket | None]: ...
def stats(self) -> CaptureStats | None: ...
def close(self) -> None: ...
@dataclass
class RawSocketCapture:
local_ip: str
fallback_reason: str = ""
def __post_init__(self) -> None:
if sys.platform != "win32":
raise RuntimeError("the raw socket capture backend is only available on Windows")
self.name = "windows_raw"
self.detail = self.local_ip
self.socket = open_raw_udp_socket(self.local_ip)
def packets(self) -> Iterator[ParsedIpUdpPacket | None]:
return read_packets(self.socket)
def stats(self) -> CaptureStats | None:
return None
def close(self) -> None:
try:
self.socket.ioctl(socket.SIO_RCVALL, socket.RCVALL_OFF)
except Exception:
pass
self.socket.close()
def open_capture_backend(local_ip: str, requested: str = "auto") -> CaptureBackend:
if requested not in {"auto", "libpcap", "raw"}:
raise ValueError(f"unknown capture backend: {requested}")
if requested in {"auto", "libpcap"}:
try:
capture = open_libpcap_capture(local_ip)
except LibpcapUnavailable as exc:
if requested != "auto" or sys.platform != "win32":
raise
return RawSocketCapture(local_ip, fallback_reason=str(exc))
capture.name = "npcap" if sys.platform == "win32" else "libpcap"
capture.detail = capture.device
capture.fallback_reason = ""
return capture
return RawSocketCapture(local_ip)

View file

@ -0,0 +1,404 @@
from __future__ import annotations
import ctypes
import ctypes.util
import os
import socket
import struct
import sys
from dataclasses import dataclass
from pathlib import Path
from typing import Iterator
from nte_history_exporter.live_capture.windows_raw import ParsedIpUdpPacket, parse_ipv4_udp_packet
PCAP_ERRBUF_SIZE = 256
SNAP_LENGTH = 65535
CAPTURE_BUFFER_SIZE = 16 * 1024 * 1024
READ_TIMEOUT_MS = 100
DLT_NULL = 0
DLT_EN10MB = 1
DLT_RAW_BSD = 12
DLT_RAW = 101
DLT_LOOP = 108
DLT_LINUX_SLL = 113
DLT_LINUX_SLL2 = 276
class LibpcapUnavailable(RuntimeError):
pass
class _SockAddr(ctypes.Structure):
_fields_ = [("sa_family", ctypes.c_ushort), ("sa_data", ctypes.c_ubyte * 14)]
class _PcapAddr(ctypes.Structure):
pass
_PcapAddrPtr = ctypes.POINTER(_PcapAddr)
_PcapAddr._fields_ = [
("next", _PcapAddrPtr),
("addr", ctypes.POINTER(_SockAddr)),
("netmask", ctypes.POINTER(_SockAddr)),
("broadaddr", ctypes.POINTER(_SockAddr)),
("dstaddr", ctypes.POINTER(_SockAddr)),
]
class _PcapIf(ctypes.Structure):
pass
_PcapIfPtr = ctypes.POINTER(_PcapIf)
_PcapIf._fields_ = [
("next", _PcapIfPtr),
("name", ctypes.c_char_p),
("description", ctypes.c_char_p),
("addresses", _PcapAddrPtr),
("flags", ctypes.c_uint),
]
class _Timeval(ctypes.Structure):
_fields_ = [("tv_sec", ctypes.c_long), ("tv_usec", ctypes.c_long)]
class _PcapPacketHeader(ctypes.Structure):
_fields_ = [("ts", _Timeval), ("caplen", ctypes.c_uint), ("length", ctypes.c_uint)]
class _BpfProgram(ctypes.Structure):
_fields_ = [("bf_len", ctypes.c_uint), ("bf_insns", ctypes.c_void_p)]
class _PcapStat(ctypes.Structure):
# Npcap adds ps_capt after the three standard libpcap counters. Keeping the
# extra field is harmless on platforms that only populate the first three.
_fields_ = [
("ps_recv", ctypes.c_uint),
("ps_drop", ctypes.c_uint),
("ps_ifdrop", ctypes.c_uint),
("ps_capt", ctypes.c_uint),
]
@dataclass
class CaptureStats:
received: int
dropped: int
interface_dropped: int
def _windows_npcap_directory() -> Path:
return Path(os.environ.get("SystemRoot", r"C:\Windows")) / "System32" / "Npcap"
def _load_library() -> ctypes.CDLL:
candidates: list[str] = []
dll_directory = None
if sys.platform == "win32":
npcap_dir = _windows_npcap_directory()
if npcap_dir.is_dir() and hasattr(os, "add_dll_directory"):
dll_directory = os.add_dll_directory(str(npcap_dir))
candidates.extend([str(npcap_dir / "wpcap.dll"), "wpcap.dll"])
elif sys.platform == "darwin":
candidates.extend(
[
ctypes.util.find_library("pcap") or "",
"/usr/lib/libpcap.A.dylib",
"/usr/lib/libpcap.dylib",
]
)
else:
candidates.extend([ctypes.util.find_library("pcap") or "", "libpcap.so.1", "libpcap.so"])
errors = []
try:
for candidate in candidates:
if not candidate:
continue
try:
return ctypes.CDLL(candidate)
except OSError as exc:
errors.append(str(exc))
finally:
if dll_directory is not None:
dll_directory.close()
detail = f": {errors[-1]}" if errors else ""
if sys.platform == "win32":
raise LibpcapUnavailable(f"Npcap is not installed or could not be loaded{detail}")
raise LibpcapUnavailable(f"libpcap is not installed or could not be loaded{detail}")
def _configure_api(lib: ctypes.CDLL) -> None:
lib.pcap_findalldevs.argtypes = [ctypes.POINTER(_PcapIfPtr), ctypes.c_char_p]
lib.pcap_findalldevs.restype = ctypes.c_int
lib.pcap_freealldevs.argtypes = [_PcapIfPtr]
lib.pcap_create.argtypes = [ctypes.c_char_p, ctypes.c_char_p]
lib.pcap_create.restype = ctypes.c_void_p
lib.pcap_set_snaplen.argtypes = [ctypes.c_void_p, ctypes.c_int]
lib.pcap_set_promisc.argtypes = [ctypes.c_void_p, ctypes.c_int]
lib.pcap_set_timeout.argtypes = [ctypes.c_void_p, ctypes.c_int]
lib.pcap_set_buffer_size.argtypes = [ctypes.c_void_p, ctypes.c_int]
lib.pcap_activate.argtypes = [ctypes.c_void_p]
lib.pcap_activate.restype = ctypes.c_int
lib.pcap_close.argtypes = [ctypes.c_void_p]
lib.pcap_geterr.argtypes = [ctypes.c_void_p]
lib.pcap_geterr.restype = ctypes.c_char_p
lib.pcap_datalink.argtypes = [ctypes.c_void_p]
lib.pcap_datalink.restype = ctypes.c_int
lib.pcap_compile.argtypes = [
ctypes.c_void_p,
ctypes.POINTER(_BpfProgram),
ctypes.c_char_p,
ctypes.c_int,
ctypes.c_uint,
]
lib.pcap_setfilter.argtypes = [ctypes.c_void_p, ctypes.POINTER(_BpfProgram)]
lib.pcap_freecode.argtypes = [ctypes.POINTER(_BpfProgram)]
lib.pcap_next_ex.argtypes = [
ctypes.c_void_p,
ctypes.POINTER(ctypes.POINTER(_PcapPacketHeader)),
ctypes.POINTER(ctypes.POINTER(ctypes.c_ubyte)),
]
lib.pcap_next_ex.restype = ctypes.c_int
lib.pcap_stats.argtypes = [ctypes.c_void_p, ctypes.POINTER(_PcapStat)]
lib.pcap_stats.restype = ctypes.c_int
def _pcap_error(lib: ctypes.CDLL, handle: ctypes.c_void_p) -> str:
raw = lib.pcap_geterr(handle)
return raw.decode("utf-8", errors="replace") if raw else "unknown libpcap error"
def _ipv4_from_sockaddr(address: ctypes.POINTER(_SockAddr)) -> str | None:
if not address:
return None
raw = ctypes.string_at(address, 16)
family = raw[1] if sys.platform == "darwin" else struct.unpack_from("=H", raw, 0)[0]
return socket.inet_ntoa(raw[4:8]) if family == socket.AF_INET else None
def _find_windows_device(local_ip: str) -> str | None:
if sys.platform != "win32":
return None
from ctypes import wintypes
class IpAddressString(ctypes.Structure):
_fields_ = [("value", ctypes.c_char * 16)]
class IpMaskString(ctypes.Structure):
_fields_ = [("value", ctypes.c_char * 16)]
class IpAddrString(ctypes.Structure):
pass
IpAddrStringPtr = ctypes.POINTER(IpAddrString)
IpAddrString._fields_ = [
("next", IpAddrStringPtr),
("ip_address", IpAddressString),
("ip_mask", IpMaskString),
("context", wintypes.DWORD),
]
class IpAdapterInfo(ctypes.Structure):
pass
IpAdapterInfoPtr = ctypes.POINTER(IpAdapterInfo)
IpAdapterInfo._fields_ = [
("next", IpAdapterInfoPtr),
("combo_index", wintypes.DWORD),
("adapter_name", ctypes.c_char * 260),
("description", ctypes.c_char * 132),
("address_length", wintypes.UINT),
("address", ctypes.c_ubyte * 8),
("index", wintypes.DWORD),
("adapter_type", wintypes.UINT),
("dhcp_enabled", wintypes.UINT),
("current_ip_address", IpAddrStringPtr),
("ip_address_list", IpAddrString),
("gateway_list", IpAddrString),
("dhcp_server", IpAddrString),
("have_wins", wintypes.BOOL),
("primary_wins_server", IpAddrString),
("secondary_wins_server", IpAddrString),
("lease_obtained", ctypes.c_longlong),
("lease_expires", ctypes.c_longlong),
]
ip_helper = ctypes.WinDLL("iphlpapi")
ip_helper.GetAdaptersInfo.argtypes = [IpAdapterInfoPtr, ctypes.POINTER(wintypes.ULONG)]
ip_helper.GetAdaptersInfo.restype = wintypes.DWORD
size = wintypes.ULONG(0)
if ip_helper.GetAdaptersInfo(None, ctypes.byref(size)) not in (0, 111):
return None
buffer = ctypes.create_string_buffer(size.value)
adapter = ctypes.cast(buffer, IpAdapterInfoPtr)
if ip_helper.GetAdaptersInfo(adapter, ctypes.byref(size)) != 0:
return None
while adapter:
info = adapter.contents
address = ctypes.pointer(info.ip_address_list)
while address:
ip = address.contents.ip_address.value.decode("ascii", errors="ignore")
if ip == local_ip:
adapter_name = info.adapter_name.decode("ascii", errors="ignore")
return rf"\Device\NPF_{adapter_name}"
address = address.contents.next
adapter = info.next
return None
def _find_device(lib: ctypes.CDLL, local_ip: str) -> str:
windows_device = _find_windows_device(local_ip)
if windows_device:
return windows_device
devices = _PcapIfPtr()
errbuf = ctypes.create_string_buffer(PCAP_ERRBUF_SIZE)
if lib.pcap_findalldevs(ctypes.byref(devices), errbuf) != 0:
raise LibpcapUnavailable(errbuf.value.decode("utf-8", errors="replace"))
try:
device = devices
available = []
while device:
entry = device.contents
name = entry.name.decode("utf-8", errors="replace")
addresses = entry.addresses
while addresses:
ip = _ipv4_from_sockaddr(addresses.contents.addr)
if ip:
available.append((name, ip))
if ip == local_ip:
return name
addresses = addresses.contents.next
device = entry.next
finally:
lib.pcap_freealldevs(devices)
details = ", ".join(f"{name}={ip}" for name, ip in available) or "no IPv4 capture devices"
raise LibpcapUnavailable(f"no libpcap device owns local address {local_ip}; found {details}")
def _extract_ipv4_frame(frame: bytes, datalink: int) -> bytes | None:
if datalink == DLT_EN10MB:
if len(frame) < 14:
return None
offset = 14
ether_type = struct.unpack_from("!H", frame, 12)[0]
while ether_type in (0x8100, 0x88A8, 0x9100):
if len(frame) < offset + 4:
return None
ether_type = struct.unpack_from("!H", frame, offset + 2)[0]
offset += 4
return frame[offset:] if ether_type == 0x0800 else None
if datalink in (DLT_RAW_BSD, DLT_RAW):
return frame
if datalink in (DLT_NULL, DLT_LOOP):
if len(frame) < 4:
return None
family_native = struct.unpack_from("=I", frame, 0)[0]
family_network = struct.unpack_from("!I", frame, 0)[0]
return frame[4:] if socket.AF_INET in (family_native, family_network) else None
if datalink == DLT_LINUX_SLL:
if len(frame) < 16 or struct.unpack_from("!H", frame, 14)[0] != 0x0800:
return None
return frame[16:]
if datalink == DLT_LINUX_SLL2:
if len(frame) < 20 or struct.unpack_from("!H", frame, 0)[0] != 0x0800:
return None
return frame[20:]
raise LibpcapUnavailable(f"unsupported libpcap link-layer type {datalink}")
class LibpcapCapture:
def __init__(self, local_ip: str) -> None:
self.lib = _load_library()
_configure_api(self.lib)
self.device = _find_device(self.lib, local_ip)
self.handle = ctypes.c_void_p()
errbuf = ctypes.create_string_buffer(PCAP_ERRBUF_SIZE)
handle = self.lib.pcap_create(self.device.encode("utf-8"), errbuf)
if not handle:
raise LibpcapUnavailable(errbuf.value.decode("utf-8", errors="replace"))
self.handle = ctypes.c_void_p(handle)
try:
self.lib.pcap_set_snaplen(self.handle, SNAP_LENGTH)
self.lib.pcap_set_promisc(self.handle, 0)
self.lib.pcap_set_timeout(self.handle, READ_TIMEOUT_MS)
self.lib.pcap_set_buffer_size(self.handle, CAPTURE_BUFFER_SIZE)
activation = self.lib.pcap_activate(self.handle)
if activation < 0:
raise LibpcapUnavailable(_pcap_error(self.lib, self.handle))
program = _BpfProgram()
filter_expression = f"udp and host {local_ip}".encode("ascii")
if self.lib.pcap_compile(self.handle, ctypes.byref(program), filter_expression, 1, 0xFFFFFFFF) != 0:
raise LibpcapUnavailable(_pcap_error(self.lib, self.handle))
try:
if self.lib.pcap_setfilter(self.handle, ctypes.byref(program)) != 0:
raise LibpcapUnavailable(_pcap_error(self.lib, self.handle))
finally:
self.lib.pcap_freecode(ctypes.byref(program))
self.datalink = self.lib.pcap_datalink(self.handle)
except Exception:
self.close()
raise
def packets(self) -> Iterator[ParsedIpUdpPacket | None]:
header = ctypes.POINTER(_PcapPacketHeader)()
packet_data = ctypes.POINTER(ctypes.c_ubyte)()
while self.handle:
result = self.lib.pcap_next_ex(self.handle, ctypes.byref(header), ctypes.byref(packet_data))
if result == 0:
yield None
continue
if result == -2:
return
if result < 0:
raise RuntimeError(_pcap_error(self.lib, self.handle))
frame = ctypes.string_at(packet_data, header.contents.caplen)
ipv4_packet = _extract_ipv4_frame(frame, self.datalink)
if ipv4_packet is None:
continue
packet = parse_ipv4_udp_packet(ipv4_packet)
if packet is not None:
yield packet
def stats(self) -> CaptureStats | None:
if not self.handle:
return None
stats = _PcapStat()
if self.lib.pcap_stats(self.handle, ctypes.byref(stats)) != 0:
return None
return CaptureStats(stats.ps_recv, stats.ps_drop, stats.ps_ifdrop)
def close(self) -> None:
if self.handle:
self.lib.pcap_close(self.handle)
self.handle = ctypes.c_void_p()
def open_libpcap_capture(local_ip: str) -> LibpcapCapture:
return LibpcapCapture(local_ip)

View file

@ -1,8 +1,8 @@
from __future__ import annotations
import json
import msvcrt
import socket
import subprocess
import sys
import time
from datetime import datetime
from pathlib import Path
@ -12,8 +12,10 @@ from nte_history_exporter.constants import POOL_META
from nte_history_exporter.decoder.boundary import annotate_groups, select_continuous_run_from_page_1
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.session import LiveHistorySession, UdpPacket
from nte_history_exporter.live_capture.windows_raw import detect_local_ipv4, open_raw_udp_socket, read_packets
from nte_history_exporter.live_capture.stop_key import StopKeyMonitor
from nte_history_exporter.live_capture.windows_raw import detect_local_ipv4
EXPORT_PREFIXES = {
@ -23,7 +25,7 @@ EXPORT_PREFIXES = {
}
def copy_to_clipboard(text: str) -> None:
def copy_to_clipboard(text: str) -> bool:
try:
import tkinter as tk
@ -33,54 +35,99 @@ def copy_to_clipboard(text: str) -> None:
root.clipboard_append(text)
root.update()
root.destroy()
return
return True
except Exception:
pass
import subprocess
commands = []
if sys.platform == "win32":
commands.append(["clip"])
elif sys.platform == "darwin":
commands.append(["pbcopy"])
else:
commands.extend([["wl-copy"], ["xclip", "-selection", "clipboard"]])
subprocess.run(["clip"], input=text, text=True, check=True)
for command in commands:
try:
subprocess.run(command, input=text, text=True, check=True)
return True
except (FileNotFoundError, subprocess.CalledProcessError):
continue
return False
def run_live_capture(
*,
interface_ip: str | None = None,
capture_backend: str = "auto",
copy_clipboard: bool = True,
write_debug_csv: bool = False,
) -> dict:
local_ip = interface_ip or detect_local_ipv4()
session = LiveHistorySession(local_ip)
sock = open_raw_udp_socket(local_ip)
console.print_live_instructions(local_ip)
capture = open_capture_backend(local_ip, capture_backend)
console.print_live_instructions(local_ip, capture.name, capture.detail)
if capture.fallback_reason:
console.print_capture_fallback(capture.fallback_reason)
reported_missing_pages: dict[str, tuple[int, ...]] = {}
try:
for packet in read_packets(sock):
if msvcrt.kbhit():
msvcrt.getch()
break
if packet is None:
continue
matched = session.process_packet(
UdpPacket(
timestamp=time.time(),
src_ip=packet.src_ip,
dst_ip=packet.dst_ip,
src_port=packet.src_port,
dst_port=packet.dst_port,
payload=packet.payload,
with StopKeyMonitor() as stop_key:
for packet in capture.packets():
if stop_key.pressed():
break
if packet is None:
continue
pair_count_before = len(session.pairs)
matched = session.process_packet(
UdpPacket(
timestamp=time.time(),
src_ip=packet.src_ip,
dst_ip=packet.dst_ip,
src_port=packet.src_port,
dst_port=packet.dst_port,
payload=packet.payload,
)
)
)
if matched:
kind = session.pairs[-1][7] if session.pairs else "permanent"
label = POOL_META.get(kind, POOL_META["permanent"])["name"]
console.print_page_captured(label, session.last_page_seen)
if matched:
affected_kinds = []
for pair in session.pairs[pair_count_before:]:
page = pair[0]
kind = pair[7] if len(pair) > 7 else "permanent"
label = POOL_META.get(kind, POOL_META["permanent"])["name"]
was_replacement = any(
existing[0] == page
and (existing[7] if len(existing) > 7 else "permanent") == kind
for existing in session.pairs[:pair_count_before]
)
console.print_page_captured(label, page, recaptured=was_replacement)
if kind not in affected_kinds:
affected_kinds.append(kind)
for kind in affected_kinds:
label = POOL_META.get(kind, POOL_META["permanent"])["name"]
missing_pages = tuple(session.missing_pages(kind))
previously_missing = reported_missing_pages.get(kind, ())
if missing_pages and missing_pages != previously_missing:
reasons = {
page: session.missing_page_reason(kind, page)
for page in missing_pages
}
console.print_missing_pages(label, list(missing_pages), reasons)
elif previously_missing and not missing_pages:
console.print_page_gap_recovered(label)
reported_missing_pages[kind] = missing_pages
finally:
try:
sock.ioctl(socket.SIO_RCVALL, socket.RCVALL_OFF)
except Exception:
pass
sock.close()
stats = capture.stats()
capture.close()
if stats:
console.print_capture_stats(
stats.received,
stats.dropped,
stats.interface_dropped,
)
exports = []
for kind in session.kinds_seen():
@ -138,8 +185,10 @@ def run_live_capture(
console.print_note(f"Export written: {item['json_path']}")
if copy_clipboard and len(exports) == 1:
copy_to_clipboard(exports[0]["payload"])
console.print_success("Export copied to clipboard - paste it straight into your tracker.")
if copy_to_clipboard(exports[0]["payload"]):
console.print_success("Export copied to clipboard - paste it straight into your tracker.")
else:
console.print_note("Clipboard tool unavailable; use the JSON file shown above.")
elif len(exports) > 1:
console.print_note("Multiple banners captured; clipboard copy skipped so one export")
console.print_note("does not overwrite another.")

View file

@ -13,10 +13,10 @@ from nte_history_exporter.decoder.arc import (
)
from nte_history_exporter.decoder.boundary import select_continuous_run_from_page_1
from nte_history_exporter.decoder.protocol import (
decode_response_records,
history_request_kind,
is_history_request,
request_page,
response_contains_history_marker,
)
from nte_history_exporter.decoder.run import build_rows_from_pairs
@ -42,6 +42,8 @@ class PendingRequest:
dst_ip: str
src_port: int
dst_port: int
response_candidates: int = 0
response_candidate_lengths: tuple[int, ...] = ()
class LiveHistorySession:
@ -52,6 +54,42 @@ class LiveHistorySession:
self.packet_count = 0
self.last_match_time: float | None = None
self.last_page_seen: int | None = None
self.last_capture_was_replacement = False
self.requested_pages: dict[str, set[int]] = {}
self.unanswered_pages: dict[str, dict[int, str]] = {}
def _mark_unanswered(self, request: PendingRequest) -> None:
if request.response_candidates:
lengths = ", ".join(str(length) for length in request.response_candidate_lengths)
reason = (
f"{request.response_candidates} matching inbound UDP packet(s) captured "
f"but not recognized as history response (lengths: {lengths})"
)
else:
reason = "request captured; no matching response page was captured"
self.unanswered_pages.setdefault(request.kind, {})[request.page] = reason
def _queue_request(self, request: PendingRequest) -> None:
# The game may pipeline several page requests before responses arrive,
# and one response can contain multiple five-record pages. Keep that
# queue intact; only replace an earlier duplicate request for the same
# page, such as a recovery pass after reopening the history board.
retained = deque()
for pending in self.pending:
same_stream = (
pending.src_ip == request.src_ip
and pending.dst_ip == request.dst_ip
and pending.src_port == request.src_port
and pending.dst_port == request.dst_port
and pending.kind == request.kind
)
if same_stream and (request.page == 1 or pending.page == request.page):
self._mark_unanswered(pending)
else:
retained.append(pending)
self.pending = retained
self.requested_pages.setdefault(request.kind, set()).add(request.page)
self.pending.append(request)
def process_packet(self, packet: UdpPacket) -> bool:
self.packet_count += 1
@ -69,7 +107,7 @@ class LiveHistorySession:
src_port=packet.src_port,
dst_port=packet.dst_port,
)
self.pending.append(req)
self._queue_request(req)
self.last_page_seen = req.page
return False
@ -86,47 +124,82 @@ class LiveHistorySession:
src_port=packet.src_port,
dst_port=packet.dst_port,
)
self.pending.append(req)
self._queue_request(req)
self.last_page_seen = req.page
return False
if packet.dst_ip != self.local_ip or len(packet.payload) < 100:
return False
is_monopoly_response = response_contains_history_marker(packet.payload)
is_arc_response = bool(parse_arc_response(packet.payload))
if not is_monopoly_response and not is_arc_response:
return False
for req in self.pending:
if req.kind == "arc_miracle_box" and not is_arc_response:
continue
if req.kind != "arc_miracle_box" and not is_monopoly_response:
continue
connection_candidates = [
req
for req in self.pending
if (
packet.src_ip == req.dst_ip
and packet.dst_ip == req.src_ip
and packet.src_port == req.dst_port
and packet.dst_port == req.src_port
):
self.pending.remove(req)
self.pairs.append(
(
req.page,
req.offset,
req.request_msg,
req.request_time,
self.packet_count,
packet.timestamp,
packet.payload,
req.kind,
)
)
self.last_match_time = packet.timestamp
self.last_page_seen = req.page
return True
)
]
if not connection_candidates:
return False
return False
monopoly_records = decode_response_records(packet.payload)
arc_records = parse_arc_response(packet.payload) if not monopoly_records else []
if monopoly_records:
candidates = [req for req in connection_candidates if req.kind != "arc_miracle_box"]
records = monopoly_records
elif arc_records:
candidates = [req for req in connection_candidates if req.kind == "arc_miracle_box"]
records = arc_records
else:
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),
)
return False
page_count = max(1, (len(records) + 4) // 5)
if len(records) < 5:
selected = [candidates[-1]]
else:
selected = candidates[:page_count]
for page_slice, req in enumerate(selected):
self.pending.remove(req)
self.unanswered_pages.setdefault(req.kind, {}).pop(req.page, None)
slice_start = page_slice * 5
slice_count = min(5, len(records) - slice_start)
if slice_count <= 0:
break
self.last_capture_was_replacement = any(
pair[0] == req.page and (pair[7] if len(pair) > 7 else "permanent") == req.kind
for pair in self.pairs
)
self.pairs.append(
(
req.page,
req.offset,
req.request_msg,
req.request_time,
self.packet_count,
packet.timestamp,
packet.payload,
req.kind,
slice_start,
slice_count,
)
)
self.last_match_time = packet.timestamp
self.last_page_seen = req.page
return bool(selected)
def kinds_seen(self) -> list[str]:
seen = []
@ -139,6 +212,30 @@ class LiveHistorySession:
def pairs_for_kind(self, kind: str) -> list[tuple]:
return [pair for pair in self.pairs if (pair[7] if len(pair) > 7 else "permanent") == kind]
def missing_pages(self, kind: str) -> list[int]:
pages = {pair[0] for pair in self.pairs_for_kind(kind)}
if not pages:
return []
return [page for page in range(1, max(pages) + 1) if page not in pages]
def missing_page_reason(self, kind: str, page: int) -> str:
unanswered = self.unanswered_pages.get(kind, {})
if page in unanswered:
return unanswered[page]
pending = next(
(request for request in self.pending if request.kind == kind and request.page == page),
None,
)
if pending and pending.response_candidates:
lengths = ", ".join(str(length) for length in pending.response_candidate_lengths)
return (
f"{pending.response_candidates} matching inbound UDP packet(s) captured "
f"but not recognized as history response (lengths: {lengths})"
)
if page in self.requested_pages.get(kind, set()):
return "request captured; no matching response page was captured"
return "request was not captured"
def build_rows(self, kind: str | None = None) -> list[dict[str, Any]]:
if kind == "arc_miracle_box":
best_run, _warnings = select_continuous_arc_run(self.pairs_for_kind(kind))

View file

@ -0,0 +1,49 @@
from __future__ import annotations
import os
import select
import sys
class StopKeyMonitor:
def __init__(self) -> None:
self._windows = os.name == "nt"
self._fd: int | None = None
self._original_terminal = None
def __enter__(self) -> StopKeyMonitor:
if self._windows or not sys.stdin.isatty():
return self
import termios
import tty
self._fd = sys.stdin.fileno()
self._original_terminal = termios.tcgetattr(self._fd)
tty.setcbreak(self._fd)
return self
def pressed(self) -> bool:
if self._windows:
import msvcrt
if not msvcrt.kbhit():
return False
msvcrt.getch()
return True
if self._fd is None:
return False
readable, _, _ = select.select([self._fd], [], [], 0)
if not readable:
return False
os.read(self._fd, 1)
return True
def __exit__(self, _exc_type, _exc, _tb) -> None:
if self._fd is None or self._original_terminal is None:
return
import termios
termios.tcsetattr(self._fd, termios.TCSADRAIN, self._original_terminal)

View file

@ -6,6 +6,8 @@ from dataclasses import dataclass
from nte_history_exporter.live_capture.session import UdpPacket
RECEIVE_BUFFER_SIZE = 4 * 1024 * 1024
@dataclass
class ParsedIpUdpPacket:
@ -63,6 +65,7 @@ def parse_ipv4_udp_packet(data: bytes) -> ParsedIpUdpPacket | None:
def open_raw_udp_socket(local_ip: str) -> socket.socket:
sock = socket.socket(socket.AF_INET, socket.SOCK_RAW, socket.IPPROTO_IP)
sock.bind((local_ip, 0))
sock.setsockopt(socket.SOL_SOCKET, socket.SO_RCVBUF, RECEIVE_BUFFER_SIZE)
sock.setsockopt(socket.IPPROTO_IP, socket.IP_HDRINCL, 1)
sock.ioctl(socket.SIO_RCVALL, socket.RCVALL_ON)
sock.settimeout(0.5)