Merge Wave 2: thread-safety fixes (signals, fd ownership, handler epilogue, logging, scanner leak)
CI / lint (push) Successful in 1m30s
CI / sanitizers (undefined) (push) Successful in 59s
CI / sanitizers (address) (push) Successful in 1m6s
CI / fuzz-build (push) Successful in 28s
CI / coverage (push) Successful in 49s
CI / build-and-test (push) Successful in 4m31s
CI / valgrind (push) Successful in 3m10s
CI / lint (push) Successful in 1m30s
CI / sanitizers (undefined) (push) Successful in 59s
CI / sanitizers (address) (push) Successful in 1m6s
CI / fuzz-build (push) Successful in 28s
CI / coverage (push) Successful in 49s
CI / build-and-test (push) Successful in 4m31s
CI / valgrind (push) Successful in 3m10s
This commit is contained in:
@@ -1376,9 +1376,11 @@ static bool cli_handle_io_options(CliParseCtx* ctx) {
|
||||
return true;
|
||||
}
|
||||
if (config->log_file) {
|
||||
/* Detach the logger before closing: log I/O may be in flight and must
|
||||
never touch a freed FILE*. */
|
||||
log_set_file(NULL);
|
||||
fclose(config->log_file);
|
||||
config->log_file = NULL;
|
||||
log_set_file(NULL);
|
||||
}
|
||||
FILE* lf = fopen(ctx->argv[++ctx->i], "a");
|
||||
if (!lf) {
|
||||
|
||||
@@ -587,12 +587,16 @@ void directory_scanner_destroy(DirectoryScanner* scanner) {
|
||||
|
||||
static Chunk* chunk_data_to_chunk(ArrayList* chunk_data) {
|
||||
void** chunk_items = array_list_to_array(chunk_data);
|
||||
if (!chunk_items)
|
||||
if (!chunk_items) {
|
||||
array_list_delete(chunk_data);
|
||||
return NULL;
|
||||
}
|
||||
Chunk* chunk = chunk_create((File**)chunk_items, chunk_data->size);
|
||||
free(chunk_items);
|
||||
if (!chunk)
|
||||
if (!chunk) {
|
||||
array_list_delete(chunk_data);
|
||||
return NULL;
|
||||
}
|
||||
chunk_data->item_destroyer = NULL;
|
||||
array_list_delete(chunk_data);
|
||||
return chunk;
|
||||
|
||||
+68
-81
@@ -23,6 +23,7 @@
|
||||
#include <string.h>
|
||||
#include <unistd.h>
|
||||
#include <errno.h>
|
||||
#include <sys/socket.h>
|
||||
#include <sys/stat.h>
|
||||
#include <openssl/x509.h>
|
||||
|
||||
@@ -502,12 +503,16 @@ void handler(int file_descriptor) {
|
||||
gate_ctx.ssl = ssl;
|
||||
gate_ctx.fd = file_descriptor;
|
||||
gate_ctx.super_mode_override = -1;
|
||||
Config* config = config_receive_with_validate(file_descriptor, server_module_gate, &gate_ctx);
|
||||
/* All teardown state starts empty so the single `done` epilogue is safe to
|
||||
* reach from any error path (including before the config frame arrives). */
|
||||
Config* config = NULL;
|
||||
PipelineContextReceiver* context = NULL;
|
||||
char* joined_destination = NULL;
|
||||
bool charset_ready = false;
|
||||
config = config_receive_with_validate(file_descriptor, server_module_gate, &gate_ctx);
|
||||
if (config == NULL) {
|
||||
log_message(LOG_LEVEL_ERROR, "Failed to receive config");
|
||||
close(file_descriptor);
|
||||
protocol_session_unbind();
|
||||
return;
|
||||
goto done;
|
||||
}
|
||||
/* Apply the super-mode veto the gate decided on (operator --no-super, or a
|
||||
* daemon module without the `client owner = yes` opt-in) exactly once, so
|
||||
@@ -519,23 +524,15 @@ void handler(int file_descriptor) {
|
||||
protocol_set_8_bit_output(config->eight_bit_output);
|
||||
if (!authorized_root) {
|
||||
log_message(LOG_LEVEL_ERROR, "No server-side destination root configured");
|
||||
config_delete(config);
|
||||
close(file_descriptor);
|
||||
protocol_session_unbind();
|
||||
return;
|
||||
goto done;
|
||||
}
|
||||
if (!allow_unauthenticated && ssl == NULL) {
|
||||
log_message(LOG_LEVEL_ERROR, "Rejected unauthenticated plaintext connection");
|
||||
config_delete(config);
|
||||
close(file_descriptor);
|
||||
protocol_session_unbind();
|
||||
return;
|
||||
goto done;
|
||||
}
|
||||
if (ssl && required_client_cn && !tls_client_identity_allowed(ssl)) {
|
||||
log_message(LOG_LEVEL_ERROR, "Rejected TLS client with unauthorized identity");
|
||||
config_delete(config);
|
||||
close(file_descriptor);
|
||||
return;
|
||||
goto done;
|
||||
}
|
||||
/* Daemon mode: the module's root is the authorized root (installed by
|
||||
server_module_gate), and the client's destination is a MODULE-RELATIVE
|
||||
@@ -545,13 +542,9 @@ void handler(int file_descriptor) {
|
||||
if (g_daemon_conf && config->receive_root_directory && config->receive_root_directory[0] == '/') {
|
||||
log_message(LOG_LEVEL_ERROR, "Rejected absolute daemon destination (must be relative to the "
|
||||
"selected module root)");
|
||||
config_delete(config);
|
||||
close(file_descriptor);
|
||||
protocol_session_unbind();
|
||||
return;
|
||||
goto done;
|
||||
}
|
||||
char* destination = config->receive_root_directory;
|
||||
char* joined_destination = NULL;
|
||||
if (destination && destination[0] != '/')
|
||||
joined_destination = path_cat(authorized_root, destination);
|
||||
if (joined_destination)
|
||||
@@ -560,19 +553,16 @@ void handler(int file_descriptor) {
|
||||
!path_is_within(authorized_root, destination)) {
|
||||
log_message(LOG_LEVEL_ERROR, "Rejected destination outside authorized root");
|
||||
free(joined_destination);
|
||||
config_delete(config);
|
||||
close(file_descriptor);
|
||||
return;
|
||||
joined_destination = NULL;
|
||||
goto done;
|
||||
}
|
||||
if (joined_destination) {
|
||||
free(config->receive_root_directory);
|
||||
config->receive_root_directory = joined_destination;
|
||||
joined_destination = NULL;
|
||||
}
|
||||
if (!config->receive_root_directory) {
|
||||
config_delete(config);
|
||||
close(file_descriptor);
|
||||
protocol_session_unbind();
|
||||
return;
|
||||
goto done;
|
||||
}
|
||||
config->use_delete = config->use_delete && allow_delete;
|
||||
/* --iconv (protocol 2.16.0): install the receiver-side wire->local conversion
|
||||
@@ -581,13 +571,13 @@ void handler(int file_descriptor) {
|
||||
any) may override the local charset; a spec the client is known to have
|
||||
validated cannot fail here unless the server's override names an
|
||||
unsupported charset. */
|
||||
if (config->iconv_spec && !charset_wire_init_receiver(config->iconv_spec, server_iconv_spec)) {
|
||||
if (config->iconv_spec) {
|
||||
if (!charset_wire_init_receiver(config->iconv_spec, server_iconv_spec)) {
|
||||
log_message(LOG_LEVEL_ERROR,
|
||||
"--iconv: unsupported charset conversion requested (LOCAL[,REMOTE])");
|
||||
config_delete(config);
|
||||
close(file_descriptor);
|
||||
protocol_session_unbind();
|
||||
return;
|
||||
goto done;
|
||||
}
|
||||
charset_ready = true;
|
||||
}
|
||||
/* --delete-missing-args deletes destination mirrors receiver-side, so it is
|
||||
deletion and stays gated by the same --allow-delete server policy. When
|
||||
@@ -602,10 +592,7 @@ void handler(int file_descriptor) {
|
||||
log_message(LOG_LEVEL_ERROR, "destination root is not available: %s",
|
||||
escaped_root ? escaped_root : "<allocation failed>");
|
||||
free(escaped_root);
|
||||
config_delete(config);
|
||||
close(file_descriptor);
|
||||
protocol_session_unbind();
|
||||
return;
|
||||
goto done;
|
||||
}
|
||||
/* A --delay-updates transfer stages under a private 0700 directory inside
|
||||
the receive root. Create it up front (wiping leftovers of any previously
|
||||
@@ -614,11 +601,7 @@ void handler(int file_descriptor) {
|
||||
config->delay_context = delay_updates_context_create(config->receive_root_directory);
|
||||
if (!config->delay_context || !delay_updates_prepare(config->delay_context)) {
|
||||
log_message(LOG_LEVEL_ERROR, "Failed to initialize --delay-updates staging area");
|
||||
delay_updates_cleanup(config->delay_context);
|
||||
config_delete(config);
|
||||
close(file_descriptor);
|
||||
protocol_session_unbind();
|
||||
return;
|
||||
goto done;
|
||||
}
|
||||
}
|
||||
/* Preserve the negotiated identity policy for the fd-relative ownership
|
||||
@@ -628,10 +611,7 @@ void handler(int file_descriptor) {
|
||||
rather than silently applying the wrong ownership policy. */
|
||||
if (!identity_set_active(config)) {
|
||||
log_message(LOG_LEVEL_ERROR, "Failed to activate identity policy");
|
||||
config_delete(config);
|
||||
close(file_descriptor);
|
||||
protocol_session_unbind();
|
||||
return;
|
||||
goto done;
|
||||
}
|
||||
/* Persist the negotiated --keep-dirlinks policy once, here at config-accept,
|
||||
before any multithreaded receiver/writer threads are spawned, so the
|
||||
@@ -662,38 +642,25 @@ void handler(int file_descriptor) {
|
||||
if (!motd_send(file_descriptor, motd ? motd : "")) {
|
||||
free(motd);
|
||||
log_message(LOG_LEVEL_ERROR, "Failed to send daemon MOTD");
|
||||
config_delete(config);
|
||||
close(file_descriptor);
|
||||
protocol_session_unbind();
|
||||
identity_clear_active();
|
||||
return;
|
||||
goto done;
|
||||
}
|
||||
free(motd);
|
||||
}
|
||||
if (config->use_multithreading) {
|
||||
Queue* q = queue_create(100, file_destroy);
|
||||
if (q == NULL) {
|
||||
config_delete(config);
|
||||
close(file_descriptor);
|
||||
protocol_session_unbind();
|
||||
identity_clear_active();
|
||||
return;
|
||||
}
|
||||
PipelineContextReceiver* context =
|
||||
pipeline_context_receiver_create(config, q, file_descriptor, ssl);
|
||||
if (q == NULL)
|
||||
goto done;
|
||||
context = pipeline_context_receiver_create(config, q, file_descriptor, ssl);
|
||||
if (context == NULL) {
|
||||
queue_destroy(q);
|
||||
config_delete(config);
|
||||
close(file_descriptor);
|
||||
protocol_session_unbind();
|
||||
identity_clear_active();
|
||||
return;
|
||||
goto done;
|
||||
}
|
||||
protocol_session_set_max_alloc(&context->session, config->max_alloc);
|
||||
atomic_store(&context->session.total_allocated_bytes,
|
||||
atomic_load(&session.total_allocated_bytes));
|
||||
pipeline_context_receiver_set_queue_byte_limit(context, RECEIVER_QUEUE_MAX_BYTES);
|
||||
thrd_t receiver, writer;
|
||||
thrd_t receiver = {0};
|
||||
thrd_t writer = {0};
|
||||
bool receiver_created = thrd_create(&receiver, receive_thread, context) == thrd_success;
|
||||
bool writer_created = false;
|
||||
if (receiver_created)
|
||||
@@ -706,17 +673,19 @@ void handler(int file_descriptor) {
|
||||
cnd_broadcast(&context->condition_not_full);
|
||||
cnd_broadcast(&context->condition_not_empty);
|
||||
mtx_unlock(&context->mutex);
|
||||
close(file_descriptor);
|
||||
/* Unblock a worker parked in socket I/O without closing the fd (the
|
||||
* child owns the single close). shutdown() only affects sockets; for
|
||||
* the --stdio pipe the receiver's per-message poll timeout still
|
||||
* bounds the join, so do nothing there rather than close a descriptor
|
||||
* another thread may still be using. */
|
||||
struct stat fd_stat;
|
||||
if (fstat(file_descriptor, &fd_stat) == 0 && S_ISSOCK(fd_stat.st_mode))
|
||||
shutdown(file_descriptor, SHUT_RDWR);
|
||||
thrd_join(receiver, NULL);
|
||||
} else {
|
||||
close(file_descriptor);
|
||||
}
|
||||
if (writer_created)
|
||||
thrd_join(writer, NULL);
|
||||
pipeline_context_receiver_destroy(context);
|
||||
protocol_session_unbind();
|
||||
identity_clear_active();
|
||||
return;
|
||||
goto done;
|
||||
}
|
||||
int receiver_result;
|
||||
int writer_result;
|
||||
@@ -760,21 +729,36 @@ void handler(int file_descriptor) {
|
||||
} else {
|
||||
send_status(file_descriptor, STATUS_ERROR);
|
||||
}
|
||||
if (!transfer_ok) {
|
||||
if (!transfer_ok)
|
||||
log_message(LOG_LEVEL_ERROR, "Transfer failed");
|
||||
if (config->delay_updates && config->delay_context)
|
||||
delay_updates_cleanup(config->delay_context);
|
||||
}
|
||||
pipeline_context_receiver_destroy(context);
|
||||
} else {
|
||||
if (receiver_receive_files(config, file_descriptor) != 0)
|
||||
log_message(LOG_LEVEL_ERROR, "Transfer failed");
|
||||
config_delete(config);
|
||||
}
|
||||
protocol_session_unbind();
|
||||
identity_clear_active();
|
||||
|
||||
done:
|
||||
/* Single cleanup epilogue: every error path jumps here, so the iconv
|
||||
* receiver conversion is released, the identity snapshot cleared, the
|
||||
* protocol session unbound and the config freed exactly once. The
|
||||
* connection fd is deliberately NOT closed here -- the child functions own
|
||||
* its single close (plain_child_fn / tls_child_fn), and the --stdio call
|
||||
* site must leave stdin/stdout open. */
|
||||
if (charset_ready)
|
||||
charset_wire_free();
|
||||
close(file_descriptor);
|
||||
/* The delay-updates staging tree is released by config_delete (which the
|
||||
branch below always reaches), so it is cleaned exactly once. */
|
||||
identity_clear_active();
|
||||
protocol_session_unbind();
|
||||
if (context != NULL) {
|
||||
/* context owns both the config and the queue it was created with. */
|
||||
pipeline_context_receiver_destroy(context);
|
||||
context = NULL;
|
||||
config = NULL;
|
||||
} else {
|
||||
config_delete(config);
|
||||
config = NULL;
|
||||
}
|
||||
free(joined_destination);
|
||||
}
|
||||
|
||||
#ifndef FASTSYNC_SERVER_AS_LIB
|
||||
@@ -969,6 +953,9 @@ int main(int argc, char* argv[]) {
|
||||
return 1;
|
||||
}
|
||||
io_set_fds(STDIN_FILENO, STDOUT_FILENO);
|
||||
/* handler() does not own the stdio fds: it never closes its descriptor
|
||||
* argument, so STDIN/STDOUT stay open for this (single-shot) SSH session
|
||||
* and are released by process exit. */
|
||||
handler(STDIN_FILENO);
|
||||
release_authorization();
|
||||
server_cli_options_free(&opts);
|
||||
|
||||
+7
-4
@@ -693,6 +693,9 @@ void config_delete(Config* config) {
|
||||
if (config == NULL)
|
||||
return;
|
||||
if (config->log_file) {
|
||||
/* The logging subsystem borrows this FILE*; detach it before closing so a
|
||||
* concurrent log call can never touch the freed handle. */
|
||||
log_set_file(NULL);
|
||||
fclose(config->log_file);
|
||||
config->log_file = NULL;
|
||||
}
|
||||
@@ -1428,7 +1431,7 @@ Config* config_receive_with_validate(int file_descriptor, ConfigValidateFunc val
|
||||
goto error;
|
||||
if (strcmp(config->version, PROTOCOL_VERSION) != 0) {
|
||||
char* escaped_version = output_escape(config->version, false);
|
||||
fprintf(stderr, "Protocol version mismatch: client=%s, server=%s\n",
|
||||
log_message(LOG_LEVEL_ERROR, "Protocol version mismatch: client=%s, server=%s",
|
||||
escaped_version ? escaped_version : "<allocation failed>", PROTOCOL_VERSION);
|
||||
free(escaped_version);
|
||||
send_status(file_descriptor, STATUS_ERROR);
|
||||
@@ -1455,14 +1458,14 @@ Config* config_receive_with_validate(int file_descriptor, ConfigValidateFunc val
|
||||
if (config->compress_choice[0] != '\0' && strcmp(config->compress_choice, "zstd") != 0 &&
|
||||
strcmp(config->compress_choice, "none") != 0) {
|
||||
char* escaped_choice = output_escape(config->compress_choice, config->eight_bit_output);
|
||||
fprintf(stderr, "Unsupported compression choice: %s\n",
|
||||
log_message(LOG_LEVEL_ERROR, "Unsupported compression choice: %s",
|
||||
escaped_choice ? escaped_choice : "<allocation failed>");
|
||||
free(escaped_choice);
|
||||
send_status(file_descriptor, STATUS_ERROR);
|
||||
goto error;
|
||||
}
|
||||
if (!validate_received_config(config)) {
|
||||
fprintf(stderr, "Invalid configuration received from client\n");
|
||||
log_message(LOG_LEVEL_ERROR, "Invalid configuration received from client");
|
||||
send_status(file_descriptor, STATUS_ERROR);
|
||||
goto error;
|
||||
}
|
||||
@@ -1476,7 +1479,7 @@ Config* config_receive_with_validate(int file_descriptor, ConfigValidateFunc val
|
||||
* the CONFIG_VALIDATE_ALREADY_TERMINATED sentinel, so no second status is
|
||||
* written. */
|
||||
if (rejection != CONFIG_VALIDATE_ALREADY_TERMINATED) {
|
||||
fprintf(stderr, "%s\n", rejection);
|
||||
log_message(LOG_LEVEL_ERROR, "%s", rejection);
|
||||
send_status(file_descriptor, STATUS_ERROR);
|
||||
}
|
||||
goto error;
|
||||
|
||||
+72
-26
@@ -3,7 +3,9 @@
|
||||
#include <stdbool.h>
|
||||
#include <stdarg.h>
|
||||
#include <stdio.h>
|
||||
#include <stdlib.h>
|
||||
#include <string.h>
|
||||
#include <threads.h>
|
||||
#include <time.h>
|
||||
|
||||
static const char* log_level_strings[] = {"DEBUG", "INFO", "WARN", "ERROR"};
|
||||
@@ -15,6 +17,18 @@ static FILE* log_fp = NULL;
|
||||
static _Thread_local bool eight_bit_output;
|
||||
static LogStderrMode stderr_mode = LOG_STDERR_ERRORS;
|
||||
|
||||
/* Serializes access to log_fp and makes each emitted line atomic: the
|
||||
* timestamp prefix, formatted body, and trailing newline are written as one
|
||||
* critical section so concurrent threads cannot interleave partial lines.
|
||||
* Initialized lazily (matching the protocol.c bw_mutex idiom) because logging
|
||||
* can happen before main() installs any synchronization. */
|
||||
static mtx_t log_mutex;
|
||||
static once_flag log_mutex_once = ONCE_FLAG_INIT;
|
||||
|
||||
static void log_mutex_init(void) {
|
||||
mtx_init(&log_mutex, mtx_plain);
|
||||
}
|
||||
|
||||
void set_log_level(LogLevel level) {
|
||||
current_log_level = level;
|
||||
}
|
||||
@@ -41,7 +55,10 @@ uint32_t get_log_info_flags(void) {
|
||||
}
|
||||
|
||||
void log_set_file(FILE* fp) {
|
||||
call_once(&log_mutex_once, log_mutex_init);
|
||||
mtx_lock(&log_mutex);
|
||||
log_fp = fp;
|
||||
mtx_unlock(&log_mutex);
|
||||
}
|
||||
|
||||
void log_set_8_bit_output(bool enabled) {
|
||||
@@ -60,13 +77,48 @@ LogStderrMode log_get_stderr_mode(void) {
|
||||
return stderr_mode;
|
||||
}
|
||||
|
||||
static inline void write_message(FILE* dest_io, LogLevel log_level, struct tm t, const char* format,
|
||||
/* Format one complete log line (timestamp prefix + body + newline) into a
|
||||
* freshly allocated buffer. This is pure CPU/malloc work and must happen
|
||||
* OUTSIDE the log mutex: the mutex only guards the log_fp pointer, so a
|
||||
* stalled stderr/stdout pipe cannot block every logging thread. Returns NULL
|
||||
* on allocation/formatting failure. */
|
||||
static char* format_log_line(LogLevel log_level, const struct tm* t, const char* format,
|
||||
va_list args) {
|
||||
fprintf(dest_io, "%04d-%02d-%02d %02d:%02d:%02d [%s]: ", t.tm_year + 1900, t.tm_mon + 1,
|
||||
t.tm_mday, t.tm_hour, t.tm_min, t.tm_sec, log_level_strings[log_level]);
|
||||
char prefix[64];
|
||||
int prefix_len = snprintf(
|
||||
prefix, sizeof(prefix), "%04d-%02d-%02d %02d:%02d:%02d [%s]: ", t->tm_year + 1900,
|
||||
t->tm_mon + 1, t->tm_mday, t->tm_hour, t->tm_min, t->tm_sec, log_level_strings[log_level]);
|
||||
if (prefix_len < 0 || prefix_len >= (int)sizeof(prefix))
|
||||
return NULL;
|
||||
va_list copy;
|
||||
va_copy(copy, args);
|
||||
int body_len = vsnprintf(NULL, 0, format, copy);
|
||||
va_end(copy);
|
||||
if (body_len < 0)
|
||||
return NULL;
|
||||
size_t total = (size_t)prefix_len + (size_t)body_len;
|
||||
char* line = malloc(total + 2); /* body bytes + '\n' + NUL */
|
||||
if (!line)
|
||||
return NULL;
|
||||
memcpy(line, prefix, (size_t)prefix_len);
|
||||
vsnprintf(line + prefix_len, (size_t)body_len + 1, format, args);
|
||||
line[total] = '\n';
|
||||
line[total + 1] = '\0';
|
||||
return line;
|
||||
}
|
||||
|
||||
vfprintf(dest_io, format, args);
|
||||
fprintf(dest_io, "\n");
|
||||
/* Write an already-formatted line to the console and, if configured, the log
|
||||
* file. Only the log_fp pointer is read under the mutex (so log_set_file /
|
||||
* config_delete cannot free it while it is in use); the single console fputs
|
||||
* runs unlocked but is internally atomic per stdio stream. */
|
||||
static void emit_log_line(FILE* console, const char* line) {
|
||||
fputs(line, console);
|
||||
call_once(&log_mutex_once, log_mutex_init);
|
||||
mtx_lock(&log_mutex);
|
||||
FILE* file = log_fp;
|
||||
if (file)
|
||||
fputs(line, file);
|
||||
mtx_unlock(&log_mutex);
|
||||
}
|
||||
|
||||
void log_message(LogLevel log_level, const char* format, ...) {
|
||||
@@ -86,14 +138,12 @@ void log_message(LogLevel log_level, const char* format, ...) {
|
||||
|
||||
va_list args;
|
||||
va_start(args, format);
|
||||
write_message(dest_io, log_level, t, format, args);
|
||||
char* line = format_log_line(log_level, &t, format, args);
|
||||
va_end(args);
|
||||
|
||||
if (log_fp) {
|
||||
va_start(args, format);
|
||||
write_message(log_fp, log_level, t, format, args);
|
||||
va_end(args);
|
||||
}
|
||||
if (!line)
|
||||
return;
|
||||
emit_log_line(dest_io, line);
|
||||
free(line);
|
||||
}
|
||||
|
||||
void log_debug_message(LogDebugFlag flag, const char* format, ...) {
|
||||
@@ -107,14 +157,12 @@ void log_debug_message(LogDebugFlag flag, const char* format, ...) {
|
||||
|
||||
va_list args;
|
||||
va_start(args, format);
|
||||
write_message(stdout, LOG_LEVEL_DEBUG, t, format, args);
|
||||
char* line = format_log_line(LOG_LEVEL_DEBUG, &t, format, args);
|
||||
va_end(args);
|
||||
|
||||
if (log_fp) {
|
||||
va_start(args, format);
|
||||
write_message(log_fp, LOG_LEVEL_DEBUG, t, format, args);
|
||||
va_end(args);
|
||||
}
|
||||
if (!line)
|
||||
return;
|
||||
emit_log_line(stdout, line);
|
||||
free(line);
|
||||
}
|
||||
|
||||
void log_info_message(LogInfoFlag flag, const char* format, ...) {
|
||||
@@ -129,14 +177,12 @@ void log_info_message(LogInfoFlag flag, const char* format, ...) {
|
||||
|
||||
va_list args;
|
||||
va_start(args, format);
|
||||
write_message(stdout, LOG_LEVEL_INFO, t, format, args);
|
||||
char* line = format_log_line(LOG_LEVEL_INFO, &t, format, args);
|
||||
va_end(args);
|
||||
|
||||
if (log_fp) {
|
||||
va_start(args, format);
|
||||
write_message(log_fp, LOG_LEVEL_INFO, t, format, args);
|
||||
va_end(args);
|
||||
}
|
||||
if (!line)
|
||||
return;
|
||||
emit_log_line(stdout, line);
|
||||
free(line);
|
||||
}
|
||||
|
||||
void log_perror(const char* context) {
|
||||
|
||||
@@ -152,9 +152,18 @@ static void accept_loop(Server* server, void (*child_fn)(int, void*), void* chil
|
||||
log_message(LOG_LEVEL_INFO, "%s", log_fmt);
|
||||
pid_t pid = fork();
|
||||
if (pid == 0) {
|
||||
/* Connection children must not run the parent's global cleanup(): it
|
||||
* frees state (credentials / daemon conf) that the child's worker
|
||||
* threads may still be reading and closes fd numbers the child could
|
||||
* already have reused. Reset the inherited handlers so a signal
|
||||
* terminates the child directly; SIGCHLD is reset too since a child
|
||||
* must never reap the parent's children. This runs before the child
|
||||
* spawns any thread, so it cannot race one. */
|
||||
signal(SIGINT, SIG_DFL);
|
||||
signal(SIGTERM, SIG_DFL);
|
||||
signal(SIGCHLD, SIG_DFL);
|
||||
close(server->file_descriptor);
|
||||
child_fn(fd, child_ctx);
|
||||
close(fd);
|
||||
_exit(0);
|
||||
} else if (pid > 0) {
|
||||
g_active_connections++;
|
||||
@@ -169,6 +178,9 @@ struct plain_ctx {
|
||||
|
||||
static void plain_child_fn(int fd, void* ctx) {
|
||||
((struct plain_ctx*)ctx)->handler(fd);
|
||||
/* handler() never closes the connection fd; the child owns its single
|
||||
* close here after the handler has fully torn down. */
|
||||
close(fd);
|
||||
}
|
||||
|
||||
bool server_listen(Server* server, void (*handler)(int file_descriptor)) {
|
||||
|
||||
@@ -192,13 +192,18 @@ static void tls_child_fn(int fd, void* arg) {
|
||||
SSL* ssl = wrap_fd_with_ssl(fd, ctx->ssl_ctx, true, NULL);
|
||||
if (!ssl) {
|
||||
io_set_ssl(NULL);
|
||||
close(fd);
|
||||
return;
|
||||
}
|
||||
io_set_ssl(ssl);
|
||||
ctx->handler(fd);
|
||||
/* Shut the TLS layer down before releasing the fd: handler() no longer
|
||||
* closes it, so SSL_shutdown still has a valid socket. The child owns the
|
||||
* single fd close, performed last. */
|
||||
SSL_shutdown(ssl);
|
||||
SSL_free(ssl);
|
||||
io_set_ssl(NULL);
|
||||
close(fd);
|
||||
}
|
||||
|
||||
bool server_listen_tls(Server* server, void (*handler)(int file_descriptor)) {
|
||||
|
||||
@@ -1,7 +1,9 @@
|
||||
#include "test_log.h"
|
||||
#include "log.h"
|
||||
#include "test_utils.h"
|
||||
#include <fcntl.h>
|
||||
#include <string.h>
|
||||
#include <threads.h>
|
||||
#include <unistd.h>
|
||||
|
||||
/* Test default log level: WARNING and ERROR should print, DEBUG and INFO should not.
|
||||
@@ -147,6 +149,118 @@ static void test_log_debug_enabled_matches_gate() {
|
||||
set_log_debug_flags(LOG_DEBUG_ALL);
|
||||
}
|
||||
|
||||
#define LOG_CONCURRENCY_THREADS 8
|
||||
#define LOG_CONCURRENCY_LINES 250
|
||||
|
||||
typedef struct {
|
||||
int id;
|
||||
} LogConcurrencyArg;
|
||||
|
||||
static int log_concurrency_worker(void* context) {
|
||||
LogConcurrencyArg* arg = context;
|
||||
for (int i = 0; i < LOG_CONCURRENCY_LINES; i++) {
|
||||
log_message(LOG_LEVEL_WARNING, "worker %d line %d", arg->id, i);
|
||||
}
|
||||
return 0;
|
||||
}
|
||||
|
||||
static int count_substring(const char* haystack, const char* needle) {
|
||||
int count = 0;
|
||||
size_t needle_length = strlen(needle);
|
||||
const char* cursor = haystack;
|
||||
while ((cursor = strstr(cursor, needle)) != NULL) {
|
||||
count++;
|
||||
cursor += needle_length;
|
||||
}
|
||||
return count;
|
||||
}
|
||||
|
||||
/* Concurrent log_message() calls from many threads must never interleave a
|
||||
* single line: every emitted line has exactly one timestamp prefix and one
|
||||
* body. Before write_message() was serialized, the three separate fprintf
|
||||
* calls (prefix, body, newline) let lines tear. */
|
||||
static void test_log_concurrent_no_torn_lines(void) {
|
||||
FILE* fp = tmpfile();
|
||||
EXPECT_NOT_NULL(fp);
|
||||
|
||||
/* Mute the console mirror so the workers don't flood the test output. */
|
||||
fflush(stdout);
|
||||
fflush(stderr);
|
||||
int saved_stdout = dup(STDOUT_FILENO);
|
||||
int saved_stderr = dup(STDERR_FILENO);
|
||||
int null_fd = open("/dev/null", O_WRONLY);
|
||||
EXPECT_TRUE(saved_stdout >= 0);
|
||||
EXPECT_TRUE(saved_stderr >= 0);
|
||||
EXPECT_TRUE(null_fd >= 0);
|
||||
EXPECT_TRUE(dup2(null_fd, STDOUT_FILENO) >= 0);
|
||||
EXPECT_TRUE(dup2(null_fd, STDERR_FILENO) >= 0);
|
||||
close(null_fd);
|
||||
|
||||
set_log_level(LOG_LEVEL_WARNING);
|
||||
log_set_stderr_mode(LOG_STDERR_ERRORS);
|
||||
log_set_file(fp);
|
||||
|
||||
thrd_t threads[LOG_CONCURRENCY_THREADS];
|
||||
LogConcurrencyArg args[LOG_CONCURRENCY_THREADS];
|
||||
int created = 0;
|
||||
for (int i = 0; i < LOG_CONCURRENCY_THREADS; i++) {
|
||||
args[i].id = i;
|
||||
if (thrd_create(&threads[i], log_concurrency_worker, &args[i]) != thrd_success)
|
||||
break;
|
||||
created++;
|
||||
}
|
||||
for (int i = 0; i < created; i++) {
|
||||
thrd_join(threads[i], NULL);
|
||||
}
|
||||
|
||||
log_set_file(NULL);
|
||||
fflush(fp);
|
||||
|
||||
fflush(stdout);
|
||||
fflush(stderr);
|
||||
dup2(saved_stdout, STDOUT_FILENO);
|
||||
dup2(saved_stderr, STDERR_FILENO);
|
||||
close(saved_stdout);
|
||||
close(saved_stderr);
|
||||
|
||||
rewind(fp);
|
||||
char line[512];
|
||||
int total_lines = 0;
|
||||
int malformed_lines = 0;
|
||||
bool saw_missing_newline = false;
|
||||
while (fgets(line, sizeof(line), fp) != NULL) {
|
||||
size_t length = strlen(line);
|
||||
if (length == 0 || line[length - 1] != '\n')
|
||||
saw_missing_newline = true;
|
||||
if (strncmp(line, "20", 2) != 0 || count_substring(line, "[WARN]: worker ") != 1)
|
||||
malformed_lines++;
|
||||
total_lines++;
|
||||
}
|
||||
fclose(fp);
|
||||
|
||||
EXPECT_EQ_INT(created, LOG_CONCURRENCY_THREADS);
|
||||
EXPECT_FALSE(saw_missing_newline);
|
||||
EXPECT_EQ_INT(malformed_lines, 0);
|
||||
EXPECT_EQ_INT(total_lines, LOG_CONCURRENCY_THREADS * LOG_CONCURRENCY_LINES);
|
||||
log_set_stderr_mode(LOG_STDERR_ERRORS);
|
||||
}
|
||||
|
||||
/* Detaching the logger from a FILE* before it is closed must leave the logging
|
||||
* subsystem safe: later calls must not touch the freed handle. */
|
||||
static void test_log_set_file_null_before_fclose(void) {
|
||||
FILE* fp = tmpfile();
|
||||
EXPECT_NOT_NULL(fp);
|
||||
|
||||
set_log_level(LOG_LEVEL_ERROR);
|
||||
log_set_file(fp);
|
||||
log_message(LOG_LEVEL_ERROR, "line before detach");
|
||||
log_set_file(NULL);
|
||||
fclose(fp);
|
||||
|
||||
log_message(LOG_LEVEL_ERROR, "line after close");
|
||||
EXPECT_TRUE(true);
|
||||
}
|
||||
|
||||
void test_log() {
|
||||
test_log_message_debug();
|
||||
test_log_message_info();
|
||||
@@ -159,4 +273,6 @@ void test_log() {
|
||||
test_log_stderr_mode_all();
|
||||
test_log_message_formats();
|
||||
test_log_debug_enabled_matches_gate();
|
||||
test_log_concurrent_no_torn_lines();
|
||||
test_log_set_file_null_before_fclose();
|
||||
}
|
||||
|
||||
@@ -1317,6 +1317,58 @@ static void test_scanner_captures_directory_times() {
|
||||
rmdir(root);
|
||||
}
|
||||
|
||||
/* Ownership guard for chunk_data_to_chunk(): a returned Chunk owns its File
|
||||
* objects, so destroying the chunk must free them exactly once and the scanner
|
||||
* must never free them again. chunk_size = 1 forces the mid-directory
|
||||
* conversion branch (chunk_data_size > chunk_size) for every file, and the
|
||||
* chunk is destroyed immediately, catching a double free / use-after-free under
|
||||
* ASan if ownership transfer regressed.
|
||||
*
|
||||
* The failure path (array_list_to_array() or chunk_create() returning NULL) is
|
||||
* not reachable from a unit test: both allocate through protocol_alloc(), and
|
||||
* each allocation they perform is no larger than the array_list allocations
|
||||
* that already succeeded while building the list (array_list_to_array() copies
|
||||
* exactly `size` pointers, which never exceeds the capacity just grown, and
|
||||
* sizeof(Chunk) is far below the initial 100-entry item array). Binding a
|
||||
* small --max-alloc session therefore always fails *before* this function, not
|
||||
* inside it, so fault injection cannot isolate these paths. */
|
||||
static void test_scanner_chunk_ownership() {
|
||||
const char* dir = "test_scan_ownership";
|
||||
const char* file1 = "test_scan_ownership/a.txt";
|
||||
const char* file2 = "test_scan_ownership/b.txt";
|
||||
const char* file3 = "test_scan_ownership/c.txt";
|
||||
|
||||
EXPECT_EQ_INT(mkdir(dir, 0755), 0);
|
||||
create_test_file(file1, "aaaa");
|
||||
create_test_file(file2, "bbbb");
|
||||
create_test_file(file3, "cccc");
|
||||
|
||||
ScannerOptions options = {0};
|
||||
options.chunk_size = 1;
|
||||
DirectoryScanner* scanner = directory_scanner_create_with_options(dir, &options);
|
||||
EXPECT_NOT_NULL(scanner);
|
||||
|
||||
int chunks = 0;
|
||||
int files = 0;
|
||||
Chunk* chunk;
|
||||
while ((chunk = directory_scanner_next(scanner)) != NULL) {
|
||||
chunks++;
|
||||
files += chunk->element_count;
|
||||
EXPECT_EQ_INT(chunk->element_count, 1);
|
||||
chunk_destroy(chunk);
|
||||
EXPECT_FALSE(directory_scanner_failed(scanner));
|
||||
}
|
||||
EXPECT_EQ_INT(files, 3);
|
||||
EXPECT_EQ_INT(chunks, 3);
|
||||
EXPECT_FALSE(directory_scanner_failed(scanner));
|
||||
|
||||
directory_scanner_destroy(scanner);
|
||||
unlink(file1);
|
||||
unlink(file2);
|
||||
unlink(file3);
|
||||
rmdir(dir);
|
||||
}
|
||||
|
||||
void test_scanner() {
|
||||
test_scanner_single_file();
|
||||
test_scanner_multiple_files();
|
||||
@@ -1353,4 +1405,5 @@ void test_scanner() {
|
||||
test_dirs_files_from();
|
||||
test_files_from_relative_send_path();
|
||||
test_scanner_captures_directory_times();
|
||||
test_scanner_chunk_ownership();
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user