From 7d9cbad8c5500715d6de8155ddbd2e90e7811539 Mon Sep 17 00:00:00 2001 From: evgeny Date: Wed, 29 Jul 2026 16:53:30 +0300 Subject: [PATCH] media_delivery: add architecture description and implementation plan --- src/media_delivery/_desc.txt | 47 +++ src/media_delivery/_plan.md | 597 +++++++++++++++++++++++++++++++++++ 2 files changed, 644 insertions(+) create mode 100644 src/media_delivery/_desc.txt create mode 100644 src/media_delivery/_plan.md diff --git a/src/media_delivery/_desc.txt b/src/media_delivery/_desc.txt new file mode 100644 index 00000000..61cc74e2 --- /dev/null +++ b/src/media_delivery/_desc.txt @@ -0,0 +1,47 @@ +Архитектура распространения медиа + +1. src (создатель медиа) - подписывает, разделяет на блоки, создает UUID, отправляет в сообщении информацию о всех блоках для возможности скачивания и валидации (проверки сумм блоков и подписи файла). +2. req узел - кто хочет получить медиа. + + +основаня идея распространения контента +- каждый узел выбирает удобный (по пигну) на старте. выбирает следующий если таймаут (2 сек) + и устанавливает к нему прямое подключение. уведомляет суперузел что он обслуживает этот узел. + Суперузлы хранят в памяти список узлов (ll_queue) которые они обслуживают. Если соединение рвётся - узел автоматически выпадает из списка обслуживаемых. + +- req-узел когда принял полный блок - уведомляет суперузел об этом. + +- суперузлы между собой реплицируют таблицу с узлами для скачивания. + Также если подключен автор контекнта то суперузел автору контента отправляет обновления о доступности его блоков. + + На всех узлах есть таблица с узлами для скачивания блоков. + ID primary/autincrement | UUID блока | id группы | id ноды | timestamp когда добавлено-обновлено + Суперузлы в нее сохраняют весь контент групп в которых они состоят. обычные узлы сохраняют только свой контент. + Если нет суперузлов или они недоступны - то роль суперузла выполняет автор контента. + + + + +Репликация между суперузлами. + +Каждый суперузел хранит указатель (id) по какой элемент синхронизировано с каждым другим суперузлом. +отправляет злементы и ждёт подтверждения (не более 4 в очереди). переповтор - 10 сек без ответа. +Суперузлы мониторят появление других суперузлов в bgp таблице. когда узел появился пробуют установить прямое подключение (без посредников и без обратного подключения - это сделает другой суперузел сам). +Если не получается 10 сек - отключаемся и следующая попытка через час. +то есть как только видят повяление - пробуют установить прямую связь каждый с каждым. +Прямое подключение - это независимый автономный процесс. +Для этого нужен коллбэк из bgp - появился новый узел и обновился узел. Его надо добавить если нет с возможностью цепочки подписчиков (как коллбэки статусов в etcp). + +при появлении суперузлата также начинаем репликацию таблицы блоков - отдельный асинхронный процесс. +для каждого суперузла свой контекст (таймаут, состояние протокола). отправляем обновление и ждём подтверждение с таймаутом (2 сек). таумаут постепенно увеличиваем до 10 минут если нет подтверждений и сбрасываем если появилось. + + +Алгоритм. +при получении сообщения (в зависимости от настроек кеширования) узлы получившие сообщение с медиа хотят получить вложение. Порядок действий: + - выбирают суперузел с минимальным rtt и отправляют ему запрос - хочу скачать медиа. дай список узлов с которых качаем. + В запросе - UUID запрашиваемых блоков (обязательно один файл можно несколько блоков для него), id src узла. + запрос направляется в суперузел изи автору контента - если нет суперузлов или они недоступны. таймаут 2 сек. + далее выбираем лучшие узлы (не более 10 штук по пингу) и пробуем к ним установить подключение (прямое или через промежуточный узел). + перебираем пытаясь подключиться, пока не наберем нужное количество (по количеству блоков, но не более 10 штук) + к кому подключились - запрашиваем блоки также с выбором следующего при таймауте или неудаче. + отправляем им запрос и пытаемся скачать. diff --git a/src/media_delivery/_plan.md b/src/media_delivery/_plan.md new file mode 100644 index 00000000..f68f87ed --- /dev/null +++ b/src/media_delivery/_plan.md @@ -0,0 +1,597 @@ +# Архитектура media_delivery — План реализации + +## Транспортная модель + +- **Все сообщения** (управление + данные блоков) — через `etcp_router` (`etcp_route_send`) +- Перед передачей блоков устанавливается прямое соединение через `conn_mgr` (минимизация задержек) +- После установки `etcp_router` автоматически маршрутизирует через прямое соединение + +## Что уже есть в коде (используем, не дублируем) + +| Компонент | Для чего | +|---|---| +| `peers_*.node_type=4` + `adm_tags="supernode=yes"` | Определение суперузла (задаётся админом в чатгуи) | +| `member_sync` | Синхронизация мемберов (включая adm_tags) | +| `conn_mgr` | Прямые подключения между узлами (direct/reverse/indirect) | +| `etcp_router` (ETCP_RT_ID_MEDIA_DELIVERY=0x07) | Надёжная доставка сервисных сообщений (уже забинжено) | +| `media_index` | Регистрация блоков в SQLite (уже реализовано) | +| `etcp_add_conn_status_cbk` | Мониторинг статуса ETCP-соединений | + +## Новые модули + +``` +src/media_delivery/ + media_delivery.h/c — расширен: контекст, протокол, init/destroy + media_delivery_proto.h — packed-структуры протокольных пакетов + media_download.h/c — логика скачивания блоков +``` + +## Этап 1: BGP-коллбэк в `topo_group` + +Уведомление о появлении/обновлении/удалении узлов. Нужен для обнаружения новых суперузлов. + +**Файлы:** `src/routing_layer/topo_group.h`, `src/routing_layer/topo_group.c` + +Добавить цепочку подписчиков (аналогично `etcp_status_cbk_entry`): + +```c +typedef void (*topo_node_event_fn)(struct TOPO_GROUP* group, uint64_t node_id, + int event, void* arg); +// event: 0=NEW, 1=UPDATE, 2=REMOVE + +void topo_group_add_node_cbk(struct TOPO_GROUP* group, topo_node_event_fn fn, void* arg); +void topo_group_remove_node_cbk(struct TOPO_GROUP* group, topo_node_event_fn fn, void* arg); +``` + +Вызывать в `topo_group_process_nodeinfo` при создании/обновлении `TOPO_NODEQ`, и в `topo_group_process_withdraw` при удалении. + +--- + +## Этап 2: SQLite-таблицы доступности блоков + +Таблицы в `topo_sqlite_db` (уже открыт в `utun_instance`). **Создаются всегда** при `media_delivery_init`, независимо от `is_supernode`. У обычного узла просто пустые. + +### `block_availability` — какие узлы имеют блоки + +```sql +CREATE TABLE IF NOT EXISTS block_availability ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + block_uuid BLOB NOT NULL, -- 16 байт + group_id INTEGER NOT NULL, -- идентификатор группы + node_id INTEGER NOT NULL, -- узел-обладатель блока + chunk INTEGER NOT NULL, -- номер чанка в файле + timestamp INTEGER NOT NULL -- когда добавлено/обновлено (unix epoch) +); +CREATE UNIQUE INDEX IF NOT EXISTS idx_ba_uuid_node + ON block_availability(block_uuid, node_id); +CREATE INDEX IF NOT EXISTS idx_ba_id ON block_availability(id); +CREATE INDEX IF NOT EXISTS idx_ba_node_id ON block_availability(node_id); +``` + +### `super_sync` — последний полученный ID от пира (хранит только принимающая сторона) + +```sql +CREATE TABLE IF NOT EXISTS super_sync ( + peer_node_id INTEGER PRIMARY KEY, -- другой суперузел + last_recv_id INTEGER NOT NULL DEFAULT 0 -- последний ID из БАЗЫ ПИРА, полученный мной +); + +-- last_recv_id — это НЕ локальный ID, а id в таблице block_availability пира. +-- Используется чтобы сказать пиру при SUPER_HELLO: «я принял твои записи до ID X». +-- Пиру это говорит с какого id начинать следующую отправку. +``` + +Отправитель не хранит состояние — он stateless. При подключении обмениваются `SUPER_HELLO`: +A говорит «я получил от тебя до ID X», B говорит «я получил от тебя до ID Y». +Каждая сторона начинает отправку с `(чужой last_recv_id) + 1`. + +### Операции над `block_availability` (прямые SQL-запросы) + +| Операция | SQL | +|---|---| +| Добавить/обновить блок | `INSERT OR REPLACE INTO block_availability(block_uuid,group_id,node_id,chunk,timestamp) VALUES(?,?,?,?,?)` | +| Удалить все блоки узла | `DELETE FROM block_availability WHERE node_id=?` | +| Удалить чужие блоки (выключение суперузла) | `DELETE FROM block_availability WHERE node_id != ?` | +| Удалить блоки группы | `DELETE FROM block_availability WHERE group_id=?` | +| Найти узлы с блоком | `SELECT node_id, timestamp FROM block_availability WHERE block_uuid=? ORDER BY timestamp DESC` | +| Получить записи с ID > N | `SELECT * FROM block_availability WHERE id > ? ORDER BY id ASC LIMIT ?` | +| sync_ptr пира | `SELECT last_recv_id FROM super_sync WHERE peer_node_id=?` — возвращает ID из **чужой базы** | +| обновить sync_ptr | `INSERT OR REPLACE INTO super_sync(peer_node_id,last_recv_id) VALUES(?,?)` — сохраняет ID из **чужой базы** | + +**Плюсы:** не теряем данные при рестарте, не нужен отдельный модуль, проще отладка через `sqlite3`. +**Ограничение:** все SQL-запросы только из uasync-потока (уже так). + +--- + +## Этап 2.1: Коллбэк изменения свойств узла в `member_sync` + +Многоподписочный коллбэк для отслеживания изменений `adm_tags` любого узла (включая себя). +Нужен `media_delivery` для динамического включения/выключения режима суперузла. + +**Файлы:** `src/chat/member_sync.h`, `src/chat/member_sync.c` + +```c +typedef void (*node_props_changed_fn)(uint64_t node_id, const char* adm_tags, void* arg); + +void member_sync_add_props_cbk(node_props_changed_fn fn, void* arg); +void member_sync_remove_props_cbk(node_props_changed_fn fn, void* arg); +``` + +Вызывается при получении обновлённых `adm_tags` через `MSG_ITEM_UPDATE` (в `member_sync.c` в функции обработки `adm_tags` для любого узла). + +`media_delivery` подписывается и в коллбэке: +- Если `node_id == self` и изменился факт `supernode=yes` → переключить `is_supernode` (см. Init) +- Если `node_id != self`: проверить стал/перестал быть суперузлом → запустить/остановить репликацию с ним + +--- + +## Этап 3: `media_delivery_proto.h` + +Packed-структуры подкоманд протокола (первый байт данных после etcp_router): + +```c +enum { + MEDIA_SUBCMD_SERVE_REG = 0x01, // req-узел→суперузел: обслуживай меня + MEDIA_SUBCMD_SERVE_ACK = 0x02, // суперузел→req-узел: подтверждение + MEDIA_SUBCMD_SERVE_LEAVE = 0x03, // req-узел→суперузел: отключение + MEDIA_SUBCMD_QUERY = 0x04, // req-узел→суперузел: кто имеет блоки? + MEDIA_SUBCMD_QUERY_RESP = 0x05, // суперузел→req-узел: список узлов + MEDIA_SUBCMD_HAVE_BLOCK = 0x06, // узел→суперузел: у меня есть блок + MEDIA_SUBCMD_HAVE_BLOCK_ACK = 0x07, // суперузел→узел: подтверждение приёма + MEDIA_SUBCMD_SUPER_REPL = 0x08, // суперузел↔суперузел: репликация записей + MEDIA_SUBCMD_SUPER_ACK = 0x09, // подтверждение репликации + MEDIA_SUBCMD_BLOCK_REQ = 0x0A, // req-узел→блок-холдер: дай блок + MEDIA_SUBCMD_BLOCK_CHUNK = 0x0B, // блок-холдер→req-узел: чанк данных (1KB) + MEDIA_SUBCMD_BLOCK_DONE = 0x0C, // блок-холдер→req-узел: блок передан полностью + MEDIA_SUBCMD_CANCEL = 0x0D, // req-узел→блок-холдер: отмена передачи блока + MEDIA_SUBCMD_SUPER_HELLO = 0x0E, // суперузел↔суперузел: рукопожатие при подключении +}; +``` + +Packed-структуры пакетов: +- `MEDIA_SERVE_REG`: `subcmd:1, group_id:8` +- `MEDIA_SERVE_ACK`: `subcmd:1, group_id:8, status:1` +- `MEDIA_SERVE_LEAVE`: `subcmd:1, group_id:8` +- `MEDIA_QUERY`: `subcmd:1, group_id:8, media_id:16, num_blocks:2, block_ids[num*16]` + +- `MEDIA_QUERY_RESP`: `subcmd:1, num_entries:2, entries[num × {node_id:8, rtt:2, block_id:16, chunk:4}]` + +- `MEDIA_HAVE_BLOCK`: `subcmd:1, group_id:8, media_id:16, block_id:16, chunk:4, timestamp:8, node_sign:64` +- `MEDIA_HAVE_BLOCK_ACK`: `subcmd:1, media_id:16, block_id:16, status:1` (0=OK, 1=dup/already) +- `MEDIA_SUPER_REPL`: `subcmd:1, seq:4, num_entries:2, entries[num × repl_entry]` +- `MEDIA_SUPER_ACK`: `subcmd:1, ack_seq:4` +- `MEDIA_BLOCK_REQ`: `subcmd:1, media_id:16, block_id:16, chunk:4, offset:8` +- `MEDIA_BLOCK_CHUNK`: `subcmd:1, media_id:16, block_id:16, chunk:4, offset:4, data_len:2, data[data_len]` + +- `MEDIA_BLOCK_DONE`: `subcmd:1, media_id:16, block_id:16, chunk:4, total_size:4, block_sig:64` +- `MEDIA_CANCEL`: `subcmd:1, media_id:16, block_id:16, chunk:4` +- `MEDIA_SUPER_HELLO`: `subcmd:1, last_recv_id:8` + + +**Файл:** `src/media_delivery/media_delivery_proto.h` + +--- + +## Этап 4: Расширение `media_delivery_ctx` и init/destroy + +```c +#define MEDIA_MAX_REPL_INFLIGHT 4 +#define MEDIA_REPL_TIMEOUT_TB 20000 // начальный таймаут 2s (0.1ms) +#define MEDIA_REPL_TIMEOUT_MAX_TB 6000000 // макс 10 min +#define MEDIA_MAX_DOWNLOADS 10 +#define MEDIA_QUERY_TIMEOUT_TB 20000 // таймаут QUERY 2s +#define MEDIA_HAVE_BLOCK_TIMEOUT_TB 20000 // таймаут HAVE_BLOCK_ACK 2s +#define MEDIA_RECONNECT_COOLDOWN_TB 360000000 // 1 час (0.1ms) + +struct media_super_peer { + struct ll_entry ll; // индекс по peer_node_id (8 байт) + uint64_t peer_node_id; + uint64_t peer_last_recv_id; // что пир говорит он получил от нас (из SUPER_HELLO / SUPER_ACK) + uint32_t timeout_tb; // текущий таймаут (растёт при ошибках) + uint8_t inflight_count; // пакетов в полёте (макс 4) + uint8_t connected; // 1 = прямое подключение установлено + uint8_t hello_done; // 1 = SUPER_HELLO обмен завершён, можно реплицировать + void* timeout_timer; // таймер ретрансмита + void* connect_timer; // таймер переподключения (1 час при ошибке) +}; + +struct media_served_node { + struct ll_entry ll; // индекс по node_id (8 байт) + uint64_t node_id; + uint64_t group_id; + int64_t joined_at; // timestamp регистрации +}; + +struct media_download_peer { + uint64_t node_id; + uint8_t connected; // 1 = conn_mgr подключён + int num_blocks; // сколько блоков назначено этому узлу + struct { + uint8_t block_id[16]; + uint8_t started; // 1 = BLOCK_REQ отправлен, стрим активен + uint8_t received; // 1 = BLOCK_DONE получен + uint8_t validated; // 1 = подпись проверена + } blocks[64]; +}; + +struct media_download { + struct ll_entry ll; // индекс по media_id (16 байт) + uint8_t media_id[16]; + uint64_t group_id; + char dest_path[1024]; // путь для сохранения файла (как в media_index) + char media_base[512]; // базовый путь media-файлов + int num_blocks; + uint8_t* block_ids; // num_blocks * 16 + uint8_t* block_sigs; // num_blocks * 64 + uint8_t content_hash[32]; + int64_t file_size; + int64_t block_size; + int blocks_received; + int blocks_validated; + int num_peers; + struct media_download_peer peers[10]; + uint8_t active; // 1 = загрузка в процессе + uint8_t assembled; // 1 = файл собран + int err; // 0=OK, <0=ошибка + void* timeout_timer; // общий таймаут загрузки + void (*done_cb)(void* arg, int err); + void* done_arg; +}; + +struct media_delivery_ctx { + uint64_t self_node_id; + int initialized; + uint8_t is_supernode; // из adm_tags в peers_* + sqlite3* db; // = inst->topo_sqlite_db (кешируем для удобства) + struct ll_queue* served_nodes; // media_served_node — только суперузел + struct ll_queue* super_peers; // media_super_peer — только суперузел + struct ll_queue* downloads; // media_download — активные загрузки + struct UTUN_INSTANCE* inst; + void* bgp_cbk_handle; // handle для topo_group_remove_node_cbk +}; +``` + +**Init:** +1. Создать таблицы `block_availability` + `super_sync` (CREATE TABLE IF NOT EXISTS) — **всегда**, независимо от `is_supernode` +2. Прочитать `adm_tags` из `peers_*` → если `supernode=yes` → `is_supernode=1` +3. Инициализировать `downloads` +4. Подписаться на `topo_group_add_node_cbk` для всех групп (BGP-события) +5. Подписаться на `member_sync_add_props_cbk` для отслеживания изменений `adm_tags` +6. Подписаться на `etcp_add_conn_status_cbk` + +**Динамическое переключение суперузла** (через коллбэк `node_props_changed`, см. Этап 2.1): + +- **Включение** (`is_supernode` стал `1`): + - Найти других суперузлов (как при старте) + - Для каждого: `conn_mgr_connect_node` → `SUPER_HELLO` → репликация (см. 6.1, `last_recv_id` из `super_sync`) +- **Выключение** (`is_supernode` стал `0`): + - Очистить чужие записи: `DELETE FROM block_availability WHERE node_id != self` + - Остановить репликацию со всеми суперузлами, очистить `super_peers` + - `super_sync` не трогаем — пригодится если снова станем суперузлом + - `served_nodes` очищается автоматически по дисконнекту req-узлов + +**Файлы:** `src/media_delivery/media_delivery.h`, `src/media_delivery/media_delivery.c` + +--- + +## Этап 5: Обработчик `media_delivery_etcp_recv_cb` + +Диспатч по подкомандам: + +| Подкоманда | Действие суперузла | Действие req-узла | +|---|---|---| +| `SERVE_REG` | Добавить в served_nodes, ответить SERVE_ACK | — | +| `SERVE_ACK` | — | Отметить суперузел как обслуживающий | +| `SERVE_LEAVE` | Удалить из served_nodes | — | +| `QUERY` | SQL: `SELECT node_id, timestamp FROM block_availability WHERE block_uuid=?` → QUERY_RESP | — | +| `QUERY_RESP` | — | Получить список узлов, начать подключение | +| `HAVE_BLOCK` | SQL: `INSERT OR REPLACE` → ответ `HAVE_BLOCK_ACK` → репликация + уведомление автора | — | +| `HAVE_BLOCK_ACK` | — | Подтверждение получено, сброс таймера. Если dup/already — повторно не шлём | +| `CANCEL` | — | Остановить активный стрим для media_id/block_id, закрыть fd | +| `SUPER_REPL` | SQL: `INSERT OR REPLACE` для каждой записи, обновить `last_recv_id` в `super_sync`, ответить `SUPER_ACK(max_id)` | — | +| `SUPER_ACK` | Сдвинуть `peer_last_recv_id = ack_seq`, сбросить таймаут, отправить следующую пачку | — | +| `SUPER_HELLO` | Записать `hello.last_recv_id` в `peer->peer_last_recv_id` (это ID в нашей базе). Отправить ответный `SUPER_HELLO(мой_last_recv_id для этого пира)`. Начать репликацию с `peer_last_recv_id + 1` | — | +| `BLOCK_REQ` | — | Открыть файл на чтение, начать стрим чанков BLOCK_CHUNK | +| `BLOCK_CHUNK` | — | Записать чанк во временный файл блока | +| `BLOCK_DONE` | — | Закрыть файл, проверить подпись, отметить полученным | + +--- + +## Этап 6: Логика суперузла — обслуживание и репликация + +### 6.1 Старт суперузла + +При `is_supernode=1` (и при старте, и при динамическом включении): +- Таблицы `block_availability` + `super_sync` уже созданы в `media_delivery_init` +- Для каждой группы где мы член: найти других суперузлов + - `SELECT node_id FROM peers_ WHERE node_type=4 AND node_id != self` +- Для каждого другого суперузла: `conn_mgr_connect_node` +- При успешном подключении обменяться `SUPER_HELLO` (обе стороны): + - A → B: `SUPER_HELLO(A.last_recv_id_for_B)` + - B → A: `SUPER_HELLO(B.last_recv_id_for_A)` + - `last_recv_id` для пира берётся из `super_sync`: `SELECT last_recv_id FROM super_sync WHERE peer_node_id=?` + - После обмена каждая сторона начинает отправку с `peer_last_recv_id + 1` + +При динамическом **включении** (был обычный узел, стал суперузел): +- `last_recv_id` для каждого пира — из `super_sync` (мог остаться от прошлого включения). + Если записи нет — 0. Не сбрасываем принудительно. +- Отправляем `SUPER_HELLO(last_recv_id)`, начинаем репликацию с `last_recv_id + 1` + +При динамическом **выключении** (был суперузел, стал обычный): +- `DELETE FROM block_availability WHERE node_id != self` — чистим чужие записи +- Останавливаем репликацию со всеми, очищаем `super_peers` +- `super_sync` сохраняем (last_recv_id остаётся) +- `served_nodes` очищается автоматически по дисконнекту req-узлов + +### 6.2 Обнаружение через BGP-коллбэк и `node_props_changed` + +Два источника — разные назначения, без дублирования: + +- **BGP-коллбэк** (`TOPO_NODE_EVENT_NEW/UPDATE`): узел появился в BGP-таблице группы (можем до него маршрутизировать). + Проверяем `node_type=4` (суперузел) → `conn_mgr_connect_node` → `SUPER_HELLO` → репликация. + BGP не даёт `adm_tags` напрямую, только `node_type`. + +- **`node_props_changed`**: изменились `adm_tags` любого узла (включая себя). Это основной источник для отслеживания переключения роли. + Проверяем изменение `supernode=yes` в `adm_tags` (strstr). + - `node_id == self` → dynamic switch `is_supernode` + - `node_id != self`: стал суперузлом → `conn_mgr_connect_node`. Перестал → disconnect. + +- **Устранение дублей**: если узел уже в `super_peers`, второй вызов `conn_mgr_connect_node` вернёт `CONN_MGR_ERR_ALREADY_CONNECTED` — игнорируем. + +При `TOPO_NODE_EVENT_REMOVE`: +- Удалить из `super_peers`, очистить контекст +- `DELETE FROM block_availability WHERE node_id=?` (записи ушедшего узла) + +### 6.3 Репликация (per-peer media_super_peer) + +Отправитель stateless, приёмник хранит `last_recv_id` в `super_sync`. + +**Рукопожатие при подключении:** +``` +A → B: SUPER_HELLO(last_recv_id_A_from_B) // «я принял твои (B) записи вплоть до ID X в ТВОЕЙ базе» +B → A: SUPER_HELLO(last_recv_id_B_from_A) // «я принял твои (A) записи вплоть до ID Y в ТВОЕЙ базе» +Обе стороны: peer_last_recv_id = значение из HELLO (это ID в СВОЕЙ базе, с которого начинать отправку) +A начинает отправку с Y+1 (записи своей базы с id > Y), B начинает с X+1 +``` + +**Отправка репликации (sender):** +``` +media_super_repl_send(peer): + rows = SQL("SELECT id, block_uuid, group_id, node_id, chunk, timestamp + FROM block_availability WHERE id > ? ORDER BY id LIMIT ?", + peer->peer_last_recv_id, 4) + if rows пуст: return + отправить MEDIA_SUPER_REPL с seq = max(id) из rows + peer->inflight_count++ + если inflight_count < 4 и есть ещё: отправить ещё + запустить таймер ретрансмита (peer->timeout_tb) +``` + +**Приём репликации (receiver):** +``` +на MEDIA_SUPER_REPL: + для каждой записи: INSERT OR REPLACE в block_availability + обновить last_recv_id = max(id из записей пира) в super_sync + // last_recv_id — это ID в ЧУЖОЙ базе (пира) + ответить SUPER_ACK(max_id_пира) + // подтверждаем ID пира (не локальный) +``` + +**Обработка SUPER_ACK (sender):** +``` +на MEDIA_SUPER_ACK(ack_seq): + // ack_seq — ID нашей записи в НАШЕЙ базе, подтверждённый пиром + peer->peer_last_recv_id = ack_seq + peer->inflight_count-- + peer->timeout_tb = MEDIA_REPL_TIMEOUT_TB // сброс таймаута + если есть ещё записи (id > ack_seq): отправить следующую пачку + иначе: остановить таймер + +таймер ретрансмита (сработал): + повторно отправить все неподтверждённые записи + peer->timeout_tb = min(timeout_tb * 2, MEDIA_REPL_TIMEOUT_MAX_TB) + +таймаут SUPER_HELLO (2 сек без ответа): + закрыть соединение, запустить таймер переподключения MEDIA_RECONNECT_COOLDOWN_TB (1 час) + +при дисконнекте: + inflight_count = 0 + запустить таймер переподключения MEDIA_RECONNECT_COOLDOWN_TB (1 час) +``` + +### 6.4 При новом HAVE_BLOCK + +Если есть подключённые суперузлы с `inflight_count < 4`: +- Сразу отправить новую запись через `SUPER_REPL` + +### 6.5 Обслуживание req-узлов + +- `served_nodes` — ll_queue с индексом по node_id +- При `etcp_conn_status_cbk` с DOWN/DELETE: удалить из served_nodes +- При `SERVE_REG`: добавить, ответить `SERVE_ACK` +- При `SERVE_LEAVE`: удалить + +**Когда req-узел отправляет `SERVE_REG`:** сразу после успешного `conn_mgr_connect_node` к суперузлу (перед отправкой `MEDIA_QUERY`). Но `MEDIA_HAVE_BLOCK` принимается суперузлом и без предварительного `SERVE_REG` — регистрация опциональна, нужна для мониторинга и упрощённого реконнекта. + +**Subcmd не по роли:** если не-суперузел получает `QUERY`, `HAVE_BLOCK`, `SUPER_REPL` или `SUPER_HELLO` — ответить `MEDIA_ERR_NOT_SUPERNODE` (status=0xFE в следующем байте после subcmd). Если суперузел получает `QUERY_RESP` или `HAVE_BLOCK_ACK` — игнорировать. + +--- + +## Этап 7: `media_download.h/c` — скачивание блоков + +### 7.1 Запуск загрузки + +```c +int media_download_start(struct UTUN_INSTANCE* inst, uint64_t group_id, + const struct media_index_result* result, + const char* dest_path, const char* media_base, + void (*done_cb)(void* arg, int err), void* done_arg); +``` + +1. Создать `media_download`, заполнить из `result` +2. Найти суперузлов: `SELECT node_id FROM peers_ WHERE node_type=4` для всех групп +3. Для каждого суперузла получить RTT (из topo_node) +4. Выбрать суперузел с минимальным RTT +5. Отправить `MEDIA_QUERY` через `etcp_router` +6. Поставить таймаут `MEDIA_QUERY_TIMEOUT_TB` (2 сек) +7. При таймауте → следующий суперузел + +### 7.2 Обработка QUERY_RESP + +1. Получить массив `{node_id, rtt, block_id, chunk}` из ответа +2. Сгруппировать по node_id +3. Сортировать по RTT, взять до 10 лучших +4. Для каждого: `conn_mgr_connect_node(mgr, node_id, ...)` +5. При успешном подключении → `media_download_send_block_requests(dl, peer_index)` + +### 7.3 Стриминг блока чанками по 1KB с backpressure + +Передача блока (10-25MB) идёт потоком чанков по ~1KB каждый. +Congestion control через встроенный backpressure `etcp_router_on_send_ready`: +отправитель шлёт чанки пока `inflight` не заполнится, затем ждёт коллбэк «можно отправлять» +и продолжает стрим. + +Состояние одного стрима (блок-холдер → req-узел): +```c +struct media_block_stream { + uint8_t media_id[16]; + uint8_t block_id[16]; + uint32_t chunk; + uint64_t offset; // текущее смещение в блоке + uint32_t bytes_sent; // всего отправлено байт + uint32_t total_size; // ожидаемый размер блока + uint8_t done; // 1 = стрим завершён + void* waiter; // handle etcp_router_on_send_ready + FILE* file; // fd исходного файла + uint64_t dst_node_id; + uint64_t group_id; +}; +``` + +**Отправка (холдер):** +1. При `MEDIA_BLOCK_REQ` — создать `media_block_stream`, открыть файл на чтение +2. Вызвать `etcp_router_on_send_ready` для регистрации backpressure waiter +3. В коллбэке `on_send_ready`: читать чанк из файла (до 1024 байт), отправить `MEDIA_BLOCK_CHUNK` +4. Если `etcp_route_send` вернул быстрый успех — отправить следующий чанк сразу +5. Если очередь заполнена — waiter сработает когда можно отправлять дальше +6. Когда весь блок отправлен — отправить `MEDIA_BLOCK_DONE` с `block_sig` и `total_size` +7. Закрыть файл, освободить стрим + +**Приём (req-узел):** +1. При `MEDIA_BLOCK_REQ` — создать writer для блока: открыть временный файл по `dest_path + offset` +2. При `MEDIA_BLOCK_CHUNK` — записать данные файл, сдвинуть offset +3. При `MEDIA_BLOCK_DONE` — проверить что `total_size` совпадает с записанным, закрыть файл +4. Верифицировать `block_sig` (Ed25519 от подписи в `block_sigs`) +5. Если подпись верна: отметить блок полученным, `blocks_received++` +6. Если подпись не совпала: удалить файл, запросить блок с другого узла +7. Проверить `blocks_received == num_blocks` → сборка файла (см. 7.5) + +### 7.4 Распределение блоков по узлам + +- Для каждого подключённого узла: распределить неполученные блоки равномерно +- Отправить `MEDIA_BLOCK_REQ` для блоков назначенных этому узлу +- Один узел может параллельно стримить несколько блоков (отдельный стрим на каждый) +- Если узел дисконнектится — его неначатые блоки перераспределяются на других + +### 7.5 Сборка файла и обновление БД + +При получении каждого блока: +1. Проверить: `blocks_received == num_blocks`? +2. Если да — склеить блоки в порядке chunk в целевой файл: + ``` + for (n = 0; n < num_blocks; n++) { + прочитать temp_файл блока n + записать в dest_path на offset = n * block_size + удалить temp_файл + } + ``` +3. Проверить `content_hash` собранного файла (SHA256) +4. Если хеш совпал: + - Удалить записи отдельных блоков из `media_files` + - Вставить одну запись с `chunk=0`, `chunk_size=file_size` (цельный файл) + - Для каждого полученного и проверенного блока — отправить `MEDIA_HAVE_BLOCK` суперузлу + (с подтверждением и перезапросом — см. 7.6) + - Вызвать `done_cb(0)` +5. Если хеш не совпал — удалить собранный файл, перезапросить проблемные блоки + +### 7.6 HAVE_BLOCK с подтверждением и перезапросом + +Отправка `MEDIA_HAVE_BLOCK` — с обязательным подтверждением от суперузла: + +``` +req-узел → суперузел: MEDIA_HAVE_BLOCK + запустить таймер HAVE_BLOCK_TIMEOUT_TB (2 сек) + +суперузел → req-узел: MEDIA_HAVE_BLOCK_ACK(status=0 OK, 1 dup) + если OK: INSERT в block_availability → репликация → уведомление автора + если dup: запись уже есть, ничего не делаем + +req-узел получил ACK: + сброс таймера, блок считается доставленным суперузлу + +req-узел таймаут / ошибка: + выбрать другой суперузел (следующий по RTT из списка) + отправить MEDIA_HAVE_BLOCK повторно + если все суперузлы исчерпаны → done_cb(ERR) +``` + +Константа: `HAVE_BLOCK_TIMEOUT_TB = 20000` (2 сек). + +### 7.7 Отмена передачи (MEDIA_SUBCMD_CANCEL) + +Из GUI или при ошибке req-узел может отменить активную передачу блока: + +``` +req-узел → блок-холдер: MEDIA_CANCEL(media_id, block_id, chunk) + +блок-холдер при получении CANCEL: + остановить стрим для указанного блока + закрыть fd + отменить waiter (etcp_router_cancel_send_ready) + освободить media_block_stream +``` + +Также при дисконнекте блок-холдера или req-узла — все активные стримы отменяются автоматически (через `etcp_conn_status_cbk`). + +### 7.8 Обработка ошибок + +- Таймаут BLOCK_REQ (нет начала стрима за 10 сек) → перезапросить с другого узла +- Стрим stalled (нет чанков > 30 сек) → отправить CANCEL, запросить блок с другого узла +- BLOCK_DONE с неверной подписью → удалить файл, запросить с другого узла +- Все узлы для блока исчерпаны → `done_cb(ERR)` +- Дисконнект узла → перераспределить его блоки на других, отменить активные стримы с него +- Общий таймаут загрузки → `done_cb(ERR_TIMEOUT)` +- HAVE_BLOCK без ACK (2 сек) → переключиться на следующий суперузел + +--- + +## Этап 8: Интеграция + +**Файлы для изменения:** +- `src/Makefile.am` — добавить `media_delivery/media_download.c` +- `src/media_delivery/media_delivery.c` — полностью переписать обработчик +- `src/media_delivery/media_delivery.h` — расширить структуру +- `src/routing_layer/topo_group.h` + `topo_group.c` — добавить BGP-коллбэк +- `src/chat/member_sync.h` + `member_sync.c` — добавить `node_props_changed` коллбэк + +**Уже готово:** +- `media_delivery_init` вызывается в `utun_instance.c:650` +- `media_delivery_destroy` вызывается в `utun_instance.c:514` +- `ETCP_RT_ID_MEDIA_DELIVERY` забинжен в `media_delivery_init` + +--- + +## Решённые вопросы + +1. **Формат передачи блоков:** стриминг чанками по ~1KB с congestion control через `etcp_router_on_send_ready` (вариант А). Отправитель шлёт чанки пока inflight позволяет, затем ждёт backpressure-коллбэк. Завершается `MEDIA_BLOCK_DONE` с `block_sig` (Ed25519) и `total_size`. + +2. **Валидация при приёме:** каждый блок подписан Ed25519. При `BLOCK_DONE` проверяем `block_sig` против подписи из `media_index_result.block_sigs`. Если не совпала — удаляем файл, запрашиваем блок с другого узла. + +3. **Уведомление автора контента:** один раз, при успешной вставке `HAVE_BLOCK` в `block_availability` (после ответа `HAVE_BLOCK_ACK`). Без переповторов. Если что-то потерялось — req-узел перезапросит у другого суперузла. HAVE_BLOCK всегда с подтверждением: таймаут 2 сек → переключение на следующий суперузел. + +4. **Файловая система:** путь назначения (`dest_path`) и `media_base` передаются в `media_download_start` (точно как в `media_index_register_async`). Временные файлы блоков сохраняются рядом с dest_path с суффиксом `.chunk_N`. + +5. **Сборка файла:** в момент получения очередного блока проверяем `blocks_validated == num_blocks`. Если все на месте — склеиваем чанки по порядку в `dest_path`, проверяем `content_hash` (SHA256). При успехе удаляем записи отдельных блоков из `media_files`, вставляем одну запись с `chunk_size=file_size` (цельный файл).