Browse Source

media_delivery: BLOCK_PROCESSING + relay streaming + per-file download limit

- block_availability: status column (0=processing, 1=completed)
- md_ba_complete_block(): DELETE old + check exists + INSERT new id
- BLOCK_PROCESSING subcmd (0x10) — INSERT with new id, no REPLACE/IGNORE
- BLOCK_RELAY_FULL subcmd (0x11) — redirect when source/relay at capacity
- relay_block_ctx: downstream forwarding from chunk file
- md_file_load: per-file source download limit (max_downloads_per_file)
- RELAY_FULL from source node when at capacity
- relay_full handler: retry existing peers, not just new ones
- sync.md: full documentation
topo_upd
evgeny 2 months ago
parent
commit
de9e467b39
  1. 597
      src/media_delivery/_plan.md
  2. 364
      src/media_delivery/media_delivery.c
  3. 71
      src/media_delivery/media_delivery.h
  4. 13
      src/media_delivery/media_delivery_proto.h
  5. 183
      src/media_delivery/media_download.c
  6. 4
      src/media_delivery/media_download.h
  7. 251
      src/media_delivery/sync.md
  8. 108
      tests/test_media_delivery_sql.c

597
src/media_delivery/_plan.md

@ -1,597 +0,0 @@
# Архитектура 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]`
<!-- src_node_id не передаётся — суперузел ищет по block_uuid независимо от автора.
src_node_id известен из сообщения чата (media_index_result), req-узел хранит его в media_download. -->
- `MEDIA_QUERY_RESP`: `subcmd:1, num_entries:2, entries[num × {node_id:8, rtt:2, block_id:16, chunk:4}]`
<!-- rtt — cumulative_rtt из topo_node суперузла до узла-обладателя (0.1ms).
req-узел может использовать для сортировки или перемерить сам через route_ping. -->
- `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]`
<!-- data_len — реальный размер (1..1024 байт). offset — смещение в блоке. -->
- `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`
<!-- last_recv_id — ID в МОЕЙ базе (отправителя HELLO) до которого пир подтвердил приём.
Получатель HELLO использует это значение чтобы знать с какого id начинать отправку. -->
**Файл:** `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_<ch> 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_<ch> 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` (цельный файл).

364
src/media_delivery/media_delivery.c

@ -50,6 +50,8 @@ static const char* md_subcmd_name(uint8_t sc) {
case MEDIA_SUBCMD_BLOCK_DONE: return "BLOCK_DONE";
case MEDIA_SUBCMD_CANCEL: return "CANCEL";
case MEDIA_SUBCMD_SUPER_HELLO: return "SUPER_HELLO";
case MEDIA_SUBCMD_BLOCK_PROCESSING: return "BLOCK_PROCESSING";
case MEDIA_SUBCMD_BLOCK_RELAY_FULL: return "BLOCK_RELAY_FULL";
default: return "???";
}
}
@ -57,10 +59,10 @@ static const char* md_subcmd_name(uint8_t sc) {
/* ── SQL helper: insert/update block_availability ── */
static int md_ba_insert(sqlite3* db, const uint8_t* block_uuid, uint64_t group_id,
uint64_t node_id, uint32_t chunk, int64_t ts) {
uint64_t node_id, uint32_t chunk, int64_t ts, int status) {
const char* sql =
"INSERT OR REPLACE INTO block_availability(block_uuid,group_id,node_id,chunk,timestamp)"
" VALUES(?,?,?,?,?)";
"INSERT INTO block_availability(block_uuid,group_id,node_id,chunk,timestamp,status)"
" VALUES(?,?,?,?,?,?)";
sqlite3_stmt* st = NULL;
if (sqlite3_prepare_v2(db, sql, -1, &st, NULL) != SQLITE_OK) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: ba_insert prepare failed: %s", MD_ID, sqlite3_errmsg(db));
@ -71,6 +73,7 @@ static int md_ba_insert(sqlite3* db, const uint8_t* block_uuid, uint64_t group_i
sqlite3_bind_int64(st, 3, (sqlite3_int64)node_id);
sqlite3_bind_int(st, 4, (int)chunk);
sqlite3_bind_int64(st, 5, (sqlite3_int64)ts);
sqlite3_bind_int(st, 6, status);
int rc = sqlite3_step(st);
sqlite3_finalize(st);
if (rc != SQLITE_DONE) {
@ -80,6 +83,39 @@ static int md_ba_insert(sqlite3* db, const uint8_t* block_uuid, uint64_t group_i
return 0;
}
/* ── complete block: DELETE processing records + INSERT completed with new id ── */
static int md_ba_complete_block(sqlite3* db, const uint8_t* block_uuid,
uint64_t group_id, uint64_t node_id,
uint32_t chunk, int64_t ts) {
/* 1. DELETE processing records (status=0) for this (block_uuid, node_id) */
const char* del_sql = "DELETE FROM block_availability WHERE block_uuid=? AND node_id=? AND status=0";
sqlite3_stmt* st = NULL;
if (sqlite3_prepare_v2(db, del_sql, -1, &st, NULL) == SQLITE_OK) {
sqlite3_bind_blob(st, 1, block_uuid, 16, SQLITE_STATIC);
sqlite3_bind_int64(st, 2, (sqlite3_int64)node_id);
sqlite3_step(st);
sqlite3_finalize(st);
}
/* 2. check if completed already exists — avoid duplicate insert */
const char* chk_sql = "SELECT 1 FROM block_availability WHERE block_uuid=? AND node_id=? AND status=1";
if (sqlite3_prepare_v2(db, chk_sql, -1, &st, NULL) == SQLITE_OK) {
sqlite3_bind_blob(st, 1, block_uuid, 16, SQLITE_STATIC);
sqlite3_bind_int64(st, 2, (sqlite3_int64)node_id);
int exists = (sqlite3_step(st) == SQLITE_ROW);
sqlite3_finalize(st);
if (exists) {
DEBUG_DEBUG(DEBUG_CATEGORY_DEBUG, "%s: ba_complete_block: completed already exists for node=0x%016llx, skipping",
MD_ID, (unsigned long long)node_id);
return 0;
}
}
/* 3. INSERT completed record (status=1, new id via AUTOINCREMENT) */
return md_ba_insert(db, block_uuid, group_id, node_id, chunk, ts, 1);
}
static int md_ba_find_nodes(sqlite3* db, const uint8_t* block_uuid,
uint64_t* out_node_ids, int max) {
const char* sql =
@ -103,7 +139,7 @@ static int md_ba_find_nodes(sqlite3* db, const uint8_t* block_uuid,
static int md_ba_get_since(sqlite3* db, uint64_t since_id,
uint8_t* out_buf, int max_bytes) {
const char* sql =
"SELECT id,block_uuid,group_id,node_id,chunk,timestamp"
"SELECT id,block_uuid,group_id,node_id,chunk,timestamp,status"
" FROM block_availability WHERE id > ? ORDER BY id ASC LIMIT 4";
sqlite3_stmt* st = NULL;
if (sqlite3_prepare_v2(db, sql, -1, &st, NULL) != SQLITE_OK) {
@ -113,7 +149,7 @@ static int md_ba_get_since(sqlite3* db, uint64_t since_id,
sqlite3_bind_int64(st, 1, (sqlite3_int64)since_id);
int off = 0;
while (sqlite3_step(st) == SQLITE_ROW) {
if (off + 52 > max_bytes) break;
if (off + 56 > max_bytes) break;
int64_t id = sqlite3_column_int64(st, 0);
memcpy(out_buf + off, &id, 8); off += 8;
memcpy(out_buf + off, sqlite3_column_blob(st, 1), 16); off += 16;
@ -121,6 +157,7 @@ static int md_ba_get_since(sqlite3* db, uint64_t since_id,
int64_t nid = sqlite3_column_int64(st, 3); memcpy(out_buf + off, &nid, 8); off += 8;
int32_t ch = sqlite3_column_int(st, 4); memcpy(out_buf + off, &ch, 4); off += 4;
int64_t ts = sqlite3_column_int64(st, 5); memcpy(out_buf + off, &ts, 8); off += 8;
int32_t sts = sqlite3_column_int(st, 6); memcpy(out_buf + off, &sts, 4); off += 4;
}
sqlite3_finalize(st);
return off;
@ -191,6 +228,98 @@ static void md_super_peer_remove(struct media_delivery_ctx* md, uint64_t node_id
queue_entry_free(e);
}
/* ── relay block context helpers ── */
struct relay_block_ctx* md_relay_find(struct media_delivery_ctx* md, const uint8_t* block_id) {
if (!md->relay_blocks) return NULL;
struct ll_entry* e = queue_find_data_by_index(md->relay_blocks, block_id);
return e ? (struct relay_block_ctx*)e->data : NULL;
}
struct relay_block_ctx* md_relay_add(struct media_delivery_ctx* md,
const uint8_t* block_id,
const uint8_t* media_id,
uint32_t chunk) {
struct relay_block_ctx* rc = md_relay_find(md, block_id);
if (rc) return rc;
if (!md->relay_blocks) {
md->relay_blocks = queue_new(md->inst->ua, 64, offsetof(struct relay_block_ctx, block_id), 16, "md_relay");
if (!md->relay_blocks) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: queue_new(relay_blocks) failed", MD_ID); return NULL; }
}
struct ll_entry* qe = queue_entry_new(sizeof(struct relay_block_ctx));
if (!qe) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: queue_entry_new(relay_block) failed", MD_ID); return NULL; }
rc = (struct relay_block_ctx*)qe->data;
memset(rc, 0, sizeof(*rc));
memcpy(rc->block_id, block_id, 16);
memcpy(rc->media_id, media_id, 16);
rc->chunk = chunk;
rc->inst = md->inst;
memcpy(qe->data, block_id, 16);
queue_data_put_with_index(md->relay_blocks, qe);
return rc;
}
void md_relay_remove(struct media_delivery_ctx* md, const uint8_t* block_id) {
if (!md->relay_blocks) return;
struct ll_entry* e = queue_find_data_by_index(md->relay_blocks, block_id);
if (!e) return;
queue_remove_data(md->relay_blocks, e);
queue_entry_free(e);
}
/* ── file_load helpers (per-file source download tracking) ── */
struct md_file_load* md_file_load_find(struct media_delivery_ctx* md, const uint8_t* media_id) {
if (!md->file_loads) return NULL;
struct ll_entry* e = queue_find_data_by_index(md->file_loads, media_id);
return e ? (struct md_file_load*)e->data : NULL;
}
int md_file_load_inc(struct media_delivery_ctx* md, const uint8_t* media_id, uint64_t node_id) {
if (!md->file_loads) {
md->file_loads = queue_new(md->inst->ua, 64, offsetof(struct md_file_load, media_id), 16, "md_fld");
if (!md->file_loads) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: queue_new(file_loads) failed", MD_ID); return -1; }
}
struct md_file_load* fl = md_file_load_find(md, media_id);
if (!fl) {
struct ll_entry* qe = queue_entry_new(sizeof(struct md_file_load));
if (!qe) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: queue_entry_new(file_load) failed", MD_ID); return -1; }
fl = (struct md_file_load*)qe->data;
memset(fl, 0, sizeof(*fl));
memcpy(fl->media_id, media_id, 16);
memcpy(qe->data, media_id, 16);
queue_data_put_with_index(md->file_loads, qe);
}
if (fl->active_downloads < 10)
fl->downloader_ids[fl->active_downloads] = node_id;
fl->active_downloads++;
DEBUG_DEBUG(DEBUG_CATEGORY_DEBUG, "%s: file_load inc media=%02x%02x... node=0x%016llx count=%d",
MD_ID, media_id[0], media_id[1], (unsigned long long)node_id, fl->active_downloads);
return fl->active_downloads;
}
void md_file_load_dec(struct media_delivery_ctx* md, const uint8_t* media_id, uint64_t node_id) {
struct md_file_load* fl = md_file_load_find(md, media_id);
if (!fl) return;
if (fl->active_downloads > 0) fl->active_downloads--;
/* remove node_id from downloader_ids array */
for (int i = 0; i < fl->active_downloads && i < 10; i++) {
if (fl->downloader_ids[i] == node_id) {
for (int j = i; j < fl->active_downloads && j < 9; j++)
fl->downloader_ids[j] = fl->downloader_ids[j + 1];
fl->downloader_ids[fl->active_downloads < 10 ? fl->active_downloads : 9] = 0;
break;
}
}
DEBUG_DEBUG(DEBUG_CATEGORY_DEBUG, "%s: file_load dec media=%02x%02x... node=0x%016llx count=%d",
MD_ID, media_id[0], media_id[1], (unsigned long long)node_id, fl->active_downloads);
/* remove entry when count reaches 0 */
if (fl->active_downloads == 0 && md->file_loads) {
struct ll_entry* e = queue_find_data_by_index(md->file_loads, media_id);
if (e) { queue_remove_data(md->file_loads, e); queue_entry_free(e); }
}
}
/* ── replication timers ── */
static void md_repl_retrans_cb(void* arg) {
@ -229,7 +358,7 @@ static void md_super_repl_send(struct media_delivery_ctx* md, struct media_super
if (!md->db || !peer->connected || !peer->hello_done) return;
while (peer->inflight_count < MEDIA_MAX_REPL_INFLIGHT) {
uint8_t buf[512];
uint8_t buf[560];
int nb = md_ba_get_since(md->db, peer->peer_last_recv_id, buf, (int)sizeof(buf));
if (nb == 0) {
if (peer->timeout_timer) uasync_cancel_timeout(md->inst->ua, peer->timeout_timer);
@ -238,17 +367,17 @@ static void md_super_repl_send(struct media_delivery_ctx* md, struct media_super
}
/* build packet: subcmd + seq + num + raw rows */
uint8_t pkt[520];
uint8_t pkt[600];
int off = 0;
pkt[off++] = MEDIA_SUBCMD_SUPER_REPL;
int num_entries = nb / 52;
int num_entries = nb / 56;
int64_t max_id = 0;
const uint8_t* rp = buf;
for (int i = 0; i < num_entries; i++) {
int64_t id; memcpy(&id, rp, 8);
if (id > max_id) max_id = id;
rp += 52;
rp += 56;
}
uint32_t seq = (uint32_t)max_id; memcpy(pkt + off, &seq, 4); off += 4;
@ -374,7 +503,7 @@ static void md_handle_have_block(struct media_delivery_ctx* md, uint64_t from_no
if (!md->db) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: HAVE_BLOCK but no DB", MD_ID); return; }
struct media_pkt_have_block* hb = (struct media_pkt_have_block*)data;
int rc = md_ba_insert(md->db, hb->block_id, hb->group_id, from_node, hb->chunk, hb->timestamp);
int rc = md_ba_complete_block(md->db, hb->block_id, hb->group_id, from_node, hb->chunk, hb->timestamp);
uint8_t ack[sizeof(struct media_pkt_have_block_ack)];
struct media_pkt_have_block_ack* ha = (struct media_pkt_have_block_ack*)ack;
@ -401,6 +530,34 @@ static void md_handle_have_block(struct media_delivery_ctx* md, uint64_t from_no
}
}
/* ── handle BLOCK_PROCESSING: node started downloading a block, insert as processing (status=0) ── */
static void md_handle_block_processing(struct media_delivery_ctx* md, uint64_t from_node,
const uint8_t* data, size_t len) {
if (len < MEDIA_HAVE_BLOCK_SIZE) return;
if (!md->db) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: BLOCK_PROCESSING but no DB", MD_ID); return; }
struct media_pkt_have_block* hb = (struct media_pkt_have_block*)data;
/* INSERT processing (status=0) with new id */
int rc = md_ba_insert(md->db, hb->block_id, hb->group_id, from_node, hb->chunk, hb->timestamp, 0);
if (rc != 0) {
DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "%s: BLOCK_PROCESSING insert failed for node=0x%016llx (duplicate/constraint?)",
MD_ID, (unsigned long long)from_node);
}
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: BLOCK_PROCESSING from 0x%016llx block=%02x%02x...",
MD_ID, (unsigned long long)from_node, hb->block_id[0], hb->block_id[1]);
/* replicate to connected super-peers */
struct ll_entry* se = md->super_peers ? md->super_peers->head : NULL;
while (se) {
struct media_super_peer* sp = (struct media_super_peer*)se->data;
if (sp->connected && sp->hello_done && sp->inflight_count < MEDIA_MAX_REPL_INFLIGHT)
md_super_repl_send(md, sp);
se = se->next;
}
}
static void md_handle_super_repl(struct media_delivery_ctx* md, uint64_t from_node,
const uint8_t* data, size_t len) {
if (len < MEDIA_SUPER_REPL_HDR_SIZE) return;
@ -418,7 +575,23 @@ static void md_handle_super_repl(struct media_delivery_ctx* md, uint64_t from_no
int64_t nid; memcpy(&nid, ep, 8); ep += 8;
int32_t ch; memcpy(&ch, ep, 4); ep += 4;
int64_t ts; memcpy(&ts, ep, 8); ep += 8;
md_ba_insert(md->db, block_uuid, (uint64_t)gid, (uint64_t)nid, (uint32_t)ch, ts);
int32_t st; memcpy(&st, ep, 4); ep += 4;
if (st == 1) {
md_ba_complete_block(md->db, block_uuid, (uint64_t)gid, (uint64_t)nid, (uint32_t)ch, ts);
} else {
/* processing: INSERT OR IGNORE (дубликаты при репликации — норма) */
const char* ign = "INSERT OR IGNORE INTO block_availability(block_uuid,group_id,node_id,chunk,timestamp,status) VALUES(?,?,?,?,?,0)";
sqlite3_stmt* igs = NULL;
if (sqlite3_prepare_v2(md->db, ign, -1, &igs, NULL) == SQLITE_OK) {
sqlite3_bind_blob(igs, 1, block_uuid, 16, SQLITE_STATIC);
sqlite3_bind_int64(igs, 2, (sqlite3_int64)gid);
sqlite3_bind_int64(igs, 3, (sqlite3_int64)nid);
sqlite3_bind_int(igs, 4, (int)ch);
sqlite3_bind_int64(igs, 5, (sqlite3_int64)ts);
sqlite3_step(igs);
sqlite3_finalize(igs);
}
}
}
md_ss_set(md->db, from_node, (uint64_t)max_id);
@ -578,6 +751,7 @@ static void stream_send_chunk_cb(struct ll_queue* q, void* arg) {
/* block_sig is zero — receiver will verify vs block_sigs from media_index_result */
md_send(md->inst, sc->group_id, sc->dst_node_id, done_pkt, MEDIA_BLOCK_DONE_SIZE);
md_file_load_dec(md, sc->media_id, sc->dst_node_id);
media_delivery_stream_done(md->inst);
fclose(sc->file);
u_free(sc);
@ -591,6 +765,74 @@ static void stream_send_chunk_cb(struct ll_queue* q, void* arg) {
}
}
/* ── relay chunk forwarding ── */
struct relay_fwd_ctx {
struct media_delivery_ctx* md;
uint64_t dst_node_id;
uint64_t group_id;
struct relay_block_ctx* rc;
int ds_idx; // индекс в rc->downstream[]
struct queue_waiter_handle waiter;
};
static void md_relay_fwd_cb(struct ll_queue* q, void* arg) {
(void)q;
struct relay_fwd_ctx* fc = (struct relay_fwd_ctx*)arg;
struct relay_block_ctx* rc = fc->rc;
struct relay_downstream* ds = &rc->downstream[fc->ds_idx];
if (ds->sent_offset >= rc->file_offset) { u_free(fc); return; }
FILE* f = fopen(rc->chunk_file, "rb");
if (!f) { u_free(fc); return; }
fseeko(f, (off_t)ds->sent_offset, SEEK_SET);
size_t to_read = rc->file_offset - ds->sent_offset;
if (to_read > 1024) to_read = 1024;
uint8_t buf[1024];
size_t rd = fread(buf, 1, to_read, f);
fclose(f);
if (rd == 0) { u_free(fc); return; }
uint8_t pkt[MEDIA_BLOCK_CHUNK_HDR_SIZE + 1024];
struct media_pkt_block_chunk* ch = (struct media_pkt_block_chunk*)pkt;
memset(ch, 0, MEDIA_BLOCK_CHUNK_HDR_SIZE);
ch->subcmd = MEDIA_SUBCMD_BLOCK_CHUNK;
memcpy(ch->media_id, rc->media_id, 16);
memcpy(ch->block_id, rc->block_id, 16);
ch->chunk = rc->chunk;
ch->offset = (uint32_t)ds->sent_offset;
ch->data_len = (uint16_t)rd;
memcpy(pkt + MEDIA_BLOCK_CHUNK_HDR_SIZE, buf, rd);
size_t pkt_len = MEDIA_BLOCK_CHUNK_HDR_SIZE + rd;
if (md_send(fc->md->inst, fc->group_id, fc->dst_node_id, pkt, pkt_len) == 0) {
ds->sent_offset += rd;
}
if (ds->sent_offset < rc->file_offset) {
etcp_router_on_send_ready(fc->md->inst, fc->group_id,
fc->dst_node_id, ETCP_RT_ID_MEDIA_DELIVERY,
&fc->waiter, md_relay_fwd_cb, fc);
} else {
u_free(fc);
}
}
static void md_relay_catchup(struct media_delivery_ctx* md, struct relay_block_ctx* rc,
uint64_t dst_node_id, uint64_t group_id, int ds_idx) {
struct relay_downstream* ds = &rc->downstream[ds_idx];
if (rc->file_offset == 0 || ds->sent_offset >= rc->file_offset) return;
struct relay_fwd_ctx* fc = u_calloc(1, sizeof(*fc));
if (!fc) return;
fc->md = md; fc->dst_node_id = dst_node_id; fc->group_id = group_id;
fc->rc = rc; fc->ds_idx = ds_idx;
memset(&fc->waiter, 0, sizeof(fc->waiter));
md_relay_fwd_cb(NULL, fc);
}
static void md_handle_block_req(struct media_delivery_ctx* md, uint64_t from_node,
const uint8_t* data, size_t len) {
if (len < MEDIA_BLOCK_REQ_SIZE) return;
@ -637,7 +879,42 @@ static void md_handle_block_req(struct media_delivery_ctx* md, uint64_t from_nod
if (sqlite3_prepare_v2(db, sql, -1, &st, NULL) != SQLITE_OK) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: BLOCK_REQ prepare failed: %s", MD_ID, sqlite3_errmsg(db)); md->active_streams--; return; }
sqlite3_bind_blob(st, 1, req->block_id, 16, SQLITE_STATIC);
sqlite3_bind_int64(st, 2, (sqlite3_int64)md->inst->node_id);
if (sqlite3_step(st) != SQLITE_ROW) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: BLOCK_REQ block_id=%02x%02x... not found in media_files", MD_ID, req->block_id[0], req->block_id[1]); sqlite3_finalize(st); md->active_streams--; return; }
if (sqlite3_step(st) != SQLITE_ROW) {
sqlite3_finalize(st);
/* fallback: check relay_blocks — maybe we are downloading this block and can relay */
struct relay_block_ctx* rc = md_relay_find(md, req->block_id);
if (rc && rc->downstream_count < MD_MAX_RELAY_DOWNSTREAM) {
int dsi = rc->downstream_count++;
rc->downstream[dsi].node_id = from_node;
rc->downstream[dsi].group_id = group_id;
rc->downstream[dsi].sent_offset = 0;
memset(&rc->downstream[dsi].waiter, 0, sizeof(rc->downstream[dsi].waiter));
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: relay downstream[%d] added node=0x%016llx file_offset=%llu",
MD_ID, dsi, (unsigned long long)from_node, (unsigned long long)rc->file_offset);
md_relay_catchup(md, rc, from_node, group_id, dsi);
md->active_streams--; return;
}
if (rc) {
/* relay full — send RELAY_FULL with downstream list */
uint8_t rfp[256]; int roff = 0;
rfp[roff++] = MEDIA_SUBCMD_BLOCK_RELAY_FULL;
memcpy(rfp + roff, req->media_id, 16); roff += 16;
memcpy(rfp + roff, req->block_id, 16); roff += 16;
uint32_t ck = req->chunk; memcpy(rfp + roff, &ck, 4); roff += 4;
uint16_t nrn = (uint16_t)rc->downstream_count;
memcpy(rfp + roff, &nrn, 2); roff += 2;
for (int i = 0; i < rc->downstream_count; i++) {
memcpy(rfp + roff, &rc->downstream[i].node_id, 8); roff += 8;
}
md_send(md->inst, group_id, from_node, rfp, (size_t)roff);
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: relay full for block=%02x%02x..., sent %d downstream nodes to 0x%016llx",
MD_ID, req->block_id[0], req->block_id[1], rc->downstream_count, (unsigned long long)from_node);
} else {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: BLOCK_REQ block_id=%02x%02x... not found in media_files",
MD_ID, req->block_id[0], req->block_id[1]);
}
md->active_streams--; return;
}
const char* loc = (const char*)sqlite3_column_text(st, 0);
int64_t chunk_size = sqlite3_column_int64(st, 1);
@ -674,8 +951,34 @@ static void md_handle_block_req(struct media_delivery_ctx* md, uint64_t from_nod
md->active_streams--; /* undo pre-increment */
return;
}
/* per-file download limit: redirect to current downloaders if at capacity */
if (md->max_downloads_per_file > 0) {
struct md_file_load* fl = md_file_load_find(md, req->media_id);
if (fl && fl->active_downloads >= (int)md->max_downloads_per_file) {
/* source is at capacity — send RELAY_FULL with current downloader list */
uint8_t rfp[256]; int roff = 0;
rfp[roff++] = MEDIA_SUBCMD_BLOCK_RELAY_FULL;
memcpy(rfp + roff, req->media_id, 16); roff += 16;
memcpy(rfp + roff, req->block_id, 16); roff += 16;
uint32_t ck = req->chunk; memcpy(rfp + roff, &ck, 4); roff += 4;
uint16_t nrn = (uint16_t)fl->active_downloads;
memcpy(rfp + roff, &nrn, 2); roff += 2;
for (int i = 0; i < fl->active_downloads && i < 10; i++) {
memcpy(rfp + roff, &fl->downloader_ids[i], 8); roff += 8;
}
md_send(md->inst, group_id, from_node, rfp, (size_t)roff);
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: source at capacity (%d/%d) for media=%02x%02x..., redirecting to %d nodes",
MD_ID, fl->active_downloads, md->max_downloads_per_file,
req->media_id[0], req->media_id[1], fl->active_downloads);
fclose(f);
md->active_streams--; return;
}
}
md_file_load_inc(md, req->media_id, from_node);
struct stream_ctx* sc = u_calloc(1, sizeof(*sc));
if (!sc) { fclose(f); md->active_streams--; return; }
if (!sc) { fclose(f); md->active_streams--; md_file_load_dec(md, req->media_id, from_node); return; }
sc->md = md;
sc->group_id = group_id;
sc->dst_node_id = from_node;
@ -697,6 +1000,28 @@ static void md_handle_block_req(struct media_delivery_ctx* md, uint64_t from_nod
stream_send_chunk_cb(NULL, sc);
}
/* ── handle RELAY_FULL: relay node is overloaded, retry one of its downstream nodes ── */
static void md_handle_block_relay_full(struct media_delivery_ctx* md, uint64_t from_node,
const uint8_t* data, size_t len) {
if (len < MEDIA_RELAY_FULL_HDR_SIZE) return;
struct media_pkt_block_relay_full* rf = (struct media_pkt_block_relay_full*)data;
uint16_t nrn = rf->num_relay_nodes;
if (nrn > MD_MAX_RELAY_DOWNSTREAM) nrn = MD_MAX_RELAY_DOWNSTREAM;
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: RELAY_FULL from 0x%016llx block=%02x%02x... nodes=%u",
MD_ID, (unsigned long long)from_node, rf->block_id[0], rf->block_id[1], nrn);
/* forward to media_download layer to retry from one of the listed nodes */
const uint8_t* rp = data + MEDIA_RELAY_FULL_HDR_SIZE;
uint64_t node_ids[MD_MAX_RELAY_DOWNSTREAM];
for (int i = 0; i < (int)nrn && i < MD_MAX_RELAY_DOWNSTREAM; i++) {
memcpy(&node_ids[i], rp, 8); rp += 8;
}
media_download_handle_relay_full(md->inst, rf->media_id, rf->block_id,
rf->chunk, node_ids, (int)nrn);
}
/* ── decrement stream counter (called when stream finished/cancelled internally) ── */
void media_delivery_stream_done(struct UTUN_INSTANCE* inst) {
@ -894,6 +1219,10 @@ static void md_etcp_recv(struct ETCP_CONN* conn, struct ll_entry* entry) {
case MEDIA_SUBCMD_HAVE_BLOCK:
md_handle_have_block(md, from_node, data, len);
break;
case MEDIA_SUBCMD_BLOCK_PROCESSING:
if (md->is_supernode) md_handle_block_processing(md, from_node, data, len);
else DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "%s: BLOCK_PROCESSING ignored — not supernode", MD_ID);
break;
case MEDIA_SUBCMD_SUPER_REPL:
if (md->is_supernode) md_handle_super_repl(md, from_node, data, len);
else DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "%s: SUPER_REPL ignored — not supernode", MD_ID);
@ -920,6 +1249,9 @@ static void md_etcp_recv(struct ETCP_CONN* conn, struct ll_entry* entry) {
case MEDIA_SUBCMD_BLOCK_OVERLOADED:
media_download_handle_overloaded(inst, data, len);
break;
case MEDIA_SUBCMD_BLOCK_RELAY_FULL:
md_handle_block_relay_full(md, from_node, data, len);
break;
case MEDIA_SUBCMD_BLOCK_CHUNK:
media_download_handle_chunk(inst, data, len);
break;
@ -948,7 +1280,8 @@ static int media_delivery_create_tables(struct UTUN_INSTANCE* inst) {
" group_id INTEGER NOT NULL,"
" node_id INTEGER NOT NULL,"
" chunk INTEGER NOT NULL,"
" timestamp INTEGER NOT NULL);";
" timestamp INTEGER NOT NULL,"
" status INTEGER NOT NULL DEFAULT 1);"; // 0=processing, 1=completed
if (sqlite3_exec(db, sql_ba, NULL, NULL, NULL) != SQLITE_OK) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: CREATE block_availability failed: %s", MD_ID, sqlite3_errmsg(db));
return -1;
@ -980,6 +1313,7 @@ int media_delivery_init(struct UTUN_INSTANCE* inst) {
md->self_node_id = inst->node_id;
md->inst = inst;
md->db = inst->topo_sqlite_db;
md->max_downloads_per_file = 3;
if (media_delivery_create_tables(inst) != 0) return -1;
@ -1050,6 +1384,8 @@ void media_delivery_destroy(struct UTUN_INSTANCE* inst) {
}
if (md->served_nodes) { queue_free(md->served_nodes); md->served_nodes = NULL; }
if (md->downloads) { queue_free(md->downloads); md->downloads = NULL; }
if (md->relay_blocks) { queue_free(md->relay_blocks); md->relay_blocks = NULL; }
if (md->file_loads) { queue_free(md->file_loads); md->file_loads = NULL; }
memset(md, 0, sizeof(*md));
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: destroyed", MD_ID);

71
src/media_delivery/media_delivery.h

@ -38,6 +38,42 @@ struct media_super_peer {
struct media_delivery_ctx* md; // обратная ссылка
};
/* per-file download tracking (source node admission control) */
struct md_file_load {
struct ll_entry ll; // индекс по media_id (16 байт)
uint8_t media_id[16];
int active_downloads; // текущее число исходящих стримов
uint64_t downloader_ids[10]; // кто сейчас качает (для редиректа)
};
/*
* block_availability: жизненный цикл записей и гарантия уникальности id
*
* Каждое событие создаёт новую запись с новым id (AUTOINCREMENT).
* Старая запись удаляется — id НЕ переиспользуется.
* Это критично для репликации: курсор super_sync.last_recv_id
* всегда монотонно растёт и не пропускает события.
*
* BLOCK_PROCESSING (subcmd 0x10):
* md_ba_insert(... , status=0) — чистый INSERT, новый id=X.
*
* HAVE_BLOCK (subcmd 0x06):
* md_ba_complete_block():
* 1. DELETE FROM block_availability WHERE block_uuid=? AND node_id=? AND status=0
* (удаляет processing-запись id=X, если была)
* 2. Если completed уже есть → no-op (возврат 0)
* 3. INSERT ... status=1 — новый id=Y.
*
* SUPER_REPL (subcmd 0x08):
* Принимающая сторона:
* - status=0: INSERT OR IGNORE (дубликаты при репликации — норма)
* - status=1: md_ba_complete_block() — та же логика что HAVE_BLOCK
*
* Итог: processing id=X удаляется, completed вставляется с новым id=Y>X.
* Репликация через md_ba_get_since(id > last_recv_id) всегда видит Y.
*/
struct media_served_node {
struct ll_entry ll; // индекс по node_id (8 байт)
uint64_t node_id;
@ -45,17 +81,41 @@ struct media_served_node {
int64_t joined_at; // timestamp регистрации
};
#define MD_MAX_RELAY_DOWNSTREAM 5
struct relay_downstream {
uint64_t node_id;
uint64_t group_id;
uint64_t sent_offset; // сколько байт уже отправлено этому пиру
struct queue_waiter_handle waiter;
};
struct relay_block_ctx {
struct ll_entry ll; // индекс по block_id (16 байт)
uint8_t block_id[16];
uint8_t media_id[16];
uint32_t chunk;
uint64_t file_offset; // всего байт получено от апстрима (размер .chunk_N)
char chunk_file[2048]; // путь к .chunk_N файлу для relay-чтения
struct relay_downstream downstream[MD_MAX_RELAY_DOWNSTREAM];
int downstream_count;
struct UTUN_INSTANCE* inst;
};
struct media_delivery_ctx {
uint64_t self_node_id;
int initialized;
uint8_t is_supernode; // из adm_tags в peers_*
uint8_t active_streams; // текущее число исходящих стримов (admission)
uint8_t max_downloads_per_file; // лимит параллельных загрузок с одного файла (source node), default 3
uint8_t stream_completed; // 1 = хотя бы один стрим завершился (для тестов)
uint8_t streams_started; // монотонный счётчик запущенных стримов (для тестов)
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 ll_queue* relay_blocks; // relay_block_ctx — relay в процессе загрузки
struct ll_queue* file_loads; // md_file_load — per-file download count (source node)
struct UTUN_INSTANCE* inst;
};
@ -65,6 +125,17 @@ int media_delivery_bind(struct UTUN_INSTANCE* inst);
void media_delivery_set_supernode(struct UTUN_INSTANCE* inst, int is_supernode);
void media_delivery_stream_done(struct UTUN_INSTANCE* inst);
/* relay block context helpers (used by media_download.c) */
struct relay_block_ctx* md_relay_find(struct media_delivery_ctx* md, const uint8_t* block_id);
struct relay_block_ctx* md_relay_add(struct media_delivery_ctx* md,
const uint8_t* block_id, const uint8_t* media_id, uint32_t chunk);
void md_relay_remove(struct media_delivery_ctx* md, const uint8_t* block_id);
/* file_load helpers (per-file source download tracking) */
struct md_file_load* md_file_load_find(struct media_delivery_ctx* md, const uint8_t* media_id);
int md_file_load_inc(struct media_delivery_ctx* md, const uint8_t* media_id, uint64_t node_id);
void md_file_load_dec(struct media_delivery_ctx* md, const uint8_t* media_id, uint64_t node_id);
#ifdef __cplusplus
}
#endif

13
src/media_delivery/media_delivery_proto.h

@ -26,6 +26,8 @@ enum {
MEDIA_SUBCMD_CANCEL = 0x0D, // req-узел→блок-холдер: отмена передачи блока
MEDIA_SUBCMD_SUPER_HELLO = 0x0E, // суперузел↔суперузел: рукопожатие при подключении
MEDIA_SUBCMD_BLOCK_OVERLOADED = 0x0F, // автор→req-узел: перегружен, попробуй позже
MEDIA_SUBCMD_BLOCK_PROCESSING = 0x10, // узел→суперузел: начал качать блок (status=processing)
MEDIA_SUBCMD_BLOCK_RELAY_FULL = 0x11, // relay-узел→req-узел: перегружен, вот список кому транслирую
};
/* ── ответные статусы ── */
@ -164,6 +166,17 @@ struct media_pkt_block_overloaded {
uint32_t retry_after_ms; // рекомендованная задержка перед повтором (ms), 0 = default 2000
};
struct media_pkt_block_relay_full {
uint8_t subcmd; // MEDIA_SUBCMD_BLOCK_RELAY_FULL
uint8_t media_id[16];
uint8_t block_id[16];
uint32_t chunk;
uint16_t num_relay_nodes; // 1..5
// далее: relay_node_ids[num_relay_nodes * 8]
};
#define MEDIA_RELAY_FULL_HDR_SIZE (sizeof(struct media_pkt_block_relay_full)) // 1+16+16+4+2=39
#pragma pack(pop)
#ifdef __cplusplus

183
src/media_delivery/media_download.c

@ -89,6 +89,31 @@ static int md_dl_send_block_req(struct UTUN_INSTANCE* inst, uint64_t dst,
return md_dl_send(inst, dst, (const uint8_t*)&req, sizeof(req));
}
/* ── forward decl ── */
static void md_dl_send_block_processing(struct UTUN_INSTANCE* inst, struct media_download* dl, int bi);
/* ── send BLOCK_REQ + create relay context + notify supernode ── */
static int md_dl_start_block(struct UTUN_INSTANCE* inst, uint64_t dst,
struct media_download* dl, const uint8_t* block_id,
int chunk_idx, int peer_local_bi) {
/* find global block index */
int gbi = -1;
for (int i = 0; i < dl->num_blocks; i++) {
if (memcmp(dl->block_ids + i * 16, block_id, 16) == 0) { gbi = i; break; }
}
if (gbi >= 0) {
/* create relay context before first chunk arrives */
char chunk_file[2048];
snprintf(chunk_file, sizeof(chunk_file), "%s.chunk_%d", dl->dest_path, gbi);
struct relay_block_ctx* rc = md_relay_add(&inst->md, block_id, dl->media_id, (uint32_t)gbi);
if (rc) snprintf(rc->chunk_file, sizeof(rc->chunk_file), "%s", chunk_file);
/* tell supernode we're downloading this block */
md_dl_send_block_processing(inst, dl, gbi);
}
return md_dl_send_block_req(inst, dst, dl, block_id, chunk_idx, peer_local_bi);
}
/* ── send HAVE_BLOCK to supernode ── */
static void md_dl_send_have_block(struct UTUN_INSTANCE* inst, struct media_download* dl, int bi) {
@ -119,6 +144,36 @@ static void md_dl_send_have_block(struct UTUN_INSTANCE* inst, struct media_downl
md_dl_send(inst, super, (const uint8_t*)&hb, sizeof(hb));
}
/* ── send BLOCK_PROCESSING to supernode ── */
static void md_dl_send_block_processing(struct UTUN_INSTANCE* inst, struct media_download* dl, int bi) {
uint64_t super = dl->super_count > 0 ? dl->super_nodes[dl->super_current] : 0;
if (!super) { DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: BLOCK_PROCESSING but no supernode", MDL_ID); return; }
struct media_pkt_have_block hb;
memset(&hb, 0, sizeof(hb));
hb.subcmd = MEDIA_SUBCMD_BLOCK_PROCESSING;
hb.group_id = dl->group_id;
memcpy(hb.media_id, dl->media_id, 16);
memcpy(hb.block_id, dl->block_ids + bi * 16, 16);
hb.chunk = (uint32_t)bi;
hb.timestamp = (int64_t)time(NULL);
uint8_t smsg[64]; size_t soff = 0;
memcpy(smsg + soff, dl->block_ids + bi * 16, 16); soff += 16;
memcpy(smsg + soff, &hb.chunk, 4); soff += 4;
memcpy(smsg + soff, &hb.timestamp, 8); soff += 8;
uint64_t self = inst->node_id; memcpy(smsg + soff, &self, 8); soff += 8;
if (sc_ed25519_sign(inst->my_ed25519_privkey, smsg, soff, hb.node_sign) != SC_OK) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: Ed25519 sign failed for BLOCK_PROCESSING", MDL_ID);
return;
}
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: BLOCK_PROCESSING to super 0x%016llx block=%d",
MDL_ID, (unsigned long long)super, bi);
md_dl_send(inst, super, (const uint8_t*)&hb, sizeof(hb));
}
/* ── conn_mgr callback ── */
void md_dl_conn_cb(struct CONN_MGR_HANDLE* h, uint64_t node_id, uint64_t group_id,
@ -142,8 +197,8 @@ void md_dl_conn_cb(struct CONN_MGR_HANDLE* h, uint64_t node_id, uint64_t group_i
if (!dl->peers[pi].blocks[bi].started) {
dl->peers[pi].blocks[bi].started = 1;
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: conn_cb BLOCK_REQ: peer-local bi=%d",
MDL_ID, bi);
md_dl_send_block_req(dl->inst, node_id, dl, dl->peers[pi].blocks[bi].block_id, bi, bi);
MDL_ID, bi);
md_dl_start_block(dl->inst, node_id, dl, dl->peers[pi].blocks[bi].block_id, bi, bi);
break;
}
}
@ -300,7 +355,7 @@ void media_download_handle_query_resp(struct UTUN_INSTANCE* inst,
for (int bi = 0; bi < dl->peers[pi].num_blocks; bi++) {
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: direct BLOCK_REQ: peer-local bi=%d peer_num_blocks=%d",
MDL_ID, bi, dl->peers[pi].num_blocks);
md_dl_send_block_req(inst, dl->peers[pi].node_id, dl, dl->peers[pi].blocks[bi].block_id, bi, bi);
md_dl_start_block(inst, dl->peers[pi].node_id, dl, dl->peers[pi].blocks[bi].block_id, bi, bi);
}
}
}
@ -335,6 +390,41 @@ void media_download_handle_chunk(struct UTUN_INSTANCE* inst,
MDL_ID, len, (size_t)MEDIA_BLOCK_CHUNK_HDR_SIZE, (size_t)wlen);
}
fclose(f);
/* forward to relay downstreams */
struct relay_block_ctx* rc = md_relay_find(&inst->md, ch->block_id);
if (rc) {
rc->file_offset += wlen;
for (int i = 0; i < rc->downstream_count; i++) {
struct relay_downstream* ds = &rc->downstream[i];
if (ds->sent_offset >= rc->file_offset) continue;
/* read from chunk file and send catch-up to this downstream */
FILE* rf = fopen(tmp, "rb");
if (!rf) continue;
while (ds->sent_offset < rc->file_offset) {
fseeko(rf, (off_t)ds->sent_offset, SEEK_SET);
size_t to_read = rc->file_offset - ds->sent_offset;
if (to_read > 1024) to_read = 1024;
uint8_t rbuf[1024];
size_t rd = fread(rbuf, 1, to_read, rf);
if (rd == 0) break;
uint8_t rpkt[MEDIA_BLOCK_CHUNK_HDR_SIZE + 1024];
struct media_pkt_block_chunk* rch = (struct media_pkt_block_chunk*)rpkt;
memset(rch, 0, MEDIA_BLOCK_CHUNK_HDR_SIZE);
rch->subcmd = MEDIA_SUBCMD_BLOCK_CHUNK;
memcpy(rch->media_id, ch->media_id, 16);
memcpy(rch->block_id, ch->block_id, 16);
rch->chunk = ch->chunk;
rch->offset = (uint32_t)ds->sent_offset;
rch->data_len = (uint16_t)rd;
memcpy(rpkt + MEDIA_BLOCK_CHUNK_HDR_SIZE, rbuf, rd);
if (md_dl_send(inst, ds->node_id, rpkt, MEDIA_BLOCK_CHUNK_HDR_SIZE + rd) == 0)
ds->sent_offset += rd;
else break;
}
fclose(rf);
}
}
}
/* ── handle incoming BLOCK_DONE ── */
@ -407,6 +497,26 @@ void media_download_handle_done(struct UTUN_INSTANCE* inst,
/* send HAVE_BLOCK to supernode */
md_dl_send_have_block(inst, dl, bi);
/* forward BLOCK_DONE to relay downstreams */
{
struct relay_block_ctx* rc = md_relay_find(&inst->md, dl->block_ids + bi * 16);
if (rc) {
for (int i = 0; i < rc->downstream_count; i++) {
struct media_pkt_block_done rd;
memset(&rd, 0, sizeof(rd));
rd.subcmd = MEDIA_SUBCMD_BLOCK_DONE;
memcpy(rd.media_id, dl->media_id, 16);
memcpy(rd.block_id, dl->block_ids + bi * 16, 16);
rd.chunk = (uint32_t)bi;
rd.total_size = (uint32_t)rc->file_offset;
md_dl_send(inst, rc->downstream[i].node_id, (const uint8_t*)&rd, sizeof(rd));
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: relay BLOCK_DONE forwarded to 0x%016llx block=%d",
MDL_ID, (unsigned long long)rc->downstream[i].node_id, bi);
}
md_relay_remove(&inst->md, dl->block_ids + bi * 16);
}
}
/* send next block from the same peer (backpressure: one at a time) */
for (int pi = 0; pi < dl->num_peers; pi++) {
for (int bj = 0; bj < dl->peers[pi].num_blocks; bj++) {
@ -416,8 +526,8 @@ void media_download_handle_done(struct UTUN_INSTANCE* inst,
dl->peers[pi].blocks[bk].started = 1;
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: BLOCK_DONE → next req: peer[%d] bi=%d",
MDL_ID, pi, bk);
md_dl_send_block_req(inst, dl->peers[pi].node_id, dl,
dl->peers[pi].blocks[bk].block_id, bk, bk);
md_dl_start_block(inst, dl->peers[pi].node_id, dl,
dl->peers[pi].blocks[bk].block_id, bk, bk);
goto blockreq_sent;
}
}
@ -592,3 +702,66 @@ int media_download_cancel(struct UTUN_INSTANCE* inst,
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: download cancelled", MDL_ID);
return 0;
}
/* ── handle RELAY_FULL from relay node: retry block request from one of the listed nodes ── */
void media_download_handle_relay_full(struct UTUN_INSTANCE* inst,
const uint8_t* media_id, const uint8_t* block_id,
uint32_t chunk, const uint64_t* node_ids, int num_nodes) {
struct media_download* dl = md_dl_find(inst, media_id);
if (!dl || !dl->active) { DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "%s: RELAY_FULL — no active download", MDL_ID); return; }
int bi = -1;
for (int i = 0; i < dl->num_blocks; i++) {
if (memcmp(dl->block_ids + i * 16, block_id, 16) == 0) { bi = i; break; }
}
if (bi < 0) { DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "%s: RELAY_FULL — unknown block", MDL_ID); return; }
/* try each listed node — if already a peer, retry block; otherwise add new */
for (int ni = 0; ni < num_nodes && ni < MD_MAX_RELAY_DOWNSTREAM; ni++) {
int found_pi = -1;
for (int j = 0; j < dl->num_peers; j++) {
if (dl->peers[j].node_id == node_ids[ni]) { found_pi = j; break; }
}
if (found_pi >= 0) {
/* already a peer — reset block started flag and retry */
for (int bj = 0; bj < dl->peers[found_pi].num_blocks; bj++) {
if (memcmp(dl->peers[found_pi].blocks[bj].block_id, block_id, 16) == 0) {
dl->peers[found_pi].blocks[bj].started = 0;
dl->peers[found_pi].blocks[bj].received = 0;
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: RELAY_FULL → retry existing peer[%d] node=0x%016llx bi=%d",
MDL_ID, found_pi, (unsigned long long)node_ids[ni], bj);
md_dl_start_block(inst, node_ids[ni], dl, block_id, (int)chunk, bj);
return;
}
}
/* peer doesn't have this block — add it */
int nb = dl->peers[found_pi].num_blocks;
if (nb >= MD_MAX_BLOCKS_PER_PEER) continue;
memcpy(dl->peers[found_pi].blocks[nb].block_id, block_id, 16);
dl->peers[found_pi].blocks[nb].started = 0;
dl->peers[found_pi].blocks[nb].received = 0;
dl->peers[found_pi].blocks[nb].validated = 0;
dl->peers[found_pi].num_blocks++;
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: RELAY_FULL → add block to existing peer[%d] bi=%d",
MDL_ID, found_pi, nb);
md_dl_start_block(inst, node_ids[ni], dl, block_id, (int)chunk, nb);
return;
}
/* new peer */
int pi = dl->num_peers;
if (pi >= 10) break;
dl->peers[pi].node_id = node_ids[ni];
dl->peers[pi].connected = 1;
dl->peers[pi].num_blocks = 1;
memcpy(dl->peers[pi].blocks[0].block_id, block_id, 16);
dl->peers[pi].blocks[0].started = 0;
dl->peers[pi].blocks[0].received = 0;
dl->peers[pi].blocks[0].validated = 0;
dl->num_peers++;
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: RELAY_FULL → new peer[%d] node=0x%016llx",
MDL_ID, pi, (unsigned long long)node_ids[ni]);
md_dl_start_block(inst, node_ids[ni], dl, block_id, (int)chunk, 0);
break;
}
}

4
src/media_delivery/media_download.h

@ -105,6 +105,10 @@ void media_download_handle_done(struct UTUN_INSTANCE* inst,
void media_download_handle_overloaded(struct UTUN_INSTANCE* inst,
const uint8_t* data, size_t len);
void media_download_handle_relay_full(struct UTUN_INSTANCE* inst,
const uint8_t* media_id, const uint8_t* block_id,
uint32_t chunk, const uint64_t* node_ids, int num_nodes);
/* ── conn_mgr callback (non-static for test visibility) ── */
struct CONN_MGR_HANDLE;
enum conn_mgr_event;

251
src/media_delivery/sync.md

@ -0,0 +1,251 @@
# Репликация block_availability между суперузлами
## 1. Таблицы в SQLite
### `block_availability`
Хранит информацию о том, какие узлы имеют какие блоки:
```
id INTEGER PRIMARY KEY AUTOINCREMENT
block_uuid BLOB — UUID блока (16 байт)
group_id INTEGER — id группы
node_id INTEGER — id узла-владельца блока
chunk INTEGER — номер чанка
timestamp INTEGER — unix-время добавления/обновления
```
Индексы: `block_uuid+node_id` (unique), `id`, `node_id`.
### `super_sync`
Хранит указатель, до какого `id` синхронизировано с каждым другим суперузлом:
```
peer_node_id INTEGER PRIMARY KEY — id пира-суперузла
last_recv_id INTEGER — максимальный id, подтверждённый пиром (что он принял от нас)
```
Обе таблицы создаются в `media_delivery_create_tables()` (media_delivery.c:940-971).
На всех узлах есть `block_availability`. Суперузлы сохраняют в неё весь контент групп, в которых состоят. Обычные узлы — только свой контент. Если нет суперузлов или они недоступны, роль суперузла выполняет автор контента.
## 2. Обнаружение других суперузлов
Суперузлы определяются через свойство `supernode=yes` в `adm_tags` таблицы `peers_<channel>`.
Механизмы обнаружения:
- **При старте** (`md_super_start`, стр. 730): суперузел ищет в БД все узлы с `node_type=4` в своей chat-группе и инициирует прямое подключение к каждому через `conn_mgr_open_invite`.
- **При появлении нового узла в BGP** (`md_on_bgp_node`, стр. 790): проверяет `node_type=4` в `peers_<channel>` — если да, вызывает `md_super_connect`.
- **При изменении свойств** (`md_on_props_changed`, стр. 825): если `adm_tags` изменились на `supernode=yes` — подключается; если на `no` — удаляет пира.
Как только суперузел видит появление другого суперузла в BGP-таблице — пробует установить прямую связь. Прямое подключение — независимый автономный процесс (без посредников и без обратного подключения — это сделает другой суперузел сам). Если не получается 10 сек — отключаемся и следующая попытка через час.
## 3. Протокол репликации
### Handshake: SUPER_HELLO (subcmd 0x0E)
После установки прямого подключения инициатор отправляет `SUPER_HELLO` (`md_super_conn_cb`, стр. 484) с полем `last_recv_id` — до какого `id` он уже получил записи от этого пира.
При получении `SUPER_HELLO` (`md_handle_super_hello`, стр. 454):
- Создаёт/обновляет `media_super_peer` в памяти
- Устанавливает `connected=1`, `hello_done=1`
- Отправляет ответный `SUPER_HELLO` со своим `last_recv_id`
- Обнуляет `inflight_count` и запускает репликацию через `md_super_repl_send`
### Репликация: SUPER_REPL (subcmd 0x08)
`md_super_repl_send` (стр. 228):
- Запрашивает из `block_availability` записи с `id > peer->peer_last_recv_id`
- Порциями по 4 записи (константа `MEDIA_MAX_REPL_INFLIGHT = 4`)
- Каждая запись: 52 байта (id:8, block_uuid:16, group_id:8, node_id:8, chunk:4, timestamp:8)
- Пакет: `subcmd` + `seq` (макс id в пачке) + `num_entries` + сырые строки
- Инкрементит `inflight_count`
### Подтверждение: SUPER_ACK (subcmd 0x09)
При получении `SUPER_ACK` (`md_handle_super_ack`, стр. 436):
- Обновляется `peer_last_recv_id = ack_seq`
- Декрементится `inflight_count`
- Таймаут сбрасывается на начальное значение (`MEDIA_REPL_TIMEOUT_TB = 2с`)
- Запускается следующая порция репликации (если есть записи)
### Приём репликации (`md_handle_super_repl`, стр. 404):
- Разбирает каждую запись, вызывает `md_ba_insert` (INSERT OR REPLACE в `block_availability`)
- Сохраняет `max_id` в `super_sync` через `md_ss_set(db, from_node, max_id)`
- Отправляет `SUPER_ACK` с подтверждённым `max_id`
## 4. Таймауты и retransmit
- **Начальный таймаут:** 2 секунды (`MEDIA_REPL_TIMEOUT_TB`)
- **Exponential backoff:** при отсутствии ACK таймаут удваивается (до 10 минут — `MEDIA_REPL_TIMEOUT_MAX_TB`)
- **При успешном ACK:** таймаут сбрасывается на 2 секунды
- **При неудачном подключении:** повтор через 1 час (`MEDIA_RECONNECT_COOLDOWN_TB`)
- Механика: `md_repl_schedule_retrans` → `md_repl_retrans_cb` → `md_super_repl_send`
## 5. Источники записей в block_availability
Записи попадают в таблицу из трёх источников:
1. **HAVE_BLOCK** — когда req-узел принял полный блок, он уведомляет свой суперузел (`md_handle_have_block`, стр. 371)
2. **SUPER_REPL** — репликация от другого суперузла
3. **Локально** — при регистрации своего медиа через `media_index`
При получении HAVE_BLOCK суперузел немедленно реплицирует новые записи всем подключённым super-пирам (`md_handle_have_block`, стр. 398-400).
## 6. Кто есть суперузел
- `node_type=4` в `peers_<channel>` + `adm_tags` содержит `supernode=yes`
- `md->is_supernode` устанавливается при инициализации и может меняться через `media_delivery_set_supernode` / `md_on_props_changed`
- Если узел перестаёт быть суперузлом — вызывается `md_super_stop`: удаляются все чужие блоки из `block_availability`, очищается очередь свер-пиров
## 7. Обслуживаемые узлы (served_nodes)
- Обычный узел при старте загрузки отправляет суперузлу `SERVE_REG` (subcmd 0x01)
- Суперузел сохраняет его в `served_nodes` (ll_queue с индексом по node_id)
- Если соединение рвётся — узел автоматически выпадает из списка (`md_on_conn_status`)
- Также автор контента получает от суперузла обновления о доступности его блоков
## 8. Очистка
- **При удалении узла из BGP** (`md_on_bgp_node`, TOPO_NODE_EVENT_REMOVE): удаляются все его записи из `block_availability`, свер-пир удаляется
- **При разрыве соединения** (`md_on_conn_status`, ETCP_CONN_STATUS_DOWN/DELETE): свер-пир помечается как отключённый (`connected=0, hello_done=0, inflight_count=0`), таймеры отменяются; served_node удаляется
- **При отключении режима суперузла** (`md_super_stop`): удаляются все чужие блоки (`WHERE node_id != self`), очищается очередь свер-пиров
## 9. Диаграмма потоков
```
[Узел A: обычный] [Суперузел S1] [Суперузел S2]
| | |
|── HAVE_BLOCK ─────────────────>| |
| |── md_ba_insert (локально) |
| |── SUPER_REPL ────────────>|
| | |── md_ba_insert (локально)
| |<─────────── SUPER_ACK ────|
| |── md_ss_set(S2, max_id) |
| | |
|── QUERY ──────────────────────>| |
|<───── QUERY_RESP (node list) ──| |
| | |
|── BLOCK_REQ ─────────────────────────────────────────────>|
|<───── BLOCK_CHUNK (stream) ───────────────────────────────|
|<───── BLOCK_DONE ─────────────────────────────────────────|
| | |
|── HAVE_BLOCK ─────────────────>| |
| |── SUPER_REPL ────────────>|
```
## 10. Константы
| Константа | Значение | Описание |
|-----------|----------|----------|
| `MEDIA_MAX_REPL_INFLIGHT` | 4 | Макс пакетов в полёте на одного пира |
| `MEDIA_REPL_TIMEOUT_TB` | 20000 | Начальный таймаут 2с (0.1ms units) |
| `MEDIA_REPL_TIMEOUT_MAX_TB` | 6000000 | Макс таймаут 10 мин |
| `MEDIA_RECONNECT_COOLDOWN_TB` | 360000000 | Повтор подключения через 1 час |
| `MEDIA_HELLO_TIMEOUT_TB` | 20000 | Таймаут SUPER_HELLO 2с |
## 11. BLOCK_PROCESSING: регистрация блока в процессе загрузки
Когда узел начинает скачивать блок, он отправляет суперузлу `BLOCK_PROCESSING` (subcmd 0x10).
Суперузел вставляет запись в `block_availability` со `status=0` (processing) через `INSERT OR IGNORE`
— не перезаписывает существующие completed-записи (status=1).
Когда блок полностью скачан, узел отправляет `HAVE_BLOCK` (subcmd 0x06).
Суперузел вызывает `INSERT OR REPLACE` со `status=1` — processing-запись заменяется на completed.
Записи с `status=0` реплицируются между суперузлами так же, как и completed (поле `status` включено
в сериализацию `SUPER_REPL`, размер строки: 56 байт вместо 52).
### BLOCK_PROCESSING vs HAVE_BLOCK
| Операция | subcmd | SQL | Статус | id |
|----------|--------|-----|--------|----|
| Начал качать | 0x10 | INSERT (status=0) | processing | новый X |
| Скачал | 0x06 | DELETE (status=0) + проверка + INSERT (status=1) | completed | новый Y>X |
**Гарантия:** каждый переход создаёт новую запись с новым id. Старая удаляется — id не переиспользуется.
Если completed уже существует — HAVE_BLOCK игнорируется (no-op).
Если processing пытается вставиться поверх completed — `SQLITE_CONSTRAINT`, логируется WARN.
### Схема block_availability (обновлённая)
```
id INTEGER PRIMARY KEY AUTOINCREMENT — монотонный, не переиспользуется
block_uuid BLOB — UUID блока (16 байт)
group_id INTEGER — id группы
node_id INTEGER — id узла-владельца блока
chunk INTEGER — номер чанка
timestamp INTEGER — unix-время добавления/обновления
status INTEGER — 0=processing (качается), 1=completed (скачан)
```
Уникальный индекс: `(block_uuid, node_id)` — предотвращает две completed-записи для одного блока.
## 12. Relay-стриминг: передача блока в процессе загрузки
### Принцип
Если узел C хочет скачать блок, а узел B уже качает этот блок от узла A:
1. C видит B в QUERY_RESP (благодаря BLOCK_PROCESSING, status=0)
2. C отправляет BLOCK_REQ к B
3. B принимает запрос, добавляет C в `relay_downstream` для этого блока
4. B получает BLOCK_CHUNK от A → пишет в файл → форвардит C
5. B получает BLOCK_DONE от A → форвардит C → удаляет relay-контекст
### Flow control
- Relay-узел всегда отстаёт от апстрима: он не может слать быстрее, чем получает
- Каждый downstream хранит `sent_offset` — сколько байт уже отправлено
- При получении нового чанка: читается файл, отправляется отставание (до 1KB за раз)
- При ошибке отправки — чанк пропускается, ретрай на следующем чанке или в BLOCK_DONE
- Бэкпрессур через `etcp_router_on_send_ready` для асинхронной досылки
### Переполнение relay (RELAY_FULL)
- Максимум 5 downstream-узлов на блок (`MD_MAX_RELAY_DOWNSTREAM`)
- При переполнении B отправляет C пакет `RELAY_FULL` (subcmd 0x11)
со списком node_id текущих downstream-узлов
- C извлекает список и пробует запросить блок у одного из этих узлов
### Структуры данных
**relay_block_ctx** — контекст relay для одного блока (индекс по block_id в `md->relay_blocks`):
- `block_id[16]`, `media_id[16]`, `chunk` — идентификаторы блока
- `chunk_file[2048]` — путь к `.chunk_N` файлу (источник данных для форвардинга)
- `file_offset` — всего получено байт от апстрима (размер файла)
- `downstream[5]` — массив relay_downstream
**relay_downstream** — один пир, запросивший relay:
- `node_id`, `group_id` — кому слать
- `sent_offset` — сколько байт уже отправлено этому пиру
- `waiter` — handle для бэкпрессур-досылки
### Диаграмма relay-потока
```
[Узел A: источник] [Узел B: relay] [Узел C: req]
| | |
| |<── BLOCK_REQ ────────────|
| |── relay downstream add |
| |── catchup send (если есть)|
| | |
|── BLOCK_CHUNK ──────────>| |
| |── fwrite (в файл) |
| |── fread → BLOCK_CHUNK ──>|
| | |
|── BLOCK_DONE ───────────>| |
| |── fwd BLOCK_DONE ───────>|
| |── md_relay_remove |
```
### Жизненный цикл relay_block_ctx
1. **Создание** — в `md_dl_start_block` при отправке первого BLOCK_REQ
2. **Обновление file_offset** — в `media_download_handle_chunk` при каждом чанке
3. **Форвардинг downstream** — в `media_download_handle_chunk` после записи
4. **Завершение** — в `media_download_handle_done` после валидации: форвард BLOCK_DONE + удаление
5. **Отмена** — при `media_download_cancel`: контекст удаляется вместе с `md->relay_blocks`
### Константы relay
| Константа | Значение | Описание |
|-----------|----------|----------|
| `MD_MAX_RELAY_DOWNSTREAM` | 5 | Макс число downstream-узлов на блок |
| `MEDIA_RELAY_FULL_HDR_SIZE` | 39 | Размер заголовка RELAY_FULL (1+16+16+4+2) |

108
tests/test_media_delivery_sql.c

@ -48,7 +48,8 @@ static int create_tables(sqlite3* db) {
" group_id INTEGER NOT NULL,"
" node_id INTEGER NOT NULL,"
" chunk INTEGER NOT NULL,"
" timestamp INTEGER NOT NULL)";
" timestamp INTEGER NOT NULL,"
" status INTEGER NOT NULL DEFAULT 1)"; // 0=processing, 1=completed
int rc = sqlite3_exec(db, sql_ba, NULL, NULL, NULL);
if (rc != SQLITE_OK) return -1;
sqlite3_exec(db, "CREATE UNIQUE INDEX IF NOT EXISTS idx_ba_uuid_node ON block_availability(block_uuid, node_id)", NULL, NULL, NULL);
@ -220,6 +221,110 @@ static void test_ba_get_since(void) {
}
}
/* ── status column tests ── */
static void test_ba_status(void) {
/* processing→completed: DELETE old + INSERT new, ids don't reuse */
TEST("ba_complete: new id > processing id"); {
sqlite3* db = make_db(); create_tables(db);
uint8_t uuid[16]; for (int i = 0; i < 16; i++) uuid[i] = (uint8_t)(0xAA + i);
sqlite3_stmt* s = NULL;
/* 1. INSERT processing (status=0) */
sqlite3_prepare_v2(db, "INSERT INTO block_availability(block_uuid,group_id,node_id,chunk,timestamp,status) VALUES(?,1,0xAA,0,100,0)", -1, &s, NULL);
sqlite3_bind_blob(s, 1, uuid, 16, SQLITE_STATIC);
sqlite3_step(s); sqlite3_finalize(s);
/* 2. get processing id */
int64_t pid = -1;
sqlite3_prepare_v2(db, "SELECT id FROM block_availability WHERE node_id=0xAA AND status=0", -1, &s, NULL);
if (sqlite3_step(s) == SQLITE_ROW) pid = sqlite3_column_int64(s, 0);
sqlite3_finalize(s);
/* 3. DELETE processing + INSERT completed (new id) */
sqlite3_exec(db, "DELETE FROM block_availability WHERE block_uuid=x'AAABACADAEAFB0B1B2B3B4B5B6B7B8B9' AND node_id=0xAA AND status=0", NULL, NULL, NULL);
sqlite3_prepare_v2(db, "INSERT INTO block_availability(block_uuid,group_id,node_id,chunk,timestamp,status) VALUES(?,1,0xAA,0,200,1)", -1, &s, NULL);
sqlite3_bind_blob(s, 1, uuid, 16, SQLITE_STATIC);
sqlite3_step(s); sqlite3_finalize(s);
/* 4. get completed id */
int64_t cid = -1;
sqlite3_prepare_v2(db, "SELECT id,status FROM block_availability WHERE node_id=0xAA", -1, &s, NULL);
int st = -1;
if (sqlite3_step(s) == SQLITE_ROW) { cid = sqlite3_column_int64(s, 0); st = sqlite3_column_int(s, 1); }
sqlite3_finalize(s);
int n = count_rows(db, "block_availability");
if (n == 1 && st == 1 && cid > pid) OK(); else FAIL("rows=%d status=%d pid=%lld cid=%lld (expected 1 row, status=1, cid>pid)", n, st, (long long)pid, (long long)cid);
sqlite3_close(db);
}
/* completed already exists → no-op, no duplicate */
TEST("ba_complete: no-op when completed exists"); {
sqlite3* db = make_db(); create_tables(db);
uint8_t uuid[16]; for (int i = 0; i < 16; i++) uuid[i] = (uint8_t)(0xBB + i);
sqlite3_stmt* s = NULL;
/* 1. INSERT completed */
sqlite3_prepare_v2(db, "INSERT INTO block_availability(block_uuid,group_id,node_id,chunk,timestamp,status) VALUES(?,1,0xBB,0,300,1)", -1, &s, NULL);
sqlite3_bind_blob(s, 1, uuid, 16, SQLITE_STATIC);
sqlite3_step(s); sqlite3_finalize(s);
int64_t id1 = -1;
sqlite3_prepare_v2(db, "SELECT id FROM block_availability WHERE node_id=0xBB", -1, &s, NULL);
if (sqlite3_step(s) == SQLITE_ROW) id1 = sqlite3_column_int64(s, 0);
sqlite3_finalize(s);
/* 2. DELETE processing (none) + check completed exists + skip INSERT */
sqlite3_exec(db, "DELETE FROM block_availability WHERE block_uuid=x'BBBCBDBEBFC0C1C2C3C4C5C6C7C8C9' AND node_id=0xBB AND status=0", NULL, NULL, NULL);
int exists = 0;
sqlite3_prepare_v2(db, "SELECT 1 FROM block_availability WHERE block_uuid=? AND node_id=0xBB AND status=1", -1, &s, NULL);
sqlite3_bind_blob(s, 1, uuid, 16, SQLITE_STATIC);
exists = (sqlite3_step(s) == SQLITE_ROW);
sqlite3_finalize(s);
/* 3. Don't insert if exists */
int64_t id2 = -1;
sqlite3_prepare_v2(db, "SELECT id FROM block_availability WHERE node_id=0xBB", -1, &s, NULL);
if (sqlite3_step(s) == SQLITE_ROW) id2 = sqlite3_column_int64(s, 0);
sqlite3_finalize(s);
int n = count_rows(db, "block_availability");
if (n == 1 && exists && id2 == id1) OK(); else FAIL("rows=%d exists=%d id1=%lld id2=%lld (expected 1 row, same id)", n, exists, (long long)id1, (long long)id2);
sqlite3_close(db);
}
/* processing after completed INSERT → CONSTRAINT violation */
TEST("ba_processing after completed: constraint violation"); {
sqlite3* db = make_db(); create_tables(db);
uint8_t uuid[16]; for (int i = 0; i < 16; i++) uuid[i] = (uint8_t)(0xCC + i);
sqlite3_stmt* s = NULL;
/* 1. INSERT completed */
sqlite3_prepare_v2(db, "INSERT INTO block_availability(block_uuid,group_id,node_id,chunk,timestamp,status) VALUES(?,1,0xCC,0,400,1)", -1, &s, NULL);
sqlite3_bind_blob(s, 1, uuid, 16, SQLITE_STATIC);
sqlite3_step(s); sqlite3_finalize(s);
/* 2. try INSERT processing (same uuid+node) → must fail with unique constraint */
int rc = SQLITE_OK;
sqlite3_prepare_v2(db, "INSERT INTO block_availability(block_uuid,group_id,node_id,chunk,timestamp,status) VALUES(?,1,0xCC,0,500,0)", -1, &s, NULL);
sqlite3_bind_blob(s, 1, uuid, 16, SQLITE_STATIC);
rc = sqlite3_step(s);
sqlite3_finalize(s);
/* 3. completed is still there, unchanged */
int n = count_rows(db, "block_availability");
int64_t ts_check = -1;
sqlite3_prepare_v2(db, "SELECT timestamp FROM block_availability WHERE node_id=0xCC", -1, &s, NULL);
if (sqlite3_step(s) == SQLITE_ROW) ts_check = sqlite3_column_int64(s, 0);
sqlite3_finalize(s);
if (n == 1 && rc == SQLITE_CONSTRAINT && ts_check == 400) OK();
else FAIL("rows=%d rc=%d ts=%lld (expected 1 row, CONSTRAINT, ts=400)", n, rc, (long long)ts_check);
sqlite3_close(db);
}
}
static void test_ss_ops(void) {
TEST("super_sync insert+select"); {
sqlite3* db = make_db(); create_tables(db);
@ -352,6 +457,7 @@ int main(void) {
test_ba_insert();
test_ba_find();
test_ba_get_since();
test_ba_status();
test_ss_ops();
test_ba_delete();
test_proto_sizes();

Loading…
Cancel
Save