diff --git a/src/dm/dm_mailbox.c b/src/dm/dm_mailbox.c index 864ba44a..3e64e5f6 100644 --- a/src/dm/dm_mailbox.c +++ b/src/dm/dm_mailbox.c @@ -1,353 +1,406 @@ /* - * dm_mailbox.c — offline-хранение DM-сообщений на storage-узлах - * - * Storage-узел хранит зашифрованные блобы в таблице dm_mail и отдаёт их - * получателю по PULL, удаляя после ACK. Отправитель кладёт блобы по PUT. - * Содержимое для storage-узла непрозрачно (E2E-шифрование + подпись автора). + * dm_mailbox.c — durable очередь суперузла. Сообщения и квитанции непрозрачны; + * приём/удаление атомарны, роли и область маршрута проверяются на каждом PUT. */ - #include "dm_mailbox.h" #include "dm_core.h" -#include "../chat/chat_core.h" #include "../utun_instance.h" #include "../ntp_time.h" #include "../routing_layer/etcp_router.h" #include "../routing_layer/topo_node_sqlite.h" -#include "../transport_layer/etcp_api.h" #include "../transport_layer/etcp.h" #include "../../lib/mem.h" #include "../../lib/ll_queue.h" #include "../../lib/debug_config.h" -#include "../../lib/json_flat.h" -#include "../../lib/platform_compat.h" - -#include -#include #include +#include +#include #define MB_ID "dm_mailbox" +#define MB_PUT 1 /* [recipient:8][каноническое тело] */ +#define MB_DELIVER 2 /* [каноническое тело] */ +#define MB_RECEIPT 3 /* [подписанная квитанция получателя] */ +#define MB_STORED 4 /* [conv:8][seq:8] — durable PUT подтверждён суперузлом */ -/* Подкоманды протокола mailbox (первый байт payload после ROUTER_SVC-заголовка). */ -#define MB_SUBCMD_PUT 0x01 /* отправитель → storage: [recipient:8][sender:8][msg] */ -#define MB_SUBCMD_PULL 0x02 /* получатель → storage: [recipient:8][sender:8][since_seq:8] */ -#define MB_SUBCMD_PULL_RESP 0x03 /* storage → получатель: [recipient:8][count:2][msg]* */ -#define MB_SUBCMD_ACK 0x04 /* получатель → storage: [recipient:8][sender:8][up_to_seq:8] */ - -#define MB_TTL_DAYS 7 -#define MB_MAX_BATCH 32 - -/* Per-instance состояние модуля. */ struct dm_mb_state { struct UTUN_INSTANCE* inst; - sqlite3* db; - int initialized; - dm_mail_deliver_fn deliver_fn; - void* deliver_arg; + sqlite3* db; + void* timer; }; -static struct dm_mb_state* mb_of(struct UTUN_INSTANCE* inst) { - return inst ? inst->dm_mailbox : NULL; -} - -void dm_mailbox_set_deliver_cb(struct UTUN_INSTANCE* inst, dm_mail_deliver_fn fn, void* arg) { - struct dm_mb_state* mb = mb_of(inst); - if (!mb) return; - mb->deliver_fn = fn; - mb->deliver_arg = arg; +/* Проверить роль в конкретной общей группе; storage не заменяет supernode. */ +static int mb_scope(struct dm_mb_state* mb, uint64_t gid, uint64_t super, uint64_t a, uint64_t b) { + if (!gid) return 0; + char ch[64]; + snprintf(ch, sizeof(ch), "%llu", (unsigned long long)gid); + return topo_node_sqlite_member_in_channel(mb->db, ch, super) && + topo_node_sqlite_get_node_type(mb->db, ch, super) == 4 && + topo_node_sqlite_member_in_channel(mb->db, ch, a) && + topo_node_sqlite_member_in_channel(mb->db, ch, b); } -/* Выполнить SQL без результата. Возвращает 0/ошибку. */ -static int mb_exec(struct dm_mb_state* mb, const char* sql) { - char* err = NULL; - if (sqlite3_exec(mb->db, sql, NULL, NULL, &err) != SQLITE_OK) { - DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: sql failed: %s (%s)", MB_ID, sql, err ? err : "?"); - if (err) sqlite3_free(err); +/* Подписать контроль и передать владение пакетом router, учитывая backpressure. */ +static int mb_send(struct dm_mb_state* mb, uint64_t gid, uint64_t dst, uint8_t cmd, const uint8_t* body, size_t len) { + if (!etcp_router_send_q_has_room(mb->inst, gid, dst, ETCP_RT_ID_DM_MAILBOX)) { + DEBUG_DEBUG(DEBUG_CATEGORY_DM, "%s: backpressure peer=%llu", MB_ID, (unsigned long long)dst); return -1; } - return 0; + struct ll_entry* e = queue_entry_new(0); + if (!e) { DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: packet allocation failed", MB_ID); return -1; } + e->dgram = u_malloc(len + 2); + if (!e->dgram) { + DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: payload allocation failed", MB_ID); + queue_entry_free(e); + return -1; + } + e->dgram[0] = ETCP_RT_ID_DM_MAILBOX; + e->dgram[1] = cmd; + memcpy(e->dgram + 2, body, len); + e->len = (uint16_t)(len + 2); + int rc = etcp_route_send(mb->inst, gid, dst, e, 1, ROUTER_FLAG_SIGNED); + if (rc) DEBUG_WARN(DEBUG_CATEGORY_DM, "%s: route failed peer=%llu group=%llu cmd=%u rc=%d", MB_ID, + (unsigned long long)dst, (unsigned long long)gid, cmd, rc); + return rc; } -/* Прочитать seq из канонического тела сообщения (смещение 8). */ -static uint64_t mb_msg_seq(const uint8_t* msg) { - uint64_t seq; memcpy(&seq, msg + 8, 8); return seq; +/* Проверять каждый SQL результат, в том числе границы транзакций. */ +static int mb_exec(struct dm_mb_state* mb, const char* sql) { + int rc = sqlite3_exec(mb->db, sql, NULL, NULL, NULL); + if (rc != SQLITE_OK) DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: SQL rc=%d error=%s statement=%s", + MB_ID, rc, sqlite3_errmsg(mb->db), sql); + return rc == SQLITE_OK ? 0 : -1; } -/* Создать таблицу dm_mail (на любом узле; фактически используется storage-узлами). */ -static void mb_create_tables(struct dm_mb_state* mb) { - mb_exec(mb, "CREATE TABLE IF NOT EXISTS dm_mail (" - " recipient INTEGER NOT NULL," - " sender INTEGER NOT NULL," - " seq INTEGER NOT NULL," - " msg BLOB NOT NULL," - " ttl INTEGER NOT NULL," - " PRIMARY KEY (recipient, sender, seq))"); +/* Принять неизменяемое тело; подтверждать только durable запись или точный повтор. */ +static void mb_put_received(struct dm_mb_state* mb, uint64_t gid, uint64_t source, const uint8_t* p, size_t len) { + if (len < 8) { DEBUG_WARN(DEBUG_CATEGORY_DM, "%s: short PUT bytes=%zu", MB_ID, len); return; } + uint64_t recipient, author, seq; + memcpy(&recipient, p, 8); + p += 8; + len -= 8; + if (dm_verify_message(mb->inst, recipient, p, len) != 0) return; + memcpy(&seq, p + 8, 8); + memcpy(&author, p + 24, 8); + if (source != author || !mb_scope(mb, gid, mb->inst->node_id, author, recipient)) { + DEBUG_WARN(DEBUG_CATEGORY_DM, "%s: unauthorized PUT source=%llu recipient=%llu group=%llu", MB_ID, + (unsigned long long)source, (unsigned long long)recipient, (unsigned long long)gid); + return; + } + sqlite3_stmt* st = NULL; + int rc = sqlite3_prepare_v2(mb->db, + "SELECT receipt FROM dm_mail_receipts WHERE recipient=? AND sender=? AND seq=?", -1, &st, NULL); + if (rc != SQLITE_OK) goto fail; + sqlite3_bind_int64(st, 1, (sqlite3_int64)recipient); + sqlite3_bind_int64(st, 2, (sqlite3_int64)author); + sqlite3_bind_int64(st, 3, (sqlite3_int64)seq); + rc = sqlite3_step(st); + if (rc == SQLITE_ROW) { + uint8_t receipt[DM_RECEIPT_SIZE]; + if (sqlite3_column_bytes(st, 0) != DM_RECEIPT_SIZE) goto fail; + memcpy(receipt, sqlite3_column_blob(st, 0), sizeof(receipt)); + sqlite3_finalize(st); + if (dm_verify_receipt(mb->inst, receipt, p, len) == 0) + mb_send(mb, gid, author, MB_RECEIPT, receipt, sizeof(receipt)); + return; /* Поздний PUT не возрождает доставленную задачу. */ + } + if (rc != SQLITE_DONE) goto fail; + sqlite3_finalize(st); + st = NULL; + rc = sqlite3_prepare_v2(mb->db, + "INSERT INTO dm_pending_messages(recipient,sender,seq,msg,group_id,retry_at) VALUES(?,?,?,?,?,0)" + " ON CONFLICT(recipient,sender,seq) DO NOTHING", -1, &st, NULL); + if (rc != SQLITE_OK) goto fail; + sqlite3_bind_int64(st, 1, (sqlite3_int64)recipient); + sqlite3_bind_int64(st, 2, (sqlite3_int64)author); + sqlite3_bind_int64(st, 3, (sqlite3_int64)seq); + sqlite3_bind_blob(st, 4, p, (int)len, SQLITE_STATIC); + sqlite3_bind_int64(st, 5, (sqlite3_int64)gid); + rc = sqlite3_step(st); + if (rc != SQLITE_DONE) goto fail; + sqlite3_finalize(st); + st = NULL; + rc = sqlite3_prepare_v2(mb->db, + "SELECT msg FROM dm_pending_messages WHERE recipient=? AND sender=? AND seq=?", -1, &st, NULL); + if (rc != SQLITE_OK) goto fail; + sqlite3_bind_int64(st, 1, (sqlite3_int64)recipient); + sqlite3_bind_int64(st, 2, (sqlite3_int64)author); + sqlite3_bind_int64(st, 3, (sqlite3_int64)seq); + rc = sqlite3_step(st); + if (rc != SQLITE_ROW || sqlite3_column_bytes(st, 0) != (int)len || + memcmp(sqlite3_column_blob(st, 0), p, len)) goto fail; + sqlite3_finalize(st); + mb_send(mb, gid, author, MB_STORED, p, 16); + DEBUG_INFO(DEBUG_CATEGORY_DM, "%s: stored sender=%llu recipient=%llu seq=%llu group=%llu", MB_ID, + (unsigned long long)author, (unsigned long long)recipient, (unsigned long long)seq, (unsigned long long)gid); + return; +fail: + DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: PUT database/conflict rc=%d error=%s", MB_ID, rc, sqlite3_errmsg(mb->db)); + sqlite3_finalize(st); } -/* Отправить payload подкомандой subcmd в dst через etcp_router (mode=0, своё E2E). */ -static int mb_route_send(struct dm_mb_state* mb, uint64_t group_id, uint64_t dst, uint8_t subcmd, - const uint8_t* body, size_t body_len) { - 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; } - entry->dgram[0] = ETCP_RT_ID_DM_MAILBOX; - memcpy(entry->dgram + 1, &subcmd, 1); - if (body_len) memcpy(entry->dgram + 1 + 1, body, body_len); - entry->len = (uint16_t)(1 + 1 + body_len); - int rc = etcp_route_send(mb->inst, group_id, dst, entry, 1, 0); - if (rc != 0) { queue_dgram_free(entry); queue_entry_free(entry); return -1; } - return 0; +/* Сначала durable квитанция, затем удаление точной копии в той же транзакции. */ +static void mb_receipt_received(struct dm_mb_state* mb, uint64_t gid, uint64_t source, const uint8_t* r) { + if (dm_verify_receipt(mb->inst, r, NULL, 0) != 0) return; + uint64_t author, recipient, seq; + memcpy(&seq, r + 16, 8); + memcpy(&author, r + 24, 8); + memcpy(&recipient, r + 32, 8); + if (author == mb->inst->node_id) { + dm_accept_receipt(mb->inst, r); + return; + } + if ((source != author && source != recipient) || !mb_scope(mb, gid, mb->inst->node_id, author, recipient)) { + DEBUG_WARN(DEBUG_CATEGORY_DM, "%s: receipt outside authorized scope source=%llu group=%llu", MB_ID, + (unsigned long long)source, (unsigned long long)gid); + return; + } + if (mb_exec(mb, "BEGIN IMMEDIATE") != 0) return; + sqlite3_stmt* st = NULL; + int rc = sqlite3_prepare_v2(mb->db, + "SELECT msg FROM dm_pending_messages WHERE recipient=? AND sender=? AND seq=?", -1, &st, NULL); + if (rc != SQLITE_OK) goto fail; + sqlite3_bind_int64(st, 1, (sqlite3_int64)recipient); + sqlite3_bind_int64(st, 2, (sqlite3_int64)author); + sqlite3_bind_int64(st, 3, (sqlite3_int64)seq); + rc = sqlite3_step(st); + if (rc == SQLITE_ROW) { + if (dm_verify_receipt(mb->inst, r, sqlite3_column_blob(st, 0), (size_t)sqlite3_column_bytes(st, 0)) != 0) goto fail; + } else if (rc != SQLITE_DONE) goto fail; + sqlite3_finalize(st); + st = NULL; + rc = sqlite3_prepare_v2(mb->db, + "INSERT INTO dm_mail_receipts(recipient,sender,seq,receipt,group_id) VALUES(?,?,?,?,?)" + " ON CONFLICT(recipient,sender,seq) DO NOTHING", -1, &st, NULL); + if (rc != SQLITE_OK) goto fail; + sqlite3_bind_int64(st, 1, (sqlite3_int64)recipient); + sqlite3_bind_int64(st, 2, (sqlite3_int64)author); + sqlite3_bind_int64(st, 3, (sqlite3_int64)seq); + sqlite3_bind_blob(st, 4, r, DM_RECEIPT_SIZE, SQLITE_STATIC); + sqlite3_bind_int64(st, 5, (sqlite3_int64)gid); + rc = sqlite3_step(st); + if (rc != SQLITE_DONE) goto fail; + sqlite3_finalize(st); + st = NULL; + rc = sqlite3_prepare_v2(mb->db, "DELETE FROM dm_pending_messages WHERE recipient=? AND sender=? AND seq=?", -1, &st, NULL); + if (rc != SQLITE_OK) goto fail; + sqlite3_bind_int64(st, 1, (sqlite3_int64)recipient); + sqlite3_bind_int64(st, 2, (sqlite3_int64)author); + sqlite3_bind_int64(st, 3, (sqlite3_int64)seq); + rc = sqlite3_step(st); + if (rc != SQLITE_DONE) goto fail; + int deleted = sqlite3_changes(mb->db); + sqlite3_finalize(st); + st = NULL; + if (mb_exec(mb, "COMMIT") != 0) goto fail; + if (deleted) DEBUG_INFO(DEBUG_CATEGORY_DM, "%s: delivered copy removed sender=%llu recipient=%llu seq=%llu", MB_ID, + (unsigned long long)author, (unsigned long long)recipient, (unsigned long long)seq); + mb_send(mb, gid, author, MB_RECEIPT, r, DM_RECEIPT_SIZE); + return; +fail: + DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: receipt commit failed rc=%d error=%s", MB_ID, rc, sqlite3_errmsg(mb->db)); + sqlite3_finalize(st); + mb_exec(mb, "ROLLBACK"); } -/* ── обработка входящих mailbox-пакетов (etcp_router 0x35) ── */ - -static void mb_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) { - if (!entry || !conn || !entry->dgram || entry->len < ROUTER_SVC_PAYLOAD_OFF + 1) { +/* Никаких диапазонных ACK: одна квитанция подтверждает одно тело. */ +static void mb_recv(struct ETCP_CONN* conn, struct ll_entry* entry) { + if (!conn || !entry || !entry->dgram || entry->len <= ROUTER_SVC_PAYLOAD_OFF) { if (entry) { queue_dgram_free(entry); queue_entry_free(entry); } return; } - struct UTUN_INSTANCE* inst = conn->instance; - struct dm_mb_state* mb = mb_of(inst); - if (!inst || !mb || !mb->db) { queue_dgram_free(entry); queue_entry_free(entry); return; } - - const uint8_t* d = entry->dgram; - uint64_t group_id; memcpy(&group_id, d + ROUTER_SVC_GROUP_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 == MB_SUBCMD_PUT) { - /* [subcmd][recipient:8][sender:8][msg] — storage: сохранить блоб */ - if (plen < 1 + 8 + 8 + 16) { queue_dgram_free(entry); queue_entry_free(entry); return; } - uint64_t recipient; memcpy(&recipient, p + 1, 8); - uint64_t sender; memcpy(&sender, p + 9, 8); - const uint8_t* msg = p + 17; - size_t mlen = plen - 17; - uint64_t seq = mb_msg_seq(msg); - int64_t ttl = (int64_t)ntp_time_get_seconds(inst) + MB_TTL_DAYS * 86400; - sqlite3_stmt* st = NULL; - if (sqlite3_prepare_v2(mb->db, - "INSERT OR IGNORE INTO dm_mail(recipient,sender,seq,msg,ttl) VALUES(?,?,?,?,?)", - -1, &st, NULL) == SQLITE_OK) { - sqlite3_bind_int64(st, 1, (sqlite3_int64)recipient); - sqlite3_bind_int64(st, 2, (sqlite3_int64)sender); - sqlite3_bind_int64(st, 3, (sqlite3_int64)seq); - sqlite3_bind_blob(st, 4, msg, (int)mlen, SQLITE_STATIC); - sqlite3_bind_int64(st, 5, ttl); - int rc = sqlite3_step(st); - DEBUG_INFO(DEBUG_CATEGORY_DM, "%s: PUT recipient=0x%016llx sender=0x%016llx seq=%llu rc=%d", - MB_ID, (unsigned long long)recipient, (unsigned long long)sender, - (unsigned long long)seq, rc); + struct dm_mb_state* mb = conn->instance->dm_mailbox; + if (!mb) { DEBUG_WARN(DEBUG_CATEGORY_DM, "%s: receive after stop", MB_ID); goto done; } + uint64_t gid, source; + memcpy(&gid, entry->dgram + ROUTER_SVC_GROUP_OFF, 8); + memcpy(&source, entry->dgram + ROUTER_SVC_SRC_OFF, 8); + const uint8_t* p = entry->dgram + ROUTER_SVC_PAYLOAD_OFF; + size_t len = entry->len - ROUTER_SVC_PAYLOAD_OFF - 1; + if (!(entry->dgram[ROUTER_SVC_FLAGS_OFF] & ROUTER_FLAG_SIGNED)) { + DEBUG_WARN(DEBUG_CATEGORY_DM, "%s: unsigned control source=%llu", MB_ID, (unsigned long long)source); + goto done; + } + if (p[0] == MB_PUT) mb_put_received(mb, gid, source, p + 1, len); + else if (p[0] == MB_DELIVER) { + if (len < DM_MSG_FIXED_HDR) { DEBUG_WARN(DEBUG_CATEGORY_DM, "%s: short delivery", MB_ID); goto done; } + uint64_t author; + memcpy(&author, p + 25, 8); + if (!mb_scope(mb, gid, source, author, mb->inst->node_id)) { + DEBUG_WARN(DEBUG_CATEGORY_DM, "%s: delivery from non-supernode source=%llu", MB_ID, (unsigned long long)source); + goto done; } - if (st) sqlite3_finalize(st); - } else if (subcmd == MB_SUBCMD_PULL) { - /* [subcmd][recipient:8][sender:8][since_seq:8] — storage: вернуть очередь */ - if (plen < 1 + 8 + 8 + 8) { queue_dgram_free(entry); queue_entry_free(entry); return; } - uint64_t recipient; memcpy(&recipient, p + 1, 8); - uint64_t sender; memcpy(&sender, p + 9, 8); - uint64_t since_seq; memcpy(&since_seq, p + 17, 8); - - /* прочитать блобы (до MB_MAX_BATCH) в буфер: [recipient:8][sender:8][count:2][msg]* */ - size_t buf_cap = 18 + (size_t)MB_MAX_BATCH * 4096; - uint8_t* body = u_malloc(buf_cap); - if (!body) { queue_dgram_free(entry); queue_entry_free(entry); return; } - memcpy(body, &recipient, 8); - memcpy(body + 8, &sender, 8); - uint16_t count = 0; - size_t boff = 18; + uint8_t receipt[DM_RECEIPT_SIZE]; + if (dm_accept_message(mb->inst, p + 1, len, receipt) == 0) + mb_send(mb, gid, source, MB_RECEIPT, receipt, sizeof(receipt)); + } else if (p[0] == MB_RECEIPT && len == DM_RECEIPT_SIZE) mb_receipt_received(mb, gid, source, p + 1); + else if (p[0] == MB_STORED && len == 16) { + uint64_t conv, seq; + memcpy(&conv, p + 1, 8); + memcpy(&seq, p + 9, 8); + DEBUG_DEBUG(DEBUG_CATEGORY_DM, "%s: custody confirmed super=%llu conv=%llu seq=%llu (await recipient)", MB_ID, + (unsigned long long)source, (unsigned long long)conv, (unsigned long long)seq); + } else DEBUG_WARN(DEBUG_CATEGORY_DM, "%s: invalid command=%u bytes=%zu", MB_ID, p[0], len); +done: + queue_dgram_free(entry); + queue_entry_free(entry); +} +/* Продвигать ограниченную порцию; SQLite statement закрыт до вызова транспорта. */ +static void mb_tick(void* arg) { + struct dm_mb_state* mb = arg; + mb->timer = NULL; + int64_t now = (int64_t)ntp_time_get_seconds(mb->inst); + for (int i = 0; i < 32; i++) { sqlite3_stmt* st = NULL; - if (sqlite3_prepare_v2(mb->db, - "SELECT msg, length(msg) FROM dm_mail WHERE recipient=? AND sender=? AND seq>? ORDER BY seq LIMIT ?", - -1, &st, NULL) == SQLITE_OK) { - sqlite3_bind_int64(st, 1, (sqlite3_int64)recipient); - sqlite3_bind_int64(st, 2, (sqlite3_int64)sender); - sqlite3_bind_int64(st, 3, (sqlite3_int64)since_seq); - sqlite3_bind_int(st, 4, MB_MAX_BATCH); - while (sqlite3_step(st) == SQLITE_ROW) { - const uint8_t* msg = (const uint8_t*)sqlite3_column_blob(st, 0); - int mlen = sqlite3_column_int(st, 1); - if (!msg || mlen <= 0) continue; - if (boff + (size_t)mlen > buf_cap) break; - memcpy(body + boff, msg, (size_t)mlen); - boff += (size_t)mlen; - count++; - } - } - if (st) sqlite3_finalize(st); - memcpy(body + 16, &count, 2); - DEBUG_INFO(DEBUG_CATEGORY_DM, "%s: PULL recipient=0x%016llx sender=0x%016llx since=%llu → %u msgs", - MB_ID, (unsigned long long)recipient, (unsigned long long)sender, - (unsigned long long)since_seq, count); - if (count > 0) - mb_route_send(mb, group_id, recipient, MB_SUBCMD_PULL_RESP, body, boff); - u_free(body); - } else if (subcmd == MB_SUBCMD_PULL_RESP) { - /* [subcmd][recipient:8][sender:8][count:2][msg]* — получатель: передать в dm_core + ACK storage */ - if (plen < 1 + 8 + 8 + 2) { queue_dgram_free(entry); queue_entry_free(entry); return; } - uint64_t storage; memcpy(&storage, d + ROUTER_SVC_SRC_OFF, 8); - uint64_t recipient; memcpy(&recipient, p + 1, 8); - uint64_t sender; memcpy(&sender, p + 9, 8); - uint16_t count; memcpy(&count, p + 17, 2); - const uint8_t* m = p + 19; - size_t remaining = plen - 19; - uint64_t max_seq = 0; - for (int i = 0; i < count; i++) { - size_t mlen = 0; - if (!dm_msg_len(m, remaining, &mlen)) break; - uint64_t seq = mb_msg_seq(m); - if (seq > max_seq) max_seq = seq; - if (mb->deliver_fn) mb->deliver_fn(m, mlen, mb->deliver_arg); - m += mlen; remaining -= mlen; - } - /* подтвердить storage-узлу приём — удалить доставленное (recipient/sender/up_to_seq) */ - if (count > 0 && max_seq > 0) { - uint8_t ack[24]; - memcpy(ack, &recipient, 8); - memcpy(ack + 8, &sender, 8); - memcpy(ack + 16, &max_seq, 8); - mb_route_send(mb, group_id, storage, MB_SUBCMD_ACK, ack, sizeof(ack)); - DEBUG_INFO(DEBUG_CATEGORY_DM, "%s: PULL_RESP recipient=0x%016llx sender=0x%016llx count=%u → ACK up_to=%llu", - MB_ID, (unsigned long long)recipient, (unsigned long long)sender, count, - (unsigned long long)max_seq); + int rc = sqlite3_prepare_v2(mb->db, + "SELECT recipient,sender,seq,msg,group_id FROM dm_pending_messages" + " WHERE retry_at<=? ORDER BY retry_at,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; + uint64_t recipient = (uint64_t)sqlite3_column_int64(st, 0), author = (uint64_t)sqlite3_column_int64(st, 1); + uint64_t seq = (uint64_t)sqlite3_column_int64(st, 2), gid = (uint64_t)sqlite3_column_int64(st, 4); + size_t len = (size_t)sqlite3_column_bytes(st, 3); + uint8_t* msg = u_malloc(len); + if (!msg) { DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: delivery allocation failed", MB_ID); sqlite3_finalize(st); break; } + memcpy(msg, sqlite3_column_blob(st, 3), len); + sqlite3_finalize(st); + st = NULL; + rc = sqlite3_prepare_v2(mb->db, + "UPDATE dm_pending_messages SET retry_at=? WHERE recipient=? AND sender=? AND seq=?", -1, &st, NULL); + if (rc == SQLITE_OK) { + sqlite3_bind_int64(st, 1, now + 5); + sqlite3_bind_int64(st, 2, (sqlite3_int64)recipient); + sqlite3_bind_int64(st, 3, (sqlite3_int64)author); + sqlite3_bind_int64(st, 4, (sqlite3_int64)seq); + rc = sqlite3_step(st); } - } else if (subcmd == MB_SUBCMD_ACK) { - /* [subcmd][recipient:8][sender:8][up_to_seq:8] — storage: удалить доставленное */ - if (plen < 1 + 8 + 8 + 8) { queue_dgram_free(entry); queue_entry_free(entry); return; } - uint64_t recipient; memcpy(&recipient, p + 1, 8); - uint64_t sender; memcpy(&sender, p + 9, 8); - uint64_t up_to_seq; memcpy(&up_to_seq, p + 17, 8); - sqlite3_stmt* st = NULL; - if (sqlite3_prepare_v2(mb->db, - "DELETE FROM dm_mail WHERE recipient=? AND sender=? AND seq<=?", - -1, &st, NULL) == SQLITE_OK) { - sqlite3_bind_int64(st, 1, (sqlite3_int64)recipient); - sqlite3_bind_int64(st, 2, (sqlite3_int64)sender); - sqlite3_bind_int64(st, 3, (sqlite3_int64)up_to_seq); - sqlite3_step(st); + sqlite3_finalize(st); + st = NULL; + if (rc != SQLITE_DONE) { u_free(msg); goto fail; } + uint64_t route_gid = gid; + if (dm_route_group(mb->inst, recipient, gid, &route_gid) && + mb_scope(mb, route_gid, mb->inst->node_id, author, recipient)) { + mb_send(mb, route_gid, recipient, MB_DELIVER, msg, len); + DEBUG_DEBUG(DEBUG_CATEGORY_DM, "%s: pending delivery sender=%llu recipient=%llu seq=%llu group=%llu", MB_ID, + (unsigned long long)author, (unsigned long long)recipient, (unsigned long long)seq, + (unsigned long long)route_gid); } - if (st) sqlite3_finalize(st); - DEBUG_INFO(DEBUG_CATEGORY_DM, "%s: ACK recipient=0x%016llx sender=0x%016llx up_to=%llu", - MB_ID, (unsigned long long)recipient, (unsigned long long)sender, - (unsigned long long)up_to_seq); - } else { - DEBUG_WARN(DEBUG_CATEGORY_DM, "%s: unknown subcmd=%02x", MB_ID, subcmd); + u_free(msg); + continue; +fail: + DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: pending scan rc=%d error=%s", MB_ID, rc, sqlite3_errmsg(mb->db)); + sqlite3_finalize(st); + break; } - - queue_dgram_free(entry); queue_entry_free(entry); + mb->timer = uasync_set_timeout(mb->inst->ua, 10000, mb, mb_tick, "dm_mailbox"); + if (!mb->timer) DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: retry timer allocation failed", MB_ID); } -/* ── Жизненный цикл ── */ - +/* Новые таблицы протокола; прежний формат и диапазонные ACK не используются. */ int dm_mailbox_init(struct UTUN_INSTANCE* inst) { - if (!inst) return -1; - if (inst->dm_mailbox && inst->dm_mailbox->initialized) return 0; + if (!inst || !inst->topo_sqlite_db) { DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: init without core DB", MB_ID); return -1; } + if (inst->dm_mailbox) return 0; struct dm_mb_state* mb = u_calloc(1, sizeof(*mb)); - if (!mb) return -1; - inst->dm_mailbox = mb; + if (!mb) { DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: state allocation failed", MB_ID); return -1; } mb->inst = inst; - mb->db = chat_core_get_db(inst); - if (!mb->db) { DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: no chat db", MB_ID); u_free(mb); inst->dm_mailbox = NULL; return -1; } - mb_create_tables(mb); - if (etcp_router_bind(inst, ETCP_RT_ID_DM_MAILBOX, mb_recv_cb) != 0) { - DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: etcp_router_bind failed", MB_ID); - u_free(mb); inst->dm_mailbox = NULL; return -1; + mb->db = inst->topo_sqlite_db; + if (mb_exec(mb, "CREATE TABLE IF NOT EXISTS dm_pending_messages (recipient INTEGER NOT NULL,sender INTEGER NOT NULL," + "seq INTEGER NOT NULL,msg BLOB NOT NULL,group_id INTEGER NOT NULL,retry_at INTEGER NOT NULL DEFAULT 0," + "PRIMARY KEY(recipient,sender,seq))") != 0 || + mb_exec(mb, "CREATE TABLE IF NOT EXISTS dm_mail_receipts (recipient INTEGER NOT NULL,sender INTEGER NOT NULL," + "seq INTEGER NOT NULL,receipt BLOB NOT NULL,group_id INTEGER NOT NULL," + "PRIMARY KEY(recipient,sender,seq))") != 0) goto fail; + if (etcp_router_bind(inst, ETCP_RT_ID_DM_MAILBOX, mb_recv) != 0) goto fail; + inst->dm_mailbox = mb; + mb->timer = uasync_set_timeout(inst->ua, 10000, mb, mb_tick, "dm_mailbox"); + if (!mb->timer) { + etcp_router_unbind(inst, ETCP_RT_ID_DM_MAILBOX); + inst->dm_mailbox = NULL; + goto fail; } - mb->initialized = 1; - DEBUG_INFO(DEBUG_CATEGORY_DM, "%s: initialized", MB_ID); + DEBUG_INFO(DEBUG_CATEGORY_DM, "%s: initialized persistent supernode queue", MB_ID); return 0; +fail: + DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: initialization failed", MB_ID); + u_free(mb); + return -1; } +/* Отменить таймер до освобождения owner, durable копии оставить в БД. */ void dm_mailbox_destroy(struct UTUN_INSTANCE* inst) { - struct dm_mb_state* mb = mb_of(inst); - if (!inst || !mb || !mb->initialized) return; + if (!inst || !inst->dm_mailbox) return; + struct dm_mb_state* mb = inst->dm_mailbox; + if (mb->timer) uasync_cancel_timeout(inst->ua, mb->timer); etcp_router_unbind(inst, ETCP_RT_ID_DM_MAILBOX); - u_free(mb); inst->dm_mailbox = NULL; - DEBUG_INFO(DEBUG_CATEGORY_DM, "%s: destroyed", MB_ID); + DEBUG_INFO(DEBUG_CATEGORY_DM, "%s: destroyed, pending retained", MB_ID); + u_free(mb); } -/* ── Публичные операции ── */ - -/* Собрать storage-узлы (storage=yes в adm_tags) из общих групп без дублей. */ -int dm_mailbox_get_storage_nodes(struct UTUN_INSTANCE* inst, uint64_t* out, int max) { - struct dm_mb_state* mb = mb_of(inst); - if (!inst || !out || max <= 0 || !mb || !mb->db) return 0; - uint64_t my_id = inst->node_id; - uint64_t* ch_ids = NULL; - int ch_count = 0; - if (topo_node_sqlite_get_member_channels(mb->db, my_id, &ch_ids, &ch_count) != 0) return 0; - - int found = 0; - char buf[256]; - for (int c = 0; c < ch_count && found < max; c++) { - char ch[64]; snprintf(ch, sizeof(ch), "%llu", (unsigned long long)ch_ids[c]); - char tbl[80]; snprintf(tbl, sizeof(tbl), "peers_%s", ch); - char sql[256]; snprintf(sql, sizeof(sql), - "SELECT node_id, COALESCE(adm_tags,'') FROM \"%s\"", tbl); +/* Найти суперузлы в общей группе обеих сторон; для PUT выбрать только один. */ +static int mb_to_supernodes(struct dm_mb_state* mb, uint64_t recipient, uint64_t author, + uint8_t cmd, const uint8_t* body, size_t len) { + uint64_t* ids = NULL; + int count = 0, sent = 0; + if (topo_node_sqlite_get_member_channels(mb->db, mb->inst->node_id, &ids, &count) != 0) { + DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: membership lookup failed", MB_ID); + return -1; + } + for (int i = 0; i < count; i++) { + char ch[64], sql[256]; + snprintf(ch, sizeof(ch), "%llu", (unsigned long long)ids[i]); + if (!topo_node_sqlite_member_in_channel(mb->db, ch, author) || + !topo_node_sqlite_member_in_channel(mb->db, ch, recipient)) continue; + snprintf(sql, sizeof(sql), "SELECT node_id FROM \"peers_%s\" WHERE node_type=4 AND deleted=0 ORDER BY node_id", ch); sqlite3_stmt* st = NULL; - if (sqlite3_prepare_v2(mb->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 == my_id) continue; - const char* tags = (const char*)sqlite3_column_text(st, 1); - if (tags && json_flat_get(tags, "storage", buf, sizeof(buf)) == 0 - && strcmp(buf, "yes") == 0) { - int dup = 0; - for (int i = 0; i < found; i++) if (out[i] == nid) { dup = 1; break; } - if (!dup) out[found++] = nid; + if (sqlite3_prepare_v2(mb->db, sql, -1, &st, NULL) != SQLITE_OK) { + DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: supernode query failed: %s", MB_ID, sqlite3_errmsg(mb->db)); + continue; + } + int rc; + while ((rc = sqlite3_step(st)) == SQLITE_ROW) { + uint64_t node = (uint64_t)sqlite3_column_int64(st, 0); + if (node == mb->inst->node_id || node == recipient || node == author) continue; + struct ETCP_ROUTER_PATH path = etcp_router_get_path(mb->inst, ids[i], node, ETCP_RT_ID_DM_MAILBOX); + if (!path.next_hop_node_id) continue; + if (!mb_send(mb, ids[i], node, cmd, body, len)) { + sent++; + if (cmd == MB_PUT) break; } } + if (rc != SQLITE_ROW && rc != SQLITE_DONE) + DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: supernode scan failed rc=%d", MB_ID, rc); sqlite3_finalize(st); + if (sent && cmd == MB_PUT) break; } - if (ch_ids) u_free(ch_ids); - return found; + u_free(ids); + return sent; } -/* Отправитель: разослать зашифрованное сообщение на все storage-узлы. */ -int dm_mailbox_put(struct UTUN_INSTANCE* inst, uint64_t recipient, uint64_t sender, - const uint8_t* msg, size_t msg_len) { - struct dm_mb_state* mb = mb_of(inst); - if (!inst || !mb || !msg || msg_len == 0) return -1; - - uint64_t nodes[8]; - int n = dm_mailbox_get_storage_nodes(inst, nodes, 8); - if (n == 0) { - DEBUG_WARN(DEBUG_CATEGORY_DM, "%s: no storage nodes for offline delivery", MB_ID); +/* Исходящее тело остаётся у A до подписанного подтверждения B. */ +int dm_mailbox_put(struct UTUN_INSTANCE* inst, uint64_t recipient, uint64_t sender, const uint8_t* msg, size_t len) { + struct dm_mb_state* mb = inst ? inst->dm_mailbox : NULL; + if (!mb || sender != inst->node_id || dm_verify_message(inst, recipient, msg, len) != 0) { + DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: invalid local PUT", MB_ID); return -1; } - - /* body = [recipient:8][sender:8][msg] */ - size_t body_len = 8 + 8 + msg_len; - uint8_t* body = u_malloc(body_len); - if (!body) return -1; + uint8_t body[8 + DM_MSG_FIXED_HDR + 255 + 2 + DM_DATA_MAX + 16 + DM_MSG_SIG_SIZE]; memcpy(body, &recipient, 8); - memcpy(body + 8, &sender, 8); - memcpy(body + 16, msg, msg_len); - - int sent = 0; - for (int i = 0; i < n; i++) { - if (mb_route_send(mb, 0, nodes[i], MB_SUBCMD_PUT, body, body_len) == 0) sent++; - } - u_free(body); - DEBUG_INFO(DEBUG_CATEGORY_DM, "%s: put → %d/%d storage nodes", MB_ID, sent, n); - return sent > 0 ? 0 : -1; + memcpy(body + 8, msg, len); + int n = mb_to_supernodes(mb, recipient, sender, MB_PUT, body, len + 8); + if (n <= 0) DEBUG_DEBUG(DEBUG_CATEGORY_DM, "%s: no reachable supernode; local outbox retained recipient=%llu", MB_ID, + (unsigned long long)recipient); + return n > 0 ? 0 : -1; } -/* Получатель: запросить почту от отправителя со всех storage-узлов. */ -void dm_mailbox_pull(struct UTUN_INSTANCE* inst, uint64_t peer_node_id, uint64_t since_seq) { - struct dm_mb_state* mb = mb_of(inst); - if (!inst || !mb) return; - uint64_t nodes[8]; - int n = dm_mailbox_get_storage_nodes(inst, nodes, 8); - if (n == 0) return; - - uint8_t body[24]; - memcpy(body, &inst->node_id, 8); - memcpy(body + 8, &peer_node_id, 8); - memcpy(body + 16, &since_seq, 8); - for (int i = 0; i < n; i++) - mb_route_send(mb, 0, nodes[i], MB_SUBCMD_PULL, body, sizeof(body)); - DEBUG_INFO(DEBUG_CATEGORY_DM, "%s: pull from %d storage nodes (peer=0x%016llx since=%llu)", - MB_ID, n, (unsigned long long)peer_node_id, (unsigned long long)since_seq); +/* Копии могут остаться после смены маршрута/суперузла — удалить каждую доступную. */ +void dm_mailbox_receipt(struct UTUN_INSTANCE* inst, const uint8_t* receipt) { + struct dm_mb_state* mb = inst ? inst->dm_mailbox : NULL; + if (!mb || dm_verify_receipt(inst, receipt, NULL, 0) != 0) return; + uint64_t author, recipient; + memcpy(&author, receipt + 24, 8); + memcpy(&recipient, receipt + 32, 8); + mb_to_supernodes(mb, recipient, author, MB_RECEIPT, receipt, DM_RECEIPT_SIZE); } diff --git a/src/dm/dm_mailbox.h b/src/dm/dm_mailbox.h index d8edd0d1..b4c53ecb 100644 --- a/src/dm/dm_mailbox.h +++ b/src/dm/dm_mailbox.h @@ -1,12 +1,12 @@ /* - * dm_mailbox.h — offline-хранение DM-сообщений на storage-узлах + * dm_mailbox.h — очередь offline-сообщений на суперузлах * - * Когда пир недоступен (нет прямой связи), отправитель кладёт зашифрованное - * сообщение на storage-узлы (мемберы общих групп с флагом storage=yes). - * Получатель при выходе в online вытягивает накопившуюся почту, проверяет - * подпись, расшифровывает и вставляет в dm_messages (dedup по seq). + * Суперузел общей группы хранит непрозрачное подписанное тело и повторяет + * доставку при появлении маршрута до получателя, независимо от отправителя. + * Копия удаляется только по квитанции получателя после commit в dm_core. + * Квитанция сохраняется отдельно для повторов PUT и доставки отправителю. * - * Сервис ETCP_RT_ID_DM_MAILBOX (0x35), подкоманды PUT/PULL/PULL_RESP/ACK. + * Сервис ETCP_RT_ID_DM_MAILBOX (0x35); контрольные пакеты подписаны router. */ #ifndef DM_MAILBOX_H @@ -25,23 +25,15 @@ extern "C" { int dm_mailbox_init(struct UTUN_INSTANCE* inst); void dm_mailbox_destroy(struct UTUN_INSTANCE* inst); -/* Коллбэк доставки вытянутых сообщений в dm_core (msg = каноническое тело - * DM-сообщения, включая подпись). Вызывается из uasync-потока. */ -typedef void (*dm_mail_deliver_fn)(const uint8_t* msg, size_t len, void* arg); -void dm_mailbox_set_deliver_cb(struct UTUN_INSTANCE* inst, dm_mail_deliver_fn fn, void* arg); - -/* Отправитель: положить зашифрованное сообщение на storage-узлы. msg — уже - * собранное каноническое тело (зашифровано + подписано). 0=ок, <0=ошибка. */ +/* Поставить PUT выбранному доступному суперузлу общей с recipient группы. + * 0 означает принятие router, хранение подтверждается отдельно на проводе. + * До квитанции получателя источник сохраняет задачу в durable outbox. */ int dm_mailbox_put(struct UTUN_INSTANCE* inst, uint64_t recipient, uint64_t sender, const uint8_t* msg, size_t msg_len); -/* Получатель: вытянуть почту от конкретного отправителя (seq > since_seq) со - * всех известных storage-узлов. Вызывается из uasync-потока. */ -void dm_mailbox_pull(struct UTUN_INSTANCE* inst, uint64_t peer_node_id, uint64_t since_seq); - -/* Собрать node_id storage-узлов (storage=yes в adm_tags) из общих групп. До max - * штук. Возвращает количество найденных. out — буфер на max*8 байт. */ -int dm_mailbox_get_storage_nodes(struct UTUN_INSTANCE* inst, uint64_t* out, int max); +/* Передать точную квитанцию всем доступным суперузлам общей группы сторон. + * receipt — DM_RECEIPT_SIZE байт (см. dm_core.h); вызов только из uasync. */ +void dm_mailbox_receipt(struct UTUN_INSTANCE* inst, const uint8_t* receipt); #ifdef __cplusplus }