diff --git a/src/server/server.c b/src/server/server.c index c96612d..71a3ab6 100644 --- a/src/server/server.c +++ b/src/server/server.c @@ -4,6 +4,7 @@ #include "log.h" #include "metadata.h" #include "multiprocessing.h" +#include "protocol.h" #include "queue.h" #include "receiver.h" #include "transport_tcp.h" @@ -26,6 +27,13 @@ static bool allow_delete; static bool allow_unauthenticated; static const char* required_client_cn; +/* Aggregate payload bytes the multithreaded receiver may buffer ahead of the + slow disk writer. Receiving one more chunk adds up to ~2 * MAX_CHUNK_SIZE + of transient wire/decompression buffers on top of the queued payloads, so + this ceiling keeps total per-connection receive memory (decompressed and + per-file copied chunk buffers included) within MAX_CONNECTION_MEMORY. */ +#define RECEIVER_QUEUE_MAX_BYTES (MAX_CONNECTION_MEMORY - 2 * MAX_CHUNK_SIZE) + static bool tls_client_identity_allowed(SSL* ssl) { if (!ssl || !required_client_cn) return false; @@ -314,6 +322,7 @@ void handler(int file_descriptor) { protocol_session_set_max_alloc(&context->session, config->max_alloc); atomic_store(&context->session.total_allocated_bytes, atomic_load(&session.total_allocated_bytes)); + pipeline_context_receiver_set_queue_byte_limit(context, RECEIVER_QUEUE_MAX_BYTES); thrd_t receiver, writer; bool receiver_created = thrd_create(&receiver, receive_thread, context) == thrd_success; bool writer_created = false; diff --git a/src/shared/file.c b/src/shared/file.c index f8a7f84..cd87726 100644 --- a/src/shared/file.c +++ b/src/shared/file.c @@ -348,10 +348,25 @@ static bool file_to_disk_secure_impl(const char* path, const void* data, if (newer) { ok = true; } else { - if (!sparse || data_size == 0 || ftruncate(fd, (off_t)data_size) == 0) + /* In-place overwrites: pre-size sparse targets and always trim the + file to the new payload length afterwards so shorter payloads can + never leave stale trailing bytes from a previous version. */ + if (sparse && data_size > 0) + ok = ftruncate(fd, (off_t)data_size) == 0; + if (ok || !sparse || data_size == 0) ok = write_all(fd, data, data_size); - if (ok && metadata) - ok = file_restore_metadata_fd(fd, metadata, preserve_executability); + if (ok) + ok = ftruncate(fd, (off_t)data_size) == 0; + /* Normalize the mode: apply the metadata-derived safe mode when the + sender supplied metadata (setuid/setgid/sticky are never honored); + otherwise fall back to a safe default so dangerous bits on an + existing destination cannot survive an overwrite. */ + if (ok) { + if (metadata) + ok = file_restore_metadata_fd(fd, metadata, preserve_executability); + else if (fchmod(fd, S_IRUSR | S_IWUSR | S_IRGRP | S_IROTH) != 0) + ok = false; + } if (ok && use_fsync) ok = fsync(fd) == 0; } diff --git a/src/shared/multiprocessing.c b/src/shared/multiprocessing.c index ace0c9d..06c16f0 100644 --- a/src/shared/multiprocessing.c +++ b/src/shared/multiprocessing.c @@ -111,6 +111,8 @@ PipelineContextReceiver* pipeline_context_receiver_create(Config* config, Queue* protocol_session_init(&context->session, file_descriptor, file_descriptor); protocol_session_set_ssl(&context->session, ssl); context->receiver_done = false; + context->queued_bytes = 0; + context->max_queue_bytes = 0; atomic_init(&context->cancelled, false); int init = 0; if (mtx_init(&context->mutex, mtx_plain) != thrd_success) @@ -147,14 +149,73 @@ void pipeline_context_receiver_destroy(PipelineContextReceiver* context) { free(context); } +void pipeline_context_receiver_set_queue_byte_limit(PipelineContextReceiver* context, + size_t max_bytes) { + if (context == NULL) + return; + mtx_lock(&context->mutex); + context->max_queue_bytes = max_bytes; + context->queued_bytes = 0; + cnd_broadcast(&context->condition_not_full); + mtx_unlock(&context->mutex); +} + +void pipeline_context_receiver_note_bytes_released(PipelineContextReceiver* context, + size_t released_bytes) { + if (context == NULL || context->max_queue_bytes == 0 || released_bytes == 0) + return; + mtx_lock(&context->mutex); + if (released_bytes >= context->queued_bytes) + context->queued_bytes = 0; + else + context->queued_bytes -= released_bytes; + cnd_signal(&context->condition_not_full); + mtx_unlock(&context->mutex); +} + +bool pipeline_context_receiver_enqueue_file(PipelineContextReceiver* context, File* file) { + if (context == NULL || file == NULL) + return false; + size_t file_bytes = file->data ? file->data->size : 0; + mtx_lock(&context->mutex); + while (!atomic_load(&context->cancelled)) { + bool blocked_by_count = queue_is_full(context->queue); + bool blocked_by_budget = false; + if (context->max_queue_bytes > 0) { + size_t budget = context->max_queue_bytes; + size_t used = context->queued_bytes; + if (used >= budget) { + blocked_by_budget = true; + } else if (file_bytes > budget - used) { + /* A single payload larger than the whole budget (not possible with + the per-file receive cap) is only admitted to an empty pipeline so + the wait can never deadlock. */ + blocked_by_budget = used != 0; + } + } + if (!blocked_by_count && !blocked_by_budget) + break; + cnd_wait(&context->condition_not_full, &context->mutex); + } + if (atomic_load(&context->cancelled)) { + mtx_unlock(&context->mutex); + file_destroy(file); + return false; + } + if (!queue_enqueue(context->queue, file)) { + mtx_unlock(&context->mutex); + file_destroy(file); + return false; + } + context->queued_bytes += file_bytes; + cnd_signal(&context->condition_not_empty); + mtx_unlock(&context->mutex); + return true; +} + static bool receiver_enqueue_file(File* file, void* context_pointer) { - PipelineContextReceiver* context = context_pointer; - if (queue_enqueue_multithreaded_cancel(context->queue, file, &context->mutex, - &context->condition_not_empty, - &context->condition_not_full, &context->cancelled)) - return true; - file_destroy(file); - return false; + PipelineContextReceiver* context = (PipelineContextReceiver*)context_pointer; + return pipeline_context_receiver_enqueue_file(context, file); } static void receiver_thread_fail(PipelineContextReceiver* context) { @@ -215,11 +276,13 @@ int write_thread(void* pipeline_context) { protocol_session_unbind(); return thrd_success; } + size_t file_bytes = file->data ? file->data->size : 0; FileSaveResult result = FILE_SAVE_SKIPPED; if (save_to_disk) { result = file_save_to_disk_full(root_directory, file, context->config); if (result == FILE_SAVE_ERROR) { file_destroy(file); + pipeline_context_receiver_note_bytes_released(context, file_bytes); mtx_lock(&context->mutex); atomic_store(&context->cancelled, true); context->receiver_done = true; @@ -236,6 +299,7 @@ int write_thread(void* pipeline_context) { if (context->config->remove_source_files && !receiver_outcomes_append(&context->outcomes, (unsigned char)result)) { file_destroy(file); + pipeline_context_receiver_note_bytes_released(context, file_bytes); mtx_lock(&context->mutex); atomic_store(&context->cancelled, true); context->receiver_done = true; @@ -247,5 +311,6 @@ int write_thread(void* pipeline_context) { return thrd_error; } file_destroy(file); + pipeline_context_receiver_note_bytes_released(context, file_bytes); } } diff --git a/src/shared/multiprocessing.h b/src/shared/multiprocessing.h index d0e01ea..28f2fe4 100644 --- a/src/shared/multiprocessing.h +++ b/src/shared/multiprocessing.h @@ -47,6 +47,14 @@ typedef struct PipelineContextReceiver { cnd_t condition_not_empty; bool receiver_done; atomic_bool cancelled; + /* Aggregate payload bytes that have been received but not yet released by + the disk writer (queued or in the writer's hand). Guarded by `mutex`. + When `max_queue_bytes` is non-zero the receiver blocks before enqueuing + once this total would exceed it, so decompressed/copied file payloads + buffered ahead of a slow disk writer respect the per-connection memory + budget instead of growing without bound. */ + size_t queued_bytes; + size_t max_queue_bytes; } PipelineContextReceiver; PipelineContextSender* pipeline_context_sender_create(Config* config, Queue* queue_scanner, @@ -55,6 +63,18 @@ void pipeline_context_sender_destroy(PipelineContextSender* context); PipelineContextReceiver* pipeline_context_receiver_create(Config* config, Queue* queue_receiver, int file_descriptor, SSL* ssl); void pipeline_context_receiver_destroy(PipelineContextReceiver* context); +/* Bound the bytes buffered ahead of the disk writer (see max_queue_bytes). */ +void pipeline_context_receiver_set_queue_byte_limit(PipelineContextReceiver* context, + size_t max_bytes); +/* Blocking enqueue used by the receive pipeline sink. Blocks while the queue + is full by element count or when adding `file` would push queued_bytes over + the configured byte limit; waits until the disk writer releases bytes. + Takes ownership of `file` on success and destroys it on failure/cancel. */ +bool pipeline_context_receiver_enqueue_file(PipelineContextReceiver* context, File* file); +/* Account for `released_bytes` of payload memory that has been freed by the + disk writer, unblocking a receiver that is waiting on the byte limit. */ +void pipeline_context_receiver_note_bytes_released(PipelineContextReceiver* context, + size_t released_bytes); int receive_thread(void* pipeline_context); int write_thread(void* pipeline_context); #endif diff --git a/src/shared/protocol.c b/src/shared/protocol.c index 9c30709..5780b79 100644 --- a/src/shared/protocol.c +++ b/src/shared/protocol.c @@ -14,7 +14,6 @@ #define RECEIVE_TIMEOUT_SEC 60 /* 60 second per-message timeout */ #define SEND_TIMEOUT_SEC 60 -#define MAX_CONNECTION_MEMORY (256ULL * 1024 * 1024) /* bounded cumulative receive budget */ static __thread int io_read_fd = -1; static __thread int io_write_fd = -1; diff --git a/src/shared/protocol.h b/src/shared/protocol.h index ff79e25..2b4f66f 100644 --- a/src/shared/protocol.h +++ b/src/shared/protocol.h @@ -28,6 +28,10 @@ #define DEFAULT_MAX_ALLOC (1ULL * 1024 * 1024 * 1024) /* Server policy ceiling for a client-provided allocation limit. */ #define MAX_SERVER_ALLOC (256ULL * 1024 * 1024) +/* Bounded cumulative per-connection receive budget. In-flight wire buffers, + decompression buffers and queued (not yet written) file payloads for a + connection must stay within this ceiling. */ +#define MAX_CONNECTION_MEMORY (256ULL * 1024 * 1024) typedef struct ssl_st SSL; diff --git a/tests/test_file.c b/tests/test_file.c index 0d416be..5308e76 100644 --- a/tests/test_file.c +++ b/tests/test_file.c @@ -5,6 +5,7 @@ #include "utils.h" #include "protocol.h" #include "test_utils.h" +#include #include #include #include @@ -730,6 +731,161 @@ static void test_file_send_single_calls_metadata_and_path() { } } +static void test_inplace_overwrite_clears_special_mode_bits() { + const char* root = "test_inplace_tmp"; + const char* path = "test_inplace_tmp/priv.txt"; + const char* content = "olddata"; + unlink(path); + rmdir(root); + EXPECT_EQ_INT(mkdir(root, 0700), 0); + + /* Create a destination carrying setuid + sticky bits. */ + int fd = open(path, O_WRONLY | O_CREAT | O_TRUNC | O_CLOEXEC, 0644); + EXPECT_TRUE(fd >= 0); + // cppcheck-suppress knownConditionTrueFalse + if (fd < 0) { + rmdir(root); + return; + } + EXPECT_EQ_INT((int)write(fd, content, strlen(content)), (int)strlen(content)); + EXPECT_EQ_INT(fchmod(fd, S_ISUID | S_ISVTX | 0755), 0); + EXPECT_EQ_INT(close(fd), 0); + + /* Overwrite in place without metadata: the mode must be normalized to a + safe default (0644) and the setuid/sticky bits must be gone. */ + File* f = file_create("priv.txt"); + EXPECT_NOT_NULL(f); + const char* new_content = "newdata"; + f->data->data = malloc(strlen(new_content)); + EXPECT_NOT_NULL(f->data->data); + memcpy(f->data->data, new_content, strlen(new_content)); + f->data->size = strlen(new_content); + + Config* cfg = config_create(); + EXPECT_NOT_NULL(cfg); + cfg->inplace = true; + EXPECT_TRUE(file_save_to_disk(root, f, cfg)); + file_destroy(f); + config_delete(cfg); + + struct stat st; + EXPECT_EQ_INT(stat(path, &st), 0); + EXPECT_EQ_INT((int)(st.st_mode & (S_ISUID | S_ISGID | S_ISVTX)), 0); + EXPECT_EQ_INT((int)(st.st_mode & 0777), 0644); + FILE* stream = fopen(path, "rb"); + char buf[16] = {0}; + EXPECT_NOT_NULL(stream); + // cppcheck-suppress knownConditionTrueFalse + if (stream) { + size_t nread = fread(buf, 1, sizeof(buf) - 1, stream); + fclose(stream); + EXPECT_EQ_INT((int)nread, (int)strlen(new_content)); + } + EXPECT_EQ_STR(buf, new_content); + + unlink(path); + rmdir(root); +} + +static void test_inplace_overwrite_metadata_strips_special_bits() { + const char* root = "test_inplace_meta_tmp"; + const char* path = "test_inplace_meta_tmp/meta.txt"; + const char* source = "test_inplace_meta_source.txt"; + unlink(path); + unlink(source); + rmdir(root); + EXPECT_EQ_INT(mkdir(root, 0700), 0); + + /* Existing destination with setuid+sticky set. */ + int fd = open(path, O_WRONLY | O_CREAT | O_TRUNC | O_CLOEXEC, 0644); + EXPECT_TRUE(fd >= 0); + // cppcheck-suppress knownConditionTrueFalse + if (fd < 0) { + rmdir(root); + return; + } + EXPECT_EQ_INT((int)write(fd, "olddata", 7), 7); + EXPECT_EQ_INT(fchmod(fd, S_ISUID | S_ISVTX | 0755), 0); + EXPECT_EQ_INT(close(fd), 0); + + /* Build source metadata carrying a plain executable mode (no specials). */ + EXPECT_TRUE(file_write_to_disk(source, "source", 6, false, false)); + EXPECT_EQ_INT(chmod(source, 0755), 0); + struct stat source_st; + EXPECT_EQ_INT(stat(source, &source_st), 0); + + File* f = file_create("meta.txt"); + EXPECT_NOT_NULL(f); + const char* new_content = "meta"; + f->data->data = malloc(strlen(new_content)); + EXPECT_NOT_NULL(f->data->data); + memcpy(f->data->data, new_content, strlen(new_content)); + f->data->size = strlen(new_content); + f->metadata = file_metadata_create(&source_st); + EXPECT_NOT_NULL(f->metadata); + + Config* cfg = config_create(); + EXPECT_NOT_NULL(cfg); + cfg->inplace = true; + EXPECT_TRUE(file_save_to_disk(root, f, cfg)); + file_destroy(f); + config_delete(cfg); + unlink(source); + + struct stat st; + EXPECT_EQ_INT(stat(path, &st), 0); + /* Metadata-derived mode is applied and never includes setuid/setgid/sticky. */ + EXPECT_EQ_INT((int)(st.st_mode & (S_ISUID | S_ISGID | S_ISVTX)), 0); + EXPECT_EQ_INT((int)(st.st_mode & 0777), 0755); + + unlink(path); + rmdir(root); +} + +static void test_inplace_overwrite_truncates_shorter_payload() { + const char* root = "test_inplace_trunc_tmp"; + const char* path = "test_inplace_trunc_tmp/big.txt"; + unlink(path); + rmdir(root); + EXPECT_EQ_INT(mkdir(root, 0700), 0); + + const char* old_content = "0123456789abcdef"; /* 16 bytes */ + EXPECT_TRUE(file_write_to_disk(path, old_content, strlen(old_content), false, false)); + + File* f = file_create("big.txt"); + EXPECT_NOT_NULL(f); + const char* new_content = "hi"; + f->data->data = malloc(strlen(new_content)); + EXPECT_NOT_NULL(f->data->data); + memcpy(f->data->data, new_content, strlen(new_content)); + f->data->size = strlen(new_content); + + Config* cfg = config_create(); + EXPECT_NOT_NULL(cfg); + cfg->inplace = true; + EXPECT_TRUE(file_save_to_disk(root, f, cfg)); + file_destroy(f); + config_delete(cfg); + + /* A shorter payload must truncate the file: no stale trailing bytes. */ + struct stat st; + EXPECT_EQ_INT(stat(path, &st), 0); + EXPECT_EQ_INT((int)st.st_size, (int)strlen(new_content)); + FILE* stream = fopen(path, "rb"); + char buf[32] = {0}; + EXPECT_NOT_NULL(stream); + // cppcheck-suppress knownConditionTrueFalse + if (stream) { + size_t nread = fread(buf, 1, sizeof(buf) - 1, stream); + fclose(stream); + EXPECT_EQ_INT((int)nread, (int)strlen(new_content)); + } + EXPECT_EQ_STR(buf, new_content); + + unlink(path); + rmdir(root); +} + void test_file() { test_file_create(); test_file_destroy_null(); @@ -762,4 +918,7 @@ void test_file() { test_file_send_single_calls_metadata_and_path(); } test_file_metadata_create(); + test_inplace_overwrite_clears_special_mode_bits(); + test_inplace_overwrite_metadata_strips_special_bits(); + test_inplace_overwrite_truncates_shorter_payload(); } diff --git a/tests/test_multiprocessing.c b/tests/test_multiprocessing.c index fdb38ba..a12b2ee 100644 --- a/tests/test_multiprocessing.c +++ b/tests/test_multiprocessing.c @@ -250,6 +250,87 @@ static void test_write_thread_done() { config_delete(cfg); } +typedef struct { + PipelineContextReceiver* context; + File* file; + atomic_bool* done; + atomic_bool* result; +} ByteBudgetEnqueueArg; + +static int byte_budget_enqueue_worker(void* arg) { + ByteBudgetEnqueueArg* worker = arg; + bool ok = pipeline_context_receiver_enqueue_file(worker->context, worker->file); + atomic_store(worker->result, ok); + atomic_store(worker->done, true); + return thrd_success; +} + +/* A receiver must not buffer more decompressed/copied payload bytes ahead of + the (slow) disk writer than the configured byte budget: an enqueue that + would exceed the budget blocks until the writer releases bytes. */ +static void test_receiver_enqueue_byte_budget() { + Config* cfg = config_create(); + EXPECT_NOT_NULL(cfg); + free(cfg->version); + cfg->version = str_dup(PROTOCOL_VERSION); + cfg->send_directory = str_dup("/src"); + cfg->receive_root_directory = str_dup("/dst"); + cfg->save_to_disk = false; + + Queue* q = queue_create(16, file_destroy); + EXPECT_NOT_NULL(q); + PipelineContextReceiver* ctx = pipeline_context_receiver_create(cfg, q, -1, NULL); + EXPECT_NOT_NULL(ctx); + pipeline_context_receiver_set_queue_byte_limit(ctx, 3000); + ctx->receiver_done = false; + + File* first = file_create("budget_file_1"); + EXPECT_NOT_NULL(first); + first->data->size = 2000; + EXPECT_TRUE(pipeline_context_receiver_enqueue_file(ctx, first)); + EXPECT_EQ_INT((int)ctx->queued_bytes, 2000); + + /* Second 2000-byte payload would push the pipeline to 4000 > 3000 budget, + so the enqueue must block until the first payload is released. */ + File* second = file_create("budget_file_2"); + EXPECT_NOT_NULL(second); + second->data->size = 2000; + atomic_bool done; + atomic_bool result; + atomic_init(&done, false); + atomic_init(&result, false); + ByteBudgetEnqueueArg arg = {ctx, second, &done, &result}; + thrd_t enqueuer; + EXPECT_EQ_INT(thrd_create(&enqueuer, byte_budget_enqueue_worker, &arg), thrd_success); + + /* Give a broken (unbounded) implementation every chance to enqueue. */ + struct timespec wait = {0, 200 * 1000000L}; + thrd_sleep(&wait, NULL); + EXPECT_FALSE(atomic_load(&done)); + EXPECT_EQ_INT((int)ctx->queued_bytes, 2000); /* budget still honored */ + + /* Simulate the disk writer: dequeue + destroy + release the first file. */ + File* drained = queue_dequeue_multithreaded(q, &ctx->mutex, &ctx->condition_not_empty, + &ctx->condition_not_full, &ctx->receiver_done); + EXPECT_NOT_NULL(drained); + file_destroy(drained); + pipeline_context_receiver_note_bytes_released(ctx, 2000); + EXPECT_EQ_INT((int)ctx->queued_bytes, 0); + + EXPECT_EQ_INT(thrd_join(enqueuer, NULL), thrd_success); + EXPECT_TRUE(atomic_load(&done)); + EXPECT_TRUE(atomic_load(&result)); + EXPECT_EQ_INT((int)ctx->queued_bytes, 2000); /* second payload now in flight */ + + /* Tear down: the second file is still queued and is freed by queue_destroy. */ + mtx_destroy(&ctx->mutex); + cnd_destroy(&ctx->condition_not_full); + cnd_destroy(&ctx->condition_not_empty); + free(ctx); + queue_destroy(q); + config_delete(cfg); +} + void test_multiprocessing() { test_sender_create_destroy(); test_receiver_create_destroy(); @@ -261,4 +342,5 @@ void test_multiprocessing() { test_receive_thread_failure_wakes_writer(); } test_write_thread_done(); + test_receiver_enqueue_byte_budget(); }