Add remove-source-files option
CI / lint (pull_request) Successful in 11s
CI / sanitizers (undefined) (pull_request) Successful in 37s
CI / sanitizers (address) (pull_request) Successful in 38s
CI / fuzz-build (pull_request) Successful in 15s
CI / coverage (pull_request) Successful in 32s
CI / build-and-test (pull_request) Successful in 1m15s
CI / valgrind (pull_request) Successful in 33s

This commit is contained in:
2026-09-03 17:33:52 +02:00
parent 190fc5d300
commit 3b85658024
10 changed files with 138 additions and 6 deletions
+1 -1
View File
@@ -62,7 +62,7 @@ This document maps rsync's full feature set to FastSync's current implementation
| `-@`, `--modify-window=NUM` | Mod-time comparison accuracy | ❌ Not Implemented | |
| `--existing` | Skip creating new files on receiver | ❌ Not Implemented | |
| `--ignore-existing` | Skip updating existing files | ❌ Not Implemented | |
| `--remove-source-files` | Sender removes synced files | ❌ Not Implemented | |
| `--remove-source-files` | Sender removes regular files after confirmed transfer | ✅ Implemented | |
## 4. Directory Options
+1
View File
@@ -137,6 +137,7 @@ typedef struct {
/* Options that map directly onto a Config field with no side effects. */
static const OptionEntry OPTION_TABLE[] = {
{"--dry-run", "-n", OPT_FLAG, offsetof(Config, dry_run)},
{"--remove-source-files", NULL, OPT_FLAG, offsetof(Config, remove_source_files)},
{"--delete", NULL, OPT_FLAG, offsetof(Config, use_delete)},
{"--incremental", NULL, OPT_FLAG, offsetof(Config, use_incremental)},
{"--delta", NULL, OPT_FLAG, offsetof(Config, use_delta)},
+58 -5
View File
@@ -21,6 +21,7 @@
#include <string.h>
#include <threads.h>
#include <time.h>
#include <sys/stat.h>
#include <unistd.h>
#define STREAM_THRESHOLD (64ULL * 1024 * 1024)
@@ -104,6 +105,31 @@ static bool finalize_transfer(Client* client) {
receive_status(client->file_descriptor, &status) && status == STATUS_OK;
}
/* Remove only regular source files after the receiver confirms the whole transfer. */
static void remove_transferred_sources(const Config* config, ArrayList* paths) {
if (!config->remove_source_files || !paths)
return;
for (int i = 0; i < paths->size; i++) {
const char* path = paths->items[i];
struct stat st;
if (lstat(path, &st) != 0 || !S_ISREG(st.st_mode))
continue;
if (unlink(path) != 0)
log_message(LOG_LEVEL_WARNING, "Could not remove source file %s", path);
}
}
static bool remember_source_file(ArrayList* paths, const File* file) {
if (!paths || !file || !file->path)
return true;
char* path = str_dup(file->path);
if (!path || !array_list_add(paths, path)) {
free(path);
return false;
}
return true;
}
static void mark_sender_done(PipelineContextSender* context) {
mtx_lock(&context->mutex_progress);
context->sender_done = true;
@@ -343,7 +369,8 @@ static int send_single_file(Client* client, File* file, Config* config, bool use
return 0;
}
int send_chunk(Client* client, Chunk* chunk, Config* config) {
static int send_chunk_with_removal(Client* client, Chunk* chunk, Config* config,
ArrayList* remove_sources) {
if (config->use_chunk_serialization) {
if (!send_status(client->file_descriptor, STATUS_CHUNK))
return -1;
@@ -360,6 +387,12 @@ int send_chunk(Client* client, Chunk* chunk, Config* config) {
return -1;
}
data_destroy(data);
if (remove_sources) {
for (int i = 0; i < chunk->element_count; i++) {
if (!remember_source_file(remove_sources, chunk->items[i]))
return -1;
}
}
return 0;
}
@@ -375,10 +408,16 @@ int send_chunk(Client* client, Chunk* chunk, Config* config) {
continue;
if (rc < 0)
return -1;
if (rc == 0 && !remember_source_file(remove_sources, f))
return -1;
}
return 0;
}
int send_chunk(Client* client, Chunk* chunk, Config* config) {
return send_chunk_with_removal(client, chunk, config, NULL);
}
static int send_chunks_multithreaded(void* pipeline_context) {
PipelineContextSender* context = (PipelineContextSender*)pipeline_context;
Client* client = connect_transfer_client(context->config);
@@ -419,6 +458,8 @@ static int send_chunks_multithreaded(void* pipeline_context) {
goto send_fail;
}
bool ok = finalize_transfer(client);
if (ok)
remove_transferred_sources(context->config, context->remove_source_files);
disconnect_transfer_client(client);
mark_sender_done(context);
protocol_session_unbind();
@@ -431,7 +472,8 @@ static int send_chunks_multithreaded(void* pipeline_context) {
protocol_session_unbind();
return thrd_error;
}
if (send_chunk(client, current_chunk, context->config) != 0) {
if (send_chunk_with_removal(client, current_chunk, context->config,
context->remove_source_files) != 0) {
log_message(LOG_LEVEL_ERROR, "unexpected error while sending chunk");
chunk_destroy(current_chunk);
pipeline_cancel(context);
@@ -603,12 +645,16 @@ int send_files(Config* config) {
int ret = 1;
DirectoryScanner* scanner = NULL;
ArrayList* manifest = NULL;
ArrayList* remove_sources = NULL;
if (!config_send(client->file_descriptor, config))
goto send_fail;
ScannerOptions scanner_options = scanner_options_from_config(config, 0);
scanner = directory_scanner_create_with_options(config->send_directory, &scanner_options);
manifest = create_transfer_manifest(config);
if (!scanner || (config->use_delete && !manifest))
if (config->remove_source_files)
remove_sources = array_list_create(free);
if (!scanner || (config->use_delete && !manifest) ||
(config->remove_source_files && !remove_sources))
goto send_fail;
Chunk* current_chunk;
unsigned long long total_bytes = 0;
@@ -642,7 +688,7 @@ int send_files(Config* config) {
goto send_fail;
}
}
if (send_chunk(client, current_chunk, config) != 0) {
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)
@@ -672,6 +718,8 @@ int send_files(Config* config) {
manifest = NULL;
}
bool ok = finalize_transfer(client);
if (ok)
remove_transferred_sources(config, remove_sources);
if (config->show_progress)
print_transfer_progress(total_bytes, start, "Done.\n");
if (config->stats) {
@@ -687,6 +735,8 @@ send_fail:
here even on success without --delete, fixing a pre-existing leak. */
if (manifest)
array_list_delete(manifest);
if (remove_sources)
array_list_delete(remove_sources);
if (scanner)
directory_scanner_destroy(scanner);
disconnect_transfer_client(client);
@@ -731,7 +781,10 @@ int send_files_multithreaded(Config** config_ptr) {
*config_ptr = NULL; /* context now owns config through all remaining paths */
if (config->use_delete)
context->manifest = array_list_create(free);
if (config->use_delete && !context->manifest) {
if (config->remove_source_files)
context->remove_source_files = array_list_create(free);
if ((config->use_delete && !context->manifest) ||
(config->remove_source_files && !context->remove_source_files)) {
pipeline_context_sender_destroy(context);
return 1;
}
+1
View File
@@ -18,6 +18,7 @@ void print_usage(void) {
printf(" -z [level] Alias for -c\n");
printf(" -a, --archive Archive mode (-c -m -M)\n");
printf(" -n, --dry-run Show what would be transferred\n");
printf(" --remove-source-files Remove regular source files after successful transfer\n");
printf(" -p <port> SSH port (default: 22)\n");
printf(" --progress Show transfer progress\n");
printf(" --delete Delete files on receiver not in source\n");
+1
View File
@@ -19,6 +19,7 @@ static void config_set_defaults(Config* config) {
config->use_metadata = false;
config->show_progress = false;
config->dry_run = false;
config->remove_source_files = false;
config->use_delete = false;
config->compression_level = 5;
config->use_sendfile = false;
+1
View File
@@ -20,6 +20,7 @@ typedef struct Config {
bool use_metadata;
bool show_progress;
bool dry_run;
bool remove_source_files;
bool use_delete;
int compression_level;
unsigned long long chunk_size;
+3
View File
@@ -26,6 +26,7 @@ PipelineContextSender* pipeline_context_sender_create(Config* config, Queue* que
context->scanner_done = false;
context->loader_done = false;
context->manifest = NULL;
context->remove_source_files = NULL;
context->progress_bytes = 0;
context->sender_done = false;
atomic_init(&context->cancelled, false);
@@ -76,6 +77,8 @@ void pipeline_context_sender_destroy(PipelineContextSender* context) {
if (context->manifest) {
array_list_delete(context->manifest);
}
if (context->remove_source_files)
array_list_delete(context->remove_source_files);
config_delete(context->config);
queue_destroy(context->queue_scanner);
queue_destroy(context->queue_loader);
+1
View File
@@ -24,6 +24,7 @@ typedef struct {
cnd_t condition_not_empty_loader;
bool loader_done;
ArrayList* manifest;
ArrayList* remove_source_files;
mtx_t mutex_progress;
unsigned long long progress_bytes;
bool sender_done;
+48
View File
@@ -36,6 +36,54 @@ class TestDryRun:
assert "Dry run:" in result.stdout, f"No dry run output: {result.stdout[:200]}"
class TestRemoveSourceFiles:
def test_removes_only_transferred_regular_files(self, shared_server):
source = os.path.join(TEST_DATA_DIR, "remove_source")
dest = os.path.join(TEST_DATA_DIR, "remove_dest")
clean_dir(source)
clean_dir(dest)
with open(os.path.join(source, "one.txt"), "wb") as f:
f.write(b"one")
with open(os.path.join(source, "two.txt"), "wb") as f:
f.write(b"two")
os.makedirs(os.path.join(source, "directory"))
os.symlink("one.txt", os.path.join(source, "link.txt"))
result, _ = run_client(source, dest, flags=["--remove-source-files", "-m"],
port=shared_server.port)
assert result.returncode == 0, f"Remove-source sync failed: {result.stderr[:200]}"
assert not os.path.exists(os.path.join(source, "one.txt"))
assert not os.path.exists(os.path.join(source, "two.txt"))
assert os.path.isdir(os.path.join(source, "directory"))
assert os.path.islink(os.path.join(source, "link.txt"))
def test_dry_run_preserves_source_files(self):
source = os.path.join(TEST_DATA_DIR, "remove_dry_source")
dest = os.path.join(TEST_DATA_DIR, "remove_dry_dest")
clean_dir(source)
clean_dir(dest)
source_file = os.path.join(source, "file.txt")
with open(source_file, "wb") as f:
f.write(b"keep")
result, _ = run_client(source, dest, flags=["--remove-source-files", "--dry-run"])
assert result.returncode == 0
assert os.path.isfile(source_file)
def test_failed_connection_preserves_source_files(self):
source = os.path.join(TEST_DATA_DIR, "remove_failed_source")
dest = os.path.join(TEST_DATA_DIR, "remove_failed_dest")
clean_dir(source)
clean_dir(dest)
source_file = os.path.join(source, "file.txt")
with open(source_file, "wb") as f:
f.write(b"keep after failure")
result, _ = run_client(source, dest, flags=["--remove-source-files"], port=1)
assert result.returncode != 0
assert os.path.isfile(source_file)
class TestArchiveMode:
def test_archive_mode(self, shared_server):
clean_dir(DEST_DIR)
+23
View File
@@ -107,6 +107,27 @@ static void test_cli_dry_run() {
config_delete(cfg);
}
static void test_cli_remove_source_files() {
Config* cfg = config_create();
EXPECT_NOT_NULL(cfg);
EXPECT_FALSE(cfg->remove_source_files);
cfg->remove_source_files = true;
EXPECT_TRUE(cfg->remove_source_files);
config_delete(cfg);
}
static void test_parse_args_remove_source_files() {
Config* cfg = config_create();
char* argv[] = {"fastsync", "--remove-source-files", "/src", "/dst"};
int positional_args[2];
int positional_count = 0;
int ret = parse_args(cfg, 4, argv, positional_args, &positional_count);
EXPECT_EQ_INT(ret, 0);
EXPECT_TRUE(cfg->remove_source_files);
config_delete(cfg);
}
/* Test that --delete sets use_delete */
static void test_cli_delete_flag() {
Config* cfg = config_create();
@@ -347,6 +368,8 @@ void test_client_cli() {
test_cli_help();
test_cli_archive_flags();
test_cli_dry_run();
test_cli_remove_source_files();
test_parse_args_remove_source_files();
test_cli_delete_flag();
test_cli_exclude_patterns();
test_parse_args_help();