67 changed files with 3024 additions and 1268 deletions
@ -0,0 +1,128 @@
|
||||
# Разбор полного прогона 2026-09-27 |
||||
|
||||
Исходный результат `./check.sh` вне песочницы: 92 passed, 5 failed, 1 skipped (98 тестов). |
||||
Песочница запрещает сетевые сокеты; её результаты не использовались для оценки протокола. |
||||
|
||||
Промежуточный прогон: 96 passed, 1 failed, 1 skipped; оставался `test_chat_join_e2e`. |
||||
После KEY_REGISTER_ACK и исправления невыровненных чтений повторный полный |
||||
`make check -j4`: **97 passed, 0 failed, 1 skipped (98)**. |
||||
Лог: `/tmp/utun-ack-final-check.log`. Пропущенный `test_auto_socket_dynamic` |
||||
запущен отдельно с root в изолированной сети и упал — подробности ниже. |
||||
Внешние `check-proxy/check-burst/check-load` не выполнялись: цепочка root-проверок |
||||
остановилась на этом тесте. |
||||
|
||||
`test_etcp_lifecycle`, `test_etcp_link_stress`, `test_conn_mgr_phases` дополнительно |
||||
прошли ASan/UBSan/LeakSanitizer. Инструментированы тесты и изменённые транспортные/ |
||||
маршрутные модули, остальные объекты библиотек взяты из обычной сборки. |
||||
Логи: `/tmp/ncd-late-asan.log`, `/tmp/link-role-asan.log`, `/tmp/cm-phases-asan.log`. |
||||
|
||||
## Исправленные причины |
||||
|
||||
- `test_conn_mgr_phases`: после TIMEOUT фазы DIRECT NCD навсегда подавлял UP, |
||||
хотя conn_mgr сохранял handle для REVERSE. В логе физический ETCP уже UP, |
||||
но conn_mgr продолжал DIRECT_REQ до своего таймаута. TIMEOUT теперь обозначает |
||||
истечение срока первой попытки, CLOSED — окончательное завершение. Удерживаемый |
||||
handle получает поздний UP. Отдельный lifecycle-регрессионный тест воспроизводит |
||||
TIMEOUT → INIT/UP (одно уведомление) → DOWN → UP; до исправления он падал. |
||||
- `test_etcp_link_stress`: ожидал reset всей сессии при коллизии дополнительного |
||||
линка. Проверка заменена на сохранение эпох, отсутствие reset и дублей линков, |
||||
доставку и ACK сообщений через новый линк. Прежний тест забирал сырые фрагменты |
||||
вместо normalizer и мог считать служебные пакеты полезным трафиком. Теперь он |
||||
использует отдельный ETCP service и проверяет каждый байт и порядок сообщений. |
||||
Эта проверка выявила ошибку реализации: при принятии встречного INIT сохранённый |
||||
исходящий линк не менял `is_server`. Роль обновляется при отправке INIT_RESPONSE. |
||||
- `test_bgp_route_exchange`: после разрыва B–C recovery успешно строил C–A, |
||||
нарушая предположение теста о фиксированной цепочке. Теперь A и C принимают |
||||
только ключ B: все прежние проверки withdraw/restore сохранены, обход запрещён |
||||
условиями тестовой сети. Автоматический обход отдельно проверяет `test_topo_recovery`. |
||||
- `test_tcp_io`: адрес 192.0.2.1:81 в окружении устанавливал TCP-соединение, |
||||
поэтому connect timeout закономерно не срабатывал. Теперь используется локальный |
||||
listener с заполненной accept-очередью, проверяется ETIMEDOUT. Перед уничтожением |
||||
uasync тест завершает отложенное освобождение TCP и проверяет баланс таймеров. |
||||
|
||||
Диагностика тестов включается `UTUN_TEST_DEBUG=1`; обычный запуск не включает DEBUG. |
||||
Остаточные DATA после удаления соединения могут быть undecryptable и сами по себе |
||||
не означают ошибку транспорта. |
||||
|
||||
## Исправлено: гонка регистрации приглашения |
||||
|
||||
`test_chat_join_e2e`, happy path A != C, падал из-за отсутствия подтверждения KEY_REGISTER. |
||||
Раньше `chat_invite_build_link()` вызывал `chat_join_register_key()` и сразу возвращал ссылку. |
||||
Успех `etcp_route_send()` означает постановку в очередь. Он не подтверждает обработку |
||||
ключа посредником C. Параллельное подключение J–C может завершиться раньше доставки A–C. |
||||
|
||||
Доказательство: `/tmp/join-e2e-debug.log`, один ключ `90c4b081833450a0`: |
||||
|
||||
- 01:39:02.437851 — A пишет `join-key registered`, фактически ключ только поставлен в отправку. |
||||
- 01:39:02.446527 — C отвергает JOIN_INFO_REQ: ключ ещё неизвестен. |
||||
- 01:39:02.449724 — C сохраняет этот ключ, но J уже получил терминальный ERROR. |
||||
|
||||
Теперь C отвечает KEY_REGISTER_ACK только после сохранения ключа. A сопоставляет |
||||
channel/target/key и выдаёт ссылку после успешного ACK и локального сохранения ключа. |
||||
Отрицательный ACK завершает запрос ошибкой. Чужие, неподписанные, незашифрованные, |
||||
повторные и запоздалые подтверждения не могут завершить другой запрос. |
||||
|
||||
Контракт асинхронных API `chat_join_register_key` / `chat_invite_build_link`: |
||||
NULL означает ошибку запуска без callback; принятый запрос даёт ровно один callback. |
||||
Обычное завершение происходит после возврата API, включая self. Cancel/destroy |
||||
завершают запрос синхронно с CANCELLED. Handle недействителен после callback. |
||||
Для удалённого узла дедлайн 5 секунд; доставку и повтор потерянных пакетов обеспечивает |
||||
надёжный router. Повторная регистрация того же ключа тем же инвайтером обновляет TTL; |
||||
замена инвайтера для существующего ключа запрещена. Отмена может оставить ключ на C |
||||
до истечения TTL, но A не авторизует по нему вход, поскольку локально такой ключ |
||||
сохраняется только при успешном завершении. |
||||
|
||||
GUI получает request_id в событии и игнорирует устаревшие результаты. Новый запрос |
||||
отменяет предыдущее ожидание GUI. Headless связывает ответ с JSON id и отменяет все |
||||
ожидания отключившегося клиента. Ни один из путей не запускает вложенный event loop. |
||||
|
||||
Дополнительные исправления: |
||||
|
||||
- Устранён double-free при ошибке `etcp_route_send`: пакет уже освобождён роутером. |
||||
- Headless прекращает разбор пакета после `quit`; последующая команда не создаёт |
||||
запрос на уже закрытом клиенте. При destroy отменяются cleanup_timer и запросы, |
||||
удаляются клиентские сокеты и освобождается контекст. |
||||
- UBSan обнаружил невыровненное чтение group_id из `d + 1` в chat_sync. Оно и три |
||||
аналогичных диагностических чтения заменены на memcpy. |
||||
- Тест chat_join теперь останавливает вручную запущенный member_sync; |
||||
проверяется баланс выделения/освобождения таймеров. |
||||
|
||||
Проверки: |
||||
|
||||
- `test_chat_join`: 31/31, включая задержанный/отсутствующий ACK, таймаут, |
||||
неверные channel/peer/key/flags/размер/status, дубликат, отказ хранения из-за |
||||
конфликта инвайтера, отмену, shutdown, параллельные ключи и ошибку отправки. |
||||
- `test_chat_join_e2e`: 3/3 (A≠C, A=C, неверный ключ). C принудительно откладывает |
||||
регистрацию на 200 мс; проверяется отсутствие ссылки до обработки регистрации. |
||||
Дополнительно проверяются подавление старого GUI-запроса и реальный headless |
||||
control socket: параллельные ответы, валидность ссылок, quit и команда после него. |
||||
- ASan/UBSan: оба теста проходят. LeakSanitizer и баланс таймеров проверены в |
||||
test_chat_join; E2E использует fork/_exit, поэтому сам по себе не проверяет утечки |
||||
при завершении дочерних процессов. Инструментированы изменённые chat-модули и |
||||
transport/routing, остальные объекты взяты из обычной библиотеки. |
||||
Логи: `/tmp/utun-ack-join-asan.log`, `/tmp/utun-ack-e2e-asan.log`. |
||||
|
||||
## Отдельный root-тест auto_socket_dynamic |
||||
|
||||
Запуск в `sudo unshare --net` с поднятым loopback закончился TIMEOUT через 30 секунд. |
||||
Лог: `/tmp/utun-ack-root-checks.log`. Это не PASS основного набора: в нём тест пропущен. |
||||
|
||||
В конфиге теста включён auto_sockets, но не отключён стандартный фильтр |
||||
`auto_socket_skip_no_default_route=1`. Dummy-интерфейсы теста не имеют default route, |
||||
а `auto_socket.c:reconcile_iface()` исключает такие интерфейсы. Лог показывает |
||||
`no links created`, затем TIMEOUT. Для локального сценария тест должен явно задать |
||||
`auto_socket_skip_no_default_route=no` и проверить создание сокетов до трафика. |
||||
|
||||
Кроме того, тест сохраняет handle общего дедлайна, traffic_send и traffic_monitor |
||||
в одну переменную timeout_handle; фазовые таймеры вообще не сохраняются. При выходе |
||||
не отменяются все таймеры; зафиксирован остаток 9 таймеров (665/656). В обработчиках |
||||
трафика освобождается entry без queue_dgram_free. Нужна отдельная доработка fixture: |
||||
владение таймерами и payload, backpressure отправителя, завершение deferred cleanup. |
||||
Производственный auto_socket в рамках исправления приглашений не менялся. |
||||
|
||||
## Попутная проблема диагностики |
||||
|
||||
`lib/debug_config.c:debug_parse_config()` требует `:`, хотя документация описывает `=`. |
||||
Кроме того, функция всё ещё вычисляет индекс через битовую маску, хотя категории |
||||
теперь являются индексами. Исправление в этот набор не включено; тесты используют |
||||
`debug_set_category_level()` напрямую. |
||||
File diff suppressed because it is too large
Load Diff
@ -1,141 +1,102 @@
|
||||
# ETCP Router (etcp_router) |
||||
|
||||
## 1. Назначение |
||||
Сервисный слой маршрутизации поверх ETCP — упрощённый TCP с восстановлением порядка доставки, дедупликацией и контролем inflight. Мультиплексирует до 256 сервисов (`svc_id`) на одном ETCP-соединении. Поддерживает transit (многошаговую маршрутизацию через промежуточные узлы), ретрансмиты, minRTT-измерения и подписи Ed25519. |
||||
|
||||
## 2. Как пользоваться |
||||
|
||||
### Инициализация |
||||
```c |
||||
// В utun_instance_init(): |
||||
etcp_router_init(inst); |
||||
|
||||
// Зарегистрировать обработчик сервиса (svc_id 0..255) |
||||
etcp_router_bind(inst, svc_id, my_service_callback); |
||||
``` |
||||
|
||||
### Отправка |
||||
```c |
||||
// Простой способ: сформировать ll_entry с svc_id в первом байте |
||||
struct ll_entry* e = queue_entry_new(0); |
||||
e->dgram = u_malloc(1 + payload_len); |
||||
e->dgram[0] = svc_id; |
||||
memcpy(e->dgram + 1, payload, payload_len); |
||||
e->len = 1 + payload_len; |
||||
etcp_route_send(inst, group_id, dst_node_id, e, force, mode); // mode: 0 / ROUTE_CRYPTO_SIGN / ROUTE_CRYPTO_ENCRYPT |
||||
|
||||
// Или через существующий ROUTER_CONN: |
||||
struct ETCP_ROUTER_CONN* rconn = etcp_router_conn_get(inst, group_id, remote_node_id, svc_id); |
||||
etcp_router_conn_send(rconn, data, len, mode); |
||||
``` |
||||
|
||||
### Приём |
||||
```c |
||||
void my_service_callback(struct ETCP_CONN* conn, struct ll_entry* entry) { |
||||
if (!entry) { /* соединение закрыто */ return; } |
||||
uint8_t svc_id = entry->dgram[0]; |
||||
uint64_t from = *(uint64_t*)(entry->dgram + ROUTER_SVC_SRC_OFF); // remote_node_id |
||||
uint8_t rx_flags = entry->dgram[ROUTER_SVC_FLAGS_OFF]; // был ли SIGNED/ENCRYPTED на wire |
||||
// ... данные начиная с entry->dgram[ROUTER_SVC_PAYLOAD_OFF] ... |
||||
queue_dgram_free(entry); |
||||
queue_entry_free(entry); |
||||
} |
||||
``` |
||||
|
||||
### Backpressure |
||||
```c |
||||
// На send_q (inflight переполнен): |
||||
struct queue_waiter_handle h; |
||||
etcp_router_on_send_ready(inst, group_id, node_id, svc_id, &h, my_ready_cb, my_arg); |
||||
// Когда send_q освободится — вызовется my_ready_cb |
||||
``` |
||||
|
||||
### Подпись / шифрование (модуль route_crypto) |
||||
```c |
||||
// mode — битовая маска, применяется в начале функции отправки (route_crypto_encode), |
||||
// при приёме проверяется/расшифровывается перед вызовом коллбэка (route_crypto_decode). |
||||
#include "route_crypto.h" |
||||
etcp_route_send(inst, group_id, dst, e, force, ROUTE_CRYPTO_SIGN | ROUTE_CRYPTO_ENCRYPT); |
||||
``` |
||||
|
||||
## 3. API |
||||
|
||||
### Ключевые структуры |
||||
|
||||
**SVC_ROUTE_HDR** (33 байта) — заголовок пакета: |
||||
`cmd(1) + group_id(8) + dst_node_id(8) + src_node_id(8) + seq(4) + svc_id(1) + flags(1) + timestamp(2)` |
||||
|
||||
Флаги: `ROUTER_FLAG_START` (0x80), `ROUTER_FLAG_RST` (0x40), `ROUTER_FLAG_SIGNED` (0x08), `ROUTER_FLAG_ENCRYPTED` (0x04), `ROUTER_FLAG_CLOSE` (0x02). |
||||
|
||||
**Формат доставки сервису** (router → callback): |
||||
`[svc_id:1][src_node_id:8][dst_node_id:8][rx_flags:1][payload...]` (ROUTER_SVC_HDR_SIZE = 18). |
||||
|
||||
**ETCP_ROUTER_CONN** — состояние логического подключения (group_id + remote_node_id + svc_id): |
||||
- `tx_seq`, `rx_seq`, `tx_acked` — seq-нумерация для порядка и inflight |
||||
- `recv_q` — reorder-очередь с хеш-индексом по seq (восстановление порядка) |
||||
- `send_q` — очередь ожидания при полном inflight (backpressure) |
||||
- `send_waiter` — waiter на send_input_q (по 1 пакету, round-robin между сервисами) |
||||
- `watchdog_timer` — медленная проверка инварианта drain (страховка от заклинивания) |
||||
- `inflight_q` — копии финальных отправленных пакетов для ретрансмита (хеш по seq) |
||||
- `incoming_q` — FIFO между сетевым приёмом и recv_q (защита от гонок) |
||||
- `rtt`, `rtt_jitter`, `minrtt` — измерения задержки |
||||
- `inflight_limit` — текущий лимит пакетов в полёте (minrtt_probe снижает до 4) |
||||
- `no_ack_count` — последовательные ретрансмиты без ACK (при 17 — закрытие) |
||||
- `start_sent`, `peer_sync_done` — синхронизация после (пере)подключения |
||||
- `closed` — флаг закрытия (игнорирование таймеров) |
||||
|
||||
**ROUTER_INFLIGHT** — запись в inflight_q: `seq, last_sent_tb, send_count, dgram*` (финальный wire-пакет) |
||||
|
||||
**TRANSIT_QUEUE** — per (group_id, src, dst) пара на промежуточном узле: |
||||
- `q` — FIFO транзитных пакетов |
||||
- `waiter` — backpressure на send_input_q next_hop'а |
||||
- `conn` — ETCP_CONN следующего шага |
||||
|
||||
### Внутренняя архитектура |
||||
``` |
||||
Приём: сеть → incoming_q (FIFO) → router_incoming_q_cb → recv_q (хеш по seq) |
||||
↓ |
||||
router_try_assembly → deliver → route_crypto_decode → cb |
||||
|
||||
Отправка: router_enqueue_send → router_build_packet (seq + route_crypto_encode) → send_q (FIFO, backpressure) |
||||
router_send_kick → waiter на send_input_q → router_send_drain_cb (по 1 пакету, round-robin) |
||||
→ etcp_send + inflight_q (копия финального пакета) |
||||
↓ |
||||
router_track_inflight_state + retrans_schedule |
||||
ACK: периодический (10ms) + idle (500ms) → router_send_ack(rx_seq) |
||||
``` |
||||
|
||||
### Константы |
||||
| Константа | Значение | Описание | |
||||
|---|---|---| |
||||
| `ROUTER_MAX_INFLIGHT` | 256 | Макс. пакетов в полёте | |
||||
| `ROUTER_ACK_INTERVAL_TB` | 100 | ACK интервал 10ms (0.1ms timebase) | |
||||
| `ROUTER_ACK_IDLE_TB` | 5000 | Idle таймаут 500ms | |
||||
| `ROUTER_RETRANS_TIMEOUT_TB` | 3000 | Таймаут ретрансмита 300ms | |
||||
| `ROUTER_NO_ACK_MAX_RETRANS` | 17 | Макс. ретрансмитов без ACK (≈5s) | |
||||
| `ROUTER_MAX_SEND_Q_PACKETS` | 64 | Порог backpressure send_q | |
||||
| `ROUTER_SEND_WATCHDOG_TB` | 5000 | Watchdog проверка инварианта drain (500ms) | |
||||
|
||||
### Функции |
||||
| Функция | Описание | |
||||
|---|---| |
||||
| `etcp_router_init(inst)` | Инициализация: bind ETCP_ID_SVC_ROUTE, создание router_conns | |
||||
| `etcp_router_destroy(inst)` | Деинициализация: unbind + close_all + free | |
||||
| `etcp_router_bind(inst, svc_id, cb)` | Зарегистрировать обработчик сервиса | |
||||
| `etcp_router_unbind(inst, svc_id)` | Удалить обработчик сервиса | |
||||
| `etcp_route_send(inst, group, dst, entry, force, mode)` | Отправить пакет (loopback/transit/direct); mode = ROUTE_CRYPTO_SIGN/ENCRYPT | |
||||
| `etcp_router_conn_get(inst, group, remote, svc_id)` | Найти или создать ROUTER_CONN | |
||||
| `etcp_router_conn_send(rconn, data, len, mode)` | Отправить данные с авто-seq и inflight-контролем | |
||||
| `etcp_router_conn_close(rconn)` | Закрыть seq-подключение (CLOSE + уведомление сервиса) | |
||||
| `etcp_router_conn_close_all_for_node(inst, group, node_id)` | Закрыть все conn к узлу | |
||||
| `etcp_router_conn_restart(inst, group, remote, svc_id)` | Сброс состояния (перезапуск удалённой стороны) | |
||||
| `etcp_router_pause_retrans_for_node(inst, node_id)` | Сбросить ретрансмиты без закрытия (для conn reinit) | |
||||
| `etcp_router_set_max_inflight(rconn, new_max)` | Установить рабочий max_inflight | |
||||
| `etcp_router_on_send_ready/cancel_send_ready(inst, group, node, svc, h, cb, arg)` | Backpressure на send_q | |
||||
| `etcp_router_transit_queues_destroy(conn)` | Удалить все транзитные очереди ETCP_CONN | |
||||
|
||||
### Модуль route_crypto (пред/пост-обработка пакета) |
||||
| Функция | Описание | |
||||
|---|---| |
||||
| `route_crypto_encode(inst, peer, base, len, mode, &out, &len)` | Шифрование/подпись в начале отправки → финальный пакет | |
||||
| `route_crypto_decode(inst, wire, len, &out, &len, &rx_flags)` | Проверка подписи/расшифровка перед коллбэком; ошибка → дроп | |
||||
# ETCP Router: протокол и проверка транспорта |
||||
|
||||
Логическое соединение определяется `(group_id, remote_node_id, svc_id)`. Роутер обеспечивает порядок, дедупликацию и повторную отправку между конечными узлами; физический ETCP и транзитные узлы могут меняться. Формат изменён без обратной совместимости. |
||||
|
||||
## Wire-формат |
||||
|
||||
Packed `SVC_ROUTE_HDR`, 57 байт: |
||||
|
||||
`cmd:1 | group:8 | dst:8 | src:8 | seq:4 | svc:1 | flags:1 | timestamp:2 | reset_id:8 | peer_reset_id:8 | challenge:8` |
||||
|
||||
Числовое представление полей сохраняет существующий формат проекта (native endian). Отдельной версии/согласования формата нет; узлы должны обновляться совместно. |
||||
|
||||
`reset_id` — случайный ненулевой ID локального экземпляра соединения; `peer_reset_id` — ID получателя. DATA, ACK и CLOSE принимаются только при совпадении обоих ID с текущей парой. Смена физического транспорта пару не меняет. |
||||
|
||||
| Пакет | Флаги | Payload | Поля | |
||||
|---|---|---|---| |
||||
| HELLO | START (0x80) | нет | собственный ID, peer=0, challenge=0 | |
||||
| CHALLENGE | RST (0x40) | нет | оба ID, случайный ненулевой challenge | |
||||
| CONFIRM | START|RST | нет | оба ID, эхо challenge | |
||||
| DATA | 0 либо SIGNED/ENCRYPTED | непустой | seq, оба ID, challenge=0 | |
||||
| ACK | 0 | нет | seq=следующий ожидаемый, оба ID, challenge=0 | |
||||
| CLOSE | CLOSE (0x02) | нет | оба ID, challenge=0 | |
||||
|
||||
Незнакомый ID сам по себе не сбрасывает рабочую сессию. Получатель выдаёт свежий challenge и меняет peer ID только после соответствующего CONFIRM. CHALLENGE не меняет рабочую сессию: отправляется CONFIRM и, если нужно, встречный challenge. Потерянные сообщения handshake восстанавливаются повторением HELLO через 300 ms; выдача новых challenge ограничена тем же интервалом. |
||||
|
||||
Первичное согласование сохраняет накопленные исходящие данные. Подтверждённый рестарт пира сбрасывает seq/очереди и уведомляет сервис; собственный ID сохраняется. Явный `etcp_router_conn_restart` создаёт новый собственный ID. Сброс выполняется до уведомления, поэтому данные, отправленные сервисом из callback, сохраняются. Сервис должен заново выполнить прикладную синхронизацию после уведомления; гарантии доставки не продолжаются через явный сброс сессии. |
||||
|
||||
Challenge защищает от случайных старых/переставленных управляющих пакетов, но не является end-to-end аутентификацией против злонамеренного транзитного узла. Управление опирается на защиту физических ETCP-соединений. Подпись/шифрование DATA включаются вызывающим сервисом через `ROUTE_CRYPTO_SIGN` / `ROUTE_CRYPTO_ENCRYPT`. |
||||
|
||||
## Очереди и подтверждения |
||||
|
||||
- `send_q`: plaintext + seq + crypto mode, до 64 пакетов при обычной отправке. Кодирование происходит при отправке после handshake, с актуальной парой ID. |
||||
- `tx_seq`: следующий выделяемый seq; `tx_sent`: граница фактически отправленного; `tx_acked`: граница подтверждённого. |
||||
- `inflight_q`: финальные wire-копии, обычно до 256 пакетов. Копия создаётся до передачи в ETCP; ошибка не теряет pending-пакет. |
||||
- ACK принимается только в модульном интервале `(tx_acked, tx_sent]`. Дубликаты, старые и выходящие за границу ACK не продлевают ожидание прогресса. Удаляются существующие inflight-записи, а не перебираются все номера до произвольного ACK. |
||||
- Новые DATA и ретрансмиты идут по одному через waiter физической `send_input_q`. При смене пути отменяется waiter на очереди, где он действительно был зарегистрирован, включая уже запланированный callback. |
||||
- Каждые 300 ms без прогресса планируется проход повторной отправки. После 17 попыток без ACK соединение закрывается. Отсутствующий маршрут проверяется каждые 20 ms; очередь сохраняется. Пока первичный handshake не завершён, HELLO повторяется без отдельного таймаута. |
||||
- Watchdog 500 ms восстанавливает отправку и переносит waiter при смене заблокированного пути. |
||||
- Сброс физического ETCP сохраняет router inflight и seq. Восстановление самого ETCP-handshake остаётся обязанностью транспортного слоя. |
||||
- Crypto decode/проверка подписи выполняются до помещения в reorder и изменения rx_seq. Повреждённый пакет не подтверждается как доставленный. |
||||
- DATA выдаются сервису строго по seq; входящие дубликаты повторно проверяются после извлечения из очереди. |
||||
- Wire-пакет ограничен `PKTNORM_MAX_DGRAM_SIZE` (16384), включая заголовок и crypto overhead. Пустой payload запрещён: он используется для управления. |
||||
|
||||
## Сервисный API |
||||
|
||||
`etcp_router_conn_send(rconn, data, len, mode)` копирует payload; 0 означает принятие в очередь, а не end-to-end доставку. При backpressure производитель ждёт через `etcp_router_on_send_ready`; отменяет ожидание через `etcp_router_cancel_send_ready`. Handle должен жить до callback/отмены. Режим `force` низкоуровневого `etcp_route_send` обходит порог очереди и не предназначен для непрерывного потока. |
||||
|
||||
`etcp_router_bind` регистрирует callback, владеющий переданным ll_entry. Формат доставки — 26 байт заголовка: |
||||
|
||||
`svc:1 | src:8 | dst:8 | rx_flags:1 | group:8 | payload...` |
||||
|
||||
Читать невыравненные числовые поля нужно через `memcpy`. Заголовок без payload означает уведомление о закрытии/сбросе. После обработки освобождаются `queue_dgram_free(entry)` и `queue_entry_free(entry)`. |
||||
|
||||
Callback может закрыть соединение. Освобождение send_q отложено за пределы callback производителя; отправка проверяет, не закрылся/перезапустился ли rconn. При уничтожении реестр удаляется ровно один раз для каждого соединения. |
||||
|
||||
## Проверки и диагностика |
||||
|
||||
`tests/test_etcp_router_faults` перехватывает SVC_ROUTE между настоящими роутерами, очередями и UASYNC-таймерами. Физическая очередь тестовая: нижележащий ETCP не скрывает потери от роутера. Проверяются: |
||||
|
||||
1. Потери DATA/ACK, перестановка, дубликаты, четыре потока (два направления × два сервиса), backpressure и ограниченность физической очереди. |
||||
2. Потеря HELLO/CHALLENGE/CONFIRM в обоих направлениях. |
||||
3. Старые, повторные, слишком большие ACK и ACK ещё не отправленного пакета. |
||||
4. Сохранение inflight при сбросе транспорта. |
||||
5. Старые DATA/ACK/CLOSE/HELLO/CONFIRM после смены сессии. |
||||
6. Переполнение uint32 seq. |
||||
7. Закрытие после смены пути при зарегистрированном и отложенном waiter. |
||||
8. Закрытие из callback приёма. |
||||
9. Смена маршрута при занятой старой очереди без новой отправки. |
||||
10. Исчезновение и восстановление маршрута. |
||||
11. Таймаут при непрерывных повторных ACK без прогресса. |
||||
12. Одновременный рестарт обоих концов. |
||||
13. Граничные размеры payload, crypto overhead, SIZE_MAX. |
||||
14. Закрытие из callback производителя при освобождении send_q. |
||||
|
||||
Успех передачи требует совпадения содержимого, порядка и количества, финальных ACK, пустых send/inflight. Завершение проверяет освобождение таймеров. Unit-набор включает 34 проверки, в том числе криптографию и повторную доставку после неверной подписи. |
||||
|
||||
Сетевой `test_etcp_router` использует UDP/dummynet, backpressure и сброс с восстановлением ETCP-handshake. `test_etcp_router_reconnect` проверяет a–b–c и перезапуск b, по 500 пакетов до/после, с ожиданием последних ACK в каждой фазе и фиксированным seed. |
||||
|
||||
Переменная `ROUTER_TEST_LOG=/tmp/router.log` включает подробный файл. Рабочая категория — `etcp_route`; DEBUG показывает handshake, пару эпох, ретрансмиты и смену состояния; поштучные ACK/DATA оставлены на TRACE. Ошибочные ветки не скрыты. В fault-тесте WARN/ERROR для намеренно неверных ACK/размеров и искусственного таймаута ожидаемы. |
||||
|
||||
Предыдущий этап проверки роутера: 14 fault-сценариев и 34 unit-проверки прошли; fault-набор также прошёл ASan/UBSan/LeakSanitizer (инструментированы router, route_crypto и ll_queue, остальная библиотека обычная). UDP-тест: 440/440 доставлено и подтверждено. Транзитный: 1000/1000 доставлено и подтверждено. Это не доказательство отсутствия всех ошибок стека. |
||||
|
||||
## Lifecycle транспорта и восстановление |
||||
|
||||
ETCP имеет единую точку финального освобождения после DELETE, завершения callback и освобождения последней ссылки. Запрос закрытия идемпотентен и сразу запрещает новые send/ref_take. |
||||
|
||||
NCD-entry удерживает одну ссылку ETCP. При DELETE запись сначала исключается из реестра, затем владельцы получают CLOSED; старые handles безопасно закрываются после создания нового соединения с тем же node_id. Закрытые callbacks/handles удаляются после внешней рассылки. Отложенный начальный UP отменяется при DOWN/close. При удалении сокета обход использует снимок ссылок ETCP, поскольку DOWN может удалить NCD-entry. |
||||
|
||||
Recovery удерживает TOPO_NODE для каждого уникального кандидата и освобождает ссылку при любом завершении. При успехе BGP сначала получает собственный NCD-handle. `test_topo_recovery` удаляет B и проверяет самостоятельное подключение A↔C и доставку 1000 сообщений без тестового reconnect. |
||||
|
||||
`etcp_session.c` согласует эпохи физического транспорта через HELLO/CHALLENGE/CONFIRM поверх аутентифицированного UDP/STCP. Каждый DATA/ACK содержит пару эпох; старые кадры отбрасываются до обработки. Сброс потока не генерирует ложный DOWN. REINIT повторно запускает обмен таблицами BGP. Формат и таймауты описаны в `doc/etcp_protocol.txt`. |
||||
|
||||
Дополнительные проверки: |
||||
- `test_etcp_lifecycle`: DELETE/refcount, повторное закрытие, закрытие из callback, FIN_WAIT, отмена начального UP, удаление сокета, владение recovery. |
||||
- `test_etcp_session`: реальные crypto/ingress/таймеры; reset любой стороны, одновременный reset, потери всех видов control, повторный reset, отказ одного из двух линков, timeout, старые DATA/ACK/control, неверные размеры. |
||||
- `test_etcp_router_tcp`: тот же контроль содержимого, порядка, дубликатов, backpressure и финальных ACK по STCP; обе стороны сбрасываются одновременно. |
||||
|
||||
Диагностика по запросу: `SESSION_TEST_LOG=/tmp/session.log`, `LIFECYCLE_TEST_LOG=/tmp/lifecycle.log`, `ROUTER_TEST_LOG=/tmp/router.log`, `BGP_TEST_LOG=/tmp/bgp.log`. `ROUTER_TEST_RESET=client|server|both` выбирает сторону reset; `ROUTER_TEST_TCP=1` включает STCP. Ожидаемые ошибки отрицательных сценариев не скрываются. |
||||
|
||||
Дополнительно исправлен монитор conn_mgr: наличие маршрута через посредника не означает прямого соединения. Фоновый PING выбирает транспорт по node_id конечного peer и живой линк именно этого транспорта. Проверка BGP-треугольника выдерживает полный период фоновых проб и требует отсутствия ошибок входящих пакетов до внесения потерь. На старом monitor эта проверка воспроизводимо падает; с исправлением проходит. |
||||
|
||||
Итоговая проверка текущих исправлений: 16 регрессионных executable-тестов прошли (router unit/faults/lifecycle/session, UDP, TCP, reconnect, recovery, STCP traffic, NCD, три conn_mgr, normalizer/ETCP, BGP triangle). Семь наборов дополнительно прошли ASan/UBSan/LeakSanitizer: lifecycle, session, router faults, UDP router, TCP router, autonomous recovery, BGP triangle. Инструментированы ETCP core/API/connections/session, NCD, router, recovery, topo_group/connect, conn_mgr core/monitor и ll_queue; остальные библиотечные объекты обычные. Отрицательные тесты намеренно вызывают ошибки; после принудительного удаления транспорта возможен отказ приёма запоздалой UDP-датаграммы без существующего линка. Финальный BGP-прогон, включая фоновые пробы и cleanup таймеров, завершился без диагностических ошибок. |
||||
|
||||
@ -0,0 +1,124 @@
|
||||
#include "etcp_session.h" |
||||
#include "etcp.h" |
||||
#include "etcp_connections.h" |
||||
#include "../lib/debug_config.h" |
||||
#include "../lib/platform_compat.h" |
||||
#include <string.h> |
||||
|
||||
#define SESSION_RETRY_TB 5000 |
||||
#define SESSION_TIMEOUT_TB 100000 |
||||
|
||||
static uint64_t read_id(const uint8_t* p) { uint64_t v; memcpy(&v,p,8); return be64toh(v); } |
||||
static void write_id(uint8_t* p, uint64_t v) { v = htobe64(v); memcpy(p,&v,8); } |
||||
void etcp_session_encode(uint8_t* data, uint8_t code, uint64_t sender, uint64_t target) { |
||||
data[0] = code; write_id(data+1,sender); write_id(data+9,target); |
||||
} |
||||
|
||||
static void session_send(struct ETCP_CONN* c, uint8_t code, uint64_t peer, uint64_t cookie) { |
||||
for (struct ETCP_LINK* link = c->links; link; link = link->next) { |
||||
if (!link->initialized) continue; |
||||
struct ETCP_DGRAM* pkt = memory_pool_alloc(c->instance->pkt_pool); |
||||
if (!pkt) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP,"[%s] session: packet allocation failed",c->log_name); return; } |
||||
pkt->link = link; pkt->data_len = ETCP_SESSION_CONTROL_SIZE; pkt->noencrypt_len = 0; |
||||
etcp_session_encode(pkt->data,code,c->reset_id,peer); write_id(pkt->data+17,cookie); |
||||
if (etcp_encrypt_send(pkt) < 0) |
||||
DEBUG_WARN(DEBUG_CATEGORY_ETCP,"[%s] session: send failed code=%u link=%u",c->log_name,code,link->local_link_id); |
||||
memory_pool_free(c->instance->pkt_pool,pkt); |
||||
} |
||||
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP,"[%s] session TX code=%u local=%016llx peer=%016llx cookie=%016llx", |
||||
c->log_name,code,(unsigned long long)c->reset_id,(unsigned long long)peer,(unsigned long long)cookie); |
||||
} |
||||
|
||||
void etcp_session_cancel(struct ETCP_CONN* c) { |
||||
if (c->session_timer) { uasync_cancel_timeout(c->instance->ua,c->session_timer); c->session_timer = NULL; } |
||||
c->session_candidate = c->session_challenge = 0; |
||||
} |
||||
|
||||
static void session_tick(void* arg) { |
||||
struct ETCP_CONN* c = arg; |
||||
c->session_timer = NULL; |
||||
if (c->close_requested) return; |
||||
if (get_time_tb() - c->session_started >= SESSION_TIMEOUT_TB) { |
||||
DEBUG_WARN(DEBUG_CATEGORY_ETCP,"[%s] session confirmation timeout required=%u candidate=%016llx",c->log_name, |
||||
c->session_required,(unsigned long long)c->session_candidate); |
||||
etcp_session_cancel(c); |
||||
if (c->session_required) etcp_connection_close(c); |
||||
return; |
||||
} |
||||
if (c->session_required) session_send(c,ETCP_SESSION_HELLO,c->peer_reset_id,0); |
||||
if (c->session_candidate) session_send(c,ETCP_SESSION_CHALLENGE,c->session_candidate,c->session_challenge); |
||||
c->session_timer = uasync_set_timeout(c->instance->ua,SESSION_RETRY_TB,c,session_tick,"etcp_session"); |
||||
} |
||||
|
||||
static void session_schedule(struct ETCP_CONN* c) { |
||||
if (c->session_timer) return; |
||||
c->session_started = get_time_tb(); |
||||
c->session_timer = uasync_set_timeout(c->instance->ua,1,c,session_tick,"etcp_session"); |
||||
} |
||||
|
||||
void etcp_session_start(struct ETCP_CONN* c) { |
||||
etcp_session_cancel(c); c->session_required = 1; |
||||
DEBUG_INFO(DEBUG_CATEGORY_ETCP,"[%s] session negotiation started local=%016llx peer=%016llx",c->log_name, |
||||
(unsigned long long)c->reset_id,(unsigned long long)c->peer_reset_id); |
||||
session_schedule(c); |
||||
} |
||||
|
||||
void etcp_session_observe(struct ETCP_CONN* c, uint64_t peer_epoch) { |
||||
if (!peer_epoch || c->close_requested) return; |
||||
if (c->session_candidate != peer_epoch) { |
||||
c->session_candidate = peer_epoch; |
||||
do { |
||||
if (random_bytes((uint8_t*)&c->session_challenge,8) != 0) { |
||||
c->session_candidate = 0; |
||||
DEBUG_ERROR(DEBUG_CATEGORY_ETCP,"[%s] session challenge RNG failed",c->log_name); return; |
||||
} |
||||
} while (!c->session_challenge); |
||||
} |
||||
session_schedule(c); |
||||
} |
||||
|
||||
int etcp_session_receive(struct ETCP_CONN* c, const uint8_t* data, size_t len) { |
||||
if (!len || data[0] < ETCP_SESSION_HELLO || data[0] > ETCP_SESSION_CONFIRM) return 0; |
||||
if (len != ETCP_SESSION_CONTROL_SIZE || c->close_requested) { |
||||
DEBUG_WARN(DEBUG_CATEGORY_ETCP,"[%s] invalid session control size=%zu",c->log_name,len); return 1; |
||||
} |
||||
uint64_t peer = read_id(data+1), target = read_id(data+9), cookie = read_id(data+17); |
||||
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP,"[%s] session RX code=%u peer=%016llx target=%016llx cookie=%016llx", |
||||
c->log_name,data[0],(unsigned long long)peer,(unsigned long long)target,(unsigned long long)cookie); |
||||
if (!peer) { DEBUG_WARN(DEBUG_CATEGORY_ETCP,"[%s] session: zero peer epoch",c->log_name); return 1; } |
||||
if (data[0] == ETCP_SESSION_HELLO) { |
||||
etcp_session_observe(c,peer); |
||||
if (c->session_challenge) session_send(c,ETCP_SESSION_CHALLENGE,peer,c->session_challenge); |
||||
} else if (target != c->reset_id || !cookie) { |
||||
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP,"[%s] session: stale target or empty cookie",c->log_name); |
||||
} else if (data[0] == ETCP_SESSION_CHALLENGE) { |
||||
session_send(c,ETCP_SESSION_CONFIRM,peer,cookie); |
||||
if (c->session_required || peer != c->peer_reset_id) { |
||||
etcp_session_observe(c,peer); |
||||
if (c->session_challenge) session_send(c,ETCP_SESSION_CHALLENGE,peer,c->session_challenge); |
||||
} |
||||
} else if (peer == c->session_candidate && cookie == c->session_challenge) { |
||||
etcp_session_cancel(c); |
||||
if (peer != c->peer_reset_id) etcp_conn_reinit_id(c,"confirmed peer epoch",c->reset_id); |
||||
c->peer_reset_id = peer; c->session_required = 0; |
||||
DEBUG_INFO(DEBUG_CATEGORY_ETCP,"[%s] session confirmed local=%016llx peer=%016llx",c->log_name, |
||||
(unsigned long long)c->reset_id,(unsigned long long)peer); |
||||
etcp_conn_ready(c); |
||||
} else { |
||||
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP,"[%s] session: stale confirmation",c->log_name); |
||||
} |
||||
return 1; |
||||
} |
||||
|
||||
int etcp_session_accept_data(struct ETCP_CONN* c, const uint8_t* data, size_t len) { |
||||
if (len <= ETCP_SESSION_HEADER_SIZE || data[0] != ETCP_SESSION_DATA || data[ETCP_SESSION_HEADER_SIZE] > 1) { |
||||
DEBUG_WARN(DEBUG_CATEGORY_ETCP,"[%s] invalid session data size=%zu",c->log_name,len); return 0; |
||||
} |
||||
uint64_t peer = read_id(data+1), target = read_id(data+9); |
||||
if (peer != c->peer_reset_id || target != c->reset_id || c->reinit_pending || !c->initialized) { |
||||
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP,"[%s] discard stale data peer=%016llx target=%016llx pending=%u",c->log_name, |
||||
(unsigned long long)peer,(unsigned long long)target,c->reinit_pending); |
||||
return 0; |
||||
} |
||||
return 1; |
||||
} |
||||
@ -0,0 +1,21 @@
|
||||
#ifndef ETCP_SESSION_H |
||||
#define ETCP_SESSION_H |
||||
#include <stddef.h> |
||||
#include <stdint.h> |
||||
struct ETCP_CONN; |
||||
|
||||
#define ETCP_SESSION_DATA 0x09 |
||||
#define ETCP_SESSION_HELLO 0x0a |
||||
#define ETCP_SESSION_CHALLENGE 0x0b |
||||
#define ETCP_SESSION_CONFIRM 0x0c |
||||
#define ETCP_SESSION_HEADER_SIZE 17 |
||||
#define ETCP_SESSION_CONTROL_SIZE 25 |
||||
|
||||
/* Authenticated link control, independent of the reliable stream and its normalizer. */ |
||||
void etcp_session_start(struct ETCP_CONN* conn); |
||||
void etcp_session_observe(struct ETCP_CONN* conn, uint64_t peer_epoch); |
||||
void etcp_session_cancel(struct ETCP_CONN* conn); |
||||
int etcp_session_receive(struct ETCP_CONN* conn, const uint8_t* data, size_t len); |
||||
void etcp_session_encode(uint8_t* data, uint8_t code, uint64_t sender, uint64_t target); |
||||
int etcp_session_accept_data(struct ETCP_CONN* conn, const uint8_t* data, size_t len); |
||||
#endif |
||||
@ -0,0 +1,182 @@
|
||||
/* Lifecycle regressions: real ETCP queues/pools/event loop, no network required. */ |
||||
#include <stdio.h> |
||||
#include <stdlib.h> |
||||
#include "etcp.h" |
||||
#include "etcp_api.h" |
||||
#include "node_conn_direct.h" |
||||
#include "utun_instance.h" |
||||
#include "config_parser.h" |
||||
#include "topo_group.h" |
||||
#include "topo_node.h" |
||||
#include "topo_recovery.h" |
||||
#include "../lib/debug_config.h" |
||||
#include "../lib/mem.h" |
||||
|
||||
#define CHECK(x) do { if (!(x)) { fprintf(stderr,"FAIL %s:%d: %s\n",__func__,__LINE__,#x); exit(1); } } while (0) |
||||
static int deleted, victim_called, release_on_delete; |
||||
static void victim(struct ETCP_CONN* c, int ev, void* arg) { (void)c; (void)ev; (void)arg; victim_called++; } |
||||
static void remove_next(struct ETCP_CONN* c, int ev, void* arg) { |
||||
(void)ev; (void)arg; etcp_conn_remove_cbk(c, victim, NULL); etcp_conn_remove_cbk(c, remove_next, NULL); |
||||
} |
||||
static void on_delete(struct ETCP_CONN* c, int ev, void* arg) { |
||||
(void)ev; (void)arg; deleted++; |
||||
if (release_on_delete) etcp_conn_ref_free(c); |
||||
CHECK(!c->free_scheduled); // phase 1 still owns the object
|
||||
etcp_connection_close(c); |
||||
} |
||||
static void close_twice(struct ETCP_CONN* c, int ev, void* arg) { |
||||
(void)ev; (void)arg; |
||||
etcp_connection_close(c); void* token = c->close_token; CHECK(token); |
||||
etcp_connection_close(c); CHECK(c->close_token == token); |
||||
struct ll_entry* e = queue_entry_new(0); |
||||
CHECK(etcp_send(c,e) == -1); queue_entry_free(e); |
||||
} |
||||
static struct NODE_CONN_DIRECT *owner1, *owner2; |
||||
static int owner1_calls, owner2_calls, closed_calls, action; |
||||
static void owner_cb(struct NODE_CONN_DIRECT* h, enum ncd_event event, void* arg) { |
||||
int which = (int)(intptr_t)arg; |
||||
if (event == NCD_EVENT_CLOSED) { CHECK(node_conn_direct_get_conn(h) == NULL); closed_calls++; return; } |
||||
if (which == 1) owner1_calls++; else owner2_calls++; |
||||
if (which == 2 && event == NCD_EVENT_DOWN) { |
||||
if (action >= 1) { node_conn_direct_force_close(owner1); owner1 = NULL; } |
||||
if (action >= 2) { node_conn_direct_force_close(h); owner2 = NULL; } |
||||
} |
||||
} |
||||
static struct ETCP_CONN* ready_conn(struct UTUN_INSTANCE* inst) { |
||||
struct ETCP_CONN* c = etcp_connection_create(inst,"ncd_lifecycle"); CHECK(c); |
||||
c->peer_node_id = 42; c->links_up = 1; etcp_conn_ready(c); |
||||
queue_set_callback(c->send_input_q,NULL,NULL); // физического линка нет; проверяем lifecycle управляющих отправок
|
||||
return c; |
||||
} |
||||
static void ncd_cases(struct UTUN_INSTANCE* inst) { |
||||
for (int scenario=0; scenario<9; scenario++) { |
||||
owner1 = owner2 = NULL; owner1_calls = owner2_calls = closed_calls = 0; action = 0; |
||||
struct ETCP_CONN* c = ready_conn(inst); |
||||
CHECK(node_conn_direct_open(inst,42,owner_cb,(void*)1,&owner1,NULL) == NCD_REUSED); |
||||
CHECK(node_conn_direct_open(inst,42,owner_cb,(void*)2,&owner2,NULL) == NCD_REUSED); |
||||
CHECK(c->ref_count == 1); // одна ссылка на entry, не на handle
|
||||
if (scenario < 2) { |
||||
etcp_connection_close(c); uasync_poll(inst->ua,0); |
||||
CHECK(closed_calls == 2 && owner1_calls == 1 && owner2_calls == 1); // DOWN, без запоздалого UP
|
||||
CHECK(!inst->ncd_registry && !node_conn_direct_get_conn(owner1)); |
||||
struct ETCP_CONN* fresh = ready_conn(inst); |
||||
struct NODE_CONN_DIRECT* h = NULL; |
||||
CHECK(node_conn_direct_open(inst,42,NULL,NULL,&h,NULL) == NCD_REUSED); |
||||
CHECK(node_conn_direct_get_conn(h) == fresh); |
||||
if (scenario) { node_conn_direct_force_close(owner1); node_conn_direct_force_close(owner2); } |
||||
else { node_conn_direct_close(owner1); node_conn_direct_close(owner2); } |
||||
node_conn_direct_force_close(h); |
||||
} else if (scenario < 4) { |
||||
action = scenario - 1; c->links_up = 0; |
||||
etcp_cbk_fire(c,ETCP_CBK_EVENT_DOWN); |
||||
CHECK(owner2_calls == 1 && owner1_calls == 0); |
||||
if (owner2) node_conn_direct_force_close(owner2); |
||||
} else if (scenario == 4) { |
||||
node_conn_direct_close(owner1); node_conn_direct_close(owner2); |
||||
CHECK(c->fin_wait && inst->ncd_registry); |
||||
struct NODE_CONN_DIRECT* h = NULL; |
||||
CHECK(node_conn_direct_open(inst,42,owner_cb,(void*)1,&h,NULL) == NCD_REUSED); |
||||
CHECK(!c->fin_wait); uasync_poll(inst->ua,0); CHECK(owner1_calls == 1 && !owner2_calls); |
||||
node_conn_direct_force_close(h); |
||||
} else if (scenario == 5) { |
||||
node_conn_direct_force_close(owner1); node_conn_direct_force_close(owner2); |
||||
uasync_poll(inst->ua,0); CHECK(!owner1_calls && !owner2_calls && !closed_calls); |
||||
} else if (scenario == 8) { |
||||
struct ETCP_SOCKET sock = {0}; sock.instance = inst; |
||||
struct ETCP_LINK* link = u_calloc(1,sizeof(*link)); CHECK(link); |
||||
link->etcp = c; link->conn = &sock; link->initialized = link->link_status = 1; |
||||
c->links = link; action = 2; |
||||
CHECK(ncd_remove_socket_links(inst,&sock) == 1); CHECK(!owner1 && !owner2); |
||||
} else if (scenario == 7) { |
||||
uasync_poll(inst->ua,0); CHECK(owner1_calls == 1 && owner2_calls == 1); |
||||
etcp_conn_reinit_id(c,"lifecycle reset",c->reset_id); |
||||
CHECK(c->links_up == 1 && owner1_calls == 1 && owner2_calls == 1); |
||||
CHECK(!c->close_requested && node_conn_direct_get_conn(owner1) == c); |
||||
etcp_conn_ready(c); CHECK(!c->reinit_pending); |
||||
node_conn_direct_force_close(owner1); node_conn_direct_force_close(owner2); |
||||
} else { |
||||
c->links_up = 0; etcp_cbk_fire(c,ETCP_CBK_EVENT_DOWN); |
||||
uasync_poll(inst->ua,0); CHECK(owner1_calls == 1 && owner2_calls == 1); // pending initial UP cancelled
|
||||
node_conn_direct_force_close(owner1); node_conn_direct_force_close(owner2); |
||||
} |
||||
uasync_poll(inst->ua,0); |
||||
CHECK(inst->ncd_registry == NULL && inst->connections->count == 0); |
||||
printf("PASS NCD scenario %d\n",scenario+1); |
||||
} |
||||
} |
||||
static void recovery_cases(struct UTUN_INSTANCE* inst) { |
||||
struct TOPO_GROUPS groups = {0}; struct TOPO_GROUP group = {0}; struct TOPO_GROUP_NODE node = {0}; |
||||
groups.instance = inst; inst->topo_groups = &groups; group.instance = inst; node.node_id = 43; |
||||
groups.node_registry = queue_new(inst->ua,16,0,8,"recovery_registry"); |
||||
for (int scenario=0;scenario<3;scenario++) { |
||||
struct TOPO_NODE* ni = u_calloc(1,sizeof(*ni)); CHECK(ni); ni->node_id = node.node_id; |
||||
CHECK(topo_node_registry_store(&groups,ni) == ni && ni->group_ref_count == 1); |
||||
topo_recovery_add_node(&group,&node,44); topo_recovery_add_node(&group,&node,44); |
||||
CHECK(ni->group_ref_count == 2 && group.recovery_list->count == 1); |
||||
if (scenario == 1) { topo_recovery_add_node(&group,&node,45); CHECK(ni->group_ref_count == 3); } |
||||
topo_node_registry_unref(&groups,node.node_id); CHECK(topo_node_registry_find(&groups,node.node_id) == ni); |
||||
if (scenario == 0) topo_recovery_cancel_for_node(&group,node.node_id); |
||||
else if (scenario == 1) topo_recovery_cancel_all(&group); |
||||
else { topo_recovery_start(&group); topo_recovery_cancel_all(&group); } |
||||
uasync_poll(inst->ua,0); |
||||
CHECK(!group.recovery_list && !groups.node_registry->count && !inst->ncd_registry && !inst->connections->count); |
||||
printf("PASS recovery ownership scenario %d\n",scenario+1); |
||||
} |
||||
queue_free(groups.node_registry); inst->topo_groups = NULL; |
||||
} |
||||
static int late_events[4]; |
||||
static void late_up_cb(struct NODE_CONN_DIRECT* h, enum ncd_event ev, void* arg) { |
||||
(void)h; (void)arg; late_events[ev]++; |
||||
} |
||||
static void ncd_late_up_case(struct UTUN_INSTANCE* inst) { |
||||
struct ETCP_CONN* c = ready_conn(inst); c->links_up = 0; |
||||
struct NODE_CONN_DIRECT* h = NULL; |
||||
CHECK(node_conn_direct_open(inst,42,late_up_cb,NULL,&h,NULL) == NCD_REUSED); |
||||
uasync_poll(inst->ua,10); CHECK(late_events[NCD_EVENT_TIMEOUT] == 1 && !late_events[NCD_EVENT_UP]); |
||||
c->links_up = 1; |
||||
etcp_cbk_fire(c,ETCP_CBK_EVENT_INIT); etcp_cbk_fire(c,ETCP_CBK_EVENT_UP); |
||||
CHECK(late_events[NCD_EVENT_UP] == 1); // INIT + UP должны дать ровно одно уведомление
|
||||
c->links_up = 0; etcp_cbk_fire(c,ETCP_CBK_EVENT_DOWN); |
||||
c->links_up = 1; etcp_cbk_fire(c,ETCP_CBK_EVENT_UP); |
||||
CHECK(late_events[NCD_EVENT_DOWN] == 1 && late_events[NCD_EVENT_UP] == 2); |
||||
node_conn_direct_force_close(h); uasync_poll(inst->ua,0); |
||||
CHECK(!inst->ncd_registry && !inst->connections->count); |
||||
puts("PASS NCD late UP after timeout and subsequent reconnect"); |
||||
} |
||||
int main(void) { |
||||
debug_config_init(); debug_set_level(DEBUG_LEVEL_WARN); |
||||
const char* log = getenv("LIFECYCLE_TEST_LOG"); |
||||
if (log) { debug_enable_file_output(log,1); debug_set_category_level(DEBUG_CATEGORY_CONNECTION,DEBUG_LEVEL_DEBUG); } |
||||
size_t baseline = u_get_allocated_count(); |
||||
struct UTUN_INSTANCE inst = {0}; struct utun_config cfg = {0}; |
||||
inst.ua = uasync_create(); inst.config = &cfg; |
||||
inst.connections = queue_new(inst.ua,16,0,8,"test_connections"); |
||||
for (int scenario=0; scenario<4; scenario++) { |
||||
deleted = victim_called = 0; release_on_delete = scenario == 0; |
||||
struct ETCP_CONN* c = etcp_connection_create(&inst,"lifecycle"); CHECK(c); |
||||
etcp_conn_add_cbk(c,on_delete,NULL,ETCP_CBK_EVENT_DELETE); |
||||
if (scenario < 2) CHECK(etcp_conn_ref_take(c) == 0); |
||||
if (scenario == 2) { |
||||
etcp_conn_add_cbk(c,close_twice,NULL,ETCP_CBK_EVENT_REINIT); |
||||
etcp_cbk_fire(c,ETCP_CBK_EVENT_REINIT); |
||||
CHECK(c->close_requested && deleted == 0); |
||||
} |
||||
if (scenario == 3) { |
||||
etcp_conn_add_cbk(c,victim,NULL,ETCP_CBK_EVENT_REINIT); |
||||
etcp_conn_add_cbk(c,remove_next,NULL,ETCP_CBK_EVENT_REINIT); |
||||
etcp_cbk_fire(c,ETCP_CBK_EVENT_REINIT); CHECK(!victim_called); |
||||
} |
||||
etcp_connection_close(c); CHECK(deleted == 1 && c->delete_complete && !c->close_token); |
||||
if (scenario == 1) { |
||||
uasync_poll(inst.ua,0); CHECK(c->ref_count == 1 && !c->free_scheduled); |
||||
etcp_conn_ref_free(c); CHECK(c->free_scheduled); |
||||
} |
||||
uasync_poll(inst.ua,0); CHECK(inst.connections->count == 0); |
||||
printf("PASS lifecycle scenario %d\n",scenario+1); |
||||
} |
||||
ncd_cases(&inst); ncd_late_up_case(&inst); recovery_cases(&inst); |
||||
queue_free(inst.connections); uasync_poll(inst.ua,0); |
||||
CHECK(inst.ua->timer_alloc_count == inst.ua->timer_free_count); |
||||
uasync_destroy(inst.ua,0); CHECK(u_get_allocated_count() == baseline); |
||||
puts("All ETCP lifecycle scenarios passed"); return 0; |
||||
} |
||||
@ -0,0 +1,361 @@
|
||||
/* Deterministic SVC_ROUTE fault injection. Real router/queues/timers, no UDP:
|
||||
* lower ETCP must not conceal the router's own loss/reorder/restart bugs. */ |
||||
#include <stdio.h> |
||||
#include <stdlib.h> |
||||
#include <string.h> |
||||
#include "etcp.h" |
||||
#include "pkt_normalizer.h" |
||||
#include "etcp_router.h" |
||||
#include "utun_instance.h" |
||||
#include "../lib/debug_config.h" |
||||
#include "../lib/mem.h" |
||||
|
||||
#define SVC 0x42 |
||||
#define COUNT 160 |
||||
#define REQUIRE(x) do { if (!(x)) { fprintf(stderr, "FAIL %s:%d: %s\n", __func__, __LINE__, #x); return -1; } } while (0) |
||||
|
||||
struct endpoint { |
||||
struct UTUN_INSTANCE inst; |
||||
struct ETCP_CONN conn, alternate; |
||||
struct ll_entry* route; |
||||
unsigned sent[2], received[2], notified[2]; |
||||
}; |
||||
static struct endpoint ep[2]; |
||||
static struct UASYNC* ua; |
||||
static int failed, drop_data, drop_ack, drop_hello[2], drop_challenge[2], drop_confirm[2], drop_all_ack; |
||||
static int reorder, duplicated, dropped_data, dropped_ack, handshake_drops, closed_in_handler; |
||||
static unsigned max_wire; |
||||
static struct ll_entry *held, *saved_data, *saved_ack, *saved_hello, *saved_confirm; |
||||
|
||||
static void release(struct ll_entry* e) { if (e) { queue_dgram_free(e); queue_entry_free(e); } } |
||||
static struct ll_entry* copy_packet(const struct ll_entry* e) { |
||||
struct ll_entry* copy = ll_alloc_lldgram(e->len); |
||||
if (!copy) abort(); |
||||
memcpy(copy->dgram, e->dgram, e->len); copy->len = e->len; |
||||
return copy; |
||||
} |
||||
static struct ETCP_ROUTER_CONN* rc(int side, int svc) { |
||||
return etcp_router_conn_get(&ep[side].inst, 0, ep[1-side].inst.node_id, SVC + svc); |
||||
} |
||||
static void receive(struct ETCP_CONN* conn, struct ll_entry* e) { |
||||
int side = conn->instance == &ep[0].inst ? 0 : 1; |
||||
unsigned svc = e->dgram[0] - SVC; |
||||
uint64_t src, dst, group; |
||||
memcpy(&src, e->dgram + ROUTER_SVC_SRC_OFF, 8); memcpy(&dst, e->dgram + ROUTER_SVC_DST_OFF, 8); |
||||
memcpy(&group, e->dgram + ROUTER_SVC_GROUP_OFF, 8); |
||||
if (svc >= 2 || src != ep[1-side].inst.node_id || dst != ep[side].inst.node_id || group) { |
||||
fprintf(stderr, "FAIL delivery identity side=%d svc=%u\n", side, svc); failed = 1; release(e); return; |
||||
} |
||||
if (e->len == ROUTER_SVC_HDR_SIZE) { ep[side].notified[svc]++; release(e); return; } |
||||
uint32_t number = 0; |
||||
if (e->len < ROUTER_SVC_HDR_SIZE + 4) { failed = 1; release(e); return; } |
||||
memcpy(&number, e->dgram + ROUTER_SVC_PAYLOAD_OFF, 4); |
||||
if (number != ep[side].received[svc]) { |
||||
fprintf(stderr, "FAIL delivery sequence side=%d svc=%u got=%u expected=%u\n", side, svc, number, ep[side].received[svc]); |
||||
failed = 1; |
||||
} |
||||
size_t len = 4 + number % 181; |
||||
if (e->len != ROUTER_SVC_HDR_SIZE + len) failed = 1; |
||||
for (size_t i = 4; i < len; i++) if (e->dgram[ROUTER_SVC_PAYLOAD_OFF + i] != (uint8_t)(number ^ i ^ svc)) failed = 1; |
||||
ep[side].received[svc]++; |
||||
release(e); |
||||
if (closed_in_handler) { closed_in_handler = 0; etcp_router_conn_close(rc(side, svc)); } |
||||
} |
||||
static void inject(int side, struct ll_entry* e) { |
||||
ep[side].inst.api_bindings.callbacks[ETCP_RT_ID_SVC_ROUTE](&ep[side].conn, e); |
||||
} |
||||
static void wire(int from, struct ll_queue* q) { |
||||
if ((unsigned)q->count > max_wire) max_wire = q->count; |
||||
struct ll_entry* e = queue_data_get(q); |
||||
if (!e) return; |
||||
if (e->len < SVC_ROUTE_HDR_SIZE) { release(e); return; } |
||||
struct SVC_ROUTE_HDR* h = (struct SVC_ROUTE_HDR*)e->dgram; |
||||
int data = e->len > SVC_ROUTE_HDR_SIZE; |
||||
if (!data && h->flags == ROUTER_FLAG_START && !saved_hello && from == 0) saved_hello = copy_packet(e); |
||||
if (!data && h->flags == (ROUTER_FLAG_START | ROUTER_FLAG_RST) && !saved_confirm && from == 0) saved_confirm = copy_packet(e); |
||||
int* drop = NULL; |
||||
if (!data && h->flags == ROUTER_FLAG_START) drop = &drop_hello[from]; |
||||
if (!data && h->flags == ROUTER_FLAG_RST) drop = &drop_challenge[from]; |
||||
if (!data && h->flags == (ROUTER_FLAG_START | ROUTER_FLAG_RST)) drop = &drop_confirm[from]; |
||||
if (drop && *drop) { --*drop; handshake_drops++; release(e); return; } |
||||
if (data && from == 0 && h->svc_id == SVC) { |
||||
if (!saved_data) saved_data = copy_packet(e); |
||||
if (drop_data && h->seq == 2) { drop_data--; dropped_data++; release(e); return; } |
||||
if (reorder && h->seq == 3 && !held) { held = e; return; } |
||||
if (reorder && h->seq == 4 && held) { |
||||
struct ll_entry* dup = copy_packet(e); |
||||
inject(1, e); inject(1, dup); inject(1, held); held = NULL; reorder = 0; duplicated++; |
||||
return; |
||||
} |
||||
} |
||||
if (!data && h->flags == 0 && from == 1) { |
||||
if (!saved_ack) saved_ack = copy_packet(e); |
||||
if (drop_all_ack || drop_ack) { if (drop_ack) drop_ack--; dropped_ack++; release(e); return; } |
||||
} |
||||
inject(1-from, e); |
||||
} |
||||
static void step(void) { |
||||
uasync_poll(ua, 1); |
||||
for (int i = 0; i < 2; i++) { wire(i, ep[i].conn.send_input_q); wire(i, ep[i].alternate.send_input_q); } |
||||
} |
||||
static int setup(void) { |
||||
memset(ep, 0, sizeof(ep)); failed = 0; ua = uasync_create(); REQUIRE(ua); |
||||
drop_data = drop_ack = drop_all_ack = reorder = duplicated = dropped_data = dropped_ack = handshake_drops = closed_in_handler = 0; |
||||
memset(drop_hello, 0, sizeof(drop_hello)); memset(drop_challenge, 0, sizeof(drop_challenge)); memset(drop_confirm, 0, sizeof(drop_confirm)); |
||||
max_wire = 0; held = saved_data = saved_ack = saved_hello = saved_confirm = NULL; |
||||
for (int i = 0; i < 2; i++) { |
||||
ep[i].inst.ua = ua; ep[i].inst.node_id = i + 1; |
||||
REQUIRE(etcp_router_init(&ep[i].inst) == 0); |
||||
REQUIRE(etcp_router_bind(&ep[i].inst, SVC, receive) == 0); |
||||
REQUIRE(etcp_router_bind(&ep[i].inst, SVC+1, receive) == 0); |
||||
ep[i].conn.instance = ep[i].alternate.instance = &ep[i].inst; |
||||
ep[i].conn.peer_node_id = ep[i].alternate.peer_node_id = 2-i; |
||||
ep[i].conn.send_input_q = queue_new(ua, 0, 0, 0, "wire"); |
||||
ep[i].alternate.send_input_q = queue_new(ua, 0, 0, 0, "alternate_wire"); |
||||
queue_set_waiter_defer(ep[i].conn.send_input_q, 1); queue_set_waiter_defer(ep[i].alternate.send_input_q, 1); |
||||
ep[i].inst.connections = queue_new(ua, 16, 0, 8, "connections"); |
||||
ep[i].route = queue_entry_new(sizeof(struct conn_queue_entry)); |
||||
struct conn_queue_entry* ce = (struct conn_queue_entry*)ep[i].route->data; |
||||
ce->peer_node_id = 2-i; ce->conn = &ep[i].conn; |
||||
queue_data_put_with_index(ep[i].inst.connections, ep[i].route); |
||||
} |
||||
return 0; |
||||
} |
||||
static void cleanup(void) { |
||||
for (int i = 0; i < 2; i++) etcp_router_destroy(&ep[i].inst); |
||||
uasync_poll(ua, 0); |
||||
for (int i = 0; i < 2; i++) { |
||||
struct ll_entry* e; |
||||
while ((e = queue_data_get(ep[i].conn.send_input_q))) release(e); |
||||
while ((e = queue_data_get(ep[i].alternate.send_input_q))) release(e); |
||||
queue_free(ep[i].conn.send_input_q); queue_free(ep[i].alternate.send_input_q); |
||||
release(queue_data_get(ep[i].inst.connections)); queue_free(ep[i].inst.connections); |
||||
} |
||||
release(held); release(saved_data); release(saved_ack); release(saved_hello); release(saved_confirm); |
||||
uasync_poll(ua, 0); |
||||
if (ua->timer_alloc_count != ua->timer_free_count) { |
||||
fprintf(stderr, "FAIL teardown: live timers=%llu\n", (unsigned long long)(ua->timer_alloc_count - ua->timer_free_count)); |
||||
failed = 1; |
||||
} |
||||
uasync_destroy(ua, 0); |
||||
} |
||||
static int send_one(int side, int svc) { |
||||
uint32_t n = ep[side].sent[svc]; uint8_t data[185]; size_t len = 4 + n % 181; |
||||
memcpy(data, &n, 4); |
||||
for (size_t i = 4; i < len; i++) data[i] = n ^ i ^ svc; |
||||
if (etcp_router_conn_send(rc(side, svc), data, len, 0) != 0) return -1; |
||||
ep[side].sent[svc]++; return 0; |
||||
} |
||||
static int settled(int side, int svc) { |
||||
struct ETCP_ROUTER_CONN* c = rc(side, svc); |
||||
return c->tx_acked == c->tx_seq && !c->send_q->count && !c->inflight_q->count; |
||||
} |
||||
static int wait_settled(int side, int svc) { |
||||
uint64_t deadline = get_time_tb() + 40000; |
||||
while (!settled(side, svc) && !failed && get_time_tb() < deadline) step(); |
||||
REQUIRE(!failed); REQUIRE(settled(side, svc)); |
||||
REQUIRE(ep[1-side].received[svc] == ep[side].sent[svc]); |
||||
return 0; |
||||
} |
||||
static int transfer(void) { |
||||
uint64_t deadline = get_time_tb() + 40000; |
||||
int complete = 0, backpressure = 0; |
||||
while (!complete && !failed && get_time_tb() < deadline) { |
||||
complete = 1; |
||||
for (int side = 0; side < 2; side++) for (int svc = 0; svc < 2; svc++) { |
||||
while (ep[side].sent[svc] < COUNT) if (send_one(side, svc) != 0) { backpressure++; break; } |
||||
if (ep[side].sent[svc] < COUNT || !settled(side, svc)) complete = 0; |
||||
} |
||||
step(); |
||||
} |
||||
REQUIRE(!failed); REQUIRE(complete); REQUIRE(backpressure); |
||||
for (int side = 0; side < 2; side++) for (int svc = 0; svc < 2; svc++) { |
||||
REQUIRE(ep[side].received[svc] == COUNT); REQUIRE(ep[side].notified[svc] == 0); |
||||
} |
||||
printf(" delivered=%u backpressure=%d data_drop=%d ack_drop=%d reorder=%d retrans=%u wire_max=%u\n", |
||||
4*COUNT, backpressure, dropped_data, dropped_ack, duplicated, rc(0,0)->c_retrans_done, max_wire); |
||||
return 0; |
||||
} |
||||
static int loss_reorder(void) { |
||||
REQUIRE(setup() == 0); drop_data = 1; drop_ack = 2; reorder = 1; |
||||
REQUIRE(transfer() == 0); REQUIRE(dropped_data == 1 && dropped_ack == 2 && duplicated == 1); |
||||
REQUIRE(rc(0,0)->c_retrans_done > 0); REQUIRE(max_wire < 20); |
||||
cleanup(); return 0; |
||||
} |
||||
static int handshake_loss(void) { |
||||
REQUIRE(setup() == 0); |
||||
for (int i = 0; i < 2; i++) drop_hello[i] = drop_challenge[i] = drop_confirm[i] = 1; |
||||
REQUIRE(transfer() == 0); REQUIRE(handshake_drops == 6); |
||||
cleanup(); return 0; |
||||
} |
||||
static int ack_bounds(void) { |
||||
REQUIRE(setup() == 0); REQUIRE(send_one(0,0) == 0); REQUIRE(wait_settled(0,0) == 0); REQUIRE(saved_ack); |
||||
struct ETCP_ROUTER_CONN* c = rc(0,0); uint64_t progress = c->last_ack_changed_tb; c->no_ack_count = 9; |
||||
uint32_t invalid[] = {0, 1, 2, 0x7fffffffU, 0x80000001U}; |
||||
for (unsigned i = 0; i < sizeof(invalid)/sizeof(invalid[0]); i++) { |
||||
struct ll_entry* e = copy_packet(saved_ack); ((struct SVC_ROUTE_HDR*)e->dgram)->seq = invalid[i]; inject(0,e); |
||||
REQUIRE(c->tx_acked == 1 && c->no_ack_count == 9 && c->last_ack_changed_tb == progress); |
||||
} |
||||
REQUIRE(send_one(0,0) == 0); // queued, not transmitted until deferred waiter runs
|
||||
struct ll_entry* e = copy_packet(saved_ack); ((struct SVC_ROUTE_HDR*)e->dgram)->seq = 2; inject(0,e); |
||||
REQUIRE(c->tx_acked == 1); REQUIRE(wait_settled(0,0) == 0); |
||||
cleanup(); return 0; |
||||
} |
||||
static int reinit_retains(void) { |
||||
REQUIRE(setup() == 0); REQUIRE(send_one(0,0) == 0); REQUIRE(wait_settled(0,0) == 0); |
||||
REQUIRE(send_one(0,0) == 0); uasync_poll(ua,0); |
||||
struct ETCP_ROUTER_CONN* c = rc(0,0); REQUIRE(c->inflight_q->count == 1); |
||||
uint64_t epoch = c->reset_id; |
||||
etcp_router_pause_retrans_for_node(&ep[0].inst, ep[1].inst.node_id); |
||||
struct ll_entry* e; while ((e = queue_data_get(ep[0].conn.send_input_q))) release(e); |
||||
REQUIRE(c->inflight_q->count == 1 && c->retrans_timer && c->reset_id == epoch); |
||||
REQUIRE(wait_settled(0,0) == 0); REQUIRE(c->c_retrans_done > 0); |
||||
REQUIRE(ep[0].notified[0] == 0 && ep[1].notified[0] == 0); |
||||
cleanup(); return 0; |
||||
} |
||||
static int stale_session(void) { |
||||
REQUIRE(setup() == 0); REQUIRE(send_one(0,0) == 0); REQUIRE(wait_settled(0,0) == 0); |
||||
uint64_t previous = rc(0,0)->reset_id; |
||||
etcp_router_conn_restart(&ep[0].inst, 0, ep[1].inst.node_id, SVC); |
||||
REQUIRE(rc(0,0)->reset_id != previous); |
||||
REQUIRE(send_one(0,0) == 0); REQUIRE(wait_settled(0,0) == 0); |
||||
REQUIRE(ep[0].notified[0] == 1 && ep[1].notified[0] == 1); |
||||
uint64_t epoch = rc(1,0)->peer_reset_id; |
||||
inject(1,copy_packet(saved_data)); inject(0,copy_packet(saved_ack)); |
||||
inject(1,copy_packet(saved_hello)); inject(1,copy_packet(saved_confirm)); |
||||
struct ll_entry* e = copy_packet(saved_ack); ((struct SVC_ROUTE_HDR*)e->dgram)->flags = ROUTER_FLAG_CLOSE; inject(0,e); |
||||
for (int i = 0; i < 50; i++) step(); |
||||
REQUIRE(rc(1,0)->peer_reset_id == epoch && ep[1].received[0] == 2); |
||||
REQUIRE(ep[0].notified[0] == 1 && ep[1].notified[0] == 1); |
||||
REQUIRE(send_one(0,0) == 0); REQUIRE(wait_settled(0,0) == 0); |
||||
cleanup(); return 0; |
||||
} |
||||
static int wraparound(void) { |
||||
REQUIRE(setup() == 0); REQUIRE(send_one(0,0) == 0); REQUIRE(wait_settled(0,0) == 0); |
||||
rc(0,0)->tx_seq = rc(0,0)->tx_sent = rc(0,0)->tx_acked = UINT32_MAX-2; |
||||
rc(1,0)->rx_seq = rc(1,0)->last_sent_ack_seq = UINT32_MAX-2; |
||||
for (int i = 0; i < 7; i++) REQUIRE(send_one(0,0) == 0); |
||||
REQUIRE(wait_settled(0,0) == 0); REQUIRE(rc(0,0)->tx_acked == 4); |
||||
cleanup(); return 0; |
||||
} |
||||
static int waiter_lifetime(void) { |
||||
for (int deferred = 0; deferred < 2; deferred++) { |
||||
REQUIRE(setup() == 0); REQUIRE(send_one(0,0) == 0); REQUIRE(wait_settled(0,0) == 0); |
||||
struct ll_queue* old = ep[0].conn.send_input_q; |
||||
if (!deferred) queue_data_put(old, queue_entry_new(0)); |
||||
REQUIRE(send_one(0,0) == 0); |
||||
struct ETCP_ROUTER_CONN* c = rc(0,0); |
||||
REQUIRE(c->send_waiter_q == old); |
||||
REQUIRE(deferred ? c->send_waiter.call_soon_id != NULL : c->send_waiter.internal != NULL); |
||||
((struct conn_queue_entry*)ep[0].route->data)->conn = &ep[0].alternate; |
||||
etcp_router_conn_close(c); uasync_poll(ua,0); |
||||
struct ll_entry* e; while ((e = queue_data_get(old))) release(e); |
||||
REQUIRE(old->waiter_head == NULL); |
||||
cleanup(); |
||||
} |
||||
return 0; |
||||
} |
||||
static int close_from_handler(void) { |
||||
REQUIRE(setup() == 0); closed_in_handler = 1; REQUIRE(send_one(0,0) == 0); |
||||
uint64_t deadline = get_time_tb()+10000; |
||||
while (!ep[1].received[0] && get_time_tb()<deadline) step(); |
||||
REQUIRE(ep[1].received[0] == 1); REQUIRE(!failed); |
||||
cleanup(); return 0; |
||||
} |
||||
static void producer_closes(struct ll_queue* q, void* arg) { |
||||
(void)q; |
||||
struct ETCP_ROUTER_CONN* c = arg; |
||||
etcp_router_conn_close(c); |
||||
} |
||||
static int close_from_producer(void) { |
||||
REQUIRE(setup() == 0); REQUIRE(send_one(0,0) == 0); REQUIRE(wait_settled(0,0) == 0); |
||||
struct ETCP_ROUTER_CONN* c = rc(0,0); |
||||
for (int i=0; i<ROUTER_MAX_SEND_Q_PACKETS; i++) REQUIRE(send_one(0,0) == 0); |
||||
struct queue_waiter_handle waiter = {0}; |
||||
etcp_router_on_send_ready(&ep[0].inst, 0, ep[1].inst.node_id, SVC, &waiter, producer_closes, c); |
||||
REQUIRE(waiter.internal); |
||||
step(); |
||||
REQUIRE(ep[0].notified[0] == 1 && ep[0].inst.router_conns->count == 0); |
||||
REQUIRE(!waiter.internal && !waiter.call_soon_id); |
||||
cleanup(); return 0; |
||||
} |
||||
static int route_switch(void) { |
||||
REQUIRE(setup() == 0); REQUIRE(send_one(0,0) == 0); REQUIRE(wait_settled(0,0) == 0); |
||||
struct ll_queue* old = ep[0].conn.send_input_q; |
||||
queue_data_put(old, queue_entry_new(0)); REQUIRE(send_one(0,0) == 0); |
||||
REQUIRE(rc(0,0)->send_waiter.internal); |
||||
((struct conn_queue_entry*)ep[0].route->data)->conn = &ep[0].alternate; |
||||
uint64_t deadline = get_time_tb()+20000; |
||||
while (!settled(0,0) && get_time_tb()<deadline) { |
||||
uasync_poll(ua,1); wire(0,ep[0].alternate.send_input_q); wire(1,ep[1].conn.send_input_q); |
||||
} |
||||
REQUIRE(settled(0,0) && ep[1].received[0] == 2); REQUIRE(old->waiter_head == NULL); |
||||
cleanup(); return 0; |
||||
} |
||||
static int no_route_restore(void) { |
||||
REQUIRE(setup() == 0); REQUIRE(send_one(0,0) == 0); REQUIRE(wait_settled(0,0) == 0); |
||||
struct conn_queue_entry* route = (struct conn_queue_entry*)ep[0].route->data; |
||||
route->conn = NULL; REQUIRE(send_one(0,0) == 0); |
||||
REQUIRE(rc(0,0)->no_route && rc(0,0)->send_q->count == 1); |
||||
for (int i=0; i<30; i++) step(); |
||||
route->conn = &ep[0].alternate; REQUIRE(wait_settled(0,0) == 0); |
||||
REQUIRE(ep[0].notified[0] == 0 && ep[1].notified[0] == 0); |
||||
cleanup(); return 0; |
||||
} |
||||
static int timeout_with_duplicate_ack(void) { |
||||
REQUIRE(setup() == 0); drop_all_ack = 1; REQUIRE(send_one(0,0) == 0); |
||||
uint64_t deadline = get_time_tb()+70000, last_dup = 0; |
||||
while (!ep[0].notified[0] && get_time_tb()<deadline) { |
||||
if (saved_ack && get_time_tb()-last_dup >= 500) { |
||||
struct ll_entry* e = copy_packet(saved_ack); ((struct SVC_ROUTE_HDR*)e->dgram)->seq = 0; |
||||
inject(0,e); last_dup = get_time_tb(); |
||||
} |
||||
step(); |
||||
} |
||||
REQUIRE(ep[0].notified[0] == 1); REQUIRE(ep[1].received[0] == 1); |
||||
REQUIRE(dropped_ack >= ROUTER_NO_ACK_MAX_RETRANS-1); |
||||
REQUIRE(ep[0].inst.router_conns->count == 0); |
||||
cleanup(); return 0; |
||||
} |
||||
static int simultaneous_restart(void) { |
||||
REQUIRE(setup() == 0); REQUIRE(send_one(0,0) == 0); REQUIRE(send_one(1,0) == 0); |
||||
REQUIRE(wait_settled(0,0) == 0); REQUIRE(wait_settled(1,0) == 0); |
||||
for (int i=0; i<2; i++) etcp_router_conn_restart(&ep[i].inst, 0, ep[1-i].inst.node_id, SVC); |
||||
REQUIRE(send_one(0,0) == 0); REQUIRE(send_one(1,0) == 0); |
||||
REQUIRE(wait_settled(0,0) == 0); REQUIRE(wait_settled(1,0) == 0); |
||||
REQUIRE(ep[0].notified[0] == 1 && ep[1].notified[0] == 1); |
||||
cleanup(); return 0; |
||||
} |
||||
static int size_bounds(void) { |
||||
REQUIRE(setup() == 0); |
||||
uint8_t* data = u_calloc(1, PKTNORM_MAX_DGRAM_SIZE); |
||||
REQUIRE(data); |
||||
struct ETCP_ROUTER_CONN* c = rc(0,0); |
||||
REQUIRE(etcp_router_conn_send(c, data, 0, 0) == -1); |
||||
REQUIRE(etcp_router_conn_send(c, data, PKTNORM_MAX_DGRAM_SIZE - SVC_ROUTE_HDR_SIZE + 1, 0) == -1); |
||||
REQUIRE(etcp_router_conn_send(c, data, PKTNORM_MAX_DGRAM_SIZE - SVC_ROUTE_HDR_SIZE, ROUTER_FLAG_SIGNED) == -1); |
||||
REQUIRE(etcp_router_conn_send(c, data, SIZE_MAX, 0) == -1); |
||||
REQUIRE(c->tx_seq == 0 && c->send_q->count == 0); |
||||
REQUIRE(etcp_router_conn_send(c, data, PKTNORM_MAX_DGRAM_SIZE - SVC_ROUTE_HDR_SIZE, 0) == 0); |
||||
u_free(data); cleanup(); return 0; |
||||
} |
||||
|
||||
int main(void) { |
||||
debug_config_init(); debug_set_categories(DEBUG_CATEGORY_ALL); debug_set_level(DEBUG_LEVEL_WARN); |
||||
const char* log = getenv("ROUTER_TEST_LOG"); |
||||
if (log) { debug_enable_file_output(log, 1); debug_set_category_level(DEBUG_CATEGORY_ETCPROUTE, DEBUG_LEVEL_DEBUG); } |
||||
struct { const char* name; int (*run)(void); } tests[] = { |
||||
{"loss/reorder/duplicate and duplex multiplexing", loss_reorder}, {"lost HELLO/challenge/confirm", handshake_loss}, |
||||
{"ACK bounds and no false progress", ack_bounds}, {"transport reinit retains inflight", reinit_retains}, |
||||
{"restart rejects old DATA/ACK/CLOSE/HELLO/confirm", stale_session}, {"sequence wraparound", wraparound}, |
||||
{"waiter lifetime after route switch", waiter_lifetime}, {"service closes during delivery", close_from_handler}, {"route switch while old next hop blocked", route_switch}, |
||||
{"no route then restore without data loss", no_route_restore}, {"duplicate ACK cannot prevent timeout", timeout_with_duplicate_ack}, |
||||
{"simultaneous endpoint restart", simultaneous_restart}, {"payload size boundaries", size_bounds}, {"producer closes during send queue drain", close_from_producer} |
||||
}; |
||||
for (unsigned i=0; i<sizeof(tests)/sizeof(tests[0]); i++) { |
||||
printf("TEST %u: %s\n", i+1, tests[i].name); fflush(stdout); |
||||
if (tests[i].run() != 0 || failed) return 1; |
||||
puts("PASS"); |
||||
} |
||||
puts("All router fault scenarios passed"); return 0; |
||||
} |
||||
@ -0,0 +1,3 @@
|
||||
/* Same content/order/backpressure checks over real STCP, with simultaneous reset. */ |
||||
#define ROUTER_TEST_TCP_DEFAULT 1 |
||||
#include "test_etcp_router.c" |
||||
@ -0,0 +1,162 @@
|
||||
/* Real ETCP lifecycle, AES-CCM and common ingress; deterministic encrypted wire faults. */ |
||||
#include <stdio.h> |
||||
#include <stdlib.h> |
||||
#include <string.h> |
||||
#include "etcp.h" |
||||
#include "etcp_api.h" |
||||
#include "etcp_connections.h" |
||||
#include "etcp_session.h" |
||||
#include "config_parser.h" |
||||
#include "../lib/debug_config.h" |
||||
#include "../lib/mem.h" |
||||
|
||||
#define CHECK(x) do { if (!(x)) { fprintf(stderr,"FAIL %s:%d: %s\n",__func__,__LINE__,#x); exit(1); } } while (0) |
||||
struct endpoint { struct UTUN_INSTANCE inst; struct utun_config cfg; struct ETCP_CONN* c; struct ETCP_SOCKET sock; struct ETCP_LINK links[2]; }; |
||||
struct frame { int from; size_t len; uint8_t bytes[PACKET_DATA_SIZE]; }; |
||||
static struct endpoint ep[2]; |
||||
static struct UASYNC* ua; |
||||
static struct frame frames[4096], saved[4]; |
||||
static unsigned head, tail, drops, reinit[2], down[2]; |
||||
static int drop[2][3], blackhole, dead_primary; |
||||
|
||||
static ssize_t capture(socket_t fd, const void* data, size_t len, const struct sockaddr* addr, socklen_t alen, |
||||
struct ETCP_LINK* link, void* arg) { |
||||
(void)fd; (void)addr; (void)alen; |
||||
int from = (int)(intptr_t)arg; |
||||
CHECK(len <= sizeof(frames[0].bytes)); |
||||
if (blackhole || (dead_primary && link->local_link_id == 1)) { drops++; return len; } |
||||
CHECK(tail-head < 4096); |
||||
struct frame* f = &frames[tail++ % 4096]; f->from = from; f->len = len; memcpy(f->bytes,data,len); |
||||
return len; |
||||
} |
||||
static void event(struct ETCP_CONN* c, int ev, void* arg) { |
||||
(void)c; int i = (int)(intptr_t)arg; |
||||
if (ev == ETCP_CBK_EVENT_REINIT) reinit[i]++; |
||||
if (ev == ETCP_CBK_EVENT_DOWN) down[i]++; |
||||
} |
||||
static void deliver(const struct frame* f, int allow_drop) { |
||||
int to = 1-f->from; |
||||
struct ETCP_DGRAM* p = memory_pool_alloc(ep[to].inst.pkt_pool); CHECK(p); |
||||
size_t len = 0; |
||||
CHECK(sc_decrypt(&ep[to].c->crypto_ctx,f->bytes,f->len,(uint8_t*)&p->timestamp,&len) == SC_OK); |
||||
CHECK(len >= 4); |
||||
int kind = p->data[0] - ETCP_SESSION_HELLO; |
||||
if (allow_drop && kind >= 0 && kind < 3 && drop[f->from][kind]) { |
||||
drop[f->from][kind]--; drops++; memory_pool_free(ep[to].inst.pkt_pool,p); return; |
||||
} |
||||
if (kind == 0 && !saved[0].len) saved[0] = *f; |
||||
if (kind == 2 && !saved[1].len) saved[1] = *f; |
||||
CHECK(etcp_packet_decrypted(&ep[to].sock,p,&ep[to].links[0],len) == 0); |
||||
} |
||||
static void pump(void) { |
||||
uasync_poll(ua,1); |
||||
unsigned end = tail; // callbacks can append replies
|
||||
while (head < end) { struct frame f = frames[head++ % 4096]; deliver(&f,1); } |
||||
} |
||||
static void converge(void) { |
||||
uint64_t deadline = get_time_tb()+40000; |
||||
while (get_time_tb() < deadline) { |
||||
pump(); |
||||
if (!ep[0].c->reinit_pending && !ep[1].c->reinit_pending && !ep[0].c->session_timer && !ep[1].c->session_timer && head == tail) break; |
||||
} |
||||
for (int i=0;i<2;i++) { |
||||
CHECK(ep[i].c->initialized && !ep[i].c->session_required && !ep[i].c->reinit_pending); |
||||
CHECK(!ep[i].c->session_timer && !ep[i].c->close_requested && !down[i]); |
||||
CHECK(ep[i].c->peer_reset_id == ep[1-i].c->reset_id); |
||||
} |
||||
} |
||||
static void setup(int multi) { |
||||
memset(ep,0,sizeof(ep)); memset(drop,0,sizeof(drop)); memset(saved,0,sizeof(saved)); |
||||
memset(reinit,0,sizeof(reinit)); memset(down,0,sizeof(down)); head=tail=drops=0; blackhole=dead_primary=0; |
||||
ua = uasync_create(); CHECK(ua); |
||||
for (int i=0;i<2;i++) { |
||||
ep[i].inst.ua=ua; ep[i].inst.node_id=i+1; ep[i].inst.config=&ep[i].cfg; |
||||
ep[i].inst.connections=queue_new(ua,16,0,8,"session_test"); |
||||
ep[i].inst.pkt_pool=memory_pool_init(sizeof(struct ETCP_DGRAM)+PACKET_DATA_SIZE,"session_packets"); |
||||
CHECK(sc_generate_keypair(&ep[i].inst.my_keys) == SC_OK); |
||||
ep[i].c=etcp_connection_create(&ep[i].inst,"session_test"); CHECK(ep[i].c); |
||||
ep[i].c->peer_node_id=2-i; ep[i].c->links_up=1; |
||||
etcp_conn_ready(ep[i].c); |
||||
etcp_conn_add_cbk(ep[i].c,event,(void*)(intptr_t)i,ETCP_CBK_EVENT_REINIT|ETCP_CBK_EVENT_DOWN); |
||||
CHECK(sc_init_ctx(&ep[i].c->crypto_ctx,&ep[i].inst.my_keys) == SC_OK); |
||||
ep[i].sock.instance=&ep[i].inst; |
||||
for (int j=0;j<=multi;j++) { |
||||
struct ETCP_LINK* l=&ep[i].links[j]; |
||||
l->etcp=ep[i].c; l->conn=&ep[i].sock; l->local_link_id=j+1; |
||||
l->initialized=l->link_status=l->recv_keepalive=l->remote_keepalive=1; |
||||
l->link_state=LINK_STATE_CONNECTED; l->mtu=1280; l->remote_addr.ss_family=AF_INET; |
||||
l->send_hook=capture; l->send_hook_ctx=(void*)(intptr_t)i; |
||||
} |
||||
if (multi) ep[i].links[0].next=&ep[i].links[1]; |
||||
ep[i].c->links=&ep[i].links[0]; |
||||
} |
||||
for (int i=0;i<2;i++) { |
||||
ep[i].c->peer_reset_id=ep[1-i].c->reset_id; |
||||
CHECK(sc_set_peer_public_key(&ep[i].c->crypto_ctx,ep[1-i].inst.my_keys.public_key,0) == SC_OK); |
||||
} |
||||
} |
||||
static void teardown(void) { |
||||
for (int i=0;i<2;i++) { |
||||
ep[i].c->links=NULL; // fixture owns synthetic authenticated links
|
||||
etcp_connection_close(ep[i].c); |
||||
} |
||||
uasync_poll(ua,0); |
||||
for (int i=0;i<2;i++) { queue_free(ep[i].inst.connections); memory_pool_destroy(ep[i].inst.pkt_pool); } |
||||
uasync_poll(ua,0); CHECK(ua->timer_alloc_count == ua->timer_free_count); uasync_destroy(ua,0); |
||||
} |
||||
static struct frame stream_packet(int data) { |
||||
struct ETCP_DGRAM* p=memory_pool_alloc(ep[0].inst.pkt_pool); CHECK(p); |
||||
p->link=&ep[0].links[0]; p->noencrypt_len=0; p->data_len=8; memset(p->data,0,8); p->data[0]=data ? 0 : 1; if (data) p->data[1]=99; |
||||
CHECK(etcp_encrypt_send(p)>0); memory_pool_free(ep[0].inst.pkt_pool,p); |
||||
CHECK(tail == head+1); struct frame f=frames[head++ % 4096]; return f; |
||||
} |
||||
int main(void) { |
||||
debug_config_init(); debug_set_level(DEBUG_LEVEL_WARN); |
||||
const char* log=getenv("SESSION_TEST_LOG"); |
||||
if (log) { debug_enable_file_output(log,1); debug_set_category_level(DEBUG_CATEGORY_ETCP,DEBUG_LEVEL_DEBUG); } |
||||
size_t baseline=u_get_allocated_count(); |
||||
for (int scenario=0;scenario<9;scenario++) { |
||||
setup(scenario == 5); |
||||
saved[2]=stream_packet(0); saved[3]=stream_packet(1); |
||||
if (scenario == 2 || scenario == 3) for (int i=0;i<2;i++) for (int j=0;j<3;j++) drop[i][j]=1; |
||||
if (scenario == 5) dead_primary=1; |
||||
etcp_conn_fatal_reinit(ep[scenario == 1 ? 1 : 0].c,"fault test"); |
||||
if (scenario == 3) etcp_conn_fatal_reinit(ep[1].c,"simultaneous reset"); |
||||
if (scenario == 4) { pump(); etcp_conn_fatal_reinit(ep[0].c,"second reset before confirmation"); } |
||||
if (scenario == 6) { |
||||
blackhole=1; ep[0].c->session_started=get_time_tb()-100000; |
||||
CHECK(etcp_conn_ref_take(ep[0].c)==0); |
||||
ep[0].c->links=NULL; // timeout closes a real connection, links belong to fixture
|
||||
while (ep[0].c->state!=2) pump(); |
||||
CHECK(!ep[0].c->session_timer && ep[0].c->delete_complete); |
||||
etcp_conn_ref_free(ep[0].c); ep[0].c=NULL; |
||||
ep[1].c->links=NULL; etcp_connection_close(ep[1].c); uasync_poll(ua,0); |
||||
for (int i=0;i<2;i++) { queue_free(ep[i].inst.connections); memory_pool_destroy(ep[i].inst.pkt_pool); } |
||||
uasync_poll(ua,0); CHECK(ua->timer_alloc_count==ua->timer_free_count); uasync_destroy(ua,0); |
||||
} else { |
||||
converge(); |
||||
if (scenario==2 || scenario==3 || scenario==5) CHECK(drops>0); |
||||
uint64_t peer=ep[1].c->peer_reset_id, local=ep[0].c->reset_id; |
||||
unsigned before0=reinit[0], before1=reinit[1]; |
||||
ep[1].links[0].last_recv_local_time=123; |
||||
deliver(&saved[2],0); deliver(&saved[3],0); CHECK(ep[1].links[0].last_recv_local_time==123); // stale ACK never reaches stream/liveness
|
||||
if (scenario==7) { |
||||
if (saved[0].len) deliver(&saved[0],0); |
||||
if (saved[1].len) { deliver(&saved[1],0); deliver(&saved[1],0); } |
||||
converge(); |
||||
} |
||||
if (scenario==8) { |
||||
uint8_t short_frame[25]={ETCP_SESSION_CONFIRM}; |
||||
CHECK(etcp_session_receive(ep[0].c,short_frame,1)); |
||||
CHECK(!etcp_session_accept_data(ep[0].c,short_frame,0)); |
||||
etcp_conn_apply_peer_reset_id(ep[0].c,local); // stale/unproven physical INIT cannot switch peer
|
||||
CHECK(ep[0].c->peer_reset_id==ep[1].c->reset_id); |
||||
} |
||||
CHECK(ep[1].c->peer_reset_id==peer && ep[0].c->reset_id==local); |
||||
CHECK(reinit[0]==before0 && reinit[1]==before1); |
||||
teardown(); |
||||
} |
||||
CHECK(u_get_allocated_count()==baseline); printf("PASS session scenario %d\n",scenario+1); |
||||
} |
||||
puts("All ETCP session scenarios passed"); return 0; |
||||
} |
||||
@ -0,0 +1,3 @@
|
||||
/* Real a-b-c failure: no restart/reconnect help from the test. */ |
||||
#define ROUTER_TEST_RECOVERY_ONLY 1 |
||||
#include "test_etcp_router_reconnect.c" |
||||
Loading…
Reference in new issue