fix: bound 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 37s
CI / fuzz-build (pull_request) Successful in 15s
CI / coverage (pull_request) Successful in 32s
CI / build-and-test (pull_request) Successful in 1m15s
CI / valgrind (pull_request) Successful in 34s
CI / lint (pull_request) Successful in 11s
CI / sanitizers (address) (pull_request) Successful in 38s
CI / sanitizers (undefined) (pull_request) Successful in 37s
CI / fuzz-build (pull_request) Successful in 15s
CI / coverage (pull_request) Successful in 32s
CI / build-and-test (pull_request) Successful in 1m15s
CI / valgrind (pull_request) Successful in 34s
This commit is contained in:
+15
-4
@@ -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]);
|
||||
|
||||
@@ -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)
|
||||
@@ -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));
|
||||
|
||||
@@ -4,6 +4,8 @@
|
||||
#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);
|
||||
|
||||
@@ -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 */
|
||||
|
||||
@@ -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);
|
||||
|
||||
Reference in New Issue
Block a user