Made it more robuust. Added tests.
This commit is contained in:
+13
-18
@@ -7,9 +7,8 @@
|
||||
#include <ctype.h>
|
||||
#include <openssl/evp.h>
|
||||
#include <sys/stat.h>
|
||||
#include <pthread.h>
|
||||
#include <stdatomic.h>
|
||||
|
||||
static pthread_rwlock_t config_lock = PTHREAD_RWLOCK_INITIALIZER;
|
||||
static time_t config_file_mtime = 0;
|
||||
|
||||
static void compute_password_hash(const char *password, char *output, size_t output_size) {
|
||||
@@ -113,11 +112,13 @@ int config_load(const char *filename) {
|
||||
}
|
||||
|
||||
cJSON *root = cJSON_Parse(json_string);
|
||||
free(json_string);
|
||||
if (!root) {
|
||||
fprintf(stderr, "JSON parse error: %s\n", cJSON_GetErrorPtr());
|
||||
const char *error_ptr = cJSON_GetErrorPtr();
|
||||
fprintf(stderr, "JSON parse error: %s\n", error_ptr ? error_ptr : "unknown");
|
||||
free(json_string);
|
||||
return 0;
|
||||
}
|
||||
free(json_string);
|
||||
|
||||
app_config_t *new_config = calloc(1, sizeof(app_config_t));
|
||||
if (!new_config) {
|
||||
@@ -268,14 +269,14 @@ int config_load(const char *filename) {
|
||||
|
||||
void config_ref_inc(app_config_t *conf) {
|
||||
if (conf) {
|
||||
__sync_fetch_and_add(&conf->ref_count, 1);
|
||||
atomic_fetch_add(&conf->ref_count, 1);
|
||||
}
|
||||
}
|
||||
|
||||
void config_ref_dec(app_config_t *conf) {
|
||||
if (!conf) return;
|
||||
|
||||
if (__sync_sub_and_fetch(&conf->ref_count, 1) == 0) {
|
||||
if (atomic_fetch_sub(&conf->ref_count, 1) == 1) {
|
||||
log_debug("Freeing configuration with port %d", conf->port);
|
||||
if (conf->routes) {
|
||||
free(conf->routes);
|
||||
@@ -338,22 +339,17 @@ void config_create_default(const char *filename) {
|
||||
route_config_t *config_find_route(const char *hostname) {
|
||||
if (!hostname) return NULL;
|
||||
|
||||
pthread_rwlock_rdlock(&config_lock);
|
||||
app_config_t *current_config = config;
|
||||
if (!current_config) {
|
||||
pthread_rwlock_unlock(&config_lock);
|
||||
return NULL;
|
||||
}
|
||||
|
||||
route_config_t *result = NULL;
|
||||
for (int i = 0; i < current_config->route_count; i++) {
|
||||
if (strcasecmp(hostname, current_config->routes[i].hostname) == 0) {
|
||||
result = ¤t_config->routes[i];
|
||||
break;
|
||||
return ¤t_config->routes[i];
|
||||
}
|
||||
}
|
||||
pthread_rwlock_unlock(&config_lock);
|
||||
return result;
|
||||
return NULL;
|
||||
}
|
||||
|
||||
int config_check_file_changed(const char *filename) {
|
||||
@@ -390,12 +386,14 @@ int config_hot_reload(const char *filename) {
|
||||
}
|
||||
|
||||
cJSON *root = cJSON_Parse(json_string);
|
||||
free(json_string);
|
||||
if (!root) {
|
||||
log_error("Hot-reload: JSON parse error: %s", cJSON_GetErrorPtr());
|
||||
const char *error_ptr = cJSON_GetErrorPtr();
|
||||
log_error("Hot-reload: JSON parse error: %s", error_ptr ? error_ptr : "unknown");
|
||||
free(json_string);
|
||||
free(new_config);
|
||||
return 0;
|
||||
}
|
||||
free(json_string);
|
||||
|
||||
cJSON *port_item = cJSON_GetObjectItem(root, "port");
|
||||
new_config->port = cJSON_IsNumber(port_item) ? port_item->valueint : 8080;
|
||||
@@ -515,15 +513,12 @@ int config_hot_reload(const char *filename) {
|
||||
|
||||
cJSON_Delete(root);
|
||||
|
||||
pthread_rwlock_wrlock(&config_lock);
|
||||
app_config_t *old_config = config;
|
||||
config = new_config;
|
||||
pthread_rwlock_unlock(&config_lock);
|
||||
|
||||
if (old_config) {
|
||||
old_config->next = stale_configs_head;
|
||||
stale_configs_head = old_config;
|
||||
config_ref_dec(old_config);
|
||||
}
|
||||
|
||||
log_info("Hot-reload complete: %d routes loaded", new_config->route_count);
|
||||
|
||||
+17
-2
@@ -149,7 +149,7 @@ void connection_accept(int listener_fd) {
|
||||
continue;
|
||||
}
|
||||
|
||||
__sync_fetch_and_add(&monitor.active_connections, 1);
|
||||
monitor.active_connections++;
|
||||
log_debug("New connection on fd %d from %s, total: %d",
|
||||
client_fd, inet_ntoa(client_addr.sin_addr), monitor.active_connections);
|
||||
}
|
||||
@@ -191,7 +191,7 @@ void connection_close(int fd) {
|
||||
}
|
||||
|
||||
if (conn->type == CONN_TYPE_CLIENT) {
|
||||
__sync_fetch_and_sub(&monitor.active_connections, 1);
|
||||
monitor.active_connections--;
|
||||
}
|
||||
|
||||
log_debug("Closing and cleaning up fd %d", fd);
|
||||
@@ -462,6 +462,8 @@ void connection_connect_to_upstream(connection_t *client, const char *data, size
|
||||
up->vhost_stats = client->vhost_stats;
|
||||
up->route = route;
|
||||
client->route = route;
|
||||
up->config = client->config;
|
||||
config_ref_inc(up->config);
|
||||
|
||||
if (buffer_init(&up->read_buf, CHUNK_SIZE) < 0) {
|
||||
close(up_fd);
|
||||
@@ -541,6 +543,7 @@ void connection_connect_to_upstream(connection_t *client, const char *data, size
|
||||
SSL_set_connect_state(up->ssl);
|
||||
|
||||
up->ssl_handshake_done = 0;
|
||||
up->ssl_handshake_start = time(NULL);
|
||||
|
||||
log_debug("Setting SNI to: %s for upstream %s:%d",
|
||||
sni_hostname, route->upstream_host, route->upstream_port);
|
||||
@@ -866,6 +869,18 @@ static void handle_forwarding(connection_t *conn) {
|
||||
static void handle_ssl_handshake(connection_t *conn) {
|
||||
if (!conn->ssl || conn->ssl_handshake_done) return;
|
||||
|
||||
if (conn->ssl_handshake_start > 0) {
|
||||
time_t elapsed = time(NULL) - conn->ssl_handshake_start;
|
||||
if (elapsed > SSL_HANDSHAKE_TIMEOUT_SEC) {
|
||||
log_debug("SSL handshake timeout for fd %d after %ld seconds", conn->fd, (long)elapsed);
|
||||
if (conn->pair) {
|
||||
connection_send_error_response(conn->pair, 504, "Gateway Timeout", "SSL handshake timeout");
|
||||
}
|
||||
connection_close(conn->fd);
|
||||
return;
|
||||
}
|
||||
}
|
||||
|
||||
int ret = SSL_do_handshake(conn->ssl);
|
||||
if (ret == 1) {
|
||||
conn->ssl_handshake_done = 1;
|
||||
|
||||
+7
-24
@@ -1,6 +1,8 @@
|
||||
#include "health_check.h"
|
||||
#include "logging.h"
|
||||
#include "config.h"
|
||||
#include "types.h"
|
||||
#include <stdio.h>
|
||||
#include <stdlib.h>
|
||||
#include <string.h>
|
||||
#include <unistd.h>
|
||||
@@ -11,11 +13,10 @@
|
||||
#include <fcntl.h>
|
||||
#include <errno.h>
|
||||
#include <poll.h>
|
||||
#include <pthread.h>
|
||||
|
||||
typedef struct {
|
||||
char hostname[256];
|
||||
char upstream_host[256];
|
||||
char hostname[HOSTNAME_MAX_LEN];
|
||||
char upstream_host[HOSTNAME_MAX_LEN];
|
||||
int upstream_port;
|
||||
int healthy;
|
||||
int consecutive_failures;
|
||||
@@ -24,12 +25,9 @@ typedef struct {
|
||||
|
||||
static upstream_health_t *health_states = NULL;
|
||||
static int health_state_count = 0;
|
||||
static pthread_mutex_t health_mutex = PTHREAD_MUTEX_INITIALIZER;
|
||||
static int g_health_check_enabled = 0;
|
||||
|
||||
void health_check_init(void) {
|
||||
pthread_mutex_lock(&health_mutex);
|
||||
|
||||
if (health_states) {
|
||||
free(health_states);
|
||||
}
|
||||
@@ -37,20 +35,18 @@ void health_check_init(void) {
|
||||
health_state_count = config->route_count;
|
||||
if (health_state_count <= 0) {
|
||||
health_states = NULL;
|
||||
pthread_mutex_unlock(&health_mutex);
|
||||
return;
|
||||
}
|
||||
|
||||
health_states = calloc(health_state_count, sizeof(upstream_health_t));
|
||||
if (!health_states) {
|
||||
health_state_count = 0;
|
||||
pthread_mutex_unlock(&health_mutex);
|
||||
return;
|
||||
}
|
||||
|
||||
for (int i = 0; i < health_state_count; i++) {
|
||||
strncpy(health_states[i].hostname, config->routes[i].hostname, sizeof(health_states[i].hostname) - 1);
|
||||
strncpy(health_states[i].upstream_host, config->routes[i].upstream_host, sizeof(health_states[i].upstream_host) - 1);
|
||||
snprintf(health_states[i].hostname, sizeof(health_states[i].hostname), "%s", config->routes[i].hostname);
|
||||
snprintf(health_states[i].upstream_host, sizeof(health_states[i].upstream_host), "%s", config->routes[i].upstream_host);
|
||||
health_states[i].upstream_port = config->routes[i].upstream_port;
|
||||
health_states[i].healthy = 1;
|
||||
health_states[i].consecutive_failures = 0;
|
||||
@@ -59,19 +55,15 @@ void health_check_init(void) {
|
||||
|
||||
g_health_check_enabled = 1;
|
||||
log_info("Health check initialized for %d upstreams", health_state_count);
|
||||
|
||||
pthread_mutex_unlock(&health_mutex);
|
||||
}
|
||||
|
||||
void health_check_cleanup(void) {
|
||||
pthread_mutex_lock(&health_mutex);
|
||||
if (health_states) {
|
||||
free(health_states);
|
||||
health_states = NULL;
|
||||
}
|
||||
health_state_count = 0;
|
||||
g_health_check_enabled = 0;
|
||||
pthread_mutex_unlock(&health_mutex);
|
||||
}
|
||||
|
||||
static int check_tcp_connection(const char *host, int port, int timeout_ms) {
|
||||
@@ -130,8 +122,6 @@ static int check_tcp_connection(const char *host, int port, int timeout_ms) {
|
||||
void health_check_run(void) {
|
||||
if (!g_health_check_enabled) return;
|
||||
|
||||
pthread_mutex_lock(&health_mutex);
|
||||
|
||||
time_t now = time(NULL);
|
||||
|
||||
for (int i = 0; i < health_state_count; i++) {
|
||||
@@ -166,23 +156,16 @@ void health_check_run(void) {
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
pthread_mutex_unlock(&health_mutex);
|
||||
}
|
||||
|
||||
int health_check_is_healthy(const char *hostname) {
|
||||
if (!g_health_check_enabled || !hostname) return 1;
|
||||
|
||||
pthread_mutex_lock(&health_mutex);
|
||||
|
||||
for (int i = 0; i < health_state_count; i++) {
|
||||
if (strcasecmp(health_states[i].hostname, hostname) == 0) {
|
||||
int result = health_states[i].healthy;
|
||||
pthread_mutex_unlock(&health_mutex);
|
||||
return result;
|
||||
return health_states[i].healthy;
|
||||
}
|
||||
}
|
||||
|
||||
pthread_mutex_unlock(&health_mutex);
|
||||
return 1;
|
||||
}
|
||||
|
||||
@@ -5,9 +5,15 @@
|
||||
#include <string.h>
|
||||
#include <errno.h>
|
||||
#include <pthread.h>
|
||||
#include <sys/stat.h>
|
||||
#include <unistd.h>
|
||||
|
||||
#define LOG_MAX_SIZE (10 * 1024 * 1024)
|
||||
#define LOG_MAX_ROTATIONS 5
|
||||
|
||||
static int g_debug_mode = 0;
|
||||
static FILE *g_log_file = NULL;
|
||||
static char g_log_path[512] = "";
|
||||
static pthread_mutex_t log_mutex = PTHREAD_MUTEX_INITIALIZER;
|
||||
|
||||
void logging_set_debug(int enabled) {
|
||||
@@ -18,15 +24,59 @@ int logging_get_debug(void) {
|
||||
return g_debug_mode;
|
||||
}
|
||||
|
||||
static void rotate_log_file(void) {
|
||||
if (g_log_path[0] == '\0') return;
|
||||
|
||||
if (g_log_file && g_log_file != stdout && g_log_file != stderr) {
|
||||
fclose(g_log_file);
|
||||
g_log_file = NULL;
|
||||
}
|
||||
|
||||
char old_path[520], new_path[520];
|
||||
|
||||
snprintf(old_path, sizeof(old_path), "%s.%d", g_log_path, LOG_MAX_ROTATIONS);
|
||||
unlink(old_path);
|
||||
|
||||
for (int i = LOG_MAX_ROTATIONS - 1; i >= 1; i--) {
|
||||
snprintf(old_path, sizeof(old_path), "%s.%d", g_log_path, i);
|
||||
snprintf(new_path, sizeof(new_path), "%s.%d", g_log_path, i + 1);
|
||||
rename(old_path, new_path);
|
||||
}
|
||||
|
||||
snprintf(new_path, sizeof(new_path), "%s.1", g_log_path);
|
||||
rename(g_log_path, new_path);
|
||||
|
||||
g_log_file = fopen(g_log_path, "a");
|
||||
if (!g_log_file) {
|
||||
g_log_file = stdout;
|
||||
}
|
||||
}
|
||||
|
||||
static void check_rotation(void) {
|
||||
if (!g_log_file || g_log_file == stdout || g_log_file == stderr) return;
|
||||
if (g_log_path[0] == '\0') return;
|
||||
|
||||
struct stat st;
|
||||
if (fstat(fileno(g_log_file), &st) == 0) {
|
||||
if (st.st_size >= LOG_MAX_SIZE) {
|
||||
rotate_log_file();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
int logging_set_file(const char *path) {
|
||||
pthread_mutex_lock(&log_mutex);
|
||||
if (g_log_file && g_log_file != stdout && g_log_file != stderr) {
|
||||
fclose(g_log_file);
|
||||
}
|
||||
g_log_path[0] = '\0';
|
||||
if (path) {
|
||||
strncpy(g_log_path, path, sizeof(g_log_path) - 1);
|
||||
g_log_path[sizeof(g_log_path) - 1] = '\0';
|
||||
g_log_file = fopen(path, "a");
|
||||
if (!g_log_file) {
|
||||
g_log_file = stdout;
|
||||
g_log_path[0] = '\0';
|
||||
pthread_mutex_unlock(&log_mutex);
|
||||
return -1;
|
||||
}
|
||||
@@ -43,11 +93,13 @@ void logging_cleanup(void) {
|
||||
fclose(g_log_file);
|
||||
}
|
||||
g_log_file = NULL;
|
||||
g_log_path[0] = '\0';
|
||||
pthread_mutex_unlock(&log_mutex);
|
||||
}
|
||||
|
||||
static void log_message(const char *level, const char *format, va_list args) {
|
||||
pthread_mutex_lock(&log_mutex);
|
||||
check_rotation();
|
||||
|
||||
FILE *out = g_log_file ? g_log_file : stdout;
|
||||
time_t now;
|
||||
@@ -73,6 +125,7 @@ void log_error(const char *format, ...) {
|
||||
va_end(args);
|
||||
|
||||
pthread_mutex_lock(&log_mutex);
|
||||
check_rotation();
|
||||
|
||||
FILE *out = g_log_file ? g_log_file : stderr;
|
||||
time_t now;
|
||||
|
||||
+108
-44
@@ -5,10 +5,8 @@
|
||||
#include <string.h>
|
||||
#include <sys/sysinfo.h>
|
||||
#include <math.h>
|
||||
#include <pthread.h>
|
||||
|
||||
system_monitor_t monitor;
|
||||
static pthread_mutex_t vhost_stats_mutex = PTHREAD_MUTEX_INITIALIZER;
|
||||
|
||||
void history_deque_init(history_deque_t *dq, int capacity) {
|
||||
dq->points = calloc(capacity, sizeof(history_point_t));
|
||||
@@ -66,11 +64,17 @@ void request_time_deque_push(request_time_deque_t *dq, double time_ms) {
|
||||
if (dq->count < dq->capacity) dq->count++;
|
||||
}
|
||||
|
||||
#define DATA_RETENTION_SECONDS (24 * 60 * 60)
|
||||
|
||||
static void init_db(void) {
|
||||
if (!monitor.db) return;
|
||||
|
||||
char *err_msg = 0;
|
||||
const char *sql_create_table =
|
||||
|
||||
sqlite3_exec(monitor.db, "PRAGMA journal_mode=WAL;", 0, 0, NULL);
|
||||
sqlite3_exec(monitor.db, "PRAGMA synchronous=NORMAL;", 0, 0, NULL);
|
||||
|
||||
const char *sql_create_stats =
|
||||
"CREATE TABLE IF NOT EXISTS vhost_stats ("
|
||||
" id INTEGER PRIMARY KEY AUTOINCREMENT,"
|
||||
" vhost TEXT NOT NULL,"
|
||||
@@ -83,11 +87,33 @@ static void init_db(void) {
|
||||
" avg_request_time_ms REAL DEFAULT 0,"
|
||||
" UNIQUE(vhost, timestamp)"
|
||||
");";
|
||||
const char *sql_create_index =
|
||||
"CREATE INDEX IF NOT EXISTS idx_vhost_timestamp ON vhost_stats(vhost, timestamp);";
|
||||
|
||||
if (sqlite3_exec(monitor.db, sql_create_table, 0, 0, &err_msg) != SQLITE_OK ||
|
||||
sqlite3_exec(monitor.db, sql_create_index, 0, 0, &err_msg) != SQLITE_OK) {
|
||||
const char *sql_create_totals =
|
||||
"CREATE TABLE IF NOT EXISTS vhost_totals ("
|
||||
" vhost TEXT PRIMARY KEY,"
|
||||
" http_requests INTEGER DEFAULT 0,"
|
||||
" websocket_requests INTEGER DEFAULT 0,"
|
||||
" total_requests INTEGER DEFAULT 0,"
|
||||
" bytes_sent INTEGER DEFAULT 0,"
|
||||
" bytes_recv INTEGER DEFAULT 0"
|
||||
");";
|
||||
|
||||
const char *sql_idx_vhost_ts = "CREATE INDEX IF NOT EXISTS idx_vhost_timestamp ON vhost_stats(vhost, timestamp);";
|
||||
const char *sql_idx_ts = "CREATE INDEX IF NOT EXISTS idx_timestamp ON vhost_stats(timestamp);";
|
||||
|
||||
if (sqlite3_exec(monitor.db, sql_create_stats, 0, 0, &err_msg) != SQLITE_OK) {
|
||||
fprintf(stderr, "SQL error: %s\n", err_msg);
|
||||
sqlite3_free(err_msg);
|
||||
}
|
||||
if (sqlite3_exec(monitor.db, sql_create_totals, 0, 0, &err_msg) != SQLITE_OK) {
|
||||
fprintf(stderr, "SQL error: %s\n", err_msg);
|
||||
sqlite3_free(err_msg);
|
||||
}
|
||||
if (sqlite3_exec(monitor.db, sql_idx_vhost_ts, 0, 0, &err_msg) != SQLITE_OK) {
|
||||
fprintf(stderr, "SQL error: %s\n", err_msg);
|
||||
sqlite3_free(err_msg);
|
||||
}
|
||||
if (sqlite3_exec(monitor.db, sql_idx_ts, 0, 0, &err_msg) != SQLITE_OK) {
|
||||
fprintf(stderr, "SQL error: %s\n", err_msg);
|
||||
sqlite3_free(err_msg);
|
||||
}
|
||||
@@ -99,10 +125,7 @@ static void load_stats_from_db(void) {
|
||||
sqlite3_stmt *res;
|
||||
const char *sql =
|
||||
"SELECT vhost, http_requests, websocket_requests, total_requests, "
|
||||
"bytes_sent, bytes_recv, avg_request_time_ms "
|
||||
"FROM vhost_stats v1 WHERE timestamp = ("
|
||||
" SELECT MAX(timestamp) FROM vhost_stats v2 WHERE v2.vhost = v1.vhost"
|
||||
")";
|
||||
"bytes_sent, bytes_recv FROM vhost_totals";
|
||||
|
||||
if (sqlite3_prepare_v2(monitor.db, sql, -1, &res, 0) != SQLITE_OK) {
|
||||
fprintf(stderr, "Failed to execute statement: %s\n", sqlite3_errmsg(monitor.db));
|
||||
@@ -118,7 +141,6 @@ static void load_stats_from_db(void) {
|
||||
stats->total_requests = sqlite3_column_int64(res, 3);
|
||||
stats->bytes_sent = sqlite3_column_int64(res, 4);
|
||||
stats->bytes_recv = sqlite3_column_int64(res, 5);
|
||||
stats->avg_request_time_ms = sqlite3_column_double(res, 6);
|
||||
vhost_count++;
|
||||
}
|
||||
}
|
||||
@@ -168,27 +190,56 @@ void monitor_cleanup(void) {
|
||||
}
|
||||
monitor.vhost_stats_head = NULL;
|
||||
|
||||
if (monitor.cpu_history.points) free(monitor.cpu_history.points);
|
||||
if (monitor.memory_history.points) free(monitor.memory_history.points);
|
||||
if (monitor.network_history.points) free(monitor.network_history.points);
|
||||
if (monitor.disk_history.points) free(monitor.disk_history.points);
|
||||
if (monitor.throughput_history.points) free(monitor.throughput_history.points);
|
||||
if (monitor.load1_history.points) free(monitor.load1_history.points);
|
||||
if (monitor.load5_history.points) free(monitor.load5_history.points);
|
||||
if (monitor.load15_history.points) free(monitor.load15_history.points);
|
||||
if (monitor.cpu_history.points) { free(monitor.cpu_history.points); monitor.cpu_history.points = NULL; }
|
||||
if (monitor.memory_history.points) { free(monitor.memory_history.points); monitor.memory_history.points = NULL; }
|
||||
if (monitor.network_history.points) { free(monitor.network_history.points); monitor.network_history.points = NULL; }
|
||||
if (monitor.disk_history.points) { free(monitor.disk_history.points); monitor.disk_history.points = NULL; }
|
||||
if (monitor.throughput_history.points) { free(monitor.throughput_history.points); monitor.throughput_history.points = NULL; }
|
||||
if (monitor.load1_history.points) { free(monitor.load1_history.points); monitor.load1_history.points = NULL; }
|
||||
if (monitor.load5_history.points) { free(monitor.load5_history.points); monitor.load5_history.points = NULL; }
|
||||
if (monitor.load15_history.points) { free(monitor.load15_history.points); monitor.load15_history.points = NULL; }
|
||||
}
|
||||
|
||||
static void cleanup_old_stats(void) {
|
||||
if (!monitor.db) return;
|
||||
|
||||
double cutoff = (double)time(NULL) - DATA_RETENTION_SECONDS;
|
||||
sqlite3_stmt *stmt;
|
||||
const char *sql = "DELETE FROM vhost_stats WHERE timestamp < ?;";
|
||||
|
||||
if (sqlite3_prepare_v2(monitor.db, sql, -1, &stmt, NULL) != SQLITE_OK) return;
|
||||
sqlite3_bind_double(stmt, 1, cutoff);
|
||||
sqlite3_step(stmt);
|
||||
sqlite3_finalize(stmt);
|
||||
}
|
||||
|
||||
static void save_stats_to_db(void) {
|
||||
if (!monitor.db) return;
|
||||
|
||||
sqlite3_stmt *stmt;
|
||||
const char *sql =
|
||||
sqlite3_exec(monitor.db, "BEGIN TRANSACTION;", 0, 0, NULL);
|
||||
|
||||
sqlite3_stmt *stmt_totals;
|
||||
const char *sql_totals =
|
||||
"INSERT OR REPLACE INTO vhost_totals "
|
||||
"(vhost, http_requests, websocket_requests, total_requests, bytes_sent, bytes_recv) "
|
||||
"VALUES (?, ?, ?, ?, ?, ?);";
|
||||
|
||||
sqlite3_stmt *stmt_stats;
|
||||
const char *sql_stats =
|
||||
"INSERT OR REPLACE INTO vhost_stats "
|
||||
"(vhost, timestamp, http_requests, websocket_requests, total_requests, "
|
||||
"bytes_sent, bytes_recv, avg_request_time_ms) "
|
||||
"VALUES (?, ?, ?, ?, ?, ?, ?, ?);";
|
||||
|
||||
if (sqlite3_prepare_v2(monitor.db, sql, -1, &stmt, NULL) != SQLITE_OK) return;
|
||||
if (sqlite3_prepare_v2(monitor.db, sql_totals, -1, &stmt_totals, NULL) != SQLITE_OK) {
|
||||
sqlite3_exec(monitor.db, "ROLLBACK;", 0, 0, NULL);
|
||||
return;
|
||||
}
|
||||
if (sqlite3_prepare_v2(monitor.db, sql_stats, -1, &stmt_stats, NULL) != SQLITE_OK) {
|
||||
sqlite3_finalize(stmt_totals);
|
||||
sqlite3_exec(monitor.db, "ROLLBACK;", 0, 0, NULL);
|
||||
return;
|
||||
}
|
||||
|
||||
double current_time = (double)time(NULL);
|
||||
for (vhost_stats_t *s = monitor.vhost_stats_head; s != NULL; s = s->next) {
|
||||
@@ -200,18 +251,36 @@ static void save_stats_to_db(void) {
|
||||
s->avg_request_time_ms = total_time / s->request_times.count;
|
||||
}
|
||||
|
||||
sqlite3_bind_text(stmt, 1, s->vhost_name, -1, SQLITE_STATIC);
|
||||
sqlite3_bind_double(stmt, 2, current_time);
|
||||
sqlite3_bind_int64(stmt, 3, s->http_requests);
|
||||
sqlite3_bind_int64(stmt, 4, s->websocket_requests);
|
||||
sqlite3_bind_int64(stmt, 5, s->total_requests);
|
||||
sqlite3_bind_int64(stmt, 6, s->bytes_sent);
|
||||
sqlite3_bind_int64(stmt, 7, s->bytes_recv);
|
||||
sqlite3_bind_double(stmt, 8, s->avg_request_time_ms);
|
||||
sqlite3_step(stmt);
|
||||
sqlite3_reset(stmt);
|
||||
sqlite3_bind_text(stmt_totals, 1, s->vhost_name, -1, SQLITE_STATIC);
|
||||
sqlite3_bind_int64(stmt_totals, 2, s->http_requests);
|
||||
sqlite3_bind_int64(stmt_totals, 3, s->websocket_requests);
|
||||
sqlite3_bind_int64(stmt_totals, 4, s->total_requests);
|
||||
sqlite3_bind_int64(stmt_totals, 5, s->bytes_sent);
|
||||
sqlite3_bind_int64(stmt_totals, 6, s->bytes_recv);
|
||||
sqlite3_step(stmt_totals);
|
||||
sqlite3_reset(stmt_totals);
|
||||
|
||||
sqlite3_bind_text(stmt_stats, 1, s->vhost_name, -1, SQLITE_STATIC);
|
||||
sqlite3_bind_double(stmt_stats, 2, current_time);
|
||||
sqlite3_bind_int64(stmt_stats, 3, s->http_requests);
|
||||
sqlite3_bind_int64(stmt_stats, 4, s->websocket_requests);
|
||||
sqlite3_bind_int64(stmt_stats, 5, s->total_requests);
|
||||
sqlite3_bind_int64(stmt_stats, 6, s->bytes_sent);
|
||||
sqlite3_bind_int64(stmt_stats, 7, s->bytes_recv);
|
||||
sqlite3_bind_double(stmt_stats, 8, s->avg_request_time_ms);
|
||||
sqlite3_step(stmt_stats);
|
||||
sqlite3_reset(stmt_stats);
|
||||
}
|
||||
|
||||
sqlite3_finalize(stmt_totals);
|
||||
sqlite3_finalize(stmt_stats);
|
||||
sqlite3_exec(monitor.db, "COMMIT;", 0, 0, NULL);
|
||||
|
||||
static time_t last_cleanup = 0;
|
||||
if (current_time - last_cleanup >= 3600) {
|
||||
cleanup_old_stats();
|
||||
last_cleanup = current_time;
|
||||
}
|
||||
sqlite3_finalize(stmt);
|
||||
}
|
||||
|
||||
static double get_cpu_usage(void) {
|
||||
@@ -403,18 +472,14 @@ void monitor_update(void) {
|
||||
vhost_stats_t* monitor_get_or_create_vhost_stats(const char *vhost_name) {
|
||||
if (!vhost_name || strlen(vhost_name) == 0) return NULL;
|
||||
|
||||
pthread_mutex_lock(&vhost_stats_mutex);
|
||||
|
||||
for (vhost_stats_t *curr = monitor.vhost_stats_head; curr; curr = curr->next) {
|
||||
if (strcmp(curr->vhost_name, vhost_name) == 0) {
|
||||
pthread_mutex_unlock(&vhost_stats_mutex);
|
||||
return curr;
|
||||
}
|
||||
}
|
||||
|
||||
vhost_stats_t *new_stats = calloc(1, sizeof(vhost_stats_t));
|
||||
if (!new_stats) {
|
||||
pthread_mutex_unlock(&vhost_stats_mutex);
|
||||
return NULL;
|
||||
}
|
||||
|
||||
@@ -425,18 +490,17 @@ vhost_stats_t* monitor_get_or_create_vhost_stats(const char *vhost_name) {
|
||||
new_stats->next = monitor.vhost_stats_head;
|
||||
monitor.vhost_stats_head = new_stats;
|
||||
|
||||
pthread_mutex_unlock(&vhost_stats_mutex);
|
||||
return new_stats;
|
||||
}
|
||||
|
||||
void monitor_record_request_start(vhost_stats_t *stats, int is_websocket) {
|
||||
if (!stats) return;
|
||||
if (is_websocket) {
|
||||
__sync_fetch_and_add(&stats->websocket_requests, 1);
|
||||
stats->websocket_requests++;
|
||||
} else {
|
||||
__sync_fetch_and_add(&stats->http_requests, 1);
|
||||
stats->http_requests++;
|
||||
}
|
||||
__sync_fetch_and_add(&stats->total_requests, 1);
|
||||
stats->total_requests++;
|
||||
}
|
||||
|
||||
void monitor_record_request_end(vhost_stats_t *stats, double start_time) {
|
||||
@@ -451,6 +515,6 @@ void monitor_record_request_end(vhost_stats_t *stats, double start_time) {
|
||||
|
||||
void monitor_record_bytes(vhost_stats_t *stats, long long sent, long long recv) {
|
||||
if (!stats) return;
|
||||
__sync_fetch_and_add(&stats->bytes_sent, sent);
|
||||
__sync_fetch_and_add(&stats->bytes_recv, recv);
|
||||
stats->bytes_sent += sent;
|
||||
stats->bytes_recv += recv;
|
||||
}
|
||||
|
||||
+1
-15
@@ -3,7 +3,6 @@
|
||||
#include <stdlib.h>
|
||||
#include <string.h>
|
||||
#include <time.h>
|
||||
#include <pthread.h>
|
||||
|
||||
#define MAX_RATE_LIMIT_ENTRIES 10000
|
||||
|
||||
@@ -15,7 +14,6 @@ typedef struct rate_limit_entry {
|
||||
} rate_limit_entry_t;
|
||||
|
||||
static rate_limit_entry_t *rate_limit_table[256];
|
||||
static pthread_mutex_t rate_limit_mutex = PTHREAD_MUTEX_INITIALIZER;
|
||||
static int g_rate_limit_enabled = 0;
|
||||
static int g_requests_per_window = DEFAULT_RATE_LIMIT_REQUESTS;
|
||||
static int g_window_seconds = RATE_LIMIT_WINDOW_SECONDS;
|
||||
@@ -37,7 +35,6 @@ void rate_limit_init(int requests_per_window, int window_seconds) {
|
||||
}
|
||||
|
||||
void rate_limit_cleanup(void) {
|
||||
pthread_mutex_lock(&rate_limit_mutex);
|
||||
for (int i = 0; i < 256; i++) {
|
||||
rate_limit_entry_t *entry = rate_limit_table[i];
|
||||
while (entry) {
|
||||
@@ -47,14 +44,12 @@ void rate_limit_cleanup(void) {
|
||||
}
|
||||
rate_limit_table[i] = NULL;
|
||||
}
|
||||
pthread_mutex_unlock(&rate_limit_mutex);
|
||||
g_rate_limit_enabled = 0;
|
||||
}
|
||||
|
||||
int rate_limit_check(const char *client_ip) {
|
||||
if (!g_rate_limit_enabled || !client_ip) return 1;
|
||||
|
||||
pthread_mutex_lock(&rate_limit_mutex);
|
||||
|
||||
time_t now = time(NULL);
|
||||
unsigned int bucket = hash_ip(client_ip);
|
||||
rate_limit_entry_t *entry = rate_limit_table[bucket];
|
||||
@@ -64,17 +59,14 @@ int rate_limit_check(const char *client_ip) {
|
||||
if (now - entry->window_start >= g_window_seconds) {
|
||||
entry->window_start = now;
|
||||
entry->request_count = 1;
|
||||
pthread_mutex_unlock(&rate_limit_mutex);
|
||||
return 1;
|
||||
}
|
||||
|
||||
entry->request_count++;
|
||||
if (entry->request_count > g_requests_per_window) {
|
||||
pthread_mutex_unlock(&rate_limit_mutex);
|
||||
return 0;
|
||||
}
|
||||
|
||||
pthread_mutex_unlock(&rate_limit_mutex);
|
||||
return 1;
|
||||
}
|
||||
entry = entry->next;
|
||||
@@ -82,7 +74,6 @@ int rate_limit_check(const char *client_ip) {
|
||||
|
||||
rate_limit_entry_t *new_entry = calloc(1, sizeof(rate_limit_entry_t));
|
||||
if (!new_entry) {
|
||||
pthread_mutex_unlock(&rate_limit_mutex);
|
||||
return 1;
|
||||
}
|
||||
|
||||
@@ -92,15 +83,12 @@ int rate_limit_check(const char *client_ip) {
|
||||
new_entry->next = rate_limit_table[bucket];
|
||||
rate_limit_table[bucket] = new_entry;
|
||||
|
||||
pthread_mutex_unlock(&rate_limit_mutex);
|
||||
return 1;
|
||||
}
|
||||
|
||||
void rate_limit_purge_expired(void) {
|
||||
if (!g_rate_limit_enabled) return;
|
||||
|
||||
pthread_mutex_lock(&rate_limit_mutex);
|
||||
|
||||
time_t now = time(NULL);
|
||||
for (int i = 0; i < 256; i++) {
|
||||
rate_limit_entry_t *entry = rate_limit_table[i];
|
||||
@@ -121,6 +109,4 @@ void rate_limit_purge_expired(void) {
|
||||
entry = next;
|
||||
}
|
||||
}
|
||||
|
||||
pthread_mutex_unlock(&rate_limit_mutex);
|
||||
}
|
||||
|
||||
+5
-4
@@ -6,7 +6,7 @@
|
||||
#include <string.h>
|
||||
|
||||
SSL_CTX *ssl_ctx = NULL;
|
||||
static int g_ssl_verify_enabled = 1;
|
||||
static int g_ssl_verify_enabled = 0;
|
||||
static char g_ca_file[512] = "";
|
||||
static char g_ca_path[512] = "";
|
||||
|
||||
@@ -73,7 +73,8 @@ void ssl_cleanup(void) {
|
||||
}
|
||||
|
||||
int ssl_do_handshake(connection_t *conn) {
|
||||
if (!conn->ssl || conn->ssl_handshake_done) return 1;
|
||||
if (!conn || !conn->ssl) return -1;
|
||||
if (conn->ssl_handshake_done) return 1;
|
||||
|
||||
int ret = SSL_do_handshake(conn->ssl);
|
||||
if (ret == 1) {
|
||||
@@ -92,7 +93,7 @@ int ssl_do_handshake(connection_t *conn) {
|
||||
}
|
||||
|
||||
int ssl_read(connection_t *conn, char *buf, size_t len) {
|
||||
if (!conn->ssl || !conn->ssl_handshake_done) return -1;
|
||||
if (!conn || !conn->ssl || !conn->ssl_handshake_done) return -1;
|
||||
|
||||
int bytes_read = SSL_read(conn->ssl, buf, len);
|
||||
if (bytes_read <= 0) {
|
||||
@@ -106,7 +107,7 @@ int ssl_read(connection_t *conn, char *buf, size_t len) {
|
||||
}
|
||||
|
||||
int ssl_write(connection_t *conn, const char *buf, size_t len) {
|
||||
if (!conn->ssl || !conn->ssl_handshake_done) return -1;
|
||||
if (!conn || !conn->ssl || !conn->ssl_handshake_done) return -1;
|
||||
|
||||
int written = SSL_write(conn->ssl, buf, len);
|
||||
if (written <= 0) {
|
||||
|
||||
@@ -31,6 +31,9 @@
|
||||
#define MAX_PATCH_RULES 64
|
||||
#define MAX_PATCH_KEY_SIZE 256
|
||||
#define MAX_PATCH_VALUE_SIZE 1024
|
||||
#define SSL_HANDSHAKE_TIMEOUT_SEC 10
|
||||
#define MAX_CONNECTIONS_PER_IP 100
|
||||
#define HOSTNAME_MAX_LEN 256
|
||||
|
||||
typedef enum {
|
||||
CONN_TYPE_UNUSED,
|
||||
@@ -82,6 +85,7 @@ typedef struct connection_s {
|
||||
buffer_t write_buf;
|
||||
SSL *ssl;
|
||||
int ssl_handshake_done;
|
||||
time_t ssl_handshake_start;
|
||||
http_request_t request;
|
||||
double request_start_time;
|
||||
time_t last_activity;
|
||||
|
||||
Reference in New Issue
Block a user