From 3b85658024e49d353c39341f3a71d04212937f0b Mon Sep 17 00:00:00 2001 From: TapTap Date: Thu, 3 Sep 2026 17:33:52 +0200 Subject: [PATCH] Add remove-source-files option --- RSYNC_COMPAT.md | 2 +- src/client/client_cli.c | 1 + src/client/client_send.c | 63 +++++++++++++++++++++++++++--- src/client/usage.c | 1 + src/shared/config.c | 1 + src/shared/config.h | 1 + src/shared/multiprocessing.c | 3 ++ src/shared/multiprocessing.h | 1 + tests/integration/test_features.py | 48 +++++++++++++++++++++++ tests/test_client_cli.c | 23 +++++++++++ 10 files changed, 138 insertions(+), 6 deletions(-) diff --git a/RSYNC_COMPAT.md b/RSYNC_COMPAT.md index 4b236ec..ac792ce 100644 --- a/RSYNC_COMPAT.md +++ b/RSYNC_COMPAT.md @@ -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 diff --git a/src/client/client_cli.c b/src/client/client_cli.c index 191cc07..71f0f38 100644 --- a/src/client/client_cli.c +++ b/src/client/client_cli.c @@ -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)}, diff --git a/src/client/client_send.c b/src/client/client_send.c index 0b4ca5c..bf19384 100644 --- a/src/client/client_send.c +++ b/src/client/client_send.c @@ -21,6 +21,7 @@ #include #include #include +#include #include #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; } diff --git a/src/client/usage.c b/src/client/usage.c index 2f800af..3f1e469 100644 --- a/src/client/usage.c +++ b/src/client/usage.c @@ -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 SSH port (default: 22)\n"); printf(" --progress Show transfer progress\n"); printf(" --delete Delete files on receiver not in source\n"); diff --git a/src/shared/config.c b/src/shared/config.c index efc2708..01c790d 100644 --- a/src/shared/config.c +++ b/src/shared/config.c @@ -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; diff --git a/src/shared/config.h b/src/shared/config.h index 4a218b0..f2939e2 100644 --- a/src/shared/config.h +++ b/src/shared/config.h @@ -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; diff --git a/src/shared/multiprocessing.c b/src/shared/multiprocessing.c index 6b677f3..47dd5fc 100644 --- a/src/shared/multiprocessing.c +++ b/src/shared/multiprocessing.c @@ -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); diff --git a/src/shared/multiprocessing.h b/src/shared/multiprocessing.h index 37f09d5..6fd2f31 100644 --- a/src/shared/multiprocessing.h +++ b/src/shared/multiprocessing.h @@ -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; diff --git a/tests/integration/test_features.py b/tests/integration/test_features.py index 312add4..ee3690c 100644 --- a/tests/integration/test_features.py +++ b/tests/integration/test_features.py @@ -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) diff --git a/tests/test_client_cli.c b/tests/test_client_cli.c index 656fd40..f1a66f9 100644 --- a/tests/test_client_cli.c +++ b/tests/test_client_cli.c @@ -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();