diff --git a/src/shared/config.c b/src/shared/config.c index 584c61a..4816e06 100644 --- a/src/shared/config.c +++ b/src/shared/config.c @@ -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,8 +1431,8 @@ 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", - escaped_version ? escaped_version : "", PROTOCOL_VERSION); + log_message(LOG_LEVEL_ERROR, "Protocol version mismatch: client=%s, server=%s", + escaped_version ? escaped_version : "", PROTOCOL_VERSION); free(escaped_version); send_status(file_descriptor, STATUS_ERROR); goto 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", - escaped_choice ? escaped_choice : ""); + log_message(LOG_LEVEL_ERROR, "Unsupported compression choice: %s", + escaped_choice ? escaped_choice : ""); 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; diff --git a/src/shared/log.c b/src/shared/log.c index d013bbe..d0e1f0a 100644 --- a/src/shared/log.c +++ b/src/shared/log.c @@ -4,6 +4,7 @@ #include #include #include +#include #include static const char* log_level_strings[] = {"DEBUG", "INFO", "WARN", "ERROR"}; @@ -15,6 +16,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 +54,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) { @@ -79,6 +95,9 @@ void log_message(LogLevel log_level, const char* format, ...) { if (!localtime_r(&now, &t)) return; + call_once(&log_mutex_once, log_mutex_init); + mtx_lock(&log_mutex); + FILE* dest_io = stdout; if (stderr_mode == LOG_STDERR_ALL || log_level == LOG_LEVEL_ERROR) { dest_io = stderr; @@ -94,6 +113,8 @@ void log_message(LogLevel log_level, const char* format, ...) { write_message(log_fp, log_level, t, format, args); va_end(args); } + + mtx_unlock(&log_mutex); } void log_debug_message(LogDebugFlag flag, const char* format, ...) { @@ -105,6 +126,9 @@ void log_debug_message(LogDebugFlag flag, const char* format, ...) { if (!localtime_r(&now, &t)) return; + call_once(&log_mutex_once, log_mutex_init); + mtx_lock(&log_mutex); + va_list args; va_start(args, format); write_message(stdout, LOG_LEVEL_DEBUG, t, format, args); @@ -115,6 +139,8 @@ void log_debug_message(LogDebugFlag flag, const char* format, ...) { write_message(log_fp, LOG_LEVEL_DEBUG, t, format, args); va_end(args); } + + mtx_unlock(&log_mutex); } void log_info_message(LogInfoFlag flag, const char* format, ...) { @@ -127,6 +153,9 @@ void log_info_message(LogInfoFlag flag, const char* format, ...) { if (!localtime_r(&now, &t)) return; + call_once(&log_mutex_once, log_mutex_init); + mtx_lock(&log_mutex); + va_list args; va_start(args, format); write_message(stdout, LOG_LEVEL_INFO, t, format, args); @@ -137,6 +166,8 @@ void log_info_message(LogInfoFlag flag, const char* format, ...) { write_message(log_fp, LOG_LEVEL_INFO, t, format, args); va_end(args); } + + mtx_unlock(&log_mutex); } void log_perror(const char* context) { diff --git a/tests/test_log.c b/tests/test_log.c index f9caf19..b5acedd 100644 --- a/tests/test_log.c +++ b/tests/test_log.c @@ -1,7 +1,9 @@ #include "test_log.h" #include "log.h" #include "test_utils.h" +#include #include +#include #include /* 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(); }