diff --git a/src/media_delivery/file_transfer.c b/src/media_delivery/file_transfer.c new file mode 100644 index 00000000..600ce3ba --- /dev/null +++ b/src/media_delivery/file_transfer.c @@ -0,0 +1,331 @@ +/* file_transfer — signed router stream с ограниченной очередью и CM-владением. */ +#include +#include +#include +#include + +#include "../../lib/mem.h" +#include "../../lib/platform_compat.h" +#include "../../lib/debug_config.h" +#include "file_transfer.h" +#include "../utun_instance.h" +#include "../routing_layer/conn_mgr.h" +#include "../routing_layer/etcp_router.h" +#include "../transport_layer/etcp.h" + +#define FT_REQUEST 1 /* tid:16, object:16, size:8, offset:8 */ +#define FT_DATA 2 /* tid:16, offset:8, bytes */ +#define FT_DONE 3 /* tid:16 */ +#define FT_REFUSE 4 /* tid:16 */ +#define FT_CHUNK 32768u +#define FT_MAX_OPS 16 + +struct file_transfer { + struct UTUN_INSTANCE* inst; + file_transfer_lookup_fn lookup; + void* arg; + struct file_transfer_op* ops; + unsigned count; + int closing; +}; + +struct file_transfer_op { + struct file_transfer_op* next; + struct file_transfer* ft; + uint64_t group, peer, size, offset, progress_tb, request_tb; + uint8_t tid[16], object[16]; + int receiving, up; + FILE* file; + struct CONN_MGR_HANDLE* cm; + struct queue_waiter_handle waiter; + void* timer; + void* pump; + file_transfer_done_fn done; + void* arg; +}; + +/* Очистить все локальные триггеры до callback владельца. */ +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); + 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) { + DEBUG_ERROR(DEBUG_CATEGORY_MEDIA, "file_transfer: close failed peer=%llu: %s", + (unsigned long long)op->peer, strerror(errno)); + error = -1; + } + struct file_transfer_op** p = &ft->ops; + while (*p && *p != op) p = &(*p)->next; + if (*p) *p = op->next; + ft->count--; + DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "file_transfer: finish peer=%llu receive=%d bytes=%llu/%llu error=%d active=%u", + (unsigned long long)op->peer, op->receiving, (unsigned long long)op->offset, + (unsigned long long)op->size, error, ft->count); + file_transfer_done_fn done = op->done; + void* arg = op->arg; + u_free(op); + if (done) done(arg, error); +} + +/* Отправить подписанную кодограмму, не превышая порог очереди. */ +static int ft_send(struct file_transfer* ft, uint64_t group, uint64_t peer, + uint8_t cmd, const uint8_t* data, size_t size) { + struct ETCP_ROUTER_CONN* rc = etcp_router_conn_get(ft->inst, group, peer, ETCP_RT_ID_FILE_TRANSFER); + if (!rc) { DEBUG_ERROR(DEBUG_CATEGORY_MEDIA, "file_transfer: router allocation failed"); return -1; } + queue_set_threshold(rc->send_q, 0, 0); + etcp_router_set_max_inflight(rc, 16); + if (!etcp_router_send_q_has_room(ft->inst, group, peer, ETCP_RT_ID_FILE_TRANSFER)) return 1; + uint8_t packet[1 + 24 + FT_CHUNK]; + packet[0] = cmd; + memcpy(packet + 1, data, size); + int error = etcp_router_conn_send(rc, packet, size + 1, ROUTER_FLAG_SIGNED); + if (error) DEBUG_WARN(DEBUG_CATEGORY_MEDIA, "file_transfer: send failed peer=%llu cmd=%u rc=%d", + (unsigned long long)peer, cmd, error); + return error ? -1 : 0; +} + +/* Повторный REQUEST возобновляет приём со следующего неполученного байта. */ +static void ft_request(struct file_transfer_op* op) { + uint8_t body[48]; + memcpy(body, op->tid, 16); + memcpy(body + 16, op->object, 16); + memcpy(body + 32, &op->size, 8); + memcpy(body + 40, &op->offset, 8); + if (ft_send(op->ft, op->group, op->peer, FT_REQUEST, body, sizeof(body)) == 0) + op->request_tb = get_time_tb(); +} + +static void ft_pump(void* arg); + +/* Waiter может сработать синхронно; продолжение всегда через call_soon. */ +static void ft_ready(struct ll_queue* q, void* arg) { + (void)q; + struct file_transfer_op* op = arg; + if (op->pump) return; + op->pump = uasync_call_soon(op->ft->inst->ua, op, ft_pump); + if (!op->pump) DEBUG_ERROR(DEBUG_CATEGORY_MEDIA, "file_transfer: pump timer allocation failed"); +} + +/* Ограничить работу одного callback и количество байтов в router. */ +static void ft_pump(void* arg) { + struct file_transfer_op* op = 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++) { + 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; + } + uint8_t body[24 + FT_CHUNK]; + size_t n = op->size - op->offset < FT_CHUNK ? (size_t)(op->size - op->offset) : FT_CHUNK; + memcpy(body, op->tid, 16); + memcpy(body + 16, &op->offset, 8); + if (fseeko(op->file, (off_t)op->offset, SEEK_SET) != 0 || fread(body + 24, 1, n, op->file) != n) { + DEBUG_ERROR(DEBUG_CATEGORY_MEDIA, "file_transfer: read failed peer=%llu offset=%llu: %s", + (unsigned long long)op->peer, (unsigned long long)op->offset, strerror(errno)); + ft_finish(op, -1); + return; + } + int error = ft_send(ft, op->group, op->peer, FT_DATA, body, 24 + n); + if (error) return; + op->offset += n; + op->progress_tb = get_time_tb(); + } + if (op->offset < op->size) ft_ready(NULL, op); +} + +/* DOWN/TIMEOUT сохраняет CM handle; таймер проверяет позднее восстановление. */ +static void ft_connected(struct CONN_MGR_HANDLE* h, uint64_t peer, uint64_t group, + enum conn_mgr_event event, void* arg) { + (void)h; (void)peer; (void)group; + struct file_transfer_op* op = arg; + op->up = event == CONN_EVENT_UP; + DEBUG_DEBUG(DEBUG_CATEGORY_MEDIA, "file_transfer: CM event=%d peer=%llu receive=%d offset=%llu", + event, (unsigned long long)op->peer, op->receiving, (unsigned long long)op->offset); + if (op->up) { + if (op->receiving) ft_request(op); + else ft_ready(NULL, op); + } +} + +/* Неактивная операция завершается; durable владелец решает, когда повторять. */ +static void ft_tick(void* arg) { + struct file_transfer_op* op = arg; + op->timer = NULL; + uint64_t now = get_time_tb(); + if (now - op->progress_tb >= 300000) { + DEBUG_WARN(DEBUG_CATEGORY_MEDIA, "file_transfer: inactivity timeout peer=%llu offset=%llu/%llu", + (unsigned long long)op->peer, (unsigned long long)op->offset, (unsigned long long)op->size); + ft_finish(op, -1); + return; + } + if (op->up) { + if (op->receiving && now - op->request_tb >= 50000 && now - op->progress_tb >= 50000) ft_request(op); + else if (!op->receiving && !op->pump && !op->waiter.internal) ft_ready(NULL, op); + } + op->timer = uasync_set_timeout(op->ft->inst->ua, 10000, op, ft_tick, "file_transfer"); + if (!op->timer) { DEBUG_ERROR(DEBUG_CATEGORY_MEDIA, "file_transfer: watchdog allocation failed"); ft_finish(op, -1); } +} + +/* Добавить операцию до conn_mgr_open; callbacks CM используют уже живой owner. */ +static struct file_transfer_op* ft_open(struct file_transfer* ft, uint64_t group, uint64_t peer, + const uint8_t object[16], uint64_t size, const char* path, int receiving) { + if (!ft || ft->closing || !group || !peer || !object || !path || !size || size > INT64_MAX || ft->count >= FT_MAX_OPS) { + DEBUG_WARN(DEBUG_CATEGORY_MEDIA, "file_transfer: invalid request or operation limit size=%llu", + (unsigned long long)size); + return NULL; + } + struct file_transfer_op* op = u_calloc(1, sizeof(*op)); + if (!op) { DEBUG_ERROR(DEBUG_CATEGORY_MEDIA, "file_transfer: operation allocation failed"); return NULL; } + op->ft = ft; op->group = group; op->peer = peer; op->size = size; op->receiving = receiving; + memcpy(op->object, object, 16); + op->file = fopen(path, receiving ? "wb" : "rb"); + if (!op->file) { + DEBUG_ERROR(DEBUG_CATEGORY_MEDIA, "file_transfer: open failed path=%s: %s", path, strerror(errno)); + u_free(op); + return NULL; + } + op->next = ft->ops; ft->ops = op; ft->count++; + op->progress_tb = get_time_tb(); + op->timer = uasync_set_timeout(ft->inst->ua, 10000, op, ft_tick, "file_transfer"); + if (!op->timer || !etcp_router_conn_get(ft->inst, group, peer, ETCP_RT_ID_FILE_TRANSFER) || + conn_mgr_open(ft->inst, group, peer, ft_connected, op, &op->cm) != 0) { + DEBUG_ERROR(DEBUG_CATEGORY_MEDIA, "file_transfer: CM/timer launch failed peer=%llu", (unsigned long long)peer); + ft_finish(op, -1); + return NULL; + } + DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "file_transfer: start peer=%llu receive=%d size=%llu active=%u", + (unsigned long long)peer, receiving, (unsigned long long)size, ft->count); + return op; +} + +static struct file_transfer_op* ft_find(struct file_transfer* ft, uint64_t group, uint64_t peer, + const uint8_t tid[16], int receiving) { + for (struct file_transfer_op* op = ft->ops; op; op = op->next) + if (op->group == group && op->peer == peer && op->receiving == receiving && !memcmp(op->tid, tid, 16)) return op; + return NULL; +} + +/* REQUEST проверяется владельцем файла при каждом повторе, включая TTL/права. */ +static void ft_serve(struct file_transfer* ft, uint64_t group, uint64_t peer, const uint8_t* p) { + uint64_t size, offset; + memcpy(&size, p + 32, 8); + memcpy(&offset, p + 40, 8); + char path[1024]; + if (!size || offset >= size || ft->lookup(ft->arg, group, peer, p + 16, size, path, sizeof(path)) != 0) { + DEBUG_WARN(DEBUG_CATEGORY_MEDIA, "file_transfer: request refused peer=%llu size=%llu offset=%llu", + (unsigned long long)peer, (unsigned long long)size, (unsigned long long)offset); + ft_send(ft, group, peer, FT_REFUSE, p, 16); + return; + } + struct file_transfer_op* op = ft_find(ft, group, peer, p, 0); + if (op && (op->size != size || memcmp(op->object, p + 16, 16))) { + DEBUG_WARN(DEBUG_CATEGORY_MEDIA, "file_transfer: conflicting transfer ID peer=%llu", (unsigned long long)peer); + return; + } + if (!op) { + op = ft_open(ft, group, peer, p + 16, size, path, 0); + if (!op) { ft_send(ft, group, peer, FT_REFUSE, p, 16); return; } + memcpy(op->tid, p, 16); + } + op->offset = offset; + op->progress_tb = get_time_tb(); + if (op->up) ft_ready(NULL, op); +} + +/* Signed source, transfer ID и точные границы отделяют операции и повторы. */ +static void ft_recv(struct ETCP_CONN* conn, struct ll_entry* e) { + if (!e || !conn || !e->dgram || e->len <= ROUTER_SVC_PAYLOAD_OFF) goto done; + struct file_transfer* ft = conn->instance->file_transfer; + if (!ft || ft->closing) goto done; + uint64_t group, peer; + memcpy(&group, e->dgram + ROUTER_SVC_GROUP_OFF, 8); + memcpy(&peer, e->dgram + ROUTER_SVC_SRC_OFF, 8); + if (!(e->dgram[ROUTER_SVC_FLAGS_OFF] & ROUTER_FLAG_SIGNED)) { + DEBUG_WARN(DEBUG_CATEGORY_MEDIA, "file_transfer: unsigned packet peer=%llu", (unsigned long long)peer); + goto done; + } + const uint8_t* p = e->dgram + ROUTER_SVC_PAYLOAD_OFF; + uint8_t cmd = *p++; + size_t len = e->len - ROUTER_SVC_PAYLOAD_OFF - 1; + if (cmd == FT_REQUEST && len == 48) ft_serve(ft, group, peer, p); + else if (cmd == FT_DONE && len == 16) { + struct file_transfer_op* op = ft_find(ft, group, peer, p, 0); + if (op && op->offset == op->size) ft_finish(op, 0); + } else if (cmd == FT_REFUSE && len == 16) { + struct file_transfer_op* op = ft_find(ft, group, peer, p, 1); + if (op) { DEBUG_WARN(DEBUG_CATEGORY_MEDIA, "file_transfer: source refused peer=%llu", (unsigned long long)peer); ft_finish(op, -1); } + } else if (cmd == FT_DATA && len > 24 && len <= 24 + FT_CHUNK) { + struct file_transfer_op* op = ft_find(ft, group, peer, p, 1); + if (!op) goto done; /* Поздний пакет завершённой/отменённой операции. */ + uint64_t offset; + memcpy(&offset, p + 16, 8); + size_t n = len - 24; + if (offset < op->offset && n <= op->offset - offset) goto done; + if (offset != op->offset || n > op->size - op->offset) { + DEBUG_WARN(DEBUG_CATEGORY_MEDIA, "file_transfer: wrong range peer=%llu offset=%llu expected=%llu bytes=%zu", + (unsigned long long)peer, (unsigned long long)offset, (unsigned long long)op->offset, n); + ft_request(op); + goto done; + } + if (fwrite(p + 24, 1, n, op->file) != n) { + DEBUG_ERROR(DEBUG_CATEGORY_MEDIA, "file_transfer: write failed peer=%llu: %s", (unsigned long long)peer, strerror(errno)); + ft_finish(op, -1); + goto done; + } + op->offset += n; + op->progress_tb = get_time_tb(); + if (op->offset == op->size) { + ft_send(ft, group, peer, FT_DONE, op->tid, 16); + ft_finish(op, 0); + } + } else DEBUG_WARN(DEBUG_CATEGORY_MEDIA, "file_transfer: malformed packet cmd=%u bytes=%zu", cmd, len); +done: + if (e) { queue_dgram_free(e); queue_entry_free(e); } +} + +struct file_transfer* file_transfer_create(struct UTUN_INSTANCE* inst, file_transfer_lookup_fn lookup, void* arg) { + if (!inst || !lookup || inst->file_transfer) { + DEBUG_ERROR(DEBUG_CATEGORY_MEDIA, "file_transfer: invalid initialization"); + return NULL; + } + struct file_transfer* ft = u_calloc(1, sizeof(*ft)); + if (!ft) { DEBUG_ERROR(DEBUG_CATEGORY_MEDIA, "file_transfer: owner allocation failed"); return NULL; } + ft->inst = inst; ft->lookup = lookup; ft->arg = arg; + if (etcp_router_bind(inst, ETCP_RT_ID_FILE_TRANSFER, ft_recv) != 0) { u_free(ft); return NULL; } + inst->file_transfer = ft; + DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "file_transfer: initialized max_operations=%u", FT_MAX_OPS); + return ft; +} + +void file_transfer_destroy(struct file_transfer* ft) { + if (!ft) return; + ft->closing = 1; + etcp_router_unbind(ft->inst, ETCP_RT_ID_FILE_TRANSFER); + while (ft->ops) ft_finish(ft->ops, -2); + ft->inst->file_transfer = NULL; + DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "file_transfer: destroyed"); + u_free(ft); +} + +struct file_transfer_op* file_transfer_pull(struct file_transfer* ft, uint64_t group, uint64_t peer, + const uint8_t object[16], uint64_t size, const char* path, + file_transfer_done_fn done, void* arg) { + uint8_t tid[16]; + if (!done || RAND_bytes(tid, sizeof(tid)) != 1) { + DEBUG_ERROR(DEBUG_CATEGORY_MEDIA, "file_transfer: transfer ID generation/arguments failed"); + return NULL; + } + struct file_transfer_op* op = ft_open(ft, group, peer, object, size, path, 1); + if (!op) return NULL; + memcpy(op->tid, tid, sizeof(tid)); + op->done = done; op->arg = arg; + return op; +} + +void file_transfer_cancel(struct file_transfer_op* op) { + if (op) ft_finish(op, -2); +} diff --git a/src/media_delivery/file_transfer.h b/src/media_delivery/file_transfer.h new file mode 100644 index 00000000..89de0e2b --- /dev/null +++ b/src/media_delivery/file_transfer.h @@ -0,0 +1,32 @@ +/* file_transfer — передача неизменяемого файла между явно заданными узлами. + * Не знает о чатах, ключах содержимого, custody или БД. lookup авторизует запрос + * и возвращает путь; получатель проверяет содержимое в своём done. + * Каждая операция владеет CM handle до завершения/отмены, включая DOWN/TIMEOUT. + * Все API/callbacks в uasync; destroy вызывается вне callbacks, до удаления групп. + * Файл назначения принадлежит вызывающему; при ошибке его нужно отбросить. + */ +#ifndef FILE_TRANSFER_H +#define FILE_TRANSFER_H + +#include +#include + +struct UTUN_INSTANCE; +struct file_transfer; +struct file_transfer_op; + +typedef int (*file_transfer_lookup_fn)(void* arg, uint64_t group, uint64_t peer, + const uint8_t object[16], uint64_t size, char* path, size_t cap); +typedef void (*file_transfer_done_fn)(void* arg, int error); + +/* 0 в done означает точный приём size байт и закрытие файла, а не custody ACK. + * Возврат NULL означает отказ запуска, done не вызывается. cancel вызывает done(-2). + * Handle перестаёт существовать перед вызовом done, его нужно обнулить в callback. */ +struct file_transfer* file_transfer_create(struct UTUN_INSTANCE* inst, file_transfer_lookup_fn lookup, void* arg); +void file_transfer_destroy(struct file_transfer* ft); +struct file_transfer_op* file_transfer_pull(struct file_transfer* ft, uint64_t group, uint64_t peer, + const uint8_t object[16], uint64_t size, const char* destination, + file_transfer_done_fn done, void* arg); +void file_transfer_cancel(struct file_transfer_op* op); + +#endif