diff --git a/doc/service_lifecycle.md b/doc/service_lifecycle.md new file mode 100644 index 00000000..7408d890 --- /dev/null +++ b/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-сервиса. diff --git a/src/chat/chat_channel.c b/src/chat/chat_channel.c index 733d401f..aaaee7ad 100644 --- a/src/chat/chat_channel.c +++ b/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); diff --git a/src/chat/chat_core.c b/src/chat/chat_core.c index c2bb1f47..6d9c153e 100644 --- a/src/chat/chat_core.c +++ b/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; diff --git a/src/chat/chat_core.h b/src/chat/chat_core.h index 2431e965..410924a7 100644 --- a/src/chat/chat_core.h +++ b/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); diff --git a/src/chat/chat_core_priv.h b/src/chat/chat_core_priv.h index c4a3bf7e..2047b7aa 100644 --- a/src/chat/chat_core_priv.h +++ b/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) ── */ diff --git a/src/chat/chat_msg.c b/src/chat/chat_msg.c index e445757c..c2391663 100644 --- a/src/chat/chat_msg.c +++ b/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); diff --git a/src/chat/chat_sync.c b/src/chat/chat_sync.c index ee9778a4..d92f0966 100644 --- a/src/chat/chat_sync.c +++ b/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; diff --git a/src/chat/chat_whisper.c b/src/chat/chat_whisper.c index 3a8268a4..7eb52752 100644 --- a/src/chat/chat_whisper.c +++ b/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; diff --git a/src/chat/db_sync.c b/src/chat/db_sync.c index d88bc7cb..9ee34083 100644 --- a/src/chat/db_sync.c +++ b/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); diff --git a/src/chat/db_sync.h b/src/chat/db_sync.h index d25213da..3609284b 100644 --- a/src/chat/db_sync.h +++ b/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 diff --git a/src/media_async/media_async.c b/src/media_async/media_async.c index f2033306..b3ebd781 100644 --- a/src/media_async/media_async.c +++ b/src/media_async/media_async.c @@ -19,70 +19,87 @@ #include #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; - ma_work_fn work; - void* data; - ma_done_fn done; - void* arg; + 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) ─── */ diff --git a/src/media_async/media_async.h b/src/media_async/media_async.h index 69a204c9..f833ba04 100644 --- a/src/media_async/media_async.h +++ b/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, diff --git a/src/media_delivery/media_delivery.c b/src/media_delivery/media_delivery.c index 1f54f699..a923a413 100644 --- a/src/media_delivery/media_delivery.c +++ b/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; diff --git a/src/media_delivery/media_delivery.h b/src/media_delivery/media_delivery.h index f44fd0bc..076cd9be 100644 --- a/src/media_delivery/media_delivery.h +++ b/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_* diff --git a/src/media_delivery/media_index.c b/src/media_delivery/media_index.c index e5e143d5..bdc592a3 100644 --- a/src/media_delivery/media_index.c +++ b/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, diff --git a/src/routing_layer/etcp_router.c b/src/routing_layer/etcp_router.c index a46ab9bf..f8a0d8e0 100644 --- a/src/routing_layer/etcp_router.c +++ b/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); + } + } +} diff --git a/src/routing_layer/etcp_router.h b/src/routing_layer/etcp_router.h index b1e0a022..ee712ca3 100644 --- a/src/routing_layer/etcp_router.h +++ b/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 в группе и уведомить сервис. diff --git a/src/routing_layer/routing.c b/src/routing_layer/routing.c index bf143992..190c9f97 100644 --- a/src/routing_layer/routing.c +++ b/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) { diff --git a/src/routing_layer/topo_group.c b/src/routing_layer/topo_group.c index 28c62f23..21ddba5d 100644 --- a/src/routing_layer/topo_group.c +++ b/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,9 +650,10 @@ 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]; - snprintf(db_file, sizeof(db_file), "%s/chats.db", instance->config->global.db_path); + { + 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); if (rc == SQLITE_OK && instance->topo_sqlite_db) { @@ -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)) { diff --git a/src/routing_layer/topo_group.h b/src/routing_layer/topo_group.h index ee1d9c06..4b7ce53f 100644 --- a/src/routing_layer/topo_group.h +++ b/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); diff --git a/src/routing_layer/topo_group_invite.c b/src/routing_layer/topo_group_invite.c index b0931ba7..a8bdd35b 100644 --- a/src/routing_layer/topo_group_invite.c +++ b/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; diff --git a/src/routing_layer/topo_group_invite.h b/src/routing_layer/topo_group_invite.h index 34f9feb8..5ae837d3 100644 --- a/src/routing_layer/topo_group_invite.h +++ b/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); diff --git a/src/routing_layer/topo_node.c b/src/routing_layer/topo_node.c index f5e3e0e3..930f2108 100644 --- a/src/routing_layer/topo_node.c +++ b/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; } diff --git a/src/routing_layer/topo_node.h b/src/routing_layer/topo_node.h index 2b33ed74..5873a007 100644 --- a/src/routing_layer/topo_node.h +++ b/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); diff --git a/src/transport_layer/etcp_connections.c b/src/transport_layer/etcp_connections.c index 9b44c16a..70db1336 100644 --- a/src/transport_layer/etcp_connections.c +++ b/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; diff --git a/src/utun_instance.c b/src/utun_instance.c index 12cd7f7b..60ff7b72 100644 --- a/src/utun_instance.c +++ b/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)"); - } + // Generic replication is a core service; channel instances belong to chat. + if (!instance->db_sync && db_sync_init(instance) < 0) return -1; - // 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() - // 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 diff --git a/src/utun_instance.h b/src/utun_instance.h index 5b1e4564..ce2393be 100644 --- a/src/utun_instance.h +++ b/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); diff --git a/src/video/video.c b/src/video/video.c index 810bf872..35591c86 100644 --- a/src/video/video.c +++ b/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; } diff --git a/src/video/video.h b/src/video/video.h index 6b13e724..098bc3e6 100644 --- a/src/video/video.h +++ b/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); diff --git a/tests/Makefile.am b/tests/Makefile.am index 438c042f..d1c9456c 100644 --- a/tests/Makefile.am +++ b/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) diff --git a/tests/test_chat_join_e2e.c b/tests/test_chat_join_e2e.c index 97da5d0f..6e437daa 100644 --- a/tests/test_chat_join_e2e.c +++ b/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); diff --git a/tests/test_conn_mgr_phases.c b/tests/test_conn_mgr_phases.c index 03a5bb10..9a24ee4f 100644 --- a/tests/test_conn_mgr_phases.c +++ b/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; diff --git a/tests/test_etcp_100_packets.c b/tests/test_etcp_100_packets.c index 7c87c70b..d7269739 100644 --- a/tests/test_etcp_100_packets.c +++ b/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; } diff --git a/tests/test_etcp_api.c b/tests/test_etcp_api.c index 3fdca458..0f8b4a97 100644 --- a/tests/test_etcp_api.c +++ b/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; } diff --git a/tests/test_etcp_ping.c b/tests/test_etcp_ping.c index 9faa6cfa..26058f9e 100644 --- a/tests/test_etcp_ping.c +++ b/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; diff --git a/tests/test_etcp_simple_traffic.c b/tests/test_etcp_simple_traffic.c index a97cd783..3cf71f64 100644 --- a/tests/test_etcp_simple_traffic.c +++ b/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; diff --git a/tests/test_group_recovery.c b/tests/test_group_recovery.c index 4c1ff85f..0a2b5df4 100644 --- a/tests/test_group_recovery.c +++ b/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. */ diff --git a/tests/test_invite_group_create.c b/tests/test_invite_group_create.c index 2eef1eea..e237c1dc 100644 --- a/tests/test_invite_group_create.c +++ b/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 ── */ { diff --git a/tests/test_media_delivery_chat.c b/tests/test_media_delivery_chat.c index 29153d74..a55f179a 100644 --- a/tests/test_media_delivery_chat.c +++ b/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]); } diff --git a/tests/test_media_delivery_full.c b/tests/test_media_delivery_full.c index e43db8c6..5527f4e5 100644 --- a/tests/test_media_delivery_full.c +++ b/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; } } diff --git a/tests/test_media_delivery_integration.c b/tests/test_media_delivery_integration.c index a6c5004c..968390fc 100644 --- a/tests/test_media_delivery_integration.c +++ b/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); diff --git a/tests/test_nat_detection.c b/tests/test_nat_detection.c index c8db1f66..a895501e 100644 --- a/tests/test_nat_detection.c +++ b/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"); diff --git a/tests/test_pkt_normalizer_etcp.c b/tests/test_pkt_normalizer_etcp.c index d9c12a61..edbac38c 100644 --- a/tests/test_pkt_normalizer_etcp.c +++ b/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"); diff --git a/tests/test_route_ping.c b/tests/test_route_ping.c index 6bbb6961..008aa5db 100644 --- a/tests/test_route_ping.c +++ b/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; } diff --git a/tests/test_services.c b/tests/test_services.c new file mode 100644 index 00000000..e6431f9e --- /dev/null +++ b/tests/test_services.c @@ -0,0 +1,173 @@ +#include +#include +#include +#include +#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; +} diff --git a/tools/chatgui-android/libutun_lite/instance_lite.c b/tools/chatgui-android/libutun_lite/instance_lite.c index bf35920e..7cd9848a 100644 --- a/tools/chatgui-android/libutun_lite/instance_lite.c +++ b/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); diff --git a/tools/chatgui/transport/utun_node.cpp b/tools/chatgui/transport/utun_node.cpp index 73427276..0203a84f 100644 --- a/tools/chatgui/transport/utun_node.cpp +++ b/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");