Files
FastSync/src/shared/multiprocessing.h
T
TapTap 9749c7878c fix(p7-times): dir-time entries only record (never create dirs); bound/chunk dir-time frames; harden list add; docs+tests
Review fixes for Phase 7 Wave D.

#1 (HIGH): STATUS_DIR_TIMES entries no longer create directories. A new
receiver-only File.dir_time_only flag marks dir-time entries; file_save_to_disk_full
short-circuits them as FILE_SAVE_SKIPPED before any device/dir branch, so the sink
still accumulates metadata into the deferred DirTimeList but creates nothing. Empty
source dirs stay untransferred (-a), -m/--prune-empty-dirs semantics are preserved,
and a pre-existing regular file/symlink at an empty-dir mirror path no longer aborts
the transfer. dir_time_list_apply fstatat()s the leaf (AT_SYMLINK_NOFOLLOW) and skips
absent/non-directory paths QUIETLY; only a real existing directory is stamped.
Also initialize File.dir_time_only in file_create() (uninitialised garbage otherwise).

#2 (MED): send_dir_times() chunks entries into repeated STATUS_DIR_TIMES frames of at
most MAX_MANIFEST_ENTRIES, matching the receiver's per-frame bound; the tautological
> INT_MAX check is gone.

#3 (LOW): dir_time_list_add() assigns each grown array right after its realloc (no
dangling) and advances capacity only after both succeed.

#4 (LOW): RSYNC_COMPAT.md -- STATUS_MKDIR carries metadata, dir times are transmitted
via STATUS_DIR_TIMES and applied at the end, empty dirs are still never created; -m
rationale, -O row and Wave D notes updated. Summary counts untouched.

#5 (LOW): integration tests for the three #1 scenarios (empty-dir non-creation under
-a and -a -m, collision non-abort), scanner test now covers empty-dir capture, and
test_file_restore_symlink_metadata asserts the positive apply path when supported.

PROTOCOL_VERSION stays 2.17.0; config-frame layout unchanged.
2026-09-12 11:06:05 +02:00

134 lines
6.3 KiB
C

#ifndef MULTIPROCESSING_H
#define MULTIPROCESSING_H
#include <threads.h>
#include <stdatomic.h>
#include "array_list.h"
#include "config.h"
#include "file.h"
#include "protocol.h"
#include "queue.h"
#include "receiver.h"
#include "stop_condition.h"
#include <openssl/ssl.h>
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;
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;
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;
PipelineContextSender* pipeline_context_sender_create(Config* config, Queue* queue_scanner,
Queue* queue_loader);
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