diff --git a/src/dm/dm_media.c b/src/dm/dm_media.c new file mode 100644 index 00000000..91bed9da --- /dev/null +++ b/src/dm/dm_media.c @@ -0,0 +1,593 @@ +/* dm_media — E2E файлы бесед и отдельное ограниченное хранилище M. */ +#include +#include +#include +#include +#include +#include +#include + +#include "../../lib/mem.h" +#include "../../lib/platform_compat.h" +#include "../../lib/debug_config.h" +#include "../../lib/json_flat.h" +#include "dm_media_priv.h" +#include "dm_crypto.h" +#include "dm_mailbox.h" +#include "dm_mailbox_media.h" +#include "../utun_instance.h" +#include "../ntp_time.h" +#include "../chat/chat_event.h" +#include "../chat/chat_setting.h" +#include "../routing_layer/etcp_router.h" +#include "../routing_layer/topo_node_sqlite.h" +#include "../transport_layer/secure_channel.h" +#include "../transport_layer/etcp.h" +#include "../media_delivery/file_transfer.h" + +enum dm_file_work { DM_FILE_PREPARE, DM_FILE_RECEIVE, DM_FILE_CUSTODY }; +struct dm_file_job; +struct dm_media_state { + struct UTUN_INSTANCE* inst; + char root[768]; + char pending[768]; + struct dm_file_job* jobs; + void* timer; + int closing; +}; + +struct dm_file_job { + struct dm_file_job* next; + struct dm_media_state* media; + enum dm_file_work work; + struct file_transfer_op* transfer; + uint8_t id[16], key[32], hash[32]; + uint64_t author, size; + char conv[64], name[256]; + char input[1024], cipher[1024], plain[1024], temp[1024]; + int error; +}; + +sqlite3_stmt* dm_media_prepare(struct UTUN_INSTANCE* inst, const char* sql) { + sqlite3_stmt* st = NULL; + if (sqlite3_prepare_v2(inst->topo_sqlite_db, sql, -1, &st, NULL) != SQLITE_OK) + DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: SQL prepare failed: %s statement=%s", sqlite3_errmsg(inst->topo_sqlite_db), sql); + return st; +} + +/* Шаг записи всегда завершает statement; SQL-ошибки нельзя скрывать ACK. */ +int dm_media_step(struct UTUN_INSTANCE* inst, sqlite3_stmt* st) { + if (!st) return -1; + int rc = sqlite3_step(st); + if (rc != SQLITE_DONE) DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: SQL step rc=%d error=%s statement=%s", + rc, sqlite3_errmsg(inst->topo_sqlite_db), sqlite3_sql(st)); + sqlite3_finalize(st); + return rc == SQLITE_DONE ? 0 : -1; +} + +int dm_media_exec(struct UTUN_INSTANCE* inst, const char* sql) { + int rc = sqlite3_exec(inst->topo_sqlite_db, sql, NULL, NULL, NULL); + if (rc != SQLITE_OK) DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: SQL rc=%d error=%s statement=%s", + rc, sqlite3_errmsg(inst->topo_sqlite_db), sql); + return rc == SQLITE_OK ? 0 : -1; +} + +const uint8_t* dm_message_media(const uint8_t* msg, size_t len) { + size_t actual; + if (!dm_msg_len(msg, len, &actual) || actual != len || msg[32] != 4 || memcmp(msg + 33, "file", 4)) return NULL; + return msg + len - DM_MSG_SIG_SIZE - DM_MEDIA_DESC_SIZE; +} + +uint64_t dm_media_cipher_size(uint64_t size) { + return size + (size ? (size + DM_MEDIA_PLAIN_SIZE - 1) / DM_MEDIA_PLAIN_SIZE : 1) * DM_TAG_SIZE; +} + +/* Пути зависят только от UUID, имя пользователя не попадает в файловый путь. */ +static int dm_file_path(struct dm_media_state* media, int custody, const uint8_t id[16], const char* suffix, + char* path, size_t cap) { + char hex[33]; + for (int i = 0; i < 16; i++) snprintf(hex + 2 * i, 3, "%02x", id[i]); + int n = snprintf(path, cap, "%s/%s%s", custody ? media->pending : media->root, hex, suffix); + if (n < 0 || (size_t)n >= cap) { DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: file path too long"); return -1; } + return 0; +} + +/* fflush не гарантирует сохранение: fsync/_commit нужен до публикации и ACK. */ +static int dm_file_sync(FILE* file) { + if (fflush(file) != 0) goto fail; +#ifdef _WIN32 + if (_commit(_fileno(file)) != 0) goto fail; +#else + if (fsync(fileno(file)) != 0) goto fail; +#endif + return 0; +fail: + DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: file sync failed: %s", strerror(errno)); + return -1; +} + +/* Сделать создание/удаление имён durable на POSIX. */ +static int dm_dir_sync(const char* path) { +#ifndef _WIN32 + char parent[1024]; + snprintf(parent, sizeof(parent), "%s", path); + char* slash = strrchr(parent, '/'); + if (!slash) { DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: path without parent: %s", path); return -1; } + *slash = '\0'; + int fd = open(parent, O_RDONLY | O_DIRECTORY); + if (fd < 0) { DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: directory open failed: %s: %s", parent, strerror(errno)); return -1; } + int rc = fsync(fd); + if (rc) DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: directory sync failed: %s: %s", parent, strerror(errno)); + close(fd); + return rc ? -1 : 0; +#else + (void)path; + return 0; +#endif +} + +/* Windows MoveFileEx обеспечивает замену и write-through существующего артефакта. */ +static int dm_file_publish(const char* temp, const char* path) { +#ifdef _WIN32 + if (!MoveFileExA(temp, path, MOVEFILE_REPLACE_EXISTING | MOVEFILE_WRITE_THROUGH)) { + DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: publish failed %s -> %s winerr=%lu", temp, path, GetLastError()); + return -1; + } +#else + if (rename(temp, path) != 0) { + DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: publish failed %s -> %s: %s", temp, path, strerror(errno)); + return -1; + } +#endif + return dm_dir_sync(path); +} + +static int dm_file_remove(const char* path) { + if (remove(path) != 0 && errno != ENOENT) { + DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: remove failed path=%s: %s", path, strerror(errno)); + return -1; + } + return 0; +} + +/* Подписанная роль и живые мемберы в одной области; тип NAT не даёт полномочий. */ +int dm_media_scope(struct UTUN_INSTANCE* inst, uint64_t group, uint64_t node, + const char* role, uint64_t author, uint64_t recipient) { + if (!group) return 0; + char sql[256]; + snprintf(sql, sizeof(sql), "SELECT node_id,adm_tags FROM \"peers_%llu\" WHERE deleted=0 AND node_id IN (?,?,?)", + (unsigned long long)group); + sqlite3_stmt* st = dm_media_prepare(inst, sql); + if (!st) return 0; + sqlite3_bind_int64(st, 1, (sqlite3_int64)node); + sqlite3_bind_int64(st, 2, (sqlite3_int64)author); + sqlite3_bind_int64(st, 3, (sqlite3_int64)recipient); + int has_a = 0, has_b = 0, allowed = 0, rc; + while ((rc = sqlite3_step(st)) == SQLITE_ROW) { + uint64_t id = (uint64_t)sqlite3_column_int64(st, 0); + if (id == author) has_a = 1; + if (id == recipient) has_b = 1; + if (id == node) { + const char* tags = (const char*)sqlite3_column_text(st, 1); + char value[16]; + allowed = !role || (tags && json_flat_get(tags, role, value, sizeof(value)) == 0 && !strcmp(value, "yes")); + } + } + if (rc != SQLITE_DONE) DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: role scan failed rc=%d", rc); + sqlite3_finalize(st); + return rc == SQLITE_DONE && allowed && has_a && has_b; +} + +void dm_media_notify(struct UTUN_INSTANCE* inst, uint64_t conv) { + char id[64]; + int n = snprintf(id, sizeof(id), "%llu", (unsigned long long)conv); + uint8_t evt[1 + 64 + 8]; + evt[0] = (uint8_t)n; + memcpy(evt + 1, id, (size_t)n); + memcpy(evt + 1 + n, &inst->node_id, 8); + chat_event_post(inst, CHAT_EVT_DM_MSG_RECEIVED, evt, n + 9); +} + +/* Квитанция подписывает UUID и хеш всего сообщения с descriptor. */ +int dm_media_verify_receipt(struct UTUN_INSTANCE* inst, const uint8_t* r, const uint8_t* msg, size_t len) { + uint64_t conv, author, recipient, seq; + if (!inst || !r || memcmp(r, "DMFILE01", 8)) goto invalid; + memcpy(&conv, r + 8, 8); memcpy(&seq, r + 16, 8); + memcpy(&author, r + 24, 8); memcpy(&recipient, r + 32, 8); + if (!author || !recipient || author == recipient || !seq || seq > INT64_MAX || + conv != dm_derive_conv_id(author, recipient)) goto invalid; + sqlite3_stmt* st = dm_media_prepare(inst, "SELECT ed25519_pubkey FROM nodes WHERE node_id=?"); + if (!st) return -1; + sqlite3_bind_int64(st, 1, (sqlite3_int64)recipient); + uint8_t key[32]; + int found = sqlite3_step(st) == SQLITE_ROW && sqlite3_column_bytes(st, 0) == 32; + if (found) memcpy(key, sqlite3_column_blob(st, 0), sizeof(key)); + sqlite3_finalize(st); + if (!found || sc_ed25519_verify(key, r, 88, r + 88) != SC_OK) goto invalid; + if (msg) { + const uint8_t* desc = dm_message_media(msg, len); + uint8_t hash[32]; + if (!desc || memcmp(r + 8, msg, 16) || memcmp(r + 24, msg + 24, 8) || memcmp(r + 40, desc, 16)) goto invalid; + SHA256(msg, len, hash); + if (memcmp(hash, r + 56, 32)) goto invalid; + } + return 0; +invalid: + DEBUG_WARN(DEBUG_CATEGORY_DM, "dm_media: rejected media receipt body_bytes=%zu", len); + return -1; +} + +/* Транзакция dm_core связывает историю, очередь и состояние файла. */ +int dm_media_track(struct UTUN_INSTANCE* inst, int dir, const uint8_t* msg, size_t len) { + const uint8_t* desc = dm_message_media(msg, len); + if (!desc) return 0; + struct dm_media_state* media = inst->dm_media; + if (!media || media->closing) { DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: file received while owner stopped"); return -1; } + uint64_t conv, seq, author, size, recipient, group; + memcpy(&conv, msg, 8); memcpy(&seq, msg + 8, 8); memcpy(&author, msg + 24, 8); memcpy(&size, desc + 16, 8); + char id[64], path[1024]; + snprintf(id, sizeof(id), "%llu", (unsigned long long)conv); + sqlite3_stmt* st = dm_media_prepare(inst, "SELECT peer_node_id,group_id FROM dm_conversations WHERE conv_id=?"); + if (!st) return -1; + sqlite3_bind_text(st, 1, id, -1, SQLITE_STATIC); + int rc = sqlite3_step(st); + if (rc != SQLITE_ROW) { DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: conversation absent conv=%s rc=%d", id, rc); sqlite3_finalize(st); return -1; } + recipient = dir ? (uint64_t)sqlite3_column_int64(st, 0) : inst->node_id; + group = (uint64_t)sqlite3_column_int64(st, 1); + sqlite3_finalize(st); + if (dm_file_path(media, 0, desc, ".file", path, sizeof(path))) return -1; + if (dir) { + struct stat sb; + if (stat(path, &sb) || (uint64_t)sb.st_size != size) { + DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: outgoing file not prepared path=%s size=%llu", path, (unsigned long long)size); + return -1; + } + } + st = dm_media_prepare(inst, "INSERT INTO dm_files(id,conv_id,dir,seq,author,recipient,plain_size,cipher_hash,body,group_id," + "local_path,ready,holder) VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?) ON CONFLICT(id) DO NOTHING"); + if (!st) return -1; + sqlite3_bind_blob(st, 1, desc, 16, SQLITE_STATIC); + sqlite3_bind_text(st, 2, id, -1, SQLITE_STATIC); sqlite3_bind_int(st, 3, dir); + sqlite3_bind_int64(st, 4, (sqlite3_int64)seq); sqlite3_bind_int64(st, 5, (sqlite3_int64)author); + sqlite3_bind_int64(st, 6, (sqlite3_int64)recipient); sqlite3_bind_int64(st, 7, (sqlite3_int64)size); + sqlite3_bind_blob(st, 8, desc + 24, 32, SQLITE_STATIC); sqlite3_bind_blob(st, 9, msg, (int)len, SQLITE_STATIC); + sqlite3_bind_int64(st, 10, (sqlite3_int64)group); sqlite3_bind_text(st, 11, path, -1, SQLITE_STATIC); + sqlite3_bind_int(st, 12, dir ? 1 : 0); sqlite3_bind_int64(st, 13, (sqlite3_int64)author); + if (dm_media_step(inst, st)) return -1; + st = dm_media_prepare(inst, "SELECT body FROM dm_files WHERE id=?"); + if (!st) return -1; + sqlite3_bind_blob(st, 1, desc, 16, SQLITE_STATIC); + rc = sqlite3_step(st); + int exact = rc == SQLITE_ROW && sqlite3_column_bytes(st, 0) == (int)len && !memcmp(sqlite3_column_blob(st, 0), msg, len); + sqlite3_finalize(st); + if (!exact) DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: conflicting UUID/body conv=%s seq=%llu", id, (unsigned long long)seq); + return exact ? 0 : -1; +} + +/* Квитанция endpoint удаляет только media pending, история/собственный файл остаются. */ +int dm_media_accept_receipt(struct UTUN_INSTANCE* inst, const uint8_t* r) { + if (dm_media_verify_receipt(inst, r, NULL, 0)) return -1; + uint64_t author; + memcpy(&author, r + 24, 8); + if (author != inst->node_id) { DEBUG_WARN(DEBUG_CATEGORY_DM, "dm_media: receipt for another author"); return -1; } + sqlite3_stmt* st = dm_media_prepare(inst, "SELECT body,receipt FROM dm_files WHERE id=? AND dir=1"); + if (!st) return -1; + sqlite3_bind_blob(st, 1, r + 40, 16, SQLITE_STATIC); + int rc = sqlite3_step(st); + if (rc != SQLITE_ROW || dm_media_verify_receipt(inst, r, sqlite3_column_blob(st, 0), (size_t)sqlite3_column_bytes(st, 0))) { + DEBUG_WARN(DEBUG_CATEGORY_DM, "dm_media: receipt does not match outgoing artifact rc=%d", rc); + sqlite3_finalize(st); + return -1; + } + int already = sqlite3_column_bytes(st, 1) == DM_MEDIA_RECEIPT_SIZE; + sqlite3_finalize(st); + st = dm_media_prepare(inst, "UPDATE dm_files SET receipt=? WHERE id=? AND dir=1"); + if (!st) return -1; + 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 -1; + if (!already) { + uint64_t conv; + memcpy(&conv, r + 8, 8); + DEBUG_INFO(DEBUG_CATEGORY_DM, "dm_media: recipient confirmed file conv=%llu", (unsigned long long)conv); + dm_media_notify(inst, conv); + } + return 0; +} + +static void dm_media_control(struct UTUN_INSTANCE* inst, uint64_t group, uint64_t source, + uint8_t cmd, const uint8_t* body, size_t len); + +/* Local S/M тоже используют тот же проверяемый обработчик, без искусственного CM к себе. */ +int dm_media_send(struct UTUN_INSTANCE* inst, uint64_t group, uint64_t peer, + uint8_t cmd, const uint8_t* body, size_t len) { + if (len > DM_MEDIA_BODY_MAX + 24 || !inst->dm_media || inst->dm_media->closing) { + DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: invalid/stopped control send bytes=%zu", len); + return -1; + } + if (peer == inst->node_id) { dm_media_control(inst, group, peer, cmd, body, len); return 0; } + struct ETCP_ROUTER_CONN* rc = etcp_router_conn_get(inst, group, peer, ETCP_RT_ID_DM_MEDIA); + if (!rc) { DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: router channel allocation failed"); return -1; } + if (!etcp_router_send_q_has_room(inst, group, peer, ETCP_RT_ID_DM_MEDIA)) return -1; + uint8_t packet[1 + DM_MEDIA_BODY_MAX + 24]; + packet[0] = cmd; + memcpy(packet + 1, body, len); + int error = etcp_router_conn_send(rc, packet, len + 1, ROUTER_FLAG_SIGNED); + if (error) DEBUG_WARN(DEBUG_CATEGORY_DM, "dm_media: control send failed peer=%llu cmd=%u rc=%d", + (unsigned long long)peer, cmd, error); + return error; +} + +static struct dm_file_job* dm_job_find(struct dm_media_state* media, const uint8_t id[16], int custody) { + for (struct dm_file_job* job = media->jobs; job; job = job->next) + if ((job->work == DM_FILE_CUSTODY) == custody && !memcmp(job->id, id, 16)) return job; + return NULL; +} + +static void dm_job_free(struct dm_file_job* job) { + struct dm_file_job** p = &job->media->jobs; + while (*p && *p != job) p = &(*p)->next; + if (*p) *p = job->next; + OPENSSL_cleanse(job->key, sizeof(job->key)); + u_free(job); +} + +/* Копирование своего source фиксирует один снимок до шифрования. */ +static int dm_file_copy(const char* source, const char* dest) { + FILE* in = fopen(source, "rb"); + FILE* out = in ? fopen(dest, "wb") : NULL; + int error = -1; + if (!in || !out) { DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: copy open failed %s -> %s: %s", source, dest, strerror(errno)); goto done; } + uint8_t buf[32768]; + uint64_t size = 0; + size_t n; + while ((n = fread(buf, 1, sizeof(buf), in))) { + size += n; + if (size > DM_MEDIA_MAX_SIZE || fwrite(buf, 1, n, out) != n) { + DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: copy size/write failed bytes=%llu", (unsigned long long)size); + goto done; + } + } + if (ferror(in)) { DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: source read failed: %s", strerror(errno)); goto done; } + error = dm_file_sync(out); +done: + if (in && fclose(in)) { DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: source close failed"); error = -1; } + if (out && fclose(out)) { DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: copy close failed"); error = -1; } + return error; +} + +/* Worker не обращается к БД/сети/контексту instance. */ +static void dm_file_work(void* arg) { + struct dm_file_job* job = arg; + job->error = -1; + if (job->work == DM_FILE_PREPARE && dm_file_copy(job->input, job->plain)) return; + FILE* in = fopen(job->work == DM_FILE_PREPARE ? job->plain : job->cipher, "rb"); + FILE* out = job->work == DM_FILE_CUSTODY ? NULL : fopen(job->temp, "wb"); + if (!in || (job->work != DM_FILE_CUSTODY && !out)) { + DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: worker open failed: %s", strerror(errno)); + goto done; + } + if (job->work == DM_FILE_PREPARE) { + job->error = dm_media_encrypt_stream(job->key, job->author, job->id, in, out, &job->size, job->hash); + } else if (job->work == DM_FILE_RECEIVE) { + job->error = dm_media_decrypt_stream(job->key, job->author, job->id, job->size, job->hash, in, out); + } else { + SHA256_CTX hash; + SHA256_Init(&hash); + uint8_t buf[32768], digest[32]; + uint64_t size = 0; + size_t n; + while ((n = fread(buf, 1, sizeof(buf), in))) { size += n; SHA256_Update(&hash, buf, n); } + SHA256_Final(digest, &hash); + if (ferror(in) || size != dm_media_cipher_size(job->size) || memcmp(digest, job->hash, 32)) { + DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: custody hash/size/read failed bytes=%llu expected=%llu", + (unsigned long long)size, (unsigned long long)dm_media_cipher_size(job->size)); + } else job->error = 0; + if (!job->error) { + FILE* durable = fopen(job->cipher, "rb+"); + if (!durable) { DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: custody sync open failed"); job->error = -1; } + else { job->error = dm_file_sync(durable); if (fclose(durable)) job->error = -1; } + } + } + if (!job->error && out) job->error = dm_file_sync(out); +done: + if (in && fclose(in)) { DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: worker input close failed"); job->error = -1; } + if (out && fclose(out)) { DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: worker output close failed"); job->error = -1; } +} + +/* Только uasync публикует проверенный файл и изменяет durable состояние. */ +static void dm_file_done(void* arg, int error) { + struct dm_file_job* job = arg; + struct dm_media_state* media = job->media; + struct UTUN_INSTANCE* inst = media->inst; + int custody = job->work == DM_FILE_CUSTODY; + if (error || job->error || media->closing) goto fail; + if (job->work == DM_FILE_PREPARE) { + if (dm_file_publish(job->temp, job->cipher) || dm_dir_sync(job->plain)) goto fail; + uint8_t data[DM_MEDIA_DESC_SIZE + 256]; + memcpy(data, job->id, 16); memcpy(data + 16, &job->size, 8); memcpy(data + 24, job->hash, 32); + size_t n = strlen(job->name); + memcpy(data + DM_MEDIA_DESC_SIZE, job->name, n); + if (dm_send(inst, job->conv, "file", data, (uint32_t)(DM_MEDIA_DESC_SIZE + n))) goto fail; + DEBUG_INFO(DEBUG_CATEGORY_DM, "dm_media: prepared conv=%s bytes=%llu name=%s", job->conv, + (unsigned long long)job->size, job->name); + dm_job_free(job); + return; + } + const char* sql = custody ? "SELECT body,supernode,group_id,ready,expires FROM dm_custody WHERE id=?" : + "SELECT body,author,group_id,ready,0 FROM dm_files WHERE id=? AND dir=0"; + sqlite3_stmt* st = dm_media_prepare(inst, sql); + if (!st) goto fail; + sqlite3_bind_blob(st, 1, job->id, 16, SQLITE_STATIC); + int rc = sqlite3_step(st); + uint8_t body[DM_MEDIA_BODY_MAX]; + size_t len = rc == SQLITE_ROW ? (size_t)sqlite3_column_bytes(st, 0) : 0; + uint64_t peer = rc == SQLITE_ROW ? (uint64_t)sqlite3_column_int64(st, 1) : 0; + uint64_t group = rc == SQLITE_ROW ? (uint64_t)sqlite3_column_int64(st, 2) : 0; + int ready = rc == SQLITE_ROW ? sqlite3_column_int(st, 3) : -1; + uint64_t expires = rc == SQLITE_ROW ? (uint64_t)sqlite3_column_int64(st, 4) : 0; + if (len <= sizeof(body) && len) memcpy(body, sqlite3_column_blob(st, 0), len); + sqlite3_finalize(st); + if (ready < 0 || !len || len > sizeof(body) || (custody && expires <= ntp_time_get_seconds(inst))) { + DEBUG_WARN(DEBUG_CATEGORY_DM, "dm_media: completed task no longer admissible custody=%d ready=%d expiry=%llu", + custody, ready, (unsigned long long)expires); + goto fail; + } + uint8_t receipt[DM_MEDIA_RECEIPT_SIZE]; + if (!custody) { + memcpy(receipt, "DMFILE01", 8); memcpy(receipt + 8, body, 16); + memcpy(receipt + 24, body + 24, 8); memcpy(receipt + 32, &inst->node_id, 8); + memcpy(receipt + 40, job->id, 16); SHA256(body, len, receipt + 56); + if (sc_ed25519_sign(inst->my_ed25519_privkey, receipt, 88, receipt + 88) != SC_OK) { + DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: media receipt signing failed"); + goto fail; + } + } + char path[1024]; + if (dm_file_path(media, custody, job->id, custody ? ".cipher" : ".file", path, sizeof(path)) || + dm_file_publish(custody ? job->cipher : job->temp, path)) goto fail; + st = dm_media_prepare(inst, custody ? "UPDATE dm_custody SET ready=1 WHERE id=? AND ready=0" : + "UPDATE dm_files SET ready=1,downloading=0,receipt=? WHERE id=? AND ready<=0"); + if (!st) goto fail; + int idx = 1; + if (!custody) sqlite3_bind_blob(st, idx++, receipt, sizeof(receipt), SQLITE_STATIC); + sqlite3_bind_blob(st, idx, job->id, 16, SQLITE_STATIC); + if (dm_media_step(inst, st) || sqlite3_changes(inst->topo_sqlite_db) != 1) goto fail; + if (custody) { + uint8_t stored[24]; memcpy(stored, job->id, 16); memcpy(stored + 16, &expires, 8); + dm_media_send(inst, group, peer, DM_MEDIA_STORED, stored, sizeof(stored)); + } else { + dm_file_remove(job->cipher); + dm_media_send(inst, group, peer, DM_MEDIA_RECEIPT, receipt, sizeof(receipt)); + dm_mailbox_media_receipt(inst, receipt); + uint64_t conv; memcpy(&conv, body, 8); dm_media_notify(inst, conv); + } + DEBUG_INFO(DEBUG_CATEGORY_DM, "dm_media: file committed custody=%d peer=%llu bytes=%llu expires=%llu", custody, + (unsigned long long)peer, (unsigned long long)job->size, (unsigned long long)expires); + dm_job_free(job); + return; +fail: + DEBUG_WARN(DEBUG_CATEGORY_DM, "dm_media: file task failed/cancelled work=%d error=%d worker=%d closing=%d", + job->work, error, job->error, media->closing); + if (job->work == DM_FILE_PREPARE) { dm_file_remove(job->plain); dm_file_remove(job->cipher); } + else if (!custody) { + sqlite3_stmt* st = dm_media_prepare(inst, "UPDATE dm_files SET downloading=0 WHERE id=? AND dir=0"); + if (st) { sqlite3_bind_blob(st, 1, job->id, 16, SQLITE_STATIC); dm_media_step(inst, st); } + } + dm_file_remove(job->temp); + if (job->work != DM_FILE_PREPARE) dm_file_remove(job->cipher); + dm_job_free(job); +} + +static void dm_file_received(void* arg, int error) { + struct dm_file_job* job = arg; + job->transfer = NULL; + if (error || job->media->closing) { dm_file_done(job, error ? error : -2); return; } + media_async_submit(job->media->inst->media_async, job->media->inst->ua, dm_file_work, job, dm_file_done, job); +} + +/* Получать файл может endpoint или M; обе операции владеют одним CM handle в transfer. */ +static void dm_file_pull(struct dm_media_state* media, int custody, uint64_t group, uint64_t peer, const uint8_t* desc, + uint64_t author) { + if (dm_job_find(media, desc, custody)) return; + struct dm_file_job* job = u_calloc(1, sizeof(*job)); + if (!job) { DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: download allocation failed"); return; } + job->media = media; job->work = custody ? DM_FILE_CUSTODY : DM_FILE_RECEIVE; job->author = author; + memcpy(job->id, desc, 16); memcpy(&job->size, desc + 16, 8); memcpy(job->hash, desc + 24, 32); + if (dm_file_path(media, custody, job->id, ".download", job->cipher, sizeof(job->cipher)) || + dm_file_path(media, custody, job->id, ".decrypt", job->temp, sizeof(job->temp))) { u_free(job); return; } + if (!custody) { + sqlite3_stmt* st = dm_media_prepare(media->inst, "SELECT x25519_pubkey FROM nodes WHERE node_id=?"); + if (!st) { u_free(job); return; } + sqlite3_bind_int64(st, 1, (sqlite3_int64)author); + int rc = sqlite3_step(st); + uint8_t pub[32]; + int found = rc == SQLITE_ROW && sqlite3_column_bytes(st, 0) == sizeof(pub); + if (found) memcpy(pub, sqlite3_column_blob(st, 0), sizeof(pub)); + sqlite3_finalize(st); + if (!found || dm_derive_content_key(media->inst->my_keys.private_key, pub, job->key)) { + DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: target key unavailable author=%llu", (unsigned long long)author); + u_free(job); + return; + } + } + job->next = media->jobs; media->jobs = job; + job->transfer = file_transfer_pull(media->inst->file_transfer, group, peer, job->id, dm_media_cipher_size(job->size), + job->cipher, dm_file_received, job); + if (!job->transfer) { dm_file_done(job, -1); return; } + if (!custody) { + sqlite3_stmt* st = dm_media_prepare(media->inst, "UPDATE dm_files SET downloading=1 WHERE id=? AND dir=0"); + if (st) { sqlite3_bind_blob(st, 1, job->id, 16, SQLITE_STATIC); dm_media_step(media->inst, st); } + } +} + +/* Общий авторизатор сервера transfer: A отдаёт B/storage, M — только B по custody. */ +static int dm_file_lookup(void* arg, uint64_t group, uint64_t peer, const uint8_t id[16], uint64_t size, + char* path, size_t cap) { + struct dm_media_state* media = arg; + struct UTUN_INSTANCE* inst = media->inst; + if (media->closing) return -1; + for (int custody = 0; custody <= 1; custody++) { + sqlite3_stmt* st = dm_media_prepare(inst, custody ? + "SELECT author,recipient,plain_size,ready,expires FROM dm_custody WHERE id=?" : + "SELECT author,recipient,plain_size,ready,0 FROM dm_files WHERE id=? AND dir=1"); + if (!st) return -1; + sqlite3_bind_blob(st, 1, id, 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; + uint64_t plain = rc == SQLITE_ROW ? (uint64_t)sqlite3_column_int64(st, 2) : 0; + int ready = rc == SQLITE_ROW ? sqlite3_column_int(st, 3) : 0; + uint64_t expiry = rc == SQLITE_ROW ? (uint64_t)sqlite3_column_int64(st, 4) : 0; + sqlite3_finalize(st); + if (!author || ready != 1 || size != dm_media_cipher_size(plain) || + (custody && expiry <= ntp_time_get_seconds(inst))) continue; + if ((peer == recipient && dm_media_scope(inst, group, peer, NULL, author, recipient)) || + (!custody && dm_media_scope(inst, group, peer, "storage", author, recipient))) + return dm_file_path(media, custody, id, ".cipher", path, cap); + } + DEBUG_WARN(DEBUG_CATEGORY_DM, "dm_media: file serve denied peer=%llu bytes=%llu", (unsigned long long)peer, (unsigned long long)size); + return -1; +} + +/* Путь исходника копируется; caller может освободить строку после возврата. */ +int dm_send_file(struct UTUN_INSTANCE* inst, const char* conv, const char* path, const char* name) { + struct dm_media_state* media = inst ? inst->dm_media : NULL; + if (!media || media->closing || !inst->media_async || !conv || !path || !name || !*name || strlen(name) > 255 || + strchr(name, '/') || strchr(name, '\\') || !strcmp(name, ".") || !strcmp(name, "..") || strlen(path) >= 1024) { + DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: invalid file send arguments/state"); + return -1; + } + unsigned preparing = 0; + for (struct dm_file_job* job = media->jobs; job; job = job->next) if (job->work == DM_FILE_PREPARE) preparing++; + if (preparing >= 4) { DEBUG_WARN(DEBUG_CATEGORY_DM, "dm_media: preparation limit reached count=%u", preparing); return -1; } + struct dm_file_job* job = u_calloc(1, sizeof(*job)); + if (!job) { DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: preparation allocation failed"); return -1; } + job->media = media; job->work = DM_FILE_PREPARE; job->author = inst->node_id; + snprintf(job->conv, sizeof(job->conv), "%s", conv); snprintf(job->name, sizeof(job->name), "%s", name); + snprintf(job->input, sizeof(job->input), "%s", path); + sqlite3_stmt* st = dm_media_prepare(inst, "SELECT peer_x25519 FROM dm_conversations WHERE conv_id=?"); + if (!st) { u_free(job); return -1; } + sqlite3_bind_text(st, 1, conv, -1, SQLITE_STATIC); + uint8_t pub[32]; + int rc = sqlite3_step(st), found = rc == SQLITE_ROW && sqlite3_column_bytes(st, 0) == sizeof(pub); + if (found) memcpy(pub, sqlite3_column_blob(st, 0), sizeof(pub)); + sqlite3_finalize(st); + if (!found || dm_derive_content_key(inst->my_keys.private_key, pub, job->key) || RAND_bytes(job->id, sizeof(job->id)) != 1 || + dm_file_path(media, 0, job->id, ".file", job->plain, sizeof(job->plain)) || + dm_file_path(media, 0, job->id, ".cipher", job->cipher, sizeof(job->cipher)) || + dm_file_path(media, 0, job->id, ".encrypt", job->temp, sizeof(job->temp))) { + DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: preparation keys/paths failed conv=%s", conv); + OPENSSL_cleanse(job->key, sizeof(job->key)); u_free(job); return -1; + } + job->next = media->jobs; media->jobs = job; + DEBUG_INFO(DEBUG_CATEGORY_DM, "dm_media: preparing conv=%s name=%s", conv, name); + media_async_submit(inst->media_async, inst->ua, dm_file_work, job, dm_file_done, job); + return 0; +} + +void dm_send_file_trampoline(void* arg) { + struct dm_file_req* req = arg; + if (!req) return; + dm_send_file(req->inst, req->conv, req->path, req->name); + u_free(req); +}