From d7a12b24bf316dd7195100ff79bf4abc23cc4617 Mon Sep 17 00:00:00 2001 From: evgeny Date: Thu, 1 Oct 2026 15:52:05 +0300 Subject: [PATCH] Coordinate offline DM media custody and durable deletion on supernodes --- src/dm/dm_mailbox_media.c | 318 ++++++++++++++++++++++++++++++++++++++ 1 file changed, 318 insertions(+) create mode 100644 src/dm/dm_mailbox_media.c diff --git a/src/dm/dm_mailbox_media.c b/src/dm/dm_mailbox_media.c new file mode 100644 index 00000000..234ee781 --- /dev/null +++ b/src/dm/dm_mailbox_media.c @@ -0,0 +1,318 @@ +/* dm_mailbox_media — S хранит координатор отдельно от opaque custody на M. */ +#include +#include +#include + +#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); +}