Browse Source

Deliver offline DM through durable supernode custody

master
evgeny 2 days ago
parent
commit
5800902a42
  1. 623
      src/dm/dm_mailbox.c
  2. 32
      src/dm/dm_mailbox.h

623
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 <string.h>
#include <stdlib.h>
#include <sqlite3.h>
#include <string.h>
#include <stdio.h>
#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);
}

32
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
}

Loading…
Cancel
Save