From 69fe7f3c9fde6f4d0118e860e5b61d88513c1068 Mon Sep 17 00:00:00 2001 From: TapTap Date: Sun, 13 Sep 2026 04:06:30 +0200 Subject: [PATCH] perf(send,scanner): byte-bound sender queues; drop redundant stat --- src/client/client_send.c | 30 +++++--- src/client/scanner.c | 10 ++- src/shared/multiprocessing.c | 81 +++++++++++++++++++++ src/shared/multiprocessing.h | 26 +++++++ tests/test_multiprocessing.c | 119 +++++++++++++++++++++++++++++++ tests/test_scanner.c | 134 +++++++++++++++++++++++++++++++++++ 6 files changed, 389 insertions(+), 11 deletions(-) diff --git a/src/client/client_send.c b/src/client/client_send.c index a0bb9dc..bfec60a 100644 --- a/src/client/client_send.c +++ b/src/client/client_send.c @@ -37,6 +37,13 @@ #define STREAM_THRESHOLD (64ULL * 1024 * 1024) +/* Aggregate loaded payload bytes the sender may buffer across the loader queue + and the chunk in flight. Sending one chunk adds up to ~2 * MAX_CHUNK_SIZE of + transient serialize/compress buffers on top of the queued payloads, so this + ceiling keeps total pipeline memory within MAX_CONNECTION_MEMORY (mirrors the + receiver's RECEIVER_QUEUE_MAX_BYTES). */ +#define SENDER_QUEUE_MAX_BYTES (MAX_CONNECTION_MEMORY - 2 * MAX_CHUNK_SIZE) + /* Forward declaration for progress-reporting thread used in multithreaded send. */ static int progress_thread_fn(void* arg); @@ -1493,6 +1500,10 @@ static int send_chunks_multithreaded(void* pipeline_context) { protocol_session_unbind(); return thrd_error; } + /* Payload bytes this chunk was charged for on the loader's byte budget. + Computed before destruction and released after the memory is actually + freed, so a loader blocked on the budget wakes only once room exists. */ + size_t queued_payload = pipeline_context_sender_chunk_bytes(current_chunk); unsigned long long chunk_bytes = 0; int chunk_files = 0; for (int i = 0; i < current_chunk->element_count; i++) { @@ -1507,6 +1518,7 @@ static int send_chunks_multithreaded(void* pipeline_context) { context->progress_bytes = context->total_bytes; mtx_unlock(&context->mutex_progress); chunk_destroy(current_chunk); + pipeline_context_sender_note_bytes_released(context, queued_payload); } /* Completion tail: reached on natural exhaustion or an early stop deadline. @@ -1728,11 +1740,7 @@ static int load_files_multithreaded(void* pipeline_context) { } } } - if (!queue_enqueue_multithreaded_cancel(context->queue_loader, chunk, &context->mutex_loader, - &context->condition_not_empty_loader, - &context->condition_not_full_loader, - &context->cancelled)) { - chunk_destroy(chunk); + if (!pipeline_context_sender_enqueue_chunk(context, chunk)) { pipeline_cancel(context); protocol_session_unbind(); return thrd_error; @@ -2197,8 +2205,13 @@ int send_files_multithreaded(Config** config_ptr) { unsigned long long available_memory = pages > 0 && page_size > 0 ? (unsigned long long)pages * (unsigned long long)page_size : 512ULL * 1024 * 1024; - unsigned long long avg_file_size = 1024 * 1024; - int qsize = (int)(available_memory / avg_file_size); + /* Size the chunk queues from the actual chunk size rather than a fixed 1 MiB + average: a chunk holds roughly `chunk_size` bytes of file data, so counting + chunks at 1 MiB over-estimated the queue capacity by up to 10x. The byte + budget below is the authoritative bound; this count keeps the unloaded + chunks waiting in the scanner queue bounded too. */ + unsigned long long chunk_size = config->chunk_size > 0 ? config->chunk_size : DEFAULT_CHUNK_SIZE; + int qsize = (int)(available_memory / chunk_size); if (qsize < 10) qsize = 10; if (qsize > 1000) @@ -2223,7 +2236,8 @@ int send_files_multithreaded(Config** config_ptr) { } context->missing_args = missing_args; missing_args = NULL; /* owned by the context from here on */ - *config_ptr = NULL; /* context now owns config through all remaining paths */ + pipeline_context_sender_set_queue_byte_limit(context, SENDER_QUEUE_MAX_BYTES); + *config_ptr = NULL; /* context now owns config through all remaining paths */ struct timespec now_mono; if (clock_gettime(CLOCK_MONOTONIC, &now_mono) != 0) { now_mono.tv_sec = 0; diff --git a/src/client/scanner.c b/src/client/scanner.c index 6fd3407..7db3f58 100644 --- a/src/client/scanner.c +++ b/src/client/scanner.c @@ -388,9 +388,13 @@ static int scanner_inspect_entry(const ScannerOptions* options, const char* sour goto apply_filters; regular: - if (stat(entry->path, &entry->stats) != 0) - goto skip; - entry->is_directory = S_ISDIR(entry->stats.st_mode); + /* Not a symlink: the lstat() above already described this entry, and lstat + and stat are identical for every non-symlink, so reuse that result instead + of issuing a redundant stat() on the scanner hot path. stat() is still + used on the dereference paths above/below for actual symlinks (copy-links, + safe/copy-unsafe links, and -k symlinks-to-directories). */ + entry->stats = link_stats; + entry->is_directory = S_ISDIR(link_stats.st_mode); if (entry->is_directory) return 1; diff --git a/src/shared/multiprocessing.c b/src/shared/multiprocessing.c index b93c9f0..697e461 100644 --- a/src/shared/multiprocessing.c +++ b/src/shared/multiprocessing.c @@ -11,6 +11,7 @@ #include "protocol.h" #include "queue.h" #include "utils.h" +#include #include #include #include @@ -26,6 +27,8 @@ PipelineContextSender* pipeline_context_sender_create(Config* config, Queue* que context->queue_loader = queue_loader; context->scanner_done = false; context->loader_done = false; + context->queued_bytes = 0; + context->max_queue_bytes = 0; context->manifest = NULL; context->excluded_paths = NULL; context->missing_args = NULL; @@ -97,6 +100,84 @@ fail: return NULL; } +void pipeline_context_sender_set_queue_byte_limit(PipelineContextSender* context, + size_t max_bytes) { + if (context == NULL) + return; + mtx_lock(&context->mutex_loader); + context->max_queue_bytes = max_bytes; + context->queued_bytes = 0; + cnd_broadcast(&context->condition_not_full_loader); + mtx_unlock(&context->mutex_loader); +} + +size_t pipeline_context_sender_chunk_bytes(const Chunk* chunk) { + if (chunk == NULL || chunk->items == NULL) + return 0; + size_t total = 0; + for (int i = 0; i < chunk->element_count; i++) { + const File* file = chunk->items[i]; + if (file == NULL || file->data == NULL || file->data->data == NULL) + continue; + if (file->data->size > SIZE_MAX - total) + return SIZE_MAX; + total += file->data->size; + } + return total; +} + +void pipeline_context_sender_note_bytes_released(PipelineContextSender* context, + size_t released_bytes) { + if (context == NULL || context->max_queue_bytes == 0 || released_bytes == 0) + return; + mtx_lock(&context->mutex_loader); + if (released_bytes >= context->queued_bytes) + context->queued_bytes = 0; + else + context->queued_bytes -= released_bytes; + cnd_signal(&context->condition_not_full_loader); + mtx_unlock(&context->mutex_loader); +} + +bool pipeline_context_sender_enqueue_chunk(PipelineContextSender* context, Chunk* chunk) { + if (context == NULL || chunk == NULL) + return false; + size_t chunk_bytes = pipeline_context_sender_chunk_bytes(chunk); + mtx_lock(&context->mutex_loader); + while (!atomic_load(&context->cancelled)) { + bool blocked_by_count = queue_is_full(context->queue_loader); + 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 (chunk_bytes > budget - used) { + /* A single payload larger than the whole budget 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_loader, &context->mutex_loader); + } + if (atomic_load(&context->cancelled)) { + mtx_unlock(&context->mutex_loader); + chunk_destroy(chunk); + return false; + } + if (!queue_enqueue(context->queue_loader, chunk)) { + mtx_unlock(&context->mutex_loader); + chunk_destroy(chunk); + return false; + } + context->queued_bytes += chunk_bytes; + cnd_signal(&context->condition_not_empty_loader); + mtx_unlock(&context->mutex_loader); + return true; +} + void pipeline_context_sender_destroy(PipelineContextSender* context) { if (context->manifest) { array_list_delete(context->manifest); diff --git a/src/shared/multiprocessing.h b/src/shared/multiprocessing.h index 5cbf236..73162a4 100644 --- a/src/shared/multiprocessing.h +++ b/src/shared/multiprocessing.h @@ -5,6 +5,7 @@ #include #include "array_list.h" +#include "chunk.h" #include "config.h" #include "file.h" #include "protocol.h" @@ -25,6 +26,15 @@ typedef struct { cnd_t condition_not_full_loader; cnd_t condition_not_empty_loader; bool loader_done; + /* Aggregate loaded payload bytes queued on queue_loader but not yet released + by the sender. Guarded by `mutex_loader`. When `max_queue_bytes` is + non-zero the loader blocks before enqueueing a chunk that would push this + total over it, so the sender buffers a bounded number of bytes rather than + an unbounded count of chunks that may each be up to chunk_size (or a single + file) in size. Files streamed straight from disk by sendfile hold no + payload, so only in-memory (`data->data`) payloads are counted. */ + size_t queued_bytes; + size_t max_queue_bytes; ArrayList* manifest; /* Protected prefixes (paths the source scan excluded by user rules) sent with the keep-set manifest so --delete leaves them alone unless @@ -113,6 +123,22 @@ typedef struct PipelineContextReceiver { PipelineContextSender* pipeline_context_sender_create(Config* config, Queue* queue_scanner, Queue* queue_loader); void pipeline_context_sender_destroy(PipelineContextSender* context); +/* Bound the loaded payload bytes the sender may buffer ahead of the network + writer (see max_queue_bytes). */ +void pipeline_context_sender_set_queue_byte_limit(PipelineContextSender* context, size_t max_bytes); +/* Total payload bytes a chunk currently holds in memory (loaded file data + only; zero for entries with no payload or data streamed from disk). */ +size_t pipeline_context_sender_chunk_bytes(const Chunk* chunk); +/* Blocking enqueue used by the sender's loader stage. Blocks while + queue_loader is full by element count or when adding `chunk` would push the + queued payload bytes over the configured byte limit; waits until the sender + releases bytes. Takes ownership of `chunk` on success and destroys it on + failure/cancel. */ +bool pipeline_context_sender_enqueue_chunk(PipelineContextSender* context, Chunk* chunk); +/* Account for `released_bytes` of payload memory that the sender freed after + destroying a chunk, unblocking a loader waiting on the byte limit. */ +void pipeline_context_sender_note_bytes_released(PipelineContextSender* context, + size_t released_bytes); PipelineContextReceiver* pipeline_context_receiver_create(Config* config, Queue* queue_receiver, int file_descriptor, SSL* ssl); void pipeline_context_receiver_destroy(PipelineContextReceiver* context); diff --git a/tests/test_multiprocessing.c b/tests/test_multiprocessing.c index fb0241d..427b5c6 100644 --- a/tests/test_multiprocessing.c +++ b/tests/test_multiprocessing.c @@ -332,6 +332,123 @@ static void test_receiver_enqueue_byte_budget() { config_delete(cfg); } +/* Build a one-file chunk that appears to hold `bytes` of loaded payload by + handing it a real buffer of that size (the byte accounting counts only + in-memory `data->data`, mirroring sendfile's streamed chunks). */ +static Chunk* make_loaded_chunk(const char* name, size_t bytes) { + File* file = file_create(name); + if (!file) + return NULL; + void* buffer = malloc(bytes > 0 ? bytes : 1); + if (!buffer) { + file_destroy(file); + return NULL; + } + file->data->data = buffer; + file->data->size = bytes; + File* items[1] = {file}; + Chunk* chunk = chunk_create(items, 1); + if (!chunk) + file_destroy(file); + return chunk; +} + +/* Only loaded (in-memory) payload is charged: a chunk whose files have no + buffer (e.g. sendfile streams the bytes from disk) accounts for zero. */ +static void test_sender_chunk_bytes_accounting() { + EXPECT_EQ_INT((int)pipeline_context_sender_chunk_bytes(NULL), 0); + + File* streamed = file_create("sender_account_streamed"); + EXPECT_NOT_NULL(streamed); + streamed->data->size = 4096; /* declared size, but no in-memory buffer */ + File* streamed_items[1] = {streamed}; + Chunk* streamed_chunk = chunk_create(streamed_items, 1); + EXPECT_NOT_NULL(streamed_chunk); + EXPECT_EQ_INT((int)pipeline_context_sender_chunk_bytes(streamed_chunk), 0); + chunk_destroy(streamed_chunk); + + Chunk* loaded = make_loaded_chunk("sender_account_loaded", 2000); + EXPECT_NOT_NULL(loaded); + EXPECT_EQ_INT((int)pipeline_context_sender_chunk_bytes(loaded), 2000); + chunk_destroy(loaded); +} + +typedef struct { + PipelineContextSender* context; + Chunk* chunk; + atomic_bool* done; + atomic_bool* result; +} SenderByteBudgetArg; + +static int sender_byte_budget_worker(void* arg) { + SenderByteBudgetArg* worker = arg; + bool ok = pipeline_context_sender_enqueue_chunk(worker->context, worker->chunk); + atomic_store(worker->result, ok); + atomic_store(worker->done, true); + return thrd_success; +} + +/* The sender's loader stage must not buffer more loaded payload bytes ahead of + the network writer than the configured byte budget: an enqueue that would + exceed the budget blocks until the sender releases bytes. */ +static void test_sender_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"); + + Queue* q_scanner = queue_create(16, chunk_destroy); + Queue* q_loader = queue_create(16, chunk_destroy); + EXPECT_NOT_NULL(q_scanner); + EXPECT_NOT_NULL(q_loader); + PipelineContextSender* ctx = pipeline_context_sender_create(cfg, q_scanner, q_loader); + EXPECT_NOT_NULL(ctx); + pipeline_context_sender_set_queue_byte_limit(ctx, 3000); + + Chunk* first = make_loaded_chunk("sender_budget_1", 2000); + EXPECT_NOT_NULL(first); + EXPECT_TRUE(pipeline_context_sender_enqueue_chunk(ctx, first)); + EXPECT_EQ_INT((int)ctx->queued_bytes, 2000); + + /* A second 2000-byte chunk would push the pipeline to 4000 > 3000 budget, so + its enqueue must block until the first chunk's bytes are released. */ + Chunk* second = make_loaded_chunk("sender_budget_2", 2000); + EXPECT_NOT_NULL(second); + atomic_bool done; + atomic_bool result; + atomic_init(&done, false); + atomic_init(&result, false); + SenderByteBudgetArg arg = {ctx, second, &done, &result}; + thrd_t enqueuer; + EXPECT_EQ_INT(thrd_create(&enqueuer, sender_byte_budget_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 sender: dequeue + destroy the first chunk, then release its + bytes. Only the post-join state (below) is deterministic. */ + Chunk* drained = + queue_dequeue_multithreaded(q_loader, &ctx->mutex_loader, &ctx->condition_not_empty_loader, + &ctx->condition_not_full_loader, &ctx->loader_done); + EXPECT_NOT_NULL(drained); + chunk_destroy(drained); + pipeline_context_sender_note_bytes_released(ctx, 2000); + + 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 */ + + /* pipeline_context_sender_destroy frees the still-queued second chunk and + owns cfg/q_scanner/q_loader from here on. */ + pipeline_context_sender_destroy(ctx); +} + void test_multiprocessing() { test_sender_create_destroy(); test_receiver_create_destroy(); @@ -344,4 +461,6 @@ void test_multiprocessing() { } test_write_thread_done(); test_receiver_enqueue_byte_budget(); + test_sender_chunk_bytes_accounting(); + test_sender_enqueue_byte_budget(); } diff --git a/tests/test_scanner.c b/tests/test_scanner.c index 3ba35ed..9c3bcb6 100644 --- a/tests/test_scanner.c +++ b/tests/test_scanner.c @@ -1369,6 +1369,139 @@ static void test_scanner_chunk_ownership() { rmdir(dir); } +/* The scanner derives entry type from a single lstat() for non-symlinks + * (regular files and directories) and only calls stat() to dereference real + * symlinks. Guard the regular-file/directory/symlink distinction across the + * default (symlinks skipped), --copy-links (dereferenced) and -l (carried) + * modes so the lstat/stat reuse cannot misclassify entries. */ +static void test_scanner_entry_classification() { + const char* root = "test_scan_classify"; + const char* sub = "test_scan_classify/sub"; + const char* file = "test_scan_classify/file.txt"; + const char* nested = "test_scan_classify/sub/nested.txt"; + const char* link_file = "test_scan_classify/link_file"; + const char* link_dir = "test_scan_classify/link_dir"; + + EXPECT_EQ_INT(mkdir(root, 0755), 0); + EXPECT_EQ_INT(mkdir(sub, 0755), 0); + create_test_file(file, "hello"); /* 5 bytes */ + create_test_file(nested, "nested"); /* 6 bytes */ + EXPECT_EQ_INT(symlink("file.txt", link_file), 0); + EXPECT_EQ_INT(symlink("sub", link_dir), 0); + + /* Default: no link option -> symlinks are skipped entirely; regular files and + directories (descended, not emitted) are classified as before. */ + { + ScannerOptions options = {0}; + DirectoryScanner* scanner = directory_scanner_create_with_options(root, &options); + EXPECT_NOT_NULL(scanner); + size_t root_len = strlen(root); + bool file_ok = false, nested_ok = false, link_seen = false; + Chunk* chunk; + while ((chunk = directory_scanner_next(scanner)) != NULL) { + for (int i = 0; i < chunk->element_count; i++) { + const File* f = chunk->items[i]; + const char* rel = f->path + root_len; + if (*rel == '/') + rel++; + if (strcmp(rel, "file.txt") == 0) { + file_ok = !f->is_dir && !f->is_symlink && f->data->size == 5; + } else if (strcmp(rel, "sub/nested.txt") == 0) { + nested_ok = !f->is_dir && !f->is_symlink && f->data->size == 6; + } else { + link_seen = true; + } + } + chunk_destroy(chunk); + } + EXPECT_FALSE(directory_scanner_failed(scanner)); + EXPECT_TRUE(file_ok); + EXPECT_TRUE(nested_ok); + EXPECT_FALSE(link_seen); + directory_scanner_destroy(scanner); + } + + /* --copy-links: symlinks are dereferenced. A link to a file becomes a + regular file with the referent's size; a link to a directory is traversed. */ + { + ScannerOptions options = {0}; + options.copy_links = true; + DirectoryScanner* scanner = directory_scanner_create_with_options(root, &options); + EXPECT_NOT_NULL(scanner); + size_t root_len = strlen(root); + bool file_ok = false, nested_ok = false, link_file_ok = false; + bool link_dir_nested_ok = false, symlink_leaked = false; + Chunk* chunk; + while ((chunk = directory_scanner_next(scanner)) != NULL) { + for (int i = 0; i < chunk->element_count; i++) { + const File* f = chunk->items[i]; + const char* rel = f->path + root_len; + if (*rel == '/') + rel++; + if (f->is_symlink) + symlink_leaked = true; + if (strcmp(rel, "file.txt") == 0) + file_ok = !f->is_dir && f->data->size == 5; + else if (strcmp(rel, "sub/nested.txt") == 0) + nested_ok = !f->is_dir && f->data->size == 6; + else if (strcmp(rel, "link_file") == 0) + link_file_ok = !f->is_dir && f->data->size == 5; + else if (strcmp(rel, "link_dir/nested.txt") == 0) + link_dir_nested_ok = !f->is_dir && f->data->size == 6; + } + chunk_destroy(chunk); + } + EXPECT_FALSE(directory_scanner_failed(scanner)); + EXPECT_TRUE(file_ok); + EXPECT_TRUE(nested_ok); + EXPECT_TRUE(link_file_ok); + EXPECT_TRUE(link_dir_nested_ok); + EXPECT_FALSE(symlink_leaked); + directory_scanner_destroy(scanner); + } + + /* -l (--links): symlinks are carried through as symlinks, not dereferenced. */ + { + ScannerOptions options = {0}; + options.follow_symlinks = true; + DirectoryScanner* scanner = directory_scanner_create_with_options(root, &options); + EXPECT_NOT_NULL(scanner); + size_t root_len = strlen(root); + bool file_ok = false, link_file_ok = false, link_dir_ok = false, leaked_dir = false; + Chunk* chunk; + while ((chunk = directory_scanner_next(scanner)) != NULL) { + for (int i = 0; i < chunk->element_count; i++) { + const File* f = chunk->items[i]; + const char* rel = f->path + root_len; + if (*rel == '/') + rel++; + if (strcmp(rel, "file.txt") == 0) + file_ok = !f->is_dir && !f->is_symlink && f->data->size == 5; + else if (strcmp(rel, "link_file") == 0) + link_file_ok = f->is_symlink && !f->is_dir; + else if (strcmp(rel, "link_dir") == 0) + link_dir_ok = f->is_symlink && !f->is_dir; + else if (strcmp(rel, "link_dir/nested.txt") == 0) + leaked_dir = true; + } + chunk_destroy(chunk); + } + EXPECT_FALSE(directory_scanner_failed(scanner)); + EXPECT_TRUE(file_ok); + EXPECT_TRUE(link_file_ok); + EXPECT_TRUE(link_dir_ok); + EXPECT_FALSE(leaked_dir); + directory_scanner_destroy(scanner); + } + + unlink(link_file); + unlink(link_dir); + unlink(nested); + unlink(file); + rmdir(sub); + rmdir(root); +} + void test_scanner() { test_scanner_single_file(); test_scanner_multiple_files(); @@ -1406,4 +1539,5 @@ void test_scanner() { test_files_from_relative_send_path(); test_scanner_captures_directory_times(); test_scanner_chunk_ownership(); + test_scanner_entry_classification(); }