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.
 
 
 
 
 
 

173 lines
9.9 KiB

#include <assert.h>
#include <string.h>
#include <stdlib.h>
#include <unistd.h>
#include "utun_instance.h"
#include "config_parser.h"
#include "topo_group.h"
#include "topo_node_sqlite.h"
#include "route_lib.h"
#include "etcp.h"
#include "etcp_router.h"
#include "media_delivery/media_index.h"
#include "media_delivery/media_delivery_proto.h"
#include "node_conn_direct.h"
#include "chat/chat_core.h"
#include "chat/db_sync.h"
#include "media_async/media_async.h"
#include "../lib/debug_config.h"
struct job { int worked, done, error; };
static void work(void* arg) { ((struct job*)arg)->worked = 1; }
static void done(void* arg, int err) { struct job* job = arg; job->done++; job->error = err; }
static void ready(struct ll_queue* queue, void* arg) { (void)queue; (*(int*)arg)++; }
/* Pending transit shares the physical connection but belongs to its group. */
static void pending_transit(struct ETCP_CONN* conn, uint64_t group_id, int* calls) {
if (!conn->transit_queues) conn->transit_queues = queue_new(conn->instance->ua, 16, 0, 24, "test_transit");
assert(conn->transit_queues);
struct TRANSIT_QUEUE* transit = (struct TRANSIT_QUEUE*)queue_entry_new(sizeof(*transit) - sizeof(struct ll_entry));
assert(transit); transit->group_id = group_id; transit->conn = conn;
transit->src_node_id = 123; transit->dst_node_id = conn->peer_node_id;
transit->q = queue_new(conn->instance->ua, 0, 0, 0, "test_transit_packets"); assert(transit->q);
assert(queue_data_put_with_index(conn->transit_queues, &transit->ll) == 0);
struct ll_entry* packet = ll_alloc_lldgram(1); assert(packet); packet->dgram[0] = 0; packet->len = 1;
assert(queue_data_put(transit->q, packet) == 0);
assert(queue_waiter_wait(conn->send_input_q, &transit->waiter, ready, calls) == 0);
}
static struct TOPO_GROUP* chat_group(struct UTUN_INSTANCE* inst) {
for (struct ll_entry* e = inst->topo_groups->group_list->head; e; e = e->next) {
struct TOPO_GROUP* group = (struct TOPO_GROUP*)e;
if (group->group_type == TOPO_GROUP_TYPE_CHAT) return group;
}
return NULL;
}
int main(void) {
debug_config_init(); debug_set_level(DEBUG_LEVEL_INFO);
struct UASYNC* ua = uasync_create(); assert(ua);
struct UTUN_INSTANCE* inst = utun_instance_create_from_str(ua,
"[global]\nmy_private_key=704f2e012c8fa8768130cb0f988a997dccb628372bc5ceccacc78dcbfec5916f\n"
"my_public_key=b3193173def895bd0fcea6f86af077c7d77216f10395275f627ac18242ec0f01\n"
"[routing]\nmy_subnet=192.168.42.0/24\n[server: udp]\naddr=127.0.0.1:0\ntype=public\n"
"[chatserver]\ngroup_autoconnect=0\nstorage_autoload=0\n");
assert(inst && !topo_groups_get_default(inst->topo_groups) && !inst->chat_core);
inst->config->global.db_sync_enabled = 0;
assert(utun_core_start(inst) == 0 && utun_core_start(inst) == 0);
struct ETCP_SOCKET* socket = inst->etcp_sockets;
sqlite3* db = inst->topo_sqlite_db;
struct TOPO_NODE* identity = topo_node_registry_find(inst->topo_groups, inst->node_id);
assert(identity && identity->v4_addrs);
assert(chat_service_start(inst) == 0 && chat_service_start(inst) == 0);
assert(!topo_groups_get_default(inst->topo_groups));
chat_core_create_channel_auto(inst, "lifecycle");
struct TOPO_GROUP* chat = chat_group(inst); assert(chat);
uint64_t chat_id = chat->group_id;
assert(utun_service_start(inst) == 0 && utun_service_start(inst) == 0);
assert(inst->rt->count > 0);
struct TOPO_GROUP* vpn = topo_groups_get_default(inst->topo_groups); assert(vpn);
struct SC_MYKEYS keys; assert(sc_generate_keypair(&keys) == SC_OK);
struct TOPO_ADDR4 address = { .addr = {127, 0, 0, 1}, .port = 9, .protocol = TOPO_PROTO_UDP };
struct TOPO_NODE node = { .v4_addrs = &address };
memcpy(node.public_key, keys.public_key, 32); node.node_id = sc_derive_node_id_from_pubkey(keys.public_key);
uint8_t sig[64] = {1};
assert(topo_node_sqlite_member_block_put(db, chat->channel_id, node.node_id, sig, 1, sig, 1,
keys.public_key, keys.public_key, "{}", NULL, 0, sig) == 0);
struct NODE_CONN_DIRECT* seed = NULL;
assert(node_conn_direct_open_node(inst, node.node_id, NULL, NULL, &seed, &node, NULL) >= 0);
struct ETCP_CONN* conn = node_conn_direct_get_conn(seed); assert(conn);
struct TOPO_PEER_REQUEST *chat_request = NULL, *vpn_request = NULL;
assert(topo_group_peer_open(chat, node.node_id, &chat_request) == 0);
assert(topo_group_peer_open(vpn, node.node_id, &vpn_request) == 0);
node_conn_direct_close(seed);
utun_service_stop(inst); utun_service_stop(inst);
assert(!conn->close_requested && topo_group_peer_conn(chat_request) == conn);
assert(topo_group_peer_phase(vpn_request) == TOPO_PEER_FAILED && inst->rt->count == 0);
topo_group_peer_close(vpn_request);
assert(inst->etcp_sockets == socket && inst->topo_sqlite_db == db && inst->chat_started);
assert(topo_node_registry_find(inst->topo_groups, inst->node_id) == identity);
assert(topo_node_update_my_addresses(inst) >= 0);
assert(utun_service_start(inst) == 0);
vpn = topo_groups_get_default(inst->topo_groups);
assert(topo_group_peer_open(vpn, node.node_id, &vpn_request) == 0);
/* Stop a real stream while its router send queue is blocked. */
char path[] = "/tmp/utun-service-stream-XXXXXX";
int fd = mkstemp(path); assert(fd >= 0);
char bytes[8192] = {0}; assert(write(fd, bytes, sizeof(bytes)) == sizeof(bytes)); close(fd);
chat_core_save_ui_state(inst, "media_base", "/tmp");
uint8_t media_id[16] = {1}, block_id[16] = {2}, hash[32] = {3};
assert(media_index_register_downloaded(db, media_id, block_id, hash, chat->channel_id,
strrchr(path, '/') + 1, inst->node_id, sizeof(bytes), sizeof(bytes), 0, 0) == 0);
struct ETCP_ROUTER_CONN* routed = etcp_router_conn_get(inst, chat_id, node.node_id, ETCP_RT_ID_MEDIA_DELIVERY);
assert(routed);
struct queue_waiter_handle waiter = {0}; int ready_calls = 0;
struct ll_queue* old_queue = routed->send_q;
queue_set_threshold(old_queue, -1, 0);
etcp_router_on_send_ready(inst, chat_id, node.node_id, ETCP_RT_ID_MEDIA_DELIVERY, &waiter, ready, &ready_calls);
assert(waiter.internal);
etcp_router_conn_restart(inst, chat_id, node.node_id, ETCP_RT_ID_MEDIA_DELIVERY);
assert(routed->send_q != old_queue);
queue_set_threshold(routed->send_q, -1, 0);
etcp_router_on_send_ready(inst, chat_id, node.node_id, ETCP_RT_ID_MEDIA_DELIVERY, &waiter, ready, &ready_calls);
assert(!old_queue->waiter_head && waiter.internal);
etcp_router_cancel_send_ready(inst, chat_id, node.node_id, ETCP_RT_ID_MEDIA_DELIVERY, &waiter);
assert(!routed->send_q->waiter_head && !waiter.internal);
queue_set_threshold(routed->send_q, 0, 0);
routed->send_q->waiter_defer = 1;
etcp_router_on_send_ready(inst, chat_id, node.node_id, ETCP_RT_ID_MEDIA_DELIVERY, &waiter, ready, &ready_calls);
assert(waiter.call_soon_id);
etcp_router_cancel_send_ready(inst, chat_id, node.node_id, ETCP_RT_ID_MEDIA_DELIVERY, &waiter);
struct media_pkt_block_req request = { .subcmd = MEDIA_SUBCMD_BLOCK_REQ, .group_id = chat_id };
memcpy(request.media_id, media_id, 16); memcpy(request.block_id, block_id, 16);
struct ll_entry* packet = ll_alloc_lldgram(ROUTER_SVC_PAYLOAD_OFF + sizeof(request)); assert(packet);
memset(packet->dgram, 0, ROUTER_SVC_PAYLOAD_OFF);
packet->dgram[0] = ETCP_RT_ID_MEDIA_DELIVERY;
memcpy(packet->dgram + ROUTER_SVC_SRC_OFF, &node.node_id, 8);
memcpy(packet->dgram + ROUTER_SVC_GROUP_OFF, &chat_id, 8);
memcpy(packet->dgram + ROUTER_SVC_PAYLOAD_OFF, &request, sizeof(request));
packet->len = ROUTER_SVC_PAYLOAD_OFF + sizeof(request);
inst->router_bindings.callbacks[ETCP_RT_ID_MEDIA_DELIVERY](conn, packet);
assert(inst->md.streams && inst->md.active_streams == 1);
struct job cancelled = {0};
media_async_submit(inst->media_async, ua, work, &cancelled, done, &cancelled);
queue_set_threshold(conn->send_input_q, -1, 0);
pending_transit(conn, chat_id, &ready_calls);
pending_transit(conn, vpn->group_id, &ready_calls);
chat_service_stop(inst); chat_service_stop(inst);
assert(queue_entry_count(conn->transit_queues) == 1);
assert(((struct TRANSIT_QUEUE*)conn->transit_queues->head)->group_id == vpn->group_id);
assert(!inst->md.streams && routed->closed); unlink(path);
assert(cancelled.worked && cancelled.done == 1 && cancelled.error == MEDIA_ASYNC_CANCELLED);
assert(!chat_group(inst) && !inst->chat_core && !inst->chat_sync && !inst->media_async && !inst->md.initialized);
assert(!db_sync_instance_find(inst, chat_id));
assert(!conn->close_requested && topo_group_peer_conn(vpn_request) == conn);
assert(topo_group_peer_phase(chat_request) == TOPO_PEER_FAILED);
topo_group_peer_close(chat_request);
/* Remove the other group's pending transit before processing the event loop. */
etcp_router_close_group(inst, vpn->group_id);
assert(queue_entry_count(conn->transit_queues) == 0);
queue_set_threshold(conn->send_input_q, 0, 0);
for (int i = 0; i < 3; i++) {
assert(chat_service_start(inst) == 0 && chat_group(inst)->group_id == chat_id);
assert(db_sync_instance_find(inst, chat_id));
chat_service_stop(inst);
assert(inst->etcp_sockets == socket && inst->topo_sqlite_db == db && inst->utun_started);
uasync_poll(ua, 0);
}
topo_group_peer_close(vpn_request);
utun_service_stop(inst);
struct job completed = {0};
struct media_async* jobs = media_async_create(); assert(jobs);
media_async_submit(jobs, ua, work, &completed, done, &completed);
while (!completed.done) uasync_poll(ua, 100);
assert(completed.worked && completed.done == 1 && !completed.error);
media_async_destroy(jobs);
utun_instance_destroy(inst); uasync_poll(ua, 0);
assert(!ready_calls);
assert(ua->timer_alloc_count == ua->timer_free_count && ua->socket_alloc_count == ua->socket_free_count);
uasync_destroy(ua, 0);
return 0;
}