fix: enforce max alloc across protocol workers
CI / lint (pull_request) Successful in 12s
CI / sanitizers (address) (pull_request) Successful in 36s
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 33s
CI / lint (pull_request) Successful in 12s
CI / sanitizers (address) (pull_request) Successful in 36s
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 33s
This commit is contained in:
@@ -456,6 +456,7 @@ static int send_chunks_multithreaded(void* pipeline_context) {
|
|||||||
|
|
||||||
static int scan_directory_multithreaded(void* pipeline_context) {
|
static int scan_directory_multithreaded(void* pipeline_context) {
|
||||||
PipelineContextSender* context = (PipelineContextSender*)pipeline_context;
|
PipelineContextSender* context = (PipelineContextSender*)pipeline_context;
|
||||||
|
protocol_session_bind(&context->allocation_session);
|
||||||
ScannerOptions options = scanner_options_from_config(context->config, 4);
|
ScannerOptions options = scanner_options_from_config(context->config, 4);
|
||||||
ParallelScanner* scanner =
|
ParallelScanner* scanner =
|
||||||
parallel_scanner_create_with_options(context->config->send_directory, &options);
|
parallel_scanner_create_with_options(context->config->send_directory, &options);
|
||||||
@@ -464,6 +465,7 @@ static int scan_directory_multithreaded(void* pipeline_context) {
|
|||||||
if (scanner == NULL) {
|
if (scanner == NULL) {
|
||||||
log_message(LOG_LEVEL_ERROR, "Failed to create parallel scanner");
|
log_message(LOG_LEVEL_ERROR, "Failed to create parallel scanner");
|
||||||
pipeline_cancel(context);
|
pipeline_cancel(context);
|
||||||
|
protocol_session_unbind();
|
||||||
return thrd_error;
|
return thrd_error;
|
||||||
}
|
}
|
||||||
while ((current_chunk = parallel_scanner_next(scanner)) != NULL) {
|
while ((current_chunk = parallel_scanner_next(scanner)) != NULL) {
|
||||||
@@ -475,6 +477,7 @@ static int scan_directory_multithreaded(void* pipeline_context) {
|
|||||||
pipeline_cancel(context);
|
pipeline_cancel(context);
|
||||||
chunk_destroy(current_chunk);
|
chunk_destroy(current_chunk);
|
||||||
parallel_scanner_destroy(scanner);
|
parallel_scanner_destroy(scanner);
|
||||||
|
protocol_session_unbind();
|
||||||
return thrd_error;
|
return thrd_error;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -485,6 +488,7 @@ static int scan_directory_multithreaded(void* pipeline_context) {
|
|||||||
chunk_destroy(current_chunk);
|
chunk_destroy(current_chunk);
|
||||||
pipeline_cancel(context);
|
pipeline_cancel(context);
|
||||||
parallel_scanner_destroy(scanner);
|
parallel_scanner_destroy(scanner);
|
||||||
|
protocol_session_unbind();
|
||||||
return thrd_error;
|
return thrd_error;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -496,6 +500,7 @@ static int scan_directory_multithreaded(void* pipeline_context) {
|
|||||||
cnd_broadcast(&context->condition_not_full_scanner);
|
cnd_broadcast(&context->condition_not_full_scanner);
|
||||||
mtx_unlock(&context->mutex_scanner);
|
mtx_unlock(&context->mutex_scanner);
|
||||||
pipeline_cancel(context);
|
pipeline_cancel(context);
|
||||||
|
protocol_session_unbind();
|
||||||
return thrd_error;
|
return thrd_error;
|
||||||
}
|
}
|
||||||
mtx_lock(&context->mutex_scanner);
|
mtx_lock(&context->mutex_scanner);
|
||||||
@@ -504,11 +509,13 @@ static int scan_directory_multithreaded(void* pipeline_context) {
|
|||||||
mtx_unlock(&context->mutex_scanner);
|
mtx_unlock(&context->mutex_scanner);
|
||||||
|
|
||||||
parallel_scanner_destroy(scanner);
|
parallel_scanner_destroy(scanner);
|
||||||
|
protocol_session_unbind();
|
||||||
return thrd_success;
|
return thrd_success;
|
||||||
}
|
}
|
||||||
|
|
||||||
static int load_files_multithreaded(void* pipeline_context) {
|
static int load_files_multithreaded(void* pipeline_context) {
|
||||||
PipelineContextSender* context = (PipelineContextSender*)pipeline_context;
|
PipelineContextSender* context = (PipelineContextSender*)pipeline_context;
|
||||||
|
protocol_session_bind(&context->allocation_session);
|
||||||
while (true) {
|
while (true) {
|
||||||
Chunk* chunk = queue_dequeue_multithreaded(
|
Chunk* chunk = queue_dequeue_multithreaded(
|
||||||
context->queue_scanner, &context->mutex_scanner, &context->condition_not_empty_scanner,
|
context->queue_scanner, &context->mutex_scanner, &context->condition_not_empty_scanner,
|
||||||
@@ -518,6 +525,7 @@ static int load_files_multithreaded(void* pipeline_context) {
|
|||||||
context->loader_done = true;
|
context->loader_done = true;
|
||||||
cnd_signal(&context->condition_not_empty_loader);
|
cnd_signal(&context->condition_not_empty_loader);
|
||||||
mtx_unlock(&context->mutex_loader);
|
mtx_unlock(&context->mutex_loader);
|
||||||
|
protocol_session_unbind();
|
||||||
return thrd_success;
|
return thrd_success;
|
||||||
}
|
}
|
||||||
if (!context->config->use_sendfile) {
|
if (!context->config->use_sendfile) {
|
||||||
@@ -529,6 +537,7 @@ static int load_files_multithreaded(void* pipeline_context) {
|
|||||||
log_message(LOG_LEVEL_ERROR, "Failed to load file data");
|
log_message(LOG_LEVEL_ERROR, "Failed to load file data");
|
||||||
chunk_destroy(chunk);
|
chunk_destroy(chunk);
|
||||||
pipeline_cancel(context);
|
pipeline_cancel(context);
|
||||||
|
protocol_session_unbind();
|
||||||
return thrd_error;
|
return thrd_error;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -539,6 +548,7 @@ static int load_files_multithreaded(void* pipeline_context) {
|
|||||||
&context->cancelled)) {
|
&context->cancelled)) {
|
||||||
chunk_destroy(chunk);
|
chunk_destroy(chunk);
|
||||||
pipeline_cancel(context);
|
pipeline_cancel(context);
|
||||||
|
protocol_session_unbind();
|
||||||
return thrd_error;
|
return thrd_error;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -294,6 +294,7 @@ void handler(int file_descriptor) {
|
|||||||
protocol_session_unbind();
|
protocol_session_unbind();
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
protocol_session_set_max_alloc(&context->session, config->max_alloc);
|
||||||
context->session.total_allocated_bytes = session.total_allocated_bytes;
|
context->session.total_allocated_bytes = session.total_allocated_bytes;
|
||||||
thrd_t receiver, writer;
|
thrd_t receiver, writer;
|
||||||
bool receiver_created = thrd_create(&receiver, receive_thread, context) == thrd_success;
|
bool receiver_created = thrd_create(&receiver, receive_thread, context) == thrd_success;
|
||||||
|
|||||||
@@ -262,6 +262,8 @@ static bool receive_core_fields(int fd, Config* c) {
|
|||||||
int value;
|
int value;
|
||||||
if (!receive_n_data(fd, &c->max_alloc, sizeof(c->max_alloc)) || c->max_alloc == 0)
|
if (!receive_n_data(fd, &c->max_alloc, sizeof(c->max_alloc)) || c->max_alloc == 0)
|
||||||
return false;
|
return false;
|
||||||
|
if (c->max_alloc > MAX_SERVER_ALLOC)
|
||||||
|
c->max_alloc = MAX_SERVER_ALLOC;
|
||||||
protocol_session_set_max_alloc(NULL, c->max_alloc);
|
protocol_session_set_max_alloc(NULL, c->max_alloc);
|
||||||
c->send_directory = receive_str(fd);
|
c->send_directory = receive_str(fd);
|
||||||
c->receive_root_directory = receive_str(fd);
|
c->receive_root_directory = receive_str(fd);
|
||||||
|
|||||||
+1
-1
@@ -129,7 +129,7 @@ typedef struct Config {
|
|||||||
char* compress_choice;
|
char* compress_choice;
|
||||||
} Config;
|
} Config;
|
||||||
|
|
||||||
#define PROTOCOL_VERSION "2.2.0"
|
#define PROTOCOL_VERSION "2.3.0"
|
||||||
#define DEFAULT_CHUNK_SIZE (10 * 1024 * 1024)
|
#define DEFAULT_CHUNK_SIZE (10 * 1024 * 1024)
|
||||||
|
|
||||||
Config* config_create(void);
|
Config* config_create(void);
|
||||||
|
|||||||
@@ -29,6 +29,8 @@ PipelineContextSender* pipeline_context_sender_create(Config* config, Queue* que
|
|||||||
context->progress_bytes = 0;
|
context->progress_bytes = 0;
|
||||||
context->sender_done = false;
|
context->sender_done = false;
|
||||||
atomic_init(&context->cancelled, false);
|
atomic_init(&context->cancelled, false);
|
||||||
|
protocol_session_init(&context->allocation_session, -1, -1);
|
||||||
|
protocol_session_set_max_alloc(&context->allocation_session, config->max_alloc);
|
||||||
int init = 0;
|
int init = 0;
|
||||||
if (mtx_init(&context->mutex_scanner, mtx_plain) != thrd_success)
|
if (mtx_init(&context->mutex_scanner, mtx_plain) != thrd_success)
|
||||||
goto fail;
|
goto fail;
|
||||||
|
|||||||
@@ -28,6 +28,7 @@ typedef struct {
|
|||||||
unsigned long long progress_bytes;
|
unsigned long long progress_bytes;
|
||||||
bool sender_done;
|
bool sender_done;
|
||||||
atomic_bool cancelled;
|
atomic_bool cancelled;
|
||||||
|
ProtocolSession allocation_session;
|
||||||
} PipelineContextSender;
|
} PipelineContextSender;
|
||||||
|
|
||||||
typedef struct PipelineContextReceiver {
|
typedef struct PipelineContextReceiver {
|
||||||
|
|||||||
@@ -19,6 +19,8 @@
|
|||||||
/* Aggregate bytes retained by one received deletion manifest. */
|
/* Aggregate bytes retained by one received deletion manifest. */
|
||||||
#define MAX_MANIFEST_BYTES (16ULL * 1024 * 1024)
|
#define MAX_MANIFEST_BYTES (16ULL * 1024 * 1024)
|
||||||
#define DEFAULT_MAX_ALLOC (1ULL * 1024 * 1024 * 1024)
|
#define DEFAULT_MAX_ALLOC (1ULL * 1024 * 1024 * 1024)
|
||||||
|
/* Server policy ceiling for a client-provided allocation limit. */
|
||||||
|
#define MAX_SERVER_ALLOC (256ULL * 1024 * 1024)
|
||||||
|
|
||||||
typedef struct ssl_st SSL;
|
typedef struct ssl_st SSL;
|
||||||
|
|
||||||
|
|||||||
+4
-3
@@ -95,6 +95,7 @@ static void test_pipeline_sender_lifecycle() {
|
|||||||
EXPECT_EQ_INT(pcs->queue_loader->capacity, 15);
|
EXPECT_EQ_INT(pcs->queue_loader->capacity, 15);
|
||||||
EXPECT_FALSE(pcs->scanner_done);
|
EXPECT_FALSE(pcs->scanner_done);
|
||||||
EXPECT_FALSE(pcs->loader_done);
|
EXPECT_FALSE(pcs->loader_done);
|
||||||
|
EXPECT_EQ_INT((int)pcs->allocation_session.max_alloc, (int)cfg->max_alloc);
|
||||||
|
|
||||||
pipeline_context_sender_destroy(pcs);
|
pipeline_context_sender_destroy(pcs);
|
||||||
}
|
}
|
||||||
@@ -126,7 +127,7 @@ static void test_config_send_receive() {
|
|||||||
send_cfg->use_metadata = true;
|
send_cfg->use_metadata = true;
|
||||||
send_cfg->compression_level = 5;
|
send_cfg->compression_level = 5;
|
||||||
send_cfg->chunk_size = 1024;
|
send_cfg->chunk_size = 1024;
|
||||||
send_cfg->max_alloc = 8ULL * 1024 * 1024;
|
send_cfg->max_alloc = MAX_SERVER_ALLOC + 1;
|
||||||
|
|
||||||
/* Use socketpair for bidirectional communication */
|
/* Use socketpair for bidirectional communication */
|
||||||
int p[2];
|
int p[2];
|
||||||
@@ -161,7 +162,7 @@ static void test_config_send_receive() {
|
|||||||
ok = false;
|
ok = false;
|
||||||
if (recv_cfg->chunk_size != 1024)
|
if (recv_cfg->chunk_size != 1024)
|
||||||
ok = false;
|
ok = false;
|
||||||
if (recv_cfg->max_alloc != 8ULL * 1024 * 1024)
|
if (recv_cfg->max_alloc != MAX_SERVER_ALLOC)
|
||||||
ok = false;
|
ok = false;
|
||||||
}
|
}
|
||||||
config_delete(recv_cfg);
|
config_delete(recv_cfg);
|
||||||
@@ -192,7 +193,7 @@ static void test_config_send_receive_version_mismatch() {
|
|||||||
Config* cfg = config_create();
|
Config* cfg = config_create();
|
||||||
EXPECT_NOT_NULL(cfg);
|
EXPECT_NOT_NULL(cfg);
|
||||||
free(cfg->version);
|
free(cfg->version);
|
||||||
cfg->version = str_dup("0.0");
|
cfg->version = str_dup("2.2.0");
|
||||||
cfg->send_directory = str_dup("/src");
|
cfg->send_directory = str_dup("/src");
|
||||||
cfg->receive_root_directory = str_dup("/dst");
|
cfg->receive_root_directory = str_dup("/dst");
|
||||||
|
|
||||||
|
|||||||
@@ -202,6 +202,19 @@ static void test_max_alloc_rejects_single_buffer() {
|
|||||||
close(p[1]);
|
close(p[1]);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
static void test_max_alloc_allows_configured_buffer() {
|
||||||
|
ProtocolSession session;
|
||||||
|
protocol_session_init(&session, -1, -1);
|
||||||
|
protocol_session_set_max_alloc(&session, 4);
|
||||||
|
protocol_session_bind(&session);
|
||||||
|
void* allowed = protocol_alloc(4);
|
||||||
|
const void* rejected = protocol_alloc(5);
|
||||||
|
EXPECT_NOT_NULL(allowed);
|
||||||
|
EXPECT_NULL(rejected);
|
||||||
|
free(allowed);
|
||||||
|
protocol_session_unbind();
|
||||||
|
}
|
||||||
|
|
||||||
void test_protocol() {
|
void test_protocol() {
|
||||||
test_send_receive_n_data();
|
test_send_receive_n_data();
|
||||||
test_send_receive_n_data_zero();
|
test_send_receive_n_data_zero();
|
||||||
@@ -214,4 +227,5 @@ void test_protocol() {
|
|||||||
test_receive_n_data_truncated();
|
test_receive_n_data_truncated();
|
||||||
test_receive_str_truncated();
|
test_receive_str_truncated();
|
||||||
test_max_alloc_rejects_single_buffer();
|
test_max_alloc_rejects_single_buffer();
|
||||||
|
test_max_alloc_allows_configured_buffer();
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user