Files
SDR-Rover/experiments/lab033_priority_channel_scheduler.py
LittleSam129 c486039053 Split experiments from tests
The tests/ directory held 50 laboratory programs and no tests. They model
channels, run hundreds of repetitions and write CSV, PNG and reports;
calling that a test suite blocked introducing a real one, because any
pytest run would have collected the labs and re-executed every
experiment.

- move all 50 lab programs to experiments/ with git mv, preserving history
- rewrite the 38 cross-imports between labs from tests.labNNN to
  experiments.labNNN
- leave tests/ empty for actual fast checks of protocol/
- point quick_gate and the hook at the new layout and add experiments/ to
  the syntax sweep
- update the paths quoted in the Lab042 specification and the verifier
  agent definition

This also defuses the import-time work finding without touching 41 files:
the labs still create directories and write files on import, but nothing
imports them now except the gate, which does so deliberately.

Gate passes: syntax clean, protocol imports, 15 lab modules import, 2
functional suites run.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-10 14:34:58 +03:00

1044 lines
47 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""Lab033: one shared link queue for control, telemetry, and FEC video.
The experiment deliberately models one abstract, error-free, non-preemptive
serialization resource. Both logical directions consume the same configured
bitrate; direction-switch time and a physical duplexing scheme are outside the
scope of this laboratory.
"""
from __future__ import annotations
import csv
from dataclasses import asdict, dataclass
from pathlib import Path
import struct
import subprocess
from typing import Iterable
import cv2
import matplotlib
import numpy as np
matplotlib.use("Agg")
import matplotlib.pyplot as plt
from protocol.link_packet import (
Direction,
HEADER_FORMAT,
HEADER_SIZE,
LinkPacket,
LinkPacketCRCError,
TrafficClass,
decode_link_packet,
encode_link_packet,
)
from protocol.packet_erasure_fec import (
decode_fec_block,
decode_outer_symbol,
)
from protocol.priority_scheduler import (
QueuedPacket,
ScheduleResult,
SchedulerMode,
schedule_packets,
transmission_duration_seconds,
)
from protocol.video_packet import (
CompositeReassembler,
decode_packet as decode_inner_packet,
)
from experiments.lab028_video_packetization import (
COMPOSITE_FPS,
SOURCE_VIDEO_PATH,
EncodedComposite,
VideoMetadata,
load_video_profile,
)
from experiments.lab029_packet_channel_simulation import prepare_profiles
from experiments.lab030_packet_erasure_fec import (
FECBlockPlan,
SourcePacket,
TransmissionUnit,
prepare_source_packets,
)
from experiments.lab032_fec_parameter_sweep import (
SweepMode,
build_parameter_units,
)
OUTPUT_DIRECTORY = Path("data/processed/lab033")
SUMMARY_CSV_PATH = OUTPUT_DIRECTORY / "lab033_summary.csv"
CLASS_CSV_PATH = OUTPUT_DIRECTORY / "lab033_class_metrics.csv"
REPORT_PATH = OUTPUT_DIRECTORY / "lab033_report.txt"
CONTROL_PLOT_PATH = OUTPUT_DIRECTORY / "lab033_control_delay.png"
EMERGENCY_PLOT_PATH = OUTPUT_DIRECTORY / "lab033_emergency_delay.png"
VIDEO_PLOT_PATH = OUTPUT_DIRECTORY / "lab033_video_delay.png"
QUEUE_PLOT_PATH = OUTPUT_DIRECTORY / "lab033_queue_length.png"
DEADLINE_PLOT_PATH = OUTPUT_DIRECTORY / "lab033_deadline_misses.png"
COMPARISON_PLOT_PATH = OUTPUT_DIRECTORY / "lab033_scheduler_comparison.png"
PLOT_PATHS = (
CONTROL_PLOT_PATH,
EMERGENCY_PLOT_PATH,
VIDEO_PLOT_PATH,
QUEUE_PLOT_PATH,
DEADLINE_PLOT_PATH,
COMPARISON_PLOT_PATH,
)
VIDEO_MODE = SweepMode("12+3", 12, 3, "12+3")
VIDEO_PAYLOAD_SIZE = 512
CHANNEL_RATES_KBPS = (300.0, 260.0, 230.0)
SCHEDULER_MODES = (
SchedulerMode.FIFO,
SchedulerMode.STRICT_PRIORITY,
SchedulerMode.LATEST_STATE,
)
CLASS_NAMES = {
TrafficClass.EMERGENCY: "аварийная команда",
TrafficClass.CONTROL: "обычные команды",
TrafficClass.TELEMETRY: "телеметрия",
TrafficClass.VIDEO: "видео",
}
MODE_NAMES = {
SchedulerMode.FIFO: "FIFO",
SchedulerMode.STRICT_PRIORITY: "Строгий приоритет",
SchedulerMode.LATEST_STATE: "Приоритет + замена",
}
STREAM_EMERGENCY = 1
STREAM_CONTROL = 2
STREAM_TELEMETRY = 3
STREAM_VIDEO = 4
CONTROL_PERIOD_US = 50_000
TELEMETRY_PERIOD_US = 100_000
EMERGENCY_TIME_US = 10_000_000
TIME_EPSILON_SECONDS = 1e-9
@dataclass(frozen=True)
class Workload:
metadata: VideoMetadata
composites: tuple[EncodedComposite, ...]
packets: tuple[QueuedPacket, ...]
video_units: tuple[TransmissionUnit, ...]
video_blocks: tuple[FECBlockPlan, ...]
total_jpeg_bytes: int
@dataclass(frozen=True)
class ClassMetrics:
channel_kbps: float
scheduler: str
traffic_class: str
offered_load_kbps: float
created_packets: int
transmitted_packets: int
replaced_packets: int
remaining_at_source_end_packets: int
mean_start_delay_ms: float
p95_start_delay_ms: float
max_start_delay_ms: float
mean_total_delay_ms: float
p95_total_delay_ms: float
max_total_delay_ms: float
deadline_misses: int
deadline_miss_fraction: float
mean_blocking_delay_ms: float
max_blocking_delay_ms: float
@dataclass(frozen=True)
class SimulationSummary:
channel_kbps: float
scheduler: str
offered_load_kbps: float
emergency_load_kbps: float
control_load_kbps: float
telemetry_load_kbps: float
video_load_kbps: float
service_share_percent: float
offered_to_capacity_ratio: float
transmitted_packets: int
transmitted_bytes: int
source_duration_seconds: float
drain_end_seconds: float
additional_drain_seconds: float
queue_at_source_end_packets: int
queue_at_source_end_bytes: float
mean_queue_packets: float
max_queue_packets: int
mean_queue_bytes: float
max_queue_bytes: int
control_mean_age_ms: float
control_p95_age_ms: float
control_max_age_ms: float
control_gaps_over_100ms: int
control_max_receive_gap_ms: float
control_replaced: int
emergency_start_delay_ms: float
emergency_total_delay_ms: float
emergency_deadline_met: bool
emergency_blocker_class: str
emergency_blocker_size_bytes: int
emergency_blocker_sequence: int
emergency_blocking_delay_ms: float
telemetry_mean_age_ms: float
telemetry_p95_age_ms: float
telemetry_max_age_ms: float
telemetry_misses_500ms: int
telemetry_miss_fraction: float
telemetry_replaced: int
video_published_at_source_end_fraction: float
video_published_after_drain_fraction: float
video_final_recovery_fraction: float
video_mean_publication_delay_ms: float
video_p95_publication_delay_ms: float
video_max_publication_delay_ms: float
video_mean_display_age_ms: float
video_p95_display_age_ms: float
video_max_display_age_ms: float
video_age_over_1s_fraction: float
video_late_frame_count: int
video_max_late_run: int
@dataclass(frozen=True)
class ScenarioResult:
summary: SimulationSummary
classes: tuple[ClassMetrics, ...]
schedule: ScheduleResult
publication_times: dict[int, float]
@dataclass(frozen=True)
class FunctionalTestResult:
name: str
passed: bool
detail: str
def percentile(values: Iterable[float], q: float) -> float:
values = tuple(values)
return float(np.percentile(values, q)) if values else 0.0
def deterministic_payload(prefix: bytes, sequence: int, size: int) -> bytes:
"""Build an exact-size reproducible educational payload."""
head = prefix[:8].ljust(8, b"\0") + struct.pack("!I", sequence)
tail = bytes((sequence * 37 + index * 17) & 0xFF for index in range(size))
return (head + tail)[:size]
def _packet(
traffic_class: TrafficClass,
direction: Direction,
stream_id: int,
sequence: int,
generation_time_us: int,
deadline_ms: int,
payload: bytes,
) -> LinkPacket:
return LinkPacket(
traffic_class=traffic_class,
direction=direction,
stream_id=stream_id,
sequence_number=sequence,
generation_time_us=generation_time_us,
deadline_ms=deadline_ms,
payload=payload,
)
def build_workload() -> Workload:
metadata, composites = load_video_profile(SOURCE_VIDEO_PATH)
profile = prepare_profiles(composites)[VIDEO_PAYLOAD_SIZE]
source_packets: tuple[SourcePacket, ...] = prepare_source_packets(profile)
video_units, video_blocks = build_parameter_units(source_packets, VIDEO_MODE)
candidates: list[LinkPacket] = []
duration_us = int(round(metadata.duration_seconds * 1_000_000.0))
for sequence, generation_us in enumerate(range(0, duration_us, CONTROL_PERIOD_US)):
candidates.append(
_packet(
TrafficClass.CONTROL,
Direction.GROUND_TO_ROVER,
STREAM_CONTROL,
sequence,
generation_us,
100,
deterministic_payload(b"CONTROL", sequence, 32),
)
)
for sequence, generation_us in enumerate(range(0, duration_us, TELEMETRY_PERIOD_US)):
candidates.append(
_packet(
TrafficClass.TELEMETRY,
Direction.ROVER_TO_GROUND,
STREAM_TELEMETRY,
sequence,
generation_us,
500,
deterministic_payload(b"TELEM", sequence, 64),
)
)
if EMERGENCY_TIME_US < duration_us:
candidates.append(
_packet(
TrafficClass.EMERGENCY,
Direction.GROUND_TO_ROVER,
STREAM_EMERGENCY,
0,
EMERGENCY_TIME_US,
50,
deterministic_payload(b"E-STOP", 0, 32),
)
)
for unit in video_units:
candidates.append(
_packet(
TrafficClass.VIDEO,
Direction.ROVER_TO_GROUND,
STREAM_VIDEO,
unit.sequence_index,
int(round(unit.generation_time_seconds * 1_000_000.0)),
0,
unit.wire_packet,
)
)
# FIFO coincidence rule: class priority, stream_id, then sequence number.
candidates.sort(
key=lambda packet: (
packet.generation_time_us,
int(packet.traffic_class),
packet.stream_id,
packet.sequence_number,
)
)
packets = tuple(
QueuedPacket(packet=packet, arrival_order=index)
for index, packet in enumerate(candidates)
)
return Workload(
metadata=metadata,
composites=tuple(composites),
packets=packets,
video_units=video_units,
video_blocks=video_blocks,
total_jpeg_bytes=sum(
len(frame.base_jpeg) + len(frame.roi_jpeg) for frame in composites
),
)
def _class_items(workload: Workload, traffic_class: TrafficClass) -> tuple[QueuedPacket, ...]:
return tuple(
item for item in workload.packets
if item.packet.traffic_class is traffic_class
)
def _wire_bytes(item: QueuedPacket) -> int:
return HEADER_SIZE + len(item.packet.payload)
def queue_metrics(
schedule: ScheduleResult,
source_end: float,
) -> tuple[float, int, float, int, int, float]:
"""Time-weighted in-system queue metrics over the source interval."""
intervals: list[tuple[float, float, int]] = []
for sent in schedule.transmitted:
intervals.append(
(
sent.queued.generation_time_seconds,
sent.end_seconds,
_wire_bytes(sent.queued),
)
)
for replacement in schedule.replacements:
intervals.append(
(
replacement.removed.generation_time_seconds,
replacement.replacement.generation_time_seconds,
_wire_bytes(replacement.removed),
)
)
packet_area = 0.0
byte_area = 0.0
events: dict[float, list[int]] = {}
for start, end, size in intervals:
clipped_start = max(0.0, start)
clipped_end = min(source_end, end)
if clipped_end <= clipped_start + TIME_EPSILON_SECONDS:
continue
packet_area += clipped_end - clipped_start
byte_area += (clipped_end - clipped_start) * size
delta = events.setdefault(clipped_start, [0, 0])
delta[0] += 1
delta[1] += size
delta = events.setdefault(clipped_end, [0, 0])
delta[0] -= 1
delta[1] -= size
packets_now = bytes_now = max_packets = max_bytes = 0
for time_seconds in sorted(events):
packets_now += events[time_seconds][0]
bytes_now += events[time_seconds][1]
max_packets = max(max_packets, packets_now)
max_bytes = max(max_bytes, bytes_now)
remaining = [
sent for sent in schedule.transmitted
if sent.queued.generation_time_seconds <= source_end + TIME_EPSILON_SECONDS
and sent.end_seconds > source_end + TIME_EPSILON_SECONDS
]
remaining_bytes = 0.0
for sent in remaining:
if sent.start_seconds >= source_end:
remaining_bytes += _wire_bytes(sent.queued)
else:
remaining_bytes += (
sent.end_seconds - source_end
) * schedule.channel_bitrate_bps / 8.0
return (
packet_area / source_end,
max_packets,
byte_area / source_end,
max_bytes,
len(remaining),
remaining_bytes,
)
def receive_video(
schedule: ScheduleResult,
frame_count: int,
) -> tuple[dict[int, float], int]:
"""Unwrap both CRC layers and publish only complete BASE+ROI frames."""
receiver = CompositeReassembler()
publication_times: dict[int, float] = {}
block_symbols: dict[int, list[bytes]] = {}
decoded_blocks: set[int] = set()
for sent in schedule.transmitted:
link = decode_link_packet(sent.wire_packet)
if link.traffic_class is not TrafficClass.VIDEO:
continue
outer = decode_outer_symbol(link.payload)
symbols = block_symbols.setdefault(outer.block_id, [])
symbols.append(link.payload)
if not outer.is_parity:
completed = receiver.ingest(outer.data)
if completed is not None:
publication_times[completed.composite_frame_id] = sent.end_seconds
if outer.block_id not in decoded_blocks and len(symbols) >= outer.source_count:
decoded = decode_fec_block(tuple(symbols))
for inner in decoded.source_packets:
decode_inner_packet(inner)
decoded_blocks.add(outer.block_id)
if len(publication_times) != frame_count:
raise AssertionError(
f"video receiver published {len(publication_times)} of {frame_count} frames"
)
return publication_times, len(decoded_blocks)
def display_age_metrics(
publication_times: dict[int, float],
drain_end: float,
) -> tuple[float, float, float, float]:
events = sorted((time, frame_id) for frame_id, time in publication_times.items())
sample_times = np.arange(0.0, drain_end + 0.005, 0.01)
ages = []
event_index = 0
last_frame_id: int | None = None
for sample_time in sample_times:
while event_index < len(events) and events[event_index][0] <= sample_time:
last_frame_id = events[event_index][1]
event_index += 1
generation = 0.0 if last_frame_id is None else last_frame_id / COMPOSITE_FPS
ages.append(max(0.0, sample_time - generation))
return (
float(np.mean(ages)),
percentile(ages, 95),
max(ages, default=0.0),
sum(age > 1.0 for age in ages) / len(ages) if ages else 0.0,
)
def max_true_run(flags: Iterable[bool]) -> int:
longest = current = 0
for flag in flags:
current = current + 1 if flag else 0
longest = max(longest, current)
return longest
def class_metrics(
workload: Workload,
schedule: ScheduleResult,
channel_kbps: float,
source_end: float,
) -> tuple[ClassMetrics, ...]:
rows = []
for traffic_class in TrafficClass:
created = _class_items(workload, traffic_class)
sent = tuple(
record for record in schedule.transmitted
if record.queued.packet.traffic_class is traffic_class
)
replaced = tuple(
item for item in schedule.replacements
if item.removed.packet.traffic_class is traffic_class
)
start_delays = [
record.start_seconds - record.queued.generation_time_seconds
for record in sent
]
total_delays = [
record.end_seconds - record.queued.generation_time_seconds
for record in sent
]
blocking = [record.blocking_delay_seconds for record in sent]
deadline_misses = sum(
delay > record.queued.packet.deadline_ms / 1000.0 + TIME_EPSILON_SECONDS
for delay, record in zip(total_delays, sent)
if record.queued.packet.deadline_ms > 0
)
deadline_population = sum(
record.queued.packet.deadline_ms > 0 for record in sent
)
rows.append(
ClassMetrics(
channel_kbps=channel_kbps,
scheduler=schedule.mode.value,
traffic_class=traffic_class.name.lower(),
offered_load_kbps=sum(_wire_bytes(item) for item in created) * 8.0 / source_end / 1000.0,
created_packets=len(created),
transmitted_packets=len(sent),
replaced_packets=len(replaced),
remaining_at_source_end_packets=sum(
record.end_seconds > source_end + TIME_EPSILON_SECONDS for record in sent
),
mean_start_delay_ms=float(np.mean(start_delays)) * 1000.0 if start_delays else 0.0,
p95_start_delay_ms=percentile(start_delays, 95) * 1000.0,
max_start_delay_ms=max(start_delays, default=0.0) * 1000.0,
mean_total_delay_ms=float(np.mean(total_delays)) * 1000.0 if total_delays else 0.0,
p95_total_delay_ms=percentile(total_delays, 95) * 1000.0,
max_total_delay_ms=max(total_delays, default=0.0) * 1000.0,
deadline_misses=deadline_misses,
deadline_miss_fraction=deadline_misses / deadline_population if deadline_population else 0.0,
mean_blocking_delay_ms=float(np.mean(blocking)) * 1000.0 if blocking else 0.0,
max_blocking_delay_ms=max(blocking, default=0.0) * 1000.0,
)
)
return tuple(rows)
def simulate(workload: Workload, channel_kbps: float, mode: SchedulerMode) -> ScenarioResult:
bitrate = channel_kbps * 1000.0
source_end = workload.metadata.duration_seconds
schedule = schedule_packets(workload.packets, mode, bitrate)
publication_times, decoded_blocks = receive_video(
schedule, len(workload.composites)
)
if decoded_blocks != len(workload.video_blocks):
raise AssertionError("not every FEC block passed Lab030/Lab028 decoding")
classes = class_metrics(workload, schedule, channel_kbps, source_end)
by_class = {
TrafficClass[row.traffic_class.upper()]: row for row in classes
}
sent_by_class = {
traffic_class: tuple(
record for record in schedule.transmitted
if record.queued.packet.traffic_class is traffic_class
)
for traffic_class in TrafficClass
}
load_by_class = {
traffic_class: sum(_wire_bytes(item) for item in _class_items(workload, traffic_class))
* 8.0 / source_end / 1000.0
for traffic_class in TrafficClass
}
offered_load = sum(load_by_class.values())
useful_bytes = workload.total_jpeg_bytes + sum(
len(item.packet.payload)
for item in workload.packets
if item.packet.traffic_class is not TrafficClass.VIDEO
)
proposed_wire_bytes = sum(_wire_bytes(item) for item in workload.packets)
service_share = 100.0 * (proposed_wire_bytes - useful_bytes) / proposed_wire_bytes
drain_end = max(record.end_seconds for record in schedule.transmitted)
mean_qp, max_qp, mean_qb, max_qb, queue_end_packets, queue_end_bytes = queue_metrics(
schedule, source_end
)
control = sent_by_class[TrafficClass.CONTROL]
control_ages = [record.end_seconds - record.queued.generation_time_seconds for record in control]
control_receive_times = [record.end_seconds for record in control]
control_gaps = [right - left for left, right in zip(control_receive_times, control_receive_times[1:])]
telemetry = sent_by_class[TrafficClass.TELEMETRY]
telemetry_ages = [record.end_seconds - record.queued.generation_time_seconds for record in telemetry]
emergency = sent_by_class[TrafficClass.EMERGENCY]
if len(emergency) != 1:
raise AssertionError("exactly one emergency command must be transmitted")
emergency_record = emergency[0]
blocker = emergency_record.blocked_by
publication_delays = [
publication_times[index] - index / COMPOSITE_FPS
for index in range(len(workload.composites))
]
late_flags = [delay > 1.0 for delay in publication_delays]
mean_age, p95_age, max_age, age_over_one = display_age_metrics(
publication_times, drain_end
)
summary = SimulationSummary(
channel_kbps=channel_kbps,
scheduler=mode.value,
offered_load_kbps=offered_load,
emergency_load_kbps=load_by_class[TrafficClass.EMERGENCY],
control_load_kbps=load_by_class[TrafficClass.CONTROL],
telemetry_load_kbps=load_by_class[TrafficClass.TELEMETRY],
video_load_kbps=load_by_class[TrafficClass.VIDEO],
service_share_percent=service_share,
offered_to_capacity_ratio=offered_load / channel_kbps,
transmitted_packets=len(schedule.transmitted),
transmitted_bytes=sum(len(record.wire_packet) for record in schedule.transmitted),
source_duration_seconds=source_end,
drain_end_seconds=drain_end,
additional_drain_seconds=max(0.0, drain_end - source_end),
queue_at_source_end_packets=queue_end_packets,
queue_at_source_end_bytes=queue_end_bytes,
mean_queue_packets=mean_qp,
max_queue_packets=max_qp,
mean_queue_bytes=mean_qb,
max_queue_bytes=max_qb,
control_mean_age_ms=float(np.mean(control_ages)) * 1000.0,
control_p95_age_ms=percentile(control_ages, 95) * 1000.0,
control_max_age_ms=max(control_ages) * 1000.0,
control_gaps_over_100ms=sum(gap > 0.1 + TIME_EPSILON_SECONDS for gap in control_gaps),
control_max_receive_gap_ms=max(control_gaps, default=0.0) * 1000.0,
control_replaced=by_class[TrafficClass.CONTROL].replaced_packets,
emergency_start_delay_ms=(emergency_record.start_seconds - emergency_record.queued.generation_time_seconds) * 1000.0,
emergency_total_delay_ms=(emergency_record.end_seconds - emergency_record.queued.generation_time_seconds) * 1000.0,
emergency_deadline_met=(emergency_record.end_seconds - emergency_record.queued.generation_time_seconds) <= 0.05 + TIME_EPSILON_SECONDS,
emergency_blocker_class=blocker.packet.traffic_class.name.lower() if blocker else "none",
emergency_blocker_size_bytes=_wire_bytes(blocker) if blocker else 0,
emergency_blocker_sequence=blocker.packet.sequence_number if blocker else -1,
emergency_blocking_delay_ms=emergency_record.blocking_delay_seconds * 1000.0,
telemetry_mean_age_ms=float(np.mean(telemetry_ages)) * 1000.0,
telemetry_p95_age_ms=percentile(telemetry_ages, 95) * 1000.0,
telemetry_max_age_ms=max(telemetry_ages) * 1000.0,
telemetry_misses_500ms=sum(age > 0.5 + TIME_EPSILON_SECONDS for age in telemetry_ages),
telemetry_miss_fraction=sum(age > 0.5 + TIME_EPSILON_SECONDS for age in telemetry_ages) / len(telemetry_ages),
telemetry_replaced=by_class[TrafficClass.TELEMETRY].replaced_packets,
video_published_at_source_end_fraction=sum(time <= source_end + TIME_EPSILON_SECONDS for time in publication_times.values()) / len(workload.composites),
video_published_after_drain_fraction=sum(time <= drain_end + TIME_EPSILON_SECONDS for time in publication_times.values()) / len(workload.composites),
video_final_recovery_fraction=len(publication_times) / len(workload.composites),
video_mean_publication_delay_ms=float(np.mean(publication_delays)) * 1000.0,
video_p95_publication_delay_ms=percentile(publication_delays, 95) * 1000.0,
video_max_publication_delay_ms=max(publication_delays) * 1000.0,
video_mean_display_age_ms=mean_age * 1000.0,
video_p95_display_age_ms=p95_age * 1000.0,
video_max_display_age_ms=max_age * 1000.0,
video_age_over_1s_fraction=age_over_one,
video_late_frame_count=sum(late_flags),
video_max_late_run=max_true_run(late_flags),
)
return ScenarioResult(summary, classes, schedule, publication_times)
def run_experiment(workload: Workload) -> tuple[ScenarioResult, ...]:
return tuple(
simulate(workload, rate, mode)
for rate in CHANNEL_RATES_KBPS
for mode in SCHEDULER_MODES
)
def _synthetic_item(
traffic_class: TrafficClass,
sequence: int,
generation_us: int,
arrival_order: int,
payload_size: int = 32,
stream_id: int | None = None,
) -> QueuedPacket:
direction = (
Direction.GROUND_TO_ROVER
if traffic_class in (TrafficClass.EMERGENCY, TrafficClass.CONTROL)
else Direction.ROVER_TO_GROUND
)
return QueuedPacket(
_packet(
traffic_class,
direction,
int(traffic_class) if stream_id is None else stream_id,
sequence,
generation_us,
100,
deterministic_payload(b"TEST", sequence, payload_size),
),
arrival_order,
)
def run_functional_tests(
workload: Workload,
results: tuple[ScenarioResult, ...],
) -> tuple[FunctionalTestResult, ...]:
checks: list[tuple[str, callable]] = []
def add(name: str):
def decorator(function):
checks.append((name, function))
return function
return decorator
sample = next(item for item in workload.packets if item.packet.traffic_class is TrafficClass.VIDEO)
@add("01_header_roundtrip")
def _():
decoded = decode_link_packet(encode_link_packet(sample.packet))
assert decoded == LinkPacket(**{**asdict(sample.packet), "packet_crc32": decoded.packet_crc32})
@add("02_header_crc_detection")
def _():
wire = bytearray(encode_link_packet(sample.packet)); wire[11] ^= 1
try: decode_link_packet(wire)
except LinkPacketCRCError: return
raise AssertionError("header corruption was accepted")
@add("03_payload_crc_detection")
def _():
wire = bytearray(encode_link_packet(sample.packet)); wire[-1] ^= 1
try: decode_link_packet(wire)
except LinkPacketCRCError: return
raise AssertionError("payload corruption was accepted")
@add("04_video_wrapper_byte_exact")
def _():
assert decode_link_packet(encode_link_packet(sample.packet)).payload == sample.packet.payload
@add("05_lab030_crc")
def _():
decode_outer_symbol(decode_link_packet(encode_link_packet(sample.packet)).payload)
@add("06_lab028_crc")
def _():
systematic = next(unit for unit in workload.video_units if not unit.is_parity)
decode_inner_packet(decode_outer_symbol(systematic.wire_packet).data)
@add("07_fifo_global_order")
def _():
items = (_synthetic_item(TrafficClass.VIDEO, 0, 0, 0, 1000), _synthetic_item(TrafficClass.EMERGENCY, 0, 1, 1))
order = schedule_packets(items, SchedulerMode.FIFO, 100_000).transmitted
assert [item.queued.arrival_order for item in order] == [0, 1]
@add("08_strict_priority")
def _():
items = tuple(_synthetic_item(cls, 0, 0, index) for index, cls in enumerate((TrafficClass.VIDEO, TrafficClass.TELEMETRY, TrafficClass.CONTROL, TrafficClass.EMERGENCY)))
order = schedule_packets(items, SchedulerMode.STRICT_PRIORITY, 100_000).transmitted
assert [item.queued.packet.traffic_class for item in order] == list(TrafficClass)
@add("09_non_preemptive")
def _():
items = (_synthetic_item(TrafficClass.VIDEO, 0, 0, 0, 1000), _synthetic_item(TrafficClass.EMERGENCY, 0, 1, 1))
order = schedule_packets(items, SchedulerMode.STRICT_PRIORITY, 100_000).transmitted
assert order[0].queued.packet.traffic_class is TrafficClass.VIDEO
assert order[1].start_seconds >= order[0].end_seconds - TIME_EPSILON_SECONDS
@add("10_intra_class_sequence")
def _():
items = tuple(_synthetic_item(TrafficClass.CONTROL, index, 0, index) for index in range(4))
order = schedule_packets(items, SchedulerMode.STRICT_PRIORITY, 100_000).transmitted
assert [item.queued.packet.sequence_number for item in order] == list(range(4))
@add("11_latest_state_scope")
def _():
items = (
_synthetic_item(TrafficClass.VIDEO, 0, 0, 0, 1000),
_synthetic_item(TrafficClass.CONTROL, 0, 1, 1, stream_id=20),
_synthetic_item(TrafficClass.CONTROL, 1, 2, 2, stream_id=20),
_synthetic_item(TrafficClass.CONTROL, 2, 3, 3, stream_id=21),
_synthetic_item(TrafficClass.TELEMETRY, 0, 4, 4, stream_id=30),
_synthetic_item(TrafficClass.TELEMETRY, 1, 5, 5, stream_id=30),
)
result = schedule_packets(items, SchedulerMode.LATEST_STATE, 100_000)
removed = {(x.removed.packet.traffic_class, x.removed.packet.stream_id, x.removed.packet.sequence_number) for x in result.replacements}
assert removed == {(TrafficClass.CONTROL, 20, 0), (TrafficClass.TELEMETRY, 30, 0)}
@add("12_active_state_not_removed")
def _():
items = (_synthetic_item(TrafficClass.CONTROL, 0, 0, 0, 1000), _synthetic_item(TrafficClass.CONTROL, 1, 1, 1))
result = schedule_packets(items, SchedulerMode.LATEST_STATE, 100_000)
assert len(result.transmitted) == 2 and not result.replacements
@add("13_emergency_never_replaced")
def _():
items = (_synthetic_item(TrafficClass.VIDEO, 0, 0, 0, 1000), _synthetic_item(TrafficClass.EMERGENCY, 0, 1, 1), _synthetic_item(TrafficClass.EMERGENCY, 1, 2, 2))
result = schedule_packets(items, SchedulerMode.LATEST_STATE, 100_000)
assert sum(x.queued.packet.traffic_class is TrafficClass.EMERGENCY for x in result.transmitted) == 2
@add("14_300kbps_all_video")
def _():
assert all(r.summary.video_final_recovery_fraction == 1.0 for r in results if r.summary.channel_kbps == 300.0)
@add("15_overload_visible")
def _():
low = [r.summary for r in results if r.summary.channel_kbps == 230.0]
assert any(r.queue_at_source_end_packets > 0 and r.additional_drain_seconds > 0 for r in low)
@add("16_exact_transmission_duration")
def _():
assert abs(transmission_duration_seconds(123, 230_000) - 984 / 230_000) < 1e-15
@add("17_total_bits_accounting")
def _():
for result in results:
assert result.summary.transmitted_bytes * 8 == sum(len(x.wire_packet) * 8 for x in result.schedule.transmitted)
@add("18_reproducible_schedule")
def _():
original = results[0].schedule
repeated = schedule_packets(workload.packets, original.mode, original.channel_bitrate_bps)
assert [(x.queued.arrival_order, x.start_seconds, x.end_seconds) for x in original.transmitted] == [(x.queued.arrival_order, x.start_seconds, x.end_seconds) for x in repeated.transmitted]
@add("19_no_video_control_isolated")
def _():
items = tuple(_synthetic_item(TrafficClass.CONTROL, index, index * CONTROL_PERIOD_US, index) for index in range(20))
result = schedule_packets(items, SchedulerMode.STRICT_PRIORITY, 300_000)
assert max(x.end_seconds - x.queued.generation_time_seconds for x in result.transmitted) < 0.1
@add("20_latest_state_reaches_receiver")
def _():
latest = next(r for r in results if r.summary.channel_kbps == 230.0 and r.summary.scheduler == SchedulerMode.LATEST_STATE.value)
for cls in (TrafficClass.CONTROL, TrafficClass.TELEMETRY):
created_last = max(item.packet.sequence_number for item in _class_items(workload, cls))
sent_last = max(item.queued.packet.sequence_number for item in latest.schedule.transmitted if item.queued.packet.traffic_class is cls)
assert sent_last == created_last
output = []
for name, function in checks:
try:
function()
output.append(FunctionalTestResult(name, True, "PASS"))
except Exception as error:
output.append(FunctionalTestResult(name, False, f"{type(error).__name__}: {error}"))
if not all(item.passed for item in output):
failed = ", ".join(item.name for item in output if not item.passed)
raise AssertionError(f"functional checks failed: {failed}")
return tuple(output)
def save_csv(results: tuple[ScenarioResult, ...]) -> None:
OUTPUT_DIRECTORY.mkdir(parents=True, exist_ok=True)
with SUMMARY_CSV_PATH.open("w", newline="", encoding="utf-8") as file:
writer = csv.DictWriter(file, fieldnames=list(SimulationSummary.__dataclass_fields__))
writer.writeheader()
writer.writerows(asdict(result.summary) for result in results)
with CLASS_CSV_PATH.open("w", newline="", encoding="utf-8") as file:
writer = csv.DictWriter(file, fieldnames=list(ClassMetrics.__dataclass_fields__))
writer.writeheader()
writer.writerows(asdict(row) for result in results for row in result.classes)
def _grouped_plot(
results: tuple[ScenarioResult, ...],
values,
ylabel: str,
title: str,
path: Path,
) -> None:
x = np.arange(len(CHANNEL_RATES_KBPS)); width = 0.24
fig, axis = plt.subplots(figsize=(10, 5.5))
for index, mode in enumerate(SCHEDULER_MODES):
subset = [r for r in results if r.summary.scheduler == mode.value]
axis.bar(x + (index - 1) * width, [values(r.summary) for r in subset], width, label=MODE_NAMES[mode])
axis.set_xticks(x, [f"{rate:.0f}" for rate in CHANNEL_RATES_KBPS])
axis.set_xlabel("Скорость общего канала, кбит/с")
axis.set_ylabel(ylabel); axis.set_title(title); axis.grid(axis="y", alpha=0.3); axis.legend()
fig.tight_layout(); fig.savefig(path, dpi=150); plt.close(fig)
def save_plots(results: tuple[ScenarioResult, ...]) -> None:
_grouped_plot(results, lambda x: x.control_p95_age_ms, "P95 полной задержки, мс", "Задержка обычных команд", CONTROL_PLOT_PATH)
_grouped_plot(results, lambda x: x.emergency_total_delay_ms, "Полная задержка, мс", "Задержка аварийной команды", EMERGENCY_PLOT_PATH)
_grouped_plot(results, lambda x: x.video_p95_publication_delay_ms, "P95 задержки, мс", "Задержка публикации видео", VIDEO_PLOT_PATH)
_grouped_plot(results, lambda x: x.max_queue_packets, "Максимум, пакетов", "Длина общей очереди", QUEUE_PLOT_PATH)
_grouped_plot(results, lambda x: sum(row.deadline_misses for row in next(r.classes for r in results if r.summary is x)), "Число превышений", "Превышения допустимых сроков", DEADLINE_PLOT_PATH)
fig, axes = plt.subplots(1, 2, figsize=(12, 5))
for mode in SCHEDULER_MODES:
subset = [r.summary for r in results if r.summary.scheduler == mode.value]
axes[0].plot(CHANNEL_RATES_KBPS, [r.control_p95_age_ms for r in subset], marker="o", label=MODE_NAMES[mode])
axes[1].plot(CHANNEL_RATES_KBPS, [r.additional_drain_seconds for r in subset], marker="o", label=MODE_NAMES[mode])
axes[0].set_ylabel("P95 команд, мс"); axes[1].set_ylabel("Освобождение после видео, с")
for axis in axes:
axis.set_xlabel("Скорость, кбит/с"); axis.grid(alpha=0.3); axis.invert_xaxis()
axes[0].legend(); fig.suptitle("Сравнение планировщиков"); fig.tight_layout(); fig.savefig(COMPARISON_PLOT_PATH, dpi=150); plt.close(fig)
def write_report(
workload: Workload,
results: tuple[ScenarioResult, ...],
tests: tuple[FunctionalTestResult, ...],
) -> None:
first = results[0].summary
git_root = subprocess.run(
("git", "rev-parse", "--show-toplevel"),
check=True,
capture_output=True,
text=True,
encoding="utf-8",
).stdout.strip()
git_branch = subprocess.run(
("git", "branch", "--show-current"),
check=True,
capture_output=True,
text=True,
encoding="utf-8",
).stdout.strip()
git_head = subprocess.run(
("git", "rev-parse", "HEAD"),
check=True,
capture_output=True,
text=True,
encoding="utf-8",
).stdout.strip()
git_status = subprocess.run(
("git", "status", "--short", "--branch"),
check=True,
capture_output=True,
text=True,
encoding="utf-8",
).stdout.rstrip()
external_video_kbps = (
sum(len(unit.wire_packet) for unit in workload.video_units)
* 8.0
/ workload.metadata.duration_seconds
/ 1000.0
)
lines = [
"Lab033. Общий пакет канала и приоритетное обслуживание",
"",
"1. Git и конфигурация",
f"- Git-корень: {git_root}.",
f"- Ветка: {git_branch}; HEAD: {git_head}.",
f"- Исходное видео: {workload.metadata.duration_seconds:.6f} с, {len(workload.composites)} составных кадров, {COMPOSITE_FPS:.0f} кадра/с.",
"- BASE 240x135 grayscale JPEG Q23; ROI 320x180 grayscale JPEG Q33.",
"- Lab028: payload 512 байт; Lab030: FEC 12+3; Lab031: D=1.",
f"- Фактический внешний видеопоток до Lab033: {external_video_kbps:.3f} кбит/с; с общим заголовком: {first.video_load_kbps:.3f} кбит/с.",
f"- Аварийная команда: {first.emergency_load_kbps:.6f}; обычные команды: {first.control_load_kbps:.3f}; телеметрия: {first.telemetry_load_kbps:.3f} кбит/с.",
f"- Общая предложенная нагрузка: {first.offered_load_kbps:.3f} кбит/с; суммарная служебная доля: {first.service_share_percent:.3f}%.",
"",
"2. Формат общего заголовка",
f"- struct: {HEADER_FORMAT}; размер: {HEADER_SIZE} байта; network byte order, неявного padding нет.",
"- Смещения: magic[4]@0, version:u8@4, traffic_class:u8@5, direction:u8@6, flags:u8@7, stream_id:u16@8, sequence_number:u32@10, generation_time_us:u64@14, deadline_ms:u16@22, payload_length:u32@24, packet_crc32:u32@28.",
"- traffic_class: 1=EMERGENCY, 2=CONTROL, 3=TELEMETRY, 4=VIDEO; direction: 1=GROUND_TO_ROVER, 2=ROVER_TO_GROUND; flags в версии 1 равен нулю.",
"- CRC32 защищает заголовок с нулевым packet_crc32 и всю полезную нагрузку.",
"",
"3. Девять сочетаний скорости и планировщика",
"speed | scheduler | offered/capacity | control P95 ms | emergency ms | telemetry P95 ms | video P95 ms | queue@end packets | drain extra s | final video",
]
for result in results:
row = result.summary
lines.append(
f"{row.channel_kbps:.0f} | {MODE_NAMES[SchedulerMode(row.scheduler)]} | {row.offered_to_capacity_ratio:.4f} | {row.control_p95_age_ms:.3f} | {row.emergency_total_delay_ms:.3f} | {row.telemetry_p95_age_ms:.3f} | {row.video_p95_publication_delay_ms:.3f} | {row.queue_at_source_end_packets} | {row.additional_drain_seconds:.6f} | {row.video_final_recovery_fraction:.6f}"
)
lines.extend([
"",
"4. Подробные показатели",
])
for result in results:
row = result.summary
class_rows = {item.traffic_class: item for item in result.classes}
lines.extend([
f"[{row.channel_kbps:.0f} кбит/с; {MODE_NAMES[SchedulerMode(row.scheduler)]}]",
f"- Передано {row.transmitted_packets} пакетов, {row.transmitted_bytes} байт; средняя/максимальная очередь {row.mean_queue_packets:.3f}/{row.max_queue_packets} пакетов и {row.mean_queue_bytes:.1f}/{row.max_queue_bytes} байт.",
f"- На конце видео: {row.queue_at_source_end_packets} пакетов, {row.queue_at_source_end_bytes:.1f} байт; полное освобождение {row.drain_end_seconds:.6f} с, дополнительно {row.additional_drain_seconds:.6f} с.",
f"- Команды: mean/P95/max {row.control_mean_age_ms:.3f}/{row.control_p95_age_ms:.3f}/{row.control_max_age_ms:.3f} мс; превышений 100 мс {class_rows['control'].deadline_misses}; промежутков >100 мс {row.control_gaps_over_100ms}; max gap {row.control_max_receive_gap_ms:.3f} мс; заменено {row.control_replaced}.",
f"- Аварийная команда: start/total {row.emergency_start_delay_ms:.3f}/{row.emergency_total_delay_ms:.3f} мс; deadline 50 мс {'соблюдён' if row.emergency_deadline_met else 'НАРУШЕН'}; блокировал {row.emergency_blocker_class}, {row.emergency_blocker_size_bytes} байт, seq={row.emergency_blocker_sequence}, ожидание {row.emergency_blocking_delay_ms:.3f} мс.",
f"- Телеметрия: mean/P95/max {row.telemetry_mean_age_ms:.3f}/{row.telemetry_p95_age_ms:.3f}/{row.telemetry_max_age_ms:.3f} мс; превышений 500 мс {row.telemetry_misses_500ms} ({row.telemetry_miss_fraction:.6f}); заменено {row.telemetry_replaced}.",
f"- Видео: опубликовано к концу источника {row.video_published_at_source_end_fraction:.6f}, после drain {row.video_published_after_drain_fraction:.6f}, итог {row.video_final_recovery_fraction:.6f}; delay mean/P95/max {row.video_mean_publication_delay_ms:.3f}/{row.video_p95_publication_delay_ms:.3f}/{row.video_max_publication_delay_ms:.3f} мс.",
f"- Возраст изображения mean/P95/max {row.video_mean_display_age_ms:.3f}/{row.video_p95_display_age_ms:.3f}/{row.video_max_display_age_ms:.3f} мс; доля >1 с {row.video_age_over_1s_fraction:.6f}; поздних кадров {row.video_late_frame_count}, max серия {row.video_max_late_run}.",
])
lines.extend([
"",
"5. Интерпретация планировщиков",
"- FIFO задерживает команды, потому что все ранее поступившие внешние видеопакеты остаются перед ними независимо от срочности.",
"- Строгий приоритет после каждого окончания пакета выбирает управление раньше телеметрии и видео, поэтому сокращает очередь команд.",
"- Уже начатая передача не прерывается: модель сериализует целый пакет как атомарную единицу и не задаёт фрагментацию общего пакета или возобновление передачи.",
"- Поэтому нижняя граница задержки срочной команды включает остаток времени самого длинного пакета, уже занявшего ресурс.",
"- Замена старого состояния новым не тратит канал на запоздалую команду, которая к моменту приёма уже не отражает актуальное управление.",
"- Строгий приоритет перераспределяет порядок, но не уменьшает объём видео; при нагрузке выше пропускной способности видеопакеты продолжают накапливаться.",
"- Длинная очередь видео увеличивает возраст изображения даже без потерь: полные, но старые кадры публикуются позже времени формирования.",
"- Восстановление 100% кадров после drain подтверждает целостность, но не пригодность задержек канала для управления в реальном времени.",
"",
"6. Допущения модели",
"- Один общий абстрактный ресурс и одна расчётная скорость для обоих направлений; переключение направления не задерживает передачу.",
"- Не выбран физический полудуплекс/дуплекс, TDD/FDD, модуляция, ретранслятор или физическое разделение восходящего и нисходящего каналов.",
"- Ошибки и помехи отсутствуют; полностью переданный пакет доставляется без повреждений; видеопакеты не отбрасываются по возрасту.",
"- После конца исходного видео новые данные не создаются, а очередь полностью освобождается.",
"- При совпадении времени FIFO использует traffic_class, stream_id и sequence_number как устойчивое правило разрешения совпадений.",
"- Средняя очередь рассчитана по времени на основном интервале и включает пакет, находящийся в передатчике; остаток байтов активного пакета учитывается пропорционально недопереданному времени.",
"",
"7. Функциональные проверки",
])
lines.extend(f"- {'PASS' if item.passed else 'FAIL'} {item.name}: {item.detail}" for item in tests)
lines.extend([
"",
"8. Созданные файлы и итоговый Git status",
"- protocol/link_packet.py",
"- protocol/priority_scheduler.py",
"- experiments/lab033_priority_channel_scheduler.py",
"- data/processed/lab033/lab033_summary.csv",
"- data/processed/lab033/lab033_class_metrics.csv",
"- data/processed/lab033/lab033_report.txt",
])
lines.extend(f"- data/processed/lab033/{path.name}" for path in PLOT_PATHS)
lines.extend([
"- Lab028-Lab032 и их сохранённые результаты не изменены.",
"- Lab033 не добавлена в индекс и не закоммичена.",
"",
git_status,
])
REPORT_PATH.write_text("\n".join(lines) + "\n", encoding="utf-8")
def validate_outputs() -> None:
for path, expected_rows in ((SUMMARY_CSV_PATH, 9), (CLASS_CSV_PATH, 36)):
with path.open("r", newline="", encoding="utf-8") as file:
rows = list(csv.DictReader(file))
if len(rows) != expected_rows or not rows or any(set(rows[0]) != set(row) for row in rows):
raise AssertionError(f"invalid CSV output: {path}")
text = REPORT_PATH.read_text(encoding="utf-8")
if "Lab033" not in text or "Функциональные проверки" not in text:
raise AssertionError("report validation failed")
for path in PLOT_PATHS:
image = cv2.imread(str(path), cv2.IMREAD_UNCHANGED)
if image is None or image.size == 0:
raise AssertionError(f"OpenCV could not read {path}")
def main() -> None:
workload = build_workload()
results = run_experiment(workload)
tests = run_functional_tests(workload, results)
save_csv(results)
save_plots(results)
write_report(workload, results, tests)
validate_outputs()
print(
f"Lab033 complete: {len(results)} scenarios, "
f"{len(tests)} functional checks, "
f"{len(workload.composites)} video frames"
)
if __name__ == "__main__":
main()