pltobing's picture
feat(readiness): materialize full-word timing consensus
d6e2489
Raw History Blame Contribute Delete
32.1 kB
# License: Apache-2.0 License
# Created by: Patrick Lumbantobing, VertoX-AI
# Copyright (c) 2026 VertoX-AI. All rights reserved.
#
# This work is licensed under the Apache-2.0 License.
# To view a copy of this license, visit
# https://www.apache.org/licenses/LICENSE-2.0
"""Materialize deterministic full-word timing from accepted alignment evidence.
This module consumes the repository's backend-neutral alignment-result shape and
frozen reference metadata. It deliberately does not import or execute any
aligner. Source-character spans, Decimal timestamps, provenance, and compact
JSON serialization are kept together so callers can reuse the capability for
other accepted evidence archives without changing backend adapters.
"""
from __future__ import annotations
import csv
import fnmatch
import hashlib
import json
import math
import re
import unicodedata
import zipfile
from collections.abc import Mapping, Sequence
from dataclasses import dataclass
from decimal import Decimal, InvalidOperation
from pathlib import Path
from typing import Any, cast
NORMALIZATION_ID = "sr-v1-nfc-ascii-lexical"
TOKENIZER_PATTERN = r"[A-Za-z]+(?:['’-][A-Za-z]+)*|[0-9]+(?:[.,:/-][0-9]+)*"
_TOKENIZER = re.compile(TOKENIZER_PATTERN)
_ASCII_LETTERS_DIGITS = frozenset(
"abcdefghijklmnopqrstuvwxyzABCDEFGHIJKLMNOPQRSTUVWXYZ0123456789"
)
class MaterializationError(ValueError):
"""Identify an invalid source archive, mapping, or output invariant."""
@dataclass(frozen=True, slots=True)
class TokenSpan:
"""Describe one readiness token in normalized and source coordinates.
Parameters
----------
text: Exact lexical token surface from normalized text.
normalized_start: Inclusive normalized-text offset.
normalized_end: Exclusive normalized-text offset.
raw_start: Inclusive authoritative source offset.
raw_end: Exclusive authoritative source offset.
"""
text: str
normalized_start: int
normalized_end: int
raw_start: int
raw_end: int
@dataclass(frozen=True, slots=True)
class NormalizedSource:
"""Retain normalized text and raw source spans for each output character.
Parameters
----------
text: NFC-normalized, whitespace-collapsed source text.
raw_spans: Half-open raw source spans corresponding one-to-one to text.
"""
text: str
raw_spans: tuple[tuple[int, int], ...]
@dataclass(frozen=True, slots=True)
class MaterializationResult:
"""Identify the files emitted by a successful materialization.
Parameters
----------
jsonl_path: Path to the deterministic full-word JSONL artifact.
manifest_path: Path to its exact-key manifest artifact.
case_count: Number of materialized ordered cases.
lexical_token_count: Total number of emitted lexical words.
accepted_backend_observations: Counts by backend.
"""
jsonl_path: Path
manifest_path: Path
case_count: int
lexical_token_count: int
accepted_backend_observations: tuple[tuple[str, int], ...]
def normalize_readiness_text(source: str) -> NormalizedSource:
"""Normalize authoritative text while preserving source-character evidence.
Parameters
----------
source: Raw authoritative transcript string.
Returns
-------
Normalized text and one raw half-open span for every normalized character.
Raises
------
MaterializationError: If the source is empty after normalization.
Notes
-----
NFC is applied incrementally so a composed character receives the span of
all raw code points that produced it. Unicode whitespace runs collapse to
one ASCII space and retain the complete run span.
"""
if not source:
raise MaterializationError("authoritative transcript is empty")
nfc_chars: list[str] = []
nfc_spans: list[tuple[int, int]] = []
previous = ""
for raw_index in range(len(source)):
candidate = unicodedata.normalize("NFC", source[: raw_index + 1])
common = 0
limit = min(len(previous), len(candidate))
while common < limit and previous[common] == candidate[common]:
common += 1
if common < len(previous):
affected = nfc_spans[common:] + [(raw_index, raw_index + 1)]
if not affected:
raise MaterializationError("failed to retain NFC source span")
span = (
min(item[0] for item in affected),
max(item[1] for item in affected),
)
del nfc_chars[common:]
del nfc_spans[common:]
nfc_chars.extend(candidate[common:])
nfc_spans.extend([span] * len(candidate[common:]))
else:
suffix = candidate[common:]
nfc_chars.extend(suffix)
nfc_spans.extend([(raw_index, raw_index + 1)] * len(suffix))
previous = candidate
collapsed_chars: list[str] = []
collapsed_spans: list[tuple[int, int]] = []
index = 0
while index < len(nfc_chars):
if not nfc_chars[index].isspace():
collapsed_chars.append(nfc_chars[index])
collapsed_spans.append(nfc_spans[index])
index += 1
continue
run_start = index
index += 1
while index < len(nfc_chars) and nfc_chars[index].isspace():
index += 1
run = nfc_spans[run_start:index]
collapsed_chars.append(" ")
collapsed_spans.append(
(min(item[0] for item in run), max(item[1] for item in run))
)
normalized = "".join(collapsed_chars).strip()
if not normalized:
raise MaterializationError("authoritative transcript has no lexical content")
left_trim = len("".join(collapsed_chars)) - len("".join(collapsed_chars).lstrip())
right_trim = len("".join(collapsed_chars).rstrip())
return NormalizedSource(
normalized,
tuple(collapsed_spans[left_trim:right_trim]),
)
def tokenize_readiness_text(source: str) -> tuple[TokenSpan, ...]:
"""Extract exact Stage-2 lexical tokens and their raw source spans.
Parameters
----------
source: Raw authoritative transcript string.
Returns
-------
Ordered lexical token spans using the frozen readiness regex.
Raises
------
MaterializationError: If normalization produces no lexical tokens.
"""
normalized = normalize_readiness_text(source)
spans: list[TokenSpan] = []
for match in _TOKENIZER.finditer(normalized.text):
raw = normalized.raw_spans[match.start() : match.end()]
if not raw:
raise MaterializationError("token has no source-character span")
spans.append(
TokenSpan(
match.group(0),
match.start(),
match.end(),
min(item[0] for item in raw),
max(item[1] for item in raw),
)
)
if not spans:
raise MaterializationError("normalized transcript has no lexical tokens")
return tuple(spans)
def reconcile_token_word_indices(
source: str,
token: TokenSpan,
alignment_words: Sequence[Mapping[str, Any]],
) -> tuple[int, ...]:
"""Map a readiness token to covering backend-neutral alignment words.
Parameters
----------
source: Raw authoritative transcript.
token: Token span produced by :func:`tokenize_readiness_text`.
alignment_words: Backend-neutral word records with ``lexical_word`` spans.
Returns
-------
Ordered zero-based indices of uniquely covering alignment words.
Raises
------
MaterializationError: If coverage is incomplete, overlapping, or invalid.
Invariants
----------
Mapping is based on source-character coverage, not text search or ordinal
position. Multiple readiness tokens may therefore share one aligner word.
"""
relevant = {
index
for index in range(token.raw_start, token.raw_end)
if index < len(source) and source[index] in _ASCII_LETTERS_DIGITS
}
if not relevant:
raise MaterializationError("token has no ASCII letter/digit source evidence")
candidates: list[int] = []
covered_by: dict[int, list[int]] = {index: [] for index in relevant}
for index, word in enumerate(alignment_words):
lexical = word.get("lexical_word")
if not isinstance(lexical, Mapping):
raise MaterializationError("alignment word lacks lexical span")
start = lexical.get("start_char")
end = lexical.get("end_char")
if (
type(start) is not int
or type(end) is not int
or not 0 <= start < end <= len(source)
):
raise MaterializationError("alignment lexical span is invalid")
covered = relevant.intersection(range(start, end))
if covered:
candidates.append(index)
for character in covered:
covered_by[character].append(index)
if not candidates or any(len(indices) != 1 for indices in covered_by.values()):
raise MaterializationError(
"token source characters are ambiguously or incompletely covered"
)
return tuple(candidates)
def decimal_consensus(values: Sequence[Decimal]) -> Decimal:
"""Return the exact Decimal median of finite nonnegative timestamps.
Parameters
----------
values: At least one finite, nonnegative Decimal timestamp.
Returns
-------
Median timestamp without conversion through binary floating point.
Raises
------
MaterializationError: If a value is invalid or no values are supplied.
"""
if not values:
raise MaterializationError("cannot compute consensus without timestamps")
ordered = sorted(values)
if any(not value.is_finite() or value < 0 for value in ordered):
raise MaterializationError(
"consensus timestamps must be finite and nonnegative"
)
middle = len(ordered) // 2
if len(ordered) % 2:
return ordered[middle]
return (ordered[middle - 1] + ordered[middle]) / Decimal(2)
def confidence_label(values: Sequence[Decimal]) -> str:
"""Classify backend agreement using the frozen Decimal thresholds.
Parameters
----------
values: Valid backend end times for one lexical token.
Returns
-------
One of ``UNRESOLVED``, ``SINGLE_BACKEND``, ``HIGH``, ``MEDIUM``, or ``LOW``.
Raises
------
MaterializationError: If a supplied timestamp is invalid.
"""
if any(not value.is_finite() or value < 0 for value in values):
raise MaterializationError(
"confidence timestamps must be finite and nonnegative"
)
count = len(values)
if count == 0:
return "UNRESOLVED"
if count == 1:
return "SINGLE_BACKEND"
spread = max(values) - min(values)
if count == 3 and spread <= Decimal("0.100"):
return "HIGH"
if spread <= Decimal("0.250"):
return "MEDIUM"
return "LOW"
def compute_producer_implementation_sha256(
repository_root: Path, relative_paths: Sequence[str]
) -> str:
"""Hash producer files in the contract's path/NUL/bytes canonical form.
Parameters
----------
repository_root: Repository directory containing the listed files.
relative_paths: Repository-relative production Python paths.
Returns
-------
SHA-256 hex digest over sorted path records.
Raises
------
MaterializationError: If a path is not relative, missing, or duplicated.
"""
if len(set(relative_paths)) != len(relative_paths):
raise MaterializationError("producer implementation paths must be unique")
digest = hashlib.sha256()
for relative in sorted(relative_paths):
path = Path(relative)
if path.is_absolute() or path.parts[0] == "..":
raise MaterializationError(
f"producer path is not repository-relative: {relative}"
)
file_path = repository_root / path
if not file_path.is_file():
raise MaterializationError(
f"producer implementation file is missing: {relative}"
)
digest.update(relative.encode("utf-8"))
digest.update(b"\0")
digest.update(file_path.read_bytes())
digest.update(b"\0")
return digest.hexdigest()
def materialize_full_word_timing(
source_archive: Path,
readiness_archive: Path,
source_manifest_path: Path,
producer_implementation_sha256: str,
producer_config_sha256: str,
output_dir: Path,
) -> MaterializationResult:
"""Materialize full-word timing from immutable accepted evidence archives.
Parameters
----------
source_archive: Frozen W1-T review ZIP containing accepted capsules.
readiness_archive: Frozen Stage-2 reference ZIP containing ordered cases.
source_manifest_path: Manifest binding accepted member paths and hashes.
producer_implementation_sha256: Frozen producer implementation digest.
producer_config_sha256: Frozen producer configuration digest.
output_dir: Empty directory receiving the two materialized artifacts.
Returns
-------
Paths, counts, and backend observation totals for the emitted artifacts.
Raises
------
MaterializationError: If any frozen identity, mapping, provenance, timing,
schema, or deterministic output invariant fails.
"""
if output_dir.exists() and any(output_dir.iterdir()):
raise MaterializationError("materialization output directory must be empty")
output_dir.mkdir(parents=True, exist_ok=True)
manifest = _read_json_file(source_manifest_path)
if _sha256(source_archive.read_bytes()) != manifest["source_review_zip_sha256"]:
raise MaterializationError("source review archive hash does not match manifest")
with (
zipfile.ZipFile(source_archive) as source_zip,
zipfile.ZipFile(readiness_archive) as readiness_zip,
):
cases = _read_cases(readiness_zip)
source_w1t_member = manifest["frozen_w1t_timing_member"]
source_w1t_bytes = source_zip.read(source_w1t_member)
expected_w1t_sha = manifest["frozen_w1t_timing_sha256"]
if _sha256(source_w1t_bytes) != expected_w1t_sha:
raise MaterializationError(
"frozen W1-T timing hash does not match manifest"
)
readiness_w1t_bytes = readiness_zip.read(
"inputs/W1_T_SEMANTIC_READINESS_TIMING.jsonl"
)
if readiness_w1t_bytes != source_w1t_bytes:
raise MaterializationError("readiness and source W1-T timing bytes differ")
frozen_rows = _read_jsonl(source_w1t_bytes)
frozen_by_id = {row["audio_id"]: row for row in frozen_rows}
if set(frozen_by_id) != {case["audio_id"] for case in cases}:
raise MaterializationError("frozen W1-T cases do not match readiness cases")
rows, counts = _build_rows(source_zip, manifest, cases, frozen_by_id)
jsonl_bytes = b"".join(_json_bytes(row) + b"\n" for row in rows)
jsonl_path = output_dir / "W1_T_FULL_WORD_TIMING.jsonl"
manifest_path = output_dir / "W1_T_FULL_WORD_TIMING_MANIFEST.json"
jsonl_path.write_bytes(jsonl_bytes)
output_manifest = {
"case_count": len(rows),
"evaluation_only": True,
"full_word_timing_file_sha256": _sha256(jsonl_bytes),
"human_measured": False,
"ordered_ids_sha256": _ordered_ids_sha256([case["audio_id"] for case in cases]),
"producer_config_sha256": producer_config_sha256,
"producer_implementation_sha256": producer_implementation_sha256,
"schema_version": "sst.w1t.full-word-timing-manifest.v1",
"source_w1t_timing_sha256": manifest["frozen_w1t_timing_sha256"],
}
manifest_path.write_bytes(_json_bytes(output_manifest) + b"\n")
return MaterializationResult(
jsonl_path,
manifest_path,
len(rows),
sum(len(row["words"]) for row in rows),
tuple(sorted(counts.items())),
)
def publish_materialized_outputs(
materialized_dir: Path, readiness_root: Path
) -> tuple[tuple[str, str], ...]:
"""Publish the two validated artifacts without overwriting different bytes.
Parameters
----------
materialized_dir: Directory containing both exact output filenames.
readiness_root: External readiness-input directory receiving the files.
Returns
-------
Ordered filename/status pairs, using ``PUBLISHED`` or
``ALREADY_PRESENT_IDENTICAL``.
Raises
------
MaterializationError: If a source artifact is missing or a target has
different bytes. No differing target is modified.
"""
filenames = (
"W1_T_FULL_WORD_TIMING.jsonl",
"W1_T_FULL_WORD_TIMING_MANIFEST.json",
)
sources = {name: materialized_dir / name for name in filenames}
if any(not path.is_file() for path in sources.values()):
raise MaterializationError("both materialized output files are required")
readiness_root.mkdir(parents=True, exist_ok=True)
statuses: list[tuple[str, str]] = []
for name in filenames:
source_bytes = sources[name].read_bytes()
target = readiness_root / name
if target.exists():
if not target.is_file() or target.read_bytes() != source_bytes:
raise MaterializationError(f"target conflict: {target}")
statuses.append((name, "ALREADY_PRESENT_IDENTICAL"))
continue
target.write_bytes(source_bytes)
statuses.append((name, "PUBLISHED"))
return tuple(statuses)
def _read_cases(readiness_zip: zipfile.ZipFile) -> list[dict[str, str]]:
"""Read and validate the authoritative ordered English metadata rows."""
raw = readiness_zip.read("inputs/STAGE2_36_ENGLISH_ONLY.csv").decode("utf-8")
cases = list(csv.DictReader(raw.splitlines()))
required = {"audio_id", "en_transcript"}
if len(cases) != 36 or not all(required <= set(case) for case in cases):
raise MaterializationError(
"readiness archive must contain exactly 36 English cases"
)
ids = [case["audio_id"] for case in cases]
if len(set(ids)) != len(ids):
raise MaterializationError("readiness case IDs must be unique")
expected = "6f5d20ce8a4b86eb0b51d65cdb463e040165432a589208baa967df0cb016936c"
if _ordered_ids_sha256(ids) != expected:
raise MaterializationError(
"ordered readiness IDs do not match the frozen cohort"
)
return cases
def _build_rows(
source_zip: zipfile.ZipFile,
source_manifest: Mapping[str, Any],
cases: Sequence[Mapping[str, str]],
frozen_by_id: Mapping[str, Mapping[str, Any]],
) -> tuple[list[dict[str, Any]], dict[str, int]]:
"""Build all output rows after validating accepted capsules and mappings."""
accepted = source_manifest["accepted_full_word_evidence"]
output_rows: list[dict[str, Any]] = []
counts: dict[str, int] = {}
for case in cases:
audio_id = case["audio_id"]
transcript = case["en_transcript"]
normalized = normalize_readiness_text(transcript)
tokens = tokenize_readiness_text(transcript)
raw_sha = _sha256(transcript.encode("utf-8"))
normalized_sha = _sha256(normalized.text.encode("utf-8"))
frozen = frozen_by_id[audio_id]
if frozen["original_en_transcript_sha256"] != raw_sha:
raise MaterializationError(
f"source transcript identity failed for {audio_id}"
)
if frozen["w1_normalized_transcript_sha256"] != normalized_sha:
raise MaterializationError(
f"normalized transcript identity failed for {audio_id}"
)
observations_by_token: list[list[dict[str, Any]]] = [[] for _ in tokens]
for backend, backend_spec in sorted(accepted.items()):
member_sha = backend_spec["member_sha256_by_audio_id"].get(audio_id)
if member_sha is None:
continue
member = _find_member(source_zip, backend_spec["member_pattern"], audio_id)
member_bytes = source_zip.read(member)
if _sha256(member_bytes) != member_sha:
raise MaterializationError(
f"accepted member hash failed for {backend}/{audio_id}"
)
capsule = _read_json_bytes(member_bytes)
backend_words = _validate_capsule(
capsule, backend, audio_id, raw_sha, transcript
)
for token_index, token in enumerate(tokens):
mapped = reconcile_token_word_indices(transcript, token, backend_words)
end_time = _decimal(backend_words[mapped[-1]]["end_s"])
observation = {
"backend": backend,
"backend_version": capsule["result"]["provenance"].get(
"backend_version"
),
"model_id": capsule["result"]["provenance"].get("model_id"),
"model_revision": capsule["result"]["provenance"].get(
"model_revision"
),
"config_sha256": capsule["result"]["provenance"].get(
"config_sha256"
),
"source_artifact_sha256": member_sha,
"end_time_s": end_time,
"valid": True,
}
observations_by_token[token_index].append(observation)
counts[backend] = counts.get(backend, 0) + len(tokens)
words: list[dict[str, Any]] = []
for index, (token, observations) in enumerate(
zip(tokens, observations_by_token, strict=True)
):
observations.sort(key=lambda observation: observation["backend"])
valid_times = [observation["end_time_s"] for observation in observations]
if len(valid_times) < 2:
raise MaterializationError(
f"fewer than two accepted observations for {audio_id}"
)
words.append(
{
"backend_observations": observations,
"confidence": confidence_label(valid_times),
"consensus_end_time_s": decimal_consensus(valid_times),
"token_sha256": _sha256(token.text.encode("utf-8")),
"valid_backend_count": len(valid_times),
"word_index": index + 1,
}
)
_validate_frozen_reference(audio_id, tokens, words, frozen)
output_rows.append(
{
"audio_id": audio_id,
"evaluation_only": True,
"human_measured": False,
"lexical_word_count": len(tokens),
"normalization_id": NORMALIZATION_ID,
"normalized_transcript_sha256": normalized_sha,
"schema_version": "sst.w1t.full-word-timing.v1",
"source_transcript_sha256": raw_sha,
"timing_source": "W1_T_FORCED_ALIGNMENT_CONSENSUS",
"words": words,
}
)
return output_rows, counts
def _validate_frozen_reference(
audio_id: str,
tokens: Sequence[TokenSpan],
words: Sequence[Mapping[str, Any]],
frozen: Mapping[str, Any],
) -> None:
"""Require generated timing to reproduce one frozen semantic-ready token.
Parameters
----------
audio_id: Case identity used in deterministic failure messages.
tokens: Complete ordered Stage-2 lexical sequence with source spans.
words: Complete generated full-word timing sequence for the case.
frozen: Frozen W1 semantic-ready reference row for the same case.
Returns
-------
``None`` after all four frozen reproduction invariants pass exactly.
Raises
------
MaterializationError: If the frozen source mapping is malformed or
ambiguous, or if exact Decimal timing, backend count, confidence,
or one-based Stage-2 lexical index differs from the frozen row.
"""
mapping = frozen.get("mapping")
if not isinstance(mapping, Mapping):
raise MaterializationError(f"frozen source mapping is invalid for {audio_id}")
start = mapping.get("start_char")
end = mapping.get("end_char")
if type(start) is not int or type(end) is not int or start < 0 or end <= start:
raise MaterializationError(f"frozen source mapping is invalid for {audio_id}")
# Source identity, rather than W1 ordinal position, locates the Stage-2 token.
matching_indices = [
index
for index, token in enumerate(tokens)
if token.raw_start == start and token.raw_end == end
]
if len(matching_indices) != 1:
raise MaterializationError(
f"frozen source mapping is not unique in Stage-2 tokens for {audio_id}"
)
token_index = matching_indices[0]
generated = words[token_index]
if generated["consensus_end_time_s"] != _decimal(
frozen.get("reference_semantic_ready_time_s")
):
raise MaterializationError(f"frozen W1 timing mismatch for {audio_id}")
if generated["valid_backend_count"] != frozen.get("valid_backend_count"):
raise MaterializationError(f"frozen W1 backend-count mismatch for {audio_id}")
if generated["confidence"] != frozen.get("confidence"):
raise MaterializationError(f"frozen W1 confidence mismatch for {audio_id}")
if generated["word_index"] != frozen.get("reference_semantic_ready_word_count"):
raise MaterializationError(f"frozen W1 lexical-index mismatch for {audio_id}")
def _validate_capsule(
capsule: Mapping[str, Any],
backend: str,
audio_id: str,
transcript_sha: str,
transcript: str,
) -> Sequence[Mapping[str, Any]]:
"""Validate one accepted backend-neutral evidence capsule."""
if (
capsule.get("audio_id") != audio_id
or capsule.get("original_en_transcript_sha256") != transcript_sha
):
raise MaterializationError(f"capsule identity failed for {backend}/{audio_id}")
if capsule.get("reference_text_authoritative") is not True:
raise MaterializationError(
f"capsule is not reference-text authoritative for {backend}/{audio_id}"
)
result = capsule.get("result")
if not isinstance(result, Mapping) or not isinstance(result.get("words"), Sequence):
raise MaterializationError(
f"capsule lacks backend-neutral words for {backend}/{audio_id}"
)
provenance = result.get("provenance")
if not isinstance(provenance, Mapping) or provenance.get("backend_name") != backend:
raise MaterializationError(
f"capsule provenance backend mismatch for {backend}/{audio_id}"
)
previous_end = Decimal(0)
words = result["words"]
for word in words:
if not isinstance(word, Mapping):
raise MaterializationError("backend word is not an object")
lexical = word.get("lexical_word")
if not isinstance(lexical, Mapping):
raise MaterializationError("backend word has no lexical span")
start = lexical.get("start_char")
end = lexical.get("end_char")
if (
type(start) is not int
or type(end) is not int
or not 0 <= start < end <= len(transcript)
):
raise MaterializationError("backend lexical span is out of source bounds")
if transcript[start:end] != lexical.get("text"):
raise MaterializationError(
"backend lexical text differs from authoritative source"
)
start_time = _decimal(word.get("start_s"))
end_time = _decimal(word.get("end_s"))
if start_time < 0 or end_time < start_time or start_time < previous_end:
raise MaterializationError(
"backend word times are negative or non-monotonic"
)
previous_end = end_time
if not words:
raise MaterializationError("backend capsule contains no words")
return cast(Sequence[Mapping[str, Any]], words)
def _find_member(zip_file: zipfile.ZipFile, pattern: str, audio_id: str) -> str:
"""Find exactly one accepted JSON member for an audio ID and pattern."""
matches = [
name
for name in zip_file.namelist()
if fnmatch.fnmatch(name, pattern) and Path(name).stem == audio_id
]
if len(matches) != 1:
raise MaterializationError(
f"expected one accepted member for {audio_id}, found {len(matches)}"
)
return matches[0]
def _read_json_file(path: Path) -> Mapping[str, Any]:
"""Read a UTF-8 JSON object from a filesystem path."""
value = _read_json_bytes(path.read_bytes())
if not isinstance(value, Mapping):
raise MaterializationError(f"JSON manifest is not an object: {path}")
return value
def _read_json_bytes(raw: bytes) -> Any:
"""Parse JSON while retaining source numeric values as Decimal instances."""
try:
return json.loads(raw.decode("utf-8"), parse_int=int, parse_float=Decimal)
except (UnicodeDecodeError, json.JSONDecodeError) as error:
raise MaterializationError("source JSON is malformed") from error
def _read_jsonl(raw: bytes) -> list[Mapping[str, Any]]:
"""Parse non-empty JSONL rows as objects with Decimal numeric values."""
rows: list[Mapping[str, Any]] = []
for line in raw.decode("utf-8").splitlines():
value = _read_json_bytes(line.encode("utf-8"))
if not isinstance(value, Mapping):
raise MaterializationError("JSONL row is not an object")
rows.append(value)
return rows
def _decimal(value: Any) -> Decimal:
"""Convert an already parsed numeric value to a finite Decimal."""
if isinstance(value, Decimal):
result = value
elif isinstance(value, (int, str)) and not isinstance(value, bool):
try:
result = Decimal(value)
except InvalidOperation as error:
raise MaterializationError("timestamp is not Decimal-compatible") from error
else:
raise MaterializationError("timestamp must be a JSON number")
if not result.is_finite():
raise MaterializationError("timestamp must be finite")
return result
def _json_bytes(value: Any) -> bytes:
"""Serialize JSON-compatible data compactly with Decimal numeric literals."""
if isinstance(value, Decimal):
if not value.is_finite():
raise MaterializationError("cannot serialize non-finite Decimal")
return format(value, "f").encode("ascii")
if value is None or isinstance(value, (str, int, float, bool)):
if isinstance(value, float) and not math.isfinite(value):
raise MaterializationError("cannot serialize non-finite float")
return json.dumps(value, ensure_ascii=False, separators=(",", ":")).encode(
"utf-8"
)
if isinstance(value, Mapping):
parts = []
for key in sorted(value):
parts.append(_json_bytes(str(key)) + b":" + _json_bytes(value[key]))
return b"{" + b",".join(parts) + b"}"
if isinstance(value, Sequence) and not isinstance(value, (str, bytes, bytearray)):
return b"[" + b",".join(_json_bytes(item) for item in value) + b"]"
raise MaterializationError(
f"value is not JSON serializable: {type(value).__name__}"
)
def _sha256(raw: bytes) -> str:
"""Return the SHA-256 digest for raw bytes."""
return hashlib.sha256(raw).hexdigest()
def _ordered_ids_sha256(audio_ids: Sequence[str]) -> str:
"""Hash ordered IDs using the frozen no-final-LF convention."""
return _sha256("\n".join(audio_ids).encode("utf-8"))