feat: add compression worker threads
CI / lint (pull_request) Successful in 11s
CI / sanitizers (address) (pull_request) Successful in 38s
CI / sanitizers (undefined) (pull_request) Successful in 38s
CI / fuzz-build (pull_request) Successful in 15s
CI / coverage (pull_request) Successful in 31s
CI / build-and-test (pull_request) Successful in 1m16s
CI / valgrind (pull_request) Successful in 33s

This commit is contained in:
2026-09-03 18:15:14 +02:00
parent 190fc5d300
commit f443b2986d
16 changed files with 136 additions and 10 deletions
+1
View File
@@ -339,6 +339,7 @@ features without changing the meaning of ordinary compatibility options.
| `-m` | Enable the multithreaded scanner/loader/sender pipeline. | | `-m` | Enable the multithreaded scanner/loader/sender pipeline. |
| `-c [level]`, `-z [level]` | Enable streaming zstd compression, levels 1-22. | | `-c [level]`, `-z [level]` | Enable streaming zstd compression, levels 1-22. |
| `--compress-level <n>` | Set the zstd compression level. | | `--compress-level <n>` | Set the zstd compression level. |
| `--compress-threads <n>` | Use `n` zstd compression workers. Requires compression and a zstd build with threaded support; the setting affects sender CPU work only. |
| `--chunk-size <bytes>` | Set the transfer chunk size. | | `--chunk-size <bytes>` | Set the transfer chunk size. |
| `-s` | Enable FastSync chunk serialization. | | `-s` | Enable FastSync chunk serialization. |
| `-f`, `--sendfile` | Use TCP `sendfile()` zero-copy transfer. Incompatible with compression and chunk serialization. | | `-f`, `--sendfile` | Use TCP `sendfile()` zero-copy transfer. Incompatible with compression and chunk serialization. |
+11
View File
@@ -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 parse_args(Config* config, int argc, char* argv[], int* positional_args,
int* positional_count) { int* positional_count) {
for (int i = 1; i < argc; i++) { 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]); const OptionEntry* entry = find_table_option(argv[i]);
if (entry) { if (entry) {
if (entry->kind != OPT_FLAG) { 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"); log_message(LOG_LEVEL_ERROR, "--compress-level must be between 1 and 22");
return -1; 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] == '-') { } else if (argv[i][0] == '-') {
fprintf(stderr, "Unknown option: %s\n", argv[i]); fprintf(stderr, "Unknown option: %s\n", argv[i]);
print_usage(); print_usage();
+14 -8
View File
@@ -235,7 +235,8 @@ static int send_delta(Client* client, File* file, DeltaSignature* sig, Config* c
Data* to_send = delta_data; Data* to_send = delta_data;
if (config->use_compression) { 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); data_destroy(delta_data);
if (!to_send) if (!to_send)
return send_status(client->file_descriptor, STATUS_NEXT) ? 1 : -1; 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; 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). // 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)) if (!send_status(fd, STATUS_NEXT))
return false; 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). // 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 return send_file_direct_sendfile(file, client->file_descriptor, config->use_metadata) ? 0
: -1; : -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 ? 0
: -1; : -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) // 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; DeltaSignature* sig = NULL;
int rc = incremental_check(client, file, config, &sig); int rc = incremental_check(client, file, config, &sig);
if (rc < 0) { if (rc < 0) {
@@ -338,7 +342,8 @@ static int send_single_file(Client* client, File* file, Config* config, bool use
return -1; 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 -1;
return 0; return 0;
} }
@@ -349,7 +354,8 @@ int send_chunk(Client* client, Chunk* chunk, Config* config) {
return -1; return -1;
Data* data; Data* data;
if (config->use_compression) { 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 { } else {
data = chunk_serialize(chunk, config->use_metadata); data = chunk_serialize(chunk, config->use_metadata);
} }
+4
View File
@@ -15,6 +15,10 @@ bool validate_config(const Config* config) {
"(chunk serialization)"); "(chunk serialization)");
return false; 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) { if (config->transport == TRANSPORT_SSH && config->use_sendfile) {
log_message(LOG_LEVEL_ERROR, "-f/--sendfile is not supported with SSH transport"); log_message(LOG_LEVEL_ERROR, "-f/--sendfile is not supported with SSH transport");
return false; return false;
+1
View File
@@ -69,6 +69,7 @@ void print_usage(void) {
printf(" -S, --sparse Handle sparse files efficiently\n"); printf(" -S, --sparse Handle sparse files efficiently\n");
printf(" --inplace Update files in-place (no temp+rename)\n"); printf(" --inplace Update files in-place (no temp+rename)\n");
printf(" --compress-level <n> Compression level (default: 5)\n"); printf(" --compress-level <n> Compression level (default: 5)\n");
printf(" --compress-threads <n> Compression worker threads (requires zstd threaded support)\n");
printf(" --help Show this help\n"); printf(" --help Show this help\n");
printf(" -V, --version Show version\n"); printf(" -V, --version Show version\n");
} }
+6 -1
View File
@@ -290,11 +290,16 @@ Chunk* chunk_deserialize(Data* data, bool use_metadata) {
} }
Data* chunk_compress(Chunk* chunk, int compression_level, 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"); log_message(LOG_LEVEL_DEBUG, "Starting to compress chunk");
Data* serialized = chunk_serialize(chunk, use_metadata); Data* serialized = chunk_serialize(chunk, use_metadata);
if (serialized == NULL) if (serialized == NULL)
return NULL; return NULL;
Data* compressed = data_compress(serialized, compression_level); Data* compressed = data_compress_with_threads(serialized, compression_level, compression_threads);
data_destroy(serialized); data_destroy(serialized);
if (compressed == NULL) if (compressed == NULL)
return NULL; return NULL;
+2
View File
@@ -19,6 +19,8 @@ void chunk_destroy(void* chunk);
Data* chunk_serialize(Chunk* chunk, bool use_metadata); Data* chunk_serialize(Chunk* chunk, bool use_metadata);
Chunk* chunk_deserialize(Data* data, 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(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); Chunk* receive_chunk_data(int fd, const Config* config);
#endif #endif
+28
View File
@@ -28,6 +28,14 @@ bool compression_should_skip(const char* path) {
} }
Data* data_compress(Data* data_to_compress, int compression_level) { 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"); log_message(LOG_LEVEL_DEBUG, "Starting to compress data");
size_t dst_size = ZSTD_compressBound(data_to_compress->size); size_t dst_size = ZSTD_compressBound(data_to_compress->size);
Data* compressed_data = data_create_empty(dst_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; 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_inBuffer input = {data_to_compress->data, data_to_compress->size, 0};
ZSTD_outBuffer output = {compressed_data->data, dst_size, 0}; ZSTD_outBuffer output = {compressed_data->data, dst_size, 0};
+2
View File
@@ -5,6 +5,8 @@
#include <stdbool.h> #include <stdbool.h>
Data* data_compress(Data* data_to_compress, int compression_level); 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(Data* compressed_data);
Data* data_decompress_limited(Data* compressed_data, size_t maximum_size); Data* data_decompress_limited(Data* compressed_data, size_t maximum_size);
bool compression_should_skip(const char* path); bool compression_should_skip(const char* path);
+1
View File
@@ -21,6 +21,7 @@ static void config_set_defaults(Config* config) {
config->dry_run = false; config->dry_run = false;
config->use_delete = false; config->use_delete = false;
config->compression_level = 5; config->compression_level = 5;
config->compression_threads = 0;
config->use_sendfile = false; config->use_sendfile = false;
config->chunk_size = DEFAULT_CHUNK_SIZE; config->chunk_size = DEFAULT_CHUNK_SIZE;
config->ssh_port = 22; config->ssh_port = 22;
+1
View File
@@ -22,6 +22,7 @@ typedef struct Config {
bool dry_run; bool dry_run;
bool use_delete; bool use_delete;
int compression_level; int compression_level;
int compression_threads;
unsigned long long chunk_size; unsigned long long chunk_size;
int ssh_port; int ssh_port;
TransportType transport; TransportType transport;
+9 -1
View File
@@ -19,12 +19,20 @@
bool file_send_single_calls(File* file, int file_descriptor, bool use_metadata, bool file_send_single_calls(File* file, int file_descriptor, bool use_metadata,
int compression_level, bool send_path) { 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)) if (!file || !file->path || !file->data || (file->data->size != 0 && !file->data->data))
return false; return false;
const Data* data_to_send = file->data; const Data* data_to_send = file->data;
Data* compressed_data = NULL; Data* compressed_data = NULL;
if (compression_level > 0 && !compression_should_skip(file->path)) { 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) { if (compressed_data == NULL) {
log_message(LOG_LEVEL_ERROR, "Failed to compress file data"); log_message(LOG_LEVEL_ERROR, "Failed to compress file data");
return false; return false;
+3
View File
@@ -8,6 +8,9 @@
bool file_send_single_calls(File* file, int file_descriptor, bool use_metadata, bool file_send_single_calls(File* file, int file_descriptor, bool use_metadata,
int compression_level, bool send_path); 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 file_send_sendfile(File* file, int file_descriptor, bool use_metadata, int compression_level,
bool send_path); bool send_path);
+5
View File
@@ -71,6 +71,11 @@ class TestTCPFlags:
r = _run_tcp_test("Compression (-c)", shared_server.port, ["-c"]) r = _run_tcp_test("Compression (-c)", shared_server.port, ["-c"])
assert r["status"] == "Success", r["error"] 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): def test_chunk_serialization(self, shared_server):
r = _run_tcp_test("Chunk Serialization (-s)", shared_server.port, ["-s"]) r = _run_tcp_test("Chunk Serialization (-s)", shared_server.port, ["-s"])
assert r["status"] == "Success", r["error"] assert r["status"] == "Success", r["error"]
+31
View File
@@ -36,6 +36,12 @@ static void test_validate_config_incompatible_options() {
cfg->use_incremental = true; cfg->use_incremental = true;
cfg->use_chunk_serialization = true; cfg->use_chunk_serialization = true;
EXPECT_FALSE(validate_config(cfg)); 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); config_delete(cfg);
} }
@@ -247,6 +253,30 @@ static void test_parse_args_valid_compression_level() {
config_delete(cfg); 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 */ /* Test parse_args unknown option returns error */
static void test_parse_args_unknown_option() { static void test_parse_args_unknown_option() {
Config* cfg = config_create(); Config* cfg = config_create();
@@ -357,6 +387,7 @@ void test_client_cli() {
test_parse_args_invalid_server_port(); test_parse_args_invalid_server_port();
test_parse_args_invalid_compression_level(); test_parse_args_invalid_compression_level();
test_parse_args_valid_compression_level(); test_parse_args_valid_compression_level();
test_parse_args_compression_threads();
test_parse_args_unknown_option(); test_parse_args_unknown_option();
test_parse_args_rejects_unimplemented_options(); test_parse_args_rejects_unimplemented_options();
test_parse_args_archive(); test_parse_args_archive();
+17
View File
@@ -55,6 +55,22 @@ static void test_data_compress_decompress_large() {
data_destroy(decompressed); 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() { static void test_chunk_compress_decompress_roundtrip() {
char* path1 = "temp_comp_test_1.txt"; char* path1 = "temp_comp_test_1.txt";
char* content1 = "chunk compression test file 1"; char* content1 = "chunk compression test file 1";
@@ -115,5 +131,6 @@ static void test_chunk_compress_decompress_roundtrip() {
void test_compression() { void test_compression() {
test_data_compress_decompress_roundtrip(); test_data_compress_decompress_roundtrip();
test_data_compress_decompress_large(); test_data_compress_decompress_large();
test_data_compress_with_threads_roundtrip();
test_chunk_compress_decompress_roundtrip(); test_chunk_compress_decompress_roundtrip();
} }