Add safe snapshot and segment assembly

Add envelope-aware snapshot and segment assembly for structured Monopoly and Arc history records.
This commit is contained in:
Golumpa 2026-07-15 19:48:11 +01:00
parent a5312520e7
commit 73bb100cd6
8 changed files with 601 additions and 23 deletions

View file

@ -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]

View file

@ -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)

View file

@ -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

View file

@ -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("<Q", data, pos)[0]

View file

@ -29,7 +29,10 @@ FIELDNAMES = [
"dice_raw_u32",
"decoder_mode",
"structured_protocol_view",
"structured_generation_index",
"structured_pool_id",
"structured_assembly",
"structured_assembly_warning_count",
"reward_type",
"reward_id",
"reward_name",