Browse Source
- etcp_router.h/c: new flag ROUTER_FLAG_ENCRYPTED (0x04), etcp_route_send_encrypted(), inflight/retrans/send_q support - utun_instance.h: e2e_ctx_cache[8] — per-peer sc_context_t cache with derived session key - topo_node.c: topo_node_sign_self() — single point for Ed25519 self-signing of NODEINFO - topo_node.c: fix stale x25519_self_sig after address update (topo_node_sign_self call) - nat_detection.c: fix stale x25519_self_sig after in-place socket meta modification - test_etcp_router_unit: 5 E2E tests (basic round-trip, no NODEINFO drop, tampered, cache reuse, ENCRYPTED+SIGNED combined) - builds clean, 70/70 tests passtopo_upd
44 changed files with 2586 additions and 145 deletions
@ -0,0 +1,259 @@
|
||||
#include <stdlib.h> |
||||
#include <string.h> |
||||
#include "../lib/platform_compat.h" |
||||
#include "../lib/debug_config.h" |
||||
#include "../lib/mem.h" |
||||
#include "../lib/u_async.h" |
||||
#include "etcp_api.h" |
||||
#include "etcp_connections.h" |
||||
#include "etcp.h" |
||||
#include "topo_node.h" |
||||
#include "topo_group.h" |
||||
#include "broadcast.h" |
||||
|
||||
struct broadcast_msg { |
||||
uint64_t local_time_tb; // offset 0
|
||||
uint8_t uuid[BROADCAST_UUID_SIZE]; // offset 8 — queue index key
|
||||
uint16_t data_len; // offset 24
|
||||
uint8_t data[]; // offset 26
|
||||
} __attribute__((packed)); |
||||
|
||||
#define BMSG_UUID_OFFSET 8 |
||||
|
||||
struct broadcast_ctx { |
||||
struct TOPO_GROUP* group; |
||||
struct ll_queue* messages; |
||||
struct broadcast_cbk_entry* cbks; |
||||
void* cleanup_timer; |
||||
uint32_t send_seq; |
||||
}; |
||||
|
||||
static void broadcast_cleanup_timer_cb(void* arg); |
||||
|
||||
int broadcast_init(struct TOPO_GROUP* group) { |
||||
if (!group) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "broadcast_init: null group"); return -1; } |
||||
if (group->broadcast) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "broadcast_init: already initialized grp=%016llx", (unsigned long long)group->group_id); return -1; } |
||||
|
||||
struct broadcast_ctx* ctx = u_calloc(1, sizeof(struct broadcast_ctx)); |
||||
if (!ctx) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "broadcast_init: alloc failed"); return -1; } |
||||
ctx->group = group; |
||||
ctx->send_seq = 0; |
||||
|
||||
ctx->messages = queue_new(group->instance->ua, 64, BMSG_UUID_OFFSET, BROADCAST_UUID_SIZE, "broadcast_msg"); |
||||
if (!ctx->messages) { u_free(ctx); DEBUG_ERROR(DEBUG_CATEGORY_BGP, "broadcast_init: queue_new failed"); return -1; } |
||||
|
||||
group->broadcast = ctx; |
||||
|
||||
ctx->cleanup_timer = uasync_set_timeout(group->instance->ua, BROADCAST_CLEANUP_INTERVAL_TB, ctx, broadcast_cleanup_timer_cb, "broadcast_cleanup_timer"); |
||||
if (!ctx->cleanup_timer) DEBUG_ERROR(DEBUG_CATEGORY_BGP, "broadcast_init: cleanup_timer failed grp=%016llx", (unsigned long long)group->group_id); |
||||
|
||||
DEBUG_INFO(DEBUG_CATEGORY_BGP, "broadcast_init ok grp=%016llx", (unsigned long long)group->group_id); |
||||
return 0; |
||||
} |
||||
|
||||
void broadcast_destroy(struct TOPO_GROUP* group) { |
||||
if (!group || !group->broadcast) return; |
||||
struct broadcast_ctx* ctx = group->broadcast; |
||||
DEBUG_INFO(DEBUG_CATEGORY_BGP, "broadcast_destroy grp=%016llx", (unsigned long long)group->group_id); |
||||
|
||||
if (ctx->cleanup_timer) { uasync_cancel_timeout(group->instance->ua, ctx->cleanup_timer); ctx->cleanup_timer = NULL; } |
||||
|
||||
while (ctx->cbks) { struct broadcast_cbk_entry* c = ctx->cbks; ctx->cbks = c->next; u_free(c); } |
||||
|
||||
if (ctx->messages) { |
||||
struct ll_entry* e; |
||||
while ((e = queue_data_get(ctx->messages)) != NULL) queue_entry_free(e); |
||||
queue_free(ctx->messages); ctx->messages = NULL; |
||||
} |
||||
|
||||
u_free(ctx); group->broadcast = NULL; |
||||
} |
||||
|
||||
int broadcast_send(struct TOPO_GROUP* group, const uint8_t* data, uint16_t data_len) { |
||||
if (!group || !group->broadcast || !data || data_len > BROADCAST_MAX_DATA) return -1; |
||||
struct broadcast_ctx* ctx = group->broadcast; |
||||
|
||||
uint8_t uuid[BROADCAST_UUID_SIZE]; |
||||
{ |
||||
uint64_t tb = get_time_tb(); |
||||
uint32_t seq = ctx->send_seq++; |
||||
uint64_t ptr = (uint64_t)(uintptr_t)ctx; |
||||
memset(uuid, 0, BROADCAST_UUID_SIZE); |
||||
memcpy(uuid, &tb, 8); memcpy(uuid + 8, &seq, 4); memcpy(uuid + 12, &ptr, 4); |
||||
} |
||||
|
||||
{ |
||||
size_t sz = 26 + (size_t)data_len; // 8(time) + 16(uuid) + 2(data_len) + data
|
||||
struct ll_entry* e = queue_entry_new(sz); |
||||
if (!e) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "broadcast_send: alloc failed"); return -1; } |
||||
struct broadcast_msg* msg = (struct broadcast_msg*)e->data; |
||||
msg->local_time_tb = get_time_tb(); |
||||
memcpy(msg->uuid, uuid, BROADCAST_UUID_SIZE); |
||||
msg->data_len = data_len; |
||||
memcpy(msg->data, data, data_len); |
||||
queue_data_put_with_index(ctx->messages, e); |
||||
} |
||||
|
||||
if (group->senders_list) { |
||||
struct ll_entry* se = group->senders_list->head; |
||||
while (se) { |
||||
struct TOPO_GROUP_CONN_ITEM* item = (struct TOPO_GROUP_CONN_ITEM*)se->data; |
||||
if (item && item->conn) { |
||||
uint16_t total = 1 + 8 + BROADCAST_UUID_SIZE + 2 + data_len; |
||||
uint8_t* pkt = u_malloc(total); |
||||
if (!pkt) { se = se->next; continue; } |
||||
pkt[0] = ETCP_ID_BROADCAST; |
||||
memcpy(pkt + 1, &group->group_id, 8); |
||||
memcpy(pkt + 9, uuid, BROADCAST_UUID_SIZE); |
||||
pkt[25] = (uint8_t)(data_len & 0xFF); |
||||
pkt[26] = (uint8_t)((data_len >> 8) & 0xFF); |
||||
memcpy(pkt + 27, data, data_len); |
||||
struct ll_entry* ee = queue_entry_new(0); |
||||
if (!ee) { u_free(pkt); se = se->next; continue; } |
||||
ee->dgram = pkt; ee->len = total; |
||||
if (etcp_send(item->conn, ee) != 0) { u_free(pkt); queue_entry_free(ee); } |
||||
} |
||||
se = se->next; |
||||
} |
||||
} |
||||
return 0; |
||||
} |
||||
|
||||
void broadcast_add_cbk(struct TOPO_GROUP* group, broadcast_recv_fn fn, void* arg) { |
||||
if (!group || !group->broadcast || !fn) return; |
||||
struct broadcast_ctx* ctx = group->broadcast; |
||||
struct broadcast_cbk_entry* e = u_malloc(sizeof(*e)); |
||||
if (!e) return; |
||||
e->fn = fn; e->arg = arg; |
||||
e->next = ctx->cbks; |
||||
ctx->cbks = e; |
||||
} |
||||
|
||||
void broadcast_remove_cbk(struct TOPO_GROUP* group, broadcast_recv_fn fn, void* arg) { |
||||
if (!group || !group->broadcast || !fn) return; |
||||
struct broadcast_ctx* ctx = group->broadcast; |
||||
struct broadcast_cbk_entry** p = &ctx->cbks; |
||||
while (*p) { |
||||
if ((*p)->fn == fn && (*p)->arg == arg) { |
||||
struct broadcast_cbk_entry* r = *p; |
||||
*p = r->next; u_free(r); |
||||
return; |
||||
} |
||||
p = &(*p)->next; |
||||
} |
||||
} |
||||
|
||||
static void broadcast_fire_cbks(struct broadcast_ctx* ctx, const uint8_t* uuid, const uint8_t* data, uint16_t data_len) { |
||||
struct broadcast_cbk_entry* c = ctx->cbks; |
||||
while (c) { c->fn(uuid, data, data_len, c->arg); c = c->next; } |
||||
} |
||||
|
||||
static void broadcast_cleanup_timer_cb(void* arg) { |
||||
struct broadcast_ctx* ctx = (struct broadcast_ctx*)arg; |
||||
if (!ctx || !ctx->group) return; |
||||
uint64_t now_tb = get_time_tb(); |
||||
|
||||
int removed = 0; |
||||
while (1) { |
||||
struct ll_entry* e = queue_data_get(ctx->messages); |
||||
if (!e) break; |
||||
struct broadcast_msg* msg = (struct broadcast_msg*)e->data; |
||||
if (now_tb - msg->local_time_tb > BROADCAST_TTL_TB) { |
||||
queue_entry_free(e); removed++; |
||||
} else { |
||||
queue_data_put_first(ctx->messages, e); |
||||
break; |
||||
} |
||||
} |
||||
if (removed > 0) DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "broadcast_cleanup: removed %d stale msgs grp=%016llx", removed, (unsigned long long)ctx->group->group_id); |
||||
|
||||
ctx->cleanup_timer = uasync_set_timeout(ctx->group->instance->ua, BROADCAST_CLEANUP_INTERVAL_TB, ctx, broadcast_cleanup_timer_cb, "broadcast_cleanup_timer"); |
||||
} |
||||
|
||||
void broadcast_recv(struct TOPO_GROUP* group, struct ETCP_CONN* from_conn, const uint8_t* payload, size_t payload_len) { |
||||
if (!group || !group->broadcast || !payload || payload_len < BROADCAST_UUID_SIZE + 2) return; |
||||
struct broadcast_ctx* ctx = group->broadcast; |
||||
|
||||
const uint8_t* uuid = payload; |
||||
uint16_t data_len = (uint16_t)payload[BROADCAST_UUID_SIZE] | ((uint16_t)payload[BROADCAST_UUID_SIZE + 1] << 8); |
||||
const uint8_t* data = payload + BROADCAST_UUID_SIZE + 2; |
||||
if (data_len > BROADCAST_MAX_DATA || BROADCAST_UUID_SIZE + 2 + (size_t)data_len > payload_len) { |
||||
DEBUG_WARN(DEBUG_CATEGORY_BGP, "broadcast_recv: invalid size data_len=%u payload_len=%zu", data_len, payload_len); |
||||
return; |
||||
} |
||||
|
||||
if (queue_find_data_by_index(ctx->messages, uuid) != NULL) return; |
||||
|
||||
{ |
||||
size_t sz = 26 + (size_t)data_len; |
||||
struct ll_entry* e = queue_entry_new(sz); |
||||
if (!e) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "broadcast_recv: alloc failed"); return; } |
||||
struct broadcast_msg* msg = (struct broadcast_msg*)e->data; |
||||
msg->local_time_tb = get_time_tb(); |
||||
memcpy(msg->uuid, uuid, BROADCAST_UUID_SIZE); |
||||
msg->data_len = data_len; |
||||
memcpy(msg->data, data, data_len); |
||||
queue_data_put_with_index(ctx->messages, e); |
||||
} |
||||
|
||||
if (group->senders_list) { |
||||
struct ll_entry* se = group->senders_list->head; |
||||
while (se) { |
||||
struct TOPO_GROUP_CONN_ITEM* item = (struct TOPO_GROUP_CONN_ITEM*)se->data; |
||||
if (item && item->conn && item->conn != from_conn) { |
||||
uint16_t total = 1 + 8 + BROADCAST_UUID_SIZE + 2 + data_len; |
||||
uint8_t* pkt = u_malloc(total); |
||||
if (!pkt) { se = se->next; continue; } |
||||
pkt[0] = ETCP_ID_BROADCAST; |
||||
memcpy(pkt + 1, &group->group_id, 8); |
||||
memcpy(pkt + 9, uuid, BROADCAST_UUID_SIZE); |
||||
pkt[25] = (uint8_t)(data_len & 0xFF); |
||||
pkt[26] = (uint8_t)((data_len >> 8) & 0xFF); |
||||
memcpy(pkt + 27, data, data_len); |
||||
struct ll_entry* ee = queue_entry_new(0); |
||||
if (!ee) { u_free(pkt); se = se->next; continue; } |
||||
ee->dgram = pkt; ee->len = total; |
||||
if (etcp_send(item->conn, ee) != 0) { u_free(pkt); queue_entry_free(ee); } |
||||
} |
||||
se = se->next; |
||||
} |
||||
} |
||||
|
||||
broadcast_fire_cbks(ctx, uuid, data, data_len); |
||||
} |
||||
|
||||
static void broadcast_recv_dispatcher(struct ETCP_CONN* from_conn, struct ll_entry* entry) { |
||||
if (!from_conn || !entry || entry->len < 1 + 8) { |
||||
if (entry) { queue_dgram_free(entry); queue_entry_free(entry); } |
||||
return; |
||||
} |
||||
struct UTUN_INSTANCE* instance = from_conn->instance; |
||||
if (!instance) { queue_dgram_free(entry); queue_entry_free(entry); return; } |
||||
|
||||
uint8_t* dgram = entry->dgram; |
||||
size_t dgram_len = entry->len; |
||||
if (dgram[0] != ETCP_ID_BROADCAST) { queue_dgram_free(entry); queue_entry_free(entry); return; } |
||||
|
||||
uint64_t group_id; |
||||
memcpy(&group_id, dgram + 1, 8); |
||||
size_t payload_len = dgram_len - 1 - 8; |
||||
const uint8_t* payload = dgram + 1 + 8; |
||||
|
||||
if (instance->topo_groups) { |
||||
struct TOPO_GROUP* group = topo_groups_find(instance->topo_groups, group_id); |
||||
if (group) broadcast_recv(group, from_conn, payload, payload_len); |
||||
} |
||||
queue_dgram_free(entry); queue_entry_free(entry); |
||||
} |
||||
|
||||
int broadcast_init_instance(struct UTUN_INSTANCE* instance) { |
||||
if (!instance) return -1; |
||||
DEBUG_INFO(DEBUG_CATEGORY_BGP, "broadcast_init_instance: binding ETCP_ID_BROADCAST"); |
||||
return etcp_bind(instance, ETCP_ID_BROADCAST, broadcast_recv_dispatcher); |
||||
} |
||||
|
||||
void broadcast_destroy_instance(struct UTUN_INSTANCE* instance) { |
||||
if (!instance) return; |
||||
DEBUG_INFO(DEBUG_CATEGORY_BGP, "broadcast_destroy_instance: unbinding ETCP_ID_BROADCAST"); |
||||
etcp_unbind(instance, ETCP_ID_BROADCAST); |
||||
} |
||||
@ -0,0 +1,41 @@
|
||||
#ifndef BROADCAST_H |
||||
#define BROADCAST_H |
||||
|
||||
#ifdef __cplusplus |
||||
extern "C" { |
||||
#endif |
||||
|
||||
#include <stdint.h> |
||||
#include <stddef.h> |
||||
|
||||
struct TOPO_GROUP; |
||||
struct UTUN_INSTANCE; |
||||
|
||||
#define BROADCAST_MAX_DATA 1600 |
||||
#define BROADCAST_UUID_SIZE 16 |
||||
#define BROADCAST_CLEANUP_INTERVAL_TB 600000 // 1 минута
|
||||
#define BROADCAST_TTL_TB 1200000 // 2 минуты
|
||||
|
||||
typedef void (*broadcast_recv_fn)(const uint8_t* uuid, const uint8_t* data, uint16_t data_len, void* arg); |
||||
|
||||
struct broadcast_cbk_entry { |
||||
broadcast_recv_fn fn; |
||||
void* arg; |
||||
struct broadcast_cbk_entry* next; |
||||
}; |
||||
|
||||
struct broadcast_ctx; |
||||
|
||||
int broadcast_init(struct TOPO_GROUP* group); |
||||
void broadcast_destroy(struct TOPO_GROUP* group); |
||||
int broadcast_send(struct TOPO_GROUP* group, const uint8_t* data, uint16_t data_len); |
||||
void broadcast_add_cbk(struct TOPO_GROUP* group, broadcast_recv_fn fn, void* arg); |
||||
void broadcast_remove_cbk(struct TOPO_GROUP* group, broadcast_recv_fn fn, void* arg); |
||||
|
||||
int broadcast_init_instance(struct UTUN_INSTANCE* instance); |
||||
void broadcast_destroy_instance(struct UTUN_INSTANCE* instance); |
||||
|
||||
#ifdef __cplusplus |
||||
} |
||||
#endif |
||||
#endif |
||||
@ -0,0 +1,516 @@
|
||||
/*
|
||||
* chat_whisper.c — Whisper speech-to-text интеграция для голосовых сообщений |
||||
* |
||||
* Поток данных: |
||||
* on_msg_inserted (content_type == "voice") |
||||
* → md_auto_download (скачиваем .wav/.opus) |
||||
* → chat_whisper_transcribe_async |
||||
* → [worker thread]: PCM decode → resample 16kHz → whisper_full → текст |
||||
* → [uasync callback]: chat_core_submit_message(ct="voice_transcription") |
||||
*/ |
||||
|
||||
#include "chat_whisper.h" |
||||
#include "chat_setting.h" |
||||
|
||||
#include "../media_async/media_async.h" |
||||
#include "../../lib/u_async.h" |
||||
#include "../../lib/mem.h" |
||||
#include "../../lib/platform_compat.h" |
||||
#include "../../lib/debug_config.h" |
||||
|
||||
#include <stdlib.h> |
||||
#include <string.h> |
||||
#include <stdio.h> |
||||
#include <math.h> |
||||
|
||||
#ifdef HAVE_WHISPER |
||||
#include <whisper.h> |
||||
#endif |
||||
|
||||
#define CW_ID "chat_whisper" |
||||
|
||||
/* ─── состояние ─── */ |
||||
|
||||
static struct whisper_context* g_whisper_ctx = NULL; |
||||
static int g_initialized = 0; |
||||
|
||||
/* очередь транскрипций: whisper не thread-safe, обрабатываем по одной */ |
||||
struct wh_job { |
||||
char audio_path[1024]; |
||||
char channel_id[64]; |
||||
uint64_t reply_to_ts; |
||||
uint64_t reply_to_node; |
||||
chat_whisper_done_fn done_cb; |
||||
void* done_arg; |
||||
struct media_async* ma; |
||||
struct UASYNC* ua; |
||||
}; |
||||
|
||||
static struct wh_job* g_pending_job = NULL; |
||||
static int g_processing = 0; |
||||
|
||||
/* ─── wav reader ─── */ |
||||
|
||||
struct wav_pcm { |
||||
float* samples; /* float32 PCM моно */ |
||||
int n_samples; /* количество семплов */ |
||||
int sample_rate; /* исходный sample rate */ |
||||
}; |
||||
|
||||
static int read_u16_le(FILE* f, uint16_t* v) { return fread(v, 2, 1, f) == 1; } |
||||
static int read_u32_le(FILE* f, uint32_t* v) { return fread(v, 4, 1, f) == 1; } |
||||
|
||||
static struct wav_pcm* wav_read(const char* path) { |
||||
FILE* f = fopen(path, "rb"); |
||||
if (!f) { DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "%s: cannot open %s", CW_ID, path); return NULL; } |
||||
|
||||
char riff[4], wave[4]; |
||||
if (fread(riff, 1, 4, f) != 4 || memcmp(riff, "RIFF", 4) != 0) { fclose(f); DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "%s: no RIFF header", CW_ID); return NULL; } |
||||
uint32_t file_size; read_u32_le(f, &file_size); |
||||
if (fread(wave, 1, 4, f) != 4 || memcmp(wave, "WAVE", 4) != 0) { fclose(f); DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "%s: no WAVE", CW_ID); return NULL; } |
||||
|
||||
uint16_t audio_format = 0, channels = 0; |
||||
uint32_t sample_rate = 0, byte_rate = 0, data_size = 0; |
||||
uint16_t block_align = 0, bits_per_sample = 0; |
||||
|
||||
for (;;) { |
||||
char chunk_id[4]; size_t rd = fread(chunk_id, 1, 4, f); |
||||
if (rd < 4) break; |
||||
uint32_t chunk_size; read_u32_le(f, &chunk_size); |
||||
|
||||
if (memcmp(chunk_id, "fmt ", 4) == 0) { |
||||
read_u16_le(f, &audio_format); read_u16_le(f, &channels); |
||||
read_u32_le(f, &sample_rate); read_u32_le(f, &byte_rate); |
||||
read_u16_le(f, &block_align); read_u16_le(f, &bits_per_sample); |
||||
if (chunk_size > 16) fseek(f, chunk_size - 16, SEEK_CUR); |
||||
} else if (memcmp(chunk_id, "data", 4) == 0) { |
||||
data_size = chunk_size; |
||||
break; |
||||
} else { |
||||
fseek(f, chunk_size, SEEK_CUR); |
||||
} |
||||
} |
||||
|
||||
if (audio_format != 1) { fclose(f); DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "%s: not PCM (fmt=%u)", CW_ID, audio_format); return NULL; } |
||||
if (channels < 1 || channels > 2) { fclose(f); DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "%s: bad channels %u", CW_ID, channels); return NULL; } |
||||
if (bits_per_sample != 16) { fclose(f); DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "%s: not 16-bit PCM", CW_ID); return NULL; } |
||||
if (sample_rate < 8000 || sample_rate > 48000) { fclose(f); DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "%s: bad sample_rate %u", CW_ID, sample_rate); return NULL; } |
||||
|
||||
int frame_size = (int)(bits_per_sample / 8) * channels; |
||||
int total_frames = (int)(data_size / frame_size); |
||||
|
||||
struct wav_pcm* pcm = u_calloc(1, sizeof(*pcm)); |
||||
if (!pcm) { fclose(f); return NULL; } |
||||
pcm->samples = u_malloc((size_t)total_frames * sizeof(float)); |
||||
if (!pcm->samples) { fclose(f); u_free(pcm); return NULL; } |
||||
pcm->sample_rate = (int)sample_rate; |
||||
|
||||
pcm->n_samples = 0; |
||||
for (int i = 0; i < total_frames; i++) { |
||||
int16_t raw[2] = {0, 0}; |
||||
if (fread(raw, frame_size < 4 ? 2 : frame_size, 1, f) != 1) break; |
||||
float s = (float)raw[0]; |
||||
if (channels == 2) s = (s + (float)raw[1]) * 0.5f; |
||||
pcm->samples[pcm->n_samples++] = s / 32768.0f; |
||||
} |
||||
|
||||
fclose(f); |
||||
DEBUG_DEBUG(DEBUG_CATEGORY_DEBUG, "%s: wav_read %s: %d samples %dHz mono", CW_ID, path, pcm->n_samples, pcm->sample_rate); |
||||
return pcm; |
||||
} |
||||
|
||||
static void wav_free(struct wav_pcm* pcm) { |
||||
if (!pcm) return; |
||||
u_free(pcm->samples); |
||||
u_free(pcm); |
||||
} |
||||
|
||||
/* ─── opus frame reader (простой header: "OPUS" + sample_rate(4LE) + channels(4LE) + frames...) ─── */ |
||||
|
||||
#include "opus_codec.h" |
||||
|
||||
struct opus_pcm { |
||||
float* samples; |
||||
int n_samples; |
||||
int sample_rate; |
||||
}; |
||||
|
||||
static struct opus_pcm* opus_read(const char* path) { |
||||
FILE* f = fopen(path, "rb"); |
||||
if (!f) { DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "%s: cannot open %s", CW_ID, path); return NULL; } |
||||
|
||||
char magic[4]; |
||||
if (fread(magic, 1, 4, f) != 4 || memcmp(magic, "OPUS", 4) != 0) { fclose(f); return NULL; } |
||||
|
||||
uint32_t sr, ch; |
||||
if (fread(&sr, 4, 1, f) != 1 || fread(&ch, 4, 1, f) != 1) { fclose(f); DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "%s: bad opus header", CW_ID); return NULL; } |
||||
|
||||
opus_codec_decoder_t* dec = opus_codec_decoder_create((int)sr, (int)ch); |
||||
if (!dec) { fclose(f); DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "%s: opus decoder create failed", CW_ID); return NULL; } |
||||
|
||||
int frame_samples = (int)sr * 20 / 1000; |
||||
size_t alloc_samples = 0; |
||||
float* buffer = NULL; |
||||
int total_samples = 0; |
||||
|
||||
for (;;) { |
||||
uint8_t len_buf[2]; |
||||
if (fread(len_buf, 1, 2, f) != 2) break; |
||||
uint16_t pkt_len = (uint16_t)((len_buf[0] << 8) | len_buf[1]); |
||||
if (pkt_len == 0 || pkt_len > 16384) break; |
||||
|
||||
uint8_t* pkt = u_malloc(pkt_len); |
||||
if (!pkt) break; |
||||
if (fread(pkt, 1, pkt_len, f) != pkt_len) { u_free(pkt); break; } |
||||
|
||||
if (total_samples + frame_samples > (int)alloc_samples) { |
||||
alloc_samples = (total_samples + frame_samples + 48000) & ~1023; |
||||
float* nb = u_realloc(buffer, alloc_samples * sizeof(float)); |
||||
if (!nb) { u_free(pkt); break; } |
||||
buffer = nb; |
||||
} |
||||
|
||||
int16_t pcm_buf[5760]; /* max frame at 48kHz stereo */ |
||||
int decoded = opus_codec_decode(dec, pkt, (int)pkt_len, pcm_buf, frame_samples); |
||||
u_free(pkt); |
||||
if (decoded <= 0) continue; |
||||
|
||||
for (int k = 0; k < decoded; k++) |
||||
buffer[total_samples + k] = (float)pcm_buf[k] / 32768.0f; |
||||
total_samples += decoded; |
||||
} |
||||
|
||||
fclose(f); |
||||
opus_codec_decoder_destroy(dec); |
||||
|
||||
struct opus_pcm* pcm = u_calloc(1, sizeof(*pcm)); |
||||
if (!pcm) { u_free(buffer); return NULL; } |
||||
pcm->samples = buffer; |
||||
pcm->n_samples = total_samples; |
||||
pcm->sample_rate = (int)sr; |
||||
DEBUG_DEBUG(DEBUG_CATEGORY_DEBUG, "%s: opus_read %s: %d samples %dHz", CW_ID, path, total_samples, (int)sr); |
||||
return pcm; |
||||
} |
||||
|
||||
static void opus_free(struct opus_pcm* pcm) { |
||||
if (!pcm) return; |
||||
u_free(pcm->samples); |
||||
u_free(pcm); |
||||
} |
||||
|
||||
/* ─── ресемпл в 16kHz float32 моно (линейная интерполяция) ─── */ |
||||
|
||||
static float* resample_16k(const float* src, int n_src, int src_rate, int* n_dst_out) { |
||||
if (src_rate == 16000) { |
||||
float* dst = u_malloc((size_t)n_src * sizeof(float)); |
||||
if (dst) { memcpy(dst, src, (size_t)n_src * sizeof(float)); *n_dst_out = n_src; } |
||||
return dst; |
||||
} |
||||
|
||||
double ratio = 16000.0 / (double)src_rate; |
||||
int n_dst = (int)((double)n_src * ratio) + 16; |
||||
float* dst = u_malloc((size_t)n_dst * sizeof(float)); |
||||
if (!dst) { *n_dst_out = 0; return NULL; } |
||||
|
||||
for (int i = 0; i < n_dst; i++) { |
||||
double pos = (double)i / ratio; |
||||
int idx = (int)pos; |
||||
double frac = pos - (double)idx; |
||||
float v0 = (idx < n_src) ? src[idx] : 0.0f; |
||||
float v1 = (idx + 1 < n_src) ? src[idx + 1] : v0; |
||||
dst[i] = v0 + (float)((v1 - v0) * frac); |
||||
} |
||||
|
||||
*n_dst_out = n_dst; |
||||
return dst; |
||||
} |
||||
|
||||
/* ─── worker thread ─── */ |
||||
|
||||
struct wh_work_ctx; |
||||
static void chat_whisper_transcribe_async_impl( |
||||
struct media_async* ma, struct UASYNC* ua, |
||||
const char* audio_path, const char* channel_id, |
||||
uint64_t reply_to_ts, uint64_t reply_to_node, |
||||
chat_whisper_done_fn done_cb, void* done_arg); |
||||
|
||||
struct wh_work_ctx { |
||||
struct wh_job job; |
||||
char* text; /* результат транскрипции (u_malloc) */ |
||||
int err; |
||||
struct UASYNC* ua; |
||||
}; |
||||
|
||||
static void wh_work_fn(void* raw) { |
||||
struct wh_work_ctx* w = (struct wh_work_ctx*)raw; |
||||
|
||||
#ifndef HAVE_WHISPER |
||||
DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "%s: whisper not compiled in", CW_ID); |
||||
w->err = -1; |
||||
return; |
||||
#else |
||||
if (!g_whisper_ctx) { |
||||
DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "%s: whisper not initialized", CW_ID); |
||||
w->err = -1; |
||||
return; |
||||
} |
||||
|
||||
/* 1. Load PCM from file */ |
||||
float* pcm_samples = NULL; |
||||
int pcm_n = 0; |
||||
int pcm_rate = 0; |
||||
|
||||
/* try wav first, then opus */ |
||||
struct wav_pcm* wpcm = wav_read(w->job.audio_path); |
||||
if (wpcm) { |
||||
pcm_samples = wpcm->samples; pcm_n = wpcm->n_samples; pcm_rate = wpcm->sample_rate; |
||||
wpcm->samples = NULL; wav_free(wpcm); |
||||
} else { |
||||
struct opus_pcm* opcm = opus_read(w->job.audio_path); |
||||
if (opcm) { |
||||
pcm_samples = opcm->samples; pcm_n = opcm->n_samples; pcm_rate = opcm->sample_rate; |
||||
opcm->samples = NULL; opus_free(opcm); |
||||
} |
||||
} |
||||
|
||||
if (!pcm_samples || pcm_n <= 0) { |
||||
DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "%s: cannot decode %s", CW_ID, w->job.audio_path); |
||||
w->err = -1; |
||||
return; |
||||
} |
||||
|
||||
/* 2. Resample to 16kHz */ |
||||
int n_16k = 0; |
||||
float* samples_16k = resample_16k(pcm_samples, pcm_n, pcm_rate, &n_16k); |
||||
u_free(pcm_samples); |
||||
if (!samples_16k || n_16k <= 0) { |
||||
DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "%s: resample failed", CW_ID); |
||||
w->err = -1; |
||||
return; |
||||
} |
||||
|
||||
DEBUG_DEBUG(DEBUG_CATEGORY_DEBUG, "%s: %d samples @ 16kHz, running whisper_full", CW_ID, n_16k); |
||||
|
||||
/* 3. Whisper */ |
||||
struct whisper_full_params wparams = whisper_full_default_params(WHISPER_SAMPLING_GREEDY); |
||||
wparams.language = "ru"; |
||||
wparams.n_threads = 4; |
||||
wparams.no_timestamps = 1; |
||||
wparams.single_segment = 1; |
||||
|
||||
int ret = whisper_full(g_whisper_ctx, wparams, samples_16k, n_16k); |
||||
u_free(samples_16k); |
||||
|
||||
if (ret != 0) { |
||||
DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "%s: whisper_full failed ret=%d", CW_ID, ret); |
||||
w->err = -1; |
||||
return; |
||||
} |
||||
|
||||
/* 4. Collect text */ |
||||
int n_seg = whisper_full_n_segments(g_whisper_ctx); |
||||
if (n_seg <= 0) { |
||||
DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "%s: whisper returned no segments", CW_ID); |
||||
w->text = u_strdup(""); |
||||
} else { |
||||
size_t total = 0; |
||||
for (int i = 0; i < n_seg; i++) { |
||||
const char* seg = whisper_full_get_segment_text(g_whisper_ctx, i); |
||||
if (seg) total += strlen(seg); |
||||
} |
||||
w->text = u_malloc(total + 1); |
||||
if (w->text) { |
||||
w->text[0] = '\0'; |
||||
for (int i = 0; i < n_seg; i++) { |
||||
const char* seg = whisper_full_get_segment_text(g_whisper_ctx, i); |
||||
if (seg) strcat(w->text, seg); |
||||
} |
||||
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: transcribed [%s]", CW_ID, w->text); |
||||
} |
||||
} |
||||
|
||||
w->err = 0; |
||||
#endif /* HAVE_WHISPER */ |
||||
} |
||||
|
||||
static void wh_done_fn(void* raw, int err) { |
||||
(void)err; |
||||
struct wh_work_ctx* w = (struct wh_work_ctx*)raw; |
||||
|
||||
g_processing = 0; |
||||
|
||||
if (w->job.done_cb) { |
||||
w->job.done_cb(w->job.done_arg, w->job.channel_id, |
||||
w->text, w->err, |
||||
w->job.reply_to_ts, w->job.reply_to_node); |
||||
} |
||||
|
||||
u_free(w->text); |
||||
u_free(w); |
||||
|
||||
/* запустить следующую из очереди */ |
||||
if (g_pending_job) { |
||||
struct wh_job* nj = g_pending_job; |
||||
g_pending_job = NULL; |
||||
chat_whisper_transcribe_async_impl(nj->ma, nj->ua, nj->audio_path, nj->channel_id, |
||||
nj->reply_to_ts, nj->reply_to_node, |
||||
nj->done_cb, nj->done_arg); |
||||
u_free(nj); |
||||
} |
||||
} |
||||
|
||||
/* ─── авто-поиск модели ─── */ |
||||
|
||||
static int find_model_path(char* out, size_t out_sz) { |
||||
const char* setting = chat_setting_get("whisper_model_path"); |
||||
if (setting && setting[0] != '\0') { |
||||
/* если путь относительный, не проверяем доступ — вернём как есть */ |
||||
if (access(setting, R_OK) == 0) { snprintf(out, out_sz, "%s", setting); return 0; } |
||||
/* если задан явно но не доступен — ошибка */ |
||||
DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "%s: model_path from setting not accessible: %s", CW_ID, setting); |
||||
} |
||||
|
||||
const char* search_paths[] = { |
||||
"ggml-base.bin", |
||||
"/usr/share/utun/models/ggml-base.bin", |
||||
NULL |
||||
}; |
||||
const char* home = getenv("HOME"); |
||||
char home_path[512]; |
||||
if (home) { snprintf(home_path, sizeof(home_path), "%s/.local/share/utun/models/ggml-base.bin", home); } |
||||
|
||||
for (int i = 0; search_paths[i]; i++) { |
||||
if (access(search_paths[i], R_OK) == 0) { snprintf(out, out_sz, "%s", search_paths[i]); return 0; } |
||||
} |
||||
if (home && access(home_path, R_OK) == 0) { snprintf(out, out_sz, "%s", home_path); return 0; } |
||||
|
||||
return -1; |
||||
} |
||||
|
||||
/* ─── публичный API ─── */ |
||||
|
||||
int chat_whisper_init(void) { |
||||
if (g_initialized) return 0; |
||||
g_initialized = 1; |
||||
|
||||
#ifndef HAVE_WHISPER |
||||
DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "%s: compiled without whisper support", CW_ID); |
||||
return -1; |
||||
#else |
||||
if (!chat_setting_get_int("whisper_enabled", 0)) { |
||||
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: whisper disabled in settings", CW_ID); |
||||
return -1; |
||||
} |
||||
|
||||
char model_path[512]; |
||||
if (find_model_path(model_path, sizeof(model_path)) != 0) { |
||||
DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "%s: model not found (set whisper_model_path)", CW_ID); |
||||
return -1; |
||||
} |
||||
|
||||
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: loading model %s..." , CW_ID, model_path); |
||||
struct whisper_context_params cparams = whisper_context_default_params(); |
||||
cparams.use_gpu = 0; |
||||
g_whisper_ctx = whisper_init_from_file_with_params(model_path, cparams); |
||||
|
||||
if (!g_whisper_ctx) { |
||||
DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "%s: whisper_init_from_file failed for %s", CW_ID, model_path); |
||||
return -1; |
||||
} |
||||
|
||||
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: model loaded successfully", CW_ID); |
||||
g_chat_whisper_trigger = chat_whisper_transcribe_async; |
||||
return 0; |
||||
#endif /* HAVE_WHISPER */ |
||||
} |
||||
|
||||
int chat_whisper_available(void) { |
||||
#ifdef HAVE_WHISPER |
||||
return g_whisper_ctx != NULL && chat_setting_get_int("whisper_enabled", 0); |
||||
#else |
||||
return 0; |
||||
#endif |
||||
} |
||||
|
||||
void chat_whisper_destroy(void) { |
||||
#ifdef HAVE_WHISPER |
||||
if (g_whisper_ctx) { |
||||
whisper_free(g_whisper_ctx); |
||||
g_whisper_ctx = NULL; |
||||
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: model unloaded", CW_ID); |
||||
} |
||||
#endif |
||||
/* отменяем ожидающие задачи */ |
||||
if (g_pending_job) { u_free(g_pending_job); g_pending_job = NULL; } |
||||
g_initialized = 0; |
||||
g_processing = 0; |
||||
} |
||||
|
||||
void chat_whisper_transcribe_async( |
||||
struct media_async* ma, struct UASYNC* ua, |
||||
const char* audio_path, const char* channel_id, |
||||
uint64_t reply_to_ts, uint64_t reply_to_node, |
||||
chat_whisper_done_fn done_cb, void* done_arg) |
||||
{ |
||||
chat_whisper_transcribe_async_impl(ma, ua, audio_path, channel_id, |
||||
reply_to_ts, reply_to_node, |
||||
done_cb, done_arg); |
||||
} |
||||
|
||||
static void chat_whisper_transcribe_async_impl( |
||||
struct media_async* ma, struct UASYNC* ua, |
||||
const char* audio_path, const char* channel_id, |
||||
uint64_t reply_to_ts, uint64_t reply_to_node, |
||||
chat_whisper_done_fn done_cb, void* done_arg) |
||||
{ |
||||
if (!audio_path || !channel_id || !done_cb || !ma || !ua) { |
||||
DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "%s: invalid args", CW_ID); |
||||
if (done_cb) done_cb(done_arg, channel_id ? channel_id : "", NULL, -1, reply_to_ts, reply_to_node); |
||||
return; |
||||
} |
||||
|
||||
if (!chat_whisper_available()) { |
||||
int rc = chat_whisper_init(); |
||||
if (rc != 0) { |
||||
DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "%s: whisper not available, skip transcription", CW_ID); |
||||
done_cb(done_arg, channel_id, NULL, -1, reply_to_ts, reply_to_node); |
||||
return; |
||||
} |
||||
} |
||||
|
||||
if (g_processing) { |
||||
/* поставить в очередь */ |
||||
if (g_pending_job) { u_free(g_pending_job); } /* заменяем */ |
||||
g_pending_job = u_malloc(sizeof(*g_pending_job)); |
||||
if (!g_pending_job) { done_cb(done_arg, channel_id, NULL, -1, reply_to_ts, reply_to_node); return; } |
||||
snprintf(g_pending_job->audio_path, sizeof(g_pending_job->audio_path), "%s", audio_path); |
||||
snprintf(g_pending_job->channel_id, sizeof(g_pending_job->channel_id), "%s", channel_id); |
||||
g_pending_job->reply_to_ts = reply_to_ts; |
||||
g_pending_job->reply_to_node = reply_to_node; |
||||
g_pending_job->done_cb = done_cb; |
||||
g_pending_job->done_arg = done_arg; |
||||
g_pending_job->ma = ma; |
||||
g_pending_job->ua = ua; |
||||
DEBUG_DEBUG(DEBUG_CATEGORY_DEBUG, "%s: queued transcription for ch=%s", CW_ID, channel_id); |
||||
return; |
||||
} |
||||
|
||||
g_processing = 1; |
||||
|
||||
struct wh_work_ctx* w = u_calloc(1, sizeof(*w)); |
||||
if (!w) { g_processing = 0; done_cb(done_arg, channel_id, NULL, -1, reply_to_ts, reply_to_node); return; } |
||||
|
||||
snprintf(w->job.audio_path, sizeof(w->job.audio_path), "%s", audio_path); |
||||
snprintf(w->job.channel_id, sizeof(w->job.channel_id), "%s", channel_id); |
||||
w->job.reply_to_ts = reply_to_ts; |
||||
w->job.reply_to_node = reply_to_node; |
||||
w->job.done_cb = done_cb; |
||||
w->job.done_arg = done_arg; |
||||
w->job.ma = ma; |
||||
w->job.ua = ua; |
||||
w->ua = ua; |
||||
w->err = 0; |
||||
|
||||
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: starting transcription ch=%s file=%s", CW_ID, channel_id, audio_path); |
||||
media_async_submit(ma, ua, wh_work_fn, w, wh_done_fn, w); |
||||
} |
||||
@ -0,0 +1,52 @@
|
||||
/*
|
||||
* chat_whisper.h — Whisper speech-to-text транскрипция для голосовых сообщений |
||||
* |
||||
* Интеграция с whisper.cpp C API (libwhisper): |
||||
* - Детекция доступности whisper в системе |
||||
* - Ленивая загрузка модели при первом голосовом сообщении |
||||
* - Асинхронная транскрипция через media_async (thread-per-task) |
||||
* - Сериализация: одна транскрипция за раз (whisper не thread-safe) |
||||
* - Авто-поиск модели: настройка → ~/.local/share/utun/models/ → /usr/share/utun/models/ |
||||
* |
||||
* chat_msg.c вызывает транскрипцию через глобальный function pointer |
||||
* g_chat_whisper_trigger, который устанавливается при инициализации whisper. |
||||
* Если whisper не скомпилирован/недоступен — pointer = NULL, вызов игнорируется. |
||||
*/ |
||||
#ifndef CHAT_WHISPER_H |
||||
#define CHAT_WHISPER_H |
||||
|
||||
#include <stdint.h> |
||||
#include <stddef.h> |
||||
|
||||
#ifdef __cplusplus |
||||
extern "C" { |
||||
#endif |
||||
|
||||
struct media_async; |
||||
struct UASYNC; |
||||
|
||||
/* ── жизненный цикл ── */ |
||||
|
||||
int chat_whisper_init(void); /* ленивая загрузка модели */ |
||||
int chat_whisper_available(void); /* модель загружена и готова */ |
||||
void chat_whisper_destroy(void); /* выгрузка модели */ |
||||
|
||||
/* ── асинхронная транскрипция ── */ |
||||
|
||||
typedef void (*chat_whisper_done_fn)(void* arg, const char* channel_id, |
||||
const char* text, int err, |
||||
uint64_t reply_ts, uint64_t reply_node); |
||||
|
||||
typedef void (*chat_whisper_trigger_fn)( |
||||
struct media_async* ma, struct UASYNC* ua, |
||||
const char* audio_path, const char* channel_id, |
||||
uint64_t reply_ts, uint64_t reply_node, |
||||
chat_whisper_done_fn done_cb, void* done_arg); |
||||
|
||||
/* Глобальный триггер: устанавливается chat_whisper_init(), вызывается из chat_msg.c */ |
||||
extern chat_whisper_trigger_fn g_chat_whisper_trigger; |
||||
|
||||
#ifdef __cplusplus |
||||
} |
||||
#endif |
||||
#endif /* CHAT_WHISPER_H */ |
||||
@ -0,0 +1,277 @@
|
||||
#include <stdio.h> |
||||
#include <stdlib.h> |
||||
#include <string.h> |
||||
#include "../lib/platform_compat.h" |
||||
#include "test_utils.h" |
||||
#ifdef _WIN32 |
||||
#include <windows.h> |
||||
#else |
||||
#include <unistd.h> |
||||
#endif |
||||
|
||||
#include "etcp.h" |
||||
#include "etcp_connections.h" |
||||
#include "../src/config_parser.h" |
||||
#include "../src/config_updater.h" |
||||
#include "../src/utun_instance.h" |
||||
#include "routing.h" |
||||
#include "topo_group.h" |
||||
#include "topo_node.h" |
||||
#include "../src/tun_if.h" |
||||
#include "secure_channel.h" |
||||
#include "../lib/u_async.h" |
||||
#include "../lib/debug_config.h" |
||||
#include "../src/broadcast.h" |
||||
|
||||
#define TEST_TIMEOUT_TB 10000 |
||||
#define POLL_INTERVAL_MS 1 |
||||
|
||||
static struct UTUN_INSTANCE* inst[5]; |
||||
static struct UASYNC* ua; |
||||
static uint64_t nid[5]; |
||||
static int test_phase = 0; |
||||
static void* test_timeout_id = NULL; |
||||
|
||||
static int recv_count[5]; |
||||
static uint8_t recv_data[5][256]; |
||||
static uint16_t recv_data_len[5]; |
||||
|
||||
static void broadcast_cb(const uint8_t* uuid, const uint8_t* data, uint16_t data_len, void* arg) { |
||||
(void)uuid; |
||||
int idx = (int)(intptr_t)arg; |
||||
recv_count[idx]++; |
||||
if (data_len <= 256) { recv_data_len[idx] = data_len; memcpy(recv_data[idx], data, data_len); } |
||||
} |
||||
|
||||
static void test_timeout_cb(void* arg) { |
||||
(void)arg; |
||||
if (test_phase == 0) test_phase = 2; |
||||
} |
||||
|
||||
#define FAIL(msg) do { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "FAIL: %s", msg); test_phase = 2; goto cleanup; } while(0) |
||||
#define ASSERT(cond, msg) do { if (!(cond)) { FAIL(msg); } } while(0) |
||||
|
||||
struct ncfg { |
||||
uint64_t node_id; const char* tun_ip; const char* priv_hex; const char* pub_hex; |
||||
int srv_cnt; struct { const char* name; int port; } srvs[3]; |
||||
int cli_cnt; struct { const char* name; const char* peer_pubkey_hex; int lnk_cnt; struct { const char* local_srv; int remote_port; } lnks[2]; } clis[2]; |
||||
}; |
||||
|
||||
static void build_cfg(char* buf, size_t size, const struct ncfg* c) { |
||||
int off = snprintf(buf, size, |
||||
"[global]\n" |
||||
"my_private_key=%s\n" |
||||
"my_public_key=%s\n" |
||||
"tun_ip=%s\n" |
||||
"tun_ifname=tun99\n" |
||||
"keepalive_timeout=200\n" |
||||
"keepalive_interval=20\n" |
||||
"keepalive_adaptive=0\n" |
||||
"\n", |
||||
c->priv_hex, c->pub_hex, c->tun_ip); |
||||
for (int i = 0; i < c->srv_cnt; i++) |
||||
off += snprintf(buf + off, size - off, "[server: %s]\naddr=127.0.0.1:%d\ntype=public\n\n", c->srvs[i].name, c->srvs[i].port); |
||||
for (int i = 0; i < c->cli_cnt; i++) { |
||||
off += snprintf(buf + off, size - off, "[client: %s]\nkeepalive=1\npeer_public_key=%s\n", c->clis[i].name, c->clis[i].peer_pubkey_hex); |
||||
for (int j = 0; j < c->clis[i].lnk_cnt; j++) |
||||
off += snprintf(buf + off, size - off, "link=%s:127.0.0.1:%d\n", c->clis[i].lnks[j].local_srv, c->clis[i].lnks[j].remote_port); |
||||
off += snprintf(buf + off, size - off, "\n"); |
||||
} |
||||
off += snprintf(buf + off, size - off, "[allowed_keys]\nallow_all=1\n"); |
||||
} |
||||
|
||||
static int peer_in_nodes(int inst_idx, uint64_t node_id) { |
||||
struct TOPO_GROUP* g = topo_groups_get_default(inst[inst_idx]->topo_groups); |
||||
return g && topo_node_find_by_id(g, node_id) != NULL; |
||||
} |
||||
|
||||
static struct ETCP_LINK* find_client_link(struct UTUN_INSTANCE* ins, const char* cname) { |
||||
if (!ins || !ins->connections) return NULL; |
||||
struct ll_entry* e = ins->connections->head; |
||||
while (e) { |
||||
struct conn_queue_entry* ce = (struct conn_queue_entry*)e->data; |
||||
if (!ce->conn->name || strcmp(ce->conn->name, cname) != 0) { e = e->next; continue; } |
||||
struct ETCP_LINK* l = ce->conn->links; |
||||
while (l) { if (l->is_server == 0) return l; l = l->next; } |
||||
e = e->next; |
||||
} |
||||
return NULL; |
||||
} |
||||
|
||||
static int count_initialized_links(void) { |
||||
int n = 0; |
||||
for (int i = 0; i < 5; i++) { |
||||
if (!inst[i] || !inst[i]->connections) continue; |
||||
struct ll_entry* e = inst[i]->connections->head; |
||||
while (e) { |
||||
struct conn_queue_entry* ce = (struct conn_queue_entry*)e->data; |
||||
struct ETCP_LINK* l = ce->conn->links; |
||||
while (l) { if (l->initialized) n++; l = l->next; } |
||||
e = e->next; |
||||
} |
||||
} |
||||
return n; |
||||
} |
||||
|
||||
static int cond_links_init(void) { return count_initialized_links() >= 5; } |
||||
|
||||
static int cond_e_in_a(void) { return peer_in_nodes(0, nid[4]); } |
||||
static int cond_d_in_a(void) { return peer_in_nodes(0, nid[3]); } |
||||
static int cond_c_in_a(void) { return peer_in_nodes(0, nid[2]); } |
||||
static int cond_all_bgp(void) { return cond_c_in_a() && cond_d_in_a() && cond_e_in_a(); } |
||||
|
||||
static int wait_for(const char* desc, int (*cond)(void), int timeout_tb) { |
||||
uint64_t start = get_time_tb(); |
||||
while (!cond() && (get_time_tb() - start) < (uint64_t)timeout_tb && test_phase == 0) |
||||
uasync_poll(ua, POLL_INTERVAL_MS); |
||||
if (!cond() && test_phase == 0) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "wait_for timeout: %s", desc); return 0; } |
||||
return test_phase == 0; |
||||
} |
||||
|
||||
static int all_recv(int sender_idx, int expect) { |
||||
for (int i = 0; i < 5; i++) { |
||||
if (i == sender_idx) continue; |
||||
if (recv_count[i] != expect) return 0; |
||||
} |
||||
return 1; |
||||
} |
||||
|
||||
static int wait_recv(const char* desc, int sender_idx, int expect, int timeout_tb) { |
||||
uint64_t start = get_time_tb(); |
||||
while (!all_recv(sender_idx, expect) && (get_time_tb() - start) < (uint64_t)timeout_tb && test_phase == 0) |
||||
uasync_poll(ua, POLL_INTERVAL_MS); |
||||
if (!all_recv(sender_idx, expect) && test_phase == 0) { |
||||
DEBUG_ERROR(DEBUG_CATEGORY_BGP, "wait_recv timeout: %s", desc); |
||||
return 0; |
||||
} |
||||
return test_phase == 0; |
||||
} |
||||
|
||||
int main(void) { |
||||
debug_config_init(); |
||||
debug_set_level(DEBUG_LEVEL_ERROR); |
||||
debug_set_categories(DEBUG_CATEGORY_BGP); |
||||
debug_set_categories(DEBUG_CATEGORY_BGP); |
||||
utun_instance_set_tun_init_enabled(0); |
||||
|
||||
struct SC_MYKEYS keys[5]; |
||||
char pub_hex[5][SC_PUBKEY_SIZE * 2 + 1]; |
||||
char priv_hex[5][SC_PRIVKEY_SIZE * 2 + 1]; |
||||
for (int i = 0; i < 5; i++) { |
||||
sc_generate_keypair(&keys[i]); |
||||
bytes_to_hex(keys[i].public_key, SC_PUBKEY_SIZE, pub_hex[i], sizeof(pub_hex[i])); |
||||
bytes_to_hex(keys[i].private_key, SC_PRIVKEY_SIZE, priv_hex[i], sizeof(priv_hex[i])); |
||||
} |
||||
|
||||
int base = 43000 + (getpid() % 10000); |
||||
int p_ab_a = base++, p_ab_b = base++; |
||||
int p_ac_a = base++, p_ac_c = base++; |
||||
int p_bc_b = base++, p_bc_c = base++; |
||||
int p_cd_c = base++, p_cd_d = base++; |
||||
int p_de_d = base++, p_de_e = base++; |
||||
|
||||
/* A: servers for B and C, clients to B and C */ |
||||
char cfg_a[2048]; build_cfg(cfg_a, sizeof(cfg_a), &(struct ncfg){ |
||||
.priv_hex = priv_hex[0], .pub_hex = pub_hex[0], .tun_ip = "10.200.0.1/24", |
||||
.srv_cnt = 2, .srvs = {{"a_srv_b", p_ab_a}, {"a_srv_c", p_ac_a}}, |
||||
.cli_cnt = 2, .clis = { |
||||
{"to_b", pub_hex[1], 1, {{"a_srv_b", p_ab_b}}}, |
||||
{"to_c", pub_hex[2], 1, {{"a_srv_c", p_ac_c}}}, |
||||
} |
||||
}); |
||||
|
||||
/* B: server for A, server for C, client to C */ |
||||
char cfg_b[2048]; build_cfg(cfg_b, sizeof(cfg_b), &(struct ncfg){ |
||||
.priv_hex = priv_hex[1], .pub_hex = pub_hex[1], .tun_ip = "10.200.0.2/24", |
||||
.srv_cnt = 2, .srvs = {{"b_srv", p_ab_b}, {"b_srv_c", p_bc_b}}, |
||||
.cli_cnt = 1, .clis = {{"to_c", pub_hex[2], 1, {{"b_srv_c", p_bc_c}}}} |
||||
}); |
||||
|
||||
/* C: servers for A, B, D; client to D */ |
||||
char cfg_c[2048]; build_cfg(cfg_c, sizeof(cfg_c), &(struct ncfg){ |
||||
.priv_hex = priv_hex[2], .pub_hex = pub_hex[2], .tun_ip = "10.200.0.3/24", |
||||
.srv_cnt = 3, .srvs = {{"c_srv_a", p_ac_c}, {"c_srv_b", p_bc_c}, {"c_srv_d", p_cd_c}}, |
||||
.cli_cnt = 1, .clis = {{"to_d", pub_hex[3], 1, {{"c_srv_d", p_cd_d}}}} |
||||
}); |
||||
|
||||
/* D: server for C, server for E; client to E */ |
||||
char cfg_d[2048]; build_cfg(cfg_d, sizeof(cfg_d), &(struct ncfg){ |
||||
.priv_hex = priv_hex[3], .pub_hex = pub_hex[3], .tun_ip = "10.200.0.4/24", |
||||
.srv_cnt = 2, .srvs = {{"d_srv", p_cd_d}, {"d_srv_e", p_de_d}}, |
||||
.cli_cnt = 1, .clis = {{"to_e", pub_hex[4], 1, {{"d_srv_e", p_de_e}}}} |
||||
}); |
||||
|
||||
/* E: server for D, no clients */ |
||||
char cfg_e[2048]; build_cfg(cfg_e, sizeof(cfg_e), &(struct ncfg){ |
||||
.priv_hex = priv_hex[4], .pub_hex = pub_hex[4], .tun_ip = "10.200.0.5/24", |
||||
.srv_cnt = 1, .srvs = {{"e_srv", p_de_e}}, |
||||
.cli_cnt = 0, .clis = {} |
||||
}); |
||||
|
||||
ua = uasync_create(); ASSERT(ua, "uasync_create"); |
||||
|
||||
inst[0] = utun_instance_create_from_str(ua, cfg_a); ASSERT(inst[0], "inst A"); |
||||
inst[1] = utun_instance_create_from_str(ua, cfg_b); ASSERT(inst[1], "inst B"); |
||||
inst[2] = utun_instance_create_from_str(ua, cfg_c); ASSERT(inst[2], "inst C"); |
||||
inst[3] = utun_instance_create_from_str(ua, cfg_d); ASSERT(inst[3], "inst D"); |
||||
inst[4] = utun_instance_create_from_str(ua, cfg_e); ASSERT(inst[4], "inst E"); |
||||
for (int i = 0; i < 5; i++) { ASSERT(utun_instance_init(inst[i]) == 0, "init"); nid[i] = inst[i]->node_id; } |
||||
|
||||
test_timeout_id = uasync_set_timeout(ua, TEST_TIMEOUT_TB, NULL, test_timeout_cb, "test_timeout"); |
||||
|
||||
/* Phase 1: BGP sync */ |
||||
DEBUG_INFO(DEBUG_CATEGORY_BGP, "=== PHASE 1: links init ==="); |
||||
wait_for("links init", cond_links_init, 5000); |
||||
DEBUG_INFO(DEBUG_CATEGORY_BGP, "=== PHASE 1: BGP sync ==="); |
||||
wait_for("all nodes visible", cond_all_bgp, 5000); |
||||
ASSERT(peer_in_nodes(0, nid[1]), "B missing in A"); |
||||
ASSERT(peer_in_nodes(0, nid[2]), "C missing in A"); |
||||
ASSERT(peer_in_nodes(0, nid[3]), "D missing in A"); |
||||
ASSERT(peer_in_nodes(0, nid[4]), "E missing in A"); |
||||
DEBUG_INFO(DEBUG_CATEGORY_BGP, "Phase 1 PASSED: all nodes visible"); |
||||
|
||||
/* Phase 2: broadcast from A */ |
||||
DEBUG_INFO(DEBUG_CATEGORY_BGP, "=== PHASE 2: broadcast ==="); |
||||
for (int i = 0; i < 5; i++) { |
||||
struct TOPO_GROUP* g = topo_groups_get_default(inst[i]->topo_groups); |
||||
ASSERT(g, "no default group"); |
||||
broadcast_add_cbk(g, broadcast_cb, (void*)(intptr_t)i); |
||||
} |
||||
|
||||
{ struct TOPO_GROUP* g = topo_groups_get_default(inst[0]->topo_groups); |
||||
const uint8_t* msg = (const uint8_t*)"hello_broadcast"; |
||||
ASSERT(broadcast_send(g, msg, 15) == 0, "broadcast_send failed"); } |
||||
|
||||
wait_recv("broadcast received", 0, 1, 5000); |
||||
for (int i = 1; i < 5; i++) { |
||||
ASSERT(recv_count[i] == 1, "wrong recv count"); |
||||
ASSERT(recv_data_len[i] == 15, "wrong data len"); |
||||
ASSERT(memcmp(recv_data[i], "hello_broadcast", 15) == 0, "wrong data"); |
||||
} |
||||
DEBUG_INFO(DEBUG_CATEGORY_BGP, "Phase 2 PASSED: all 4 nodes received broadcast"); |
||||
|
||||
/* Phase 3: second broadcast — dedup check */ |
||||
DEBUG_INFO(DEBUG_CATEGORY_BGP, "=== PHASE 3: second broadcast ==="); |
||||
{ struct TOPO_GROUP* g = topo_groups_get_default(inst[0]->topo_groups); |
||||
const uint8_t* msg = (const uint8_t*)"second_msg"; |
||||
ASSERT(broadcast_send(g, msg, 10) == 0, "broadcast_send2 failed"); } |
||||
|
||||
wait_recv("broadcast2 received", 0, 2, 5000); |
||||
for (int i = 1; i < 5; i++) ASSERT(recv_count[i] == 2, "wrong recv2 count"); |
||||
DEBUG_INFO(DEBUG_CATEGORY_BGP, "Phase 3 PASSED: all 4 nodes received second broadcast, no dups"); |
||||
|
||||
test_phase = 1; |
||||
|
||||
cleanup: |
||||
for (int i = 0; i < 5; i++) { |
||||
if (inst[i]) { |
||||
struct TOPO_GROUP* g = topo_groups_get_default(inst[i]->topo_groups); |
||||
if (g) broadcast_remove_cbk(g, broadcast_cb, (void*)(intptr_t)i); |
||||
} |
||||
} |
||||
if (test_timeout_id) uasync_cancel_timeout(ua, test_timeout_id); |
||||
for (int i = 0; i < 5; i++) { if (inst[i]) { inst[i]->running = 0; utun_instance_destroy(inst[i]); inst[i] = NULL; } } |
||||
if (ua) { uasync_destroy(ua, 0); ua = NULL; } |
||||
printf("=== %s ===\n", test_phase == 1 ? "TEST PASSED" : "TEST FAILED"); |
||||
return test_phase == 1 ? 0 : 1; |
||||
} |
||||
@ -0,0 +1,69 @@
|
||||
#include "flagpainter.h" |
||||
#include <QPainterPath> |
||||
#include <cmath> |
||||
|
||||
void FlagPainter::draw(QPainter* p, uint8_t flags, int x, int y, int sz) { |
||||
if (!flags) return; |
||||
int cx = x; |
||||
if (flags & MEMBER_FLAG_SUPERNODE) { drawLightning(p, cx, y, sz); cx += sz + 2; } |
||||
if (flags & MEMBER_FLAG_ADMIN) { drawStar(p, cx, y, sz); cx += sz + 2; } |
||||
if (flags & MEMBER_FLAG_MODER) { drawShield(p, cx, y, sz); cx += sz + 2; } |
||||
} |
||||
|
||||
void FlagPainter::drawLightning(QPainter* p, int x, int y, int sz) { |
||||
p->save(); |
||||
p->setRenderHint(QPainter::Antialiasing); |
||||
QPainterPath path; |
||||
qreal cx = x + sz / 2.0, t = y + 1.0, b = y + sz - 1.0; |
||||
path.moveTo(cx, t); |
||||
path.lineTo(cx - sz * 0.3, t + sz * 0.55); |
||||
path.lineTo(cx - sz * 0.05, t + sz * 0.5); |
||||
path.lineTo(cx + sz * 0.15, b); |
||||
path.lineTo(cx + sz * 0.1, t + sz * 0.55); |
||||
path.lineTo(cx + sz * 0.35, t + sz * 0.4); |
||||
path.closeSubpath(); |
||||
p->setPen(Qt::NoPen); |
||||
p->setBrush(QColor("#FFC107")); |
||||
p->drawPath(path); |
||||
p->restore(); |
||||
} |
||||
|
||||
void FlagPainter::drawStar(QPainter* p, int x, int y, int sz) { |
||||
p->save(); |
||||
p->setRenderHint(QPainter::Antialiasing); |
||||
QPainterPath path; |
||||
qreal cx = x + sz / 2.0, cy = y + sz / 2.0, r = sz / 2.0 - 1.0; |
||||
for (int i = 0; i < 10; i++) { |
||||
qreal a = (i * 36.0 - 90.0) * 3.14159265 / 180.0; |
||||
qreal rr = (i & 1) ? r * 0.42 : r; |
||||
qreal px = cx + rr * cos(a), py = cy + rr * sin(a); |
||||
if (i == 0) path.moveTo(px, py); |
||||
else path.lineTo(px, py); |
||||
} |
||||
path.closeSubpath(); |
||||
p->setPen(QPen(QColor("#E65100"), 0.5)); |
||||
p->setBrush(QColor("#FFD700")); |
||||
p->drawPath(path); |
||||
p->restore(); |
||||
} |
||||
|
||||
void FlagPainter::drawShield(QPainter* p, int x, int y, int sz) { |
||||
p->save(); |
||||
p->setRenderHint(QPainter::Antialiasing); |
||||
QPainterPath path; |
||||
qreal cx = x + sz / 2.0, t = y + 1.0, b = y + sz - 1.0; |
||||
qreal w = sz * 0.38; |
||||
qreal notch = sz * 0.22; |
||||
path.moveTo(cx - w, t); |
||||
path.lineTo(cx - w, t + notch); |
||||
path.lineTo(cx - w * 0.3, t + notch + 1); |
||||
path.lineTo(cx, b); |
||||
path.lineTo(cx + w * 0.3, t + notch + 1); |
||||
path.lineTo(cx + w, t + notch); |
||||
path.lineTo(cx + w, t); |
||||
path.closeSubpath(); |
||||
p->setPen(QPen(QColor("#1565C0"), 0.8)); |
||||
p->setBrush(QColor("#2196F3")); |
||||
p->drawPath(path); |
||||
p->restore(); |
||||
} |
||||
@ -0,0 +1,22 @@
|
||||
#ifndef FLAGPAINTER_H |
||||
#define FLAGPAINTER_H |
||||
|
||||
#include <QPainter> |
||||
#include <QColor> |
||||
#include <stdint.h> |
||||
|
||||
#define MEMBER_FLAG_SUPERNODE 0x01 |
||||
#define MEMBER_FLAG_ADMIN 0x02 |
||||
#define MEMBER_FLAG_MODER 0x04 |
||||
|
||||
class FlagPainter { |
||||
public: |
||||
static void draw(QPainter* p, uint8_t flags, int x, int y, int sz); |
||||
|
||||
private: |
||||
static void drawLightning(QPainter* p, int x, int y, int sz); |
||||
static void drawStar(QPainter* p, int x, int y, int sz); |
||||
static void drawShield(QPainter* p, int x, int y, int sz); |
||||
}; |
||||
|
||||
#endif |
||||
Loading…
Reference in new issue