From 41fc5902be8e86cf0573b9540d98464a9c5ae951 Mon Sep 17 00:00:00 2001 From: evgeny Date: Thu, 1 Oct 2026 15:47:05 +0300 Subject: [PATCH] Implement durable DM downloads, media-node custody, quota and expiry --- src/dm/dm_media.c | 401 ++++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 401 insertions(+) diff --git a/src/dm/dm_media.c b/src/dm/dm_media.c index 91bed9da..c1addcdb 100644 --- a/src/dm/dm_media.c +++ b/src/dm/dm_media.c @@ -591,3 +591,404 @@ void dm_send_file_trampoline(void* arg) { dm_send_file(req->inst, req->conv, req->path, req->name); u_free(req); } + +/* STORE сначала резервирует место и фиксирует expiry; повтор не продлевает TTL. */ +static void dm_custody_store(struct UTUN_INSTANCE* inst, uint64_t group, uint64_t source, const uint8_t* p, size_t len) { + if (len < 8) { DEBUG_WARN(DEBUG_CATEGORY_DM, "dm_media: short STORE"); return; } + uint64_t recipient, author, size; + memcpy(&recipient, p, 8); p += 8; len -= 8; + const uint8_t* desc = dm_message_media(p, len); + if (!desc || dm_verify_message(inst, recipient, p, len)) return; + memcpy(&author, p + 24, 8); memcpy(&size, desc + 16, 8); + if (!dm_media_scope(inst, group, source, "supernode", author, recipient) || + !dm_media_scope(inst, group, inst->node_id, "storage", author, recipient)) { + DEBUG_WARN(DEBUG_CATEGORY_DM, "dm_media: unauthorized STORE source=%llu group=%llu", + (unsigned long long)source, (unsigned long long)group); + return; + } + uint8_t hash[32]; SHA256(p, len, hash); + sqlite3_stmt* st = dm_media_prepare(inst, "SELECT body_hash,receipt FROM dm_custody_done WHERE id=?"); + if (!st) return; + sqlite3_bind_blob(st, 1, desc, 16, SQLITE_STATIC); + int rc = sqlite3_step(st); + if (rc == SQLITE_ROW) { + int exact = sqlite3_column_bytes(st, 0) == 32 && !memcmp(sqlite3_column_blob(st, 0), hash, 32); + uint8_t receipt[DM_MEDIA_RECEIPT_SIZE]; + int delivered = sqlite3_column_bytes(st, 1) == sizeof(receipt); + if (delivered) memcpy(receipt, sqlite3_column_blob(st, 1), sizeof(receipt)); + sqlite3_finalize(st); + if (!exact) { DEBUG_WARN(DEBUG_CATEGORY_DM, "dm_media: STORE conflicts with retired UUID"); return; } + dm_media_send(inst, group, source, delivered ? DM_MEDIA_RECEIPT : DM_MEDIA_EXPIRED, + delivered ? receipt : desc, delivered ? sizeof(receipt) : 16); + return; + } + sqlite3_finalize(st); + if (rc != SQLITE_DONE) { DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: tombstone lookup failed rc=%d", rc); return; } + st = dm_media_prepare(inst, "SELECT body,ready,expires FROM dm_custody WHERE id=?"); + if (!st) return; + sqlite3_bind_blob(st, 1, desc, 16, SQLITE_STATIC); + rc = sqlite3_step(st); + if (rc == SQLITE_ROW) { + int exact = sqlite3_column_bytes(st, 0) == (int)len && !memcmp(sqlite3_column_blob(st, 0), p, len); + int ready = sqlite3_column_int(st, 1); + uint64_t expiry = (uint64_t)sqlite3_column_int64(st, 2); + sqlite3_finalize(st); + if (!exact) { DEBUG_WARN(DEBUG_CATEGORY_DM, "dm_media: conflicting custody body"); return; } + if (ready == 1 && expiry > ntp_time_get_seconds(inst)) { + uint8_t stored[24]; memcpy(stored, desc, 16); memcpy(stored + 16, &expiry, 8); + dm_media_send(inst, group, source, DM_MEDIA_STORED, stored, sizeof(stored)); + } + return; + } + sqlite3_finalize(st); + if (rc != SQLITE_DONE) { DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: custody lookup failed rc=%d", rc); return; } + st = dm_media_prepare(inst, "SELECT COALESCE(sum(cipher_size),0) FROM dm_custody"); + if (!st) return; + rc = sqlite3_step(st); + uint64_t used = rc == SQLITE_ROW ? (uint64_t)sqlite3_column_int64(st, 0) : UINT64_MAX; + sqlite3_finalize(st); + uint64_t quota = (uint64_t)chat_setting_get_int(inst, "dm_media_storage_mb", 1024) * 1024 * 1024; + uint64_t cipher_size = dm_media_cipher_size(size); + if (used > quota || cipher_size > quota - used) { + DEBUG_WARN(DEBUG_CATEGORY_DM, "dm_media: custody quota refused used=%llu requested=%llu quota=%llu", + (unsigned long long)used, (unsigned long long)cipher_size, (unsigned long long)quota); + uint8_t reject[17]; memcpy(reject, desc, 16); reject[16] = 1; + dm_media_send(inst, group, source, DM_MEDIA_REJECT, reject, sizeof(reject)); + return; + } + uint64_t expires = ntp_time_get_seconds(inst) + (uint64_t)chat_setting_get_int(inst, "dm_media_ttl_sec", 604800); + st = dm_media_prepare(inst, "INSERT INTO dm_custody(id,author,recipient,plain_size,cipher_size,body,body_hash," + "group_id,supernode,expires) VALUES(?,?,?,?,?,?,?,?,?,?)"); + if (!st) return; + 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)size); + sqlite3_bind_int64(st, 5, (sqlite3_int64)cipher_size); sqlite3_bind_blob(st, 6, p, (int)len, SQLITE_STATIC); + sqlite3_bind_blob(st, 7, hash, 32, SQLITE_STATIC); sqlite3_bind_int64(st, 8, (sqlite3_int64)group); + sqlite3_bind_int64(st, 9, (sqlite3_int64)source); sqlite3_bind_int64(st, 10, (sqlite3_int64)expires); + if (!dm_media_step(inst, st)) + DEBUG_INFO(DEBUG_CATEGORY_DM, "dm_media: custody reserved author=%llu recipient=%llu bytes=%llu expires=%llu", + (unsigned long long)author, (unsigned long long)recipient, (unsigned long long)cipher_size, (unsigned long long)expires); +} + +/* Удаление подтверждается только после durable удаления файла и копии из БД. */ +static void dm_custody_remove(struct dm_media_state* media, const uint8_t id[16]) { + struct UTUN_INSTANCE* inst = media->inst; + if (dm_job_find(media, id, 1)) return; + sqlite3_stmt* st = dm_media_prepare(inst, "SELECT body_hash,receipt,supernode,group_id FROM dm_custody WHERE id=? AND ready=-1"); + if (!st) return; + sqlite3_bind_blob(st, 1, id, 16, SQLITE_STATIC); + int rc = sqlite3_step(st); + uint8_t hash[32], receipt[DM_MEDIA_RECEIPT_SIZE]; + if (rc != SQLITE_ROW || sqlite3_column_bytes(st, 0) != sizeof(hash)) { + DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: delete state missing/corrupt rc=%d", rc); + sqlite3_finalize(st); return; + } + memcpy(hash, sqlite3_column_blob(st, 0), sizeof(hash)); + int delivered = sqlite3_column_bytes(st, 1) == sizeof(receipt); + if (delivered) memcpy(receipt, sqlite3_column_blob(st, 1), sizeof(receipt)); + uint64_t super = (uint64_t)sqlite3_column_int64(st, 2), group = (uint64_t)sqlite3_column_int64(st, 3); + sqlite3_finalize(st); + char path[1024]; + const char* suffixes[] = {".cipher", ".download"}; + for (unsigned i = 0; i < sizeof(suffixes) / sizeof(suffixes[0]); i++) + if (dm_file_path(media, 1, id, suffixes[i], path, sizeof(path)) || dm_file_remove(path)) return; + if (dm_dir_sync(path) || dm_media_exec(inst, "BEGIN IMMEDIATE")) return; + st = dm_media_prepare(inst, "INSERT INTO dm_custody_done(id,body_hash,receipt) VALUES(?,?,?)" + " ON CONFLICT(id) DO UPDATE SET receipt=COALESCE(excluded.receipt,receipt)"); + if (!st) goto rollback; + sqlite3_bind_blob(st, 1, id, 16, SQLITE_STATIC); sqlite3_bind_blob(st, 2, hash, 32, SQLITE_STATIC); + if (delivered) sqlite3_bind_blob(st, 3, receipt, sizeof(receipt), SQLITE_STATIC); else sqlite3_bind_null(st, 3); + if (dm_media_step(inst, st)) goto rollback; + st = dm_media_prepare(inst, "DELETE FROM dm_custody WHERE id=? AND ready=-1"); + if (!st) goto rollback; + sqlite3_bind_blob(st, 1, id, 16, SQLITE_STATIC); + if (dm_media_step(inst, st) || dm_media_exec(inst, "COMMIT")) goto rollback; + dm_media_send(inst, group, super, delivered ? DM_MEDIA_DELETE_ACK : DM_MEDIA_EXPIRED, id, 16); + DEBUG_INFO(DEBUG_CATEGORY_DM, "dm_media: custody copy removed reason=%s super=%llu", + delivered ? "delivered" : "expired", (unsigned long long)super); + return; +rollback: + dm_media_exec(inst, "ROLLBACK"); +} + +/* Подпись B позволяет любому авторизованному S завершить удаление точного артефакта. */ +static void dm_custody_delete(struct UTUN_INSTANCE* inst, uint64_t group, uint64_t source, const uint8_t* r) { + 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, source, "supernode", author, recipient) || + !dm_media_scope(inst, group, inst->node_id, "storage", author, recipient)) { + DEBUG_WARN(DEBUG_CATEGORY_DM, "dm_media: unauthorized DELETE source=%llu", (unsigned long long)source); return; + } + sqlite3_stmt* st = dm_media_prepare(inst, "SELECT body FROM dm_custody WHERE id=?"); + if (!st) return; + sqlite3_bind_blob(st, 1, r + 40, 16, SQLITE_STATIC); + int rc = sqlite3_step(st); + if (rc == SQLITE_ROW) { + int valid = !dm_media_verify_receipt(inst, r, sqlite3_column_blob(st, 0), (size_t)sqlite3_column_bytes(st, 0)); + sqlite3_finalize(st); + if (!valid) return; + st = dm_media_prepare(inst, "UPDATE dm_custody SET ready=-1,receipt=?,retry_at=0 WHERE id=?"); + 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; + file_transfer_cancel_object(inst->file_transfer, r + 40); + dm_custody_remove(inst->dm_media, r + 40); + return; + } + sqlite3_finalize(st); + if (rc != SQLITE_DONE) { DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: delete lookup failed rc=%d", rc); return; } + st = dm_media_prepare(inst, "SELECT body_hash FROM dm_custody_done WHERE id=?"); + if (!st) return; + sqlite3_bind_blob(st, 1, r + 40, 16, SQLITE_STATIC); + rc = sqlite3_step(st); + int exact = rc == SQLITE_DONE || (rc == SQLITE_ROW && sqlite3_column_bytes(st, 0) == 32 && + !memcmp(sqlite3_column_blob(st, 0), r + 56, 32)); + sqlite3_finalize(st); + if (!exact) { DEBUG_WARN(DEBUG_CATEGORY_DM, "dm_media: delete conflicts with retired artifact"); return; } + st = dm_media_prepare(inst, "INSERT INTO dm_custody_done(id,body_hash,receipt) VALUES(?,?,?)" + " ON CONFLICT(id) DO UPDATE SET receipt=excluded.receipt"); + if (!st) return; + sqlite3_bind_blob(st, 1, r + 40, 16, SQLITE_STATIC); sqlite3_bind_blob(st, 2, r + 56, 32, SQLITE_STATIC); + sqlite3_bind_blob(st, 3, r, DM_MEDIA_RECEIPT_SIZE, SQLITE_STATIC); + if (!dm_media_step(inst, st)) dm_media_send(inst, group, source, DM_MEDIA_DELETE_ACK, r + 40, 16); +} + +/* OFFER не меняет идентичность беседы. Чужой holder допустим только от S. */ +static void dm_media_offer(struct UTUN_INSTANCE* inst, uint64_t group, uint64_t source, const uint8_t* p, size_t len) { + if (len <= 24) { DEBUG_WARN(DEBUG_CATEGORY_DM, "dm_media: short OFFER"); return; } + uint64_t recipient, holder, expiry, author; + memcpy(&recipient, p, 8); memcpy(&holder, p + 8, 8); memcpy(&expiry, p + 16, 8); + p += 24; len -= 24; + const uint8_t* desc = dm_message_media(p, len); + if (recipient != inst->node_id || !desc || dm_verify_message(inst, recipient, p, len)) return; + memcpy(&author, p + 24, 8); + if ((source == author && (holder != author || expiry)) || + (source != author && !dm_media_scope(inst, group, source, "supernode", author, recipient)) || + (holder != author && !dm_media_scope(inst, group, holder, "storage", author, recipient))) { + DEBUG_WARN(DEBUG_CATEGORY_DM, "dm_media: unauthorized OFFER source=%llu holder=%llu", + (unsigned long long)source, (unsigned long long)holder); return; + } + uint8_t receipt[DM_RECEIPT_SIZE]; + if (dm_accept_message(inst, p, len, receipt)) return; + dm_mailbox_receipt(inst, receipt); + sqlite3_stmt* st = dm_media_prepare(inst, "UPDATE dm_files SET holder=?,expires=?,group_id=?,retry_at=0 WHERE id=? AND ready<=0"); + if (!st) return; + sqlite3_bind_int64(st, 1, (sqlite3_int64)holder); sqlite3_bind_int64(st, 2, (sqlite3_int64)expiry); + sqlite3_bind_int64(st, 3, (sqlite3_int64)group); sqlite3_bind_blob(st, 4, desc, 16, SQLITE_STATIC); + dm_media_step(inst, st); +} + +/* Строгие границы control, независимые подписи DM и квитанций. */ +static void dm_media_control(struct UTUN_INSTANCE* inst, uint64_t group, uint64_t source, + uint8_t cmd, const uint8_t* p, size_t len) { + if (!inst->dm_media || inst->dm_media->closing) return; + if (cmd == DM_MEDIA_OFFER) dm_media_offer(inst, group, source, p, len); + else if (cmd == DM_MEDIA_STORE) dm_custody_store(inst, group, source, p, len); + else if (cmd == DM_MEDIA_DELETE && len == DM_MEDIA_RECEIPT_SIZE) dm_custody_delete(inst, group, source, p); + else if (cmd == DM_MEDIA_RECEIPT && len == DM_MEDIA_RECEIPT_SIZE) { + uint64_t author; memcpy(&author, p + 24, 8); + if (author == inst->node_id) { + if (!dm_media_accept_receipt(inst, p)) dm_mailbox_media_receipt(inst, p); + } else dm_mailbox_media_recv(dm_mailbox_get_media(inst), group, source, cmd, p, len); + } else if ((cmd == DM_MEDIA_STORED && len == 24) || (cmd == DM_MEDIA_DELETE_ACK && len == 16) || + (cmd == DM_MEDIA_EXPIRED && len == 16) || (cmd == DM_MEDIA_REJECT && len == 17)) { + dm_mailbox_media_recv(dm_mailbox_get_media(inst), group, source, cmd, p, len); + if (cmd == DM_MEDIA_EXPIRED) { + sqlite3_stmt* st = dm_media_prepare(inst, "SELECT author,recipient FROM dm_files WHERE id=? AND dir=0 AND ready<=0"); + if (!st) return; + sqlite3_bind_blob(st, 1, p, 16, SQLITE_STATIC); + int rc = sqlite3_step(st); + uint64_t author = rc == SQLITE_ROW ? (uint64_t)sqlite3_column_int64(st, 0) : 0; + uint64_t recipient = rc == SQLITE_ROW ? (uint64_t)sqlite3_column_int64(st, 1) : 0; + sqlite3_finalize(st); + if (author && dm_media_scope(inst, group, source, "supernode", author, recipient)) { + st = dm_media_prepare(inst, "UPDATE dm_files SET ready=-1,expires=0 WHERE id=? AND ready<=0"); + if (st) { sqlite3_bind_blob(st, 1, p, 16, SQLITE_STATIC); dm_media_step(inst, st); } + dm_media_notify(inst, dm_derive_conv_id(author, recipient)); + } + } + } else DEBUG_WARN(DEBUG_CATEGORY_DM, "dm_media: invalid control cmd=%u bytes=%zu", cmd, len); +} + +static void dm_media_recv(struct ETCP_CONN* conn, struct ll_entry* e) { + if (!conn || !e || !e->dgram || e->len <= ROUTER_SVC_PAYLOAD_OFF) goto done; + uint64_t group, source; + memcpy(&group, e->dgram + ROUTER_SVC_GROUP_OFF, 8); memcpy(&source, e->dgram + ROUTER_SVC_SRC_OFF, 8); + if (!(e->dgram[ROUTER_SVC_FLAGS_OFF] & ROUTER_FLAG_SIGNED)) { + DEBUG_WARN(DEBUG_CATEGORY_DM, "dm_media: unsigned control source=%llu", (unsigned long long)source); goto done; + } + const uint8_t* p = e->dgram + ROUTER_SVC_PAYLOAD_OFF; + dm_media_control(conn->instance, group, source, p[0], p + 1, e->len - ROUTER_SVC_PAYLOAD_OFF - 1); +done: + if (e) { queue_dgram_free(e); queue_entry_free(e); } +} + +/* Текст и файл имеют независимые durable retries; наличие text ACK не останавливает медиа. */ +static void dm_files_tick(struct dm_media_state* media, uint64_t now) { + struct UTUN_INSTANCE* inst = media->inst; + for (int i = 0; i < 8; i++) { + sqlite3_stmt* st = dm_media_prepare(inst, "SELECT id,dir,author,recipient,group_id,body,ready,holder,expires,receipt" + " FROM dm_files WHERE retry_at<=? ORDER BY retry_at,seq LIMIT 1"); + if (!st) return; + 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_media: local files scan failed rc=%d", rc); sqlite3_finalize(st); return; + } + uint8_t id[16], body[DM_MEDIA_BODY_MAX], receipt[DM_MEDIA_RECEIPT_SIZE]; + memcpy(id, sqlite3_column_blob(st, 0), 16); + int dir = sqlite3_column_int(st, 1), ready = sqlite3_column_int(st, 6); + uint64_t author = (uint64_t)sqlite3_column_int64(st, 2), recipient = (uint64_t)sqlite3_column_int64(st, 3); + uint64_t group = (uint64_t)sqlite3_column_int64(st, 4), holder = (uint64_t)sqlite3_column_int64(st, 7); + uint64_t expiry = (uint64_t)sqlite3_column_int64(st, 8); + size_t len = (size_t)sqlite3_column_bytes(st, 5); + int delivered = sqlite3_column_bytes(st, 9) == sizeof(receipt); + if (len <= sizeof(body)) memcpy(body, sqlite3_column_blob(st, 5), len); + if (delivered) memcpy(receipt, sqlite3_column_blob(st, 9), sizeof(receipt)); + sqlite3_finalize(st); + st = dm_media_prepare(inst, "UPDATE dm_files SET retry_at=? WHERE id=?"); + if (!st) return; + sqlite3_bind_int64(st, 1, (sqlite3_int64)(now + (dir && delivered ? 3600 : 5))); + sqlite3_bind_blob(st, 2, id, 16, SQLITE_STATIC); + if (dm_media_step(inst, st)) return; + const uint8_t* desc = len <= sizeof(body) ? dm_message_media(body, len) : NULL; + if (!desc) { DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: corrupt saved file descriptor"); continue; } + if (dir && !delivered) { + uint64_t gid = group; + if (dm_route_group(inst, recipient, group, &gid)) { + uint8_t offer[24 + DM_MEDIA_BODY_MAX]; + memcpy(offer, &recipient, 8); memcpy(offer + 8, &author, 8); memset(offer + 16, 0, 8); + memcpy(offer + 24, body, len); + dm_media_send(inst, gid, recipient, DM_MEDIA_OFFER, offer, len + 24); + } else dm_mailbox_put(inst, recipient, author, body, len); + } else if (!dir && ready == 1 && delivered) { + uint64_t gid = group; + if (dm_route_group(inst, author, group, &gid)) dm_media_send(inst, gid, author, DM_MEDIA_RECEIPT, receipt, sizeof(receipt)); + dm_mailbox_media_receipt(inst, receipt); + } else if (!dir && ready <= 0 && !dm_job_find(media, id, 0)) { + uint64_t peer = author, gid = group; + if (!dm_route_group(inst, author, group, &gid)) { + peer = holder; + if (holder == author || (expiry && expiry <= now) || !dm_route_group(inst, holder, group, &gid)) continue; + } + dm_file_pull(media, 0, gid, peer, desc, author); + } + } +} + +/* Периодическая уборка TTL не вытесняет подтверждённые файлы раньше их expiry. */ +static void dm_custody_tick(struct dm_media_state* media, uint64_t now) { + struct UTUN_INSTANCE* inst = media->inst; + for (int i = 0; i < 8; i++) { + sqlite3_stmt* st = dm_media_prepare(inst, "SELECT id,body,author,supernode,group_id,expires,ready FROM dm_custody" + " WHERE retry_at<=? ORDER BY retry_at LIMIT 1"); + if (!st) return; + 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_media: custody scan failed rc=%d", rc); sqlite3_finalize(st); return; + } + uint8_t id[16], body[DM_MEDIA_BODY_MAX]; + memcpy(id, sqlite3_column_blob(st, 0), 16); + size_t len = (size_t)sqlite3_column_bytes(st, 1); + if (len <= sizeof(body)) memcpy(body, sqlite3_column_blob(st, 1), len); + uint64_t author = (uint64_t)sqlite3_column_int64(st, 2), super = (uint64_t)sqlite3_column_int64(st, 3); + uint64_t group = (uint64_t)sqlite3_column_int64(st, 4), expires = (uint64_t)sqlite3_column_int64(st, 5); + int ready = sqlite3_column_int(st, 6); + sqlite3_finalize(st); + st = dm_media_prepare(inst, "UPDATE dm_custody SET retry_at=?,ready=CASE WHEN expires<=? THEN -1 ELSE ready END WHERE id=?"); + if (!st) return; + sqlite3_bind_int64(st, 1, (sqlite3_int64)(now + 5)); sqlite3_bind_int64(st, 2, (sqlite3_int64)now); + sqlite3_bind_blob(st, 3, id, 16, SQLITE_STATIC); + if (dm_media_step(inst, st)) return; + if (expires <= now || ready < 0) { + file_transfer_cancel_object(inst->file_transfer, id); + dm_custody_remove(media, id); + } else if (!ready) { + const uint8_t* desc = len <= sizeof(body) ? dm_message_media(body, len) : NULL; + uint64_t gid = group; + if (desc && dm_route_group(inst, author, group, &gid)) dm_file_pull(media, 1, gid, author, desc, author); + } else { + uint8_t stored[24]; memcpy(stored, id, 16); memcpy(stored + 16, &expires, 8); + dm_media_send(inst, group, super, DM_MEDIA_STORED, stored, sizeof(stored)); + } + } +} + +static void dm_media_tick(void* arg) { + struct dm_media_state* media = arg; + media->timer = NULL; + uint64_t now = ntp_time_get_seconds(media->inst); + dm_files_tick(media, now); + dm_custody_tick(media, now); + media->timer = uasync_set_timeout(media->inst->ua, 10000, media, dm_media_tick, "dm_media"); + if (!media->timer) DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: retry timer allocation failed"); +} + +int dm_media_init(struct UTUN_INSTANCE* inst) { + if (!inst || !inst->dm || !inst->media_async || inst->dm_media) { + DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: invalid initialization state"); return -1; + } + struct dm_media_state* media = u_calloc(1, sizeof(*media)); + if (!media) { DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: owner allocation failed"); return -1; } + media->inst = inst; + const char* base = inst->config->global.db_path; + if (!base[0]) base = "."; + int n = snprintf(media->root, sizeof(media->root), "%s/dm_media", base); + int m = snprintf(media->pending, sizeof(media->pending), "%s/dm_pending", base); + if (n < 0 || (size_t)n >= sizeof(media->root) || m < 0 || (size_t)m >= sizeof(media->pending) || + (utun_mkdir(media->root, 0700) && errno != EEXIST) || (utun_mkdir(media->pending, 0700) && errno != EEXIST)) { + DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: cannot create file roots: %s", strerror(errno)); goto fail; + } + if (dm_media_exec(inst, "CREATE TABLE IF NOT EXISTS dm_files(id BLOB PRIMARY KEY,conv_id TEXT NOT NULL,dir INTEGER NOT NULL," + "seq INTEGER NOT NULL,author INTEGER NOT NULL,recipient INTEGER NOT NULL,plain_size INTEGER NOT NULL," + "cipher_hash BLOB NOT NULL,body BLOB NOT NULL,group_id INTEGER NOT NULL,local_path TEXT NOT NULL," + "ready INTEGER NOT NULL DEFAULT 0,holder INTEGER NOT NULL,expires INTEGER NOT NULL DEFAULT 0," + "receipt BLOB,downloading INTEGER NOT NULL DEFAULT 0,retry_at INTEGER NOT NULL DEFAULT 0," + "UNIQUE(conv_id,dir,seq))") || + dm_media_exec(inst, "CREATE TABLE IF NOT EXISTS dm_custody(id BLOB PRIMARY KEY,author INTEGER NOT NULL,recipient INTEGER NOT NULL," + "plain_size INTEGER NOT NULL,cipher_size INTEGER NOT NULL,body BLOB NOT NULL,body_hash BLOB NOT NULL," + "group_id INTEGER NOT NULL,supernode INTEGER NOT NULL,expires INTEGER NOT NULL," + "ready INTEGER NOT NULL DEFAULT 0,receipt BLOB,retry_at INTEGER NOT NULL DEFAULT 0)") || + dm_media_exec(inst, "CREATE TABLE IF NOT EXISTS dm_custody_done(id BLOB PRIMARY KEY,body_hash BLOB NOT NULL,receipt BLOB)") || + dm_media_exec(inst, "UPDATE dm_files SET downloading=0,retry_at=0") || dm_media_exec(inst, "UPDATE dm_custody SET retry_at=0")) goto fail; + inst->dm_media = media; + if (!file_transfer_create(inst, dm_file_lookup, media) || etcp_router_bind(inst, ETCP_RT_ID_DM_MEDIA, dm_media_recv)) goto fail; + media->timer = uasync_set_timeout(inst->ua, 10000, media, dm_media_tick, "dm_media"); + if (!media->timer) goto fail; + DEBUG_INFO(DEBUG_CATEGORY_DM, "dm_media: initialized root=%s custody=%s ttl=%ds quota=%dMiB", media->root, media->pending, + chat_setting_get_int(inst, "dm_media_ttl_sec", 604800), chat_setting_get_int(inst, "dm_media_storage_mb", 1024)); + return 0; +fail: + DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: initialization failed"); + if (inst->dm_media) dm_media_destroy(inst); else u_free(media); + return -1; +} + +/* Снять producers до ожидания workers; callbacks отмены ещё имеют живого owner. */ +void dm_media_quiesce(struct UTUN_INSTANCE* inst) { + struct dm_media_state* media = inst ? inst->dm_media : NULL; + if (!media || media->closing) return; + media->closing = 1; + if (media->timer) uasync_cancel_timeout(inst->ua, media->timer); + media->timer = NULL; + etcp_router_unbind(inst, ETCP_RT_ID_DM_MEDIA); + file_transfer_destroy(inst->file_transfer); + DEBUG_INFO(DEBUG_CATEGORY_DM, "dm_media: quiesced, waiting for workers"); +} + +void dm_media_destroy(struct UTUN_INSTANCE* inst) { + if (!inst || !inst->dm_media) return; + dm_media_quiesce(inst); + struct dm_media_state* media = inst->dm_media; + if (media->jobs) { + DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: destroy before media_async workers drained"); + return; + } + inst->dm_media = NULL; + DEBUG_INFO(DEBUG_CATEGORY_DM, "dm_media: destroyed, durable files/jobs retained"); + u_free(media); +}