Initial Awawawa
This commit is contained in:
1091
backend/src/client.c
Normal file
1091
backend/src/client.c
Normal file
File diff suppressed because it is too large
Load Diff
13
backend/src/client.h
Normal file
13
backend/src/client.h
Normal file
@@ -0,0 +1,13 @@
|
||||
#ifndef PB_CLIENT_H
|
||||
#define PB_CLIENT_H
|
||||
|
||||
#include "config.h"
|
||||
|
||||
// Runs the client role: keeps an SSH session to the target, reconnecting with
|
||||
// backoff, until SIGTERM/SIGINT.
|
||||
int client_run(const struct pb_config *cfg);
|
||||
|
||||
// Resolves the identity file (IdentityFile, root's key, host key) into out.
|
||||
int client_identity(const struct pb_config *cfg, char *out, size_t outsz);
|
||||
|
||||
#endif
|
||||
129
backend/src/config.c
Normal file
129
backend/src/config.c
Normal file
@@ -0,0 +1,129 @@
|
||||
#include <ctype.h>
|
||||
#include <stdio.h>
|
||||
#include <stdlib.h>
|
||||
#include <string.h>
|
||||
#include <strings.h>
|
||||
|
||||
#include "config.h"
|
||||
#include "log.h"
|
||||
|
||||
void config_defaults(struct pb_config *c)
|
||||
{
|
||||
memset(c, 0, sizeof(*c));
|
||||
c->role = ROLE_CLIENT;
|
||||
snprintf(c->db_path, sizeof(c->db_path), "%s", PB_DEFAULT_DB);
|
||||
snprintf(c->run_dir, sizeof(c->run_dir), "%s", PB_RUN_DIR);
|
||||
c->log_level = LOG_LVL_INFO;
|
||||
c->target_port = 2222;
|
||||
snprintf(c->known_hosts, sizeof(c->known_hosts), "/etc/patchbay/known_hosts");
|
||||
snprintf(c->ssh_user, sizeof(c->ssh_user), "patchbay");
|
||||
c->services_interval = 30;
|
||||
c->ssh_port = 0; // 0 = use target_port
|
||||
c->hub_port = 7701;
|
||||
snprintf(c->sshd_host_key, sizeof(c->sshd_host_key), "/etc/patchbay/ssh_host_ed25519_key");
|
||||
snprintf(c->authorized_keys, sizeof(c->authorized_keys), "/etc/patchbay/authorized_keys");
|
||||
snprintf(c->daemon_path, sizeof(c->daemon_path), "/usr/local/sbin/patchbayd");
|
||||
}
|
||||
|
||||
static char *trim(char *s)
|
||||
{
|
||||
while (isspace((unsigned char)*s))
|
||||
s++;
|
||||
char *e = s + strlen(s);
|
||||
while (e > s && isspace((unsigned char)e[-1]))
|
||||
*--e = '\0';
|
||||
return s;
|
||||
}
|
||||
|
||||
static void set_str(char *dst, size_t n, const char *v)
|
||||
{
|
||||
snprintf(dst, n, "%s", v);
|
||||
}
|
||||
|
||||
#define STR(field) set_str(c->field, sizeof(c->field), val)
|
||||
|
||||
int config_parse_line(struct pb_config *c, char *line)
|
||||
{
|
||||
char *s = trim(line);
|
||||
if (*s == '\0' || *s == '#' || *s == ';' || *s == '[')
|
||||
return 0;
|
||||
|
||||
char *eq = strchr(s, '=');
|
||||
if (!eq)
|
||||
return -1;
|
||||
*eq = '\0';
|
||||
char *key = trim(s);
|
||||
char *val = trim(eq + 1);
|
||||
|
||||
// Allow optional quoting so values with leading/trailing spaces survive.
|
||||
size_t vl = strlen(val);
|
||||
if (vl >= 2 && val[0] == '"' && val[vl - 1] == '"') {
|
||||
val[vl - 1] = '\0';
|
||||
val++;
|
||||
}
|
||||
|
||||
if (!strcasecmp(key, "Role")) {
|
||||
if (!strcasecmp(val, "target"))
|
||||
c->role = ROLE_TARGET;
|
||||
else if (!strcasecmp(val, "client"))
|
||||
c->role = ROLE_CLIENT;
|
||||
else
|
||||
return -1;
|
||||
} else if (!strcasecmp(key, "Database")) {
|
||||
STR(db_path);
|
||||
} else if (!strcasecmp(key, "RunDir")) {
|
||||
STR(run_dir);
|
||||
} else if (!strcasecmp(key, "LogLevel")) {
|
||||
if (!strcasecmp(val, "error")) c->log_level = LOG_LVL_ERR;
|
||||
else if (!strcasecmp(val, "warn")) c->log_level = LOG_LVL_WARN;
|
||||
else if (!strcasecmp(val, "info")) c->log_level = LOG_LVL_INFO;
|
||||
else if (!strcasecmp(val, "debug")) c->log_level = LOG_LVL_DEBUG;
|
||||
else return -1;
|
||||
} else if (!strcasecmp(key, "TargetHost")) {
|
||||
STR(target_host);
|
||||
} else if (!strcasecmp(key, "TargetPort")) {
|
||||
c->target_port = atoi(val);
|
||||
} else if (!strcasecmp(key, "IdentityFile")) {
|
||||
STR(identity_file);
|
||||
} else if (!strcasecmp(key, "IdentityPassphrase")) {
|
||||
STR(identity_pass);
|
||||
} else if (!strcasecmp(key, "KnownHosts")) {
|
||||
STR(known_hosts);
|
||||
} else if (!strcasecmp(key, "SSHUser")) {
|
||||
STR(ssh_user);
|
||||
} else if (!strcasecmp(key, "ServicesInterval")) {
|
||||
c->services_interval = atoi(val);
|
||||
} else if (!strcasecmp(key, "HubPort")) {
|
||||
c->hub_port = atoi(val);
|
||||
} else if (!strcasecmp(key, "SSHHostKey")) {
|
||||
STR(sshd_host_key);
|
||||
} else if (!strcasecmp(key, "AuthorizedKeys")) {
|
||||
STR(authorized_keys);
|
||||
} else if (!strcasecmp(key, "DaemonPath")) {
|
||||
STR(daemon_path);
|
||||
}
|
||||
// Anything else belongs to the web frontend.
|
||||
return 0;
|
||||
}
|
||||
|
||||
int config_load(struct pb_config *c, const char *path)
|
||||
{
|
||||
FILE *f = fopen(path, "r");
|
||||
if (!f)
|
||||
return -1;
|
||||
|
||||
char line[1024];
|
||||
int lineno = 0;
|
||||
while (fgets(line, sizeof(line), f)) {
|
||||
lineno++;
|
||||
if (config_parse_line(c, line) < 0)
|
||||
log_warn("%s:%d: invalid line ignored", path, lineno);
|
||||
}
|
||||
fclose(f);
|
||||
|
||||
if (c->ssh_port == 0)
|
||||
c->ssh_port = c->target_port;
|
||||
if (c->services_interval < 5)
|
||||
c->services_interval = 5;
|
||||
return 0;
|
||||
}
|
||||
39
backend/src/config.h
Normal file
39
backend/src/config.h
Normal file
@@ -0,0 +1,39 @@
|
||||
#ifndef PB_CONFIG_H
|
||||
#define PB_CONFIG_H
|
||||
|
||||
#define PB_DEFAULT_CONF "/etc/patchbay/patchbay.conf"
|
||||
#define PB_DEFAULT_DB "/etc/patchbay/patchbay.db"
|
||||
#define PB_RUN_DIR "/run/patchbay"
|
||||
|
||||
enum pb_role { ROLE_CLIENT, ROLE_TARGET };
|
||||
|
||||
struct pb_config {
|
||||
enum pb_role role;
|
||||
char db_path[256];
|
||||
char run_dir[256];
|
||||
int log_level;
|
||||
|
||||
// client role
|
||||
char target_host[256];
|
||||
int target_port;
|
||||
char identity_file[256]; // empty = auto (root key, then host key)
|
||||
char identity_pass[256];
|
||||
char known_hosts[256];
|
||||
char ssh_user[64];
|
||||
int services_interval; // seconds between Services reports
|
||||
|
||||
// target role
|
||||
int ssh_port; // dedicated sshd port (same value clients use as TargetPort)
|
||||
int hub_port; // loopback port data channels are forwarded to
|
||||
char sshd_host_key[256];
|
||||
char authorized_keys[256];
|
||||
char daemon_path[256]; // forced command written into authorized_keys
|
||||
};
|
||||
|
||||
void config_defaults(struct pb_config *c);
|
||||
// Returns 0 on success, -1 if the file cannot be read; unknown keys are ignored
|
||||
// because the web frontend shares the file.
|
||||
int config_load(struct pb_config *c, const char *path);
|
||||
int config_parse_line(struct pb_config *c, char *line);
|
||||
|
||||
#endif
|
||||
54
backend/src/db.c
Normal file
54
backend/src/db.c
Normal file
@@ -0,0 +1,54 @@
|
||||
#include <stdlib.h>
|
||||
#include <sys/stat.h>
|
||||
|
||||
#include "db.h"
|
||||
#include "log.h"
|
||||
#include "schema.h" // generated from schema.sql by the Makefile
|
||||
|
||||
int db_exec(sqlite3 *db, const char *sql)
|
||||
{
|
||||
char *err = NULL;
|
||||
if (sqlite3_exec(db, sql, NULL, NULL, &err) != SQLITE_OK) {
|
||||
log_err("sqlite: %s", err ? err : "unknown error");
|
||||
sqlite3_free(err);
|
||||
return -1;
|
||||
}
|
||||
return 0;
|
||||
}
|
||||
|
||||
sqlite3 *db_open(const char *path)
|
||||
{
|
||||
sqlite3 *db;
|
||||
mode_t old = umask(0077); // the DB holds password hashes
|
||||
int rc = sqlite3_open(path, &db);
|
||||
umask(old);
|
||||
if (rc != SQLITE_OK) {
|
||||
log_err("cannot open database %s: %s", path, sqlite3_errmsg(db));
|
||||
sqlite3_close(db);
|
||||
return NULL;
|
||||
}
|
||||
sqlite3_busy_timeout(db, 5000);
|
||||
if (db_exec(db, pb_schema_sql) < 0 || db_exec(db, "PRAGMA foreign_keys = ON;") < 0) {
|
||||
sqlite3_close(db);
|
||||
return NULL;
|
||||
}
|
||||
// Migration for databases from before tunnel nodes; fails harmlessly if present.
|
||||
sqlite3_exec(db, "ALTER TABLE nodes ADD COLUMN iface TEXT NOT NULL DEFAULT ''", NULL, NULL, NULL);
|
||||
return db;
|
||||
}
|
||||
|
||||
long db_setting_int(sqlite3 *db, const char *key, long def)
|
||||
{
|
||||
sqlite3_stmt *st;
|
||||
long v = def;
|
||||
if (sqlite3_prepare_v2(db, "SELECT value FROM settings WHERE key = ?", -1, &st, NULL) != SQLITE_OK)
|
||||
return def;
|
||||
sqlite3_bind_text(st, 1, key, -1, SQLITE_STATIC);
|
||||
if (sqlite3_step(st) == SQLITE_ROW) {
|
||||
const char *s = (const char *)sqlite3_column_text(st, 0);
|
||||
if (s)
|
||||
v = strtol(s, NULL, 10);
|
||||
}
|
||||
sqlite3_finalize(st);
|
||||
return v;
|
||||
}
|
||||
13
backend/src/db.h
Normal file
13
backend/src/db.h
Normal file
@@ -0,0 +1,13 @@
|
||||
#ifndef PB_DB_H
|
||||
#define PB_DB_H
|
||||
|
||||
#include <sqlite3.h>
|
||||
|
||||
// Opens (and creates/migrates) the database with WAL and a busy timeout so the
|
||||
// daemon and the web frontend can write concurrently.
|
||||
sqlite3 *db_open(const char *path);
|
||||
int db_exec(sqlite3 *db, const char *sql);
|
||||
// Integer setting with fallback.
|
||||
long db_setting_int(sqlite3 *db, const char *key, long def);
|
||||
|
||||
#endif
|
||||
1724
backend/src/hub.c
Normal file
1724
backend/src/hub.c
Normal file
File diff suppressed because it is too large
Load Diff
14
backend/src/hub.h
Normal file
14
backend/src/hub.h
Normal file
@@ -0,0 +1,14 @@
|
||||
#ifndef PB_HUB_H
|
||||
#define PB_HUB_H
|
||||
|
||||
#include <sqlite3.h>
|
||||
|
||||
#include "config.h"
|
||||
|
||||
// Runs the target role until SIGTERM/SIGINT. SIGHUP reloads the patch graph.
|
||||
int hub_run(const struct pb_config *cfg);
|
||||
|
||||
// Rewrites the sshd authorized_keys file from the clients table.
|
||||
int hub_write_authorized_keys(const struct pb_config *cfg, sqlite3 *db);
|
||||
|
||||
#endif
|
||||
31
backend/src/json.c
Normal file
31
backend/src/json.c
Normal file
@@ -0,0 +1,31 @@
|
||||
#include <stdio.h>
|
||||
|
||||
#include "json.h"
|
||||
|
||||
int json_str(struct buf *b, const char *s)
|
||||
{
|
||||
if (buf_append(b, "\"", 1) < 0)
|
||||
return -1;
|
||||
for (; *s; s++) {
|
||||
unsigned char c = (unsigned char)*s;
|
||||
char esc[8];
|
||||
const char *out = NULL;
|
||||
size_t n = 0;
|
||||
switch (c) {
|
||||
case '"': out = "\\\""; n = 2; break;
|
||||
case '\\': out = "\\\\"; n = 2; break;
|
||||
case '\n': out = "\\n"; n = 2; break;
|
||||
case '\r': out = "\\r"; n = 2; break;
|
||||
case '\t': out = "\\t"; n = 2; break;
|
||||
default:
|
||||
if (c < 0x20) {
|
||||
snprintf(esc, sizeof(esc), "\\u%04x", c);
|
||||
out = esc;
|
||||
n = 6;
|
||||
}
|
||||
}
|
||||
if (out ? buf_append(b, out, n) : buf_append(b, s, 1))
|
||||
return -1;
|
||||
}
|
||||
return buf_append(b, "\"", 1);
|
||||
}
|
||||
9
backend/src/json.h
Normal file
9
backend/src/json.h
Normal file
@@ -0,0 +1,9 @@
|
||||
#ifndef PB_JSON_H
|
||||
#define PB_JSON_H
|
||||
|
||||
#include "proto.h"
|
||||
|
||||
// Appends s as a quoted, escaped JSON string.
|
||||
int json_str(struct buf *b, const char *s);
|
||||
|
||||
#endif
|
||||
33
backend/src/log.c
Normal file
33
backend/src/log.c
Normal file
@@ -0,0 +1,33 @@
|
||||
#include <stdarg.h>
|
||||
#include <stdio.h>
|
||||
#include <time.h>
|
||||
|
||||
#include "log.h"
|
||||
|
||||
static int cur_level = LOG_LVL_INFO;
|
||||
static const char *names[] = { "error", "warn", "info", "debug" };
|
||||
|
||||
void log_set_level(int level)
|
||||
{
|
||||
cur_level = level;
|
||||
}
|
||||
|
||||
void log_msg(int level, const char *fmt, ...)
|
||||
{
|
||||
if (level > cur_level)
|
||||
return;
|
||||
|
||||
// Service managers add their own timestamps, but plain foreground runs do not.
|
||||
char ts[32];
|
||||
time_t now = time(NULL);
|
||||
struct tm tm;
|
||||
localtime_r(&now, &tm);
|
||||
strftime(ts, sizeof(ts), "%Y-%m-%d %H:%M:%S", &tm);
|
||||
|
||||
va_list ap;
|
||||
va_start(ap, fmt);
|
||||
fprintf(stderr, "%s patchbayd[%s]: ", ts, names[level]);
|
||||
vfprintf(stderr, fmt, ap);
|
||||
fputc('\n', stderr);
|
||||
va_end(ap);
|
||||
}
|
||||
14
backend/src/log.h
Normal file
14
backend/src/log.h
Normal file
@@ -0,0 +1,14 @@
|
||||
#ifndef PB_LOG_H
|
||||
#define PB_LOG_H
|
||||
|
||||
enum { LOG_LVL_ERR, LOG_LVL_WARN, LOG_LVL_INFO, LOG_LVL_DEBUG };
|
||||
|
||||
void log_set_level(int level);
|
||||
void log_msg(int level, const char *fmt, ...) __attribute__((format(printf, 2, 3)));
|
||||
|
||||
#define log_err(...) log_msg(LOG_LVL_ERR, __VA_ARGS__)
|
||||
#define log_warn(...) log_msg(LOG_LVL_WARN, __VA_ARGS__)
|
||||
#define log_info(...) log_msg(LOG_LVL_INFO, __VA_ARGS__)
|
||||
#define log_debug(...) log_msg(LOG_LVL_DEBUG, __VA_ARGS__)
|
||||
|
||||
#endif
|
||||
120
backend/src/main.c
Normal file
120
backend/src/main.c
Normal file
@@ -0,0 +1,120 @@
|
||||
#include <getopt.h>
|
||||
#include <stdio.h>
|
||||
#include <stdlib.h>
|
||||
#include <string.h>
|
||||
|
||||
#include "client.h"
|
||||
#include "config.h"
|
||||
#include "db.h"
|
||||
#include "hub.h"
|
||||
#include "log.h"
|
||||
#include "proto.h"
|
||||
#include "relay.h"
|
||||
|
||||
static void usage(void)
|
||||
{
|
||||
fprintf(stderr,
|
||||
"usage: patchbayd [-c config] [-v] run in the role set by Role=\n"
|
||||
" patchbayd --relay <client_id> forced command used by the PatchBay sshd\n"
|
||||
" patchbayd --pubkey print this client's public key\n"
|
||||
" patchbayd --write-keys regenerate authorized_keys (target)\n"
|
||||
" patchbayd --print-config print effective settings as KEY=value\n"
|
||||
" patchbayd --version\n");
|
||||
}
|
||||
|
||||
static int print_pubkey(const struct pb_config *cfg)
|
||||
{
|
||||
char ident[256], pub[300];
|
||||
if (client_identity(cfg, ident, sizeof(ident)) < 0) {
|
||||
fprintf(stderr, "no SSH identity found\n");
|
||||
return 1;
|
||||
}
|
||||
snprintf(pub, sizeof(pub), "%s.pub", ident);
|
||||
FILE *f = fopen(pub, "r");
|
||||
if (!f) {
|
||||
fprintf(stderr, "cannot read %s; create it with: ssh-keygen -y -f %s > %s\n", pub, ident, pub);
|
||||
return 1;
|
||||
}
|
||||
char line[8192];
|
||||
while (fgets(line, sizeof(line), f))
|
||||
fputs(line, stdout);
|
||||
fclose(f);
|
||||
return 0;
|
||||
}
|
||||
|
||||
static void print_config(const struct pb_config *cfg)
|
||||
{
|
||||
printf("Role=%s\n", cfg->role == ROLE_TARGET ? "target" : "client");
|
||||
printf("Database=%s\n", cfg->db_path);
|
||||
printf("RunDir=%s\n", cfg->run_dir);
|
||||
printf("TargetPort=%d\n", cfg->target_port);
|
||||
printf("HubPort=%d\n", cfg->hub_port);
|
||||
printf("SSHUser=%s\n", cfg->ssh_user);
|
||||
printf("SSHHostKey=%s\n", cfg->sshd_host_key);
|
||||
printf("AuthorizedKeys=%s\n", cfg->authorized_keys);
|
||||
printf("DaemonPath=%s\n", cfg->daemon_path);
|
||||
}
|
||||
|
||||
int main(int argc, char **argv)
|
||||
{
|
||||
const char *conf_path = PB_DEFAULT_CONF;
|
||||
const char *relay_id = NULL;
|
||||
int verbose = 0, mode = 0;
|
||||
enum { M_RUN, M_RELAY, M_PUBKEY, M_KEYS, M_PRINT };
|
||||
|
||||
static const struct option opts[] = {
|
||||
{ "config", required_argument, NULL, 'c' },
|
||||
{ "relay", required_argument, NULL, 'r' },
|
||||
{ "pubkey", no_argument, NULL, 'p' },
|
||||
{ "write-keys", no_argument, NULL, 'k' },
|
||||
{ "print-config", no_argument, NULL, 'P' },
|
||||
{ "verbose", no_argument, NULL, 'v' },
|
||||
{ "version", no_argument, NULL, 'V' },
|
||||
{ "help", no_argument, NULL, 'h' },
|
||||
{ NULL, 0, NULL, 0 },
|
||||
};
|
||||
int o;
|
||||
while ((o = getopt_long(argc, argv, "c:vh", opts, NULL)) != -1) {
|
||||
switch (o) {
|
||||
case 'c': conf_path = optarg; break;
|
||||
case 'r': mode = M_RELAY; relay_id = optarg; break;
|
||||
case 'p': mode = M_PUBKEY; break;
|
||||
case 'k': mode = M_KEYS; break;
|
||||
case 'P': mode = M_PRINT; break;
|
||||
case 'v': verbose = 1; break;
|
||||
case 'V': printf("patchbayd %s\n", PB_VERSION); return 0;
|
||||
default: usage(); return o == 'h' ? 0 : 2;
|
||||
}
|
||||
}
|
||||
|
||||
struct pb_config cfg;
|
||||
config_defaults(&cfg);
|
||||
// The relay runs as the unprivileged SSH user and may not be able to read
|
||||
// the root-only config; defaults are enough for it.
|
||||
if (config_load(&cfg, conf_path) < 0 && mode != M_RELAY) {
|
||||
fprintf(stderr, "cannot read %s\n", conf_path);
|
||||
return 1;
|
||||
}
|
||||
log_set_level(verbose ? LOG_LVL_DEBUG : cfg.log_level);
|
||||
|
||||
switch (mode) {
|
||||
case M_RELAY:
|
||||
return relay_run(&cfg, relay_id);
|
||||
case M_PUBKEY:
|
||||
return print_pubkey(&cfg);
|
||||
case M_PRINT:
|
||||
print_config(&cfg);
|
||||
return 0;
|
||||
case M_KEYS: {
|
||||
sqlite3 *db = db_open(cfg.db_path);
|
||||
if (!db)
|
||||
return 1;
|
||||
int rc = hub_write_authorized_keys(&cfg, db) < 0 ? 1 : 0;
|
||||
sqlite3_close(db);
|
||||
return rc;
|
||||
}
|
||||
}
|
||||
|
||||
log_info("patchbayd %s starting as %s", PB_VERSION, cfg.role == ROLE_TARGET ? "target" : "client");
|
||||
return cfg.role == ROLE_TARGET ? hub_run(&cfg) : client_run(&cfg);
|
||||
}
|
||||
216
backend/src/net.c
Normal file
216
backend/src/net.c
Normal file
@@ -0,0 +1,216 @@
|
||||
#include <arpa/inet.h>
|
||||
#include <errno.h>
|
||||
#include <fcntl.h>
|
||||
#include <netdb.h>
|
||||
#include <netinet/in.h>
|
||||
#include <net/if.h>
|
||||
#include <netinet/tcp.h>
|
||||
#include <stdio.h>
|
||||
#include <string.h>
|
||||
#include <strings.h>
|
||||
#include <sys/stat.h>
|
||||
#include <sys/un.h>
|
||||
#include <unistd.h>
|
||||
|
||||
#include "net.h"
|
||||
|
||||
int proto_parse(const char *s)
|
||||
{
|
||||
if (!strcasecmp(s, "tcp"))
|
||||
return PROTO_TCP;
|
||||
if (!strcasecmp(s, "udp"))
|
||||
return PROTO_UDP;
|
||||
return -1;
|
||||
}
|
||||
|
||||
const char *proto_name(int proto)
|
||||
{
|
||||
return proto == PROTO_UDP ? "udp" : "tcp";
|
||||
}
|
||||
|
||||
int set_nonblock(int fd)
|
||||
{
|
||||
int fl = fcntl(fd, F_GETFL);
|
||||
if (fl < 0)
|
||||
return -1;
|
||||
return fcntl(fd, F_SETFL, fl | O_NONBLOCK);
|
||||
}
|
||||
|
||||
void tune_stream(int fd)
|
||||
{
|
||||
int one = 1;
|
||||
setsockopt(fd, IPPROTO_TCP, TCP_NODELAY, &one, sizeof(one));
|
||||
setsockopt(fd, SOL_SOCKET, SO_KEEPALIVE, &one, sizeof(one));
|
||||
}
|
||||
|
||||
static int resolve(const char *host, int port, int proto, int passive, struct addrinfo **res)
|
||||
{
|
||||
struct addrinfo hints;
|
||||
memset(&hints, 0, sizeof(hints));
|
||||
hints.ai_family = AF_UNSPEC;
|
||||
hints.ai_socktype = proto == PROTO_UDP ? SOCK_DGRAM : SOCK_STREAM;
|
||||
hints.ai_flags = passive ? AI_PASSIVE : 0;
|
||||
char portstr[16];
|
||||
snprintf(portstr, sizeof(portstr), "%d", port);
|
||||
if (host && (!*host || !strcmp(host, "*")))
|
||||
host = NULL;
|
||||
int rc = getaddrinfo(host, portstr, &hints, res);
|
||||
if (rc != 0) {
|
||||
errno = rc == EAI_SYSTEM ? errno : EADDRNOTAVAIL;
|
||||
return -1;
|
||||
}
|
||||
return 0;
|
||||
}
|
||||
|
||||
int iface_valid(const char *iface)
|
||||
{
|
||||
size_t n = strlen(iface);
|
||||
if (n == 0 || n >= IFNAMSIZ)
|
||||
return 0;
|
||||
for (const char *p = iface; *p; p++)
|
||||
if (!((*p >= 'a' && *p <= 'z') || (*p >= 'A' && *p <= 'Z') || (*p >= '0' && *p <= '9') ||
|
||||
*p == '_' || *p == '.' || *p == '-'))
|
||||
return 0;
|
||||
return 1;
|
||||
}
|
||||
|
||||
static int bind_device(int fd, const char *iface)
|
||||
{
|
||||
if (!iface || !*iface || !strcmp(iface, "-"))
|
||||
return 0;
|
||||
return setsockopt(fd, SOL_SOCKET, SO_BINDTODEVICE, iface, (socklen_t)strlen(iface));
|
||||
}
|
||||
|
||||
int net_listen(const char *addr, int port, int proto, const char *iface)
|
||||
{
|
||||
struct addrinfo *res;
|
||||
if (resolve(addr, port, proto, 1, &res) < 0)
|
||||
return -1;
|
||||
|
||||
int fd = -1, err = 0;
|
||||
for (struct addrinfo *ai = res; ai; ai = ai->ai_next) {
|
||||
fd = socket(ai->ai_family, ai->ai_socktype | SOCK_CLOEXEC | SOCK_NONBLOCK, 0);
|
||||
if (fd < 0) {
|
||||
err = errno;
|
||||
continue;
|
||||
}
|
||||
int one = 1;
|
||||
setsockopt(fd, SOL_SOCKET, SO_REUSEADDR, &one, sizeof(one));
|
||||
if (bind_device(fd, iface) == 0 && bind(fd, ai->ai_addr, ai->ai_addrlen) == 0 &&
|
||||
(proto == PROTO_UDP || listen(fd, 256) == 0))
|
||||
break;
|
||||
err = errno;
|
||||
close(fd);
|
||||
fd = -1;
|
||||
}
|
||||
freeaddrinfo(res);
|
||||
if (fd < 0)
|
||||
errno = err;
|
||||
return fd;
|
||||
}
|
||||
|
||||
int net_connect(const char *host, int port, int proto, const char *iface, int *in_progress)
|
||||
{
|
||||
struct addrinfo *res;
|
||||
*in_progress = 0;
|
||||
if (resolve(host, port, proto, 0, &res) < 0)
|
||||
return -1;
|
||||
|
||||
int fd = -1, err = 0;
|
||||
for (struct addrinfo *ai = res; ai; ai = ai->ai_next) {
|
||||
fd = socket(ai->ai_family, ai->ai_socktype | SOCK_CLOEXEC | SOCK_NONBLOCK, 0);
|
||||
if (fd < 0) {
|
||||
err = errno;
|
||||
continue;
|
||||
}
|
||||
if (bind_device(fd, iface) < 0) {
|
||||
err = errno;
|
||||
close(fd);
|
||||
fd = -1;
|
||||
continue;
|
||||
}
|
||||
if (connect(fd, ai->ai_addr, ai->ai_addrlen) == 0)
|
||||
break;
|
||||
if (errno == EINPROGRESS) {
|
||||
*in_progress = 1;
|
||||
break;
|
||||
}
|
||||
err = errno;
|
||||
close(fd);
|
||||
fd = -1;
|
||||
}
|
||||
freeaddrinfo(res);
|
||||
if (fd < 0)
|
||||
errno = err;
|
||||
else if (proto == PROTO_TCP)
|
||||
tune_stream(fd);
|
||||
return fd;
|
||||
}
|
||||
|
||||
int net_listen_unix(const char *path, int mode)
|
||||
{
|
||||
struct sockaddr_un sa;
|
||||
if (strlen(path) >= sizeof(sa.sun_path)) {
|
||||
errno = ENAMETOOLONG;
|
||||
return -1;
|
||||
}
|
||||
int fd = socket(AF_UNIX, SOCK_STREAM | SOCK_CLOEXEC | SOCK_NONBLOCK, 0);
|
||||
if (fd < 0)
|
||||
return -1;
|
||||
memset(&sa, 0, sizeof(sa));
|
||||
sa.sun_family = AF_UNIX;
|
||||
strcpy(sa.sun_path, path);
|
||||
unlink(path);
|
||||
// umask covers the window between bind and chmod.
|
||||
mode_t old = umask(0177);
|
||||
int rc = bind(fd, (struct sockaddr *)&sa, sizeof(sa));
|
||||
umask(old);
|
||||
if (rc < 0 || chmod(path, (mode_t)mode) < 0 || listen(fd, 64) < 0) {
|
||||
int e = errno;
|
||||
close(fd);
|
||||
errno = e;
|
||||
return -1;
|
||||
}
|
||||
return fd;
|
||||
}
|
||||
|
||||
int net_connect_unix(const char *path)
|
||||
{
|
||||
struct sockaddr_un sa;
|
||||
if (strlen(path) >= sizeof(sa.sun_path)) {
|
||||
errno = ENAMETOOLONG;
|
||||
return -1;
|
||||
}
|
||||
int fd = socket(AF_UNIX, SOCK_STREAM | SOCK_CLOEXEC, 0);
|
||||
if (fd < 0)
|
||||
return -1;
|
||||
memset(&sa, 0, sizeof(sa));
|
||||
sa.sun_family = AF_UNIX;
|
||||
strcpy(sa.sun_path, path);
|
||||
if (connect(fd, (struct sockaddr *)&sa, sizeof(sa)) < 0) {
|
||||
int e = errno;
|
||||
close(fd);
|
||||
errno = e;
|
||||
return -1;
|
||||
}
|
||||
return fd;
|
||||
}
|
||||
|
||||
void sockaddr_str(const struct sockaddr *sa, char *out, size_t outsz)
|
||||
{
|
||||
char host[INET6_ADDRSTRLEN] = "?";
|
||||
int port = 0;
|
||||
if (sa->sa_family == AF_INET) {
|
||||
const struct sockaddr_in *s4 = (const struct sockaddr_in *)sa;
|
||||
inet_ntop(AF_INET, &s4->sin_addr, host, sizeof(host));
|
||||
port = ntohs(s4->sin_port);
|
||||
snprintf(out, outsz, "%s:%d", host, port);
|
||||
} else if (sa->sa_family == AF_INET6) {
|
||||
const struct sockaddr_in6 *s6 = (const struct sockaddr_in6 *)sa;
|
||||
inet_ntop(AF_INET6, &s6->sin6_addr, host, sizeof(host));
|
||||
port = ntohs(s6->sin6_port);
|
||||
snprintf(out, outsz, "[%s]:%d", host, port);
|
||||
} else {
|
||||
snprintf(out, outsz, "?");
|
||||
}
|
||||
}
|
||||
29
backend/src/net.h
Normal file
29
backend/src/net.h
Normal file
@@ -0,0 +1,29 @@
|
||||
#ifndef PB_NET_H
|
||||
#define PB_NET_H
|
||||
|
||||
#include <sys/socket.h>
|
||||
|
||||
enum pb_proto { PROTO_TCP, PROTO_UDP };
|
||||
|
||||
int proto_parse(const char *s); // -1 on unknown
|
||||
const char *proto_name(int proto);
|
||||
|
||||
int set_nonblock(int fd);
|
||||
void tune_stream(int fd);
|
||||
|
||||
// Binds a listening (TCP) or bound (UDP) socket. iface (NULL/"" = any) pins the
|
||||
// socket to one interface with SO_BINDTODEVICE. Returns fd or -1 with errno.
|
||||
int net_listen(const char *addr, int port, int proto, const char *iface);
|
||||
// Starts a non-blocking connect, optionally pinned to iface. Returns fd or -1;
|
||||
// *in_progress set if pending.
|
||||
int net_connect(const char *host, int port, int proto, const char *iface, int *in_progress);
|
||||
int net_listen_unix(const char *path, int mode);
|
||||
int net_connect_unix(const char *path);
|
||||
|
||||
// 1 if iface is a plausible interface name (also safe as a protocol field).
|
||||
int iface_valid(const char *iface);
|
||||
|
||||
// Formats a sockaddr as "addr:port" (IPv6 in brackets).
|
||||
void sockaddr_str(const struct sockaddr *sa, char *out, size_t outsz);
|
||||
|
||||
#endif
|
||||
152
backend/src/proto.c
Normal file
152
backend/src/proto.c
Normal file
@@ -0,0 +1,152 @@
|
||||
#include <stdarg.h>
|
||||
#include <stdio.h>
|
||||
#include <stdlib.h>
|
||||
#include <string.h>
|
||||
#include <sys/random.h>
|
||||
|
||||
#include "proto.h"
|
||||
|
||||
/* Byte buffer */
|
||||
|
||||
void buf_init(struct buf *b)
|
||||
{
|
||||
b->data = NULL;
|
||||
b->len = 0;
|
||||
b->cap = 0;
|
||||
}
|
||||
|
||||
void buf_free(struct buf *b)
|
||||
{
|
||||
free(b->data);
|
||||
buf_init(b);
|
||||
}
|
||||
|
||||
static int buf_reserve(struct buf *b, size_t extra)
|
||||
{
|
||||
if (b->len + extra <= b->cap)
|
||||
return 0;
|
||||
size_t cap = b->cap ? b->cap : 4096;
|
||||
while (cap < b->len + extra)
|
||||
cap *= 2;
|
||||
char *p = realloc(b->data, cap);
|
||||
if (!p)
|
||||
return -1;
|
||||
b->data = p;
|
||||
b->cap = cap;
|
||||
return 0;
|
||||
}
|
||||
|
||||
int buf_append(struct buf *b, const void *p, size_t n)
|
||||
{
|
||||
if (buf_reserve(b, n) < 0)
|
||||
return -1;
|
||||
memcpy(b->data + b->len, p, n);
|
||||
b->len += n;
|
||||
return 0;
|
||||
}
|
||||
|
||||
void buf_consume(struct buf *b, size_t n)
|
||||
{
|
||||
if (n >= b->len) {
|
||||
b->len = 0;
|
||||
return;
|
||||
}
|
||||
memmove(b->data, b->data + n, b->len - n);
|
||||
b->len -= n;
|
||||
}
|
||||
|
||||
int buf_printf(struct buf *b, const char *fmt, ...)
|
||||
{
|
||||
va_list ap;
|
||||
va_start(ap, fmt);
|
||||
int n = vsnprintf(NULL, 0, fmt, ap);
|
||||
va_end(ap);
|
||||
if (n < 0 || buf_reserve(b, (size_t)n + 1) < 0)
|
||||
return -1;
|
||||
va_start(ap, fmt);
|
||||
vsnprintf(b->data + b->len, (size_t)n + 1, fmt, ap);
|
||||
va_end(ap);
|
||||
b->len += (size_t)n;
|
||||
return n;
|
||||
}
|
||||
|
||||
/* Line protocol helpers */
|
||||
|
||||
int line_next(struct buf *b, char *out, size_t outsz)
|
||||
{
|
||||
char *nl = b->len ? memchr(b->data, '\n', b->len) : NULL;
|
||||
if (!nl)
|
||||
return b->len >= PB_MAX_LINE ? -1 : 0;
|
||||
|
||||
size_t n = (size_t)(nl - b->data);
|
||||
if (n >= outsz || n >= PB_MAX_LINE)
|
||||
return -1;
|
||||
memcpy(out, b->data, n);
|
||||
if (n > 0 && out[n - 1] == '\r')
|
||||
n--;
|
||||
out[n] = '\0';
|
||||
buf_consume(b, (size_t)(nl - b->data) + 1);
|
||||
return 1;
|
||||
}
|
||||
|
||||
int line_split(char *line, char **fields, int max)
|
||||
{
|
||||
int n = 0;
|
||||
char *save = NULL;
|
||||
for (char *t = strtok_r(line, " ", &save); t && n < max; t = strtok_r(NULL, " ", &save))
|
||||
fields[n++] = t;
|
||||
return n;
|
||||
}
|
||||
|
||||
void field_sanitise(char *s)
|
||||
{
|
||||
if (!*s) {
|
||||
// Empty fields would shift the field count on the other side.
|
||||
s[0] = '-';
|
||||
s[1] = '\0';
|
||||
return;
|
||||
}
|
||||
for (; *s; s++) {
|
||||
unsigned char ch = (unsigned char)*s;
|
||||
if (ch <= ' ' || ch >= 127)
|
||||
*s = '_';
|
||||
}
|
||||
}
|
||||
|
||||
/* UDP framing */
|
||||
|
||||
int udp_frame_append(struct buf *b, const void *payload, size_t n)
|
||||
{
|
||||
if (n > 0xffff)
|
||||
return -1;
|
||||
unsigned char hdr[2] = { (unsigned char)(n >> 8), (unsigned char)(n & 0xff) };
|
||||
if (buf_append(b, hdr, 2) < 0 || buf_append(b, payload, n) < 0)
|
||||
return -1;
|
||||
return 0;
|
||||
}
|
||||
|
||||
int udp_frame_peek(const struct buf *b, const char **payload)
|
||||
{
|
||||
if (b->len < 2)
|
||||
return 0;
|
||||
const unsigned char *p = (const unsigned char *)b->data;
|
||||
size_t n = ((size_t)p[0] << 8) | p[1];
|
||||
if (b->len < 2 + n)
|
||||
return 0;
|
||||
*payload = b->data + 2;
|
||||
// A zero-length datagram is valid; signal it with -2 to keep 0 = incomplete.
|
||||
return n ? (int)n : -2;
|
||||
}
|
||||
|
||||
int random_token(char *out, size_t outsz)
|
||||
{
|
||||
unsigned char raw[PB_TOKEN_LEN / 2];
|
||||
if (outsz < PB_TOKEN_LEN + 1)
|
||||
return -1;
|
||||
if (getrandom(raw, sizeof(raw), 0) != (ssize_t)sizeof(raw))
|
||||
return -1;
|
||||
for (size_t i = 0; i < sizeof(raw); i++)
|
||||
sprintf(out + i * 2, "%02x", raw[i]);
|
||||
out[PB_TOKEN_LEN] = '\0';
|
||||
return 0;
|
||||
}
|
||||
73
backend/src/proto.h
Normal file
73
backend/src/proto.h
Normal file
@@ -0,0 +1,73 @@
|
||||
#ifndef PB_PROTO_H
|
||||
#define PB_PROTO_H
|
||||
|
||||
#include <stddef.h>
|
||||
#include <stdint.h>
|
||||
|
||||
/* Control protocol
|
||||
*
|
||||
* One text line per message, fields separated by single spaces, no field may
|
||||
* contain whitespace. Carried over the SSH exec channel whose forced command
|
||||
* (patchbayd --relay) bridges it to the hub socket on the target.
|
||||
*
|
||||
* hub -> client: HELLO <token> <hub_port>
|
||||
* SINKS-BEGIN / SINK <sink_id> <proto> <bind> <port> <iface|-> / SINKS-END
|
||||
* OPEN <conn_id> <proto> <host> <port> <iface|->
|
||||
* SVC-REQ, PING
|
||||
* client -> hub: HELLO <version> <hostname>
|
||||
* SVC-BEGIN / SVC <proto> <addr> <port> <pid> <process> / SVC-END
|
||||
* IF-BEGIN / IF <name> <ipv4/prefix|-> / IF-END (tun interfaces)
|
||||
* SINKSTATE <sink_id> <ok|err> <reason>
|
||||
* PONG
|
||||
*
|
||||
* <iface> pins a socket to an interface (tunnel nodes, SO_BINDTODEVICE).
|
||||
*
|
||||
* Data channels are direct-tcpip channels to 127.0.0.1:HubPort on the target.
|
||||
* They start with "PB1 <token> SINK <sink_id>\n" (a client sink accepted a
|
||||
* connection) or "PB1 <token> OPEN <conn_id>\n" (answer to OPEN), followed by
|
||||
* raw stream bytes. UDP flows use length-prefixed frames, see udp_frame_*.
|
||||
*/
|
||||
|
||||
#define PB_VERSION "0.1.0"
|
||||
#define PB_TOKEN_LEN 32 // hex chars
|
||||
#define PB_MAX_LINE 1024
|
||||
#define PB_MAX_FIELDS 8
|
||||
#define PB_UDP_MAX 65507
|
||||
|
||||
/* Byte buffer */
|
||||
|
||||
struct buf {
|
||||
char *data;
|
||||
size_t len;
|
||||
size_t cap;
|
||||
};
|
||||
|
||||
void buf_init(struct buf *b);
|
||||
void buf_free(struct buf *b);
|
||||
int buf_append(struct buf *b, const void *p, size_t n);
|
||||
void buf_consume(struct buf *b, size_t n);
|
||||
int buf_printf(struct buf *b, const char *fmt, ...) __attribute__((format(printf, 2, 3)));
|
||||
|
||||
/* Line protocol helpers */
|
||||
|
||||
// Extracts the next complete line (without '\n', '\r' stripped) from b into
|
||||
// out. Returns 1 if a line was extracted, 0 if incomplete, -1 if a line exceeds
|
||||
// PB_MAX_LINE.
|
||||
int line_next(struct buf *b, char *out, size_t outsz);
|
||||
|
||||
// Splits a line in place on spaces. Returns number of fields.
|
||||
int line_split(char *line, char **fields, int max);
|
||||
|
||||
// Replaces characters that are not safe for a protocol field with '_'.
|
||||
void field_sanitise(char *s);
|
||||
|
||||
/* UDP framing: 2 byte big endian length + payload */
|
||||
|
||||
int udp_frame_append(struct buf *b, const void *payload, size_t n);
|
||||
// Returns payload length and sets *payload if a full frame is at the start of
|
||||
// b->data, 0 if incomplete. Caller consumes 2 + length bytes afterwards.
|
||||
int udp_frame_peek(const struct buf *b, const char **payload);
|
||||
|
||||
int random_token(char *out, size_t outsz);
|
||||
|
||||
#endif
|
||||
66
backend/src/relay.c
Normal file
66
backend/src/relay.c
Normal file
@@ -0,0 +1,66 @@
|
||||
#include <errno.h>
|
||||
#include <poll.h>
|
||||
#include <stdio.h>
|
||||
#include <stdlib.h>
|
||||
#include <string.h>
|
||||
#include <unistd.h>
|
||||
|
||||
#include "log.h"
|
||||
#include "net.h"
|
||||
#include "relay.h"
|
||||
|
||||
// Copies whatever is readable from one fd to the other. 0 on EOF/error.
|
||||
static int copy_once(int from, int to)
|
||||
{
|
||||
char buf[16384];
|
||||
ssize_t n = read(from, buf, sizeof(buf));
|
||||
if (n <= 0)
|
||||
return 0;
|
||||
for (ssize_t off = 0; off < n;) {
|
||||
ssize_t w = write(to, buf + off, (size_t)(n - off));
|
||||
if (w < 0) {
|
||||
if (errno == EINTR)
|
||||
continue;
|
||||
return 0;
|
||||
}
|
||||
off += w;
|
||||
}
|
||||
return 1;
|
||||
}
|
||||
|
||||
int relay_run(const struct pb_config *cfg, const char *client_id)
|
||||
{
|
||||
char path[300];
|
||||
snprintf(path, sizeof(path), "%s/hub.sock", cfg->run_dir);
|
||||
int fd = net_connect_unix(path);
|
||||
if (fd < 0) {
|
||||
log_err("relay: cannot reach hub at %s: %s", path, strerror(errno));
|
||||
return 1;
|
||||
}
|
||||
|
||||
// sshd sets SSH_CONNECTION="client_ip client_port server_ip server_port".
|
||||
char addr[64] = "-";
|
||||
const char *conn = getenv("SSH_CONNECTION");
|
||||
if (conn && sscanf(conn, "%63s", addr) != 1)
|
||||
strcpy(addr, "-");
|
||||
|
||||
char hello[128];
|
||||
int n = snprintf(hello, sizeof(hello), "RELAY %d %s\n", atoi(client_id), addr);
|
||||
if (write(fd, hello, (size_t)n) != n)
|
||||
return 1;
|
||||
|
||||
struct pollfd p[2] = { { .fd = 0, .events = POLLIN }, { .fd = fd, .events = POLLIN } };
|
||||
for (;;) {
|
||||
if (poll(p, 2, -1) < 0) {
|
||||
if (errno == EINTR)
|
||||
continue;
|
||||
break;
|
||||
}
|
||||
if (p[0].revents && !copy_once(0, fd))
|
||||
break;
|
||||
if (p[1].revents && !copy_once(fd, 1))
|
||||
break;
|
||||
}
|
||||
close(fd);
|
||||
return 0;
|
||||
}
|
||||
10
backend/src/relay.h
Normal file
10
backend/src/relay.h
Normal file
@@ -0,0 +1,10 @@
|
||||
#ifndef PB_RELAY_H
|
||||
#define PB_RELAY_H
|
||||
|
||||
#include "config.h"
|
||||
|
||||
// Forced command of every client key in authorized_keys: bridges the SSH exec
|
||||
// channel (stdin/stdout) to the hub socket, announcing the client id.
|
||||
int relay_run(const struct pb_config *cfg, const char *client_id);
|
||||
|
||||
#endif
|
||||
213
backend/src/services.c
Normal file
213
backend/src/services.c
Normal file
@@ -0,0 +1,213 @@
|
||||
#include <arpa/inet.h>
|
||||
#include <ctype.h>
|
||||
#include <dirent.h>
|
||||
#include <ifaddrs.h>
|
||||
#include <net/if.h>
|
||||
#include <stdio.h>
|
||||
#include <stdlib.h>
|
||||
#include <string.h>
|
||||
#include <unistd.h>
|
||||
|
||||
#include "net.h"
|
||||
#include "services.h"
|
||||
|
||||
#define TCP_LISTEN 0x0A
|
||||
#define TCP_CLOSE 0x07
|
||||
|
||||
int services_parse_line(const char *line, int v6, int udp, struct svc_entry *e)
|
||||
{
|
||||
char local[64], remote[64];
|
||||
unsigned int state;
|
||||
unsigned long inode;
|
||||
|
||||
// sl local rem st tx:rx tr:when retrnsmt uid timeout inode
|
||||
if (sscanf(line, " %*d: %63s %63s %x %*s %*s %*s %*u %*u %lu",
|
||||
local, remote, &state, &inode) != 4)
|
||||
return -1;
|
||||
|
||||
char *lport = strchr(local, ':');
|
||||
char *rport = strchr(remote, ':');
|
||||
if (!lport || !rport)
|
||||
return -1;
|
||||
*lport++ = '\0';
|
||||
|
||||
if (udp) {
|
||||
// Unconnected UDP sockets have no remote port and sit in TCP_CLOSE.
|
||||
if (state != TCP_CLOSE || strtoul(rport + 1, NULL, 16) != 0)
|
||||
return 0;
|
||||
} else if (state != TCP_LISTEN) {
|
||||
return 0;
|
||||
}
|
||||
|
||||
memset(e, 0, sizeof(*e));
|
||||
e->proto = udp ? PROTO_UDP : PROTO_TCP;
|
||||
e->port = (int)strtoul(lport, NULL, 16);
|
||||
e->inode = inode;
|
||||
|
||||
// The kernel prints the raw 32 bit words, so copying them back in host
|
||||
// order restores the network byte order address.
|
||||
if (v6) {
|
||||
struct in6_addr a;
|
||||
if (strlen(local) != 32)
|
||||
return -1;
|
||||
for (int i = 0; i < 4; i++) {
|
||||
char word[9];
|
||||
memcpy(word, local + i * 8, 8);
|
||||
word[8] = '\0';
|
||||
unsigned int w = (unsigned int)strtoul(word, NULL, 16);
|
||||
memcpy((char *)&a + i * 4, &w, 4);
|
||||
}
|
||||
inet_ntop(AF_INET6, &a, e->addr, sizeof(e->addr));
|
||||
} else {
|
||||
struct in_addr a;
|
||||
unsigned int w = (unsigned int)strtoul(local, NULL, 16);
|
||||
memcpy(&a, &w, 4);
|
||||
inet_ntop(AF_INET, &a, e->addr, sizeof(e->addr));
|
||||
}
|
||||
return 1;
|
||||
}
|
||||
|
||||
static int scan_file(const char *path, int v6, int udp, struct svc_entry **arr, int *n, int *cap)
|
||||
{
|
||||
FILE *f = fopen(path, "r");
|
||||
if (!f)
|
||||
return 0; // e.g. no IPv6
|
||||
char line[512];
|
||||
if (!fgets(line, sizeof(line), f)) { // header
|
||||
fclose(f);
|
||||
return 0;
|
||||
}
|
||||
while (fgets(line, sizeof(line), f)) {
|
||||
struct svc_entry e;
|
||||
if (services_parse_line(line, v6, udp, &e) != 1)
|
||||
continue;
|
||||
if (*n == *cap) {
|
||||
int nc = *cap ? *cap * 2 : 64;
|
||||
struct svc_entry *p = realloc(*arr, (size_t)nc * sizeof(*p));
|
||||
if (!p) {
|
||||
fclose(f);
|
||||
return -1;
|
||||
}
|
||||
*arr = p;
|
||||
*cap = nc;
|
||||
}
|
||||
(*arr)[(*n)++] = e;
|
||||
}
|
||||
fclose(f);
|
||||
return 0;
|
||||
}
|
||||
|
||||
static void map_pids(struct svc_entry *arr, int n)
|
||||
{
|
||||
DIR *proc = opendir("/proc");
|
||||
if (!proc)
|
||||
return;
|
||||
struct dirent *de;
|
||||
while ((de = readdir(proc))) {
|
||||
if (!isdigit((unsigned char)de->d_name[0]))
|
||||
continue;
|
||||
int pid = atoi(de->d_name);
|
||||
char fddir[64];
|
||||
snprintf(fddir, sizeof(fddir), "/proc/%d/fd", pid);
|
||||
DIR *fds = opendir(fddir);
|
||||
if (!fds)
|
||||
continue;
|
||||
struct dirent *fe;
|
||||
while ((fe = readdir(fds))) {
|
||||
char lpath[320], target[64];
|
||||
snprintf(lpath, sizeof(lpath), "%s/%s", fddir, fe->d_name);
|
||||
ssize_t l = readlink(lpath, target, sizeof(target) - 1);
|
||||
if (l <= 0)
|
||||
continue;
|
||||
target[l] = '\0';
|
||||
unsigned long inode;
|
||||
if (sscanf(target, "socket:[%lu]", &inode) != 1)
|
||||
continue;
|
||||
for (int i = 0; i < n; i++) {
|
||||
if (arr[i].inode != inode || arr[i].pid)
|
||||
continue;
|
||||
arr[i].pid = pid;
|
||||
char cpath[64];
|
||||
snprintf(cpath, sizeof(cpath), "/proc/%d/comm", pid);
|
||||
FILE *cf = fopen(cpath, "r");
|
||||
if (cf) {
|
||||
if (fgets(arr[i].process, sizeof(arr[i].process), cf))
|
||||
arr[i].process[strcspn(arr[i].process, "\n")] = '\0';
|
||||
fclose(cf);
|
||||
}
|
||||
}
|
||||
}
|
||||
closedir(fds);
|
||||
}
|
||||
closedir(proc);
|
||||
}
|
||||
|
||||
int services_scan(struct svc_entry **out)
|
||||
{
|
||||
struct svc_entry *arr = NULL;
|
||||
int n = 0, cap = 0;
|
||||
if (scan_file("/proc/net/tcp", 0, 0, &arr, &n, &cap) < 0 ||
|
||||
scan_file("/proc/net/tcp6", 1, 0, &arr, &n, &cap) < 0 ||
|
||||
scan_file("/proc/net/udp", 0, 1, &arr, &n, &cap) < 0 ||
|
||||
scan_file("/proc/net/udp6", 1, 1, &arr, &n, &cap) < 0) {
|
||||
free(arr);
|
||||
return -1;
|
||||
}
|
||||
map_pids(arr, n);
|
||||
|
||||
// SO_REUSEPORT and multiple workers produce duplicates; keep the first.
|
||||
int m = 0;
|
||||
for (int i = 0; i < n; i++) {
|
||||
int dup = 0;
|
||||
for (int j = 0; j < m && !dup; j++)
|
||||
dup = arr[j].proto == arr[i].proto && arr[j].port == arr[i].port &&
|
||||
!strcmp(arr[j].addr, arr[i].addr);
|
||||
if (!dup)
|
||||
arr[m++] = arr[i];
|
||||
}
|
||||
for (int i = 0; i < m; i++)
|
||||
if (!arr[i].process[0])
|
||||
strcpy(arr[i].process, "-");
|
||||
*out = arr;
|
||||
return m;
|
||||
}
|
||||
|
||||
int ifaces_scan(struct iface_entry **out)
|
||||
{
|
||||
struct ifaddrs *ifa;
|
||||
if (getifaddrs(&ifa) < 0)
|
||||
return -1;
|
||||
struct iface_entry *arr = NULL;
|
||||
int n = 0, cap = 0;
|
||||
for (struct ifaddrs *i = ifa; i; i = i->ifa_next) {
|
||||
if (strncmp(i->ifa_name, "tun", 3))
|
||||
continue;
|
||||
int k;
|
||||
for (k = 0; k < n && strcmp(arr[k].name, i->ifa_name); k++)
|
||||
;
|
||||
if (k == n) {
|
||||
if (n == cap) {
|
||||
cap = cap ? cap * 2 : 8;
|
||||
struct iface_entry *p = realloc(arr, (size_t)cap * sizeof(*p));
|
||||
if (!p)
|
||||
break;
|
||||
arr = p;
|
||||
}
|
||||
snprintf(arr[n].name, sizeof(arr[n].name), "%s", i->ifa_name);
|
||||
strcpy(arr[n].addr, "-");
|
||||
n++;
|
||||
}
|
||||
// First IPv4 address wins; interfaces without one keep "-".
|
||||
if (i->ifa_addr && i->ifa_addr->sa_family == AF_INET && !strcmp(arr[k].addr, "-")) {
|
||||
char a[INET_ADDRSTRLEN];
|
||||
int prefix = 0;
|
||||
inet_ntop(AF_INET, &((struct sockaddr_in *)i->ifa_addr)->sin_addr, a, sizeof(a));
|
||||
if (i->ifa_netmask)
|
||||
prefix = __builtin_popcount(((struct sockaddr_in *)i->ifa_netmask)->sin_addr.s_addr);
|
||||
snprintf(arr[k].addr, sizeof(arr[k].addr), "%s/%d", a, prefix);
|
||||
}
|
||||
}
|
||||
freeifaddrs(ifa);
|
||||
*out = arr;
|
||||
return n;
|
||||
}
|
||||
30
backend/src/services.h
Normal file
30
backend/src/services.h
Normal file
@@ -0,0 +1,30 @@
|
||||
#ifndef PB_SERVICES_H
|
||||
#define PB_SERVICES_H
|
||||
|
||||
struct svc_entry {
|
||||
int proto;
|
||||
char addr[64];
|
||||
int port;
|
||||
int pid; // 0 if unknown
|
||||
char process[32];
|
||||
unsigned long inode;
|
||||
};
|
||||
|
||||
// Collects listening TCP and bound, unconnected UDP sockets with their owning
|
||||
// processes from /proc. Returns count (>= 0) and a malloc'd array in *out.
|
||||
int services_scan(struct svc_entry **out);
|
||||
|
||||
// Parses one data line of /proc/net/{tcp,udp}[6]. Returns 1 if it describes a
|
||||
// listening socket, 0 if not, -1 on parse error.
|
||||
int services_parse_line(const char *line, int v6, int udp, struct svc_entry *e);
|
||||
|
||||
struct iface_entry {
|
||||
char name[16];
|
||||
char addr[64]; // "a.b.c.d/prefix" or "-" without IPv4
|
||||
};
|
||||
|
||||
// Lists tun* interfaces with their IPv4 addresses (one entry per interface
|
||||
// without an address). Returns count and a malloc'd array in *out.
|
||||
int ifaces_scan(struct iface_entry **out);
|
||||
|
||||
#endif
|
||||
141
backend/src/stats.c
Normal file
141
backend/src/stats.c
Normal file
@@ -0,0 +1,141 @@
|
||||
#include <stdlib.h>
|
||||
#include <string.h>
|
||||
|
||||
#include "db.h"
|
||||
#include "log.h"
|
||||
#include "stats.h"
|
||||
|
||||
// Entries are allocated individually and never freed, so connections can keep
|
||||
// pointers to them across reloads.
|
||||
static struct sink_stats **all;
|
||||
static int nall, capall;
|
||||
|
||||
struct sink_stats *stats_get(int sink_id)
|
||||
{
|
||||
for (int i = 0; i < nall; i++)
|
||||
if (all[i]->sink_id == sink_id)
|
||||
return all[i];
|
||||
if (nall == capall) {
|
||||
int nc = capall ? capall * 2 : 32;
|
||||
struct sink_stats **p = realloc(all, (size_t)nc * sizeof(*p));
|
||||
if (!p)
|
||||
return NULL;
|
||||
all = p;
|
||||
capall = nc;
|
||||
}
|
||||
struct sink_stats *s = calloc(1, sizeof(*s));
|
||||
if (!s)
|
||||
return NULL;
|
||||
s->sink_id = sink_id;
|
||||
all[nall++] = s;
|
||||
return s;
|
||||
}
|
||||
|
||||
void stats_add(struct sink_stats *s, uint64_t in, uint64_t out)
|
||||
{
|
||||
s->sec_in += in;
|
||||
s->sec_out += out;
|
||||
s->min_in += in;
|
||||
s->min_out += out;
|
||||
s->total_in += in;
|
||||
s->total_out += out;
|
||||
}
|
||||
|
||||
void stats_conn_open(struct sink_stats *s)
|
||||
{
|
||||
s->active++;
|
||||
s->min_conns++;
|
||||
}
|
||||
|
||||
void stats_conn_close(struct sink_stats *s)
|
||||
{
|
||||
if (s->active > 0)
|
||||
s->active--;
|
||||
}
|
||||
|
||||
void stats_tick_second(double dt)
|
||||
{
|
||||
if (dt <= 0)
|
||||
dt = 1;
|
||||
for (int i = 0; i < nall; i++) {
|
||||
all[i]->rate_in = (double)all[i]->sec_in / dt;
|
||||
all[i]->rate_out = (double)all[i]->sec_out / dt;
|
||||
all[i]->sec_in = all[i]->sec_out = 0;
|
||||
}
|
||||
}
|
||||
|
||||
static void upsert(sqlite3_stmt *st, int sink, const char *tier, long ts,
|
||||
uint64_t in, uint64_t out, uint32_t conns)
|
||||
{
|
||||
sqlite3_reset(st);
|
||||
sqlite3_bind_int(st, 1, sink);
|
||||
sqlite3_bind_text(st, 2, tier, -1, SQLITE_STATIC);
|
||||
sqlite3_bind_int64(st, 3, ts);
|
||||
sqlite3_bind_int64(st, 4, (sqlite3_int64)in);
|
||||
sqlite3_bind_int64(st, 5, (sqlite3_int64)out);
|
||||
sqlite3_bind_int(st, 6, (int)conns);
|
||||
if (sqlite3_step(st) != SQLITE_DONE)
|
||||
log_warn("stats insert failed");
|
||||
}
|
||||
|
||||
static void prune(sqlite3 *db, const char *tier, long cutoff)
|
||||
{
|
||||
sqlite3_stmt *st;
|
||||
if (sqlite3_prepare_v2(db, "DELETE FROM stats WHERE tier = ? AND ts < ?", -1, &st, NULL) != SQLITE_OK)
|
||||
return;
|
||||
sqlite3_bind_text(st, 1, tier, -1, SQLITE_STATIC);
|
||||
sqlite3_bind_int64(st, 2, cutoff);
|
||||
sqlite3_step(st);
|
||||
sqlite3_finalize(st);
|
||||
}
|
||||
|
||||
void stats_flush(sqlite3 *db, long now)
|
||||
{
|
||||
long minute = now - now % 60, hour = now - now % 3600, day = now - now % 86400;
|
||||
sqlite3_stmt *st;
|
||||
const char *sql =
|
||||
"INSERT INTO stats (sink_id, tier, ts, bytes_in, bytes_out, conns) VALUES (?, ?, ?, ?, ?, ?) "
|
||||
"ON CONFLICT (sink_id, tier, ts) DO UPDATE SET "
|
||||
"bytes_in = bytes_in + excluded.bytes_in, bytes_out = bytes_out + excluded.bytes_out, "
|
||||
"conns = conns + excluded.conns";
|
||||
if (sqlite3_prepare_v2(db, sql, -1, &st, NULL) != SQLITE_OK) {
|
||||
log_err("stats: %s", sqlite3_errmsg(db));
|
||||
return;
|
||||
}
|
||||
db_exec(db, "BEGIN");
|
||||
for (int i = 0; i < nall; i++) {
|
||||
struct sink_stats *s = all[i];
|
||||
if (!s->min_in && !s->min_out && !s->min_conns)
|
||||
continue;
|
||||
upsert(st, s->sink_id, "m", minute, s->min_in, s->min_out, s->min_conns);
|
||||
upsert(st, s->sink_id, "h", hour, s->min_in, s->min_out, s->min_conns);
|
||||
upsert(st, s->sink_id, "d", day, s->min_in, s->min_out, s->min_conns);
|
||||
s->min_in = s->min_out = 0;
|
||||
s->min_conns = 0;
|
||||
}
|
||||
sqlite3_finalize(st);
|
||||
|
||||
long mh = db_setting_int(db, "stats_minute_hours", 48);
|
||||
long hd = db_setting_int(db, "stats_hour_days", 90);
|
||||
long dd = db_setting_int(db, "stats_day_days", 0);
|
||||
if (mh > 0)
|
||||
prune(db, "m", now - mh * 3600);
|
||||
if (hd > 0)
|
||||
prune(db, "h", now - hd * 86400);
|
||||
if (dd > 0)
|
||||
prune(db, "d", now - dd * 86400);
|
||||
db_exec(db, "COMMIT");
|
||||
}
|
||||
|
||||
int stats_live_json(struct buf *b)
|
||||
{
|
||||
buf_append(b, "{", 1);
|
||||
for (int i = 0; i < nall; i++) {
|
||||
struct sink_stats *s = all[i];
|
||||
buf_printf(b, "%s\"%d\":{\"in\":%.0f,\"out\":%.0f,\"active\":%d,"
|
||||
"\"total_in\":%llu,\"total_out\":%llu}",
|
||||
i ? "," : "", s->sink_id, s->rate_in, s->rate_out, s->active,
|
||||
(unsigned long long)s->total_in, (unsigned long long)s->total_out);
|
||||
}
|
||||
return buf_append(b, "}", 1);
|
||||
}
|
||||
33
backend/src/stats.h
Normal file
33
backend/src/stats.h
Normal file
@@ -0,0 +1,33 @@
|
||||
#ifndef PB_STATS_H
|
||||
#define PB_STATS_H
|
||||
|
||||
#include <sqlite3.h>
|
||||
#include <stdint.h>
|
||||
|
||||
#include "proto.h"
|
||||
|
||||
// Traffic counters per connection (= per sink node), kept on the target.
|
||||
struct sink_stats {
|
||||
int sink_id;
|
||||
uint64_t sec_in, sec_out; // current second
|
||||
uint64_t min_in, min_out; // current minute
|
||||
uint32_t min_conns;
|
||||
double rate_in, rate_out; // bytes/s over the last second
|
||||
int active; // open connections
|
||||
uint64_t total_in, total_out; // since daemon start
|
||||
};
|
||||
|
||||
struct sink_stats *stats_get(int sink_id);
|
||||
void stats_add(struct sink_stats *s, uint64_t in, uint64_t out);
|
||||
void stats_conn_open(struct sink_stats *s);
|
||||
void stats_conn_close(struct sink_stats *s);
|
||||
|
||||
void stats_tick_second(double dt);
|
||||
// Writes the minute's counters (also rolled into hour and day rows) and
|
||||
// prunes rows past their retention.
|
||||
void stats_flush(sqlite3 *db, long now);
|
||||
|
||||
// {"<sink_id>": {"in": rate, "out": rate, "active": n, "total_in": .., ...}, ...}
|
||||
int stats_live_json(struct buf *b);
|
||||
|
||||
#endif
|
||||
145
backend/tests/test_main.c
Normal file
145
backend/tests/test_main.c
Normal file
@@ -0,0 +1,145 @@
|
||||
/* Minimal unit test runner: ./test_runner [test_name ...] */
|
||||
|
||||
#include <stdio.h>
|
||||
#include <stdlib.h>
|
||||
#include <string.h>
|
||||
|
||||
#include "../src/config.h"
|
||||
#include "../src/json.h"
|
||||
#include "../src/net.h"
|
||||
#include "../src/proto.h"
|
||||
#include "../src/services.h"
|
||||
|
||||
static int failures;
|
||||
|
||||
#define CHECK(cond) do { \
|
||||
if (!(cond)) { \
|
||||
fprintf(stderr, " %s:%d: CHECK(%s) failed\n", __FILE__, __LINE__, #cond); \
|
||||
failures++; \
|
||||
} \
|
||||
} while (0)
|
||||
|
||||
/* Tests */
|
||||
|
||||
static void test_config(void)
|
||||
{
|
||||
struct pb_config c;
|
||||
config_defaults(&c);
|
||||
char l1[] = "Role = target";
|
||||
char l2[] = " TargetPort=2200 # not a comment, atoi stops";
|
||||
char l3[] = "# comment";
|
||||
char l4[] = "IdentityFile = \" /etc/x \"";
|
||||
char l5[] = "SysopPassword = secret";
|
||||
char l6[] = "garbage";
|
||||
CHECK(config_parse_line(&c, l1) == 0 && c.role == ROLE_TARGET);
|
||||
CHECK(config_parse_line(&c, l2) == 0 && c.target_port == 2200);
|
||||
CHECK(config_parse_line(&c, l3) == 0);
|
||||
CHECK(config_parse_line(&c, l4) == 0 && !strcmp(c.identity_file, " /etc/x "));
|
||||
CHECK(config_parse_line(&c, l5) == 0);
|
||||
CHECK(config_parse_line(&c, l6) == -1);
|
||||
}
|
||||
|
||||
static void test_lines(void)
|
||||
{
|
||||
struct buf b;
|
||||
buf_init(&b);
|
||||
char out[PB_MAX_LINE];
|
||||
buf_append(&b, "HELLO abc 7701\r\nSINK 1 tcp", 26);
|
||||
CHECK(line_next(&b, out, sizeof(out)) == 1 && !strcmp(out, "HELLO abc 7701"));
|
||||
CHECK(line_next(&b, out, sizeof(out)) == 0);
|
||||
buf_append(&b, " 0.0.0.0 80\n", 12);
|
||||
CHECK(line_next(&b, out, sizeof(out)) == 1);
|
||||
char *f[PB_MAX_FIELDS];
|
||||
CHECK(line_split(out, f, PB_MAX_FIELDS) == 5 && !strcmp(f[4], "80"));
|
||||
buf_free(&b);
|
||||
|
||||
char s1[16] = "a b\tc";
|
||||
field_sanitise(s1);
|
||||
CHECK(!strcmp(s1, "a_b_c"));
|
||||
char s2[4] = "";
|
||||
field_sanitise(s2);
|
||||
CHECK(!strcmp(s2, "-"));
|
||||
}
|
||||
|
||||
static void test_udp_frames(void)
|
||||
{
|
||||
struct buf b;
|
||||
buf_init(&b);
|
||||
const char *p;
|
||||
CHECK(udp_frame_append(&b, "hello", 5) == 0);
|
||||
CHECK(udp_frame_append(&b, "", 0) == 0);
|
||||
CHECK(udp_frame_peek(&b, &p) == 5 && !memcmp(p, "hello", 5));
|
||||
buf_consume(&b, 7);
|
||||
CHECK(udp_frame_peek(&b, &p) == -2);
|
||||
buf_consume(&b, 2);
|
||||
CHECK(udp_frame_peek(&b, &p) == 0);
|
||||
buf_append(&b, "\x00\x03" "ab", 4);
|
||||
CHECK(udp_frame_peek(&b, &p) == 0); // incomplete
|
||||
buf_free(&b);
|
||||
}
|
||||
|
||||
static void test_services_parse(void)
|
||||
{
|
||||
struct svc_entry e;
|
||||
// 127.0.0.1:22 LISTEN, as printed on a little endian machine.
|
||||
const char *tcp = " 0: 0100007F:0016 00000000:0000 0A 00000000:00000000 00:00000000 00000000 0 0 12345 1 0000000000000000 100 0 0 10 0";
|
||||
CHECK(services_parse_line(tcp, 0, 0, &e) == 1);
|
||||
CHECK(e.port == 22 && e.inode == 12345 && e.proto == PROTO_TCP);
|
||||
// Established connection is skipped.
|
||||
const char *est = " 1: 0100007F:0016 0100007F:D431 01 00000000:00000000 00:00000000 00000000 0 0 222 1 0000000000000000 100 0 0 10 0";
|
||||
CHECK(services_parse_line(est, 0, 0, &e) == 0);
|
||||
// :: port 53 UDP
|
||||
const char *udp6 = " 10: 00000000000000000000000000000000:0035 00000000000000000000000000000000:0000 07 00000000:00000000 00:00000000 00000000 0 0 999 2 0000000000000000 0";
|
||||
CHECK(services_parse_line(udp6, 1, 1, &e) == 1);
|
||||
CHECK(e.port == 53 && !strcmp(e.addr, "::") && e.proto == PROTO_UDP);
|
||||
CHECK(services_parse_line("bogus", 0, 0, &e) == -1);
|
||||
}
|
||||
|
||||
static void test_json(void)
|
||||
{
|
||||
struct buf b;
|
||||
buf_init(&b);
|
||||
json_str(&b, "a\"b\\c\n\x01");
|
||||
buf_append(&b, "", 1);
|
||||
CHECK(!strcmp(b.data, "\"a\\\"b\\\\c\\n\\u0001\""));
|
||||
buf_free(&b);
|
||||
}
|
||||
|
||||
static void test_token(void)
|
||||
{
|
||||
char a[PB_TOKEN_LEN + 1], b[PB_TOKEN_LEN + 1];
|
||||
CHECK(random_token(a, sizeof(a)) == 0 && strlen(a) == PB_TOKEN_LEN);
|
||||
CHECK(random_token(b, sizeof(b)) == 0 && strcmp(a, b));
|
||||
}
|
||||
|
||||
/* Runner */
|
||||
|
||||
static const struct { const char *name; void (*fn)(void); } tests[] = {
|
||||
{ "config", test_config },
|
||||
{ "lines", test_lines },
|
||||
{ "udp_frames", test_udp_frames },
|
||||
{ "services_parse", test_services_parse },
|
||||
{ "json", test_json },
|
||||
{ "token", test_token },
|
||||
};
|
||||
|
||||
int main(int argc, char **argv)
|
||||
{
|
||||
int ran = 0;
|
||||
for (size_t i = 0; i < sizeof(tests) / sizeof(tests[0]); i++) {
|
||||
int want = argc < 2;
|
||||
for (int a = 1; a < argc; a++)
|
||||
want |= !strcmp(argv[a], tests[i].name);
|
||||
if (!want)
|
||||
continue;
|
||||
int before = failures;
|
||||
tests[i].fn();
|
||||
printf("%s %s\n", failures == before ? "ok " : "FAIL", tests[i].name);
|
||||
ran++;
|
||||
}
|
||||
if (!ran) {
|
||||
fprintf(stderr, "no matching tests\n");
|
||||
return 2;
|
||||
}
|
||||
return failures ? 1 : 0;
|
||||
}
|
||||
Reference in New Issue
Block a user