#!/usr/bin/env python3 # SPDX-License-Identifier: GPL-3.0-or-later """Partial-read and deadline tests for Phase-1.0AR.""" from __future__ import annotations import argparse from pathlib import Path import sys import unittest PARSER = argparse.ArgumentParser() PARSER.add_argument("--root", type=Path, required=True) ROOT = PARSER.parse_args().root sys.path.insert(0, str(ROOT / "tools")) from phase10aq_worker_result_record import * # noqa: E402,F403 from phase10ar_result_channel_model import * # noqa: E402,F403 def setup() -> tuple[WorkerPrecommit, bytes, ChannelPlan]: precommit = WorkerPrecommit(b"a" * 16, b"b" * 16, 101, 202, 7, 4096) raw = encode_record(precommit, STATUS_SUCCESS, 4096) return precommit, raw, ChannelPlan(precommit, 101, 7, b"b" * 16) class ResultChannelTests(unittest.TestCase): def test_every_split_point_assembles_exact_record(self) -> None: _, raw, plan = setup() for split in range(1, RECORD_SIZE): events = (FakeReadEvent(DATA, raw[:split]), FakeReadEvent(DATA, raw[split:])) with self.subTest(split=split): outcome = receive_record(plan, events) self.assertTrue(outcome.accepted) self.assertFalse(outcome.eof_is_success or outcome.live_transport_present) def test_byte_at_a_time_is_bounded_and_accepted(self) -> None: _, raw, plan = setup() events = tuple(FakeReadEvent(DATA, bytes((value,))) for value in raw) outcome = receive_record(plan, events) self.assertTrue(outcome.accepted) self.assertEqual(outcome.buffered, RECORD_SIZE) def test_eof_deadline_and_silent_incomplete_are_never_success(self) -> None: _, raw, plan = setup() scripts = ((FakeReadEvent(DATA, raw[:50]), FakeReadEvent(EOF)), (FakeReadEvent(DATA, raw[:50]), FakeReadEvent(DEADLINE)), (FakeReadEvent(DATA, raw[:50]),)) for events in scripts: with self.subTest(events=events): outcome = receive_record(plan, events) self.assertFalse(outcome.accepted or outcome.eof_is_success) self.assertTrue(outcome.containment_required) def test_overflow_and_digest_or_identity_mismatch_require_containment(self) -> None: precommit, raw, plan = setup() overflow = receive_record(plan, (FakeReadEvent(DATA, raw[:100]), FakeReadEvent(DATA, raw[100:] + b"x"),)) self.assertEqual(overflow.classification, "OFFLINE_CHANNEL_OVERFLOW") damaged = bytearray(raw) damaged[10] ^= 1 rejected = receive_record(plan, (FakeReadEvent(DATA, bytes(damaged)),)) self.assertEqual(rejected.classification, "OFFLINE_RECORD_REJECTED") other = WorkerPrecommit(precommit.attempt_id, b"c" * 16, 101, 202, 7, 4096) other_plan = ChannelPlan(other, 101, 7, b"c" * 16) rejected = receive_record(other_plan, (FakeReadEvent(DATA, raw),)) self.assertFalse(rejected.accepted) def test_writer_precommit_must_be_exclusive_and_exact(self) -> None: precommit, _, _ = setup() with self.assertRaises(ChannelModelError): ChannelPlan(precommit, 102, 7, b"b" * 16) with self.assertRaises(ChannelModelError): ChannelPlan(precommit, 101, 7, b"b" * 16, False) def test_deadline_preempts_crossing_read(self) -> None: _, raw, plan = setup() events = (FakeReadEvent(DATA, raw[:64], ticks=256), FakeReadEvent(DATA, raw[64:])) outcome = receive_record(plan, events) self.assertEqual(outcome.classification, "OFFLINE_CHANNEL_DEADLINE") self.assertEqual(outcome.buffered, 64) def test_extra_event_after_boundary_and_invalid_event_are_rejected(self) -> None: _, raw, plan = setup() with self.assertRaises(ChannelModelError): receive_record(plan, (FakeReadEvent(DATA, raw), FakeReadEvent(EOF))) with self.assertRaises(ChannelModelError): FakeReadEvent(DATA, b"") if __name__ == "__main__": unittest.main(argv=[__file__])