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.
1074 lines
45 KiB
1074 lines
45 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"; |
|
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) { |
|
const char* sql = |
|
"INSERT OR REPLACE INTO block_availability(block_uuid,group_id,node_id,chunk,timestamp)" |
|
" 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); |
|
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; |
|
} |
|
|
|
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" |
|
" 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 + 52 > 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; |
|
} |
|
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); |
|
} |
|
|
|
/* ── 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[512]; |
|
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[520]; |
|
int off = 0; |
|
pkt[off++] = MEDIA_SUBCMD_SUPER_REPL; |
|
|
|
int num_entries = nb / 52; |
|
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 += 52; |
|
} |
|
|
|
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_DEBUG, "%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[3] = { MEDIA_SUBCMD_SERVE_ACK, 0, 0 }; |
|
memcpy(ack + 1, &pkt->group_id, 8); |
|
(void)from_node; |
|
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%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_DEBUG, "%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_insert(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_DEBUG, "%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; |
|
} |
|
} |
|
} |
|
|
|
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; |
|
md_ba_insert(md->db, block_uuid, (uint64_t)gid, (uint64_t)nid, (uint32_t)ch, ts); |
|
} |
|
|
|
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_DEBUG, "%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_DEBUG, "%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_DEBUG, "%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_DEBUG, "%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_DEBUG, "%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); |
|
media_delivery_stream_done(md->inst); |
|
fclose(sc->file); |
|
u_free(sc); |
|
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%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); |
|
} |
|
} |
|
|
|
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_DEBUG, "%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_DEBUG, "%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) { 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]); sqlite3_finalize(st); 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; |
|
} |
|
struct stream_ctx* sc = u_calloc(1, sizeof(*sc)); |
|
if (!sc) { fclose(f); md->active_streams--; 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_DEBUG, "%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); |
|
} |
|
|
|
/* ── 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_DEBUG, "%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; |
|
|
|
if (status == ETCP_CONN_STATUS_DOWN || status == ETCP_CONN_STATUS_DELETE) { |
|
/* remove from served_nodes */ |
|
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); } |
|
} |
|
/* mark super_peer as disconnected */ |
|
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; } |
|
} |
|
} |
|
} |
|
|
|
/* ── 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_DEBUG, "%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_DEBUG, "%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_DEBUG, "%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_SUPER_REPL: |
|
if (md->is_supernode) md_handle_super_repl(md, from_node, data, len); |
|
else DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "%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_DEBUG, "%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_DEBUG, "%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_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_DEBUG, "%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);"; |
|
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; |
|
|
|
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; } |
|
|
|
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); |
|
} |
|
}
|
|
|