diff --git a/src/client/client_send.c b/src/client/client_send.c index 0b4ca5c..41c54ed 100644 --- a/src/client/client_send.c +++ b/src/client/client_send.c @@ -456,6 +456,7 @@ static int send_chunks_multithreaded(void* pipeline_context) { static int scan_directory_multithreaded(void* pipeline_context) { PipelineContextSender* context = (PipelineContextSender*)pipeline_context; + protocol_session_bind(&context->allocation_session); ScannerOptions options = scanner_options_from_config(context->config, 4); ParallelScanner* scanner = 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) { log_message(LOG_LEVEL_ERROR, "Failed to create parallel scanner"); pipeline_cancel(context); + protocol_session_unbind(); return thrd_error; } while ((current_chunk = parallel_scanner_next(scanner)) != NULL) { @@ -475,6 +477,7 @@ static int scan_directory_multithreaded(void* pipeline_context) { pipeline_cancel(context); chunk_destroy(current_chunk); parallel_scanner_destroy(scanner); + protocol_session_unbind(); return thrd_error; } } @@ -485,6 +488,7 @@ static int scan_directory_multithreaded(void* pipeline_context) { chunk_destroy(current_chunk); pipeline_cancel(context); parallel_scanner_destroy(scanner); + protocol_session_unbind(); return thrd_error; } } @@ -496,6 +500,7 @@ static int scan_directory_multithreaded(void* pipeline_context) { cnd_broadcast(&context->condition_not_full_scanner); mtx_unlock(&context->mutex_scanner); pipeline_cancel(context); + protocol_session_unbind(); return thrd_error; } mtx_lock(&context->mutex_scanner); @@ -504,11 +509,13 @@ static int scan_directory_multithreaded(void* pipeline_context) { mtx_unlock(&context->mutex_scanner); parallel_scanner_destroy(scanner); + protocol_session_unbind(); return thrd_success; } static int load_files_multithreaded(void* pipeline_context) { PipelineContextSender* context = (PipelineContextSender*)pipeline_context; + protocol_session_bind(&context->allocation_session); while (true) { Chunk* chunk = queue_dequeue_multithreaded( 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; cnd_signal(&context->condition_not_empty_loader); mtx_unlock(&context->mutex_loader); + protocol_session_unbind(); return thrd_success; } 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"); chunk_destroy(chunk); pipeline_cancel(context); + protocol_session_unbind(); return thrd_error; } } @@ -539,6 +548,7 @@ static int load_files_multithreaded(void* pipeline_context) { &context->cancelled)) { chunk_destroy(chunk); pipeline_cancel(context); + protocol_session_unbind(); return thrd_error; } } diff --git a/src/server/server.c b/src/server/server.c index af5cd43..d7221ef 100644 --- a/src/server/server.c +++ b/src/server/server.c @@ -294,6 +294,7 @@ void handler(int file_descriptor) { protocol_session_unbind(); return; } + protocol_session_set_max_alloc(&context->session, config->max_alloc); context->session.total_allocated_bytes = session.total_allocated_bytes; thrd_t receiver, writer; bool receiver_created = thrd_create(&receiver, receive_thread, context) == thrd_success; diff --git a/src/shared/config.c b/src/shared/config.c index 968b80d..ac3bcb6 100644 --- a/src/shared/config.c +++ b/src/shared/config.c @@ -262,6 +262,8 @@ static bool receive_core_fields(int fd, Config* c) { int value; if (!receive_n_data(fd, &c->max_alloc, sizeof(c->max_alloc)) || c->max_alloc == 0) return false; + if (c->max_alloc > MAX_SERVER_ALLOC) + c->max_alloc = MAX_SERVER_ALLOC; protocol_session_set_max_alloc(NULL, c->max_alloc); c->send_directory = receive_str(fd); c->receive_root_directory = receive_str(fd); diff --git a/src/shared/config.h b/src/shared/config.h index a473e49..29e2698 100644 --- a/src/shared/config.h +++ b/src/shared/config.h @@ -129,7 +129,7 @@ typedef struct Config { char* compress_choice; } Config; -#define PROTOCOL_VERSION "2.2.0" +#define PROTOCOL_VERSION "2.3.0" #define DEFAULT_CHUNK_SIZE (10 * 1024 * 1024) Config* config_create(void); diff --git a/src/shared/multiprocessing.c b/src/shared/multiprocessing.c index 6b677f3..f58e946 100644 --- a/src/shared/multiprocessing.c +++ b/src/shared/multiprocessing.c @@ -29,6 +29,8 @@ PipelineContextSender* pipeline_context_sender_create(Config* config, Queue* que context->progress_bytes = 0; context->sender_done = 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; if (mtx_init(&context->mutex_scanner, mtx_plain) != thrd_success) goto fail; diff --git a/src/shared/multiprocessing.h b/src/shared/multiprocessing.h index 37f09d5..3fd6dc8 100644 --- a/src/shared/multiprocessing.h +++ b/src/shared/multiprocessing.h @@ -28,6 +28,7 @@ typedef struct { unsigned long long progress_bytes; bool sender_done; atomic_bool cancelled; + ProtocolSession allocation_session; } PipelineContextSender; typedef struct PipelineContextReceiver { diff --git a/src/shared/protocol.h b/src/shared/protocol.h index 908b745..9022e22 100644 --- a/src/shared/protocol.h +++ b/src/shared/protocol.h @@ -19,6 +19,8 @@ /* Aggregate bytes retained by one received deletion manifest. */ #define MAX_MANIFEST_BYTES (16ULL * 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; diff --git a/tests/test_config.c b/tests/test_config.c index 4a7f093..f31b20f 100644 --- a/tests/test_config.c +++ b/tests/test_config.c @@ -95,6 +95,7 @@ static void test_pipeline_sender_lifecycle() { EXPECT_EQ_INT(pcs->queue_loader->capacity, 15); EXPECT_FALSE(pcs->scanner_done); EXPECT_FALSE(pcs->loader_done); + EXPECT_EQ_INT((int)pcs->allocation_session.max_alloc, (int)cfg->max_alloc); pipeline_context_sender_destroy(pcs); } @@ -126,7 +127,7 @@ static void test_config_send_receive() { send_cfg->use_metadata = true; send_cfg->compression_level = 5; 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 */ int p[2]; @@ -161,7 +162,7 @@ static void test_config_send_receive() { ok = false; if (recv_cfg->chunk_size != 1024) ok = false; - if (recv_cfg->max_alloc != 8ULL * 1024 * 1024) + if (recv_cfg->max_alloc != MAX_SERVER_ALLOC) ok = false; } config_delete(recv_cfg); @@ -192,7 +193,7 @@ static void test_config_send_receive_version_mismatch() { Config* cfg = config_create(); EXPECT_NOT_NULL(cfg); free(cfg->version); - cfg->version = str_dup("0.0"); + cfg->version = str_dup("2.2.0"); cfg->send_directory = str_dup("/src"); cfg->receive_root_directory = str_dup("/dst"); diff --git a/tests/test_protocol.c b/tests/test_protocol.c index fc0583a..fb325fa 100644 --- a/tests/test_protocol.c +++ b/tests/test_protocol.c @@ -202,6 +202,19 @@ static void test_max_alloc_rejects_single_buffer() { 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() { test_send_receive_n_data(); test_send_receive_n_data_zero(); @@ -214,4 +227,5 @@ void test_protocol() { test_receive_n_data_truncated(); test_receive_str_truncated(); test_max_alloc_rejects_single_buffer(); + test_max_alloc_allows_configured_buffer(); }