Merge remote-tracking branch 'origin/feat/compress-threads' into dev
This commit is contained in:
@@ -346,6 +346,7 @@ features without changing the meaning of ordinary compatibility options.
|
||||
| `--zc <alg>` | Alias for `--compress-choice`. FastSync supports `zstd` and `none`. |
|
||||
| `--zl <n>` | Alias for `--compress-level`. |
|
||||
| `--skip-compress <list>` | Skip compression for comma-separated suffixes; incompatible with `-s`. |
|
||||
| `--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. |
|
||||
| `-s` | Enable FastSync chunk serialization. |
|
||||
| `-f`, `--sendfile` | Use TCP `sendfile()` zero-copy transfer. Incompatible with compression and chunk serialization. |
|
||||
|
||||
@@ -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 : "<allocation failed>");
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -90,6 +90,7 @@ void print_usage(void) {
|
||||
printf(" --compress-level <n> Compression level (default: 5)\n");
|
||||
printf(" --zl <n> Alias for --compress-level\n");
|
||||
printf(" --skip-compress=LIST Skip compression for comma-separated suffixes\n");
|
||||
printf(" --compress-threads <n> Compression worker threads (requires zstd threaded support)\n");
|
||||
printf(" --help Show this help\n");
|
||||
printf(" -V, --version Show version\n");
|
||||
}
|
||||
|
||||
+7
-2
@@ -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;
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -6,6 +6,7 @@
|
||||
#include <stdint.h>
|
||||
#include <string.h>
|
||||
#include <strings.h>
|
||||
#include <unistd.h>
|
||||
#include <zstd.h>
|
||||
|
||||
#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};
|
||||
|
||||
|
||||
@@ -4,7 +4,11 @@
|
||||
#include "data.h"
|
||||
#include <stdbool.h>
|
||||
|
||||
#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);
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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"]
|
||||
|
||||
@@ -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();
|
||||
|
||||
@@ -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();
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user