From 3499baf80b475fa9a1ad23b099af44867dcd5257 Mon Sep 17 00:00:00 2001 From: TapTap Date: Sun, 13 Sep 2026 07:20:28 +0200 Subject: [PATCH 1/2] build: explicit CMake targets; move receiver pipeline out of shared --- CMakeLists.txt | 184 ++++++++++++++++++++---- src/server/receiver_pipeline.c | 254 +++++++++++++++++++++++++++++++++ src/server/receiver_pipeline.h | 68 +++++++++ src/server/server.c | 2 +- src/shared/multiprocessing.c | 246 ------------------------------- src/shared/multiprocessing.h | 52 ------- tests/test_config.c | 1 + tests/test_multiprocessing.c | 1 + 8 files changed, 482 insertions(+), 326 deletions(-) create mode 100644 src/server/receiver_pipeline.c create mode 100644 src/server/receiver_pipeline.h diff --git a/CMakeLists.txt b/CMakeLists.txt index f384f62..e61b0fa 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -70,36 +70,105 @@ endif() find_package(OpenSSL REQUIRED) -file(GLOB SHARED_SRCS "src/shared/*.c") -set(FILE_STORE_SRCS "${CMAKE_CURRENT_SOURCE_DIR}/src/shared/file_store.c") -list(REMOVE_ITEM SHARED_SRCS ${FILE_STORE_SRCS}) -file(GLOB SERVER_SRCS "src/server/*.c") -set(SERVER_RECEIVER_SRCS src/server/receiver.c) -file(GLOB CLIENT_SRCS "src/client/*.c") +# --- Explicit source lists --- +# The shared library is self-contained: it must never depend on the client or +# server modules. In particular, the receiver pipeline (receive_thread / +# write_thread) lives under src/server, not here, so the client executable can +# link the shared library without pulling in any server code. +set(SHARED_SRCS + src/shared/array_list.c + src/shared/batch.c + src/shared/charset.c + src/shared/checksum.c + src/shared/chmod.c + src/shared/chunk.c + src/shared/compression.c + src/shared/config.c + src/shared/credentials.c + src/shared/daemon_conf.c + src/shared/data.c + src/shared/delay_updates.c + src/shared/delta.c + src/shared/file.c + src/shared/file_list.c + src/shared/file_receive.c + src/shared/file_send.c + src/shared/file_store.c + src/shared/filter.c + src/shared/hardlink.c + src/shared/identity.c + src/shared/log.c + src/shared/metadata.c + src/shared/motd.c + src/shared/multiprocessing.c + src/shared/protocol.c + src/shared/queue.c + src/shared/stop_condition.c + src/shared/transport_ssh.c + src/shared/transport_tcp.c + src/shared/transport_tls.c + src/shared/utils.c + src/shared/xattr.c +) + +# Server implementation (no main): the receiver read/write pipeline plus the +# CLI parser. The server executable adds its own main (server.c). +set(SERVER_CORE_SRCS + src/server/receiver.c + src/server/receiver_pipeline.c + src/server/server_cli.c +) +set(SERVER_MAIN_SRCS src/server/server.c) + +# Client implementation (no main): everything except the CLI entry point. +set(CLIENT_CORE_SRCS + src/client/change_list.c + src/client/client_send.c + src/client/client_validation.c + src/client/scanner.c + src/client/usage.c +) +set(CLIENT_MAIN_SRCS src/client/client_cli.c) + +# --- Library targets --- +add_library(fastsync_shared STATIC ${SHARED_SRCS}) +target_include_directories(fastsync_shared PUBLIC src/shared) +target_link_libraries(fastsync_shared PUBLIC Threads::Threads ${ZSTD_LIBRARY} OpenSSL::SSL + OpenSSL::Crypto xxhash) + +add_library(fastsync_client_core STATIC ${CLIENT_CORE_SRCS}) +target_include_directories(fastsync_client_core PUBLIC src/client) +target_link_libraries(fastsync_client_core PUBLIC fastsync_shared) + +add_library(fastsync_server_core STATIC ${SERVER_CORE_SRCS}) +target_include_directories(fastsync_server_core PUBLIC src/server) +target_link_libraries(fastsync_server_core PUBLIC fastsync_shared) # --- Main executables --- -add_executable(server ${SERVER_SRCS} ${SHARED_SRCS} ${FILE_STORE_SRCS}) -target_include_directories(server PRIVATE src/shared src/server src/client) -target_link_libraries(server PRIVATE Threads::Threads ${ZSTD_LIBRARY} OpenSSL::SSL OpenSSL::Crypto xxhash) +# The client links only the shared library and its own core; it deliberately +# does NOT get src/server on its include path nor compile receiver.c. +add_executable(server ${SERVER_MAIN_SRCS}) +target_link_libraries(server PRIVATE fastsync_server_core) -add_executable(client ${CLIENT_SRCS} ${SHARED_SRCS} ${FILE_STORE_SRCS} ${SERVER_RECEIVER_SRCS}) -target_include_directories(client PRIVATE src/shared src/server src/client) -target_link_libraries(client PRIVATE Threads::Threads ${ZSTD_LIBRARY} OpenSSL::SSL OpenSSL::Crypto xxhash) +add_executable(client ${CLIENT_MAIN_SRCS}) +target_link_libraries(client PRIVATE fastsync_client_core) # --- Production hardening --- # Each compile flag is probed so a compiler/architecture that lacks it still # configures cleanly. _FORTIFY_SOURCE is guarded separately because it only # works in an optimising build. xxHash is a static archive built by -# FetchContent, so it must be position-independent for the -pie link. +# FetchContent, so it must be position-independent for the -pie link; the same +# applies to the first-party static libraries linked into the -pie binaries. if(HARDENING_ACTIVE) - set_target_properties(xxhash PROPERTIES POSITION_INDEPENDENT_CODE ON) + set_target_properties(xxhash fastsync_shared fastsync_server_core fastsync_client_core + PROPERTIES POSITION_INDEPENDENT_CODE ON) include(CheckCCompilerFlag) foreach(flag -fstack-protector-strong -fstack-clash-protection -fPIE) string(MAKE_C_IDENTIFIER "HARDEN_${flag}" _harden_var) check_c_compiler_flag("${flag}" ${_harden_var}) endforeach() check_c_compiler_flag("-D_FORTIFY_SOURCE=2" HARDEN_FORTIFY_SOURCE) - foreach(target server client) + foreach(target fastsync_shared fastsync_server_core fastsync_client_core server client) foreach(flag -fstack-protector-strong -fstack-clash-protection -fPIE) string(MAKE_C_IDENTIFIER "HARDEN_${flag}" _harden_var) if(${_harden_var}) @@ -109,6 +178,8 @@ if(HARDENING_ACTIVE) if(HARDEN_FORTIFY_SOURCE) target_compile_options(${target} PRIVATE -D_FORTIFY_SOURCE=2) endif() + endforeach() + foreach(target server client) target_link_options(${target} PRIVATE -pie -Wl,-z,relro -Wl,-z,now -Wl,-z,noexecstack) endforeach() endif() @@ -116,16 +187,59 @@ endif() # --- Testing --- enable_testing() -# Common test libraries -set(TEST_LIBS Threads::Threads ${ZSTD_LIBRARY} OpenSSL::SSL OpenSSL::Crypto xxhash) -set(TEST_INCLUDES tests src/shared src/server src/client) +# --- Unit tests --- +# The monolithic test binary exercises both client and server code, so it is +# the one place that legitimately sees both include directories and links both +# core libraries. client_cli.c is compiled here directly (with the test build +# define) rather than linked from fastsync_client_core so its test-only shims +# and the absence of main() are preserved. +set(TEST_SRCS + tests/runner.c + tests/test_array_list.c + tests/test_batch.c + tests/test_change_list.c + tests/test_checksum.c + tests/test_chunk.c + tests/test_client_cli.c + tests/test_compression.c + tests/test_config.c + tests/test_credentials.c + tests/test_daemon_conf.c + tests/test_data.c + tests/test_delay_updates.c + tests/test_delta.c + tests/test_file.c + tests/test_file_list.c + tests/test_file_sendfile.c + tests/test_fuzz_smoke.c + tests/test_glob.c + tests/test_hardlink.c + tests/test_iconv.c + tests/test_log.c + tests/test_metadata.c + tests/test_motd.c + tests/test_multiprocessing.c + tests/test_property.c + tests/test_protocol.c + tests/test_queue.c + tests/test_receiver_timeout.c + tests/test_robustness.c + tests/test_scanner.c + tests/test_server.c + tests/test_server_cli.c + tests/test_shared_utils.c + tests/test_stop.c + tests/test_stress.c + tests/test_transport_ssh.c + tests/test_transport_tcp.c + tests/test_transport_tls.c + tests/test_xattr.c +) -# Monolithic test binary (backward compatible) -file(GLOB TEST_SRCS "tests/test_*.c" "tests/runner.c") -add_executable(tests ${TEST_SRCS} ${SHARED_SRCS} ${FILE_STORE_SRCS} ${SERVER_RECEIVER_SRCS} src/client/scanner.c src/client/change_list.c src/client/client_cli.c src/client/client_validation.c src/client/usage.c src/server/server_cli.c) -target_include_directories(tests PRIVATE ${TEST_INCLUDES}) +add_executable(tests ${TEST_SRCS} src/client/client_cli.c) +target_include_directories(tests PRIVATE tests) target_compile_definitions(tests PRIVATE FASTSYNC_TEST_BUILD) -target_link_libraries(tests PRIVATE ${TEST_LIBS}) +target_link_libraries(tests PRIVATE fastsync_server_core fastsync_client_core) add_test(NAME unit_all COMMAND tests) # --- Fuzz targets (requires clang) --- @@ -134,13 +248,29 @@ if(ENABLE_FUZZ) if(NOT CMAKE_C_COMPILER_ID MATCHES "Clang") message(FATAL_ERROR "ENABLE_FUZZ requires Clang (compiler is ${CMAKE_C_COMPILER_ID})") endif() - file(GLOB FUZZ_SRCS "tests/fuzz/*.c") + set(FUZZ_SRCS + tests/fuzz/fuzz_chunk_deserialize.c + tests/fuzz/fuzz_compress_decompress.c + tests/fuzz/fuzz_config_receive.c + tests/fuzz/fuzz_delta_deserialize.c + tests/fuzz/fuzz_delta_signature_deserialize.c + tests/fuzz/fuzz_glob_match.c + tests/fuzz/fuzz_identity_parse.c + tests/fuzz/fuzz_manifest.c + tests/fuzz/fuzz_metadata_from_buf.c + tests/fuzz/fuzz_protocol_framing.c + tests/fuzz/fuzz_xattr_block.c + ) + # Compile the sources under test directly so libFuzzer's coverage + # instrumentation sees them (static libraries would be uninstrumented). + set(FUZZ_CORE_SRCS ${SHARED_SRCS} src/server/receiver.c src/server/receiver_pipeline.c) foreach(FUZZ_SRC ${FUZZ_SRCS}) get_filename_component(FUZZ_NAME ${FUZZ_SRC} NAME_WE) - add_executable(${FUZZ_NAME} ${FUZZ_SRC} ${SHARED_SRCS} ${FILE_STORE_SRCS} ${SERVER_RECEIVER_SRCS}) - target_include_directories(${FUZZ_NAME} PRIVATE ${TEST_INCLUDES}) + add_executable(${FUZZ_NAME} ${FUZZ_SRC} ${FUZZ_CORE_SRCS}) + target_include_directories(${FUZZ_NAME} PRIVATE tests src/shared src/server) target_compile_options(${FUZZ_NAME} PRIVATE -fsanitize=fuzzer,address,undefined -fno-omit-frame-pointer) target_link_options(${FUZZ_NAME} PRIVATE -fsanitize=fuzzer,address,undefined) - target_link_libraries(${FUZZ_NAME} PRIVATE ${TEST_LIBS}) + target_link_libraries(${FUZZ_NAME} PRIVATE Threads::Threads ${ZSTD_LIBRARY} OpenSSL::SSL + OpenSSL::Crypto xxhash) endforeach() endif() diff --git a/src/server/receiver_pipeline.c b/src/server/receiver_pipeline.c new file mode 100644 index 0000000..abe8719 --- /dev/null +++ b/src/server/receiver_pipeline.c @@ -0,0 +1,254 @@ +#include "receiver_pipeline.h" + +#include "log.h" +#include "protocol.h" +#include "queue.h" +#include "utils.h" +#include +#include +#include + +PipelineContextReceiver* pipeline_context_receiver_create(Config* config, Queue* queue, + int file_descriptor, SSL* ssl) { + PipelineContextReceiver* context = malloc(sizeof(PipelineContextReceiver)); + if (context == NULL) + return NULL; + context->config = config; + context->queue = queue; + context->file_descriptor = file_descriptor; + context->ssl = ssl; + context->outcomes.entries = NULL; + context->outcomes.count = 0; + context->outcomes.capacity = 0; + dir_time_list_init(&context->dir_times); + 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; + context->deferred_manifest = NULL; + atomic_init(&context->cancelled, false); + int init = 0; + if (mtx_init(&context->mutex, mtx_plain) != thrd_success) + goto fail; + init++; + if (cnd_init(&context->condition_not_full) != thrd_success) + goto fail; + init++; + if (cnd_init(&context->condition_not_empty) != thrd_success) + goto fail; + // cppcheck-suppress unreadVariable + init++; + return context; + +fail: + log_perror("Error initializing synchronization objects"); + if (init >= 3) + cnd_destroy(&context->condition_not_empty); + if (init >= 2) + cnd_destroy(&context->condition_not_full); + if (init >= 1) + mtx_destroy(&context->mutex); + free(context); + return NULL; +} + +void pipeline_context_receiver_destroy(PipelineContextReceiver* context) { + config_delete(context->config); + if (context->deferred_manifest) + delete_manifest_free(context->deferred_manifest); + queue_destroy(context->queue); + receiver_outcomes_destroy(&context->outcomes); + dir_time_list_free(&context->dir_times); + mtx_destroy(&context->mutex); + cnd_destroy(&context->condition_not_full); + cnd_destroy(&context->condition_not_empty); + 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 = (PipelineContextReceiver*)context_pointer; + return pipeline_context_receiver_enqueue_file(context, file); +} + +static void receiver_thread_fail(PipelineContextReceiver* context) { + mtx_lock(&context->mutex); + atomic_store(&context->cancelled, true); + context->receiver_done = true; + cnd_broadcast(&context->condition_not_empty); + cnd_broadcast(&context->condition_not_full); + mtx_unlock(&context->mutex); +} + +int receive_thread(void* pipeline_context) { + PipelineContextReceiver* context = (PipelineContextReceiver*)pipeline_context; + protocol_session_bind(&context->session); + mtx_lock(&context->mutex); + int file_descriptor = context->file_descriptor; + const Config* config = context->config; + mtx_unlock(&context->mutex); + + ReceiverSink sink = {receiver_enqueue_file, context, false, false, NULL}; + if (receiver_process_pending((Config*)config, file_descriptor, &sink, + &context->deferred_manifest) != 0) { + receiver_thread_fail(context); + protocol_session_unbind(); + return thrd_error; + } + mtx_lock(&context->mutex); + context->receiver_done = true; + cnd_signal(&context->condition_not_empty); + mtx_unlock(&context->mutex); + protocol_session_unbind(); + return thrd_success; +} + +int write_thread(void* pipeline_context) { + PipelineContextReceiver* context = (PipelineContextReceiver*)pipeline_context; + protocol_session_bind(&context->session); + mtx_lock(&context->mutex); + bool save_to_disk = context->config->save_to_disk; + char* root_directory = str_dup(context->config->receive_root_directory); + mtx_unlock(&context->mutex); + if (save_to_disk && !root_directory) { + mtx_lock(&context->mutex); + atomic_store(&context->cancelled, true); + context->receiver_done = true; + cnd_broadcast(&context->condition_not_full); + cnd_broadcast(&context->condition_not_empty); + mtx_unlock(&context->mutex); + protocol_session_unbind(); + return thrd_error; + } + + while (true) { + File* file = + queue_dequeue_multithreaded(context->queue, &context->mutex, &context->condition_not_empty, + &context->condition_not_full, &context->receiver_done); + if (file == NULL) { + free(root_directory); + 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; + cnd_broadcast(&context->condition_not_full); + cnd_broadcast(&context->condition_not_empty); + mtx_unlock(&context->mutex); + free(root_directory); + protocol_session_unbind(); + return thrd_error; + } + } + /* P7 Wave D: a directory's times are never applied inline (a later child + write would clobber them); accumulate the metadata here and let the + caller apply it once every writer has drained. */ + if (result != FILE_SAVE_ERROR && file->is_dir && file->metadata && + dir_times_should_capture(context->config) && + !dir_time_list_add(&context->dir_times, file->path, file->metadata)) { + 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; + cnd_broadcast(&context->condition_not_full); + cnd_broadcast(&context->condition_not_empty); + mtx_unlock(&context->mutex); + free(root_directory); + protocol_session_unbind(); + return thrd_error; + } + /* Record the per-file outcome so a --remove-source-files sender learns + which sources were actually written versus skipped on the receiver. + Explicit directory entries and recreated device/special nodes have no + source and are never acknowledged (mirrors receiver.c). */ + if (context->config->remove_source_files && !file->is_dir && !file->is_special && !file->skip && + !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; + cnd_broadcast(&context->condition_not_full); + cnd_broadcast(&context->condition_not_empty); + mtx_unlock(&context->mutex); + free(root_directory); + protocol_session_unbind(); + return thrd_error; + } + file_destroy(file); + pipeline_context_receiver_note_bytes_released(context, file_bytes); + } +} diff --git a/src/server/receiver_pipeline.h b/src/server/receiver_pipeline.h new file mode 100644 index 0000000..7aef579 --- /dev/null +++ b/src/server/receiver_pipeline.h @@ -0,0 +1,68 @@ +#ifndef RECEIVER_PIPELINE_H +#define RECEIVER_PIPELINE_H + +#include +#include +#include + +#include "config.h" +#include "file.h" +#include "file_receive.h" +#include "protocol.h" +#include "queue.h" +#include "receiver.h" +#include + +typedef struct PipelineContextReceiver { + Queue* queue; + Config* config; + int file_descriptor; + SSL* ssl; + ProtocolSession session; + ReceiverOutcomes outcomes; + mtx_t mutex; + cnd_t condition_not_full; + 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; + /* Keep-set manifest for the commit-style (late) deletion + (--delete/--delete-after/--delete-delay). receive_thread parses the whole + protocol stream but hands the manifest here instead of deleting while the + disk writer may still be draining; the caller (server.c) commits the + deletion after both threads have joined, so no extra is removed unless the + transfer truly succeeded. NULL in the early delete modes (which delete at + the manifest). */ + DeleteManifest* deferred_manifest; + /* P7 Wave D: directory metadata collected by write_thread from received + directory entries. Only write_thread mutates it (before it joins); the + caller (server.c) applies it after the delete/delay-updates phase. */ + DirTimeList dir_times; +} PipelineContextReceiver; + +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/server/server.c b/src/server/server.c index da3e65d..494a840 100644 --- a/src/server/server.c +++ b/src/server/server.c @@ -7,10 +7,10 @@ #include "identity.h" #include "log.h" #include "motd.h" -#include "multiprocessing.h" #include "protocol.h" #include "queue.h" #include "receiver.h" +#include "receiver_pipeline.h" #include "server_cli.h" #include "transport_tcp.h" #include "transport_tls.h" diff --git a/src/shared/multiprocessing.c b/src/shared/multiprocessing.c index 14ad46e..2155837 100644 --- a/src/shared/multiprocessing.c +++ b/src/shared/multiprocessing.c @@ -1,5 +1,4 @@ #include "multiprocessing.h" -#include "receiver.h" #include "array_list.h" #include "chunk.h" @@ -206,248 +205,3 @@ void pipeline_context_sender_destroy(PipelineContextSender* context) { mtx_destroy(&context->mutex_progress); free(context); } - -PipelineContextReceiver* pipeline_context_receiver_create(Config* config, Queue* queue, - int file_descriptor, SSL* ssl) { - PipelineContextReceiver* context = malloc(sizeof(PipelineContextReceiver)); - if (context == NULL) - return NULL; - context->config = config; - context->queue = queue; - context->file_descriptor = file_descriptor; - context->ssl = ssl; - context->outcomes.entries = NULL; - context->outcomes.count = 0; - context->outcomes.capacity = 0; - dir_time_list_init(&context->dir_times); - 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; - context->deferred_manifest = NULL; - atomic_init(&context->cancelled, false); - int init = 0; - if (mtx_init(&context->mutex, mtx_plain) != thrd_success) - goto fail; - init++; - if (cnd_init(&context->condition_not_full) != thrd_success) - goto fail; - init++; - if (cnd_init(&context->condition_not_empty) != thrd_success) - goto fail; - // cppcheck-suppress unreadVariable - init++; - return context; - -fail: - log_perror("Error initializing synchronization objects"); - if (init >= 3) - cnd_destroy(&context->condition_not_empty); - if (init >= 2) - cnd_destroy(&context->condition_not_full); - if (init >= 1) - mtx_destroy(&context->mutex); - free(context); - return NULL; -} - -void pipeline_context_receiver_destroy(PipelineContextReceiver* context) { - config_delete(context->config); - if (context->deferred_manifest) - delete_manifest_free(context->deferred_manifest); - queue_destroy(context->queue); - receiver_outcomes_destroy(&context->outcomes); - dir_time_list_free(&context->dir_times); - mtx_destroy(&context->mutex); - cnd_destroy(&context->condition_not_full); - cnd_destroy(&context->condition_not_empty); - 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 = (PipelineContextReceiver*)context_pointer; - return pipeline_context_receiver_enqueue_file(context, file); -} - -static void receiver_thread_fail(PipelineContextReceiver* context) { - mtx_lock(&context->mutex); - atomic_store(&context->cancelled, true); - context->receiver_done = true; - cnd_broadcast(&context->condition_not_empty); - cnd_broadcast(&context->condition_not_full); - mtx_unlock(&context->mutex); -} - -int receive_thread(void* pipeline_context) { - PipelineContextReceiver* context = (PipelineContextReceiver*)pipeline_context; - protocol_session_bind(&context->session); - mtx_lock(&context->mutex); - int file_descriptor = context->file_descriptor; - const Config* config = context->config; - mtx_unlock(&context->mutex); - - ReceiverSink sink = {receiver_enqueue_file, context, false, false, NULL}; - if (receiver_process_pending((Config*)config, file_descriptor, &sink, - &context->deferred_manifest) != 0) { - receiver_thread_fail(context); - protocol_session_unbind(); - return thrd_error; - } - mtx_lock(&context->mutex); - context->receiver_done = true; - cnd_signal(&context->condition_not_empty); - mtx_unlock(&context->mutex); - protocol_session_unbind(); - return thrd_success; -} - -int write_thread(void* pipeline_context) { - PipelineContextReceiver* context = (PipelineContextReceiver*)pipeline_context; - protocol_session_bind(&context->session); - mtx_lock(&context->mutex); - bool save_to_disk = context->config->save_to_disk; - char* root_directory = str_dup(context->config->receive_root_directory); - mtx_unlock(&context->mutex); - if (save_to_disk && !root_directory) { - mtx_lock(&context->mutex); - atomic_store(&context->cancelled, true); - context->receiver_done = true; - cnd_broadcast(&context->condition_not_full); - cnd_broadcast(&context->condition_not_empty); - mtx_unlock(&context->mutex); - protocol_session_unbind(); - return thrd_error; - } - - while (true) { - File* file = - queue_dequeue_multithreaded(context->queue, &context->mutex, &context->condition_not_empty, - &context->condition_not_full, &context->receiver_done); - if (file == NULL) { - free(root_directory); - 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; - cnd_broadcast(&context->condition_not_full); - cnd_broadcast(&context->condition_not_empty); - mtx_unlock(&context->mutex); - free(root_directory); - protocol_session_unbind(); - return thrd_error; - } - } - /* P7 Wave D: a directory's times are never applied inline (a later child - write would clobber them); accumulate the metadata here and let the - caller apply it once every writer has drained. */ - if (result != FILE_SAVE_ERROR && file->is_dir && file->metadata && - dir_times_should_capture(context->config) && - !dir_time_list_add(&context->dir_times, file->path, file->metadata)) { - 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; - cnd_broadcast(&context->condition_not_full); - cnd_broadcast(&context->condition_not_empty); - mtx_unlock(&context->mutex); - free(root_directory); - protocol_session_unbind(); - return thrd_error; - } - /* Record the per-file outcome so a --remove-source-files sender learns - which sources were actually written versus skipped on the receiver. - Explicit directory entries and recreated device/special nodes have no - source and are never acknowledged (mirrors receiver.c). */ - if (context->config->remove_source_files && !file->is_dir && !file->is_special && !file->skip && - !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; - cnd_broadcast(&context->condition_not_full); - cnd_broadcast(&context->condition_not_empty); - mtx_unlock(&context->mutex); - free(root_directory); - protocol_session_unbind(); - 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 6de2e98..4bf1fac 100644 --- a/src/shared/multiprocessing.h +++ b/src/shared/multiprocessing.h @@ -10,7 +10,6 @@ #include "file.h" #include "protocol.h" #include "queue.h" -#include "receiver.h" #include "stop_condition.h" #include @@ -86,40 +85,6 @@ typedef struct { bool dir_entries_mutex_init; } PipelineContextSender; -typedef struct PipelineContextReceiver { - Queue* queue; - Config* config; - int file_descriptor; - SSL* ssl; - ProtocolSession session; - ReceiverOutcomes outcomes; - mtx_t mutex; - cnd_t condition_not_full; - 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; - /* Keep-set manifest for the commit-style (late) deletion - (--delete/--delete-after/--delete-delay). receive_thread parses the whole - protocol stream but hands the manifest here instead of deleting while the - disk writer may still be draining; the caller (server.c) commits the - deletion after both threads have joined, so no extra is removed unless the - transfer truly succeeded. NULL in the early delete modes (which delete at - the manifest). */ - DeleteManifest* deferred_manifest; - /* P7 Wave D: directory metadata collected by write_thread from received - directory entries. Only write_thread mutates it (before it joins); the - caller (server.c) applies it after the delete/delay-updates phase. */ - DirTimeList dir_times; -} PipelineContextReceiver; - /* `config` is borrowed and must outlive the context: destroy does NOT free it, so the caller owns it and frees it with config_delete() afterwards. */ PipelineContextSender* pipeline_context_sender_create(Config* config, Queue* queue_scanner, @@ -141,21 +106,4 @@ bool pipeline_context_sender_enqueue_chunk(PipelineContextSender* context, Chunk 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); -/* 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/tests/test_config.c b/tests/test_config.c index c9e672e..7689cf5 100644 --- a/tests/test_config.c +++ b/tests/test_config.c @@ -4,6 +4,7 @@ #include "multiprocessing.h" #include "protocol.h" #include "queue.h" +#include "receiver_pipeline.h" #include "test_utils.h" #include "utils.h" #include diff --git a/tests/test_multiprocessing.c b/tests/test_multiprocessing.c index a051897..1a0e6f8 100644 --- a/tests/test_multiprocessing.c +++ b/tests/test_multiprocessing.c @@ -3,6 +3,7 @@ #include "config.h" #include "protocol.h" #include "queue.h" +#include "receiver_pipeline.h" #include "utils.h" #include "test_utils.h" #include From 0155902d95b03ab5778357ce0101b4e644e2b0ae Mon Sep 17 00:00:00 2001 From: TapTap Date: Sun, 13 Sep 2026 07:30:28 +0200 Subject: [PATCH 2/2] docs: update stale PipelineContextReceiver reference --- src/shared/delay_updates.h | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/src/shared/delay_updates.h b/src/shared/delay_updates.h index 225c19b..9e03311 100644 --- a/src/shared/delay_updates.h +++ b/src/shared/delay_updates.h @@ -18,8 +18,9 @@ typedef struct { /* Receiver-side --delay-updates staging registry. All successfully written files land under a private staging directory inside the receive root and are atomically renamed into their final destination only at the very end of the - transfer. A single PipelineContextReceiver has exactly one writer thread, - but the registry is still mutex-protected so the same object can be safely + transfer. A single receiver pipeline (see src/server/receiver_pipeline.h) + has exactly one writer thread, but the registry is still mutex-protected so + the same object can be safely shared with the publish/cleanup phase that runs after the threads join. */ typedef struct DelayUpdatesContext { char* root_directory; /* receive root the staging dir lives under */