From cc072cceb2d442ad06b60c552d2c13ccc099e4f3 Mon Sep 17 00:00:00 2001 From: evgeny Date: Thu, 1 Oct 2026 14:14:46 +0300 Subject: [PATCH] Persist DM outbox and verify exact delivery receipts --- src/dm/dm_core.c | 866 ++++++++++++++++++++++------------------------- src/dm/dm_core.h | 21 +- 2 files changed, 425 insertions(+), 462 deletions(-) diff --git a/src/dm/dm_core.c b/src/dm/dm_core.c index f998fdab..e5198e71 100644 --- a/src/dm/dm_core.c +++ b/src/dm/dm_core.c @@ -30,20 +30,20 @@ #include #include #include +#include #define DM_ID "dm_core" -#define DM_DATA_MAX 1024 /* макс размер открытого текста сообщения (v1 — текст) */ /* Подкоманды протокола DM (первый байт payload после ROUTER_SVC-заголовка). */ #define DM_SUBCMD_MSG 0x01 /* [msg] — сообщение */ -#define DM_SUBCMD_ACK 0x02 /* [conv_id:8][seq:8] — подтверждение приёма */ -#define DM_SUBCMD_HELLO 0x03 /* [conv_id:8][last_out_seq:8][last_in_seq:8] — catch-up при connect */ +#define DM_SUBCMD_ACK 0x02 /* Подписанная квитанция получателя. */ /* Per-instance состояние модуля. */ struct dm_state { struct UTUN_INSTANCE* inst; sqlite3* db; int initialized; + void* timer; }; /* ── вспомогательные ── */ @@ -64,8 +64,8 @@ static int dm_exec(struct dm_state* dm, const char* sql) { } /* Создать таблицы бесед и сообщений (общие для всех узлов). */ -static void dm_create_tables(struct dm_state* dm) { - dm_exec(dm, "CREATE TABLE IF NOT EXISTS dm_conversations (" +static int dm_create_tables(struct dm_state* dm) { + if (dm_exec(dm, "CREATE TABLE IF NOT EXISTS dm_conversations (" " conv_id TEXT PRIMARY KEY," " peer_node_id INTEGER NOT NULL," " peer_x25519 BLOB NOT NULL," @@ -75,8 +75,8 @@ static void dm_create_tables(struct dm_state* dm) { " last_out_seq INTEGER DEFAULT 0," " last_in_seq INTEGER DEFAULT 0," " created_at INTEGER DEFAULT 0," - " last_ts INTEGER DEFAULT 0)"); - dm_exec(dm, "CREATE TABLE IF NOT EXISTS dm_messages (" + " last_ts INTEGER DEFAULT 0)") != 0) return -1; + if (dm_exec(dm, "CREATE TABLE IF NOT EXISTS dm_messages (" " conv_id TEXT NOT NULL," " dir INTEGER NOT NULL," /* 0=входящее, 1=исходящее */ " seq INTEGER NOT NULL," @@ -85,7 +85,9 @@ static void dm_create_tables(struct dm_state* dm) { " ct TEXT DEFAULT ''," " data BLOB," /* зашифрованное тело (ciphertext+tag) */ " sig BLOB," /* Ed25519 автора */ - " PRIMARY KEY (conv_id, dir, seq))"); + " PRIMARY KEY (conv_id, dir, seq))") != 0) return -1; + return dm_exec(dm, "CREATE TABLE IF NOT EXISTS dm_outbox (conv_id TEXT NOT NULL, seq INTEGER NOT NULL," + " body BLOB NOT NULL, retry_at INTEGER NOT NULL DEFAULT 0, PRIMARY KEY(conv_id,seq))"); } /* Прочитать оба pubkey узла (x25519 + ed25519) из таблицы nodes. 0=ок, <0=нет. */ @@ -93,19 +95,23 @@ static int dm_node_pubkeys(struct dm_state* dm, uint64_t node_id, uint8_t x25519 sqlite3_stmt* st = NULL; if (sqlite3_prepare_v2(dm->db, "SELECT x25519_pubkey, ed25519_pubkey FROM nodes WHERE node_id=?", - -1, &st, NULL) != SQLITE_OK) return -1; + -1, &st, NULL) != SQLITE_OK) { + DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: node keys query failed: %s", DM_ID, sqlite3_errmsg(dm->db)); + return -1; + } sqlite3_bind_int64(st, 1, (sqlite3_int64)node_id); int rc = -1; if (sqlite3_step(st) == SQLITE_ROW) { const uint8_t* x = (const uint8_t*)sqlite3_column_blob(st, 0); const uint8_t* e = (const uint8_t*)sqlite3_column_blob(st, 1); - if (x && e && sqlite3_column_bytes(st, 0) >= 32 && sqlite3_column_bytes(st, 1) >= 32) { + if (x && e && sqlite3_column_bytes(st, 0) == 32 && sqlite3_column_bytes(st, 1) == 32) { memcpy(x25519, x, 32); memcpy(ed25519, e, 32); rc = 0; } } sqlite3_finalize(st); + if (rc) DEBUG_WARN(DEBUG_CATEGORY_DM, "%s: missing node keys peer=%llu", DM_ID, (unsigned long long)node_id); return rc; } @@ -129,68 +135,59 @@ static int dm_shared_group_id(struct dm_state* dm, uint64_t peer_id, uint64_t* o return shared; } -/* Есть ли общая группа с пиром (anti-spam: принимаем DM только от знакомых из общих каналов). */ -static int dm_share_group(struct dm_state* dm, uint64_t peer_id) { - uint64_t gid = 0; - return dm_shared_group_id(dm, peer_id, &gid); -} - -/* Собрать node_id всех мемберов общих групп (без self и дублей). До max штук. - * Используется для offline-pull: потенциальные отправители почты = мемберы общих групп. */ -static int dm_shared_peers(struct dm_state* dm, uint64_t* out, int max) { - if (!dm || !out || max <= 0) return 0; - uint64_t* ch_ids = NULL; - int ch_count = 0; - if (topo_node_sqlite_get_member_channels(dm->db, dm->inst->node_id, &ch_ids, &ch_count) != 0) +/* Выбрать общую область с действительно живым маршрутом, включая BGP-транзит. */ +int dm_route_group(struct UTUN_INSTANCE* inst, uint64_t peer, uint64_t preferred, uint64_t* group_id) { + if (!inst || !inst->topo_sqlite_db || !group_id) { + DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: invalid route lookup", DM_ID); return 0; + } + uint64_t* ids = NULL; + int count = 0; + if (topo_node_sqlite_get_member_channels(inst->topo_sqlite_db, inst->node_id, &ids, &count) != 0) { + DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: route membership lookup failed", DM_ID); + return 0; + } int found = 0; - for (int c = 0; c < ch_count && found < max; c++) { - char tbl[80]; snprintf(tbl, sizeof(tbl), "peers_%llu", (unsigned long long)ch_ids[c]); - char sql[256]; snprintf(sql, sizeof(sql), "SELECT node_id FROM \"%s\"", tbl); - sqlite3_stmt* st = NULL; - if (sqlite3_prepare_v2(dm->db, sql, -1, &st, NULL) != SQLITE_OK) continue; - while (sqlite3_step(st) == SQLITE_ROW && found < max) { - uint64_t nid = (uint64_t)sqlite3_column_int64(st, 0); - if (nid == dm->inst->node_id) continue; - int dup = 0; - for (int i = 0; i < found; i++) if (out[i] == nid) { dup = 1; break; } - if (!dup) out[found++] = nid; - } - sqlite3_finalize(st); + for (int i = -1; i < count; i++) { + uint64_t gid = i < 0 ? preferred : ids[i]; + if (!gid) continue; + char ch[64]; + snprintf(ch, sizeof(ch), "%llu", (unsigned long long)gid); + if (!topo_node_sqlite_member_in_channel(inst->topo_sqlite_db, ch, inst->node_id) || + !topo_node_sqlite_member_in_channel(inst->topo_sqlite_db, ch, peer)) continue; + struct ETCP_ROUTER_PATH path = etcp_router_get_path(inst, gid, peer, ETCP_RT_ID_DM); + if (!path.next_hop_node_id) continue; + *group_id = gid; + found = 1; + break; } - if (ch_ids) u_free(ch_ids); + u_free(ids); return found; } -/* Доступен ли пир для прямой доставки: живой ETCP-линк или присутствие в BGP-группе. */ -static int dm_peer_reachable(struct dm_state* dm, uint64_t peer_id) { - struct ETCP_CONN* conn = instance_find_conn(dm->inst, peer_id); - if (conn && conn->links_up && conn->initialized) return 1; - if (dm->inst->topo_groups && dm->inst->topo_groups->group_list) { - struct ll_entry* e = dm->inst->topo_groups->group_list->head; - while (e) { - struct TOPO_GROUP* g = (struct TOPO_GROUP*)e; - if (g->group_type == TOPO_GROUP_TYPE_CHAT && topo_node_find_by_id(g, peer_id)) return 1; - e = e->next; - } - } - return 0; -} - -/* Отправить payload подкомандой subcmd в dst через etcp_router (mode=0, своё E2E). */ +/* Router владеет entry после вызова, в том числе при ошибке. */ static int dm_route_send(struct dm_state* dm, uint64_t group_id, uint64_t dst, uint8_t subcmd, const uint8_t* body, size_t body_len) { + if (!etcp_router_send_q_has_room(dm->inst, group_id, dst, ETCP_RT_ID_DM)) { + DEBUG_DEBUG(DEBUG_CATEGORY_DM, "%s: send deferred by backpressure peer=%llu", DM_ID, (unsigned long long)dst); + return -1; + } struct ll_entry* entry = queue_entry_new(0); - if (!entry) return -1; - entry->dgram = u_malloc(1 + 1 + body_len); - if (!entry->dgram) { queue_entry_free(entry); return -1; } + if (!entry) { DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: packet allocation failed", DM_ID); return -1; } + entry->dgram = u_malloc(2 + body_len); + if (!entry->dgram) { + DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: payload allocation failed", DM_ID); + queue_entry_free(entry); + return -1; + } entry->dgram[0] = ETCP_RT_ID_DM; entry->dgram[1] = subcmd; if (body_len) memcpy(entry->dgram + 2, body, body_len); entry->len = (uint16_t)(2 + body_len); - int rc = etcp_route_send(dm->inst, group_id, dst, entry, 1, 0); - if (rc != 0) { queue_dgram_free(entry); queue_entry_free(entry); return -1; } - return 0; + int rc = etcp_route_send(dm->inst, group_id, dst, entry, 1, ROUTER_FLAG_SIGNED); + if (rc) DEBUG_WARN(DEBUG_CATEGORY_DM, "%s: route send failed peer=%llu cmd=%u rc=%d", + DM_ID, (unsigned long long)dst, subcmd, rc); + return rc; } /* ── работа с беседой ── */ @@ -213,7 +210,10 @@ static int dm_conv_load(struct dm_state* dm, const char* conv_id, struct dm_conv if (sqlite3_prepare_v2(dm->db, "SELECT peer_node_id, peer_x25519, peer_ed25519, COALESCE(peer_name,'')," " COALESCE(group_id,0), last_out_seq, last_in_seq FROM dm_conversations WHERE conv_id=?", - -1, &st, NULL) != SQLITE_OK) return -1; + -1, &st, NULL) != SQLITE_OK) { + DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: conversation query failed: %s", DM_ID, sqlite3_errmsg(dm->db)); + return -1; + } sqlite3_bind_text(st, 1, conv_id, -1, SQLITE_STATIC); int rc = -1; if (sqlite3_step(st) == SQLITE_ROW) { @@ -222,8 +222,13 @@ static int dm_conv_load(struct dm_state* dm, const char* conv_id, struct dm_conv c->peer_node_id = (uint64_t)sqlite3_column_int64(st, 0); const uint8_t* x = (const uint8_t*)sqlite3_column_blob(st, 1); const uint8_t* e = (const uint8_t*)sqlite3_column_blob(st, 2); - if (x && sqlite3_column_bytes(st, 1) >= 32) memcpy(c->peer_x25519, x, 32); - if (e && sqlite3_column_bytes(st, 2) >= 32) memcpy(c->peer_ed25519, e, 32); + if (!x || !e || sqlite3_column_bytes(st, 1) != 32 || sqlite3_column_bytes(st, 2) != 32) { + DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: damaged conversation keys conv=%s", DM_ID, conv_id); + sqlite3_finalize(st); + return -1; + } + memcpy(c->peer_x25519, x, 32); + memcpy(c->peer_ed25519, e, 32); const char* name = (const char*)sqlite3_column_text(st, 3); if (name) snprintf(c->peer_name, sizeof(c->peer_name), "%s", name); c->group_id = (uint64_t)sqlite3_column_int64(st, 4); @@ -235,8 +240,7 @@ static int dm_conv_load(struct dm_state* dm, const char* conv_id, struct dm_conv uint64_t gid = 0; if (dm_shared_group_id(dm, c->peer_node_id, &gid)) c->group_id = gid; } - dm_derive_content_key(dm->inst->my_keys.private_key, c->peer_x25519, c->content_key); - rc = 0; + rc = dm_derive_content_key(dm->inst->my_keys.private_key, c->peer_x25519, c->content_key); } sqlite3_finalize(st); return rc; @@ -253,7 +257,10 @@ static int dm_conv_save(struct dm_state* dm, const struct dm_conv* c) { " peer_x25519=excluded.peer_x25519, peer_ed25519=excluded.peer_ed25519," " peer_name=excluded.peer_name, group_id=excluded.group_id," " last_out_seq=excluded.last_out_seq, last_in_seq=excluded.last_in_seq", - -1, &st, NULL) != SQLITE_OK) return -1; + -1, &st, NULL) != SQLITE_OK) { + DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: conversation save prepare failed: %s", DM_ID, sqlite3_errmsg(dm->db)); + return -1; + } sqlite3_bind_text(st, 1, c->conv_id, -1, SQLITE_STATIC); sqlite3_bind_int64(st, 2, (sqlite3_int64)c->peer_node_id); sqlite3_bind_blob(st, 3, c->peer_x25519, 32, SQLITE_STATIC); @@ -264,6 +271,7 @@ static int dm_conv_save(struct dm_state* dm, const struct dm_conv* c) { sqlite3_bind_int64(st, 8, (sqlite3_int64)c->last_in_seq); sqlite3_bind_int64(st, 9, (sqlite3_int64)ntp_time_get_seconds(dm->inst)); int rc = sqlite3_step(st) == SQLITE_DONE ? 0 : -1; + if (rc) DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: conversation save failed conv=%s: %s", DM_ID, c->conv_id, sqlite3_errmsg(dm->db)); sqlite3_finalize(st); return rc; } @@ -275,7 +283,10 @@ static int dm_conv_save(struct dm_state* dm, const struct dm_conv* c) { static int dm_build_body(struct dm_state* dm, const struct dm_conv* c, uint64_t seq, uint64_t ts, const char* ct, const uint8_t* data, uint32_t data_len, uint8_t** out, size_t* out_len) { - if (!c || !ct || (!data && data_len) || data_len > DM_DATA_MAX) return -1; + if (!c || !ct || (!data && data_len) || data_len > DM_DATA_MAX || strlen(ct) > 255) { + DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: invalid message size=%u", DM_ID, data_len); + return -1; + } uint8_t ct_len = (uint8_t)strnlen(ct, 255); uint64_t conv_num = (uint64_t)strtoull(c->conv_id, NULL, 10); @@ -307,7 +318,12 @@ static int dm_build_body(struct dm_state* dm, const struct dm_conv* c, uint64_t memcpy(body + off, &el, 2); off += 2; memcpy(body + off, enc, enc_len); off += enc_len; /* подпись по всем полям до sig */ - sc_ed25519_sign(dm->inst->my_ed25519_privkey, body, off, body + off); + if (sc_ed25519_sign(dm->inst->my_ed25519_privkey, body, off, body + off) != SC_OK) { + DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: message signing failed", DM_ID); + u_free(enc); + u_free(body); + return -1; + } u_free(enc); *out = body; @@ -315,353 +331,309 @@ static int dm_build_body(struct dm_state* dm, const struct dm_conv* c, uint64_t return 0; } -/* ── обработка входящего сообщения ── */ - -/* Обработать каноническое тело входящего сообщения (из live-доставки или mailbox): - * проверить подпись, расшифровать, dedup по seq, вставить, отослать ACK, событие GUI. */ -static void dm_on_msg(struct dm_state* dm, const uint8_t* msg, size_t mlen) { - if (!msg || mlen < DM_MSG_FIXED_HDR + DM_MSG_SIG_SIZE) return; - uint64_t conv_num; memcpy(&conv_num, msg, 8); - uint64_t seq; memcpy(&seq, msg + 8, 8); - uint64_t ts; memcpy(&ts, msg + 16, 8); - uint64_t author; memcpy(&author, msg + 24, 8); - uint8_t ct_len = msg[32]; - const uint8_t* ct = msg + 33; - uint16_t enc_len; memcpy(&enc_len, msg + 33 + ct_len, 2); - const uint8_t* enc = msg + 35 + ct_len; - const uint8_t* sig = msg + mlen - DM_MSG_SIG_SIZE; +/* Проверить тело перед любой записью: подпись связывает conv, автора и адресата. */ +int dm_verify_message(struct UTUN_INSTANCE* inst, uint64_t recipient, const uint8_t* msg, size_t len) { + struct dm_state view = { .inst = inst, .db = inst ? inst->topo_sqlite_db : NULL }; + size_t actual = 0; + uint64_t conv, seq, author; + if (!view.db || !dm_msg_len(msg, len, &actual) || actual != len) goto invalid; + memcpy(&conv, msg, 8); + memcpy(&seq, msg + 8, 8); + memcpy(&author, msg + 24, 8); + if (!recipient || !author || author == recipient || !seq || seq > INT64_MAX || + conv != dm_derive_conv_id(author, recipient)) goto invalid; + uint8_t x[32], ed[32]; + if (dm_node_pubkeys(&view, author, x, ed) != 0 || + sc_ed25519_verify(ed, msg, len - DM_MSG_SIG_SIZE, msg + len - DM_MSG_SIG_SIZE) != SC_OK) goto invalid; + return 0; +invalid: + DEBUG_WARN(DEBUG_CATEGORY_DM, "%s: rejected message recipient=%llu bytes=%zu", DM_ID, + (unsigned long long)recipient, len); + return -1; +} - if (author == dm->inst->node_id) return; /* своё (эхо) — игнор */ +/* Сохранить тело; точный повтор допустим, другое тело с тем же seq запрещено. Транзакция у вызывающего. */ +static int dm_store_message(struct dm_state* dm, const char* conv, int dir, const uint8_t* msg, size_t len) { + uint64_t seq, ts, author; + memcpy(&seq, msg + 8, 8); + memcpy(&ts, msg + 16, 8); + memcpy(&author, msg + 24, 8); + uint16_t enc_len; + uint8_t ctl = msg[32]; + memcpy(&enc_len, msg + 33 + ctl, 2); + sqlite3_stmt* st = NULL; + int rc = sqlite3_prepare_v2(dm->db, + "INSERT INTO dm_messages(conv_id,dir,seq,ts,author,ct,data,sig) VALUES(?,?,?,?,?,?,?,?)" + " ON CONFLICT(conv_id,dir,seq) DO NOTHING", -1, &st, NULL); + if (rc != SQLITE_OK) goto fail; + sqlite3_bind_text(st, 1, conv, -1, SQLITE_STATIC); + sqlite3_bind_int(st, 2, dir); + sqlite3_bind_int64(st, 3, (sqlite3_int64)seq); + sqlite3_bind_int64(st, 4, (sqlite3_int64)ts); + sqlite3_bind_int64(st, 5, (sqlite3_int64)author); + sqlite3_bind_text(st, 6, (const char*)msg + 33, ctl, SQLITE_STATIC); + sqlite3_bind_blob(st, 7, msg + 35 + ctl, enc_len, SQLITE_STATIC); + sqlite3_bind_blob(st, 8, msg + len - DM_MSG_SIG_SIZE, DM_MSG_SIG_SIZE, SQLITE_STATIC); + rc = sqlite3_step(st); + if (rc != SQLITE_DONE) goto fail; + int inserted = sqlite3_changes(dm->db); + sqlite3_finalize(st); + st = NULL; + if (inserted) return 1; + rc = sqlite3_prepare_v2(dm->db, "SELECT sig FROM dm_messages WHERE conv_id=? AND dir=? AND seq=?", -1, &st, NULL); + if (rc != SQLITE_OK) goto fail; + sqlite3_bind_text(st, 1, conv, -1, SQLITE_STATIC); + sqlite3_bind_int(st, 2, dir); + sqlite3_bind_int64(st, 3, (sqlite3_int64)seq); + rc = sqlite3_step(st); + if (rc != SQLITE_ROW || sqlite3_column_bytes(st, 0) != DM_MSG_SIG_SIZE || + memcmp(sqlite3_column_blob(st, 0), msg + len - DM_MSG_SIG_SIZE, DM_MSG_SIG_SIZE)) goto fail; + sqlite3_finalize(st); + return 0; +fail: + DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: message write/conflict conv=%s dir=%d seq=%llu rc=%d sql=%s", + DM_ID, conv, dir, (unsigned long long)seq, rc, sqlite3_errmsg(dm->db)); + sqlite3_finalize(st); + return -1; +} +/* Принять сообщение и создать подписанную квитанцию только после успешного commit. */ +int dm_accept_message(struct UTUN_INSTANCE* inst, const uint8_t* msg, size_t len, uint8_t receipt[DM_RECEIPT_SIZE]) { + struct dm_state* dm = dm_of(inst); + if (!dm || !dm->initialized || !receipt) { + DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: receive without initialized owner", DM_ID); + return -1; + } + if (dm_verify_message(inst, inst->node_id, msg, len) != 0) return -1; + uint64_t conv, seq, author; + memcpy(&conv, msg, 8); + memcpy(&seq, msg + 8, 8); + memcpy(&author, msg + 24, 8); char conv_id[64]; - snprintf(conv_id, sizeof(conv_id), "%llu", (unsigned long long)conv_num); - + snprintf(conv_id, sizeof(conv_id), "%llu", (unsigned long long)conv); struct dm_conv c; - if (dm_conv_load(dm, conv_id, &c) != 0) { - /* первое сообщение: автоприём при общей группе + совпадение conv_id */ - if (dm_derive_conv_id(dm->inst->node_id, author) != conv_num) { - DEBUG_WARN(DEBUG_CATEGORY_DM, "%s: conv_id mismatch from 0x%016llx", DM_ID, (unsigned long long)author); - return; - } - uint64_t gid = 0; - if (!dm_shared_group_id(dm, author, &gid)) { - DEBUG_WARN(DEBUG_CATEGORY_DM, "%s: no shared group with 0x%016llx — reject DM", DM_ID, (unsigned long long)author); - return; - } + int created = dm_conv_load(dm, conv_id, &c) != 0; + if (created) { memset(&c, 0, sizeof(c)); snprintf(c.conv_id, sizeof(c.conv_id), "%s", conv_id); c.peer_node_id = author; - c.group_id = gid; - if (dm_node_pubkeys(dm, author, c.peer_x25519, c.peer_ed25519) != 0) { - DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: no pubkeys for 0x%016llx", DM_ID, (unsigned long long)author); - return; - } - chat_core_get_node_name(dm->inst, author, c.peer_name, sizeof(c.peer_name)); - dm_derive_content_key(dm->inst->my_keys.private_key, c.peer_x25519, c.content_key); - dm_conv_save(dm, &c); - DEBUG_INFO(DEBUG_CATEGORY_DM, "%s: auto-created conversation conv=%s peer=0x%016llx name=%s group=%llu", - DM_ID, conv_id, (unsigned long long)author, c.peer_name, (unsigned long long)c.group_id); - - /* событие GUI: появилась новая беседа (чтобы GUI добавил её в список) */ - { - size_t cl = strlen(conv_id); - uint8_t evt[1 + 64]; - evt[0] = (uint8_t)cl; - memcpy(evt + 1, conv_id, cl); - chat_event_post(dm->inst, CHAT_EVT_DM_CONV_UPDATED, evt, 1 + (int)cl); + if (!dm_shared_group_id(dm, author, &c.group_id)) { + DEBUG_WARN(DEBUG_CATEGORY_DM, "%s: unknown correspondent peer=%llu", DM_ID, (unsigned long long)author); + return -1; } + if (dm_node_pubkeys(dm, author, c.peer_x25519, c.peer_ed25519) != 0 || + dm_derive_content_key(inst->my_keys.private_key, c.peer_x25519, c.content_key) != 0) return -1; + chat_core_get_node_name(inst, author, c.peer_name, sizeof(c.peer_name)); } - - /* проверка подписи автора по всем полям до sig */ - if (sc_ed25519_verify(c.peer_ed25519, msg, mlen - DM_MSG_SIG_SIZE, sig) != SC_OK) { - DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: invalid signature conv=%s seq=%llu", DM_ID, conv_id, (unsigned long long)seq); - return; - } - - /* dedup: вставляем только новое */ - if (seq <= c.last_in_seq) { - /* уже есть — только подтвердить */ - uint8_t ack[16]; memcpy(ack, &conv_num, 8); memcpy(ack + 8, &seq, 8); - dm_route_send(dm, c.group_id, author, DM_SUBCMD_ACK, ack, sizeof(ack)); - return; + if (c.peer_node_id != author) { + DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: conversation author mismatch conv=%s", DM_ID, conv_id); + return -1; } - - /* расшифровать (nonce детерминирован по conv_id+author+seq) */ - uint8_t nonce[DM_NONCE_SIZE]; - dm_build_nonce(conv_num, author, seq, nonce); - uint8_t* plain = u_malloc(enc_len > DM_TAG_SIZE ? enc_len - DM_TAG_SIZE : 0); - if (!plain) return; + uint16_t enc_len; + memcpy(&enc_len, msg + 33 + msg[32], 2); + uint8_t nonce[DM_NONCE_SIZE], plain[DM_DATA_MAX + 1]; size_t plain_len = 0; - if (dm_decrypt(c.content_key, nonce, enc, enc_len, plain, &plain_len) != 0) { - u_free(plain); - DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: decrypt failed conv=%s seq=%llu", DM_ID, conv_id, (unsigned long long)seq); - return; + dm_build_nonce(conv, author, seq, nonce); + if (dm_decrypt(c.content_key, nonce, msg + 35 + msg[32], enc_len, plain, &plain_len) != 0) return -1; + if (dm_exec(dm, "BEGIN IMMEDIATE") != 0) return -1; + int inserted = dm_store_message(dm, conv_id, 0, msg, len); + if (seq > c.last_in_seq) c.last_in_seq = seq; /* Только статистика; не курсор dedup. */ + if (inserted < 0 || dm_conv_save(dm, &c) != 0 || dm_exec(dm, "COMMIT") != 0) { + dm_exec(dm, "ROLLBACK"); + return -1; } - - /* вставить в dm_messages (dir=0, входящее) */ - sqlite3_stmt* st = NULL; - if (sqlite3_prepare_v2(dm->db, - "INSERT OR IGNORE INTO dm_messages(conv_id,dir,seq,ts,author,ct,data,sig)" - " VALUES(?,0,?,?,?,?,?,?)", - -1, &st, NULL) == SQLITE_OK) { - sqlite3_bind_text(st, 1, conv_id, -1, SQLITE_STATIC); - sqlite3_bind_int64(st, 2, (sqlite3_int64)seq); - sqlite3_bind_int64(st, 3, (sqlite3_int64)ts); - sqlite3_bind_int64(st, 4, (sqlite3_int64)author); - sqlite3_bind_text(st, 5, (const char*)ct, ct_len, SQLITE_STATIC); - sqlite3_bind_blob(st, 6, enc, enc_len, SQLITE_STATIC); - sqlite3_bind_blob(st, 7, sig, DM_MSG_SIG_SIZE, SQLITE_STATIC); - sqlite3_step(st); + memcpy(receipt, "DMACK001", 8); + memcpy(receipt + 8, msg, 16); + memcpy(receipt + 24, &author, 8); + memcpy(receipt + 32, &inst->node_id, 8); + SHA256(msg, len, receipt + 40); + if (sc_ed25519_sign(inst->my_ed25519_privkey, receipt, 72, receipt + 72) != SC_OK) { + DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: receipt signing failed conv=%s seq=%llu", DM_ID, conv_id, (unsigned long long)seq); + return -1; } - if (st) sqlite3_finalize(st); - - c.last_in_seq = seq; - dm_conv_save(dm, &c); - - /* подтвердить приём (нужно и для mailbox-доставки → storage удалит по ACK) */ - uint8_t ack[16]; memcpy(ack, &conv_num, 8); memcpy(ack + 8, &seq, 8); - dm_route_send(dm, c.group_id, author, DM_SUBCMD_ACK, ack, sizeof(ack)); - - DEBUG_INFO(DEBUG_CATEGORY_DM, "%s: recv conv=%s seq=%llu author=0x%016llx ct=%.*s len=%zu", - DM_ID, conv_id, (unsigned long long)seq, (unsigned long long)author, ct_len, ct, plain_len); - - /* событие GUI: [conv_id_len:1][conv_id][author_node_id:8] */ - { + DEBUG_DEBUG(DEBUG_CATEGORY_DM, "%s: committed receive conv=%s seq=%llu new=%d plain=%zu", DM_ID, + conv_id, (unsigned long long)seq, inserted, plain_len); + if (inserted) { size_t cl = strlen(conv_id); uint8_t evt[1 + 64 + 8]; evt[0] = (uint8_t)cl; memcpy(evt + 1, conv_id, cl); memcpy(evt + 1 + cl, &author, 8); - chat_event_post(dm->inst, CHAT_EVT_DM_MSG_RECEIVED, evt, 1 + (int)cl + 8); + if (created) chat_event_post(inst, CHAT_EVT_DM_CONV_UPDATED, evt, 1 + (int)cl); + chat_event_post(inst, CHAT_EVT_DM_MSG_RECEIVED, evt, 1 + (int)cl + 8); + DEBUG_INFO(DEBUG_CATEGORY_DM, "%s: received conv=%s seq=%llu author=%llu", DM_ID, conv_id, + (unsigned long long)seq, (unsigned long long)author); } - u_free(plain); + return 0; } -/* ── обработка HELLO (catch-up) ── */ - -/* Пересобрать каноническое тело из сохранённых полей (для catch-up resend). */ -static int dm_rebuild_body(uint64_t conv_num, uint64_t seq, uint64_t ts, uint64_t author, - const char* ct, const uint8_t* enc, int enc_len, const uint8_t sig[64], - uint8_t** out, size_t* out_len) { - uint8_t ctl = (uint8_t)strnlen(ct, 255); - size_t body_len = DM_MSG_FIXED_HDR + ctl + 2 + (size_t)enc_len + DM_MSG_SIG_SIZE; - uint8_t* body = u_malloc(body_len); - if (!body) return -1; - size_t off = 0; - memcpy(body + off, &conv_num, 8); off += 8; - memcpy(body + off, &seq, 8); off += 8; - memcpy(body + off, &ts, 8); off += 8; - memcpy(body + off, &author, 8); off += 8; - body[off++] = ctl; - memcpy(body + off, ct, ctl); off += ctl; - uint16_t el = (uint16_t)enc_len; - memcpy(body + off, &el, 2); off += 2; - memcpy(body + off, enc, enc_len); off += enc_len; - memcpy(body + off, sig, DM_MSG_SIG_SIZE); - *out = body; - *out_len = body_len; +/* Квитанцию может пересылать суперузел, но подписывает её только получатель. */ +int dm_verify_receipt(struct UTUN_INSTANCE* inst, const uint8_t receipt[DM_RECEIPT_SIZE], const uint8_t* msg, size_t len) { + struct dm_state view = { .inst = inst, .db = inst ? inst->topo_sqlite_db : NULL }; + uint64_t conv, seq, author, recipient; + if (!view.db || !receipt || memcmp(receipt, "DMACK001", 8)) goto invalid; + memcpy(&conv, receipt + 8, 8); + memcpy(&seq, receipt + 16, 8); + memcpy(&author, receipt + 24, 8); + memcpy(&recipient, receipt + 32, 8); + if (!author || !recipient || author == recipient || !seq || seq > INT64_MAX || + conv != dm_derive_conv_id(author, recipient)) goto invalid; + uint8_t x[32], ed[32]; + if (dm_node_pubkeys(&view, recipient, x, ed) != 0 || + sc_ed25519_verify(ed, receipt, 72, receipt + 72) != SC_OK) goto invalid; + if (msg) { + uint8_t digest[32]; + size_t actual = 0; + if (!dm_msg_len(msg, len, &actual) || actual != len || memcmp(receipt + 8, msg, 16) || + memcmp(receipt + 24, msg + 24, 8)) goto invalid; + SHA256(msg, len, digest); + if (memcmp(digest, receipt + 40, 32)) goto invalid; + } return 0; +invalid: + DEBUG_WARN(DEBUG_CATEGORY_DM, "%s: invalid delivery receipt bytes=%zu", DM_ID, len); + return -1; } -/* На connect обе стороны шлют HELLO со своими счётчиками; получив HELLO, досылаем - * пиру недостающие ему наши сообщения (его last_in_seq+1 .. наш last_out_seq). */ -static void dm_on_hello(struct dm_state* dm, uint64_t conv_num, uint64_t from_peer, uint64_t peer_out_seq, uint64_t peer_in_seq) { +/* Удалить исходящую задачу по точной подписанной квитанции; историю беседы оставить. */ +int dm_accept_receipt(struct UTUN_INSTANCE* inst, const uint8_t receipt[DM_RECEIPT_SIZE]) { + struct dm_state* dm = dm_of(inst); + if (!dm || !dm->initialized || dm_verify_receipt(inst, receipt, NULL, 0) != 0) return -1; + uint64_t author, conv, seq; + memcpy(&author, receipt + 24, 8); + memcpy(&conv, receipt + 8, 8); + memcpy(&seq, receipt + 16, 8); + if (author != inst->node_id) { + DEBUG_WARN(DEBUG_CATEGORY_DM, "%s: receipt for another sender=%llu", DM_ID, (unsigned long long)author); + return -1; + } char conv_id[64]; - snprintf(conv_id, sizeof(conv_id), "%llu", (unsigned long long)conv_num); - - struct dm_conv c; - if (dm_conv_load(dm, conv_id, &c) != 0) { - /* Беседы нет локально (первое сообщение ушло, пока мы были offline): - * автосоздать её (anti-spam: conv_id + общая группа) и ответить HELLO, - * чтобы отправитель дослал пропущенные сообщения. */ - if (dm_derive_conv_id(dm->inst->node_id, from_peer) != conv_num) { - DEBUG_WARN(DEBUG_CATEGORY_DM, "%s: hello conv_id mismatch from 0x%016llx", DM_ID, (unsigned long long)from_peer); - return; - } - uint64_t gid = 0; - if (!dm_shared_group_id(dm, from_peer, &gid)) { - DEBUG_WARN(DEBUG_CATEGORY_DM, "%s: hello no shared group with 0x%016llx — ignore", DM_ID, (unsigned long long)from_peer); - return; - } - memset(&c, 0, sizeof(c)); - snprintf(c.conv_id, sizeof(c.conv_id), "%s", conv_id); - c.peer_node_id = from_peer; - c.group_id = gid; - if (dm_node_pubkeys(dm, from_peer, c.peer_x25519, c.peer_ed25519) != 0) { - DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: hello no pubkeys for 0x%016llx", DM_ID, (unsigned long long)from_peer); - return; - } - chat_core_get_node_name(dm->inst, from_peer, c.peer_name, sizeof(c.peer_name)); - dm_derive_content_key(dm->inst->my_keys.private_key, c.peer_x25519, c.content_key); - dm_conv_save(dm, &c); - DEBUG_INFO(DEBUG_CATEGORY_DM, "%s: hello auto-created conversation conv=%s peer=0x%016llx name=%s group=%llu", - DM_ID, conv_id, (unsigned long long)from_peer, c.peer_name, (unsigned long long)c.group_id); - - { - size_t cl = strlen(conv_id); - uint8_t evt[1 + 64]; - evt[0] = (uint8_t)cl; - memcpy(evt + 1, conv_id, cl); - chat_event_post(dm->inst, CHAT_EVT_DM_CONV_UPDATED, evt, 1 + (int)cl); - } - - /* ответить HELLO (out=0,in=0) — отправитель увидит peer_in_seq < last_out_seq и дослал сообщения */ - uint8_t hello[24]; - uint64_t z = 0; - memcpy(hello, &conv_num, 8); - memcpy(hello + 8, &z, 8); - memcpy(hello + 16, &z, 8); - if (dm_route_send(dm, 0, from_peer, DM_SUBCMD_HELLO, hello, sizeof(hello)) != 0) - DEBUG_WARN(DEBUG_CATEGORY_DM, "%s: hello reply send failed peer=0x%016llx", DM_ID, (unsigned long long)from_peer); - else - DEBUG_INFO(DEBUG_CATEGORY_DM, "%s: hello reply → 0x%016llx conv=%s (auto-created)", - DM_ID, (unsigned long long)from_peer, conv_id); - return; + snprintf(conv_id, sizeof(conv_id), "%llu", (unsigned long long)conv); + sqlite3_stmt* st = NULL; + int rc = sqlite3_prepare_v2(dm->db, "SELECT body FROM dm_outbox WHERE conv_id=? AND seq=?", -1, &st, NULL); + if (rc != SQLITE_OK) goto fail; + sqlite3_bind_text(st, 1, conv_id, -1, SQLITE_STATIC); + sqlite3_bind_int64(st, 2, (sqlite3_int64)seq); + rc = sqlite3_step(st); + if (rc == SQLITE_DONE) { sqlite3_finalize(st); return 0; } + if (rc != SQLITE_ROW) goto fail; + if (dm_verify_receipt(inst, receipt, sqlite3_column_blob(st, 0), (size_t)sqlite3_column_bytes(st, 0)) != 0) { + sqlite3_finalize(st); + return -1; } + sqlite3_finalize(st); + st = NULL; + rc = sqlite3_prepare_v2(dm->db, "DELETE FROM dm_outbox WHERE conv_id=? AND seq=?", -1, &st, NULL); + if (rc != SQLITE_OK) goto fail; + sqlite3_bind_text(st, 1, conv_id, -1, SQLITE_STATIC); + sqlite3_bind_int64(st, 2, (sqlite3_int64)seq); + rc = sqlite3_step(st); + if (rc != SQLITE_DONE) goto fail; + sqlite3_finalize(st); + DEBUG_INFO(DEBUG_CATEGORY_DM, "%s: delivery confirmed conv=%s seq=%llu", DM_ID, conv_id, (unsigned long long)seq); + return 0; +fail: + DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: receipt database error rc=%d sql=%s", DM_ID, rc, sqlite3_errmsg(dm->db)); + sqlite3_finalize(st); + return -1; +} - if (peer_in_seq < c.last_out_seq) { +/* Ограниченная порция durable outbox; retry_at обеспечивает продвижение всей очереди. */ +static void dm_pump(struct dm_state* dm) { + int64_t now = (int64_t)ntp_time_get_seconds(dm->inst); + for (int i = 0; i < 32; i++) { sqlite3_stmt* st = NULL; - if (sqlite3_prepare_v2(dm->db, - "SELECT seq, ts, author, ct, data, sig FROM dm_messages" - " WHERE conv_id=? AND dir=1 AND seq>? ORDER BY seq", - -1, &st, NULL) == SQLITE_OK) { - sqlite3_bind_text(st, 1, conv_id, -1, SQLITE_STATIC); - sqlite3_bind_int64(st, 2, (sqlite3_int64)peer_in_seq); - while (sqlite3_step(st) == SQLITE_ROW) { - uint64_t seq = (uint64_t)sqlite3_column_int64(st, 0); - uint64_t ts = (uint64_t)sqlite3_column_int64(st, 1); - uint64_t author = (uint64_t)sqlite3_column_int64(st, 2); - const char* ct = (const char*)sqlite3_column_text(st, 3); - const uint8_t* enc = (const uint8_t*)sqlite3_column_blob(st, 4); - int enc_len = sqlite3_column_bytes(st, 4); - const uint8_t* sig = (const uint8_t*)sqlite3_column_blob(st, 5); - uint8_t* body = NULL; size_t body_len = 0; - if (dm_rebuild_body(conv_num, seq, ts, author, ct, enc, enc_len, sig, &body, &body_len) == 0) { - dm_route_send(dm, c.group_id, from_peer, DM_SUBCMD_MSG, body, body_len); - u_free(body); - } - } + int rc = sqlite3_prepare_v2(dm->db, + "SELECT o.conv_id,o.seq,o.body,c.peer_node_id,c.group_id FROM dm_outbox o" + " JOIN dm_conversations c USING(conv_id) WHERE o.retry_at<=? ORDER BY o.retry_at,o.seq LIMIT 1", -1, &st, NULL); + if (rc != SQLITE_OK) goto fail; + sqlite3_bind_int64(st, 1, now); + rc = sqlite3_step(st); + if (rc == SQLITE_DONE) { sqlite3_finalize(st); break; } + if (rc != SQLITE_ROW) goto fail; + char conv[64]; + snprintf(conv, sizeof(conv), "%s", (const char*)sqlite3_column_text(st, 0)); + uint64_t seq = (uint64_t)sqlite3_column_int64(st, 1); + size_t len = (size_t)sqlite3_column_bytes(st, 2); + uint8_t* body = u_malloc(len); + if (!body) { + DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: outbox allocation failed bytes=%zu", DM_ID, len); + sqlite3_finalize(st); + break; + } + memcpy(body, sqlite3_column_blob(st, 2), len); + uint64_t peer = (uint64_t)sqlite3_column_int64(st, 3), gid = (uint64_t)sqlite3_column_int64(st, 4); + sqlite3_finalize(st); + st = NULL; + rc = sqlite3_prepare_v2(dm->db, "UPDATE dm_outbox SET retry_at=? WHERE conv_id=? AND seq=?", -1, &st, NULL); + if (rc == SQLITE_OK) { + sqlite3_bind_int64(st, 1, now + 5); + sqlite3_bind_text(st, 2, conv, -1, SQLITE_STATIC); + sqlite3_bind_int64(st, 3, (sqlite3_int64)seq); + rc = sqlite3_step(st); + } + sqlite3_finalize(st); + st = NULL; + if (rc != SQLITE_DONE) { u_free(body); goto fail; } + if (dm_route_group(dm->inst, peer, gid, &gid)) { + if (!dm_route_send(dm, gid, peer, DM_SUBCMD_MSG, body, len)) + DEBUG_DEBUG(DEBUG_CATEGORY_DM, "%s: outbox direct conv=%s seq=%llu group=%llu", DM_ID, + conv, (unsigned long long)seq, (unsigned long long)gid); + } else { + dm_mailbox_put(dm->inst, peer, dm->inst->node_id, body, len); } - if (st) sqlite3_finalize(st); + u_free(body); + continue; +fail: + DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: outbox scan failed rc=%d sql=%s", DM_ID, rc, sqlite3_errmsg(dm->db)); + sqlite3_finalize(st); + break; } - DEBUG_INFO(DEBUG_CATEGORY_DM, "%s: hello from 0x%016llx conv=%s out=%llu in=%llu", - DM_ID, (unsigned long long)from_peer, conv_id, - (unsigned long long)peer_out_seq, (unsigned long long)peer_in_seq); } -/* ── обработка ACK ── */ - -/* Получено подтверждение доставки от пира — логируем (надёжность гарантирует etcp_router). */ -static void dm_on_ack(uint64_t conv_num, uint64_t seq) { - DEBUG_INFO(DEBUG_CATEGORY_DM, "%s: ack conv=%llu seq=%llu", DM_ID, - (unsigned long long)conv_num, (unsigned long long)seq); +/* Durable-задачи переживают stop; таймер является только локальным триггером. */ +static void dm_tick(void* arg) { + struct dm_state* dm = arg; + dm->timer = NULL; + dm_pump(dm); + dm->timer = uasync_set_timeout(dm->inst->ua, 10000, dm, dm_tick, "dm_outbox"); + if (!dm->timer) DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: retry timer allocation failed", DM_ID); } -/* ── etcp_router recv (0x34) ── */ - +/* Live-доставка использует тот же commit/receipt, что и доставка с суперузла. */ static void dm_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) { - if (!entry || !conn || !entry->dgram || entry->len < ROUTER_SVC_PAYLOAD_OFF + 1) { + if (!entry || !conn || !entry->dgram || entry->len <= ROUTER_SVC_PAYLOAD_OFF) { if (entry) { queue_dgram_free(entry); queue_entry_free(entry); } - return; + return; /* Уведомление router о закрытии канала. */ } struct dm_state* dm = dm_of(conn->instance); - if (!dm || !dm->initialized) { queue_dgram_free(entry); queue_entry_free(entry); return; } - - const uint8_t* d = entry->dgram; - uint64_t from_node; memcpy(&from_node, d + ROUTER_SVC_SRC_OFF, 8); - const uint8_t* p = d + ROUTER_SVC_PAYLOAD_OFF; - size_t plen = entry->len - ROUTER_SVC_PAYLOAD_OFF; - uint8_t subcmd = p[0]; - - if (subcmd == DM_SUBCMD_MSG) { - size_t mlen = 0; - if (dm_msg_len(p + 1, plen - 1, &mlen)) - dm_on_msg(dm, p + 1, mlen); - } else if (subcmd == DM_SUBCMD_ACK) { - if (plen >= 1 + 16) { - uint64_t conv_num; memcpy(&conv_num, p + 1, 8); - uint64_t seq; memcpy(&seq, p + 9, 8); - dm_on_ack(conv_num, seq); - } - } else if (subcmd == DM_SUBCMD_HELLO) { - if (plen >= 1 + 24) { - uint64_t conv_num; memcpy(&conv_num, p + 1, 8); - uint64_t out_seq; memcpy(&out_seq, p + 9, 8); - uint64_t in_seq; memcpy(&in_seq, p + 17, 8); - dm_on_hello(dm, conv_num, from_node, out_seq, in_seq); - } - } else { - DEBUG_WARN(DEBUG_CATEGORY_DM, "%s: unknown subcmd=%02x", DM_ID, subcmd); + if (!dm || !dm->initialized) { + DEBUG_WARN(DEBUG_CATEGORY_DM, "%s: receive while stopped", DM_ID); + goto done; } - queue_dgram_free(entry); queue_entry_free(entry); -} - -/* ── connection status ── */ - -/* При поднятии соединения с пиром: HELLO (catch-up) + вытянуть mailbox. */ -static void dm_on_conn_status(struct ETCP_CONN* conn, int status, void* arg) { - struct UTUN_INSTANCE* inst = (struct UTUN_INSTANCE*)arg; - struct dm_state* dm = dm_of(inst); - if (!conn || status != ETCP_CONN_STATUS_UP || !dm) return; - uint64_t peer = conn->peer_node_id; - if (!peer) return; - - /* На подъём ЛЮБОГО соединения: HELLO для беседы с подключившимся пиром + - * PULL mailbox для бесед, чей пир не подключён напрямую (offline-доставка). */ - uint64_t conv_peers[64]; - int nconv = 0; - sqlite3_stmt* st = NULL; - if (sqlite3_prepare_v2(dm->db, - "SELECT conv_id, peer_node_id, last_out_seq, last_in_seq FROM dm_conversations", - -1, &st, NULL) == SQLITE_OK) { - while (sqlite3_step(st) == SQLITE_ROW) { - const char* conv_id = (const char*)sqlite3_column_text(st, 0); - uint64_t conv_peer = (uint64_t)sqlite3_column_int64(st, 1); - uint64_t out_seq = (uint64_t)sqlite3_column_int64(st, 2); - uint64_t in_seq = (uint64_t)sqlite3_column_int64(st, 3); - uint64_t conv_num = strtoull(conv_id, NULL, 10); - - if (nconv < (int)(sizeof(conv_peers) / sizeof(conv_peers[0]))) - conv_peers[nconv++] = conv_peer; - - if (conv_peer == peer) { - uint8_t hello[24]; - memcpy(hello, &conv_num, 8); - memcpy(hello + 8, &out_seq, 8); - memcpy(hello + 16, &in_seq, 8); - /* group_id берём из строки беседы — для простоты HELLO по 0, маршрутизация глобальная */ - if (dm_route_send(dm, 0, peer, DM_SUBCMD_HELLO, hello, sizeof(hello)) != 0) - DEBUG_WARN(DEBUG_CATEGORY_DM, "%s: hello send failed conv=%s peer=0x%016llx", - DM_ID, conv_id, (unsigned long long)peer); - else - DEBUG_INFO(DEBUG_CATEGORY_DM, "%s: hello → 0x%016llx conv=%s out=%llu in=%llu", - DM_ID, (unsigned long long)peer, conv_id, - (unsigned long long)out_seq, (unsigned long long)in_seq); - } - - /* пир не подключён напрямую — тянем mailbox со storage-узлов */ - struct ETCP_CONN* pc = instance_find_conn(inst, conv_peer); - if (!pc || !pc->links_up || !pc->initialized) { - dm_mailbox_pull(inst, conv_peer, in_seq); - } - } - sqlite3_finalize(st); + const uint8_t* p = entry->dgram + ROUTER_SVC_PAYLOAD_OFF; + size_t len = entry->len - ROUTER_SVC_PAYLOAD_OFF - 1; + uint64_t source, gid; + memcpy(&source, entry->dgram + ROUTER_SVC_SRC_OFF, 8); + memcpy(&gid, entry->dgram + ROUTER_SVC_GROUP_OFF, 8); + if (!(entry->dgram[ROUTER_SVC_FLAGS_OFF] & ROUTER_FLAG_SIGNED)) { + DEBUG_WARN(DEBUG_CATEGORY_DM, "%s: unsigned route packet source=%llu", DM_ID, (unsigned long long)source); + goto done; } - - /* Первое сообщение могло прийти, пока мы были offline: мемберы общих групп, - * с которыми ещё нет беседы, тоже опрашиваем mailbox (отправитель мог уйти - * offline, тогда HELLO-восстановление недоступно). */ - uint64_t peers[128]; - int np = dm_shared_peers(dm, peers, (int)(sizeof(peers) / sizeof(peers[0]))); - for (int i = 0; i < np; i++) { - uint64_t p = peers[i]; - struct ETCP_CONN* pc = instance_find_conn(inst, p); - if (pc && pc->links_up && pc->initialized) continue; /* online — HELLO-путь */ - int has_conv = 0; - for (int j = 0; j < nconv; j++) if (conv_peers[j] == p) { has_conv = 1; break; } - if (has_conv) continue; /* уже обработан циклом выше */ - dm_mailbox_pull(inst, p, 0); + if (p[0] == DM_SUBCMD_MSG) { + uint8_t receipt[DM_RECEIPT_SIZE]; + if (dm_accept_message(dm->inst, p + 1, len, receipt) == 0) { + dm_route_send(dm, gid, source, DM_SUBCMD_ACK, receipt, sizeof(receipt)); + dm_mailbox_receipt(dm->inst, receipt); + } + } else if (p[0] == DM_SUBCMD_ACK && len == DM_RECEIPT_SIZE) { + if (dm_accept_receipt(dm->inst, p + 1) == 0) dm_mailbox_receipt(dm->inst, p + 1); + } else { + DEBUG_WARN(DEBUG_CATEGORY_DM, "%s: invalid command=%u bytes=%zu", DM_ID, p[0], len); } +done: + queue_dgram_free(entry); + queue_entry_free(entry); } /* ── публичный API ── */ @@ -725,123 +697,97 @@ int dm_start(struct UTUN_INSTANCE* inst, uint64_t peer_node_id, return 0; } -/* Отправить сообщение в беседу: зашифровать, подписать, сохранить, доставить. */ +/* Атомарно сохранить seq, историю и неизменяемое тело исходящей задачи. */ int dm_send(struct UTUN_INSTANCE* inst, const char* conv_id, const char* content_type, const uint8_t* data, uint32_t data_len) { - if (!inst || !conv_id || !content_type || (!data && data_len)) return -1; struct dm_state* dm = dm_of(inst); - if (!dm || !dm->initialized) return -1; - + if (!dm || !dm->initialized || !conv_id || !content_type || (!data && data_len)) { + DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: send invalid arguments/state", DM_ID); + return -1; + } struct dm_conv c; - if (dm_conv_load(dm, conv_id, &c) != 0) { - DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: no conversation %s", DM_ID, conv_id); + if (dm_conv_load(dm, conv_id, &c) != 0 || c.last_out_seq >= INT64_MAX) { + DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: send conversation missing/seq exhausted conv=%s", DM_ID, conv_id); return -1; } - uint64_t seq = c.last_out_seq + 1; uint64_t ts = (uint64_t)(ntp_time_get_us(inst) / 1000); - uint8_t* body = NULL; - size_t body_len = 0; - if (dm_build_body(dm, &c, seq, ts, content_type, data, data_len, &body, &body_len) != 0) { - DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: build body failed", DM_ID); - return -1; - } - - /* сохранить локально (dir=1, исходящее) */ + size_t len = 0; + if (dm_build_body(dm, &c, seq, ts, content_type, data, data_len, &body, &len) != 0) return -1; sqlite3_stmt* st = NULL; - if (sqlite3_prepare_v2(dm->db, - "INSERT INTO dm_messages(conv_id,dir,seq,ts,author,ct,data,sig)" - " VALUES(?,1,?,?,?,?,?,?)", - -1, &st, NULL) == SQLITE_OK) { + if (dm_exec(dm, "BEGIN IMMEDIATE") != 0) { u_free(body); return -1; } + c.last_out_seq = seq; + if (dm_store_message(dm, conv_id, 1, body, len) != 1 || dm_conv_save(dm, &c) != 0) goto rollback; + int rc = sqlite3_prepare_v2(dm->db, "INSERT INTO dm_outbox(conv_id,seq,body) VALUES(?,?,?)", -1, &st, NULL); + if (rc == SQLITE_OK) { sqlite3_bind_text(st, 1, conv_id, -1, SQLITE_STATIC); sqlite3_bind_int64(st, 2, (sqlite3_int64)seq); - sqlite3_bind_int64(st, 3, (sqlite3_int64)ts); - sqlite3_bind_int64(st, 4, (sqlite3_int64)inst->node_id); - sqlite3_bind_text(st, 5, content_type, -1, SQLITE_STATIC); - /* data/sig уже внутри body: data_enc по смещению 35+ct_len, sig в конце */ - uint8_t ctl = (uint8_t)strnlen(content_type, 255); - const uint8_t* enc = body + 35 + ctl; - uint16_t enc_len; memcpy(&enc_len, body + 33 + ctl, 2); - sqlite3_bind_blob(st, 6, enc, enc_len, SQLITE_STATIC); - sqlite3_bind_blob(st, 7, body + body_len - DM_MSG_SIG_SIZE, DM_MSG_SIG_SIZE, SQLITE_STATIC); - sqlite3_step(st); + sqlite3_bind_blob(st, 3, body, (int)len, SQLITE_STATIC); + rc = sqlite3_step(st); } - if (st) sqlite3_finalize(st); - - c.last_out_seq = seq; - dm_conv_save(dm, &c); - - /* доставка: доступен → прямой send, иначе → mailbox */ - if (dm_peer_reachable(dm, c.peer_node_id)) { - if (dm_route_send(dm, c.group_id, c.peer_node_id, DM_SUBCMD_MSG, body, body_len) != 0) - DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: direct send failed conv=%s seq=%llu → 0x%016llx", - DM_ID, conv_id, (unsigned long long)seq, (unsigned long long)c.peer_node_id); - else - DEBUG_INFO(DEBUG_CATEGORY_DM, "%s: sent direct conv=%s seq=%llu → 0x%016llx", - DM_ID, conv_id, (unsigned long long)seq, (unsigned long long)c.peer_node_id); - } else { - if (dm_mailbox_put(inst, c.peer_node_id, inst->node_id, body, body_len) != 0) - DEBUG_WARN(DEBUG_CATEGORY_DM, "%s: mailbox put failed conv=%s seq=%llu (peer offline, no storage)", - DM_ID, conv_id, (unsigned long long)seq); - else - DEBUG_INFO(DEBUG_CATEGORY_DM, "%s: peer offline conv=%s seq=%llu → mailbox", - DM_ID, conv_id, (unsigned long long)seq); + sqlite3_finalize(st); + st = NULL; + if (rc != SQLITE_DONE) { + DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: outbox insert failed conv=%s seq=%llu rc=%d sql=%s", DM_ID, + conv_id, (unsigned long long)seq, rc, sqlite3_errmsg(dm->db)); + goto rollback; } + if (dm_exec(dm, "COMMIT") != 0) goto rollback; u_free(body); - - /* событие GUI (своё сообщение тоже, чтобы GUI обновил список) */ - { - size_t cl = strlen(conv_id); - uint8_t evt[1 + 64 + 8]; - evt[0] = (uint8_t)cl; - memcpy(evt + 1, conv_id, cl); - uint64_t me = inst->node_id; - memcpy(evt + 1 + cl, &me, 8); - chat_event_post(inst, CHAT_EVT_DM_MSG_RECEIVED, evt, 1 + (int)cl + 8); - } + DEBUG_INFO(DEBUG_CATEGORY_DM, "%s: queued conv=%s seq=%llu peer=%llu bytes=%u", DM_ID, conv_id, + (unsigned long long)seq, (unsigned long long)c.peer_node_id, data_len); + dm_pump(dm); + size_t cl = strlen(conv_id); + uint8_t evt[1 + 64 + 8]; + evt[0] = (uint8_t)cl; + memcpy(evt + 1, conv_id, cl); + memcpy(evt + 1 + cl, &inst->node_id, 8); + chat_event_post(inst, CHAT_EVT_DM_MSG_RECEIVED, evt, 1 + (int)cl + 8); return 0; +rollback: + dm_exec(dm, "ROLLBACK"); + u_free(body); + return -1; } -/* Вставка сообщения, полученного через mailbox (вызывается dm_mailbox по PULL_RESP). - * Обрабатывается тем же путём, что и live-сообщение (verify/decrypt/dedup/insert). */ -void dm_core_inject_message(const uint8_t* msg, size_t len, void* arg) { - struct UTUN_INSTANCE* inst = (struct UTUN_INSTANCE*)arg; - struct dm_state* dm = dm_of(inst); - if (!dm) return; - dm_on_msg(dm, msg, len); -} - -/* ── жизненный цикл ── */ - +/* Таблицы принадлежат DM, соединение с SQLite — общему ядру. */ int dm_core_init(struct UTUN_INSTANCE* inst) { - if (!inst) return -1; + if (!inst) { DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: init without instance", DM_ID); return -1; } if (inst->dm && inst->dm->initialized) return 0; struct dm_state* dm = u_calloc(1, sizeof(*dm)); - if (!dm) return -1; + if (!dm) { DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: state allocation failed", DM_ID); return -1; } inst->dm = dm; dm->inst = inst; - dm->db = chat_core_get_db(inst); - if (!dm->db) { DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: no chat db", DM_ID); u_free(dm); inst->dm = NULL; return -1; } - dm_create_tables(dm); - if (etcp_router_bind(inst, ETCP_RT_ID_DM, dm_recv_cb) != 0) { - DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: etcp_router_bind failed", DM_ID); - u_free(dm); inst->dm = NULL; return -1; - } - etcp_add_conn_status_cbk(inst, dm_on_conn_status, inst); - dm_mailbox_set_deliver_cb(inst, dm_core_inject_message, inst); + dm->db = inst->topo_sqlite_db; + if (!dm->db || dm_create_tables(dm) != 0) goto fail; + if (etcp_router_bind(inst, ETCP_RT_ID_DM, dm_recv_cb) != 0) goto fail; dm->initialized = 1; - DEBUG_INFO(DEBUG_CATEGORY_DM, "%s: initialized", DM_ID); + dm->timer = uasync_set_timeout(inst->ua, 10000, dm, dm_tick, "dm_outbox"); + if (!dm->timer) { + etcp_router_unbind(inst, ETCP_RT_ID_DM); + goto fail; + } + DEBUG_INFO(DEBUG_CATEGORY_DM, "%s: initialized persistent outbox", DM_ID); return 0; +fail: + DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: initialization failed db=%p", DM_ID, (void*)dm->db); + u_free(dm); + inst->dm = NULL; + return -1; } +/* Снять триггеры перед освобождением контекста; pending остаётся в БД. */ void dm_core_destroy(struct UTUN_INSTANCE* inst) { - if (!inst || !inst->dm || !inst->dm->initialized) return; - etcp_remove_conn_status_cbk(inst, dm_on_conn_status, inst); + struct dm_state* dm = dm_of(inst); + if (!dm) return; + if (dm->timer) uasync_cancel_timeout(inst->ua, dm->timer); + dm->timer = NULL; + dm->initialized = 0; etcp_router_unbind(inst, ETCP_RT_ID_DM); - u_free(inst->dm); + DEBUG_INFO(DEBUG_CATEGORY_DM, "%s: destroyed, durable outbox preserved", DM_ID); + u_free(dm); inst->dm = NULL; - DEBUG_INFO(DEBUG_CATEGORY_DM, "%s: destroyed", DM_ID); } /* ── JSON для GUI/headless ── */ diff --git a/src/dm/dm_core.h b/src/dm/dm_core.h index 6c3a94f6..b9b9e2d6 100644 --- a/src/dm/dm_core.h +++ b/src/dm/dm_core.h @@ -5,7 +5,7 @@ * - conv_id и content_key детерминированно выводятся обеими сторонами * из своих ключей и pubkey пира (см. dm_crypto.h) — без переговоров. * - сообщения идут через etcp_router (прямое соединение или релей), - * при недоступности пира — через dm_mailbox (storage-узлы). + * при недоступности пира — через dm_mailbox (суперузлы). * - содержимое шифруется E2E (AES-256-CCM), запись подписывается Ed25519. * * Два независимых направленных потока (я→пир и пир→я) с монотонным seq. @@ -29,21 +29,38 @@ extern "C" { * data_enc — AES-256-CCM (шифротекст + tag), sig — Ed25519 автора по всем полям до sig. */ #define DM_MSG_FIXED_HDR 33 /* conv_id+seq+ts+author(32) + ct_len(1) — до ct */ #define DM_MSG_SIG_SIZE 64 +#define DM_DATA_MAX 1024 +#define DM_RECEIPT_SIZE 136 /* Вычислить полную длину канонического тела из его заголовка. * Возвращает 1 и *out_len, либо 0 если буфер мал/повреждён. */ static inline int dm_msg_len(const uint8_t* msg, size_t avail, size_t* out_len) { - if (!msg || avail < DM_MSG_FIXED_HDR) return 0; + if (!msg || !out_len || avail < DM_MSG_FIXED_HDR) return 0; uint8_t ct_len = msg[32]; if (avail < (size_t)DM_MSG_FIXED_HDR + ct_len + 2) return 0; uint16_t data_len; memcpy(&data_len, msg + DM_MSG_FIXED_HDR + ct_len, 2); + if (data_len < 16 || data_len > DM_DATA_MAX + 16) return 0; size_t total = (size_t)DM_MSG_FIXED_HDR + ct_len + 2 + data_len + DM_MSG_SIG_SIZE; if (avail < total) return 0; *out_len = total; return 1; } +/* Общие операции протокола DM; все вызовы в uasync-потоке. + * group_id — только область маршрутизации, идентичность беседы от неё не зависит. + * route_group выбирает общую группу с живым путём, возвращает 1/0. + * verify_message проверяет каноническое тело и подпись автора для recipient. + * accept_message возвращает 0 только после commit (включая точный повтор). + * Квитанция: ["DMACK001":8][conv:8][seq:8][author:8][recipient:8][SHA256(msg):32][sig:64]. + * verify_receipt проверяет подпись получателя и при msg!=NULL точное соответствие телу. + * accept_receipt удаляет только соответствующее исходящее pending-сообщение. */ +int dm_route_group(struct UTUN_INSTANCE* inst, uint64_t peer, uint64_t preferred, uint64_t* group_id); +int dm_verify_message(struct UTUN_INSTANCE* inst, uint64_t recipient, const uint8_t* msg, size_t len); +int dm_accept_message(struct UTUN_INSTANCE* inst, const uint8_t* msg, size_t len, uint8_t receipt[DM_RECEIPT_SIZE]); +int dm_verify_receipt(struct UTUN_INSTANCE* inst, const uint8_t receipt[DM_RECEIPT_SIZE], const uint8_t* msg, size_t len); +int dm_accept_receipt(struct UTUN_INSTANCE* inst, const uint8_t receipt[DM_RECEIPT_SIZE]); + /* Жизненный цикл: инициализация после chat_core_init, destroy перед chat_core_destroy. */ int dm_core_init(struct UTUN_INSTANCE* inst); void dm_core_destroy(struct UTUN_INSTANCE* inst);