Browse Source

Coordinate offline DM media custody and durable deletion on supernodes

master
evgeny 1 day ago
parent
commit
d7a12b24bf
  1. 318
      src/dm/dm_mailbox_media.c

318
src/dm/dm_mailbox_media.c

@ -0,0 +1,318 @@
/* dm_mailbox_media — S хранит координатор отдельно от opaque custody на M. */
#include <stdio.h>
#include <string.h>
#include <openssl/sha.h>
#include "../../lib/mem.h"
#include "../../lib/debug_config.h"
#include "../../lib/json_flat.h"
#include "dm_media_priv.h"
#include "dm_mailbox_media.h"
#include "../utun_instance.h"
#include "../ntp_time.h"
#include "../routing_layer/etcp_router.h"
#include "../routing_layer/topo_node_sqlite.h"
struct dm_mb_media {
struct UTUN_INSTANCE* inst;
void* timer;
};
/* Регистрация идёт в транзакции PUT, до выдачи какого-либо задания M. */
int dm_mailbox_media_register(struct dm_mb_media* media, uint64_t group, uint64_t recipient,
const uint8_t* body, size_t len) {
const uint8_t* desc = dm_message_media(body, len);
if (!desc) return 0;
if (!media) { DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_mailbox_media: registration without owner"); return -1; }
struct UTUN_INSTANCE* inst = media->inst;
uint64_t author, seq;
memcpy(&author, body + 24, 8); memcpy(&seq, body + 8, 8);
uint8_t hash[32]; SHA256(body, len, hash);
sqlite3_stmt* st = dm_media_prepare(inst, "INSERT INTO dm_media_jobs(id,author,recipient,seq,body,body_hash,group_id)"
" VALUES(?,?,?,?,?,?,?) ON CONFLICT(id) DO NOTHING");
if (!st) return -1;
sqlite3_bind_blob(st, 1, desc, 16, SQLITE_STATIC); sqlite3_bind_int64(st, 2, (sqlite3_int64)author);
sqlite3_bind_int64(st, 3, (sqlite3_int64)recipient); sqlite3_bind_int64(st, 4, (sqlite3_int64)seq);
sqlite3_bind_blob(st, 5, body, (int)len, SQLITE_STATIC); sqlite3_bind_blob(st, 6, hash, 32, SQLITE_STATIC);
sqlite3_bind_int64(st, 7, (sqlite3_int64)group);
if (dm_media_step(inst, st)) return -1;
int inserted = sqlite3_changes(inst->topo_sqlite_db);
st = dm_media_prepare(inst, "SELECT body_hash FROM dm_media_jobs WHERE id=?");
if (!st) return -1;
sqlite3_bind_blob(st, 1, desc, 16, SQLITE_STATIC);
int rc = sqlite3_step(st);
int exact = rc == SQLITE_ROW && sqlite3_column_bytes(st, 0) == 32 && !memcmp(sqlite3_column_blob(st, 0), hash, 32);
sqlite3_finalize(st);
if (!exact) { DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_mailbox_media: conflicting media UUID/body"); return -1; }
if (inserted) DEBUG_INFO(DEBUG_CATEGORY_DM, "dm_mailbox_media: registered author=%llu recipient=%llu seq=%llu group=%llu",
(unsigned long long)author, (unsigned long long)recipient, (unsigned long long)seq, (unsigned long long)group);
return 0;
}
/* Список S используется только для пересылки подписанного подтверждения B. */
void dm_mailbox_media_receipt(struct UTUN_INSTANCE* inst, const uint8_t* receipt) {
if (dm_media_verify_receipt(inst, receipt, NULL, 0)) return;
uint64_t author, recipient;
memcpy(&author, receipt + 24, 8); memcpy(&recipient, receipt + 32, 8);
uint64_t* groups = NULL;
int count = 0;
if (topo_node_sqlite_get_member_channels(inst->topo_sqlite_db, inst->node_id, &groups, &count)) {
DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_mailbox_media: common group lookup failed"); return;
}
for (int i = 0; i < count; i++) {
char sql[192];
snprintf(sql, sizeof(sql), "SELECT node_id FROM \"peers_%llu\" WHERE deleted=0", (unsigned long long)groups[i]);
sqlite3_stmt* st = dm_media_prepare(inst, sql);
if (!st) continue;
uint64_t peers[64];
unsigned n = 0;
int rc;
while ((rc = sqlite3_step(st)) == SQLITE_ROW && n < 64) {
uint64_t node = (uint64_t)sqlite3_column_int64(st, 0);
if (node != author && node != recipient && dm_media_scope(inst, groups[i], node, "supernode", author, recipient))
peers[n++] = node;
}
if (rc != SQLITE_ROW && rc != SQLITE_DONE) DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_mailbox_media: supernode scan failed rc=%d", rc);
sqlite3_finalize(st);
for (unsigned j = 0; j < n; j++)
if (peers[j] == inst->node_id || etcp_router_get_path(inst, groups[i], peers[j], ETCP_RT_ID_DM_MEDIA).next_hop_node_id)
dm_media_send(inst, groups[i], peers[j], DM_MEDIA_RECEIPT, receipt, DM_MEDIA_RECEIPT_SIZE);
}
u_free(groups);
}
/* Подтверждение медиадоставки сохраняется до начала DELETE и никогда не отменяется STORED. */
static void dm_job_receipt(struct dm_mb_media* media, uint64_t group, uint64_t source, const uint8_t* r) {
struct UTUN_INSTANCE* inst = media->inst;
if (dm_media_verify_receipt(inst, r, NULL, 0)) return;
uint64_t author, recipient;
memcpy(&author, r + 24, 8); memcpy(&recipient, r + 32, 8);
if (!dm_media_scope(inst, group, inst->node_id, "supernode", author, recipient)) {
DEBUG_WARN(DEBUG_CATEGORY_DM, "dm_mailbox_media: receipt outside S scope source=%llu", (unsigned long long)source); return;
}
sqlite3_stmt* st = dm_media_prepare(inst, "SELECT body_hash,receipt FROM dm_media_jobs WHERE id=?");
if (!st) return;
sqlite3_bind_blob(st, 1, r + 40, 16, SQLITE_STATIC);
int rc = sqlite3_step(st);
if (rc == SQLITE_DONE) { sqlite3_finalize(st); return; } /* Нет выданных копий на этом S. */
int exact = rc == SQLITE_ROW && sqlite3_column_bytes(st, 0) == 32 && !memcmp(sqlite3_column_blob(st, 0), r + 56, 32);
int already = rc == SQLITE_ROW && sqlite3_column_bytes(st, 1) == DM_MEDIA_RECEIPT_SIZE;
sqlite3_finalize(st);
if (!exact) { DEBUG_WARN(DEBUG_CATEGORY_DM, "dm_mailbox_media: receipt body hash mismatch rc=%d", rc); return; }
st = dm_media_prepare(inst, "UPDATE dm_media_jobs SET receipt=?,retry_at=0 WHERE id=? AND receipt IS NULL");
if (!st) return;
sqlite3_bind_blob(st, 1, r, DM_MEDIA_RECEIPT_SIZE, SQLITE_STATIC); sqlite3_bind_blob(st, 2, r + 40, 16, SQLITE_STATIC);
if (dm_media_step(inst, st)) return;
if (!already) {
DEBUG_INFO(DEBUG_CATEGORY_DM, "dm_mailbox_media: recipient confirmed file; delete jobs activated source=%llu",
(unsigned long long)source);
dm_media_send(inst, group, author, DM_MEDIA_RECEIPT, r, DM_MEDIA_RECEIPT_SIZE);
}
}
/* Control от M принимается только для durable зарегистрированного назначения. */
void dm_mailbox_media_recv(struct dm_mb_media* media, uint64_t group, uint64_t source,
uint8_t cmd, const uint8_t* p, size_t len) {
if (!media) { DEBUG_WARN(DEBUG_CATEGORY_DM, "dm_mailbox_media: control after stop"); return; }
struct UTUN_INSTANCE* inst = media->inst;
if (cmd == DM_MEDIA_RECEIPT && len == DM_MEDIA_RECEIPT_SIZE) { dm_job_receipt(media, group, source, p); return; }
sqlite3_stmt* st = dm_media_prepare(inst, "SELECT j.author,j.recipient,j.receipt,j.holder FROM dm_media_jobs j"
" JOIN dm_media_holders h ON h.id=j.id WHERE j.id=? AND h.node_id=? AND h.group_id=?");
if (!st) return;
sqlite3_bind_blob(st, 1, p, 16, SQLITE_STATIC); sqlite3_bind_int64(st, 2, (sqlite3_int64)source);
sqlite3_bind_int64(st, 3, (sqlite3_int64)group);
int rc = sqlite3_step(st);
if (rc == SQLITE_DONE) { sqlite3_finalize(st); return; }
if (rc != SQLITE_ROW) { DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_mailbox_media: location lookup failed rc=%d", rc); sqlite3_finalize(st); return; }
uint64_t author = (uint64_t)sqlite3_column_int64(st, 0), recipient = (uint64_t)sqlite3_column_int64(st, 1);
uint64_t holder = (uint64_t)sqlite3_column_int64(st, 3);
uint8_t receipt[DM_MEDIA_RECEIPT_SIZE];
int delivered = sqlite3_column_bytes(st, 2) == sizeof(receipt);
if (delivered) memcpy(receipt, sqlite3_column_blob(st, 2), sizeof(receipt));
sqlite3_finalize(st);
if (cmd == DM_MEDIA_STORED && len == 24) {
uint64_t expiry;
memcpy(&expiry, p + 16, 8);
if (delivered) { dm_media_send(inst, group, source, DM_MEDIA_DELETE, receipt, sizeof(receipt)); return; }
if (holder != source || expiry <= ntp_time_get_seconds(inst) ||
!dm_media_scope(inst, group, source, "storage", author, recipient)) {
DEBUG_WARN(DEBUG_CATEGORY_DM, "dm_mailbox_media: stale/unauthorized STORED source=%llu expiry=%llu",
(unsigned long long)source, (unsigned long long)expiry); return;
}
st = dm_media_prepare(inst, "UPDATE dm_media_jobs SET stored=1,expires=?,retry_at=0 WHERE id=? AND receipt IS NULL AND stored=0");
if (!st) return;
sqlite3_bind_int64(st, 1, (sqlite3_int64)expiry); sqlite3_bind_blob(st, 2, p, 16, SQLITE_STATIC);
if (dm_media_step(inst, st)) return;
if (sqlite3_changes(inst->topo_sqlite_db))
DEBUG_INFO(DEBUG_CATEGORY_DM, "dm_mailbox_media: custody committed M=%llu expires=%llu",
(unsigned long long)source, (unsigned long long)expiry);
} else if ((cmd == DM_MEDIA_DELETE_ACK || cmd == DM_MEDIA_EXPIRED) && len == 16) {
st = dm_media_prepare(inst, "UPDATE dm_media_holders SET released=1 WHERE id=? AND node_id=?");
if (!st) return;
sqlite3_bind_blob(st, 1, p, 16, SQLITE_STATIC); sqlite3_bind_int64(st, 2, (sqlite3_int64)source);
if (dm_media_step(inst, st)) return;
if (!delivered && cmd == DM_MEDIA_EXPIRED) {
st = dm_media_prepare(inst, "UPDATE dm_media_jobs SET holder=0,stored=0,expired=1,retry_at=0 WHERE id=? AND holder=?");
if (!st) return;
sqlite3_bind_blob(st, 1, p, 16, SQLITE_STATIC); sqlite3_bind_int64(st, 2, (sqlite3_int64)source);
dm_media_step(inst, st);
}
} else if (cmd == DM_MEDIA_REJECT && len == 17 && !delivered) {
DEBUG_WARN(DEBUG_CATEGORY_DM, "dm_mailbox_media: storage refused M=%llu reason=%u", (unsigned long long)source, p[16]);
st = dm_media_prepare(inst, "UPDATE dm_media_holders SET released=1 WHERE id=? AND node_id=?");
if (!st) return;
sqlite3_bind_blob(st, 1, p, 16, SQLITE_STATIC); sqlite3_bind_int64(st, 2, (sqlite3_int64)source);
if (dm_media_step(inst, st)) return;
st = dm_media_prepare(inst, "UPDATE dm_media_jobs SET holder=0,stored=0,retry_at=0 WHERE id=? AND holder=?");
if (!st) return;
sqlite3_bind_blob(st, 1, p, 16, SQLITE_STATIC); sqlite3_bind_int64(st, 2, (sqlite3_int64)source);
dm_media_step(inst, st);
} else DEBUG_WARN(DEBUG_CATEGORY_DM, "dm_mailbox_media: unexpected control cmd=%u bytes=%zu", cmd, len);
}
/* Выбрать один достижимый storage, исключив A/B и прежние отказы/истёкшие назначения. */
static uint64_t dm_job_storage(struct dm_mb_media* media, uint64_t group, const uint8_t id[16], uint64_t a, uint64_t b) {
struct UTUN_INSTANCE* inst = media->inst;
char sql[256];
snprintf(sql, sizeof(sql), "SELECT node_id,adm_tags FROM \"peers_%llu\" WHERE deleted=0 AND node_id NOT IN"
" (SELECT node_id FROM dm_media_holders WHERE id=?) ORDER BY node_id", (unsigned long long)group);
sqlite3_stmt* st = dm_media_prepare(inst, sql);
if (!st) return 0;
sqlite3_bind_blob(st, 1, id, 16, SQLITE_STATIC);
uint64_t found = 0;
int rc;
while ((rc = sqlite3_step(st)) == SQLITE_ROW) {
uint64_t node = (uint64_t)sqlite3_column_int64(st, 0);
const char* tags = (const char*)sqlite3_column_text(st, 1);
char role[16];
if (node == a || node == b || !tags || json_flat_get(tags, "storage", role, sizeof(role)) || strcmp(role, "yes")) continue;
if (node == inst->node_id || etcp_router_get_path(inst, group, node, ETCP_RT_ID_DM_MEDIA).next_hop_node_id) { found = node; break; }
}
if (rc != SQLITE_ROW && rc != SQLITE_DONE) DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_mailbox_media: storage scan failed rc=%d", rc);
sqlite3_finalize(st);
return found;
}
/* Удаление каждого когда-либо выданного назначения повторяется до signed DELETE_ACK. */
static void dm_job_cleanup(struct dm_mb_media* media, const uint8_t id[16], const uint8_t* receipt) {
struct UTUN_INSTANCE* inst = media->inst;
sqlite3_stmt* st = dm_media_prepare(inst, "SELECT node_id,group_id FROM dm_media_holders WHERE id=? AND released=0 LIMIT 16");
if (!st) return;
sqlite3_bind_blob(st, 1, id, 16, SQLITE_STATIC);
uint64_t peers[16], groups[16];
unsigned count = 0;
int rc;
while ((rc = sqlite3_step(st)) == SQLITE_ROW && count < 16) {
peers[count] = (uint64_t)sqlite3_column_int64(st, 0); groups[count++] = (uint64_t)sqlite3_column_int64(st, 1);
}
sqlite3_finalize(st);
if (rc != SQLITE_ROW && rc != SQLITE_DONE) { DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_mailbox_media: cleanup scan failed rc=%d", rc); return; }
for (unsigned i = 0; i < count; i++) dm_media_send(inst, groups[i], peers[i], DM_MEDIA_DELETE, receipt, DM_MEDIA_RECEIPT_SIZE);
if (!count) {
st = dm_media_prepare(inst, "UPDATE dm_media_jobs SET body=NULL WHERE id=? AND receipt IS NOT NULL AND body IS NOT NULL");
if (!st) return;
sqlite3_bind_blob(st, 1, id, 16, SQLITE_STATIC);
if (!dm_media_step(inst, st) && sqlite3_changes(inst->topo_sqlite_db))
DEBUG_INFO(DEBUG_CATEGORY_DM, "dm_mailbox_media: all custody copies deleted; message body removed");
}
}
/* Snapshot statement закрывается до control send, включая локальное совмещение S/M. */
static void dm_mailbox_media_tick(void* arg) {
struct dm_mb_media* media = arg;
struct UTUN_INSTANCE* inst = media->inst;
media->timer = NULL;
uint64_t now = ntp_time_get_seconds(inst);
for (int i = 0; i < 8; i++) {
sqlite3_stmt* st = dm_media_prepare(inst, "SELECT id,author,recipient,body,group_id,holder,stored,expires,receipt,expired"
" FROM dm_media_jobs WHERE retry_at<=? ORDER BY retry_at,seq LIMIT 1");
if (!st) break;
sqlite3_bind_int64(st, 1, (sqlite3_int64)now);
int rc = sqlite3_step(st);
if (rc == SQLITE_DONE) { sqlite3_finalize(st); break; }
if (rc != SQLITE_ROW || sqlite3_column_bytes(st, 0) != 16) {
DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_mailbox_media: job scan failed rc=%d", rc); sqlite3_finalize(st); break;
}
uint8_t id[16], body[DM_MEDIA_BODY_MAX], receipt[DM_MEDIA_RECEIPT_SIZE];
memcpy(id, sqlite3_column_blob(st, 0), 16);
uint64_t author = (uint64_t)sqlite3_column_int64(st, 1), recipient = (uint64_t)sqlite3_column_int64(st, 2);
size_t len = (size_t)sqlite3_column_bytes(st, 3);
if (len <= sizeof(body) && len) memcpy(body, sqlite3_column_blob(st, 3), len);
uint64_t group = (uint64_t)sqlite3_column_int64(st, 4), holder = (uint64_t)sqlite3_column_int64(st, 5);
int stored = sqlite3_column_int(st, 6), expired = sqlite3_column_int(st, 9);
uint64_t expires = (uint64_t)sqlite3_column_int64(st, 7);
int delivered = sqlite3_column_bytes(st, 8) == sizeof(receipt);
if (delivered) memcpy(receipt, sqlite3_column_blob(st, 8), sizeof(receipt));
sqlite3_finalize(st);
st = dm_media_prepare(inst, "UPDATE dm_media_jobs SET retry_at=? WHERE id=?");
if (!st) break;
sqlite3_bind_int64(st, 1, (sqlite3_int64)(now + 5)); sqlite3_bind_blob(st, 2, id, 16, SQLITE_STATIC);
if (dm_media_step(inst, st)) break;
if (delivered) {
dm_job_cleanup(media, id, receipt);
uint64_t gid = group;
if (dm_route_group(inst, author, group, &gid)) dm_media_send(inst, gid, author, DM_MEDIA_RECEIPT, receipt, sizeof(receipt));
continue;
}
if (!len || len > sizeof(body) || !dm_message_media(body, len)) {
DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_mailbox_media: corrupt saved descriptor"); continue;
}
if (!dm_media_scope(inst, group, inst->node_id, "supernode", author, recipient)) continue;
uint64_t gid = group;
if (dm_route_group(inst, recipient, group, &gid)) {
uint8_t offer[24 + DM_MEDIA_BODY_MAX];
uint64_t location = stored && expires > now ? holder : author;
uint64_t expiry = stored ? expires : 0;
memcpy(offer, &recipient, 8); memcpy(offer + 8, &location, 8); memcpy(offer + 16, &expiry, 8);
memcpy(offer + 24, body, len);
dm_media_send(inst, gid, recipient, DM_MEDIA_OFFER, offer, len + 24);
if (expired || (stored && expires <= now)) dm_media_send(inst, gid, recipient, DM_MEDIA_EXPIRED, id, 16);
} else {
if (!holder) {
holder = dm_job_storage(media, group, id, author, recipient);
if (!holder) continue;
if (dm_media_exec(inst, "BEGIN IMMEDIATE")) continue;
st = dm_media_prepare(inst, "INSERT INTO dm_media_holders(id,node_id,group_id) VALUES(?,?,?)");
if (!st) { dm_media_exec(inst, "ROLLBACK"); continue; }
sqlite3_bind_blob(st, 1, id, 16, SQLITE_STATIC); sqlite3_bind_int64(st, 2, (sqlite3_int64)holder);
sqlite3_bind_int64(st, 3, (sqlite3_int64)group);
if (dm_media_step(inst, st)) { dm_media_exec(inst, "ROLLBACK"); continue; }
st = dm_media_prepare(inst, "UPDATE dm_media_jobs SET holder=?,stored=0 WHERE id=? AND receipt IS NULL");
if (!st) { dm_media_exec(inst, "ROLLBACK"); continue; }
sqlite3_bind_int64(st, 1, (sqlite3_int64)holder); sqlite3_bind_blob(st, 2, id, 16, SQLITE_STATIC);
if (dm_media_step(inst, st) || dm_media_exec(inst, "COMMIT")) { dm_media_exec(inst, "ROLLBACK"); continue; }
DEBUG_INFO(DEBUG_CATEGORY_DM, "dm_mailbox_media: assigned storage M=%llu", (unsigned long long)holder);
}
if (!stored) {
uint8_t put[8 + DM_MEDIA_BODY_MAX]; memcpy(put, &recipient, 8); memcpy(put + 8, body, len);
dm_media_send(inst, group, holder, DM_MEDIA_STORE, put, len + 8);
}
}
}
media->timer = uasync_set_timeout(inst->ua, 10000, media, dm_mailbox_media_tick, "dm_mailbox_media");
if (!media->timer) DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_mailbox_media: retry timer allocation failed");
}
struct dm_mb_media* dm_mailbox_media_create(struct UTUN_INSTANCE* inst) {
if (dm_media_exec(inst, "CREATE TABLE IF NOT EXISTS dm_media_jobs(id BLOB PRIMARY KEY,author INTEGER NOT NULL,"
"recipient INTEGER NOT NULL,seq INTEGER NOT NULL,body BLOB,body_hash BLOB NOT NULL,group_id INTEGER NOT NULL,"
"holder INTEGER NOT NULL DEFAULT 0,stored INTEGER NOT NULL DEFAULT 0,expires INTEGER NOT NULL DEFAULT 0,"
"expired INTEGER NOT NULL DEFAULT 0,receipt BLOB,retry_at INTEGER NOT NULL DEFAULT 0)") ||
dm_media_exec(inst, "CREATE TABLE IF NOT EXISTS dm_media_holders(id BLOB NOT NULL,node_id INTEGER NOT NULL,"
"group_id INTEGER NOT NULL,released INTEGER NOT NULL DEFAULT 0,PRIMARY KEY(id,node_id))") ||
dm_media_exec(inst, "UPDATE dm_media_jobs SET retry_at=0")) return NULL;
struct dm_mb_media* media = u_calloc(1, sizeof(*media));
if (!media) { DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_mailbox_media: owner allocation failed"); return NULL; }
media->inst = inst;
media->timer = uasync_set_timeout(inst->ua, 10000, media, dm_mailbox_media_tick, "dm_mailbox_media");
if (!media->timer) { DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_mailbox_media: timer allocation failed"); u_free(media); return NULL; }
DEBUG_INFO(DEBUG_CATEGORY_DM, "dm_mailbox_media: initialized persistent coordinator");
return media;
}
void dm_mailbox_media_destroy(struct dm_mb_media* media) {
if (!media) return;
if (media->timer) uasync_cancel_timeout(media->inst->ua, media->timer);
DEBUG_INFO(DEBUG_CATEGORY_DM, "dm_mailbox_media: destroyed, pending custody/delete jobs retained");
u_free(media);
}
Loading…
Cancel
Save