Files
chimera-gfx-Public/tools/phase10ar_result_channel_model.py
Chimera GFX release export a6037502d7
phase0-ci / build-and-audit (push) Successful in 2m14s
Publish Chimera GFX source
2026-09-03 03:27:14 +02:00

106 lines
3.9 KiB
Python

#!/usr/bin/env python3
# SPDX-License-Identifier: GPL-3.0-or-later
"""Transport-free partial-read/deadline model for one worker result record."""
from __future__ import annotations
from dataclasses import dataclass
from phase10aq_worker_result_record import (
RECORD_SIZE, WorkerPrecommit, WorkerRecord, WorkerRecordError, parse_record,
)
MAX_CHUNKS = 128
MAX_TICKS = 256
DATA = "DATA"
EOF = "EOF"
DEADLINE = "DEADLINE"
EVENTS = {DATA, EOF, DEADLINE}
class ChannelModelError(RuntimeError):
"""The synthetic channel script or boundary is invalid."""
@dataclass(frozen=True)
class FakeReadEvent:
kind: str
data: bytes = b""
ticks: int = 1
def __post_init__(self) -> None:
if self.kind not in EVENTS or type(self.data) is not bytes:
raise ChannelModelError("read event is invalid")
if (self.kind == DATA) != bool(self.data) or len(self.data) > RECORD_SIZE:
raise ChannelModelError("read event payload is invalid")
if not isinstance(self.ticks, int) or isinstance(self.ticks, bool) \
or not 1 <= self.ticks <= MAX_TICKS:
raise ChannelModelError("read event ticks are invalid")
@dataclass(frozen=True)
class ChannelPlan:
precommit: WorkerPrecommit
writer_pid: int
writer_generation: int
writer_nonce: bytes
exclusive_writer: bool = True
def __post_init__(self) -> None:
if type(self.precommit) is not WorkerPrecommit \
or self.writer_pid != self.precommit.worker_pid \
or self.writer_generation != self.precommit.generation \
or self.writer_nonce != self.precommit.nonce \
or self.exclusive_writer is not True:
raise ChannelModelError("exclusive writer precommit differs")
@dataclass(frozen=True)
class ChannelOutcome:
classification: str
accepted: bool
buffered: int
record: WorkerRecord | None
containment_required: bool
eof_is_success: bool = False
live_transport_present: bool = False
def receive_record(plan: ChannelPlan,
events: tuple[FakeReadEvent, ...]) -> ChannelOutcome:
"""Assemble one exact record from supplied bytes; no I/O is performed."""
if type(plan) is not ChannelPlan or not isinstance(events, tuple) \
or not 1 <= len(events) <= MAX_CHUNKS \
or any(type(item) is not FakeReadEvent for item in events):
raise ChannelModelError("channel model boundary is invalid")
buffer = bytearray()
ticks = 0
for index, event in enumerate(events):
if ticks + event.ticks > MAX_TICKS:
return ChannelOutcome("OFFLINE_CHANNEL_DEADLINE", False,
len(buffer), None, True)
ticks += event.ticks
if event.kind == DEADLINE:
return ChannelOutcome("OFFLINE_CHANNEL_DEADLINE", False,
len(buffer), None, True)
if event.kind == EOF:
return ChannelOutcome("OFFLINE_EOF_BEFORE_COMPLETE_RECORD", False,
len(buffer), None, True)
if len(buffer) + len(event.data) > RECORD_SIZE:
return ChannelOutcome("OFFLINE_CHANNEL_OVERFLOW", False,
len(buffer), None, True)
buffer.extend(event.data)
if len(buffer) == RECORD_SIZE:
if index != len(events) - 1:
raise ChannelModelError("events remain after record boundary")
try:
record = parse_record(bytes(buffer), plan.precommit)
except WorkerRecordError:
return ChannelOutcome("OFFLINE_RECORD_REJECTED", False,
RECORD_SIZE, None, True)
return ChannelOutcome("OFFLINE_EXACT_RECORD_ACCEPTED", True,
RECORD_SIZE, record, False)
return ChannelOutcome("OFFLINE_INCOMPLETE_WITHOUT_DEADLINE", False,
len(buffer), None, True)