Browse Source

Add bounded signed file transfers with per-operation CM ownership

master
evgeny 2 days ago
parent
commit
f39ba840b8
  1. 331
      src/media_delivery/file_transfer.c
  2. 32
      src/media_delivery/file_transfer.h

331
src/media_delivery/file_transfer.c

@ -0,0 +1,331 @@
/* file_transfer — signed router stream с ограниченной очередью и CM-владением. */
#include <stdio.h>
#include <string.h>
#include <errno.h>
#include <openssl/rand.h>
#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);
}

32
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 <stdint.h>
#include <stddef.h>
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
Loading…
Cancel
Save