From 01fcefdec553b3b67212e669989ea7f204a53728 Mon Sep 17 00:00:00 2001 From: evgeny Date: Mon, 28 Sep 2026 12:01:08 +0300 Subject: [PATCH] docs: clarify core, group and chat header contracts --- src/chat/chat_core.h | 58 +++++++++++----- src/chat/chat_join.h | 14 ++-- src/chat/chat_sync.h | 37 ++++------- src/chat/db_sync.h | 65 ++++++------------ src/chat/member_sync.h | 55 +++++---------- src/chat/merkle_sync.h | 21 ++++-- src/chat/merkle_tree.h | 15 ++++- src/media_async/media_async.h | 18 +++-- src/routing_layer/conn_mgr.h | 60 +++++------------ src/routing_layer/etcp_router.h | 33 ++++++--- src/routing_layer/topo_group.h | 92 +++++++++++--------------- src/routing_layer/topo_group_connect.h | 24 ++++--- src/routing_layer/topo_group_invite.h | 37 ++++------- src/routing_layer/topo_node.h | 72 ++++++++++++-------- src/routing_layer/topo_recovery.h | 7 ++ src/transport_layer/etcp.h | 49 +++++++------- src/transport_layer/etcp_api.h | 46 +++++-------- src/transport_layer/node_conn_direct.h | 31 +++++---- src/utun_instance.h | 76 ++++++++++++--------- 19 files changed, 396 insertions(+), 414 deletions(-) diff --git a/src/chat/chat_core.h b/src/chat/chat_core.h index 410924a7..3a0bffe2 100644 --- a/src/chat/chat_core.h +++ b/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-потоке. * Все функции 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. - * Stop cancels delivery and waits for native media workers before releasing chat state. */ +/* Запустить чат поверх core_start; UTUN не требуется. Повтор допустим; 0/-1. */ int chat_service_start(struct UTUN_INSTANCE* inst); +/* Остановить чат, сохранив ядро, UTUN и БД. Повтор допустим; вызывать вне callbacks чата. + * Отменяет доставку и ждёт native media workers; длительная задача может задержать stop. */ 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); +/* Освободить chat_core после остановки остальных частей чата; общая БД остаётся у ядра. */ void chat_core_destroy(struct UTUN_INSTANCE* inst); /* Backfill: регистрирует уже скачанные медиафайлы (лежат на диске, но отсутствуют @@ -36,33 +42,36 @@ void chat_media_backfill(struct UTUN_INSTANCE* inst); void chat_media_autodownload_backfill(struct UTUN_INSTANCE* inst); /* Регистрирует триггер: после первого 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_destroy(struct UTUN_INSTANCE* inst); +/* Заимствованная БД текущего chat_core; NULL до init/после destroy. */ struct sqlite3* chat_core_get_db(struct UTUN_INSTANCE* inst); +/* 1 — chat_core инициализирован, 0 — нет; не проверяет доступность сети. */ int chat_core_is_initialized(struct UTUN_INSTANCE* inst); /* ── Отправка сообщения (GUI → uasync) ── */ struct chat_msg_submit { struct UTUN_INSTANCE* inst; - char channel_id[64]; - char content_type[32]; - char media_src[1024]; - char media_dest[1024]; - uint8_t media_copy; + char channel_id[64]; /* десятичный ID канала */ + char content_type[32]; /* тип содержимого сообщения */ + char media_src[1024]; /* исходный файл вложения */ + char media_dest[1024]; /* путь подготовленного локального файла */ + uint8_t media_copy; /* 1 = скопировать исходный файл при регистрации */ uint8_t media_video; /* 1 = подготовить как видео (пробинг+транскод) */ int duration_ms; /* видео: длительность (0 = неизвестно) */ int width; /* видео: ширина */ int height; /* видео: высота */ - uint8_t* data; - uint32_t data_len; - uint64_t timestamp; + uint8_t* data; /* содержимое сообщения, заимствовано обычным submit */ + uint32_t data_len; /* размер data в байтах */ + uint64_t timestamp; /* входное время; отправка назначает новый timestamp через db_sync */ uint64_t reply_to_ts; /* 0 = не ответ */ uint64_t reply_to_node_id; /* node_id автора исходного сообщения */ }; +/* Подписать, сохранить и разослать сообщение. req/data остаются у вызывающего. */ void chat_core_submit_message(struct UTUN_INSTANCE* inst, struct chat_msg_submit* req); /* Отправка медиа-сообщения (вложение): src->dest (copy) или видео (transcode). */ @@ -70,10 +79,18 @@ void chat_core_submit_media_message(struct UTUN_INSTANCE* inst, struct chat_msg_ /* ── DB-операции (оставлены для интроспекции) ── */ +/* Число локальных сообщений; 0 также при отсутствии БД или ошибке запроса. */ 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); +/* Список 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); +/* Список участников: 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); +/* Ключи 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); /* ── Создание канала (GUI → uasync) ── */ @@ -92,9 +109,12 @@ struct chat_channel_create { 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); +/* Создать канал по готовым ключам и метаданным req; 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_auto(struct UTUN_INSTANCE* inst, const char* name); 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); +/* Сохранить имя узла и обновить собственные member-записи каналов. */ 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]; }; void chat_core_update_my_name_trampoline(void* arg); +/* Обновить адреса topo_node и собственные member-записи каналов. */ 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); /* Сохранение 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}] */ 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}] - * count=0 → все, offset=0 → сначала */ +/* Сообщения канала в JSON: [{id,ts,author_id,author_name,content_type,data,local_attrs}]. + * От новых к старым; count<=0 — до 100 записей, 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); -/* Мемберы канала с полным состоянием в JSON: - * [{node_id,name,online,connected,x25519,ed25519,addrs:[{ip,port,proto,rtt}]}] */ +/* Мемберы в JSON: [{node_id,name,online,connected,adm_tags,x25519,ed25519,addrs}]. + * 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); /* Имя узла из таблицы nodes */ diff --git a/src/chat/chat_join.h b/src/chat/chat_join.h index 473d2a22..0c7ae1d8 100644 --- a/src/chat/chat_join.h +++ b/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 * │ KEY_REGISTER(key,ch) │ │ @@ -100,13 +102,13 @@ typedef void (*chat_join_registered_fn)(void* arg, int result); /* Полный набор мембера джойнера (без подписи дерева — её ставит инвайтер). */ struct join_member_data { - uint64_t node_id; + uint64_t node_id; /* derive(x25519), идентичность нового участника */ const uint8_t* x25519; /* 32 */ const uint8_t* ed25519; /* 32 */ const uint8_t* join_sig; /* 64, может быть NULL */ - uint64_t join_ts; + uint64_t join_ts; /* время подписи вступления, Unix seconds */ const uint8_t* update_sig; /* 64, может быть NULL */ - uint64_t update_ts; + uint64_t update_ts; /* версия подписанного блока участника */ 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). */ int chat_join_init(struct UTUN_INSTANCE* inst); -/* Деинициализация. */ +/* Снять обработчик, освободить ключи; ожидающие регистрации получают CANCELLED. */ void chat_join_destroy(struct UTUN_INSTANCE* inst); /* Локальное приглашение (инвайтер == connection). 0=stored, -1=ошибка. */ diff --git a/src/chat/chat_sync.h b/src/chat/chat_sync.h index 6c4fbee3..afda2438 100644 --- a/src/chat/chat_sync.h +++ b/src/chat/chat_sync.h @@ -1,25 +1,10 @@ /* - * chat_sync.h — оркестратор синхронизации чата по ETCP - * - * При поднятии / разрыве ETCP-соединений: - * - обновляет online-статус пиров в БД и GUI - * - запускает member_sync (синхронизацию участников каналов) через merkle_sync - * - управляет кешем каналов (периодический 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) — вызывается при завершении + * chat_sync — связь чат-сервиса с транспортом и группами. + * Отслеживает пиров, запускает member_sync, поддерживает кеш каналов и статус в GUI. + * Вход по ссылке держит отдельный NCD handle до результата join, а не только до UP. + * JOIN_INFO_REQ → JOIN_INFO_RESP → JOIN_REQUEST → JOIN_READY; правила — в chat_join.h. + * Жизненным циклом управляет chat_service_start/stop. Вызовы — в потоке uasync, + * кроме явно отмеченной обёртки, которая копирует аргументы и постит работу в этот поток. */ #ifndef CHAT_SYNC_H @@ -65,15 +50,17 @@ struct UASYNC; /* ── Public API ── */ +/* Создать контекст, обработчики join и member_sync. 0/-1; повтор допустим. */ int chat_sync_init(struct UTUN_INSTANCE* inst); +/* Отменить ожидания join, синхронизацию и таймеры; не удаляет данные каналов. */ 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); /* Подключиться к пиру по данным invite-ссылки (вызывается из uasync-потока). 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, const uint8_t* pubkey_bin, 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). Принимает pubkey и адреса целевого узла. Если соединения ещё нет — строит TOPO_NODE - с адресами, запускает conn_mgr для подключения, и по conn_up отправляет CS_MSG_CHANNEL_INVITE. + с адресами, удерживает NCD для join, и по UP отправляет CS_MSG_CHANNEL_INVITE. Если соединение уже есть — отправляет CHANNEL_INVITE сразу. Копирует pubkey_bin и addrs_data (вызывающий может освободить после возврата). Вызов из любого потока — внутри постит в uasync. */ @@ -112,7 +99,7 @@ void chat_sync_join_channel(struct UTUN_INSTANCE* inst, const char* ch_id, uint6 Вызывается из auto_socket при появлении/изменении сетевой связности. */ 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. */ void chat_sync_on_channel_deleted(struct UTUN_INSTANCE* inst, const char* ch_id); diff --git a/src/chat/db_sync.h b/src/chat/db_sync.h index 3609284b..8edb49f4 100644 --- a/src/chat/db_sync.h +++ b/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) -// -// Назначение: децентрализованная реплицируемая таблица JSON-записей между всеми узлами сети. -// Модуль поддерживает несколько независимых инстансов (таблиц), каждый идентифицируется парой (name, id). -// -// Каждая запись ОБЯЗАТЕЛЬНО содержит Ed25519-подпись автора. -// chain_hash[pos] = SHA256(chain_hash[pos-1] || id || timestamp || author || author_signature[64]) -// Первичный ключ: (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, путь: /sync -// - Синхронизация идёт через etcp_api (service ID 0x20), только прямые P2P соединения +/* db_sync — репликация подписанных записей SQLite по прямым ETCP-соединениям. + * Каждый DB_SYNC_INSTANCE связывает таблицу с group_id. Различия находятся по + * цепочке хешей; новые записи рассылаются PUSH. Ключ и порядок записей: + * (timestamp, author_signature). Подпись автора покрывает timestamp || data. + * Ядро создаёт модуль, чат включает его и регистрирует таблицы каналов. + * Обычно используется общая SQLite ядра; remove отключает таблицу от обмена, + * сохраняя данные. Все вызовы и callbacks выполняются в потоке uasync. */ #ifndef DB_SYNC_H #define DB_SYNC_H @@ -90,32 +58,37 @@ struct DB_SYNC_INSTANCE; // Ed25519 signature size #define DB_SIG_SIZE 64 -// Global lifecycle +// Создать контекст и учесть db_sync_enabled из конфига. 0/-1. int db_sync_init(struct UTUN_INSTANCE* inst); +// Включить ранее созданный модуль независимо от флага конфига. Повтор допустим; 0/-1. int db_sync_enable(struct UTUN_INSTANCE* inst); +// Отменить таймеры/обмен, снять обработчики; общую SQLite не закрывает. 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); +// Освободить экземпляр синхронизации; сама таблица остаётся в БД. 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); // group_id инстанса (для chat-таблиц равен числовому channel_id). Возвращает 0 при si==NULL. uint64_t db_sync_instance_group_id(struct DB_SYNC_INSTANCE* si); // Data operations (per-instance) // 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. 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 char* local_attrs); 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); +// Зарезервировать следующий timestamp перед подписью; растёт даже в пределах одной ms. 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). -// Returns number of records passed to callback. +// Обход по (timestamp, author_signature), начиная с offset; limit=0 — без ограничения. +// Возвращает число записей или -1. Данные callback заимствованы на время вызова. 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 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); // 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), // 0 для live-PUSH и локальных вставок. typedef void (*db_sync_insert_cb)(struct DB_SYNC_INSTANCE* si, diff --git a/src/chat/member_sync.h b/src/chat/member_sync.h index 55688be8..ada046e6 100644 --- a/src/chat/member_sync.h +++ b/src/chat/member_sync.h @@ -1,38 +1,11 @@ /* - * ── member_sync — синхронизация участников чата ── - * - * Тонкая прослойка над merkle_sync, адаптированная под мемберов каналов. - * Данные хранятся в таблицах: nodes, node_addresses, peers_. - * Хеш включает все реплицируемые поля, длины строк и целые в big-endian (см. member_sync_doc.md). - * - * ── Использование ── - * - * // разово: инициализация (из 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); + * member_sync — подписанные записи участников каналов и их обмен через merkle_sync. + * Запись содержит блок участника, блок владельца и подпись дерева приглашений. + * Версии блоков сравниваются отдельно; запись и Merkle-хеш меняются атомарно. + * Локальный put запускает распространение после commit; start ждёт сверки с пиром. + * Участие в CHAT-группе определяется локальной таблицей peers_. + * Сетевые адреса принадлежат topo_node и не входят в блок участника. + * Модуль живёт внутри чат-сервиса; все операции и callbacks — в потоке uasync. */ #ifndef MEMBER_SYNC_H @@ -57,15 +30,17 @@ struct sqlite3; /* UTF-8 JSON bytes, excluding NUL. A complete record fits in one Merkle page. */ #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); -/* Распарсенный рекорд мембера (два подписанных блока). */ +/* Представление записи; все указатели заимствованы на время вызова. */ struct ms_member_rec { - uint64_t node_id; + uint64_t node_id; /* derive(x25519), ключ записи участника */ const uint8_t* x25519; /* 32 байта, обязательно */ const uint8_t* ed25519; /* 32 байта, обязательно */ const uint8_t* join_sig; /* 64 байта, может быть NULL */ - uint64_t join_ts; + uint64_t join_ts; /* эпоха вступления, Unix seconds */ const uint8_t* update_sig; /* 64 байта, может быть NULL */ uint64_t update_ts; /* ver блока мембера */ 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 с коллбэками для мемберов - * (ETCP сервис 0x31), регистрирует _on_node_updated для пересчёта дерева - * при обновлении node info. + * (ETCP_RT_ID_MEMBER_SYNC), проверяет локальные записи, восстанавливает деревья + * и подписывается на события CHAT-групп. 0/-1. */ 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 любого узла (включая себя). * 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); diff --git a/src/chat/merkle_sync.h b/src/chat/merkle_sync.h index 76ed7691..7d882c8d 100644 --- a/src/chat/merkle_sync.h +++ b/src/chat/merkle_sync.h @@ -1,3 +1,8 @@ +/* merkle_sync — обмен различиями между наборами записей по Merkle-дереву. + * Движок сверяет хеши и передаёт страницы изменившихся листьев; модель данных + * через data_ops читает записи, проверяет подписи/права и сохраняет изменения. + * Сессия относится к (пир, namespace), отправка учитывает backpressure. + * Все операции — в потоке uasync; один движок на UTUN_INSTANCE. */ #ifndef MERKLE_SYNC_H #define MERKLE_SYNC_H @@ -19,7 +24,7 @@ struct UTUN_INSTANCE; * SHA256 листа обновляется моделью в порядке unsigned key; родителей считает merkle_tree. * Все операции выполняются в uasync-потоке экземпляра. */ struct merkle_sync_data_ops { - merkle_leaf_hash_fn update_bucket_hash; + merkle_leaf_hash_fn update_bucket_hash; // хеш содержимого одного листа /* Одна страница листа: after_valid=0 начинает обход, иначе только key > after. * len: ёмкость -> размер; next — последний ключ; more — есть продолжение. * Пустой результат разрешён только с more=0. Ошибка не кодируется пустой страницей. */ @@ -28,12 +33,12 @@ struct merkle_sync_data_ops { /* Проверить весь формат, применить записи и обновить дерево в одной DB-транзакции. * <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 (*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); - /* Optional pull transaction. All three hooks must be supplied together. - * abort_pull releases the context after both success and failure. - * Staged records must not affect published hashes before commit_pull. */ + /* Необязательный отложенный приём: все четыре callback задаются вместе. + * stage_page накапливает страницы без изменения опубликованных хешей; + * commit_pull применяет результат; abort_pull освобождает pull и после успеха. */ int (*begin_pull)(void* ctx, const char* ns, uint64_t peer, uint64_t round, void** pull); int (*commit_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); }; +/* result — MT_OK или MT_ERR_*; ns заимствован на время callback. */ 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, const struct merkle_sync_data_ops* ops, void* data_ctx); +/* Отменить сессии без done callbacks, снять обработчики и освободить движок. */ void merkle_sync_destroy(struct UTUN_INSTANCE* inst); /* Подключить namespace к текущему ETCP-соединению и запросить контрольную точку. @@ -56,6 +64,7 @@ void merkle_sync_destroy(struct UTUN_INSTANCE* inst); * Callback может отменить сессию; уничтожение всего экземпляра должно быть отложенным. */ 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); +/* Отменить все сессии namespace / пира без callbacks. */ void merkle_sync_cancel_ns(struct UTUN_INSTANCE* inst, const char* ns); 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. * Возвращает 1=изменилось, 0=no-op, <0=ошибка. */ 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); +/* Прочитать хеш узла дерева в 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]); #endif diff --git a/src/chat/merkle_tree.h b/src/chat/merkle_tree.h index 95dcde65..307d9e11 100644 --- a/src/chat/merkle_tree.h +++ b/src/chat/merkle_tree.h @@ -1,3 +1,7 @@ +/* merkle_tree — производный индекс хешей в SQLite для сравнения записей модели. + * В каждом namespace дерево SHA256: корень (уровень 0), 5 уровней по 32 ветви. + * Лист выбирается по старшим 25 битам ключа; его содержимое хеширует модель. + * Пустое поддерево имеет нулевой хеш. БД принадлежит вызывающему. */ #ifndef MERKLE_TREE_H #define MERKLE_TREE_H @@ -9,16 +13,23 @@ #define MT_BUCKETS 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); -/* Дерево — производный индекс. При открытии очищается и восстанавливается из записей модели. */ +/* Пересоздать пустой индекс; вызывающий затем восстанавливает его из модели. 0/-1. */ 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]); +/* Прочитать 32 дочерних хеша узла уровня 0..4; отсутствующие — нули. 0/-1. */ int merkle_tree_children(sqlite3* db, const char* ns, uint8_t level, uint64_t prefix, 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); +/* Оставить старшие level*5 бит ключа, остальные обнулить; уровень 0 — корень. */ uint64_t merkle_sync_level_prefix(uint64_t key, uint8_t level); +/* Число байтов для старших level*5 бит; 0 для корня или недопустимого уровня. */ uint8_t merkle_sync_prefix_bytes(uint8_t level); #endif diff --git a/src/media_async/media_async.h b/src/media_async/media_async.h index f833ba04..9170d459 100644 --- a/src/media_async/media_async.h +++ b/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 #define MEDIA_ASYNC_H @@ -13,24 +16,29 @@ struct UASYNC; struct media_async; -typedef void (*ma_work_fn)(void* data); -typedef void (*ma_done_fn)(void* arg, int err); +typedef void (*ma_work_fn)(void* data); // рабочий поток; не обращаться к состоянию uasync +typedef void (*ma_done_fn)(void* arg, int err); // 0=work завершён, -1=ошибка запуска/таймера, -2=отмена +/* Создать владельца задач; NULL при ошибке выделения памяти. */ struct media_async* media_async_create(void); #define MEDIA_ASYNC_CANCELLED (-2) -/* Owning event-loop thread only, outside done callbacks. Joins workers, then calls - * each pending done with CANCELLED while caller-owned state is still alive. */ +/* Вызывать вне done. Ждёт работы, вызывает ожидающие done и освобождает ma; NULL допустим. */ 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, ma_work_fn work, void* data, 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_sign_block(const uint8_t* ed25519_privkey, const uint8_t* data, size_t len, uint8_t sig_out[64]); int ma_copy_file(const char* src, const char* dst); +/* Размер файла в байтах (int; для файлов, помещающихся в этот тип), -1 при ошибке stat. */ int ma_file_size(const char* path); +/* Заполнить 16 случайных байт идентификатора. */ void ma_uuid(uint8_t uuid_out[16]); /* Рекурсивный обход каталога: для каждого regular-файла вызывается cb(full_path, size). diff --git a/src/routing_layer/conn_mgr.h b/src/routing_layer/conn_mgr.h index 58808f2c..7b9300db 100644 --- a/src/routing_layer/conn_mgr.h +++ b/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 #define CONN_MGR_H @@ -15,44 +20,9 @@ struct UTUN_INSTANCE; struct ETCP_CONN; 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 { CONN_EVENT_UP = 0, /* связь есть. handle жив, можно отправлять данные */ - CONN_EVENT_DOWN = 1, /* временный обрыв. handle жив, само восстановится. НЕ закрывать */ + CONN_EVENT_DOWN = 1, /* связь потеряна; handle остаётся у владельца, успех восстановления не гарантирован */ CONN_EVENT_TIMEOUT = 2, /* подключение не удалось. handle жив, закройте сами */ }; @@ -69,19 +39,20 @@ typedef void (*conn_mgr_cb_t)(struct CONN_MGR_HANDLE* h, /* ═══════ lifecycle ═══════ */ +/* Создать менеджер группы и его фоновые проверки; NULL при ошибке. */ struct CONN_MGR* conn_mgr_init(struct TOPO_GROUP* group); +/* Остановить проверки и освободить менеджер. Все клиентские handles закрыть ДО destroy. */ 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); /* ═══════ единая точка входа ═══════ */ /** * Единственная публичная функция подключения — 3-фазное (DIRECT→REVERSE→INDIRECT). - * Используется медиа-доставкой. Пиры/группы подключаются через node_conn_direct, - * invite — через topo_group_invite. - * - * Ошибка/таймаут: CONN_EVENT_TIMEOUT. - * node_id и group_id ВСЕГДА валидны в коллбэке (даже при TIMEOUT). + * Возврат 0 — запрос принят, -1 — ошибка запуска. cb обязателен. + * Ошибка до запуска может возвращаться без callback. В callback h может быть NULL. + * Успешно выданный handle закрывает вызывающий, в том числе после TIMEOUT. */ int conn_mgr_open(struct UTUN_INSTANCE* inst, uint64_t group_id, @@ -89,15 +60,16 @@ int conn_mgr_open(struct UTUN_INSTANCE* inst, conn_mgr_cb_t cb, void* cb_arg, struct CONN_MGR_HANDLE** out_handle); -/** Закрыть handle. Последний handle рвёт соединение. */ +/** Закрыть свой handle. Последний освобождает ресурсы CM; чужие NCD-владельцы сохраняют транспорт. */ 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); -/** Отправить данные (для INDIRECT — через посредников). Забирает владение e. */ +/** Передать e->dgram через маршрутизатор CM. Освобождает entry при ненулевом h; + * буфер dgram отдельно не освобождает. При h==NULL entry остаётся у вызывающего. */ int conn_mgr_send(struct CONN_MGR_HANDLE* h, struct ll_entry* e); /* ═══════ мониторинг (вызывается извне) ═══════ */ diff --git a/src/routing_layer/etcp_router.h b/src/routing_layer/etcp_router.h index ee712ca3..8ea480eb 100644 --- a/src/routing_layer/etcp_router.h +++ b/src/routing_layer/etcp_router.h @@ -2,6 +2,9 @@ // Маршрутизирует сервисные пакеты до целевой ноды, на промежуточных нодах ретранслирует // Упрощённый TCP поверх ETCP: восстановление порядка, дедупликация, ретрансмиты, inflight-контроль. // Подпись/шифрование вынесены в автономный модуль route_crypto (encode/decode). +// Логический канал задаётся (group_id, node_id, svc_id), физический транспорт берётся у топологии. +// Сервисы регистрируют приём через bind, отправляют через route_send и соблюдают backpressure. +// Все операции и callbacks — в потоке uasync. Закрытие группы удаляет её каналы и ожидающий транзит. // // Формат 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] @@ -77,8 +80,8 @@ struct TRANSIT_QUEUE { uint64_t group_id; // = data[0..7] uint64_t src_node_id; // = data[8..15] uint64_t dst_node_id; // = data[16..23] - struct ll_queue* q; // FIFO транзитных пакетов (ll_entry с dgram), начиная с data[16] - struct queue_waiter_handle waiter; // backpressure waiter на bp-очередь (normalizer->input или tx_queue) + struct ll_queue* q; // FIFO ещё не переданных транспорту пакетов + struct queue_waiter_handle waiter; // ожидание свободной conn->send_input_q struct ETCP_CONN* conn; // next_hop (для etcp_send в drain_cb) }; @@ -133,9 +136,9 @@ struct ETCP_ROUTER_CONN { uint64_t reset_id; // собственный идентификатор; не меняется при рестарте пира uint64_t peer_reset_id; // подтверждённый идентификатор пира (0 = handshake не завершён) uint64_t pending_peer_id; // кандидат, ещё не имеющий права менять состояние сессии - uint64_t pending_challenge; - uint64_t challenge_sent_tb; - void* handshake_timer; + uint64_t pending_challenge; // случайный запрос подтверждения кандидату + uint64_t challenge_sent_tb; // время отправки запроса, 0.1 ms + void* handshake_timer; // таймер повторов согласования uint8_t start_sent; // получен первый ACK данных 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 int etcp_router_init(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_unbind(struct UTUN_INSTANCE* inst, uint8_t svc_id); // Отправить сервисный пакет (авто-conn, seq, inflight-контроль через send_q). // Принимает владение entry во всех исходах, включая ошибку. +// entry->dgram = svc_id[1] || payload; 0=принято в очередь, <0=ошибка, не подтверждение доставки. +// force разрешает превысить лимит числа пакетов в send_q. // 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); -// Найти/создать состояние seq-подключения по (group_id, remote_node_id, svc_id) +// Найти/создать логический канал. Указатель заимствован до close/destroy, NULL при ошибке. 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); // Есть ли сейчас физический маршрут до узла (прямой, 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); // Отправить данные с авто-seq и контролем inflight // 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, const uint8_t* data, size_t len, int mode); -// Закрыть seq-подключение +// Закрыть канал и уведомить сервис. После вызова rconn использовать нельзя. void etcp_router_conn_close(struct ETCP_ROUTER_CONN* rconn); // Установить рабочий max_inflight, пересчитывает inflight_limit и state machine 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, struct queue_waiter_handle* h, 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, 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-очередей 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_conn_close_all_for_node(struct UTUN_INSTANCE* inst, uint64_t group_id, uint64_t remote_node_id); // Начать новую локальную сессию для конкретного peer+svc в группе и уведомить сервис. diff --git a/src/routing_layer/topo_group.h b/src/routing_layer/topo_group.h index 4b7ce53f..f0339b6c 100644 --- a/src/routing_layer/topo_group.h +++ b/src/routing_layer/topo_group.h @@ -1,32 +1,13 @@ /** * @file topo_group.h - * @brief BGP-подобный обмен топологией узлов между пирами через ETCP. + * @brief Групповые сессии и обмен маршрутами поверх общего транспорта NCD. * - * Идея: модуль автоматически выстраивает карту маршрутизации между узлами, - * используя только пассивное наблюдение (подписка на события) - * Можно иметь много групп узлов. Группа - это список узлов. - * Для каждой группы выстраивается своя независимая таблица маршрутизации. - * Один узел может входить в любое число групп. - - * Механика: - * - Узлы обмениваются информацией друг о друге (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_). База открывается в topo_groups_init() - * если в конфиге задан db_path. + * У каждой группы своя топология; идентичности и адреса узлов общие (topo_node.h). + * UTUN обменивается также подсетями, CHAT допускает только локально известных участников. + * Для подключения используйте peer_open/close: группа владеет NCD с начала попытки. + * JOIN/ACCEPT согласует сессию, NODEINFO/WITHDRAW поддерживают её маршруты. + * Потерю путей обрабатывает topo_recovery, выбор CHAT-пиров — topo_group_connect. + * Все операции — в потоке uasync. Сервисы создают и удаляют свои группы явно. */ #ifndef TOPO_GROUP_H #define TOPO_GROUP_H @@ -170,19 +151,19 @@ struct TOPOMSG_ERR_GROUP_MISMATCH { } __attribute__((packed)); struct TOPO_GROUP_CONN_ITEM { - struct TOPO_GROUP* group; - uint64_t node_id; + struct TOPO_GROUP* group; // владелец сессии + uint64_t node_id; // идентичность пира, включая CONNECTING struct ETCP_CONN* conn; // NULL до UP, в состоянии CONNECTING struct NODE_CONN_DIRECT* handle; // владение BGP этим conn (NCD-handle) - struct TOPO_PEER_REQUEST* requests; + struct TOPO_PEER_REQUEST* requests; // запросы инициаторов; каждый закрывает свой uint64_t progress; // время последнего прогресса обмена (timebase) uint8_t retained; // участие запрошено постоянным владельцем или достигло READY uint64_t local_epoch, peer_epoch; // поколения сторон текущего присоединения uint8_t accepted; // JOIN_ACCEPT подтвердил local_epoch - uint8_t tx_failed, transport_down; - struct topo_tx_item *tx_head, *tx_tail; - struct queue_waiter_handle tx_waiter; - void* tx_wake; + uint8_t tx_failed, transport_down; // ошибка отправки / потеря транспорта + struct topo_tx_item *tx_head, *tx_tail; // FIFO команд, ещё не переданных транспорту + struct queue_waiter_handle tx_waiter; // ожидание свободной send_input_q + void* tx_wake; // отложенный запуск отправки uint8_t table_received; // TABLE_COMPLETE текущей сессии принят 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* nodes; // TOPO_GROUP_NODE{ll_entry,node_id,paths,...} — per-group узлы, хеш-индекс по node_id(8B) 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 struct CONN_MGR* conn_mgr; // менеджер соединений для этой группы struct TOPO_RECOVERY_CTX* recovery; // один последовательный recovery на группу @@ -215,7 +196,7 @@ struct TOPO_GROUP { struct radio_ctx* radio; // radio (PTT walkie-talkie) per-group context uint8_t radio_active; // 1 = я слушаю рацию этого канала (TOPO_FLAG_RADIO в NODEINFO) 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 и конфигом * @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); /** - * @brief Возвращает группу по умолчанию (group_id=TOPO_GROUP_UTUN=0x8000000000000000). + * @brief Возвращает UTUN-группу или NULL, если сервис не запущен. Указатель заимствован. */ 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); @@ -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); /** - * @brief Добавляет conn в senders_list (если нет), отправляет запрос таблицы (nodeinfo). - * - * Вызывается при ETCP on_up. + * @brief Присоединяет существующий транспорт к группе и начинает JOIN/ACCEPT. + * Группа приобретает собственный NCD handle; повтор не сбрасывает сессию. 0/-1. + * Для обычного инициатора используйте отменяемый topo_group_peer_open(). */ 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/удаление группы отсоединяет запрос окончательно. * Запрос остаётся валиден после удаления группы (phase=FAILED), caller обязан * закрыть его. Все операции выполняются в потоке единственного 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); +/* Освободить свой запрос; после вызова request недействителен. NULL допустим. */ 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 */ +/* Завершить всю сессию этой пары: LEAVE, удаление путей, отсоединение всех запросов. */ 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); +/* Время последнего прогресса в единицах 0.1 ms; 0, если сессии нет. */ uint64_t topo_group_peer_progress(const struct TOPO_PEER_REQUEST* request); -/** - * @brief Удаляет conn из senders_list, очищает paths во всех nodes, отправляет withdraw если node unreachable. - * - * Вызывается при ETCP on_down. - */ +/* Причина завершения групповой сессии определяет необходимость recovery. */ enum topo_group_remove_reason { TOPO_REMOVE_TRANSPORT_DOWN, TOPO_REMOVE_LOCAL_LEAVE, @@ -331,25 +314,29 @@ enum topo_group_remove_reason { TOPO_REMOVE_MEMBER_INVALID, 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); /** * @brief Обрабатывает пакет NODEINFO. * - * Проверка версии, обновление или создание TOPO_NODEQ, добавление пути, + * Проверка подписи и timestamp, обновление TOPO_GROUP_NODE, добавление пути, * вставка в роутинг, broadcast если не max hops. * * @return 0 при успехе */ 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); +/* Отсоединить пира и запланировать LEAVE. Чужой handle не забирает. 0/-1. */ int topo_group_leave_peer(struct TOPO_GROUP* group, struct NODE_CONN_DIRECT* handle); /** * @brief Обрабатывает WITHDRAW. * - * Удаляет node из роутинга и nodes, broadcast withdraw. + * Удаляет затронутые пути; узел удаляется только без оставшихся путей. * * @return 0 при успехе */ @@ -369,11 +356,12 @@ void topo_group_send_withdraw(struct TOPO_GROUP* group, uint64_t node_id); /** * @brief Поиск оптимального ETCP соединения для указанного node_id. * - * Перебирает paths узла, выбирает путь с минимальным hop_count. + * Выбирает минимум hop_count среди живых путей. Если живых нет, возвращает + * лучший сохранённый путь: ненулевой результат сам по себе не гарантирует UP. * * @param group указатель на TOPO_GROUP * @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); diff --git a/src/routing_layer/topo_group_connect.h b/src/routing_layer/topo_group_connect.h index 273bcfc4..461dd985 100644 --- a/src/routing_layer/topo_group_connect.h +++ b/src/routing_layer/topo_group_connect.h @@ -1,11 +1,10 @@ -/* CHAT connection policy over cancellable group requests. - * Phase 1: historical connected peers in parallel (up to 128). - * Phase 2: sequential supernode/public peers. Phase 3: local peers. - * Each node is tried once per cycle. CONNECTING: 2s; SYNCING: 5s without progress. - * Goal: one READY supernode or three READY nonmobile peers. After exhaustion, - * retry after 1s (or standby burst). Existing READY group sessions survive - * planner restart/stop. Manual attempts are tracked and cancelled with the group. - * active_count is derived from the current group sessions, not a separate counter. +/* topo_group_connect — выбор пиров и попытки подключения CHAT-группы. + * Сначала ранее подключённые пиры (до 128 параллельно), затем по одному + * суперузлы/публичные узлы, затем локальные. Цель — READY суперузел или три + * READY немобильных пира. CONNECTING: 2 s; SYNCING: 5 s без прогресса. + * После исчерпания кандидатов новый цикл через 1 s или следующий standby burst. + * Планировщик владеет запросами, группа — NCD. Остановка отменяет попытки, + * включая ручные, но сохраняет готовые сессии. Все вызовы — в потоке uasync. */ #ifndef TOPO_GROUP_CONNECT_H #define TOPO_GROUP_CONNECT_H @@ -15,15 +14,22 @@ struct TOPO_GROUP; struct ETCP_CONN; +/* Запустить новый автоматический цикл вместо прежнего. 0/-1. */ int topo_group_connect_init(struct TOPO_GROUP* group); +/* Отменить запросы, таймеры и отложенные вызовы планировщика. */ void topo_group_connect_destroy(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_down(struct TOPO_GROUP* group, struct ETCP_CONN* conn); +/* Число текущих READY-пиров группы. */ int topo_group_connect_active_count(struct TOPO_GROUP* group); +/* Перезапустить автоматический цикл, только если нет READY-пиров. */ 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); #endif diff --git a/src/routing_layer/topo_group_invite.h b/src/routing_layer/topo_group_invite.h index 5ae837d3..37569656 100644 --- a/src/routing_layer/topo_group_invite.h +++ b/src/routing_layer/topo_group_invite.h @@ -1,23 +1,11 @@ /** * @file topo_group_invite.h - * @brief Invite/join к каналу через прямое (ncd) подключение. + * @brief Получение описания канала по прямому NCD-транспорту (INVITE_INFO). * - * Модуль содержит всю логику работы с invite-ссылками: прямое подключение к - * приглашающему узлу через node_conn_direct (ncd) и проверку членства - * (INVITE_INFO_REQ/RESP) с бутстрапом криптографической идентичности канала. - * - * Поток джойнера (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 больше от него - * не зависит. + * Запрашивает и проверяет подписанные сведения канала, готовит локальную группу. + * Успешный ответ не добавляет участника на других узлах и не означает READY группы. + * Вход по пользовательской ссылке и добавление участника описаны в chat_sync.h/chat_join.h. + * Все операции — в потоке uasync; временные запросы отменяются при остановке чата. */ #ifndef 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_EVENT_JOIN 0 -#define TGI_EVENT_TIMEOUT 1 +#define TGI_EVENT_JOIN 0 /* описание канала принято, локальная инфраструктура готова */ +#define TGI_EVENT_TIMEOUT 1 /* попытка завершилась ошибкой или таймаутом */ 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) к приглашающему + проверка членства. * 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, struct TOPO_NODE* ni, uint64_t node_id, tgi_cb_t cb, void* cb_arg); /** - * Инвайтера сторона: прямое подключение (ncd) к целевому узлу для последующей - * отправки CHANNEL_INVITE (через chat_sync, по ETCP_CONN_STATUS_UP). - * Без INVITE_INFO-проверки членства. ni может быть NULL — тогда адреса грузятся - * из БД. Группа владеет NCD; инициатор держит отменяемый групповой запрос. + * Запустить групповую попытку к уже известному участнику, при необходимости + * передав адреса ni (заимствованы; NULL — взять известные). 0/-1. + * Применяется обычная проверка CHAT-членства. Сам CHANNEL_INVITE не отправляет. */ int topo_group_invite_to_channel(struct UTUN_INSTANCE* inst, uint64_t group_id, struct TOPO_NODE* ni, uint64_t node_id); diff --git a/src/routing_layer/topo_node.h b/src/routing_layer/topo_node.h index 5873a007..3249cb69 100644 --- a/src/routing_layer/topo_node.h +++ b/src/routing_layer/topo_node.h @@ -2,21 +2,11 @@ * @file topo_node.h * @brief Модель данных узла сети — структуры, сериализация, wire-формат. * - * Этот модуль определяет: - * - Как узел представлен в памяти (TOPO_NODE, TOPO_GROUP_NODE, адреса, подсети) - * - Как узел передаётся по сети (packed-структуры TOPOMSG_*) - * - Как узел сериализуется/десериализуется (topo_node_serialize/deserialize) - * - * Хранение: - * - Идентичность узла (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 + * TOPO_NODE — общие ключи, имя и адреса; реестр хранит указатель с подсчётом ссылок. + * Подписанная запись заменяется целиком при более свежем timestamp, вместе с подписью. + * TOPO_GROUP_NODE — пути, подсети и состояние узла в конкретной группе. + * TOPOMSG_NODE — сетевой формат. Модуль сериализует, подписывает и проверяет записи. + * Работа с реестром и группами выполняется в потоке uasync экземпляра. */ #ifndef TOPO_NODE_H #define TOPO_NODE_H @@ -87,9 +77,11 @@ struct UTUN_INSTANCE; #define TOPO_GROUP_TYPE_UTUN 1 #define TOPO_GROUP_TYPE_CHAT 2 +/* Результаты проб адресов. *_time — timebase (0.1 ms), *_rtt — те же единицы; + * *_status — PROBE_RESULT_*, probe_status — PROBE_STATUS_*. */ struct TOPO_CONNECTIVITY { uint8_t probe_status; - uint8_t pending_count; + uint8_t pending_count; // незавершённые пробы uint64_t probe_start_time; uint16_t interface_min_rtt; uint16_t nat_min_rtt; @@ -103,7 +95,7 @@ struct TOPO_CONNECTIVITY { uint64_t real_probe_time; uint64_t ping_req_time; uint64_t last_ping_time; - void* probe_list; + void* probe_list; // принадлежащие узлу контексты текущих проб void* probe_deferred_wait; /* контекст отложенной пробы (standby_wait, Android) */ }; @@ -171,17 +163,17 @@ struct TOPO_ADDR_REALITY_OPTS { 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; } -/** Глобальная идентичность узла. Встроен в TOPO_GROUPS->node_registry как queue entry. */ +/** Общая запись узла; node_registry хранит указатель на неё в отдельном ll_entry. */ struct TOPO_NODE { struct ll_entry ll; - uint32_t group_ref_count; + uint32_t group_ref_count; // ссылки групп, ядра и временных пользователей uint64_t node_id; uint8_t ver; uint64_t timestamp; // 0 = неподписанные bootstrap-сведения 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_activity; // CLIENT_ACTIVITY_STANDBY/ACTIVE - char* node_name; + char* node_name; // собственная строка; освобождается с записью struct TOPO_SOCKMETA4* v4_sock_meta; struct TOPO_ADDR4* v4_addrs; struct TOPO_SOCKMETA6* v6_sock_meta; @@ -196,9 +188,9 @@ struct TOPO_NODESUBNETS { struct TOPO_NODEPATH { struct ll_entry ll; - struct ETCP_CONN* conn; - uint8_t hop_count; - uint16_t cumulative_rtt; + struct ETCP_CONN* conn; // заимствованный next-hop; транспорт удерживает сессия + uint8_t hop_count; // число node_id в массиве сразу после структуры + uint16_t cumulative_rtt; // RTT оставшейся цепочки, без локального линка (0.1 ms) }; /** Узел в контексте группы (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 uint8_t conn_mgr_intermediariy_count; // количество посредников struct TOPO_CONNECTIVITY connectivity; // результаты ping-проб (interface/nat/real) - uint8_t conn_presence; // NODE_CONN_* — какие подключения есть в принципе - uint8_t conn_up; // NODE_CONN_* — какие из них подняты + uint8_t conn_presence; // маска NCONN_*: известные виды подключения + uint8_t conn_up; // маска NCONN_*: работающие виды подключения uint8_t radio; // 1 = узел слушает рацию канала (TOPO_FLAG_RADIO) - struct NODE_CONN_DIRECT* handle; // активный ncd-handle прямого соединения к узлу + struct NODE_CONN_DIRECT* handle; // не владеет транспортом; NCD принадлежит сессии группы }; // API — глобальный реестр TOPO_NODE +/* Заимствованный указатель; ref нужен, если запись должна пережить удаление из группы. */ 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); +/* Добавить/снять ссылку; последняя unref удаляет запись из реестра и освобождает её. */ 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); +/* Освободить отдельную запись и её списки; зарегистрированную запись отпускать через unref. */ void topo_node_destroy(struct TOPO_GROUPS* groups, struct TOPO_NODE* ni); // API — per-group +/* Заимствованный узел группы, NULL если не найден. */ 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); /* Полное удаление узла из группы: отменяет влётные пробы, освобождает paths/subnets, * убирает из group->nodes и освобождает сам 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); +/* Собрать NODEINFO в out; длина или -1. Входные структуры остаются у вызывающего. */ int topo_node_serialize(struct TOPO_NODE* ni, struct TOPO_GROUP_NODE* nq, uint64_t group_id, uint8_t flags, 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, uint64_t** out_hop_list, uint8_t* out_hop_count, 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); +/* Обновить локальный узел группы и его анонс. 0/-1. */ 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_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_dump_all(struct TOPO_GROUP* group); 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); +/* Обновить время ping и минимальные RTT во всех группах и БД; rtt в единицах 0.1 ms. */ 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); -/** Build canonical message for Ed25519 signature: x25519_pubkey || name || client_type || client_activity || addresses */ +/* Лимиты буфера подписываемых данных и полного сетевого NODEINFO. */ #define TOPO_SIG_MSG_MAX_SIZE 2048 #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); -/** 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); +/* Проверить timestamp, привязку node_id к ключу и подпись. 0/-1. */ 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; } diff --git a/src/routing_layer/topo_recovery.h b/src/routing_layer/topo_recovery.h index 78e370d5..206d6ccb 100644 --- a/src/routing_layer/topo_recovery.h +++ b/src/routing_layer/topo_recovery.h @@ -1,3 +1,6 @@ +/* topo_recovery — восстановление путей группы после потери транспорта. + * Собирает потерянные узлы, пробует кандидатов через групповые запросы и ждёт + * нужных маршрутов через READY-пиров. Контекст принадлежит группе; работа — в uasync. */ #ifndef TOPO_RECOVERY_H #define TOPO_RECOVERY_H @@ -30,9 +33,13 @@ struct TOPO_RECOVERY_CTX; * NODEINFO/смена состояния лишь планируют проверку через call_soon: текущий * BGP callback завершается до изменения попытки или освобождения контекста. * Остановка группы отменяет 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_start(struct TOPO_GROUP* group); +/* Отложить проверку прогресса; повторные уведомления объединяются. */ void topo_recovery_changed(struct TOPO_GROUP* group); +/* Отменить попытку и освободить контекст; готовые сессии остаются у группы. */ void topo_recovery_cancel_all(struct TOPO_GROUP* group); #ifdef __cplusplus diff --git a/src/transport_layer/etcp.h b/src/transport_layer/etcp.h index c8c60008..53b8c0dd 100644 --- a/src/transport_layer/etcp.h +++ b/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 #define ETCP_H @@ -14,8 +18,6 @@ extern "C" { #include "pkt_normalizer.h" struct stcp_link; // forward declaration -// In struct ETCP_CONN, add: -//struct pn_pair* normalizer; // Forward declarations struct UTUN_INSTANCE; @@ -23,12 +25,13 @@ struct ETCP_CONN; struct etcp_cbk_entry; // defined in etcp_api.h struct UASYNC; +/* Младшие 16 бит локального времени в единицах 0.1 ms (циклический timestamp). */ uint16_t get_current_timestamp(void); // ETCP packet section types (from protocol spec) #define ETCP_SECTION_PAYLOAD 0x00 // Data payload #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_RESP 0x08 // Measurement response (burst result) #define ETCP_SECTION_FILLER 0x09 // Filler/dummy data (discarded by receiver) @@ -99,8 +102,8 @@ struct ACK_PACKET { // ETCP connection structure (refactored) struct ETCP_CONN { - // State: 0=not ready, 1=ready (indexed in instance->connections), 2=deleted - int state; // 0=pending, 1=ready, 2=deleted (phase 1 of close done) + // state описывает объект в реестре; доступность транспорта определяется links_up. + int state; // 0=pending, 1=индексирован по peer_node_id, 2=удалён (phase 1 завершена) int ref_count; // External reference count. >0 blocks deferred resource free. // Take/free via etcp_conn_ref_take()/etcp_conn_ref_free(). @@ -108,9 +111,9 @@ struct ETCP_CONN { 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_queue* conn_queue; // указатель на очередь где лежим (pending или connections) + struct ll_queue* conn_queue; // заимствованная очередь, содержащая conn_queue_entry // Links (channels) - linked list struct ETCP_LINK* links; @@ -146,7 +149,7 @@ 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* 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) @@ -198,8 +201,8 @@ struct ETCP_CONN { uint8_t initialized; // 0 - только созданный ETCP, 1 - хотя бы один линк проинициалзирован (обмен ключами произведен) uint64_t reset_id; // локальная эпоха потока (меняется только при локальном reset) uint64_t peer_reset_id; // подтверждённая эпоха потока пира - uint64_t session_candidate, session_challenge, session_started; - void* session_timer; + uint64_t session_candidate, session_challenge, session_started; // эпоха-кандидат, challenge, время начала (0.1 ms) + void* session_timer; // таймер согласования транспортной сессии 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 links_up; // 0 - канал не готов для передачи, 1 - канал готов для передачи (хотя бы один линк не down) @@ -207,10 +210,10 @@ struct ETCP_CONN { int callbacks_running; // счётчик вложенности колбэк-цепочек (>0 — внутри цепочки, etcp_connection_close откладывается) uint8_t close_requested; // новые операции запрещены, phase 1 может ждать выхода из callback - uint8_t close_running; - uint8_t delete_complete; - uint8_t free_scheduled; - void* close_token; + uint8_t close_running; // защита от повторного входа в close + uint8_t delete_complete; // detach и рассылка DELETE завершены + uint8_t free_scheduled; // освобождение уже поставлено в uasync + void* close_token; // отложенный close при вызове внутри callback // Unified callback chain with event mask (init/reinit/up/down/node_changed) struct etcp_cbk_entry* cbks; @@ -235,30 +238,28 @@ struct ETCP_CONN { #define RTT_CB_PERIOD_TB 600000 // 1 минута в 0.1ms -// Functions +// Создать транспортный объект; жизненным циклом прикладного соединения управляет NCD. struct ETCP_CONN* etcp_connection_create(struct UTUN_INSTANCE* instance, char* name); /** * @brief Закрыть соединение: фаза-1 detach + отложенная фаза-2 очистка. * * Порядок коллбэков при закрытии (строгая последовательность): - * 1) если соединение ещё UP (links_up != 0) — финальный per-connection DOWN - * (ETCP_CBK_EVENT_DOWN), затем links_up = 0; + * 1) если соединение ещё UP — links_up = 0 и финальный per-connection DOWN; * 2) state = 2; - * 3) conn-status DELETE (ETCP_CONN_STATUS_DELETE) — единственный статус при - * state==2 (по нему подписчики снимают conn из своих списков); + * 3) per-connection DELETE, затем instance-level DELETE; * 4) cleanup: таймеры, routing, линки, очередь, deferred u_free. * - * После state==2 per-connection события (etcp_cbk_fire) и conn-status - * NEW/UP/DOWN не вызываются. + * При вызове внутри callback закрытие откладывается, close_requested ставится сразу. + * После state==2 допустим только DELETE; ref_count задерживает освобождение памяти. */ 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). * @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() * will detach but defer resource cleanup until all references are released. diff --git a/src/transport_layer/etcp_api.h b/src/transport_layer/etcp_api.h index 5c21b2e5..772bfc13 100644 --- a/src/transport_layer/etcp_api.h +++ b/src/transport_layer/etcp_api.h @@ -1,17 +1,11 @@ /** * @file etcp_api.h - * @brief API для приёма-передачи пакетов через ETCP (per-instance bindings) - * - * Основные функции: - * - etcp_send() - отправить пакет в очередь normalizer - * - etcp_bind() - подписаться на пакеты с определенным ID - * - etcp_int_recv() - коллбэк для сбора пакетов из всех подключений - * !!! для обычной отправки-приёма между узлами используем более универсальный модушь etcp_router. - * это в первую очередь - апи для использования etcp-router-ом - * - * Формат кодограмм: - * cmd = 0 - пакет для передачи адресату - * cmd = 1 - модуль обмена роутинг-таблицами + * @brief Прямой обмен кодограммами ETCP и подписки на события транспорта. + * + * bind выбирает обработчик по первому байту cmd, send передаёт кодограмму соседу. + * Для доставки через промежуточные узлы используется etcp_router.h. + * Все вызовы — в потоке uasync. Перед отправкой ждать UP и свободной send_input_q. + * Владение entry различается: send забирает его только при успехе, recv — всегда. */ #ifndef ETCP_API_H @@ -86,11 +80,12 @@ extern "C" { #define ETCP_RT_ID_CALL 0x36 // call — P2P аудио-звонок // Connection status events (instance-level callback). -// Строгая последовательность: NEW -> UP <-> DOWN -> DELETE. +// NEW начинает жизнь объекта, UP/DOWN отражают доступность, DELETE завершает её. // NEW — один раз при создании соединения. // UP — переход links_up 0→1 (UP после UP не вызывается). // DOWN — переход links_up 1→0 (DOWN после DOWN не вызывается). // DELETE — один раз при закрытии, финальный; единственный статус, допустимый при state==2. +// REINIT — новая транспортная сессия; не означает потерю физического пути. #define ETCP_CONN_STATUS_NEW 0 // соединение создано #define ETCP_CONN_STATUS_UP 1 // соединение поднялось #define ETCP_CONN_STATUS_DOWN 2 // соединение упало @@ -115,16 +110,10 @@ struct etcp_status_cbk_entry { * @param conn ETCP соединение от которого получен пакет * @param entry Элемент очереди с данными пакета * - * @note Коллбэк должен освободить entry через queue_entry_free() - * и dgram через queue_dgram_free() после обработки - * - * @note Формат кодограммы для etcp_route_send: - * 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]. + * @note После обработки освободить сначала dgram через queue_dgram_free(), + * затем entry через queue_entry_free(). conn заимствован на время вызова. + * @note Raw bind получает cmd || payload. Router bind получает заголовок + * ROUTER_SVC_* и payload со смещения ROUTER_SVC_PAYLOAD_OFF (etcp_router.h). */ typedef void (*etcp_recv_fn)(struct ETCP_CONN* conn, struct ll_entry* entry); @@ -134,7 +123,7 @@ struct UTUN_INSTANCE; struct TOPO_GROUP_NODE; // Per-connection callback event types (bitmask). -// Вызываются только при state != 2 (после DELETE события не идут). +// При state==2 допустим только терминальный ETCP_CBK_EVENT_DELETE. // UP/DOWN зеркалят conn-status UP/DOWN; при закрытии (если соединение было UP) // отправляется один финальный DOWN до state=2. #define ETCP_CBK_EVENT_INIT (1 << 0) @@ -160,8 +149,7 @@ struct etcp_cbk_entry { * * @note Вызываются все коллбэки соединения, у которых event_mask * пересекается с переданным event. - * @note Не вызывается при state==2 (соединение удалено) — после DELETE - * per-connection события не рассылаются. + * @note При state==2 допускается только DELETE; прочие события игнорируются. */ 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 id Идентификатор пакета - * @return 0 при успехе, -1 если binding не найден + * @return 0 при успехе (в том числе без прежнего binding), -1 при неверном instance */ 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 Пользовательский аргумент * * @note status = ETCP_CONN_STATUS_NEW / UP / DOWN / DELETE / REINIT. - * Строгая последовательность NEW -> UP <-> DOWN -> DELETE; DELETE — - * единственный статус, отправляемый при state==2 (при закрытии). + * DELETE — единственный статус при state==2. Финальный DOWN при close + * доставляется per-connection подписчикам; instance-подписчики получают DELETE. */ void etcp_add_conn_status_cbk(struct UTUN_INSTANCE* inst, etcp_conn_status_fn fn, void* arg); diff --git a/src/transport_layer/node_conn_direct.h b/src/transport_layer/node_conn_direct.h index 3de79a3c..25303a78 100644 --- a/src/transport_layer/node_conn_direct.h +++ b/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). - * Несколько handle'ов могут разделять одно ETCP_CONN — закрытие происходит только - * когда все handle'ы закрыты (refcounting). При закрытии последнего handle - * используется graceful shutdown: CLOSE/KEEP_ALIVE протокол (fin_wait 5 сек). + * Каждый владелец открывает свой handle; несколько handle разделяют ETCP_CONN. + * close снимает только своё владение. Последний handle закрывает неготовый транспорт + * сразу; для UP начинает CLOSE/KEEP_ALIVE (до 5 s), пир может сохранить соединение. * * Кто использует: * - conn_mgr (Connection Manager) — создаёт handle'ы для DIR/REV/IND-соединений - * - chat_core / topo_group — связь с узлами канала/группы + * - topo_group — владеет транспортом групповой сессии с начала CONNECTING + * - протокол join — временный транспорт до добавления участника * - любой модуль, кому нужно надёжное ETCP-соединение с конкретным node_id * * Переходы UP/DOWN/TIMEOUT/CLOSED доставляются синхронно. Закрытие любого handle * из callback безопасно. Начальное UP при open готового conn откладывается через uasync_call_soon. + * Все операции — в потоке uasync. NCD не проверяет членство и готовность группы. */ +#ifndef NODE_CONN_DIRECT_H +#define NODE_CONN_DIRECT_H #include @@ -48,7 +49,8 @@ typedef void (*ncd_callback)(struct NODE_CONN_DIRECT* h, enum ncd_event event, v * Node info ищется через node_registry, fallback — SQLite. * cb вызывается при изменении статуса / таймауте. * Если 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, 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 ETCP_SOCKET* specific_sock); /* NULL — авто-подбор */ +/* Освободить handle и отменить его callbacks; другие владельцы сохраняют транспорт. */ 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); -/* Force close без fin_wait: немедленно удаляет ETCP-коллбэки и, если это - * последний handle, закрывает соединение и освобождает ncd_entry. - * Используется при уничтожении CM entry — гарантирует, что после вызова - * никакие NCD-коллбэки не доставят события в освобождённую память. */ +/* Как close, но последний handle закрывает транспорт сразу, без fin_wait. + * При наличии других handle общий транспорт и их callbacks сохраняются. */ void node_conn_direct_force_close(struct NODE_CONN_DIRECT* h); /* Сменить или сбросить (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); +/* Идентификатор узла; 0 для NULL. */ 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); /* ─── Протокол CLOSE / KEEP_ALIVE (ETCP_RT_ID_NCD_CONTROL = 0x12) ─── */ diff --git a/src/utun_instance.h b/src/utun_instance.h index ce2393be..ffaf2c85 100644 --- a/src/utun_instance.h +++ b/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 #define UTUN_INSTANCE_H @@ -75,11 +80,11 @@ struct conn_queue_entry { struct ETCP_CONN* conn; }; -// Handle for config-based connections via node_conn_direct +// Постоянный запрос UTUN-сервиса к пиру из конфигурации; транспортом владеет группа. struct CONFIG_CONN_HANDLE { uint64_t node_id; char name[MAX_CONN_NAME_LEN]; - struct TOPO_PEER_REQUEST* request; + struct TOPO_PEER_REQUEST* request; // принадлежит этой записи; закрывается при stop/reload struct CONFIG_CONN_HANDLE* next; }; @@ -109,23 +114,23 @@ struct peer_sleep_cbk_entry { struct peer_sleep_cbk_entry* next; }; -// uTun instance configuration +// Состояние одного узла. Вложенные модули освобождаются через их lifecycle API. struct UTUN_INSTANCE { - uint8_t core_started, utun_started, chat_started; + uint8_t core_started, utun_started, chat_started; // успешный запуск ядра / UTUN / чата // Identification char name[MAX_CONN_NAME_LEN]; // Instance name from config - // Configuration (moved from utun_state) + // Конфигурация принадлежит экземпляру после успешного create. struct utun_config *config; // TUN interface struct tun_if* tun; - // Route subnets (for cleanup on shutdown) + // Заимствованный список системных маршрутов из config; нужен для stop UTUN. struct CFG_ROUTE_ENTRY* route_subnets; - struct ROUTE_TABLE* rt; + struct ROUTE_TABLE* rt; // таблица маршрутов подсетей struct TOPO_GROUPS* topo_groups; // Groups module for topology exchange sqlite3* topo_sqlite_db; // Shared SQLite DB (nodes/channels/peers) 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) - // Main async context + // Заимствованный цикл событий: создаёт и уничтожает вызывающий. struct UASYNC* ua; // State - int running; + int running; // продолжать utun_instance_run; не флаг готовности сервисов // Connections (очередь всех подключений для instance) 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) struct memory_pool* data_pool;// для входных-выходных данных пакета - struct memory_pool* pkt_pool; - struct memory_pool* ack_pool; + struct memory_pool* pkt_pool; // ETCP_DGRAM и буферы сетевых пакетов + struct memory_pool* ack_pool; // ACK_PACKET // Active sockets (UDP + TCP, is_tcp=1 flag) 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) struct ETCP_ROUTER_BINDINGS router_bindings; 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) - // Chat/DM subsystem (per-instance contexts; were global singletons) - struct chat_core_ctx* chat_core; // chat_core.c (was g_cc) - struct chat_sync* chat_sync; // chat_sync.c (was g_cs) - struct join_key_entry* join_keys; // chat_join.c (was g_keys; собственный init/destroy) - struct chat_join_registration* join_registrations; - int join_stopping; - struct chat_invite_build_req* invite_gui_request; - struct dm_state* dm; // dm/dm_core.c (was g_dm) - struct dm_mb_state* dm_mailbox; // dm/dm_mailbox.c (was g_mb) + // Контексты чат-сервиса; освобождаются при chat_service_stop. + struct chat_core_ctx* chat_core; // каналы, сообщения и доступ к общей БД + struct chat_sync* chat_sync; // join и синхронизация каналов + struct join_key_entry* join_keys; // зарегистрированные ключи приглашений + struct chat_join_registration* join_registrations; // ожидания KEY_REGISTER_ACK + int join_stopping; // запрещает новые регистрации при teardown + struct chat_invite_build_req* invite_gui_request; // текущее построение ссылки для GUI + struct dm_state* dm; // личные сообщения + struct dm_mb_state* dm_mailbox; // хранилище личных сообщений для получателей struct call_ctx* call; // call/call.c — P2P audio call 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_headless* radio_headless; // radio/radio_headless.c — audio stream socket - struct chat_setting_state chat_settings; // chat_setting.c (was globals) - chat_event_handler_fn chat_event_handler; // chat_event.c (was g_handler) - void* headless; // chat_headless_control.c (was g_hc) + struct chat_setting_state chat_settings; // текущие настройки чата + chat_event_handler_fn chat_event_handler; // получатель событий чата + void* headless; // контекст TCP-управления headless-чатом // E2E encryption cache — per-peer sc_context_t with derived session key #define E2E_CTX_CACHE_SIZE 8 @@ -253,7 +258,7 @@ struct UTUN_INSTANCE { struct NTP_TIME ntp; struct NTP_NODE_TIME ntp_node; - // Config-based connection handles (node_conn_direct) + // Постоянные групповые запросы UTUN к пирам из конфигурации. struct CONFIG_CONN_HANDLE* config_conn_handles; // 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) сделан }; -// Functions +// Создать ядро по файлу конфигурации; NULL при ошибке. Сервисы ещё не запущены. 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); +// Как create, но конфигурация передана строкой; строку не сохраняет. 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); -/* create constructs the core; start enables its runtime. Services are independent. - * All lifecycle calls run on the owning uasync thread, outside service callbacks. - * Stop is idempotent. destroy stops both services before releasing core resources. */ +/* Запустить общие сетевые службы ядра. Повтор допустим; 0/-1. */ int utun_core_start(struct UTUN_INSTANCE* instance); +/* Запустить UTUN-группу, TUN/NAT по конфигу и запросы к пирам. Нужен core_start; 0/-1. */ int utun_service_start(struct UTUN_INSTANCE* instance); +/* Остановить UTUN, сохранив ядро и чат. Повтор допустим. */ void utun_service_stop(struct UTUN_INSTANCE* instance); -/* Convenience: core + UTUN + configured chat. */ +/* Запустить ядро, UTUN и включённый в конфиге чат. 0/-1. */ 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); +/* Выполнять цикл событий до stop(). */ void utun_instance_run(struct UTUN_INSTANCE *instance); +/* Снять running и разбудить цикл; ресурсы освобождает destroy(), а не stop(). */ void utun_instance_stop(struct UTUN_INSTANCE *instance); +/* Глобальные переключатели создания TUN/топологии для тестов и встраивания. */ void utun_instance_set_tun_init_enabled(int enabled); void utun_instance_set_topo_group_enabled(int enabled); 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 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. - Returns ETCP_CONN* ready for etcp_send(), or NULL if not found. */ +/* Найти транспорт по node_id в реестрах ETCP/TCP. Указатель заимствован. + * NULL — не найден; ненулевой результат не гарантирует UP или готовность группы. */ static inline struct ETCP_CONN* instance_find_conn(struct UTUN_INSTANCE* inst, uint64_t node_id) { if (inst->connections) { struct ll_entry* e = queue_find_data_by_index(inst->connections, (const uint8_t*)&node_id);