Browse Source

docs: clarify core, group and chat header contracts

proxy
evgeny 3 days ago
parent
commit
01fcefdec5
  1. 58
      src/chat/chat_core.h
  2. 14
      src/chat/chat_join.h
  3. 37
      src/chat/chat_sync.h
  4. 65
      src/chat/db_sync.h
  5. 55
      src/chat/member_sync.h
  6. 21
      src/chat/merkle_sync.h
  7. 15
      src/chat/merkle_tree.h
  8. 18
      src/media_async/media_async.h
  9. 60
      src/routing_layer/conn_mgr.h
  10. 33
      src/routing_layer/etcp_router.h
  11. 92
      src/routing_layer/topo_group.h
  12. 24
      src/routing_layer/topo_group_connect.h
  13. 37
      src/routing_layer/topo_group_invite.h
  14. 72
      src/routing_layer/topo_node.h
  15. 7
      src/routing_layer/topo_recovery.h
  16. 49
      src/transport_layer/etcp.h
  17. 44
      src/transport_layer/etcp_api.h
  18. 31
      src/transport_layer/node_conn_direct.h
  19. 76
      src/utun_instance.h

58
src/chat/chat_core.h

@ -1,5 +1,7 @@
/* /*
* chat_core.h — центральный API чата в потоке uasync * chat_core — API каналов, сообщений, профиля и настроек одного узла.
* chat_service_start/stop управляют всем чат-сервисом поверх работающего ядра.
* Данные хранятся в общей SQLite; события для интерфейса доставляет chat_event.
* *
* Все DB-операции и сетевые функции выполняются в uasync-потоке. * Все DB-операции и сетевые функции выполняются в uasync-потоке.
* Все функции per-instance: первым аргументом struct UTUN_INSTANCE*. * Все функции per-instance: первым аргументом struct UTUN_INSTANCE*.
@ -19,12 +21,16 @@ struct UTUN_INSTANCE;
/* ── Жизненный цикл ── */ /* ── Жизненный цикл ── */
/* Run on the owning uasync thread, outside callbacks. Stop preserves core/UTUN/DB. /* Запустить чат поверх core_start; UTUN не требуется. Повтор допустим; 0/-1. */
* Stop cancels delivery and waits for native media workers before releasing chat state. */
int chat_service_start(struct UTUN_INSTANCE* inst); int chat_service_start(struct UTUN_INSTANCE* inst);
/* Остановить чат, сохранив ядро, UTUN и БД. Повтор допустим; вызывать вне callbacks чата.
* Отменяет доставку и ждёт native media workers; длительная задача может задержать stop. */
void chat_service_stop(struct UTUN_INSTANCE* inst); void chat_service_stop(struct UTUN_INSTANCE* inst);
/* Внутренняя часть запуска: открыть/подключить БД и загрузить каналы. 0/-1.
* Для запуска всего чат-сервиса использовать chat_service_start(). */
int chat_core_init(struct UTUN_INSTANCE* inst, const char* db_path); int chat_core_init(struct UTUN_INSTANCE* inst, const char* db_path);
/* Освободить chat_core после остановки остальных частей чата; общая БД остаётся у ядра. */
void chat_core_destroy(struct UTUN_INSTANCE* inst); void chat_core_destroy(struct UTUN_INSTANCE* inst);
/* Backfill: регистрирует уже скачанные медиафайлы (лежат на диске, но отсутствуют /* Backfill: регистрирует уже скачанные медиафайлы (лежат на диске, но отсутствуют
@ -36,33 +42,36 @@ void chat_media_backfill(struct UTUN_INSTANCE* inst);
void chat_media_autodownload_backfill(struct UTUN_INSTANCE* inst); void chat_media_autodownload_backfill(struct UTUN_INSTANCE* inst);
/* Регистрирует триггер: после первого SYNC_DONE (+ дебаунс 2с) анонсирует локальные /* Регистрирует триггер: после первого SYNC_DONE (+ дебаунс 2с) анонсирует локальные
* блоки суперузлам и запускает chat_media_autodownload_backfill. Один раз за процесс. */ * блоки суперузлам и запускает chat_media_autodownload_backfill. Один раз за запуск чата. */
void chat_media_startup_backfill_init(struct UTUN_INSTANCE* inst); void chat_media_startup_backfill_init(struct UTUN_INSTANCE* inst);
void chat_media_startup_backfill_destroy(struct UTUN_INSTANCE* inst); void chat_media_startup_backfill_destroy(struct UTUN_INSTANCE* inst);
/* Заимствованная БД текущего chat_core; NULL до init/после destroy. */
struct sqlite3* chat_core_get_db(struct UTUN_INSTANCE* inst); struct sqlite3* chat_core_get_db(struct UTUN_INSTANCE* inst);
/* 1 — chat_core инициализирован, 0 — нет; не проверяет доступность сети. */
int chat_core_is_initialized(struct UTUN_INSTANCE* inst); int chat_core_is_initialized(struct UTUN_INSTANCE* inst);
/* ── Отправка сообщения (GUI → uasync) ── */ /* ── Отправка сообщения (GUI → uasync) ── */
struct chat_msg_submit { struct chat_msg_submit {
struct UTUN_INSTANCE* inst; struct UTUN_INSTANCE* inst;
char channel_id[64]; char channel_id[64]; /* десятичный ID канала */
char content_type[32]; char content_type[32]; /* тип содержимого сообщения */
char media_src[1024]; char media_src[1024]; /* исходный файл вложения */
char media_dest[1024]; char media_dest[1024]; /* путь подготовленного локального файла */
uint8_t media_copy; uint8_t media_copy; /* 1 = скопировать исходный файл при регистрации */
uint8_t media_video; /* 1 = подготовить как видео (пробинг+транскод) */ uint8_t media_video; /* 1 = подготовить как видео (пробинг+транскод) */
int duration_ms; /* видео: длительность (0 = неизвестно) */ int duration_ms; /* видео: длительность (0 = неизвестно) */
int width; /* видео: ширина */ int width; /* видео: ширина */
int height; /* видео: высота */ int height; /* видео: высота */
uint8_t* data; uint8_t* data; /* содержимое сообщения, заимствовано обычным submit */
uint32_t data_len; uint32_t data_len; /* размер data в байтах */
uint64_t timestamp; uint64_t timestamp; /* входное время; отправка назначает новый timestamp через db_sync */
uint64_t reply_to_ts; /* 0 = не ответ */ uint64_t reply_to_ts; /* 0 = не ответ */
uint64_t reply_to_node_id; /* node_id автора исходного сообщения */ uint64_t reply_to_node_id; /* node_id автора исходного сообщения */
}; };
/* Подписать, сохранить и разослать сообщение. req/data остаются у вызывающего. */
void chat_core_submit_message(struct UTUN_INSTANCE* inst, struct chat_msg_submit* req); void chat_core_submit_message(struct UTUN_INSTANCE* inst, struct chat_msg_submit* req);
/* Отправка медиа-сообщения (вложение): src->dest (copy) или видео (transcode). */ /* Отправка медиа-сообщения (вложение): src->dest (copy) или видео (transcode). */
@ -70,10 +79,18 @@ void chat_core_submit_media_message(struct UTUN_INSTANCE* inst, struct chat_msg_
/* ── DB-операции (оставлены для интроспекции) ── */ /* ── DB-операции (оставлены для интроспекции) ── */
/* Число локальных сообщений; 0 также при отсутствии БД или ошибке запроса. */
uint32_t chat_core_count(struct UTUN_INSTANCE* inst, const char* ch_id); uint32_t chat_core_count(struct UTUN_INSTANCE* inst, const char* ch_id);
/* Хеш записи по позиции с нуля в порядке (timestamp, node_id); hash_out — 32 байта.
* 0 = запрос выполнен (нет записи — нулевой хеш), -1 = ошибка. */
int chat_core_chain_hash_at(struct UTUN_INSTANCE* inst, const char* ch_id, uint32_t pos, uint8_t* hash_out); int chat_core_chain_hash_at(struct UTUN_INSTANCE* inst, const char* ch_id, uint32_t pos, uint8_t* hash_out);
/* Список ID: count:u16, затем [длина:u8, строка]; buf_size >= 2. 0/-1.
* Списки ограничены буфером; out_len — записанные байты, числа в порядке байтов хоста. */
int chat_core_list_channels(struct UTUN_INSTANCE* inst, uint8_t* buf, size_t buf_size, size_t* out_len); int chat_core_list_channels(struct UTUN_INSTANCE* inst, uint8_t* buf, size_t buf_size, size_t* out_len);
/* Список участников: count:u16, затем node_id:u64; buf_size >= 2. 0/-1. */
int chat_core_list_peers(struct UTUN_INSTANCE* inst, const char* ch_id, uint8_t* buf, size_t buf_size, size_t* out_len); int chat_core_list_peers(struct UTUN_INSTANCE* inst, const char* ch_id, uint8_t* buf, size_t buf_size, size_t* out_len);
/* Ключи X25519/Ed25519 (по 32 байта) и список адресов из БД; buf_size >= 65.
* out_len — размер результата; 0 = прочитано, -1 = нет узла или ошибка. */
int chat_core_load_nodeinfo(struct UTUN_INSTANCE* inst, uint64_t node_id, uint8_t* buf, size_t buf_size, size_t* out_len); int chat_core_load_nodeinfo(struct UTUN_INSTANCE* inst, uint64_t node_id, uint8_t* buf, size_t buf_size, size_t* out_len);
/* ── Создание канала (GUI → uasync) ── */ /* ── Создание канала (GUI → uasync) ── */
@ -92,9 +109,12 @@ struct chat_channel_create {
struct chat_create_auto_arg { struct UTUN_INSTANCE* inst; char name[128]; }; struct chat_create_auto_arg { struct UTUN_INSTANCE* inst; char name[128]; };
/* Подготовить локальные таблицы, группу и синхронизацию канала. Членства на других узлах не даёт. */
void chat_core_ensure_channel_ready(struct UTUN_INSTANCE* inst, const char* ch_id); void chat_core_ensure_channel_ready(struct UTUN_INSTANCE* inst, const char* ch_id);
/* Создать канал по готовым ключам и метаданным req; req остаётся у вызывающего. */
void chat_core_create_channel(struct UTUN_INSTANCE* inst, struct chat_channel_create* req); void chat_core_create_channel(struct UTUN_INSTANCE* inst, struct chat_channel_create* req);
void chat_core_create_channel_trampoline(void* arg); void chat_core_create_channel_trampoline(void* arg);
/* Сгенерировать ключи и создать канал, добавив себя владельцем. */
void chat_core_create_channel_auto(struct UTUN_INSTANCE* inst, const char* name); void chat_core_create_channel_auto(struct UTUN_INSTANCE* inst, const char* name);
void chat_core_create_channel_auto_trampoline(void* arg); void chat_core_create_channel_auto_trampoline(void* arg);
@ -113,13 +133,17 @@ void chat_core_delete_channel_trampoline(void* arg);
/* ── Утилиты ── */ /* ── Утилиты ── */
/* Задать локальный ID в контексте чата; ключи и идентичность ядра не меняет. */
void chat_core_set_my_node_id(struct UTUN_INSTANCE* inst, uint64_t node_id); void chat_core_set_my_node_id(struct UTUN_INSTANCE* inst, uint64_t node_id);
/* Сохранить имя узла и обновить собственные member-записи каналов. */
void chat_core_update_my_name(struct UTUN_INSTANCE* inst, const char* name); void chat_core_update_my_name(struct UTUN_INSTANCE* inst, const char* name);
struct update_my_name_arg { struct UTUN_INSTANCE* inst; char name[128]; }; struct update_my_name_arg { struct UTUN_INSTANCE* inst; char name[128]; };
void chat_core_update_my_name_trampoline(void* arg); void chat_core_update_my_name_trampoline(void* arg);
/* Обновить адреса topo_node и собственные member-записи каналов. */
void chat_core_sync_my_addresses(struct UTUN_INSTANCE* inst); void chat_core_sync_my_addresses(struct UTUN_INSTANCE* inst);
/* Единое обновление своего мембера: адреса + bump update_ts + подпись (с адресами) + синк */ /* Обновить своё userinfo во всех каналах: увеличить update_ts, подписать и запустить sync.
* Адреса обновляются отдельно через topo_node; подпись member-блока их не содержит. */
void chat_core_update_my_member(struct UTUN_INSTANCE* inst); void chat_core_update_my_member(struct UTUN_INSTANCE* inst);
/* Сохранение key-value в ui_state (вызывается из uasync-потока) */ /* Сохранение key-value в ui_state (вызывается из uasync-потока) */
@ -203,13 +227,13 @@ void chat_core_collect_member_detail_trampoline(void* arg);
/* Список каналов с метаданными в JSON: [{id,name,owner_id,owner_name,peers,msgs,online}] */ /* Список каналов с метаданными в JSON: [{id,name,owner_id,owner_name,peers,msgs,online}] */
int chat_core_get_channels_json(struct UTUN_INSTANCE* inst, char* buf, size_t buf_size, size_t* out_len); int chat_core_get_channels_json(struct UTUN_INSTANCE* inst, char* buf, size_t buf_size, size_t* out_len);
/* Сообщения канала в JSON: [{id,ts,author_id,author_name,content_type,data}] /* Сообщения канала в JSON: [{id,ts,author_id,author_name,content_type,data,local_attrs}].
* count=0 → все, offset=0 → сначала */ * От новых к старым; count<=0 — до 100 записей, offset — сколько пропустить. */
int chat_core_get_messages_json(struct UTUN_INSTANCE* inst, const char* ch_id, int count, int offset, int chat_core_get_messages_json(struct UTUN_INSTANCE* inst, const char* ch_id, int count, int offset,
char* buf, size_t buf_size, size_t* out_len); char* buf, size_t buf_size, size_t* out_len);
/* Мемберы канала с полным состоянием в JSON: /* Мемберы в JSON: [{node_id,name,online,connected,adm_tags,x25519,ed25519,addrs}].
* [{node_id,name,online,connected,x25519,ed25519,addrs:[{ip,port,proto,rtt}]}] */ * connected означает наличие транспорта в реестре, а не его UP или готовность группы. */
int chat_core_get_members_json(struct UTUN_INSTANCE* inst, const char* ch_id, char* buf, size_t buf_size, size_t* out_len); int chat_core_get_members_json(struct UTUN_INSTANCE* inst, const char* ch_id, char* buf, size_t buf_size, size_t* out_len);
/* Имя узла из таблицы nodes */ /* Имя узла из таблицы nodes */

14
src/chat/chat_join.h

@ -1,7 +1,9 @@
/* /*
* chat_join.h — авторизация входа в канал (join) через connection-узел и инвайтера * chat_join — добавление участника канала через узел подключения C и инвайтера A.
* Модуль регистрирует invite-ключи, пересылает запрос к A и сохраняет подписанного
* участника. Прямой обмен с новым узлом J выполняет chat_sync; Merkle разносит запись.
* *
* Flow (замена старого CHANNEL_INFO_REQ/CHANNEL_JOIN/WELCOME): * Протокол (A и C могут быть одним узлом):
* *
* инвайтер A connection-узел C джойнер J * инвайтер A connection-узел C джойнер J
* │ KEY_REGISTER(key,ch) │ │ * │ KEY_REGISTER(key,ch) │ │
@ -100,13 +102,13 @@ typedef void (*chat_join_registered_fn)(void* arg, int result);
/* Полный набор мембера джойнера (без подписи дерева — её ставит инвайтер). */ /* Полный набор мембера джойнера (без подписи дерева — её ставит инвайтер). */
struct join_member_data { struct join_member_data {
uint64_t node_id; uint64_t node_id; /* derive(x25519), идентичность нового участника */
const uint8_t* x25519; /* 32 */ const uint8_t* x25519; /* 32 */
const uint8_t* ed25519; /* 32 */ const uint8_t* ed25519; /* 32 */
const uint8_t* join_sig; /* 64, может быть NULL */ const uint8_t* join_sig; /* 64, может быть NULL */
uint64_t join_ts; uint64_t join_ts; /* время подписи вступления, Unix seconds */
const uint8_t* update_sig; /* 64, может быть NULL */ const uint8_t* update_sig; /* 64, может быть NULL */
uint64_t update_ts; uint64_t update_ts; /* версия подписанного блока участника */
char userinfo[256]; /* JSON, null-terminated */ char userinfo[256]; /* JSON, null-terminated */
}; };
@ -120,7 +122,7 @@ int chat_join_parse_member(const uint8_t* data, size_t len, struct join_member_
/* Инициализация: etcp_router_bind(ETCP_RT_ID_JOIN). */ /* Инициализация: etcp_router_bind(ETCP_RT_ID_JOIN). */
int chat_join_init(struct UTUN_INSTANCE* inst); int chat_join_init(struct UTUN_INSTANCE* inst);
/* Деинициализация. */ /* Снять обработчик, освободить ключи; ожидающие регистрации получают CANCELLED. */
void chat_join_destroy(struct UTUN_INSTANCE* inst); void chat_join_destroy(struct UTUN_INSTANCE* inst);
/* Локальное приглашение (инвайтер == connection). 0=stored, -1=ошибка. */ /* Локальное приглашение (инвайтер == connection). 0=stored, -1=ошибка. */

37
src/chat/chat_sync.h

@ -1,25 +1,10 @@
/* /*
* chat_sync.h — оркестратор синхронизации чата по ETCP * chat_sync — связь чат-сервиса с транспортом и группами.
* * Отслеживает пиров, запускает member_sync, поддерживает кеш каналов и статус в GUI.
* При поднятии / разрыве ETCP-соединений: * Вход по ссылке держит отдельный NCD handle до результата join, а не только до UP.
* - обновляет online-статус пиров в БД и GUI * JOIN_INFO_REQ → JOIN_INFO_RESP → JOIN_REQUEST → JOIN_READY; правила — в chat_join.h.
* - запускает member_sync (синхронизацию участников каналов) через merkle_sync * Жизненным циклом управляет chat_service_start/stop. Вызовы — в потоке uasync,
* - управляет кешем каналов (периодический refresh из БД, 30s) * кроме явно отмеченной обёртки, которая копирует аргументы и постит работу в этот поток.
*
* Обслуживает протокол входа в канал (join) между джойнером, connection-узлом и инвайтером:
* JOIN_INFO_REQ → JOIN_INFO_RESP → JOIN_REQUEST → (форвард инвайтеру) → JOIN_READY.
* Полное описание flow и дерева приглашений — в chat_join.h.
*
* Обрабатывает invite-ссылки (chat_sync_connect_from_invite):
* - сохраняет pubkey, join_key и адреса invite-узла в БД (nodes, node_addresses)
* - создаёт инфраструктуру канала, использует NCD для прямого подключения
* - после поднятия ETCP-соединения запускает join-протокол
*
* Использование:
* 1. chat_sync_init(inst) — вызывается при старте, биндит ETCP_RT_ID_CHAT_SYNC (0x30)
* 2. chat_sync_connect_from_invite(ch_id, node_id, pubkey, addrs, count) — вход по invite
* 3. chat_sync_join_channel(inst, ch_id, target_node_id) — вход через уже подключённый узел
* 4. chat_sync_destroy(inst) — вызывается при завершении
*/ */
#ifndef CHAT_SYNC_H #ifndef CHAT_SYNC_H
@ -65,15 +50,17 @@ struct UASYNC;
/* ── Public API ── */ /* ── Public API ── */
/* Создать контекст, обработчики join и member_sync. 0/-1; повтор допустим. */
int chat_sync_init(struct UTUN_INSTANCE* inst); int chat_sync_init(struct UTUN_INSTANCE* inst);
/* Отменить ожидания join, синхронизацию и таймеры; не удаляет данные каналов. */
void chat_sync_destroy(struct UTUN_INSTANCE* inst); void chat_sync_destroy(struct UTUN_INSTANCE* inst);
/* Прямое ETCP-подключение к узлу (pubkey+адреса из БД) */ /* Сейчас пустая функция. Для подключения использовать chat_core_connect_node с каналом. */
void chat_sync_connect_node(struct UTUN_INSTANCE* inst, uint64_t node_id); void chat_sync_connect_node(struct UTUN_INSTANCE* inst, uint64_t node_id);
/* Подключиться к пиру по данным invite-ссылки (вызывается из uasync-потока). /* Подключиться к пиру по данным invite-ссылки (вызывается из uasync-потока).
join_key — 64-битный ключ входа из ссылки (регистрируется инвайтером на connection-узле). join_key — 64-битный ключ входа из ссылки (регистрируется инвайтером на connection-узле).
password может быть NULL — тогда отправляется без пароля (совместимость с v1) */ password может быть NULL — тогда передаётся пустой пароль. */
void chat_sync_connect_from_invite(struct UTUN_INSTANCE* inst, uint64_t channel_id, uint64_t node_id, void chat_sync_connect_from_invite(struct UTUN_INSTANCE* inst, uint64_t channel_id, uint64_t node_id,
const uint8_t* pubkey_bin, const uint8_t* pubkey_bin,
const uint8_t* addrs_data, int addr_count, const uint8_t* addrs_data, int addr_count,
@ -82,7 +69,7 @@ void chat_sync_connect_from_invite(struct UTUN_INSTANCE* inst, uint64_t channel_
/* Пригласить target_node_id в канал ch_id с адресами из invite-ссылки (inviter → joiner). /* Пригласить target_node_id в канал ch_id с адресами из invite-ссылки (inviter → joiner).
Принимает pubkey и адреса целевого узла. Если соединения ещё нет — строит TOPO_NODE Принимает pubkey и адреса целевого узла. Если соединения ещё нет — строит TOPO_NODE
с адресами, запускает conn_mgr для подключения, и по conn_up отправляет CS_MSG_CHANNEL_INVITE. с адресами, удерживает NCD для join, и по UP отправляет CS_MSG_CHANNEL_INVITE.
Если соединение уже есть — отправляет CHANNEL_INVITE сразу. Если соединение уже есть — отправляет CHANNEL_INVITE сразу.
Копирует pubkey_bin и addrs_data (вызывающий может освободить после возврата). Копирует pubkey_bin и addrs_data (вызывающий может освободить после возврата).
Вызов из любого потока — внутри постит в uasync. */ Вызов из любого потока — внутри постит в uasync. */
@ -112,7 +99,7 @@ void chat_sync_join_channel(struct UTUN_INSTANCE* inst, const char* ch_id, uint6
Вызывается из auto_socket при появлении/изменении сетевой связности. */ Вызывается из auto_socket при появлении/изменении сетевой связности. */
void chat_sync_retry_channels_on_socket_change(struct UTUN_INSTANCE* inst); void chat_sync_retry_channels_on_socket_change(struct UTUN_INSTANCE* inst);
/* Немедленно убрать канал из кеша g_cs->channels (локальное удаление канала), /* Немедленно убрать канал из кеша chat_sync (локальное удаление канала),
* чтобы cs_on_conn_status не запускал member_sync для удалённого ns до следующего refresh. */ * чтобы cs_on_conn_status не запускал member_sync для удалённого ns до следующего refresh. */
void chat_sync_on_channel_deleted(struct UTUN_INSTANCE* inst, const char* ch_id); void chat_sync_on_channel_deleted(struct UTUN_INSTANCE* inst, const char* ch_id);

65
src/chat/db_sync.h

@ -1,42 +1,10 @@
// db_sync.h — Distributed content-addressed table with SQLite + peer sync via etcp_api (direct P2P) /* db_sync — репликация подписанных записей SQLite по прямым ETCP-соединениям.
// * Каждый DB_SYNC_INSTANCE связывает таблицу с group_id. Различия находятся по
// Назначение: децентрализованная реплицируемая таблица JSON-записей между всеми узлами сети. * цепочке хешей; новые записи рассылаются PUSH. Ключ и порядок записей:
// Модуль поддерживает несколько независимых инстансов (таблиц), каждый идентифицируется парой (name, id). * (timestamp, author_signature). Подпись автора покрывает timestamp || data.
// * Ядро создаёт модуль, чат включает его и регистрирует таблицы каналов.
// Каждая запись ОБЯЗАТЕЛЬНО содержит Ed25519-подпись автора. * Обычно используется общая SQLite ядра; remove отключает таблицу от обмена,
// chain_hash[pos] = SHA256(chain_hash[pos-1] || id || timestamp || author || author_signature[64]) * сохраняя данные. Все вызовы и callbacks выполняются в потоке uasync. */
// Первичный ключ: (timestamp, author) — тот же автор в ту же ms = дубликат.
// Упорядочение: ORDER BY timestamp, author.
//
// Использование:
// 1. Включить в конфиге: db_sync_enabled = 1
// Опционально: db_sync_ttl = 86400 (по умолчанию)
// 2. db_sync_init() — вызывается автоматически при старте utun_instance
// 3. struct DB_SYNC_INSTANCE* si = db_sync_instance_add(inst, "chats", 1);
// Создаёт/регистрирует инстанс. При поднятии соединений автоматически запускает sync.
// 4. uint64_t ts = db_sync_next_timestamp(si);
// sig = Ed25519(ts || my_node_id || json_data)
// db_sync_insert_signed(si, json_data, len, sig, 64, ts) — вставить запись.
// sig=NULL → ошибка.
// 5. db_sync_count(si) — количество записей в локальной БД для этого инстанса
// 6. db_sync_instance_remove(si) — удалить инстанс (таблица БД не удаляется)
// 7. db_sync_destroy() — вызывается автоматически при завершении
//
// Синхронизация:
// - При поднятии ETCP-соединения с пиром для каждого инстанса запускается полная синхронизация
// - Каждое сообщение содержит group_id (uint64, binary), что позволяет маршрутизировать сообщения к нужному инстансу на приёмной стороне
// - Если пир не имеет инстанса с таким хешем — возвращает DB_MSG_ERROR(0x01)
// - Если пир добавляет инстанс позже — сам инициирует sync в нашу сторону
// - Протокол: сравнение chain_hash для поиска расхождений (7 типов сообщений)
// - Новые записи немедленно рассылаются подключённым пирам через PUSH
//
// Wire-формат записи (SEND_DATA/PUSH): [id:8][ts:8][author:8][dlen:4][data][sig_len:1=64][sig:64]
//
// Нюансы:
// - Записи не редактируются и не удаляются явно — только TTL-очистка (per-instance)
// - Дубликаты определяются по (timestamp, author)
// - БД хранится в SQLite, путь: <db_path>/sync
// - Синхронизация идёт через etcp_api (service ID 0x20), только прямые P2P соединения
#ifndef DB_SYNC_H #ifndef DB_SYNC_H
#define DB_SYNC_H #define DB_SYNC_H
@ -90,32 +58,37 @@ struct DB_SYNC_INSTANCE;
// Ed25519 signature size // Ed25519 signature size
#define DB_SIG_SIZE 64 #define DB_SIG_SIZE 64
// Global lifecycle // Создать контекст и учесть db_sync_enabled из конфига. 0/-1.
int db_sync_init(struct UTUN_INSTANCE* inst); int db_sync_init(struct UTUN_INSTANCE* inst);
// Включить ранее созданный модуль независимо от флага конфига. Повтор допустим; 0/-1.
int db_sync_enable(struct UTUN_INSTANCE* inst); int db_sync_enable(struct UTUN_INSTANCE* inst);
// Отменить таймеры/обмен, снять обработчики; общую SQLite не закрывает.
void db_sync_destroy(struct UTUN_INSTANCE* inst); void db_sync_destroy(struct UTUN_INSTANCE* inst);
// Instance management // Зарегистрировать/найти таблицу для group_id. auto_sync запускает начальный обмен; NULL при ошибке.
struct DB_SYNC_INSTANCE* db_sync_instance_add(struct UTUN_INSTANCE* inst, const char* table_name, uint64_t group_id, int auto_sync); struct DB_SYNC_INSTANCE* db_sync_instance_add(struct UTUN_INSTANCE* inst, const char* table_name, uint64_t group_id, int auto_sync);
// Освободить экземпляр синхронизации; сама таблица остаётся в БД.
void db_sync_instance_remove(struct DB_SYNC_INSTANCE* si); void db_sync_instance_remove(struct DB_SYNC_INSTANCE* si);
// Найти живой инстанс по group_id (поиск в текущем массиве db->instances — указатель всегда актуален). // Заимствованный указатель по group_id; действителен до remove/destroy, NULL если нет.
struct DB_SYNC_INSTANCE* db_sync_instance_find(struct UTUN_INSTANCE* inst, uint64_t group_id); struct DB_SYNC_INSTANCE* db_sync_instance_find(struct UTUN_INSTANCE* inst, uint64_t group_id);
// group_id инстанса (для chat-таблиц равен числовому channel_id). Возвращает 0 при si==NULL. // group_id инстанса (для chat-таблиц равен числовому channel_id). Возвращает 0 при si==NULL.
uint64_t db_sync_instance_group_id(struct DB_SYNC_INSTANCE* si); uint64_t db_sync_instance_group_id(struct DB_SYNC_INSTANCE* si);
// Data operations (per-instance) // Data operations (per-instance)
// db_sync_insert_signed: sig MUST be non-NULL, 64 bytes. // db_sync_insert_signed: sig MUST be non-NULL, 64 bytes.
// sig = Ed25519(ts[8] || author[8] || json). // sig = Ed25519(ts[8] || json); ts записан в порядке байтов хоста.
// ts — call db_sync_next_timestamp(si) before signing to reserve monotonically increasing timestamp. // ts — call db_sync_next_timestamp(si) before signing to reserve monotonically increasing timestamp.
int db_sync_insert_signed(struct DB_SYNC_INSTANCE* si, const char* json_data, size_t len, int db_sync_insert_signed(struct DB_SYNC_INSTANCE* si, const char* json_data, size_t len,
const uint8_t* sig, size_t sig_len, uint64_t ts, const uint8_t* sig, size_t sig_len, uint64_t ts,
const char* local_attrs); const char* local_attrs);
uint32_t db_sync_count(struct DB_SYNC_INSTANCE* si); uint32_t db_sync_count(struct DB_SYNC_INSTANCE* si);
// Последний известный/зарезервированный timestamp в миллисекундах Unix; 0 для NULL.
uint64_t db_sync_get_last_timestamp(struct DB_SYNC_INSTANCE* si); uint64_t db_sync_get_last_timestamp(struct DB_SYNC_INSTANCE* si);
// Зарезервировать следующий timestamp перед подписью; растёт даже в пределах одной ms.
uint64_t db_sync_next_timestamp(struct DB_SYNC_INSTANCE* si); uint64_t db_sync_next_timestamp(struct DB_SYNC_INSTANCE* si);
// Select: iterate records ordered by (timestamp, author), starting at offset, max limit (0=unlimited). // Обход по (timestamp, author_signature), начиная с offset; limit=0 — без ограничения.
// Returns number of records passed to callback. // Возвращает число записей или -1. Данные callback заимствованы на время вызова.
typedef void (*db_sync_select_cb)(void* arg, uint64_t id, uint64_t timestamp, typedef void (*db_sync_select_cb)(void* arg, uint64_t id, uint64_t timestamp,
const char* data, size_t data_len, uint64_t author, const char* data, size_t data_len, uint64_t author,
const uint8_t* author_sig, size_t sig_len, const uint8_t* author_sig, size_t sig_len,
@ -124,7 +97,7 @@ int db_sync_select(struct DB_SYNC_INSTANCE* si, uint32_t offset, uint32_t l
db_sync_select_cb cb, void* arg); db_sync_select_cb cb, void* arg);
// Insert callback: fired after local or peer insert succeeds. // Insert callback: fired after local or peer insert succeeds.
// author_node_id = self for local inserts, peer node_id for remote. // author_node_id — исходный автор записи, не обязательно переславший её пир.
// initial_sync = 1 если запись пришла в рамках bulk-обмена SEND_DATA (initial/re-sync), // initial_sync = 1 если запись пришла в рамках bulk-обмена SEND_DATA (initial/re-sync),
// 0 для live-PUSH и локальных вставок. // 0 для live-PUSH и локальных вставок.
typedef void (*db_sync_insert_cb)(struct DB_SYNC_INSTANCE* si, typedef void (*db_sync_insert_cb)(struct DB_SYNC_INSTANCE* si,

55
src/chat/member_sync.h

@ -1,38 +1,11 @@
/* /*
* ── member_sync — синхронизация участников чата ── * member_sync — подписанные записи участников каналов и их обмен через merkle_sync.
* * Запись содержит блок участника, блок владельца и подпись дерева приглашений.
* Тонкая прослойка над merkle_sync, адаптированная под мемберов каналов. * Версии блоков сравниваются отдельно; запись и Merkle-хеш меняются атомарно.
* Данные хранятся в таблицах: nodes, node_addresses, peers_<channel_id>. * Локальный put запускает распространение после commit; start ждёт сверки с пиром.
* Хеш включает все реплицируемые поля, длины строк и целые в big-endian (см. member_sync_doc.md). * Участие в CHAT-группе определяется локальной таблицей peers_<channel_id>.
* * Сетевые адреса принадлежат topo_node и не входят в блок участника.
* ── Использование ── * Модуль живёт внутри чат-сервиса; все операции и callbacks — в потоке uasync.
*
* // разово: инициализация (из chat_sync_init)
* member_sync_init(inst);
*
* // асинхронный запуск синхронизации (из cs_on_conn_up)
* member_sync_start(inst, peer, ch_id, on_done, ch);
*
* // коллбэк: синхронизация завершена
* static void on_done(uint64_t peer, const char* ch_id, int result, void* arg) {
* struct channel_cache* ch = (struct channel_cache*)arg;
* if (result == MT_OK) ch->synced = CS_SYNC_DONE;
* }
*
* // добавление/обновление мембера (из chat_core_create_channel, cs_handle_welcome)
* member_sync_put(inst, ch_id, node_id, x25519, ed25519,
* join_sig, join_ts, update_sig, update_ts, userinfo,
* adm_tags, adm_tags_sig, storage);
* // запись и дерево изменены атомарно; после commit запускается автоматическая синхронизация
*
* // отмена синхронизации (коллбэк НЕ вызывается)
* member_sync_cancel(inst, peer, ch_id);
*
* // количество мемберов в канале
* int n = member_sync_count(inst, ch_id);
*
* // разово: завершение (из chat_sync_destroy)
* member_sync_destroy(inst);
*/ */
#ifndef MEMBER_SYNC_H #ifndef MEMBER_SYNC_H
@ -57,15 +30,17 @@ struct sqlite3;
/* UTF-8 JSON bytes, excluding NUL. A complete record fits in one Merkle page. */ /* UTF-8 JSON bytes, excluding NUL. A complete record fits in one Merkle page. */
#define MS_ADM_TAGS_MAX 8192 #define MS_ADM_TAGS_MAX 8192
/* Байты для подписи владельца: tags без NUL, затем node_id[8] в порядке байтов хоста.
* Возвращает длину или -1, если данные/буфер не подходят. */
int member_sync_build_owner_msg(uint64_t node_id, const char* tags, uint8_t* out, size_t capacity); int member_sync_build_owner_msg(uint64_t node_id, const char* tags, uint8_t* out, size_t capacity);
/* Распарсенный рекорд мембера (два подписанных блока). */ /* Представление записи; все указатели заимствованы на время вызова. */
struct ms_member_rec { struct ms_member_rec {
uint64_t node_id; uint64_t node_id; /* derive(x25519), ключ записи участника */
const uint8_t* x25519; /* 32 байта, обязательно */ const uint8_t* x25519; /* 32 байта, обязательно */
const uint8_t* ed25519; /* 32 байта, обязательно */ const uint8_t* ed25519; /* 32 байта, обязательно */
const uint8_t* join_sig; /* 64 байта, может быть NULL */ const uint8_t* join_sig; /* 64 байта, может быть NULL */
uint64_t join_ts; uint64_t join_ts; /* эпоха вступления, Unix seconds */
const uint8_t* update_sig; /* 64 байта, может быть NULL */ const uint8_t* update_sig; /* 64 байта, может быть NULL */
uint64_t update_ts; /* ver блока мембера */ uint64_t update_ts; /* ver блока мембера */
const char* userinfo; /* JSON {"name":...}, может быть NULL */ const char* userinfo; /* JSON {"name":...}, может быть NULL */
@ -109,8 +84,8 @@ void member_sync_remove_apply_cbk(struct UTUN_INSTANCE* inst, member_apply_cbk_f
/* /*
* Инициализировать модуль: создаёт merkle_sync с коллбэками для мемберов * Инициализировать модуль: создаёт merkle_sync с коллбэками для мемберов
* (ETCP сервис 0x31), регистрирует _on_node_updated для пересчёта дерева * (ETCP_RT_ID_MEMBER_SYNC), проверяет локальные записи, восстанавливает деревья
* при обновлении node info. * и подписывается на события CHAT-групп. 0/-1.
*/ */
int member_sync_init(struct UTUN_INSTANCE* inst); int member_sync_init(struct UTUN_INSTANCE* inst);
@ -222,7 +197,7 @@ int member_sync_count(struct UTUN_INSTANCE* inst, const char* ch_id);
/* /*
* Коллбэк: изменились adm_tags любого узла (включая себя). * Коллбэк: изменились adm_tags любого узла (включая себя).
* channel_id — канал (namespace), в котором произошло изменение. * channel_id — канал (namespace), в котором произошло изменение.
* Вызывается при успешной обработке MSG_ITEM_UPDATE с adm_tags. * Уведомляет об изменении свойств после commit и при обновлении состояния узла в BGP.
* Многоподписочный — можно добавить несколько подписчиков. * Многоподписочный — можно добавить несколько подписчиков.
*/ */
typedef void (*node_props_changed_fn)(uint64_t node_id, const char* adm_tags, const char* channel_id, void* arg); typedef void (*node_props_changed_fn)(uint64_t node_id, const char* adm_tags, const char* channel_id, void* arg);

21
src/chat/merkle_sync.h

@ -1,3 +1,8 @@
/* merkle_sync — обмен различиями между наборами записей по Merkle-дереву.
* Движок сверяет хеши и передаёт страницы изменившихся листьев; модель данных
* через data_ops читает записи, проверяет подписи/права и сохраняет изменения.
* Сессия относится к (пир, namespace), отправка учитывает backpressure.
* Все операции — в потоке uasync; один движок на UTUN_INSTANCE. */
#ifndef MERKLE_SYNC_H #ifndef MERKLE_SYNC_H
#define MERKLE_SYNC_H #define MERKLE_SYNC_H
@ -19,7 +24,7 @@ struct UTUN_INSTANCE;
* SHA256 листа обновляется моделью в порядке unsigned key; родителей считает merkle_tree. * SHA256 листа обновляется моделью в порядке unsigned key; родителей считает merkle_tree.
* Все операции выполняются в uasync-потоке экземпляра. */ * Все операции выполняются в uasync-потоке экземпляра. */
struct merkle_sync_data_ops { struct merkle_sync_data_ops {
merkle_leaf_hash_fn update_bucket_hash; merkle_leaf_hash_fn update_bucket_hash; // хеш содержимого одного листа
/* Одна страница листа: after_valid=0 начинает обход, иначе только key > after. /* Одна страница листа: after_valid=0 начинает обход, иначе только key > after.
* len: ёмкость -> размер; next — последний ключ; more — есть продолжение. * len: ёмкость -> размер; next — последний ключ; more — есть продолжение.
* Пустой результат разрешён только с more=0. Ошибка не кодируется пустой страницей. */ * Пустой результат разрешён только с more=0. Ошибка не кодируется пустой страницей. */
@ -28,12 +33,12 @@ struct merkle_sync_data_ops {
/* Проверить весь формат, применить записи и обновить дерево в одной DB-транзакции. /* Проверить весь формат, применить записи и обновить дерево в одной DB-транзакции.
* <0 — ошибка. Никаких broadcast/send-back из этого callback. */ * <0 — ошибка. Никаких broadcast/send-back из этого callback. */
int (*apply_items)(void* ctx, const char* ns, uint64_t from_peer, const uint8_t* data, size_t len); int (*apply_items)(void* ctx, const char* ns, uint64_t from_peer, const uint8_t* data, size_t len);
int (*validate_peer)(void* ctx, const char* ns, uint64_t peer); int (*validate_peer)(void* ctx, const char* ns, uint64_t peer); // необязательная проверка: только 1 разрешает приём
/* Проверить отложенные зависимости модели после обоих проходов, до подтверждения корня. */ /* Проверить отложенные зависимости модели после обоих проходов, до подтверждения корня. */
int (*finish)(void* ctx, const char* ns); int (*finish)(void* ctx, const char* ns);
/* Optional pull transaction. All three hooks must be supplied together. /* Необязательный отложенный приём: все четыре callback задаются вместе.
* abort_pull releases the context after both success and failure. * stage_page накапливает страницы без изменения опубликованных хешей;
* Staged records must not affect published hashes before commit_pull. */ * commit_pull применяет результат; abort_pull освобождает pull и после успеха. */
int (*begin_pull)(void* ctx, const char* ns, uint64_t peer, uint64_t round, void** pull); int (*begin_pull)(void* ctx, const char* ns, uint64_t peer, uint64_t round, void** pull);
int (*commit_pull)(void* ctx, void* pull); int (*commit_pull)(void* ctx, void* pull);
void (*abort_pull)(void* ctx, void* pull); void (*abort_pull)(void* ctx, void* pull);
@ -41,10 +46,13 @@ struct merkle_sync_data_ops {
uint64_t next, int more, const uint8_t* data, size_t len); uint64_t next, int more, const uint8_t* data, size_t len);
}; };
/* result — MT_OK или MT_ERR_*; ns заимствован на время callback. */
typedef void (*merkle_sync_done_cb)(uint64_t peer, const char* ns, int result, void* arg); typedef void (*merkle_sync_done_cb)(uint64_t peer, const char* ns, int result, void* arg);
/* Зарегистрировать сервис; нужны ua и общая SQLite. ops/data_ctx живут до destroy. 0/-1. */
int merkle_sync_init(struct UTUN_INSTANCE* inst, uint8_t svc_id, int merkle_sync_init(struct UTUN_INSTANCE* inst, uint8_t svc_id,
const struct merkle_sync_data_ops* ops, void* data_ctx); const struct merkle_sync_data_ops* ops, void* data_ctx);
/* Отменить сессии без done callbacks, снять обработчики и освободить движок. */
void merkle_sync_destroy(struct UTUN_INSTANCE* inst); void merkle_sync_destroy(struct UTUN_INSTANCE* inst);
/* Подключить namespace к текущему ETCP-соединению и запросить контрольную точку. /* Подключить namespace к текущему ETCP-соединению и запросить контрольную точку.
@ -56,6 +64,7 @@ void merkle_sync_destroy(struct UTUN_INSTANCE* inst);
* Callback может отменить сессию; уничтожение всего экземпляра должно быть отложенным. */ * Callback может отменить сессию; уничтожение всего экземпляра должно быть отложенным. */
int merkle_sync_start(struct UTUN_INSTANCE* inst, uint64_t peer, const char* ns, merkle_sync_done_cb cb, void* arg); int merkle_sync_start(struct UTUN_INSTANCE* inst, uint64_t peer, const char* ns, merkle_sync_done_cb cb, void* arg);
void merkle_sync_cancel(struct UTUN_INSTANCE* inst, uint64_t peer, const char* ns); void merkle_sync_cancel(struct UTUN_INSTANCE* inst, uint64_t peer, const char* ns);
/* Отменить все сессии namespace / пира без callbacks. */
void merkle_sync_cancel_ns(struct UTUN_INSTANCE* inst, const char* ns); void merkle_sync_cancel_ns(struct UTUN_INSTANCE* inst, const char* ns);
void merkle_sync_cancel_peer(struct UTUN_INSTANCE* inst, uint64_t peer); void merkle_sync_cancel_peer(struct UTUN_INSTANCE* inst, uint64_t peer);
@ -64,7 +73,9 @@ void merkle_sync_cancel_peer(struct UTUN_INSTANCE* inst, uint64_t peer);
* Без внешней транзакции пересчёт сам создаёт savepoint и уведомляет после commit. * Без внешней транзакции пересчёт сам создаёт savepoint и уведомляет после commit.
* Возвращает 1=изменилось, 0=no-op, <0=ошибка. */ * Возвращает 1=изменилось, 0=no-op, <0=ошибка. */
int merkle_sync_recompute_path(struct UTUN_INSTANCE* inst, const char* ns, uint64_t key); int merkle_sync_recompute_path(struct UTUN_INSTANCE* inst, const char* ns, uint64_t key);
/* После commit уведомить сессии namespace о локальном изменении дерева. */
void merkle_sync_changed(struct UTUN_INSTANCE* inst, const char* ns); void merkle_sync_changed(struct UTUN_INSTANCE* inst, const char* ns);
/* Прочитать хеш узла дерева в hash. 0=успех, <0=ошибка. */
int merkle_sync_read_hash(struct UTUN_INSTANCE* inst, const char* ns, uint8_t level, uint64_t prefix, uint8_t hash[MT_HASH_SIZE]); int merkle_sync_read_hash(struct UTUN_INSTANCE* inst, const char* ns, uint8_t level, uint64_t prefix, uint8_t hash[MT_HASH_SIZE]);
#endif #endif

15
src/chat/merkle_tree.h

@ -1,3 +1,7 @@
/* merkle_tree — производный индекс хешей в SQLite для сравнения записей модели.
* В каждом namespace дерево SHA256: корень (уровень 0), 5 уровней по 32 ветви.
* Лист выбирается по старшим 25 битам ключа; его содержимое хеширует модель.
* Пустое поддерево имеет нулевой хеш. БД принадлежит вызывающему. */
#ifndef MERKLE_TREE_H #ifndef MERKLE_TREE_H
#define MERKLE_TREE_H #define MERKLE_TREE_H
@ -9,16 +13,23 @@
#define MT_BUCKETS 32 #define MT_BUCKETS 32
#define MT_HASH_SIZE 32 #define MT_HASH_SIZE 32
/* Добавить записи листа в уже начатый hash через EVP_DigestUpdate; не завершать hash.
* Возвращает число записей (0 = пустой лист), <0 при ошибке. */
typedef int (*merkle_leaf_hash_fn)(void* ctx, const char* ns, uint8_t level, uint64_t prefix, EVP_MD_CTX* hash); typedef int (*merkle_leaf_hash_fn)(void* ctx, const char* ns, uint8_t level, uint64_t prefix, EVP_MD_CTX* hash);
/* Дерево — производный индекс. При открытии очищается и восстанавливается из записей модели. */ /* Пересоздать пустой индекс; вызывающий затем восстанавливает его из модели. 0/-1. */
int merkle_tree_init(sqlite3* db); int merkle_tree_init(sqlite3* db);
/* Прочитать хеш уровня 0..5; отсутствующий узел даёт нулевой хеш. 0/-1. */
int merkle_tree_get(sqlite3* db, const char* ns, uint8_t level, uint64_t prefix, uint8_t hash[MT_HASH_SIZE]); int merkle_tree_get(sqlite3* db, const char* ns, uint8_t level, uint64_t prefix, uint8_t hash[MT_HASH_SIZE]);
/* Прочитать 32 дочерних хеша узла уровня 0..4; отсутствующие — нули. 0/-1. */
int merkle_tree_children(sqlite3* db, const char* ns, uint8_t level, uint64_t prefix, int merkle_tree_children(sqlite3* db, const char* ns, uint8_t level, uint64_t prefix,
uint8_t hashes[MT_BUCKETS][MT_HASH_SIZE]); uint8_t hashes[MT_BUCKETS][MT_HASH_SIZE]);
/* Работает внутри транзакции вызывающей стороны. Сначала лист, затем родители до корня. */ /* Пересчитать лист ключа и родителей в транзакции вызывающего.
* 1 = индекс изменился, 0 = прежний хеш, -1 = ошибка (транзакцию нужно откатить). */
int merkle_tree_update(sqlite3* db, const char* ns, uint64_t key, merkle_leaf_hash_fn leaf_hash, void* ctx); int merkle_tree_update(sqlite3* db, const char* ns, uint64_t key, merkle_leaf_hash_fn leaf_hash, void* ctx);
/* Оставить старшие level*5 бит ключа, остальные обнулить; уровень 0 — корень. */
uint64_t merkle_sync_level_prefix(uint64_t key, uint8_t level); uint64_t merkle_sync_level_prefix(uint64_t key, uint8_t level);
/* Число байтов для старших level*5 бит; 0 для корня или недопустимого уровня. */
uint8_t merkle_sync_prefix_bytes(uint8_t level); uint8_t merkle_sync_prefix_bytes(uint8_t level);
#endif #endif

18
src/media_async/media_async.h

@ -1,4 +1,7 @@
// media_async.h — асинхронный движок thread-per-task + примитивы (хеш, подпись, копия файла) /* media_async — владелец фоновых задач обработки файлов и криптографии.
* Каждая work выполняется в отдельном pthread; done — в потоке uasync.
* create/submit/destroy вызываются из этого потока. destroy ждёт workers и
* завершает ожидающие done с CANCELLED, пока данные вызывающего ещё существуют. */
#ifndef MEDIA_ASYNC_H #ifndef MEDIA_ASYNC_H
#define MEDIA_ASYNC_H #define MEDIA_ASYNC_H
@ -13,24 +16,29 @@ struct UASYNC;
struct media_async; struct media_async;
typedef void (*ma_work_fn)(void* data); typedef void (*ma_work_fn)(void* data); // рабочий поток; не обращаться к состоянию uasync
typedef void (*ma_done_fn)(void* arg, int err); typedef void (*ma_done_fn)(void* arg, int err); // 0=work завершён, -1=ошибка запуска/таймера, -2=отмена
/* Создать владельца задач; NULL при ошибке выделения памяти. */
struct media_async* media_async_create(void); struct media_async* media_async_create(void);
#define MEDIA_ASYNC_CANCELLED (-2) #define MEDIA_ASYNC_CANCELLED (-2)
/* Owning event-loop thread only, outside done callbacks. Joins workers, then calls /* Вызывать вне done. Ждёт работы, вызывает ожидающие done и освобождает ma; NULL допустим. */
* each pending done with CANCELLED while caller-owned state is still alive. */
void media_async_destroy(struct media_async* ma); void media_async_destroy(struct media_async* ma);
/* Передать задачу; data/arg остаются у вызывающего до done. При ошибке done
* может быть вызван синхронно. Ошибку самой work передают через data, а не err. */
void media_async_submit(struct media_async* ma, struct UASYNC* ua, void media_async_submit(struct media_async* ma, struct UASYNC* ua,
ma_work_fn work, void* data, ma_work_fn work, void* data,
ma_done_fn done, void* arg); ma_done_fn done, void* arg);
/* Синхронные примитивы для worker: SHA256 файла, Ed25519-подпись, копирование. 0/-1. */
int ma_sha256_file(const char* path, uint8_t hash_out[32]); int ma_sha256_file(const char* path, uint8_t hash_out[32]);
int ma_sign_block(const uint8_t* ed25519_privkey, int ma_sign_block(const uint8_t* ed25519_privkey,
const uint8_t* data, size_t len, uint8_t sig_out[64]); const uint8_t* data, size_t len, uint8_t sig_out[64]);
int ma_copy_file(const char* src, const char* dst); int ma_copy_file(const char* src, const char* dst);
/* Размер файла в байтах (int; для файлов, помещающихся в этот тип), -1 при ошибке stat. */
int ma_file_size(const char* path); int ma_file_size(const char* path);
/* Заполнить 16 случайных байт идентификатора. */
void ma_uuid(uint8_t uuid_out[16]); void ma_uuid(uint8_t uuid_out[16]);
/* Рекурсивный обход каталога: для каждого regular-файла вызывается cb(full_path, size). /* Рекурсивный обход каталога: для каждого regular-файла вызывается cb(full_path, size).

60
src/routing_layer/conn_mgr.h

@ -1,3 +1,8 @@
/* conn_mgr — соединение с узлом группы напрямую, в обратную сторону или через посредника.
* Один менеджер принадлежит одной группе; open выдаёт отдельный handle владельцу.
* DIRECT использует NCD, REVERSE просит пира подключиться, INDIRECT выбирает посредника.
* Используется сервисами для доставки данных. Групповые сессии открывают через
* topo_group_peer_open(). Все операции — в потоке uasync; handle закрывает вызывающий. */
#ifndef CONN_MGR_H #ifndef CONN_MGR_H
#define CONN_MGR_H #define CONN_MGR_H
@ -15,44 +20,9 @@ struct UTUN_INSTANCE;
struct ETCP_CONN; struct ETCP_CONN;
struct TOPO_GROUP; struct TOPO_GROUP;
/**
* @file conn_mgr.h
* @brief Подключение к узлам P2P-сети. Handle-based API.
*
* Подключает к удалегному узлу либо напрямую, либо через посредника с прямыми адресами.
* При открытии выдаёт handle. и держит подключение открытым пока не закроешь handle.
* Отправлять можно обычным send через роутер - он сам разберётся как отправить,
* либо send этого моделя - он точно отправит через этот мост.
*
* Один CONN_MGR = одна TOPO_GROUP. Создаётся conn_mgr_init(group), живёт пока
* существует группа. Внутри: кеш entries (по одной на целевой узел), список
* handle'ов (по одному на каждый вызов open), фоновые ping/idle/мониторинг.
*
* === Единая точка входа: conn_mgr_open ===
*
* 3-фазное подключение (DIRECT→REVERSE→INDIRECT).
* Коллбэк: CONN_EVENT_UP (успех) или CONN_EVENT_TIMEOUT (провал).
*
* === 3 фазы ===
*
* 1. DIRECT: node_conn_direct на все адреса, NAT-фильтрация, LAN broadcast ping.
* 2. REVERSE: у нас прямой IP — шлём DIRECT_REQ цели через BGP.
* 3. INDIRECT: через посредника с минимальной суммой RTT (etcp_router).
*
* === Фоновые процессы ===
*
* bg_ping: раз в ~100ms один узел группы за тик. Поддерживает свежие RTT.
* candidate_ping: раз в ~2с, удаляет протухших кандидатов.
*
* === Refcounting и закрытие ===
*
* Несколько handle на один ETCP_CONN. conn_mgr_close удаляет ОДИН handle.
* Последний handle: DIRECT/REVERSE — DISCONNECT + NCD CLOSE.
*/
enum conn_mgr_event { enum conn_mgr_event {
CONN_EVENT_UP = 0, /* связь есть. handle жив, можно отправлять данные */ CONN_EVENT_UP = 0, /* связь есть. handle жив, можно отправлять данные */
CONN_EVENT_DOWN = 1, /* временный обрыв. handle жив, само восстановится. НЕ закрывать */ CONN_EVENT_DOWN = 1, /* связь потеряна; handle остаётся у владельца, успех восстановления не гарантирован */
CONN_EVENT_TIMEOUT = 2, /* подключение не удалось. handle жив, закройте сами */ CONN_EVENT_TIMEOUT = 2, /* подключение не удалось. handle жив, закройте сами */
}; };
@ -69,19 +39,20 @@ typedef void (*conn_mgr_cb_t)(struct CONN_MGR_HANDLE* h,
/* ═══════ lifecycle ═══════ */ /* ═══════ lifecycle ═══════ */
/* Создать менеджер группы и его фоновые проверки; NULL при ошибке. */
struct CONN_MGR* conn_mgr_init(struct TOPO_GROUP* group); struct CONN_MGR* conn_mgr_init(struct TOPO_GROUP* group);
/* Остановить проверки и освободить менеджер. Все клиентские handles закрыть ДО destroy. */
void conn_mgr_destroy(struct CONN_MGR* mgr); void conn_mgr_destroy(struct CONN_MGR* mgr);
/* Приём управляющего пакета от etcp_router; забирает entry. */
void conn_mgr_router_recv_handler(struct ETCP_CONN* conn, struct ll_entry* entry); void conn_mgr_router_recv_handler(struct ETCP_CONN* conn, struct ll_entry* entry);
/* ═══════ единая точка входа ═══════ */ /* ═══════ единая точка входа ═══════ */
/** /**
* Единственная публичная функция подключения — 3-фазное (DIRECT→REVERSE→INDIRECT). * Единственная публичная функция подключения — 3-фазное (DIRECT→REVERSE→INDIRECT).
* Используется медиа-доставкой. Пиры/группы подключаются через node_conn_direct, * Возврат 0 — запрос принят, -1 — ошибка запуска. cb обязателен.
* invite — через topo_group_invite. * Ошибка до запуска может возвращаться без callback. В callback h может быть NULL.
* * Успешно выданный handle закрывает вызывающий, в том числе после TIMEOUT.
* Ошибка/таймаут: CONN_EVENT_TIMEOUT.
* node_id и group_id ВСЕГДА валидны в коллбэке (даже при TIMEOUT).
*/ */
int conn_mgr_open(struct UTUN_INSTANCE* inst, int conn_mgr_open(struct UTUN_INSTANCE* inst,
uint64_t group_id, uint64_t group_id,
@ -89,15 +60,16 @@ int conn_mgr_open(struct UTUN_INSTANCE* inst,
conn_mgr_cb_t cb, void* cb_arg, conn_mgr_cb_t cb, void* cb_arg,
struct CONN_MGR_HANDLE** out_handle); struct CONN_MGR_HANDLE** out_handle);
/** Закрыть handle. Последний handle рвёт соединение. */ /** Закрыть свой handle. Последний освобождает ресурсы CM; чужие NCD-владельцы сохраняют транспорт. */
void conn_mgr_close(struct CONN_MGR_HANDLE* h); void conn_mgr_close(struct CONN_MGR_HANDLE* h);
/* ═══════ доступ к соединению ═══════ */ /* ═══════ доступ к соединению ═══════ */
/** ETCP_CONN для DIRECT/REVERSE, NULL для INDIRECT. */ /** Заимствованный ETCP_CONN от NCD, если он есть; NULL для INDIRECT или отсоединённого handle. */
struct ETCP_CONN* conn_mgr_get_conn(struct CONN_MGR_HANDLE* h); struct ETCP_CONN* conn_mgr_get_conn(struct CONN_MGR_HANDLE* h);
/** Отправить данные (для INDIRECT — через посредников). Забирает владение e. */ /** Передать e->dgram через маршрутизатор CM. Освобождает entry при ненулевом h;
* буфер dgram отдельно не освобождает. При h==NULL entry остаётся у вызывающего. */
int conn_mgr_send(struct CONN_MGR_HANDLE* h, struct ll_entry* e); int conn_mgr_send(struct CONN_MGR_HANDLE* h, struct ll_entry* e);
/* ═══════ мониторинг (вызывается извне) ═══════ */ /* ═══════ мониторинг (вызывается извне) ═══════ */

33
src/routing_layer/etcp_router.h

@ -2,6 +2,9 @@
// Маршрутизирует сервисные пакеты до целевой ноды, на промежуточных нодах ретранслирует // Маршрутизирует сервисные пакеты до целевой ноды, на промежуточных нодах ретранслирует
// Упрощённый TCP поверх ETCP: восстановление порядка, дедупликация, ретрансмиты, inflight-контроль. // Упрощённый TCP поверх ETCP: восстановление порядка, дедупликация, ретрансмиты, inflight-контроль.
// Подпись/шифрование вынесены в автономный модуль route_crypto (encode/decode). // Подпись/шифрование вынесены в автономный модуль route_crypto (encode/decode).
// Логический канал задаётся (group_id, node_id, svc_id), физический транспорт берётся у топологии.
// Сервисы регистрируют приём через bind, отправляют через route_send и соблюдают backpressure.
// Все операции и callbacks — в потоке uasync. Закрытие группы удаляет её каналы и ожидающий транзит.
// //
// Формат SVC_ROUTE пакета: // Формат SVC_ROUTE пакета:
// [cmd:1] [group_id:8] [dst_node_id:8] [src_node_id:8] [seq:4] [svc_id:1] [flags:1] [timestamp:2] [reset_id:8] // [cmd:1] [group_id:8] [dst_node_id:8] [src_node_id:8] [seq:4] [svc_id:1] [flags:1] [timestamp:2] [reset_id:8]
@ -77,8 +80,8 @@ struct TRANSIT_QUEUE {
uint64_t group_id; // = data[0..7] uint64_t group_id; // = data[0..7]
uint64_t src_node_id; // = data[8..15] uint64_t src_node_id; // = data[8..15]
uint64_t dst_node_id; // = data[16..23] uint64_t dst_node_id; // = data[16..23]
struct ll_queue* q; // FIFO транзитных пакетов (ll_entry с dgram), начиная с data[16] struct ll_queue* q; // FIFO ещё не переданных транспорту пакетов
struct queue_waiter_handle waiter; // backpressure waiter на bp-очередь (normalizer->input или tx_queue) struct queue_waiter_handle waiter; // ожидание свободной conn->send_input_q
struct ETCP_CONN* conn; // next_hop (для etcp_send в drain_cb) struct ETCP_CONN* conn; // next_hop (для etcp_send в drain_cb)
}; };
@ -133,9 +136,9 @@ struct ETCP_ROUTER_CONN {
uint64_t reset_id; // собственный идентификатор; не меняется при рестарте пира uint64_t reset_id; // собственный идентификатор; не меняется при рестарте пира
uint64_t peer_reset_id; // подтверждённый идентификатор пира (0 = handshake не завершён) uint64_t peer_reset_id; // подтверждённый идентификатор пира (0 = handshake не завершён)
uint64_t pending_peer_id; // кандидат, ещё не имеющий права менять состояние сессии uint64_t pending_peer_id; // кандидат, ещё не имеющий права менять состояние сессии
uint64_t pending_challenge; uint64_t pending_challenge; // случайный запрос подтверждения кандидату
uint64_t challenge_sent_tb; uint64_t challenge_sent_tb; // время отправки запроса, 0.1 ms
void* handshake_timer; void* handshake_timer; // таймер повторов согласования
uint8_t start_sent; // получен первый ACK данных uint8_t start_sent; // получен первый ACK данных
uint8_t peer_sync_done; // 1 = завершено подтверждение peer ID uint8_t peer_sync_done; // 1 = завершено подтверждение peer ID
@ -203,41 +206,48 @@ struct ETCP_ROUTER_BINDINGS {
// Инициализация: etcp_bind(ETCP_ID_SVC_ROUTE) + очистка bindings + создание router_conns // Инициализация: etcp_bind(ETCP_ID_SVC_ROUTE) + очистка bindings + создание router_conns
int etcp_router_init(struct UTUN_INSTANCE* inst); int etcp_router_init(struct UTUN_INSTANCE* inst);
// Деинициализация // Закрыть каналы, снять обработчик и освободить маршрутизатор.
void etcp_router_destroy(struct UTUN_INSTANCE* inst); void etcp_router_destroy(struct UTUN_INSTANCE* inst);
// Зарегистрировать обработчик сервиса // Зарегистрировать обработчик сервиса
int etcp_router_bind(struct UTUN_INSTANCE* inst, uint8_t svc_id, etcp_recv_fn callback); int etcp_router_bind(struct UTUN_INSTANCE* inst, uint8_t svc_id, etcp_recv_fn callback);
// Снять обработчик; само по себе не закрывает логические каналы сервиса.
int etcp_router_unbind(struct UTUN_INSTANCE* inst, uint8_t svc_id); int etcp_router_unbind(struct UTUN_INSTANCE* inst, uint8_t svc_id);
// Отправить сервисный пакет (авто-conn, seq, inflight-контроль через send_q). // Отправить сервисный пакет (авто-conn, seq, inflight-контроль через send_q).
// Принимает владение entry во всех исходах, включая ошибку. // Принимает владение entry во всех исходах, включая ошибку.
// entry->dgram = svc_id[1] || payload; 0=принято в очередь, <0=ошибка, не подтверждение доставки.
// force разрешает превысить лимит числа пакетов в send_q.
// mode — битовая маска ROUTE_CRYPTO_SIGN / ROUTE_CRYPTO_ENCRYPT (см. route_crypto.h), 0 = обычный пакет. // mode — битовая маска ROUTE_CRYPTO_SIGN / ROUTE_CRYPTO_ENCRYPT (см. route_crypto.h), 0 = обычный пакет.
int etcp_route_send(struct UTUN_INSTANCE* inst, uint64_t group_id, uint64_t dst_node_id, struct ll_entry* entry, int force, int mode); int etcp_route_send(struct UTUN_INSTANCE* inst, uint64_t group_id, uint64_t dst_node_id, struct ll_entry* entry, int force, int mode);
// Найти/создать состояние seq-подключения по (group_id, remote_node_id, svc_id) // Найти/создать логический канал. Указатель заимствован до close/destroy, NULL при ошибке.
struct ETCP_ROUTER_CONN* etcp_router_conn_get(struct UTUN_INSTANCE* inst, struct ETCP_ROUTER_CONN* etcp_router_conn_get(struct UTUN_INSTANCE* inst,
uint64_t group_id, uint64_t remote_node_id, uint8_t svc_id); uint64_t group_id, uint64_t remote_node_id, uint8_t svc_id);
// Есть ли сейчас физический маршрут до узла (прямой, indirect-посредник или глобальный). // Есть ли сейчас физический маршрут до узла (прямой, indirect-посредник или глобальный).
// 1 = есть ETCP-соединение для отправки, 0 = нет (передача встанет в no_route). // 1 = найден транспорт (для self всегда 1); это не гарантия UP/готовности логического канала.
int etcp_router_has_route(struct UTUN_INSTANCE* inst, uint64_t group_id, uint64_t dst_node_id); int etcp_router_has_route(struct UTUN_INSTANCE* inst, uint64_t group_id, uint64_t dst_node_id);
// Отправить данные с авто-seq и контролем inflight // Отправить данные с авто-seq и контролем inflight
// data: payload без svc_id, mode: ROUTE_CRYPTO_SIGN / ROUTE_CRYPTO_ENCRYPT (0 = обычный) // data: payload без svc_id, mode: ROUTE_CRYPTO_SIGN / ROUTE_CRYPTO_ENCRYPT (0 = обычный)
// Копирует data; 0=принято в очередь, <0=ошибка. Доставка подтверждается отдельно протоколом.
int etcp_router_conn_send(struct ETCP_ROUTER_CONN* rconn, int etcp_router_conn_send(struct ETCP_ROUTER_CONN* rconn,
const uint8_t* data, size_t len, int mode); const uint8_t* data, size_t len, int mode);
// Закрыть seq-подключение // Закрыть канал и уведомить сервис. После вызова rconn использовать нельзя.
void etcp_router_conn_close(struct ETCP_ROUTER_CONN* rconn); void etcp_router_conn_close(struct ETCP_ROUTER_CONN* rconn);
// Установить рабочий max_inflight, пересчитывает inflight_limit и state machine // Установить рабочий max_inflight, пересчитывает inflight_limit и state machine
void etcp_router_set_max_inflight(struct ETCP_ROUTER_CONN* rconn, uint16_t new_max); void etcp_router_set_max_inflight(struct ETCP_ROUTER_CONN* rconn, uint16_t new_max);
// Backpressure: зарегистрировать/отменить waiter на send_q очереди // Дождаться порога send_q. h должен быть обнулён и жить до callback/отмены.
// Callback может выполниться внутри вызова. Повторная регистрация отменяет прежнюю.
// Если канала нет, ожидание не регистрируется. Перед освобождением arg отменить waiter.
void etcp_router_on_send_ready(struct UTUN_INSTANCE* inst, uint64_t group_id, uint64_t node_id, uint8_t svc_id, void etcp_router_on_send_ready(struct UTUN_INSTANCE* inst, uint64_t group_id, uint64_t node_id, uint8_t svc_id,
struct queue_waiter_handle* h, struct queue_waiter_handle* h,
queue_threshold_callback_fn callback, void* arg); queue_threshold_callback_fn callback, void* arg);
// Отмена относится к исходной очереди, даже если логический канал уже перезапущен.
void etcp_router_cancel_send_ready(struct UTUN_INSTANCE* inst, uint64_t group_id, uint64_t node_id, uint8_t svc_id, void etcp_router_cancel_send_ready(struct UTUN_INSTANCE* inst, uint64_t group_id, uint64_t node_id, uint8_t svc_id,
struct queue_waiter_handle* h); struct queue_waiter_handle* h);
@ -252,8 +262,9 @@ int etcp_router_send_q_has_room(struct UTUN_INSTANCE* inst, uint64_t group_id, u
// Используется при etcp_conn_reinit чтобы избежать гонки с очисткой ETCP-очередей // Используется при etcp_conn_reinit чтобы избежать гонки с очисткой ETCP-очередей
void etcp_router_pause_retrans_for_node(struct UTUN_INSTANCE* inst, uint64_t remote_node_id); void etcp_router_pause_retrans_for_node(struct UTUN_INSTANCE* inst, uint64_t remote_node_id);
// Закрыть все router_conn для указанного (group_id, remote_node_id) (peer умер) // Закрыть все каналы и ожидающий транзит группы. Сначала сервисы отменяют свои producers/waiters.
void etcp_router_close_group(struct UTUN_INSTANCE* inst, uint64_t group_id); void etcp_router_close_group(struct UTUN_INSTANCE* inst, uint64_t group_id);
// Закрыть логические каналы одного пира в группе.
void etcp_router_conn_close_all_for_node(struct UTUN_INSTANCE* inst, uint64_t group_id, uint64_t remote_node_id); void etcp_router_conn_close_all_for_node(struct UTUN_INSTANCE* inst, uint64_t group_id, uint64_t remote_node_id);
// Начать новую локальную сессию для конкретного peer+svc в группе и уведомить сервис. // Начать новую локальную сессию для конкретного peer+svc в группе и уведомить сервис.

92
src/routing_layer/topo_group.h

@ -1,32 +1,13 @@
/** /**
* @file topo_group.h * @file topo_group.h
* @brief BGP-подобный обмен топологией узлов между пирами через ETCP. * @brief Групповые сессии и обмен маршрутами поверх общего транспорта NCD.
* *
* Идея: модуль автоматически выстраивает карту маршрутизации между узлами, * У каждой группы своя топология; идентичности и адреса узлов общие (topo_node.h).
* используя только пассивное наблюдение (подписка на события) * UTUN обменивается также подсетями, CHAT допускает только локально известных участников.
* Можно иметь много групп узлов. Группа - это список узлов. * Для подключения используйте peer_open/close: группа владеет NCD с начала попытки.
* Для каждой группы выстраивается своя независимая таблица маршрутизации. * JOIN/ACCEPT согласует сессию, NODEINFO/WITHDRAW поддерживают её маршруты.
* Один узел может входить в любое число групп. * Потерю путей обрабатывает topo_recovery, выбор CHAT-пиров — topo_group_connect.
* Все операции — в потоке uasync. Сервисы создают и удаляют свои группы явно.
* Механика:
* - Узлы обмениваются информацией друг о друге (pubkey, адреса, подсети)
* через NODEINFO/WITHDRAW сообщения, маршрутизируемые по ETCP.
* - Каждый узел хранит полную таблицу известных узлов (структура TOPO_NODEQ).
* - При подключении нового пира — полная синхронизация таблицы.
* - При отключении пира — withdraw всех узлов, достижимых только через него.
*
* Группы (TOPO_GROUP):
* - TOPO_GROUP_TYPE_UTUN — основная VPN-сеть (обмен маршрутами, подсетями)
* - TOPO_GROUP_TYPE_CHAT — чат-группа (только узлы, без подсетей)
* - Каждая группа изолирована: узлы из utun-группы не видны в чат-группе и наоборот
*
* Дополнительные функции:
* - NAT-детекция: выделена в отдельный модуль nat_detection.h
* - Поиск оптимального маршрута до узла (topo_group_find_conn_for_node)
*
* Хранение: все узлы сохраняются в SQLite (таблицы nodes, node_addresses,
* channels, peers_<channel_id>). База открывается в topo_groups_init()
* если в конфиге задан db_path.
*/ */
#ifndef TOPO_GROUP_H #ifndef TOPO_GROUP_H
#define TOPO_GROUP_H #define TOPO_GROUP_H
@ -170,19 +151,19 @@ struct TOPOMSG_ERR_GROUP_MISMATCH {
} __attribute__((packed)); } __attribute__((packed));
struct TOPO_GROUP_CONN_ITEM { struct TOPO_GROUP_CONN_ITEM {
struct TOPO_GROUP* group; struct TOPO_GROUP* group; // владелец сессии
uint64_t node_id; uint64_t node_id; // идентичность пира, включая CONNECTING
struct ETCP_CONN* conn; // NULL до UP, в состоянии CONNECTING struct ETCP_CONN* conn; // NULL до UP, в состоянии CONNECTING
struct NODE_CONN_DIRECT* handle; // владение BGP этим conn (NCD-handle) struct NODE_CONN_DIRECT* handle; // владение BGP этим conn (NCD-handle)
struct TOPO_PEER_REQUEST* requests; struct TOPO_PEER_REQUEST* requests; // запросы инициаторов; каждый закрывает свой
uint64_t progress; // время последнего прогресса обмена (timebase) uint64_t progress; // время последнего прогресса обмена (timebase)
uint8_t retained; // участие запрошено постоянным владельцем или достигло READY uint8_t retained; // участие запрошено постоянным владельцем или достигло READY
uint64_t local_epoch, peer_epoch; // поколения сторон текущего присоединения uint64_t local_epoch, peer_epoch; // поколения сторон текущего присоединения
uint8_t accepted; // JOIN_ACCEPT подтвердил local_epoch uint8_t accepted; // JOIN_ACCEPT подтвердил local_epoch
uint8_t tx_failed, transport_down; uint8_t tx_failed, transport_down; // ошибка отправки / потеря транспорта
struct topo_tx_item *tx_head, *tx_tail; struct topo_tx_item *tx_head, *tx_tail; // FIFO команд, ещё не переданных транспорту
struct queue_waiter_handle tx_waiter; struct queue_waiter_handle tx_waiter; // ожидание свободной send_input_q
void* tx_wake; void* tx_wake; // отложенный запуск отправки
uint8_t table_received; // TABLE_COMPLETE текущей сессии принят uint8_t table_received; // TABLE_COMPLETE текущей сессии принят
uint8_t table_sent; // ответ на запрос пира целиком передан транспорту uint8_t table_sent; // ответ на запрос пира целиком передан транспорту
}; };
@ -206,7 +187,7 @@ struct TOPO_GROUP {
struct ll_queue* senders_list; // ll_entry.data = TOPO_GROUP_CONN_ITEM; включает CONNECTING с conn=NULL struct ll_queue* senders_list; // ll_entry.data = TOPO_GROUP_CONN_ITEM; включает CONNECTING с conn=NULL
struct ll_queue* nodes; // TOPO_GROUP_NODE{ll_entry,node_id,paths,...} — per-group узлы, хеш-индекс по node_id(8B) struct ll_queue* nodes; // TOPO_GROUP_NODE{ll_entry,node_id,paths,...} — per-group узлы, хеш-индекс по node_id(8B)
struct TOPO_GROUP_NODE* local_node; // свой узел в этой группе struct TOPO_GROUP_NODE* local_node; // свой узел в этой группе
uint8_t ed25519_public_key[SC_PUBKEY_SIZE]; uint8_t ed25519_public_key[SC_PUBKEY_SIZE]; // публичный ключ локального узла
char channel_id[64]; // channel_id для групп типа CHAT char channel_id[64]; // channel_id для групп типа CHAT
struct CONN_MGR* conn_mgr; // менеджер соединений для этой группы struct CONN_MGR* conn_mgr; // менеджер соединений для этой группы
struct TOPO_RECOVERY_CTX* recovery; // один последовательный recovery на группу struct TOPO_RECOVERY_CTX* recovery; // один последовательный recovery на группу
@ -215,7 +196,7 @@ struct TOPO_GROUP {
struct radio_ctx* radio; // radio (PTT walkie-talkie) per-group context struct radio_ctx* radio; // radio (PTT walkie-talkie) per-group context
uint8_t radio_active; // 1 = я слушаю рацию этого канала (TOPO_FLAG_RADIO в NODEINFO) uint8_t radio_active; // 1 = я слушаю рацию этого канала (TOPO_FLAG_RADIO в NODEINFO)
struct topo_node_cbk_entry* node_cbks; // цепочка подписчиков на события узлов struct topo_node_cbk_entry* node_cbks; // цепочка подписчиков на события узлов
uint8_t stopping; uint8_t stopping; // teardown начат; новые попытки запрещены
}; };
/** /**
@ -235,7 +216,8 @@ struct TOPO_GROUPS {
}; };
/** /**
* @brief Создаёт контейнер групп и группу utun по умолчанию (group_id=1). * @brief Создаёт реестр и локальную идентичность; сервисных групп пока нет.
* Открывает общую SQLite: db_path/chats.db, без db_path — :memory:.
* *
* @param instance экземпляр utun с node_id и конфигом * @param instance экземпляр utun с node_id и конфигом
* @return TOPO_GROUPS или NULL при ошибке * @return TOPO_GROUPS или NULL при ошибке
@ -257,12 +239,12 @@ void topo_group_on_activity_change(struct UTUN_INSTANCE* instance, int active, v
void topo_groups_destroy(struct UTUN_INSTANCE* instance); void topo_groups_destroy(struct UTUN_INSTANCE* instance);
/** /**
* @brief Возвращает группу по умолчанию (group_id=TOPO_GROUP_UTUN=0x8000000000000000). * @brief Возвращает UTUN-группу или NULL, если сервис не запущен. Указатель заимствован.
*/ */
struct TOPO_GROUP* topo_groups_get_default(struct TOPO_GROUPS* g); struct TOPO_GROUP* topo_groups_get_default(struct TOPO_GROUPS* g);
/** /**
* @brief Находит группу по id. * @brief Заимствованный указатель на группу по id; NULL, если её нет.
*/ */
struct TOPO_GROUP* topo_groups_find(struct TOPO_GROUPS* g, uint64_t group_id); struct TOPO_GROUP* topo_groups_find(struct TOPO_GROUPS* g, uint64_t group_id);
@ -287,9 +269,9 @@ struct TOPO_GROUP* topo_groups_create_group(struct TOPO_GROUPS* g, uint64_t grou
void topo_groups_remove_group(struct TOPO_GROUPS* g, uint64_t group_id); void topo_groups_remove_group(struct TOPO_GROUPS* g, uint64_t group_id);
/** /**
* @brief Добавляет conn в senders_list (если нет), отправляет запрос таблицы (nodeinfo). * @brief Присоединяет существующий транспорт к группе и начинает JOIN/ACCEPT.
* * Группа приобретает собственный NCD handle; повтор не сбрасывает сессию. 0/-1.
* Вызывается при ETCP on_up. * Для обычного инициатора используйте отменяемый topo_group_peer_open().
*/ */
int topo_group_new_conn(struct TOPO_GROUP* group, struct ETCP_CONN* conn); int topo_group_new_conn(struct TOPO_GROUP* group, struct ETCP_CONN* conn);
@ -311,19 +293,20 @@ enum topo_peer_phase { TOPO_PEER_FAILED, TOPO_PEER_CONNECTING, TOPO_PEER_SYNCING
* CLOSED/leave/удаление группы отсоединяет запрос окончательно. * CLOSED/leave/удаление группы отсоединяет запрос окончательно.
* Запрос остаётся валиден после удаления группы (phase=FAILED), caller обязан * Запрос остаётся валиден после удаления группы (phase=FAILED), caller обязан
* закрыть его. Все операции выполняются в потоке единственного uasync. * закрыть его. Все операции выполняются в потоке единственного uasync.
* progress — timebase последнего принятого NODEINFO/шага обмена, не polling time. */ * progress — timebase последнего принятого NODEINFO/шага обмена, не polling time.
* open: 0 = запрос создан в *out, -1 = ошибка. */
int topo_group_peer_open(struct TOPO_GROUP* group, uint64_t node_id, struct TOPO_PEER_REQUEST** out); int topo_group_peer_open(struct TOPO_GROUP* group, uint64_t node_id, struct TOPO_PEER_REQUEST** out);
/* Освободить свой запрос; после вызова request недействителен. NULL допустим. */
void topo_group_peer_close(struct TOPO_PEER_REQUEST* request); void topo_group_peer_close(struct TOPO_PEER_REQUEST* request);
struct ETCP_CONN* topo_group_peer_conn(const struct TOPO_PEER_REQUEST* request); /* borrowed, including CONNECTING */ struct ETCP_CONN* topo_group_peer_conn(const struct TOPO_PEER_REQUEST* request); /* borrowed, including CONNECTING */
/* Завершить всю сессию этой пары: LEAVE, удаление путей, отсоединение всех запросов. */
void topo_group_peer_leave(struct TOPO_GROUP* group, uint64_t node_id); void topo_group_peer_leave(struct TOPO_GROUP* group, uint64_t node_id);
/* Текущая фаза; NULL или отсоединённый запрос дают FAILED. */
enum topo_peer_phase topo_group_peer_phase(const struct TOPO_PEER_REQUEST* request); enum topo_peer_phase topo_group_peer_phase(const struct TOPO_PEER_REQUEST* request);
/* Время последнего прогресса в единицах 0.1 ms; 0, если сессии нет. */
uint64_t topo_group_peer_progress(const struct TOPO_PEER_REQUEST* request); uint64_t topo_group_peer_progress(const struct TOPO_PEER_REQUEST* request);
/** /* Причина завершения групповой сессии определяет необходимость recovery. */
* @brief Удаляет conn из senders_list, очищает paths во всех nodes, отправляет withdraw если node unreachable.
*
* Вызывается при ETCP on_down.
*/
enum topo_group_remove_reason { enum topo_group_remove_reason {
TOPO_REMOVE_TRANSPORT_DOWN, TOPO_REMOVE_TRANSPORT_DOWN,
TOPO_REMOVE_LOCAL_LEAVE, TOPO_REMOVE_LOCAL_LEAVE,
@ -331,25 +314,29 @@ enum topo_group_remove_reason {
TOPO_REMOVE_MEMBER_INVALID, TOPO_REMOVE_MEMBER_INVALID,
TOPO_REMOVE_REJECTED TOPO_REMOVE_REJECTED
}; };
/* Только потеря транспорта запускает recovery каскадно потерянных маршрутов. */ /* Сбросить сессию и удалить пути через conn, разослать WITHDRAW для потерянных узлов.
* При DOWN живые запросы сохраняют NCD для повторного UP. Только потеря транспорта
* запускает recovery каскадно потерянных маршрутов. */
void topo_group_remove_conn(struct TOPO_GROUP* group, struct ETCP_CONN* conn, enum topo_group_remove_reason reason); void topo_group_remove_conn(struct TOPO_GROUP* group, struct ETCP_CONN* conn, enum topo_group_remove_reason reason);
/** /**
* @brief Обрабатывает пакет NODEINFO. * @brief Обрабатывает пакет NODEINFO.
* *
* Проверка версии, обновление или создание TOPO_NODEQ, добавление пути, * Проверка подписи и timestamp, обновление TOPO_GROUP_NODE, добавление пути,
* вставка в роутинг, broadcast если не max hops. * вставка в роутинг, broadcast если не max hops.
* *
* @return 0 при успехе * @return 0 при успехе
*/ */
int topo_group_process_nodeinfo(struct TOPO_GROUP* group, struct ETCP_CONN* from, const uint8_t* data, size_t len); int topo_group_process_nodeinfo(struct TOPO_GROUP* group, struct ETCP_CONN* from, const uint8_t* data, size_t len);
/* Присоединить UP-транспорт; переданный handle остаётся у вызывающего. 0/-1. */
int topo_group_join_peer(struct TOPO_GROUP* group, struct NODE_CONN_DIRECT* handle); int topo_group_join_peer(struct TOPO_GROUP* group, struct NODE_CONN_DIRECT* handle);
/* Отсоединить пира и запланировать LEAVE. Чужой handle не забирает. 0/-1. */
int topo_group_leave_peer(struct TOPO_GROUP* group, struct NODE_CONN_DIRECT* handle); int topo_group_leave_peer(struct TOPO_GROUP* group, struct NODE_CONN_DIRECT* handle);
/** /**
* @brief Обрабатывает WITHDRAW. * @brief Обрабатывает WITHDRAW.
* *
* Удаляет node из роутинга и nodes, broadcast withdraw. * Удаляет затронутые пути; узел удаляется только без оставшихся путей.
* *
* @return 0 при успехе * @return 0 при успехе
*/ */
@ -369,11 +356,12 @@ void topo_group_send_withdraw(struct TOPO_GROUP* group, uint64_t node_id);
/** /**
* @brief Поиск оптимального ETCP соединения для указанного node_id. * @brief Поиск оптимального ETCP соединения для указанного node_id.
* *
* Перебирает paths узла, выбирает путь с минимальным hop_count. * Выбирает минимум hop_count среди живых путей. Если живых нет, возвращает
* лучший сохранённый путь: ненулевой результат сам по себе не гарантирует UP.
* *
* @param group указатель на TOPO_GROUP * @param group указатель на TOPO_GROUP
* @param node_id целевой узел * @param node_id целевой узел
* @return оптимальный ETCP_CONN* или NULL * @return заимствованный ETCP_CONN* или NULL
*/ */
struct ETCP_CONN* topo_group_find_conn_for_node(struct TOPO_GROUP* group, uint64_t node_id); struct ETCP_CONN* topo_group_find_conn_for_node(struct TOPO_GROUP* group, uint64_t node_id);

24
src/routing_layer/topo_group_connect.h

@ -1,11 +1,10 @@
/* CHAT connection policy over cancellable group requests. /* topo_group_connect — выбор пиров и попытки подключения CHAT-группы.
* Phase 1: historical connected peers in parallel (up to 128). * Сначала ранее подключённые пиры (до 128 параллельно), затем по одному
* Phase 2: sequential supernode/public peers. Phase 3: local peers. * суперузлы/публичные узлы, затем локальные. Цель — READY суперузел или три
* Each node is tried once per cycle. CONNECTING: 2s; SYNCING: 5s without progress. * READY немобильных пира. CONNECTING: 2 s; SYNCING: 5 s без прогресса.
* Goal: one READY supernode or three READY nonmobile peers. After exhaustion, * После исчерпания кандидатов новый цикл через 1 s или следующий standby burst.
* retry after 1s (or standby burst). Existing READY group sessions survive * Планировщик владеет запросами, группа — NCD. Остановка отменяет попытки,
* planner restart/stop. Manual attempts are tracked and cancelled with the group. * включая ручные, но сохраняет готовые сессии. Все вызовы — в потоке uasync.
* active_count is derived from the current group sessions, not a separate counter.
*/ */
#ifndef TOPO_GROUP_CONNECT_H #ifndef TOPO_GROUP_CONNECT_H
#define TOPO_GROUP_CONNECT_H #define TOPO_GROUP_CONNECT_H
@ -15,15 +14,22 @@
struct TOPO_GROUP; struct TOPO_GROUP;
struct ETCP_CONN; struct ETCP_CONN;
/* Запустить новый автоматический цикл вместо прежнего. 0/-1. */
int topo_group_connect_init(struct TOPO_GROUP* group); int topo_group_connect_init(struct TOPO_GROUP* group);
/* Отменить запросы, таймеры и отложенные вызовы планировщика. */
void topo_group_connect_destroy(struct TOPO_GROUP* group); void topo_group_connect_destroy(struct TOPO_GROUP* group);
/* Запланировать проверку изменившихся сессий. */
void topo_group_connect_changed(struct TOPO_GROUP* group); void topo_group_connect_changed(struct TOPO_GROUP* group);
/* Учесть появление пира; для READY записать connected в БД. */
void topo_group_connect_on_up(struct TOPO_GROUP* group, struct ETCP_CONN* conn); void topo_group_connect_on_up(struct TOPO_GROUP* group, struct ETCP_CONN* conn);
/* Учесть потерю пира и продолжить подбор кандидатов. */
void topo_group_connect_on_down(struct TOPO_GROUP* group, struct ETCP_CONN* conn); void topo_group_connect_on_down(struct TOPO_GROUP* group, struct ETCP_CONN* conn);
/* Число текущих READY-пиров группы. */
int topo_group_connect_active_count(struct TOPO_GROUP* group); int topo_group_connect_active_count(struct TOPO_GROUP* group);
/* Перезапустить автоматический цикл, только если нет READY-пиров. */
void topo_group_connect_restart(struct TOPO_GROUP* group); void topo_group_connect_restart(struct TOPO_GROUP* group);
/* One manual request; shares the group session, no NCD ownership handoff on UP. */ /* Одна ручная попытка, общий запрос к сессии группы. 0=принято, -1=ошибка. */
int topo_group_connect_node_once(struct TOPO_GROUP* group, uint64_t node_id); int topo_group_connect_node_once(struct TOPO_GROUP* group, uint64_t node_id);
#endif #endif

37
src/routing_layer/topo_group_invite.h

@ -1,23 +1,11 @@
/** /**
* @file topo_group_invite.h * @file topo_group_invite.h
* @brief Invite/join к каналу через прямое (ncd) подключение. * @brief Получение описания канала по прямому NCD-транспорту (INVITE_INFO).
* *
* Модуль содержит всю логику работы с invite-ссылками: прямое подключение к * Запрашивает и проверяет подписанные сведения канала, готовит локальную группу.
* приглашающему узлу через node_conn_direct (ncd) и проверку членства * Успешный ответ не добавляет участника на других узлах и не означает READY группы.
* (INVITE_INFO_REQ/RESP) с бутстрапом криптографической идентичности канала. * Вход по пользовательской ссылке и добавление участника описаны в chat_sync.h/chat_join.h.
* * Все операции — в потоке uasync; временные запросы отменяются при остановке чата.
* Поток джойнера (topo_group_invite_join):
* 1. Создаёт/находит группу, сохраняет TOPO_NODE приглашающего.
* 2. Открывает прямое соединение (ncd) и шлёт INVITE_INFO_REQ.
* 3. Получает INVITE_INFO_RESP — проверяет ch_sig, сохраняет канал в БД,
* поднимает инфраструктуру группы (chat_core_ensure_channel_ready).
* 4. Доставляет TGI_EVENT_JOIN.
*
* Инвайтера сторона: tgi_handle_info_req отвечает на INVITE_INFO_REQ —
* проверяет членство в БД и возвращает полную информацию о канале.
*
* conn_mgr при этом используется ТОЛЬКО для медиа — invite больше от него
* не зависит.
*/ */
#ifndef TOPO_GROUP_INVITE_H #ifndef TOPO_GROUP_INVITE_H
#define TOPO_GROUP_INVITE_H #define TOPO_GROUP_INVITE_H
@ -37,8 +25,8 @@ struct NODE_CONN_DIRECT;
#define TGI_SUBCMD_INFO_RESP 0x02 #define TGI_SUBCMD_INFO_RESP 0x02
/* события коллбэка джойнера */ /* события коллбэка джойнера */
#define TGI_EVENT_JOIN 0 #define TGI_EVENT_JOIN 0 /* описание канала принято, локальная инфраструктура готова */
#define TGI_EVENT_TIMEOUT 1 #define TGI_EVENT_TIMEOUT 1 /* попытка завершилась ошибкой или таймаутом */
typedef void (*tgi_cb_t)(uint64_t node_id, uint64_t group_id, int event, void* arg); typedef void (*tgi_cb_t)(uint64_t node_id, uint64_t group_id, int event, void* arg);
@ -48,17 +36,18 @@ void topo_group_invite_init(struct UTUN_INSTANCE* inst);
/** /**
* Джойнер: прямое подключение (ncd) к приглашающему + проверка членства. * Джойнер: прямое подключение (ncd) к приглашающему + проверка членства.
* ni — TOPO_NODE приглашающего (pubkey + адреса), не владеем. * ni — TOPO_NODE приглашающего (pubkey + адреса), не владеем.
* При успехе вызывается cb(TGI_EVENT_JOIN), при провале cb(TGI_EVENT_TIMEOUT). * 0 — попытка начата, -1 — ошибка. До создания попытки ошибка не вызывает cb;
* после создания ошибки дают TGI_EVENT_TIMEOUT, иногда до возврата функции.
* Успешный обмен даёт TGI_EVENT_JOIN; cancel_all завершает ожидание без cb.
*/ */
int topo_group_invite_join(struct UTUN_INSTANCE* inst, uint64_t group_id, int topo_group_invite_join(struct UTUN_INSTANCE* inst, uint64_t group_id,
struct TOPO_NODE* ni, uint64_t node_id, struct TOPO_NODE* ni, uint64_t node_id,
tgi_cb_t cb, void* cb_arg); tgi_cb_t cb, void* cb_arg);
/** /**
* Инвайтера сторона: прямое подключение (ncd) к целевому узлу для последующей * Запустить групповую попытку к уже известному участнику, при необходимости
* отправки CHANNEL_INVITE (через chat_sync, по ETCP_CONN_STATUS_UP). * передав адреса ni (заимствованы; NULL — взять известные). 0/-1.
* Без INVITE_INFO-проверки членства. ni может быть NULL — тогда адреса грузятся * Применяется обычная проверка CHAT-членства. Сам CHANNEL_INVITE не отправляет.
* из БД. Группа владеет NCD; инициатор держит отменяемый групповой запрос.
*/ */
int topo_group_invite_to_channel(struct UTUN_INSTANCE* inst, uint64_t group_id, int topo_group_invite_to_channel(struct UTUN_INSTANCE* inst, uint64_t group_id,
struct TOPO_NODE* ni, uint64_t node_id); struct TOPO_NODE* ni, uint64_t node_id);

72
src/routing_layer/topo_node.h

@ -2,21 +2,11 @@
* @file topo_node.h * @file topo_node.h
* @brief Модель данных узла сети — структуры, сериализация, wire-формат. * @brief Модель данных узла сети — структуры, сериализация, wire-формат.
* *
* Этот модуль определяет: * TOPO_NODE — общие ключи, имя и адреса; реестр хранит указатель с подсчётом ссылок.
* - Как узел представлен в памяти (TOPO_NODE, TOPO_GROUP_NODE, адреса, подсети) * Подписанная запись заменяется целиком при более свежем timestamp, вместе с подписью.
* - Как узел передаётся по сети (packed-структуры TOPOMSG_*) * TOPO_GROUP_NODE — пути, подсети и состояние узла в конкретной группе.
* - Как узел сериализуется/десериализуется (topo_node_serialize/deserialize) * TOPOMSG_NODE — сетевой формат. Модуль сериализует, подписывает и проверяет записи.
* * Работа с реестром и группами выполняется в потоке uasync экземпляра.
* Хранение:
* - Идентичность узла (pubkeys, адреса, имя) — в глобальном реестре TOPO_GROUPS->node_registry
* Единственный экземпляр TOPO_NODE, разделяется между группами через group_ref_count.
* - Per-group данные (node_id, paths, subnets, connectivity) — в TOPO_GROUP->nodes (очередь TOPO_GROUP_NODE)
*
* Основные понятия:
* TOPO_NODE — глобальная идентичность узла. Встроен в node_registry как queue entry.
* TOPO_GROUP_NODE — узел в контексте группы: node_id + paths, hop_list, связность
* TOPO_NODEPATH — путь до узла через конкретное ETCP-соединение
* TOPOMSG_NODE — wire-формат, передаваемый по NODEINFO
*/ */
#ifndef TOPO_NODE_H #ifndef TOPO_NODE_H
#define TOPO_NODE_H #define TOPO_NODE_H
@ -87,9 +77,11 @@ struct UTUN_INSTANCE;
#define TOPO_GROUP_TYPE_UTUN 1 #define TOPO_GROUP_TYPE_UTUN 1
#define TOPO_GROUP_TYPE_CHAT 2 #define TOPO_GROUP_TYPE_CHAT 2
/* Результаты проб адресов. *_time — timebase (0.1 ms), *_rtt — те же единицы;
* *_status — PROBE_RESULT_*, probe_status — PROBE_STATUS_*. */
struct TOPO_CONNECTIVITY { struct TOPO_CONNECTIVITY {
uint8_t probe_status; uint8_t probe_status;
uint8_t pending_count; uint8_t pending_count; // незавершённые пробы
uint64_t probe_start_time; uint64_t probe_start_time;
uint16_t interface_min_rtt; uint16_t interface_min_rtt;
uint16_t nat_min_rtt; uint16_t nat_min_rtt;
@ -103,7 +95,7 @@ struct TOPO_CONNECTIVITY {
uint64_t real_probe_time; uint64_t real_probe_time;
uint64_t ping_req_time; uint64_t ping_req_time;
uint64_t last_ping_time; uint64_t last_ping_time;
void* probe_list; void* probe_list; // принадлежащие узлу контексты текущих проб
void* probe_deferred_wait; /* контекст отложенной пробы (standby_wait, Android) */ void* probe_deferred_wait; /* контекст отложенной пробы (standby_wait, Android) */
}; };
@ -171,17 +163,17 @@ struct TOPO_ADDR_REALITY_OPTS {
struct _topo_head { struct _topo_head* next; }; struct _topo_head { struct _topo_head* next; };
static inline int topo_list_count(struct _topo_head* head) { int n = 0; while (head) { n++; head = head->next; } return n; } static inline int topo_list_count(struct _topo_head* head) { int n = 0; while (head) { n++; head = head->next; } return n; }
/** Глобальная идентичность узла. Встроен в TOPO_GROUPS->node_registry как queue entry. */ /** Общая запись узла; node_registry хранит указатель на неё в отдельном ll_entry. */
struct TOPO_NODE { struct TOPO_NODE {
struct ll_entry ll; struct ll_entry ll;
uint32_t group_ref_count; uint32_t group_ref_count; // ссылки групп, ядра и временных пользователей
uint64_t node_id; uint8_t ver; uint64_t node_id; uint8_t ver;
uint64_t timestamp; // 0 = неподписанные bootstrap-сведения uint64_t timestamp; // 0 = неподписанные bootstrap-сведения
uint8_t public_key[SC_PUBKEY_SIZE], ed25519_public_key[SC_PUBKEY_SIZE]; uint8_t public_key[SC_PUBKEY_SIZE], ed25519_public_key[SC_PUBKEY_SIZE];
uint8_t x25519_self_sig[64]; uint8_t x25519_self_sig[64]; // Ed25519-подпись общей записи
uint8_t client_type; // CLIENT_TYPE_SERVER/DESKTOP/MOBILE uint8_t client_type; // CLIENT_TYPE_SERVER/DESKTOP/MOBILE
uint8_t client_activity; // CLIENT_ACTIVITY_STANDBY/ACTIVE uint8_t client_activity; // CLIENT_ACTIVITY_STANDBY/ACTIVE
char* node_name; char* node_name; // собственная строка; освобождается с записью
struct TOPO_SOCKMETA4* v4_sock_meta; struct TOPO_SOCKMETA4* v4_sock_meta;
struct TOPO_ADDR4* v4_addrs; struct TOPO_ADDR4* v4_addrs;
struct TOPO_SOCKMETA6* v6_sock_meta; struct TOPO_SOCKMETA6* v6_sock_meta;
@ -196,9 +188,9 @@ struct TOPO_NODESUBNETS {
struct TOPO_NODEPATH { struct TOPO_NODEPATH {
struct ll_entry ll; struct ll_entry ll;
struct ETCP_CONN* conn; struct ETCP_CONN* conn; // заимствованный next-hop; транспорт удерживает сессия
uint8_t hop_count; uint8_t hop_count; // число node_id в массиве сразу после структуры
uint16_t cumulative_rtt; uint16_t cumulative_rtt; // RTT оставшейся цепочки, без локального линка (0.1 ms)
}; };
/** Узел в контексте группы (per-group данные). Хранится в TOPO_GROUP->nodes (ll_queue). */ /** Узел в контексте группы (per-group данные). Хранится в TOPO_GROUP->nodes (ll_queue). */
@ -213,52 +205,74 @@ struct TOPO_GROUP_NODE {
uint64_t conn_mgr_intermediaries[CONN_MGR_MAX_INTERMEDIARIES]; // посредники для INDIRECT uint64_t conn_mgr_intermediaries[CONN_MGR_MAX_INTERMEDIARIES]; // посредники для INDIRECT
uint8_t conn_mgr_intermediariy_count; // количество посредников uint8_t conn_mgr_intermediariy_count; // количество посредников
struct TOPO_CONNECTIVITY connectivity; // результаты ping-проб (interface/nat/real) struct TOPO_CONNECTIVITY connectivity; // результаты ping-проб (interface/nat/real)
uint8_t conn_presence; // NODE_CONN_* — какие подключения есть в принципе uint8_t conn_presence; // маска NCONN_*: известные виды подключения
uint8_t conn_up; // NODE_CONN_* — какие из них подняты uint8_t conn_up; // маска NCONN_*: работающие виды подключения
uint8_t radio; // 1 = узел слушает рацию канала (TOPO_FLAG_RADIO) uint8_t radio; // 1 = узел слушает рацию канала (TOPO_FLAG_RADIO)
struct NODE_CONN_DIRECT* handle; // активный ncd-handle прямого соединения к узлу struct NODE_CONN_DIRECT* handle; // не владеет транспортом; NCD принадлежит сессии группы
}; };
// API — глобальный реестр TOPO_NODE // API — глобальный реестр TOPO_NODE
/* Заимствованный указатель; ref нужен, если запись должна пережить удаление из группы. */
struct TOPO_NODE* topo_node_registry_find(struct TOPO_GROUPS* groups, uint64_t node_id); struct TOPO_NODE* topo_node_registry_find(struct TOPO_GROUPS* groups, uint64_t node_id);
/* Принять ni и вернуть запись с одной новой ссылкой. Может вернуть существующий объект.
* При NULL входной ni остаётся у вызывающего. Подписанные записи сравниваются по timestamp. */
struct TOPO_NODE* topo_node_registry_store(struct TOPO_GROUPS* groups, struct TOPO_NODE* ni); struct TOPO_NODE* topo_node_registry_store(struct TOPO_GROUPS* groups, struct TOPO_NODE* ni);
/* Добавить/снять ссылку; последняя unref удаляет запись из реестра и освобождает её. */
void topo_node_registry_ref(struct TOPO_GROUPS* groups, uint64_t node_id); void topo_node_registry_ref(struct TOPO_GROUPS* groups, uint64_t node_id);
void topo_node_registry_unref(struct TOPO_GROUPS* groups, uint64_t node_id); void topo_node_registry_unref(struct TOPO_GROUPS* groups, uint64_t node_id);
/* Освободить отдельную запись и её списки; зарегистрированную запись отпускать через unref. */
void topo_node_destroy(struct TOPO_GROUPS* groups, struct TOPO_NODE* ni); void topo_node_destroy(struct TOPO_GROUPS* groups, struct TOPO_NODE* ni);
// API — per-group // API — per-group
/* Заимствованный узел группы, NULL если не найден. */
struct TOPO_GROUP_NODE* topo_node_find_by_id(struct TOPO_GROUP* group, uint64_t node_id); struct TOPO_GROUP_NODE* topo_node_find_by_id(struct TOPO_GROUP* group, uint64_t node_id);
/* Отпустить ссылку реестра и подсети; nq и paths не освобождает. */
void topo_nodeq_free_group_fields(struct TOPO_GROUPS* groups, struct TOPO_GROUP_NODE* nq); void topo_nodeq_free_group_fields(struct TOPO_GROUPS* groups, struct TOPO_GROUP_NODE* nq);
/* Полное удаление узла из группы: отменяет влётные пробы, освобождает paths/subnets, /* Полное удаление узла из группы: отменяет влётные пробы, освобождает paths/subnets,
* убирает из group->nodes и освобождает сам nq. Единственный корректный способ удаления узла. */ * убирает из group->nodes и освобождает сам nq. Единственный корректный способ удаления узла. */
void topo_nodeq_remove_node(struct TOPO_GROUP* group, struct TOPO_GROUP_NODE* nq); void topo_nodeq_remove_node(struct TOPO_GROUP* group, struct TOPO_GROUP_NODE* nq);
/* Размер переменной части NODEINFO по счётчикам заголовка. */
int topo_node_dyn_size(const struct TOPOMSG_NODE* msg); int topo_node_dyn_size(const struct TOPOMSG_NODE* msg);
/* Собрать NODEINFO в out; длина или -1. Входные структуры остаются у вызывающего. */
int topo_node_serialize(struct TOPO_NODE* ni, struct TOPO_GROUP_NODE* nq, int topo_node_serialize(struct TOPO_NODE* ni, struct TOPO_GROUP_NODE* nq,
uint64_t group_id, uint8_t flags, uint64_t group_id, uint8_t flags,
uint8_t* out, size_t out_max, uint16_t cumulative_rtt); uint8_t* out, size_t out_max, uint16_t cumulative_rtt);
/* Разобрать NODEINFO; 0/-1. Выходные объекты выделены отдельно и передаются вызывающему. */
int topo_node_deserialize(struct TOPO_GROUP* group, const uint8_t* data, size_t len, struct TOPO_NODE** out_ni, struct TOPO_NODESUBNETS** out_subnets, int topo_node_deserialize(struct TOPO_GROUP* group, const uint8_t* data, size_t len, struct TOPO_NODE** out_ni, struct TOPO_NODESUBNETS** out_subnets,
uint64_t** out_hop_list, uint8_t* out_hop_count, uint64_t** out_hop_list, uint8_t* out_hop_count,
uint16_t* out_cumulative_rtt); uint16_t* out_cumulative_rtt);
/* Заимствованный список hops: кратчайший живой путь, иначе первый сохранённый.
* out_count обязателен; NULL означает отсутствие пути. */
uint64_t* topo_node_best_hop_list(struct TOPO_GROUP_NODE* nq, uint8_t* out_count, uint16_t* out_rtt); uint64_t* topo_node_best_hop_list(struct TOPO_GROUP_NODE* nq, uint8_t* out_count, uint16_t* out_rtt);
/* Обновить локальный узел группы и его анонс. 0/-1. */
int topo_group_update_my_nodeinfo(struct UTUN_INSTANCE* instance, struct TOPO_GROUP* group); int topo_group_update_my_nodeinfo(struct UTUN_INSTANCE* instance, struct TOPO_GROUP* group);
/* Создать подписанную локальную идентичность с отдельной ссылкой ядра. 0/-1. */
int topo_node_init_self(struct UTUN_INSTANCE* instance); int topo_node_init_self(struct UTUN_INSTANCE* instance);
int topo_node_update_my_addresses(struct UTUN_INSTANCE* instance); /* 1=changed, 0=unchanged, -1=error */ int topo_node_update_my_addresses(struct UTUN_INSTANCE* instance); /* 1=changed, 0=unchanged, -1=error */
/* Обработчик изменений локального сокета: обновляет адреса узла. */
void topo_node_on_socket_changed(struct ETCP_SOCKET* sock, int event, void* arg); void topo_node_on_socket_changed(struct ETCP_SOCKET* sock, int event, void* arg);
/* Вывести диагностику узлов группы в лог / текстовый буфер. */
void topo_node_dump_all(struct TOPO_GROUP* group); void topo_node_dump_all(struct TOPO_GROUP* group);
int topo_node_format_all(struct TOPO_GROUP* group, char* buf, size_t buf_size); int topo_node_format_all(struct TOPO_GROUP* group, char* buf, size_t buf_size);
/* 1, если хотя бы в одной группе пора повторить ping этого узла; иначе 0. */
int topo_node_ping_request_cbk(struct TOPO_GROUPS* groups, uint64_t node_id); int topo_node_ping_request_cbk(struct TOPO_GROUPS* groups, uint64_t node_id);
/* Обновить время ping и минимальные RTT во всех группах и БД; rtt в единицах 0.1 ms. */
void topo_node_ping_update_rtt(struct TOPO_GROUPS* groups, uint64_t node_id, uint16_t rtt); void topo_node_ping_update_rtt(struct TOPO_GROUPS* groups, uint64_t node_id, uint16_t rtt);
/* Минимальный суммарный RTT живого пути в единицах 0.1 ms; UINT16_MAX, если неизвестен. */
uint16_t topo_get_chain_rtt(struct TOPO_GROUP_NODE* nq); uint16_t topo_get_chain_rtt(struct TOPO_GROUP_NODE* nq);
/** Build canonical message for Ed25519 signature: x25519_pubkey || name || client_type || client_activity || addresses */ /* Лимиты буфера подписываемых данных и полного сетевого NODEINFO. */
#define TOPO_SIG_MSG_MAX_SIZE 2048 #define TOPO_SIG_MSG_MAX_SIZE 2048
#define TOPO_NODE_WIRE_MAX_SIZE 16384 #define TOPO_NODE_WIRE_MAX_SIZE 16384
/* Канонические байты подписи: node_id, ver, timestamp, X25519-ключ, имя,
* тип/активность и полные списки сокетов/адресов/Reality. Длина или -1. */
int topo_node_build_sig_msg(struct TOPO_NODE* ni, uint8_t* buf, size_t buf_size); int topo_node_build_sig_msg(struct TOPO_NODE* ni, uint8_t* buf, size_t buf_size);
/** Sign self NODEINFO with Ed25519: build_sig_msg + sc_ed25519_sign → ni->x25519_self_sig */ /* Увеличить timestamp, подписать локальный NODEINFO и сохранить снимок в SQLite. 0/-1. */
int topo_node_sign_self(struct UTUN_INSTANCE* instance, struct TOPO_NODE* ni); int topo_node_sign_self(struct UTUN_INSTANCE* instance, struct TOPO_NODE* ni);
/* Проверить timestamp, привязку node_id к ключу и подпись. 0/-1. */
int topo_node_verify(const struct TOPO_NODE* ni); int topo_node_verify(const struct TOPO_NODE* ni);
static inline const struct TOPO_SOCKMETA4* topo_v4_sock_meta(const struct TOPO_NODE* ni) { return ni->v4_sock_meta; } static inline const struct TOPO_SOCKMETA4* topo_v4_sock_meta(const struct TOPO_NODE* ni) { return ni->v4_sock_meta; }

7
src/routing_layer/topo_recovery.h

@ -1,3 +1,6 @@
/* topo_recovery — восстановление путей группы после потери транспорта.
* Собирает потерянные узлы, пробует кандидатов через групповые запросы и ждёт
* нужных маршрутов через READY-пиров. Контекст принадлежит группе; работа — в uasync. */
#ifndef TOPO_RECOVERY_H #ifndef TOPO_RECOVERY_H
#define TOPO_RECOVERY_H #define TOPO_RECOVERY_H
@ -30,9 +33,13 @@ struct TOPO_RECOVERY_CTX;
* NODEINFO/смена состояния лишь планируют проверку через call_soon: текущий * NODEINFO/смена состояния лишь планируют проверку через call_soon: текущий
* BGP callback завершается до изменения попытки или освобождения контекста. * BGP callback завершается до изменения попытки или освобождения контекста.
* Остановка группы отменяет recovery до освобождения сессий и таблицы узлов. */ * Остановка группы отменяет recovery до освобождения сессий и таблицы узлов. */
/* Добавить потерянную цель до освобождения её записи; rtt — в единицах 0.1 ms. */
void topo_recovery_add_node(struct TOPO_GROUP* group, uint64_t node_id, uint64_t next_hop, uint16_t rtt); void topo_recovery_add_node(struct TOPO_GROUP* group, uint64_t node_id, uint64_t next_hop, uint16_t rtt);
/* Запустить собранные попытки после удаления нерабочих путей. */
void topo_recovery_start(struct TOPO_GROUP* group); void topo_recovery_start(struct TOPO_GROUP* group);
/* Отложить проверку прогресса; повторные уведомления объединяются. */
void topo_recovery_changed(struct TOPO_GROUP* group); void topo_recovery_changed(struct TOPO_GROUP* group);
/* Отменить попытку и освободить контекст; готовые сессии остаются у группы. */
void topo_recovery_cancel_all(struct TOPO_GROUP* group); void topo_recovery_cancel_all(struct TOPO_GROUP* group);
#ifdef __cplusplus #ifdef __cplusplus

49
src/transport_layer/etcp.h

@ -1,4 +1,8 @@
// etcp.h - ETCP Protocol Header (refactored based on etcp_protocol.txt) /* etcp — надёжный зашифрованный транспорт с несколькими линками к одному пиру.
* Хранит очереди, ACK, повторы и состояние транспортной сессии; normalizer
* разбивает/собирает кодограммы. Сервисы используют etcp_api.h или etcp_router.h.
* Владельцы прямых подключений работают через NCD. Функции этого заголовка —
* внутренний API стека, выполняемый в потоке uasync экземпляра. */
#ifndef ETCP_H #ifndef ETCP_H
#define ETCP_H #define ETCP_H
@ -14,8 +18,6 @@ extern "C" {
#include "pkt_normalizer.h" #include "pkt_normalizer.h"
struct stcp_link; // forward declaration struct stcp_link; // forward declaration
// In struct ETCP_CONN, add:
//struct pn_pair* normalizer;
// Forward declarations // Forward declarations
struct UTUN_INSTANCE; struct UTUN_INSTANCE;
@ -23,12 +25,13 @@ struct ETCP_CONN;
struct etcp_cbk_entry; // defined in etcp_api.h struct etcp_cbk_entry; // defined in etcp_api.h
struct UASYNC; struct UASYNC;
/* Младшие 16 бит локального времени в единицах 0.1 ms (циклический timestamp). */
uint16_t get_current_timestamp(void); uint16_t get_current_timestamp(void);
// ETCP packet section types (from protocol spec) // ETCP packet section types (from protocol spec)
#define ETCP_SECTION_PAYLOAD 0x00 // Data payload #define ETCP_SECTION_PAYLOAD 0x00 // Data payload
#define ETCP_SECTION_ACK 0x01 // ACK section #define ETCP_SECTION_ACK 0x01 // ACK section
#define ETCP_SECTION_TIMESTAMP 0x06 // Channel timestamp (example, adjust if needed) #define ETCP_SECTION_TIMESTAMP 0x06 // Channel timestamp
#define ETCP_SECTION_MEAS_TS 0x07 // Measurement timestamp for bandwidth (burst packet) #define ETCP_SECTION_MEAS_TS 0x07 // Measurement timestamp for bandwidth (burst packet)
#define ETCP_SECTION_MEAS_RESP 0x08 // Measurement response (burst result) #define ETCP_SECTION_MEAS_RESP 0x08 // Measurement response (burst result)
#define ETCP_SECTION_FILLER 0x09 // Filler/dummy data (discarded by receiver) #define ETCP_SECTION_FILLER 0x09 // Filler/dummy data (discarded by receiver)
@ -99,8 +102,8 @@ struct ACK_PACKET {
// ETCP connection structure (refactored) // ETCP connection structure (refactored)
struct ETCP_CONN { struct ETCP_CONN {
// State: 0=not ready, 1=ready (indexed in instance->connections), 2=deleted // state описывает объект в реестре; доступность транспорта определяется links_up.
int state; // 0=pending, 1=ready, 2=deleted (phase 1 of close done) int state; // 0=pending, 1=индексирован по peer_node_id, 2=удалён (phase 1 завершена)
int ref_count; // External reference count. >0 blocks deferred resource free. int ref_count; // External reference count. >0 blocks deferred resource free.
// Take/free via etcp_conn_ref_take()/etcp_conn_ref_free(). // Take/free via etcp_conn_ref_take()/etcp_conn_ref_free().
@ -108,9 +111,9 @@ struct ETCP_CONN {
struct UTUN_INSTANCE* instance; struct UTUN_INSTANCE* instance;
// Queue entries in instance->connections or instance->pending_connections // Запись в реестре instance->connections; известный пир индексируется по node_id.
struct ll_entry* conn_queue_entry; // entry в очереди instance struct ll_entry* conn_queue_entry; // entry в очереди instance
struct ll_queue* conn_queue; // указатель на очередь где лежим (pending или connections) struct ll_queue* conn_queue; // заимствованная очередь, содержащая conn_queue_entry
// Links (channels) - linked list // Links (channels) - linked list
struct ETCP_LINK* links; struct ETCP_LINK* links;
@ -146,7 +149,7 @@ struct ETCP_CONN {
void (*link_ready_for_send_fn)(struct ETCP_CONN*);// функцию которую должен вызвать драйвер линка при готовности линка принимать данные void (*link_ready_for_send_fn)(struct ETCP_CONN*);// функцию которую должен вызвать драйвер линка при готовности линка принимать данные
struct ll_queue* output_queue; // Assembled outgoing packets (storage: ETCP_FRAGMENT / rx_pool) struct ll_queue* output_queue; // Assembled outgoing packets (storage: ETCP_FRAGMENT / rx_pool)
struct ll_queue* transit_queues; // hash по (src_node_id:8, dst_node_id:8) — транзитные очереди struct ll_queue* transit_queues; // ожидающий транзит, ключ (group_id:8, src_node_id:8, dst_node_id:8)
struct ll_queue* send_input_q; // единая входная очередь отправки (normalizer->input или tx_queue) struct ll_queue* send_input_q; // единая входная очередь отправки (normalizer->input или tx_queue)
@ -198,8 +201,8 @@ struct ETCP_CONN {
uint8_t initialized; // 0 - только созданный ETCP, 1 - хотя бы один линк проинициалзирован (обмен ключами произведен) uint8_t initialized; // 0 - только созданный ETCP, 1 - хотя бы один линк проинициалзирован (обмен ключами произведен)
uint64_t reset_id; // локальная эпоха потока (меняется только при локальном reset) uint64_t reset_id; // локальная эпоха потока (меняется только при локальном reset)
uint64_t peer_reset_id; // подтверждённая эпоха потока пира uint64_t peer_reset_id; // подтверждённая эпоха потока пира
uint64_t session_candidate, session_challenge, session_started; uint64_t session_candidate, session_challenge, session_started; // эпоха-кандидат, challenge, время начала (0.1 ms)
void* session_timer; void* session_timer; // таймер согласования транспортной сессии
uint8_t session_required; // только подтверждение challenge может завершить локальный reset uint8_t session_required; // только подтверждение challenge может завершить локальный reset
uint8_t tx_state; // 0 - n/a, 1 - data_wait (queues empty), 2 - link_wait (link busy) uint8_t tx_state; // 0 - n/a, 1 - data_wait (queues empty), 2 - link_wait (link busy)
uint8_t links_up; // 0 - канал не готов для передачи, 1 - канал готов для передачи (хотя бы один линк не down) uint8_t links_up; // 0 - канал не готов для передачи, 1 - канал готов для передачи (хотя бы один линк не down)
@ -207,10 +210,10 @@ struct ETCP_CONN {
int callbacks_running; // счётчик вложенности колбэк-цепочек (>0 — внутри цепочки, etcp_connection_close откладывается) int callbacks_running; // счётчик вложенности колбэк-цепочек (>0 — внутри цепочки, etcp_connection_close откладывается)
uint8_t close_requested; // новые операции запрещены, phase 1 может ждать выхода из callback uint8_t close_requested; // новые операции запрещены, phase 1 может ждать выхода из callback
uint8_t close_running; uint8_t close_running; // защита от повторного входа в close
uint8_t delete_complete; uint8_t delete_complete; // detach и рассылка DELETE завершены
uint8_t free_scheduled; uint8_t free_scheduled; // освобождение уже поставлено в uasync
void* close_token; void* close_token; // отложенный close при вызове внутри callback
// Unified callback chain with event mask (init/reinit/up/down/node_changed) // Unified callback chain with event mask (init/reinit/up/down/node_changed)
struct etcp_cbk_entry* cbks; struct etcp_cbk_entry* cbks;
@ -235,30 +238,28 @@ struct ETCP_CONN {
#define RTT_CB_PERIOD_TB 600000 // 1 минута в 0.1ms #define RTT_CB_PERIOD_TB 600000 // 1 минута в 0.1ms
// Functions // Создать транспортный объект; жизненным циклом прикладного соединения управляет NCD.
struct ETCP_CONN* etcp_connection_create(struct UTUN_INSTANCE* instance, char* name); struct ETCP_CONN* etcp_connection_create(struct UTUN_INSTANCE* instance, char* name);
/** /**
* @brief Закрыть соединение: фаза-1 detach + отложенная фаза-2 очистка. * @brief Закрыть соединение: фаза-1 detach + отложенная фаза-2 очистка.
* *
* Порядок коллбэков при закрытии (строгая последовательность): * Порядок коллбэков при закрытии (строгая последовательность):
* 1) если соединение ещё UP (links_up != 0) — финальный per-connection DOWN * 1) если соединение ещё UP — links_up = 0 и финальный per-connection DOWN;
* (ETCP_CBK_EVENT_DOWN), затем links_up = 0;
* 2) state = 2; * 2) state = 2;
* 3) conn-status DELETE (ETCP_CONN_STATUS_DELETE) — единственный статус при * 3) per-connection DELETE, затем instance-level DELETE;
* state==2 (по нему подписчики снимают conn из своих списков);
* 4) cleanup: таймеры, routing, линки, очередь, deferred u_free. * 4) cleanup: таймеры, routing, линки, очередь, deferred u_free.
* *
* После state==2 per-connection события (etcp_cbk_fire) и conn-status * При вызове внутри callback закрытие откладывается, close_requested ставится сразу.
* NEW/UP/DOWN не вызываются. * После state==2 допустим только DELETE; ref_count задерживает освобождение памяти.
*/ */
void etcp_connection_close(struct ETCP_CONN* etcp); void etcp_connection_close(struct ETCP_CONN* etcp);
void etcp_conn_queue_set_ready(struct ETCP_CONN* conn); // move pending->connections, fire ready cbks void etcp_conn_queue_set_ready(struct ETCP_CONN* conn); // переиндексировать запись по известному peer_node_id
/** /**
* @brief Take a reference on the connection (blocks deferred resource free). * @brief Take a reference on the connection (blocks deferred resource free).
* @param conn connection * @param conn connection
* @return 0 on success, -1 if conn is NULL or already in state 2 (deleted). * @return 0 on success, -1 if conn is NULL, closing or already deleted (state 2).
* *
* Increments ref_count. While ref_count > 0, etcp_connection_close() * Increments ref_count. While ref_count > 0, etcp_connection_close()
* will detach but defer resource cleanup until all references are released. * will detach but defer resource cleanup until all references are released.

44
src/transport_layer/etcp_api.h

@ -1,17 +1,11 @@
/** /**
* @file etcp_api.h * @file etcp_api.h
* @brief API для приёма-передачи пакетов через ETCP (per-instance bindings) * @brief Прямой обмен кодограммами ETCP и подписки на события транспорта.
* *
* Основные функции: * bind выбирает обработчик по первому байту cmd, send передаёт кодограмму соседу.
* - etcp_send() - отправить пакет в очередь normalizer * Для доставки через промежуточные узлы используется etcp_router.h.
* - etcp_bind() - подписаться на пакеты с определенным ID * Все вызовы — в потоке uasync. Перед отправкой ждать UP и свободной send_input_q.
* - etcp_int_recv() - коллбэк для сбора пакетов из всех подключений * Владение entry различается: send забирает его только при успехе, recv — всегда.
* !!! для обычной отправки-приёма между узлами используем более универсальный модушь etcp_router.
* это в первую очередь - апи для использования etcp-router-ом
*
* Формат кодограмм: <cmd 1 byte> <data ... n bytes>
* cmd = 0 - пакет для передачи адресату
* cmd = 1 - модуль обмена роутинг-таблицами
*/ */
#ifndef ETCP_API_H #ifndef ETCP_API_H
@ -86,11 +80,12 @@ extern "C" {
#define ETCP_RT_ID_CALL 0x36 // call — P2P аудио-звонок #define ETCP_RT_ID_CALL 0x36 // call — P2P аудио-звонок
// Connection status events (instance-level callback). // Connection status events (instance-level callback).
// Строгая последовательность: NEW -> UP <-> DOWN -> DELETE. // NEW начинает жизнь объекта, UP/DOWN отражают доступность, DELETE завершает её.
// NEW — один раз при создании соединения. // NEW — один раз при создании соединения.
// UP — переход links_up 0→1 (UP после UP не вызывается). // UP — переход links_up 0→1 (UP после UP не вызывается).
// DOWN — переход links_up 1→0 (DOWN после DOWN не вызывается). // DOWN — переход links_up 1→0 (DOWN после DOWN не вызывается).
// DELETE — один раз при закрытии, финальный; единственный статус, допустимый при state==2. // DELETE — один раз при закрытии, финальный; единственный статус, допустимый при state==2.
// REINIT — новая транспортная сессия; не означает потерю физического пути.
#define ETCP_CONN_STATUS_NEW 0 // соединение создано #define ETCP_CONN_STATUS_NEW 0 // соединение создано
#define ETCP_CONN_STATUS_UP 1 // соединение поднялось #define ETCP_CONN_STATUS_UP 1 // соединение поднялось
#define ETCP_CONN_STATUS_DOWN 2 // соединение упало #define ETCP_CONN_STATUS_DOWN 2 // соединение упало
@ -115,16 +110,10 @@ struct etcp_status_cbk_entry {
* @param conn ETCP соединение от которого получен пакет * @param conn ETCP соединение от которого получен пакет
* @param entry Элемент очереди с данными пакета * @param entry Элемент очереди с данными пакета
* *
* @note Коллбэк должен освободить entry через queue_entry_free() * @note После обработки освободить сначала dgram через queue_dgram_free(),
* и dgram через queue_dgram_free() после обработки * затем entry через queue_entry_free(). conn заимствован на время вызова.
* * @note Raw bind получает cmd || payload. Router bind получает заголовок
* @note Формат кодограммы для etcp_route_send: * ROUTER_SVC_* и payload со смещения ROUTER_SVC_PAYLOAD_OFF (etcp_router.h).
* entry->dgram[0] = svc_id (напр. ETCP_RT_ID_MEDIA_DELIVERY=0x07).
* При приёме router_deliver предваряет payload байтом svc_id:
* entry->dgram[0]=svc_id, entry->dgram[1..]=payload (исходный dgram).
* Если payload содержит cmd+subcmd (двухуровневый, как conn_mgr),
* subcmd в entry->dgram[2]. Если одноуровневый (media_delivery),
* subcmd в entry->dgram[1].
*/ */
typedef void (*etcp_recv_fn)(struct ETCP_CONN* conn, struct ll_entry* entry); typedef void (*etcp_recv_fn)(struct ETCP_CONN* conn, struct ll_entry* entry);
@ -134,7 +123,7 @@ struct UTUN_INSTANCE;
struct TOPO_GROUP_NODE; struct TOPO_GROUP_NODE;
// Per-connection callback event types (bitmask). // Per-connection callback event types (bitmask).
// Вызываются только при state != 2 (после DELETE события не идут). // При state==2 допустим только терминальный ETCP_CBK_EVENT_DELETE.
// UP/DOWN зеркалят conn-status UP/DOWN; при закрытии (если соединение было UP) // UP/DOWN зеркалят conn-status UP/DOWN; при закрытии (если соединение было UP)
// отправляется один финальный DOWN до state=2. // отправляется один финальный DOWN до state=2.
#define ETCP_CBK_EVENT_INIT (1 << 0) #define ETCP_CBK_EVENT_INIT (1 << 0)
@ -160,8 +149,7 @@ struct etcp_cbk_entry {
* *
* @note Вызываются все коллбэки соединения, у которых event_mask * @note Вызываются все коллбэки соединения, у которых event_mask
* пересекается с переданным event. * пересекается с переданным event.
* @note Не вызывается при state==2 (соединение удалено) — после DELETE * @note При state==2 допускается только DELETE; прочие события игнорируются.
* per-connection события не рассылаются.
*/ */
void etcp_cbk_fire(struct ETCP_CONN* conn, int event); void etcp_cbk_fire(struct ETCP_CONN* conn, int event);
@ -329,7 +317,7 @@ int etcp_bind(struct UTUN_INSTANCE* inst, uint8_t id, etcp_recv_fn callback);
* *
* @param inst UTUN instance * @param inst UTUN instance
* @param id Идентификатор пакета * @param id Идентификатор пакета
* @return 0 при успехе, -1 если binding не найден * @return 0 при успехе (в том числе без прежнего binding), -1 при неверном instance
*/ */
int etcp_unbind(struct UTUN_INSTANCE* inst, uint8_t id); int etcp_unbind(struct UTUN_INSTANCE* inst, uint8_t id);
@ -341,8 +329,8 @@ int etcp_unbind(struct UTUN_INSTANCE* inst, uint8_t id);
* @param arg Пользовательский аргумент * @param arg Пользовательский аргумент
* *
* @note status = ETCP_CONN_STATUS_NEW / UP / DOWN / DELETE / REINIT. * @note status = ETCP_CONN_STATUS_NEW / UP / DOWN / DELETE / REINIT.
* Строгая последовательность NEW -> UP <-> DOWN -> DELETE; DELETE — * DELETE — единственный статус при state==2. Финальный DOWN при close
* единственный статус, отправляемый при state==2 (при закрытии). * доставляется per-connection подписчикам; instance-подписчики получают DELETE.
*/ */
void etcp_add_conn_status_cbk(struct UTUN_INSTANCE* inst, etcp_conn_status_fn fn, void* arg); void etcp_add_conn_status_cbk(struct UTUN_INSTANCE* inst, etcp_conn_status_fn fn, void* arg);

31
src/transport_layer/node_conn_direct.h

@ -1,22 +1,23 @@
#ifndef NODE_CONN_DIRECT_H
#define NODE_CONN_DIRECT_H
/* /*
* node_conn_direct.h — handle-based прямое подключение к удалённому узлу (node_id) * node_conn_direct — общий прямой транспорт к узлу по node_id.
* *
* Модуль управляет ETCP-соединениями через непрозрачные handle'ы (NODE_CONN_DIRECT). * Модуль управляет ETCP-соединениями через непрозрачные handle'ы (NODE_CONN_DIRECT).
* Несколько handle'ов могут разделять одно ETCP_CONN — закрытие происходит только * Каждый владелец открывает свой handle; несколько handle разделяют ETCP_CONN.
* когда все handle'ы закрыты (refcounting). При закрытии последнего handle * close снимает только своё владение. Последний handle закрывает неготовый транспорт
* используется graceful shutdown: CLOSE/KEEP_ALIVE протокол (fin_wait 5 сек). * сразу; для UP начинает CLOSE/KEEP_ALIVE (до 5 s), пир может сохранить соединение.
* *
* Кто использует: * Кто использует:
* - conn_mgr (Connection Manager) — создаёт handle'ы для DIR/REV/IND-соединений * - conn_mgr (Connection Manager) — создаёт handle'ы для DIR/REV/IND-соединений
* - chat_core / topo_group — связь с узлами канала/группы * - topo_group — владеет транспортом групповой сессии с начала CONNECTING
* - протокол join — временный транспорт до добавления участника
* - любой модуль, кому нужно надёжное ETCP-соединение с конкретным node_id * - любой модуль, кому нужно надёжное ETCP-соединение с конкретным node_id
* *
* Переходы UP/DOWN/TIMEOUT/CLOSED доставляются синхронно. Закрытие любого handle * Переходы UP/DOWN/TIMEOUT/CLOSED доставляются синхронно. Закрытие любого handle
* из callback безопасно. Начальное UP при open готового conn откладывается через uasync_call_soon. * из callback безопасно. Начальное UP при open готового conn откладывается через uasync_call_soon.
* Все операции — в потоке uasync. NCD не проверяет членство и готовность группы.
*/ */
#ifndef NODE_CONN_DIRECT_H
#define NODE_CONN_DIRECT_H
#include <stdint.h> #include <stdint.h>
@ -48,7 +49,8 @@ typedef void (*ncd_callback)(struct NODE_CONN_DIRECT* h, enum ncd_event event, v
* Node info ищется через node_registry, fallback — SQLite. * Node info ищется через node_registry, fallback — SQLite.
* cb вызывается при изменении статуса / таймауте. * cb вызывается при изменении статуса / таймауте.
* Если conn уже готов — cb(NCD_EVENT_UP) через uasync_call_soon (не синхронно). * Если conn уже готов — cb(NCD_EVENT_UP) через uasync_call_soon (не синхронно).
* Пока не придёт событие UP - отправляения могут теряться. Поэтому перед первой отправкой всегда ждём UP. * До первой отправки дождаться UP. Результат: NCD_NEW/NCD_REUSED или NCD_ERR.
* При успехе вызывающий обязан закрыть *out_handle, в том числе после CLOSED.
*/ */
int node_conn_direct_open(struct UTUN_INSTANCE* inst, uint64_t node_id, int node_conn_direct_open(struct UTUN_INSTANCE* inst, uint64_t node_id,
ncd_callback cb, void* cb_arg, ncd_callback cb, void* cb_arg,
@ -65,16 +67,15 @@ int node_conn_direct_open_node(struct UTUN_INSTANCE* inst, uint64_t node_id,
struct TOPO_NODE* ni, /* публичный ключ + адреса пира (не владеем) */ struct TOPO_NODE* ni, /* публичный ключ + адреса пира (не владеем) */
struct ETCP_SOCKET* specific_sock); /* NULL — авто-подбор */ struct ETCP_SOCKET* specific_sock); /* NULL — авто-подбор */
/* Освободить handle и отменить его callbacks; другие владельцы сохраняют транспорт. */
void node_conn_direct_close(struct NODE_CONN_DIRECT* h); void node_conn_direct_close(struct NODE_CONN_DIRECT* h);
/* Применить принятую запись из общего реестра к уже существующему соединению. /* Применить принятую запись из общего реестра к уже существующему соединению.
* Не создаёт соединение и не приобретает handle. Возвращает число новых линков. */ * Не создаёт соединение и не приобретает handle. Число новых линков или -1 при ошибке. */
int node_conn_direct_update_node(struct UTUN_INSTANCE* inst, uint64_t node_id); int node_conn_direct_update_node(struct UTUN_INSTANCE* inst, uint64_t node_id);
/* Force close без fin_wait: немедленно удаляет ETCP-коллбэки и, если это /* Как close, но последний handle закрывает транспорт сразу, без fin_wait.
* последний handle, закрывает соединение и освобождает ncd_entry. * При наличии других handle общий транспорт и их callbacks сохраняются. */
* Используется при уничтожении CM entry — гарантирует, что после вызова
* никакие NCD-коллбэки не доставят события в освобождённую память. */
void node_conn_direct_force_close(struct NODE_CONN_DIRECT* h); void node_conn_direct_force_close(struct NODE_CONN_DIRECT* h);
/* Сменить или сбросить (cb=NULL) callback на уже открытом handle. /* Сменить или сбросить (cb=NULL) callback на уже открытом handle.
@ -86,8 +87,10 @@ void node_conn_direct_set_callback(struct NODE_CONN_DIRECT* h, ncd_callback cb,
* должен забыть указатель. Соединение не трогается. */ * должен забыть указатель. Соединение не трогается. */
void node_conn_direct_transfer(struct NODE_CONN_DIRECT* h, ncd_callback new_cb, void* new_cb_arg); void node_conn_direct_transfer(struct NODE_CONN_DIRECT* h, ncd_callback new_cb, void* new_cb_arg);
/* Идентификатор узла; 0 для NULL. */
uint64_t node_conn_direct_node_id(const struct NODE_CONN_DIRECT* h); uint64_t node_conn_direct_node_id(const struct NODE_CONN_DIRECT* h);
/* Заимствованный транспорт, в том числе до UP; NULL после CLOSED. Не закрывать напрямую. */
struct ETCP_CONN* node_conn_direct_get_conn(struct NODE_CONN_DIRECT* h); struct ETCP_CONN* node_conn_direct_get_conn(struct NODE_CONN_DIRECT* h);
/* ─── Протокол CLOSE / KEEP_ALIVE (ETCP_RT_ID_NCD_CONTROL = 0x12) ─── */ /* ─── Протокол CLOSE / KEEP_ALIVE (ETCP_RT_ID_NCD_CONTROL = 0x12) ─── */

76
src/utun_instance.h

@ -1,3 +1,8 @@
/* utun_instance — владелец ядра и сервисов одного узла.
* Ядро содержит транспорт, маршрутизатор, идентичность и общую БД.
* create → core_start; UTUN и чат запускаются отдельно через service_start/stop.
* Все операции жизненного цикла — в потоке переданного uasync, вне callbacks сервисов.
* destroy останавливает сервисы и освобождает ядро, но не сам uasync. */
#ifndef UTUN_INSTANCE_H #ifndef UTUN_INSTANCE_H
#define UTUN_INSTANCE_H #define UTUN_INSTANCE_H
@ -75,11 +80,11 @@ struct conn_queue_entry {
struct ETCP_CONN* conn; struct ETCP_CONN* conn;
}; };
// Handle for config-based connections via node_conn_direct // Постоянный запрос UTUN-сервиса к пиру из конфигурации; транспортом владеет группа.
struct CONFIG_CONN_HANDLE { struct CONFIG_CONN_HANDLE {
uint64_t node_id; uint64_t node_id;
char name[MAX_CONN_NAME_LEN]; char name[MAX_CONN_NAME_LEN];
struct TOPO_PEER_REQUEST* request; struct TOPO_PEER_REQUEST* request; // принадлежит этой записи; закрывается при stop/reload
struct CONFIG_CONN_HANDLE* next; struct CONFIG_CONN_HANDLE* next;
}; };
@ -109,23 +114,23 @@ struct peer_sleep_cbk_entry {
struct peer_sleep_cbk_entry* next; struct peer_sleep_cbk_entry* next;
}; };
// uTun instance configuration // Состояние одного узла. Вложенные модули освобождаются через их lifecycle API.
struct UTUN_INSTANCE { struct UTUN_INSTANCE {
uint8_t core_started, utun_started, chat_started; uint8_t core_started, utun_started, chat_started; // успешный запуск ядра / UTUN / чата
// Identification // Identification
char name[MAX_CONN_NAME_LEN]; // Instance name from config char name[MAX_CONN_NAME_LEN]; // Instance name from config
// Configuration (moved from utun_state) // Конфигурация принадлежит экземпляру после успешного create.
struct utun_config *config; struct utun_config *config;
// TUN interface // TUN interface
struct tun_if* tun; struct tun_if* tun;
// Route subnets (for cleanup on shutdown) // Заимствованный список системных маршрутов из config; нужен для stop UTUN.
struct CFG_ROUTE_ENTRY* route_subnets; struct CFG_ROUTE_ENTRY* route_subnets;
struct ROUTE_TABLE* rt; struct ROUTE_TABLE* rt; // таблица маршрутов подсетей
struct TOPO_GROUPS* topo_groups; // Groups module for topology exchange struct TOPO_GROUPS* topo_groups; // Groups module for topology exchange
sqlite3* topo_sqlite_db; // Shared SQLite DB (nodes/channels/peers) sqlite3* topo_sqlite_db; // Shared SQLite DB (nodes/channels/peers)
struct NAT_DETECTION* nat_det; // NAT detection module struct NAT_DETECTION* nat_det; // NAT detection module
@ -141,11 +146,11 @@ struct UTUN_INSTANCE {
uint8_t next_socket_id; // Counter for unique socket IDs (0-255) uint8_t next_socket_id; // Counter for unique socket IDs (0-255)
// Main async context // Заимствованный цикл событий: создаёт и уничтожает вызывающий.
struct UASYNC* ua; struct UASYNC* ua;
// State // State
int running; int running; // продолжать utun_instance_run; не флаг готовности сервисов
// Connections (очередь всех подключений для instance) // Connections (очередь всех подключений для instance)
struct ll_queue* connections; // indexed by peer_node_id (0=pending), data=conn_queue_entry struct ll_queue* connections; // indexed by peer_node_id (0=pending), data=conn_queue_entry
@ -160,8 +165,8 @@ struct UTUN_INSTANCE {
void* test_user_ptr; // Generic user pointer (used by tests) void* test_user_ptr; // Generic user pointer (used by tests)
struct memory_pool* data_pool;// для входных-выходных данных пакета struct memory_pool* data_pool;// для входных-выходных данных пакета
struct memory_pool* pkt_pool; struct memory_pool* pkt_pool; // ETCP_DGRAM и буферы сетевых пакетов
struct memory_pool* ack_pool; struct memory_pool* ack_pool; // ACK_PACKET
// Active sockets (UDP + TCP, is_tcp=1 flag) // Active sockets (UDP + TCP, is_tcp=1 flag)
struct ETCP_SOCKET* etcp_sockets;// linked-list (UDP + TCP) struct ETCP_SOCKET* etcp_sockets;// linked-list (UDP + TCP)
@ -205,25 +210,25 @@ struct UTUN_INSTANCE {
// etcp_router bindings и seq-connections (per-instance service routing) // etcp_router bindings и seq-connections (per-instance service routing)
struct ETCP_ROUTER_BINDINGS router_bindings; struct ETCP_ROUTER_BINDINGS router_bindings;
struct ll_queue* router_conns; struct ll_queue* router_conns;
struct CONN_MGR* conn_mgr; // Connection Manager (может быть NULL) struct CONN_MGR* conn_mgr; // заимствован у UTUN-группы; NULL при остановленном UTUN
struct DB_SYNC* db_sync; // Distributed DB sync (может быть NULL) struct DB_SYNC* db_sync; // Distributed DB sync (может быть NULL)
// Chat/DM subsystem (per-instance contexts; were global singletons) // Контексты чат-сервиса; освобождаются при chat_service_stop.
struct chat_core_ctx* chat_core; // chat_core.c (was g_cc) struct chat_core_ctx* chat_core; // каналы, сообщения и доступ к общей БД
struct chat_sync* chat_sync; // chat_sync.c (was g_cs) struct chat_sync* chat_sync; // join и синхронизация каналов
struct join_key_entry* join_keys; // chat_join.c (was g_keys; собственный init/destroy) struct join_key_entry* join_keys; // зарегистрированные ключи приглашений
struct chat_join_registration* join_registrations; struct chat_join_registration* join_registrations; // ожидания KEY_REGISTER_ACK
int join_stopping; int join_stopping; // запрещает новые регистрации при teardown
struct chat_invite_build_req* invite_gui_request; struct chat_invite_build_req* invite_gui_request; // текущее построение ссылки для GUI
struct dm_state* dm; // dm/dm_core.c (was g_dm) struct dm_state* dm; // личные сообщения
struct dm_mb_state* dm_mailbox; // dm/dm_mailbox.c (was g_mb) struct dm_mb_state* dm_mailbox; // хранилище личных сообщений для получателей
struct call_ctx* call; // call/call.c — P2P audio call struct call_ctx* call; // call/call.c — P2P audio call
void* call_audio; // call/call_headless.c — audio stream socket void* call_audio; // call/call_headless.c — audio stream socket
struct radio_instance* radio; // radio/radio.c — PTT walkie-talkie (frame cb + counters) struct radio_instance* radio; // radio/radio.c — PTT walkie-talkie (frame cb + counters)
struct radio_headless* radio_headless; // radio/radio_headless.c — audio stream socket struct radio_headless* radio_headless; // radio/radio_headless.c — audio stream socket
struct chat_setting_state chat_settings; // chat_setting.c (was globals) struct chat_setting_state chat_settings; // текущие настройки чата
chat_event_handler_fn chat_event_handler; // chat_event.c (was g_handler) chat_event_handler_fn chat_event_handler; // получатель событий чата
void* headless; // chat_headless_control.c (was g_hc) void* headless; // контекст TCP-управления headless-чатом
// E2E encryption cache — per-peer sc_context_t with derived session key // E2E encryption cache — per-peer sc_context_t with derived session key
#define E2E_CTX_CACHE_SIZE 8 #define E2E_CTX_CACHE_SIZE 8
@ -253,7 +258,7 @@ struct UTUN_INSTANCE {
struct NTP_TIME ntp; struct NTP_TIME ntp;
struct NTP_NODE_TIME ntp_node; struct NTP_NODE_TIME ntp_node;
// Config-based connection handles (node_conn_direct) // Постоянные групповые запросы UTUN к пирам из конфигурации.
struct CONFIG_CONN_HANDLE* config_conn_handles; struct CONFIG_CONN_HANDLE* config_conn_handles;
// Per-instance NCD state (replaces global statics in node_conn_direct.c) // Per-instance NCD state (replaces global statics in node_conn_direct.c)
@ -261,22 +266,29 @@ struct UTUN_INSTANCE {
uint8_t ncd_control_bound; // 1 = etcp_bind(ETCP_RT_ID_NCD_CONTROL) сделан uint8_t ncd_control_bound; // 1 = etcp_bind(ETCP_RT_ID_NCD_CONTROL) сделан
}; };
// Functions // Создать ядро по файлу конфигурации; NULL при ошибке. Сервисы ещё не запущены.
struct UTUN_INSTANCE* utun_instance_create(struct UASYNC* ua, const char* config_file); struct UTUN_INSTANCE* utun_instance_create(struct UASYNC* ua, const char* config_file);
// При успехе забирает config; при ошибке config остаётся у вызывающего.
struct UTUN_INSTANCE* utun_instance_create_from_config(struct UASYNC* ua, struct utun_config* config); struct UTUN_INSTANCE* utun_instance_create_from_config(struct UASYNC* ua, struct utun_config* config);
// Как create, но конфигурация передана строкой; строку не сохраняет.
struct UTUN_INSTANCE* utun_instance_create_from_str(struct UASYNC* ua, const char* config_text); struct UTUN_INSTANCE* utun_instance_create_from_str(struct UASYNC* ua, const char* config_text);
// Остановить сервисы и освободить instance. Переданный uasync остаётся у вызывающего.
void utun_instance_destroy(struct UTUN_INSTANCE* instance); void utun_instance_destroy(struct UTUN_INSTANCE* instance);
/* create constructs the core; start enables its runtime. Services are independent. /* Запустить общие сетевые службы ядра. Повтор допустим; 0/-1. */
* All lifecycle calls run on the owning uasync thread, outside service callbacks.
* Stop is idempotent. destroy stops both services before releasing core resources. */
int utun_core_start(struct UTUN_INSTANCE* instance); int utun_core_start(struct UTUN_INSTANCE* instance);
/* Запустить UTUN-группу, TUN/NAT по конфигу и запросы к пирам. Нужен core_start; 0/-1. */
int utun_service_start(struct UTUN_INSTANCE* instance); int utun_service_start(struct UTUN_INSTANCE* instance);
/* Остановить UTUN, сохранив ядро и чат. Повтор допустим. */
void utun_service_stop(struct UTUN_INSTANCE* instance); void utun_service_stop(struct UTUN_INSTANCE* instance);
/* Convenience: core + UTUN + configured chat. */ /* Запустить ядро, UTUN и включённый в конфиге чат. 0/-1. */
int utun_instance_init(struct UTUN_INSTANCE *instance); int utun_instance_init(struct UTUN_INSTANCE *instance);
/* Перечитать конфигурацию, при необходимости пересоздать экземпляр; использовать возвращённый указатель. */
struct UTUN_INSTANCE *utun_instance_reload(struct UTUN_INSTANCE *instance, struct UASYNC *ua, const char *config_file); struct UTUN_INSTANCE *utun_instance_reload(struct UTUN_INSTANCE *instance, struct UASYNC *ua, const char *config_file);
/* Выполнять цикл событий до stop(). */
void utun_instance_run(struct UTUN_INSTANCE *instance); void utun_instance_run(struct UTUN_INSTANCE *instance);
/* Снять running и разбудить цикл; ресурсы освобождает destroy(), а не stop(). */
void utun_instance_stop(struct UTUN_INSTANCE *instance); void utun_instance_stop(struct UTUN_INSTANCE *instance);
/* Глобальные переключатели создания TUN/топологии для тестов и встраивания. */
void utun_instance_set_tun_init_enabled(int enabled); void utun_instance_set_tun_init_enabled(int enabled);
void utun_instance_set_topo_group_enabled(int enabled); void utun_instance_set_topo_group_enabled(int enabled);
void utun_set_client_activity(struct UTUN_INSTANCE* instance, int active); void utun_set_client_activity(struct UTUN_INSTANCE* instance, int active);
@ -289,8 +301,8 @@ void utun_fire_peer_sleep_cbk(struct UTUN_INSTANCE* instance, uint64_t peer_node
// Diagnostic function for memory leak analysis // Diagnostic function for memory leak analysis
void utun_instance_diagnose_leaks(struct UTUN_INSTANCE* instance, const char* phase); void utun_instance_diagnose_leaks(struct UTUN_INSTANCE* instance, const char* phase);
/* Find active connection by node_id — searches both UDP (connections) and TCP (tcp_connections) queues. /* Найти транспорт по node_id в реестрах ETCP/TCP. Указатель заимствован.
Returns ETCP_CONN* ready for etcp_send(), or NULL if not found. */ * NULL — не найден; ненулевой результат не гарантирует UP или готовность группы. */
static inline struct ETCP_CONN* instance_find_conn(struct UTUN_INSTANCE* inst, uint64_t node_id) { static inline struct ETCP_CONN* instance_find_conn(struct UTUN_INSTANCE* inst, uint64_t node_id) {
if (inst->connections) { if (inst->connections) {
struct ll_entry* e = queue_find_data_by_index(inst->connections, (const uint8_t*)&node_id); struct ll_entry* e = queue_find_data_by_index(inst->connections, (const uint8_t*)&node_id);

Loading…
Cancel
Save