41 changed files with 2428 additions and 290 deletions
@ -0,0 +1,536 @@
|
||||
/*
|
||||
* chat_headless_control.c — TCP control socket for headless chat CLI |
||||
* |
||||
* JSON line protocol over TCP. Integrated into uasync event loop. |
||||
* Default: localhost-only, port from config. |
||||
*/ |
||||
#include "chat_headless_control.h" |
||||
#include "invite_link.h" |
||||
#include "chat_core.h" |
||||
#include "chat_sync.h" |
||||
#include "chat_event.h" |
||||
#include "../utun_instance.h" |
||||
#include "../ntp_time.h" |
||||
#include "../transport_layer/etcp.h" |
||||
#include "../transport_layer/etcp_connections.h" |
||||
#include "../../lib/u_async.h" |
||||
#include "../../lib/debug_config.h" |
||||
#include "../../lib/mem.h" |
||||
#include "../../lib/ll_queue.h" |
||||
#include "../routing_layer/topo_node_sqlite.h" |
||||
|
||||
#include <stdlib.h> |
||||
#include <string.h> |
||||
#include <stdio.h> |
||||
#include <errno.h> |
||||
#include <unistd.h> |
||||
#include <sys/socket.h> |
||||
#include <netinet/in.h> |
||||
#include <arpa/inet.h> |
||||
#include <fcntl.h> |
||||
|
||||
#define MAX_CLIENTS 16 |
||||
#define RECV_BUF_SIZE 16384 |
||||
#define SEND_BUF_SIZE 16384 |
||||
|
||||
#ifndef DEBUG_CATEGORY_HEADLESS |
||||
#define DEBUG_CATEGORY_HEADLESS 27 |
||||
#endif |
||||
|
||||
struct headless_client { |
||||
socket_t fd; |
||||
void* socket_id; |
||||
char recv_buf[RECV_BUF_SIZE]; |
||||
int recv_len; |
||||
char send_buf[SEND_BUF_SIZE]; |
||||
int send_len; |
||||
int send_offset; |
||||
int subscribed; |
||||
int closing; |
||||
struct headless_client* next; |
||||
}; |
||||
|
||||
static struct { |
||||
struct UASYNC* ua; |
||||
struct UTUN_INSTANCE* inst; |
||||
struct headless_client* clients; |
||||
int client_count; |
||||
socket_t listen_fd; |
||||
void* listen_sock_id; |
||||
int running; |
||||
} g_hc; |
||||
|
||||
/* ── Forward declarations ── */ |
||||
|
||||
static void hc_handle_command(struct headless_client* cli, const char* json); |
||||
|
||||
/* ── JSON helpers ── */ |
||||
|
||||
static int json_get_int(const char* s, const char* key, int def) { |
||||
char pat[128]; snprintf(pat, sizeof(pat), "\"%s\"", key); |
||||
const char* p = strstr(s, pat); if (!p) return def; |
||||
p += strlen(pat); |
||||
while (*p == ' ' || *p == '\t' || *p == ':') p++; |
||||
return (int)strtol(p, NULL, 10); |
||||
} |
||||
|
||||
static int json_get_str(const char* s, const char* key, char* out, size_t out_sz) { |
||||
char pat[128]; snprintf(pat, sizeof(pat), "\"%s\"", key); |
||||
const char* p = strstr(s, pat); if (!p) return -1; |
||||
p += strlen(pat); |
||||
while (*p == ' ' || *p == '\t' || *p == ':') p++; |
||||
if (*p == '"') p++; |
||||
size_t i = 0; |
||||
while (*p && *p != '"' && i < out_sz - 1) { |
||||
if (*p == '\\' && *(p + 1)) { p++; out[i++] = *(p++); } |
||||
else out[i++] = *(p++); |
||||
} |
||||
out[i] = '\0'; return 0; |
||||
} |
||||
static const char* json_get_cmd(const char* s) { |
||||
const char* p = strstr(s, "\"cmd\""); |
||||
if (!p) return NULL; |
||||
p += 5; |
||||
while (*p == ' ' || *p == '\t' || *p == ':') p++; |
||||
if (*p == '"') p++; |
||||
return p; |
||||
} |
||||
static const char* json_get_cmd_end(const char* cmd_start) { |
||||
const char* p = cmd_start; |
||||
while (*p && *p != '"') { if (*p == '\\' && *(p + 1)) p++; p++; } |
||||
return p; |
||||
} |
||||
|
||||
/* ── Send helpers ── */ |
||||
|
||||
static void cli_send(struct headless_client* cli, const char* data, size_t len) { |
||||
if (!cli || cli->closing || cli->fd < 0 || !data || len == 0) return; |
||||
if (cli->send_len > 0) { |
||||
size_t space = sizeof(cli->send_buf) - (size_t)cli->send_len; |
||||
if (len + 1 > space) return; |
||||
memcpy(cli->send_buf + cli->send_len, data, len); |
||||
cli->send_len += (int)len; |
||||
return; |
||||
} |
||||
ssize_t r = write(cli->fd, data, len); |
||||
if (r < 0 && (errno == EAGAIN || errno == EWOULDBLOCK)) { |
||||
if (len <= sizeof(cli->send_buf)) { memcpy(cli->send_buf, data, len); cli->send_len = (int)len; cli->send_offset = 0; } |
||||
if (g_hc.ua) uasync_set_socket_write(g_hc.ua, cli->socket_id, 1); |
||||
return; |
||||
} |
||||
if (r >= 0 && (size_t)r < len) { |
||||
size_t remain = len - (size_t)r; |
||||
if (remain <= sizeof(cli->send_buf)) { memcpy(cli->send_buf, data + r, remain); cli->send_len = (int)remain; cli->send_offset = 0; } |
||||
if (g_hc.ua) uasync_set_socket_write(g_hc.ua, cli->socket_id, 1); |
||||
} |
||||
} |
||||
|
||||
static void cli_send_str(struct headless_client* cli, const char* s) { |
||||
if (s) cli_send(cli, s, strlen(s)); |
||||
} |
||||
|
||||
static void send_response(struct headless_client* cli, int id, const char* ok_data, const char* error) { |
||||
char buf[4096]; int len; |
||||
if (error) len = snprintf(buf, sizeof(buf), "{\"id\":%d,\"ok\":false,\"error\":\"%s\"}\n", id, error); |
||||
else if (ok_data) len = snprintf(buf, sizeof(buf), "{\"id\":%d,\"ok\":true,\"data\":%s}\n", id, ok_data); |
||||
else len = snprintf(buf, sizeof(buf), "{\"id\":%d,\"ok\":true}\n", id); |
||||
if (len > 0) cli_send(cli, buf, (size_t)len); |
||||
} |
||||
|
||||
/* ── Broadcast event ── */ |
||||
|
||||
static void hc_broadcast_event(const char* json) { |
||||
size_t jlen = strlen(json); |
||||
char buf[4096]; int blen = snprintf(buf, sizeof(buf), "%s\n", json); |
||||
if (blen <= 0) return; |
||||
for (struct headless_client* c = g_hc.clients; c; c = c->next) |
||||
if (c->subscribed) cli_send(c, buf, (size_t)blen); |
||||
} |
||||
|
||||
/* ── chat_event callback ── */ |
||||
|
||||
static void hc_on_chat_event(int type, const uint8_t* data, int len) { |
||||
if (type == 10) return; /* skip status text */ |
||||
char json[1024]; json[0] = '\0'; |
||||
if (type == 1 && data && len >= 10) { /* MSG_RECEIVED */ |
||||
uint8_t cl = data[0]; |
||||
uint64_t author; memcpy(&author, data + 1 + cl, 8); |
||||
snprintf(json, sizeof(json), "{\"event\":\"msg\",\"ch\":\"%.*s\",\"author_id\":\"0x%016llx\"}", (int)cl, data + 1, (unsigned long long)author); |
||||
} else if (type == 5 && data && len >= 2) { /* MEMBERS_CHANGED */ |
||||
uint8_t cl = data[0]; |
||||
snprintf(json, sizeof(json), "{\"event\":\"members_changed\",\"ch\":\"%.*s\"}", (int)cl, data + 1); |
||||
} else if (type == 4 && data && len >= 2) { /* CHANNEL_UPDATED */ |
||||
uint8_t cl = data[0]; |
||||
snprintf(json, sizeof(json), "{\"event\":\"channel_updated\",\"ch\":\"%.*s\"}", (int)cl, data + 1); |
||||
} else if (type == 18 && data && len >= 2) { /* INVITE_RECEIVED */ |
||||
uint8_t cl = data[0]; |
||||
if (len >= 2 + cl) { |
||||
uint8_t nl = data[1 + cl]; |
||||
uint64_t inviter; int off = 2 + cl + nl; |
||||
if (len >= off + 8) { memcpy(&inviter, data + off, 8); |
||||
snprintf(json, sizeof(json), "{\"event\":\"invite_received\",\"ch\":\"%.*s\",\"ch_name\":\"%.*s\",\"from_id\":\"0x%016llx\"}", |
||||
(int)cl, data + 1, (int)nl, data + 2 + cl, (unsigned long long)inviter); |
||||
} |
||||
} |
||||
} |
||||
if (json[0]) hc_broadcast_event(json); |
||||
} |
||||
|
||||
/* ── Write callback ── */ |
||||
|
||||
static void client_write_callback(socket_t fd, void* arg) { |
||||
struct headless_client* cli = (struct headless_client*)arg; |
||||
if (!cli || cli->closing || fd != cli->fd) return; |
||||
if (cli->send_len <= 0) { if (g_hc.ua) uasync_set_socket_write(g_hc.ua, cli->socket_id, 0); return; } |
||||
ssize_t r = write(cli->fd, cli->send_buf + cli->send_offset, (size_t)(cli->send_len - cli->send_offset)); |
||||
if (r < 0 && (errno == EAGAIN || errno == EWOULDBLOCK)) return; |
||||
if (r <= 0) { cli->send_len = 0; cli->send_offset = 0; if (g_hc.ua) uasync_set_socket_write(g_hc.ua, cli->socket_id, 0); return; } |
||||
cli->send_offset += (int)r; |
||||
if (cli->send_offset >= cli->send_len) { cli->send_len = 0; cli->send_offset = 0; if (g_hc.ua) uasync_set_socket_write(g_hc.ua, cli->socket_id, 0); } |
||||
} |
||||
|
||||
/* ── Close client ── */ |
||||
|
||||
static void hc_close_client(struct headless_client* cli) { |
||||
if (!cli || cli->closing) return; |
||||
cli->closing = 1; |
||||
if (g_hc.ua && cli->socket_id) uasync_remove_socket_t(g_hc.ua, cli->fd); |
||||
if (cli->fd >= 0) { close(cli->fd); cli->fd = -1; } |
||||
} |
||||
|
||||
/* ── Read callback ── */ |
||||
|
||||
static void client_read_callback(socket_t fd, void* arg) { |
||||
struct headless_client* cli = (struct headless_client*)arg; |
||||
if (!cli || cli->closing) return; |
||||
char tmp[4096]; |
||||
ssize_t r = read(fd, tmp, sizeof(tmp) - 1); |
||||
if (r <= 0) { hc_close_client(cli); return; } |
||||
for (ssize_t i = 0; i < r; i++) { |
||||
if (tmp[i] == '\n') { |
||||
if (cli->recv_len > 0) { cli->recv_buf[cli->recv_len] = '\0'; hc_handle_command(cli, cli->recv_buf); cli->recv_len = 0; } |
||||
} else if (cli->recv_len < RECV_BUF_SIZE - 1) { |
||||
cli->recv_buf[cli->recv_len++] = tmp[i]; |
||||
} |
||||
} |
||||
} |
||||
|
||||
/* ── Command handlers ── */ |
||||
|
||||
static void hc_handle_ping(struct headless_client* cli, int id, const char* json) { |
||||
(void)json; |
||||
uint64_t uptime = 0; |
||||
(void)uptime; |
||||
char data[128]; snprintf(data, sizeof(data), "{\"pong\":true,\"uptime_sec\":%llu}", (unsigned long long)uptime); |
||||
send_response(cli, id, data, NULL); |
||||
} |
||||
|
||||
static void hc_handle_status(struct headless_client* cli, int id, const char* json) { |
||||
(void)json; |
||||
char buf[8192]; int off = 0; |
||||
off += snprintf(buf + off, sizeof(buf) - off, "\""); |
||||
if (!g_hc.inst) { |
||||
off += snprintf(buf + off, sizeof(buf) - off, "ERROR: no instance"); |
||||
} else { |
||||
int conn_count = g_hc.inst->connections ? queue_entry_count(g_hc.inst->connections) : 0; |
||||
struct NTP_TIME* ntp = &g_hc.inst->ntp; |
||||
off += snprintf(buf + off, sizeof(buf) - off, "NTP: enabled=%s synced=%s offset=%.1fs\\n",
|
||||
ntp->enabled ? "yes" : "no", ntp->synced ? "yes" : "no", ntp->offset_us / 1000000.0); |
||||
off += snprintf(buf + off, sizeof(buf) - off, "Connections: %d active\\n", conn_count); |
||||
uint64_t my_id = g_hc.inst->node_id; |
||||
struct ll_entry* entry = g_hc.inst->connections ? g_hc.inst->connections->head : NULL; |
||||
while (entry) { |
||||
struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data; |
||||
if (ce) { |
||||
char name[64]; chat_core_get_node_name(ce->peer_node_id, name, sizeof(name)); |
||||
off += snprintf(buf + off, sizeof(buf) - off, " [%04llX]->[%04llX] %s ETCP:%s\\n", |
||||
(unsigned long long)(my_id & 0xFFFF), (unsigned long long)(ce->peer_node_id & 0xFFFF), |
||||
name[0] ? name : "?", ce->conn && ce->conn->links_up ? "UP" : "DOWN"); |
||||
} |
||||
entry = entry->next; |
||||
} |
||||
} |
||||
if (off >= (int)sizeof(buf) - 2) off = (int)sizeof(buf) - 2; |
||||
off += snprintf(buf + off, sizeof(buf) - off, "\""); |
||||
send_response(cli, id, buf, NULL); |
||||
} |
||||
|
||||
static void hc_handle_channels(struct headless_client* cli, int id, const char* json) { |
||||
(void)json; |
||||
char buf[16384]; size_t len = 0; |
||||
if (chat_core_get_channels_json(buf, sizeof(buf), &len) < 0) |
||||
send_response(cli, id, NULL, "failed to query channels"); |
||||
else |
||||
send_response(cli, id, buf, NULL); |
||||
} |
||||
|
||||
static void hc_handle_members(struct headless_client* cli, int id, const char* json) { |
||||
char ch[64]; |
||||
if (json_get_str(json, "ch", ch, sizeof(ch)) < 0) { send_response(cli, id, NULL, "missing 'ch' param"); return; } |
||||
char buf[16384]; size_t len = 0; |
||||
if (chat_core_get_members_json(ch, buf, sizeof(buf), &len) < 0) |
||||
send_response(cli, id, NULL, "failed to query members"); |
||||
else |
||||
send_response(cli, id, buf, NULL); |
||||
} |
||||
|
||||
static void hc_handle_messages(struct headless_client* cli, int id, const char* json) { |
||||
char ch[64]; int count = 20, offset = 0; |
||||
if (json_get_str(json, "ch", ch, sizeof(ch)) < 0) { send_response(cli, id, NULL, "missing 'ch' param"); return; } |
||||
count = json_get_int(json, "count", 20); |
||||
offset = json_get_int(json, "offset", 0); |
||||
char buf[16384]; size_t len = 0; |
||||
if (chat_core_get_messages_json(ch, count, offset, buf, sizeof(buf), &len) < 0) |
||||
send_response(cli, id, NULL, "failed to query messages"); |
||||
else |
||||
send_response(cli, id, buf, NULL); |
||||
} |
||||
|
||||
static void hc_handle_send(struct headless_client* cli, int id, const char* json) { |
||||
char ch[64], ct[32], data[4096]; |
||||
if (json_get_str(json, "ch", ch, sizeof(ch)) < 0) { send_response(cli, id, NULL, "missing 'ch' param"); return; } |
||||
if (json_get_str(json, "content_type", ct, sizeof(ct)) < 0) { ct[0] = 't'; ct[1] = 'e'; ct[2] = 'x'; ct[3] = 't'; ct[4] = '\0'; } |
||||
if (json_get_str(json, "data", data, sizeof(data)) < 0) { send_response(cli, id, NULL, "missing 'data' param"); return; } |
||||
struct chat_msg_submit* req = u_calloc(1, sizeof(*req)); |
||||
if (!req) { send_response(cli, id, NULL, "out of memory"); return; } |
||||
snprintf(req->channel_id, sizeof(req->channel_id), "%s", ch); |
||||
snprintf(req->content_type, sizeof(req->content_type), "%s", ct); |
||||
req->data = (uint8_t*)u_strdup(data); |
||||
req->data_len = (uint32_t)strlen(data); |
||||
chat_core_submit_message(req); |
||||
if (req->data) u_free(req->data); |
||||
u_free(req); |
||||
send_response(cli, id, "{\"sent\":true}", NULL); |
||||
} |
||||
|
||||
static void hc_handle_invite(struct headless_client* cli, int id, const char* json) { |
||||
char ch[64]; |
||||
if (json_get_str(json, "ch", ch, sizeof(ch)) < 0) { send_response(cli, id, NULL, "missing 'ch' param"); return; } |
||||
if (!g_hc.inst) { send_response(cli, id, NULL, "no instance"); return; } |
||||
|
||||
uint64_t my_node_id = g_hc.inst->node_id; |
||||
uint8_t my_pubkey[32]; memcpy(my_pubkey, g_hc.inst->my_keys.public_key, 32); |
||||
|
||||
sqlite3* db = chat_core_get_db(); |
||||
if (!db) { send_response(cli, id, NULL, "no database"); return; } |
||||
|
||||
struct InviteData inv; memset(&inv, 0, sizeof(inv)); |
||||
inv.channelId = strtoull(ch, NULL, 10); |
||||
if (inv.channelId == 0) { send_response(cli, id, NULL, "invalid channel_id"); return; } |
||||
memcpy(inv.pubkey, my_pubkey, 32); |
||||
inv.nodeId = my_node_id; |
||||
|
||||
sqlite3_stmt* as = NULL; |
||||
sqlite3_prepare_v2(db, "SELECT family, protocol, address, port, rtt, socket_id FROM node_addresses WHERE node_id=? AND protocol IN (1,2) ORDER BY family, protocol", |
||||
-1, &as, NULL); |
||||
if (as) { |
||||
sqlite3_bind_int64(as, 1, (sqlite3_int64)my_node_id); |
||||
while (sqlite3_step(as) == SQLITE_ROW && inv.addrCount < INVITE_ADDR_MAX) { |
||||
struct InviteAddr* a = &inv.addrs[inv.addrCount++]; |
||||
int fam = sqlite3_column_int(as, 0); |
||||
a->proto = (uint8_t)sqlite3_column_int(as, 1); |
||||
const uint8_t* addr = (const uint8_t*)sqlite3_column_blob(as, 2); |
||||
int addr_len = sqlite3_column_bytes(as, 2); |
||||
a->port = (uint16_t)sqlite3_column_int(as, 3); |
||||
a->socketId = (uint8_t)sqlite3_column_int(as, 5); |
||||
a->family = fam; |
||||
memset(a->address, 0, sizeof(a->address)); |
||||
if (addr && addr_len > 0) memcpy(a->address, addr, addr_len < 16 ? (size_t)addr_len : 16); |
||||
} |
||||
sqlite3_finalize(as); |
||||
} |
||||
|
||||
if (inv.addrCount == 0) { |
||||
/* fallback: collect from sockets directly */ |
||||
struct ETCP_SOCKET* s = g_hc.inst->etcp_sockets; |
||||
while (s && inv.addrCount < INVITE_ADDR_MAX) { |
||||
if (s->local_addr.ss_family == AF_INET) { |
||||
struct sockaddr_in* sin = (struct sockaddr_in*)&s->local_addr; |
||||
struct InviteAddr* a = &inv.addrs[inv.addrCount++]; |
||||
a->family = 4; a->proto = 1; a->socketId = s->sock_id; |
||||
a->port = ntohs(sin->sin_port); memcpy(a->address, &sin->sin_addr, 4); |
||||
} |
||||
s = s->next; |
||||
} |
||||
struct TCP_SOCKET* ts = g_hc.inst->tcp_sockets; |
||||
while (ts && inv.addrCount < INVITE_ADDR_MAX) { |
||||
if (ts->interface_addr.ss_family == AF_INET) { |
||||
struct sockaddr_in* sin = (struct sockaddr_in*)&ts->interface_addr; |
||||
struct InviteAddr* a = &inv.addrs[inv.addrCount++]; |
||||
a->family = 4; a->proto = 2; a->socketId = ts->sock_id; |
||||
a->port = ntohs(sin->sin_port); memcpy(a->address, &sin->sin_addr, 4); |
||||
} |
||||
ts = ts->next; |
||||
} |
||||
} |
||||
|
||||
if (inv.addrCount == 0) { send_response(cli, id, NULL, "no addresses available to build invite link"); return; } |
||||
|
||||
char link[1024]; |
||||
int r = invite_link_encode(&inv, NULL, link, sizeof(link)); |
||||
if (r < 0) { send_response(cli, id, NULL, "failed to encode invite link"); return; } |
||||
|
||||
char resp[1280]; snprintf(resp, sizeof(resp), "{\"link\":\"%s\",\"channel_id\":%llu,\"addrs\":%d}", link, |
||||
(unsigned long long)inv.channelId, (int)inv.addrCount); |
||||
send_response(cli, id, resp, NULL); |
||||
} |
||||
|
||||
static void hc_handle_connect(struct headless_client* cli, int id, const char* json) { |
||||
char link[1024]; |
||||
if (json_get_str(json, "link", link, sizeof(link)) < 0) { send_response(cli, id, NULL, "missing 'link' param"); return; } |
||||
struct InviteData d; char err[256]; |
||||
if (invite_link_decode(link, strlen(link), &d, err, sizeof(err)) < 0) { send_response(cli, id, NULL, err); return; } |
||||
if (!g_hc.inst) { send_response(cli, id, NULL, "no instance"); return; } |
||||
|
||||
DEBUG_INFO((int)DEBUG_CATEGORY_HEADLESS, "headless: connect ch=%llu node=0x%016llx addrs=%d", |
||||
(unsigned long long)d.channelId, (unsigned long long)d.nodeId, (int)d.addrCount); |
||||
|
||||
uint8_t addrs_buf[2048]; |
||||
int addrs_len = invite_serialize_addrs(&d, addrs_buf, sizeof(addrs_buf)); |
||||
if (addrs_len <= 0) { send_response(cli, id, NULL, "serialize addrs failed"); return; } |
||||
|
||||
chat_core_ensure_channel_ready((const char*)(&(char[64]){0})); |
||||
{ char chstr[64]; snprintf(chstr, sizeof(chstr), "%llu", (unsigned long long)d.channelId); |
||||
chat_core_ensure_channel_ready(chstr); } |
||||
|
||||
chat_core_set_my_node_id(g_hc.inst->node_id); |
||||
|
||||
chat_sync_connect_from_invite(d.channelId, d.nodeId, d.pubkey, |
||||
addrs_buf, d.addrCount, addrs_len, d.password_len ? d.password : NULL); |
||||
|
||||
char resp[512]; snprintf(resp, sizeof(resp), |
||||
"{\"channel_id\":%llu,\"node_id\":\"0x%016llx\",\"addrs\":%d,\"connecting\":true}", |
||||
(unsigned long long)d.channelId, (unsigned long long)d.nodeId, (int)d.addrCount); |
||||
send_response(cli, id, resp, NULL); |
||||
} |
||||
|
||||
static void hc_handle_create_channel(struct headless_client* cli, int id, const char* json) { |
||||
char name[256]; |
||||
if (json_get_str(json, "name", name, sizeof(name)) < 0) { send_response(cli, id, NULL, "missing 'name' param"); return; } |
||||
chat_core_create_channel_auto(name); |
||||
send_response(cli, id, "{\"created\":true}", NULL); |
||||
} |
||||
|
||||
static void hc_handle_subscribe(struct headless_client* cli, int id, const char* json) { |
||||
int enable = json_get_int(json, "enable", 1); |
||||
cli->subscribed = enable; |
||||
char resp[64]; snprintf(resp, sizeof(resp), "{\"subscribed\":%s}", enable ? "true" : "false"); |
||||
send_response(cli, id, resp, NULL); |
||||
} |
||||
|
||||
static void hc_handle_quit(struct headless_client* cli, int id, const char* json) { |
||||
(void)json; |
||||
send_response(cli, id, "{\"bye\":true}", NULL); |
||||
cli_send(cli, "{\"status\":\"closing\"}\n", 20); |
||||
hc_close_client(cli); |
||||
} |
||||
|
||||
/* ── Command dispatch ── */ |
||||
|
||||
typedef void (*hc_cmd_fn)(struct headless_client* cli, int id, const char* json); |
||||
|
||||
static const char* cmd_names[] = { "ping", "status", "channels", "members", "messages", "send", "invite", "connect", "create_channel", "subscribe", "quit", NULL }; |
||||
static hc_cmd_fn cmd_handlers[] = { hc_handle_ping, hc_handle_status, hc_handle_channels, hc_handle_members, hc_handle_messages, hc_handle_send, hc_handle_invite, hc_handle_connect, hc_handle_create_channel, hc_handle_subscribe, hc_handle_quit }; |
||||
|
||||
void hc_handle_command(struct headless_client* cli, const char* json) { |
||||
if (!cli || !json) return; |
||||
int id = json_get_int(json, "id", 0); |
||||
const char* cs = json_get_cmd(json); |
||||
if (!cs) { send_response(cli, id, NULL, "missing cmd"); return; } |
||||
const char* ce = json_get_cmd_end(cs); |
||||
size_t clen = (size_t)(ce - cs); |
||||
if (clen == 0 || clen > 31) { send_response(cli, id, NULL, "invalid cmd"); return; } |
||||
char cmd[32]; memcpy(cmd, cs, clen); cmd[clen] = '\0'; |
||||
for (int i = 0; cmd_names[i]; i++) { |
||||
if (strcmp(cmd, cmd_names[i]) == 0) { cmd_handlers[i](cli, id, json); return; } |
||||
} |
||||
send_response(cli, id, NULL, "unknown command"); |
||||
} |
||||
|
||||
/* ── Accept callback ── */ |
||||
|
||||
static void accept_callback(socket_t fd, void* arg) { |
||||
(void)arg; |
||||
struct sockaddr_storage addr; socklen_t alen = sizeof(addr); |
||||
socket_t cfd = accept(fd, (struct sockaddr*)&addr, &alen); |
||||
if (cfd < 0) { if (errno != EAGAIN && errno != EWOULDBLOCK) DEBUG_ERROR((int)DEBUG_CATEGORY_HEADLESS, "headless: accept failed: %s", strerror(errno)); return; } |
||||
|
||||
int flags = fcntl(cfd, F_GETFL, 0); if (flags >= 0) fcntl(cfd, F_SETFL, flags | O_NONBLOCK); |
||||
|
||||
if (g_hc.client_count >= MAX_CLIENTS) { close(cfd); DEBUG_WARN((int)DEBUG_CATEGORY_HEADLESS, "headless: max clients reached"); return; } |
||||
|
||||
struct headless_client* cli = u_calloc(1, sizeof(*cli)); |
||||
if (!cli) { close(cfd); DEBUG_ERROR((int)DEBUG_CATEGORY_HEADLESS, "headless: failed to allocate client"); return; } |
||||
cli->fd = cfd; |
||||
|
||||
cli->socket_id = uasync_add_socket_t(g_hc.ua, cfd, client_read_callback, client_write_callback, NULL, cli); |
||||
if (!cli->socket_id) { u_free(cli); close(cfd); DEBUG_ERROR((int)DEBUG_CATEGORY_HEADLESS, "headless: failed to register client fd"); return; } |
||||
|
||||
cli->next = g_hc.clients; g_hc.clients = cli; g_hc.client_count++; |
||||
DEBUG_INFO((int)DEBUG_CATEGORY_HEADLESS, "headless: client connected fd=%d total=%d", (int)cfd, g_hc.client_count); |
||||
} |
||||
|
||||
/* ── Periodic cleanup ── */ |
||||
|
||||
static void hc_cleanup(void* arg) { |
||||
(void)arg; |
||||
struct headless_client** prev = &g_hc.clients; |
||||
while (*prev) { |
||||
if ((*prev)->closing) { |
||||
struct headless_client* dead = *prev; *prev = dead->next; |
||||
if (dead->socket_id && g_hc.ua) uasync_remove_socket_t(g_hc.ua, dead->fd); |
||||
if (dead->fd >= 0) close(dead->fd); |
||||
u_free(dead); g_hc.client_count--; |
||||
} else prev = &(*prev)->next; |
||||
} |
||||
if (g_hc.ua) uasync_set_timeout(g_hc.ua, 5000, NULL, hc_cleanup, "hc_cleanup"); |
||||
} |
||||
|
||||
/* ── Init / Destroy ── */ |
||||
|
||||
int chat_headless_control_init(struct UASYNC* ua, struct UTUN_INSTANCE* inst, |
||||
const char* bind_ip, int port) { |
||||
if (!ua || !bind_ip || port <= 0 || port > 65535) return -1; |
||||
if (g_hc.running) return 0; |
||||
|
||||
memset(&g_hc, 0, sizeof(g_hc)); |
||||
g_hc.ua = ua; g_hc.inst = inst; |
||||
|
||||
g_hc.listen_fd = socket(AF_INET, SOCK_STREAM, 0); |
||||
if (g_hc.listen_fd < 0) { DEBUG_ERROR((int)DEBUG_CATEGORY_HEADLESS, "headless: socket() failed: %s", strerror(errno)); return -1; } |
||||
|
||||
int reuse = 1; setsockopt(g_hc.listen_fd, SOL_SOCKET, SO_REUSEADDR, &reuse, sizeof(reuse)); |
||||
int flags = fcntl(g_hc.listen_fd, F_GETFL, 0); if (flags >= 0) fcntl(g_hc.listen_fd, F_SETFL, flags | O_NONBLOCK); |
||||
|
||||
struct sockaddr_in sin; memset(&sin, 0, sizeof(sin)); |
||||
sin.sin_family = AF_INET; sin.sin_port = htons((uint16_t)port); |
||||
if (inet_pton(AF_INET, bind_ip, &sin.sin_addr) != 1) { close(g_hc.listen_fd); g_hc.listen_fd = -1; return -1; } |
||||
|
||||
if (bind(g_hc.listen_fd, (struct sockaddr*)&sin, sizeof(sin)) < 0) { |
||||
DEBUG_ERROR((int)DEBUG_CATEGORY_HEADLESS, "headless: bind(%s:%d) failed: %s", bind_ip, port, strerror(errno)); |
||||
close(g_hc.listen_fd); g_hc.listen_fd = -1; return -1; |
||||
} |
||||
if (listen(g_hc.listen_fd, 5) < 0) { close(g_hc.listen_fd); g_hc.listen_fd = -1; return -1; } |
||||
|
||||
g_hc.listen_sock_id = uasync_add_socket_t(ua, g_hc.listen_fd, accept_callback, NULL, NULL, NULL); |
||||
if (!g_hc.listen_sock_id) { close(g_hc.listen_fd); g_hc.listen_fd = -1; return -1; } |
||||
|
||||
chat_event_set_handler(hc_on_chat_event); |
||||
uasync_set_timeout(ua, 5000, NULL, hc_cleanup, "hc_cleanup"); |
||||
|
||||
g_hc.running = 1; |
||||
DEBUG_INFO((int)DEBUG_CATEGORY_HEADLESS, "headless: listening on %s:%d", bind_ip, port); |
||||
return 0; |
||||
} |
||||
|
||||
void chat_headless_control_destroy(void) { |
||||
if (!g_hc.running) return; |
||||
g_hc.running = 0; |
||||
chat_event_set_handler(NULL); |
||||
if (g_hc.listen_sock_id && g_hc.ua) uasync_remove_socket_t(g_hc.ua, g_hc.listen_fd); |
||||
if (g_hc.listen_fd >= 0) { close(g_hc.listen_fd); g_hc.listen_fd = -1; } |
||||
struct headless_client* c = g_hc.clients; |
||||
while (c) { struct headless_client* n = c->next; if (c->fd >= 0) close(c->fd); u_free(c); c = n; } |
||||
g_hc.clients = NULL; g_hc.client_count = 0; |
||||
DEBUG_INFO((int)DEBUG_CATEGORY_HEADLESS, "headless: destroyed"); |
||||
} |
||||
@ -0,0 +1,25 @@
|
||||
/*
|
||||
* chat_headless_control.h — TCP control socket for headless chat CLI |
||||
* |
||||
* JSON line protocol: |
||||
* Request: {"id":N,"cmd":"<command>",...params} |
||||
* Response: {"id":N,"ok":true,"data":{...}} | {"id":N,"ok":false,"error":"..."} |
||||
* Event: {"event":"<type>",...} |
||||
* |
||||
* Commands: ping, status, channels, members, messages, send, invite, connect, |
||||
* create_channel, subscribe, quit |
||||
*/ |
||||
#ifndef CHAT_HEADLESS_CONTROL_H |
||||
#define CHAT_HEADLESS_CONTROL_H |
||||
|
||||
#include <stdint.h> |
||||
#include <stddef.h> |
||||
|
||||
struct UASYNC; |
||||
struct UTUN_INSTANCE; |
||||
|
||||
int chat_headless_control_init(struct UASYNC* ua, struct UTUN_INSTANCE* inst, |
||||
const char* bind_ip, int port); |
||||
void chat_headless_control_destroy(void); |
||||
|
||||
#endif /* CHAT_HEADLESS_CONTROL_H */ |
||||
@ -0,0 +1,176 @@
|
||||
/*
|
||||
* invite_link.c — encode/decode invite-ссылок utun://
|
||||
* |
||||
* Перенесено из tools/chatgui-android/libutun_lite/invite_link_c.c |
||||
* Формат совместим с десктопной Qt-версией (tools/chatgui/src/invite_link.cpp). |
||||
*/ |
||||
#include "invite_link.h" |
||||
#include "../transport_layer/secure_channel.h" |
||||
#include <string.h> |
||||
#include <stdlib.h> |
||||
#include <stdio.h> |
||||
|
||||
#define INVITE_PREFIX "utun://"
|
||||
#define INVITE_PREFIX_LEN 7 |
||||
|
||||
static const char base64_table[] = "ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz0123456789+/"; |
||||
|
||||
static int base64_char_val(char c) { |
||||
if (c >= 'A' && c <= 'Z') return c - 'A'; |
||||
if (c >= 'a' && c <= 'z') return c - 'a' + 26; |
||||
if (c >= '0' && c <= '9') return c - '0' + 52; |
||||
if (c == '+') return 62; |
||||
if (c == '/') return 63; |
||||
return -1; |
||||
} |
||||
|
||||
static int base64_decode(const char* in, size_t in_len, uint8_t* out, size_t out_cap) { |
||||
size_t opos = 0; int group[4]; int gi = 0; |
||||
for (size_t i = 0; i < in_len && in[i] != '='; i++) { |
||||
int v = base64_char_val(in[i]); |
||||
if (v < 0) return -1; |
||||
group[gi++] = v; |
||||
if (gi == 4) { |
||||
out[opos++] = (uint8_t)((group[0] << 2) | (group[1] >> 4)); |
||||
if (opos >= out_cap) return -1; |
||||
out[opos++] = (uint8_t)(((group[1] & 0xF) << 4) | (group[2] >> 2)); |
||||
if (opos >= out_cap) return -1; |
||||
out[opos++] = (uint8_t)(((group[2] & 0x3) << 6) | group[3]); |
||||
if (opos >= out_cap) return -1; |
||||
gi = 0; |
||||
} |
||||
} |
||||
if (gi == 2) { out[opos++] = (uint8_t)((group[0] << 2) | (group[1] >> 4)); } |
||||
else if (gi == 3) { |
||||
out[opos++] = (uint8_t)((group[0] << 2) | (group[1] >> 4)); |
||||
out[opos++] = (uint8_t)(((group[1] & 0xF) << 4) | (group[2] >> 2)); |
||||
} |
||||
return (int)opos; |
||||
} |
||||
|
||||
int invite_link_decode(const char* link, size_t link_len, struct InviteData* out, |
||||
char* error_buf, size_t error_buf_size) { |
||||
if (!link || !out) { if (error_buf && error_buf_size) snprintf(error_buf, error_buf_size, "null argument"); return -1; } |
||||
if (link_len < INVITE_PREFIX_LEN || strncmp(link, INVITE_PREFIX, INVITE_PREFIX_LEN) != 0) { |
||||
if (error_buf && error_buf_size) snprintf(error_buf, error_buf_size, "invalid prefix (expected utun://)"); return -1; |
||||
} |
||||
const char* b64 = link + INVITE_PREFIX_LEN; |
||||
size_t b64_len = link_len - INVITE_PREFIX_LEN; |
||||
uint8_t raw[4096]; |
||||
int raw_len = base64_decode(b64, b64_len, raw, sizeof(raw)); |
||||
if (raw_len < 0) { if (error_buf && error_buf_size) snprintf(error_buf, error_buf_size, "base64 decode failed"); return -1; } |
||||
if (raw_len < 11) { if (error_buf && error_buf_size) snprintf(error_buf, error_buf_size, "blob too short (%d bytes)", raw_len); return -1; } |
||||
|
||||
int off = 0; |
||||
uint8_t ver = raw[off++]; |
||||
if (ver != INVITE_LINK_VERSION && ver != 0x02) { |
||||
if (error_buf && error_buf_size) snprintf(error_buf, error_buf_size, "unsupported version 0x%02x", ver); return -1; |
||||
} |
||||
out->password[0] = '\0'; out->password_len = 0; |
||||
|
||||
if (ver == 0x02) { |
||||
if (off + 1 > raw_len) { if (error_buf && error_buf_size) snprintf(error_buf, error_buf_size, "truncated at pass_len"); return -1; } |
||||
uint8_t plen = raw[off++]; |
||||
if (plen > 0) { |
||||
if (plen > INVITE_PASS_MAX - 1) { if (error_buf && error_buf_size) snprintf(error_buf, error_buf_size, "password too long %d", plen); return -1; } |
||||
if (off + plen > raw_len) { if (error_buf && error_buf_size) snprintf(error_buf, error_buf_size, "truncated at password"); return -1; } |
||||
memcpy(out->password, raw + off, plen); out->password[plen] = '\0'; out->password_len = plen; |
||||
off += plen; |
||||
} |
||||
} |
||||
|
||||
if (off + 8 > raw_len) { if (error_buf && error_buf_size) snprintf(error_buf, error_buf_size, "truncated at channel_id"); return -1; } |
||||
uint64_t chId = 0; |
||||
for (int i = 0; i < 8; i++) chId = (chId << 8) | raw[off++]; |
||||
out->channelId = chId; |
||||
memset(out->pubkey, 0, INVITE_PUBKEY_SIZE); |
||||
out->addrCount = 0; |
||||
|
||||
int first_block = 1; |
||||
while (off + 1 <= raw_len) { |
||||
uint8_t header = raw[off++]; |
||||
int cnt = (header & 0x03) + 1; |
||||
if (cnt < 1 || cnt > 4) { if (error_buf && error_buf_size) snprintf(error_buf, error_buf_size, "invalid addr count %d", cnt); return -1; } |
||||
if (off + 32 > raw_len) { if (error_buf && error_buf_size) snprintf(error_buf, error_buf_size, "truncated at pubkey"); return -1; } |
||||
if (first_block) { memcpy(out->pubkey, raw + off, INVITE_PUBKEY_SIZE); out->nodeId = sc_derive_node_id_from_pubkey(out->pubkey); first_block = 0; } |
||||
off += 32; |
||||
for (int j = 0; j < cnt; j++) { |
||||
int is_v6 = (header >> (2 + j)) & 1; |
||||
int ip_len = is_v6 ? 16 : 4; |
||||
if (off + 2 + ip_len + 2 > raw_len) { if (error_buf && error_buf_size) snprintf(error_buf, error_buf_size, "truncated at addr %d", j); return -1; } |
||||
if (out->addrCount >= INVITE_ADDR_MAX) break; |
||||
struct InviteAddr* a = &out->addrs[out->addrCount++]; |
||||
a->socketId = raw[off++]; a->proto = raw[off++]; |
||||
a->family = is_v6 ? 6 : 4; |
||||
memcpy(a->address, raw + off, (size_t)ip_len); off += ip_len; |
||||
a->port = (uint16_t)((raw[off] << 8) | raw[off + 1]); off += 2; |
||||
} |
||||
} |
||||
return 0; |
||||
} |
||||
|
||||
int invite_serialize_addrs(const struct InviteData* data, uint8_t* buf, size_t buf_size) { |
||||
if (!data || !buf) return -1; |
||||
size_t needed = 0; |
||||
for (int i = 0; i < (int)data->addrCount; i++) needed += (data->addrs[i].family == 6) ? 21U : 9U; |
||||
if (needed > buf_size) return -1; |
||||
size_t pos = 0; |
||||
for (int i = 0; i < (int)data->addrCount; i++) { |
||||
const struct InviteAddr* a = &data->addrs[i]; |
||||
int ip_len = (a->family == 6) ? 16 : 4; |
||||
buf[pos++] = (uint8_t)a->family; buf[pos++] = a->socketId; buf[pos++] = a->proto; |
||||
memcpy(buf + pos, a->address, (size_t)ip_len); pos += ip_len; |
||||
buf[pos++] = (uint8_t)(a->port >> 8); buf[pos++] = (uint8_t)(a->port & 0xFF); |
||||
} |
||||
return (int)pos; |
||||
} |
||||
|
||||
int invite_link_encode(const struct InviteData* data, const char* password, char* out, size_t out_size) { |
||||
if (!data || !out || data->pubkey[0] == 0) return -1; |
||||
if (data->addrCount == 0) return -1; |
||||
int has_pass = (password && password[0]) ? 1 : 0; |
||||
uint8_t raw[4096]; size_t pos = 0; |
||||
raw[pos++] = has_pass ? 0x02 : INVITE_LINK_VERSION; |
||||
if (has_pass) { |
||||
size_t plen = strlen(password); |
||||
if (plen > INVITE_PASS_MAX - 1) return -1; |
||||
raw[pos++] = (uint8_t)plen; memcpy(raw + pos, password, plen); pos += plen; |
||||
} |
||||
for (int i = 7; i >= 0; i--) raw[pos++] = (uint8_t)((data->channelId >> (i * 8)) & 0xFF); |
||||
uint8_t pk[32]; memcpy(pk, data->pubkey, 32); |
||||
|
||||
int ia = 0; |
||||
while (ia < (int)data->addrCount) { |
||||
int rem = (int)data->addrCount - ia; |
||||
int cnt = (rem > 4) ? 4 : rem; |
||||
uint8_t header = (uint8_t)((cnt - 1) & 0x03); |
||||
for (int j = 0; j < cnt; j++) if (data->addrs[ia + j].family == 6) header |= (uint8_t)(1 << (2 + j)); |
||||
raw[pos++] = header; memcpy(raw + pos, pk, 32); pos += 32; |
||||
for (int j = 0; j < cnt; j++) { |
||||
const struct InviteAddr* a = &data->addrs[ia + j]; |
||||
int ip_len = (a->family == 6) ? 16 : 4; |
||||
raw[pos++] = a->socketId; raw[pos++] = a->proto; |
||||
memcpy(raw + pos, a->address, (size_t)ip_len); pos += (size_t)ip_len; |
||||
raw[pos++] = (uint8_t)(a->port >> 8); raw[pos++] = (uint8_t)(a->port & 0xFF); |
||||
} |
||||
ia += cnt; |
||||
} |
||||
|
||||
if (out_size < 12 + (pos * 4 + 2) / 3 + 1) return -1; |
||||
memcpy(out, INVITE_PREFIX, INVITE_PREFIX_LEN); |
||||
size_t opos = INVITE_PREFIX_LEN; |
||||
size_t i = 0; |
||||
while (i < pos) { |
||||
uint32_t val = (uint32_t)raw[i] << 16; |
||||
val |= (i + 1 < pos) ? (uint32_t)raw[i + 1] << 8 : 0; |
||||
val |= (i + 2 < pos) ? (uint32_t)raw[i + 2] : 0; |
||||
out[opos++] = base64_table[(val >> 18) & 0x3F]; |
||||
out[opos++] = base64_table[(val >> 12) & 0x3F]; |
||||
out[opos++] = (i + 1 < pos) ? base64_table[(val >> 6) & 0x3F] : '='; |
||||
out[opos++] = (i + 2 < pos) ? base64_table[val & 0x3F] : '='; |
||||
i += 3; |
||||
} |
||||
if (opos >= out_size) return -1; |
||||
out[opos] = '\0'; |
||||
return (int)(opos - INVITE_PREFIX_LEN); |
||||
} |
||||
@ -0,0 +1,54 @@
|
||||
/*
|
||||
* invite_link.h — invite-ссылки utun:// для каналов
|
||||
* |
||||
* Формат: utun:// + base64(version | password? | channel_id | [header | pubkey | addrs]+)
|
||||
* Совместим с десктопной (Qt) и Android (Kotlin) версиями. |
||||
*/ |
||||
#ifndef INVITE_LINK_H |
||||
#define INVITE_LINK_H |
||||
|
||||
#include <stdint.h> |
||||
#include <stddef.h> |
||||
|
||||
#define INVITE_LINK_VERSION 0x01 |
||||
#define INVITE_PUBKEY_SIZE 32 |
||||
#define INVITE_ADDR_MAX 32 |
||||
#define INVITE_PASS_MAX 128 |
||||
|
||||
#define INVITE_PROTO_UDP 0x01 |
||||
#define INVITE_PROTO_TCP 0x02 |
||||
#define INVITE_PROTO_BOTH 0x03 |
||||
|
||||
#ifdef __cplusplus |
||||
extern "C" { |
||||
#endif |
||||
|
||||
struct InviteAddr { |
||||
int family; |
||||
uint8_t address[16]; |
||||
uint8_t socketId; |
||||
uint8_t proto; |
||||
uint16_t port; |
||||
}; |
||||
|
||||
struct InviteData { |
||||
uint64_t channelId; |
||||
uint8_t pubkey[INVITE_PUBKEY_SIZE]; |
||||
uint64_t nodeId; |
||||
uint8_t addrCount; |
||||
struct InviteAddr addrs[INVITE_ADDR_MAX]; |
||||
char password[INVITE_PASS_MAX]; |
||||
uint8_t password_len; |
||||
}; |
||||
|
||||
int invite_link_decode(const char* link, size_t link_len, struct InviteData* out, |
||||
char* error_buf, size_t error_buf_size); |
||||
|
||||
int invite_link_encode(const struct InviteData* data, const char* password, char* out, size_t out_size); |
||||
|
||||
int invite_serialize_addrs(const struct InviteData* data, uint8_t* buf, size_t buf_size); |
||||
|
||||
#ifdef __cplusplus |
||||
} |
||||
#endif |
||||
#endif /* INVITE_LINK_H */ |
||||
@ -0,0 +1,109 @@
|
||||
""" |
||||
chat_client.py — async TCP JSON-RPC client for uTun headless chat control. |
||||
|
||||
Usage: |
||||
async with ChatClient(port=9999) as cli: |
||||
await cli.create_channel("Test") |
||||
channels = await cli.channels() |
||||
await cli.send(channels[0]["id"], "hello") |
||||
""" |
||||
|
||||
import asyncio |
||||
import json |
||||
|
||||
|
||||
class ChatClientError(Exception): |
||||
"""Error from chat client (timeout, refused, protocol, server error).""" |
||||
|
||||
|
||||
class ChatClient: |
||||
def __init__(self, host="127.0.0.1", port=9999, timeout=5.0): |
||||
self.host = host |
||||
self.port = port |
||||
self.timeout = timeout |
||||
self._connected = False |
||||
|
||||
async def connect(self): |
||||
self._connected = True |
||||
|
||||
async def disconnect(self): |
||||
self._connected = False |
||||
|
||||
async def __aenter__(self): |
||||
await self.connect() |
||||
return self |
||||
|
||||
async def __aexit__(self, *args): |
||||
await self.disconnect() |
||||
|
||||
# ── Internal: one-shot request ── |
||||
|
||||
async def _request(self, cmd, **params): |
||||
if not self._connected: |
||||
raise ChatClientError("not connected") |
||||
|
||||
reader, writer = await asyncio.wait_for( |
||||
asyncio.open_connection(self.host, self.port), |
||||
timeout=self.timeout, |
||||
) |
||||
|
||||
msg = {"id": 1, "cmd": cmd} |
||||
msg.update(params) |
||||
writer.write((json.dumps(msg, ensure_ascii=True) + "\n").encode()) |
||||
await writer.drain() |
||||
|
||||
try: |
||||
line = await asyncio.wait_for( |
||||
reader.readline(), timeout=self.timeout |
||||
) |
||||
except asyncio.TimeoutError: |
||||
writer.close() |
||||
raise ChatClientError(f"timeout waiting for '{cmd}' response") |
||||
|
||||
writer.close() |
||||
if not line: |
||||
raise ChatClientError(f"empty response for '{cmd}'") |
||||
|
||||
try: |
||||
obj = json.loads(line.decode("utf-8", errors="replace")) |
||||
except json.JSONDecodeError: |
||||
raw = line.decode("utf-8", errors="replace") |
||||
raise ChatClientError(f"invalid JSON for '{cmd}' len={len(raw)} tail=...{raw[-50:] if len(raw)>50 else raw}") |
||||
|
||||
if not obj.get("ok"): |
||||
raise ChatClientError(obj.get("error", "unknown error")) |
||||
|
||||
return obj.get("data", {}) |
||||
|
||||
# ── Commands ── |
||||
|
||||
async def ping(self): |
||||
return await self._request("ping") |
||||
|
||||
async def status(self): |
||||
return await self._request("status") |
||||
|
||||
async def channels(self): |
||||
return await self._request("channels") |
||||
|
||||
async def members(self, ch_id): |
||||
return await self._request("members", ch=str(ch_id)) |
||||
|
||||
async def messages(self, ch_id, count=20, offset=0): |
||||
return await self._request( |
||||
"messages", ch=str(ch_id), count=int(count), offset=int(offset) |
||||
) |
||||
|
||||
async def send(self, ch_id, text): |
||||
return await self._request( |
||||
"send", ch=str(ch_id), content_type="text", data=str(text) |
||||
) |
||||
|
||||
async def invite(self, ch_id): |
||||
return await self._request("invite", ch=str(ch_id)) |
||||
|
||||
async def connect_channel(self, link): |
||||
return await self._request("connect", link=str(link)) |
||||
|
||||
async def create_channel(self, name): |
||||
return await self._request("create_channel", name=str(name)) |
||||
@ -0,0 +1,268 @@
|
||||
#!/usr/bin/env python3 |
||||
""" |
||||
chat_integration_test.py — integration test: 2 chat nodes over UDP, sync messages. |
||||
|
||||
Запускает два процесса utun в фоне, создаёт канал, пишет 4 сообщения, |
||||
подключает второй узел по invite-ссылке, проверяет members и messages. |
||||
|
||||
Usage: |
||||
python3 tools/chat_integration_test.py |
||||
""" |
||||
|
||||
import asyncio |
||||
import json |
||||
import os |
||||
import socket |
||||
import sys |
||||
import tempfile |
||||
import time |
||||
|
||||
sys.path.insert(0, os.path.dirname(__file__)) |
||||
from chat_client import ChatClient, ChatClientError |
||||
|
||||
UTUN_BIN = os.path.join(os.path.dirname(__file__), "..", "src", "utun") |
||||
READY_TIMEOUT = 3.0 # max wait for utun ready |
||||
SYNC_TIMEOUT = 8.0 # max wait for join + sync |
||||
REQUEST_TIMEOUT = 5.0 # per-request timeout |
||||
TOTAL_TIMEOUT = 25.0 # overall test timeout |
||||
|
||||
|
||||
def find_free_port(): |
||||
s = socket.socket(socket.AF_INET, socket.SOCK_STREAM) |
||||
s.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) |
||||
s.bind(("127.0.0.1", 0)) |
||||
port = s.getsockname()[1] |
||||
s.close() |
||||
return port |
||||
|
||||
|
||||
def write_config(path, etcp_port, ctrl_port, db_subdir): |
||||
content = f"""[global] |
||||
db_path={db_subdir} |
||||
|
||||
[server: srv] |
||||
addr=127.0.0.1:{etcp_port} |
||||
type=public |
||||
|
||||
[chatserver] |
||||
db_path={db_subdir} |
||||
headless_control_bind=127.0.0.1:{ctrl_port} |
||||
|
||||
[allowed_keys] |
||||
allow_all=1 |
||||
""" |
||||
with open(path, "w") as f: |
||||
f.write(content) |
||||
|
||||
|
||||
async def kill_proc(proc, label): |
||||
if proc is None or proc.returncode is not None: |
||||
return |
||||
try: |
||||
proc.terminate() |
||||
try: |
||||
await asyncio.wait_for(proc.wait(), timeout=3.0) |
||||
except asyncio.TimeoutError: |
||||
proc.kill() |
||||
await proc.wait() |
||||
except ProcessLookupError: |
||||
pass |
||||
|
||||
|
||||
async def wait_ready(cli, timeout=READY_TIMEOUT): |
||||
deadline = time.monotonic() + timeout |
||||
while time.monotonic() < deadline: |
||||
try: |
||||
await cli.ping() |
||||
return True |
||||
except ChatClientError: |
||||
await asyncio.sleep(0.1) |
||||
return False |
||||
|
||||
|
||||
def check(name, expr, detail=""): |
||||
if not expr: |
||||
detail = f" ({detail})" if detail else "" |
||||
raise AssertionError(f"FAIL: {name}{detail}") |
||||
print(f" OK: {name}") |
||||
|
||||
|
||||
async def main(): |
||||
etcp_a = find_free_port() |
||||
etcp_b = find_free_port() |
||||
ctrl_a = find_free_port() |
||||
ctrl_b = find_free_port() |
||||
|
||||
print(f"ports: etcp={etcp_a},{etcp_b} ctrl={ctrl_a},{ctrl_b}") |
||||
|
||||
tmpdir = tempfile.mkdtemp(prefix="utun_chat_test_") |
||||
db_a = os.path.join(tmpdir, "db_a") |
||||
db_b = os.path.join(tmpdir, "db_b") |
||||
os.makedirs(db_a, exist_ok=True) |
||||
os.makedirs(db_b, exist_ok=True) |
||||
|
||||
config_a = os.path.join(tmpdir, "a.conf") |
||||
config_b = os.path.join(tmpdir, "b.conf") |
||||
write_config(config_a, etcp_a, ctrl_a, db_a) |
||||
write_config(config_b, etcp_b, ctrl_b, db_b) |
||||
|
||||
proc_a = None |
||||
proc_b = None |
||||
|
||||
try: |
||||
# ── Start utun processes ── |
||||
print("\n--- Starting utun ---") |
||||
log_a = os.path.join(tmpdir, "utun_a.log") |
||||
log_b = os.path.join(tmpdir, "utun_b.log") |
||||
pid_a = os.path.join(tmpdir, "utun_a.pid") |
||||
pid_b = os.path.join(tmpdir, "utun_b.pid") |
||||
|
||||
proc_a = await asyncio.create_subprocess_exec( |
||||
UTUN_BIN, "-f", "-p", pid_a, "-l", log_a, "-c", config_a, |
||||
stdout=asyncio.subprocess.DEVNULL, stderr=asyncio.subprocess.DEVNULL, |
||||
) |
||||
proc_b = await asyncio.create_subprocess_exec( |
||||
UTUN_BIN, "-f", "-p", pid_b, "-l", log_b, "-c", config_b, |
||||
stdout=asyncio.subprocess.DEVNULL, stderr=asyncio.subprocess.DEVNULL, |
||||
) |
||||
|
||||
print(f" proc_a pid={proc_a.pid} proc_b pid={proc_b.pid}") |
||||
|
||||
await asyncio.sleep(0.5) |
||||
|
||||
# ── Open connections ── |
||||
print("\n--- Connecting to headless control ---") |
||||
async with ChatClient(port=ctrl_a, timeout=REQUEST_TIMEOUT) as cli_a: |
||||
if not await wait_ready(cli_a): |
||||
raise RuntimeError(f"node A not ready after {READY_TIMEOUT}s") |
||||
print(f" node A ready") |
||||
|
||||
async with ChatClient(port=ctrl_b, timeout=REQUEST_TIMEOUT) as cli_b: |
||||
if not await wait_ready(cli_b): |
||||
raise RuntimeError(f"node B not ready after {READY_TIMEOUT}s") |
||||
print(f" node B ready") |
||||
|
||||
# ── Create channel ── |
||||
print("\n--- Create channel ---") |
||||
await cli_a.create_channel("TestGroup") |
||||
channels = await cli_a.channels() |
||||
check("channel created", isinstance(channels, list) and len(channels) == 1, |
||||
f"channels={json.dumps(channels)}") |
||||
ch_id = str(channels[0]["id"]) |
||||
print(f" channel_id={ch_id} name={channels[0]['name']}") |
||||
|
||||
# ── Write 4 messages ── |
||||
print("\n--- Write messages ---") |
||||
for i in range(4): |
||||
await cli_a.send(ch_id, f"Message {i+1}") |
||||
|
||||
msgs_a = await cli_a.messages(ch_id, count=10) |
||||
check("4 messages on A", isinstance(msgs_a, list) and len(msgs_a) == 4, |
||||
f"got {len(msgs_a) if isinstance(msgs_a, list) else '?'}") |
||||
for i, m in enumerate(msgs_a): |
||||
print(f" [#{i+1}] ts={m.get('ts','?')} author={m.get('author_id','?')}") |
||||
|
||||
# ── Invite ── |
||||
print("\n--- Invite link ---") |
||||
invite = await cli_a.invite(ch_id) |
||||
link = invite.get("link", "") |
||||
check("invite link", isinstance(link, str) and link.startswith("utun://"), |
||||
f"link={'...' + link[-20:] if link else 'MISSING'}") |
||||
print(f" link={link[:80]}...") |
||||
|
||||
# ── Connect node B ── |
||||
print("\n--- Connect node B ---") |
||||
conn = await cli_b.connect_channel(link) |
||||
check("connect accepted", isinstance(conn, dict) and conn.get("connecting"), |
||||
f"resp={json.dumps(conn)}") |
||||
print(f" connecting={conn.get('connecting')}") |
||||
|
||||
# Wait for join + sync |
||||
print(f"\n--- Wait sync (max {SYNC_TIMEOUT}s) ---") |
||||
deadline = time.monotonic() + SYNC_TIMEOUT |
||||
member_count = 0 |
||||
while time.monotonic() < deadline: |
||||
try: |
||||
members = await cli_b.members(ch_id) |
||||
member_count = len(members) if isinstance(members, list) else 0 |
||||
if member_count >= 2: |
||||
break |
||||
except ChatClientError: |
||||
pass |
||||
await asyncio.sleep(0.15) |
||||
check("2 members on B", member_count >= 2, |
||||
f"got {member_count} members after {SYNC_TIMEOUT}s") |
||||
|
||||
# Small pause to let db_sync fully settle |
||||
await asyncio.sleep(0.5) |
||||
|
||||
# ── Verify members on B ── |
||||
members = await cli_b.members(ch_id) |
||||
print(f"\n--- Members on B ({len(members)}) ---") |
||||
for m in members: |
||||
connected = "✓" if m.get("connected") else "✗" |
||||
online = "✓" if m.get("online") else "✗" |
||||
print(f" {m['node_id']} name={m.get('name','?')} online={online} connected={connected}") |
||||
check("owner online", any(m.get("online") for m in members), |
||||
f"members={json.dumps(members)}") |
||||
check("owner connected", any(m.get("connected") for m in members), |
||||
f"members={json.dumps(members)}") |
||||
|
||||
# ── Verify messages on B ── |
||||
if proc_b.returncode is not None: |
||||
raise RuntimeError(f"node B died with code {proc_b.returncode}") |
||||
print(f"\n--- Messages on B ---") |
||||
try: |
||||
msgs_b = await cli_b.messages(ch_id, count=10) |
||||
except ChatClientError as e: |
||||
print(f" messages error: {e}") |
||||
print(f" proc_b.returncode={proc_b.returncode}") |
||||
raise |
||||
check("4 messages on B", isinstance(msgs_b, list) and len(msgs_b) == 4, |
||||
f"got {len(msgs_b) if isinstance(msgs_b, list) else '?'}") |
||||
for i, m in enumerate(msgs_b): |
||||
print(f" [#{i+1}] ts={m.get('ts','?')} author={m.get('author_id','?')}") |
||||
|
||||
# ── Verify members on A ── |
||||
members_a = await cli_a.members(ch_id) |
||||
print(f"\n--- Members on A ({len(members_a)}) ---") |
||||
for m in members_a: |
||||
print(f" {m['node_id']} name={m.get('name','?')} online={m.get('online')} connected={m.get('connected')}") |
||||
check("2 members on A", isinstance(members_a, list) and len(members_a) >= 2, |
||||
f"got {len(members_a) if isinstance(members_a, list) else '?'}") |
||||
|
||||
print("\n=== TEST PASSED ===") |
||||
return 0 |
||||
|
||||
except Exception as e: |
||||
print(f"\n=== TEST FAILED: {e} ===", file=sys.stderr) |
||||
import traceback |
||||
traceback.print_exc() |
||||
return 1 |
||||
|
||||
finally: |
||||
print("\n--- Cleanup ---") |
||||
await kill_proc(proc_a, "proc_a") |
||||
await kill_proc(proc_b, "proc_b") |
||||
|
||||
for f in [config_a, config_b]: |
||||
try: os.unlink(f) |
||||
except OSError: pass |
||||
try: os.rmdir(db_a) |
||||
except OSError: pass |
||||
try: os.rmdir(db_b) |
||||
except OSError: pass |
||||
try: os.rmdir(tmpdir) |
||||
except OSError: pass |
||||
|
||||
print(f" temp dir cleaned: {tmpdir}") |
||||
|
||||
|
||||
if __name__ == "__main__": |
||||
async def _run(): |
||||
try: |
||||
return await asyncio.wait_for(main(), timeout=TOTAL_TIMEOUT) |
||||
except asyncio.TimeoutError: |
||||
print(f"\n=== TEST FAILED: total timeout {TOTAL_TIMEOUT}s ===", file=sys.stderr) |
||||
return 1 |
||||
sys.exit(asyncio.run(_run())) |
||||
@ -0,0 +1,23 @@
|
||||
[global] |
||||
db_path=db_a |
||||
|
||||
my_node_id=3033fd1d683b41e4 |
||||
my_private_key=289da7dc18c1fc693bc8409c6ae8125a2ba4017b77a4f343ff581842620bce4d |
||||
my_public_key=2e944d6c761a9bd5a915d5bbed973c9d22041fb66c5b4900f04903df13d27c3a |
||||
[server: srv] |
||||
addr=127.0.0.1:15001 |
||||
type=public |
||||
transport=tcp |
||||
|
||||
[chatserver] |
||||
db_path=db_a |
||||
headless_control_bind=127.0.0.1:15002 |
||||
|
||||
[debug] |
||||
debug=trace |
||||
etcp=info |
||||
general=info |
||||
member_sync=info |
||||
|
||||
[allowed_keys] |
||||
allow_all=1 |
||||
@ -0,0 +1,23 @@
|
||||
[global] |
||||
db_path=db_b |
||||
|
||||
my_node_id=0d83f84b68c25ad4 |
||||
my_private_key=30b409f91449fa87b2ad39b1dc9398e1e60d978763289041bb96914af1e5d557 |
||||
my_public_key=4de544515b606cc7af37967150e2e9bcf8e8c76fd9ce05f602d50ccdf77e6b01 |
||||
[server: srv] |
||||
addr=127.0.0.1:15003 |
||||
type=public |
||||
transport=tcp |
||||
|
||||
[chatserver] |
||||
db_path=db_b |
||||
headless_control_bind=127.0.0.1:15004 |
||||
|
||||
[debug] |
||||
debug=trace |
||||
etcp=info |
||||
general=info |
||||
member_sync=info |
||||
|
||||
[allowed_keys] |
||||
allow_all=1 |
||||
@ -0,0 +1,120 @@
|
||||
#!/bin/bash |
||||
# |
||||
# chat_tcp_test/run.sh — TCP integration test with diagnostics |
||||
# |
||||
set -eu |
||||
TESTDIR="$(cd "$(dirname "$0")" && pwd)" |
||||
cd "$TESTDIR" |
||||
|
||||
CHATCLI="python3 $TESTDIR/../chatcli" |
||||
UTUN="$TESTDIR/../../src/utun" |
||||
|
||||
# ── cleanup handler ── |
||||
cleanup() { |
||||
local rc=${1:-$?} |
||||
echo "" |
||||
echo "=== cleanup (exit=$rc) ===" |
||||
pkill -9 -f "utun.*15001\|utun.*15002\|utun.*15003\|utun.*15004" 2>/dev/null || true |
||||
sleep 0.3 |
||||
rm -f "$TESTDIR"/*.pid |
||||
if [ "$rc" -ne 0 ]; then |
||||
echo "=== a.log (last 40) ==="; tail -40 "$TESTDIR/a.log" 2>/dev/null || true |
||||
echo "=== b.log (last 40) ==="; tail -40 "$TESTDIR/b.log" 2>/dev/null || true |
||||
fi |
||||
exit "$rc" |
||||
} |
||||
trap 'cleanup 1' INT TERM |
||||
trap 'cleanup' EXIT |
||||
|
||||
die() { echo "FATAL: $*" >&2; exit 1; } |
||||
|
||||
# ── kill old + clean start ── |
||||
echo "--- Killing old processes on test ports ---" |
||||
for port in 15001 15002 15003 15004; do |
||||
fuser -k ${port}/tcp 2>/dev/null || true |
||||
done |
||||
sleep 0.5 |
||||
|
||||
rm -rf "$TESTDIR"/db_a "$TESTDIR"/db_b "$TESTDIR"/*.log "$TESTDIR"/*.pid |
||||
mkdir -p "$TESTDIR"/db_a "$TESTDIR"/db_b |
||||
|
||||
# ── start ── |
||||
echo "--- Starting utun ---" |
||||
"$UTUN" -f -p "$TESTDIR/a.pid" -l "$TESTDIR/a.log" -c "$TESTDIR/a.conf" > /dev/null 2>&1 & |
||||
PID_A=$! |
||||
"$UTUN" -f -p "$TESTDIR/b.pid" -l "$TESTDIR/b.log" -c "$TESTDIR/b.conf" > /dev/null 2>&1 & |
||||
PID_B=$! |
||||
|
||||
# ── wait ready (max 5 sec) ── |
||||
echo "--- Wait ready ---" |
||||
for i in $(seq 1 50); do |
||||
if $CHATCLI --port 15002 ping 2>/dev/null | grep -q uptime; then echo " A ready (${i}00ms)"; break; fi |
||||
sleep 0.1 |
||||
[ $i -lt 50 ] || die "node A not ready" |
||||
done |
||||
for i in $(seq 1 50); do |
||||
if $CHATCLI --port 15004 ping 2>/dev/null | grep -q uptime; then echo " B ready (${i}00ms)"; break; fi |
||||
sleep 0.1 |
||||
[ $i -lt 50 ] || die "node B not ready" |
||||
done |
||||
|
||||
# ── create channel ── |
||||
echo "--- Create channel ---" |
||||
$CHATCLI --port 15002 create "TestTCP" > /dev/null 2>&1 |
||||
sleep 0.3 |
||||
CH_ID=$($CHATCLI --port 15002 channels 2>/dev/null | python3 -c "import sys,json;print(json.load(sys.stdin)[0]['id'])") |
||||
[ -n "$CH_ID" ] || die "failed to get channel_id" |
||||
echo " ch_id=$CH_ID" |
||||
|
||||
# ── send messages ── |
||||
echo "--- Send messages ---" |
||||
for i in 1 2 3 4; do |
||||
$CHATCLI --port 15002 send "$CH_ID" "Message $i" > /dev/null 2>&1 || die "send $i failed" |
||||
done |
||||
echo " sent 4 OK" |
||||
|
||||
# ── invite ── |
||||
echo "--- Invite ---" |
||||
LINK=$($CHATCLI --port 15002 invite "$CH_ID" 2>/dev/null | python3 -c "import sys,json;print(json.load(sys.stdin).get('link',''))") |
||||
[ -n "$LINK" ] || die "invite failed" |
||||
echo " link OK" |
||||
|
||||
# ── connect B ── |
||||
echo "--- Connect B ---" |
||||
$CHATCLI --port 15004 connect "$LINK" > /dev/null 2>&1 || die "connect failed" |
||||
|
||||
# ── wait ETCP:UP (max 5 sec) ── |
||||
echo "--- Wait ETCP:UP ---" |
||||
UP=0 |
||||
for i in $(seq 1 50); do |
||||
sleep 0.1 |
||||
if $CHATCLI --port 15004 status 2>/dev/null | grep -q 'ETCP:UP'; then |
||||
echo " ETCP:UP after ${i}00ms" |
||||
UP=1; break |
||||
fi |
||||
done |
||||
[ $UP -eq 1 ] || die "no ETCP:UP after 5s" |
||||
|
||||
# ── wait members sync ── |
||||
echo "--- Wait members sync ---" |
||||
MEM_OK=0 |
||||
for i in $(seq 1 30); do |
||||
local MEM_COUNT |
||||
MEM_COUNT=$($CHATCLI --port 15004 members "$CH_ID" 2>/dev/null | python3 -c "import sys,json;print(len(json.load(sys.stdin)))" 2>/dev/null || echo 0) |
||||
if [ "$MEM_COUNT" -ge 2 ]; then |
||||
echo " members=$MEM_COUNT OK (${i}00ms)" |
||||
MEM_OK=1; break |
||||
fi |
||||
sleep 0.1 |
||||
done |
||||
[ $MEM_OK -eq 1 ] || die "members sync timeout" |
||||
|
||||
# ── messages ── |
||||
echo "--- Check messages ---" |
||||
MSG_COUNT=$($CHATCLI --port 15004 messages "$CH_ID" 10 2>/dev/null | python3 -c "import sys,json;print(len(json.load(sys.stdin)))" 2>/dev/null || echo 0) |
||||
[ "$MSG_COUNT" -eq 4 ] || die "expected 4 messages, got $MSG_COUNT" |
||||
echo " messages=$MSG_COUNT OK" |
||||
|
||||
echo "" |
||||
echo "=== TEST PASSED ===" |
||||
exit 0 |
||||
@ -0,0 +1,177 @@
|
||||
# chatcli — uTun headless chat CLI client |
||||
# |
||||
# Usage: |
||||
# chatcli [--host HOST] [--port PORT] <command> [args...] |
||||
# |
||||
# Environment: |
||||
# CHATCLI_HOST — default host (default: 127.0.0.1) |
||||
# CHATCLI_PORT — default port (default: 9999) |
||||
# |
||||
# Commands (see chatcli_commands.txt for full description): |
||||
# ping alive check |
||||
# status node info, connections, NTP |
||||
# channels list all channels |
||||
# members <ch_id> list members with full state |
||||
# messages <ch_id> [count] read last messages |
||||
# send <ch_id> <text> send text message |
||||
# invite <ch_id> create invite link |
||||
# connect <utun://...> join channel via invite link |
||||
# create <name> create new channel |
||||
# listen interactive event listener |
||||
|
||||
import sys, os, json, socket, struct |
||||
|
||||
DEFAULT_HOST = os.environ.get("CHATCLI_HOST", "127.0.0.1") |
||||
DEFAULT_PORT = int(os.environ.get("CHATCLI_PORT", "9999")) |
||||
TIMEOUT = 3 |
||||
|
||||
def _connect(): |
||||
s = socket.socket(socket.AF_INET, socket.SOCK_STREAM) |
||||
s.settimeout(TIMEOUT) |
||||
try: |
||||
s.connect((DEFAULT_HOST, DEFAULT_PORT)) |
||||
except (socket.timeout, ConnectionRefusedError, OSError) as e: |
||||
print(f"ERROR: cannot connect to {DEFAULT_HOST}:{DEFAULT_PORT} — {e}", file=sys.stderr) |
||||
sys.exit(1) |
||||
return s |
||||
|
||||
def _readline(s): |
||||
buf = b"" |
||||
while True: |
||||
try: |
||||
ch = s.recv(1) |
||||
except socket.timeout: |
||||
break |
||||
if not ch: |
||||
break |
||||
if ch == b"\n": |
||||
break |
||||
buf += ch |
||||
return buf.decode("utf-8", errors="replace") |
||||
|
||||
def _send(s, js): |
||||
s.sendall(js.encode()) |
||||
|
||||
def _req(cmd, **params): |
||||
s = _connect() |
||||
js = {"id": 1, "cmd": cmd} |
||||
js.update(params) |
||||
_send(s, json.dumps(js, ensure_ascii=True) + "\n") |
||||
resp = _readline(s) |
||||
s.close() |
||||
try: |
||||
obj = json.loads(resp) |
||||
except json.JSONDecodeError: |
||||
print(f"RAW: {resp}") |
||||
return |
||||
if obj.get("ok"): |
||||
data = obj.get("data", {}) |
||||
if isinstance(data, str): |
||||
print(data) |
||||
else: |
||||
print(json.dumps(data, indent=2, ensure_ascii=False)) |
||||
else: |
||||
print(f"ERROR: {obj.get('error', 'unknown')}", file=sys.stderr) |
||||
sys.exit(1) |
||||
|
||||
def cmd_ping(): |
||||
_req("ping") |
||||
|
||||
def cmd_status(): |
||||
_req("status") |
||||
|
||||
def cmd_channels(): |
||||
_req("channels") |
||||
|
||||
def cmd_members(ch_id): |
||||
_req("members", ch=ch_id) |
||||
|
||||
def cmd_messages(ch_id, count=None): |
||||
params = {"ch": ch_id} |
||||
if count: |
||||
params["count"] = int(count) |
||||
_req("messages", **params) |
||||
|
||||
def cmd_send(ch_id, text): |
||||
_req("send", ch=ch_id, content_type="text", data=text) |
||||
|
||||
def cmd_invite(ch_id): |
||||
_req("invite", ch=ch_id) |
||||
|
||||
def cmd_connect(link): |
||||
_req("connect", link=link) |
||||
|
||||
def cmd_create(name): |
||||
_req("create_channel", name=name) |
||||
|
||||
def cmd_listen(): |
||||
print(f"Listening on {DEFAULT_HOST}:{DEFAULT_PORT} (Ctrl+C to quit)") |
||||
s = _connect() |
||||
_send(s, json.dumps({"id": 1, "cmd": "subscribe", "enable": 1}) + "\n") |
||||
try: |
||||
while True: |
||||
line = _readline(s) |
||||
if not line: |
||||
break |
||||
try: |
||||
obj = json.loads(line) |
||||
except json.JSONDecodeError: |
||||
print(line) |
||||
continue |
||||
evt = obj.get("event", "") |
||||
if evt == "msg": |
||||
print(f"[{obj.get('ch','?')}] {obj.get('author_id','?')}: new message") |
||||
elif evt == "members_changed": |
||||
print(f"[{obj.get('ch','?')}] members changed") |
||||
elif evt == "channel_updated": |
||||
print(f"[{obj.get('ch','?')}] channel updated") |
||||
elif evt == "invite_received": |
||||
print(f"INVITE: ch={obj.get('ch','?')} name={obj.get('ch_name','?')} from={obj.get('from_id','?')}") |
||||
else: |
||||
print(json.dumps(obj, indent=2, ensure_ascii=False)) |
||||
except KeyboardInterrupt: |
||||
print("\nDisconnected.") |
||||
finally: |
||||
s.close() |
||||
|
||||
def main(): |
||||
args = sys.argv[1:] |
||||
i = 0 |
||||
while i < len(args): |
||||
if args[i] == "--host" and i + 1 < len(args): |
||||
global DEFAULT_HOST; DEFAULT_HOST = args[i + 1]; i += 2 |
||||
elif args[i] == "--port" and i + 1 < len(args): |
||||
global DEFAULT_PORT; DEFAULT_PORT = int(args[i + 1]); i += 2 |
||||
else: |
||||
break |
||||
args = args[i:] |
||||
|
||||
if not args or args[0] in ("help", "--help", "-h"): |
||||
print(__doc__) |
||||
sys.exit(0) |
||||
|
||||
cmd = args[0] |
||||
try: |
||||
if cmd == "ping": cmd_ping() |
||||
elif cmd == "status": cmd_status() |
||||
elif cmd == "channels": cmd_channels() |
||||
elif cmd == "members": cmd_members(*args[1:2] if len(args) > 1 else (_die("usage: members <ch_id>"),)) |
||||
elif cmd in ("messages","msgs"): cmd_messages(*args[1:3] if len(args) > 1 else (_die("usage: messages <ch_id> [count]"),)) |
||||
elif cmd == "send": cmd_send(*args[1:3] if len(args) > 2 else (_die("usage: send <ch_id> <text>"),)) |
||||
elif cmd == "invite": cmd_invite(*args[1:2] if len(args) > 1 else (_die("usage: invite <ch_id>"),)) |
||||
elif cmd in ("connect","join"): cmd_connect(*args[1:2] if len(args) > 1 else (_die("usage: connect <utun://...>"),)) |
||||
elif cmd == "create": cmd_create(*args[1:2] if len(args) > 1 else (_die("usage: create <name>"),)) |
||||
elif cmd == "listen": cmd_listen() |
||||
else: print(f"Unknown command: {cmd}\nUse --help for usage", file=sys.stderr); sys.exit(1) |
||||
except socket.timeout: |
||||
print(f"ERROR: timeout ({TIMEOUT}s) connecting to {DEFAULT_HOST}:{DEFAULT_PORT}", file=sys.stderr) |
||||
sys.exit(1) |
||||
except ConnectionRefusedError: |
||||
print(f"ERROR: connection refused at {DEFAULT_HOST}:{DEFAULT_PORT} — is uTun running with headless_control_bind?", file=sys.stderr) |
||||
sys.exit(1) |
||||
|
||||
def _die(msg): |
||||
print(msg, file=sys.stderr); sys.exit(1) |
||||
|
||||
if __name__ == "__main__": |
||||
main() |
||||
@ -0,0 +1,76 @@
|
||||
#!/bin/bash |
||||
# chatcli — headless chat CLI client for uTun |
||||
# |
||||
# Usage: |
||||
# chatcli [--host HOST] [--port PORT] <command> [args...] |
||||
# |
||||
# Commands: |
||||
# channels list channels |
||||
# members <ch_id> list members with full state |
||||
# messages <ch_id> [count] read last messages |
||||
# send <ch_id> <content_type> <text> send message |
||||
# invite <ch_id> create invite link |
||||
# connect <utun://...> join channel via invite link |
||||
# create <name> create new channel |
||||
# status show connection status |
||||
# listen subscribe to events (interactive) |
||||
# --help show this help |
||||
|
||||
set -e |
||||
HOST="${CHATCLI_HOST:-127.0.0.1}" |
||||
PORT="${CHATCLI_PORT:-9999}" |
||||
ID=1 |
||||
NC_TIMEOUT="${CHATCLI_TIMEOUT:-5}" |
||||
|
||||
_send() { |
||||
local json="$1" |
||||
if command -v socat &>/dev/null; then |
||||
echo "$json" | socat -t "$NC_TIMEOUT" - TCP:"$HOST":"$PORT" 2>/dev/null || true |
||||
elif command -v nc &>/dev/null; then |
||||
if nc -h 2>&1 | grep -q '^BusyBox'; then |
||||
echo "$json" | nc "$HOST" "$PORT" 2>/dev/null || true |
||||
else |
||||
echo "$json" | nc -q "$NC_TIMEOUT" "$HOST" "$PORT" 2>/dev/null || true |
||||
fi |
||||
else |
||||
echo "ERROR: socat or nc required" >&2; exit 1 |
||||
fi |
||||
} |
||||
|
||||
_request() { |
||||
local cmd="$1" data="$2" |
||||
local json='{"id":'"$ID"',"cmd":"'"$cmd"'"' |
||||
if [ -n "$data" ]; then json="$json,$data"; fi |
||||
json="$json}" |
||||
_send "$json" |
||||
ID=$((ID + 1)) |
||||
} |
||||
|
||||
_show() { local json="$1"; echo "$json" | python3 -m json.tool 2>/dev/null || echo "$json"; } |
||||
|
||||
case "${1:-}" in |
||||
channels) r=$(_request channels); _show "$r" ;; |
||||
members) r=$(_request members '"ch":"'"$2"'"'); _show "$r" ;; |
||||
messages|msgs) |
||||
cnt="${3:-20}" |
||||
r=$(_request messages '"ch":"'"$2"'","count":'"$cnt"','"offset":0'); _show "$r" ;; |
||||
send) |
||||
ct="${3:-text}"; text="$4" |
||||
r=$(_request send '"ch":"'"$2"'","content_type":"'"$ct"'","data":"'"$text"'"'); _show "$r" ;; |
||||
invite) r=$(_request invite '"ch":"'"$2"'"'); _show "$r" ;; |
||||
connect|join) r=$(_request connect '"link":"'"$2"'"'); _show "$r" ;; |
||||
create) r=$(_request create_channel '"name":"'"$2"'"'); _show "$r" ;; |
||||
status|stat) r=$(_request status); _show "$r" ;; |
||||
listen) |
||||
echo "Connecting to $HOST:$PORT (Ctrl+C to quit)..." |
||||
{ |
||||
printf '{"id":1,"cmd":"subscribe","enable":1}\n' |
||||
cat |
||||
} | nc "$HOST" "$PORT" |
||||
;; |
||||
ping) r=$(_request ping); _show "$r" ;; |
||||
--help|help|-h|"") |
||||
head -20 "$0" | sed 's/^# //' | grep -v '^#' |
||||
;; |
||||
*) echo "Unknown command: $1"; echo "Use --help for usage" >&2; exit 1 ;; |
||||
esac |
||||
@ -0,0 +1,77 @@
|
||||
chatcli — headless chat CLI for uTun |
||||
===================================== |
||||
|
||||
TCP JSON-RPC клиент для управления чатом uTun без GUI. |
||||
|
||||
Настройка сервера |
||||
----------------- |
||||
В конфиге uTun в секции [chatserver] добавить: |
||||
headless_control_bind=127.0.0.1:9999 |
||||
|
||||
Протокол |
||||
-------- |
||||
JSON line: запрос {"id":N,"cmd":"...","params"} → ответ {"id":N,"ok":true,"data":...} |
||||
При подписке (subscribe) сервер шлёт асинхронные события: {"event":"...",...} |
||||
|
||||
Переменные окружения |
||||
-------------------- |
||||
CHATCLI_HOST — хост по умолчанию (127.0.0.1) |
||||
CHATCLI_PORT — порт по умолчанию (9999) |
||||
|
||||
Команды |
||||
------- |
||||
|
||||
ping |
||||
Проверка соединения и uptime сервера. |
||||
Пример: chatcli ping |
||||
|
||||
status |
||||
Статус узла: NTP синхронизация, активные соединения, NAT. |
||||
Пример: chatcli status |
||||
|
||||
channels |
||||
Список всех каналов с метаданными: id, имя, owner, число мемберов, |
||||
число сообщений, сколько онлайн. |
||||
Пример: chatcli channels |
||||
|
||||
members <ch_id> |
||||
Список участников канала с полным состоянием: |
||||
node_id, имя, online, connected, публичные ключи (x25519/ed25519 hex), |
||||
адреса (ip, port, proto, rtt). |
||||
Пример: chatcli members 12345 |
||||
|
||||
messages <ch_id> [count] |
||||
Последние N сообщений канала (по умолчанию 20): |
||||
id, timestamp, author_id, author_name, content_type, data. |
||||
Пример: chatcli messages 12345 10 |
||||
|
||||
send <ch_id> <text> |
||||
Отправить текстовое сообщение в канал. |
||||
Пример: chatcli send 12345 "hello world" |
||||
|
||||
invite <ch_id> |
||||
Создать invite-ссылку для канала (формат utun://...). |
||||
Ссылка включает pubkey, адреса сервера и ID канала. |
||||
Пример: chatcli invite 12345 |
||||
|
||||
connect <utun://...> |
||||
Подключиться к каналу по invite-ссылке. |
||||
Декодирует ссылку, сохраняет pubkey и адреса узла в БД, |
||||
запускает процесс подключения. |
||||
Пример: chatcli connect utun://AgR0ZXN0... |
||||
|
||||
create <name> |
||||
Создать новый канал. Генерирует ключи X25519/Ed25519, |
||||
подписывает join-сообщение. |
||||
Пример: chatcli create "My Group" |
||||
|
||||
listen |
||||
Подписаться на события и слушать в реальном времени. |
||||
Получает: новые сообщения, изменения мемберов, приглашения. |
||||
Выход: Ctrl+C |
||||
Пример: chatcli listen |
||||
|
||||
Настройки хоста/порта |
||||
--------------------- |
||||
chatcli --host 192.168.1.1 --port 8888 channels |
||||
CHATCLI_HOST=10.0.0.1 CHATCLI_PORT=7777 chatcli status |
||||
Loading…
Reference in new issue