From ec02ee6ddedc260fe06ae320e220f7ec61590014 Mon Sep 17 00:00:00 2001 From: TapTap Date: Thu, 3 Sep 2026 22:22:04 +0200 Subject: [PATCH] fix: bound compression worker threads --- src/client/client_cli.c | 19 +++++++++++++++---- src/shared/compression.c | 9 +++++++-- src/shared/compression.h | 2 ++ tests/test_client_cli.c | 6 ++++++ tests/test_compression.c | 11 ++++++----- 5 files changed, 36 insertions(+), 11 deletions(-) diff --git a/src/client/client_cli.c b/src/client/client_cli.c index 65c614e..86a6a2d 100644 --- a/src/client/client_cli.c +++ b/src/client/client_cli.c @@ -1,5 +1,6 @@ #include "client_send.h" #include "client_validation.h" +#include "compression.h" #include "config.h" #include "delta.h" #include "log.h" @@ -77,6 +78,17 @@ static int set_positive_int_option(int* dest, const char* value, const char* opt return 0; } +static int set_compression_threads_option(int* dest, const char* value) { + if (set_positive_int_option(dest, value, "--compress-threads") != 0) + return -1; + if (*dest > COMPRESSION_MAX_THREADS) { + log_message(LOG_LEVEL_ERROR, "--compress-threads must be between 1 and %d", + COMPRESSION_MAX_THREADS); + return -1; + } + return 0; +} + /* Parse a string as a non-negative integer into *dest. Returns true on success, false on error. */ static int set_nonneg_int_option(int* dest, const char* value, const char* option_name) { if (!parse_nonneg_int(value, dest)) { @@ -213,8 +225,8 @@ int parse_args(Config* config, int argc, char* argv[], int* positional_args, for (int i = 1; i < argc; i++) { const char* threads_prefix = "--compress-threads="; if (strncmp(argv[i], threads_prefix, strlen(threads_prefix)) == 0) { - if (set_positive_int_option(&config->compression_threads, argv[i] + strlen(threads_prefix), - "--compress-threads") != 0) + if (set_compression_threads_option(&config->compression_threads, + argv[i] + strlen(threads_prefix)) != 0) return -1; continue; } @@ -369,8 +381,7 @@ int parse_args(Config* config, int argc, char* argv[], int* positional_args, return -1; } } else if (opt_is(argv[i], "--compress-threads", NULL) && i + 1 < argc) { - if (set_positive_int_option(&config->compression_threads, argv[++i], "--compress-threads") != - 0) + if (set_compression_threads_option(&config->compression_threads, argv[++i]) != 0) return -1; } else if (argv[i][0] == '-') { fprintf(stderr, "Unknown option: %s\n", argv[i]); diff --git a/src/shared/compression.c b/src/shared/compression.c index 5071e27..3e4791a 100644 --- a/src/shared/compression.c +++ b/src/shared/compression.c @@ -6,6 +6,7 @@ #include #include #include +#include #include #define INITIAL_DECOMPRESS_BUF_SIZE (1024 * 1024) @@ -34,7 +35,7 @@ Data* data_compress(Data* data_to_compress, int compression_level) { Data* data_compress_with_threads(Data* data_to_compress, int compression_level, int compression_threads) { if (!data_to_compress || (!data_to_compress->data && data_to_compress->size != 0) || - compression_threads < 0) + compression_threads < 0 || compression_threads > COMPRESSION_MAX_THREADS) return NULL; log_message(LOG_LEVEL_DEBUG, "Starting to compress data"); size_t dst_size = ZSTD_compressBound(data_to_compress->size); @@ -58,7 +59,11 @@ Data* data_compress_with_threads(Data* data_to_compress, int compression_level, } if (compression_threads > 0) { - zret = ZSTD_CCtx_setParameter(cctx, ZSTD_c_nbWorkers, compression_threads); + long online_cpus = sysconf(_SC_NPROCESSORS_ONLN); + int available_threads = online_cpus > 0 && online_cpus < compression_threads + ? (int)online_cpus + : compression_threads; + zret = ZSTD_CCtx_setParameter(cctx, ZSTD_c_nbWorkers, available_threads); if (ZSTD_isError(zret)) { log_message(LOG_LEVEL_ERROR, "Failed to set compression threads: %s", ZSTD_getErrorName(zret)); diff --git a/src/shared/compression.h b/src/shared/compression.h index b40a04f..3dcc90a 100644 --- a/src/shared/compression.h +++ b/src/shared/compression.h @@ -4,6 +4,8 @@ #include "data.h" #include +#define COMPRESSION_MAX_THREADS 64 + Data* data_compress(Data* data_to_compress, int compression_level); Data* data_compress_with_threads(Data* data_to_compress, int compression_level, int compression_threads); diff --git a/tests/test_client_cli.c b/tests/test_client_cli.c index cbfad83..9d90011 100644 --- a/tests/test_client_cli.c +++ b/tests/test_client_cli.c @@ -275,6 +275,12 @@ static void test_parse_args_compression_threads() { positional_count = 0; EXPECT_EQ_INT(parse_args(cfg, 4, invalid_argv, positional_args, &positional_count), -1); config_delete(cfg); + + cfg = config_create(); + char* excessive_argv[] = {"fastsync", "--compress-threads=65", "/src", "/dst"}; + positional_count = 0; + EXPECT_EQ_INT(parse_args(cfg, 4, excessive_argv, positional_args, &positional_count), -1); + config_delete(cfg); } /* Test parse_args unknown option returns error */ diff --git a/tests/test_compression.c b/tests/test_compression.c index 029463b..7fe8bc0 100644 --- a/tests/test_compression.c +++ b/tests/test_compression.c @@ -56,16 +56,17 @@ static void test_data_compress_decompress_large() { } static void test_data_compress_with_threads_roundtrip() { - const char original[] = "Threaded zstd compression test data"; - Data* input = data_create_empty(sizeof(original) - 1); + const size_t size = 8 * 1024 * 1024; + Data* input = data_create_empty(size); EXPECT_NOT_NULL(input); - memcpy(input->data, original, sizeof(original) - 1); + for (size_t i = 0; i < size; i++) + ((char*)input->data)[i] = (char)((i / 4096) % 7); Data* compressed = data_compress_with_threads(input, 3, 2); EXPECT_NOT_NULL(compressed); Data* decompressed = data_decompress(compressed); EXPECT_NOT_NULL(decompressed); - EXPECT_EQ_INT((int)decompressed->size, (int)(sizeof(original) - 1)); - EXPECT_EQ_INT(memcmp(decompressed->data, original, sizeof(original) - 1), 0); + EXPECT_EQ_INT((int)decompressed->size, (int)size); + EXPECT_EQ_INT(memcmp(decompressed->data, input->data, size), 0); data_destroy(input); data_destroy(compressed); data_destroy(decompressed);