Merge feat/p3-delete-timing: rsync delete timing (--delete-before/--delete-during/--del/--delete-after/--delete-delay)

This commit is contained in:
2026-09-06 19:58:40 +02:00
22 changed files with 934 additions and 105 deletions
+10 -2
View File
@@ -102,7 +102,11 @@ partial, alternate, and planned behavior.
| `-q, --quiet` | Suppress non-error output | | `-q, --quiet` | Suppress non-error output |
| `--progress` | Show real-time transfer speed | | `--progress` | Show real-time transfer speed |
| `-P` | Enables partial-transfer mode and progress output (partial retention is incomplete) | | `-P` | Enables partial-transfer mode and progress output (partial retention is incomplete) |
| `--delete` | Delete files on receiver not present in source | | `--delete` | Delete files on receiver not present in source (default timing: delete-after, i.e. only after the whole transfer succeeded) |
| `--delete-before` | Delete extras before the transfer starts (implies `--delete`) |
| `--delete-during`, `--del` | Delete extras once the keep-set is known, before data is applied (implies `--delete`) |
| `--delete-delay` | Delete extras only after a successful transfer (implies `--delete`) |
| `--delete-after` | Explicit delete-after timing (implies `--delete`) |
| `--exclude <pattern>` | Exclude files matching glob pattern (repeatable) | | `--exclude <pattern>` | Exclude files matching glob pattern (repeatable) |
| `--exclude-from <file>` | Read exclude patterns from a file (one per line) | | `--exclude-from <file>` | Read exclude patterns from a file (one per line) |
| `--include <pattern>` | Only transfer files matching glob pattern (repeatable, whitelist) | | `--include <pattern>` | Only transfer files matching glob pattern (repeatable, whitelist) |
@@ -386,7 +390,11 @@ option.
|---|---| |---|---|
| `-a`, `--archive` | Enable current archive preset. Full rsync archive semantics are planned. | | `-a`, `--archive` | Enable current archive preset. Full rsync archive semantics are planned. |
| `-n`, `--dry-run` | Scan and report without writing files. | | `-n`, `--dry-run` | Scan and report without writing files. |
| `--delete` | Request removal of destination entries absent from the source. The server must allow deletion. | | `--delete` | Request removal of destination entries absent from the source. The server must allow deletion. Default timing is delete-after: extras are removed only after the whole transfer succeeded. |
| `--delete-before` | Delete extras before the transfer starts (implies `--delete`). |
| `--delete-during`, `--del` | Delete extras once the keep-set manifest is known, before data is applied (implies `--delete`; early mode, same engine behaviour as `--delete-before`). |
| `--delete-delay` | Delete extras only after a successful transfer (implies `--delete`; commit mode, same behaviour as `--delete-after`). |
| `--delete-after` | Explicit delete-after timing: delete only after the transfer succeeded (implies `--delete`). |
| `--exclude <pattern>` | Exclude matching paths. Repeatable. | | `--exclude <pattern>` | Exclude matching paths. Repeatable. |
| `--include <pattern>` | Include matching paths. Repeatable. | | `--include <pattern>` | Include matching paths. Repeatable. |
| `--exclude-from <file>` | Read exclude patterns from a file. | | `--exclude-from <file>` | Read exclude patterns from a file. |
+44 -5
View File
@@ -103,17 +103,56 @@ This document maps rsync's full feature set to FastSync's current implementation
| Flag | Rsync Description | FastSync Status | Notes | | Flag | Rsync Description | FastSync Status | Notes |
|------|-------------------|-----------------|-------| |------|-------------------|-----------------|-------|
| `--delete` | Delete extraneous files from dest | ✅ Implemented | `use_delete` config field | | `--delete` | Delete extraneous files from dest | ✅ Implemented | `use_delete` config field. Deletion is always derived from the transmitted keep-set manifest of the paths the sender sent/keeps (never from unchecked input), runs through the symlink-safe walker bounded by `MAX_SERVER_DELETE_COUNT`, and skips the `.fastsync-stage` staging dir under `--delay-updates`. FastSync's default timing when no timing flag is given is **delete-after** (extras are removed only once the whole transfer succeeded) — intentionally NOT rsync's `--del`/delete-during default, to preserve FastSync's commit-style safety |
| `--delete-before` | Delete before transfer | ❌ Not Implemented | Removed because it had no effect | | `--delete-before` | Delete before transfer | ✅ Implemented | Implies `--delete`. The sender runs a full source pre-scan (paths only) and transmits the keep-set manifest BEFORE any file data; the receiver validates it, removes every destination entry not listed (bounded walk, staging-dir skip), then acks `STATUS_OK`. The sender only starts streaming after the deletion committed, or aborts if the receiver reported a deletion error. By definition the deletions already happened when a later transfer phase fails — rsync's delete-before is destructive the same way; a subsequent failure does not restore the removed files. Divergence: the keep-set is the pre-scan snapshot, so a file that appears on the source between the pre-scan and the data pass is still transferred but was not protected from deletion |
| `--del`, `--delete-during` | Delete during transfer | ❌ Not Implemented | Both flags are recognized but rejected; delete timing is not implemented | | `--del`, `--delete-during` | Delete during transfer | ✅ Implemented | Both spellings accepted; imply `--delete`. FastSync streams the source in a single directory scan and has no per-directory generator pass, so deletions cannot be interleaved per-directory the way rsync's delete-during does. `--delete-during` therefore selects the same early engine mode as `--delete-before` (manifest transmitted before any data, extras removed and acknowledged before data is applied); observable success/failure behaviour equals `--delete-before`. That is the documented divergence from rsync, where `--del` is the default meaning of `--delete` |
| `--delete-delay` | Find deletions during, delete after | ❌ Not Implemented | | | `--delete-delay` | Find deletions during, delete after | ✅ Implemented | Implies `--delete`. Commit-mode timing: extras are removed only after the whole transfer succeeded. rsync's delete-delay records the deletion list during its scan and applies it at the end; FastSync never snapshots the destination while data flows (the keep-set is the transmitted manifest and the destination is listed only at deletion time), so `--delete-delay` is implemented as the same end-of-transfer commit as `--delete-after` with identical safety. That is the documented divergence |
| `--delete-after` | Delete after transfer | ❌ Not Implemented | Removed because it had no effect | | `--delete-after` | Delete after transfer | ✅ Implemented | Implies `--delete`. The delete-after timing is also what plain `--delete` does: the keep-set manifest closes the data stream and the receiver commits the bounded deletion only after the terminal `STATUS_FINISHED` proves the whole transfer (every data frame received and stored) succeeded. A failed or aborted transfer removes nothing |
| `--delete-excluded` | Also delete excluded files | ❌ Not Implemented | Removed because it had no effect | | `--delete-excluded` | Also delete excluded files | ❌ Not Implemented | Removed because it had no effect |
| `--max-delete=NUM` | Max files to delete | ❌ Not Implemented | Removed because it had no effect | | `--max-delete=NUM` | Max files to delete | ❌ Not Implemented | Removed because it had no effect |
| `--ignore-errors` | Delete even with I/O errors | ❌ Not Implemented | | | `--ignore-errors` | Delete even with I/O errors | ❌ Not Implemented | |
| `--force` | Force deletion of non-empty dirs | ❌ Not Implemented | | | `--force` | Force deletion of non-empty dirs | ❌ Not Implemented | |
| `--prune-empty-dirs` | Prune empty dir chains | ❌ Not Implemented | Removed because it had no effect | | `--prune-empty-dirs` | Prune empty dir chains | ❌ Not Implemented | Removed because it had no effect |
**Deletion-timing implementation notes (Phase 3):** the delete flags above are
real. Two new config booleans (`delete_during`, `delete_delay`) join the already
serialized `delete_before`/`delete_after`, so the on-the-wire config layout
changed and `PROTOCOL_VERSION` was bumped **2.7.0 → 2.8.0** (peers must match).
The `STATUS_MANIFEST` frame is count-delimited and position-independent: the
receiver commits the deletion either when the manifest arrives (early modes:
`--delete-before`/`--delete-during`, which additionally acknowledge with
`STATUS_OK` before data flows) or after the terminal `STATUS_FINISHED` proves
the whole transfer succeeded (commit modes: plain `--delete`/`--delete-after`/
`--delete-delay`). Timing is chosen purely from the config, so server policy
(`--allow-delete` off) still disables deletion without deadlocking the early
manifest ack. `--delete-delay` and `--delete-during` are each implemented as
the closest safe approximation their engine mode allows; the divergences are
noted in the rows above.
Manifest size: the sender's keep-set collection (streaming or early pre-scan)
is unbounded, but the receiver rejects any manifest beyond `MAX_MANIFEST_ENTRIES`
(1 048 576 entries) / `MAX_MANIFEST_BYTES` (16 MB of paths) as a hard protocol
error. In the commit modes this only means the deletion is refused after the
data already arrived; in the NEW early modes (`--delete-before`/`--delete-during`)
the manifest is the first frame, so an oversized keep-set now aborts the whole
transfer BEFORE any data is sent (previously all data transferred and only the
deletion step failed). Keep the source tree small enough for the receiver's
manifest caps when using the early timing.
Early-delete ACK wait: after committing a large deletion (up to
`MAX_SERVER_DELETE_COUNT` unlinks) the receiver's `STATUS_OK`/`STATUS_ERROR`
reply can legitimately take much longer than a normal round trip, so the sender
waits for that single ACK with an extended explicit deadline (1 hour) instead
of the default 60 s per-message receive window. A receiver that is genuinely
gone still aborts the wait via connection close/error; the extended bound only
protects against aborting after the deletion already committed on the receiver.
Flag-conflict policy: unlike rsync's last-one-wins behaviour, every deletion
timing flag implies `--delete`, and combining a timing flag with `--no-delete`
(in either argument order) — or more than one timing flag — is rejected as a
configuration error rather than silently resolved. Note the check is
order-independent because it runs over the fully parsed config.
## 8. Metadata Preservation ## 8. Metadata Preservation
| Flag | Rsync Description | FastSync Status | Notes | | Flag | Rsync Description | FastSync Status | Notes |
+20 -20
View File
@@ -370,7 +370,6 @@ typedef enum {
OPT_POS_INT, OPT_POS_INT,
OPT_NONNEG_INT, OPT_NONNEG_INT,
OPT_ULL, OPT_ULL,
OPT_UNSUPPORTED,
} OptKind; } OptKind;
typedef struct { typedef struct {
@@ -430,7 +429,10 @@ static const OptionEntry OPTION_TABLE[] = {
{"--old-d", NULL, OPT_FLAG, offsetof(Config, dirs)}, {"--old-d", NULL, OPT_FLAG, offsetof(Config, dirs)},
{"--relative", "-R", OPT_FLAG, offsetof(Config, relative)}, {"--relative", "-R", OPT_FLAG, offsetof(Config, relative)},
{"--mkpath", NULL, OPT_FLAG, offsetof(Config, mkpath)}, {"--mkpath", NULL, OPT_FLAG, offsetof(Config, mkpath)},
{"--delete-during", "--del", OPT_UNSUPPORTED, 0}, {"--delete-before", NULL, OPT_FLAG, offsetof(Config, delete_before)},
{"--delete-during", "--del", OPT_FLAG, offsetof(Config, delete_during)},
{"--delete-delay", NULL, OPT_FLAG, offsetof(Config, delete_delay)},
{"--delete-after", NULL, OPT_FLAG, offsetof(Config, delete_after)},
{"--source-dir", NULL, OPT_STRING, offsetof(Config, send_directory)}, {"--source-dir", NULL, OPT_STRING, offsetof(Config, send_directory)},
{"--dest-dir", NULL, OPT_STRING, offsetof(Config, receive_root_directory)}, {"--dest-dir", NULL, OPT_STRING, offsetof(Config, receive_root_directory)},
@@ -547,8 +549,7 @@ static int apply_negation(Config* config, const char* arg) {
return 0; return 0;
} }
static int apply_table_option(Config* config, const OptionEntry* entry, const char* option_name, static int apply_table_option(Config* config, const OptionEntry* entry, const char* value) {
const char* value) {
if (entry->kind == OPT_NOOP) if (entry->kind == OPT_NOOP)
return 0; return 0;
void* field = (char*)config + entry->offset; void* field = (char*)config + entry->offset;
@@ -578,11 +579,6 @@ static int apply_table_option(Config* config, const OptionEntry* entry, const ch
*(unsigned long long*)field = v; *(unsigned long long*)field = v;
return 0; return 0;
} }
case OPT_UNSUPPORTED: {
const char* reason = "delete-during is not implemented";
log_message(LOG_LEVEL_ERROR, "%s: %s; refusing to ignore option", option_name, reason);
return -1;
}
} }
return -1; return -1;
} }
@@ -665,23 +661,20 @@ int parse_args(Config* config, int argc, char* argv[], int* positional_args,
if (!entry) if (!entry)
entry = find_table_option_with_equals(argv[i], &inline_value); entry = find_table_option_with_equals(argv[i], &inline_value);
if (entry) { if (entry) {
const char* option_name = argv[i];
const char* value = NULL; const char* value = NULL;
if (entry->kind != OPT_FLAG) { if (entry->kind != OPT_FLAG) {
if (entry->kind != OPT_UNSUPPORTED) { value = inline_value;
value = inline_value; if (!value && i + 1 < argc)
if (!value && i + 1 < argc) value = argv[++i];
value = argv[++i]; if (!value) {
if (!value) { log_message(LOG_LEVEL_ERROR, "missing argument for %s", entry->name);
log_message(LOG_LEVEL_ERROR, "missing argument for %s", entry->name); return -1;
return -1;
}
} }
if (strcmp(entry->name, "--compress-choice") == 0) { if (strcmp(entry->name, "--compress-choice") == 0) {
if (set_compression_choice(config, value) != 0) if (set_compression_choice(config, value) != 0)
return -1; return -1;
} else { } else {
if (apply_table_option(config, entry, option_name, value) != 0) if (apply_table_option(config, entry, value) != 0)
return -1; return -1;
if (strcmp(entry->name, "--compress-level") == 0 && if (strcmp(entry->name, "--compress-level") == 0 &&
(config->compression_level < 1 || config->compression_level > 22)) { (config->compression_level < 1 || config->compression_level > 22)) {
@@ -697,11 +690,18 @@ int parse_args(Config* config, int argc, char* argv[], int* positional_args,
config->use_metadata = true; config->use_metadata = true;
} }
} }
} else if (apply_table_option(config, entry, option_name, NULL) != 0) { } else if (apply_table_option(config, entry, NULL) != 0) {
return -1; return -1;
} }
if (entry->offset == offsetof(Config, eight_bit_output)) if (entry->offset == offsetof(Config, eight_bit_output))
protocol_set_8_bit_output(true); protocol_set_8_bit_output(true);
/* A delete-timing flag selects when --delete removes extras, so it
implies --delete exactly like the rsync options do. */
if (entry->offset == offsetof(Config, delete_before) ||
entry->offset == offsetof(Config, delete_during) ||
entry->offset == offsetof(Config, delete_delay) ||
entry->offset == offsetof(Config, delete_after))
config->use_delete = true;
continue; continue;
} }
+131 -17
View File
@@ -254,10 +254,6 @@ static void disconnect_transfer_client(Client* client) {
client_delete(client); client_delete(client);
} }
static ArrayList* create_transfer_manifest(const Config* config) {
return config->use_delete ? array_list_create(free) : NULL;
}
static bool add_chunk_to_manifest(ArrayList* manifest, const Chunk* chunk) { static bool add_chunk_to_manifest(ArrayList* manifest, const Chunk* chunk) {
if (!manifest) if (!manifest)
return true; return true;
@@ -597,6 +593,59 @@ static int send_delete_manifest(int fd, ArrayList* manifest) {
return 0; return 0;
} }
/* Transmit the keep-set manifest and wait for the receiver's verdict. Used by
--delete-before/--delete-during, where the extras are removed on the receiver
BEFORE the first byte of file data is sent: the receiver acknowledges with
STATUS_OK once the bounded delete committed, or STATUS_ERROR if it could not
(in which case the sender aborts without streaming any data). The ACK may
take much longer than an ordinary per-message round trip because the receiver
performs the whole bounded deletion walk (up to MAX_SERVER_DELETE_COUNT
unlinks) before replying, so the wait uses a generous explicit deadline
instead of the default 60 s receive window. */
#define DELETE_ACK_TIMEOUT_SEC 3600
static bool send_delete_manifest_early(Client* client, ArrayList* manifest) {
if (!client || !manifest)
return false;
if (send_delete_manifest(client->file_descriptor, manifest) != 0)
return false;
Status ack;
if (!receive_status_timed(client->file_descriptor, &ack, DELETE_ACK_TIMEOUT_SEC))
return false;
if (ack != STATUS_OK) {
log_message(LOG_LEVEL_ERROR, "Server failed to delete files before the transfer");
return false;
}
return true;
}
/* Walk the whole source tree once collecting only destination-relative wire
paths, loading and sending nothing. --delete-before/--delete-during need the
complete keep-set manifest before the first data byte, so it is built by a
dedicated pre-scan pass and transmitted early; the data pass then re-scans
with a fresh scanner. */
static bool scan_paths_only(const Config* config, const ScannerOptions* options,
ArrayList* manifest) {
DirectoryScanner* scanner =
directory_scanner_create_with_options(config->send_directory, options);
if (!scanner)
return false;
bool ok = true;
Chunk* chunk;
while ((chunk = directory_scanner_next(scanner)) != NULL) {
if (!add_chunk_to_manifest(manifest, chunk)) {
ok = false;
chunk_destroy(chunk);
break;
}
chunk_destroy(chunk);
}
if (ok && directory_scanner_failed(scanner))
ok = false;
directory_scanner_destroy(scanner);
return ok;
}
static int incremental_check(Client* client, File* file, const Config* config, static int incremental_check(Client* client, File* file, const Config* config,
DeltaSignature** out_sig) { DeltaSignature** out_sig) {
*out_sig = NULL; *out_sig = NULL;
@@ -906,6 +955,17 @@ static int send_chunks_multithreaded(void* pipeline_context) {
protocol_session_unbind(); protocol_session_unbind();
return thrd_error; return thrd_error;
} }
if (context->early_delete) {
/* The keep-set manifest was prebuilt by a path-only pre-scan. Transmit it
and wait for the receiver to delete extras before streaming any data. */
if (!send_delete_manifest_early(client, context->manifest)) {
pipeline_cancel(context);
disconnect_transfer_client(client);
mark_sender_done(context);
protocol_session_unbind();
return thrd_error;
}
}
while (true) { while (true) {
Chunk* current_chunk = queue_dequeue_multithreaded( Chunk* current_chunk = queue_dequeue_multithreaded(
@@ -919,7 +979,7 @@ static int send_chunks_multithreaded(void* pipeline_context) {
protocol_session_unbind(); protocol_session_unbind();
return thrd_error; return thrd_error;
} }
if (context->config->use_delete) { if (context->config->use_delete && !context->early_delete) {
if (send_delete_manifest(client->file_descriptor, context->manifest) != 0) if (send_delete_manifest(client->file_descriptor, context->manifest) != 0)
goto send_fail; goto send_fail;
} }
@@ -1013,7 +1073,7 @@ static int scan_directory_multithreaded(void* pipeline_context) {
failed = dirs_mode ? directory_scanner_failed(dscanner) : parallel_scanner_failed(scanner); failed = dirs_mode ? directory_scanner_failed(dscanner) : parallel_scanner_failed(scanner);
break; break;
} }
if (context->config->use_delete) { if (context->config->use_delete && !context->early_delete) {
mtx_lock(&context->mutex_scanner); mtx_lock(&context->mutex_scanner);
bool manifest_ok = add_chunk_to_manifest(context->manifest, current_chunk); bool manifest_ok = add_chunk_to_manifest(context->manifest, current_chunk);
mtx_unlock(&context->mutex_scanner); mtx_unlock(&context->mutex_scanner);
@@ -1172,19 +1232,46 @@ int send_files(Config* config) {
DirectoryScanner* scanner = NULL; DirectoryScanner* scanner = NULL;
ArrayList* manifest = NULL; ArrayList* manifest = NULL;
ArrayList* remove_sources = NULL; ArrayList* remove_sources = NULL;
bool delete_early = config->use_delete && config_delete_timing_early(config);
bool send_failed = false;
PreparedScanner prepared; PreparedScanner prepared;
memset(&prepared, 0, sizeof(prepared)); memset(&prepared, 0, sizeof(prepared));
if (!config_send(client->file_descriptor, config)) if (!config_send(client->file_descriptor, config))
goto send_fail; goto send_fail;
if (!prepare_scanner(config, 0, &prepared)) if (!prepare_scanner(config, 0, &prepared))
goto send_fail; goto send_fail;
scanner = directory_scanner_create_with_options(config->send_directory, &prepared.options);
manifest = create_transfer_manifest(config);
if (config->remove_source_files) if (config->remove_source_files)
remove_sources = array_list_create(source_file_destroy); remove_sources = array_list_create(source_file_destroy);
if (!scanner || (config->use_delete && !manifest) || if (config->remove_source_files && !remove_sources)
(config->remove_source_files && !remove_sources))
goto send_fail; goto send_fail;
/* The late-timing modes (plain --delete / --delete-after / --delete-delay)
build the manifest while streaming and send it after the last data frame.
The early modes (--delete-before/--delete-during) send it up front from a
dedicated path-only pre-scan, so no manifest is kept during the data pass. */
if (delete_early) {
/* Pass 1: collect the complete keep-set (paths only, no data loaded) and
transmit it now, before any file data. The receiver removes extras and
acks; the transfer aborts here if the deletion could not commit. */
ArrayList* early_manifest = array_list_create(free);
if (!early_manifest)
goto send_fail;
if (!scan_paths_only(config, &prepared.options, early_manifest)) {
array_list_delete(early_manifest);
goto send_fail;
}
bool early_ok = send_delete_manifest_early(client, early_manifest);
array_list_delete(early_manifest);
if (!early_ok)
goto send_fail;
} else if (config->use_delete) {
manifest = array_list_create(free);
if (!manifest)
goto send_fail;
}
scanner = directory_scanner_create_with_options(config->send_directory, &prepared.options);
if (!scanner)
goto send_fail;
Chunk* current_chunk; Chunk* current_chunk;
unsigned long long total_bytes = 0; unsigned long long total_bytes = 0;
int total_files = 0; int total_files = 0;
@@ -1196,7 +1283,7 @@ int send_files(Config* config) {
chunk_bytes += current_chunk->items[i]->data->size; chunk_bytes += current_chunk->items[i]->data->size;
total_files++; total_files++;
} }
if (!add_chunk_to_manifest(manifest, current_chunk)) { if (manifest && !add_chunk_to_manifest(manifest, current_chunk)) {
chunk_destroy(current_chunk); chunk_destroy(current_chunk);
goto send_fail; goto send_fail;
} }
@@ -1220,9 +1307,7 @@ int send_files(Config* config) {
if (send_chunk_with_removal(client, current_chunk, config, remove_sources) != 0) { if (send_chunk_with_removal(client, current_chunk, config, remove_sources) != 0) {
log_message(LOG_LEVEL_ERROR, "Failed to send chunk"); log_message(LOG_LEVEL_ERROR, "Failed to send chunk");
chunk_destroy(current_chunk); chunk_destroy(current_chunk);
if (manifest) send_failed = true;
array_list_delete(manifest);
manifest = NULL;
break; break;
} }
total_bytes += chunk_bytes; total_bytes += chunk_bytes;
@@ -1235,9 +1320,18 @@ int send_files(Config* config) {
} }
chunk_destroy(current_chunk); chunk_destroy(current_chunk);
} }
if (directory_scanner_failed(scanner) || (config->use_delete && manifest == NULL)) if (send_failed) {
if (manifest) {
array_list_delete(manifest);
manifest = NULL;
}
goto send_fail; goto send_fail;
if (config->use_delete) { }
if (directory_scanner_failed(scanner))
goto send_fail;
if (manifest) {
/* Late (commit) ordering: all file data is out; transmit the keep-set
manifest so the receiver deletes only after the transfer succeeds. */
if (send_delete_manifest(client->file_descriptor, manifest) != 0) { if (send_delete_manifest(client->file_descriptor, manifest) != 0) {
array_list_delete(manifest); array_list_delete(manifest);
manifest = NULL; manifest = NULL;
@@ -1324,8 +1418,28 @@ int send_files_multithreaded(Config** config_ptr) {
return 1; return 1;
} }
*config_ptr = NULL; /* context now owns config through all remaining paths */ *config_ptr = NULL; /* context now owns config through all remaining paths */
if (config->use_delete) if (config->use_delete) {
context->manifest = array_list_create(free); context->manifest = array_list_create(free);
if (!context->manifest) {
pipeline_context_sender_destroy(context);
return 1;
}
if (config_delete_timing_early(config)) {
/* --delete-before/--delete-during: build the complete keep-set manifest
(paths only, nothing loaded or sent) up front so the sender thread can
transmit it before the first data byte. */
PreparedScanner prepared;
memset(&prepared, 0, sizeof(prepared));
bool prebuilt = prepare_scanner(config, 4, &prepared) &&
scan_paths_only(config, &prepared.options, context->manifest);
prepared_scanner_destroy(&prepared);
if (!prebuilt) {
pipeline_context_sender_destroy(context);
return 1;
}
context->early_delete = true;
}
}
if (config->remove_source_files) if (config->remove_source_files)
context->remove_source_files = array_list_create(source_file_destroy); context->remove_source_files = array_list_create(source_file_destroy);
if ((config->use_delete && !context->manifest) || if ((config->use_delete && !context->manifest) ||
+6
View File
@@ -71,5 +71,11 @@ bool validate_config(const Config* config) {
"staging directory)"); "staging directory)");
return false; return false;
} }
if (!config_has_valid_delete_timing(config)) {
log_message(LOG_LEVEL_ERROR,
"--delete-before/--delete-during/--delete-delay/--delete-after select the delete "
"timing; at most one may be given and each implies --delete");
return false;
}
return true; return true;
} }
+13 -1
View File
@@ -24,6 +24,19 @@ void print_usage(void) {
printf(" -P Partial mode with progress (retention incomplete)\n"); printf(" -P Partial mode with progress (retention incomplete)\n");
printf(" -8, --8-bit-output Leave high-bit characters unescaped in output\n"); printf(" -8, --8-bit-output Leave high-bit characters unescaped in output\n");
printf(" --delete Delete files on receiver not in source\n"); printf(" --delete Delete files on receiver not in source\n");
printf(" (default timing: delete only after the whole\n");
printf(" transfer has succeeded)\n");
printf(" --delete-before Delete extras before the transfer starts\n");
printf(" (implies --delete)\n");
printf(" --delete-during Delete extras once the keep-set manifest is known,\n");
printf(" before the data is applied (implies --delete)\n");
printf(" --del Alias for --delete-during\n");
printf(" --delete-delay Delete extras only after a successful transfer\n");
printf(" (implies --delete)\n");
printf(" --delete-after Delete only after the whole transfer succeeded\n");
printf(" (the default --delete timing; implies --delete)\n");
printf(" Note: each timing flag implies --delete. Combining a timing flag with\n");
printf(" --no-delete (in either order) is rejected as a config error.\n");
printf(" --ignore-existing Skip files that already exist on receiver\n"); printf(" --ignore-existing Skip files that already exist on receiver\n");
printf(" --delay-updates Put updated files into place only at the end of transfer\n"); printf(" --delay-updates Put updated files into place only at the end of transfer\n");
printf(" --dirs, -d, --old-dirs, --old-d Transfer the named directory entries without\n"); printf(" --dirs, -d, --old-dirs, --old-d Transfer the named directory entries without\n");
@@ -37,7 +50,6 @@ void print_usage(void) {
printf(" parent directory is not itself listed\n"); printf(" parent directory is not itself listed\n");
printf(" --mkpath Create the destination root directory on the server when it\n"); printf(" --mkpath Create the destination root directory on the server when it\n");
printf(" does not exist yet\n"); printf(" does not exist yet\n");
printf(" --del Alias for --delete-during (not implemented)\n");
printf(" --exclude <pattern> Exclude files matching pattern\n"); printf(" --exclude <pattern> Exclude files matching pattern\n");
printf(" --include <pattern> Only include files matching pattern\n"); printf(" --include <pattern> Only include files matching pattern\n");
printf(" --exclude-from <file> Read exclude patterns from file\n"); printf(" --exclude-from <file> Read exclude patterns from file\n");
+99 -11
View File
@@ -129,20 +129,40 @@ static bool receiver_process_batch(Config* config, int file_descriptor) {
} }
int receiver_process(Config* config, int file_descriptor, const ReceiverSink* sink) { int receiver_process(Config* config, int file_descriptor, const ReceiverSink* sink) {
return receiver_process_pending(config, file_descriptor, sink, NULL);
}
/* Runs the whole receive loop. The delete manifest may legitimately arrive
either FIRST (--delete-before / --delete-during: the sender transmits the
validated keep-set before any file data) or LAST (plain --delete /
--delete-after / --delete-delay: the manifest closes the data stream). In
the early modes the receiver deletes as soon as the manifest has been read
and acknowledges with STATUS_OK so the sender only starts streaming once the
deletion has committed (or failed); in the late modes the manifest is held
and the deletion is committed only after the terminal STATUS_FINISHED proves
the whole transfer succeeded. See receiver_process_pending() for how the -m
receiver defers that commit until its disk writer has drained. */
int receiver_process_pending(Config* config, int file_descriptor, const ReceiverSink* sink,
ArrayList** pending_manifest) {
Status status; Status status;
if (!receive_status(file_descriptor, &status)) if (!receive_status(file_descriptor, &status))
return -1; return -1;
bool early_delete = config_delete_timing_early(config);
/* Parked keep-set for the late/commit timing. Every exit path below frees it
exactly once; the only exception is the successful FINISHED handoff, which
transfers ownership to *pending_manifest (used by the -m receiver). */
ArrayList* deferred_manifest = NULL;
while (status == STATUS_NEXT || status == STATUS_CHUNK || status == STATUS_CHECK || while (status == STATUS_NEXT || status == STATUS_CHUNK || status == STATUS_CHECK ||
status == STATUS_KEEPALIVE || status == STATUS_ABORT || status == STATUS_CHECK_BATCH || status == STATUS_KEEPALIVE || status == STATUS_ABORT || status == STATUS_CHECK_BATCH ||
status == STATUS_MKDIR) { status == STATUS_MKDIR || status == STATUS_MANIFEST) {
if (status == STATUS_KEEPALIVE) { if (status == STATUS_KEEPALIVE) {
if (!send_status(file_descriptor, STATUS_KEEPALIVE)) if (!send_status(file_descriptor, STATUS_KEEPALIVE))
return -1; goto fail;
goto next; goto next_status;
} }
if (status == STATUS_ABORT) { if (status == STATUS_ABORT) {
log_message(LOG_LEVEL_INFO, "Received abort from client, cleaning up"); log_message(LOG_LEVEL_INFO, "Received abort from client, cleaning up");
return -1; goto fail;
} }
if (status == STATUS_CHECK) { if (status == STATUS_CHECK) {
bool skipped; bool skipped;
@@ -155,12 +175,46 @@ int receiver_process(Config* config, int file_descriptor, const ReceiverSink* si
goto receive_error; goto receive_error;
} else if (status == STATUS_CHECK_BATCH) { } else if (status == STATUS_CHECK_BATCH) {
if (!receiver_process_batch(config, file_descriptor)) if (!receiver_process_batch(config, file_descriptor))
return -1; goto fail;
goto next; goto next_status;
} else if (status == STATUS_MKDIR) { } else if (status == STATUS_MKDIR) {
File* dir = file_receive_directory(file_descriptor); File* dir = file_receive_directory(file_descriptor);
if (!dir || !sink->store_file(dir, sink->context)) if (!dir || !sink->store_file(dir, sink->context))
goto receive_error; goto receive_error;
} else if (status == STATUS_MANIFEST) {
ArrayList* manifest = receive_manifest_entries(file_descriptor);
if (!manifest)
goto fail; /* receive_manifest_entries already sent STATUS_ERROR */
if (early_delete) {
/* --delete-before / --delete-during: the manifest is authoritative the
moment it arrives, before any file data. Delete now and acknowledge
so the sender only starts streaming once the deletion committed (or
failed). This is the rsync delete-before/delete-during window: a
later transfer failure does not restore these deletions. */
bool deletion_ok = config->use_delete ? manifest_delete_extras(config, manifest) : true;
array_list_delete(manifest);
if (!deletion_ok) {
send_status(file_descriptor, STATUS_ERROR);
goto fail;
}
if (!send_status(file_descriptor, STATUS_OK))
goto fail;
} else if (config->use_delete) {
/* Plain --delete / --delete-after / --delete-delay: hold the keep-set
and commit the deletion only after STATUS_FINISHED. */
if (deferred_manifest) {
log_message(LOG_LEVEL_ERROR, "Received a second delete manifest");
array_list_delete(deferred_manifest);
deferred_manifest = NULL;
array_list_delete(manifest);
send_status(file_descriptor, STATUS_ERROR);
goto fail;
}
deferred_manifest = manifest;
} else {
array_list_delete(manifest);
}
goto next_status;
} else { } else {
File* file = file_receive(config, file_descriptor); File* file = file_receive(config, file_descriptor);
if (!file) { if (!file) {
@@ -170,27 +224,61 @@ int receiver_process(Config* config, int file_descriptor, const ReceiverSink* si
if (!sink->store_file(file, sink->context)) if (!sink->store_file(file, sink->context))
goto receive_error; goto receive_error;
} }
next: next_status:
if (!receive_status(file_descriptor, &status)) if (!receive_status(file_descriptor, &status))
goto receive_error; goto receive_error;
} }
if (status == STATUS_MANIFEST && receive_manifest(file_descriptor, config, &status) != 0)
return -1;
if (status != STATUS_FINISHED) { if (status != STATUS_FINISHED) {
log_message(LOG_LEVEL_ERROR, "Did not receive FINISHED Status"); log_message(LOG_LEVEL_ERROR, "Did not receive FINISHED Status");
goto receive_error; goto receive_error;
} }
/* Commit-style (late) deletion: every data frame has been received and the
sender proved the whole tree with STATUS_FINISHED. The single-threaded
receiver stores files synchronously, so everything is on disk here and the
deletion can be committed before the --delay-updates publication in
send_success (the walker skips the staging dir, so staged files are never
treated as extras). The -m receiver passes `pending_manifest` because its
disk writer may still be draining; the caller commits after the writer has
joined so no extra file is removed unless the transfer is known to have
succeeded. */
if (deferred_manifest) {
if (pending_manifest) {
*pending_manifest = deferred_manifest;
deferred_manifest = NULL;
} else {
bool deletion_ok = manifest_delete_extras(config, deferred_manifest);
array_list_delete(deferred_manifest);
deferred_manifest = NULL;
if (!deletion_ok) {
send_status(file_descriptor, STATUS_ERROR);
goto fail;
}
}
}
if (sink->send_success) { if (sink->send_success) {
if (sink->send_success_frame) { if (sink->send_success_frame) {
if (!sink->send_success_frame(file_descriptor, sink->context)) if (!sink->send_success_frame(file_descriptor, sink->context))
return -1; goto fail;
} else if (!send_status(file_descriptor, STATUS_OK)) { } else if (!send_status(file_descriptor, STATUS_OK)) {
return -1; goto fail;
} }
} }
return 0; return 0;
fail:
/* Failure exits that must not (or already did) report a STATUS_ERROR. The
parked keep-set is dropped: never commit a deletion for a failed stream. */
if (deferred_manifest) {
array_list_delete(deferred_manifest);
deferred_manifest = NULL;
}
return -1;
receive_error: receive_error:
if (deferred_manifest) {
array_list_delete(deferred_manifest);
deferred_manifest = NULL;
}
if (sink->send_error) if (sink->send_error)
send_status(file_descriptor, STATUS_ERROR); send_status(file_descriptor, STATUS_ERROR);
return -1; return -1;
+8
View File
@@ -35,6 +35,14 @@ void receiver_outcomes_destroy(ReceiverOutcomes* outcomes);
bool receiver_send_final_success(int fd, const Config* config, const ReceiverOutcomes* outcomes); bool receiver_send_final_success(int fd, const Config* config, const ReceiverOutcomes* outcomes);
int receiver_process(Config* config, int file_descriptor, const ReceiverSink* sink); int receiver_process(Config* config, int file_descriptor, const ReceiverSink* sink);
/* receiver_process with an escape hatch for the commit-style (late) deletion:
when `pending_manifest` is non-NULL the receiver does NOT delete at
STATUS_FINISHED itself; instead it stores the owned keep-set manifest there
(leaving *pending_manifest untouched on early modes/errors) so the caller can
commit the deletion only after its disk writer has fully drained. Pass NULL
to keep the default behaviour (delete before the success frame). */
int receiver_process_pending(Config* config, int file_descriptor, const ReceiverSink* sink,
ArrayList** pending_manifest);
int receiver_receive_files(Config* config, int file_descriptor); int receiver_receive_files(Config* config, int file_descriptor);
#endif #endif
+14
View File
@@ -255,6 +255,20 @@ void handler(int file_descriptor) {
thrd_join(receiver, &receiver_result); thrd_join(receiver, &receiver_result);
thrd_join(writer, &writer_result); thrd_join(writer, &writer_result);
bool transfer_ok = receiver_result == thrd_success && writer_result == thrd_success; bool transfer_ok = receiver_result == thrd_success && writer_result == thrd_success;
if (transfer_ok) {
/* Commit-style (late) deletion: receive_thread handed the keep-set
manifest here instead of deleting while write_thread might still be
draining, so by now every file is on disk and the whole transfer is
known to have succeeded. Remove the extras before publishing a
--delay-updates run; the walker skips the staging directory. */
if (context->deferred_manifest) {
if (!manifest_delete_extras(config, context->deferred_manifest)) {
transfer_ok = false;
}
array_list_delete(context->deferred_manifest);
context->deferred_manifest = NULL;
}
}
if (transfer_ok) { if (transfer_ok) {
/* --delay-updates: receive_thread has finished the whole protocol stream /* --delay-updates: receive_thread has finished the whole protocol stream
(including manifest/delete handling) and write_thread has drained its (including manifest/delete handling) and write_thread has drained its
+29 -2
View File
@@ -115,6 +115,8 @@ static void config_set_defaults(Config* config) {
config->partial_dir = NULL; config->partial_dir = NULL;
config->suffix = NULL; config->suffix = NULL;
config->delete_before = false; config->delete_before = false;
config->delete_during = false;
config->delete_delay = false;
config->address = NULL; config->address = NULL;
config->bind_address = NULL; config->bind_address = NULL;
config->ipv6 = false; config->ipv6 = false;
@@ -161,12 +163,14 @@ static bool validate_received_config(const Config* config) {
valid_wire_bool(config->inplace) && valid_wire_bool(config->append) && valid_wire_bool(config->inplace) && valid_wire_bool(config->append) &&
valid_wire_bool(config->use_fsync) && valid_wire_bool(config->append_verify) && valid_wire_bool(config->use_fsync) && valid_wire_bool(config->append_verify) &&
valid_wire_bool(config->delete_excluded) && valid_wire_bool(config->delete_after) && valid_wire_bool(config->delete_excluded) && valid_wire_bool(config->delete_after) &&
valid_wire_bool(config->delete_delay) && valid_wire_bool(config->delete_during) &&
valid_wire_bool(config->relative) && valid_wire_bool(config->prune_empty_dirs) && valid_wire_bool(config->relative) && valid_wire_bool(config->prune_empty_dirs) &&
valid_wire_bool(config->delay_updates) && valid_wire_bool(config->mkpath) && valid_wire_bool(config->delay_updates) && valid_wire_bool(config->mkpath) &&
!(config->delay_updates && config->inplace) && !(config->delay_updates && config->inplace) &&
!(config->delay_updates && delay_updates_staging_name_conflict(config->backup_dir)) && !(config->delay_updates && delay_updates_staging_name_conflict(config->backup_dir)) &&
valid_wire_bool(config->partial) && valid_wire_bool(config->delete_before) && valid_wire_bool(config->partial) && valid_wire_bool(config->delete_before) &&
valid_wire_bool(config->checksum) && valid_wire_bool(config->eight_bit_output) && valid_wire_bool(config->checksum) && valid_wire_bool(config->eight_bit_output) &&
config_has_valid_delete_timing(config) &&
!(config->skip_compress_set && config->use_chunk_serialization) && !(config->skip_compress_set && config->use_chunk_serialization) &&
(!config->use_compression || (!config->use_compression ||
(config->compression_level >= 1 && config->compression_level <= 22)) && (config->compression_level >= 1 && config->compression_level <= 22)) &&
@@ -188,6 +192,26 @@ Config* config_create(void) {
return config; return config;
} }
bool config_delete_timing_early(const Config* config) {
if (!config)
return false;
return config->delete_before || config->delete_during;
}
/* A delete-timing flag is only meaningful together with --delete. At most one
of the four flags may be set; several simultaneous timings are a client bug
and are rejected on both ends. */
bool config_has_valid_delete_timing(const Config* config) {
if (!config)
return false;
if (!config->use_delete)
return !config->delete_before && !config->delete_during && !config->delete_delay &&
!config->delete_after;
int timing_count = (config->delete_before ? 1 : 0) + (config->delete_during ? 1 : 0) +
(config->delete_delay ? 1 : 0) + (config->delete_after ? 1 : 0);
return timing_count <= 1;
}
bool config_is_remote_dest(const char* s) { bool config_is_remote_dest(const char* s) {
if (s == NULL) if (s == NULL)
return false; return false;
@@ -311,7 +335,8 @@ static bool send_selection_options(int fd, const Config* c) {
send_int(fd, c->use_fsync) && send_int(fd, c->append_verify) && send_int(fd, c->use_fsync) && send_int(fd, c->append_verify) &&
send_int(fd, c->delete_excluded) && send_int(fd, c->delete_after) && send_int(fd, c->delete_excluded) && send_int(fd, c->delete_after) &&
send_n_data(fd, &c->max_delete, sizeof(c->max_delete)) && send_int(fd, c->relative) && send_n_data(fd, &c->max_delete, sizeof(c->max_delete)) && send_int(fd, c->relative) &&
send_int(fd, c->prune_empty_dirs) && send_int(fd, c->mkpath); send_int(fd, c->prune_empty_dirs) && send_int(fd, c->mkpath) &&
send_int(fd, c->delete_during) && send_int(fd, c->delete_delay);
} }
static bool send_skip_compress_options(int fd, const Config* c) { static bool send_skip_compress_options(int fd, const Config* c) {
@@ -418,7 +443,9 @@ static bool receive_selection_options(int fd, Config* c) {
return false; return false;
if (!receive_wire_bool(fd, &c->mkpath)) if (!receive_wire_bool(fd, &c->mkpath))
return false; return false;
return true; if (!receive_wire_bool(fd, &c->delete_during))
return false;
return receive_wire_bool(fd, &c->delete_delay);
} }
static bool receive_resume_options(int fd, Config* c) { static bool receive_resume_options(int fd, Config* c) {
+26 -1
View File
@@ -147,6 +147,19 @@ typedef struct Config {
// PR #179: Delete policies // PR #179: Delete policies
bool delete_before; bool delete_before;
/* rsync deletion-timing family (real from Phase 3). At most one of
delete_before / delete_during / delete_delay / delete_after may be set, and
only together with use_delete (the CLI implies --delete for each of them).
delete_before and delete_during select the EARLY engine mode: the keep-set
manifest is transmitted before any file data and extras are removed then,
acknowledged, before the first data byte. delete_delay and delete_after
select the LATE commit mode: extras are removed only after the whole
transfer has succeeded (plain --delete keeps this mode). The exact
semantics and the divergences from rsync are documented in RSYNC_COMPAT.md
and in config_delete_timing_early() below. */
bool delete_during;
bool delete_delay;
// PR #181: IPv6 and bind address // PR #181: IPv6 and bind address
char* address; char* address;
char* bind_address; char* bind_address;
@@ -174,7 +187,7 @@ typedef struct Config {
DelayUpdatesContext* delay_context; DelayUpdatesContext* delay_context;
} Config; } Config;
#define PROTOCOL_VERSION "2.7.0" #define PROTOCOL_VERSION "2.8.0"
#define DEFAULT_CHUNK_SIZE (10 * 1024 * 1024) #define DEFAULT_CHUNK_SIZE (10 * 1024 * 1024)
Config* config_create(void); Config* config_create(void);
@@ -184,4 +197,16 @@ Config* config_receive(int file_descriptor);
bool config_is_remote_dest(const char* s); bool config_is_remote_dest(const char* s);
void config_parse_ssh_dest(Config* config); void config_parse_ssh_dest(Config* config);
/* True when the negotiated delete timing performs the extra-file deletion
* BEFORE the transfer data (--delete-before / --delete-during). The flag is
* a pure function of the config and is used identically on the sender (to pick
* the manifest-first frame order) and the receiver (to delete when the early
* manifest arrives). When false the deletion is committed only after the whole
* transfer succeeded (--delete / --delete-after / --delete-delay). */
bool config_delete_timing_early(const Config* config);
/* Delete-timing sanity: with deletion enabled at most one timing flag may be
* set (none = the default delete-after commit timing); without deletion no
* timing flag may be set (each timing flag implies --delete). */
bool config_has_valid_delete_timing(const Config* config);
#endif #endif
+25 -34
View File
@@ -792,26 +792,27 @@ File* file_receive_directory(int file_descriptor) {
return file; return file;
} }
int receive_manifest(int fd, const Config* config, int* next_status) { /* Read a delete-manifest frame (the STATUS_MANIFEST leading code has already
if (!config) { been consumed): an entry count followed by that many destination-relative
send_status(fd, STATUS_ERROR); paths. The frame is self-delimiting (the count is authoritative), so the
return -1; caller decides what to do next and continues reading the following STATUS_*
} frame. Returns an owned ArrayList of validated path strings, or NULL after
int received_status = STATUS_ERROR; sending STATUS_ERROR when the frame is malformed (bad count, empty/absolute
int* status_out = next_status ? next_status : &received_status; path, path traversal, or an aggregate size beyond MAX_MANIFEST_BYTES). */
ArrayList* receive_manifest_entries(int fd) {
int count; int count;
if (!receive_int(fd, &count)) { if (!receive_int(fd, &count)) {
send_status(fd, STATUS_ERROR); send_status(fd, STATUS_ERROR);
return -1; return NULL;
} }
if (count < 0 || count > MAX_MANIFEST_ENTRIES) { if (count < 0 || count > MAX_MANIFEST_ENTRIES) {
send_status(fd, STATUS_ERROR); send_status(fd, STATUS_ERROR);
return -1; return NULL;
} }
ArrayList* manifest = array_list_create(free); ArrayList* manifest = array_list_create(free);
if (!manifest) { if (!manifest) {
send_status(fd, STATUS_ERROR); send_status(fd, STATUS_ERROR);
return -1; return NULL;
} }
size_t manifest_bytes = 0; size_t manifest_bytes = 0;
for (int i = 0; i < count; i++) { for (int i = 0; i < count; i++) {
@@ -823,32 +824,22 @@ int receive_manifest(int fd, const Config* config, int* next_status) {
free(s); free(s);
array_list_delete(manifest); array_list_delete(manifest);
send_status(fd, STATUS_ERROR); send_status(fd, STATUS_ERROR);
return -1; return NULL;
} }
} }
if (!receive_status(fd, status_out)) { return manifest;
array_list_delete(manifest); }
send_status(fd, STATUS_ERROR);
return -1; /* Remove every destination entry under the receive root that is not listed in
} `manifest`, bounded by MAX_SERVER_DELETE_COUNT, using the symlink-safe
/* Deletion is a commit operation: never perform it until the sender has delete walker. With --delay-updates the not-yet-published staging directory
completed the manifest frame successfully. */ is a direct child of the receive root and must not be treated as a set of
if (*status_out != STATUS_FINISHED || !config->use_delete) { extras. Prints a notice and returns true on success. */
array_list_delete(manifest); bool manifest_delete_extras(const Config* config, ArrayList* manifest) {
if (*status_out != STATUS_FINISHED) if (!config || !manifest)
send_status(fd, STATUS_ERROR); return false;
return *status_out == STATUS_FINISHED ? 0 : -1;
}
fprintf(stderr, "Deleting files not in manifest...\n"); fprintf(stderr, "Deleting files not in manifest...\n");
/* With --delay-updates the staged (not yet published) files live directly
under the receive root in the staging directory; the delete walker must
not treat them as extras or it would remove every staged file before it
can be published. */
const char* skip_staging = config->delay_updates ? DELAY_UPDATES_STAGING_DIR : NULL; const char* skip_staging = config->delay_updates ? DELAY_UPDATES_STAGING_DIR : NULL;
bool deletion_ok = delete_extras_limited(config->receive_root_directory, manifest, return delete_extras_limited(config->receive_root_directory, manifest, MAX_SERVER_DELETE_COUNT,
MAX_SERVER_DELETE_COUNT, skip_staging); skip_staging);
array_list_delete(manifest);
if (!deletion_ok)
send_status(fd, STATUS_ERROR);
return deletion_ok ? 0 : -1;
} }
+8 -1
View File
@@ -10,7 +10,14 @@
File* file_receive(const Config* config, int file_descriptor); File* file_receive(const Config* config, int file_descriptor);
File* file_receive_directory(int file_descriptor); File* file_receive_directory(int file_descriptor);
File* receive_incremental_check(int fd, const Config* config, bool* skipped); File* receive_incremental_check(int fd, const Config* config, bool* skipped);
int receive_manifest(int fd, const Config* config, int* next_status); /* Read a delete-manifest frame: entry count then paths (self-delimiting; the
leading STATUS_MANIFEST code has been consumed). Returns an owned path
ArrayList, or NULL after signalling STATUS_ERROR on a malformed frame. */
ArrayList* receive_manifest_entries(int fd);
/* Remove destination entries under config->receive_root_directory that are not
in `manifest` (bounded walk, staging-dir skip). The caller decides WHEN to
run it based on the negotiated delete timing. */
bool manifest_delete_extras(const Config* config, ArrayList* manifest);
/* Outcome of a single file_save_to_disk operation. The receiver needs to /* Outcome of a single file_save_to_disk operation. The receiver needs to
distinguish "written" from "skipped" so --remove-source-files can be told distinguish "written" from "skipped" so --remove-source-files can be told
+6 -1
View File
@@ -27,6 +27,7 @@ PipelineContextSender* pipeline_context_sender_create(Config* config, Queue* que
context->loader_done = false; context->loader_done = false;
context->manifest = NULL; context->manifest = NULL;
context->remove_source_files = NULL; context->remove_source_files = NULL;
context->early_delete = false;
context->total_files = 0; context->total_files = 0;
context->progress_bytes = 0; context->progress_bytes = 0;
context->total_bytes = 0; context->total_bytes = 0;
@@ -113,6 +114,7 @@ PipelineContextReceiver* pipeline_context_receiver_create(Config* config, Queue*
context->receiver_done = false; context->receiver_done = false;
context->queued_bytes = 0; context->queued_bytes = 0;
context->max_queue_bytes = 0; context->max_queue_bytes = 0;
context->deferred_manifest = NULL;
atomic_init(&context->cancelled, false); atomic_init(&context->cancelled, false);
int init = 0; int init = 0;
if (mtx_init(&context->mutex, mtx_plain) != thrd_success) if (mtx_init(&context->mutex, mtx_plain) != thrd_success)
@@ -141,6 +143,8 @@ fail:
void pipeline_context_receiver_destroy(PipelineContextReceiver* context) { void pipeline_context_receiver_destroy(PipelineContextReceiver* context) {
config_delete(context->config); config_delete(context->config);
if (context->deferred_manifest)
array_list_delete(context->deferred_manifest);
queue_destroy(context->queue); queue_destroy(context->queue);
receiver_outcomes_destroy(&context->outcomes); receiver_outcomes_destroy(&context->outcomes);
mtx_destroy(&context->mutex); mtx_destroy(&context->mutex);
@@ -236,7 +240,8 @@ int receive_thread(void* pipeline_context) {
mtx_unlock(&context->mutex); mtx_unlock(&context->mutex);
ReceiverSink sink = {receiver_enqueue_file, context, false, false, NULL}; ReceiverSink sink = {receiver_enqueue_file, context, false, false, NULL};
if (receiver_process((Config*)config, file_descriptor, &sink) != 0) { if (receiver_process_pending((Config*)config, file_descriptor, &sink,
&context->deferred_manifest) != 0) {
receiver_thread_fail(context); receiver_thread_fail(context);
protocol_session_unbind(); protocol_session_unbind();
return thrd_error; return thrd_error;
+13
View File
@@ -26,6 +26,11 @@ typedef struct {
bool loader_done; bool loader_done;
ArrayList* manifest; ArrayList* manifest;
ArrayList* remove_source_files; 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; mtx_t mutex_progress;
int total_files; int total_files;
unsigned long long progress_bytes; unsigned long long progress_bytes;
@@ -55,6 +60,14 @@ typedef struct PipelineContextReceiver {
budget instead of growing without bound. */ budget instead of growing without bound. */
size_t queued_bytes; size_t queued_bytes;
size_t max_queue_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). */
ArrayList* deferred_manifest;
} PipelineContextReceiver; } PipelineContextReceiver;
PipelineContextSender* pipeline_context_sender_create(Config* config, Queue* queue_scanner, PipelineContextSender* pipeline_context_sender_create(Config* config, Queue* queue_scanner,
+26 -2
View File
@@ -302,15 +302,25 @@ bool protocol_send_n_data(ProtocolSession* session, const void* data, size_t dat
return true; return true;
} }
bool protocol_receive_n_data_timed(ProtocolSession* session, void* data, size_t data_size,
int timeout_sec);
bool protocol_receive_n_data(ProtocolSession* session, void* data, size_t data_size) { bool protocol_receive_n_data(ProtocolSession* session, void* data, size_t data_size) {
return protocol_receive_n_data_timed(session, data, data_size, RECEIVE_TIMEOUT_SEC);
}
bool protocol_receive_n_data_timed(ProtocolSession* session, void* data, size_t data_size,
int timeout_sec) {
log_debug_message(LOG_DEBUG_IO, " Receiving n Data: %zu", data_size); log_debug_message(LOG_DEBUG_IO, " Receiving n Data: %zu", data_size);
if (!session) if (!session)
return false; return false;
int fd = session->read_fd; int fd = session->read_fd;
if (timeout_sec <= 0)
timeout_sec = RECEIVE_TIMEOUT_SEC;
struct timespec deadline; struct timespec deadline;
clock_gettime(CLOCK_MONOTONIC, &deadline); clock_gettime(CLOCK_MONOTONIC, &deadline);
deadline.tv_sec += RECEIVE_TIMEOUT_SEC; deadline.tv_sec += timeout_sec;
size_t total_bytes_received = 0; size_t total_bytes_received = 0;
short wait_events = POLLIN; short wait_events = POLLIN;
@@ -319,7 +329,7 @@ bool protocol_receive_n_data(ProtocolSession* session, void* data, size_t data_s
struct pollfd pfd = {.fd = fd, .events = wait_events}; struct pollfd pfd = {.fd = fd, .events = wait_events};
int poll_result = poll(&pfd, 1, deadline_remaining_ms(&deadline)); int poll_result = poll(&pfd, 1, deadline_remaining_ms(&deadline));
if (poll_result == 0) { if (poll_result == 0) {
log_message(LOG_LEVEL_ERROR, "Receive timeout after %ds", RECEIVE_TIMEOUT_SEC); log_message(LOG_LEVEL_ERROR, "Receive timeout after %ds", timeout_sec);
return false; return false;
} }
if (poll_result < 0) { if (poll_result < 0) {
@@ -516,6 +526,17 @@ bool protocol_receive_status(ProtocolSession* session, Status* status) {
return true; return true;
} }
/* protocol_receive_status with an explicit per-message deadline (seconds).
Used where a single reply may legitimately take far longer than the default
60 s receive window - e.g. the sender waiting for the early-delete ACK after
the receiver committed a large (up to MAX_SERVER_DELETE_COUNT) deletion. */
bool protocol_receive_status_timed(ProtocolSession* session, Status* status, int timeout_sec) {
if (!protocol_receive_n_data_timed(session, status, sizeof(Status), timeout_sec))
return false;
log_debug_message(LOG_DEBUG_PROTO, "Received Status: %s", status_to_string(*status));
return true;
}
bool send_str(int fd, const char* data) { bool send_str(int fd, const char* data) {
return protocol_send_str(legacy_session(-1, fd), data); return protocol_send_str(legacy_session(-1, fd), data);
} }
@@ -543,3 +564,6 @@ bool send_status(int fd, Status status) {
bool receive_status(int fd, Status* status) { bool receive_status(int fd, Status* status) {
return protocol_receive_status(legacy_session(fd, -1), status); return protocol_receive_status(legacy_session(fd, -1), status);
} }
bool receive_status_timed(int fd, Status* status, int timeout_sec) {
return protocol_receive_status_timed(legacy_session(fd, -1), status, timeout_sec);
}
+5
View File
@@ -112,5 +112,10 @@ bool send_int(int file_descriptor, int data);
bool receive_int(int file_descriptor, int* data); bool receive_int(int file_descriptor, int* data);
bool send_status(int file_descriptor, Status status); bool send_status(int file_descriptor, Status status);
bool receive_status(int file_descriptor, Status* status); bool receive_status(int file_descriptor, Status* status);
/* receive_status with an explicit per-message deadline in seconds, instead of
the default RECEIVE_TIMEOUT_SEC. A reply that may legitimately take longer
(e.g. the early-delete ACK after a large receiver-side deletion) must use
this so the sender does not abort after the deletion already committed. */
bool receive_status_timed(int file_descriptor, Status* status, int timeout_sec);
#endif #endif
+131
View File
@@ -1971,3 +1971,134 @@ class TestFilters:
"""--filter/-C/-F rule layer: excludes prune, ordering is first-match-wins, """--filter/-C/-F rule layer: excludes prune, ordering is first-match-wins,
the default with no matching rule is include, and legacy --exclude remains the default with no matching rule is include, and legacy --exclude remains
an independent layer.""" an independent layer."""
class TestDeleteTiming:
"""rsync deletion-timing family. --delete-before/--delete-during transmit
the keep-set manifest BEFORE any file data (the receiver deletes extras and
acks first); --delete/--delete-after/--delete-delay commit deletions only
after the whole transfer succeeded. Every timing flag implies --delete."""
def _seed(self, tag):
source = os.path.join(TEST_DATA_DIR, f"deltiming_{tag}_src")
clean_dir(source)
entries = {
"top.txt": b"top level\n",
"sub/deep.txt": b"deeply nested file\n",
}
for rel, content in entries.items():
full = os.path.join(source, rel)
os.makedirs(os.path.dirname(full), exist_ok=True)
with open(full, "wb") as fh:
fh.write(content)
return source
@pytest.mark.parametrize("flag", ["--delete-before", "--delete-during", "--del",
"--delete-after", "--delete-delay"])
@pytest.mark.parametrize("mt", [False, True])
def test_flag_removes_extras_on_success(self, flag, mt):
"""Every timing flag is accepted, implies --delete, and on a successful
transfer removes the destination extras exactly like plain --delete."""
source = self._seed("ok")
dest = os.path.join(TEST_DATA_DIR, "deltiming_ok_dst")
clean_dir(dest)
with ServerManager() as server:
server.start(extra_args=["--allow-delete"])
result, _ = run_client(source, dest, port=server.port)
assert result.returncode == 0, f"seed sync failed: {result.stderr[:200]}"
received = get_dest_received_dir(dest, source)
extra = os.path.join(received, "extra.txt")
with open(extra, "wb") as fh:
fh.write(b"should be deleted")
flags = [flag] + (["-m"] if mt else [])
result, _ = run_client(source, dest, flags=flags, port=server.port)
assert result.returncode == 0, \
f"{flag} sync failed: {(result.stderr or result.stdout)[:300]}"
assert not os.path.exists(extra), f"{flag} did not remove the extra file"
mismatches, missing = verify_transfer(source, received)
assert not missing, f"{flag} missing files: {missing}"
assert not mismatches, f"{flag} mismatched files: {mismatches}"
@pytest.mark.parametrize("flag", ["--delete-before", "--delete-during", "--del"])
@pytest.mark.parametrize("mt", [False, True])
def test_early_flags_delete_before_data(self, flag, mt):
"""--delete-before/--delete-during remove extras (and a file blocking a
destination directory) BEFORE data is applied, so a nested write that
would fail while the blocker still exists succeeds."""
source = self._seed("early")
dest = os.path.join(TEST_DATA_DIR, "deltiming_early_dst")
clean_dir(dest)
with ServerManager() as server:
server.start(extra_args=["--allow-delete"])
result, _ = run_client(source, dest, port=server.port)
assert result.returncode == 0, f"seed sync failed: {result.stderr[:200]}"
received = get_dest_received_dir(dest, source)
extra = os.path.join(received, "extra.txt")
with open(extra, "wb") as fh:
fh.write(b"extra file")
blocker = os.path.join(received, "sub")
shutil.rmtree(blocker)
with open(blocker, "wb") as fh:
fh.write(b"blocks the nested destination directory")
flags = [flag] + (["-m"] if mt else [])
result, _ = run_client(source, dest, flags=flags, port=server.port)
assert result.returncode == 0, \
f"{flag} (early delete) did not remove the blocker in time: " \
f"{(result.stderr or result.stdout)[:300]}"
assert not os.path.exists(extra), f"{flag} did not delete the extra before data"
assert _read_file(os.path.join(received, "sub", "deep.txt")) == b"deeply nested file\n", \
f"{flag}: nested file was not written after the early deletion"
@pytest.mark.parametrize("flag", ["--delete", "--delete-after", "--delete-delay"])
@pytest.mark.parametrize("mt", [False, True])
def test_late_flags_commit_only_after_success(self, flag, mt):
"""Plain --delete/--delete-after/--delete-delay defer deletion until the
whole transfer succeeds: a mid-transfer write failure must leave every
extra in place (commit-style safety). The -m receiver must also keep
the extras: the deferred keep-set is committed by the server only after
the disk-writer thread has finished, and a failing writer means the
manifest is freed, never applied."""
source = self._seed("late")
dest = os.path.join(TEST_DATA_DIR, "deltiming_late_dst")
clean_dir(dest)
with ServerManager() as server:
server.start(extra_args=["--allow-delete"])
result, _ = run_client(source, dest, port=server.port)
assert result.returncode == 0, f"seed sync failed: {result.stderr[:200]}"
received = get_dest_received_dir(dest, source)
extra = os.path.join(received, "extra.txt")
with open(extra, "wb") as fh:
fh.write(b"extra file")
blocker = os.path.join(received, "sub")
shutil.rmtree(blocker)
with open(blocker, "wb") as fh:
fh.write(b"blocks the nested destination directory")
flags = [flag] + (["-m"] if mt else [])
result, _ = run_client(source, dest, flags=flags, port=server.port)
assert result.returncode != 0, \
f"{flag} (mt={mt}) unexpectedly succeeded (deletion must be deferred)"
assert os.path.exists(extra), \
f"{flag} (mt={mt}) removed an extra although the transfer failed"
assert os.path.isfile(blocker), \
f"{flag} (mt={mt}) deleted the blocker although the transfer failed"
def test_early_flag_respected_when_server_refuses_delete(self, shared_server):
"""With an --allow-delete-less server the client's early timing still
completes (no deadlock on the pre-delete ack) and simply never deletes,
exactly like the plain server policy."""
source = self._seed("refused")
dest = os.path.join(TEST_DATA_DIR, "deltiming_refused_dst")
clean_dir(dest)
result, _ = run_client(source, dest, port=shared_server.port)
assert result.returncode == 0, f"seed sync failed: {result.stderr[:200]}"
received = get_dest_received_dir(dest, source)
extra = os.path.join(received, "extra.txt")
with open(extra, "wb") as fh:
fh.write(b"extra file")
result, _ = run_client(source, dest, flags=["--delete-before"], port=shared_server.port)
assert result.returncode == 0, \
f"--delete-before against a refuse-delete server failed: {result.stderr[:300]}"
assert os.path.exists(extra), "unauthorized delete removed an extra file"
+78 -7
View File
@@ -595,8 +595,9 @@ static void test_parse_args_relative_no_implied_mkpath() {
config_delete(cfg); config_delete(cfg);
} }
/* --del is recognized as the rsync alias, but its timing mode is not implemented. */ /* --del is accepted as the rsync alias for --delete-during: it enables
static void test_parse_args_delete_during_alias_unimplemented() { * deletion with the during (early) timing. */
static void test_parse_args_delete_during_alias() {
static const char* const options[] = {"--del", "--delete-during"}; static const char* const options[] = {"--del", "--delete-during"};
for (size_t i = 0; i < sizeof(options) / sizeof(options[0]); i++) { for (size_t i = 0; i < sizeof(options) / sizeof(options[0]); i++) {
@@ -605,12 +606,81 @@ static void test_parse_args_delete_during_alias_unimplemented() {
int positional_args[2]; int positional_args[2];
int positional_count = 0; int positional_count = 0;
EXPECT_EQ_INT(parse_args(cfg, 4, argv, positional_args, &positional_count), -1); EXPECT_EQ_INT(parse_args(cfg, 4, argv, positional_args, &positional_count), 0);
EXPECT_FALSE(cfg->use_delete); EXPECT_TRUE(cfg->use_delete);
EXPECT_TRUE(cfg->delete_during);
EXPECT_FALSE(cfg->delete_before);
EXPECT_FALSE(cfg->delete_delay);
EXPECT_FALSE(cfg->delete_after);
config_delete(cfg); config_delete(cfg);
} }
} }
/* Each rsync deletion-timing flag is accepted and implies --delete. */
static void test_parse_args_delete_timing_flags() {
Config* cfg = config_create();
char* argv[] = {"fastsync", "--delete-before", "/src", "/dst"};
int positional_args[2];
int positional_count = 0;
EXPECT_EQ_INT(parse_args(cfg, 4, argv, positional_args, &positional_count), 0);
EXPECT_TRUE(cfg->use_delete);
EXPECT_TRUE(cfg->delete_before);
config_delete(cfg);
cfg = config_create();
char* argv_after[] = {"fastsync", "--delete-after", "/src", "/dst"};
positional_count = 0;
EXPECT_EQ_INT(parse_args(cfg, 4, argv_after, positional_args, &positional_count), 0);
EXPECT_TRUE(cfg->use_delete);
EXPECT_TRUE(cfg->delete_after);
EXPECT_FALSE(cfg->delete_before);
config_delete(cfg);
cfg = config_create();
char* argv_delay[] = {"fastsync", "--delete-delay", "/src", "/dst"};
positional_count = 0;
EXPECT_EQ_INT(parse_args(cfg, 4, argv_delay, positional_args, &positional_count), 0);
EXPECT_TRUE(cfg->use_delete);
EXPECT_TRUE(cfg->delete_delay);
EXPECT_FALSE(cfg->delete_before);
EXPECT_FALSE(cfg->delete_after);
config_delete(cfg);
}
/* Two different delete-timing flags on one command line are a conflict, not a
* silent last-one-wins choice. */
static void test_parse_args_delete_timing_conflict_rejected() {
Config* cfg = config_create();
char* argv[] = {"fastsync", "--delete-before", "--delete-after", "/src", "/dst"};
int positional_args[2];
int positional_count = 0;
EXPECT_EQ_INT(parse_args(cfg, 5, argv, positional_args, &positional_count), 0);
EXPECT_TRUE(cfg->use_delete);
EXPECT_FALSE(validate_config(cfg));
config_delete(cfg);
cfg = config_create();
char* argv2[] = {"fastsync", "--delete-during", "--delete-delay", "/src", "/dst"};
positional_count = 0;
EXPECT_EQ_INT(parse_args(cfg, 5, argv2, positional_args, &positional_count), 0);
EXPECT_FALSE(validate_config(cfg));
config_delete(cfg);
}
/* A timing flag whose --delete was then negated away must be rejected: timing
* without deletion is meaningless. */
static void test_parse_args_delete_timing_without_delete_rejected() {
Config* cfg = config_create();
char* argv[] = {"fastsync", "--delete-before", "--no-delete", "/src", "/dst"};
int positional_args[2];
int positional_count = 0;
EXPECT_EQ_INT(parse_args(cfg, 5, argv, positional_args, &positional_count), 0);
EXPECT_FALSE(cfg->use_delete);
EXPECT_TRUE(cfg->delete_before);
EXPECT_FALSE(validate_config(cfg));
config_delete(cfg);
}
/* Parsed-but-unimplemented options must fail instead of being silently accepted. */ /* Parsed-but-unimplemented options must fail instead of being silently accepted. */
static void test_parse_args_rejects_unimplemented_options() { static void test_parse_args_rejects_unimplemented_options() {
static const char* const options[] = {"--silent", static const char* const options[] = {"--silent",
@@ -626,7 +696,6 @@ static void test_parse_args_rejects_unimplemented_options() {
"--append", "--append",
"--append-verify", "--append-verify",
"--delete-excluded", "--delete-excluded",
"--delete-after",
"--max-delete", "--max-delete",
"--prune-empty-dirs", "--prune-empty-dirs",
"-e", "-e",
@@ -635,7 +704,6 @@ static void test_parse_args_rejects_unimplemented_options() {
"--compare-dest", "--compare-dest",
"--copy-dest", "--copy-dest",
"--link-dest", "--link-dest",
"--delete-before",
"--address", "--address",
"--bind-address", "--bind-address",
"--ipv6", "--ipv6",
@@ -1542,7 +1610,10 @@ void test_client_cli() {
test_parse_args_unknown_option(); test_parse_args_unknown_option();
test_parse_args_dirs_aliases(); test_parse_args_dirs_aliases();
test_parse_args_relative_no_implied_mkpath(); test_parse_args_relative_no_implied_mkpath();
test_parse_args_delete_during_alias_unimplemented(); test_parse_args_delete_during_alias();
test_parse_args_delete_timing_flags();
test_parse_args_delete_timing_conflict_rejected();
test_parse_args_delete_timing_without_delete_rejected();
test_parse_args_rejects_unimplemented_options(); test_parse_args_rejects_unimplemented_options();
test_parse_args_quiet(); test_parse_args_quiet();
test_parse_args_human_readable(); test_parse_args_human_readable();
+126
View File
@@ -438,6 +438,129 @@ static void test_config_delay_updates_reserved_backup_rejected() {
config_delete(c); config_delete(c);
} }
static void test_config_delete_timing_early_helper() {
Config* cfg = config_create();
EXPECT_NOT_NULL(cfg);
EXPECT_FALSE(config_delete_timing_early(cfg));
EXPECT_TRUE(config_has_valid_delete_timing(cfg));
cfg->use_delete = true;
EXPECT_TRUE(config_has_valid_delete_timing(cfg));
EXPECT_FALSE(config_delete_timing_early(cfg));
config_delete(cfg);
cfg = config_create();
cfg->use_delete = true;
cfg->delete_before = true;
EXPECT_TRUE(config_delete_timing_early(cfg));
EXPECT_TRUE(config_has_valid_delete_timing(cfg));
config_delete(cfg);
cfg = config_create();
cfg->use_delete = true;
cfg->delete_during = true;
EXPECT_TRUE(config_delete_timing_early(cfg));
EXPECT_TRUE(config_has_valid_delete_timing(cfg));
config_delete(cfg);
cfg = config_create();
cfg->use_delete = true;
cfg->delete_delay = true;
EXPECT_FALSE(config_delete_timing_early(cfg));
EXPECT_TRUE(config_has_valid_delete_timing(cfg));
config_delete(cfg);
cfg = config_create();
cfg->use_delete = true;
cfg->delete_after = true;
EXPECT_FALSE(config_delete_timing_early(cfg));
EXPECT_TRUE(config_has_valid_delete_timing(cfg));
config_delete(cfg);
/* Two simultaneous timings are invalid. */
cfg = config_create();
cfg->use_delete = true;
cfg->delete_before = true;
cfg->delete_after = true;
EXPECT_TRUE(config_delete_timing_early(cfg));
EXPECT_FALSE(config_has_valid_delete_timing(cfg));
config_delete(cfg);
/* A timing flag without deletion is invalid. */
cfg = config_create();
cfg->delete_delay = true;
EXPECT_FALSE(config_has_valid_delete_timing(cfg));
EXPECT_FALSE(config_delete_timing_early(cfg));
config_delete(cfg);
}
/* New delete-timing fields must survive config_send/config_receive unchanged,
and a config carrying two conflicting timings must be rejected. */
static void test_config_delete_timing_wire_roundtrip() {
if (is_running_under_valgrind())
return;
struct {
bool before, during, delay, after;
} cases[] = {
{false, false, false, false}, {true, false, false, false}, {false, true, false, false},
{false, false, true, false}, {false, false, false, true},
};
for (size_t i = 0; i < sizeof(cases) / sizeof(cases[0]); i++) {
int p[2];
EXPECT_EQ_INT(socketpair(AF_UNIX, SOCK_STREAM, 0, p), 0);
pid_t pid = fork();
if (pid == 0) {
close(p[1]);
io_set_fds(p[0], p[0]);
Config* recv = config_receive(p[0]);
bool ok = recv != NULL;
if (ok) {
ok = recv->use_delete && recv->delete_before == cases[i].before &&
recv->delete_during == cases[i].during && recv->delete_delay == cases[i].delay &&
recv->delete_after == cases[i].after;
}
config_delete(recv);
close(p[0]);
_exit(ok ? 0 : 1);
} else {
close(p[0]);
io_set_fds(p[1], p[1]);
Config* send_cfg = config_create();
EXPECT_NOT_NULL(send_cfg);
send_cfg->send_directory = str_dup("/src");
send_cfg->receive_root_directory = str_dup("/dst");
send_cfg->use_delete = true;
send_cfg->delete_before = cases[i].before;
send_cfg->delete_during = cases[i].during;
send_cfg->delete_delay = cases[i].delay;
send_cfg->delete_after = cases[i].after;
bool sent = config_send(p[1], send_cfg);
int status;
waitpid(pid, &status, 0);
close(p[1]);
config_delete(send_cfg);
EXPECT_TRUE(sent);
EXPECT_TRUE(WIFEXITED(status) && WEXITSTATUS(status) == 0);
}
}
}
/* The receiver-side wire validation rejects a keep-set config with two
conflicting delete-timing flags. */
static void test_config_delete_timing_conflict_rejected() {
if (is_running_under_valgrind())
return;
Config* c = config_create();
EXPECT_NOT_NULL(c);
c->send_directory = str_dup("/src");
c->receive_root_directory = str_dup("/dst");
c->use_delete = true;
c->delete_before = true;
c->delete_delay = true;
EXPECT_FALSE(roundtrip_config_ok(c));
config_delete(c);
}
static void test_config_is_remote_dest() { static void test_config_is_remote_dest() {
/* Valid SSH-style destinations */ /* Valid SSH-style destinations */
EXPECT_TRUE(config_is_remote_dest("user@host:/path")); EXPECT_TRUE(config_is_remote_dest("user@host:/path"));
@@ -474,6 +597,9 @@ void test_config() {
test_config_string_null_vs_empty_roundtrip(); test_config_string_null_vs_empty_roundtrip();
test_config_temp_dir_roundtrip(); test_config_temp_dir_roundtrip();
test_config_delay_updates_reserved_backup_rejected(); test_config_delay_updates_reserved_backup_rejected();
test_config_delete_timing_wire_roundtrip();
test_config_delete_timing_conflict_rejected();
} }
test_config_delete_timing_early_helper();
test_config_is_remote_dest(); test_config_is_remote_dest();
} }
+20
View File
@@ -412,6 +412,25 @@ static void test_protocol_accounting_release_does_not_underflow() {
protocol_session_unbind(); protocol_session_unbind();
} }
static void test_send_receive_status_timed() {
int p[2];
EXPECT_EQ_INT(pipe(p), 0);
io_set_fds(p[0], p[1]);
io_set_bwlimit(0);
/* The extended-deadline variant must read an ordinary status just like the
default window, and must fail cleanly on EOF rather than block. */
EXPECT_TRUE(send_status(0, STATUS_OK));
Status received = -1;
EXPECT_TRUE(receive_status_timed(0, &received, 5));
EXPECT_EQ_INT((int)received, (int)STATUS_OK);
close(p[1]);
EXPECT_FALSE(receive_status_timed(0, &received, 5));
close(p[0]);
}
void test_protocol() { void test_protocol() {
test_send_receive_n_data(); test_send_receive_n_data();
test_send_receive_n_data_zero(); test_send_receive_n_data_zero();
@@ -421,6 +440,7 @@ void test_protocol() {
test_send_receive_data(); test_send_receive_data();
test_send_receive_int(); test_send_receive_int();
test_send_receive_status(); test_send_receive_status();
test_send_receive_status_timed();
test_receive_n_data_truncated(); test_receive_n_data_truncated();
test_receive_str_truncated(); test_receive_str_truncated();
test_max_alloc_rejects_single_buffer(); test_max_alloc_rejects_single_buffer();
+96 -1
View File
@@ -180,7 +180,10 @@ static void test_receive_manifest_rejects_traversal() {
io_set_fds(p[0], p[1]); io_set_fds(p[0], p[1]);
EXPECT_TRUE(send_int(p[1], 1)); EXPECT_TRUE(send_int(p[1], 1));
EXPECT_TRUE(send_str(p[1], "../outside")); EXPECT_TRUE(send_str(p[1], "../outside"));
EXPECT_EQ_INT(receive_manifest(p[0], cfg, NULL), -1); EXPECT_NULL(receive_manifest_entries(p[0]));
Status status;
EXPECT_TRUE(receive_status(p[1], &status));
EXPECT_EQ_INT(status, STATUS_ERROR);
close(p[0]); close(p[0]);
close(p[1]); close(p[1]);
config_delete(cfg); config_delete(cfg);
@@ -457,6 +460,95 @@ static void test_incremental_check_delta_oversize_reports_failure() {
} }
} }
/* Late-timing keep-set leak guard: a manifest parked by the commit path must
be freed on every error exit, never leaked. These tests drive
receiver_process_pending() through an error AFTER the manifest was parked and
are exercised under ASan/valgrind to prove the list is released. */
static Config* make_late_delete_config(const char* root) {
Config* cfg = config_create();
if (!cfg)
return NULL;
cfg->send_directory = str_dup("/src");
cfg->receive_root_directory = str_dup(root);
cfg->use_delete = true;
cfg->delete_after = true;
return cfg;
}
static int run_pending_receiver(Config* cfg, int fd, ArrayList** pending) {
ReceiverSink sink = {0};
return receiver_process_pending(cfg, fd, &sink, pending);
}
static void test_late_manifest_abort_frees_keepset() {
Config* cfg = make_late_delete_config("/tmp/fastsync_late_abort");
EXPECT_NOT_NULL(cfg);
int p[2];
EXPECT_EQ_INT(socketpair(AF_UNIX, SOCK_STREAM, 0, p), 0);
io_set_fds(p[0], p[1]);
io_set_bwlimit(0);
EXPECT_TRUE(send_status(p[1], STATUS_MANIFEST));
EXPECT_TRUE(send_int(p[1], 1));
EXPECT_TRUE(send_str(p[1], "keep.txt"));
EXPECT_TRUE(send_status(p[1], STATUS_ABORT));
ArrayList* pending = NULL;
EXPECT_EQ_INT(run_pending_receiver(cfg, p[0], &pending), -1);
EXPECT_NULL(pending);
close(p[0]);
close(p[1]);
config_delete(cfg);
}
static void test_late_manifest_eof_frees_keepset() {
Config* cfg = make_late_delete_config("/tmp/fastsync_late_eof");
EXPECT_NOT_NULL(cfg);
int p[2];
EXPECT_EQ_INT(socketpair(AF_UNIX, SOCK_STREAM, 0, p), 0);
io_set_fds(p[0], p[1]);
io_set_bwlimit(0);
EXPECT_TRUE(send_status(p[1], STATUS_MANIFEST));
EXPECT_TRUE(send_int(p[1], 1));
EXPECT_TRUE(send_str(p[1], "keep.txt"));
shutdown(p[1], SHUT_WR);
ArrayList* pending = NULL;
EXPECT_EQ_INT(run_pending_receiver(cfg, p[0], &pending), -1);
EXPECT_NULL(pending);
close(p[0]);
close(p[1]);
config_delete(cfg);
}
static void test_late_second_manifest_frees_both() {
Config* cfg = make_late_delete_config("/tmp/fastsync_late_second");
EXPECT_NOT_NULL(cfg);
int p[2];
EXPECT_EQ_INT(socketpair(AF_UNIX, SOCK_STREAM, 0, p), 0);
io_set_fds(p[0], p[1]);
io_set_bwlimit(0);
EXPECT_TRUE(send_status(p[1], STATUS_MANIFEST));
EXPECT_TRUE(send_int(p[1], 1));
EXPECT_TRUE(send_str(p[1], "first.txt"));
EXPECT_TRUE(send_status(p[1], STATUS_MANIFEST));
EXPECT_TRUE(send_int(p[1], 1));
EXPECT_TRUE(send_str(p[1], "second.txt"));
ArrayList* pending = NULL;
EXPECT_EQ_INT(run_pending_receiver(cfg, p[0], &pending), -1);
EXPECT_NULL(pending);
close(p[0]);
close(p[1]);
config_delete(cfg);
}
void test_server() { void test_server() {
if (!is_running_under_valgrind()) { if (!is_running_under_valgrind()) {
test_receive_files_finished(); test_receive_files_finished();
@@ -467,5 +559,8 @@ void test_server() {
test_incremental_check_quick_skip_by_mtime(); test_incremental_check_quick_skip_by_mtime();
test_incremental_check_size_mismatch_full_transfer(); test_incremental_check_size_mismatch_full_transfer();
test_incremental_check_delta_oversize_reports_failure(); test_incremental_check_delta_oversize_reports_failure();
test_late_manifest_abort_frees_keepset();
test_late_manifest_eof_frees_keepset();
test_late_second_manifest_frees_both();
} }
} }