#ifndef MULTIPROCESSING_H #define MULTIPROCESSING_H #include #include #include "array_list.h" #include "chunk.h" #include "config.h" #include "file.h" #include "protocol.h" #include "queue.h" #include "stop_condition.h" #include typedef struct { Config* config; Queue* queue_scanner; mtx_t mutex_scanner; cnd_t condition_not_full_scanner; cnd_t condition_not_empty_scanner; bool scanner_done; Queue* queue_loader; mtx_t mutex_loader; 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 --delete-excluded is set. NULL when not collecting. Populated by the scanner thread (parallel workers append under mutex_scanner via the scanner's exclusion sink) or, in the early modes, by the path-only pre-scan on the calling thread before the pipeline starts. */ ArrayList* excluded_paths; /* --delete-missing-args: the destination-relative mirrors of the --files-from entries that are missing under the source. Computed by the preflight on the calling thread before the pipeline starts; the sender thread transmits them in the manifest frame's third section and the receiver deletes each as an explicit request. */ ArrayList* missing_args; /* A source I/O error (unreadable directory) was recorded during the scan. Set by the pre-scan (before the threads start) or by the scanner thread under mutex_scanner; the caller turns it into a non-zero exit when --ignore-errors kept the run going. */ bool scan_had_io_error; ArrayList* remove_source_files; /* True when --delete-before/--delete-during require the keep-set manifest to be transmitted before any file data: context->manifest is then prebuilt by a path-only pre-scan on the calling thread and the pipeline scanner must not append to it. Set once before the worker threads start. */ bool early_delete; mtx_t mutex_progress; int total_files; unsigned long long progress_bytes; unsigned long long total_bytes; bool sender_done; atomic_bool cancelled; ProtocolSession allocation_session; /* Phase 6: client-only sender stop deadline, computed once before the worker * threads start and shared read-only by the scanner and the sender thread. */ StopCondition stop_condition; /* Phase 6: set when the scanner/sender reached the stop deadline before the * scan (and thus the keep-set manifest) completed naturally. When true the * completion tail must NOT transmit the partial manifest, or the receiver * would delete unscanned source mirrors. Written by the sender thread * before it reads the manifest, so no additional synchronization is needed * to suppress the manifest. */ bool scan_stopped_early; /* P7 Wave D: captured source directory times, filled by the scanner thread * (and its parallel workers, guarded by dir_entries_mutex) and drained by the * sender thread in trailing STATUS_DIR_TIMES frame(s). Owned by the * context; NULL for non-metadata transfers. */ ArrayList* dir_entries; mtx_t dir_entries_mutex; bool dir_entries_mutex_init; } PipelineContextSender; /* `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, 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); #endif