Merge remote-tracking branch 'origin/feat/remove-source-files' into dev
This commit is contained in:
+1
-1
@@ -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
|
||||
|
||||
|
||||
@@ -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)},
|
||||
|
||||
+126
-7
@@ -16,11 +16,14 @@
|
||||
#include "transport_ssh.h"
|
||||
#include "transport_tls.h"
|
||||
#include "utils.h"
|
||||
#include <fcntl.h>
|
||||
#include <limits.h>
|
||||
#include <stdio.h>
|
||||
#include <stdlib.h>
|
||||
#include <string.h>
|
||||
#include <threads.h>
|
||||
#include <time.h>
|
||||
#include <sys/stat.h>
|
||||
#include <unistd.h>
|
||||
|
||||
#define STREAM_THRESHOLD (64ULL * 1024 * 1024)
|
||||
@@ -104,6 +107,88 @@ static bool finalize_transfer(Client* client) {
|
||||
receive_status(client->file_descriptor, &status) && status == STATUS_OK;
|
||||
}
|
||||
|
||||
typedef struct {
|
||||
char* path;
|
||||
dev_t device;
|
||||
ino_t inode;
|
||||
} SourceFile;
|
||||
|
||||
static void source_file_destroy(void* item) {
|
||||
SourceFile* source = item;
|
||||
if (source) {
|
||||
free(source->path);
|
||||
free(source);
|
||||
}
|
||||
}
|
||||
|
||||
/* Remove only the same regular source file that was sent. */
|
||||
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++) {
|
||||
SourceFile* source = paths->items[i];
|
||||
const char* slash = strrchr(source->path, '/');
|
||||
const char* leaf = slash ? slash + 1 : source->path;
|
||||
char parent[PATH_MAX];
|
||||
if (slash) {
|
||||
size_t parent_length = (size_t)(slash - source->path);
|
||||
if (parent_length == 0)
|
||||
parent_length = 1;
|
||||
if (parent_length >= sizeof(parent))
|
||||
continue;
|
||||
memcpy(parent, source->path, parent_length);
|
||||
parent[parent_length] = '\0';
|
||||
} else {
|
||||
(void)snprintf(parent, sizeof(parent), ".");
|
||||
}
|
||||
|
||||
int dirfd = open(parent, O_RDONLY | O_DIRECTORY | O_CLOEXEC);
|
||||
if (dirfd < 0)
|
||||
continue;
|
||||
struct stat st;
|
||||
if (fstatat(dirfd, leaf, &st, AT_SYMLINK_NOFOLLOW) != 0 || !S_ISREG(st.st_mode) ||
|
||||
st.st_dev != source->device || st.st_ino != source->inode) {
|
||||
close(dirfd);
|
||||
continue;
|
||||
}
|
||||
if (unlinkat(dirfd, leaf, 0) != 0)
|
||||
log_message(LOG_LEVEL_WARNING, "Could not remove source file %s", source->path);
|
||||
close(dirfd);
|
||||
}
|
||||
}
|
||||
|
||||
static SourceFile* source_file_create(const File* file) {
|
||||
if (!file || !file->path)
|
||||
return NULL;
|
||||
struct stat st;
|
||||
if (lstat(file->path, &st) != 0 || !S_ISREG(st.st_mode))
|
||||
return NULL;
|
||||
SourceFile* source = malloc(sizeof(*source));
|
||||
if (!source)
|
||||
return NULL;
|
||||
source->path = str_dup(file->path);
|
||||
source->device = st.st_dev;
|
||||
source->inode = st.st_ino;
|
||||
if (!source->path) {
|
||||
source_file_destroy(source);
|
||||
return NULL;
|
||||
}
|
||||
return source;
|
||||
}
|
||||
|
||||
static bool remember_source_file(ArrayList* paths, const File* file) {
|
||||
if (!paths || !file || !file->path)
|
||||
return true;
|
||||
SourceFile* source = source_file_create(file);
|
||||
if (!source)
|
||||
return true;
|
||||
if (!array_list_add(paths, source)) {
|
||||
source_file_destroy(source);
|
||||
return false;
|
||||
}
|
||||
return true;
|
||||
}
|
||||
|
||||
static void mark_sender_done(PipelineContextSender* context) {
|
||||
mtx_lock(&context->mutex_progress);
|
||||
context->sender_done = true;
|
||||
@@ -343,8 +428,15 @@ 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 (remove_sources) {
|
||||
for (int i = 0; i < chunk->element_count; i++) {
|
||||
if (!remember_source_file(remove_sources, chunk->items[i]))
|
||||
return -1;
|
||||
}
|
||||
}
|
||||
if (!send_status(client->file_descriptor, STATUS_CHUNK))
|
||||
return -1;
|
||||
Data* data;
|
||||
@@ -370,15 +462,28 @@ int send_chunk(Client* client, Chunk* chunk, Config* config) {
|
||||
bool stream = f->data->data == NULL && f->data->size > 0;
|
||||
bool use_sendfile =
|
||||
(config->use_sendfile && !config->use_compression) || (stream && !config->use_compression);
|
||||
SourceFile* source = remove_sources ? source_file_create(f) : NULL;
|
||||
int rc = send_single_file(client, f, config, config->use_incremental, use_sendfile);
|
||||
if (rc == 1)
|
||||
if (rc == 1) {
|
||||
source_file_destroy(source);
|
||||
continue;
|
||||
if (rc < 0)
|
||||
}
|
||||
if (rc < 0) {
|
||||
source_file_destroy(source);
|
||||
return -1;
|
||||
}
|
||||
if (source && !array_list_add(remove_sources, source)) {
|
||||
source_file_destroy(source);
|
||||
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 +524,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 +538,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 +711,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(source_file_destroy);
|
||||
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 +754,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 +784,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 +801,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 +847,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(source_file_destroy);
|
||||
if ((config->use_delete && !context->manifest) ||
|
||||
(config->remove_source_files && !context->remove_source_files)) {
|
||||
pipeline_context_sender_destroy(context);
|
||||
return 1;
|
||||
}
|
||||
|
||||
@@ -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");
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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);
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -36,6 +36,85 @@ 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_single_threaded_removes_transferred_file(self, shared_server):
|
||||
source = os.path.join(TEST_DATA_DIR, "remove_single_source")
|
||||
dest = os.path.join(TEST_DATA_DIR, "remove_single_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"single threaded")
|
||||
|
||||
result, _ = run_client(source, dest, flags=["--remove-source-files"],
|
||||
port=shared_server.port)
|
||||
assert result.returncode == 0
|
||||
assert not os.path.exists(source_file)
|
||||
|
||||
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)
|
||||
|
||||
def test_incremental_skip_preserves_source_file(self, shared_server):
|
||||
source = os.path.join(TEST_DATA_DIR, "remove_skipped_source")
|
||||
dest = os.path.join(TEST_DATA_DIR, "remove_skipped_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 skip")
|
||||
|
||||
result, _ = run_client(source, dest, port=shared_server.port)
|
||||
assert result.returncode == 0
|
||||
result, _ = run_client(source, dest,
|
||||
flags=["--remove-source-files", "--incremental"],
|
||||
port=shared_server.port)
|
||||
assert result.returncode == 0
|
||||
assert os.path.isfile(source_file)
|
||||
|
||||
|
||||
class TestArchiveMode:
|
||||
def test_archive_mode(self, shared_server):
|
||||
clean_dir(DEST_DIR)
|
||||
|
||||
@@ -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();
|
||||
@@ -358,6 +379,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();
|
||||
|
||||
Reference in New Issue
Block a user