Skip to content

fix: bound and validate stream reassembly on the transport level - #31

Open
Woralem wants to merge 1 commit into
igmunv:devfrom
Woralem:fix/transport-stream-bounds
Open

fix: bound and validate stream reassembly on the transport level#31
Woralem wants to merge 1 commit into
igmunv:devfrom
Woralem:fix/transport-stream-bounds

Conversation

@Woralem

@Woralem Woralem commented Aug 25, 2026

Copy link
Copy Markdown
Contributor

fix: bound and validate stream reassembly on the transport level

Связано с: refs #18 (пункт про транспортный уровень)

Кратко

Транспорт лежит ниже проверки подписи, поэтому он обрабатывает пакеты от кого угодно, ещё до того как станет понятно, кто их прислал. Сейчас он верит заголовку пакета на слово: таблица незавершённых потоков не ограничена и никогда не чистится, chunk_count / chunk_id / size не проверяются, а окно по времени проверяется только в одну сторону. Одиночного UDP-пакета достаточно, чтобы навсегда занять память узла

PR закрывает три из четырёх транспортных дефектов из #18: границы, валидация заголовка и симметричное окно времени. Аутентификация служебных пакетов (пинги, ACK) в этот PR не входит. Я могу его сделать здесь же, но там придётся менять формат кадра

Причина

  1. Незавершённые потоки живут вечно. WAITING_STREAMS[stream_id] создаётся по первому же чанку и удаляется только когда поток собран полностью. Если прислать один чанк из ста и замолчать, запись останется в памяти до перезапуска процесса. Ни таймаута, ни лимита на количество потоков нет, а stream_id - один байт, то есть 256 записей можно занять 256 пакетами.
  2. Заголовку никто не верит, но все ему доверяют. chunk_count - 2 байта, то есть до 65535 чанков, size тоже приходит из пакета. chunk_count == 0 навсегда оставляет поток недособранным, chunk_id >= chunk_count пишет мусор в словарь, а большой size даёт неограниченный рост памяти на один stream_id. Плюс ACK отправляется до проверок, то есть узел ещё и подтверждает мусор
  3. Окно по времени односторонее. Проверка difference_seconds >= 300 отбрасывает старые пакеты, но пакет с временем из будущего даёт отрицательную разницу и проходит всегда. Такой пакет можно переиспользовать бесконечно (replay-окно ограничено только SEEN_PACKETS_WINDOW)

Что сделано

Все лимиты - константы класса, их видно в одном месте и можно переопределить в тестах

Константа Значение Зачем
STREAM_TIMEOUT ACK_TIMEOUT * ACK_RETRIES (30s) Перекрывает полный бюджет повторов отправителя: если за это время не пришло ни одного чанка, поток не соберётся никогда
MAX_WAITING_STREAMS 8 Сколько незавершённых потоков держим одновременно
MAX_STREAM_CHUNKS 4096 Верхняя граница на chunk_count из заголовка
MAX_STREAM_BYTES 1 MiB Бюджет по байтам на один поток
PACKET_MAX_AGE 300 Прежнее поведение, вынесено из тела rworker
CLOCK_SKEW_TOLERANCE 60 Допуск на расхождение часов для пакетов «из будущего»

Валидация заголовка - valid_stream_header(packet)
Вызывается в rworker до send_acknowledgment, чтобы не подтверждать заведомо битый пакет. Отбрасывает chunk_count == 0, chunk_count > MAX_STREAM_CHUNKS, chunk_id >= chunk_count и size > MAX_STREAM_BYTES, каждый случай пишется в лог

Истечение потоков - expire_waiting_streams()
У каждого потока есть deadline на time.monotonic(). Дедлайн продлевается на каждом новом чанке, просроченные потоки выбрасываются перед обработкой очередного пакета.

Учёт потоков - stream_for(packet)

  • если chunk_count не совпал с уже известным для этого stream_id - старая запись считается устаревшей и выбрасывается (переиспользование stream_id больше не склеивает два разных сообщения);
  • если таблица заполнена - вытесняется поток с самым ранним дедлайном, то есть новый трафик не выдавливается мусором, накопленным раньше;
  • запись хранит count, packets, bytes, deadline.

Бюджет по байтам
stream["bytes"] растёт только на новых chunk_id, поэтому честные повторы (retransmit) бюджет не тратят. При превышении MAX_STREAM_BYTES чанк отбрасывается.

Симметричное окно времени
difference_seconds >= PACKET_MAX_AGE or difference_seconds < -CLOCK_SKEW_TOLERANCE - теперь пакет из будущего дальше 60s тоже отбрасывается.

Проверка

Новый файл tests/test_transport.py - 7 тестов, только stdlib + levels.*, без сети и без внешних зависимостей, в том же стиле, что tests/test_pipeline.py:

ok    test_whole_stream_reassembles          (1.20s)
ok    test_bogus_chunk_count_is_refused      (1.20s)
ok    test_abandoned_stream_is_forgotten     (1.60s)
ok    test_stream_id_reuse_does_not_merge     (1.20s)
ok    test_stream_table_stays_bounded        (1.20s)
ok    test_stream_byte_budget_is_enforced    (1.20s)
ok    test_timestamp_window_is_symmetric     (1.20s)

7/7 passed

Cуществующий tests/test_pipeline.py на этой же ветке:

ok    test_wordcoder_roundtrip               (0.00s)
ok    test_roundtrip_varied_sizes            (0.41s)
ok    test_multichunk_large_message          (0.36s)
ok    test_survives_packet_loss              (9.93s)
ok    test_stream_id_wraparound              (0.52s)
ok    test_plaintext_never_hits_the_channel   (0.31s)

6/6 passed

Сборка потоков, повторы при 30% потерь и переиспользование stream_id работают как раньше. Лимиты специально выставлены с запасом относительно рабочих сценариев: 14 KiB сообщение при CHUNK_SIZE = 100 - это ~140 чанков против лимита 4096, а отправитель шлёт один поток за раз, так что MAX_WAITING_STREAMS = 8 не мешает нормальному трафику

Шо осталось с #18

@igmunv

igmunv commented Aug 26, 2026

Copy link
Copy Markdown
Owner

Транспорт лежит ниже проверки подписи

Нет. Подпись проверяется на переходном уровне, которые лежит ниже транспортного.

image

Правильно ли я вас понял?

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants