From f443b2986dbe3fe5272f7b4f8aed4fec606e3641 Mon Sep 17 00:00:00 2001 From: TapTap Date: Thu, 3 Sep 2026 18:15:14 +0200 Subject: [PATCH] feat: add compression worker threads --- README.md | 1 + src/client/client_cli.c | 11 +++++++++++ src/client/client_send.c | 22 ++++++++++++++-------- src/client/client_validation.c | 4 ++++ src/client/usage.c | 1 + src/shared/chunk.c | 7 ++++++- src/shared/chunk.h | 2 ++ src/shared/compression.c | 28 ++++++++++++++++++++++++++++ src/shared/compression.h | 2 ++ src/shared/config.c | 1 + src/shared/config.h | 1 + src/shared/file_send.c | 10 +++++++++- src/shared/file_send.h | 3 +++ tests/integration/test_tcp.py | 5 +++++ tests/test_client_cli.c | 31 +++++++++++++++++++++++++++++++ tests/test_compression.c | 17 +++++++++++++++++ 16 files changed, 136 insertions(+), 10 deletions(-) diff --git a/README.md b/README.md index fdbd9f3..953618d 100644 --- a/README.md +++ b/README.md @@ -339,6 +339,7 @@ features without changing the meaning of ordinary compatibility options. | `-m` | Enable the multithreaded scanner/loader/sender pipeline. | | `-c [level]`, `-z [level]` | Enable streaming zstd compression, levels 1-22. | | `--compress-level ` | Set the zstd compression level. | +| `--compress-threads ` | Use `n` zstd compression workers. Requires compression and a zstd build with threaded support; the setting affects sender CPU work only. | | `--chunk-size ` | Set the transfer chunk size. | | `-s` | Enable FastSync chunk serialization. | | `-f`, `--sendfile` | Use TCP `sendfile()` zero-copy transfer. Incompatible with compression and chunk serialization. | diff --git a/src/client/client_cli.c b/src/client/client_cli.c index 191cc07..65c614e 100644 --- a/src/client/client_cli.c +++ b/src/client/client_cli.c @@ -211,6 +211,13 @@ static int apply_table_option(Config* config, const OptionEntry* entry, const ch int parse_args(Config* config, int argc, char* argv[], int* positional_args, int* positional_count) { 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) + return -1; + continue; + } const OptionEntry* entry = find_table_option(argv[i]); if (entry) { if (entry->kind != OPT_FLAG) { @@ -361,6 +368,10 @@ int parse_args(Config* config, int argc, char* argv[], int* positional_args, log_message(LOG_LEVEL_ERROR, "--compress-level must be between 1 and 22"); 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) + return -1; } else if (argv[i][0] == '-') { fprintf(stderr, "Unknown option: %s\n", argv[i]); print_usage(); diff --git a/src/client/client_send.c b/src/client/client_send.c index 0b4ca5c..945deef 100644 --- a/src/client/client_send.c +++ b/src/client/client_send.c @@ -235,7 +235,8 @@ static int send_delta(Client* client, File* file, DeltaSignature* sig, Config* c Data* to_send = delta_data; if (config->use_compression) { - to_send = data_compress(delta_data, config->compression_level); + to_send = data_compress_with_threads(delta_data, config->compression_level, + config->compression_threads); data_destroy(delta_data); if (!to_send) return send_status(client->file_descriptor, STATUS_NEXT) ? 1 : -1; @@ -251,13 +252,15 @@ static int send_delta(Client* client, File* file, DeltaSignature* sig, Config* c return ok ? 0 : -1; } -typedef bool (*file_send_fn)(File*, int, bool, int, bool); +typedef bool (*file_send_fn)(File*, int, bool, int, int, bool); // Send a single file directly (non-incremental path). -static bool send_file_direct(File* file, int fd, bool use_metadata, int compression_level) { +static bool send_file_direct(File* file, int fd, bool use_metadata, int compression_level, + int compression_threads) { if (!send_status(fd, STATUS_NEXT)) return false; - return file_send_single_calls(file, fd, use_metadata, compression_level, true); + return file_send_single_calls_with_threads(file, fd, use_metadata, compression_level, + compression_threads, true); } // Send a single file directly via sendfile (non-incremental path). @@ -278,7 +281,8 @@ static int send_single_file(Client* client, File* file, Config* config, bool use return send_file_direct_sendfile(file, client->file_descriptor, config->use_metadata) ? 0 : -1; } - return send_file_direct(file, client->file_descriptor, config->use_metadata, compression_level) + return send_file_direct(file, client->file_descriptor, config->use_metadata, compression_level, + config->compression_threads) ? 0 : -1; } @@ -310,7 +314,7 @@ static int send_single_file(Client* client, File* file, Config* config, bool use } // Incremental path with single_calls (supports compression and delta) - file_send_fn send_fn = (file_send_fn)file_send_single_calls; + file_send_fn send_fn = (file_send_fn)file_send_single_calls_with_threads; DeltaSignature* sig = NULL; int rc = incremental_check(client, file, config, &sig); if (rc < 0) { @@ -338,7 +342,8 @@ static int send_single_file(Client* client, File* file, Config* config, bool use return -1; } } - if (!send_fn(file, client->file_descriptor, config->use_metadata, compression_level, false)) + if (!send_fn(file, client->file_descriptor, config->use_metadata, compression_level, + config->compression_threads, false)) return -1; return 0; } @@ -349,7 +354,8 @@ int send_chunk(Client* client, Chunk* chunk, Config* config) { return -1; Data* data; if (config->use_compression) { - data = chunk_compress(chunk, config->compression_level, config->use_metadata); + data = chunk_compress_with_threads(chunk, config->compression_level, config->use_metadata, + config->compression_threads); } else { data = chunk_serialize(chunk, config->use_metadata); } diff --git a/src/client/client_validation.c b/src/client/client_validation.c index 106580f..f91c66c 100644 --- a/src/client/client_validation.c +++ b/src/client/client_validation.c @@ -15,6 +15,10 @@ bool validate_config(const Config* config) { "(chunk serialization)"); return false; } + if (config->compression_threads > 0 && !config->use_compression) { + log_message(LOG_LEVEL_ERROR, "--compress-threads requires compression (-c or -z)"); + return false; + } if (config->transport == TRANSPORT_SSH && config->use_sendfile) { log_message(LOG_LEVEL_ERROR, "-f/--sendfile is not supported with SSH transport"); return false; diff --git a/src/client/usage.c b/src/client/usage.c index 2f800af..af6049c 100644 --- a/src/client/usage.c +++ b/src/client/usage.c @@ -69,6 +69,7 @@ void print_usage(void) { printf(" -S, --sparse Handle sparse files efficiently\n"); printf(" --inplace Update files in-place (no temp+rename)\n"); printf(" --compress-level Compression level (default: 5)\n"); + printf(" --compress-threads Compression worker threads (requires zstd threaded support)\n"); printf(" --help Show this help\n"); printf(" -V, --version Show version\n"); } diff --git a/src/shared/chunk.c b/src/shared/chunk.c index 3af1c63..7abee45 100644 --- a/src/shared/chunk.c +++ b/src/shared/chunk.c @@ -290,11 +290,16 @@ Chunk* chunk_deserialize(Data* data, bool use_metadata) { } Data* chunk_compress(Chunk* chunk, int compression_level, bool use_metadata) { + return chunk_compress_with_threads(chunk, compression_level, use_metadata, 0); +} + +Data* chunk_compress_with_threads(Chunk* chunk, int compression_level, bool use_metadata, + int compression_threads) { log_message(LOG_LEVEL_DEBUG, "Starting to compress chunk"); Data* serialized = chunk_serialize(chunk, use_metadata); if (serialized == NULL) return NULL; - Data* compressed = data_compress(serialized, compression_level); + Data* compressed = data_compress_with_threads(serialized, compression_level, compression_threads); data_destroy(serialized); if (compressed == NULL) return NULL; diff --git a/src/shared/chunk.h b/src/shared/chunk.h index dbba7b7..2c04e85 100644 --- a/src/shared/chunk.h +++ b/src/shared/chunk.h @@ -19,6 +19,8 @@ void chunk_destroy(void* chunk); Data* chunk_serialize(Chunk* chunk, bool use_metadata); Chunk* chunk_deserialize(Data* data, bool use_metadata); Data* chunk_compress(Chunk* chunk, int compression_level, bool use_metadata); +Data* chunk_compress_with_threads(Chunk* chunk, int compression_level, bool use_metadata, + int compression_threads); Chunk* receive_chunk_data(int fd, const Config* config); #endif diff --git a/src/shared/compression.c b/src/shared/compression.c index 1a42ba5..5071e27 100644 --- a/src/shared/compression.c +++ b/src/shared/compression.c @@ -28,6 +28,14 @@ bool compression_should_skip(const char* path) { } Data* data_compress(Data* data_to_compress, int compression_level) { + return data_compress_with_threads(data_to_compress, compression_level, 0); +} + +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) + return NULL; log_message(LOG_LEVEL_DEBUG, "Starting to compress data"); size_t dst_size = ZSTD_compressBound(data_to_compress->size); Data* compressed_data = data_create_empty(dst_size); @@ -49,6 +57,26 @@ Data* data_compress(Data* data_to_compress, int compression_level) { return NULL; } + if (compression_threads > 0) { + zret = ZSTD_CCtx_setParameter(cctx, ZSTD_c_nbWorkers, compression_threads); + if (ZSTD_isError(zret)) { + log_message(LOG_LEVEL_ERROR, "Failed to set compression threads: %s", + ZSTD_getErrorName(zret)); + ZSTD_freeCCtx(cctx); + data_destroy(compressed_data); + return NULL; + } + /* Streaming compression needs the source size before threaded mode can end a frame. */ + zret = ZSTD_CCtx_setPledgedSrcSize(cctx, data_to_compress->size); + if (ZSTD_isError(zret)) { + log_message(LOG_LEVEL_ERROR, "Failed to set compression source size: %s", + ZSTD_getErrorName(zret)); + ZSTD_freeCCtx(cctx); + data_destroy(compressed_data); + return NULL; + } + } + ZSTD_inBuffer input = {data_to_compress->data, data_to_compress->size, 0}; ZSTD_outBuffer output = {compressed_data->data, dst_size, 0}; diff --git a/src/shared/compression.h b/src/shared/compression.h index 5322abe..b40a04f 100644 --- a/src/shared/compression.h +++ b/src/shared/compression.h @@ -5,6 +5,8 @@ #include 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); Data* data_decompress(Data* compressed_data); Data* data_decompress_limited(Data* compressed_data, size_t maximum_size); bool compression_should_skip(const char* path); diff --git a/src/shared/config.c b/src/shared/config.c index efc2708..10a20e9 100644 --- a/src/shared/config.c +++ b/src/shared/config.c @@ -21,6 +21,7 @@ static void config_set_defaults(Config* config) { config->dry_run = false; config->use_delete = false; config->compression_level = 5; + config->compression_threads = 0; config->use_sendfile = false; config->chunk_size = DEFAULT_CHUNK_SIZE; config->ssh_port = 22; diff --git a/src/shared/config.h b/src/shared/config.h index 4a218b0..2055e29 100644 --- a/src/shared/config.h +++ b/src/shared/config.h @@ -22,6 +22,7 @@ typedef struct Config { bool dry_run; bool use_delete; int compression_level; + int compression_threads; unsigned long long chunk_size; int ssh_port; TransportType transport; diff --git a/src/shared/file_send.c b/src/shared/file_send.c index 5757ce1..ed74910 100644 --- a/src/shared/file_send.c +++ b/src/shared/file_send.c @@ -19,12 +19,20 @@ bool file_send_single_calls(File* file, int file_descriptor, bool use_metadata, int compression_level, bool send_path) { + return file_send_single_calls_with_threads(file, file_descriptor, use_metadata, compression_level, + 0, send_path); +} + +bool file_send_single_calls_with_threads(File* file, int file_descriptor, bool use_metadata, + int compression_level, int compression_threads, + bool send_path) { if (!file || !file->path || !file->data || (file->data->size != 0 && !file->data->data)) return false; const Data* data_to_send = file->data; Data* compressed_data = NULL; if (compression_level > 0 && !compression_should_skip(file->path)) { - compressed_data = data_compress(file->data, compression_level); + compressed_data = + data_compress_with_threads(file->data, compression_level, compression_threads); if (compressed_data == NULL) { log_message(LOG_LEVEL_ERROR, "Failed to compress file data"); return false; diff --git a/src/shared/file_send.h b/src/shared/file_send.h index 86c17ad..c0388af 100644 --- a/src/shared/file_send.h +++ b/src/shared/file_send.h @@ -8,6 +8,9 @@ bool file_send_single_calls(File* file, int file_descriptor, bool use_metadata, int compression_level, bool send_path); +bool file_send_single_calls_with_threads(File* file, int file_descriptor, bool use_metadata, + int compression_level, int compression_threads, + bool send_path); bool file_send_sendfile(File* file, int file_descriptor, bool use_metadata, int compression_level, bool send_path); diff --git a/tests/integration/test_tcp.py b/tests/integration/test_tcp.py index 18cc07c..ab4f405 100644 --- a/tests/integration/test_tcp.py +++ b/tests/integration/test_tcp.py @@ -71,6 +71,11 @@ class TestTCPFlags: r = _run_tcp_test("Compression (-c)", shared_server.port, ["-c"]) assert r["status"] == "Success", r["error"] + def test_compression_threads(self, shared_server): + r = _run_tcp_test("Compression threads (-c --compress-threads=2)", shared_server.port, + ["-c", "--compress-threads=2"]) + assert r["status"] == "Success", r["error"] + def test_chunk_serialization(self, shared_server): r = _run_tcp_test("Chunk Serialization (-s)", shared_server.port, ["-s"]) assert r["status"] == "Success", r["error"] diff --git a/tests/test_client_cli.c b/tests/test_client_cli.c index 656fd40..cbfad83 100644 --- a/tests/test_client_cli.c +++ b/tests/test_client_cli.c @@ -36,6 +36,12 @@ static void test_validate_config_incompatible_options() { cfg->use_incremental = true; cfg->use_chunk_serialization = true; EXPECT_FALSE(validate_config(cfg)); + cfg->use_incremental = false; + cfg->compression_threads = 2; + EXPECT_FALSE(validate_config(cfg)); + cfg->use_compression = true; + cfg->use_sendfile = false; + EXPECT_TRUE(validate_config(cfg)); config_delete(cfg); } @@ -247,6 +253,30 @@ static void test_parse_args_valid_compression_level() { config_delete(cfg); } +static void test_parse_args_compression_threads() { + Config* cfg = config_create(); + char* argv[] = {"fastsync", "--compress-threads", "4", "/src", "/dst"}; + int positional_args[2]; + int positional_count = 0; + + EXPECT_EQ_INT(parse_args(cfg, 5, argv, positional_args, &positional_count), 0); + EXPECT_EQ_INT(cfg->compression_threads, 4); + config_delete(cfg); + + cfg = config_create(); + char* equals_argv[] = {"fastsync", "--compress-threads=3", "/src", "/dst"}; + positional_count = 0; + EXPECT_EQ_INT(parse_args(cfg, 4, equals_argv, positional_args, &positional_count), 0); + EXPECT_EQ_INT(cfg->compression_threads, 3); + config_delete(cfg); + + cfg = config_create(); + char* invalid_argv[] = {"fastsync", "--compress-threads=0", "/src", "/dst"}; + positional_count = 0; + EXPECT_EQ_INT(parse_args(cfg, 4, invalid_argv, positional_args, &positional_count), -1); + config_delete(cfg); +} + /* Test parse_args unknown option returns error */ static void test_parse_args_unknown_option() { Config* cfg = config_create(); @@ -357,6 +387,7 @@ void test_client_cli() { test_parse_args_invalid_server_port(); test_parse_args_invalid_compression_level(); test_parse_args_valid_compression_level(); + test_parse_args_compression_threads(); test_parse_args_unknown_option(); test_parse_args_rejects_unimplemented_options(); test_parse_args_archive(); diff --git a/tests/test_compression.c b/tests/test_compression.c index a234e00..029463b 100644 --- a/tests/test_compression.c +++ b/tests/test_compression.c @@ -55,6 +55,22 @@ static void test_data_compress_decompress_large() { data_destroy(decompressed); } +static void test_data_compress_with_threads_roundtrip() { + const char original[] = "Threaded zstd compression test data"; + Data* input = data_create_empty(sizeof(original) - 1); + EXPECT_NOT_NULL(input); + memcpy(input->data, original, sizeof(original) - 1); + 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); + data_destroy(input); + data_destroy(compressed); + data_destroy(decompressed); +} + static void test_chunk_compress_decompress_roundtrip() { char* path1 = "temp_comp_test_1.txt"; char* content1 = "chunk compression test file 1"; @@ -115,5 +131,6 @@ static void test_chunk_compress_decompress_roundtrip() { void test_compression() { test_data_compress_decompress_roundtrip(); test_data_compress_decompress_large(); + test_data_compress_with_threads_roundtrip(); test_chunk_compress_decompress_roundtrip(); }