22 changed files with 1449 additions and 40 deletions
@ -0,0 +1,190 @@ |
|||||||
|
#ifndef CONN_MGR_H |
||||||
|
#define CONN_MGR_H |
||||||
|
|
||||||
|
#include <stdint.h> |
||||||
|
#include <stddef.h> |
||||||
|
#include "../lib/ll_queue.h" |
||||||
|
#include "etcp_api.h" |
||||||
|
#include "etcp_router.h" |
||||||
|
#include "route_node.h" |
||||||
|
|
||||||
|
struct UTUN_INSTANCE; |
||||||
|
struct ETCP_CONN; |
||||||
|
struct NODEINFO_Q; |
||||||
|
struct ROUTE_BGP; |
||||||
|
|
||||||
|
#define CONN_MGR_MAX_CANDIDATES 3 |
||||||
|
#define CONN_MGR_CANDIDATE_CACHE_MS 20000 |
||||||
|
#define CONN_MGR_CANDIDATE_STALE_TB 300000 |
||||||
|
#define CONN_MGR_CANDIDATE_PING_TB 20000 |
||||||
|
#define CONN_MGR_BG_PING_INTERVAL_TB 1000 |
||||||
|
#define CONN_MGR_BG_PING_CYCLE_MIN_TB 100000 |
||||||
|
#define CONN_MGR_IDLE_CHECK_INTERVAL_TB 10000 |
||||||
|
#define CONN_MGR_LOCAL_SCAN_ATTEMPTS 3 |
||||||
|
#define CONN_MGR_LOCAL_SCAN_TIMEOUT_MS 100 |
||||||
|
#define CONN_MGR_CONNECT_DIRECT_TIMEOUT_MS 5000 |
||||||
|
#define CONN_MGR_CONNECT_REVERSE_TIMEOUT_MS 15000 |
||||||
|
#define CONN_MGR_INTERM_EXCHANGE_TIMEOUT_MS 15000 |
||||||
|
|
||||||
|
#define CONN_MGR_OK 0 |
||||||
|
#define CONN_MGR_ERR_NOT_FOUND -1 |
||||||
|
#define CONN_MGR_ERR_NO_ADDRESSES -2 |
||||||
|
#define CONN_MGR_ERR_TIMEOUT -3 |
||||||
|
#define CONN_MGR_ERR_UNREACHABLE -4 |
||||||
|
#define CONN_MGR_ERR_REFUSED -5 |
||||||
|
#define CONN_MGR_ERR_ALREADY_CONNECTED -6 |
||||||
|
#define CONN_MGR_ERR_INTERNAL -7 |
||||||
|
|
||||||
|
enum CONN_MGR_STATE { |
||||||
|
CONN_MGR_STATE_DISCONNECTED = 0, |
||||||
|
CONN_MGR_STATE_CONNECTING = 1, |
||||||
|
CONN_MGR_STATE_CONNECTED = 2, |
||||||
|
}; |
||||||
|
|
||||||
|
enum CM_TRY_STATE { |
||||||
|
CM_TRY_NONE = 0, |
||||||
|
CM_TRY_PENDING = 1, |
||||||
|
CM_TRY_OK = 2, |
||||||
|
CM_TRY_FAILED = 3, |
||||||
|
}; |
||||||
|
|
||||||
|
#define CONN_MGR_SUBCMD_DIRECT_REQ 0x01 |
||||||
|
#define CONN_MGR_SUBCMD_DIRECT_RESP 0x02 |
||||||
|
#define CONN_MGR_SUBCMD_INTERM_EXCHANGE_REQ 0x03 |
||||||
|
#define CONN_MGR_SUBCMD_INTERM_EXCHANGE_RESP 0x04 |
||||||
|
#define CONN_MGR_SUBCMD_INTERM_SELECTED 0x05 |
||||||
|
#define CONN_MGR_SUBCMD_DISCONNECT 0x06 |
||||||
|
|
||||||
|
struct CONN_MGR_CANDIDATE { |
||||||
|
uint64_t node_id; |
||||||
|
uint16_t rtt; |
||||||
|
} __attribute__((packed)); |
||||||
|
|
||||||
|
#pragma pack(push, 1) |
||||||
|
struct CONN_MGR_DIRECT_REQ { |
||||||
|
uint8_t cmd; |
||||||
|
uint8_t subcmd; |
||||||
|
uint32_t request_id; |
||||||
|
uint8_t addr_count; |
||||||
|
}; |
||||||
|
|
||||||
|
struct CONN_MGR_DIRECT_RESP { |
||||||
|
uint8_t cmd; |
||||||
|
uint8_t subcmd; |
||||||
|
uint32_t request_id; |
||||||
|
uint8_t accepted; |
||||||
|
}; |
||||||
|
|
||||||
|
struct CONN_MGR_INTERM_EXCHANGE_REQ { |
||||||
|
uint8_t cmd; |
||||||
|
uint8_t subcmd; |
||||||
|
uint32_t request_id; |
||||||
|
uint8_t candidate_count; |
||||||
|
struct CONN_MGR_CANDIDATE candidates[4]; |
||||||
|
}; |
||||||
|
|
||||||
|
struct CONN_MGR_INTERM_EXCHANGE_RESP { |
||||||
|
uint8_t cmd; |
||||||
|
uint8_t subcmd; |
||||||
|
uint32_t request_id; |
||||||
|
uint8_t my_count; |
||||||
|
uint8_t your_count; |
||||||
|
struct CONN_MGR_CANDIDATE my_candidates[4]; |
||||||
|
struct CONN_MGR_CANDIDATE your_candidates[4]; |
||||||
|
}; |
||||||
|
|
||||||
|
struct CONN_MGR_INTERM_SELECTED { |
||||||
|
uint8_t cmd; |
||||||
|
uint8_t subcmd; |
||||||
|
uint32_t request_id; |
||||||
|
uint8_t count; |
||||||
|
struct CONN_MGR_CANDIDATE selected[3]; |
||||||
|
}; |
||||||
|
|
||||||
|
struct CONN_MGR_DISCONNECT { |
||||||
|
uint8_t cmd; |
||||||
|
uint8_t subcmd; |
||||||
|
uint64_t node_id; |
||||||
|
}; |
||||||
|
#pragma pack(pop) |
||||||
|
|
||||||
|
#define CONN_MGR_DIRECT_REQ_HDR_SIZE sizeof(struct CONN_MGR_DIRECT_REQ) |
||||||
|
#define CONN_MGR_DIRECT_RESP_HDR_SIZE sizeof(struct CONN_MGR_DIRECT_RESP) |
||||||
|
#define CONN_MGR_INTERM_EXCHANGE_REQ_SIZE sizeof(struct CONN_MGR_INTERM_EXCHANGE_REQ) |
||||||
|
#define CONN_MGR_INTERM_EXCHANGE_RESP_SIZE sizeof(struct CONN_MGR_INTERM_EXCHANGE_RESP) |
||||||
|
#define CONN_MGR_INTERM_SELECTED_SIZE sizeof(struct CONN_MGR_INTERM_SELECTED) |
||||||
|
#define CONN_MGR_DISCONNECT_SIZE sizeof(struct CONN_MGR_DISCONNECT) |
||||||
|
|
||||||
|
typedef void (*conn_mgr_connect_callback_t)(int result, uint64_t node_id, void* arg); |
||||||
|
|
||||||
|
struct cm_cb_node { |
||||||
|
conn_mgr_connect_callback_t cb; |
||||||
|
void* arg; |
||||||
|
struct cm_cb_node* next; |
||||||
|
}; |
||||||
|
|
||||||
|
struct CONN_MGR_ENTRY { |
||||||
|
uint64_t node_id; |
||||||
|
uint8_t state; |
||||||
|
uint8_t conn_type; |
||||||
|
uint8_t alien; |
||||||
|
uint32_t idle_timeout_ms; |
||||||
|
uint64_t last_traffic_tb; |
||||||
|
void* idle_timer; |
||||||
|
|
||||||
|
struct cm_cb_node* cb_list; |
||||||
|
|
||||||
|
uint64_t intermediaries[CONN_MGR_MAX_INTERMEDIARIES]; |
||||||
|
uint8_t intermediariy_count; |
||||||
|
uint8_t rr_idx; |
||||||
|
|
||||||
|
uint8_t local_scan_state; |
||||||
|
uint8_t main_connect_state; |
||||||
|
struct { |
||||||
|
uint8_t phase; |
||||||
|
void* timer; |
||||||
|
uint32_t request_id; |
||||||
|
void* conn_ctx; |
||||||
|
} main; |
||||||
|
|
||||||
|
struct CONN_MGR* mgr; |
||||||
|
}; |
||||||
|
|
||||||
|
struct CONN_MGR { |
||||||
|
struct UTUN_INSTANCE* instance; |
||||||
|
struct CONN_MGR_ENTRY* entries; |
||||||
|
size_t entry_count; |
||||||
|
size_t entry_capacity; |
||||||
|
|
||||||
|
void* bg_ping_timer; |
||||||
|
size_t bg_ping_cursor; |
||||||
|
uint64_t bg_ping_cycle_start_tb; |
||||||
|
|
||||||
|
void* candidate_ping_timer; |
||||||
|
|
||||||
|
struct CONN_MGR_CANDIDATE best_candidates[CONN_MGR_MAX_CANDIDATES]; |
||||||
|
uint8_t best_candidate_count; |
||||||
|
|
||||||
|
uint32_t next_request_id; |
||||||
|
uint8_t initialized; |
||||||
|
|
||||||
|
struct cm_reverse_pending* reverse_pending; |
||||||
|
struct cm_exchange_pending* exchange_pending; |
||||||
|
}; |
||||||
|
|
||||||
|
struct CONN_MGR* conn_mgr_init(struct UTUN_INSTANCE* instance); |
||||||
|
void conn_mgr_destroy(struct CONN_MGR* mgr); |
||||||
|
|
||||||
|
int conn_mgr_connect_node(struct CONN_MGR* mgr, uint64_t node_id, |
||||||
|
uint32_t idle_timeout_ms, |
||||||
|
conn_mgr_connect_callback_t cb, void* cb_arg); |
||||||
|
int conn_mgr_disconnect_node(struct CONN_MGR* mgr, uint64_t node_id); |
||||||
|
int conn_mgr_set_idle_timeout(struct CONN_MGR* mgr, uint64_t node_id, uint32_t timeout_ms); |
||||||
|
int conn_mgr_get_status(struct CONN_MGR* mgr, uint64_t node_id, |
||||||
|
uint8_t* out_state, uint8_t* out_conn_type); |
||||||
|
int conn_mgr_add_alien_node(struct CONN_MGR* mgr, const uint8_t* nodeinfo_data, size_t len); |
||||||
|
int conn_mgr_send(struct CONN_MGR* mgr, uint64_t node_id, struct ll_entry* entry); |
||||||
|
|
||||||
|
void conn_mgr_update_best_candidates(struct CONN_MGR* mgr, uint64_t node_id, uint16_t rtt); |
||||||
|
|
||||||
|
#endif |
||||||
@ -0,0 +1,136 @@ |
|||||||
|
/**
|
||||||
|
* @file test_conn_mgr.c |
||||||
|
* @brief Connection Manager — smoke test |
||||||
|
* |
||||||
|
* Test 1: Direct via existing ETCP conn → CONN_TYPE_DIRECT |
||||||
|
* Test 2: Status check → state/type |
||||||
|
* Test 3: Alien node → NODEINFO_Q.alien=1 |
||||||
|
* |
||||||
|
* Pending: |
||||||
|
* Test 4: New link UP callback (INIT to new peer via conn_mgr) |
||||||
|
* Test 5: Indirect through intermediary (needs full Phase 3 exchange) |
||||||
|
* Test 6: Idle timeout |
||||||
|
*/ |
||||||
|
#include <stdio.h> |
||||||
|
#include <stdlib.h> |
||||||
|
#include <string.h> |
||||||
|
#include <stdarg.h> |
||||||
|
#include "../lib/platform_compat.h" |
||||||
|
#include "test_utils.h" |
||||||
|
#ifndef _WIN32 |
||||||
|
#include <unistd.h> |
||||||
|
#endif |
||||||
|
|
||||||
|
#include "../src/etcp.h" |
||||||
|
#include "../src/etcp_connections.h" |
||||||
|
#include "../src/config_parser.h" |
||||||
|
#include "../src/config_updater.h" |
||||||
|
#include "../src/utun_instance.h" |
||||||
|
#include "../src/route_bgp.h" |
||||||
|
#include "../src/conn_mgr.h" |
||||||
|
#include "../lib/u_async.h" |
||||||
|
#include "../lib/debug_config.h" |
||||||
|
#include "../lib/mem.h" |
||||||
|
|
||||||
|
#define TIMEOUT_TB 300000 |
||||||
|
#define POLL_MS 5 |
||||||
|
#define NID_A 0xAAAA000000000001ULL |
||||||
|
#define NID_B 0xBBBB000000000002ULL |
||||||
|
|
||||||
|
static struct UTUN_INSTANCE* g_a = NULL, *g_b = NULL; |
||||||
|
static struct UASYNC* ua = NULL; |
||||||
|
static int result = 0; |
||||||
|
static void* ttimer = NULL; |
||||||
|
static char tdir[] = "/tmp/utun_cm_XXXXXX"; |
||||||
|
static char ca[256], cb[256]; |
||||||
|
static int pa = 0, pb = 0; |
||||||
|
static volatile int cdone = 0, cresult = 0; |
||||||
|
|
||||||
|
static int wf(const char* p, const char* f, ...) { |
||||||
|
va_list ap; FILE* fp = fopen(p, "w"); if (!fp) return -1; |
||||||
|
va_start(ap, f); vfprintf(fp, f, ap); va_end(ap); fclose(fp); return 0; |
||||||
|
} |
||||||
|
static char* gv(const char* p, const char* k) { |
||||||
|
struct utun_config* c = parse_config(p); if (!c) return NULL; |
||||||
|
char* r = (strcmp(k, "pub") == 0) ? u_strdup(c->global.my_public_key_hex) : u_strdup(c->global.my_private_key_hex); |
||||||
|
free_config(c); return r; |
||||||
|
} |
||||||
|
static int lks(struct UTUN_INSTANCE* i) { |
||||||
|
int n = 0; struct ETCP_CONN* c = i->connections; |
||||||
|
while (c) { struct ETCP_LINK* l = c->links; while (l) { if (l->initialized && l->link_status) n++; l = l->next; } c = c->next; } |
||||||
|
return n; |
||||||
|
} |
||||||
|
static void fail(const char* msg) { fprintf(stderr, "FAIL: %s\n", msg); fflush(stderr); result = 2; } |
||||||
|
static void test1(void* arg); |
||||||
|
static void test2(void* arg); |
||||||
|
static void test3(void* arg); |
||||||
|
|
||||||
|
static void to_cb(void* arg) { (void)arg; fprintf(stderr, "TIMEOUT\n"); result = 2; } |
||||||
|
static void ccb(int r, uint64_t id, void* arg) { |
||||||
|
(void)arg; fprintf(stderr, "connect_cb: result=%d node=0x%llx\n", r, (unsigned long long)id); fflush(stderr); |
||||||
|
cdone = 1; cresult = r; |
||||||
|
} |
||||||
|
|
||||||
|
static void test1(void* arg) { |
||||||
|
(void)arg; if (result) return; |
||||||
|
if (lks(g_a) < 1 || lks(g_b) < 1) { uasync_set_timeout(ua, 50, NULL, (timeout_callback_t)test1, "t1"); return; } |
||||||
|
fprintf(stderr, "Test 1: direct — conn_mgr_connect_node(B)\n"); fflush(stderr); |
||||||
|
conn_mgr_connect_node(g_a->conn_mgr, NID_B, 0, ccb, NULL); |
||||||
|
uasync_set_timeout(ua, 100, NULL, (timeout_callback_t)test2, "t2a"); |
||||||
|
} |
||||||
|
static void test2(void* arg) { |
||||||
|
(void)arg; if (result) return; |
||||||
|
if (!cdone) { uasync_set_timeout(ua, 100, NULL, (timeout_callback_t)test2, "t2b"); return; } |
||||||
|
if (cresult != CONN_MGR_OK) { fail("connect failed"); return; } |
||||||
|
uint8_t st, ty; conn_mgr_get_status(g_a->conn_mgr, NID_B, &st, &ty); |
||||||
|
fprintf(stderr, "Test 2: status state=%d type=%d\n", st, ty); fflush(stderr); |
||||||
|
if (st == CONN_MGR_STATE_CONNECTED && ty == CONN_TYPE_DIRECT) fprintf(stderr, " OK: DIRECT\n"); |
||||||
|
else { fail("status mismatch"); return; } |
||||||
|
test3(NULL); |
||||||
|
} |
||||||
|
static void test3(void* arg) { |
||||||
|
(void)arg; |
||||||
|
uint64_t aid = 0xEEEE000000000001ULL; |
||||||
|
struct NODEINFO ni; memset(&ni, 0, sizeof(ni)); ni.node_id = aid; ni.ver = 1; |
||||||
|
conn_mgr_add_alien_node(g_a->conn_mgr, (uint8_t*)&ni, sizeof(ni)); |
||||||
|
struct NODEINFO_Q* nq = nodeinfo_find_by_id(g_a->bgp, aid); |
||||||
|
fprintf(stderr, "Test 3: alien alien=%d\n", nq ? nq->alien : -1); fflush(stderr); |
||||||
|
if (nq && nq->alien == 1) fprintf(stderr, " OK: alien flag set\n"); |
||||||
|
else { fail("alien flag not set"); return; } |
||||||
|
fprintf(stderr, "=== ALL DONE ===\n"); fflush(stderr); |
||||||
|
result = (result == 0) ? 1 : 2; |
||||||
|
} |
||||||
|
|
||||||
|
static void setup(void) { |
||||||
|
test_mkdtemp(tdir); |
||||||
|
int base = 47000 + (getpid() % 15000); pa = base; pb = base + 1; |
||||||
|
snprintf(ca, sizeof(ca), "%s/a.conf", tdir); snprintf(cb, sizeof(cb), "%s/b.conf", tdir); |
||||||
|
wf(ca, "[global]\nmy_node_id=0x%llx\ntun_ip=10.99.0.1/24\ntun_ifname=tun99\n[server: s1]\naddr=127.0.0.1:%d\ntype=public\n[allowed_keys]\nallow_all=1\n", NID_A, pa); |
||||||
|
wf(cb, "[global]\nmy_node_id=0x%llx\ntun_ip=10.99.0.2/24\ntun_ifname=tun98\n[server: s1]\naddr=127.0.0.1:%d\ntype=public\n[allowed_keys]\nallow_all=1\n", NID_B, pb); |
||||||
|
config_ensure_keys_and_node_id(ca); config_ensure_keys_and_node_id(cb); |
||||||
|
char *p0 = gv(ca,"pub"), *r0 = gv(ca,"priv"), *p1 = gv(cb,"pub"), *r1 = gv(cb,"priv"); |
||||||
|
wf(ca, "[global]\nmy_node_id=0x%llx\nmy_private_key=%s\nmy_public_key=%s\ntun_ip=10.99.0.1/24\ntun_ifname=tun99\n[server: s1]\naddr=127.0.0.1:%d\ntype=public\n[client: to_b]\nkeepalive=1\npeer_public_key=%s\nlink=s1:127.0.0.1:%d\n[allowed_keys]\nallow_all=1\n", NID_A, r0, p0, pa, p1, pb); |
||||||
|
wf(cb, "[global]\nmy_node_id=0x%llx\nmy_private_key=%s\nmy_public_key=%s\ntun_ip=10.99.0.2/24\ntun_ifname=tun98\n[server: s1]\naddr=127.0.0.1:%d\ntype=public\n[allowed_keys]\nallow_all=1\n", NID_B, r1, p1, pb); |
||||||
|
u_free(p0); u_free(r0); u_free(p1); u_free(r1); |
||||||
|
} |
||||||
|
static void cleanup(void) { test_unlink(ca); test_unlink(cb); test_rmdir(tdir); } |
||||||
|
|
||||||
|
int main(void) { |
||||||
|
debug_config_init(); debug_set_level(DEBUG_LEVEL_WARN); |
||||||
|
utun_instance_set_tun_init_enabled(0); setup(); |
||||||
|
ua = uasync_create(); |
||||||
|
g_a = utun_instance_create(ua, ca); g_b = utun_instance_create(ua, cb); |
||||||
|
if (!g_a || !g_b) goto done; |
||||||
|
utun_instance_init(g_a); utun_instance_init(g_b); |
||||||
|
uasync_call_soon(ua, NULL, test1); |
||||||
|
ttimer = uasync_set_timeout(ua, TIMEOUT_TB, NULL, to_cb, "to"); |
||||||
|
{ int el = 0; while (!result && el < TIMEOUT_TB + 5000) { uasync_poll(ua, POLL_MS); el += POLL_MS; } } |
||||||
|
fprintf(stderr, "final result=%d\n", result); fflush(stderr); |
||||||
|
done: |
||||||
|
if (ttimer) uasync_cancel_timeout(ua, ttimer); |
||||||
|
if (g_a) { g_a->running = 0; utun_instance_destroy(g_a); } |
||||||
|
if (g_b) { g_b->running = 0; utun_instance_destroy(g_b); } |
||||||
|
if (ua) { uasync_destroy(ua, 0); ua = NULL; } |
||||||
|
cleanup(); |
||||||
|
return (result == 1) ? 0 : 1; |
||||||
|
} |
||||||
Loading…
Reference in new issue