From b47918fad5dedb4d3ded24cbf13a7dd9774126a5 Mon Sep 17 00:00:00 2001 From: evgeny Date: Thu, 1 Oct 2026 15:57:30 +0300 Subject: [PATCH] Keep media publication and durable deletion off the event loop --- src/dm/dm_media.c | 56 +++++++++++++++++++++--------- src/media_delivery/file_transfer.c | 4 +-- 2 files changed, 42 insertions(+), 18 deletions(-) diff --git a/src/dm/dm_media.c b/src/dm/dm_media.c index c1addcdb..b5c31b22 100644 --- a/src/dm/dm_media.c +++ b/src/dm/dm_media.c @@ -25,7 +25,7 @@ #include "../transport_layer/etcp.h" #include "../media_delivery/file_transfer.h" -enum dm_file_work { DM_FILE_PREPARE, DM_FILE_RECEIVE, DM_FILE_CUSTODY }; +enum dm_file_work { DM_FILE_PREPARE, DM_FILE_RECEIVE, DM_FILE_CUSTODY, DM_FILE_DELETE }; struct dm_file_job; struct dm_media_state { struct UTUN_INSTANCE* inst; @@ -44,7 +44,7 @@ struct dm_file_job { 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]; + char input[1024], cipher[1024], plain[1024], temp[1024], output[1024]; int error; }; @@ -319,7 +319,7 @@ int dm_media_send(struct UTUN_INSTANCE* inst, uint64_t group, uint64_t peer, 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; + if ((job->work == DM_FILE_CUSTODY || job->work == DM_FILE_DELETE) == custody && !memcmp(job->id, id, 16)) return job; return NULL; } @@ -359,6 +359,10 @@ done: static void dm_file_work(void* arg) { struct dm_file_job* job = arg; job->error = -1; + if (job->work == DM_FILE_DELETE) { + if (!dm_file_remove(job->cipher) && !dm_file_remove(job->temp)) job->error = dm_dir_sync(job->cipher); + return; + } 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"); @@ -392,17 +396,30 @@ static void dm_file_work(void* arg) { 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; } + if (!job->error) { + if (job->work == DM_FILE_PREPARE) { + job->error = dm_file_publish(job->temp, job->cipher); + if (!job->error) job->error = dm_dir_sync(job->plain); + } else job->error = dm_file_publish(job->work == DM_FILE_CUSTODY ? job->cipher : job->temp, job->output); + } } +static void dm_custody_finalize(struct dm_media_state* media, const uint8_t id[16]); + /* Только 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; + if (job->work == DM_FILE_DELETE) { + if (!error && !job->error && !media->closing) dm_custody_finalize(media, job->id); + else DEBUG_WARN(DEBUG_CATEGORY_DM, "dm_media: custody deletion worker failed/cancelled error=%d worker=%d", error, job->error); + dm_job_free(job); + return; + } 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); @@ -442,9 +459,6 @@ static void dm_file_done(void* arg, int error) { 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; @@ -474,6 +488,7 @@ fail: if (st) { sqlite3_bind_blob(st, 1, job->id, 16, SQLITE_STATIC); dm_media_step(inst, st); } } dm_file_remove(job->temp); + if (job->output[0]) dm_file_remove(job->output); if (job->work != DM_FILE_PREPARE) dm_file_remove(job->cipher); dm_job_free(job); } @@ -494,7 +509,8 @@ static void dm_file_pull(struct dm_media_state* media, int custody, uint64_t gro 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; } + dm_file_path(media, custody, job->id, ".decrypt", job->temp, sizeof(job->temp)) || + dm_file_path(media, custody, job->id, custody ? ".cipher" : ".file", job->output, sizeof(job->output))) { 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; } @@ -671,9 +687,8 @@ static void dm_custody_store(struct UTUN_INSTANCE* inst, uint64_t group, uint64_ } /* Удаление подтверждается только после durable удаления файла и копии из БД. */ -static void dm_custody_remove(struct dm_media_state* media, const uint8_t id[16]) { +static void dm_custody_finalize(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); @@ -688,11 +703,7 @@ static void dm_custody_remove(struct dm_media_state* media, const uint8_t id[16] 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; + if (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; @@ -711,6 +722,19 @@ rollback: dm_media_exec(inst, "ROLLBACK"); } +/* Удаление больших файлов и fsync каталога выполняются worker, а не uasync. */ +static void dm_custody_remove(struct dm_media_state* media, const uint8_t id[16]) { + if (dm_job_find(media, id, 1)) return; + struct dm_file_job* job = u_calloc(1, sizeof(*job)); + if (!job) { DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: delete worker allocation failed"); return; } + job->media = media; job->work = DM_FILE_DELETE; + memcpy(job->id, id, 16); + if (dm_file_path(media, 1, id, ".cipher", job->cipher, sizeof(job->cipher)) || + dm_file_path(media, 1, id, ".download", job->temp, sizeof(job->temp))) { u_free(job); return; } + job->next = media->jobs; media->jobs = job; + media_async_submit(media->inst->media_async, media->inst->ua, dm_file_work, job, dm_file_done, job); +} + /* Подпись 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; @@ -789,7 +813,7 @@ static void dm_media_control(struct UTUN_INSTANCE* inst, uint64_t group, uint64_ 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); + dm_media_accept_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)) { diff --git a/src/media_delivery/file_transfer.c b/src/media_delivery/file_transfer.c index e11f5769..d7c6358d 100644 --- a/src/media_delivery/file_transfer.c +++ b/src/media_delivery/file_transfer.c @@ -48,7 +48,7 @@ struct file_transfer_op { static void ft_finish(struct file_transfer_op* op, int error) { struct file_transfer* ft = op->ft; if (op->timer) uasync_cancel_timeout(ft->inst->ua, op->timer); - if (op->pump) uasync_cancel_timeout(ft->inst->ua, op->pump); + if (op->pump) uasync_call_soon_cancel(ft->inst->ua, op->pump); etcp_router_cancel_send_ready(ft->inst, op->group, op->peer, ETCP_RT_ID_FILE_TRANSFER, &op->waiter); if (op->cm) conn_mgr_close(op->cm); if (op->file && fclose(op->file) != 0) { @@ -114,7 +114,7 @@ static void ft_pump(void* arg) { op->pump = NULL; if (!op->up || op->receiving) return; struct file_transfer* ft = op->ft; - for (unsigned i = 0; i < 8 && op->offset < op->size; i++) { + for (unsigned i = 0; i < 1 && op->offset < op->size; i++) { if (!etcp_router_send_q_has_room(ft->inst, op->group, op->peer, ETCP_RT_ID_FILE_TRANSFER)) { etcp_router_on_send_ready(ft->inst, op->group, op->peer, ETCP_RT_ID_FILE_TRANSFER, &op->waiter, ft_ready, op); return;