diff --git a/README.md b/README.md index e9cb7a6..8868dc8 100644 --- a/README.md +++ b/README.md @@ -346,6 +346,7 @@ features without changing the meaning of ordinary compatibility options. | `--zc ` | Alias for `--compress-choice`. FastSync supports `zstd` and `none`. | | `--zl ` | Alias for `--compress-level`. | | `--skip-compress ` | Skip compression for comma-separated suffixes; incompatible with `-s`. | +| `--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 e7768b0..bce7d4c 100644 --- a/src/client/client_cli.c +++ b/src/client/client_cli.c @@ -1,6 +1,7 @@ #include "client_send.h" #include "client_validation.h" #include "chmod.h" +#include "compression.h" #include "config.h" #include "delta.h" #include "log.h" @@ -90,6 +91,17 @@ static int set_compression_choice(Config* config, const char* value) { 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)) { @@ -369,6 +381,13 @@ int parse_args(Config* config, int argc, char* argv[], int* positional_args, return -1; continue; } + const char* threads_prefix = "--compress-threads="; + if (strncmp(argv[i], threads_prefix, strlen(threads_prefix)) == 0) { + if (set_compression_threads_option(&config->compression_threads, + argv[i] + strlen(threads_prefix)) != 0) + return -1; + continue; + } const OptionEntry* entry = find_table_option(argv[i]); const char* inline_value = NULL; if (!entry) @@ -586,6 +605,9 @@ int parse_args(Config* config, int argc, char* argv[], int* positional_args, } else if (opt_is(argv[i], "--skip-compress", NULL) && i + 1 < argc) { if (parse_skip_compress(config, argv[++i]) != 0) return -1; + } else if (opt_is(argv[i], "--compress-threads", NULL) && i + 1 < argc) { + if (set_compression_threads_option(&config->compression_threads, argv[++i]) != 0) + return -1; } else if (argv[i][0] == '-') { char* escaped = output_escape(argv[i], false); fprintf(stderr, "Unknown option: %s\n", escaped ? escaped : ""); diff --git a/src/client/client_send.c b/src/client/client_send.c index 335a9d6..0f017c7 100644 --- a/src/client/client_send.c +++ b/src/client/client_send.c @@ -354,7 +354,8 @@ static int send_delta(Client* client, File* file, DeltaSignature* sig, Config* c int skip_count = config->skip_compress_set ? config->skip_compress_count : -1; if (config->use_compression && !compression_should_skip_with_suffixes( file->path, config->skip_compress_suffixes, skip_count)) { - 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; @@ -377,7 +378,8 @@ static bool send_file_direct(File* file, int fd, bool use_metadata, int compress return false; int skip_count = config->skip_compress_set ? config->skip_compress_count : -1; return file_send_single_calls_with_skip(file, fd, use_metadata, compression_level, true, - config->skip_compress_suffixes, skip_count); + config->skip_compress_suffixes, skip_count, + config->compression_threads); } // Send a single file directly via sendfile (non-incremental path). @@ -386,7 +388,8 @@ static bool send_file_direct_sendfile(File* file, int fd, bool use_metadata, con return false; int skip_count = config->skip_compress_set ? config->skip_compress_count : -1; return file_send_sendfile_with_skip(file, fd, use_metadata, 0, true, - config->skip_compress_suffixes, skip_count); + config->skip_compress_suffixes, skip_count, + config->compression_threads); } // Process one file in a chunk: either via incremental check or direct send. @@ -430,7 +433,8 @@ static int send_single_file(Client* client, File* file, Config* config, bool use // Fall through: send full file via sendfile (pass 0 for compression_level) int skip_count = config->skip_compress_set ? config->skip_compress_count : -1; if (!file_send_sendfile_with_skip(file, client->file_descriptor, config->use_metadata, 0, false, - config->skip_compress_suffixes, skip_count)) + config->skip_compress_suffixes, skip_count, + config->compression_threads)) return -1; return 0; } @@ -466,7 +470,7 @@ static int send_single_file(Client* client, File* file, Config* config, bool use int skip_count = config->skip_compress_set ? config->skip_compress_count : -1; if (!file_send_single_calls_with_skip(file, client->file_descriptor, config->use_metadata, compression_level, false, config->skip_compress_suffixes, - skip_count)) + skip_count, config->compression_threads)) return -1; return 0; } @@ -484,7 +488,8 @@ static int send_chunk_with_removal(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 b0bb6b9..432220c 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 6fffcff..e38b003 100644 --- a/src/client/usage.c +++ b/src/client/usage.c @@ -90,6 +90,7 @@ void print_usage(void) { printf(" --compress-level Compression level (default: 5)\n"); printf(" --zl Alias for --compress-level\n"); printf(" --skip-compress=LIST Skip compression for comma-separated suffixes\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 fb6fd72..37eb2ff 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) { - log_debug_message(LOG_DEBUG_PACK, "Starting to compress chunk"); + 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 f512278..96221b3 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) @@ -38,7 +39,15 @@ bool compression_should_skip_with_suffixes(const char* path, char* const* suffix } Data* data_compress(Data* data_to_compress, int compression_level) { - log_debug_message(LOG_DEBUG_UTIL, "Starting to compress data"); + 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 || 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); Data* compressed_data = data_create_empty(dst_size); if (compressed_data == NULL) @@ -59,6 +68,30 @@ Data* data_compress(Data* data_to_compress, int compression_level) { return NULL; } + if (compression_threads > 0) { + 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)); + 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 2cbb17a..b30622d 100644 --- a/src/shared/compression.h +++ b/src/shared/compression.h @@ -4,7 +4,11 @@ #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); 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 7a93fa2..e0cb796 100644 --- a/src/shared/config.c +++ b/src/shared/config.c @@ -24,6 +24,7 @@ static void config_set_defaults(Config* config) { config->remove_source_files = 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 6778ae1..1499efc 100644 --- a/src/shared/config.h +++ b/src/shared/config.h @@ -24,6 +24,7 @@ typedef struct Config { bool remove_source_files; 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 28bfca6..4d628bf 100644 --- a/src/shared/file_send.c +++ b/src/shared/file_send.c @@ -20,19 +20,21 @@ 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_skip(file, file_descriptor, use_metadata, compression_level, - send_path, NULL, -1); + send_path, NULL, -1, 0); } bool file_send_single_calls_with_skip(File* file, int file_descriptor, bool use_metadata, int compression_level, bool send_path, - char* const* skip_suffixes, int skip_count) { + char* const* skip_suffixes, int skip_count, + int compression_threads) { 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_with_suffixes(file->path, skip_suffixes, skip_count)) { - 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; @@ -58,17 +60,18 @@ bool file_send_single_calls_with_skip(File* file, int file_descriptor, bool use_ bool file_send_sendfile(File* file, int file_descriptor, bool use_metadata, int compression_level, bool send_path) { return file_send_sendfile_with_skip(file, file_descriptor, use_metadata, compression_level, - send_path, NULL, -1); + send_path, NULL, -1, 0); } bool file_send_sendfile_with_skip(File* file, int file_descriptor, bool use_metadata, int compression_level, bool send_path, char* const* skip_suffixes, - int skip_count) { + int skip_count, int compression_threads) { if (!file || !file->path || !file->data) return false; if (compression_level > 0) return file_send_single_calls_with_skip(file, file_descriptor, use_metadata, compression_level, - send_path, skip_suffixes, skip_count); + send_path, skip_suffixes, skip_count, + compression_threads); if (send_path && !send_str(file_descriptor, file->path)) return false; diff --git a/src/shared/file_send.h b/src/shared/file_send.h index 093195c..67585fc 100644 --- a/src/shared/file_send.h +++ b/src/shared/file_send.h @@ -10,11 +10,12 @@ 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_skip(File* file, int file_descriptor, bool use_metadata, int compression_level, bool send_path, - char* const* skip_suffixes, int skip_count); + char* const* skip_suffixes, int skip_count, + int compression_threads); bool file_send_sendfile(File* file, int file_descriptor, bool use_metadata, int compression_level, bool send_path); bool file_send_sendfile_with_skip(File* file, int file_descriptor, bool use_metadata, int compression_level, bool send_path, char* const* skip_suffixes, - int skip_count); + int skip_count, int compression_threads); #endif 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 5fbf18c..7c54a33 100644 --- a/tests/test_client_cli.c +++ b/tests/test_client_cli.c @@ -42,6 +42,12 @@ static void test_validate_config_incompatible_options() { cfg->use_incremental = false; cfg->skip_compress_set = true; EXPECT_FALSE(validate_config(cfg)); + cfg->skip_compress_set = 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); } @@ -455,6 +461,36 @@ static void test_parse_args_empty_skip_compress() { 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); + + 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 */ static void test_parse_args_unknown_option() { Config* cfg = config_create(); @@ -838,6 +874,7 @@ void test_client_cli() { test_parse_args_rejects_invalid_modify_window(); test_parse_args_skip_compress(); test_parse_args_empty_skip_compress(); + test_parse_args_compression_threads(); test_parse_args_unknown_option(); test_parse_args_rejects_unimplemented_options(); test_parse_args_quiet(); diff --git a/tests/test_compression.c b/tests/test_compression.c index 8877f13..8784046 100644 --- a/tests/test_compression.c +++ b/tests/test_compression.c @@ -63,6 +63,23 @@ static void test_skip_compress_suffix_matching() { EXPECT_FALSE(compression_should_skip_with_suffixes("archive.zip", suffixes, 0)); } +static void test_data_compress_with_threads_roundtrip() { + const size_t size = 8 * 1024 * 1024; + Data* input = data_create_empty(size); + EXPECT_NOT_NULL(input); + 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)size); + EXPECT_EQ_INT(memcmp(decompressed->data, input->data, size), 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"; @@ -124,5 +141,6 @@ void test_compression() { test_data_compress_decompress_roundtrip(); test_data_compress_decompress_large(); test_skip_compress_suffix_matching(); + test_data_compress_with_threads_roundtrip(); test_chunk_compress_decompress_roundtrip(); }