feat: implement rsync delete timing (--delete-before/--delete-during/--delete-delay/--delete-after)
Deletion timing is now real and selected by the four rsync flags plus the plain --delete default. Wire protocol bumps to 2.8.0: two new config booleans (delete_during, delete_delay) are serialized and validated, joining the existing delete_before/delete_after. - Early modes (--delete-before, --delete-during/--del): the sender pre-scans the whole tree (paths only), transmits the keep-set manifest BEFORE any file data, and the receiver removes extras and acks STATUS_OK; the sender only streams data after the deletion committed. Deletion is thus performed even if a later transfer phase fails (rsync delete-before/during are destructive by definition). FastSync streams in a single scan so it cannot interleave per-directory like rsync delete-during; --delete-during selects the same engine mode as --delete-before (documented divergence). - Late/commit modes (plain --delete, --delete-after, --delete-delay): the manifest closes the data stream and deletion is committed only after STATUS_FINISHED proves the whole transfer succeeded, preserving FastSync's commit-style safety. --delete-delay converges with --delete-after because FastSync never snapshots the destination during data flow (documented). - The STATUS_MANIFEST frame is now self-delimiting and position-independent. Single-threaded receivers delete before the success frame; the -m receiver hands the keep-set to server.c, which commits the deletion only after the disk writer thread has drained (fixes a delete-vs-in-flight-temp race). - Every timing flag implies --delete; at most one timing flag is allowed. - Each timing flag implies --delete, matching rsync; conflicts are rejected.
This commit is contained in:
+20
-20
@@ -370,7 +370,6 @@ typedef enum {
|
||||
OPT_POS_INT,
|
||||
OPT_NONNEG_INT,
|
||||
OPT_ULL,
|
||||
OPT_UNSUPPORTED,
|
||||
} OptKind;
|
||||
|
||||
typedef struct {
|
||||
@@ -430,7 +429,10 @@ static const OptionEntry OPTION_TABLE[] = {
|
||||
{"--old-d", NULL, OPT_FLAG, offsetof(Config, dirs)},
|
||||
{"--relative", "-R", OPT_FLAG, offsetof(Config, relative)},
|
||||
{"--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)},
|
||||
{"--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;
|
||||
}
|
||||
|
||||
static int apply_table_option(Config* config, const OptionEntry* entry, const char* option_name,
|
||||
const char* value) {
|
||||
static int apply_table_option(Config* config, const OptionEntry* entry, const char* value) {
|
||||
if (entry->kind == OPT_NOOP)
|
||||
return 0;
|
||||
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;
|
||||
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;
|
||||
}
|
||||
@@ -665,23 +661,20 @@ int parse_args(Config* config, int argc, char* argv[], int* positional_args,
|
||||
if (!entry)
|
||||
entry = find_table_option_with_equals(argv[i], &inline_value);
|
||||
if (entry) {
|
||||
const char* option_name = argv[i];
|
||||
const char* value = NULL;
|
||||
if (entry->kind != OPT_FLAG) {
|
||||
if (entry->kind != OPT_UNSUPPORTED) {
|
||||
value = inline_value;
|
||||
if (!value && i + 1 < argc)
|
||||
value = argv[++i];
|
||||
if (!value) {
|
||||
log_message(LOG_LEVEL_ERROR, "missing argument for %s", entry->name);
|
||||
return -1;
|
||||
}
|
||||
value = inline_value;
|
||||
if (!value && i + 1 < argc)
|
||||
value = argv[++i];
|
||||
if (!value) {
|
||||
log_message(LOG_LEVEL_ERROR, "missing argument for %s", entry->name);
|
||||
return -1;
|
||||
}
|
||||
if (strcmp(entry->name, "--compress-choice") == 0) {
|
||||
if (set_compression_choice(config, value) != 0)
|
||||
return -1;
|
||||
} else {
|
||||
if (apply_table_option(config, entry, option_name, value) != 0)
|
||||
if (apply_table_option(config, entry, value) != 0)
|
||||
return -1;
|
||||
if (strcmp(entry->name, "--compress-level") == 0 &&
|
||||
(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;
|
||||
}
|
||||
}
|
||||
} else if (apply_table_option(config, entry, option_name, NULL) != 0) {
|
||||
} else if (apply_table_option(config, entry, NULL) != 0) {
|
||||
return -1;
|
||||
}
|
||||
if (entry->offset == offsetof(Config, eight_bit_output))
|
||||
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;
|
||||
}
|
||||
|
||||
|
||||
+125
-17
@@ -254,10 +254,6 @@ static void disconnect_transfer_client(Client* 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) {
|
||||
if (!manifest)
|
||||
return true;
|
||||
@@ -597,6 +593,53 @@ static int send_delete_manifest(int fd, ArrayList* manifest) {
|
||||
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). */
|
||||
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(client->file_descriptor, &ack))
|
||||
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,
|
||||
DeltaSignature** out_sig) {
|
||||
*out_sig = NULL;
|
||||
@@ -906,6 +949,17 @@ static int send_chunks_multithreaded(void* pipeline_context) {
|
||||
protocol_session_unbind();
|
||||
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) {
|
||||
Chunk* current_chunk = queue_dequeue_multithreaded(
|
||||
@@ -919,7 +973,7 @@ static int send_chunks_multithreaded(void* pipeline_context) {
|
||||
protocol_session_unbind();
|
||||
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)
|
||||
goto send_fail;
|
||||
}
|
||||
@@ -1013,7 +1067,7 @@ static int scan_directory_multithreaded(void* pipeline_context) {
|
||||
failed = dirs_mode ? directory_scanner_failed(dscanner) : parallel_scanner_failed(scanner);
|
||||
break;
|
||||
}
|
||||
if (context->config->use_delete) {
|
||||
if (context->config->use_delete && !context->early_delete) {
|
||||
mtx_lock(&context->mutex_scanner);
|
||||
bool manifest_ok = add_chunk_to_manifest(context->manifest, current_chunk);
|
||||
mtx_unlock(&context->mutex_scanner);
|
||||
@@ -1172,19 +1226,46 @@ int send_files(Config* config) {
|
||||
DirectoryScanner* scanner = NULL;
|
||||
ArrayList* manifest = NULL;
|
||||
ArrayList* remove_sources = NULL;
|
||||
bool delete_early = config->use_delete && config_delete_timing_early(config);
|
||||
bool send_failed = false;
|
||||
PreparedScanner prepared;
|
||||
memset(&prepared, 0, sizeof(prepared));
|
||||
if (!config_send(client->file_descriptor, config))
|
||||
goto send_fail;
|
||||
if (!prepare_scanner(config, 0, &prepared))
|
||||
goto send_fail;
|
||||
scanner = directory_scanner_create_with_options(config->send_directory, &prepared.options);
|
||||
manifest = create_transfer_manifest(config);
|
||||
if (config->remove_source_files)
|
||||
remove_sources = array_list_create(source_file_destroy);
|
||||
if (!scanner || (config->use_delete && !manifest) ||
|
||||
(config->remove_source_files && !remove_sources))
|
||||
if (config->remove_source_files && !remove_sources)
|
||||
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;
|
||||
unsigned long long total_bytes = 0;
|
||||
int total_files = 0;
|
||||
@@ -1196,7 +1277,7 @@ int send_files(Config* config) {
|
||||
chunk_bytes += current_chunk->items[i]->data->size;
|
||||
total_files++;
|
||||
}
|
||||
if (!add_chunk_to_manifest(manifest, current_chunk)) {
|
||||
if (manifest && !add_chunk_to_manifest(manifest, current_chunk)) {
|
||||
chunk_destroy(current_chunk);
|
||||
goto send_fail;
|
||||
}
|
||||
@@ -1220,9 +1301,7 @@ int send_files(Config* config) {
|
||||
if (send_chunk_with_removal(client, current_chunk, config, remove_sources) != 0) {
|
||||
log_message(LOG_LEVEL_ERROR, "Failed to send chunk");
|
||||
chunk_destroy(current_chunk);
|
||||
if (manifest)
|
||||
array_list_delete(manifest);
|
||||
manifest = NULL;
|
||||
send_failed = true;
|
||||
break;
|
||||
}
|
||||
total_bytes += chunk_bytes;
|
||||
@@ -1235,9 +1314,18 @@ int send_files(Config* config) {
|
||||
}
|
||||
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;
|
||||
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) {
|
||||
array_list_delete(manifest);
|
||||
manifest = NULL;
|
||||
@@ -1324,8 +1412,28 @@ int send_files_multithreaded(Config** config_ptr) {
|
||||
return 1;
|
||||
}
|
||||
*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);
|
||||
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)
|
||||
context->remove_source_files = array_list_create(source_file_destroy);
|
||||
if ((config->use_delete && !context->manifest) ||
|
||||
|
||||
@@ -71,5 +71,11 @@ bool validate_config(const Config* config) {
|
||||
"staging directory)");
|
||||
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;
|
||||
}
|
||||
|
||||
+10
-1
@@ -24,6 +24,16 @@ void print_usage(void) {
|
||||
printf(" -P Partial mode with progress (retention incomplete)\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(" (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 is known, before\n");
|
||||
printf(" --del data is applied (alias --del; implies --delete)\n");
|
||||
printf(" --delete-delay Delete extras only after a successful transfer\n");
|
||||
printf(" (implies --delete)\n");
|
||||
printf(" --delete-after Alias of the default --delete timing: delete only\n");
|
||||
printf(" after the transfer succeeded (implies --delete)\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(" --dirs, -d, --old-dirs, --old-d Transfer the named directory entries without\n");
|
||||
@@ -37,7 +47,6 @@ void print_usage(void) {
|
||||
printf(" parent directory is not itself listed\n");
|
||||
printf(" --mkpath Create the destination root directory on the server when it\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(" --include <pattern> Only include files matching pattern\n");
|
||||
printf(" --exclude-from <file> Read exclude patterns from file\n");
|
||||
|
||||
Reference in New Issue
Block a user