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.
 
 
 
 
 
 

364 lines
19 KiB

/* test_dm_media — реальные A–S–B и S–M, независимые роли supernode/storage.
* Прямая доставка, offline custody, restart S/M при выключенном A, TTL,
* квота, точные квитанции, поздние STORE/STORED и отзыв активных операций.
*/
#include <stdio.h>
#include <string.h>
#include <stdlib.h>
#include <openssl/evp.h>
#include <openssl/sha.h>
#include "../lib/platform_compat.h"
#include "../lib/debug_config.h"
#include "../lib/mem.h"
#include "../src/utun_instance.h"
#include "../src/ntp_time.h"
#include "../src/chat/chat_core.h"
#include "../src/chat/member_sync.h"
#include "../src/chat/chat_setting.h"
#include "../src/dm/dm_core.h"
#include "../src/dm/dm_crypto.h"
#include "../src/dm/dm_media.h"
#include "../src/routing_layer/topo_node_sqlite.h"
#include "../src/routing_layer/topo_group.h"
#include "../src/routing_layer/etcp_router.h"
#include "../src/transport_layer/secure_channel.h"
#include "../src/transport_layer/etcp.h"
#include "../src/media_delivery/file_transfer.h"
#include "test_utils.h"
#define GROUP 4242424242424243ULL
#define CHANNEL "4242424242424243"
#define FILE_BYTES (3 * 32768 + 17)
enum { A, B, S, M, NODES };
struct media_test {
struct UASYNC* ua;
struct UTUN_INSTANCE* nodes[NODES];
struct SC_MYKEYS keys[NODES];
uint64_t ids[NODES];
uint8_t ed[NODES][32], ch_x[32], ch_ed[32], ch_priv[32], ch_sig[64];
char dir[512], conv[64], file[768];
uint8_t body[2048], id[16], receipt[DM_MEDIA_RECEIPT_SIZE];
size_t body_len;
unsigned phase, ticks, heartbeats, work_ticks;
int failed;
void* timer;
};
static int value(struct UTUN_INSTANCE* inst, const char* sql) {
if (!inst) return -1;
sqlite3_stmt* st = NULL;
int result = -1;
if (sqlite3_prepare_v2(inst->topo_sqlite_db, sql, -1, &st, NULL) == SQLITE_OK && sqlite3_step(st) == SQLITE_ROW)
result = sqlite3_column_int(st, 0);
sqlite3_finalize(st);
return result;
}
static int live(struct UTUN_INSTANCE* inst, uint64_t peer) {
return inst && etcp_router_get_path(inst, GROUP, peer, ETCP_RT_ID_DM).next_hop_node_id != 0;
}
static void failure(struct media_test* t, const char* reason) {
DEBUG_ERROR(DEBUG_CATEGORY_DM, "media test failed phase=%u reason=%s", t->phase, reason);
t->failed = 1;
uasync_stop(t->ua);
}
static void hex(const uint8_t* p, size_t len, char* out) {
for (size_t i = 0; i < len; i++) snprintf(out + i * 2, 3, "%02x", p[i]);
}
/* Четыре ядра используют разные БД; клиентские физические линки сходятся на S. */
static int start_node(struct media_test* t, int index, int base) {
char path[768], pub[65], priv[65], server[65], db[768];
hex(t->keys[index].public_key, 32, pub); hex(t->keys[index].private_key, 32, priv);
snprintf(db, sizeof(db), "%s/db%d", t->dir, index); utun_mkdir(db, 0700);
snprintf(path, sizeof(path), "%s/%d.conf", t->dir, index);
FILE* f = fopen(path, "w");
if (!f) return -1;
fprintf(f, "[global]\nmy_private_key=%s\nmy_public_key=%s\ndb_path=%s\n"
"[server:s1]\naddr=127.0.0.1:%d\ntype=public\n[allowed_keys]\nallow_all=1\n"
"[chatserver]\nstorage_autoload=0\n", priv, pub, db, base + index);
if (index != S) {
hex(t->keys[S].public_key, 32, server);
fprintf(f, "[client:to_s]\npeer_public_key=%s\nlink=s1:127.0.0.1:%d\n", server, base + S);
}
fclose(f);
t->nodes[index] = utun_instance_create(t->ua, path);
return t->nodes[index] ? utun_instance_init(t->nodes[index]) : -1;
}
/* Подписать мемберские и административные блоки, не подменяя авторизацию SQL-ом. */
static int setup(struct media_test* t) {
for (int i = 0; i < NODES; i++) {
struct UTUN_INSTANCE* inst = t->nodes[i];
if (topo_node_sqlite_channel_put(inst->topo_sqlite_db, CHANNEL, "dm-media-test", t->ids[A],
t->ch_x, NULL, t->ch_ed, i == A ? t->ch_priv : NULL, t->ch_sig)) return -1;
chat_core_ensure_channel_ready(inst, CHANNEL);
for (int j = 0; j < NODES; j++) {
uint8_t msg[512], join[64], adm_sig[64];
uint64_t ts = ntp_time_get_seconds(inst);
int n = member_sync_build_join_msg(t->ch_x, t->ch_ed, t->ids[j], t->keys[j].public_key, ts, msg, sizeof(msg));
if (n < 0 || sc_ed25519_sign(t->nodes[j]->my_ed25519_privkey, msg, (size_t)n, join) != SC_OK) return -1;
const char* tags = j == S ? "{\"supernode\":\"yes\",\"ver\":\"1\"}" :
j == M ? "{\"storage\":\"yes\",\"ver\":\"1\"}" : "{\"ver\":\"1\"}";
size_t off = strlen(tags);
memcpy(msg, tags, off); memcpy(msg + off, &t->ids[j], 8);
if (sc_ed25519_sign(t->ch_priv, msg, off + 8, adm_sig) != SC_OK) return -1;
char info[64]; snprintf(info, sizeof(info), "{\"name\":\"node%d\"}", j);
if (member_sync_put(inst, CHANNEL, t->ids[j], t->keys[j].public_key, t->ed[j], join, ts,
NULL, 0, info, tags, adm_sig, 0, 0, NULL) < 0) return -1;
}
if (i != S) member_sync_start(inst, t->ids[S], CHANNEL, NULL, NULL);
}
return dm_start(t->nodes[A], t->ids[B], t->keys[B].public_key, t->ed[B], "B", CHANNEL);
}
/* Снимок подписанного тела нужен для атак/replay и не заменяет UDP-передачу. */
static int snapshot(struct media_test* t, unsigned seq) {
sqlite3_stmt* st = NULL;
if (sqlite3_prepare_v2(t->nodes[A]->topo_sqlite_db, "SELECT body FROM dm_files WHERE seq=? AND dir=1", -1, &st, NULL)) return -1;
sqlite3_bind_int(st, 1, (int)seq);
int rc = sqlite3_step(st);
size_t n = rc == SQLITE_ROW ? (size_t)sqlite3_column_bytes(st, 0) : 0;
if (n && n <= sizeof(t->body)) memcpy(t->body, sqlite3_column_blob(st, 0), n);
sqlite3_finalize(st);
const uint8_t* desc = n <= sizeof(t->body) ? dm_message_media(t->body, n) : NULL;
if (!desc) return -1;
memcpy(t->id, desc, 16); t->body_len = n;
return 0;
}
/* Полученный файл сверяется побайтно; размер или ready-флаг недостаточны. */
static int file_correct(struct media_test* t) {
sqlite3_stmt* st = NULL;
if (sqlite3_prepare_v2(t->nodes[B]->topo_sqlite_db, "SELECT local_path,receipt FROM dm_files WHERE id=? AND ready=1", -1, &st, NULL)) return 0;
sqlite3_bind_blob(st, 1, t->id, 16, SQLITE_STATIC);
int rc = sqlite3_step(st);
char path[1024] = "";
if (rc == SQLITE_ROW && sqlite3_column_bytes(st, 1) == sizeof(t->receipt)) {
snprintf(path, sizeof(path), "%s", (const char*)sqlite3_column_text(st, 0));
memcpy(t->receipt, sqlite3_column_blob(st, 1), sizeof(t->receipt));
}
sqlite3_finalize(st);
FILE* f = path[0] ? fopen(path, "rb") : NULL;
if (!f) return 0;
int ok = 1;
for (unsigned i = 0; i < FILE_BYTES; i++) if (fgetc(f) != (int)((i * 73 + 19) & 255)) { ok = 0; break; }
if (ok && fgetc(f) != EOF) ok = 0;
fclose(f);
return ok;
}
/* Вход после проверки router, для проверки SQL rollback и подписанных ACK. */
static void control(struct media_test* t, int dst, uint64_t source, uint8_t cmd, const uint8_t* p, size_t len) {
struct UTUN_INSTANCE* inst = t->nodes[dst];
struct ETCP_CONN* conn = instance_find_conn(inst, t->ids[dst == S ? M : S]);
struct ll_entry* e = ll_alloc_lldgram((uint16_t)(ROUTER_SVC_PAYLOAD_OFF + 1 + len));
if (!conn || !e) { if (e) { queue_dgram_free(e); queue_entry_free(e); } failure(t, "control fixture unavailable"); return; }
e->len = (uint16_t)(ROUTER_SVC_PAYLOAD_OFF + 1 + len);
memset(e->dgram, 0, e->len);
e->dgram[0] = ETCP_RT_ID_DM_MEDIA;
memcpy(e->dgram + ROUTER_SVC_SRC_OFF, &source, 8); memcpy(e->dgram + ROUTER_SVC_DST_OFF, &inst->node_id, 8);
uint64_t group = GROUP; memcpy(e->dgram + ROUTER_SVC_GROUP_OFF, &group, 8);
e->dgram[ROUTER_SVC_FLAGS_OFF] = ROUTER_FLAG_SIGNED; e->dgram[ROUTER_SVC_PAYLOAD_OFF] = cmd;
memcpy(e->dgram + ROUTER_SVC_PAYLOAD_OFF + 1, p, len);
inst->router_bindings.callbacks[ETCP_RT_ID_DM_MEDIA](conn, e);
}
static void cancelled(void* arg, int error) {
int* calls = arg;
*calls += error == -2 ? 1 : 100;
}
/* Ошибки регистрации не расходуют seq; квота проверяется до получения байтов. */
static int protocol_checks(struct media_test* t) {
struct UTUN_INSTANCE* a = t->nodes[A];
struct UTUN_INSTANCE* m = t->nodes[M];
const uint8_t* desc = dm_message_media(t->body, t->body_len);
uint8_t input[DM_MEDIA_DESC_SIZE + ATTACHMENT_WIRE_MAX];
struct attachment_info info = { .kind = ATTACHMENT_FILE, .name = "rollback.bin" };
int metadata_len = attachment_encode(&info, input + DM_MEDIA_DESC_SIZE, ATTACHMENT_WIRE_MAX);
if (metadata_len < 0) return -1;
memcpy(input, desc, DM_MEDIA_DESC_SIZE);
if (sqlite3_exec(a->topo_sqlite_db, "CREATE TEMP TRIGGER reject_dm_file BEFORE INSERT ON dm_files"
" BEGIN SELECT RAISE(ABORT,'injected media registration failure'); END", NULL, NULL, NULL)) return -1;
int rejected = dm_send(a, t->conv, "file", input, DM_MEDIA_DESC_SIZE + metadata_len) < 0;
if (sqlite3_exec(a->topo_sqlite_db, "DROP TRIGGER reject_dm_file", NULL, NULL, NULL)) return -1;
if (!rejected || value(a, "SELECT last_out_seq FROM dm_conversations") != 1 ||
value(a, "SELECT count(*) FROM dm_messages") != 1 || value(a, "SELECT count(*) FROM dm_outbox") != 0) return -1;
char json[8192]; size_t len;
if (dm_list_messages_json(a, t->conv, 200, 0, json, sizeof(json), &len) ||
!strstr(json, "\"media_status\":\"delivered\"") || !strstr(json, "\"media_path\":") ||
!strstr(json, "\"content_type\":\"audio/opus\"") || !strstr(json, "\"duration_ms\":3000")) return -1;
uint8_t request[8 + sizeof(t->body)];
memcpy(request, &t->ids[B], 8); memcpy(request + 8, t->body, t->body_len);
uint8_t* metadata = (uint8_t*)dm_message_media(request + 8, t->body_len);
metadata[0] ^= 0x80;
uint64_t size = 2 * 1024 * 1024;
memcpy(metadata + 16, &size, 8);
if (sc_ed25519_sign(a->my_ed25519_privkey, request + 8, t->body_len - DM_MSG_SIG_SIZE,
request + 8 + t->body_len - DM_MSG_SIG_SIZE) ||
dm_verify_message(m, t->ids[B], request + 8, t->body_len) || chat_setting_set(m, "dm_media_storage_mb", "1")) return -1;
control(t, M, t->ids[S], DM_MEDIA_STORE, request, 8 + t->body_len);
if (value(m, "SELECT count(*) FROM dm_custody") != 0 || chat_setting_set(m, "dm_media_storage_mb", "1024")) return -1;
char path[768]; snprintf(path, sizeof(path), "%s/cancelled.download", t->dir);
int calls = 0;
struct file_transfer_op* op = file_transfer_pull(t->nodes[B]->file_transfer, GROUP, t->ids[A], t->id,
dm_media_cipher_size(FILE_BYTES), path, cancelled, &calls);
if (!op) return -1;
file_transfer_cancel_object(t->nodes[B]->file_transfer, t->id);
if (calls != 1) return -1;
DEBUG_INFO(DEBUG_CATEGORY_DM, "media protocol checks passed: SQL rollback, quota, JSON state, CM operation cancellation");
return 0;
}
/* Таймер служит heartbeat: большие криптозадачи не должны останавливать сценарий. */
/* Одни и те же transfer/custody гарантии для разных типов вложений. */
static int send_fixture(struct media_test* t, enum attachment_kind kind, const char* name) {
struct attachment_info info = { .kind = kind };
snprintf(info.name, sizeof(info.name), "%s", name);
if (kind == ATTACHMENT_VOICE) { info.duration_ms = 3000; memset(info.waveform, 'A', 100); }
if (kind == ATTACHMENT_VIDEO) { info.duration_ms = 3000; info.width = 640; info.height = 480; }
return dm_send_file(t->nodes[A], t->conv, t->file, &info, 0);
}
static void tick(void* arg) {
struct media_test* t = arg;
t->timer = NULL; t->ticks++; t->heartbeats++;
struct UTUN_INSTANCE* a = t->nodes[A];
struct UTUN_INSTANCE* b = t->nodes[B];
struct UTUN_INSTANCE* s = t->nodes[S];
struct UTUN_INSTANCE* m = t->nodes[M];
switch (t->phase) {
case 0:
if (instance_find_conn(a, t->ids[S]) && instance_find_conn(b, t->ids[S]) && instance_find_conn(m, t->ids[S])) {
if (setup(t)) { failure(t, "signed fixture setup"); return; }
t->phase++;
}
break;
case 1:
if (live(a, t->ids[B]) && live(m, t->ids[A])) {
if (send_fixture(t, ATTACHMENT_VOICE, "direct.opus")) { failure(t, "direct voice enqueue"); return; }
t->phase++;
}
break;
case 2:
if (value(a, "SELECT count(*) FROM dm_files WHERE dir=1 AND receipt IS NOT NULL") == 1) {
if (snapshot(t, 1) || !file_correct(t) || value(s, "SELECT count(*) FROM dm_media_jobs") != 0 ||
value(m, "SELECT count(*) FROM dm_custody") != 0) { failure(t, "direct transfer/custody isolation"); return; }
uint8_t bad[DM_MEDIA_RECEIPT_SIZE]; memcpy(bad, t->receipt, sizeof(bad)); bad[56] ^= 1;
if (!dm_media_verify_receipt(a, bad, t->body, t->body_len)) { failure(t, "forged media receipt accepted"); return; }
if (protocol_checks(t)) { failure(t, "media protocol invariants"); return; }
utun_instance_destroy(b); t->nodes[B] = NULL; t->phase++;
}
break;
case 3:
if (!live(a, t->ids[B])) {
if (send_fixture(t, ATTACHMENT_VIDEO, "offline.mp4")) { failure(t, "offline video enqueue"); return; }
t->phase++;
}
break;
case 4:
if (value(s, "SELECT count(*) FROM dm_media_jobs WHERE stored=1") == 1 && value(m, "SELECT count(*) FROM dm_custody WHERE ready=1") == 1) {
if (snapshot(t, 2)) { failure(t, "offline body snapshot"); return; }
utun_instance_destroy(a); t->nodes[A] = NULL;
utun_instance_destroy(m); t->nodes[M] = NULL;
utun_instance_destroy(s); t->nodes[S] = NULL;
int base = 48000 + getpid() % 5000;
if (start_node(t, S, base) || start_node(t, M, base) || start_node(t, B, base)) { failure(t, "custody restart"); return; }
t->phase++;
}
break;
case 5:
if (file_correct(t) && value(m, "SELECT count(*) FROM dm_custody") == 0 &&
value(s, "SELECT count(*) FROM dm_media_jobs WHERE receipt IS NOT NULL AND body IS NULL") == 1) {
uint8_t stored[24]; memcpy(stored, t->id, 16);
uint64_t expiry = ntp_time_get_seconds(s) + 100; memcpy(stored + 16, &expiry, 8);
control(t, S, t->ids[M], DM_MEDIA_STORED, stored, sizeof(stored));
uint8_t put[8 + sizeof(t->body)]; memcpy(put, &t->ids[B], 8); memcpy(put + 8, t->body, t->body_len);
control(t, M, t->ids[S], DM_MEDIA_STORE, put, 8 + t->body_len);
if (value(m, "SELECT count(*) FROM dm_custody") != 0 ||
value(s, "SELECT count(*) FROM dm_media_jobs WHERE receipt IS NULL") != 0) { failure(t, "late custody resurrected artifact"); return; }
if (start_node(t, A, 48000 + getpid() % 5000)) { failure(t, "source restart"); return; }
t->phase++;
}
break;
case 6:
if (live(a, t->ids[B]) && value(a, "SELECT count(*) FROM dm_files WHERE receipt IS NOT NULL") == 2) {
utun_instance_destroy(b); t->nodes[B] = NULL;
if (chat_setting_set(m, "dm_media_ttl_sec", "1")) { failure(t, "TTL setting"); return; }
t->phase++;
}
break;
case 7:
if (!live(a, t->ids[B])) {
if (send_fixture(t, ATTACHMENT_FILE, "expires.bin")) { failure(t, "TTL file enqueue"); return; }
t->phase++;
}
break;
case 8:
if (value(m, "SELECT count(*) FROM dm_custody_done WHERE receipt IS NULL") == 1) {
if (value(m, "SELECT count(*) FROM dm_custody") != 0 || value(s, "SELECT count(*) FROM dm_media_jobs WHERE expired=1") != 1) break;
if (snapshot(t, 3)) { failure(t, "expired body snapshot"); return; }
uint8_t put[8 + sizeof(t->body)]; memcpy(put, &t->ids[B], 8); memcpy(put + 8, t->body, t->body_len);
control(t, M, t->ids[S], DM_MEDIA_STORE, put, 8 + t->body_len);
if (value(m, "SELECT count(*) FROM dm_custody") != 0) { failure(t, "late STORE extended TTL"); return; }
t->phase++;
uasync_stop(t->ua); return;
}
break;
}
if (t->failed) return;
if (t->ticks > 6000) { failure(t, "timeout"); return; }
t->timer = uasync_set_timeout(t->ua, 100, t, tick, "dm_media_test");
}
int main(void) {
struct media_test t = {0};
debug_config_init(); debug_set_level(DEBUG_LEVEL_ERROR);
if (getenv("UTUN_TEST_DEBUG")) {
debug_set_category_level_by_name("dm", "debug"); debug_set_category_level_by_name("media", "debug");
debug_set_category_level_by_name("bgp", "info"); debug_set_category_level_by_name("etcp_route", "info");
}
utun_instance_set_tun_init_enabled(0);
snprintf(t.dir, sizeof(t.dir), "/tmp/utun_dm_media_XXXXXX");
if (test_mkdtemp(t.dir)) return 1;
for (int i = 0; i < NODES; i++) {
if (sc_generate_keypair(&t.keys[i]) || sc_derive_ed25519_pubkey(t.keys[i].private_key, t.ed[i])) return 1;
t.ids[i] = sc_derive_node_id_from_pubkey(t.keys[i].public_key);
}
struct SC_MYKEYS channel;
if (sc_generate_keypair(&channel)) return 1;
memcpy(t.ch_x, channel.public_key, 32);
if (sc_derive_ed25519_pubkey(channel.private_key, t.ch_ed)) return 1;
uint8_t channel_seed[SHA512_DIGEST_LENGTH];
SHA512(channel.private_key, 32, channel_seed);
memcpy(t.ch_priv, channel_seed, 32);
uint8_t signed_channel[256];
size_t off = sizeof(CHANNEL);
memcpy(signed_channel, CHANNEL, off);
const char name[] = "dm-media-test";
memcpy(signed_channel + off, name, sizeof(name)); off += sizeof(name);
memcpy(signed_channel + off, &t.ids[A], 8); off += 8;
memcpy(signed_channel + off, t.ch_x, 32); off += 32; memcpy(signed_channel + off, t.ch_ed, 32); off += 32;
if (sc_ed25519_sign(t.ch_priv, signed_channel, off, t.ch_sig)) return 1;
snprintf(t.conv, sizeof(t.conv), "%llu", (unsigned long long)dm_derive_conv_id(t.ids[A], t.ids[B]));
snprintf(t.file, sizeof(t.file), "%s/source.bin", t.dir);
FILE* f = fopen(t.file, "wb");
if (!f) return 1;
for (unsigned i = 0; i < FILE_BYTES; i++) fputc((int)((i * 73 + 19) & 255), f);
fclose(f);
t.ua = uasync_create();
int base = 48000 + getpid() % 5000;
if (!t.ua || start_node(&t, S, base) || start_node(&t, A, base) || start_node(&t, M, base) || start_node(&t, B, base)) return 1;
t.timer = uasync_set_timeout(t.ua, 100, &t, tick, "dm_media_test");
uasync_mainloop(t.ua);
for (int i = 0; i < NODES; i++) if (t.nodes[i]) utun_instance_destroy(t.nodes[i]);
uasync_destroy(t.ua, 0);
printf("test_dm_media: phase=%u heartbeats=%u %s artifacts=%s\n", t.phase, t.heartbeats,
!t.failed && t.phase == 9 ? "PASS" : "FAIL", t.dir);
return t.failed || t.phase != 9;
}