1 changed files with 593 additions and 0 deletions
@ -0,0 +1,593 @@
|
||||
/* dm_media — E2E файлы бесед и отдельное ограниченное хранилище M. */ |
||||
#include <stdio.h> |
||||
#include <string.h> |
||||
#include <stdlib.h> |
||||
#include <errno.h> |
||||
#include <sys/stat.h> |
||||
#include <openssl/sha.h> |
||||
#include <openssl/rand.h> |
||||
|
||||
#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); |
||||
} |
||||
Loading…
Reference in new issue