"""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()