From 5eaea9d92e027bb94c52480d3b39107e9ee4cf95 Mon Sep 17 00:00:00 2001 From: TapTap Date: Sat, 1 Aug 2026 18:12:41 +0200 Subject: [PATCH 1/4] refactor: extract CLI parser helpers, fix validation bugs - Add set_string_option(), set_positive_int_option(), set_nonneg_int_option() helpers to eliminate duplicated alloc/free/assign patterns - Replace 14 string-option branches with set_string_option() calls - Replace 7 integer-option branches with set_positive/nonneg_int_option() - Fix --info/--debug: replace atoi() with validated parse - Fix --max-size/--min-size: add strtoull() error checking - Fix --chunk-size: validate input with error message on failure --- src/client/client_cli.c | 191 +++++++++++++++++----------------------- 1 file changed, 81 insertions(+), 110 deletions(-) diff --git a/src/client/client_cli.c b/src/client/client_cli.c index de1fc3c..fac0f98 100644 --- a/src/client/client_cli.c +++ b/src/client/client_cli.c @@ -49,6 +49,37 @@ static bool parse_nonneg_int(const char* s, int* out_val) { return true; } +/* Duplicate a string argument into *dest, freeing the old value. Returns 0 on success, -1 on + * failure. */ +static int set_string_option(char** dest, const char* value, const char* option_name) { + char* dup = str_dup(value); + if (!dup) { + fprintf(stderr, "Error: memory allocation failed for %s\n", option_name); + return -1; + } + free(*dest); + *dest = dup; + return 0; +} + +/* Parse a string as a positive integer into *dest. Returns 0 on success, -1 on error. */ +static int set_positive_int_option(int* dest, const char* value, const char* option_name) { + if (!parse_positive_int(value, dest)) { + fprintf(stderr, "Error: %s must be a positive integer\n", option_name); + return -1; + } + return 0; +} + +/* Parse a string as a non-negative integer into *dest. Returns 0 on success, -1 on error. */ +static int set_nonneg_int_option(int* dest, const char* value, const char* option_name) { + if (!parse_nonneg_int(value, dest)) { + fprintf(stderr, "Error: %s must be a non-negative integer\n", option_name); + return -1; + } + return 0; +} + static void print_usage(void); static int read_patterns_from_file(const char* filepath, char*** patterns, int* count); @@ -100,9 +131,23 @@ static int parse_args(Config* config, int argc, char* argv[], int* positional_ar } config->include_patterns[config->include_count++] = dup; } else if (strcmp(argv[i], "--max-size") == 0 && i + 1 < argc) { - config->max_size = strtoull(argv[++i], NULL, 10); + char* end; + errno = 0; + unsigned long long val = strtoull(argv[++i], &end, 10); + if (errno != 0 || *end != '\0') { + fprintf(stderr, "Error: --max-size must be a non-negative integer\n"); + return -1; + } + config->max_size = val; } else if (strcmp(argv[i], "--min-size") == 0 && i + 1 < argc) { - config->min_size = strtoull(argv[++i], NULL, 10); + char* end; + errno = 0; + unsigned long long val = strtoull(argv[++i], &end, 10); + if (errno != 0 || *end != '\0') { + fprintf(stderr, "Error: --min-size must be a non-negative integer\n"); + return -1; + } + config->min_size = val; } else if (strcmp(argv[i], "--incremental") == 0) { config->use_incremental = true; } else if (strcmp(argv[i], "--delta") == 0) { @@ -132,21 +177,11 @@ static int parse_args(Config* config, int argc, char* argv[], int* positional_ar } } } else if (strcmp(argv[i], "--source-dir") == 0 && i + 1 < argc) { - char* dup = str_dup(argv[++i]); - if (!dup) { - fprintf(stderr, "Error: memory allocation failed for --source-dir\n"); + if (set_string_option(&config->send_directory, argv[++i], "--source-dir") != 0) return -1; - } - free(config->send_directory); - config->send_directory = dup; } else if (strcmp(argv[i], "--dest-dir") == 0 && i + 1 < argc) { - char* dup = str_dup(argv[++i]); - if (!dup) { - fprintf(stderr, "Error: memory allocation failed for --dest-dir\n"); + if (set_string_option(&config->receive_root_directory, argv[++i], "--dest-dir") != 0) return -1; - } - free(config->receive_root_directory); - config->receive_root_directory = dup; } else if (strcmp(argv[i], "--save-to-disk") == 0) { config->save_to_disk = true; } else if (strcmp(argv[i], "-M") == 0 || strcmp(argv[i], "--preserve") == 0) { @@ -162,13 +197,8 @@ static int parse_args(Config* config, int argc, char* argv[], int* positional_ar config->use_chunk_serialization = true; log_message(LOG_LEVEL_INFO, "Enabled Chunk Serialization"); } else if (strcmp(argv[i], "--server-host") == 0 && i + 1 < argc) { - char* dup = str_dup(argv[++i]); - if (!dup) { - fprintf(stderr, "Error: memory allocation failed for --server-host\n"); + if (set_string_option(&config->server_host, argv[++i], "--server-host") != 0) return -1; - } - free(config->server_host); - config->server_host = dup; } else if (strcmp(argv[i], "--server-port") == 0 && i + 1 < argc) { if (!parse_positive_int(argv[++i], &config->server_port)) { fprintf(stderr, "Error: invalid --server-port value: %s\n", argv[i]); @@ -191,69 +221,44 @@ static int parse_args(Config* config, int argc, char* argv[], int* positional_ar } else if (strcmp(argv[i], "--progress") == 0) { config->show_progress = true; } else if (strcmp(argv[i], "--chunk-size") == 0 && i + 1 < argc) { - unsigned long long val = strtoull(argv[++i], NULL, 10); - if (val > 0) - config->chunk_size = val; + char* end; + errno = 0; + unsigned long long val = strtoull(argv[++i], &end, 10); + if (errno != 0 || *end != '\0' || val == 0) { + fprintf(stderr, "Error: --chunk-size must be a positive integer\n"); + return -1; + } + config->chunk_size = val; } else if (strcmp(argv[i], "--tls") == 0) { config->use_tls = true; } else if (strcmp(argv[i], "--cert") == 0 && i + 1 < argc) { - char* dup = str_dup(argv[++i]); - if (!dup) { - fprintf(stderr, "Error: memory allocation failed for --cert\n"); + if (set_string_option(&config->tls_cert, argv[++i], "--cert") != 0) return -1; - } - free(config->tls_cert); - config->tls_cert = dup; } else if (strcmp(argv[i], "--key") == 0 && i + 1 < argc) { - char* dup = str_dup(argv[++i]); - if (!dup) { - fprintf(stderr, "Error: memory allocation failed for --key\n"); + if (set_string_option(&config->tls_key, argv[++i], "--key") != 0) return -1; - } - free(config->tls_key); - config->tls_key = dup; } else if (strcmp(argv[i], "--ca") == 0 && i + 1 < argc) { - char* dup = str_dup(argv[++i]); - if (!dup) { - fprintf(stderr, "Error: memory allocation failed for --ca\n"); + if (set_string_option(&config->tls_ca, argv[++i], "--ca") != 0) return -1; - } - free(config->tls_ca); - config->tls_ca = dup; } else if (strcmp(argv[i], "--timeout") == 0 && i + 1 < argc) { - int val; - if (!parse_positive_int(argv[++i], &val)) { - fprintf(stderr, "Error: --timeout must be a positive integer\n"); + if (set_positive_int_option(&config->timeout, argv[++i], "--timeout") != 0) return -1; - } - config->timeout = val; } else if (strcmp(argv[i], "--contimeout") == 0 && i + 1 < argc) { - int val; - if (!parse_positive_int(argv[++i], &val)) { - fprintf(stderr, "Error: --contimeout must be a positive integer\n"); + if (set_positive_int_option(&config->contimeout, argv[++i], "--contimeout") != 0) return -1; - } - config->contimeout = val; } else if (strcmp(argv[i], "-q") == 0 || strcmp(argv[i], "--quiet") == 0 || strcmp(argv[i], "--silent") == 0) { config->quiet = true; } else if (strcmp(argv[i], "--backup") == 0) { config->backup = true; } else if (strcmp(argv[i], "--backup-dir") == 0 && i + 1 < argc) { - char* dup = str_dup(argv[++i]); - if (!dup) { - fprintf(stderr, "Error: memory allocation failed for --backup-dir\n"); + if (set_string_option(&config->backup_dir, argv[++i], "--backup-dir") != 0) return -1; - } - free(config->backup_dir); - config->backup_dir = dup; } else if (strcmp(argv[i], "--stats") == 0) { config->stats = true; } else if (strcmp(argv[i], "--max-depth") == 0 && i + 1 < argc) { - if (!parse_nonneg_int(argv[++i], &config->max_depth)) { - fprintf(stderr, "Error: --max-depth must be a non-negative integer\n"); + if (set_nonneg_int_option(&config->max_depth, argv[++i], "--max-depth") != 0) return -1; - } } else if (strcmp(argv[i], "--log-file") == 0 && i + 1 < argc) { FILE* lf = fopen(argv[++i], "a"); if (!lf) { @@ -263,12 +268,8 @@ static int parse_args(Config* config, int argc, char* argv[], int* positional_ar config->log_file = lf; log_set_file(lf); } else if (strcmp(argv[i], "--queue-size") == 0 && i + 1 < argc) { - int val; - if (!parse_positive_int(argv[++i], &val)) { - fprintf(stderr, "Error: --queue-size must be a positive integer\n"); + if (set_positive_int_option(&config->queue_size, argv[++i], "--queue-size") != 0) return -1; - } - config->queue_size = val; } else if (strcmp(argv[i], "--exclude-from") == 0 && i + 1 < argc) { if (read_patterns_from_file(argv[++i], &config->exclude_patterns, &config->exclude_count) != 0) @@ -280,13 +281,9 @@ static int parse_args(Config* config, int argc, char* argv[], int* positional_ar } else if (strcmp(argv[i], "--partial") == 0) { config->partial = true; } else if (strcmp(argv[i], "--fastsync-server-path") == 0 && i + 1 < argc) { - char* dup = str_dup(argv[++i]); - if (!dup) { - fprintf(stderr, "Error: memory allocation failed for --fastsync-server-path\n"); + if (set_string_option(&config->fastsync_server_path, argv[++i], "--fastsync-server-path") != + 0) return -1; - } - free(config->fastsync_server_path); - config->fastsync_server_path = dup; } else if (strcmp(argv[i], "-v") == 0 || strcmp(argv[i], "--verbose") == 0) { set_log_level(LOG_LEVEL_DEBUG); } else if (strcmp(argv[i], "-l") == 0 || strcmp(argv[i], "--links") == 0) { @@ -310,15 +307,14 @@ static int parse_args(Config* config, int argc, char* argv[], int* positional_ar } else if (strcmp(argv[i], "-i") == 0 || strcmp(argv[i], "--itemize-changes") == 0) { config->itemize_changes = true; } else if (strcmp(argv[i], "--out-format") == 0 && i + 1 < argc) { - char* dup = str_dup(argv[++i]); - if (!dup) + if (set_string_option(&config->out_format, argv[++i], "--out-format") != 0) return -1; - free(config->out_format); - config->out_format = dup; } else if (strcmp(argv[i], "--info") == 0 && i + 1 < argc) { - config->info_level = atoi(argv[++i]); + if (set_nonneg_int_option(&config->info_level, argv[++i], "--info") != 0) + return -1; } else if (strcmp(argv[i], "--debug") == 0 && i + 1 < argc) { - config->debug_level = atoi(argv[++i]); + if (set_nonneg_int_option(&config->debug_level, argv[++i], "--debug") != 0) + return -1; } else if (strcmp(argv[i], "--list-only") == 0) { config->list_only = true; } else if (strcmp(argv[i], "-h") == 0 || strcmp(argv[i], "--human-readable") == 0) { @@ -336,12 +332,8 @@ static int parse_args(Config* config, int argc, char* argv[], int* positional_ar } else if (strcmp(argv[i], "--delete-after") == 0) { config->delete_after = true; } else if (strcmp(argv[i], "--max-delete") == 0 && i + 1 < argc) { - int val; - if (!parse_nonneg_int(argv[++i], &val)) { - fprintf(stderr, "Error: --max-delete must be a non-negative integer\n"); + if (set_nonneg_int_option(&config->max_delete, argv[++i], "--max-delete") != 0) return -1; - } - config->max_delete = val; } else if (strcmp(argv[i], "--filter") == 0 && i + 1 < argc) { if (!config->filters) config->filters = array_list_create(free); @@ -350,11 +342,8 @@ static int parse_args(Config* config, int argc, char* argv[], int* positional_ar return -1; array_list_add(config->filters, dup); } else if (strcmp(argv[i], "--files-from") == 0 && i + 1 < argc) { - char* dup = str_dup(argv[++i]); - if (!dup) + if (set_string_option(&config->files_from, argv[++i], "--files-from") != 0) return -1; - free(config->files_from); - config->files_from = dup; } else if (strcmp(argv[i], "--cvs-exclude") == 0) { config->cvs_exclude = true; } else if (strcmp(argv[i], "--prune-empty-dirs") == 0) { @@ -363,45 +352,27 @@ static int parse_args(Config* config, int argc, char* argv[], int* positional_ar config->relative = true; } else if (strcmp(argv[i], "-e") == 0 || strcmp(argv[i], "--rsh") == 0) { if (i + 1 < argc) { - char* dup = str_dup(argv[++i]); - if (!dup) + if (set_string_option(&config->rsh_command, argv[++i], "-e/--rsh") != 0) return -1; - free(config->rsh_command); - config->rsh_command = dup; } else { fprintf(stderr, "Error: -e/--rsh requires a command argument\n"); return -1; } } else if (strcmp(argv[i], "--rsync-path") == 0 && i + 1 < argc) { - char* dup = str_dup(argv[++i]); - if (!dup) + if (set_string_option(&config->rsync_path, argv[++i], "--rsync-path") != 0) return -1; - free(config->rsync_path); - config->rsync_path = dup; } else if (strcmp(argv[i], "--temp-dir") == 0 && i + 1 < argc) { - char* dup = str_dup(argv[++i]); - if (!dup) + if (set_string_option(&config->temp_dir, argv[++i], "--temp-dir") != 0) return -1; - free(config->temp_dir); - config->temp_dir = dup; } else if (strcmp(argv[i], "--compare-dest") == 0 && i + 1 < argc) { - char* dup = str_dup(argv[++i]); - if (!dup) + if (set_string_option(&config->compare_dest, argv[++i], "--compare-dest") != 0) return -1; - free(config->compare_dest); - config->compare_dest = dup; } else if (strcmp(argv[i], "--copy-dest") == 0 && i + 1 < argc) { - char* dup = str_dup(argv[++i]); - if (!dup) + if (set_string_option(&config->copy_dest, argv[++i], "--copy-dest") != 0) return -1; - free(config->copy_dest); - config->copy_dest = dup; } else if (strcmp(argv[i], "--link-dest") == 0 && i + 1 < argc) { - char* dup = str_dup(argv[++i]); - if (!dup) + if (set_string_option(&config->link_dest, argv[++i], "--link-dest") != 0) return -1; - free(config->link_dest); - config->link_dest = dup; } else if (argv[i][0] == '-') { fprintf(stderr, "Unknown option: %s\n", argv[i]); print_usage(); From 1064ae6c192bdf9bac0d0e35a1f587473809abe9 Mon Sep 17 00:00:00 2001 From: TapTap Date: Sat, 1 Aug 2026 18:42:02 +0200 Subject: [PATCH 2/4] fix: address PR review issues in client_cli.c - Fix FILE* leak when --log-file specified multiple times - Add strtoull endptr validation for --delta-block/--delta-max - Add range validation (1-22) for -c compression level --- src/client/client_cli.c | 22 ++++++++++++++++++++-- 1 file changed, 20 insertions(+), 2 deletions(-) diff --git a/src/client/client_cli.c b/src/client/client_cli.c index 8e9f844..5a15301 100644 --- a/src/client/client_cli.c +++ b/src/client/client_cli.c @@ -153,13 +153,25 @@ static int parse_args(Config* config, int argc, char* argv[], int* positional_ar } else if (strcmp(argv[i], "--delta") == 0) { config->use_delta = true; } else if (strcmp(argv[i], "--delta-block") == 0 && i + 1 < argc) { - unsigned long long val = strtoull(argv[++i], NULL, 10); + char* end; + errno = 0; + unsigned long long val = strtoull(argv[++i], &end, 10); + if (errno != 0 || *end != '\0') { + fprintf(stderr, "Error: --delta-block must be a positive integer\n"); + return -1; + } if (val >= DELTA_BLOCK_SIZE_MIN && val <= DELTA_BLOCK_SIZE_MAX) config->delta_block_size = (uint32_t)val; else fprintf(stderr, "Warning: --delta-block value %llu out of range, using default\n", val); } else if (strcmp(argv[i], "--delta-max") == 0 && i + 1 < argc) { - unsigned long long val = strtoull(argv[++i], NULL, 10); + char* end; + errno = 0; + unsigned long long val = strtoull(argv[++i], &end, 10); + if (errno != 0 || *end != '\0') { + fprintf(stderr, "Error: --delta-max must be a positive integer\n"); + return -1; + } if (val >= DELTA_MIN_FILE_SIZE) config->delta_max_file_size = val; else @@ -171,6 +183,10 @@ static int parse_args(Config* config, int argc, char* argv[], int* positional_ar char* end_ptr; long level = strtol(argv[i + 1], &end_ptr, 10); if (*end_ptr == '\0') { + if (level < 1 || level > 22) { + fprintf(stderr, "Error: compression level must be 1-22\n"); + return -1; + } config->compression_level = (int)level; log_message(LOG_LEVEL_INFO, "Set Compression level to %ld", level); i++; @@ -265,6 +281,8 @@ static int parse_args(Config* config, int argc, char* argv[], int* positional_ar fprintf(stderr, "Error: could not open log file '%s': %s\n", argv[i], strerror(errno)); return -1; } + if (config->log_file) + fclose(config->log_file); config->log_file = lf; log_set_file(lf); } else if (strcmp(argv[i], "--queue-size") == 0 && i + 1 < argc) { From 05a74770cad6ddb3a8ae0027f0722a72161ed7a1 Mon Sep 17 00:00:00 2001 From: TapTap Date: Sun, 2 Aug 2026 09:20:10 +0200 Subject: [PATCH 3/4] fix: address PR #200 review issues - Restore PROTOCOL_VERSION to a forward-compatible 2.1.0 and document wire format - Add NULL guard to config_delete - Add pipeline cancellation flag and cancellation-aware queue helper - Join running threads before destroying pipeline contexts on creation failure - Fix NULL dereference and memory leaks in manifest/chunk handling - Fix TLS/TCP socket fd leak on connect error paths - Add compression-level range validation (1-22) - Close previous log file before opening a new one - Use getline for unbounded pattern-file lines - Fix thread-unsafe localtime() and add log level bounds check - Fix file_load_data to clean up data on read size mismatch - Add hard ceiling to decompression buffer growth - Fix mkdir_r bounds check and restore glob comments - Add send_str NULL guard and mutex-protect bandwidth limiter - Update AGENTS.md for per-thread io_ssl contract --- AGENTS.md | 2 +- src/client/client_cli.c | 20 +++++-- src/client/client_send.c | 97 +++++++++++++++++++++++++------ src/client/scanner.c | 49 +++++++++++++--- src/server/server.c | 22 ++++++- src/shared/compression.c | 23 +++++++- src/shared/config.c | 13 +++++ src/shared/config.h | 3 +- src/shared/file.c | 4 +- src/shared/log.c | 14 +++-- src/shared/multiprocessing.c | 107 +++++++++++++++++++++++++++-------- src/shared/multiprocessing.h | 2 + src/shared/protocol.c | 15 +++++ src/shared/queue.c | 16 ++++++ src/shared/queue.h | 3 + src/shared/utils.c | 22 ++++++- 16 files changed, 344 insertions(+), 68 deletions(-) diff --git a/AGENTS.md b/AGENTS.md index 042ab34..9fc41ec 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -206,7 +206,7 @@ tea pr close --repo TapTap/FastSync ## Common pitfalls -- **`__thread` on shared SSL context**: io_ssl must NOT be thread-local — worker threads inherit the SSL context from the main thread. Use regular `static SSL* io_ssl`. +- **Per-thread SSL context**: `io_ssl` is stored per-thread (`static __thread SSL* io_ssl`). Each thread that performs protocol I/O must call `io_set_ssl()` to install its own SSL object before using `send_*` / `receive_*` primitives. The main thread's SSL context is not automatically inherited by worker threads. - **SSL WANT_READ/WANT_WRITE retry**: Always retry on `SSL_ERROR_WANT_READ` and `SSL_ERROR_WANT_WRITE` in `send_n_data`/`receive_n_data`. Removing these breaks TLS multithreaded transfers. - **clang-format version**: The CI image uses clang-format 18. Always format inside the CI Docker container for exact match. - **Merge order matters**: Merge the most comprehensive branch first, then smaller ones, to minimize conflicts when creating a combined branch. diff --git a/src/client/client_cli.c b/src/client/client_cli.c index 5a15301..5ed58bb 100644 --- a/src/client/client_cli.c +++ b/src/client/client_cli.c @@ -276,13 +276,16 @@ static int parse_args(Config* config, int argc, char* argv[], int* positional_ar if (set_nonneg_int_option(&config->max_depth, argv[++i], "--max-depth") != 0) return -1; } else if (strcmp(argv[i], "--log-file") == 0 && i + 1 < argc) { + if (config->log_file) { + fclose(config->log_file); + config->log_file = NULL; + log_set_file(NULL); + } FILE* lf = fopen(argv[++i], "a"); if (!lf) { fprintf(stderr, "Error: could not open log file '%s': %s\n", argv[i], strerror(errno)); return -1; } - if (config->log_file) - fclose(config->log_file); config->log_file = lf; log_set_file(lf); } else if (strcmp(argv[i], "--queue-size") == 0 && i + 1 < argc) { @@ -427,6 +430,10 @@ static int parse_args(Config* config, int argc, char* argv[], int* positional_ar } else if (strcmp(argv[i], "--compress-level") == 0 && i + 1 < argc) { if (set_positive_int_option(&config->compression_level, argv[++i], "--compress-level") != 0) return -1; + if (config->compression_level < 1 || config->compression_level > 22) { + fprintf(stderr, "Error: --compress-level must be between 1 and 22\n"); + return -1; + } } else if (argv[i][0] == '-') { fprintf(stderr, "Unknown option: %s\n", argv[i]); print_usage(); @@ -600,8 +607,10 @@ static int read_patterns_from_file(const char* filepath, char*** patterns, int* fprintf(stderr, "Error: could not open pattern file '%s': %s\n", filepath, strerror(errno)); return -1; } - char line[4096]; - while (fgets(line, sizeof(line), fp)) { + char* line = NULL; + size_t line_size = 0; + ssize_t n; + while ((n = getline(&line, &line_size, fp)) != -1) { char* p = line; while (*p == ' ' || *p == '\t') p++; @@ -615,6 +624,7 @@ static int read_patterns_from_file(const char* filepath, char*** patterns, int* char** tmp = realloc(*patterns, (*count + 1) * sizeof(char*)); if (!tmp) { fprintf(stderr, "Error: memory allocation failed for pattern file\n"); + free(line); fclose(fp); return -1; } @@ -622,11 +632,13 @@ static int read_patterns_from_file(const char* filepath, char*** patterns, int* char* dup = str_dup(p); if (!dup) { fprintf(stderr, "Error: memory allocation failed for pattern file\n"); + free(line); fclose(fp); return -1; } (*patterns)[(*count)++] = dup; } + free(line); fclose(fp); return 0; } diff --git a/src/client/client_send.c b/src/client/client_send.c index 620cb57..16e0e66 100644 --- a/src/client/client_send.c +++ b/src/client/client_send.c @@ -346,6 +346,7 @@ static int send_chunks_multithreaded(void* pipeline_context) { } if (send_chunk(client, current_chunk, context->config) != 0) { fprintf(stderr, "Error: unexpected error while sending chunk\n"); + chunk_destroy(current_chunk); client_disconnect(client); client_delete(client); mtx_lock(&context->mutex_progress); @@ -385,13 +386,31 @@ static int scan_directory_multithreaded(void* pipeline_context) { const char* p = current_chunk->items[i]->path; if (*p == '/') p++; - array_list_add(context->manifest, str_dup(p)); + char* manifest_entry = str_dup(p); + if (!manifest_entry) { + log_message(LOG_LEVEL_ERROR, "Failed to allocate manifest entry"); + mtx_unlock(&context->mutex_scanner); + context->cancelled = true; + cnd_broadcast(&context->condition_not_full_scanner); + cnd_broadcast(&context->condition_not_empty_scanner); + parallel_scanner_destroy(scanner); + return thrd_error; + } + array_list_add(context->manifest, manifest_entry); } mtx_unlock(&context->mutex_scanner); } - queue_enqueue_multithreaded(context->queue_scanner, current_chunk, &context->mutex_scanner, - &context->condition_not_empty_scanner, - &context->condition_not_full_scanner); + if (!queue_enqueue_multithreaded_cancel( + context->queue_scanner, current_chunk, &context->mutex_scanner, + &context->condition_not_empty_scanner, &context->condition_not_full_scanner, + &context->cancelled)) { + chunk_destroy(current_chunk); + context->cancelled = true; + cnd_broadcast(&context->condition_not_full_scanner); + cnd_broadcast(&context->condition_not_empty_scanner); + parallel_scanner_destroy(scanner); + return thrd_error; + } } mtx_lock(&context->mutex_scanner); context->scanner_done = true; @@ -427,9 +446,16 @@ static int load_files_multithreaded(void* pipeline_context) { } } } - queue_enqueue_multithreaded(context->queue_loader, chunk, &context->mutex_loader, - &context->condition_not_empty_loader, - &context->condition_not_full_loader); + if (!queue_enqueue_multithreaded_cancel(context->queue_loader, chunk, &context->mutex_loader, + &context->condition_not_empty_loader, + &context->condition_not_full_loader, + &context->cancelled)) { + chunk_destroy(chunk); + context->cancelled = true; + cnd_broadcast(&context->condition_not_full_loader); + cnd_broadcast(&context->condition_not_empty_loader); + return thrd_error; + } } } @@ -487,16 +513,20 @@ int send_files(Config* config) { client = client_create(); if (!client || !client_connect_tls(client, config->server_host, config->server_port, config->tls_cert, config->tls_key, config->tls_ca)) { - if (client) + if (client) { + client_disconnect(client); client_delete(client); + } fprintf(stderr, "Error: could not connect to server via TLS\n"); return 1; } } else { client = client_create(); if (!client || !client_connect(client, config->server_host, config->server_port)) { - if (client) + if (client) { + client_disconnect(client); client_delete(client); + } fprintf(stderr, "Error: could not connect to server\n"); return 1; } @@ -526,7 +556,17 @@ int send_files(Config* config) { const char* p = current_chunk->items[i]->path; if (*p == '/') p++; - array_list_add(manifest, str_dup(p)); + char* manifest_entry = str_dup(p); + if (!manifest_entry) { + log_message(LOG_LEVEL_ERROR, "Failed to allocate manifest entry"); + chunk_destroy(current_chunk); + array_list_delete(manifest); + directory_scanner_destroy(scanner); + client_disconnect(client); + client_delete(client); + return 1; + } + array_list_add(manifest, manifest_entry); } } if (!config->use_sendfile) { @@ -625,17 +665,42 @@ int send_files_multithreaded(Config* config) { if (config->use_delete) context->manifest = array_list_create(free); - thrd_t scanner, loader, sender, progress; - if (thrd_create(&scanner, scan_directory_multithreaded, context) != thrd_success || - thrd_create(&loader, load_files_multithreaded, context) != thrd_success || - thrd_create(&sender, send_chunks_multithreaded, context) != thrd_success) { + thrd_t scanner, loader, sender; + bool scanner_created = false; + bool loader_created = false; + bool sender_created = false; + + scanner_created = (thrd_create(&scanner, scan_directory_multithreaded, context) == thrd_success); + if (scanner_created) + loader_created = (thrd_create(&loader, load_files_multithreaded, context) == thrd_success); + if (scanner_created && loader_created) + sender_created = (thrd_create(&sender, send_chunks_multithreaded, context) == thrd_success); + + if (!scanner_created || !loader_created || !sender_created) { perror("Error creating threads.\n"); + context->cancelled = true; + context->scanner_done = true; + context->loader_done = true; + context->sender_done = true; + cnd_broadcast(&context->condition_not_full_scanner); + cnd_broadcast(&context->condition_not_empty_scanner); + cnd_broadcast(&context->condition_not_full_loader); + cnd_broadcast(&context->condition_not_empty_loader); + if (sender_created) + thrd_join(sender, NULL); + if (loader_created) + thrd_join(loader, NULL); + if (scanner_created) + thrd_join(scanner, NULL); pipeline_context_sender_destroy(context); return 1; } + thrd_t progress; + bool progress_created = false; if (config->show_progress) { - if (thrd_create(&progress, progress_thread_fn, context) != thrd_success) { + progress_created = (thrd_create(&progress, progress_thread_fn, context) == thrd_success); + if (!progress_created) { perror("Error creating progress thread.\n"); /* Non-fatal; continue without progress reporting */ } @@ -646,7 +711,7 @@ int send_files_multithreaded(Config* config) { thrd_join(loader, NULL); thrd_join(sender, &sender_result); - if (config->show_progress) { + if (progress_created) { /* Signal progress thread to exit if it hasn't already */ mtx_lock(&context->mutex_progress); context->sender_done = true; diff --git a/src/client/scanner.c b/src/client/scanner.c index dfaa50e..848f66d 100644 --- a/src/client/scanner.c +++ b/src/client/scanner.c @@ -27,10 +27,14 @@ static void dir_entry_destroy(void* item) { static DirEntry* dir_entry_create(const char* path, int depth) { DirEntry* de = malloc(sizeof(DirEntry)); - if (de) { - de->path = str_dup(path); - de->depth = depth; + if (!de) + return NULL; + de->path = str_dup(path); + if (!de->path) { + free(de); + return NULL; } + de->depth = depth; return de; } @@ -62,7 +66,18 @@ DirectoryScanner* directory_scanner_create(const char* root_directory, bool use_ scanner->safe_links = safe_links; scanner->copy_unsafe_links = copy_unsafe_links; scanner->checksum = checksum; - queue_enqueue(scanner->directories, dir_entry_create(root_directory, 0)); + DirEntry* root = dir_entry_create(root_directory, 0); + if (!root) { + queue_destroy(scanner->directories); + free(scanner); + return NULL; + } + if (!queue_enqueue(scanner->directories, root)) { + dir_entry_destroy(root); + queue_destroy(scanner->directories); + free(scanner); + return NULL; + } return scanner; } @@ -323,9 +338,28 @@ ParallelScanner* parallel_scanner_create(char* root_directory, bool use_metadata free(ps); return NULL; } - if (mtx_init(&ps->result_mutex, mtx_plain) != thrd_success || - cnd_init(&ps->result_not_empty) != thrd_success || - cnd_init(&ps->result_not_full) != thrd_success) { + int init = 0; + bool ok = true; + if (mtx_init(&ps->result_mutex, mtx_plain) != thrd_success) + ok = false; + if (ok) { + init++; + if (cnd_init(&ps->result_not_empty) != thrd_success) + ok = false; + } + if (ok) { + // cppcheck-suppress unreadVariable + init++; + if (cnd_init(&ps->result_not_full) != thrd_success) + ok = false; + } + if (!ok) { + if (init >= 3) + cnd_destroy(&ps->result_not_full); + if (init >= 2) + cnd_destroy(&ps->result_not_empty); + if (init >= 1) + mtx_destroy(&ps->result_mutex); queue_destroy(ps->result_queue); free(ps); return NULL; @@ -564,6 +598,7 @@ void parallel_scanner_destroy(ParallelScanner* ps) { return; ps->done = true; cnd_signal(&ps->result_not_empty); + cnd_broadcast(&ps->result_not_full); for (int i = 0; i < ps->num_threads; i++) thrd_join(ps->threads[i], NULL); free(ps->threads); diff --git a/src/server/server.c b/src/server/server.c index f8fdb54..1827d6c 100644 --- a/src/server/server.c +++ b/src/server/server.c @@ -135,11 +135,27 @@ void handler(int file_descriptor) { return; } thrd_t receiver, writer; - if (thrd_create(&receiver, receive_thread, context) != thrd_success || - thrd_create(&writer, write_thread, context) != thrd_success) { + bool receiver_created = false; + bool writer_created = false; + receiver_created = (thrd_create(&receiver, receive_thread, context) == thrd_success); + if (receiver_created) + writer_created = (thrd_create(&writer, write_thread, context) == thrd_success); + if (!receiver_created || !writer_created) { perror("Error creating Threads"); + if (receiver_created) { + mtx_lock(&context->mutex); + context->cancelled = true; + cnd_broadcast(&context->condition_not_full); + cnd_broadcast(&context->condition_not_empty); + mtx_unlock(&context->mutex); + close(file_descriptor); + thrd_join(receiver, NULL); + } else { + close(file_descriptor); + } + if (writer_created) + thrd_join(writer, NULL); pipeline_context_receiver_destroy(context); - close(file_descriptor); return; } thrd_join(receiver, NULL); diff --git a/src/shared/compression.c b/src/shared/compression.c index c1eb1f1..27b6da4 100644 --- a/src/shared/compression.c +++ b/src/shared/compression.c @@ -1,12 +1,13 @@ #include "compression.h" #include "data.h" #include "log.h" -#include "stdlib.h" -#include "string.h" +#include +#include #include -#include "zstd.h" +#include #define INITIAL_DECOMPRESS_BUF_SIZE (1024 * 1024) +#define MAX_DECOMPRESSED_SIZE (100ULL * 1024 * 1024) /* 100 MB hard ceiling */ static const char* SKIP_COMPRESSION_EXTENSIONS[] = {".jpg", ".jpeg", ".png", ".gif", ".mp4", ".mkv", ".zip", ".gz", ".xz", ".zst", NULL}; @@ -84,6 +85,13 @@ Data* data_decompress(Data* compressed_data) { dst_size = compressed_data->size * 3; if (dst_size < INITIAL_DECOMPRESS_BUF_SIZE) dst_size = INITIAL_DECOMPRESS_BUF_SIZE; + if (dst_size > MAX_DECOMPRESSED_SIZE) + dst_size = MAX_DECOMPRESSED_SIZE; + } + if (dst_size > MAX_DECOMPRESSED_SIZE) { + log_message(LOG_LEVEL_ERROR, "Declared decompressed size exceeds %llu bytes", + (unsigned long long)MAX_DECOMPRESSED_SIZE); + return NULL; } ZSTD_DCtx* dctx = ZSTD_createDCtx(); @@ -113,7 +121,16 @@ Data* data_decompress(Data* compressed_data) { return NULL; } if (ret > 0 && output.pos == output.size) { + if (buf_size >= MAX_DECOMPRESSED_SIZE) { + log_message(LOG_LEVEL_ERROR, "Decompressed data exceeds %llu bytes", + (unsigned long long)MAX_DECOMPRESSED_SIZE); + ZSTD_freeDCtx(dctx); + data_destroy(uncompressed_data); + return NULL; + } buf_size *= 2; + if (buf_size > MAX_DECOMPRESSED_SIZE) + buf_size = MAX_DECOMPRESSED_SIZE; void* new_data = realloc(uncompressed_data->data, buf_size); if (!new_data) { log_message(LOG_LEVEL_ERROR, "Failed to grow decompression buffer"); diff --git a/src/shared/config.c b/src/shared/config.c index 6d52cc4..c7900cc 100644 --- a/src/shared/config.c +++ b/src/shared/config.c @@ -131,6 +131,8 @@ void config_parse_ssh_dest(Config* config) { } void config_delete(Config* config) { + if (config == NULL) + return; free(config->version); free(config->send_directory); free(config->receive_root_directory); @@ -167,6 +169,16 @@ void config_delete(Config* config) { free(config); } +/* Wire format order (must match config_receive and be updated when PROTOCOL_VERSION bumps): + * version, send_directory, receive_root_directory, save_to_disk, use_multithreading, + * use_chunk_serialization, use_compression, use_metadata, compression_level, chunk_size, + * use_sendfile, use_delete, use_incremental, use_delta, delta_block_size, delta_max_file_size, + * backup, backup_dir, follow_symlinks, copy_links, safe_links, copy_unsafe_links, + * preserve_hard_links, preserve_acls, preserve_xattrs, preserve_devices, preserve_sparse, + * update, inplace, append, append_verify, delete_excluded, delete_after, max_delete, relative, + * prune_empty_dirs, temp_dir, partial, partial_dir, suffix, delete_before, checksum, + * compress_choice, status + */ bool config_send(int file_descriptor, const Config* config) { if (!send_str(file_descriptor, config->version)) return false; @@ -264,6 +276,7 @@ bool config_send(int file_descriptor, const Config* config) { return true; } +/* Wire format order: see the comment above config_send. */ Config* config_receive(int file_descriptor) { Config* config = (Config*)malloc(sizeof(Config)); if (config == NULL) diff --git a/src/shared/config.h b/src/shared/config.h index 9c78505..22359e4 100644 --- a/src/shared/config.h +++ b/src/shared/config.h @@ -128,7 +128,8 @@ typedef struct Config { char* compress_choice; } Config; -#define PROTOCOL_VERSION "1.3.0" +/* This version must be bumped whenever config_send / config_receive wire format changes. */ +#define PROTOCOL_VERSION "2.1.0" #define DEFAULT_CHUNK_SIZE (10 * 1024 * 1024) Config* config_create(void); diff --git a/src/shared/file.c b/src/shared/file.c index 2b730f8..488321c 100644 --- a/src/shared/file.c +++ b/src/shared/file.c @@ -16,7 +16,6 @@ #include "config.h" #include "data.h" #include "file.h" -#include "log.h" #include "metadata.h" #include "protocol.h" #include "utils.h" @@ -96,6 +95,9 @@ bool file_load_data(File* file) { size_t bytes_read = file_content_to_buffer(file); if (bytes_read != file->data->size) { log_message(LOG_LEVEL_ERROR, "Did not read expected amount of bytes from file"); + free(file->data->data); + file->data->data = NULL; + file->data->size = 0; return false; } return true; diff --git a/src/shared/log.c b/src/shared/log.c index 07a5d27..fbd9400 100644 --- a/src/shared/log.c +++ b/src/shared/log.c @@ -18,11 +18,15 @@ void log_set_file(FILE* fp) { void log_message(LogLevel log_level, const char* format, ...) { if (log_level < current_log_level) return; + if (log_level < 0 || log_level >= (int)(sizeof(log_level_strings) / sizeof(log_level_strings[0]))) + return; time_t now = time(NULL); - const struct tm* t = localtime(&now); + struct tm t; + if (!localtime_r(&now, &t)) + return; - fprintf(stderr, "%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]); + fprintf(stderr, "%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]); va_list args; va_start(args, format); @@ -31,8 +35,8 @@ void log_message(LogLevel log_level, const char* format, ...) { fprintf(stderr, "\n"); if (log_fp) { - fprintf(log_fp, "%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]); + fprintf(log_fp, "%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]); va_start(args, format); vfprintf(log_fp, format, args); va_end(args); diff --git a/src/shared/multiprocessing.c b/src/shared/multiprocessing.c index df03388..24a4a5c 100644 --- a/src/shared/multiprocessing.c +++ b/src/shared/multiprocessing.c @@ -26,18 +26,48 @@ PipelineContextSender* pipeline_context_sender_create(Config* config, Queue* que context->manifest = NULL; context->progress_bytes = 0; context->sender_done = false; - if (mtx_init(&context->mutex_scanner, mtx_plain) != thrd_success || - cnd_init(&context->condition_not_full_scanner) != thrd_success || - cnd_init(&context->condition_not_empty_scanner) != thrd_success || - mtx_init(&context->mutex_loader, mtx_plain) != thrd_success || - cnd_init(&context->condition_not_full_loader) != thrd_success || - mtx_init(&context->mutex_progress, mtx_plain) != thrd_success || - cnd_init(&context->condition_not_empty_loader) != thrd_success) { - perror("Error initializing synchronization objects"); - free(context); - return NULL; - } + context->cancelled = false; + int init = 0; + if (mtx_init(&context->mutex_scanner, mtx_plain) != thrd_success) + goto fail; + init++; + if (cnd_init(&context->condition_not_full_scanner) != thrd_success) + goto fail; + init++; + if (cnd_init(&context->condition_not_empty_scanner) != thrd_success) + goto fail; + init++; + if (mtx_init(&context->mutex_loader, mtx_plain) != thrd_success) + goto fail; + init++; + if (cnd_init(&context->condition_not_full_loader) != thrd_success) + goto fail; + init++; + if (cnd_init(&context->condition_not_empty_loader) != thrd_success) + goto fail; + init++; + if (mtx_init(&context->mutex_progress, mtx_plain) != thrd_success) + goto fail; + // cppcheck-suppress unreadVariable + init++; return context; + +fail: + perror("Error initializing synchronization objects"); + if (init >= 6) + cnd_destroy(&context->condition_not_empty_loader); + if (init >= 5) + cnd_destroy(&context->condition_not_full_loader); + if (init >= 4) + mtx_destroy(&context->mutex_loader); + if (init >= 3) + cnd_destroy(&context->condition_not_empty_scanner); + if (init >= 2) + cnd_destroy(&context->condition_not_full_scanner); + if (init >= 1) + mtx_destroy(&context->mutex_scanner); + free(context); + return NULL; } void pipeline_context_sender_destroy(PipelineContextSender* context) { @@ -67,14 +97,30 @@ PipelineContextReceiver* pipeline_context_receiver_create(Config* config, Queue* context->file_descriptor = file_descriptor; context->ssl = ssl; context->receiver_done = false; - if (mtx_init(&context->mutex, mtx_plain) != thrd_success || - cnd_init(&context->condition_not_full) != thrd_success || - cnd_init(&context->condition_not_empty) != thrd_success) { - perror("Error initializing synchronization objects"); - free(context); - return NULL; - } + context->cancelled = false; + int init = 0; + if (mtx_init(&context->mutex, mtx_plain) != thrd_success) + goto fail; + init++; + if (cnd_init(&context->condition_not_full) != thrd_success) + goto fail; + init++; + if (cnd_init(&context->condition_not_empty) != thrd_success) + goto fail; + // cppcheck-suppress unreadVariable + init++; return context; + +fail: + perror("Error initializing synchronization objects"); + if (init >= 3) + cnd_destroy(&context->condition_not_empty); + if (init >= 2) + cnd_destroy(&context->condition_not_full); + if (init >= 1) + mtx_destroy(&context->mutex); + free(context); + return NULL; } void pipeline_context_receiver_destroy(PipelineContextReceiver* context) { @@ -94,8 +140,13 @@ static bool receive_chunk_enqueue(int file_descriptor, PipelineContextReceiver* for (int i = 0; i < chunk->element_count; i++) { File* file = chunk->items[i]; chunk->items[i] = NULL; - queue_enqueue_multithreaded(context->queue, file, &context->mutex, - &context->condition_not_empty, &context->condition_not_full); + if (!queue_enqueue_multithreaded_cancel(context->queue, file, &context->mutex, + &context->condition_not_empty, + &context->condition_not_full, &context->cancelled)) { + file_destroy(file); + chunk_destroy(chunk); + return false; + } } chunk_destroy(chunk); return true; @@ -129,8 +180,12 @@ int receive_thread(void* pipeline_context) { if (!skipped) { if (file == NULL) return thrd_error; - queue_enqueue_multithreaded(context->queue, file, &context->mutex, - &context->condition_not_empty, &context->condition_not_full); + if (!queue_enqueue_multithreaded_cancel( + context->queue, file, &context->mutex, &context->condition_not_empty, + &context->condition_not_full, &context->cancelled)) { + file_destroy(file); + return thrd_error; + } } } else if (status == STATUS_CHUNK) { if (!receive_chunk_enqueue(file_descriptor, context)) @@ -166,8 +221,12 @@ int receive_thread(void* pipeline_context) { } else { File* file = file_receive(config, file_descriptor); if (file) { - queue_enqueue_multithreaded(context->queue, file, &context->mutex, - &context->condition_not_empty, &context->condition_not_full); + if (!queue_enqueue_multithreaded_cancel( + context->queue, file, &context->mutex, &context->condition_not_empty, + &context->condition_not_full, &context->cancelled)) { + file_destroy(file); + return thrd_error; + } } else { log_message(LOG_LEVEL_ERROR, "Failed to receive file"); return thrd_error; diff --git a/src/shared/multiprocessing.h b/src/shared/multiprocessing.h index aa53a13..443720a 100644 --- a/src/shared/multiprocessing.h +++ b/src/shared/multiprocessing.h @@ -26,6 +26,7 @@ typedef struct { mtx_t mutex_progress; unsigned long long progress_bytes; bool sender_done; + bool cancelled; } PipelineContextSender; typedef struct PipelineContextReceiver { @@ -37,6 +38,7 @@ typedef struct PipelineContextReceiver { cnd_t condition_not_full; cnd_t condition_not_empty; bool receiver_done; + bool cancelled; } PipelineContextReceiver; PipelineContextSender* pipeline_context_sender_create(Config* config, Queue* queue_scanner, diff --git a/src/shared/protocol.c b/src/shared/protocol.c index 0b32f96..0d41808 100644 --- a/src/shared/protocol.c +++ b/src/shared/protocol.c @@ -6,6 +6,7 @@ #include #include #include +#include #include #include @@ -20,6 +21,8 @@ static __thread SSL* io_ssl; static unsigned long long io_bwlimit = 0; static long long bw_tokens = 0; static struct timespec bw_last_refill = {0, 0}; +static mtx_t bw_mutex; +static once_flag bw_mutex_once = ONCE_FLAG_INIT; static __thread unsigned long long total_allocated_bytes = 0; @@ -28,15 +31,24 @@ void io_set_fds(int read_fd, int write_fd) { io_write_fd = write_fd; } +static void bw_mutex_init(void) { + mtx_init(&bw_mutex, mtx_plain); +} + void io_set_bwlimit(unsigned long long bytes_per_sec) { + call_once(&bw_mutex_once, bw_mutex_init); + mtx_lock(&bw_mutex); io_bwlimit = bytes_per_sec; bw_tokens = (long long)io_bwlimit; clock_gettime(CLOCK_MONOTONIC, &bw_last_refill); + mtx_unlock(&bw_mutex); } static void bw_throttle(size_t bytes_written) { if (io_bwlimit == 0) return; + call_once(&bw_mutex_once, bw_mutex_init); + mtx_lock(&bw_mutex); struct timespec now; clock_gettime(CLOCK_MONOTONIC, &now); @@ -61,6 +73,7 @@ static void bw_throttle(size_t bytes_written) { bw_tokens = 0; clock_gettime(CLOCK_MONOTONIC, &bw_last_refill); } + mtx_unlock(&bw_mutex); } void io_set_ssl(SSL* ssl) { @@ -177,6 +190,8 @@ static const char* status_to_string(Status status) { } bool send_str(int file_descriptor, const char* data) { + if (data == NULL) + return false; size_t size = strlen(data); if (!send_n_data(file_descriptor, &size, sizeof(size_t))) return false; diff --git a/src/shared/queue.c b/src/shared/queue.c index 14d901d..1213746 100644 --- a/src/shared/queue.c +++ b/src/shared/queue.c @@ -103,6 +103,22 @@ bool queue_enqueue_multithreaded(Queue* queue, void* item, mtx_t* mutex, cnd_t* return ok; } +bool queue_enqueue_multithreaded_cancel(Queue* queue, void* item, mtx_t* mutex, + cnd_t* condition_not_empty, cnd_t* condition_not_full, + const bool* cancelled) { + mtx_lock(mutex); + while (queue_is_full(queue) && (cancelled == NULL || !*cancelled)) + cnd_wait(condition_not_full, mutex); + if (cancelled != NULL && *cancelled) { + mtx_unlock(mutex); + return false; + } + bool ok = queue_enqueue(queue, item); + cnd_signal(condition_not_empty); + mtx_unlock(mutex); + return ok; +} + void* queue_dequeue(Queue* queue) { if (queue == NULL || queue_is_empty(queue)) { perror("ERROR: Could not dequeue from null or empty queue."); diff --git a/src/shared/queue.h b/src/shared/queue.h index 8bd6ba4..f482d86 100644 --- a/src/shared/queue.h +++ b/src/shared/queue.h @@ -20,6 +20,9 @@ bool queue_is_full(const Queue* queue); bool queue_enqueue(Queue* queue, void* item); bool queue_enqueue_multithreaded(Queue* queue, void* item, mtx_t* mutex, cnd_t* condition_not_empty, cnd_t* condition_not_full); +bool queue_enqueue_multithreaded_cancel(Queue* queue, void* item, mtx_t* mutex, + cnd_t* condition_not_empty, cnd_t* condition_not_full, + const bool* cancelled); void* queue_dequeue(Queue* queue); void* queue_dequeue_multithreaded(Queue* queue, mtx_t* mutex, cnd_t* condition_not_empty, cnd_t* condition_not_full, const bool* other_thread_done); diff --git a/src/shared/utils.c b/src/shared/utils.c index dc7f0f9..43ed1e2 100644 --- a/src/shared/utils.c +++ b/src/shared/utils.c @@ -10,11 +10,13 @@ #include bool mkdir_r(const char* path) { - char* path_duplicate = malloc(strlen(path) + 1); + size_t path_len = strlen(path); + char* path_duplicate = malloc(path_len + 1); if (!path_duplicate) return false; - memcpy(path_duplicate, path, strlen(path) + 1); - char* path_current = (char*)malloc((strlen(path) + 2) * sizeof(char)); + memcpy(path_duplicate, path, path_len + 1); + size_t capacity = path_len + 2; + char* path_current = (char*)malloc(capacity * sizeof(char)); if (!path_current) { free(path_duplicate); return false; @@ -33,6 +35,10 @@ bool mkdir_r(const char* path) { bool ok = true; while (part != NULL) { size_t part_len = strlen(part); + if ((size_t)(path_current_position - path_current) + part_len + 2 > capacity) { + ok = false; + break; + } memcpy(path_current_position, part, part_len); path_current_position += part_len; path_current_position[0] = '/'; @@ -64,10 +70,18 @@ char* str_dup(const char* string) { return new_string; } +/* Match a glob pattern against a string. Supported wildcards: + * ? matches any single character except '/'. + * * matches any sequence of characters within one path component (no '/'). + * ** matches any sequence of characters, including '/' (cross-directory). + * slash-star-star-slash is treated as a cross-directory wildcard when it appears between + * literals. + */ bool glob_match(const char* pattern, const char* str) { while (*pattern) { if (*pattern == '*') { if (*(pattern + 1) == '*') { + /* globstar: match across directories */ pattern += 2; if (*pattern == '\0') return true; @@ -80,6 +94,7 @@ bool glob_match(const char* pattern, const char* str) { } return glob_match(pattern, str); } + /* single *: match within one path component */ pattern++; while (*str && *str != '/') { if (glob_match(pattern, str)) @@ -94,6 +109,7 @@ bool glob_match(const char* pattern, const char* str) { str++; } else { if (*pattern != *str) { + /* allow literal / ** / rest to match any number of directories */ if (*pattern == '/' && *(pattern + 1) == '*' && *(pattern + 2) == '*') { const char* rest = pattern + 3; if (*rest == '/') From 8743d7f5227af8bb83b47e0241f08225d92c98bc Mon Sep 17 00:00:00 2001 From: TapTap Date: Sat, 8 Aug 2026 20:01:48 +0200 Subject: [PATCH 4/4] fix: address PR #200 review issues - Fix memory leak in delta_deserialize() on BLOCK_MATCH error path - Fix memory leak on malloc failure for literal data - Add port range validation (1-65535) for -p/--port and --server-port - Restore -V/--version flag with usage text - Use MAX_DATA_PAYLOAD_SIZE consistently (remove local MAX_DATA_SIZE) - Fix integer truncation in config_send/config_receive for delta_block_size - Add log messages for malloc failures - Extract config_set_defaults() helper to eliminate duplication - Add 10 new CLI tests for argument parsing --- CMakeLists.txt | 3 +- src/client/client_cli.c | 22 +++++- src/shared/config.c | 89 +++++------------------ src/shared/delta.c | 9 +++ src/shared/protocol.c | 5 +- tests/test_client_cli.c | 153 ++++++++++++++++++++++++++++++++++++++++ 6 files changed, 202 insertions(+), 79 deletions(-) diff --git a/CMakeLists.txt b/CMakeLists.txt index db6a15f..42c2353 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -79,8 +79,9 @@ set(TEST_INCLUDES tests src/shared src/server src/client) # Monolithic test binary (backward compatible) file(GLOB TEST_SRCS "tests/test_*.c" "tests/runner.c") -add_executable(tests ${TEST_SRCS} ${SHARED_SRCS} src/client/scanner.c) +add_executable(tests ${TEST_SRCS} ${SHARED_SRCS} src/client/scanner.c src/client/client_cli.c) target_include_directories(tests PRIVATE ${TEST_INCLUDES}) +target_compile_definitions(tests PRIVATE FASTSYNC_TEST_BUILD) target_link_libraries(tests PRIVATE ${TEST_LIBS}) add_test(NAME unit_all COMMAND tests) diff --git a/src/client/client_cli.c b/src/client/client_cli.c index 5ed58bb..711a449 100644 --- a/src/client/client_cli.c +++ b/src/client/client_cli.c @@ -12,6 +12,7 @@ #include #include +#ifndef FASTSYNC_TEST_BUILD /* Parse environment variables for source/destination directories and save-to-disk flag. */ static void parse_environment(const char** out_env_source, const char** out_env_dest, bool* out_save_to_disk) { @@ -22,6 +23,7 @@ static void parse_environment(const char** out_env_source, const char** out_env_ if (env_save && (strcmp(env_save, "true") == 0 || strcmp(env_save, "1") == 0)) *out_save_to_disk = true; } +#endif /* Parse a string as a positive integer, returning true on success. */ static bool parse_positive_int(const char* s, int* out_val) { @@ -84,12 +86,15 @@ static void print_usage(void); static int read_patterns_from_file(const char* filepath, char*** patterns, int* count); /* Parse CLI arguments into config. Returns 0 on success, -1 on error, 1 for help/clean-exit. */ -static int parse_args(Config* config, int argc, char* argv[], int* positional_args, - int* positional_count) { +int parse_args(Config* config, int argc, char* argv[], int* positional_args, + int* positional_count) { for (int i = 1; i < argc; i++) { if (strcmp(argv[i], "--help") == 0) { print_usage(); return 1; + } else if (strcmp(argv[i], "-V") == 0 || strcmp(argv[i], "--version") == 0) { + printf("fastsync version %s\n", PROTOCOL_VERSION); + return 1; } else if (strcmp(argv[i], "-a") == 0 || strcmp(argv[i], "--archive") == 0) { config->use_compression = true; config->use_multithreading = true; @@ -102,6 +107,10 @@ static int parse_args(Config* config, int argc, char* argv[], int* positional_ar fprintf(stderr, "Error: invalid --port/-p value: %s\n", argv[i]); return -1; } + if (config->ssh_port > 65535) { + fprintf(stderr, "Error: SSH port must be 1-65535\n"); + return -1; + } } else if (strcmp(argv[i], "--delete") == 0) { config->use_delete = true; } else if (strcmp(argv[i], "--exclude") == 0 && i + 1 < argc) { @@ -220,6 +229,10 @@ static int parse_args(Config* config, int argc, char* argv[], int* positional_ar fprintf(stderr, "Error: invalid --server-port value: %s\n", argv[i]); return -1; } + if (config->server_port > 65535) { + fprintf(stderr, "Error: server port must be 1-65535\n"); + return -1; + } } else if (strcmp(argv[i], "--bwlimit") == 0 && i + 1 < argc) { char* end; errno = 0; @@ -451,6 +464,7 @@ static int parse_args(Config* config, int argc, char* argv[], int* positional_ar return 0; } +#ifndef FASTSYNC_TEST_BUILD /* Validate config after parsing. Returns true if valid. */ static bool validate_config(const Config* config) { if (!config->send_directory || !config->receive_root_directory) { @@ -491,6 +505,7 @@ static bool validate_config(const Config* config) { } return true; } +#endif /* FASTSYNC_TEST_BUILD */ static void print_usage(void) { printf("Usage:\n"); @@ -599,6 +614,7 @@ static void print_usage(void) { printf(" --copy-dest Copy destination\n"); printf(" --link-dest Link destination\n"); printf(" --help Show this help\n"); + printf(" -V, --version Show version\n"); } static int read_patterns_from_file(const char* filepath, char*** patterns, int* count) { @@ -643,6 +659,7 @@ static int read_patterns_from_file(const char* filepath, char*** patterns, int* return 0; } +#ifndef FASTSYNC_TEST_BUILD int main(int argc, char* argv[]) { const char* env_source = NULL; const char* env_dest = NULL; @@ -748,3 +765,4 @@ cleanup: } return exit_code; } +#endif /* FASTSYNC_TEST_BUILD */ diff --git a/src/shared/config.c b/src/shared/config.c index c7900cc..cc8b8e6 100644 --- a/src/shared/config.c +++ b/src/shared/config.c @@ -8,10 +8,7 @@ #include #include -Config* config_create(void) { - Config* config = malloc(sizeof(Config)); - if (!config) - return NULL; +static void config_set_defaults(Config* config) { config->version = str_dup(PROTOCOL_VERSION); config->send_directory = NULL; config->receive_root_directory = NULL; @@ -101,6 +98,13 @@ Config* config_create(void) { config->server_mode = false; config->checksum = false; config->compress_choice = NULL; +} + +Config* config_create(void) { + Config* config = malloc(sizeof(Config)); + if (!config) + return NULL; + config_set_defaults(config); return config; } @@ -208,7 +212,7 @@ bool config_send(int file_descriptor, const Config* config) { return false; if (!send_int(file_descriptor, config->use_delta)) return false; - if (!send_int(file_descriptor, (int)config->delta_block_size)) + if (!send_n_data(file_descriptor, &config->delta_block_size, sizeof(config->delta_block_size))) return false; if (!send_n_data(file_descriptor, &config->delta_max_file_size, sizeof(unsigned long long))) return false; @@ -281,9 +285,11 @@ Config* config_receive(int file_descriptor) { Config* config = (Config*)malloc(sizeof(Config)); if (config == NULL) return NULL; - memset(config, 0, sizeof(*config)); + config_set_defaults(config); + free(config->version); config->version = receive_str(file_descriptor); if (!config->version) { + free(config->server_host); free(config); return NULL; } @@ -291,6 +297,7 @@ Config* config_receive(int file_descriptor) { fprintf(stderr, "Protocol version mismatch: client=%s, server=%s\n", config->version, PROTOCOL_VERSION); free(config->version); + free(config->server_host); free(config); send_status(file_descriptor, STATUS_ERROR); return NULL; @@ -298,6 +305,7 @@ Config* config_receive(int file_descriptor) { config->send_directory = receive_str(file_descriptor); if (!config->send_directory) { free(config->version); + free(config->server_host); free(config); return NULL; } @@ -305,6 +313,7 @@ Config* config_receive(int file_descriptor) { if (!config->receive_root_directory) { free(config->version); free(config->send_directory); + free(config->server_host); free(config); return NULL; } @@ -341,67 +350,10 @@ Config* config_receive(int file_descriptor) { if (!receive_int(file_descriptor, &tmp)) goto error; config->use_delta = tmp; - if (!receive_int(file_descriptor, &tmp)) + if (!receive_n_data(file_descriptor, &config->delta_block_size, sizeof(config->delta_block_size))) goto error; - config->delta_block_size = (uint32_t)tmp; if (!receive_n_data(file_descriptor, &config->delta_max_file_size, sizeof(unsigned long long))) goto error; - config->show_progress = false; - config->dry_run = false; - config->ssh_port = 22; - config->transport = TRANSPORT_TCP; - config->ssh_destination = NULL; - config->fastsync_server_path = NULL; - config->exclude_patterns = NULL; - config->exclude_count = 0; - config->include_patterns = NULL; - config->include_count = 0; - config->max_size = 0; - config->min_size = 0; - config->use_tls = false; - config->tls_cert = NULL; - config->tls_key = NULL; - config->tls_ca = NULL; - config->timeout = 30; - config->contimeout = 10; - config->quiet = false; - config->stats = false; - config->max_depth = 0; - config->log_file = NULL; - config->queue_size = 100; - config->follow_symlinks = false; - config->copy_links = false; - config->safe_links = false; - config->copy_unsafe_links = false; - config->preserve_hard_links = false; - config->preserve_acls = false; - config->preserve_xattrs = false; - config->preserve_devices = false; - config->preserve_sparse = false; - config->itemize_changes = false; - config->out_format = NULL; - config->info_level = 0; - config->debug_level = 0; - config->list_only = false; - config->human_readable = false; - config->update = false; - config->inplace = false; - config->append = false; - config->append_verify = false; - config->delete_excluded = false; - config->delete_after = false; - config->max_delete = 0; - config->filters = NULL; - config->files_from = NULL; - config->cvs_exclude = false; - config->prune_empty_dirs = false; - config->relative = false; - config->rsh_command = NULL; - config->rsync_path = NULL; - config->temp_dir = NULL; - config->compare_dest = NULL; - config->copy_dest = NULL; - config->link_dest = NULL; if (!receive_int(file_descriptor, &tmp)) goto error; config->backup = tmp; @@ -482,15 +434,6 @@ Config* config_receive(int file_descriptor) { config->compress_choice = receive_str(file_descriptor); if (config->compress_choice == NULL) goto error; - config->address = NULL; - config->bind_address = NULL; - config->ipv6 = false; - config->ipv4 = false; - config->daemon = false; - config->daemon_config = NULL; - config->server_mode = false; - config->server_host = str_dup("127.0.0.1"); - config->server_port = 8080; if (!send_status(file_descriptor, STATUS_OK)) goto error; return config; diff --git a/src/shared/delta.c b/src/shared/delta.c index 0305fdc..f9f7c84 100644 --- a/src/shared/delta.c +++ b/src/shared/delta.c @@ -388,6 +388,10 @@ Delta* delta_deserialize(const Data* data) { if (type == DELTA_OP_BLOCK_MATCH) { if (pos + sizeof(uint32_t) * 3 > data->size) { + for (uint32_t k = 0; k < i; k++) { + if (delta->instructions[k].type == DELTA_INSTR_LITERAL) + free(delta->instructions[k].literal.data); + } free(delta->instructions); free(delta); return NULL; @@ -426,6 +430,11 @@ Delta* delta_deserialize(const Data* data) { } delta->instructions[i].literal.data = malloc(lit_len); if (!delta->instructions[i].literal.data) { + log_message(LOG_LEVEL_ERROR, "Failed to allocate %u bytes for literal data", lit_len); + for (uint32_t k = 0; k < i; k++) { + if (delta->instructions[k].type == DELTA_INSTR_LITERAL) + free(delta->instructions[k].literal.data); + } free(delta->instructions); free(delta); return NULL; diff --git a/src/shared/protocol.c b/src/shared/protocol.c index 0d41808..bce9507 100644 --- a/src/shared/protocol.c +++ b/src/shared/protocol.c @@ -10,7 +10,6 @@ #include #include -#define MAX_DATA_SIZE (100ULL * 1024 * 1024) /* 100 MB max per data message */ #define RECEIVE_TIMEOUT_SEC 60 /* 60 second per-message timeout */ #define MAX_CONNECTION_MEMORY (1024ULL * 1024 * 1024) /* 1 GB total per connection */ @@ -236,9 +235,9 @@ Data* receive_data(int file_descriptor) { unsigned long long size = 0; if (!receive_n_data(file_descriptor, &size, sizeof(unsigned long long))) return NULL; - if (size > MAX_DATA_SIZE) { + if (size > MAX_DATA_PAYLOAD_SIZE) { log_message(LOG_LEVEL_ERROR, "Data size %llu exceeds maximum %llu", size, - (unsigned long long)MAX_DATA_SIZE); + (unsigned long long)MAX_DATA_PAYLOAD_SIZE); return NULL; } if (total_allocated_bytes + size > MAX_CONNECTION_MEMORY) { diff --git a/tests/test_client_cli.c b/tests/test_client_cli.c index ab44023..cf13dfe 100644 --- a/tests/test_client_cli.c +++ b/tests/test_client_cli.c @@ -6,6 +6,9 @@ #include #include +/* Declaration of parse_args from client_cli.c */ +int parse_args(Config* config, int argc, char* argv[], int* positional_args, int* positional_count); + /* Test main() with --help flag (early return path, no server connection needed) */ static void test_cli_help() { /* We can't easily call main() because it calls send_files which needs a server. @@ -80,10 +83,160 @@ static void test_cli_exclude_patterns() { config_delete(cfg); } +/* Test parse_args with --help returns 1 (clean exit) */ +static void test_parse_args_help() { + Config* cfg = config_create(); + char* argv[] = {"fastsync", "--help"}; + int positional_args[2]; + int positional_count = 0; + + int ret = parse_args(cfg, 2, argv, positional_args, &positional_count); + EXPECT_EQ_INT(ret, 1); + + config_delete(cfg); +} + +/* Test parse_args with -V/--version returns 1 */ +static void test_parse_args_version() { + Config* cfg = config_create(); + char* argv_short[] = {"fastsync", "-V"}; + char* argv_long[] = {"fastsync", "--version"}; + int positional_args[2]; + int positional_count = 0; + + int ret = parse_args(cfg, 2, argv_short, positional_args, &positional_count); + EXPECT_EQ_INT(ret, 1); + + ret = parse_args(cfg, 2, argv_long, positional_args, &positional_count); + EXPECT_EQ_INT(ret, 1); + + config_delete(cfg); +} + +/* Test parse_args with valid port */ +static void test_parse_args_valid_port() { + Config* cfg = config_create(); + char* argv[] = {"fastsync", "-p", "2222", "/src", "/dst"}; + int positional_args[2]; + int positional_count = 0; + + int ret = parse_args(cfg, 5, argv, positional_args, &positional_count); + EXPECT_EQ_INT(ret, 0); + EXPECT_EQ_INT(cfg->ssh_port, 2222); + EXPECT_EQ_INT(positional_count, 2); + + config_delete(cfg); +} + +/* Test parse_args rejects port > 65535 */ +static void test_parse_args_invalid_port() { + Config* cfg = config_create(); + char* argv[] = {"fastsync", "-p", "99999", "/src", "/dst"}; + int positional_args[2]; + int positional_count = 0; + + int ret = parse_args(cfg, 5, argv, positional_args, &positional_count); + EXPECT_EQ_INT(ret, -1); + + config_delete(cfg); +} + +/* Test parse_args rejects non-numeric port */ +static void test_parse_args_non_numeric_port() { + Config* cfg = config_create(); + char* argv[] = {"fastsync", "-p", "abc", "/src", "/dst"}; + int positional_args[2]; + int positional_count = 0; + + int ret = parse_args(cfg, 5, argv, positional_args, &positional_count); + EXPECT_EQ_INT(ret, -1); + + config_delete(cfg); +} + +/* Test parse_args rejects server port > 65535 */ +static void test_parse_args_invalid_server_port() { + Config* cfg = config_create(); + char* argv[] = {"fastsync", "--server-port", "70000", "/src", "/dst"}; + int positional_args[2]; + int positional_count = 0; + + int ret = parse_args(cfg, 5, argv, positional_args, &positional_count); + EXPECT_EQ_INT(ret, -1); + + config_delete(cfg); +} + +/* Test parse_args rejects invalid compression level */ +static void test_parse_args_invalid_compression_level() { + Config* cfg = config_create(); + char* argv[] = {"fastsync", "-c", "25", "/src", "/dst"}; + int positional_args[2]; + int positional_count = 0; + + int ret = parse_args(cfg, 5, argv, positional_args, &positional_count); + EXPECT_EQ_INT(ret, -1); + + config_delete(cfg); +} + +/* Test parse_args accepts valid compression level */ +static void test_parse_args_valid_compression_level() { + Config* cfg = config_create(); + char* argv[] = {"fastsync", "-c", "10", "/src", "/dst"}; + int positional_args[2]; + int positional_count = 0; + + int ret = parse_args(cfg, 5, argv, positional_args, &positional_count); + EXPECT_EQ_INT(ret, 0); + EXPECT_EQ_INT(cfg->compression_level, 10); + + config_delete(cfg); +} + +/* Test parse_args unknown option returns error */ +static void test_parse_args_unknown_option() { + Config* cfg = config_create(); + char* argv[] = {"fastsync", "--nonexistent", "/src", "/dst"}; + int positional_args[2]; + int positional_count = 0; + + int ret = parse_args(cfg, 4, argv, positional_args, &positional_count); + EXPECT_EQ_INT(ret, -1); + + config_delete(cfg); +} + +/* Test parse_args with --archive flag */ +static void test_parse_args_archive() { + Config* cfg = config_create(); + char* argv[] = {"fastsync", "--archive", "/src", "/dst"}; + int positional_args[2]; + int positional_count = 0; + + int ret = parse_args(cfg, 4, argv, positional_args, &positional_count); + EXPECT_EQ_INT(ret, 0); + EXPECT_TRUE(cfg->use_compression); + EXPECT_TRUE(cfg->use_multithreading); + EXPECT_TRUE(cfg->use_metadata); + + config_delete(cfg); +} + void test_client_cli() { test_cli_help(); test_cli_archive_flags(); test_cli_dry_run(); test_cli_delete_flag(); test_cli_exclude_patterns(); + test_parse_args_help(); + test_parse_args_version(); + test_parse_args_valid_port(); + test_parse_args_invalid_port(); + test_parse_args_non_numeric_port(); + test_parse_args_invalid_server_port(); + test_parse_args_invalid_compression_level(); + test_parse_args_valid_compression_level(); + test_parse_args_unknown_option(); + test_parse_args_archive(); }