Browse Source

Separate core, UTUN and chat lifecycles with owned asynchronous work

proxy
evgeny 3 days ago
parent
commit
72fee9e245
  1. 59
      doc/service_lifecycle.md
  2. 4
      src/chat/chat_channel.c
  3. 2
      src/chat/chat_core.c
  4. 5
      src/chat/chat_core.h
  5. 1
      src/chat/chat_core_priv.h
  6. 12
      src/chat/chat_msg.c
  7. 1
      src/chat/chat_sync.c
  8. 2
      src/chat/chat_whisper.c
  9. 53
      src/chat/db_sync.c
  10. 1
      src/chat/db_sync.h
  11. 93
      src/media_async/media_async.c
  12. 3
      src/media_async/media_async.h
  13. 90
      src/media_delivery/media_delivery.c
  14. 4
      src/media_delivery/media_delivery.h
  15. 2
      src/media_delivery/media_index.c
  16. 30
      src/routing_layer/etcp_router.c
  17. 1
      src/routing_layer/etcp_router.h
  18. 2
      src/routing_layer/routing.c
  19. 43
      src/routing_layer/topo_group.c
  20. 4
      src/routing_layer/topo_group.h
  21. 7
      src/routing_layer/topo_group_invite.c
  22. 2
      src/routing_layer/topo_group_invite.h
  23. 22
      src/routing_layer/topo_node.c
  24. 1
      src/routing_layer/topo_node.h
  25. 15
      src/transport_layer/etcp_connections.c
  26. 258
      src/utun_instance.c
  27. 9
      src/utun_instance.h
  28. 32
      src/video/video.c
  29. 7
      src/video/video.h
  30. 5
      tests/Makefile.am
  31. 25
      tests/test_chat_join_e2e.c
  32. 2
      tests/test_conn_mgr_phases.c
  33. 4
      tests/test_etcp_100_packets.c
  34. 4
      tests/test_etcp_api.c
  35. 4
      tests/test_etcp_ping.c
  36. 4
      tests/test_etcp_simple_traffic.c
  37. 2
      tests/test_group_recovery.c
  38. 4
      tests/test_invite_group_create.c
  39. 1
      tests/test_media_delivery_chat.c
  40. 1
      tests/test_media_delivery_full.c
  41. 1
      tests/test_media_delivery_integration.c
  42. 5
      tests/test_nat_detection.c
  43. 4
      tests/test_pkt_normalizer_etcp.c
  44. 2
      tests/test_route_ping.c
  45. 173
      tests/test_services.c
  46. 12
      tools/chatgui-android/libutun_lite/instance_lite.c
  47. 15
      tools/chatgui/transport/utun_node.cpp

59
doc/service_lifecycle.md

@ -0,0 +1,59 @@
# Жизненный цикл ядра и сервисов
Все операции выполняются в потоке единственного `UASYNC` экземпляра, вне callbacks
останавливаемых модулей. GUI передаёт команды в этот поток.
```c
struct UTUN_INSTANCE* inst = utun_instance_create_from_str(ua, config);
if (!inst || utun_core_start(inst) < 0) /* обработать ошибку */;
if (chat_service_start(inst) < 0) /* обработать ошибку */;
if (utun_service_start(inst) < 0) /* обработать ошибку */;
chat_service_stop(inst); /* UTUN продолжает работать */
chat_service_start(inst); /* загружает сохранённые каналы */
utun_service_stop(inst); /* чат продолжает работать */
utun_instance_destroy(inst);
```
`utun_instance_init()` — сокращение для запуска ядра, UTUN и чата, если чат включён
в конфигурации. Desktop GUI и Android запускают ядро и чат явно, без UTUN.
Повторный start работающего сервиса и повторный stop допустимы. Ошибка старта
сервиса откатывает созданные им ресурсы. `utun_instance_stop()` останавливает цикл
приложения; освобождение ядра выполняет `utun_instance_destroy()`.
| Владелец | Ресурсы |
|---|---|
| Ядро | Идентичность узла, реестр узлов и групп, SQLite, транспортные сокеты, NCD, маршрутизатор, общая репликация, время, мониторинг |
| UTUN | Группа UTUN, запросы соединений из конфигурации, DATA binding, TUN, системные маршруты, NAT transport |
| Чат | CHAT-группы, запросы подключения и join, экземпляры синхронизации каналов, chat API, медиа, личные сообщения, звонки и headless API |
| Группа | Сессии пиров, её маршруты и логические каналы, планировщик подключения и recovery |
| Сессия пира группы | NCD handle с момента CONNECTING, согласованные эпохи, состояние обмена таблицей, очередь отправки |
SQLite открывается ядром один раз: `db_path/chats.db`, либо `:memory:`, если путь
не задан. Остановка чата сохраняет каналы и сообщения в этой БД. В варианте
`:memory:` данные сохраняются только до уничтожения ядра.
Общий транспорт живёт, пока есть его владельцы. Закрытие запроса одной группы
не разрывает соединение, нужное другой группе. При остановке группы отменяются
её планировщики, запросы, таймеры и ожидания очередей; отправляется LEAVE,
закрываются логические каналы маршрутизатора и удаляется её ожидающий транзит
на общих соединениях. Транзит других групп сохраняется. Отложенный LEAVE может временно
удерживать транспорт самостоятельно, не обращаясь к уже удалённой группе.
Остановка чата сначала прекращает внешние команды и фоновые задания, затем
закрывает медиа и сервисы, синхронизацию и группы, после чего освобождает chat core.
`media_async_destroy()` дожидается завершения native workers и вызывает каждый
ожидающий callback с `MEDIA_ASYNC_CANCELLED`, пока его контекст ещё существует.
Результаты отменённых задач не записываются в БД. Это синхронное ожидание: длительное
транскодирование или распознавание речи может задержать stop; принудительного
прерывания этих библиотек нет.
Статус транспорта UP подтверждает восстановление физического пути. READY сессии
группы означает согласованный JOIN/ACCEPT и завершённый обмен её таблицей; именно
его ждёт recovery. REINIT меняет эпохи обмена, сохраняя рабочие маршруты. Эпохи
защищают от сообщений старой сессии и не задают срок жизни маршрута.
Проверки: `test_services` проверяет независимый stop/start, общий NCD в CONNECTING,
сохранение идентичности/сокетов/БД, повторную загрузку каналов, отмену workers,
заблокированного медиа-потока, транзита и ожиданий заменённой очереди. `test_chat_join_e2e`
проверяет join и распространение membership между узлами без UTUN-сервиса.

4
src/chat/chat_channel.c

@ -72,9 +72,7 @@ void chat_core_ensure_channel_ready(struct UTUN_INSTANCE* inst, const char* ch_i
}
if (si) {
si_register(inst, si, ch_id);
struct msg_insert_arg* ia = u_malloc(sizeof(*ia));
if (ia) { ia->inst = inst; snprintf(ia->ch_id, sizeof(ia->ch_id), "%s", ch_id); }
db_sync_set_insert_cb(si, on_msg_inserted, ia);
db_sync_set_insert_cb(si, on_msg_inserted, inst);
db_sync_instance_set_gated(si, 1);
DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "%s: channel ready ch=%s tbl=%s group=0x%016llx",
CC_ID, ch_id, tbl_msg, (unsigned long long)gid);

2
src/chat/chat_core.c

@ -80,6 +80,7 @@ void si_register(struct UTUN_INSTANCE* inst, struct DB_SYNC_INSTANCE* si, const
int chat_core_init(struct UTUN_INSTANCE* inst, const char* db_path) {
if (!inst || !db_path) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "%s: invalid args inst=%p db_path=%p", CC_ID, (void*)inst, (void*)db_path); return -1; }
if (inst->chat_core) return 0;
struct chat_core_ctx* cc = u_calloc(1, sizeof(*cc));
if (!cc) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "%s: OOM for chat_core_ctx", CC_ID); return -1; }
inst->chat_core = cc;
@ -235,6 +236,7 @@ void chat_core_destroy(struct UTUN_INSTANCE* inst) {
struct chat_core_ctx* cc = CC(inst);
if (!cc || !cc->initialized) return;
cc->initialized = 0;
member_sync_destroy(inst); /* Also releases callback lists after partial startup. */
if (cc->db && !cc->shared_db) { sqlite3_close(cc->db); }
inst->chat_core = NULL;

5
src/chat/chat_core.h

@ -19,6 +19,11 @@ 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. */
int chat_service_start(struct UTUN_INSTANCE* inst);
void chat_service_stop(struct UTUN_INSTANCE* inst);
int chat_core_init(struct UTUN_INSTANCE* inst, const char* db_path);
void chat_core_destroy(struct UTUN_INSTANCE* inst);

1
src/chat/chat_core_priv.h

@ -113,7 +113,6 @@ void si_register(struct UTUN_INSTANCE* inst, struct DB_SYNC_INSTANCE* si, const
/* ── db_sync callback (chat_msg.c) ── */
struct msg_insert_arg { struct UTUN_INSTANCE* inst; char ch_id[64]; };
void on_msg_inserted(struct DB_SYNC_INSTANCE* si, uint64_t record_ts, const char* data, size_t len, uint64_t author, int initial_sync, void* arg);
/* ── применить настройку group_autoconnect ко всем CHAT-группам (chat_channel.c) ── */

12
src/chat/chat_msg.c

@ -342,8 +342,8 @@ void chat_core_submit_media_message(struct UTUN_INSTANCE* inst, struct chat_msg_
}
if (req->media_video) {
video_transcode_start(inst->ua, req->media_src, req->media_dest,
on_video_prepared, mctx);
if (video_transcode_start(inst->media_async, inst->ua, req->media_src, req->media_dest, on_video_prepared, mctx) < 0)
on_video_prepared(mctx, -1);
return;
}
@ -842,12 +842,10 @@ static int md_auto_download(struct UTUN_INSTANCE* inst,
/* ─── db_sync callback ─── */
void on_msg_inserted(struct DB_SYNC_INSTANCE* si, uint64_t record_ts, const char* data, size_t len, uint64_t author, int initial_sync, void* arg) {
(void)si;
struct msg_insert_arg* ia = (struct msg_insert_arg*)arg;
struct UTUN_INSTANCE* inst = ia ? ia->inst : NULL;
struct UTUN_INSTANCE* inst = arg;
struct chat_core_ctx* cc = CC(inst);
if (!ia || !cc) return;
const char* ch_id = ia->ch_id;
if (!cc) return;
char ch_id[32]; snprintf(ch_id, sizeof(ch_id), "%llu", (unsigned long long)db_sync_instance_group_id(si));
uint8_t evt[73]; uint8_t cl = (uint8_t)strlen(ch_id); evt[0] = cl; memcpy(evt + 1, ch_id, cl);
memcpy(evt + 1 + cl, &author, 8);
chat_event_post(inst, CHAT_EVT_MSG_RECEIVED, evt, 1 + cl + 8);

1
src/chat/chat_sync.c

@ -740,6 +740,7 @@ static void refresh_timer_cb(void* arg) {
int chat_sync_init(struct UTUN_INSTANCE* inst) {
if (!inst) return -1;
if (inst->chat_sync) return 0;
DEBUG_TRACE(DEBUG_CATEGORY_MEMBER_SYNC, "%s: init", CS_ID);
struct chat_sync* cs = u_calloc(1, sizeof(*cs));
if (!cs) return -1;

2
src/chat/chat_whisper.c

@ -345,8 +345,8 @@ static void wh_work_fn(void* raw) {
}
static void wh_done_fn(void* raw, int err) {
(void)err;
struct wh_work_ctx* w = (struct wh_work_ctx*)raw;
if (err) w->err = err;
g_processing = 0;

53
src/chat/db_sync.c

@ -569,9 +569,10 @@ static int db_verify_author_sig(struct DB_SYNC_INSTANCE* si,
memcpy(msg + off, &ts, 8); off += 8;
if (jlen > 0) { memcpy(msg + off, json, jlen); off += jlen; }
if (sc_ed25519_verify(pubkey, msg, off, sig) != SC_OK) {
uint64_t key_prefix, sig_prefix; memcpy(&key_prefix, pubkey, 8); memcpy(&sig_prefix, sig, 8);
DEBUG_WARN(DEBUG_CATEGORY_CHAT_SYNC,
"author_sig VERIFY FAIL author=%016llx ed_pubkey=%016llx... msg_len=%zu sig=%016llx... — discarding as forgery",
(unsigned long long)author_node_id, *(uint64_t*)pubkey, off, *(uint64_t*)sig);
(unsigned long long)author_node_id, (unsigned long long)key_prefix, off, (unsigned long long)sig_prefix);
return -1;
}
return 0;
@ -808,7 +809,7 @@ static void db_handle_init_sync(struct DB_SYNC_INSTANCE* si, uint64_t src, const
{
if (len < 4) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "INIT_SYNC too short %zu from %016llx", len, (unsigned long long)src); return; }
db_sync_flush_recalc(si);
uint32_t pc = *(uint32_t*)p;
uint32_t pc; memcpy(&pc, p, 4);
uint32_t mc = db_count(si);
uint32_t tp = (pc < mc ? pc : mc);
if (tp > 0) tp--;
@ -861,7 +862,7 @@ static void db_handle_init_resp(struct DB_SYNC_INSTANCE* si, uint64_t src, const
{
if (len < 14) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "INIT_RESP too short %zu from %016llx", len, (unsigned long long)src); return; }
db_sync_flush_recalc(si);
uint32_t tp = *(uint32_t*)p;
uint32_t tp; memcpy(&tp, p, 4);
uint64_t peer_ch8; memcpy(&peer_ch8, p + 4, 8);
uint8_t has_tail = p[12];
uint8_t sc = p[13];
@ -918,7 +919,7 @@ static void db_handle_init_resp(struct DB_SYNC_INSTANCE* si, uint64_t src, const
// Mismatch: find first divergent position
uint32_t fm = tp;
for (int i = 0; i < sc && spr + 12 <= p + len; i++) {
uint32_t pos = *(uint32_t*)spr; uint64_t pch8; memcpy(&pch8, spr + 4, 8); spr += 12;
uint32_t pos; memcpy(&pos, spr, 4); uint64_t pch8; memcpy(&pch8, spr + 4, 8); spr += 12;
uint64_t mch8;
if (db_chain_hash8_at(si, pos, &mch8) == 0) {
if (mch8 != pch8 && pos < fm) fm = pos;
@ -950,7 +951,7 @@ static int si_send_data_batch(struct DB_SYNC_INSTANCE* si, uint64_t dst, uint32_
buf[off++] = DB_MSG_SEND_DATA;
memcpy(buf + off, &from, 4); off += 4;
uint16_t rc = 0;
uint16_t* rcp = (uint16_t*)(buf + off); off += 2;
uint8_t* rcp = buf + off; off += 2;
memcpy(buf + off, &vp, 4); off += 4;
memcpy(buf + off, &want_from, 4); off += 4; // NEW: request peer's data from this position, DB_WANT_FROM_NONE=none
@ -999,7 +1000,7 @@ static int si_send_data_batch(struct DB_SYNC_INSTANCE* si, uint64_t dst, uint32_
SI_SHRT(si), (unsigned long long)(dst >> 16), from, count);
}
*rcp = rc;
memcpy(rcp, &rc, 2);
db_sync_send(si, dst, buf, off);
DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "sync [%s:%04llX] → SEND_DATA: %u recs from=%u want=%u vp=%u",
SI_SHRT(si), (unsigned long long)(dst >> 16), rc, from, want_from, vp);
@ -1016,10 +1017,10 @@ static int si_parse_record(const uint8_t** pp, const uint8_t* end,
const uint8_t** rsig, int* rsiglen)
{
if (*pp + 28 > end) { DEBUG_WARN(DEBUG_CATEGORY_CHAT_SYNC, "si_parse_record: truncated header, need 28 have %td", end - *pp); return -1; }
*rid = *(uint64_t*)*pp; *pp += 8;
*rts = *(uint64_t*)*pp; *pp += 8;
*rauthor = *(uint64_t*)*pp; *pp += 8;
*rdlen = *(uint32_t*)*pp; *pp += 4;
memcpy(rid, *pp, 8); *pp += 8;
memcpy(rts, *pp, 8); *pp += 8;
memcpy(rauthor, *pp, 8); *pp += 8;
memcpy(rdlen, *pp, 4); *pp += 4;
if (*pp + *rdlen > end) { DEBUG_WARN(DEBUG_CATEGORY_CHAT_SYNC, "si_parse_record: data overrun dlen=%u have=%td", *rdlen, end - *pp); return -1; }
*rdata = (*rdlen > 0) ? *pp : NULL;
@ -1044,10 +1045,10 @@ static void db_handle_send_data(struct DB_SYNC_INSTANCE* si, uint64_t src, const
{
if (len < 14) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "sync [%s:%04llX] ← SEND_DATA: truncated len=%zu (need >=14)", SI_SHRT(si), (unsigned long long)(src >> 16), len); return; }
db_sync_flush_recalc(si);
uint32_t from = *(uint32_t*)p;
uint16_t count = *(uint16_t*)(p + 4);
uint32_t vp = *(uint32_t*)(p + 6);
uint32_t want_from = *(uint32_t*)(p + 10);
uint32_t from; memcpy(&from, p, 4);
uint16_t count; memcpy(&count, p + 4, 2);
uint32_t vp; memcpy(&vp, p + 6, 4);
uint32_t want_from; memcpy(&want_from, p + 10, 4);
struct SI_PEER* sp = si_peer_find(si, src);
@ -1110,7 +1111,7 @@ static void db_handle_sync_done(struct DB_SYNC_INSTANCE* si, uint64_t src, const
{
if (len < 12) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "SYNC_DONE too short %zu from %016llx", len, (unsigned long long)src); return; }
db_sync_flush_recalc(si);
uint32_t pc = *(uint32_t*)p;
uint32_t pc; memcpy(&pc, p, 4);
uint64_t pch8; memcpy(&pch8, p + 4, 8);
uint32_t mc = db_count(si);
uint64_t mch8 = 0; if (mc > 0) db_chain_hash8_at(si, mc - 1, &mch8);
@ -1152,8 +1153,8 @@ static void db_handle_sync_done(struct DB_SYNC_INSTANCE* si, uint64_t src, const
static void db_handle_ack_push(struct DB_SYNC_INSTANCE* si, uint64_t src, const uint8_t* p, size_t len)
{
if (len < 16) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "sync [%s:%04llX] ← PUSH ACK: truncated len=%zu (need >=16)", SI_SHRT(si), (unsigned long long)(src >> 16), len); return; }
uint64_t ts = *(uint64_t*)p;
uint64_t author = *(uint64_t*)(p + 8);
uint64_t ts; memcpy(&ts, p, 8);
uint64_t author; memcpy(&author, p + 8, 8);
DEBUG_TRACE(DEBUG_CATEGORY_CHAT_SYNC, "from=%016llx ts=%llu author=%016llx", (unsigned long long)src, (unsigned long long)ts, (unsigned long long)author);
sqlite3_stmt* stmt;
@ -1224,7 +1225,7 @@ static void db_sync_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry)
if (!db || !db->enabled) { queue_dgram_free(entry); queue_entry_free(entry); return; }
uint64_t src = conn ? conn->peer_node_id : 0;
uint64_t group_id = be64toh(*(uint64_t*)(entry->dgram + 1));
uint64_t group_id; memcpy(&group_id, entry->dgram + 1, 8); group_id = be64toh(group_id);
uint8_t type = entry->dgram[9];
const uint8_t* payload = entry->dgram + 10;
size_t plen = entry->len - 10;
@ -1596,13 +1597,13 @@ static void db_sync_peer_check_cb(void* arg)
#endif
DEBUG_TRACE(DEBUG_CATEGORY_CHAT_SYNC, "instances=%d", db->instance_count);
struct TOPO_GROUP* g = topo_groups_get_default(db->inst->topo_groups);
int total_synced = 0, total_skipped = 0, any_peers = 0;
char launched_list[256] = "";
if (g) {
{
for (int i = 0; i < db->instance_count; i++) {
struct DB_SYNC_INSTANCE* si = db->instances[i];
struct TOPO_GROUP* g = topo_groups_find(db->inst->topo_groups, si->group_id);
if (!g) continue;
if (!si->enabled) continue;
if (si->sync_gated) continue;
@ -1735,9 +1736,14 @@ int db_sync_init(struct UTUN_INSTANCE* inst)
return 0;
}
db->enabled = 1;
inst->db_sync = db;
return db_sync_enable(inst);
}
int db_sync_enable(struct UTUN_INSTANCE* inst) {
struct DB_SYNC* db = inst ? inst->db_sync : NULL;
if (!db) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "db_sync enable requires initialized core"); return -1; }
if (db->enabled) return 0;
const char* dp = inst->config->global.db_path;
if (inst->topo_sqlite_db) {
db->db = inst->topo_sqlite_db; db->shared_db = 1;
@ -1748,10 +1754,11 @@ int db_sync_init(struct UTUN_INSTANCE* inst)
else snprintf(sp, sizeof(sp), "/tmp/utun_db_sync");
if (db_sqlite_open(db, sp) != 0) {
DEBUG_WARN(DEBUG_CATEGORY_CHAT_SYNC, "SQLite open failed, sync disabled");
db->enabled = 0;
return -1;
}
}
db->enabled = 1;
etcp_bind(inst, ETCP_RT_ID_DB_SYNC, db_sync_recv_cb);
etcp_add_conn_status_cbk(inst, db_sync_on_conn_status, NULL);

1
src/chat/db_sync.h

@ -92,6 +92,7 @@ struct DB_SYNC_INSTANCE;
// Global lifecycle
int db_sync_init(struct UTUN_INSTANCE* inst);
int db_sync_enable(struct UTUN_INSTANCE* inst);
void db_sync_destroy(struct UTUN_INSTANCE* inst);
// Instance management

93
src/media_async/media_async.c

@ -19,70 +19,87 @@
#include <dirent.h>
#endif
struct media_async {
int initialized;
};
struct ma_thread_ctx;
struct media_async { struct ma_thread_ctx* jobs; int closing; };
struct ma_thread_ctx {
struct media_async* owner;
struct ma_thread_ctx* next;
struct UASYNC* ua;
pthread_t thread;
void* timer;
ma_work_fn work;
void* data;
ma_done_fn done;
void* arg;
int finished;
};
static void ma_done_trampoline(void* raw) {
struct ma_thread_ctx* ctx = (struct ma_thread_ctx*)raw;
ctx->done(ctx->arg, 0);
u_free(ctx);
static void ma_finish(struct ma_thread_ctx* job, int error) {
struct ma_thread_ctx** p = &job->owner->jobs;
while (*p && *p != job) p = &(*p)->next;
if (*p) *p = job->next;
pthread_join(job->thread, NULL);
job->done(job->arg, error);
u_free(job);
}
static void ma_poll(void* arg) {
struct ma_thread_ctx* job = arg; job->timer = NULL;
if (__atomic_load_n(&job->finished, __ATOMIC_ACQUIRE)) { ma_finish(job, 0); return; }
job->timer = uasync_set_timeout(job->ua, 100, job, ma_poll, "media_worker");
if (!job->timer) {
DEBUG_ERROR(DEBUG_CATEGORY_MEDIA, "media worker completion timer failed");
ma_finish(job, -1);
}
}
static void* ma_thread_entry(void* raw) {
struct ma_thread_ctx* ctx = (struct ma_thread_ctx*)raw;
ctx->work(ctx->data);
uasync_post(ctx->ua, ma_done_trampoline, ctx);
static void* ma_thread_entry(void* arg) {
struct ma_thread_ctx* job = arg;
job->work(job->data);
__atomic_store_n(&job->finished, 1, __ATOMIC_RELEASE);
uasync_wakeup(job->ua);
return NULL;
}
struct media_async* media_async_create(void) {
struct media_async* ma = u_malloc(sizeof(*ma));
if (!ma) return NULL;
ma->initialized = 1;
DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "media_async: created");
struct media_async* ma = u_calloc(1, sizeof(*ma));
if (!ma) DEBUG_ERROR(DEBUG_CATEGORY_MEDIA, "media worker owner allocation failed");
else DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "media_async: created");
return ma;
}
void media_async_destroy(struct media_async* ma) {
if (!ma) return;
ma->initialized = 0;
ma->closing = 1;
unsigned cancelled = 0;
while (ma->jobs) {
struct ma_thread_ctx* job = ma->jobs;
if (job->timer) uasync_cancel_timeout(job->ua, job->timer);
ma_finish(job, MEDIA_ASYNC_CANCELLED); cancelled++;
}
DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "media_async: destroyed cancelled=%u", cancelled);
u_free(ma);
DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "media_async: destroyed");
}
void media_async_submit(struct media_async* ma, struct UASYNC* ua,
ma_work_fn work, void* data,
ma_done_fn done, void* arg) {
(void)ma;
if (!ua || !work || !done) return;
struct ma_thread_ctx* ctx = u_malloc(sizeof(*ctx));
if (!ctx) { done(arg, -1); return; }
ctx->ua = ua;
ctx->work = work;
ctx->data = data;
ctx->done = done;
ctx->arg = arg;
pthread_t tid;
int rc = pthread_create(&tid, NULL, ma_thread_entry, ctx);
if (rc != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_MEDIA, "media_async: pthread_create failed rc=%d", rc);
u_free(ctx);
done(arg, -1);
ma_work_fn work, void* data, ma_done_fn done, void* arg) {
if (!ma || !ua || !work || !done || ma->closing) {
DEBUG_WARN(DEBUG_CATEGORY_MEDIA, "media task rejected: invalid arguments or stopping owner");
if (done) done(arg, MEDIA_ASYNC_CANCELLED);
return;
}
pthread_detach(tid);
struct ma_thread_ctx* job = u_calloc(1, sizeof(*job));
if (!job) { DEBUG_ERROR(DEBUG_CATEGORY_MEDIA, "media task allocation failed"); done(arg, -1); return; }
job->owner = ma; job->ua = ua; job->work = work; job->data = data; job->done = done; job->arg = arg;
job->timer = uasync_set_timeout(ua, 1, job, ma_poll, "media_worker");
if (!job->timer) { DEBUG_ERROR(DEBUG_CATEGORY_MEDIA, "media task timer failed"); u_free(job); done(arg, -1); return; }
int rc = pthread_create(&job->thread, NULL, ma_thread_entry, job);
if (rc) {
DEBUG_ERROR(DEBUG_CATEGORY_MEDIA, "media task thread failed rc=%d", rc);
uasync_cancel_timeout(ua, job->timer); u_free(job); done(arg, -1); return;
}
job->next = ma->jobs; ma->jobs = job;
}
/* ─── примитивы (вызываются из worker thread) ─── */

3
src/media_async/media_async.h

@ -17,6 +17,9 @@ typedef void (*ma_work_fn)(void* data);
typedef void (*ma_done_fn)(void* arg, int err);
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. */
void media_async_destroy(struct media_async* ma);
void media_async_submit(struct media_async* ma, struct UASYNC* ua,

90
src/media_delivery/media_delivery.c

@ -23,6 +23,7 @@
/* ── forward decls ── */
static void md_relay_cancel_block(struct media_delivery_ctx* md, struct relay_block_ctx* block);
static void md_etcp_recv(struct ETCP_CONN* conn, struct ll_entry* entry);
static void md_on_conn_status(struct ETCP_CONN* conn, int status, void* arg);
static void md_on_bgp_node(struct TOPO_GROUP* group, uint64_t node_id, int event, void* arg);
@ -266,6 +267,7 @@ void md_relay_remove(struct media_delivery_ctx* md, const uint8_t* block_id) {
if (!md->relay_blocks) return;
struct ll_entry* e = queue_find_data_by_index(md->relay_blocks, block_id);
if (!e) return;
md_relay_cancel_block(md, (struct relay_block_ctx*)e->data);
queue_remove_data(md->relay_blocks, e);
queue_entry_free(e);
}
@ -717,6 +719,7 @@ static void md_super_conn_cb(struct CONN_MGR_HANDLE* h, uint64_t node_id, uint64
#define MD_CHUNK_SIZE 1024
struct stream_ctx {
struct stream_ctx* next;
struct media_delivery_ctx* md;
uint64_t group_id;
uint64_t dst_node_id;
@ -731,18 +734,30 @@ struct stream_ctx {
struct queue_waiter_handle waiter;
};
static void stream_free(struct stream_ctx* sc) {
struct media_delivery_ctx* md = sc->md;
etcp_router_cancel_send_ready(md->inst, sc->group_id, sc->dst_node_id, ETCP_RT_ID_MEDIA_DELIVERY, &sc->waiter);
struct stream_ctx** p = &md->streams;
while (*p && *p != sc) p = &(*p)->next;
if (*p) *p = sc->next;
md_file_load_dec(md, sc->media_id, sc->dst_node_id);
media_delivery_stream_done(md->inst);
if (sc->file) fclose(sc->file);
u_free(sc);
}
static void stream_send_chunk_cb(struct ll_queue* q, void* arg) {
(void)q;
struct stream_ctx* sc = (struct stream_ctx*)arg;
struct media_delivery_ctx* md = sc->md;
if (!md || sc->remaining == 0) { u_free(sc); return; }
if (!md->initialized || sc->remaining == 0) { stream_free(sc); return; }
while (sc->remaining > 0) {
size_t to_read = sc->remaining < MD_CHUNK_SIZE ? (size_t)sc->remaining : MD_CHUNK_SIZE;
uint8_t buf[MD_CHUNK_SIZE];
fseeko(sc->file, (off_t)sc->offset, SEEK_SET);
size_t rd = fread(buf, 1, to_read, sc->file);
if (rd == 0) { u_free(sc); return; }
if (rd == 0) { DEBUG_ERROR(DEBUG_CATEGORY_MEDIA, "stream read failed"); stream_free(sc); return; }
uint8_t pkt[MEDIA_BLOCK_CHUNK_HDR_SIZE + MD_CHUNK_SIZE];
struct media_pkt_block_chunk* ch = (struct media_pkt_block_chunk*)pkt;
@ -759,10 +774,7 @@ static void stream_send_chunk_cb(struct ll_queue* q, void* arg) {
if (snd != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_MEDIA, "%s: stream chunk send FAILED rc=%d — abort stream chunk=%d to 0x%016llx",
MD_ID, snd, sc->chunk, (unsigned long long)sc->dst_node_id);
md_file_load_dec(md, sc->media_id, sc->dst_node_id);
media_delivery_stream_done(md->inst);
fclose(sc->file);
u_free(sc);
stream_free(sc);
return;
}
@ -802,10 +814,7 @@ static void stream_send_chunk_cb(struct ll_queue* q, void* arg) {
bd->total_size = (uint32_t)sc->block_data_len;
md_send(md->inst, sc->group_id, sc->dst_node_id, done_pkt, MEDIA_BLOCK_DONE_SIZE);
md_file_load_dec(md, sc->media_id, sc->dst_node_id);
media_delivery_stream_done(md->inst);
fclose(sc->file);
u_free(sc);
stream_free(sc);
DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%s: BLOCK_DONE sent chunk=%d total=%u",
MD_ID, chunk, bd->total_size);
}
@ -814,6 +823,7 @@ static void stream_send_chunk_cb(struct ll_queue* q, void* arg) {
/* ── relay chunk forwarding ── */
struct relay_fwd_ctx {
struct relay_fwd_ctx* next;
struct media_delivery_ctx* md;
uint64_t dst_node_id;
uint64_t group_id;
@ -822,16 +832,33 @@ struct relay_fwd_ctx {
struct queue_waiter_handle waiter;
};
static void relay_fwd_free(struct relay_fwd_ctx* fc) {
etcp_router_cancel_send_ready(fc->md->inst, fc->group_id, fc->dst_node_id, ETCP_RT_ID_MEDIA_DELIVERY, &fc->waiter);
struct relay_fwd_ctx** p = &fc->md->forwards;
while (*p && *p != fc) p = &(*p)->next;
if (*p) *p = fc->next;
u_free(fc);
}
static void md_relay_cancel_block(struct media_delivery_ctx* md, struct relay_block_ctx* block) {
struct relay_fwd_ctx* fc = md->forwards;
while (fc) {
struct relay_fwd_ctx* next = fc->next;
if (fc->rc == block) relay_fwd_free(fc);
fc = next;
}
}
static void md_relay_fwd_cb(struct ll_queue* q, void* arg) {
(void)q;
struct relay_fwd_ctx* fc = (struct relay_fwd_ctx*)arg;
struct relay_block_ctx* rc = fc->rc;
struct relay_downstream* ds = &rc->downstream[fc->ds_idx];
if (ds->sent_offset >= rc->file_offset) { u_free(fc); return; }
if (ds->sent_offset >= rc->file_offset) { relay_fwd_free(fc); return; }
FILE* f = fopen(rc->chunk_file, "rb");
if (!f) { u_free(fc); return; }
if (!f) { relay_fwd_free(fc); return; }
fseeko(f, (off_t)ds->sent_offset, SEEK_SET);
size_t to_read = rc->file_offset - ds->sent_offset;
@ -839,7 +866,7 @@ static void md_relay_fwd_cb(struct ll_queue* q, void* arg) {
uint8_t buf[1024];
size_t rd = fread(buf, 1, to_read, f);
fclose(f);
if (rd == 0) { u_free(fc); return; }
if (rd == 0) { relay_fwd_free(fc); return; }
uint8_t pkt[MEDIA_BLOCK_CHUNK_HDR_SIZE + 1024];
struct media_pkt_block_chunk* ch = (struct media_pkt_block_chunk*)pkt;
@ -862,7 +889,7 @@ static void md_relay_fwd_cb(struct ll_queue* q, void* arg) {
fc->dst_node_id, ETCP_RT_ID_MEDIA_DELIVERY,
&fc->waiter, md_relay_fwd_cb, fc);
} else {
u_free(fc);
relay_fwd_free(fc);
}
}
@ -873,6 +900,7 @@ static void md_relay_catchup(struct media_delivery_ctx* md, struct relay_block_c
struct relay_fwd_ctx* fc = u_calloc(1, sizeof(*fc));
if (!fc) return;
fc->next = md->forwards; md->forwards = fc;
fc->md = md; fc->dst_node_id = dst_node_id; fc->group_id = group_id;
fc->rc = rc; fc->ds_idx = ds_idx;
memset(&fc->waiter, 0, sizeof(fc->waiter));
@ -1026,7 +1054,7 @@ static void md_handle_block_req(struct media_delivery_ctx* md, uint64_t from_nod
struct stream_ctx* sc = u_calloc(1, sizeof(*sc));
if (!sc) { fclose(f); md->active_streams--; md_file_load_dec(md, req->media_id, from_node); return; }
sc->md = md;
sc->md = md; sc->next = md->streams; md->streams = sc;
sc->group_id = group_id;
sc->dst_node_id = from_node;
memcpy(sc->media_id, req->media_id, 16);
@ -1415,6 +1443,10 @@ void media_delivery_destroy(struct UTUN_INSTANCE* inst) {
if (!inst || !inst->md.initialized) return;
struct media_delivery_ctx* md = &inst->md;
md->initialized = 0;
while (md->streams) stream_free(md->streams);
while (md->forwards) relay_fwd_free(md->forwards);
etcp_router_unbind(inst, ETCP_RT_ID_MEDIA_DELIVERY);
etcp_remove_conn_status_cbk(inst, md_on_conn_status, md);
member_sync_remove_props_cbk(inst, md_on_props_changed, md);
@ -1440,18 +1472,11 @@ void media_delivery_destroy(struct UTUN_INSTANCE* inst) {
}
if (md->served_nodes) { queue_free(md->served_nodes); md->served_nodes = NULL; }
if (md->downloads) {
struct ll_entry* de = md->downloads->head;
while (de) {
struct media_download* dl = (struct media_download*)de->data;
if (dl->watchdog_timer) { uasync_cancel_timeout(md->inst->ua, dl->watchdog_timer); dl->watchdog_timer = NULL; }
for (int pi = 0; pi < dl->num_peers; pi++)
if (dl->peers[pi].cm_handle) { conn_mgr_close(dl->peers[pi].cm_handle); dl->peers[pi].cm_handle = NULL; }
if (dl->blocks) u_free(dl->blocks);
if (dl->block_sigs) u_free(dl->block_sigs);
de = de->next;
while (md->downloads->head) {
struct media_download* dl = (struct media_download*)md->downloads->head->data;
media_download_cancel(inst, dl->media_id, NULL, 0);
}
queue_free(md->downloads);
md->downloads = NULL;
queue_free(md->downloads); md->downloads = NULL;
}
if (md->relay_blocks) { queue_free(md->relay_blocks); md->relay_blocks = NULL; }
if (md->file_loads) { queue_free(md->file_loads); md->file_loads = NULL; }
@ -1464,6 +1489,19 @@ void media_delivery_cancel_channel(struct UTUN_INSTANCE* inst, uint64_t group_id
if (!inst || !inst->md.initialized) return;
struct media_delivery_ctx* md = &inst->md;
struct stream_ctx* stream = md->streams;
while (stream) {
struct stream_ctx* next = stream->next;
if (stream->group_id == group_id) stream_free(stream);
stream = next;
}
struct relay_fwd_ctx* forward = md->forwards;
while (forward) {
struct relay_fwd_ctx* next = forward->next;
if (forward->group_id == group_id) relay_fwd_free(forward);
forward = next;
}
/* 1. собрать media_id активных загрузок этого канала и отменить их */
uint8_t media_ids[MEDIA_MAX_DOWNLOADS][16];
int nm = 0;

4
src/media_delivery/media_delivery.h

@ -105,7 +105,11 @@ struct relay_block_ctx {
struct UTUN_INSTANCE* inst;
};
struct stream_ctx;
struct relay_fwd_ctx;
struct media_delivery_ctx {
struct stream_ctx* streams;
struct relay_fwd_ctx* forwards;
uint64_t self_node_id;
int initialized;
uint8_t is_supernode; // из adm_tags в peers_*

2
src/media_delivery/media_index.c

@ -294,8 +294,8 @@ static void mi_reg_work(void* data) {
}
static void mi_reg_done(void* arg, int err) {
(void)err;
struct mi_reg_ctx* ctx = (struct mi_reg_ctx*)arg;
if (err) ctx->err = err;
if (ctx->err == 0) {
int rc = media_index_commit(ctx->db, &ctx->result, ctx->node_id,

30
src/routing_layer/etcp_router.c

@ -1356,18 +1356,18 @@ void etcp_router_on_send_ready(struct UTUN_INSTANCE* inst, uint64_t group_id, ui
struct queue_waiter_handle* h,
queue_threshold_callback_fn callback, void* arg) {
if (!inst || !h) return;
etcp_router_cancel_send_ready(inst, group_id, node_id, svc_id, h);
struct ETCP_ROUTER_CONN* rconn = router_conn_find(inst, group_id, node_id, svc_id);
if (!rconn || !rconn->send_q) return;
h->defer_q = rconn->send_q; /* Keep the registration queue across route/session replacement. */
queue_waiter_wait(rconn->send_q, h, callback, arg);
}
// Backpressure: отменить waiter на send_q для (group_id, node_id, svc_id).
void etcp_router_cancel_send_ready(struct UTUN_INSTANCE* inst, uint64_t group_id, uint64_t node_id, uint8_t svc_id,
struct queue_waiter_handle* h) {
if (!inst || !h) return;
struct ETCP_ROUTER_CONN* rconn = router_conn_find(inst, group_id, node_id, svc_id);
if (!rconn || !rconn->send_q) return;
queue_waiter_cancel(rconn->send_q, h);
(void)inst; (void)group_id; (void)node_id; (void)svc_id;
if (h && (h->internal || h->call_soon_id)) queue_waiter_cancel(h->defer_q, h);
}
// Backpressure: есть ли место в send_q (без регистрации waiter). 1 — есть, 0 — нет.
@ -1466,3 +1466,25 @@ void etcp_router_conn_close_all_for_node(struct UTUN_INSTANCE* inst, uint64_t gr
for (int i = 0; i < count; i++)
router_close_and_notify((struct ETCP_ROUTER_CONN*)to_close[i]);
}
/* Service stop: discard every logical channel in a group before its routes disappear.
* Call after service producers have cancelled their waiters. */
void etcp_router_close_group(struct UTUN_INSTANCE* inst, uint64_t group_id) {
if (!inst) return;
/* Notifications may close other channels too; do not retain list cursors across them. */
while (inst->router_conns) {
struct ll_entry* entry = inst->router_conns->head;
while (entry && ((struct ETCP_ROUTER_CONN*)entry)->group_id != group_id) entry = entry->next;
if (!entry) break;
router_close_and_notify((struct ETCP_ROUTER_CONN*)entry);
}
for (struct ll_entry* e = inst->connections ? inst->connections->head : NULL; e; e = e->next) {
struct ETCP_CONN* conn = ((struct conn_queue_entry*)e->data)->conn;
if (!conn) continue;
struct ll_entry* entry = conn->transit_queues ? conn->transit_queues->head : NULL;
while (entry) {
struct TRANSIT_QUEUE* transit = (struct TRANSIT_QUEUE*)entry; entry = entry->next;
if (transit->group_id == group_id) transit_queue_destroy(conn, transit);
}
}
}

1
src/routing_layer/etcp_router.h

@ -253,6 +253,7 @@ int etcp_router_send_q_has_room(struct UTUN_INSTANCE* inst, uint64_t group_id, u
void etcp_router_pause_retrans_for_node(struct UTUN_INSTANCE* inst, uint64_t remote_node_id);
// Закрыть все router_conn для указанного (group_id, remote_node_id) (peer умер)
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 в группе и уведомить сервис.

2
src/routing_layer/routing.c

@ -281,7 +281,7 @@ void routing_destroy(struct UTUN_INSTANCE* instance) {
(unsigned long long)instance->node_id);
// Unbind DATA handler from etcp_router
etcp_router_unbind(instance, ETCP_RT_ID_DATA);
if (instance->router_bindings.callbacks[ETCP_RT_ID_DATA]) etcp_router_unbind(instance, ETCP_RT_ID_DATA);
// Clean up route table if exists
if (instance->rt) {

43
src/routing_layer/topo_group.c

@ -563,6 +563,16 @@ static void topo_group_destroy(struct TOPO_GROUP* group) {
radio_destroy(group);
DEBUG_INFO(DEBUG_CATEGORY_SYS, "[GRP_DESTROY] 2c radio_destroy done");
/* A stopped service must withdraw membership even if another group keeps ETCP alive. */
if (group->instance->running) {
struct ll_entry* current = group->senders_list->head;
while (current) {
struct TOPO_GROUP_CONN_ITEM* peer = (struct TOPO_GROUP_CONN_ITEM*)current->data;
current = current->next;
if (peer->conn && peer->conn->links_up) topo_group_leave_peer(group, peer->handle);
}
}
struct ll_entry* e;
while ((e = queue_data_get(group->senders_list)) != NULL) {
struct TOPO_GROUP_CONN_ITEM* item = (struct TOPO_GROUP_CONN_ITEM*)e->data;
@ -574,6 +584,7 @@ static void topo_group_destroy(struct TOPO_GROUP* group) {
queue_free(group->senders_list);
DEBUG_INFO(DEBUG_CATEGORY_SYS, "[GRP_DESTROY] 3 senders_list done");
etcp_router_close_group(group->instance, group->group_id);
if (group->conn_mgr) { conn_mgr_destroy(group->conn_mgr); group->conn_mgr = NULL; }
DEBUG_INFO(DEBUG_CATEGORY_SYS, "[GRP_DESTROY] 4 conn_mgr done");
@ -589,6 +600,7 @@ static void topo_group_destroy(struct TOPO_GROUP* group) {
if (nc) DEBUG_INFO(DEBUG_CATEGORY_SYS, "[GRP_DESTROY] 4b nodes cleaned (%d)", nc);
}
if (group->local_node) {
route_delete(group->instance->rt, group->local_node);
if (group->local_node->paths) { queue_free(group->local_node->paths); group->local_node->paths = NULL; }
topo_nodeq_free_group_fields(group->instance->topo_groups, group->local_node);
u_free(group->local_node);
@ -601,15 +613,16 @@ static void topo_group_destroy(struct TOPO_GROUP* group) {
DEBUG_INFO(DEBUG_CATEGORY_SYS, "[GRP_DESTROY] 5 done");
}
/* Удаляет группу из контейнера (кроме UTUN-группы по умолчанию). */
/* Удаляет группу и все принадлежащие ей ресурсы. */
void topo_groups_remove_group(struct TOPO_GROUPS* g, uint64_t group_id) {
if (!g || !g->group_list || group_id == TOPO_GROUP_UTUN) return;
if (!g || !g->group_list) return;
struct TOPO_GROUP* group = topo_groups_find(g, group_id);
if (!group) { DEBUG_WARN(DEBUG_CATEGORY_BGP, "remove_group: group %016llx not found", (unsigned long long)group_id); return; }
DEBUG_INFO(DEBUG_CATEGORY_BGP, "remove_group: id=%016llx type=%d ch=%s",
(unsigned long long)group_id, group->group_type, group->channel_id);
if (g->instance->conn_mgr == group->conn_mgr) g->instance->conn_mgr = NULL;
queue_remove_data(g->group_list, &group->ll);
topo_group_destroy(group); /* node_cbks освобождаются внутри */
queue_entry_free(&group->ll);
@ -637,8 +650,9 @@ struct TOPO_GROUPS* topo_groups_init(struct UTUN_INSTANCE* instance) {
g->v4_subnet_pool = memory_pool_init(sizeof(struct TOPO_SUBNET4), "to_v4sub");
g->v6_subnet_pool = memory_pool_init(sizeof(struct TOPO_SUBNET6), "to_v6sub");
if (instance->config && instance->config->global.db_path[0]) {
char db_file[512];
{
char db_file[512] = ":memory:";
if (instance->config && instance->config->global.db_path[0])
snprintf(db_file, sizeof(db_file), "%s/chats.db", instance->config->global.db_path);
int rc = sqlite3_open_v2(db_file, &instance->topo_sqlite_db,
SQLITE_OPEN_READWRITE | SQLITE_OPEN_CREATE | SQLITE_OPEN_FULLMUTEX, NULL);
@ -661,18 +675,7 @@ struct TOPO_GROUPS* topo_groups_init(struct UTUN_INSTANCE* instance) {
instance->topo_groups = g;
struct TOPO_GROUP* default_group = topo_group_create(instance, TOPO_GROUP_UTUN, TOPO_GROUP_TYPE_UTUN);
if (!default_group) {
memory_pool_destroy(g->v4_sock_meta_pool); memory_pool_destroy(g->v4_addr_pool);
memory_pool_destroy(g->v6_sock_meta_pool); memory_pool_destroy(g->v6_addr_pool);
memory_pool_destroy(g->v4_subnet_pool); memory_pool_destroy(g->v6_subnet_pool);
queue_free(g->node_registry); queue_free(g->group_list); u_free(g);
DEBUG_ERROR(DEBUG_CATEGORY_BGP, "failed to create default group");
return NULL;
}
queue_data_put_with_index(g->group_list, &default_group->ll);
instance->conn_mgr = default_group->conn_mgr; /* recv handler shortcut */
if (topo_node_init_self(instance) < 0) { topo_groups_destroy(instance); return NULL; }
etcp_bind(instance, ETCP_ID_TOPO_ENTRY, topo_group_receive_cbk);
etcp_add_conn_status_cbk(instance, topo_group_conn_status, g);
@ -680,7 +683,7 @@ struct TOPO_GROUPS* topo_groups_init(struct UTUN_INSTANCE* instance) {
ETCP_SOCKET_EVENT_ADDR_CHANGED | ETCP_SOCKET_EVENT_STATUS_CHANGED);
utun_add_activity_cbk(instance, topo_group_on_activity_change, NULL);
DEBUG_INFO(DEBUG_CATEGORY_BGP, "TOPO groups initialized (with default group)");
DEBUG_INFO(DEBUG_CATEGORY_BGP, "TOPO registry initialized (services create their own groups)");
return g;
}
@ -1655,8 +1658,12 @@ static void topo_leave_event(struct NODE_CONN_DIRECT* handle, enum ncd_event eve
}
static void topo_leave_ready(struct ll_queue* queue, void* arg) {
(void)queue;
struct topo_leave* leave = arg;
if (queue_entry_count(queue)) {
if (queue_waiter_wait(queue, &leave->waiter, topo_leave_ready, leave) >= 0) return;
DEBUG_ERROR(DEBUG_CATEGORY_BGP, "cannot rearm LEAVE_GROUP waiter");
topo_leave_finish(leave); return;
}
struct ETCP_CONN* conn = node_conn_direct_get_conn(leave->handle);
struct TOPO_GROUP* group = topo_groups_find(leave->instance->topo_groups, leave->group_id);
if (group && topo_group_has_sender(group, conn)) {

4
src/routing_layer/topo_group.h

@ -281,8 +281,8 @@ struct TOPO_GROUP* topo_groups_create_group(struct TOPO_GROUPS* g, uint64_t grou
* @brief Удаляет одну группу по id (разрушает conn_mgr, connect-цикл, broadcast,
* recovery, узлы и подписчиков node_cbks) и извлекает её из group_list.
*
* Используется при локальном удалении чат-канала. Группу по умолчанию
* (TOPO_GROUP_UTUN) удалять нельзя — функция это игнорирует.
* Используется при удалении канала и остановке сервиса. UTUN-группой
* управляет utun_service_start/stop.
*/
void topo_groups_remove_group(struct TOPO_GROUPS* g, uint64_t group_id);

7
src/routing_layer/topo_group_invite.c

@ -122,8 +122,7 @@ static void tgi_handle_info_req(struct ETCP_CONN* conn, const uint8_t* data, siz
char ch_id[32]; snprintf(ch_id, sizeof(ch_id), "%llu", (unsigned long long)req->group_id);
struct TOPO_GROUP* grp = topo_groups_find(inst->topo_groups, req->group_id);
struct TOPO_GROUP* sg = grp ? grp : topo_groups_get_default(inst->topo_groups);
struct TOPO_NODE* ln = sg && sg->local_node ? topo_node_registry_find(inst->topo_groups, sg->local_node->node_id) : NULL;
struct TOPO_NODE* ln = topo_node_registry_find(inst->topo_groups, inst->node_id);
const char* nm = ln ? ln->node_name : NULL; size_t nl = nm ? strlen(nm) : 0; if (nl > 255) nl = 255;
sqlite3* db = inst->topo_sqlite_db;
@ -138,7 +137,7 @@ static void tgi_handle_info_req(struct ETCP_CONN* conn, const uint8_t* data, siz
struct TGI_INVITE_RESP* resp = (struct TGI_INVITE_RESP*)rb;
resp->cmd = ETCP_RT_ID_GROUP_INVITE; resp->subcmd = TGI_SUBCMD_INFO_RESP;
resp->group_id = req->group_id; resp->is_member = 0;
if (sg) memcpy(resp->ed25519_pubkey, sg->ed25519_public_key, SC_PUBKEY_SIZE);
memcpy(resp->ed25519_pubkey, inst->my_ed25519_pubkey, SC_PUBKEY_SIZE);
resp->node_name_len = (uint8_t)nl; if (nl) memcpy(resp->node_name, nm, nl);
struct ll_entry* qe = queue_entry_new(0);
if (qe) { qe->dgram = rb; qe->len = (uint16_t)rs;
@ -207,7 +206,7 @@ static void tgi_handle_info_req(struct ETCP_CONN* conn, const uint8_t* data, siz
{ struct TGI_INVITE_RESP* resp = (struct TGI_INVITE_RESP*)rb;
resp->cmd = ETCP_RT_ID_GROUP_INVITE; resp->subcmd = TGI_SUBCMD_INFO_RESP;
resp->group_id = req->group_id; resp->is_member = 1;
if (sg) memcpy(resp->ed25519_pubkey, sg->ed25519_public_key, SC_PUBKEY_SIZE);
memcpy(resp->ed25519_pubkey, inst->my_ed25519_pubkey, SC_PUBKEY_SIZE);
resp->node_name_len = (uint8_t)nl; if (nl) memcpy(resp->node_name, nm, nl); }
uint8_t* ch = rb + TGI_INVITE_RESP_HDR_SIZE + nl;

2
src/routing_layer/topo_group_invite.h

@ -58,7 +58,7 @@ int topo_group_invite_join(struct UTUN_INSTANCE* inst, uint64_t group_id,
* Инвайтера сторона: прямое подключение (ncd) к целевому узлу для последующей
* отправки CHANNEL_INVITE (через chat_sync, по ETCP_CONN_STATUS_UP).
* Без INVITE_INFO-проверки членства. ni может быть NULL — тогда адреса грузятся
* из БД (node_conn_direct_open). Handle сохраняется в nq->handle.
* из БД. Группа владеет NCD; инициатор держит отменяемый групповой запрос.
*/
int topo_group_invite_to_channel(struct UTUN_INSTANCE* inst, uint64_t group_id,
struct TOPO_NODE* ni, uint64_t node_id);

22
src/routing_layer/topo_node.c

@ -192,6 +192,7 @@ void topo_nodeq_free_group_fields(struct TOPO_GROUPS* groups, struct TOPO_GROUP_
void topo_nodeq_remove_node(struct TOPO_GROUP* group, struct TOPO_GROUP_NODE* nq) {
if (!group || !nq) return;
route_connectivity_cancel_node(nq);
if (group->instance) route_delete(group->instance->rt, nq);
if (nq->paths) { queue_free(nq->paths); nq->paths = NULL; }
queue_remove_data(group->nodes, &nq->ll);
topo_nodeq_free_group_fields(group->instance ? group->instance->topo_groups : NULL, nq);
@ -928,6 +929,25 @@ int topo_node_format_all(struct TOPO_GROUP* group, char* buf, size_t buf_size) {
// ===== update my nodeinfo =====
/* Core owns one reference independently of all services. */
int topo_node_init_self(struct UTUN_INSTANCE* instance) {
struct TOPO_NODE* ni = u_calloc(1, sizeof(*ni));
if (!ni) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "local identity allocation failed"); return -1; }
ni->node_id = instance->node_id; ni->ver = 1;
ni->client_type = instance->client_type; ni->client_activity = instance->client_activity;
memcpy(ni->public_key, instance->my_keys.public_key, SC_PUBKEY_SIZE);
memcpy(ni->ed25519_public_key, instance->my_ed25519_pubkey, SC_PUBKEY_SIZE);
ni->node_name = u_strdup(instance->name);
if (!ni->node_name || topo_node_sign_self(instance, ni) < 0) {
DEBUG_ERROR(DEBUG_CATEGORY_BGP, "cannot initialize local identity");
topo_node_destroy(instance->topo_groups, ni); return -1;
}
if (!topo_node_registry_store(instance->topo_groups, ni)) {
DEBUG_ERROR(DEBUG_CATEGORY_BGP, "cannot register local identity"); return -1;
}
return 0;
}
/* Пересоздаёт local_node группы при изменении имени/подсетей/ключей, инкрементируя версию. */
int topo_group_update_my_nodeinfo(struct UTUN_INSTANCE* instance, struct TOPO_GROUP* group) {
if (!instance || !group) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "invalid args"); return -1; }
@ -1039,8 +1059,6 @@ int topo_group_update_my_nodeinfo(struct UTUN_INSTANCE* instance, struct TOPO_GR
/* Пересобирает списки адресов своего узла из ETCP-сокетов; при изменении — подпись, БД, рассылка. */
int topo_node_update_my_addresses(struct UTUN_INSTANCE* instance) {
if (!instance || !instance->topo_groups) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "invalid args"); return -1; }
struct TOPO_GROUP* default_group = topo_groups_get_default(instance->topo_groups);
if (!default_group || !default_group->local_node) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "no local_node"); return -1; }
struct TOPO_NODE* ni = topo_node_registry_find(instance->topo_groups, instance->node_id);
if (!ni) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "my node not in registry"); return -1; }

1
src/routing_layer/topo_node.h

@ -243,6 +243,7 @@ int topo_node_deserialize(struct TOPO_GROUP* group, const uint8_t* data, size_t
uint64_t* topo_node_best_hop_list(struct TOPO_GROUP_NODE* nq, uint8_t* out_count, uint16_t* out_rtt);
int topo_group_update_my_nodeinfo(struct UTUN_INSTANCE* instance, struct TOPO_GROUP* group);
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);

15
src/transport_layer/etcp_connections.c

@ -1899,19 +1899,7 @@ static void send_init_response(struct ETCP_SOCKET* e_sock, struct ETCP_DGRAM* pk
nat_sin->sin_family = AF_INET;
nat_sin->sin_addr.s_addr = my_ip;
nat_sin->sin_port = ((struct sockaddr_in*)&e_sock->interface_addr)->sin_port;
{ struct TOPO_GROUP* g = topo_groups_get_default(link->etcp->instance->topo_groups);
if (g) {
topo_group_update_my_nodeinfo(g->instance, g);
if (g->local_node && g->senders_list) {
struct ll_entry* se = g->senders_list->head;
while (se) {
struct TOPO_GROUP_CONN_ITEM* item = (struct TOPO_GROUP_CONN_ITEM*)se->data;
if (item && item->conn) topo_group_send_nodeinfo(g, g->local_node, item->conn);
se = se->next;
}
}
}
}
if (link->etcp->instance->topo_groups) topo_node_update_my_addresses(link->etcp->instance);
DEBUG_INFO(DEBUG_CATEGORY_NAT, "Self-detected PUBLIC: socket=%s ip=%s port=%u sock_id=%d",
e_sock->name, ip_to_str(&my_ip, AF_INET).str,
ntohs(nat_sin->sin_port), e_sock->sock_id);
@ -2890,7 +2878,6 @@ int init_connections(struct UTUN_INSTANCE* instance) {
etcp_keepalive_register(instance);
if (init_client_connections(instance, config->clients) < 0) return -1;
// Return 1 if there was a partial socket initialization error
if (socket_result == 1) return 1;
return 0;

258
src/utun_instance.c

@ -52,6 +52,7 @@
// Forward declarations
static uint32_t get_dest_ip(const uint8_t *packet, size_t len);
static int mkdir_recursive(const char* path);
// Global instance for signal handlers
static struct UTUN_INSTANCE *g_instance = NULL;
@ -193,22 +194,6 @@ static int instance_init_common(struct UTUN_INSTANCE* instance, struct UASYNC* u
return -1;
}
DEBUG_INFO(DEBUG_CATEGORY_ROUTING, "Routing module created");
if (g_tun_init_enabled && config->global.tun_enabled) {
instance->tun = tun_init(ua, config);
if (!instance->tun) {
DEBUG_ERROR(DEBUG_CATEGORY_TUN, "Failed to initialize TUN device");
return -1;
}
DEBUG_INFO(DEBUG_CATEGORY_TUN, "TUN interface initialized: %s", instance->tun->ifname);
if (config->route_subnets) {
int added = tun_route_add_all(instance->tun->ifindex, instance->tun->ifname, config->route_subnets);
DEBUG_INFO(DEBUG_CATEGORY_TUN, "Added %d system routes for TUN interface", added);
instance->route_subnets = config->route_subnets;
}
} else {
DEBUG_INFO(DEBUG_CATEGORY_TUN, "TUN initialization disabled - skipping TUN device setup");
instance->tun = NULL;
}
int socket_result = init_sockets(instance);
instance->socket_init_status = socket_result;
if (socket_result != 0) {
@ -218,17 +203,14 @@ static int instance_init_common(struct UTUN_INSTANCE* instance, struct UASYNC* u
instance->nat_det = nat_detection_create(instance);
if (!instance->nat_det) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "NAT detection creation failed"); }
if (config->global.db_path[0]) mkdir_recursive(config->global.db_path);
if (g_topo_group_enabled) {
instance->topo_groups = topo_groups_init(instance);
if (!instance->topo_groups) {
DEBUG_ERROR(DEBUG_CATEGORY_BGP, "Failed to initialize BGP module");
} else {
DEBUG_INFO(DEBUG_CATEGORY_BGP, "BGP module initialized");
topo_group_invite_init(instance);
struct TOPO_GROUP* g = topo_groups_get_default(instance->topo_groups);
if (instance->rt && g && g->local_node) {
topo_group_update_my_nodeinfo(instance, g);
}
}
} else {
DEBUG_INFO(DEBUG_CATEGORY_BGP, "topo_group initialization disabled, skipping BGP module");
@ -261,12 +243,6 @@ static int instance_init_common(struct UTUN_INSTANCE* instance, struct UASYNC* u
if (g_topo_group_enabled && instance->topo_groups)
etcp_router_bind(instance, ETCP_RT_ID_CONN_MGR, conn_mgr_router_recv_handler);
// Bind DATA handler via etcp_router (after etcp_router_init)
if (routing_bind(instance) != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_ROUTING, "Failed to bind DATA via etcp_router");
return -1;
}
// TCP proxy server (exit node, optional) — must be before tcp_proxy_client so
// tcp_proxy_client_create can overwrite the handler if it has remote mappings
#ifndef _WIN32
@ -440,20 +416,8 @@ void utun_instance_destroy(struct UTUN_INSTANCE *instance) {
}
DEBUG_INFO(DEBUG_CATEGORY_SYS, "[DESTROY] C done — control_srv");
/* Phase D: config conn handles */
{
struct CONFIG_CONN_HANDLE* ch = instance->config_conn_handles;
int ch_count = 0;
while (ch) {
struct CONFIG_CONN_HANDLE* next = ch->next;
topo_group_peer_close(ch->request);
u_free(ch); ch_count++;
ch = next;
}
if (ch_count) DEBUG_INFO(DEBUG_CATEGORY_SYS, "[DESTROY] D: closed %d config handles", ch_count);
}
instance->config_conn_handles = NULL;
DEBUG_INFO(DEBUG_CATEGORY_SYS, "[DESTROY] D done — config_conn");
chat_service_stop(instance);
utun_service_stop(instance);
/* Phase E: socket_monitor + auto_socket */
socket_monitor_destroy(instance);
@ -494,18 +458,6 @@ void utun_instance_destroy(struct UTUN_INSTANCE *instance) {
instance->etcp_sockets = NULL;
DEBUG_INFO(DEBUG_CATEGORY_SYS, "[DESTROY] G done — ETCP sockets");
/* Phase H: chat */
chat_headless_control_destroy(instance);
chat_media_startup_backfill_destroy(instance);
dm_mailbox_destroy(instance);
dm_core_destroy(instance);
call_headless_destroy(instance);
call_destroy(instance);
radio_headless_destroy(instance);
chat_sync_destroy(instance);
chat_core_destroy(instance);
DEBUG_INFO(DEBUG_CATEGORY_SYS, "[DESTROY] H done — chat");
/* Phase I: db_sync */
db_sync_destroy(instance);
DEBUG_INFO(DEBUG_CATEGORY_SYS, "[DESTROY] I done — db_sync");
@ -669,7 +621,6 @@ void utun_instance_stop(struct UTUN_INSTANCE *instance) {
instance->running = 0;
// Wakeup main loop using built-in uasync wakeup
if (instance->ua) {
memory_pool_destroy(instance->pkt_pool);
uasync_wakeup(instance->ua);
}
@ -685,8 +636,9 @@ static int mkdir_recursive(const char *path) {
return utun_mkdir(tmp, 0755);
}
int utun_instance_init(struct UTUN_INSTANCE *instance) {
if (!instance) return -1;
int utun_core_start(struct UTUN_INSTANCE *instance) {
if (!instance) { DEBUG_ERROR(DEBUG_CATEGORY_SYS, "core start: NULL instance"); return -1; }
if (instance->core_started) return 0;
struct utun_config* config = instance->config;
if (config) {
@ -697,56 +649,11 @@ int utun_instance_init(struct UTUN_INSTANCE *instance) {
instance->keepalive_interval = ka;
}
// db_sync — распределённая таблица с репликацией
db_sync_init(instance);
// chat — сообщения, каналы, P2P-синхронизация
if (instance->config->global.chatserver_enabled) {
if (!instance->config->global.db_path[0])
snprintf(instance->config->global.db_path, sizeof(instance->config->global.db_path), "/var/lib/utun");
mkdir_recursive(instance->config->global.db_path);
char db_file[512];
snprintf(db_file, sizeof(db_file), "%s/chats.db", instance->config->global.db_path);
if (chat_core_init(instance, db_file) == 0) {
chat_sync_init(instance);
chat_media_backfill(instance);
chat_media_startup_backfill_init(instance);
/* dm_mailbox до dm_core: dm_core_init ставит deliver-cb в mailbox */
dm_mailbox_init(instance);
dm_core_init(instance);
call_init(instance);
}
}
// Set TUN interface in routing module
if (instance->tun) {
routing_set_tun(instance);
DEBUG_INFO(DEBUG_CATEGORY_ROUTING, "TUN interface registered in routing module");
}
// Initialize NAT (after routing_set_tun, before connections)
if (instance->config->global.nat_enabled) {
int nat_ret = nat_transport_init(instance);
if (nat_ret != 0) {
DEBUG_WARN(DEBUG_CATEGORY_NAT, "NAT transport init failed (non-fatal)");
}
}
// Initialize media delivery
if (media_delivery_init(instance) != 0) {
DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "media_delivery init failed (non-fatal)");
}
// Initialize media async engine
instance->media_async = media_async_create();
if (!instance->media_async) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "media_async create failed");
}
// Note: TUN socket is already registered in tun_init()
// Generic replication is a core service; channel instances belong to chat.
if (!instance->db_sync && db_sync_init(instance) < 0) return -1;
// Initialize connections
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "utun_instance_init() calling init_connections() for instance %p", instance);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "core start: initializing transports instance=%p", instance);
int conn_result = init_connections(instance);
@ -774,6 +681,101 @@ int utun_instance_init(struct UTUN_INSTANCE *instance) {
}
}
// Start the main loop
instance->running = 1;
// Initialize NTP time sync (non-fatal if fails)
if (ntp_time_init(instance) != 0) {
DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "NTP time sync init failed (non-fatal)");
}
// Initialize NTP node time exchange (non-fatal)
if (ntp_node_time_init(instance) != 0) {
DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "NTP node time exchange init failed (non-fatal)");
}
broadcast_init_instance(instance);
radio_init_instance(instance);
instance->core_started = 1;
DEBUG_INFO(DEBUG_CATEGORY_SYS, "core started");
return 0;
}
/* Service lifecycle runs on the owning uasync thread, outside module callbacks. */
int utun_service_start(struct UTUN_INSTANCE* instance) {
if (!instance || !instance->core_started || !instance->topo_groups) {
DEBUG_ERROR(DEBUG_CATEGORY_SYS, "UTUN start requires a running core"); return -1;
}
if (instance->utun_started) return 0;
struct utun_config* config = instance->config;
struct UASYNC* ua = instance->ua;
struct TOPO_GROUP* group = topo_groups_create_group(instance->topo_groups, TOPO_GROUP_UTUN, TOPO_GROUP_TYPE_UTUN, NULL);
if (!group) return -1;
instance->conn_mgr = group->conn_mgr;
if (g_tun_init_enabled && config->global.tun_enabled) {
instance->tun = tun_init(ua, config);
if (!instance->tun) {
DEBUG_ERROR(DEBUG_CATEGORY_TUN, "Failed to initialize TUN device");
goto fail;
}
DEBUG_INFO(DEBUG_CATEGORY_TUN, "TUN interface initialized: %s", instance->tun->ifname);
if (config->route_subnets) {
int added = tun_route_add_all(instance->tun->ifindex, instance->tun->ifname, config->route_subnets);
DEBUG_INFO(DEBUG_CATEGORY_TUN, "Added %d system routes for TUN interface", added);
instance->route_subnets = config->route_subnets;
}
} else {
DEBUG_INFO(DEBUG_CATEGORY_TUN, "TUN initialization disabled - skipping TUN device setup");
instance->tun = NULL;
}
if (routing_bind(instance) < 0) goto fail;
if (instance->tun) routing_set_tun(instance);
if (config->global.nat_enabled && nat_transport_init(instance) < 0) goto fail;
if (init_client_connections(instance, config->clients) < 0) goto fail;
instance->utun_started = 1;
DEBUG_INFO(DEBUG_CATEGORY_SYS, "UTUN service started"); return 0;
fail:
DEBUG_ERROR(DEBUG_CATEGORY_SYS, "UTUN service start failed; rolling back");
utun_service_stop(instance); return -1;
}
void utun_service_stop(struct UTUN_INSTANCE* instance) {
if (!instance) return;
instance->utun_started = 0;
if (instance->nat_tr.initialized) nat_transport_destroy(instance);
if (instance->router_bindings.callbacks[ETCP_RT_ID_DATA]) etcp_router_unbind(instance, ETCP_RT_ID_DATA);
if (instance->tun) {
if (instance->tun->ifindex && instance->route_subnets)
tun_route_del_all(instance->tun->ifindex, instance->tun->ifname, instance->route_subnets);
tun_close(instance->tun); instance->tun = NULL;
}
instance->route_subnets = NULL;
while (instance->config_conn_handles) {
struct CONFIG_CONN_HANDLE* owner = instance->config_conn_handles;
instance->config_conn_handles = owner->next;
topo_group_peer_close(owner->request); u_free(owner);
}
if (topo_groups_get_default(instance->topo_groups)) topo_groups_remove_group(instance->topo_groups, TOPO_GROUP_UTUN);
DEBUG_INFO(DEBUG_CATEGORY_SYS, "UTUN service stopped");
}
int chat_service_start(struct UTUN_INSTANCE* instance) {
if (!instance || !instance->core_started || !instance->topo_sqlite_db) {
DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "chat start requires a running core with SQLite"); return -1;
}
if (instance->chat_started) return 0;
char db_file[512];
snprintf(db_file, sizeof(db_file), "%s/chats.db", instance->config->global.db_path);
if (db_sync_enable(instance) < 0) return -1;
topo_group_invite_init(instance);
if (chat_core_init(instance, db_file) < 0 || chat_sync_init(instance) < 0) goto fail;
instance->media_async = media_async_create();
if (!instance->media_async || media_delivery_init(instance) < 0) goto fail;
if (dm_mailbox_init(instance) < 0 || dm_core_init(instance) < 0 || call_init(instance) < 0) goto fail;
chat_media_backfill(instance);
chat_media_startup_backfill_init(instance);
// Initialize headless chat control if configured in [chatserver]
if (instance->config->global.headless_control_bind[0] != '\0') {
char ip_str[64] = "127.0.0.1"; int port = 9999;
@ -816,23 +818,43 @@ int utun_instance_init(struct UTUN_INSTANCE *instance) {
DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "Failed to init headless radio audio socket (non-fatal)");
}
// Start the main loop
instance->running = 1;
// Initialize NTP time sync (non-fatal if fails)
if (ntp_time_init(instance) != 0) {
DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "NTP time sync init failed (non-fatal)");
}
instance->chat_started = 1;
DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "chat service started"); return 0;
fail:
DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "chat service start failed; rolling back");
chat_service_stop(instance); return -1;
}
// Initialize NTP node time exchange (non-fatal)
if (ntp_node_time_init(instance) != 0) {
DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "NTP node time exchange init failed (non-fatal)");
void chat_service_stop(struct UTUN_INSTANCE* instance) {
if (!instance) return;
instance->chat_started = 0;
chat_headless_control_destroy(instance);
call_headless_destroy(instance); radio_headless_destroy(instance);
chat_media_startup_backfill_destroy(instance);
media_async_destroy(instance->media_async); instance->media_async = NULL;
media_delivery_destroy(instance);
call_destroy(instance); dm_core_destroy(instance); dm_mailbox_destroy(instance);
etcp_unbind(instance, ETCP_RT_ID_GROUP_INVITE);
topo_group_invite_cancel_all(instance);
chat_sync_destroy(instance);
struct ll_entry* entry = instance->topo_groups ? instance->topo_groups->group_list->head : NULL;
while (entry) {
struct TOPO_GROUP* group = (struct TOPO_GROUP*)entry; entry = entry->next;
if (group->group_type != TOPO_GROUP_TYPE_CHAT) continue;
struct DB_SYNC_INSTANCE* sync = db_sync_instance_find(instance, group->group_id);
if (sync) db_sync_instance_remove(sync);
topo_groups_remove_group(instance->topo_groups, group->group_id);
}
chat_core_destroy(instance);
DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "chat service stopped");
}
broadcast_init_instance(instance);
radio_init_instance(instance);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "Connections initialized successfully");
/* Application convenience entry point; embedders may start each service explicitly. */
int utun_instance_init(struct UTUN_INSTANCE* instance) {
if (utun_core_start(instance) < 0) return -1;
if (instance->topo_groups && utun_service_start(instance) < 0) return -1;
if (instance->config->global.chatserver_enabled && chat_service_start(instance) < 0) return -1;
return 0;
}
@ -1009,7 +1031,7 @@ struct UTUN_INSTANCE *utun_instance_reload(struct UTUN_INSTANCE *instance, struc
}
os = next;
}
int client_result = init_client_connections(instance, new_config->clients);
int client_result = instance->utun_started ? init_client_connections(instance, new_config->clients) : 0;
if (client_result < 0)
DEBUG_ERROR(DEBUG_CATEGORY_CONFIG, "reload: failed to replace client handles; previous connections retained");
// Reload networks: clear and repopulate

9
src/utun_instance.h

@ -111,6 +111,8 @@ struct peer_sleep_cbk_entry {
// uTun instance configuration
struct UTUN_INSTANCE {
uint8_t core_started, utun_started, chat_started;
// Identification
char name[MAX_CONN_NAME_LEN]; // Instance name from config
@ -264,6 +266,13 @@ struct UTUN_INSTANCE* utun_instance_create(struct UASYNC* ua, const char* config
struct UTUN_INSTANCE* utun_instance_create_from_config(struct UASYNC* ua, struct utun_config* config);
struct UTUN_INSTANCE* utun_instance_create_from_str(struct UASYNC* ua, const char* config_text);
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. */
int utun_core_start(struct UTUN_INSTANCE* instance);
int utun_service_start(struct UTUN_INSTANCE* instance);
void utun_service_stop(struct UTUN_INSTANCE* instance);
/* Convenience: core + UTUN + configured chat. */
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);
void utun_instance_run(struct UTUN_INSTANCE *instance);

32
src/video/video.c

@ -1,5 +1,6 @@
// video.c — пробинг и транскодирование видео через FFmpeg (опционально, HAVE_FFMPEG)
#include "video.h"
#include "../media_async/media_async.h"
#include "../../lib/debug_config.h"
#include "../../lib/mem.h"
#include "../../lib/u_async.h"
@ -44,11 +45,10 @@ int video_probe(const char* path) { (void)path; return 0; }
int video_get_info(const char* path, struct video_info* out) { (void)path; (void)out; return -1; }
int video_needs_transcode(const struct video_info* info) { (void)info; return 1; }
int video_transcode_start(struct UASYNC* ua, const char* src, const char* dst,
int video_transcode_start(struct media_async* owner, struct UASYNC* ua, const char* src, const char* dst,
void (*done)(void* arg, int err), void* arg) {
(void)ua; (void)src; (void)dst;
(void)owner; (void)ua; (void)src; (void)dst; (void)done; (void)arg;
DEBUG_ERROR(DEBUG_CATEGORY_VIDEO, "video: built without FFmpeg (HAVE_FFMPEG undefined)");
if (done) done(arg, -1);
return -1;
}
@ -525,15 +525,15 @@ static void vc_cleanup(struct vc_ctx* ctx) {
if (ctx->ifmt) { avformat_close_input(&ctx->ifmt); ctx->ifmt = NULL; }
}
static void vc_done_trampoline(void* raw) {
static void vc_done_trampoline(void* raw, int error) {
struct vc_ctx* ctx = (struct vc_ctx*)raw;
vc_active_set(0);
vc_progress_set(-1);
if (ctx->done) ctx->done(ctx->arg, ctx->err);
if (ctx->done) ctx->done(ctx->arg, error ? error : ctx->err);
u_free(ctx);
}
static void* vc_thread_entry(void* raw) {
static void vc_thread_entry(void* raw) {
struct vc_ctx* ctx = (struct vc_ctx*)raw;
struct video_info info;
@ -567,13 +567,12 @@ static void* vc_thread_entry(void* raw) {
vc_cleanup(ctx);
out:
uasync_post(ctx->ua, vc_done_trampoline, ctx);
return NULL;
return;
}
int video_transcode_start(struct UASYNC* ua, const char* src, const char* dst,
int video_transcode_start(struct media_async* owner, struct UASYNC* ua, const char* src, const char* dst,
void (*done)(void* arg, int err), void* arg) {
if (!ua || !src || !dst || !done) {
if (!owner || !ua || !src || !dst || !done) {
DEBUG_ERROR(DEBUG_CATEGORY_VIDEO, "video: transcode_start invalid args");
return -1;
}
@ -583,7 +582,7 @@ int video_transcode_start(struct UASYNC* ua, const char* src, const char* dst,
}
struct vc_ctx* ctx = u_calloc(1, sizeof(*ctx));
if (!ctx) return -1;
if (!ctx) { DEBUG_ERROR(DEBUG_CATEGORY_VIDEO, "video: transcode context allocation failed"); return -1; }
snprintf(ctx->src, sizeof(ctx->src), "%s", src);
snprintf(ctx->dst, sizeof(ctx->dst), "%s", dst);
ctx->ua = ua;
@ -597,16 +596,7 @@ int video_transcode_start(struct UASYNC* ua, const char* src, const char* dst,
vc_active_set(1);
vc_progress_set(0);
pthread_t tid;
int rc = pthread_create(&tid, NULL, vc_thread_entry, ctx);
if (rc != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_VIDEO, "video: pthread_create failed rc=%d", rc);
vc_active_set(0);
vc_progress_set(-1);
u_free(ctx);
return -1;
}
pthread_detach(tid);
media_async_submit(owner, ua, vc_thread_entry, ctx, vc_done_trampoline, ctx);
return 0;
}

7
src/video/video.h

@ -51,9 +51,12 @@ int video_needs_transcode(const struct video_info* info);
* Создаёт отдельный pthread, который либо перекодирует src->dst в целевой
* профиль, либо (если транскод не нужен) копирует src->dst как есть.
* Прогресс обновляется атомарно, читается через video_transcode_get_progress().
* По завершении done(arg, err) выполняется в uasync-потоке (через uasync_post).
* При возврате 0 done(arg, err) вызывается ровно один раз в uasync-потоке,
* в том числе при отмене owner. Callback может выполниться внутри start.
* При возврате -1 callback не вызывается, arg остаётся у вызывающего.
*/
int video_transcode_start(struct UASYNC* ua,
struct media_async;
int video_transcode_start(struct media_async* owner, struct UASYNC* ua,
const char* src, const char* dst,
void (*done)(void* arg, int err), void* arg);

5
tests/Makefile.am

@ -65,6 +65,7 @@ check_PROGRAMS = \
test_etcp_connect \
test_node_conn_direct \
test_group_ownership \
test_services \
test_group_exchange \
test_group_recovery \
test_node_snapshot \
@ -622,3 +623,7 @@ test_etcp_session_LDADD = $(LIBUTUN) $(COMMON_LIBS)
test_etcp_router_tcp_SOURCES = test_etcp_router_tcp.c
test_etcp_router_tcp_LDADD = $(LIBUTUN) $(COMMON_LIBS)
test_services_SOURCES = test_services.c
test_services_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/src/routing_layer -I$(top_srcdir)/src/transport_layer -I$(top_srcdir)/lib
test_services_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS)

25
tests/test_chat_join_e2e.c

@ -27,6 +27,7 @@
#include "../routing_layer/topo_node.h"
#include "../utun_instance.h"
#include "../transport_layer/etcp.h"
#include "../transport_layer/node_conn_direct.h"
#include "../transport_layer/secure_channel.h"
#include "../ntp_time.h"
#include "../lib/u_async.h"
@ -184,22 +185,13 @@ static void write_configs(const char* dir) {
snprintf(dbd, sizeof(dbd), "%s/db%s", dir, i == IDX_A ? "a" : i == IDX_C ? "c" : "j");
utun_mkdir(dbd, 0755);
char client[256] = "";
if (i == IDX_C && !g_sh.degraded) {
char apub[65]; to_hex(g_sh.x_pub[IDX_A], 32, apub);
snprintf(client, sizeof(client),
"[client: to_a]\nkeepalive=1\npeer_public_key=%s\nlink=s1:127.0.0.1:%d\n",
apub, g_sh.port[IDX_A]);
}
wf(path,
"[global]\ntun_ip=10.98.%d.1/24\ntun_ifname=tun%d0\ndb_path=%s/db%s\n"
"my_private_key=%s\nmy_public_key=%s\n"
"[server: s1]\naddr=127.0.0.1:%d\ntype=public\n"
"%s"
"[chatserver]\nstorage_autoload=0\n[allowed_keys]\nallow_all=1\n",
i, i, dir, i == IDX_A ? "a" : i == IDX_C ? "c" : "j",
priv, pub, g_sh.port[i], client);
priv, pub, g_sh.port[i]);
}
}
@ -448,7 +440,9 @@ static int child_main(const char* role, const char* dir, int invalid_key) {
strcmp(role, "a") == 0 ? "a" : strcmp(role, "c") == 0 ? "c" : "j");
struct UTUN_INSTANCE* inst = utun_instance_create(ua, cfg);
if (!inst) { fprintf(stderr, "%s: create failed\n", role); uasync_destroy(ua, 0); return 1; }
utun_instance_init(inst);
if (utun_core_start(inst) < 0 || chat_service_start(inst) < 0) {
utun_instance_destroy(inst); uasync_destroy(ua, 0); return 1;
}
chat_event_set_handler(inst, join_event_handler);
if (strcmp(role, "c") == 0) {
scenario_dir = dir;
@ -456,6 +450,14 @@ static int child_main(const char* role, const char* dir, int invalid_key) {
etcp_router_bind(inst, ETCP_RT_ID_JOIN, delay_registration);
}
/* Fixture bootstrap between the existing members uses core NCD, without UTUN. */
struct NODE_CONN_DIRECT* bootstrap = NULL;
if (!strcmp(role, "c") && !g_sh.degraded) {
struct TOPO_ADDR4 address = { .addr = {127, 0, 0, 1}, .port = g_sh.port[IDX_A], .protocol = TOPO_PROTO_UDP };
struct TOPO_NODE node = { .node_id = g_sh.nid[IDX_A], .v4_addrs = &address };
memcpy(node.public_key, g_sh.x_pub[IDX_A], 32);
if (node_conn_direct_open_node(inst, node.node_id, NULL, NULL, &bootstrap, &node, NULL) == NCD_ERR) return 1;
}
int rc = 1;
uint64_t ch_num = strtoull(g_sh.ch_id, NULL, 10);
sqlite3* db = inst->topo_sqlite_db;
@ -589,6 +591,7 @@ static int child_main(const char* role, const char* dir, int invalid_key) {
}
out:
node_conn_direct_close(bootstrap);
if (early_link || delayed_registration || early_bgp_membership) rc = 1;
fprintf(stderr, "%s: %s\n", role, rc == 0 ? "OK" : "FAIL");
fflush(stdout); fflush(stderr);

2
tests/test_conn_mgr_phases.c

@ -202,7 +202,7 @@ static int scenario_setup(void) {
inst_b = utun_instance_create(ua, cfg_b);
if (!inst_c || !inst_a || !inst_b) { fail("instance create"); return -1; }
if (init_connections(inst_c) != 0 || init_connections(inst_a) != 0 || init_connections(inst_b) != 0) {
if (utun_instance_init(inst_c) != 0 || utun_instance_init(inst_a) != 0 || utun_instance_init(inst_b) != 0) {
fail("init_connections"); return -1;
}
nid_a = inst_a->node_id; nid_b = inst_b->node_id; nid_c = inst_c->node_id;

4
tests/test_etcp_100_packets.c

@ -409,7 +409,7 @@ int main() {
printf("Creating server...\n");
ua = uasync_create();
server_instance = utun_instance_create(ua, server_config_path);
if (!server_instance || init_connections(server_instance) < 0) {
if (!server_instance || utun_instance_init(server_instance) < 0) {
printf("Failed to create server\n");
return 1;
}
@ -417,7 +417,7 @@ int main() {
printf("Creating client...\n");
client_instance = utun_instance_create(ua, client_config_path);
if (!client_instance || init_connections(client_instance) < 0) {
if (!client_instance || utun_instance_init(client_instance) < 0) {
printf("Failed to create client\n");
return 1;
}

4
tests/test_etcp_api.c

@ -564,7 +564,7 @@ int main() {
printf("Creating server...\n");
ua = uasync_create();
server_instance = utun_instance_create(ua, server_config_path);
if (!server_instance || init_connections(server_instance) < 0) {
if (!server_instance || utun_instance_init(server_instance) < 0) {
printf("Failed to create server\n");
return 1;
}
@ -572,7 +572,7 @@ int main() {
printf("Creating client...\n");
client_instance = utun_instance_create(ua, client_config_path);
if (!client_instance || init_connections(client_instance) < 0) {
if (!client_instance || utun_instance_init(client_instance) < 0) {
printf("Failed to create client\n");
return 1;
}

4
tests/test_etcp_ping.c

@ -166,12 +166,12 @@ int main(void) {
goto cleanup;
}
int ret;
ret = init_connections(server_instance);
ret = utun_instance_init(server_instance);
if (ret != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "Server init_connections failed with code %d", ret);
goto cleanup;
}
ret = init_connections(client_instance);
ret = utun_instance_init(client_instance);
if (ret != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "Client init_connections failed with code %d", ret);
goto cleanup;

4
tests/test_etcp_simple_traffic.c

@ -211,14 +211,14 @@ int main(void) {
ua = uasync_create();
server_instance = utun_instance_create(ua, server_config_path);
if (!server_instance || init_connections(server_instance) < 0) {
if (!server_instance || utun_instance_init(server_instance) < 0) {
fprintf(stderr, "Failed to create server\n");
return 1;
}
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "Server ready (node_id=%llx)", (unsigned long long)server_instance->node_id);
client_instance = utun_instance_create(ua, client_config_path);
if (!client_instance || init_connections(client_instance) < 0) {
if (!client_instance || utun_instance_init(client_instance) < 0) {
fprintf(stderr, "Failed to create client\n");
utun_instance_destroy(server_instance);
return 1;

2
tests/test_group_recovery.c

@ -284,7 +284,7 @@ static void persistent_request(void) {
static void manual_planner_cancel(void) {
struct fixture f; create(&f);
assert(sqlite3_open(":memory:", &f.inst->topo_sqlite_db) == SQLITE_OK);
assert(f.inst->topo_sqlite_db);
assert(topo_node_sqlite_init(f.inst->topo_sqlite_db) == 0);
f.group->group_type = TOPO_GROUP_TYPE_CHAT; strcpy(f.group->channel_id, "42");
/* The membership table is the only authorization source. */

4
tests/test_invite_group_create.c

@ -83,7 +83,7 @@ int main(void) {
struct UASYNC* ua = uasync_create();
struct UTUN_INSTANCE* a = utun_instance_create(ua, ca);
if (!a) { fail("create A"); goto clean; }
utun_instance_init(a);
utun_instance_init(a); topo_group_invite_init(a);
/* ── Add a TCP socket (simulating Android STCP server) ── */
{
@ -104,7 +104,7 @@ int main(void) {
/* Get B's node_id and pubkey from its instance */
struct UTUN_INSTANCE* b = utun_instance_create(ua, cb);
if (!b) { fail("create B"); a->running = 0; utun_instance_destroy(a); uasync_destroy(ua, 0); goto clean; }
utun_instance_init(b);
utun_instance_init(b); topo_group_invite_init(b);
/* ── Set up B as member of group 0x7777777700000001 ── */
{

1
tests/test_media_delivery_chat.c

@ -574,6 +574,7 @@ int main(void) {
g_inst[i] = utun_instance_create(g_ua, g_cfg[i]);
if (!g_inst[i]) { printf("FAIL: create n%d\n", i); goto done; }
utun_instance_init(g_inst[i]);
media_delivery_init(g_inst[i]);
member_sync_init(g_inst[i]);
}

1
tests/test_media_delivery_full.c

@ -468,6 +468,7 @@ int main(void) {
g_inst[i] = utun_instance_create(g_ua, g_cfg[i]);
if (!g_inst[i]) { printf("FAIL: create n%d\n", i); goto done; }
utun_instance_init(g_inst[i]);
media_delivery_init(g_inst[i]);
if (member_sync_init(g_inst[i]) < 0) { FAIL("member_sync_init n%d", i); goto done; }
}

1
tests/test_media_delivery_integration.c

@ -284,6 +284,7 @@ int main(void) {
g_a = utun_instance_create(g_ua, g_ca); g_b = utun_instance_create(g_ua, g_cb);
if (!g_a || !g_b) { printf("FAIL: instance create\n"); goto done; }
utun_instance_init(g_a); utun_instance_init(g_b);
media_delivery_init(g_a); media_delivery_init(g_b);
/* make A supernode */
media_delivery_set_supernode(g_a, 1);

5
tests/test_nat_detection.c

@ -212,13 +212,14 @@ int main(void) {
}
// Разрешаем NAT check для localhost (для теста)
if (topo_groups_get_default(inst_s->topo_groups)) topo_group_set_nat_check_local(topo_groups_get_default(inst_s->topo_groups), 1);
if (init_connections(inst_s) != 0 || init_connections(inst_c1) != 0 || init_connections(inst_c2) != 0) {
if (utun_instance_init(inst_s) != 0 || utun_instance_init(inst_c1) != 0 || utun_instance_init(inst_c2) != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "Failed to init connections");
goto cleanup;
}
if (topo_groups_get_default(inst_s->topo_groups)) topo_group_set_nat_check_local(topo_groups_get_default(inst_s->topo_groups), 1);
g_node_id_s = inst_s->node_id; g_node_id_c1 = inst_c1->node_id; g_node_id_c2 = inst_c2->node_id;
test_timeout_id = uasync_set_timeout(ua, TEST_TIMEOUT_MS * 10, NULL, test_timeout_cb, "test_nat");

4
tests/test_pkt_normalizer_etcp.c

@ -250,10 +250,10 @@ int main(void) {
ua = uasync_create();
server_instance = utun_instance_create(ua, server_config_path);
if (!server_instance || init_connections(server_instance) < 0) { fprintf(stderr, "Failed to create server\n"); return 1; }
if (!server_instance || utun_instance_init(server_instance) < 0) { fprintf(stderr, "Failed to create server\n"); return 1; }
client_instance = utun_instance_create(ua, client_config_path);
if (!client_instance || init_connections(client_instance) < 0) { fprintf(stderr, "Failed to create client\n"); return 1; }
if (!client_instance || utun_instance_init(client_instance) < 0) { fprintf(stderr, "Failed to create client\n"); return 1; }
packet_timeout_id = uasync_set_timeout(ua, 500, NULL, monitor_and_send, "test_monitor");
void* global_timeout_id = uasync_set_timeout(ua, TEST_TIMEOUT_TB * 10, NULL, test_timeout, "test_timeout");

2
tests/test_route_ping.c

@ -218,7 +218,7 @@ int main(void) {
goto cleanup;
}
if (init_connections(inst_a) != 0 || init_connections(inst_b) != 0 || init_connections(inst_c) != 0) {
if (utun_instance_init(inst_a) != 0 || utun_instance_init(inst_b) != 0 || utun_instance_init(inst_c) != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "Failed to init connections");
goto cleanup;
}

173
tests/test_services.c

@ -0,0 +1,173 @@
#include <assert.h>
#include <string.h>
#include <stdlib.h>
#include <unistd.h>
#include "utun_instance.h"
#include "config_parser.h"
#include "topo_group.h"
#include "topo_node_sqlite.h"
#include "route_lib.h"
#include "etcp.h"
#include "etcp_router.h"
#include "media_delivery/media_index.h"
#include "media_delivery/media_delivery_proto.h"
#include "node_conn_direct.h"
#include "chat/chat_core.h"
#include "chat/db_sync.h"
#include "media_async/media_async.h"
#include "../lib/debug_config.h"
struct job { int worked, done, error; };
static void work(void* arg) { ((struct job*)arg)->worked = 1; }
static void done(void* arg, int err) { struct job* job = arg; job->done++; job->error = err; }
static void ready(struct ll_queue* queue, void* arg) { (void)queue; (*(int*)arg)++; }
/* Pending transit shares the physical connection but belongs to its group. */
static void pending_transit(struct ETCP_CONN* conn, uint64_t group_id, int* calls) {
if (!conn->transit_queues) conn->transit_queues = queue_new(conn->instance->ua, 16, 0, 24, "test_transit");
assert(conn->transit_queues);
struct TRANSIT_QUEUE* transit = (struct TRANSIT_QUEUE*)queue_entry_new(sizeof(*transit) - sizeof(struct ll_entry));
assert(transit); transit->group_id = group_id; transit->conn = conn;
transit->src_node_id = 123; transit->dst_node_id = conn->peer_node_id;
transit->q = queue_new(conn->instance->ua, 0, 0, 0, "test_transit_packets"); assert(transit->q);
assert(queue_data_put_with_index(conn->transit_queues, &transit->ll) == 0);
struct ll_entry* packet = ll_alloc_lldgram(1); assert(packet); packet->dgram[0] = 0; packet->len = 1;
assert(queue_data_put(transit->q, packet) == 0);
assert(queue_waiter_wait(conn->send_input_q, &transit->waiter, ready, calls) == 0);
}
static struct TOPO_GROUP* chat_group(struct UTUN_INSTANCE* inst) {
for (struct ll_entry* e = inst->topo_groups->group_list->head; e; e = e->next) {
struct TOPO_GROUP* group = (struct TOPO_GROUP*)e;
if (group->group_type == TOPO_GROUP_TYPE_CHAT) return group;
}
return NULL;
}
int main(void) {
debug_config_init(); debug_set_level(DEBUG_LEVEL_INFO);
struct UASYNC* ua = uasync_create(); assert(ua);
struct UTUN_INSTANCE* inst = utun_instance_create_from_str(ua,
"[global]\nmy_private_key=704f2e012c8fa8768130cb0f988a997dccb628372bc5ceccacc78dcbfec5916f\n"
"my_public_key=b3193173def895bd0fcea6f86af077c7d77216f10395275f627ac18242ec0f01\n"
"[routing]\nmy_subnet=192.168.42.0/24\n[server: udp]\naddr=127.0.0.1:0\ntype=public\n"
"[chatserver]\ngroup_autoconnect=0\nstorage_autoload=0\n");
assert(inst && !topo_groups_get_default(inst->topo_groups) && !inst->chat_core);
inst->config->global.db_sync_enabled = 0;
assert(utun_core_start(inst) == 0 && utun_core_start(inst) == 0);
struct ETCP_SOCKET* socket = inst->etcp_sockets;
sqlite3* db = inst->topo_sqlite_db;
struct TOPO_NODE* identity = topo_node_registry_find(inst->topo_groups, inst->node_id);
assert(identity && identity->v4_addrs);
assert(chat_service_start(inst) == 0 && chat_service_start(inst) == 0);
assert(!topo_groups_get_default(inst->topo_groups));
chat_core_create_channel_auto(inst, "lifecycle");
struct TOPO_GROUP* chat = chat_group(inst); assert(chat);
uint64_t chat_id = chat->group_id;
assert(utun_service_start(inst) == 0 && utun_service_start(inst) == 0);
assert(inst->rt->count > 0);
struct TOPO_GROUP* vpn = topo_groups_get_default(inst->topo_groups); assert(vpn);
struct SC_MYKEYS keys; assert(sc_generate_keypair(&keys) == SC_OK);
struct TOPO_ADDR4 address = { .addr = {127, 0, 0, 1}, .port = 9, .protocol = TOPO_PROTO_UDP };
struct TOPO_NODE node = { .v4_addrs = &address };
memcpy(node.public_key, keys.public_key, 32); node.node_id = sc_derive_node_id_from_pubkey(keys.public_key);
uint8_t sig[64] = {1};
assert(topo_node_sqlite_member_block_put(db, chat->channel_id, node.node_id, sig, 1, sig, 1,
keys.public_key, keys.public_key, "{}", NULL, 0, sig) == 0);
struct NODE_CONN_DIRECT* seed = NULL;
assert(node_conn_direct_open_node(inst, node.node_id, NULL, NULL, &seed, &node, NULL) >= 0);
struct ETCP_CONN* conn = node_conn_direct_get_conn(seed); assert(conn);
struct TOPO_PEER_REQUEST *chat_request = NULL, *vpn_request = NULL;
assert(topo_group_peer_open(chat, node.node_id, &chat_request) == 0);
assert(topo_group_peer_open(vpn, node.node_id, &vpn_request) == 0);
node_conn_direct_close(seed);
utun_service_stop(inst); utun_service_stop(inst);
assert(!conn->close_requested && topo_group_peer_conn(chat_request) == conn);
assert(topo_group_peer_phase(vpn_request) == TOPO_PEER_FAILED && inst->rt->count == 0);
topo_group_peer_close(vpn_request);
assert(inst->etcp_sockets == socket && inst->topo_sqlite_db == db && inst->chat_started);
assert(topo_node_registry_find(inst->topo_groups, inst->node_id) == identity);
assert(topo_node_update_my_addresses(inst) >= 0);
assert(utun_service_start(inst) == 0);
vpn = topo_groups_get_default(inst->topo_groups);
assert(topo_group_peer_open(vpn, node.node_id, &vpn_request) == 0);
/* Stop a real stream while its router send queue is blocked. */
char path[] = "/tmp/utun-service-stream-XXXXXX";
int fd = mkstemp(path); assert(fd >= 0);
char bytes[8192] = {0}; assert(write(fd, bytes, sizeof(bytes)) == sizeof(bytes)); close(fd);
chat_core_save_ui_state(inst, "media_base", "/tmp");
uint8_t media_id[16] = {1}, block_id[16] = {2}, hash[32] = {3};
assert(media_index_register_downloaded(db, media_id, block_id, hash, chat->channel_id,
strrchr(path, '/') + 1, inst->node_id, sizeof(bytes), sizeof(bytes), 0, 0) == 0);
struct ETCP_ROUTER_CONN* routed = etcp_router_conn_get(inst, chat_id, node.node_id, ETCP_RT_ID_MEDIA_DELIVERY);
assert(routed);
struct queue_waiter_handle waiter = {0}; int ready_calls = 0;
struct ll_queue* old_queue = routed->send_q;
queue_set_threshold(old_queue, -1, 0);
etcp_router_on_send_ready(inst, chat_id, node.node_id, ETCP_RT_ID_MEDIA_DELIVERY, &waiter, ready, &ready_calls);
assert(waiter.internal);
etcp_router_conn_restart(inst, chat_id, node.node_id, ETCP_RT_ID_MEDIA_DELIVERY);
assert(routed->send_q != old_queue);
queue_set_threshold(routed->send_q, -1, 0);
etcp_router_on_send_ready(inst, chat_id, node.node_id, ETCP_RT_ID_MEDIA_DELIVERY, &waiter, ready, &ready_calls);
assert(!old_queue->waiter_head && waiter.internal);
etcp_router_cancel_send_ready(inst, chat_id, node.node_id, ETCP_RT_ID_MEDIA_DELIVERY, &waiter);
assert(!routed->send_q->waiter_head && !waiter.internal);
queue_set_threshold(routed->send_q, 0, 0);
routed->send_q->waiter_defer = 1;
etcp_router_on_send_ready(inst, chat_id, node.node_id, ETCP_RT_ID_MEDIA_DELIVERY, &waiter, ready, &ready_calls);
assert(waiter.call_soon_id);
etcp_router_cancel_send_ready(inst, chat_id, node.node_id, ETCP_RT_ID_MEDIA_DELIVERY, &waiter);
struct media_pkt_block_req request = { .subcmd = MEDIA_SUBCMD_BLOCK_REQ, .group_id = chat_id };
memcpy(request.media_id, media_id, 16); memcpy(request.block_id, block_id, 16);
struct ll_entry* packet = ll_alloc_lldgram(ROUTER_SVC_PAYLOAD_OFF + sizeof(request)); assert(packet);
memset(packet->dgram, 0, ROUTER_SVC_PAYLOAD_OFF);
packet->dgram[0] = ETCP_RT_ID_MEDIA_DELIVERY;
memcpy(packet->dgram + ROUTER_SVC_SRC_OFF, &node.node_id, 8);
memcpy(packet->dgram + ROUTER_SVC_GROUP_OFF, &chat_id, 8);
memcpy(packet->dgram + ROUTER_SVC_PAYLOAD_OFF, &request, sizeof(request));
packet->len = ROUTER_SVC_PAYLOAD_OFF + sizeof(request);
inst->router_bindings.callbacks[ETCP_RT_ID_MEDIA_DELIVERY](conn, packet);
assert(inst->md.streams && inst->md.active_streams == 1);
struct job cancelled = {0};
media_async_submit(inst->media_async, ua, work, &cancelled, done, &cancelled);
queue_set_threshold(conn->send_input_q, -1, 0);
pending_transit(conn, chat_id, &ready_calls);
pending_transit(conn, vpn->group_id, &ready_calls);
chat_service_stop(inst); chat_service_stop(inst);
assert(queue_entry_count(conn->transit_queues) == 1);
assert(((struct TRANSIT_QUEUE*)conn->transit_queues->head)->group_id == vpn->group_id);
assert(!inst->md.streams && routed->closed); unlink(path);
assert(cancelled.worked && cancelled.done == 1 && cancelled.error == MEDIA_ASYNC_CANCELLED);
assert(!chat_group(inst) && !inst->chat_core && !inst->chat_sync && !inst->media_async && !inst->md.initialized);
assert(!db_sync_instance_find(inst, chat_id));
assert(!conn->close_requested && topo_group_peer_conn(vpn_request) == conn);
assert(topo_group_peer_phase(chat_request) == TOPO_PEER_FAILED);
topo_group_peer_close(chat_request);
/* Remove the other group's pending transit before processing the event loop. */
etcp_router_close_group(inst, vpn->group_id);
assert(queue_entry_count(conn->transit_queues) == 0);
queue_set_threshold(conn->send_input_q, 0, 0);
for (int i = 0; i < 3; i++) {
assert(chat_service_start(inst) == 0 && chat_group(inst)->group_id == chat_id);
assert(db_sync_instance_find(inst, chat_id));
chat_service_stop(inst);
assert(inst->etcp_sockets == socket && inst->topo_sqlite_db == db && inst->utun_started);
uasync_poll(ua, 0);
}
topo_group_peer_close(vpn_request);
utun_service_stop(inst);
struct job completed = {0};
struct media_async* jobs = media_async_create(); assert(jobs);
media_async_submit(jobs, ua, work, &completed, done, &completed);
while (!completed.done) uasync_poll(ua, 100);
assert(completed.worked && completed.done == 1 && !completed.error);
media_async_destroy(jobs);
utun_instance_destroy(inst); uasync_poll(ua, 0);
assert(!ready_calls);
assert(ua->timer_alloc_count == ua->timer_free_count && ua->socket_alloc_count == ua->socket_free_count);
uasync_destroy(ua, 0);
return 0;
}

12
tools/chatgui-android/libutun_lite/instance_lite.c

@ -299,11 +299,9 @@ static void* instance_thread(void* arg) {
IL_LOGI("instance created, node_id=0x%016llx", (unsigned long long)g_inst->node_id);
/* Shared SQLite DB открывается внутри utun_instance_init() → topo_groups_init()
(topo_group.c). Двойное открытие здесь давало утечку соединения на каждый рестарт. */
/* Full init: db_sync, chat_core, chat_sync, init_connections */
if (utun_instance_init(g_inst) != 0) {
IL_LOGE("utun_instance_init failed");
/* Core owns SQLite and transport; this frontend starts only the chat service. */
if (utun_core_start(g_inst) != 0 || chat_service_start(g_inst) != 0) {
IL_LOGE("core/chat start failed");
utun_instance_destroy(g_inst);
g_inst = NULL;
standby_deinit();
@ -386,8 +384,8 @@ static void* instance_thread(void* arg) {
if (!g_inst) { IL_LOGE("poll exit: create_from_config failed"); break; }
chat_event_set_handler(g_inst, chat_event_forward);
/* Shared SQLite DB открывается внутри utun_instance_init() (topo_group.c). */
if (utun_instance_init(g_inst) != 0) { IL_LOGE("poll exit: utun_instance_init failed"); break; }
/* Recreate core and chat independently of UTUN. */
if (utun_core_start(g_inst) != 0 || chat_service_start(g_inst) != 0) { IL_LOGE("poll exit: core/chat start failed"); break; }
chat_core_sync_my_addresses(g_inst);
fire_local_sockets_event();
if (g_inst->config->global.name[0]) chat_core_update_my_name(g_inst, g_inst->config->global.name);

15
tools/chatgui/transport/utun_node.cpp

@ -213,7 +213,7 @@ void UtunNode::runLoop() {
gui_bridge_set_uasync(ua);
gui_bridge_set_inst(m_instance);
if (utun_instance_init(m_instance) != 0) {
if (utun_core_start(m_instance) != 0) {
QMetaObject::invokeMethod(this, [this] { emit error("utun_instance_init failed"); });
utun_instance_destroy(m_instance);
m_instance = nullptr;
@ -234,14 +234,11 @@ void UtunNode::runLoop() {
gui_bridge_post(type, data, len);
});
/* Initialize chat_core (DB) and chat_sync (channel/message P2P sync) */
chat_core_init(m_instance, QString(m_dbPath + "/chats.db").toUtf8().constData());
chat_sync_init(m_instance);
chat_media_backfill(m_instance);
chat_media_startup_backfill_init(m_instance);
dm_mailbox_init(m_instance);
dm_core_init(m_instance);
call_init(m_instance);
if (chat_service_start(m_instance) != 0) {
QMetaObject::invokeMethod(this, [this] { emit error("chat_service_start failed"); });
utun_instance_destroy(m_instance); m_instance = nullptr;
uasync_destroy(ua, 0); m_ua = nullptr; return;
}
chat_core_sync_my_addresses(m_instance);
DEBUG_DEBUG(DEBUG_CATEGORY_GENERAL, "chat_core + chat_sync initialized");

Loading…
Cancel
Save