diff --git a/PROJECT_LOG.md b/PROJECT_LOG.md index d08cc36..a531046 100644 --- a/PROJECT_LOG.md +++ b/PROJECT_LOG.md @@ -163,3 +163,28 @@ git version 2.51.1.windows.1 - Режим 8+3 не принят из-за отсутствия устойчивого преимущества. - Подтверждено, что продолжительные помехи нельзя компенсировать только увеличением числа проверочных пакетов. - Следующий этап посвящён совместной передаче команд, телеметрии и видео с приоритетным обслуживанием. + +--- + +# Запись 008 + +## Дата + +3 августа 2026 года + +## Тема + +Завершение Lab033: общий пакет канала и приоритетное обслуживание. + +## Выполнено + +- Добавлен общий 32-байтный заголовок канала для команд, телеметрии и видео. +- Общая предложенная нагрузка составляет 258,214 кбит/с. +- Подтверждено, что FIFO не обеспечивает требуемую задержку команд. +- Выбран строгий приоритет с заменой устаревших состояний. +- Аварийная команда имеет наивысший приоритет. +- Скорость 260 кбит/с практически полностью загружена. +- При 230 кбит/с видеоданные накапливаются и устаревают. +- Итоговое восстановление всех кадров после освобождения очереди не означает работу видео в реальном времени. +- Расчёты и графики выполнены с Matplotlib 3.11.1. +- Следующий этап посвящён отбрасыванию устаревших видеокадров. diff --git a/data/processed/lab033/lab033_class_metrics.csv b/data/processed/lab033/lab033_class_metrics.csv new file mode 100644 index 0000000..d382f7a --- /dev/null +++ b/data/processed/lab033/lab033_class_metrics.csv @@ -0,0 +1,37 @@ +channel_kbps,scheduler,traffic_class,offered_load_kbps,created_packets,transmitted_packets,replaced_packets,remaining_at_source_end_packets,mean_start_delay_ms,p95_start_delay_ms,max_start_delay_ms,mean_total_delay_ms,p95_total_delay_ms,max_total_delay_ms,deadline_misses,deadline_miss_fraction,mean_blocking_delay_ms,max_blocking_delay_ms +300.0,fifo,emergency,0.02465489566613162,1,1,0,0,0.0,0.0,0.0,1.7066666666671892,1.7066666666671892,1.7066666666671892,0,0.0,0.0,0.0 +300.0,fifo,control,10.256436597110755,416,416,0,2,104.3649927884586,236.92691666666235,294.47966666665957,106.0716594551255,238.6335833333294,296.1863333333259,203,0.4879807692307692,5.879864583331442,16.212999999996036 +300.0,fifo,telemetry,7.692327447833066,208,208,0,1,98.47181249999751,231.12033333332866,294.26699999999425,101.03181249999754,233.68033333332914,296.8269999999933,0,0.0,5.797325320511048,16.212999999996036 +300.0,fifo,video,240.24038523274479,1078,1078,0,15,128.3044211502762,254.33199999999601,309.6799999999931,143.7311125541102,264.7599999999962,325.89333333332604,0,0.0,0.0026771799628383768,0.16033333332998723 +300.0,strict_priority,emergency,0.02465489566613162,1,1,0,0,0.0,0.0,0.0,1.7066666666671892,1.7066666666671892,1.7066666666671892,0,0.0,0.0,0.0 +300.0,strict_priority,control,10.256436597110755,416,416,0,0,7.733262019229404,15.759666666666755,16.106333333334,9.439928685896334,17.46633333333346,17.813000000000745,0,0.0,7.7291594551268386,16.106333333334 +300.0,strict_priority,telemetry,7.692327447833066,208,208,0,0,9.864120192306657,17.546883333334474,17.813000000000745,12.424120192306697,20.10688333333437,20.37300000000064,0,0.0,8.149248397434599,16.106333333334 +300.0,strict_priority,video,240.24038523274479,1078,1078,0,15,135.37173098330103,269.12799999999464,323.36000000000143,150.79842238713505,279.70533333332844,333.4933333333332,0,0.0,0.0026771799628976987,0.16033333333353994 +300.0,latest_state,emergency,0.02465489566613162,1,1,0,0,0.0,0.0,0.0,1.7066666666671892,1.7066666666671892,1.7066666666671892,0,0.0,0.0,0.0 +300.0,latest_state,control,10.256436597110755,416,416,0,0,7.733262019229404,15.759666666666755,16.106333333334,9.439928685896334,17.46633333333346,17.813000000000745,0,0.0,7.7291594551268386,16.106333333334 +300.0,latest_state,telemetry,7.692327447833066,208,208,0,0,9.864120192306657,17.546883333334474,17.813000000000745,12.424120192306697,20.10688333333437,20.37300000000064,0,0.0,8.149248397434599,16.106333333334 +300.0,latest_state,video,240.24038523274479,1078,1078,0,15,135.37173098330103,269.12799999999464,323.36000000000143,150.79842238713505,279.70533333332844,333.4933333333332,0,0.0,0.0026771799628976987,0.16033333333353994 +260.0,fifo,emergency,0.02465489566613162,1,1,0,0,0.0,0.0,0.0,1.969230769230279,1.969230769230279,1.969230769230279,0,0.0,0.0,0.0 +260.0,fifo,control,10.256436597110755,416,416,0,7,167.10386353552764,331.0900769231244,458.97469230786925,169.07309430475814,333.0593076923547,460.9439230770995,298,0.7163461538461539,7.857339866886274,18.615384615390695 +260.0,fifo,telemetry,7.692327447833066,208,208,0,4,160.36350517753885,328.8285384616073,460.9439230770995,163.31735133138503,331.7823846154539,463.8977692309467,0,0.0,7.732144230791379,18.615384615390695 +260.0,fifo,video,240.24038523274479,1078,1078,0,21,175.94024482662746,343.1754384616482,473.6000000001752,193.74027336951343,357.6154384615683,492.3076923078682,0,0.0,3.3139568289080654,18.257076923131166 +260.0,strict_priority,emergency,0.02465489566613162,1,1,0,0,0.0,0.0,0.0,1.969230769230279,1.969230769230279,1.969230769230279,0,0.0,0.0,0.0 +260.0,strict_priority,control,10.256436597110755,416,416,0,0,8.313035133156095,16.784365384612144,18.69776923077815,10.282265902386632,18.753596153843198,20.66700000000843,0,0.0,8.308301405345444,18.69776923077815 +260.0,strict_priority,telemetry,7.692327447833066,208,208,0,0,10.909659023688215,19.493946153856623,20.66700000000843,13.863505177534366,22.44779230770204,23.620846153853847,0,0.0,8.930960798836377,18.69776923077815 +260.0,strict_priority,video,240.24038523274479,1078,1078,0,23,187.65972534610552,365.73226153856297,480.49230769247583,205.4597538889915,380.5630769230985,499.20000000016887,0,0.0,3.9965313971947523,18.25707692318801 +260.0,latest_state,emergency,0.02465489566613162,1,1,0,0,0.0,0.0,0.0,1.969230769230279,1.969230769230279,1.969230769230279,0,0.0,0.0,0.0 +260.0,latest_state,control,10.256436597110755,416,416,0,0,8.313035133156095,16.784365384612144,18.69776923077815,10.282265902386632,18.753596153843198,20.66700000000843,0,0.0,8.308301405345444,18.69776923077815 +260.0,latest_state,telemetry,7.692327447833066,208,208,0,0,10.909659023688215,19.493946153856623,20.66700000000843,13.863505177534366,22.44779230770204,23.620846153853847,0,0.0,8.930960798836377,18.69776923077815 +260.0,latest_state,video,240.24038523274479,1078,1078,0,23,187.65972534610552,365.73226153856297,480.49230769247583,205.4597538889915,380.5630769230985,499.20000000016887,0,0.0,3.9965313971947523,18.25707692318801 +230.0,fifo,emergency,0.02465489566613162,1,1,0,0,878.8869565216775,878.8869565216775,878.8869565216775,881.1130434781998,881.1130434781998,881.1130434781998,1,1.0,0.4173913042819777,0.4173913042819777 +230.0,fifo,control,10.256436597110755,416,416,0,42,1172.0617892976043,2238.3086956520533,2606.2956521737738,1174.2878762541263,2240.5347826085754,2608.5217391302963,414,0.9951923076923077,9.870735785906158,21.13043478255605 +230.0,fifo,telemetry,7.692327447833066,208,208,0,21,1161.6391304347285,2222.267826086846,2608.5217391302963,1164.978260869511,2225.6069565216276,2611.860869565078,177,0.8509615384615384,9.380434782561796,21.043478260802218 +230.0,fifo,video,240.24038523274479,1078,1078,0,118,1186.612664354224,2267.033969565095,2618.4808260868167,1206.7344357505297,2288.1817956520513,2639.628652173773,0,0.0,8.20556586266141,21.136565217338088 +230.0,strict_priority,emergency,0.02465489566613162,1,1,0,0,9.32173913036749,9.32173913036749,9.32173913036749,11.547826086889756,11.547826086889756,11.547826086889756,0,0.0,9.32173913036749,9.32173913036749 +230.0,strict_priority,control,10.256436597110755,416,416,0,1,10.02424749162144,19.62173913038057,21.113043478251825,12.250334448143445,21.84782608690261,23.33913043477409,0,0.0,10.018896321052877,21.113043478251825 +230.0,strict_priority,telemetry,7.692327447833066,208,208,0,0,12.175083612023139,21.772173913030322,23.33913043477409,15.51421404680579,25.111304347813142,26.678260869555714,0,0.0,9.938294314364006,21.113043478251825 +230.0,strict_priority,video,240.24038523274479,1078,1078,0,126,1273.2989271597767,2364.9321739130005,2626.272130434721,1293.4206985560818,2379.542608695629,2647.4199565216772,0,0.0,9.628813422585635,19.802565217331036 +230.0,latest_state,emergency,0.02465489566613162,1,1,0,0,9.32173913036749,9.32173913036749,9.32173913036749,11.547826086889756,11.547826086889756,11.547826086889756,0,0.0,9.32173913036749,9.32173913036749 +230.0,latest_state,control,10.256436597110755,416,416,0,1,10.02424749162144,19.62173913038057,21.113043478251825,12.250334448143445,21.84782608690261,23.33913043477409,0,0.0,10.018896321052877,21.113043478251825 +230.0,latest_state,telemetry,7.692327447833066,208,208,0,0,12.175083612023139,21.772173913030322,23.33913043477409,15.51421404680579,25.111304347813142,26.678260869555714,0,0.0,9.938294314364006,21.113043478251825 +230.0,latest_state,video,240.24038523274479,1078,1078,0,126,1273.2989271597767,2364.9321739130005,2626.272130434721,1293.4206985560818,2379.542608695629,2647.4199565216772,0,0.0,9.628813422585635,19.802565217331036 diff --git a/data/processed/lab033/lab033_control_delay.png b/data/processed/lab033/lab033_control_delay.png new file mode 100644 index 0000000..9443611 Binary files /dev/null and b/data/processed/lab033/lab033_control_delay.png differ diff --git a/data/processed/lab033/lab033_deadline_misses.png b/data/processed/lab033/lab033_deadline_misses.png new file mode 100644 index 0000000..e58e19a Binary files /dev/null and b/data/processed/lab033/lab033_deadline_misses.png differ diff --git a/data/processed/lab033/lab033_emergency_delay.png b/data/processed/lab033/lab033_emergency_delay.png new file mode 100644 index 0000000..4b0f487 Binary files /dev/null and b/data/processed/lab033/lab033_emergency_delay.png differ diff --git a/data/processed/lab033/lab033_queue_length.png b/data/processed/lab033/lab033_queue_length.png new file mode 100644 index 0000000..5956749 Binary files /dev/null and b/data/processed/lab033/lab033_queue_length.png differ diff --git a/data/processed/lab033/lab033_report.txt b/data/processed/lab033/lab033_report.txt new file mode 100644 index 0000000..fd9bc6a --- /dev/null +++ b/data/processed/lab033/lab033_report.txt @@ -0,0 +1,165 @@ +Lab033. Общий пакет канала и приоритетное обслуживание + +1. Git и конфигурация +- Git-корень: C:/Users/user/Desktop/projects/SDR_Rover. +- Ветка: main; HEAD: 3dbb7afa8332b49535d867f1eabca7bdd2b7cdde. +- Исходное видео: 20.766667 с, 63 составных кадров, 3 кадра/с. +- BASE 240x135 grayscale JPEG Q23; ROI 320x180 grayscale JPEG Q33. +- Lab028: payload 512 байт; Lab030: FEC 12+3; Lab031: D=1. +- Фактический внешний видеопоток до Lab033: 226.951 кбит/с; с общим заголовком: 240.240 кбит/с. +- Аварийная команда: 0.024655; обычные команды: 10.256; телеметрия: 7.692 кбит/с. +- Общая предложенная нагрузка: 258.214 кбит/с; суммарная служебная доля: 34.923%. + +2. Формат общего заголовка +- struct: !4sBBBBHIQHII; размер: 32 байта; 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 +300 | FIFO | 0.8607 | 238.634 | 1.707 | 233.680 | 304.579 | 18 | 0.231867 | 1.000000 +300 | Строгий приоритет | 0.8607 | 17.466 | 1.707 | 20.107 | 322.072 | 15 | 0.231867 | 1.000000 +300 | Приоритет + замена | 0.8607 | 17.466 | 1.707 | 20.107 | 322.072 | 15 | 0.231867 | 1.000000 +260 | FIFO | 0.9931 | 333.059 | 1.969 | 331.782 | 420.025 | 32 | 0.399200 | 1.000000 +260 | Строгий приоритет | 0.9931 | 18.754 | 1.969 | 22.448 | 445.846 | 23 | 0.399200 | 1.000000 +260 | Приоритет + замена | 0.9931 | 18.754 | 1.969 | 22.448 | 445.846 | 23 | 0.399200 | 1.000000 +230 | FIFO | 1.1227 | 2240.535 | 881.113 | 2225.607 | 2417.731 | 181 | 2.547420 | 1.000000 +230 | Строгий приоритет | 1.1227 | 21.848 | 11.548 | 25.111 | 2505.995 | 127 | 2.547420 | 1.000000 +230 | Приоритет + замена | 1.1227 | 21.848 | 11.548 | 25.111 | 2505.995 | 127 | 2.547420 | 1.000000 + +4. Подробные показатели +[300 кбит/с; FIFO] +- Передано 1703 пакетов, 670280 байт; средняя/максимальная очередь 10.478/22 пакетов и 4452.4/12221 байт. +- На конце видео: 18 пакетов, 8695.0 байт; полное освобождение 20.998534 с, дополнительно 0.231867 с. +- Команды: mean/P95/max 106.072/238.634/296.186 мс; превышений 100 мс 203; промежутков >100 мс 63; max gap 344.480 мс; заменено 0. +- Аварийная команда: start/total 0.000/1.707 мс; deadline 50 мс соблюдён; блокировал none, 0 байт, seq=-1, ожидание 0.000 мс. +- Телеметрия: mean/P95/max 101.032/233.680/296.827 мс; превышений 500 мс 0 (0.000000); заменено 0. +- Видео: опубликовано к концу источника 0.984127, после drain 1.000000, итог 1.000000; delay mean/P95/max 261.559/304.579/311.146 мс. +- Возраст изображения mean/P95/max 423.846/580.000/643.333 мс; доля >1 с 0.000000; поздних кадров 0, max серия 0. +[300 кбит/с; Строгий приоритет] +- Передано 1703 пакетов, 670280 байт; средняя/максимальная очередь 8.050/22 пакетов и 4452.4/12221 байт. +- На конце видео: 15 пакетов, 8695.0 байт; полное освобождение 20.998534 с, дополнительно 0.231867 с. +- Команды: mean/P95/max 9.440/17.466/17.813 мс; превышений 100 мс 0; промежутков >100 мс 0; max gap 65.760 мс; заменено 0. +- Аварийная команда: start/total 0.000/1.707 мс; deadline 50 мс соблюдён; блокировал none, 0 байт, seq=-1, ожидание 0.000 мс. +- Телеметрия: mean/P95/max 12.424/20.107/20.373 мс; превышений 500 мс 0 (0.000000); заменено 0. +- Видео: опубликовано к концу источника 0.984127, после drain 1.000000, итог 1.000000; delay mean/P95/max 276.337/322.072/330.773 мс. +- Возраст изображения mean/P95/max 437.649/593.333/660.000 мс; доля >1 с 0.000000; поздних кадров 0, max серия 0. +[300 кбит/с; Приоритет + замена] +- Передано 1703 пакетов, 670280 байт; средняя/максимальная очередь 8.050/22 пакетов и 4452.4/12221 байт. +- На конце видео: 15 пакетов, 8695.0 байт; полное освобождение 20.998534 с, дополнительно 0.231867 с. +- Команды: mean/P95/max 9.440/17.466/17.813 мс; превышений 100 мс 0; промежутков >100 мс 0; max gap 65.760 мс; заменено 0. +- Аварийная команда: start/total 0.000/1.707 мс; deadline 50 мс соблюдён; блокировал none, 0 байт, seq=-1, ожидание 0.000 мс. +- Телеметрия: mean/P95/max 12.424/20.107/20.373 мс; превышений 500 мс 0 (0.000000); заменено 0. +- Видео: опубликовано к концу источника 0.984127, после drain 1.000000, итог 1.000000; delay mean/P95/max 276.337/322.072/330.773 мс. +- Возраст изображения mean/P95/max 437.649/593.333/660.000 мс; доля >1 с 0.000000; поздних кадров 0, max серия 0. +[260 кбит/с; FIFO] +- Передано 1703 пакетов, 670280 байт; средняя/максимальная очередь 14.800/37 пакетов и 6010.5/16512 байт. +- На конце видео: 32 пакетов, 12974.0 байт; полное освобождение 21.165867 с, дополнительно 0.399200 с. +- Команды: mean/P95/max 169.073/333.059/460.944 мс; превышений 100 мс 298; промежутков >100 мс 63; max gap 384.348 мс; заменено 0. +- Аварийная команда: start/total 0.000/1.969 мс; deadline 50 мс соблюдён; блокировал none, 0 байт, seq=-1, ожидание 0.000 мс. +- Телеметрия: mean/P95/max 163.317/331.782/463.898 мс; превышений 500 мс 0 (0.000000); заменено 0. +- Видео: опубликовано к концу источника 0.984127, после drain 1.000000, итог 1.000000; delay mean/P95/max 329.072/420.025/436.185 мс. +- Возраст изображения mean/P95/max 489.627/666.667/766.667 мс; доля >1 с 0.000000; поздних кадров 0, max серия 0. +[260 кбит/с; Строгий приоритет] +- Передано 1703 пакетов, 670280 байт; средняя/максимальная очередь 10.783/28 пакетов и 6010.2/16384 байт. +- На конце видео: 23 пакетов, 12974.0 байт; полное освобождение 21.165867 с, дополнительно 0.399200 с. +- Команды: mean/P95/max 10.282/18.754/20.667 мс; превышений 100 мс 0; промежутков >100 мс 0; max gap 63.477 мс; заменено 0. +- Аварийная команда: start/total 0.000/1.969 мс; deadline 50 мс соблюдён; блокировал none, 0 байт, seq=-1, ожидание 0.000 мс. +- Телеметрия: mean/P95/max 13.864/22.448/23.621 мс; превышений 500 мс 0 (0.000000); заменено 0. +- Видео: опубликовано к концу источника 0.968254, после drain 1.000000, итог 1.000000; delay mean/P95/max 351.218/445.846/459.559 мс. +- Возраст изображения mean/P95/max 511.188/693.333/783.333 мс; доля >1 с 0.000000; поздних кадров 0, max серия 0. +[260 кбит/с; Приоритет + замена] +- Передано 1703 пакетов, 670280 байт; средняя/максимальная очередь 10.783/28 пакетов и 6010.2/16384 байт. +- На конце видео: 23 пакетов, 12974.0 байт; полное освобождение 21.165867 с, дополнительно 0.399200 с. +- Команды: mean/P95/max 10.282/18.754/20.667 мс; превышений 100 мс 0; промежутков >100 мс 0; max gap 63.477 мс; заменено 0. +- Аварийная команда: start/total 0.000/1.969 мс; deadline 50 мс соблюдён; блокировал none, 0 байт, seq=-1, ожидание 0.000 мс. +- Телеметрия: mean/P95/max 13.864/22.448/23.621 мс; превышений 500 мс 0 (0.000000); заменено 0. +- Видео: опубликовано к концу источника 0.968254, после drain 1.000000, итог 1.000000; delay mean/P95/max 351.218/445.846/459.559 мс. +- Возраст изображения mean/P95/max 511.188/693.333/783.333 мс; доля >1 с 0.000000; поздних кадров 0, max серия 0. +[230 кбит/с; FIFO] +- Передано 1703 пакетов, 670280 байт; средняя/максимальная очередь 86.915/184 пакетов и 34309.1/76497 байт. +- На конце видео: 181 пакетов, 73238.3 байт; полное освобождение 23.314087 с, дополнительно 2.547420 с. +- Команды: mean/P95/max 1174.288/2240.535/2608.522 мс; превышений 100 мс 414; промежутков >100 мс 63; max gap 427.304 мс; заменено 0. +- Аварийная команда: start/total 878.887/881.113 мс; deadline 50 мс НАРУШЕН; блокировал video, 608 байт, seq=461, ожидание 0.417 мс. +- Телеметрия: mean/P95/max 1164.978/2225.607/2611.861 мс; превышений 500 мс 177 (0.850962); заменено 0. +- Видео: опубликовано к концу источника 0.888889, после drain 1.000000, итог 1.000000; delay mean/P95/max 1347.217/2417.731/2576.186 мс. +- Возраст изображения mean/P95/max 1506.487/2588.167/2906.667 мс; доля >1 с 0.771870; поздних кадров 42, max серия 42. +[230 кбит/с; Строгий приоритет] +- Передано 1703 пакетов, 670280 байт; средняя/максимальная очередь 59.750/132 пакетов и 34307.1/75969 байт. +- На конце видео: 127 пакетов, 73238.3 байт; полное освобождение 23.314087 с, дополнительно 2.547420 с. +- Команды: mean/P95/max 12.250/21.848/23.339 мс; превышений 100 мс 0; промежутков >100 мс 0; max gap 69.009 мс; заменено 0. +- Аварийная команда: start/total 9.322/11.548 мс; deadline 50 мс соблюдён; блокировал video, 608 байт, seq=458, ожидание 9.322 мс. +- Телеметрия: mean/P95/max 15.514/25.111/26.678 мс; превышений 500 мс 0 (0.000000); заменено 0. +- Видео: опубликовано к концу источника 0.888889, после drain 1.000000, итог 1.000000; delay mean/P95/max 1445.995/2505.995/2583.977 мс. +- Возраст изображения mean/P95/max 1594.537/2678.167/2916.667 мс; доля >1 с 0.804460; поздних кадров 47, max серия 47. +[230 кбит/с; Приоритет + замена] +- Передано 1703 пакетов, 670280 байт; средняя/максимальная очередь 59.750/132 пакетов и 34307.1/75969 байт. +- На конце видео: 127 пакетов, 73238.3 байт; полное освобождение 23.314087 с, дополнительно 2.547420 с. +- Команды: mean/P95/max 12.250/21.848/23.339 мс; превышений 100 мс 0; промежутков >100 мс 0; max gap 69.009 мс; заменено 0. +- Аварийная команда: start/total 9.322/11.548 мс; deadline 50 мс соблюдён; блокировал video, 608 байт, seq=458, ожидание 9.322 мс. +- Телеметрия: mean/P95/max 15.514/25.111/26.678 мс; превышений 500 мс 0 (0.000000); заменено 0. +- Видео: опубликовано к концу источника 0.888889, после drain 1.000000, итог 1.000000; delay mean/P95/max 1445.995/2505.995/2583.977 мс. +- Возраст изображения mean/P95/max 1594.537/2678.167/2916.667 мс; доля >1 с 0.804460; поздних кадров 47, max серия 47. + +5. Интерпретация планировщиков +- FIFO задерживает команды, потому что все ранее поступившие внешние видеопакеты остаются перед ними независимо от срочности. +- Строгий приоритет после каждого окончания пакета выбирает управление раньше телеметрии и видео, поэтому сокращает очередь команд. +- Уже начатая передача не прерывается: модель сериализует целый пакет как атомарную единицу и не задаёт фрагментацию общего пакета или возобновление передачи. +- Поэтому нижняя граница задержки срочной команды включает остаток времени самого длинного пакета, уже занявшего ресурс. +- Замена старого состояния новым не тратит канал на запоздалую команду, которая к моменту приёма уже не отражает актуальное управление. +- Строгий приоритет перераспределяет порядок, но не уменьшает объём видео; при нагрузке выше пропускной способности видеопакеты продолжают накапливаться. +- Длинная очередь видео увеличивает возраст изображения даже без потерь: полные, но старые кадры публикуются позже времени формирования. +- Восстановление 100% кадров после drain подтверждает целостность, но не пригодность задержек канала для управления в реальном времени. + +6. Допущения модели +- Один общий абстрактный ресурс и одна расчётная скорость для обоих направлений; переключение направления не задерживает передачу. +- Не выбран физический полудуплекс/дуплекс, TDD/FDD, модуляция, ретранслятор или физическое разделение восходящего и нисходящего каналов. +- Ошибки и помехи отсутствуют; полностью переданный пакет доставляется без повреждений; видеопакеты не отбрасываются по возрасту. +- После конца исходного видео новые данные не создаются, а очередь полностью освобождается. +- При совпадении времени FIFO использует traffic_class, stream_id и sequence_number как устойчивое правило разрешения совпадений. +- Средняя очередь рассчитана по времени на основном интервале и включает пакет, находящийся в передатчике; остаток байтов активного пакета учитывается пропорционально недопереданному времени. + +7. Функциональные проверки +- PASS 01_header_roundtrip: PASS +- PASS 02_header_crc_detection: PASS +- PASS 03_payload_crc_detection: PASS +- PASS 04_video_wrapper_byte_exact: PASS +- PASS 05_lab030_crc: PASS +- PASS 06_lab028_crc: PASS +- PASS 07_fifo_global_order: PASS +- PASS 08_strict_priority: PASS +- PASS 09_non_preemptive: PASS +- PASS 10_intra_class_sequence: PASS +- PASS 11_latest_state_scope: PASS +- PASS 12_active_state_not_removed: PASS +- PASS 13_emergency_never_replaced: PASS +- PASS 14_300kbps_all_video: PASS +- PASS 15_overload_visible: PASS +- PASS 16_exact_transmission_duration: PASS +- PASS 17_total_bits_accounting: PASS +- PASS 18_reproducible_schedule: PASS +- PASS 19_no_video_control_isolated: PASS +- PASS 20_latest_state_reaches_receiver: PASS + +8. Созданные файлы и итоговый Git status +- protocol/link_packet.py +- protocol/priority_scheduler.py +- tests/lab033_priority_channel_scheduler.py +- data/processed/lab033/lab033_summary.csv +- data/processed/lab033/lab033_class_metrics.csv +- data/processed/lab033/lab033_report.txt +- data/processed/lab033/lab033_control_delay.png +- data/processed/lab033/lab033_emergency_delay.png +- data/processed/lab033/lab033_video_delay.png +- data/processed/lab033/lab033_queue_length.png +- data/processed/lab033/lab033_deadline_misses.png +- data/processed/lab033/lab033_scheduler_comparison.png +- Lab028-Lab032 и их сохранённые результаты не изменены. +- Lab033 не добавлена в индекс и не закоммичена. + +## main...origin/main +?? data/processed/lab033/ +?? protocol/link_packet.py +?? protocol/priority_scheduler.py +?? tests/lab033_priority_channel_scheduler.py diff --git a/data/processed/lab033/lab033_scheduler_comparison.png b/data/processed/lab033/lab033_scheduler_comparison.png new file mode 100644 index 0000000..aa1fda5 Binary files /dev/null and b/data/processed/lab033/lab033_scheduler_comparison.png differ diff --git a/data/processed/lab033/lab033_summary.csv b/data/processed/lab033/lab033_summary.csv new file mode 100644 index 0000000..1e7ec91 --- /dev/null +++ b/data/processed/lab033/lab033_summary.csv @@ -0,0 +1,10 @@ +channel_kbps,scheduler,offered_load_kbps,emergency_load_kbps,control_load_kbps,telemetry_load_kbps,video_load_kbps,service_share_percent,offered_to_capacity_ratio,transmitted_packets,transmitted_bytes,source_duration_seconds,drain_end_seconds,additional_drain_seconds,queue_at_source_end_packets,queue_at_source_end_bytes,mean_queue_packets,max_queue_packets,mean_queue_bytes,max_queue_bytes,control_mean_age_ms,control_p95_age_ms,control_max_age_ms,control_gaps_over_100ms,control_max_receive_gap_ms,control_replaced,emergency_start_delay_ms,emergency_total_delay_ms,emergency_deadline_met,emergency_blocker_class,emergency_blocker_size_bytes,emergency_blocker_sequence,emergency_blocking_delay_ms,telemetry_mean_age_ms,telemetry_p95_age_ms,telemetry_max_age_ms,telemetry_misses_500ms,telemetry_miss_fraction,telemetry_replaced,video_published_at_source_end_fraction,video_published_after_drain_fraction,video_final_recovery_fraction,video_mean_publication_delay_ms,video_p95_publication_delay_ms,video_max_publication_delay_ms,video_mean_display_age_ms,video_p95_display_age_ms,video_max_display_age_ms,video_age_over_1s_fraction,video_late_frame_count,video_max_late_run +300.0,fifo,258.21380417335473,0.02465489566613162,10.256436597110755,7.692327447833066,240.24038523274479,34.922718863758426,0.8607126805778491,1703,670280,20.766666666666666,20.99853366666666,0.23186699999999405,18,8695.012499999913,10.477643691813581,22,4452.409985633955,12221,106.0716594551255,238.6335833333294,296.1863333333259,63,344.4796666666594,0,0.0,1.7066666666671892,True,none,0,-1,0.0,101.03181249999754,233.68033333332914,296.8269999999933,0,0.0,0,0.9841269841269841,1.0,1.0,261.5593703703662,304.57866666666183,311.1463333333262,423.8457877201335,580.0,643.3333333333333,0.0,0,0 +300.0,strict_priority,258.21380417335473,0.02465489566613162,10.256436597110755,7.692327447833066,240.24038523274479,34.922718863758426,0.8607126805778491,1703,670280,20.766666666666666,20.99853366666666,0.23186699999999405,15,8695.012499999915,8.05016637239154,22,4452.387224542477,12221,9.439928685896334,17.46633333333346,17.813000000000745,0,65.7596666666671,0,0.0,1.7066666666671892,True,none,0,-1,0.0,12.424120192306697,20.10688333333437,20.37300000000064,0,0.0,0,0.9841269841269841,1.0,1.0,276.3369365079338,322.0720000000001,330.77299999999,437.6487386958594,593.3333333333334,660.0000000000001,0.0,0,0 +300.0,latest_state,258.21380417335473,0.02465489566613162,10.256436597110755,7.692327447833066,240.24038523274479,34.922718863758426,0.8607126805778491,1703,670280,20.766666666666666,20.99853366666666,0.23186699999999405,15,8695.012499999915,8.05016637239154,22,4452.387224542477,12221,9.439928685896334,17.46633333333346,17.813000000000745,0,65.7596666666671,0,0.0,1.7066666666671892,True,none,0,-1,0.0,12.424120192306697,20.10688333333437,20.37300000000064,0,0.0,0,0.9841269841269841,1.0,1.0,276.3369365079338,322.0720000000001,330.77299999999,437.6487386958594,593.3333333333334,660.0000000000001,0.0,0,0 +260.0,fifo,258.21380417335473,0.02465489566613162,10.256436597110755,7.692327447833066,240.24038523274479,34.922718863758426,0.9931300160513643,1703,670280,20.766666666666666,21.165867000000176,0.39920033333351057,32,12974.010833338622,14.800483013954153,37,6010.5496338357525,16512,169.07309430475814,333.0593076923547,460.9439230770995,63,384.34838461538055,0,0.0,1.969230769230279,True,none,0,-1,0.0,163.31735133138503,331.7823846154539,463.8977692309467,0,0.0,0,0.9841269841269841,1.0,1.0,329.07164590967284,420.0249487180816,436.18494871812175,489.62700661000974,666.6666666666679,766.6666666666692,0.0,0,0 +260.0,strict_priority,258.21380417335473,0.02465489566613162,10.256436597110755,7.692327447833066,240.24038523274479,34.922718863758426,0.9931300160513643,1703,670280,20.766666666666666,21.16586700000017,0.39920033333350347,23,12974.010833338338,10.783079034450305,28,6010.224630522928,16384,10.282265902386632,18.753596153843198,20.66700000000843,0,63.476923076922546,0,0.0,1.969230769230279,True,none,0,-1,0.0,13.863505177534366,22.44779230770204,23.620846153853847,0,0.0,0,0.9682539682539683,1.0,1.0,351.2176776557015,445.84648717956554,459.55930769245157,511.18822788794483,693.3333333333316,783.3333333333314,0.0,0,0 +260.0,latest_state,258.21380417335473,0.02465489566613162,10.256436597110755,7.692327447833066,240.24038523274479,34.922718863758426,0.9931300160513643,1703,670280,20.766666666666666,21.16586700000017,0.39920033333350347,23,12974.010833338338,10.783079034450305,28,6010.224630522928,16384,10.282265902386632,18.753596153843198,20.66700000000843,0,63.476923076922546,0,0.0,1.969230769230279,True,none,0,-1,0.0,13.863505177534366,22.44779230770204,23.620846153853847,0,0.0,0,0.9682539682539683,1.0,1.0,351.2176776557015,445.84648717956554,459.55930769245157,511.18822788794483,693.3333333333316,783.3333333333314,0.0,0,0 +230.0,fifo,258.21380417335473,0.02465489566613162,10.256436597110755,7.692327447833066,240.24038523274479,34.922718863758426,1.1226687137971945,1703,670280,20.766666666666666,23.3140869565216,2.5474202898549336,181,73238.33333333043,86.91502882266383,184,34309.08168164904,76497,1174.2878762541263,2240.5347826085754,2608.5217391302963,63,427.3043478260803,0,878.8869565216775,881.1130434781998,False,video,608,461,0.4173913042819777,1164.978260869511,2225.6069565216276,2611.860869565078,177,0.8509615384615384,0,0.8888888888888888,1.0,1.0,1347.216563146941,2417.731014492629,2576.1855072462367,1506.486563750715,2588.1666666666674,2906.6666666666697,0.7718696397941681,42,42 +230.0,strict_priority,258.21380417335473,0.02465489566613162,10.256436597110755,7.692327447833066,240.24038523274479,34.922718863758426,1.1226687137971945,1703,670280,20.766666666666666,23.314086956521678,2.547420289855012,127,73238.33333333276,59.749842696627866,132,34307.118394888734,75969,12.250334448143445,21.84782608690261,23.33913043477409,0,69.00869565217427,0,9.32173913036749,11.547826086889756,True,video,608,458,9.32173913036749,15.51421404680579,25.111304347813142,26.678260869555714,0,0.0,0,0.8888888888888888,1.0,1.0,1445.9947550034299,2505.9953623187953,2583.976811594141,1594.536878216124,2678.1666666666674,2916.666666666668,0.8044596912521441,47,47 +230.0,latest_state,258.21380417335473,0.02465489566613162,10.256436597110755,7.692327447833066,240.24038523274479,34.922718863758426,1.1226687137971945,1703,670280,20.766666666666666,23.314086956521678,2.547420289855012,127,73238.33333333276,59.749842696627866,132,34307.118394888734,75969,12.250334448143445,21.84782608690261,23.33913043477409,0,69.00869565217427,0,9.32173913036749,11.547826086889756,True,video,608,458,9.32173913036749,15.51421404680579,25.111304347813142,26.678260869555714,0,0.0,0,0.8888888888888888,1.0,1.0,1445.9947550034299,2505.9953623187953,2583.976811594141,1594.536878216124,2678.1666666666674,2916.666666666668,0.8044596912521441,47,47 diff --git a/data/processed/lab033/lab033_video_delay.png b/data/processed/lab033/lab033_video_delay.png new file mode 100644 index 0000000..80182cd Binary files /dev/null and b/data/processed/lab033/lab033_video_delay.png differ diff --git a/protocol/link_packet.py b/protocol/link_packet.py new file mode 100644 index 0000000..8a133c6 --- /dev/null +++ b/protocol/link_packet.py @@ -0,0 +1,195 @@ +"""Common Lab033 link packet shared by commands, telemetry, and video. + +The fixed 32-byte header uses network byte order and no implicit padding:: + + !4sBBBBHIQHII + + Offset Size Field + 0 4 magic (b"SLP1") + 4 1 version + 5 1 traffic_class + 6 1 direction + 7 1 flags + 8 2 stream_id + 10 4 sequence_number + 14 8 generation_time_us + 22 2 deadline_ms + 24 4 payload_length + 28 4 packet_crc32 + +``packet_crc32`` is IEEE CRC32 over the complete header with that field set +to zero, followed by the complete payload. +""" + +from __future__ import annotations + +from dataclasses import dataclass, replace +from enum import IntEnum +import struct +import zlib + + +MAGIC = b"SLP1" +VERSION = 1 +HEADER_FORMAT = "!4sBBBBHIQHII" +HEADER_SIZE = struct.calcsize(HEADER_FORMAT) +SUPPORTED_FLAGS_MASK = 0 + + +class TrafficClass(IntEnum): + """Traffic classes in descending scheduling priority.""" + + EMERGENCY = 1 + CONTROL = 2 + TELEMETRY = 3 + VIDEO = 4 + + +class Direction(IntEnum): + """Logical direction of a packet on the shared modelled resource.""" + + GROUND_TO_ROVER = 1 + ROVER_TO_GROUND = 2 + + +class LinkPacketError(ValueError): + """Base class for malformed Lab033 link packets.""" + + +class LinkPacketCRCError(LinkPacketError): + """The Lab033 packet header or payload failed CRC32 validation.""" + + +@dataclass(frozen=True) +class LinkPacket: + """Decoded common packet, including its CRC-validated payload.""" + + traffic_class: TrafficClass + direction: Direction + stream_id: int + sequence_number: int + generation_time_us: int + deadline_ms: int + payload: bytes + flags: int = 0 + version: int = VERSION + packet_crc32: int = 0 + + +def crc32(data: bytes) -> int: + """Return unsigned IEEE CRC32.""" + + return zlib.crc32(data) & 0xFFFFFFFF + + +def _validated(packet: LinkPacket) -> LinkPacket: + if not isinstance(packet, LinkPacket): + raise TypeError("packet must be LinkPacket") + try: + traffic_class = TrafficClass(packet.traffic_class) + direction = Direction(packet.direction) + except ValueError as error: + raise LinkPacketError("unsupported enumeration value") from error + if packet.version != VERSION: + raise LinkPacketError("unsupported link packet version") + if packet.flags & ~SUPPORTED_FLAGS_MASK: + raise LinkPacketError("unsupported link packet flags") + if not 0 <= packet.stream_id <= 0xFFFF: + raise LinkPacketError("stream_id is outside uint16") + if not 0 <= packet.sequence_number <= 0xFFFFFFFF: + raise LinkPacketError("sequence_number is outside uint32") + if not 0 <= packet.generation_time_us <= 0xFFFFFFFFFFFFFFFF: + raise LinkPacketError("generation_time_us is outside uint64") + if not 0 <= packet.deadline_ms <= 0xFFFF: + raise LinkPacketError("deadline_ms is outside uint16") + if not isinstance(packet.payload, (bytes, bytearray)): + raise TypeError("payload must be bytes or bytearray") + payload = bytes(packet.payload) + if len(payload) > 0xFFFFFFFF: + raise LinkPacketError("payload is outside uint32 length") + return replace( + packet, + traffic_class=traffic_class, + direction=direction, + payload=payload, + ) + + +def _pack_header(packet: LinkPacket, crc_value: int) -> bytes: + return struct.pack( + HEADER_FORMAT, + MAGIC, + packet.version, + int(packet.traffic_class), + int(packet.direction), + packet.flags, + packet.stream_id, + packet.sequence_number, + packet.generation_time_us, + packet.deadline_ms, + len(packet.payload), + crc_value, + ) + + +def encode_link_packet(packet: LinkPacket) -> bytes: + """Serialize a common packet and calculate its CRC32.""" + + packet = _validated(packet) + header_without_crc = _pack_header(packet, 0) + packet_crc32 = crc32(header_without_crc + packet.payload) + return _pack_header(packet, packet_crc32) + packet.payload + + +def decode_link_packet(wire_packet: bytes) -> LinkPacket: + """Deserialize a common packet and validate its format and CRC32.""" + + if not isinstance(wire_packet, (bytes, bytearray)): + raise TypeError("wire_packet must be bytes or bytearray") + wire_packet = bytes(wire_packet) + if len(wire_packet) < HEADER_SIZE: + raise LinkPacketError("link packet is shorter than its header") + ( + magic, + version, + traffic_class, + direction, + flags, + stream_id, + sequence_number, + generation_time_us, + deadline_ms, + payload_length, + received_crc32, + ) = struct.unpack(HEADER_FORMAT, wire_packet[:HEADER_SIZE]) + if magic != MAGIC: + raise LinkPacketError("invalid link packet magic") + if len(wire_packet) != HEADER_SIZE + payload_length: + raise LinkPacketError("payload_length does not match packet length") + try: + packet = LinkPacket( + traffic_class=TrafficClass(traffic_class), + direction=Direction(direction), + stream_id=stream_id, + sequence_number=sequence_number, + generation_time_us=generation_time_us, + deadline_ms=deadline_ms, + payload=wire_packet[HEADER_SIZE:], + flags=flags, + version=version, + packet_crc32=received_crc32, + ) + except ValueError as error: + raise LinkPacketError("unsupported enumeration value") from error + packet = _validated(packet) + calculated_crc32 = crc32(_pack_header(packet, 0) + packet.payload) + if calculated_crc32 != received_crc32: + raise LinkPacketCRCError( + "link packet CRC mismatch: " + f"received 0x{received_crc32:08X}, " + f"calculated 0x{calculated_crc32:08X}" + ) + return packet + + +assert HEADER_SIZE == 32 diff --git a/protocol/priority_scheduler.py b/protocol/priority_scheduler.py new file mode 100644 index 0000000..8b6777e --- /dev/null +++ b/protocol/priority_scheduler.py @@ -0,0 +1,222 @@ +"""Deterministic non-preemptive queue schedulers for Lab033.""" + +from __future__ import annotations + +from dataclasses import dataclass +from enum import Enum +from typing import Iterable + +from protocol.link_packet import ( + HEADER_SIZE, + LinkPacket, + TrafficClass, + encode_link_packet, +) + + +TIME_EPSILON_SECONDS = 1e-12 + + +class SchedulerMode(str, Enum): + FIFO = "fifo" + STRICT_PRIORITY = "strict_priority" + LATEST_STATE = "latest_state" + + +@dataclass(frozen=True) +class QueuedPacket: + """A packet plus a stable global arrival-order tie breaker.""" + + packet: LinkPacket + arrival_order: int + + @property + def generation_time_seconds(self) -> float: + return self.packet.generation_time_us / 1_000_000.0 + + @property + def wire_size_bytes(self) -> int: + return HEADER_SIZE + len(self.packet.payload) + + +@dataclass(frozen=True) +class Replacement: + removed: QueuedPacket + replacement: QueuedPacket + + +@dataclass(frozen=True) +class ScheduledPacket: + queued: QueuedPacket + start_seconds: float + end_seconds: float + blocked_by: QueuedPacket | None + blocking_delay_seconds: float + wire_packet: bytes + + +@dataclass(frozen=True) +class QueueSample: + time_seconds: float + packets: int + bytes: int + + +@dataclass(frozen=True) +class ScheduleResult: + mode: SchedulerMode + channel_bitrate_bps: float + transmitted: tuple[ScheduledPacket, ...] + replacements: tuple[Replacement, ...] + queue_samples: tuple[QueueSample, ...] + + +def transmission_duration_seconds( + packet_size_bytes: int, + channel_bitrate_bps: float, +) -> float: + """Return exact serialization duration for a complete packet.""" + + if packet_size_bytes < 0: + raise ValueError("packet size must not be negative") + if channel_bitrate_bps <= 0.0: + raise ValueError("channel bitrate must be positive") + return packet_size_bytes * 8.0 / channel_bitrate_bps + + +def _selection_key(item: QueuedPacket, mode: SchedulerMode) -> tuple[int, int]: + if mode is SchedulerMode.FIFO: + return (0, item.arrival_order) + return (int(item.packet.traffic_class), item.arrival_order) + + +def _enqueue( + ready: list[QueuedPacket], + item: QueuedPacket, + mode: SchedulerMode, + replacements: list[Replacement], +) -> None: + if ( + mode is SchedulerMode.LATEST_STATE + and item.packet.traffic_class + in (TrafficClass.CONTROL, TrafficClass.TELEMETRY) + ): + retained = [] + for old in ready: + if ( + old.packet.traffic_class == item.packet.traffic_class + and old.packet.stream_id == item.packet.stream_id + ): + replacements.append(Replacement(old, item)) + else: + retained.append(old) + ready[:] = retained + ready.append(item) + + +def schedule_packets( + packets: Iterable[QueuedPacket], + mode: SchedulerMode, + channel_bitrate_bps: float, +) -> ScheduleResult: + """Serialize all packets through one non-preemptive shared resource. + + The caller supplies a unique ``arrival_order``. Packets are admitted by + ``(generation_time, arrival_order)``. FIFO selects the smallest arrival + order; priority modes select the smallest ``TrafficClass`` and preserve + arrival order within a class. Latest-state replacement only touches items + still waiting in ``ready`` and therefore can never remove an active packet. + """ + + try: + mode = SchedulerMode(mode) + except ValueError as error: + raise ValueError("unsupported scheduler mode") from error + if channel_bitrate_bps <= 0.0: + raise ValueError("channel bitrate must be positive") + arrivals = sorted( + tuple(packets), + key=lambda item: ( + item.generation_time_seconds, + item.arrival_order, + ), + ) + if len({item.arrival_order for item in arrivals}) != len(arrivals): + raise ValueError("arrival_order values must be unique") + + ready: list[QueuedPacket] = [] + transmitted: list[ScheduledPacket] = [] + replacements: list[Replacement] = [] + samples: list[QueueSample] = [] + blocker_by_order: dict[int, tuple[QueuedPacket, float]] = {} + cursor = 0.0 + arrival_index = 0 + + def sample(time_seconds: float) -> None: + samples.append( + QueueSample( + time_seconds=time_seconds, + packets=len(ready), + bytes=sum(item.wire_size_bytes for item in ready), + ) + ) + + while arrival_index < len(arrivals) or ready: + if not ready: + cursor = max(cursor, arrivals[arrival_index].generation_time_seconds) + while ( + arrival_index < len(arrivals) + and arrivals[arrival_index].generation_time_seconds + <= cursor + TIME_EPSILON_SECONDS + ): + _enqueue(ready, arrivals[arrival_index], mode, replacements) + arrival_index += 1 + sample(cursor) + selected_index = min( + range(len(ready)), + key=lambda index: _selection_key(ready[index], mode), + ) + selected = ready.pop(selected_index) + start = max(cursor, selected.generation_time_seconds) + wire_packet = encode_link_packet(selected.packet) + end = start + transmission_duration_seconds( + len(wire_packet), channel_bitrate_bps + ) + + while ( + arrival_index < len(arrivals) + and arrivals[arrival_index].generation_time_seconds + <= end + TIME_EPSILON_SECONDS + ): + arriving = arrivals[arrival_index] + if arriving.generation_time_seconds > start + TIME_EPSILON_SECONDS: + blocker_by_order[arriving.arrival_order] = ( + selected, + max(0.0, end - arriving.generation_time_seconds), + ) + _enqueue(ready, arriving, mode, replacements) + arrival_index += 1 + + blocker, blocking_delay = blocker_by_order.get( + selected.arrival_order, (None, 0.0) + ) + transmitted.append( + ScheduledPacket( + queued=selected, + start_seconds=start, + end_seconds=end, + blocked_by=blocker, + blocking_delay_seconds=blocking_delay, + wire_packet=wire_packet, + ) + ) + cursor = end + sample(cursor) + + return ScheduleResult( + mode=mode, + channel_bitrate_bps=channel_bitrate_bps, + transmitted=tuple(transmitted), + replacements=tuple(replacements), + queue_samples=tuple(samples), + ) diff --git a/tests/lab033_priority_channel_scheduler.py b/tests/lab033_priority_channel_scheduler.py new file mode 100644 index 0000000..f7b5cb5 --- /dev/null +++ b/tests/lab033_priority_channel_scheduler.py @@ -0,0 +1,1043 @@ +"""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 tests.lab028_video_packetization import ( + COMPOSITE_FPS, + SOURCE_VIDEO_PATH, + EncodedComposite, + VideoMetadata, + load_video_profile, +) +from tests.lab029_packet_channel_simulation import prepare_profiles +from tests.lab030_packet_erasure_fec import ( + FECBlockPlan, + SourcePacket, + TransmissionUnit, + prepare_source_packets, +) +from tests.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", + "- tests/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()