diff --git a/docs/limitations.md b/docs/limitations.md index 01c4a2b..beaff11 100644 --- a/docs/limitations.md +++ b/docs/limitations.md @@ -13,6 +13,9 @@ Known limitations: - Structured Monopoly and Arc blocks are parsed in enrichment/fallback mode. The existing decoder remains authoritative when structured rows do not agree on record count, reward ID, and timestamp. +- Structured snapshot/segment assembly is limited to runs where every decoded + row came from structured fallback. Existing primary-decoder runs are never + reordered. Ambiguous generations retain the last proven snapshot. - Live capture prefers Npcap on Windows and automatically falls back to the built-in raw-socket backend if Npcap is unavailable. Linux and macOS require the system libpcap runtime. - The file adapter reads mitmproxy `.flows` captures for research and testing. - Npcap is Windows-only and is not redistributed with this project; Linux and macOS use their system libpcap. diff --git a/docs/packet-format.md b/docs/packet-format.md index 50d7a74..ff5df85 100644 --- a/docs/packet-format.md +++ b/docs/packet-format.md @@ -21,8 +21,29 @@ The structured parser is deliberately compatibility-gated: - Malformed, incomplete, mismatched, or ambiguous structured data is ignored; it cannot overwrite a successfully decoded primary row. +Structured protocol envelopes identify a history stream, page, query side, and +segment index. For an all-structured fallback run, the snapshot assembler: + +- orders segments by their protocol index while retaining row order inside + every segment; +- ignores exact retransmissions; +- starts a new generation when an existing segment index changes; +- replaces an older snapshot only when the new generation covers at least the + same segment range; +- merges a partial generation only when its suffix has one unique overlap with + the proven snapshot; and +- retains the proven snapshot and records an assembly warning when a merge is + ambiguous. + +Assembly never runs on a history run containing a successfully decoded primary +row. Such runs retain their existing packet/page order exactly. Timestamp-group +ordinals and UIDs are calculated only after any fallback assembly, using the +same inputs and algorithms as before. + `decoder_mode`, `structured_protocol_view`, `structured_pool_id`, -`secondary_reward_id`, and `secondary_quantity` are research/debug CSV fields. +`structured_generation_index`, `structured_assembly`, +`structured_assembly_warning_count`, `secondary_reward_id`, and +`secondary_quantity` are research/debug CSV fields. They are intentionally omitted from the public JSON export, whose format stays at version 1. diff --git a/src/nte_history_exporter/decoder/arc.py b/src/nte_history_exporter/decoder/arc.py index 9f1bdc7..b42b1e8 100644 --- a/src/nte_history_exporter/decoder/arc.py +++ b/src/nte_history_exporter/decoder/arc.py @@ -22,7 +22,12 @@ from nte_history_exporter.constants import ( from nte_history_exporter.decoder.boundary import select_continuous_run_from_page_1 from nte_history_exporter.decoder.run import fmt_packet_time from nte_history_exporter.mappings import ARC_META -from nte_history_exporter.decoder.structured_protocol import StructuredRecord, parse_structured_records +from nte_history_exporter.decoder.structured_protocol import ( + StructuredProtocolAssembler, + StructuredRecord, + parse_structured_blocks, + parse_structured_records, +) def is_arc_history_request(content: bytes) -> bool: @@ -125,7 +130,7 @@ def _arc_metadata(arc_id: str) -> dict[str, Any]: return next((meta for item_id, meta in ARC_META.items() if item_id.casefold() == folded), {}) -def _structured_arc_rows(structured_rows: list[StructuredRecord]) -> list[dict[str, Any]]: +def structured_arc_rows(structured_rows: list[StructuredRecord]) -> list[dict[str, Any]]: records = [] for structured in structured_rows: meta = _arc_metadata(structured.item_id) @@ -149,6 +154,7 @@ def _structured_arc_rows(structured_rows: list[StructuredRecord]) -> list[dict[s "decoder_mode": "structured_fallback", "structured_pool_id": structured.pool_id, "structured_protocol_view": structured.protocol_view, + "structured_generation_index": structured.generation_index, } ) return records @@ -176,10 +182,10 @@ def parse_arc_response(response: bytes) -> list[dict[str, Any]]: legacy_rows = _parse_legacy_arc_response(response) if legacy_rows: return _enrich_legacy_arc_rows(legacy_rows, structured_rows) - return _structured_arc_rows(structured_rows) + return structured_arc_rows(structured_rows) -def build_arc_rows_from_pairs(pairs: list[tuple]) -> list[dict[str, Any]]: +def _build_primary_arc_rows_from_pairs(pairs: list[tuple]) -> list[dict[str, Any]]: pool = POOL_META["arc_miracle_box"] rows: list[dict[str, Any]] = [] for pair in pairs: @@ -209,6 +215,48 @@ def build_arc_rows_from_pairs(pairs: list[tuple]) -> list[dict[str, Any]]: return rows +def build_arc_rows_from_pairs(pairs: list[tuple]) -> list[dict[str, Any]]: + primary_rows = _build_primary_arc_rows_from_pairs(pairs) + if primary_rows and any(row.get("decoder_mode") != "structured_fallback" for row in primary_rows): + return primary_rows + + assembler = StructuredProtocolAssembler() + for source_index, pair in enumerate(pairs): + assembler.add_blocks(parse_structured_blocks(pair[6], "fork", source_index=source_index)) + assembled = assembler.rows("fork") + if not assembled: + return primary_rows + + pool = POOL_META["arc_miracle_box"] + rows = [] + for row_index, (record, structured) in enumerate( + zip(structured_arc_rows(assembled), assembled), start=1 + ): + source_index = structured.source_index or 0 + pair = pairs[source_index] + page, offset, req_i, req_ts, resp_i, resp_ts, response = pair[:7] + rows.append( + { + **record, + "page": page, + "offset": offset, + "row": row_index, + "pool_group_id": pool["id"], + "pool_group_name": pool["name"], + "request_msg": req_i, + "request_time_utc": fmt_packet_time(req_ts), + "response_msg": resp_i, + "response_time_utc": fmt_packet_time(resp_ts), + "response_len": len(response), + "record_count": len(assembled), + "structured_assembly": "snapshot_segments", + "structured_assembly_warning_count": len(assembler.warnings), + } + ) + annotate_arc_groups(rows) + return rows + + def make_arc_uid(timestamp_raw: str, ordinal: int) -> str: source = "|".join([GAME_UID_PART, ARC_SYSTEM, ARC_BANNER_ID, timestamp_raw, str(ordinal)]) return hashlib.sha256(source.encode("utf-8")).hexdigest()[:32] diff --git a/src/nte_history_exporter/decoder/protocol.py b/src/nte_history_exporter/decoder/protocol.py index c890257..3fed6c7 100644 --- a/src/nte_history_exporter/decoder/protocol.py +++ b/src/nte_history_exporter/decoder/protocol.py @@ -368,7 +368,7 @@ def _reward_metadata(reward_id: str) -> dict[str, Any]: return next((meta for item_id, meta in REWARDS_BY_ID.items() if item_id.casefold() == folded), {}) -def _structured_monopoly_rows(structured_rows: list[StructuredRecord]) -> list[dict[str, Any]]: +def structured_monopoly_rows(structured_rows: list[StructuredRecord]) -> list[dict[str, Any]]: rows = [] for row_index, structured in enumerate(structured_rows, start=1): dice, dice_raw, result_type, result_source = _structured_result(structured.roll_points_raw) @@ -405,6 +405,7 @@ def _structured_monopoly_rows(structured_rows: list[StructuredRecord]) -> list[d "secondary_reward_id": structured.secondary_item_id, "secondary_quantity": structured.secondary_count, "structured_protocol_view": structured.protocol_view, + "structured_generation_index": structured.generation_index, } ) return rows @@ -425,4 +426,4 @@ def decode_response_records(response_content: bytes) -> list[dict[str, Any]]: break if heuristic_rows: return _enrich_heuristic_rows(heuristic_rows, structured_rows) - return _structured_monopoly_rows(structured_rows) + return structured_monopoly_rows(structured_rows) diff --git a/src/nte_history_exporter/decoder/run.py b/src/nte_history_exporter/decoder/run.py index c756755..1b4a646 100644 --- a/src/nte_history_exporter/decoder/run.py +++ b/src/nte_history_exporter/decoder/run.py @@ -4,7 +4,11 @@ from datetime import datetime, timezone from typing import Any from nte_history_exporter.constants import POOL_META -from nte_history_exporter.decoder.protocol import decode_response_records +from nte_history_exporter.decoder.protocol import decode_response_records, structured_monopoly_rows +from nte_history_exporter.decoder.structured_protocol import ( + StructuredProtocolAssembler, + parse_structured_blocks, +) def fmt_packet_time(ts: float | None) -> str: @@ -13,7 +17,7 @@ def fmt_packet_time(ts: float | None) -> str: return datetime.fromtimestamp(ts, timezone.utc).strftime("%H:%M:%S.%f")[:-3] -def build_rows_from_pairs(pairs: list[tuple]) -> list[dict[str, Any]]: +def _build_primary_rows_from_pairs(pairs: list[tuple]) -> list[dict[str, Any]]: rows_out: list[dict[str, Any]] = [] for pair in pairs: page, offset, req_i, req_ts, resp_i, resp_ts, response_content = pair[:7] @@ -59,3 +63,47 @@ def build_rows_from_pairs(pairs: list[tuple]) -> list[dict[str, Any]]: } ) return rows_out + + +def build_rows_from_pairs(pairs: list[tuple]) -> list[dict[str, Any]]: + primary_rows = _build_primary_rows_from_pairs(pairs) + decoded_rows = [row for row in primary_rows if row.get("record_count", 0)] + if decoded_rows and any(row.get("decoder_mode") != "structured_fallback" for row in decoded_rows): + return primary_rows + + assembler = StructuredProtocolAssembler() + for source_index, pair in enumerate(pairs): + assembler.add_blocks( + parse_structured_blocks(pair[6], "monopoly", source_index=source_index) + ) + assembled = assembler.rows("monopoly") + if not assembled: + return primary_rows + + rows_out: list[dict[str, Any]] = [] + converted = structured_monopoly_rows(assembled) + for row_index, (record, structured) in enumerate(zip(converted, assembled), start=1): + source_index = structured.source_index or 0 + pair = pairs[source_index] + page, offset, req_i, req_ts, resp_i, resp_ts, response_content = pair[:7] + kind = pair[7] if len(pair) > 7 else "permanent" + pool = POOL_META.get(kind, POOL_META["permanent"]) + rows_out.append( + { + "page": page, + "offset": offset, + "pool_group_id": pool["id"], + "pool_group_name": pool["name"], + "request_msg": req_i, + "request_time_utc": fmt_packet_time(req_ts), + "response_msg": resp_i, + "response_time_utc": fmt_packet_time(resp_ts), + "response_len": len(response_content), + "record_count": len(assembled), + **record, + "row": row_index, + "structured_assembly": "snapshot_segments", + "structured_assembly_warning_count": len(assembler.warnings), + } + ) + return rows_out diff --git a/src/nte_history_exporter/decoder/structured_protocol.py b/src/nte_history_exporter/decoder/structured_protocol.py index c541147..e93f79e 100644 --- a/src/nte_history_exporter/decoder/structured_protocol.py +++ b/src/nte_history_exporter/decoder/structured_protocol.py @@ -1,7 +1,7 @@ from __future__ import annotations import struct -from dataclasses import dataclass +from dataclasses import dataclass, field, replace from datetime import datetime, timezone from typing import Literal @@ -16,6 +16,10 @@ DOTNET_EPOCH_TICKS = 621_355_968_000_000_000 DOTNET_TICKS_PER_SECOND = 10_000_000 MIN_UNIX_SECONDS = 1_500_000_000 MAX_UNIX_SECONDS = 4_102_444_800 +PROTOCOL_CONSTANT = 0x03000000 +MONOPOLY_BLOCK_KIND = 527 +FORK_BLOCK_KIND = 5906 +MONOPOLY_ENVELOPE_FOOTER = 1_774_080 class StructuredProtocolError(ValueError): @@ -38,6 +42,132 @@ class StructuredRecord: record_end: int record_hex: str protocol_view: str + source_index: int | None = None + generation_index: int | None = None + + +@dataclass(frozen=True) +class ProtocolEnvelope: + record_type: RecordType + stream_key: str + page_index: int + query_high: bool + segment_index: int + + +@dataclass(frozen=True) +class StructuredBlock: + record_type: RecordType + marker_offset: int + declared_size: int + rows: tuple[StructuredRecord, ...] + envelope: ProtocolEnvelope | None + + +@dataclass +class _Generation: + index: int + segments: dict[int, StructuredBlock] = field(default_factory=dict) + + +@dataclass +class _Stream: + generations: list[_Generation] = field(default_factory=list) + + +class StructuredProtocolAssembler: + """Assemble retransmitted and overlapping structured snapshot segments.""" + + def __init__(self) -> None: + self._stream_order: list[str] = [] + self._streams: dict[str, _Stream] = {} + self._legacy_rows: list[StructuredRecord] = [] + self.warnings: list[dict[str, str | int]] = [] + + def add_blocks(self, blocks: list[StructuredBlock]) -> None: + for block in blocks: + self.add_block(block) + + def add_block(self, block: StructuredBlock) -> bool: + envelope = block.envelope + if envelope is None: + if not self._legacy_rows: + self._stream_order.append("__legacy__") + self._legacy_rows.extend(block.rows) + return True + + stream = self._streams.get(envelope.stream_key) + if stream is None: + stream = _Stream() + self._streams[envelope.stream_key] = stream + self._stream_order.append(envelope.stream_key) + if not stream.generations: + stream.generations.append(_Generation(0)) + generation = stream.generations[-1] + existing = generation.segments.get(envelope.segment_index) + if existing is not None: + if _block_signature(existing) == _block_signature(block): + return False + generation = _Generation(len(stream.generations)) + stream.generations.append(generation) + generation.segments[envelope.segment_index] = block + return True + + def rows(self, record_type: RecordType | None = None) -> list[StructuredRecord]: + rows: list[StructuredRecord] = [] + for stream_key in self._stream_order: + if stream_key == "__legacy__": + rows.extend( + row + for row in self._legacy_rows + if record_type is None or row.record_type == record_type + ) + continue + assembled = self._assemble_stream(stream_key, self._streams[stream_key]) + rows.extend(row for row in assembled if record_type is None or row.record_type == record_type) + return rows + + def _assemble_stream(self, stream_key: str, stream: _Stream) -> list[StructuredRecord]: + result: list[StructuredRecord] = [] + result_max_segment: int | None = None + for generation in stream.generations: + if not generation.segments: + continue + generation_rows = _generation_rows(generation) + segment_indexes = sorted(generation.segments) + generation_min = segment_indexes[0] + generation_max = segment_indexes[-1] + if not result: + result = generation_rows + result_max_segment = generation_max + continue + if generation_min == 0: + if result_max_segment is None or generation_max >= result_max_segment: + result = generation_rows + result_max_segment = generation_max + continue + merged = _partial_snapshot_merge(generation_rows, result) + if merged is not None: + result = merged + else: + self._warn(stream_key, generation, "partial snapshot cannot be merged safely") + continue + if result_max_segment is not None and generation_min > result_max_segment: + result.extend(generation_rows) + result_max_segment = generation_max + continue + self._warn(stream_key, generation, "non-zero snapshot reset cannot be merged safely") + return result + + def _warn(self, stream_key: str, generation: _Generation, message: str) -> None: + warning = { + "code": "AMBIGUOUS_STRUCTURED_SNAPSHOT", + "stream_key": stream_key, + "generation_index": generation.index, + "message": message, + } + if warning not in self.warnings: + self.warnings.append(warning) def parse_structured_records(payload: bytes, record_type: RecordType) -> list[StructuredRecord]: @@ -47,24 +177,38 @@ def parse_structured_records(payload: bytes, record_type: RecordType) -> list[St as enrichment/fallback and retain the established decoder as their primary path. """ + assembler = StructuredProtocolAssembler() + assembler.add_blocks(parse_structured_blocks(payload, record_type)) + return assembler.rows(record_type) + + +def parse_structured_blocks( + payload: bytes, + record_type: RecordType, + *, + source_index: int | None = None, +) -> list[StructuredBlock]: marker = MONOPOLY_MARKER if record_type == "monopoly" else FORK_MARKER for view_name, data in _iter_protocol_views(payload): if marker not in data: continue - records: list[StructuredRecord] = [] + blocks: list[StructuredBlock] = [] search_from = 0 while True: marker_pos = data.find(marker, search_from) if marker_pos < 0: break try: - parsed = _parse_block(data, marker_pos, record_type, marker, view_name) + parsed = _parse_block( + data, marker_pos, record_type, marker, view_name, source_index + ) except (StructuredProtocolError, UnicodeDecodeError): - parsed = [] - records.extend(parsed) + parsed = None + if parsed is not None: + blocks.append(parsed) search_from = marker_pos + len(marker) - if records: - return records + if blocks: + return blocks return [] @@ -74,7 +218,9 @@ def _parse_block( record_type: RecordType, marker: bytes, view_name: str, -) -> list[StructuredRecord]: + source_index: int | None, +) -> StructuredBlock: + envelope = _parse_protocol_envelope(record_type, data, marker_pos, view_name) pos = marker_pos + len(marker) if _byte_at(data, pos) == 0: pos += 1 @@ -92,14 +238,22 @@ def _parse_block( for _row_index in range(row_count): row_start = reader.pos if record_type == "monopoly": - record = _parse_monopoly_row(reader, row_start, view_name) + record = _parse_monopoly_row(reader, row_start, view_name, source_index) else: - record = _parse_fork_row(reader, row_start, view_name) + record = _parse_fork_row(reader, row_start, view_name, source_index) records.append(record) - return records + return StructuredBlock( + record_type=record_type, + marker_offset=marker_pos, + declared_size=declared_size, + rows=tuple(records), + envelope=envelope, + ) -def _parse_monopoly_row(reader: "_Reader", row_start: int, view_name: str) -> StructuredRecord: +def _parse_monopoly_row( + reader: "_Reader", row_start: int, view_name: str, source_index: int | None +) -> StructuredRecord: roll_points_raw = reader.u32() item_spec = reader.string() _reserved = reader.u32() @@ -127,10 +281,13 @@ def _parse_monopoly_row(reader: "_Reader", row_start: int, view_name: str) -> St secondary_count, row_start, view_name, + source_index, ) -def _parse_fork_row(reader: "_Reader", row_start: int, view_name: str) -> StructuredRecord: +def _parse_fork_row( + reader: "_Reader", row_start: int, view_name: str, source_index: int | None +) -> StructuredRecord: item_spec = reader.string() pool_id = reader.string() ticks = reader.u64() @@ -145,6 +302,7 @@ def _parse_fork_row(reader: "_Reader", row_start: int, view_name: str) -> Struct None, row_start, view_name, + source_index, ) @@ -159,6 +317,7 @@ def _make_record( secondary_count: int | None, row_start: int, view_name: str, + source_index: int | None, ) -> StructuredRecord: item_id, count = _parse_item_spec(item_spec) if not item_id: @@ -182,6 +341,7 @@ def _make_record( record_end=reader.pos, record_hex=reader.data[row_start : reader.pos].hex(), protocol_view=view_name, + source_index=source_index, ) @@ -197,6 +357,108 @@ def _parse_item_spec(value: str) -> tuple[str, int]: return value, 1 +def _parse_protocol_envelope( + record_type: RecordType, + data: bytes, + marker_pos: int, + view_name: str, +) -> ProtocolEnvelope | None: + if marker_pos == 0 or not view_name.startswith("shift8:"): + return None + if record_type == "monopoly": + if marker_pos < 26: + raise StructuredProtocolError("monopoly envelope is truncated") + protocol_constant = _relative_u32(data, marker_pos, -26) + query_raw = _relative_u32(data, marker_pos, -22) + page_raw = _relative_u32(data, marker_pos, -18) + block_kind = _relative_u32(data, marker_pos, -14) + pool_token = _relative_u32(data, marker_pos, -10) + footer = _relative_u32(data, marker_pos, -6) + if ( + protocol_constant != PROTOCOL_CONSTANT + or block_kind != MONOPOLY_BLOCK_KIND + or footer != MONOPOLY_ENVELOPE_FOOTER + ): + raise StructuredProtocolError("invalid monopoly envelope constants") + stream_key = f"monopoly:{pool_token}" + else: + if marker_pos < 17: + raise StructuredProtocolError("fork envelope is truncated") + protocol_constant = _relative_u32(data, marker_pos, -17) + query_raw = _relative_u32(data, marker_pos, -13) + page_raw = _relative_u32(data, marker_pos, -9) + block_kind = _relative_u32(data, marker_pos, -5) + if protocol_constant != PROTOCOL_CONSTANT or block_kind != FORK_BLOCK_KIND: + raise StructuredProtocolError("invalid fork envelope constants") + stream_key = "fork" + page_index = page_raw & 0x7FFFFFFF + query_high = bool(query_raw & 0x80000000) + return ProtocolEnvelope( + record_type=record_type, + stream_key=stream_key, + page_index=page_index, + query_high=query_high, + segment_index=_segment_index(page_index, query_high), + ) + + +def _segment_index(page_index: int, query_high: bool) -> int: + if query_high: + return page_index * 2 + if page_index > 0: + return page_index * 2 - 1 + raise StructuredProtocolError("low query cannot describe page zero") + + +def _row_signature(row: StructuredRecord) -> tuple: + return ( + row.record_type, + row.ticks, + row.pool_id, + row.item_id, + row.count, + row.roll_points_raw, + row.secondary_item_id, + row.secondary_count, + ) + + +def _block_signature(block: StructuredBlock) -> tuple: + return block.record_type, tuple(_row_signature(row) for row in block.rows) + + +def _generation_rows(generation: _Generation) -> list[StructuredRecord]: + rows = [] + for segment_index in sorted(generation.segments): + block = generation.segments[segment_index] + rows.extend(replace(row, generation_index=generation.index) for row in block.rows) + return rows + + +def _partial_snapshot_merge( + new_rows: list[StructuredRecord], old_rows: list[StructuredRecord] +) -> list[StructuredRecord] | None: + if not new_rows: + return list(old_rows) + if not old_rows: + return list(new_rows) + new_signatures = [_row_signature(row) for row in new_rows] + old_signatures = [_row_signature(row) for row in old_rows] + max_overlap = min(len(new_signatures), len(old_signatures)) + matches: list[tuple[int, int]] = [] + for overlap in range(max_overlap, 0, -1): + suffix = new_signatures[-overlap:] + for position in range(len(old_signatures) - overlap + 1): + if old_signatures[position : position + overlap] == suffix: + matches.append((overlap, position)) + if matches: + break + if len(matches) != 1: + return None + overlap, position = matches[0] + return [*new_rows, *old_rows[position + overlap :]] + + class _Reader: def __init__(self, data: bytes, pos: int) -> None: self.data = data @@ -271,6 +533,13 @@ def _u32_at(data: bytes, pos: int) -> int: raise StructuredProtocolError("u32 exceeds payload") from exc +def _relative_u32(data: bytes, marker_pos: int, relative_pos: int) -> int: + pos = marker_pos + relative_pos + if pos < 0: + raise StructuredProtocolError("envelope position precedes payload") + return _u32_at(data, pos) + + def _u64_at(data: bytes, pos: int) -> int: try: return struct.unpack_from(" bytes: return bytes(8) + packed.to_bytes(len(payload) + 1, "little") +def enveloped_monopoly_payload( + item_spec: str, + *, + page_index: int, + query_high: bool, + pool_token: int = 256, + shift: int = 3, +) -> bytes: + envelope = ( + PROTOCOL_CONSTANT.to_bytes(4, "little") + + ((0x80000000 if query_high else 0)).to_bytes(4, "little") + + page_index.to_bytes(4, "little") + + MONOPOLY_BLOCK_KIND.to_bytes(4, "little") + + pool_token.to_bytes(4, "little") + + MONOPOLY_ENVELOPE_FOOTER.to_bytes(4, "little") + + b"\0\0" + ) + return bit_pack_after_eight_byte_header(envelope + monopoly_payload(item_spec), shift) + + +def enveloped_fork_payload( + item_spec: str, + *, + page_index: int, + query_high: bool, + shift: int = 3, +) -> bytes: + envelope = ( + PROTOCOL_CONSTANT.to_bytes(4, "little") + + ((0x80000000 if query_high else 0)).to_bytes(4, "little") + + page_index.to_bytes(4, "little") + + FORK_BLOCK_KIND.to_bytes(4, "little") + + b"\0" + ) + return bit_pack_after_eight_byte_header(envelope + fork_payload(item_spec), shift) + + +def assembled_block(item_id: str, segment_index: int, source_index: int = 0): + block = parse_structured_blocks( + monopoly_payload(item_id), "monopoly", source_index=source_index + )[0] + envelope = ProtocolEnvelope( + record_type="monopoly", + stream_key="monopoly:256", + page_index=segment_index // 2, + query_high=segment_index % 2 == 0, + segment_index=segment_index, + ) + return replace(block, envelope=envelope) + + class StructuredProtocolTests(unittest.TestCase): def test_structured_monopoly_parser_enriches_the_existing_decoder(self): payload = monopoly_payload( @@ -122,6 +187,126 @@ class StructuredProtocolTests(unittest.TestCase): self.assertEqual(rows[0].item_id, "1003") self.assertEqual(rows[0].protocol_view, "shift8:3") + def test_protocol_envelope_exposes_stream_and_segment_identity(self): + payload = enveloped_monopoly_payload("1003", page_index=2, query_high=False) + + block = parse_structured_blocks(payload, "monopoly")[0] + + self.assertIsNotNone(block.envelope) + self.assertEqual(block.envelope.stream_key, "monopoly:256") + self.assertEqual(block.envelope.page_index, 2) + self.assertFalse(block.envelope.query_high) + self.assertEqual(block.envelope.segment_index, 3) + + def test_assembler_orders_segments_and_ignores_retransmissions(self): + assembler = StructuredProtocolAssembler() + segment_one = assembled_block("1010", 1) + assembler.add_blocks([segment_one, assembled_block("1003", 0), segment_one]) + + rows = assembler.rows("monopoly") + + self.assertEqual([row.item_id for row in rows], ["1003", "1010"]) + self.assertEqual([row.generation_index for row in rows], [0, 0]) + self.assertEqual(assembler.warnings, []) + + def test_assembler_never_deduplicates_blocks_without_envelopes(self): + assembler = StructuredProtocolAssembler() + block = parse_structured_blocks(monopoly_payload("1003"), "monopoly")[0] + + assembler.add_blocks([block, block]) + + self.assertEqual([row.item_id for row in assembler.rows("monopoly")], ["1003", "1003"]) + + def test_new_complete_generation_replaces_old_snapshot(self): + assembler = StructuredProtocolAssembler() + assembler.add_blocks([assembled_block("1003", 0), assembled_block("1010", 1)]) + assembler.add_blocks([assembled_block("1020", 0), assembled_block("1021", 1)]) + + rows = assembler.rows("monopoly") + + self.assertEqual([row.item_id for row in rows], ["1020", "1021"]) + self.assertEqual([row.generation_index for row in rows], [1, 1]) + + def test_partial_generation_merges_only_on_unique_overlap(self): + assembler = StructuredProtocolAssembler() + assembler.add_blocks( + [ + assembled_block("1003", 0), + assembled_block("1010", 1), + assembled_block("1020", 2), + ] + ) + assembler.add_blocks([assembled_block("1099", 0), assembled_block("1010", 1)]) + + rows = assembler.rows("monopoly") + + self.assertEqual([row.item_id for row in rows], ["1099", "1010", "1020"]) + self.assertEqual(assembler.warnings, []) + + def test_ambiguous_partial_generation_keeps_proven_snapshot(self): + assembler = StructuredProtocolAssembler() + assembler.add_blocks( + [ + assembled_block("1003", 0), + assembled_block("1010", 1), + assembled_block("1003", 2), + assembled_block("1010", 3), + ] + ) + assembler.add_blocks([assembled_block("1099", 0), assembled_block("1010", 1)]) + + rows = assembler.rows("monopoly") + + self.assertEqual([row.item_id for row in rows], ["1003", "1010", "1003", "1010"]) + self.assertEqual(assembler.warnings[0]["code"], "AMBIGUOUS_STRUCTURED_SNAPSHOT") + + def test_pair_assembly_is_used_only_for_all_structured_fallback(self): + segment_one = enveloped_monopoly_payload("1010", page_index=1, query_high=False) + segment_zero = enveloped_monopoly_payload("1003", page_index=0, query_high=True) + pairs = [ + (2, 8, 1, 1.0, 2, 1.1, segment_one, "permanent"), + (1, 4, 3, 1.2, 4, 1.3, segment_zero, "permanent"), + ] + + with patch("nte_history_exporter.decoder.protocol._decode_aligned_response_records", return_value=[]): + rows = build_rows_from_pairs(pairs) + annotated = annotate_groups(rows) + + self.assertEqual([row["reward_id"] for row in annotated], ["1003", "1010"]) + self.assertTrue(all(row["structured_assembly"] == "snapshot_segments" for row in annotated)) + self.assertEqual(annotated[0]["uid"], make_uid(annotated[0], 0)) + + def test_pair_assembly_cannot_reorder_successful_heuristic_rows(self): + segment_one = enveloped_monopoly_payload("1010", page_index=1, query_high=False) + segment_zero = enveloped_monopoly_payload("1003", page_index=0, query_high=True) + pairs = [ + (2, 8, 1, 1.0, 2, 1.1, segment_one, "permanent"), + (1, 4, 3, 1.2, 4, 1.3, segment_zero, "permanent"), + ] + + rows = build_rows_from_pairs(pairs) + + self.assertEqual([row["reward_id"] for row in rows], ["1010", "1003"]) + self.assertTrue(all(row["decoder_mode"] == "heuristic_enriched" for row in rows)) + self.assertTrue(all("structured_assembly" not in row for row in rows)) + + def test_arc_fallback_uses_the_same_segment_assembly_and_uid_order(self): + segment_one = enveloped_fork_payload("fork_vine", page_index=1, query_high=False) + segment_zero = enveloped_fork_payload("fork_dustbin", page_index=0, query_high=True) + pairs = [ + (2, 4, 1, 1.0, 2, 1.1, segment_one, "arc_miracle_box"), + (1, 2, 3, 1.2, 4, 1.3, segment_zero, "arc_miracle_box"), + ] + + rows = build_arc_rows_from_pairs(pairs) + + self.assertEqual([row["reward_id"] for row in rows], ["fork_dustbin", "fork_vine"]) + self.assertTrue(all(row["structured_assembly"] == "snapshot_segments" for row in rows)) + self.assertEqual( + rows[0]["uid"], + make_arc_uid(rows[0]["timestamp_raw_hex"], rows[0]["timestamp_group_ordinal"]), + ) + def test_matching_structured_data_enriches_without_replacing_heuristic_identity(self): heuristic_payload = fixture_payload("limited-points-gift-1") original = decode_response_records(heuristic_payload)[0]