You can not select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
 
 
 
 
 
 

1412 lines
61 KiB

#include "media_delivery.h"
#include "media_delivery_proto.h"
#include "media_download.h"
#include "../utun_instance.h"
#include "../routing_layer/etcp_router.h"
#include "../routing_layer/topo_group.h"
#include "../routing_layer/topo_node.h"
#include "../routing_layer/topo_node_sqlite.h"
#include "../routing_layer/conn_mgr.h"
#include "../transport_layer/etcp_api.h"
#include "../transport_layer/etcp.h"
#include "../chat/member_sync.h"
#include "../lib/debug_config.h"
#include "../lib/mem.h"
#include "../lib/u_async.h"
#include "../lib/sqlite3.h"
#include <string.h>
#include <time.h>
#define MD_ID "media_delivery"
/* ── forward decls ── */
static void md_etcp_recv(struct ETCP_CONN* conn, struct ll_entry* entry);
static void md_on_conn_status(struct ETCP_CONN* conn, int status, void* arg);
static void md_on_bgp_node(struct TOPO_GROUP* group, uint64_t node_id, int event, void* arg);
static void md_on_props_changed(uint64_t node_id, const char* adm_tags, void* arg);
static int md_super_start(struct media_delivery_ctx* md);
static void md_super_stop(struct media_delivery_ctx* md);
static void md_super_repl_send(struct media_delivery_ctx* md, struct media_super_peer* peer);
static void md_super_connect(struct media_delivery_ctx* md, uint64_t peer_node_id);
static int md_send(struct UTUN_INSTANCE* inst, uint64_t group_id, uint64_t dst_node_id, const uint8_t* data, size_t len);
static struct media_super_peer* md_super_peer_add(struct media_delivery_ctx* md, uint64_t node_id);
/* ── helpers ── */
static const char* md_subcmd_name(uint8_t sc) {
switch (sc) {
case MEDIA_SUBCMD_SERVE_REG: return "SERVE_REG";
case MEDIA_SUBCMD_SERVE_ACK: return "SERVE_ACK";
case MEDIA_SUBCMD_SERVE_LEAVE: return "SERVE_LEAVE";
case MEDIA_SUBCMD_QUERY: return "QUERY";
case MEDIA_SUBCMD_QUERY_RESP: return "QUERY_RESP";
case MEDIA_SUBCMD_HAVE_BLOCK: return "HAVE_BLOCK";
case MEDIA_SUBCMD_HAVE_BLOCK_ACK: return "HAVE_BLOCK_ACK";
case MEDIA_SUBCMD_SUPER_REPL: return "SUPER_REPL";
case MEDIA_SUBCMD_SUPER_ACK: return "SUPER_ACK";
case MEDIA_SUBCMD_BLOCK_REQ: return "BLOCK_REQ";
case MEDIA_SUBCMD_BLOCK_CHUNK: return "BLOCK_CHUNK";
case MEDIA_SUBCMD_BLOCK_DONE: return "BLOCK_DONE";
case MEDIA_SUBCMD_CANCEL: return "CANCEL";
case MEDIA_SUBCMD_SUPER_HELLO: return "SUPER_HELLO";
case MEDIA_SUBCMD_BLOCK_PROCESSING: return "BLOCK_PROCESSING";
case MEDIA_SUBCMD_BLOCK_RELAY_FULL: return "BLOCK_RELAY_FULL";
default: return "???";
}
}
/* ── SQL helper: insert/update block_availability ── */
static int md_ba_insert(sqlite3* db, const uint8_t* block_uuid, uint64_t group_id,
uint64_t node_id, uint32_t chunk, int64_t ts, int status) {
const char* sql =
"INSERT INTO block_availability(block_uuid,group_id,node_id,chunk,timestamp,status)"
" VALUES(?,?,?,?,?,?)";
sqlite3_stmt* st = NULL;
if (sqlite3_prepare_v2(db, sql, -1, &st, NULL) != SQLITE_OK) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: ba_insert prepare failed: %s", MD_ID, sqlite3_errmsg(db));
return -1;
}
sqlite3_bind_blob(st, 1, block_uuid, 16, SQLITE_STATIC);
sqlite3_bind_int64(st, 2, (sqlite3_int64)group_id);
sqlite3_bind_int64(st, 3, (sqlite3_int64)node_id);
sqlite3_bind_int(st, 4, (int)chunk);
sqlite3_bind_int64(st, 5, (sqlite3_int64)ts);
sqlite3_bind_int(st, 6, status);
int rc = sqlite3_step(st);
sqlite3_finalize(st);
if (rc != SQLITE_DONE) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: ba_insert step failed rc=%d: %s", MD_ID, rc, sqlite3_errmsg(db));
return -1;
}
return 0;
}
/* ── complete block: DELETE processing records + INSERT completed with new id ── */
static int md_ba_complete_block(sqlite3* db, const uint8_t* block_uuid,
uint64_t group_id, uint64_t node_id,
uint32_t chunk, int64_t ts) {
/* 1. DELETE processing records (status=0) for this (block_uuid, node_id) */
const char* del_sql = "DELETE FROM block_availability WHERE block_uuid=? AND node_id=? AND status=0";
sqlite3_stmt* st = NULL;
if (sqlite3_prepare_v2(db, del_sql, -1, &st, NULL) == SQLITE_OK) {
sqlite3_bind_blob(st, 1, block_uuid, 16, SQLITE_STATIC);
sqlite3_bind_int64(st, 2, (sqlite3_int64)node_id);
sqlite3_step(st);
sqlite3_finalize(st);
}
/* 2. check if completed already exists — avoid duplicate insert */
const char* chk_sql = "SELECT 1 FROM block_availability WHERE block_uuid=? AND node_id=? AND status=1";
if (sqlite3_prepare_v2(db, chk_sql, -1, &st, NULL) == SQLITE_OK) {
sqlite3_bind_blob(st, 1, block_uuid, 16, SQLITE_STATIC);
sqlite3_bind_int64(st, 2, (sqlite3_int64)node_id);
int exists = (sqlite3_step(st) == SQLITE_ROW);
sqlite3_finalize(st);
if (exists) {
DEBUG_DEBUG(DEBUG_CATEGORY_MEDIA, "%s: ba_complete_block: completed already exists for node=0x%016llx, skipping",
MD_ID, (unsigned long long)node_id);
return 0;
}
}
/* 3. INSERT completed record (status=1, new id via AUTOINCREMENT) */
return md_ba_insert(db, block_uuid, group_id, node_id, chunk, ts, 1);
}
static int md_ba_find_nodes(sqlite3* db, const uint8_t* block_uuid,
uint64_t* out_node_ids, int max) {
const char* sql =
"SELECT node_id, timestamp FROM block_availability WHERE block_uuid=? ORDER BY timestamp DESC LIMIT ?";
sqlite3_stmt* st = NULL;
if (sqlite3_prepare_v2(db, sql, -1, &st, NULL) != SQLITE_OK) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: ba_find_nodes prepare failed: %s", MD_ID, sqlite3_errmsg(db));
return 0;
}
sqlite3_bind_blob(st, 1, block_uuid, 16, SQLITE_STATIC);
sqlite3_bind_int(st, 2, max);
int n = 0;
while (n < max && sqlite3_step(st) == SQLITE_ROW) {
out_node_ids[n] = (uint64_t)sqlite3_column_int64(st, 0);
n++;
}
sqlite3_finalize(st);
return n;
}
static int md_ba_get_since(sqlite3* db, uint64_t since_id,
uint8_t* out_buf, int max_bytes) {
const char* sql =
"SELECT id,block_uuid,group_id,node_id,chunk,timestamp,status"
" FROM block_availability WHERE id > ? ORDER BY id ASC LIMIT 4";
sqlite3_stmt* st = NULL;
if (sqlite3_prepare_v2(db, sql, -1, &st, NULL) != SQLITE_OK) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: ba_get_since prepare failed: %s", MD_ID, sqlite3_errmsg(db));
return 0;
}
sqlite3_bind_int64(st, 1, (sqlite3_int64)since_id);
int off = 0;
while (sqlite3_step(st) == SQLITE_ROW) {
if (off + 56 > max_bytes) break;
int64_t id = sqlite3_column_int64(st, 0);
memcpy(out_buf + off, &id, 8); off += 8;
memcpy(out_buf + off, sqlite3_column_blob(st, 1), 16); off += 16;
int64_t gid = sqlite3_column_int64(st, 2); memcpy(out_buf + off, &gid, 8); off += 8;
int64_t nid = sqlite3_column_int64(st, 3); memcpy(out_buf + off, &nid, 8); off += 8;
int32_t ch = sqlite3_column_int(st, 4); memcpy(out_buf + off, &ch, 4); off += 4;
int64_t ts = sqlite3_column_int64(st, 5); memcpy(out_buf + off, &ts, 8); off += 8;
int32_t sts = sqlite3_column_int(st, 6); memcpy(out_buf + off, &sts, 4); off += 4;
}
sqlite3_finalize(st);
return off;
}
static uint64_t md_ss_get(sqlite3* db, uint64_t peer_node_id) {
const char* sql = "SELECT last_recv_id FROM super_sync WHERE peer_node_id=?";
sqlite3_stmt* st = NULL;
if (sqlite3_prepare_v2(db, sql, -1, &st, NULL) != SQLITE_OK) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: ss_get prepare failed: %s", MD_ID, sqlite3_errmsg(db));
return 0;
}
sqlite3_bind_int64(st, 1, (sqlite3_int64)peer_node_id);
uint64_t v = 0;
if (sqlite3_step(st) == SQLITE_ROW) v = (uint64_t)sqlite3_column_int64(st, 0);
sqlite3_finalize(st);
return v;
}
static void md_ss_set(sqlite3* db, uint64_t peer_node_id, uint64_t last_recv_id) {
const char* sql = "INSERT OR REPLACE INTO super_sync(peer_node_id,last_recv_id) VALUES(?,?)";
sqlite3_stmt* st = NULL;
if (sqlite3_prepare_v2(db, sql, -1, &st, NULL) != SQLITE_OK) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: ss_set prepare failed: %s", MD_ID, sqlite3_errmsg(db));
return;
}
sqlite3_bind_int64(st, 1, (sqlite3_int64)peer_node_id);
sqlite3_bind_int64(st, 2, (sqlite3_int64)last_recv_id);
sqlite3_step(st);
sqlite3_finalize(st);
}
/* ── super_peer find/create ── */
static struct media_super_peer* md_super_peer_find(struct media_delivery_ctx* md, uint64_t node_id) {
if (!md->super_peers) return NULL;
struct ll_entry* e = queue_find_data_by_index(md->super_peers, &node_id);
return e ? (struct media_super_peer*)e->data : NULL;
}
static struct media_super_peer* md_super_peer_add(struct media_delivery_ctx* md, uint64_t node_id) {
struct media_super_peer* sp = md_super_peer_find(md, node_id);
if (sp) return sp;
if (!md->super_peers) {
md->super_peers = queue_new(md->inst->ua, 64, offsetof(struct media_super_peer, peer_node_id), 8, "md_super");
if (!md->super_peers) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: queue_new(super_peers) failed", MD_ID); return NULL; }
}
struct ll_entry* qe = queue_entry_new(sizeof(struct media_super_peer));
if (!qe) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: queue_entry_new(super_peer) failed", MD_ID); return NULL; }
sp = (struct media_super_peer*)qe->data;
memset(sp, 0, sizeof(*sp));
sp->peer_node_id = node_id;
sp->timeout_tb = MEDIA_REPL_TIMEOUT_TB;
sp->md = md;
memcpy(qe->data, &node_id, 8);
queue_data_put_with_index(md->super_peers, qe);
return sp;
}
static void md_super_peer_remove(struct media_delivery_ctx* md, uint64_t node_id) {
if (!md->super_peers) return;
struct ll_entry* e = queue_find_data_by_index(md->super_peers, &node_id);
if (!e) return;
struct media_super_peer* sp = (struct media_super_peer*)e->data;
if (sp->timeout_timer) { uasync_cancel_timeout(md->inst->ua, sp->timeout_timer); sp->timeout_timer = NULL; }
if (sp->connect_timer) { uasync_cancel_timeout(md->inst->ua, sp->connect_timer); sp->connect_timer = NULL; }
queue_remove_data(md->super_peers, e);
queue_entry_free(e);
}
/* ── relay block context helpers ── */
struct relay_block_ctx* md_relay_find(struct media_delivery_ctx* md, const uint8_t* block_id) {
if (!md->relay_blocks) return NULL;
struct ll_entry* e = queue_find_data_by_index(md->relay_blocks, block_id);
return e ? (struct relay_block_ctx*)e->data : NULL;
}
struct relay_block_ctx* md_relay_add(struct media_delivery_ctx* md,
const uint8_t* block_id,
const uint8_t* media_id,
uint32_t chunk) {
struct relay_block_ctx* rc = md_relay_find(md, block_id);
if (rc) return rc;
if (!md->relay_blocks) {
md->relay_blocks = queue_new(md->inst->ua, 64, offsetof(struct relay_block_ctx, block_id), 16, "md_relay");
if (!md->relay_blocks) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: queue_new(relay_blocks) failed", MD_ID); return NULL; }
}
struct ll_entry* qe = queue_entry_new(sizeof(struct relay_block_ctx));
if (!qe) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: queue_entry_new(relay_block) failed", MD_ID); return NULL; }
rc = (struct relay_block_ctx*)qe->data;
memset(rc, 0, sizeof(*rc));
memcpy(rc->block_id, block_id, 16);
memcpy(rc->media_id, media_id, 16);
rc->chunk = chunk;
rc->inst = md->inst;
memcpy(qe->data, block_id, 16);
queue_data_put_with_index(md->relay_blocks, qe);
return rc;
}
void md_relay_remove(struct media_delivery_ctx* md, const uint8_t* block_id) {
if (!md->relay_blocks) return;
struct ll_entry* e = queue_find_data_by_index(md->relay_blocks, block_id);
if (!e) return;
queue_remove_data(md->relay_blocks, e);
queue_entry_free(e);
}
/* ── file_load helpers (per-file source download tracking) ── */
struct md_file_load* md_file_load_find(struct media_delivery_ctx* md, const uint8_t* media_id) {
if (!md->file_loads) return NULL;
struct ll_entry* e = queue_find_data_by_index(md->file_loads, media_id);
return e ? (struct md_file_load*)e->data : NULL;
}
int md_file_load_inc(struct media_delivery_ctx* md, const uint8_t* media_id, uint64_t node_id) {
if (!md->file_loads) {
md->file_loads = queue_new(md->inst->ua, 64, offsetof(struct md_file_load, media_id), 16, "md_fld");
if (!md->file_loads) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: queue_new(file_loads) failed", MD_ID); return -1; }
}
struct md_file_load* fl = md_file_load_find(md, media_id);
if (!fl) {
struct ll_entry* qe = queue_entry_new(sizeof(struct md_file_load));
if (!qe) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: queue_entry_new(file_load) failed", MD_ID); return -1; }
fl = (struct md_file_load*)qe->data;
memset(fl, 0, sizeof(*fl));
memcpy(fl->media_id, media_id, 16);
memcpy(qe->data, media_id, 16);
queue_data_put_with_index(md->file_loads, qe);
}
if (fl->active_downloads < 10)
fl->downloader_ids[fl->active_downloads] = node_id;
fl->active_downloads++;
DEBUG_DEBUG(DEBUG_CATEGORY_MEDIA, "%s: file_load inc media=%02x%02x... node=0x%016llx count=%d",
MD_ID, media_id[0], media_id[1], (unsigned long long)node_id, fl->active_downloads);
return fl->active_downloads;
}
void md_file_load_dec(struct media_delivery_ctx* md, const uint8_t* media_id, uint64_t node_id) {
struct md_file_load* fl = md_file_load_find(md, media_id);
if (!fl) return;
if (fl->active_downloads > 0) fl->active_downloads--;
/* remove node_id from downloader_ids array */
for (int i = 0; i < fl->active_downloads && i < 10; i++) {
if (fl->downloader_ids[i] == node_id) {
for (int j = i; j < fl->active_downloads && j < 9; j++)
fl->downloader_ids[j] = fl->downloader_ids[j + 1];
fl->downloader_ids[fl->active_downloads < 10 ? fl->active_downloads : 9] = 0;
break;
}
}
DEBUG_DEBUG(DEBUG_CATEGORY_MEDIA, "%s: file_load dec media=%02x%02x... node=0x%016llx count=%d",
MD_ID, media_id[0], media_id[1], (unsigned long long)node_id, fl->active_downloads);
/* remove entry when count reaches 0 */
if (fl->active_downloads == 0 && md->file_loads) {
struct ll_entry* e = queue_find_data_by_index(md->file_loads, media_id);
if (e) { queue_remove_data(md->file_loads, e); queue_entry_free(e); }
}
}
/* ── replication timers ── */
static void md_repl_retrans_cb(void* arg) {
struct media_super_peer* peer = (struct media_super_peer*)arg;
struct media_delivery_ctx* md = peer->md;
if (!md || !peer->connected) return;
peer->timeout_timer = NULL;
peer->timeout_tb = peer->timeout_tb < MEDIA_REPL_TIMEOUT_MAX_TB / 2
? peer->timeout_tb * 2 : peer->timeout_tb;
if (peer->timeout_tb > MEDIA_REPL_TIMEOUT_MAX_TB)
peer->timeout_tb = MEDIA_REPL_TIMEOUT_MAX_TB;
md_super_repl_send(md, peer);
}
static void md_repl_schedule_retrans(struct media_delivery_ctx* md, struct media_super_peer* peer) {
if (peer->timeout_timer) uasync_cancel_timeout(md->inst->ua, peer->timeout_timer);
peer->timeout_timer = uasync_set_timeout(md->inst->ua, peer->timeout_tb, peer,
md_repl_retrans_cb, "md_repl_retrans");
}
static void md_reconnect_cb(void* arg) {
struct media_delivery_ctx* md = (struct media_delivery_ctx*)arg;
struct media_super_peer* peer = (struct media_super_peer*)md->super_peers;
while (peer) {
if (!peer->connected && peer->connect_timer) {
peer->connect_timer = NULL;
md_super_connect(md, peer->peer_node_id);
}
peer = peer->connected ? NULL : peer;
}
}
/* ── replication send ── */
static void md_super_repl_send(struct media_delivery_ctx* md, struct media_super_peer* peer) {
if (!md->db || !peer->connected || !peer->hello_done) return;
while (peer->inflight_count < MEDIA_MAX_REPL_INFLIGHT) {
uint8_t buf[560];
int nb = md_ba_get_since(md->db, peer->peer_last_recv_id, buf, (int)sizeof(buf));
if (nb == 0) {
if (peer->timeout_timer) uasync_cancel_timeout(md->inst->ua, peer->timeout_timer);
peer->timeout_timer = NULL;
return;
}
/* build packet: subcmd + seq + num + raw rows */
uint8_t pkt[600];
int off = 0;
pkt[off++] = MEDIA_SUBCMD_SUPER_REPL;
int num_entries = nb / 56;
int64_t max_id = 0;
const uint8_t* rp = buf;
for (int i = 0; i < num_entries; i++) {
int64_t id; memcpy(&id, rp, 8);
if (id > max_id) max_id = id;
rp += 56;
}
uint32_t seq = (uint32_t)max_id; memcpy(pkt + off, &seq, 4); off += 4;
uint16_t num = (uint16_t)num_entries; memcpy(pkt + off, &num, 2); off += 2;
memcpy(pkt + off, buf, (size_t)nb); off += nb;
if (md_send(md->inst, TOPO_GROUP_UTUN, peer->peer_node_id, pkt, (size_t)off) == 0) {
peer->inflight_count++;
DEBUG_DEBUG(DEBUG_CATEGORY_MEDIA, "%s: SUPER_REPL sent to 0x%016llx seq=%u entries=%d",
MD_ID, (unsigned long long)peer->peer_node_id, seq, num_entries);
} else {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: SUPER_REPL send failed to 0x%016llx",
MD_ID, (unsigned long long)peer->peer_node_id);
return;
}
if (peer->inflight_count >= MEDIA_MAX_REPL_INFLIGHT) break;
}
if (peer->inflight_count > 0) md_repl_schedule_retrans(md, peer);
}
/* ── send via etcp_router helper ── */
static int md_send(struct UTUN_INSTANCE* inst, uint64_t group_id, uint64_t dst_node_id,
const uint8_t* data, size_t len) {
if (!inst || !data || len == 0) return -1;
struct ll_entry* e = queue_entry_new(0);
if (!e) return -1;
e->dgram = u_malloc(len + 1);
if (!e->dgram) { queue_entry_free(e); return -1; }
e->dgram[0] = ETCP_RT_ID_MEDIA_DELIVERY;
memcpy(e->dgram + 1, data, len); e->len = (uint16_t)(len + 1);
int rc = etcp_route_send(inst, group_id, dst_node_id, e, 1);
if (rc != 0) { u_free(e->dgram); queue_entry_free(e); }
return rc;
}
/* ── subcommand handlers ── */
static void md_handle_serve_reg(struct media_delivery_ctx* md, uint64_t from_node,
const uint8_t* data, size_t len) {
if (len < sizeof(struct media_pkt_serve_reg)) return;
struct media_pkt_serve_reg* pkt = (struct media_pkt_serve_reg*)data;
if (!md->served_nodes) {
md->served_nodes = queue_new(md->inst->ua, 64, offsetof(struct media_served_node, node_id), 8, "md_served");
if (!md->served_nodes) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: queue_new(served_nodes) failed", MD_ID); return; }
}
struct ll_entry* e = queue_find_data_by_index(md->served_nodes, &from_node);
if (!e) {
struct ll_entry* qe = queue_entry_new(sizeof(struct media_served_node));
if (!qe) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: queue_entry_new(served_node) failed", MD_ID); return; }
struct media_served_node* sn = (struct media_served_node*)qe->data;
memset(sn, 0, sizeof(*sn));
sn->node_id = from_node; sn->group_id = pkt->group_id;
sn->joined_at = (int64_t)time(NULL);
memcpy(qe->data, &from_node, 8);
queue_data_put_with_index(md->served_nodes, qe);
}
uint8_t ack[9] = { MEDIA_SUBCMD_SERVE_ACK, 0, 0 };
memcpy(ack + 1, &pkt->group_id, 8);
(void)from_node;
DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%s: SERVE_REG node=0x%016llx group=%016llx",
MD_ID, (unsigned long long)from_node, (unsigned long long)pkt->group_id);
}
static void md_handle_query(struct media_delivery_ctx* md, uint64_t from_node,
const uint8_t* data, size_t len) {
if (len < MEDIA_QUERY_HDR_SIZE) return;
if (!md->db) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: QUERY but no DB", MD_ID); return; }
struct media_pkt_query* q = (struct media_pkt_query*)data;
uint16_t nb = q->num_blocks;
if (nb > 64) nb = 64;
uint8_t* bids = (uint8_t*)data + MEDIA_QUERY_HDR_SIZE;
uint8_t resp[2048]; int off = 0;
resp[off++] = MEDIA_SUBCMD_QUERY_RESP;
int num_pos = off; off += 2;
int count = 0;
for (int i = 0; i < nb && off + 30 <= (int)sizeof(resp); i++) {
uint8_t* block_id = bids + i * 16;
uint64_t nodes[10]; int nf = 0;
/* 1) block_availability (supernode or HAVE_BLOCK reports) */
nf = md_ba_find_nodes(md->db, block_id, nodes, 10);
/* 2) media_files: check if we own this block */
if (nf == 0) {
const char* sql = "SELECT 1 FROM media_files WHERE block_id=? AND node_id=? LIMIT 1";
sqlite3_stmt* st = NULL;
if (sqlite3_prepare_v2(md->db, sql, -1, &st, NULL) == SQLITE_OK) {
sqlite3_bind_blob(st, 1, block_id, 16, SQLITE_STATIC);
sqlite3_bind_int64(st, 2, (sqlite3_int64)md->self_node_id);
if (sqlite3_step(st) == SQLITE_ROW) { nodes[nf++] = md->self_node_id; }
sqlite3_finalize(st);
}
}
for (int j = 0; j < nf && off + 30 <= (int)sizeof(resp); j++) {
uint16_t rtt = 0;
memcpy(resp + off, &nodes[j], 8); off += 8;
memcpy(resp + off, &rtt, 2); off += 2;
memcpy(resp + off, block_id, 16); off += 16;
uint32_t ch = q->num_blocks;
memcpy(resp + off, &ch, 4); off += 4;
count++;
}
}
uint16_t nc = (uint16_t)count;
memcpy(resp + num_pos, &nc, 2);
DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%s: QUERY from 0x%016llx → %d entries (super=%d)",
MD_ID, (unsigned long long)from_node, count, md->is_supernode);
md_send(md->inst, TOPO_GROUP_UTUN, from_node, resp, (size_t)off);
}
static void md_handle_have_block(struct media_delivery_ctx* md, uint64_t from_node,
const uint8_t* data, size_t len) {
if (len < MEDIA_HAVE_BLOCK_SIZE) return;
if (!md->db) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: HAVE_BLOCK but no DB", MD_ID); return; }
struct media_pkt_have_block* hb = (struct media_pkt_have_block*)data;
int rc = md_ba_complete_block(md->db, hb->block_id, hb->group_id, from_node, hb->chunk, hb->timestamp);
uint8_t ack[sizeof(struct media_pkt_have_block_ack)];
struct media_pkt_have_block_ack* ha = (struct media_pkt_have_block_ack*)ack;
memset(ack, 0, sizeof(ack));
ha->subcmd = MEDIA_SUBCMD_HAVE_BLOCK_ACK;
memcpy(ha->media_id, hb->media_id, 16);
memcpy(ha->block_id, hb->block_id, 16);
ha->status = (rc == 0) ? 0 : 1;
md_send(md->inst, TOPO_GROUP_UTUN, from_node, ack, sizeof(ack));
DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%s: HAVE_BLOCK from 0x%016llx rc=%d",
MD_ID, (unsigned long long)from_node, rc);
if (rc == 0) {
/* replicate to connected super-peers */
struct ll_entry* se = md->super_peers ? md->super_peers->head : NULL;
while (se) {
struct media_super_peer* sp = (struct media_super_peer*)se->data;
if (sp->connected && sp->hello_done && sp->inflight_count < MEDIA_MAX_REPL_INFLIGHT)
md_super_repl_send(md, sp);
se = se->next;
}
}
}
/* ── handle BLOCK_PROCESSING: node started downloading a block, insert as processing (status=0) ── */
static void md_handle_block_processing(struct media_delivery_ctx* md, uint64_t from_node,
const uint8_t* data, size_t len) {
if (len < MEDIA_HAVE_BLOCK_SIZE) return;
if (!md->db) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: BLOCK_PROCESSING but no DB", MD_ID); return; }
struct media_pkt_have_block* hb = (struct media_pkt_have_block*)data;
/* INSERT processing (status=0) with new id */
int rc = md_ba_insert(md->db, hb->block_id, hb->group_id, from_node, hb->chunk, hb->timestamp, 0);
if (rc != 0) {
DEBUG_WARN(DEBUG_CATEGORY_MEDIA, "%s: BLOCK_PROCESSING insert failed for node=0x%016llx (duplicate/constraint?)",
MD_ID, (unsigned long long)from_node);
}
DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%s: BLOCK_PROCESSING from 0x%016llx block=%02x%02x...",
MD_ID, (unsigned long long)from_node, hb->block_id[0], hb->block_id[1]);
/* replicate to connected super-peers */
struct ll_entry* se = md->super_peers ? md->super_peers->head : NULL;
while (se) {
struct media_super_peer* sp = (struct media_super_peer*)se->data;
if (sp->connected && sp->hello_done && sp->inflight_count < MEDIA_MAX_REPL_INFLIGHT)
md_super_repl_send(md, sp);
se = se->next;
}
}
static void md_handle_super_repl(struct media_delivery_ctx* md, uint64_t from_node,
const uint8_t* data, size_t len) {
if (len < MEDIA_SUPER_REPL_HDR_SIZE) return;
if (!md->db) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: SUPER_REPL but no DB", MD_ID); return; }
struct media_pkt_super_repl* rp = (struct media_pkt_super_repl*)data;
uint16_t num = rp->num_entries;
int64_t max_id = 0;
const uint8_t* ep = data + MEDIA_SUPER_REPL_HDR_SIZE;
for (int i = 0; i < (int)num; i++) {
int64_t id; memcpy(&id, ep, 8); ep += 8;
if (id > max_id) max_id = id;
uint8_t block_uuid[16]; memcpy(block_uuid, ep, 16); ep += 16;
int64_t gid; memcpy(&gid, ep, 8); ep += 8;
int64_t nid; memcpy(&nid, ep, 8); ep += 8;
int32_t ch; memcpy(&ch, ep, 4); ep += 4;
int64_t ts; memcpy(&ts, ep, 8); ep += 8;
int32_t st; memcpy(&st, ep, 4); ep += 4;
if (st == 1) {
md_ba_complete_block(md->db, block_uuid, (uint64_t)gid, (uint64_t)nid, (uint32_t)ch, ts);
} else {
/* processing: INSERT OR IGNORE (дубликаты при репликации — норма) */
const char* ign = "INSERT OR IGNORE INTO block_availability(block_uuid,group_id,node_id,chunk,timestamp,status) VALUES(?,?,?,?,?,0)";
sqlite3_stmt* igs = NULL;
if (sqlite3_prepare_v2(md->db, ign, -1, &igs, NULL) == SQLITE_OK) {
sqlite3_bind_blob(igs, 1, block_uuid, 16, SQLITE_STATIC);
sqlite3_bind_int64(igs, 2, (sqlite3_int64)gid);
sqlite3_bind_int64(igs, 3, (sqlite3_int64)nid);
sqlite3_bind_int(igs, 4, (int)ch);
sqlite3_bind_int64(igs, 5, (sqlite3_int64)ts);
sqlite3_step(igs);
sqlite3_finalize(igs);
}
}
}
md_ss_set(md->db, from_node, (uint64_t)max_id);
uint8_t ack[sizeof(struct media_pkt_super_ack)];
struct media_pkt_super_ack* sa = (struct media_pkt_super_ack*)ack;
sa->subcmd = MEDIA_SUBCMD_SUPER_ACK;
sa->ack_seq = (uint32_t)max_id;
md_send(md->inst, TOPO_GROUP_UTUN, from_node, ack, sizeof(ack));
DEBUG_DEBUG(DEBUG_CATEGORY_MEDIA, "%s: SUPER_REPL from 0x%016llx entries=%d max_id=%lld",
MD_ID, (unsigned long long)from_node, num, (long long)max_id);
}
static void md_handle_super_ack(struct media_delivery_ctx* md, uint64_t from_node,
const uint8_t* data, size_t len) {
if (len < sizeof(struct media_pkt_super_ack)) return;
struct media_pkt_super_ack* sa = (struct media_pkt_super_ack*)data;
struct media_super_peer* peer = md_super_peer_find(md, from_node);
if (!peer) { DEBUG_WARN(DEBUG_CATEGORY_MEDIA, "%s: SUPER_ACK from unknown peer 0x%016llx", MD_ID, (unsigned long long)from_node); return; }
peer->peer_last_recv_id = sa->ack_seq;
peer->inflight_count--;
peer->timeout_tb = MEDIA_REPL_TIMEOUT_TB;
DEBUG_DEBUG(DEBUG_CATEGORY_MEDIA, "%s: SUPER_ACK from 0x%016llx seq=%u inflight=%d",
MD_ID, (unsigned long long)from_node, sa->ack_seq, peer->inflight_count);
md_super_repl_send(md, peer);
}
static void md_handle_super_hello(struct media_delivery_ctx* md, uint64_t from_node,
const uint8_t* data, size_t len) {
if (len < sizeof(struct media_pkt_super_hello)) return;
if (!md->db) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: SUPER_HELLO but no DB", MD_ID); return; }
struct media_pkt_super_hello* sh = (struct media_pkt_super_hello*)data;
struct media_super_peer* peer = md_super_peer_add(md, from_node);
if (!peer) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: super_peer_add failed for 0x%016llx", MD_ID, (unsigned long long)from_node); return; }
peer->peer_last_recv_id = sh->last_recv_id;
peer->connected = 1;
peer->hello_done = 1;
/* send our HELLO back */
uint64_t my_last_recv = md_ss_get(md->db, from_node);
struct media_pkt_super_hello resp;
resp.subcmd = MEDIA_SUBCMD_SUPER_HELLO;
resp.last_recv_id = my_last_recv;
md_send(md->inst, TOPO_GROUP_UTUN, from_node, (const uint8_t*)&resp, sizeof(resp));
DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%s: SUPER_HELLO from 0x%016llx (peer_last_recv=%lld, my_last_recv=%lld)",
MD_ID, (unsigned long long)from_node, (long long)sh->last_recv_id, (long long)my_last_recv);
/* start replication */
peer->inflight_count = 0;
md_super_repl_send(md, peer);
}
/* ── supernode connection ── */
static void md_super_conn_cb(struct CONN_MGR_HANDLE* h, uint64_t node_id, uint64_t group_id,
enum conn_mgr_event event, void* arg) {
struct media_delivery_ctx* md = (struct media_delivery_ctx*)arg; (void)h; (void)group_id;
struct media_super_peer* peer = md_super_peer_find(md, node_id);
if (!peer) { DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: conn_cb for unknown peer 0x%016llx", MD_ID, (unsigned long long)node_id); return; }
if (event != CONN_EVENT_UP) {
DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: connect to supernode 0x%016llx failed rc=%d",
MD_ID, (unsigned long long)node_id, event);
peer->connect_timer = uasync_set_timeout(md->inst->ua,
MEDIA_RECONNECT_COOLDOWN_TB, md, md_reconnect_cb, "md_reconnect");
return;
}
peer->connected = 1;
/* send SUPER_HELLO */
uint64_t my_last = md_ss_get(md->db, node_id);
struct media_pkt_super_hello hello;
hello.subcmd = MEDIA_SUBCMD_SUPER_HELLO;
hello.last_recv_id = my_last;
if (md_send(md->inst, TOPO_GROUP_UTUN, node_id, (const uint8_t*)&hello, sizeof(hello)) == 0) {
DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%s: SUPER_HELLO sent to 0x%016llx last_recv=%lld",
MD_ID, (unsigned long long)node_id, (long long)my_last);
} else {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: SUPER_HELLO send failed to 0x%016llx",
MD_ID, (unsigned long long)node_id);
}
/* wait for HELLO response — md_handle_super_hello will start replication */
}
/* ── handle incoming BLOCK_REQ (admission control + streaming) ── */
#define MD_MAX_STREAMS 3
#define MD_CHUNK_SIZE 1024
struct stream_ctx {
struct media_delivery_ctx* md;
uint64_t group_id;
uint64_t dst_node_id;
uint8_t media_id[16];
uint8_t block_id[16];
uint32_t chunk;
FILE* file;
uint64_t offset;
uint64_t remaining;
uint64_t block_start;
size_t block_data_len;
struct queue_waiter_handle waiter;
};
static void stream_send_chunk_cb(struct ll_queue* q, void* arg) {
(void)q;
struct stream_ctx* sc = (struct stream_ctx*)arg;
struct media_delivery_ctx* md = sc->md;
if (!md || sc->remaining == 0) { u_free(sc); return; }
/* read up to MD_CHUNK_SIZE bytes */
size_t to_read = sc->remaining < MD_CHUNK_SIZE ? (size_t)sc->remaining : MD_CHUNK_SIZE;
uint8_t buf[MD_CHUNK_SIZE];
fseeko(sc->file, (off_t)sc->offset, SEEK_SET);
size_t rd = fread(buf, 1, to_read, sc->file);
if (rd == 0) { /* EOF or error */
u_free(sc); return;
}
/* build CHUNK packet */
uint8_t pkt[MEDIA_BLOCK_CHUNK_HDR_SIZE + MD_CHUNK_SIZE];
struct media_pkt_block_chunk* ch = (struct media_pkt_block_chunk*)pkt;
memset(ch, 0, MEDIA_BLOCK_CHUNK_HDR_SIZE);
ch->subcmd = MEDIA_SUBCMD_BLOCK_CHUNK;
memcpy(ch->media_id, sc->media_id, 16);
memcpy(ch->block_id, sc->block_id, 16);
ch->chunk = sc->chunk;
ch->offset = (uint32_t)(sc->offset - sc->block_start);
ch->data_len = (uint16_t)rd;
memcpy(pkt + MEDIA_BLOCK_CHUNK_HDR_SIZE, buf, rd);
md_send(md->inst, sc->group_id, sc->dst_node_id, pkt, MEDIA_BLOCK_CHUNK_HDR_SIZE + rd);
sc->offset += rd;
sc->remaining -= rd;
if (sc->remaining == 0) {
/* all data sent — send BLOCK_DONE */
uint8_t done_pkt[MEDIA_BLOCK_DONE_SIZE];
struct media_pkt_block_done* bd = (struct media_pkt_block_done*)done_pkt;
memset(bd, 0, sizeof(*bd));
bd->subcmd = MEDIA_SUBCMD_BLOCK_DONE;
memcpy(bd->media_id, sc->media_id, 16);
memcpy(bd->block_id, sc->block_id, 16);
bd->chunk = sc->chunk;
bd->total_size = (uint32_t)sc->block_data_len;
/* block_sig is zero — receiver will verify vs block_sigs from media_index_result */
md_send(md->inst, sc->group_id, sc->dst_node_id, done_pkt, MEDIA_BLOCK_DONE_SIZE);
md_file_load_dec(md, sc->media_id, sc->dst_node_id);
media_delivery_stream_done(md->inst);
fclose(sc->file);
u_free(sc);
DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%s: BLOCK_DONE sent chunk=%d total=%u",
MD_ID, sc->chunk, bd->total_size);
} else {
/* continue streaming — register waiter for backpressure */
etcp_router_on_send_ready(md->inst, sc->group_id,
sc->dst_node_id, ETCP_RT_ID_MEDIA_DELIVERY,
&sc->waiter, stream_send_chunk_cb, sc);
}
}
/* ── relay chunk forwarding ── */
struct relay_fwd_ctx {
struct media_delivery_ctx* md;
uint64_t dst_node_id;
uint64_t group_id;
struct relay_block_ctx* rc;
int ds_idx; // индекс в rc->downstream[]
struct queue_waiter_handle waiter;
};
static void md_relay_fwd_cb(struct ll_queue* q, void* arg) {
(void)q;
struct relay_fwd_ctx* fc = (struct relay_fwd_ctx*)arg;
struct relay_block_ctx* rc = fc->rc;
struct relay_downstream* ds = &rc->downstream[fc->ds_idx];
if (ds->sent_offset >= rc->file_offset) { u_free(fc); return; }
FILE* f = fopen(rc->chunk_file, "rb");
if (!f) { u_free(fc); return; }
fseeko(f, (off_t)ds->sent_offset, SEEK_SET);
size_t to_read = rc->file_offset - ds->sent_offset;
if (to_read > 1024) to_read = 1024;
uint8_t buf[1024];
size_t rd = fread(buf, 1, to_read, f);
fclose(f);
if (rd == 0) { u_free(fc); return; }
uint8_t pkt[MEDIA_BLOCK_CHUNK_HDR_SIZE + 1024];
struct media_pkt_block_chunk* ch = (struct media_pkt_block_chunk*)pkt;
memset(ch, 0, MEDIA_BLOCK_CHUNK_HDR_SIZE);
ch->subcmd = MEDIA_SUBCMD_BLOCK_CHUNK;
memcpy(ch->media_id, rc->media_id, 16);
memcpy(ch->block_id, rc->block_id, 16);
ch->chunk = rc->chunk;
ch->offset = (uint32_t)ds->sent_offset;
ch->data_len = (uint16_t)rd;
memcpy(pkt + MEDIA_BLOCK_CHUNK_HDR_SIZE, buf, rd);
size_t pkt_len = MEDIA_BLOCK_CHUNK_HDR_SIZE + rd;
if (md_send(fc->md->inst, fc->group_id, fc->dst_node_id, pkt, pkt_len) == 0) {
ds->sent_offset += rd;
}
if (ds->sent_offset < rc->file_offset) {
etcp_router_on_send_ready(fc->md->inst, fc->group_id,
fc->dst_node_id, ETCP_RT_ID_MEDIA_DELIVERY,
&fc->waiter, md_relay_fwd_cb, fc);
} else {
u_free(fc);
}
}
static void md_relay_catchup(struct media_delivery_ctx* md, struct relay_block_ctx* rc,
uint64_t dst_node_id, uint64_t group_id, int ds_idx) {
struct relay_downstream* ds = &rc->downstream[ds_idx];
if (rc->file_offset == 0 || ds->sent_offset >= rc->file_offset) return;
struct relay_fwd_ctx* fc = u_calloc(1, sizeof(*fc));
if (!fc) return;
fc->md = md; fc->dst_node_id = dst_node_id; fc->group_id = group_id;
fc->rc = rc; fc->ds_idx = ds_idx;
memset(&fc->waiter, 0, sizeof(fc->waiter));
md_relay_fwd_cb(NULL, fc);
}
static void md_handle_block_req(struct media_delivery_ctx* md, uint64_t from_node,
const uint8_t* data, size_t len) {
if (len < MEDIA_BLOCK_REQ_SIZE) return;
struct media_pkt_block_req* req = (struct media_pkt_block_req*)data;
uint64_t group_id = req->group_id ? req->group_id : TOPO_GROUP_UTUN;
DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%s: BLOCK_REQ from 0x%016llx group=%016llx block=%02x%02x... active_streams=%d",
MD_ID, (unsigned long long)from_node, (unsigned long long)group_id, req->block_id[0], req->block_id[1], md->active_streams);
int limit_reached = md->active_streams >= MD_MAX_STREAMS;
md->active_streams++;
if (limit_reached) {
md->active_streams--; /* undo */
struct media_pkt_block_overloaded ov;
ov.subcmd = MEDIA_SUBCMD_BLOCK_OVERLOADED;
memcpy(ov.media_id, req->media_id, 16);
memcpy(ov.block_id, req->block_id, 16);
ov.chunk = req->chunk;
ov.retry_after_ms = 2000;
md_send(md->inst, group_id, from_node, (const uint8_t*)&ov, sizeof(ov));
DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: BLOCK_REQ OVERLOADED (streams=%d) to 0x%016llx",
MD_ID, md->active_streams, (unsigned long long)from_node);
return;
}
/* establish direct connection if not already connected */
if (!instance_find_conn(md->inst, from_node) && group_id != TOPO_GROUP_UTUN) {
struct TOPO_GROUP* grp = topo_groups_find(md->inst->topo_groups, group_id);
if (!grp) grp = topo_groups_get_default(md->inst->topo_groups);
if (grp && grp->conn_mgr) {
DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%s: no direct conn to 0x%016llx — connecting via conn_mgr",
MD_ID, (unsigned long long)from_node);
conn_mgr_open_invite(grp->instance, grp->group_id, NULL, from_node, NULL, NULL, NULL);
}
}
/* find file in media_files DB */
sqlite3* db = md->db;
if (!db) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: BLOCK_REQ no DB", MD_ID); md->active_streams--; return; }
const char* sql = "SELECT location,chunk_size,chunk,offset,file_size FROM media_files WHERE block_id=? AND node_id=? LIMIT 1";
sqlite3_stmt* st = NULL;
if (sqlite3_prepare_v2(db, sql, -1, &st, NULL) != SQLITE_OK) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: BLOCK_REQ prepare failed: %s", MD_ID, sqlite3_errmsg(db)); md->active_streams--; return; }
sqlite3_bind_blob(st, 1, req->block_id, 16, SQLITE_STATIC);
sqlite3_bind_int64(st, 2, (sqlite3_int64)md->inst->node_id);
if (sqlite3_step(st) != SQLITE_ROW) {
sqlite3_finalize(st);
/* fallback: check relay_blocks — maybe we are downloading this block and can relay */
struct relay_block_ctx* rc = md_relay_find(md, req->block_id);
if (rc && rc->downstream_count < MD_MAX_RELAY_DOWNSTREAM) {
int dsi = rc->downstream_count++;
rc->downstream[dsi].node_id = from_node;
rc->downstream[dsi].group_id = group_id;
rc->downstream[dsi].sent_offset = 0;
memset(&rc->downstream[dsi].waiter, 0, sizeof(rc->downstream[dsi].waiter));
DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%s: relay downstream[%d] added node=0x%016llx file_offset=%llu",
MD_ID, dsi, (unsigned long long)from_node, (unsigned long long)rc->file_offset);
md_relay_catchup(md, rc, from_node, group_id, dsi);
md->active_streams--; return;
}
if (rc) {
/* relay full — send RELAY_FULL with downstream list */
uint8_t rfp[256]; int roff = 0;
rfp[roff++] = MEDIA_SUBCMD_BLOCK_RELAY_FULL;
memcpy(rfp + roff, req->media_id, 16); roff += 16;
memcpy(rfp + roff, req->block_id, 16); roff += 16;
uint32_t ck = req->chunk; memcpy(rfp + roff, &ck, 4); roff += 4;
uint16_t nrn = (uint16_t)rc->downstream_count;
memcpy(rfp + roff, &nrn, 2); roff += 2;
for (int i = 0; i < rc->downstream_count; i++) {
memcpy(rfp + roff, &rc->downstream[i].node_id, 8); roff += 8;
}
md_send(md->inst, group_id, from_node, rfp, (size_t)roff);
DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%s: relay full for block=%02x%02x..., sent %d downstream nodes to 0x%016llx",
MD_ID, req->block_id[0], req->block_id[1], rc->downstream_count, (unsigned long long)from_node);
} else {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: BLOCK_REQ block_id=%02x%02x... not found in media_files",
MD_ID, req->block_id[0], req->block_id[1]);
}
md->active_streams--; return;
}
const char* loc = (const char*)sqlite3_column_text(st, 0);
int64_t chunk_size = sqlite3_column_int64(st, 1);
int chunk = sqlite3_column_int(st, 2);
int64_t offset = sqlite3_column_int64(st, 3);
int64_t file_size = sqlite3_column_int64(st, 4);
/* compute chunk start and size */
int64_t block_start = (int64_t)chunk * chunk_size;
int64_t block_end = block_start + chunk_size;
if (block_end > file_size) block_end = file_size;
int64_t block_length = block_end - block_start;
/* find media_base */
char media_base[512] = "/tmp/utun_media";
{
sqlite3_stmt* ms = NULL;
sqlite3_prepare_v2(db, "SELECT value FROM ui_state WHERE key='media_base'", -1, &ms, NULL);
if (ms && sqlite3_step(ms) == SQLITE_ROW) {
const char* mb = (const char*)sqlite3_column_text(ms, 0);
if (mb) snprintf(media_base, sizeof(media_base), "%s", mb);
}
if (ms) sqlite3_finalize(ms);
}
(void)media_base;
char path[2048];
snprintf(path, sizeof(path), "%s/%s", media_base, loc);
FILE* f = fopen(path, "rb");
sqlite3_finalize(st);
if (!f) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: cannot open file %s for streaming", MD_ID, path);
md->active_streams--; /* undo pre-increment */
return;
}
/* per-file download limit: redirect to current downloaders if at capacity */
if (md->max_downloads_per_file > 0) {
struct md_file_load* fl = md_file_load_find(md, req->media_id);
if (fl && fl->active_downloads >= (int)md->max_downloads_per_file) {
/* source is at capacity — send RELAY_FULL with current downloader list */
uint8_t rfp[256]; int roff = 0;
rfp[roff++] = MEDIA_SUBCMD_BLOCK_RELAY_FULL;
memcpy(rfp + roff, req->media_id, 16); roff += 16;
memcpy(rfp + roff, req->block_id, 16); roff += 16;
uint32_t ck = req->chunk; memcpy(rfp + roff, &ck, 4); roff += 4;
uint16_t nrn = (uint16_t)fl->active_downloads;
memcpy(rfp + roff, &nrn, 2); roff += 2;
for (int i = 0; i < fl->active_downloads && i < 10; i++) {
memcpy(rfp + roff, &fl->downloader_ids[i], 8); roff += 8;
}
md_send(md->inst, group_id, from_node, rfp, (size_t)roff);
DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%s: source at capacity (%d/%d) for media=%02x%02x..., redirecting to %d nodes",
MD_ID, fl->active_downloads, md->max_downloads_per_file,
req->media_id[0], req->media_id[1], fl->active_downloads);
fclose(f);
md->active_streams--; return;
}
}
md_file_load_inc(md, req->media_id, from_node);
struct stream_ctx* sc = u_calloc(1, sizeof(*sc));
if (!sc) { fclose(f); md->active_streams--; md_file_load_dec(md, req->media_id, from_node); return; }
sc->md = md;
sc->group_id = group_id;
sc->dst_node_id = from_node;
memcpy(sc->media_id, req->media_id, 16);
memcpy(sc->block_id, req->block_id, 16);
sc->chunk = req->chunk;
sc->file = f;
sc->offset = (uint64_t)block_start;
sc->remaining = (uint64_t)block_length;
sc->block_start = (uint64_t)block_start;
sc->block_data_len = (size_t)block_length;
DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%s: streaming file=%s chunk=%d start=%lld len=%lld group=%016llx to 0x%016llx",
MD_ID, path, chunk, (long long)block_start, (long long)block_length, (unsigned long long)group_id, (unsigned long long)from_node);
md->streams_started++; /* monotonic — for test admission checks */
/* start streaming — first chunk immediately, then backpressure */
memset(&sc->waiter, 0, sizeof(sc->waiter));
stream_send_chunk_cb(NULL, sc);
}
/* ── handle RELAY_FULL: relay node is overloaded, retry one of its downstream nodes ── */
static void md_handle_block_relay_full(struct media_delivery_ctx* md, uint64_t from_node,
const uint8_t* data, size_t len) {
if (len < MEDIA_RELAY_FULL_HDR_SIZE) return;
struct media_pkt_block_relay_full* rf = (struct media_pkt_block_relay_full*)data;
uint16_t nrn = rf->num_relay_nodes;
if (nrn > MD_MAX_RELAY_DOWNSTREAM) nrn = MD_MAX_RELAY_DOWNSTREAM;
DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%s: RELAY_FULL from 0x%016llx block=%02x%02x... nodes=%u",
MD_ID, (unsigned long long)from_node, rf->block_id[0], rf->block_id[1], nrn);
/* forward to media_download layer to retry from one of the listed nodes */
const uint8_t* rp = data + MEDIA_RELAY_FULL_HDR_SIZE;
uint64_t node_ids[MD_MAX_RELAY_DOWNSTREAM];
for (int i = 0; i < (int)nrn && i < MD_MAX_RELAY_DOWNSTREAM; i++) {
memcpy(&node_ids[i], rp, 8); rp += 8;
}
media_download_handle_relay_full(md->inst, rf->media_id, rf->block_id,
rf->chunk, node_ids, (int)nrn);
}
/* ── decrement stream counter (called when stream finished/cancelled internally) ── */
void media_delivery_stream_done(struct UTUN_INSTANCE* inst) {
if (inst) { struct media_delivery_ctx* md = &inst->md; if (md->active_streams > 0) md->active_streams--; md->stream_completed = 1; }
}
static void md_super_connect(struct media_delivery_ctx* md, uint64_t peer_node_id) {
struct media_super_peer* peer = md_super_peer_add(md, peer_node_id);
if (!peer) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: super_peer_add failed for 0x%016llx", MD_ID, (unsigned long long)peer_node_id); return; }
struct TOPO_GROUP* grp = NULL;
/* find group containing this peer */
struct ll_entry* gle = md->inst->topo_groups->group_list->head;
while (gle) {
struct TOPO_GROUP* g = (struct TOPO_GROUP*)gle;
if (g->conn_mgr) { grp = g; break; }
gle = gle->next;
}
if (!grp || !grp->conn_mgr) {
DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: no group/conn_mgr for supernode connect", MD_ID);
return;
}
DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%s: connecting to supernode 0x%016llx",
MD_ID, (unsigned long long)peer_node_id);
conn_mgr_open_invite(grp->instance, grp->group_id, NULL, peer_node_id, md_super_conn_cb, md, NULL);
}
/* ── supernode start/stop ── */
static int md_super_start(struct media_delivery_ctx* md) {
if (!md->inst || !md->inst->topo_sqlite_db || !md->inst->topo_groups) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: super_start missing inst/db/groups", MD_ID);
return -1;
}
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: supernode START node_id=0x%016llx",
MD_ID, (unsigned long long)md->self_node_id);
struct ll_entry* gle = md->inst->topo_groups->group_list->head;
while (gle) {
struct TOPO_GROUP* g = (struct TOPO_GROUP*)gle;
if (g->group_type != TOPO_GROUP_TYPE_CHAT) { gle = gle->next; continue; }
char sql[256];
snprintf(sql, sizeof(sql),
"SELECT node_id FROM peers_%s WHERE node_type=4 AND node_id!=%llu",
g->channel_id, (unsigned long long)md->self_node_id);
sqlite3_stmt* st = NULL;
if (sqlite3_prepare_v2(md->db, sql, -1, &st, NULL) == SQLITE_OK) {
while (sqlite3_step(st) == SQLITE_ROW) {
uint64_t pid = (uint64_t)sqlite3_column_int64(st, 0);
md_super_connect(md, pid);
}
sqlite3_finalize(st);
}
gle = gle->next;
}
return 0;
}
static void md_super_stop(struct media_delivery_ctx* md) {
if (!md->db) { DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: super_stop but no DB", MD_ID); return; }
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: supernode STOP, cleaning foreign blocks", MD_ID);
/* delete foreign blocks */
char sql[128];
snprintf(sql, sizeof(sql),
"DELETE FROM block_availability WHERE node_id!=%llu",
(unsigned long long)md->self_node_id);
sqlite3_exec(md->db, sql, NULL, NULL, NULL);
/* clear super_peers */
if (md->super_peers) {
struct ll_entry* se = md->super_peers->head;
while (se) {
struct media_super_peer* sp = (struct media_super_peer*)se->data;
if (sp->timeout_timer) { uasync_cancel_timeout(md->inst->ua, sp->timeout_timer); sp->timeout_timer = NULL; }
if (sp->connect_timer) { uasync_cancel_timeout(md->inst->ua, sp->connect_timer); sp->connect_timer = NULL; }
se = se->next;
}
queue_free(md->super_peers);
md->super_peers = NULL;
}
}
/* ── callbacks ── */
static void md_on_bgp_node(struct TOPO_GROUP* group, uint64_t node_id, int event, void* arg) {
struct media_delivery_ctx* md = (struct media_delivery_ctx*)arg;
if (!md->is_supernode || !md->db) return;
if (event == TOPO_NODE_EVENT_REMOVE) {
char sql[128];
snprintf(sql, sizeof(sql),
"DELETE FROM block_availability WHERE node_id=%llu",
(unsigned long long)node_id);
sqlite3_exec(md->db, sql, NULL, NULL, NULL);
md_super_peer_remove(md, node_id);
return;
}
/* check if this node is a supernode */
struct TOPO_GROUP_NODE* nq = topo_node_find_by_id(group, node_id);
struct TOPO_NODE* ni = nq ? topo_node_registry_find(group->instance->topo_groups, nq->node_id) : NULL;
if (!nq || !ni) return;
/* for CHAT groups, check node_type via DB */
if (group->group_type == TOPO_GROUP_TYPE_CHAT && group->channel_id[0]) {
char sql[256];
snprintf(sql, sizeof(sql),
"SELECT node_type FROM peers_%s WHERE node_id=%llu",
group->channel_id, (unsigned long long)node_id);
sqlite3_stmt* st = NULL;
if (sqlite3_prepare_v2(md->db, sql, -1, &st, NULL) == SQLITE_OK) {
int nt = 0;
if (sqlite3_step(st) == SQLITE_ROW) nt = sqlite3_column_int(st, 0);
sqlite3_finalize(st);
if (nt == 4) md_super_connect(md, node_id);
}
}
}
static void md_on_props_changed(uint64_t node_id, const char* adm_tags, void* arg) {
struct media_delivery_ctx* md = (struct media_delivery_ctx*)arg;
int is_super = adm_tags && strstr(adm_tags, "supernode=yes");
if (node_id == md->self_node_id) {
media_delivery_set_supernode(md->inst, is_super);
} else {
if (is_super) md_super_connect(md, node_id);
else md_super_peer_remove(md, node_id);
}
}
static void md_on_conn_status(struct ETCP_CONN* conn, int status, void* arg) {
struct media_delivery_ctx* md = (struct media_delivery_ctx*)arg;
if (!conn) return;
uint64_t node_id = conn->peer_node_id;
switch (status) {
case ETCP_CONN_STATUS_DOWN:
case ETCP_CONN_STATUS_DELETE:
if (md->served_nodes) {
struct ll_entry* e = queue_find_data_by_index(md->served_nodes, &node_id);
if (e) { queue_remove_data(md->served_nodes, e); queue_entry_free(e); }
}
{ struct media_super_peer* sp = md_super_peer_find(md, node_id);
if (sp) {
sp->connected = 0; sp->hello_done = 0; sp->inflight_count = 0;
if (sp->timeout_timer) { uasync_cancel_timeout(md->inst->ua, sp->timeout_timer); sp->timeout_timer = NULL; }
} }
break;
default: break;
}
}
/* ── main dispatch ── */
static void md_etcp_recv(struct ETCP_CONN* conn, struct ll_entry* entry) {
if (!conn || !entry || !entry->dgram || entry->len < 1) return;
struct UTUN_INSTANCE* inst = conn->instance;
if (!inst) return;
struct media_delivery_ctx* md = &inst->md;
if (!md->initialized) { DEBUG_WARN(DEBUG_CATEGORY_MEDIA, "%s: recv before init", MD_ID); return; }
const uint8_t* data = entry->dgram;
size_t len = entry->len;
/* router_deliver prepends svc_id byte → data[0]=svc_id, data[1]=наш subcmd.
conn_mgr uses two-tier (cmd+subcmd) hence +2; we have single-tier so +1. */
if (len < 2) return;
uint8_t subcmd = data[1];
data += 1; len -= 1;
uint64_t from_node = conn->peer_node_id;
DEBUG_DEBUG(DEBUG_CATEGORY_MEDIA, "%s: recv %s(%02x) from 0x%016llx len=%zu",
MD_ID, md_subcmd_name(subcmd), subcmd, (unsigned long long)from_node, len);
switch (subcmd) {
case MEDIA_SUBCMD_SERVE_REG:
md_handle_serve_reg(md, from_node, data, len);
break;
case MEDIA_SUBCMD_SERVE_LEAVE:
DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%s: SERVE_LEAVE from 0x%016llx",
MD_ID, (unsigned long long)from_node);
if (md->served_nodes) {
struct ll_entry* e = queue_find_data_by_index(md->served_nodes, &from_node);
if (e) { queue_remove_data(md->served_nodes, e); queue_entry_free(e); }
}
break;
case MEDIA_SUBCMD_QUERY:
md_handle_query(md, from_node, data, len);
break;
case MEDIA_SUBCMD_HAVE_BLOCK:
md_handle_have_block(md, from_node, data, len);
break;
case MEDIA_SUBCMD_BLOCK_PROCESSING:
if (md->is_supernode) md_handle_block_processing(md, from_node, data, len);
else DEBUG_WARN(DEBUG_CATEGORY_MEDIA, "%s: BLOCK_PROCESSING ignored — not supernode", MD_ID);
break;
case MEDIA_SUBCMD_SUPER_REPL:
if (md->is_supernode) md_handle_super_repl(md, from_node, data, len);
else DEBUG_WARN(DEBUG_CATEGORY_MEDIA, "%s: SUPER_REPL ignored — not supernode", MD_ID);
break;
case MEDIA_SUBCMD_SUPER_ACK:
if (md->is_supernode) md_handle_super_ack(md, from_node, data, len);
else DEBUG_WARN(DEBUG_CATEGORY_MEDIA, "%s: SUPER_ACK ignored — not supernode", MD_ID);
break;
case MEDIA_SUBCMD_SUPER_HELLO:
if (md->is_supernode) md_handle_super_hello(md, from_node, data, len);
else DEBUG_WARN(DEBUG_CATEGORY_MEDIA, "%s: SUPER_HELLO ignored — not supernode", MD_ID);
break;
case MEDIA_SUBCMD_QUERY_RESP:
media_download_handle_query_resp(inst, data, len);
break;
case MEDIA_SUBCMD_HAVE_BLOCK_ACK:
/* handled by download module */
break;
case MEDIA_SUBCMD_SERVE_ACK:
break;
case MEDIA_SUBCMD_BLOCK_REQ:
md_handle_block_req(md, from_node, data, len);
break;
case MEDIA_SUBCMD_BLOCK_OVERLOADED:
media_download_handle_overloaded(inst, data, len);
break;
case MEDIA_SUBCMD_BLOCK_RELAY_FULL:
md_handle_block_relay_full(md, from_node, data, len);
break;
case MEDIA_SUBCMD_BLOCK_CHUNK:
media_download_handle_chunk(inst, data, len);
break;
case MEDIA_SUBCMD_BLOCK_DONE:
media_download_handle_done(inst, data, len);
break;
case MEDIA_SUBCMD_CANCEL:
/* TODO: handle cancel from remote */
break;
default:
DEBUG_WARN(DEBUG_CATEGORY_MEDIA, "%s: unknown subcmd 0x%02x from 0x%016llx",
MD_ID, subcmd, (unsigned long long)from_node);
}
}
/* ── table creation ── */
static int media_delivery_create_tables(struct UTUN_INSTANCE* inst) {
sqlite3* db = inst->topo_sqlite_db;
if (!db) return 0;
const char* sql_ba =
"CREATE TABLE IF NOT EXISTS block_availability ("
" id INTEGER PRIMARY KEY AUTOINCREMENT,"
" block_uuid BLOB NOT NULL,"
" group_id INTEGER NOT NULL,"
" node_id INTEGER NOT NULL,"
" chunk INTEGER NOT NULL,"
" timestamp INTEGER NOT NULL,"
" status INTEGER NOT NULL DEFAULT 1);"; // 0=processing, 1=completed
if (sqlite3_exec(db, sql_ba, NULL, NULL, NULL) != SQLITE_OK) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: CREATE block_availability failed: %s", MD_ID, sqlite3_errmsg(db));
return -1;
}
sqlite3_exec(db, "CREATE UNIQUE INDEX IF NOT EXISTS idx_ba_uuid_node ON block_availability(block_uuid, node_id);", NULL, NULL, NULL);
sqlite3_exec(db, "CREATE INDEX IF NOT EXISTS idx_ba_id ON block_availability(id);", NULL, NULL, NULL);
sqlite3_exec(db, "CREATE INDEX IF NOT EXISTS idx_ba_node_id ON block_availability(node_id);", NULL, NULL, NULL);
const char* sql_ss =
"CREATE TABLE IF NOT EXISTS super_sync ("
" peer_node_id INTEGER PRIMARY KEY,"
" last_recv_id INTEGER NOT NULL DEFAULT 0);";
if (sqlite3_exec(db, sql_ss, NULL, NULL, NULL) != SQLITE_OK) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: CREATE super_sync failed: %s", MD_ID, sqlite3_errmsg(db));
return -1;
}
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: tables created", MD_ID);
return 0;
}
/* ── public API ── */
int media_delivery_init(struct UTUN_INSTANCE* inst) {
if (!inst) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: init with NULL inst", MD_ID); return -1; }
struct media_delivery_ctx* md = &inst->md;
memset(md, 0, sizeof(*md));
md->self_node_id = inst->node_id;
md->inst = inst;
md->db = inst->topo_sqlite_db;
md->max_downloads_per_file = 3;
if (media_delivery_create_tables(inst) != 0) return -1;
if (etcp_router_bind(inst, ETCP_RT_ID_MEDIA_DELIVERY, md_etcp_recv) != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: failed to bind ETCP_RT_ID_MEDIA_DELIVERY", MD_ID);
return -1;
}
etcp_add_conn_status_cbk(inst, md_on_conn_status, md);
member_sync_add_props_cbk(md_on_props_changed, md);
/* determine initial supernode state and subscribe to BGP callbacks */
struct ll_entry* gle = inst->topo_groups ? inst->topo_groups->group_list->head : NULL;
while (gle) {
struct TOPO_GROUP* g = (struct TOPO_GROUP*)gle;
topo_group_add_node_cbk(g, md_on_bgp_node, md);
if (g->group_type == TOPO_GROUP_TYPE_CHAT && g->channel_id[0] && inst->topo_sqlite_db) {
char sql[256];
snprintf(sql, sizeof(sql),
"SELECT adm_tags FROM peers_%s WHERE node_id=%llu",
g->channel_id, (unsigned long long)inst->node_id);
sqlite3_stmt* st = NULL;
if (sqlite3_prepare_v2(inst->topo_sqlite_db, sql, -1, &st, NULL) == SQLITE_OK) {
if (sqlite3_step(st) == SQLITE_ROW) {
const char* tags = (const char*)sqlite3_column_text(st, 0);
if (tags && strstr(tags, "supernode=yes")) { md->is_supernode = 1; }
}
sqlite3_finalize(st);
}
}
gle = gle->next;
}
if (md->is_supernode) md_super_start(md);
md->initialized = 1;
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: initialized node_id=0x%016llx super=%d",
MD_ID, (unsigned long long)md->self_node_id, md->is_supernode);
return 0;
}
void media_delivery_destroy(struct UTUN_INSTANCE* inst) {
if (!inst || !inst->md.initialized) return;
struct media_delivery_ctx* md = &inst->md;
etcp_router_unbind(inst, ETCP_RT_ID_MEDIA_DELIVERY);
etcp_remove_conn_status_cbk(inst, md_on_conn_status, md);
member_sync_remove_props_cbk(md_on_props_changed, md);
/* unsubscribe from BGP callbacks */
struct ll_entry* gle = inst->topo_groups ? inst->topo_groups->group_list->head : NULL;
while (gle) {
struct TOPO_GROUP* g = (struct TOPO_GROUP*)gle;
topo_group_remove_node_cbk(g, md_on_bgp_node, md);
gle = gle->next;
}
if (md->super_peers) {
struct ll_entry* se = md->super_peers->head;
while (se) {
struct media_super_peer* sp = (struct media_super_peer*)se->data;
if (sp->timeout_timer) { uasync_cancel_timeout(md->inst->ua, sp->timeout_timer); sp->timeout_timer = NULL; }
if (sp->connect_timer) { uasync_cancel_timeout(md->inst->ua, sp->connect_timer); sp->connect_timer = NULL; }
se = se->next;
}
queue_free(md->super_peers);
md->super_peers = NULL;
}
if (md->served_nodes) { queue_free(md->served_nodes); md->served_nodes = NULL; }
if (md->downloads) { queue_free(md->downloads); md->downloads = NULL; }
if (md->relay_blocks) { queue_free(md->relay_blocks); md->relay_blocks = NULL; }
if (md->file_loads) { queue_free(md->file_loads); md->file_loads = NULL; }
memset(md, 0, sizeof(*md));
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: destroyed", MD_ID);
}
int media_delivery_bind(struct UTUN_INSTANCE* inst) {
return media_delivery_init(inst);
}
void media_delivery_set_supernode(struct UTUN_INSTANCE* inst, int is_supernode) {
if (!inst) return;
struct media_delivery_ctx* md = &inst->md;
if (is_supernode && !md->is_supernode) {
md->is_supernode = 1;
md_super_start(md);
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: supernode ENABLED", MD_ID);
} else if (!is_supernode && md->is_supernode) {
md->is_supernode = 0;
md_super_stop(md);
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: supernode DISABLED", MD_ID);
}
}