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.
667 lines
28 KiB
667 lines
28 KiB
/* |
|
* chat_msg.c — сообщения: отправка, DB-операции для chat_sync |
|
* |
|
* Вынесено из chat_core.c для уменьшения размера модуля. |
|
*/ |
|
|
|
#include "chat_core_priv.h" |
|
#include "chat_event.h" |
|
#include "chat_setting.h" |
|
|
|
#include "../utun_instance.h" |
|
#include "../transport_layer/secure_channel.h" |
|
#include "../media_delivery/media_index.h" |
|
#include "../media_delivery/media_download.h" |
|
#include "../../lib/mem.h" |
|
#include "../../lib/platform_compat.h" |
|
|
|
#include <openssl/sha.h> |
|
#include <stdlib.h> |
|
|
|
/* forward decl */ |
|
static void chat_core_submit_media_message(struct chat_msg_submit* req); |
|
|
|
/* ─── отправка сообщения (GUI → uasync) ─── */ |
|
|
|
static int msg_get_prev_chain_hash(const char* tbl, uint8_t out[32]) { |
|
char sql[128]; snprintf(sql, sizeof(sql), |
|
"SELECT chain_hash FROM \"%s\" ORDER BY timestamp DESC LIMIT 1", tbl); |
|
sqlite3_stmt* st = NULL; |
|
if (sqlite3_prepare_v2(g_cc.db, sql, -1, &st, NULL) != SQLITE_OK) { memset(out,0,32); return -1; } |
|
if (sqlite3_step(st) == SQLITE_ROW) { |
|
const void* blob = sqlite3_column_blob(st, 0); int blen = sqlite3_column_bytes(st, 0); |
|
if (blob && blen == 32) memcpy(out, blob, 32); else memset(out, 0, 32); |
|
} else { memset(out, 0, 32); } |
|
sqlite3_finalize(st); return 0; |
|
} |
|
|
|
void chat_core_submit_message(struct chat_msg_submit* req) { |
|
if (!g_cc.initialized) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: submit before init", CC_ID); return; } |
|
if (!req) return; |
|
|
|
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: SUBMIT ch=%s ct=%s len=%u", |
|
CC_ID, req->channel_id, req->content_type, req->data_len); |
|
|
|
struct DB_SYNC_INSTANCE* si = si_find(req->channel_id); |
|
if (!si) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: no db_sync instance for ch=%s", CC_ID, req->channel_id); return; } |
|
|
|
char json[4096]; |
|
snprintf(json, sizeof(json), |
|
"{\"n\":%llu,\"ch\":\"%s\",\"ct\":\"%s\",\"d\":\"%.*s\"}", |
|
(unsigned long long)g_cc.my_node_id, req->channel_id, |
|
req->content_type, (int)req->data_len, (const char*)req->data); |
|
|
|
uint64_t ts = db_sync_next_timestamp(si); |
|
uint8_t sig_msg[8192]; size_t soff = 0; |
|
memcpy(sig_msg + soff, &ts, 8); soff += 8; |
|
size_t jl = strlen(json); memcpy(sig_msg + soff, json, jl); soff += jl; |
|
uint8_t sig[64]; |
|
if (sc_ed25519_sign(g_cc.inst->my_ed25519_privkey, sig_msg, soff, sig) != SC_OK) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: Ed25519 sign failed", CC_ID); |
|
return; |
|
} |
|
|
|
int ret = db_sync_insert_signed(si, json, strlen(json), sig, 64, ts, NULL); |
|
if (ret != 0) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: db_sync_insert_signed failed ret=%d", CC_ID, ret); |
|
return; |
|
} |
|
} |
|
|
|
void chat_core_submit_trampoline(void* arg) { |
|
struct chat_msg_submit* req = (struct chat_msg_submit*)arg; |
|
if (req->media_src[0] != '\0') |
|
chat_core_submit_media_message(req); |
|
else |
|
chat_core_submit_message(req); |
|
u_free(req); |
|
} |
|
|
|
/* ─── Медиа-сообщения: async регистрация в media_files + отправка ─── */ |
|
|
|
struct media_submit_ctx { |
|
char channel_id[64]; |
|
char content_type[32]; |
|
char media_dest[512]; |
|
char media_base[512]; |
|
uint8_t msg_data[4096]; |
|
uint32_t msg_data_len; |
|
uint64_t timestamp; |
|
}; |
|
|
|
static void on_media_registered(void* arg, int err, const struct media_index_result* result) { |
|
struct media_submit_ctx* mctx = (struct media_submit_ctx*)arg; |
|
|
|
if (err || !result) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: media registration failed ch=%s err=%d", |
|
CC_ID, mctx->channel_id, err); |
|
u_free(mctx); |
|
return; |
|
} |
|
|
|
struct DB_SYNC_INSTANCE* si = si_find(mctx->channel_id); |
|
if (!si) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: no db_sync instance for ch=%s", CC_ID, mctx->channel_id); |
|
u_free(mctx); |
|
return; |
|
} |
|
|
|
/* build full_data: <base>|<file_size>|<block_size>|<num_blocks>|<media_id_hex>|<content_hash_hex>|<id0,sig0,...> */ |
|
size_t base_len = mctx->msg_data_len; |
|
int nb = result->num_blocks; |
|
size_t pairs_len = (size_t)nb * (32 + 1 + 128 + 1); /* id_hex, sig_hex, commas */ |
|
size_t extra = 128 + 32 + 64 + pairs_len + 256; |
|
size_t fdcap = base_len + extra; |
|
char* full_data = u_malloc(fdcap); |
|
if (!full_data) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: data alloc fail", CC_ID); u_free(mctx); return; } |
|
|
|
size_t off = 0; |
|
if (base_len > 0) { |
|
char b64[2048]; size_t b64len = b64_encode(mctx->msg_data, base_len, b64, sizeof(b64)); |
|
full_data[off++] = '*'; |
|
memcpy(full_data + off, b64, b64len); off += b64len; |
|
} |
|
off += snprintf(full_data + off, fdcap - off, "|%lld|%lld|%d|", |
|
(long long)result->file_size, (long long)result->block_size, nb); |
|
for (int i = 0; i < 16; i++) off += snprintf(full_data + off, fdcap - off, "%02x", result->media_id[i]); |
|
off += snprintf(full_data + off, fdcap - off, "|"); |
|
for (int i = 0; i < 32; i++) off += snprintf(full_data + off, fdcap - off, "%02x", result->content_hash[i]); |
|
off += snprintf(full_data + off, fdcap - off, "|"); |
|
for (int n = 0; n < nb; n++) { |
|
if (n > 0) full_data[off++] = ','; |
|
for (int i = 0; i < 16; i++) |
|
off += snprintf(full_data + off, fdcap - off, "%02x", result->block_ids[n * 16 + i]); |
|
full_data[off++] = ','; |
|
for (int i = 0; i < 64; i++) |
|
off += snprintf(full_data + off, fdcap - off, "%02x", result->block_sigs[n * 64 + i]); |
|
} |
|
size_t full_data_len = off; |
|
|
|
size_t json_cap = full_data_len + 256; |
|
char* json = u_malloc(json_cap); |
|
if (!json) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: json alloc fail", CC_ID); u_free(full_data); u_free(mctx); return; } |
|
snprintf(json, json_cap, |
|
"{\"n\":%llu,\"ch\":\"%s\",\"ct\":\"%s\",\"d\":\"%.*s\"}", |
|
(unsigned long long)g_cc.my_node_id, mctx->channel_id, |
|
mctx->content_type, (int)full_data_len, full_data); |
|
|
|
uint64_t ts = db_sync_next_timestamp(si); |
|
size_t sig_msg_len = 8 + strlen(json); |
|
uint8_t* sig_msg = u_malloc(sig_msg_len); |
|
if (!sig_msg) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: sig_msg alloc fail", CC_ID); u_free(json); u_free(full_data); u_free(mctx); return; } |
|
memcpy(sig_msg, &ts, 8); |
|
memcpy(sig_msg + 8, json, sig_msg_len - 8); |
|
uint8_t json_sig[64]; |
|
if (sc_ed25519_sign(g_cc.inst->my_ed25519_privkey, sig_msg, sig_msg_len, json_sig) != SC_OK) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: Ed25519 sign JSON failed", CC_ID); |
|
u_free(sig_msg); u_free(json); u_free(full_data); u_free(mctx); return; |
|
} |
|
|
|
int ret = db_sync_insert_signed(si, json, strlen(json), json_sig, 64, ts, NULL); |
|
if (ret != 0) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: media insert failed ret=%d", CC_ID, ret); |
|
} else { |
|
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: media message inserted ch=%s ts=%llu", |
|
CC_ID, mctx->channel_id, (unsigned long long)ts); |
|
const char* base = strrchr(mctx->media_dest, '/'); |
|
const char* filename = base ? base + 1 : mctx->media_dest; |
|
char attrs[320]; snprintf(attrs, sizeof(attrs), "{\"st\":\"fl\",\"fp\":\"%s\"}", filename); |
|
chat_core_update_local_attrs(mctx->channel_id, ts, json_sig, attrs); |
|
} |
|
|
|
u_free(sig_msg); u_free(json); u_free(full_data); u_free(mctx); |
|
} |
|
|
|
static void chat_core_submit_media_message(struct chat_msg_submit* req) { |
|
if (!g_cc.initialized) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: media submit before init", CC_ID); return; } |
|
if (!req || req->media_src[0] == '\0') return; |
|
|
|
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: MEDIA SUBMIT ch=%s ct=%s src=%s dst=%s copy=%d", |
|
CC_ID, req->channel_id, req->content_type, |
|
req->media_src, req->media_dest, req->media_copy); |
|
|
|
if (!g_cc.inst->media_async) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: media_async not initialized", CC_ID); |
|
return; |
|
} |
|
|
|
/* build media_base from db_path: dirname(db_path) e.g. /path/to/data */ |
|
char media_base[1024]; |
|
{ |
|
const char* last_slash = strrchr(g_cc.db_path, '/'); |
|
if (last_slash) |
|
snprintf(media_base, sizeof(media_base), "%.*s", (int)(last_slash - g_cc.db_path), g_cc.db_path); |
|
else |
|
snprintf(media_base, sizeof(media_base), "%s", g_cc.db_path); |
|
} |
|
|
|
/* copy fields for async callback (req is freed by trampoline immediately) */ |
|
struct media_submit_ctx* mctx = u_calloc(1, sizeof(*mctx)); |
|
if (!mctx) return; |
|
snprintf(mctx->channel_id, sizeof(mctx->channel_id), "%s", req->channel_id); |
|
snprintf(mctx->content_type, sizeof(mctx->content_type), "%s", req->content_type); |
|
snprintf(mctx->media_dest, sizeof(mctx->media_dest), "%s", req->media_dest); |
|
snprintf(mctx->media_base, sizeof(mctx->media_base), "%s", media_base); |
|
mctx->timestamp = req->timestamp; |
|
if (req->data && req->data_len > 0) { |
|
mctx->msg_data_len = req->data_len < sizeof(mctx->msg_data) ? req->data_len : sizeof(mctx->msg_data) - 1; |
|
memcpy(mctx->msg_data, req->data, mctx->msg_data_len); |
|
} |
|
|
|
media_index_register_async( |
|
g_cc.inst->media_async, g_cc.inst->ua, g_cc.db, |
|
g_cc.my_node_id, g_cc.inst->my_ed25519_privkey, |
|
req->channel_id, |
|
req->media_src, req->media_dest, req->media_copy ? 1 : 0, |
|
media_base, |
|
on_media_registered, mctx); |
|
} |
|
|
|
/* ─── Обновление local_attrs (для будущей приёмной стороны) ─── */ |
|
|
|
int chat_core_update_local_attrs(const char* ch_id, uint64_t ts, |
|
const uint8_t* author_sig, |
|
const char* local_attrs_json) |
|
{ |
|
if (!g_cc.initialized || !ch_id || !author_sig || !local_attrs_json) return -1; |
|
|
|
char tbl[80]; msg_table_name(ch_id, tbl, sizeof(tbl)); |
|
char sql[256]; |
|
snprintf(sql, sizeof(sql), |
|
"UPDATE \"%s\" SET local_attrs=? WHERE timestamp=? AND author_signature=?", tbl); |
|
|
|
sqlite3_stmt* stmt = NULL; |
|
if (sqlite3_prepare_v2(g_cc.db, sql, -1, &stmt, NULL) != SQLITE_OK) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: update local_attrs prep fail: %s", CC_ID, sqlite3_errmsg(g_cc.db)); |
|
return -1; |
|
} |
|
sqlite3_bind_text(stmt, 1, local_attrs_json, -1, SQLITE_STATIC); |
|
sqlite3_bind_int64(stmt, 2, (sqlite3_int64)ts); |
|
sqlite3_bind_blob(stmt, 3, author_sig, DB_SIG_SIZE, SQLITE_STATIC); |
|
int rc = sqlite3_step(stmt); |
|
sqlite3_finalize(stmt); |
|
if (rc != SQLITE_DONE) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: update local_attrs step fail rc=%d", CC_ID, rc); |
|
return -1; |
|
} |
|
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: local_attrs updated ch=%s ts=%llu", CC_ID, ch_id, (unsigned long long)ts); |
|
return 0; |
|
} |
|
|
|
/* ─── DB-операции для chat_sync ─── */ |
|
|
|
uint32_t chat_core_count(const char* ch_id) { |
|
if (!g_cc.initialized) return 0; |
|
char tbl[80]; msg_table_name(ch_id, tbl, sizeof(tbl)); |
|
char sql[128]; snprintf(sql, sizeof(sql), |
|
"SELECT COUNT(*) FROM \"%s\"", tbl); |
|
sqlite3_stmt* stmt = NULL; |
|
if (sqlite3_prepare_v2(g_cc.db, sql, -1, &stmt, NULL) != SQLITE_OK) return 0; |
|
uint32_t cnt = 0; |
|
if (sqlite3_step(stmt) == SQLITE_ROW) |
|
cnt = (uint32_t)sqlite3_column_int64(stmt, 0); |
|
sqlite3_finalize(stmt); |
|
return cnt; |
|
} |
|
|
|
int chat_core_chain_hash_at(const char* ch_id, uint32_t pos, uint8_t* hash_out) { |
|
if (!g_cc.initialized || !hash_out) return -1; |
|
char tbl[80]; msg_table_name(ch_id, tbl, sizeof(tbl)); |
|
char sql[200]; snprintf(sql, sizeof(sql), |
|
"SELECT chain_hash FROM \"%s\" ORDER BY timestamp, node_id ASC LIMIT 1 OFFSET ?", |
|
tbl); |
|
sqlite3_stmt* stmt = NULL; |
|
if (sqlite3_prepare_v2(g_cc.db, sql, -1, &stmt, NULL) != SQLITE_OK) { |
|
memset(hash_out, 0, 32); return -1; |
|
} |
|
sqlite3_bind_int64(stmt, 1, (sqlite3_int64)pos); |
|
if (sqlite3_step(stmt) == SQLITE_ROW) { |
|
const void* blob = sqlite3_column_blob(stmt, 0); |
|
int len = sqlite3_column_bytes(stmt, 0); |
|
if (blob && len == 32) memcpy(hash_out, blob, 32); |
|
else memset(hash_out, 0, 32); |
|
} else { |
|
memset(hash_out, 0, 32); |
|
} |
|
sqlite3_finalize(stmt); |
|
return 0; |
|
} |
|
|
|
int chat_core_list_channels(uint8_t* buf, size_t buf_size, size_t* out_len) { |
|
if (!g_cc.initialized || !buf || !out_len) return -1; |
|
sqlite3_stmt* stmt = NULL; |
|
if (sqlite3_prepare_v2(g_cc.db, |
|
"SELECT channel_id FROM channels ORDER BY created_at ASC", |
|
-1, &stmt, NULL) != SQLITE_OK) return -1; |
|
|
|
uint8_t* out = buf; |
|
uint8_t* start = out; |
|
out += 2; |
|
|
|
uint16_t cnt = 0; |
|
while (sqlite3_step(stmt) == SQLITE_ROW) { |
|
const char* ch_id = (const char*)sqlite3_column_text(stmt, 0); |
|
int ch_len = sqlite3_column_bytes(stmt, 0); |
|
if (!ch_id || ch_len <= 0 || ch_len > 63) continue; |
|
size_t need = (size_t)(out - buf) + 1 + (size_t)ch_len; |
|
if (need > buf_size) break; |
|
*out++ = (uint8_t)ch_len; |
|
memcpy(out, ch_id, (size_t)ch_len); out += ch_len; |
|
cnt++; |
|
} |
|
sqlite3_finalize(stmt); |
|
|
|
memcpy(start, &cnt, 2); |
|
*out_len = (size_t)(out - buf); |
|
return 0; |
|
} |
|
|
|
int chat_core_list_peers(const char* ch_id, uint8_t* buf, size_t buf_size, |
|
size_t* out_len) { |
|
if (!g_cc.initialized || !buf || !out_len) return -1; |
|
char tbl[80]; peers_table_name(ch_id, tbl, sizeof(tbl)); |
|
char sql[128]; snprintf(sql, sizeof(sql), |
|
"SELECT node_id FROM \"%s\"", tbl); |
|
sqlite3_stmt* stmt = NULL; |
|
if (sqlite3_prepare_v2(g_cc.db, sql, -1, &stmt, NULL) != SQLITE_OK) { |
|
uint16_t z = 0; memcpy(buf, &z, 2); *out_len = 2; return -1; |
|
} |
|
uint16_t cnt = 0; |
|
uint8_t* out = buf + 2; |
|
while (sqlite3_step(stmt) == SQLITE_ROW) { |
|
if ((size_t)(out - buf) + 8 > buf_size) break; |
|
uint64_t nid = (uint64_t)sqlite3_column_int64(stmt, 0); |
|
memcpy(out, &nid, 8); out += 8; cnt++; |
|
} |
|
sqlite3_finalize(stmt); |
|
memcpy(buf, &cnt, 2); |
|
*out_len = (size_t)(out - buf); |
|
return 0; |
|
} |
|
|
|
int chat_core_load_nodeinfo(uint64_t node_id, uint8_t* buf, size_t buf_size, |
|
size_t* out_len) { |
|
if (!g_cc.initialized || !buf || !out_len) return -1; |
|
|
|
sqlite3_stmt* stmt = NULL; |
|
if (sqlite3_prepare_v2(g_cc.db, |
|
"SELECT x25519_pubkey, ed25519_pubkey FROM nodes WHERE node_id=?", |
|
-1, &stmt, NULL) != SQLITE_OK) return -1; |
|
sqlite3_bind_int64(stmt, 1, (sqlite3_int64)node_id); |
|
if (sqlite3_step(stmt) != SQLITE_ROW) { sqlite3_finalize(stmt); return -1; } |
|
const uint8_t* x25519 = (const uint8_t*)sqlite3_column_blob(stmt, 0); |
|
int x25519_len = sqlite3_column_bytes(stmt, 0); |
|
const uint8_t* ed25519 = (const uint8_t*)sqlite3_column_blob(stmt, 1); |
|
int ed25519_len = sqlite3_column_bytes(stmt, 1); |
|
|
|
uint8_t* out = buf; |
|
if (x25519 && x25519_len == 32) memcpy(out, x25519, 32); |
|
else memset(out, 0, 32); |
|
out += 32; |
|
if (ed25519 && ed25519_len == 32) memcpy(out, ed25519, 32); |
|
else memset(out, 0, 32); |
|
out += 32; |
|
sqlite3_finalize(stmt); |
|
|
|
if (sqlite3_prepare_v2(g_cc.db, |
|
"SELECT family, protocol, address, port, rtt FROM node_addresses WHERE node_id=?", |
|
-1, &stmt, NULL) != SQLITE_OK) { |
|
*out++ = 0; *out_len = (size_t)(out - buf); return 0; |
|
} |
|
sqlite3_bind_int64(stmt, 1, (sqlite3_int64)node_id); |
|
uint8_t* cnt_pos = out; *out++ = 0; |
|
uint8_t addr_cnt = 0; |
|
|
|
while (sqlite3_step(stmt) == SQLITE_ROW && addr_cnt < 255) { |
|
int family = sqlite3_column_int(stmt, 0); |
|
int proto = sqlite3_column_int(stmt, 1); |
|
const uint8_t* addr = (const uint8_t*)sqlite3_column_blob(stmt, 2); |
|
int addr_len = sqlite3_column_bytes(stmt, 2); |
|
uint16_t port = (uint16_t)sqlite3_column_int(stmt, 3); |
|
int16_t rtt = (int16_t)sqlite3_column_int(stmt, 4); |
|
|
|
if ((size_t)(out - buf) + 7 + (size_t)addr_len > buf_size) break; |
|
*out++ = (uint8_t)family; |
|
*out++ = (uint8_t)proto; |
|
*out++ = (uint8_t)addr_len; |
|
if (addr_len > 0) { memcpy(out, addr, (size_t)addr_len); out += addr_len; } |
|
memcpy(out, &port, 2); out += 2; |
|
memcpy(out, &rtt, 2); out += 2; |
|
addr_cnt++; |
|
} |
|
*cnt_pos = addr_cnt; |
|
sqlite3_finalize(stmt); |
|
|
|
*out_len = (size_t)(out - buf); |
|
return 0; |
|
} |
|
|
|
/* ─── download completion ─── */ |
|
|
|
struct md_done_ctx { |
|
char channel_id[64]; |
|
uint64_t ts; |
|
int64_t msg_id; |
|
uint8_t author_sig[64]; |
|
char dest_relpath[512]; |
|
}; |
|
|
|
static void md_download_progress_cb(void* arg, int blocks_done, int num_blocks) { |
|
struct md_done_ctx* ctx = (struct md_done_ctx*)arg; |
|
uint8_t evt[256]; |
|
uint8_t cl = (uint8_t)strlen(ctx->channel_id); |
|
evt[0] = cl; memcpy(evt + 1, ctx->channel_id, cl); |
|
memcpy(evt + 1 + cl, &ctx->msg_id, 8); |
|
uint32_t bd = (uint32_t)blocks_done, nb = (uint32_t)num_blocks; |
|
memcpy(evt + 1 + cl + 8, &bd, 4); memcpy(evt + 1 + cl + 12, &nb, 4); |
|
chat_event_post(CHAT_EVT_DOWNLOAD_PROGRESS, evt, 1 + cl + 16); |
|
} |
|
|
|
static void md_download_done_cb(void* arg, int err) { |
|
struct md_done_ctx* ctx = (struct md_done_ctx*)arg; |
|
if (!err) { |
|
char attrs[1024]; |
|
snprintf(attrs, sizeof(attrs), "{\"st\":\"fl\",\"fp\":\"%s\"}", ctx->dest_relpath); |
|
chat_core_update_local_attrs(ctx->channel_id, ctx->ts, ctx->author_sig, attrs); |
|
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: media downloaded ch=%s", CC_ID, ctx->channel_id); |
|
} else { |
|
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: media download failed ch=%s err=%d", CC_ID, ctx->channel_id, err); |
|
} |
|
uint8_t evt[80]; uint8_t cl = (uint8_t)strlen(ctx->channel_id); |
|
evt[0] = cl; memcpy(evt + 1, ctx->channel_id, cl); |
|
memcpy(evt + 1 + cl, &ctx->msg_id, 8); |
|
chat_event_post(CHAT_EVT_ATTACHMENT_DOWNLOADED, evt, 1 + cl + 8); |
|
u_free(ctx); |
|
} |
|
|
|
static int md_start_download(struct UTUN_INSTANCE* inst, |
|
const char* data_str, size_t data_len, |
|
const char* ch_id, const char* base_filename, |
|
uint64_t ts, const uint8_t* author_sig, |
|
uint64_t author_node_id, int64_t msg_id) { |
|
(void)data_len; |
|
const char* p = data_str; |
|
while (*p && *p != '|') p++; |
|
if (!*p) return -1; p++; |
|
long long fsize = strtoll(p, (char**)&p, 10); |
|
if (*p != '|') return -1; p++; |
|
long long bsize = strtoll(p, (char**)&p, 10); |
|
if (*p != '|') return -1; p++; |
|
int nb = (int)strtol(p, (char**)&p, 10); |
|
while (*p && *p != '|') p++; |
|
if (!*p) return -1; p++; |
|
if (nb <= 0 || nb > 100 || fsize <= 0) return -1; |
|
|
|
struct media_index_result result; |
|
memset(&result, 0, sizeof(result)); |
|
result.file_size = fsize; result.block_size = bsize; result.num_blocks = nb; |
|
|
|
for (int i = 0; i < 16; i++) { |
|
char b[3] = { p[i*2], p[i*2+1], 0 }; |
|
result.media_id[i] = (uint8_t)strtol(b, NULL, 16); |
|
} |
|
p += 32; if (*p != '|') return -1; p++; |
|
for (int i = 0; i < 32; i++) { |
|
char b[3] = { p[i*2], p[i*2+1], 0 }; |
|
result.content_hash[i] = (uint8_t)strtol(b, NULL, 16); |
|
} |
|
p += 64; if (*p != '|') return -1; p++; |
|
|
|
result.block_ids = u_malloc((size_t)nb * 16); |
|
result.block_sigs = u_malloc((size_t)nb * 64); |
|
if (!result.block_ids || !result.block_sigs) { media_index_result_free(&result); return -1; } |
|
for (int n = 0; n < nb; n++) { |
|
for (int i = 0; i < 16; i++) { |
|
char b[3] = { p[i*2], p[i*2+1], 0 }; |
|
result.block_ids[n * 16 + i] = (uint8_t)strtol(b, NULL, 16); |
|
} |
|
p += 32; if (*p != ',') { media_index_result_free(&result); return -1; } p++; |
|
for (int i = 0; i < 64; i++) { |
|
char b[3] = { p[i*2], p[i*2+1], 0 }; |
|
result.block_sigs[n * 64 + i] = (uint8_t)strtol(b, NULL, 16); |
|
} |
|
p += 128; if (n < nb - 1 && *p == ',') p++; |
|
} |
|
|
|
char dest[1024], media_base[512]; |
|
{ |
|
const char* last_slash = strrchr(g_cc.db_path, '/'); |
|
if (last_slash) { |
|
snprintf(media_base, sizeof(media_base), "%.*s", (int)(last_slash - g_cc.db_path), g_cc.db_path); |
|
} else { |
|
snprintf(media_base, sizeof(media_base), "%s", g_cc.db_path); |
|
} |
|
} |
|
snprintf(dest, sizeof(dest), "%s/media/%s/%s", media_base, ch_id, base_filename); |
|
|
|
const char* fp_name = strrchr(dest, '/'); fp_name = fp_name ? fp_name + 1 : dest; |
|
|
|
struct md_done_ctx* ctx = u_calloc(1, sizeof(*ctx)); |
|
if (!ctx) { media_index_result_free(&result); return -1; } |
|
snprintf(ctx->channel_id, sizeof(ctx->channel_id), "%s", ch_id); |
|
ctx->ts = ts; ctx->msg_id = msg_id; |
|
snprintf(ctx->dest_relpath, sizeof(ctx->dest_relpath), "%s", fp_name); |
|
if (author_sig) memcpy(ctx->author_sig, author_sig, 64); |
|
|
|
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: download start ch=%s media=%02x%02x... blocks=%d size=%lld dest=%s", |
|
CC_ID, ch_id, result.media_id[0], result.media_id[1], nb, (long long)fsize, dest); |
|
media_download_start(inst, 0, &result, dest, media_base, author_node_id, md_download_done_cb, ctx, md_download_progress_cb, ctx); |
|
media_index_result_free(&result); |
|
return 0; |
|
} |
|
|
|
static void md_auto_download(struct UTUN_INSTANCE* inst, |
|
const char* data_str, size_t data_len, |
|
const char* ch_id, const char* base_filename, |
|
uint64_t ts, const uint8_t* author_sig, |
|
uint64_t author, int64_t msg_id) { |
|
if (!chat_setting_get_int("storage_autoload", 1)) return; |
|
uint64_t max_size = g_cc.inst->config->global.chatserver_storage_unit_size; |
|
|
|
const char* p = data_str; |
|
while (*p && *p != '|') p++; |
|
if (!*p) return; p++; |
|
long long fsize = strtoll(p, (char**)&p, 10); |
|
if (max_size > 0 && fsize > (long long)max_size) { |
|
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: media too large (%lld > %llu bytes), skip auto-download", |
|
CC_ID, fsize, (unsigned long long)max_size); |
|
return; |
|
} |
|
md_start_download(inst, data_str, data_len, ch_id, base_filename, ts, author_sig, author, msg_id); |
|
} |
|
|
|
/* ─── db_sync callback ─── */ |
|
|
|
void on_msg_inserted(struct DB_SYNC_INSTANCE* si, uint64_t record_ts, const char* data, size_t len, uint64_t author, void* arg) { |
|
(void)si; |
|
const char* ch_id = (const char*)arg; |
|
uint8_t evt[73]; uint8_t cl = (uint8_t)strlen(ch_id); evt[0] = cl; memcpy(evt + 1, ch_id, cl); |
|
memcpy(evt + 1 + cl, &author, 8); |
|
chat_event_post(CHAT_EVT_MSG_RECEIVED, evt, 1 + cl + 8); |
|
|
|
/* auto-download media if message contains media metadata */ |
|
if (data && len > 0 && author != g_cc.my_node_id) { |
|
/* data from db_sync may not be null-terminated; copy to local buffer */ |
|
char buf[4096]; size_t blen = len < sizeof(buf) - 1 ? len : sizeof(buf) - 1; |
|
memcpy(buf, data, blen); buf[blen] = '\0'; |
|
const char* ct = strstr(buf, "\"ct\":\""); |
|
const char* dpos = strstr(buf, "\"d\":\""); |
|
if (ct && dpos) { |
|
const char* d_start = dpos + 5; |
|
char body[4096]; size_t bi = 0; |
|
while (*d_start && *d_start != '\"' && bi < sizeof(body) - 1) body[bi++] = *d_start++; |
|
body[bi] = '\0'; |
|
|
|
char base_filename[256]; size_t bfi = 0; |
|
while (body[bfi] && body[bfi] != '|' && bfi < sizeof(base_filename) - 1) |
|
base_filename[bfi++] = body[bfi]; |
|
base_filename[bfi] = '\0'; |
|
/* base64-decode display prefix if marked with '*' (new format) */ |
|
if (base_filename[0] == '*' && bfi > 1) { |
|
uint8_t dec[256]; int dlen = (int)b64_decode(base_filename + 1, bfi - 1, dec, sizeof(dec)); |
|
if (dlen > 0 && dlen < (int)sizeof(base_filename)) { memcpy(base_filename, dec, (size_t)dlen); base_filename[dlen] = '\0'; } |
|
} |
|
if (bfi == 0) snprintf(base_filename, sizeof(base_filename), "file"); |
|
|
|
const uint8_t* sig_field = (const uint8_t*)strstr(buf, "\"sig\":\""); |
|
uint8_t author_sig[64] = {0}; |
|
if (sig_field) { |
|
/* read author_signature from DB — find by ts in on_msg_inserted we don't have sig |
|
so leave as zeros; use record_ts to locate the record later */ |
|
} |
|
|
|
/* We don't have author_sig in the callback; derive from body content. |
|
Instead, read it from the DB using the just-inserted record */ |
|
char tbl[80]; msg_table_name(ch_id, tbl, sizeof(tbl)); |
|
char sql[256]; |
|
snprintf(sql, sizeof(sql), "SELECT author_signature, id FROM \"%s\" WHERE timestamp=? AND node_id=? LIMIT 1", tbl); |
|
sqlite3_stmt* st = NULL; |
|
if (sqlite3_prepare_v2(g_cc.db, sql, -1, &st, NULL) == SQLITE_OK) { |
|
sqlite3_bind_int64(st, 1, (sqlite3_int64)record_ts); |
|
sqlite3_bind_int64(st, 2, (sqlite3_int64)author); |
|
int64_t msg_id = 0; |
|
if (sqlite3_step(st) == SQLITE_ROW) { |
|
const void* sblob = sqlite3_column_blob(st, 0); |
|
int sblen = sqlite3_column_bytes(st, 0); |
|
if (sblob && sblen == 64) memcpy(author_sig, sblob, 64); |
|
msg_id = sqlite3_column_int64(st, 1); |
|
} |
|
sqlite3_finalize(st); |
|
if (msg_id > 0) md_auto_download(g_cc.inst, body, bi, ch_id, base_filename, record_ts, author_sig, author, msg_id); |
|
} |
|
} |
|
} |
|
|
|
sqlite3_stmt* st = NULL; |
|
sqlite3_prepare_v2(g_cc.db, "UPDATE channels SET last_msg_at = ? WHERE channel_id = ?", -1, &st, NULL); |
|
if (st) { |
|
sqlite3_bind_int64(st, 1, (sqlite3_int64)record_ts); |
|
sqlite3_bind_text(st, 2, ch_id, -1, SQLITE_STATIC); |
|
sqlite3_step(st); sqlite3_finalize(st); |
|
} |
|
} |
|
|
|
/* ─── manual attachment download (called from JNI bridge via uasync post) ─── */ |
|
|
|
void chat_core_attachment_download(const char* channel_id, int64_t msg_id) { |
|
if (!channel_id || msg_id <= 0) return; |
|
char tbl[80]; msg_table_name(channel_id, tbl, sizeof(tbl)); |
|
char sql[256]; |
|
snprintf(sql, sizeof(sql), "SELECT data, timestamp, author_signature, node_id FROM \"%s\" WHERE id=?", tbl); |
|
sqlite3_stmt* st = NULL; |
|
if (sqlite3_prepare_v2(g_cc.db, sql, -1, &st, NULL) != SQLITE_OK) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: download sql prep failed ch=%s", CC_ID, channel_id); |
|
return; |
|
} |
|
sqlite3_bind_int64(st, 1, (sqlite3_int64)msg_id); |
|
if (sqlite3_step(st) != SQLITE_ROW) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: download msg not found ch=%s id=%lld", CC_ID, channel_id, (long long)msg_id); |
|
sqlite3_finalize(st); |
|
return; |
|
} |
|
const uint8_t* jdata = sqlite3_column_text(st, 0); |
|
int jlen = sqlite3_column_bytes(st, 0); |
|
int64_t db_ts = sqlite3_column_int64(st, 1); |
|
const uint8_t* sig_blob = (const uint8_t*)sqlite3_column_blob(st, 2); |
|
int sig_len = sqlite3_column_bytes(st, 2); |
|
int64_t author_node_id = sqlite3_column_int64(st, 3); |
|
if (!jdata || jlen <= 0 || !sig_blob || sig_len != 64) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: download bad msg data ch=%s id=%lld", CC_ID, channel_id, (long long)msg_id); |
|
sqlite3_finalize(st); |
|
return; |
|
} |
|
|
|
const char* ds = strstr((const char*)jdata, "\"d\":\""); |
|
if (!ds) { sqlite3_finalize(st); return; } |
|
const char* d_start = ds + 5; |
|
char body[4096]; size_t bi = 0; |
|
while (*d_start && *d_start != '\"' && bi < sizeof(body) - 1) body[bi++] = *d_start++; |
|
body[bi] = '\0'; |
|
|
|
char base_filename[256]; size_t bfi = 0; |
|
while (body[bfi] && body[bfi] != '|' && bfi < sizeof(base_filename) - 1) |
|
base_filename[bfi++] = body[bfi]; |
|
base_filename[bfi] = '\0'; |
|
if (base_filename[0] == '*' && bfi > 1) { |
|
uint8_t dec[256]; int dlen = (int)b64_decode(base_filename + 1, bfi - 1, dec, sizeof(dec)); |
|
if (dlen > 0 && dlen < (int)sizeof(base_filename)) { memcpy(base_filename, dec, (size_t)dlen); base_filename[dlen] = '\0'; } |
|
} |
|
if (bfi == 0) snprintf(base_filename, sizeof(base_filename), "file"); |
|
|
|
uint8_t author_sig[64]; memcpy(author_sig, sig_blob, 64); |
|
sqlite3_finalize(st); |
|
|
|
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: manual download ch=%s id=%lld file=%s", |
|
CC_ID, channel_id, (long long)msg_id, base_filename); |
|
md_start_download(g_cc.inst, body, bi, channel_id, base_filename, (uint64_t)db_ts, author_sig, (uint64_t)author_node_id, msg_id); |
|
} |
|
|
|
static void chat_core_attachment_download_trampoline_impl(void* arg) { |
|
struct attachment_dl_req* req = (struct attachment_dl_req*)arg; |
|
chat_core_attachment_download(req->channel_id, req->msg_id); |
|
u_free(req); |
|
} |
|
|
|
void chat_core_attachment_download_trampoline(void* arg) { |
|
chat_core_attachment_download_trampoline_impl(arg); |
|
}
|
|
|