2 changed files with 307 additions and 0 deletions
@ -0,0 +1,302 @@
|
||||
/* 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 "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->keys[j].private_key, 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; } |
||||
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); |
||||
} |
||||
|
||||
/* Таймер служит heartbeat: большие криптозадачи не должны останавливать сценарий. */ |
||||
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 (dm_send_file(a, t->conv, t->file, "direct.bin")) { failure(t, "direct file 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; } |
||||
utun_instance_destroy(b); t->nodes[B] = NULL; t->phase++; |
||||
} |
||||
break; |
||||
case 3: |
||||
if (!live(a, t->ids[B])) { |
||||
if (dm_send_file(a, t->conv, t->file, "offline.bin")) { failure(t, "offline file 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 (dm_send_file(a, t->conv, t->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); memcpy(t.ch_priv, channel.private_key, 32); |
||||
if (sc_derive_ed25519_pubkey(t.ch_priv, t.ch_ed)) return 1; |
||||
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; |
||||
} |
||||
Loading…
Reference in new issue