diff --git a/README.md b/README.md index fdbd9f3..6962093 100644 --- a/README.md +++ b/README.md @@ -465,7 +465,7 @@ defaults to the current directory. | ## Protocol and Security -FastSync protocol version `2.2.0` is shared by the client and server. The +FastSync protocol version `2.3.0` is shared by the client and server. The current protocol is sender-driven and includes configuration negotiation, incremental checks, checksums, manifests, keep-alives, abort handling, and FastSync-native delta messages. Client and server versions must currently diff --git a/RSYNC_COMPAT.md b/RSYNC_COMPAT.md index 4b236ec..134e3c4 100644 --- a/RSYNC_COMPAT.md +++ b/RSYNC_COMPAT.md @@ -59,7 +59,7 @@ This document maps rsync's full feature set to FastSync's current implementation | `--min-size=SIZE` | Skip files smaller than SIZE | ✅ Implemented | `min_size` in scanner | | `-I`, `--ignore-times` | Don't skip files matching size+time | ❌ Not Implemented | | | `--size-only` | Skip based on size only | ❌ Not Implemented | | -| `-@`, `--modify-window=NUM` | Mod-time comparison accuracy | ❌ Not Implemented | | +| `-@`, `--modify-window=NUM` | Mod-time comparison accuracy | ✅ Implemented | Whole-second tolerance with nanosecond-aware comparisons | | `--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 | | diff --git a/src/client/client_cli.c b/src/client/client_cli.c index 191cc07..548760c 100644 --- a/src/client/client_cli.c +++ b/src/client/client_cli.c @@ -139,6 +139,7 @@ static const OptionEntry OPTION_TABLE[] = { {"--dry-run", "-n", OPT_FLAG, offsetof(Config, dry_run)}, {"--delete", NULL, OPT_FLAG, offsetof(Config, use_delete)}, {"--incremental", NULL, OPT_FLAG, offsetof(Config, use_incremental)}, + {"--modify-window", "-@", OPT_NONNEG_INT, offsetof(Config, modify_window)}, {"--delta", NULL, OPT_FLAG, offsetof(Config, use_delta)}, {"--save-to-disk", NULL, OPT_FLAG, offsetof(Config, save_to_disk)}, {"--progress", NULL, OPT_FLAG, offsetof(Config, show_progress)}, @@ -211,6 +212,18 @@ static int apply_table_option(Config* config, const OptionEntry* entry, const ch int parse_args(Config* config, int argc, char* argv[], int* positional_args, int* positional_count) { for (int i = 1; i < argc; i++) { + const char* modify_window_prefix = "--modify-window="; + if (strncmp(argv[i], modify_window_prefix, strlen(modify_window_prefix)) == 0) { + if (set_nonneg_int_option(&config->modify_window, argv[i] + strlen(modify_window_prefix), + "--modify-window") != 0) + return -1; + continue; + } + if (strncmp(argv[i], "-@", 2) == 0 && argv[i][2] != '\0') { + if (set_nonneg_int_option(&config->modify_window, argv[i] + 2, "-@") != 0) + return -1; + continue; + } const OptionEntry* entry = find_table_option(argv[i]); if (entry) { if (entry->kind != OPT_FLAG) { diff --git a/src/client/client_send.c b/src/client/client_send.c index 0b4ca5c..e8d6d36 100644 --- a/src/client/client_send.c +++ b/src/client/client_send.c @@ -172,10 +172,13 @@ static int incremental_check(Client* client, File* file, const Config* config, return -1; unsigned long long fsize = file->data->size; long long mtime = file->metadata ? file->metadata->mtime_sec : 0; + long long mtime_nsec = file->metadata ? file->metadata->mtime_nsec : 0; if (!send_n_data(client->file_descriptor, &fsize, sizeof(fsize))) return -1; if (!send_n_data(client->file_descriptor, &mtime, sizeof(mtime))) return -1; + if (!send_n_data(client->file_descriptor, &mtime_nsec, sizeof(mtime_nsec))) + return -1; if (config->checksum) { uint64_t checksum; if (!file_checksum(file, &checksum) || diff --git a/src/client/usage.c b/src/client/usage.c index 2f800af..e9bf830 100644 --- a/src/client/usage.c +++ b/src/client/usage.c @@ -28,6 +28,7 @@ void print_usage(void) { printf(" --max-size Skip files larger than n bytes\n"); printf(" --min-size Skip files smaller than n bytes\n"); printf(" --incremental Skip files unchanged since last transfer\n"); + printf(" -@, --modify-window Modification time tolerance\n"); printf(" --delta Delta transfer for changed files (requires --incremental)\n"); printf(" --delta-block Delta block size in bytes (default: %d)\n", DELTA_BLOCK_SIZE_DEFAULT); diff --git a/src/server/receiver.c b/src/server/receiver.c index 2235cf9..c811b16 100644 --- a/src/server/receiver.c +++ b/src/server/receiver.c @@ -2,6 +2,7 @@ #include "chunk.h" #include "log.h" +#include "metadata.h" #include "protocol.h" #include "utils.h" #include @@ -37,9 +38,13 @@ static bool receiver_process_batch(Config* config, int file_descriptor) { return false; unsigned long long check_size; long long check_mtime; + long long check_mtime_nsec; if (!receive_n_data(file_descriptor, &check_size, sizeof(check_size)) || - !receive_n_data(file_descriptor, &check_mtime, sizeof(check_mtime))) { + !receive_n_data(file_descriptor, &check_mtime, sizeof(check_mtime)) || + !receive_n_data(file_descriptor, &check_mtime_nsec, sizeof(check_mtime_nsec)) || + check_mtime_nsec < 0 || check_mtime_nsec >= 1000000000LL) { free(check_path); + send_status(file_descriptor, STATUS_ERROR); return false; } if (!utils_valid_batch_path(check_path)) { @@ -60,8 +65,13 @@ static bool receiver_process_batch(Config* config, int file_descriptor) { } struct stat st; bool has_old = file_stat_secure(full_path, &st); + long long old_mtime_nsec = 0; +#ifdef __linux__ + old_mtime_nsec = st.st_mtim.tv_nsec; +#endif bool match = has_old && (unsigned long long)st.st_size == check_size && - (long long)st.st_mtime == check_mtime; + metadata_mtime_matches(st.st_mtime, old_mtime_nsec, (time_t)check_mtime, + (long)check_mtime_nsec, config->modify_window); bool sent = send_status(file_descriptor, match ? STATUS_OK : STATUS_NEXT); free(full_path); free(check_path); diff --git a/src/server/server.c b/src/server/server.c index af5cd43..aa1879f 100644 --- a/src/server/server.c +++ b/src/server/server.c @@ -2,6 +2,7 @@ #include "chunk.h" #include "file.h" #include "log.h" +#include "metadata.h" #include "multiprocessing.h" #include "queue.h" #include "receiver.h" @@ -161,6 +162,13 @@ int receive_files(Config* config, int fd) { free(check_path); return -1; } + long long check_mtime_nsec; + if (!receive_n_data(fd, &check_mtime_nsec, sizeof(check_mtime_nsec)) || + check_mtime_nsec < 0 || check_mtime_nsec >= 1000000000LL) { + free(check_path); + send_status(fd, STATUS_ERROR); + return -1; + } if (!utils_valid_batch_path(check_path)) { free(check_path); send_status(fd, STATUS_ERROR); @@ -174,8 +182,13 @@ int receive_files(Config* config, int fd) { return -1; } bool has_old = full_path && file_stat_secure(full_path, &st); + long long old_mtime_nsec = 0; +#ifdef __linux__ + old_mtime_nsec = st.st_mtim.tv_nsec; +#endif bool match = has_old && (unsigned long long)st.st_size == check_size && - (long long)st.st_mtime == check_mtime; + metadata_mtime_matches(st.st_mtime, old_mtime_nsec, (time_t)check_mtime, + (long)check_mtime_nsec, config->modify_window); bool sent = send_status(fd, match ? STATUS_OK : STATUS_NEXT); free(full_path); free(check_path); diff --git a/src/shared/config.c b/src/shared/config.c index efc2708..dfaf390 100644 --- a/src/shared/config.c +++ b/src/shared/config.c @@ -35,6 +35,7 @@ static void config_set_defaults(Config* config) { config->min_size = 0; config->use_incremental = false; config->use_delta = false; + config->modify_window = 0; config->delta_block_size = DELTA_BLOCK_SIZE_DEFAULT; config->delta_max_file_size = DELTA_MAX_FILE_SIZE; config->use_tls = false; @@ -134,7 +135,8 @@ static bool validate_received_config(const Config* config) { config->chunk_size > 0 && config->chunk_size <= MAX_CHUNK_SIZE && config->delta_block_size >= DELTA_BLOCK_SIZE_MIN && config->delta_block_size <= DELTA_BLOCK_SIZE_MAX && - config->delta_max_file_size <= DELTA_MAX_FILE_SIZE && config->max_delete >= 0; + config->delta_max_file_size <= DELTA_MAX_FILE_SIZE && config->modify_window >= 0 && + config->max_delete >= 0; } Config* config_create(void) { @@ -253,7 +255,8 @@ static bool send_resume_options(int fd, const Config* c) { return send_str(fd, c->temp_dir ? c->temp_dir : "") && send_int(fd, c->partial) && send_str(fd, c->partial_dir ? c->partial_dir : "") && send_str(fd, c->suffix ? c->suffix : "") && send_int(fd, c->delete_before) && - send_int(fd, c->checksum) && send_str(fd, c->compress_choice ? c->compress_choice : ""); + send_int(fd, c->checksum) && send_int(fd, c->modify_window) && + send_str(fd, c->compress_choice ? c->compress_choice : ""); } static bool receive_core_fields(int fd, Config* c) { @@ -329,6 +332,8 @@ static bool receive_resume_options(int fd, Config* c) { return false; if (!receive_wire_bool(fd, &c->checksum)) return false; + if (!receive_n_data(fd, &c->modify_window, sizeof(c->modify_window))) + return false; c->compress_choice = receive_str(fd); return c->compress_choice != NULL; } diff --git a/src/shared/config.h b/src/shared/config.h index 4a218b0..3ca61ed 100644 --- a/src/shared/config.h +++ b/src/shared/config.h @@ -35,6 +35,7 @@ typedef struct Config { unsigned long long min_size; bool use_incremental; bool use_delta; + int modify_window; uint32_t delta_block_size; unsigned long long delta_max_file_size; bool use_tls; @@ -128,7 +129,7 @@ typedef struct Config { char* compress_choice; } Config; -#define PROTOCOL_VERSION "2.2.0" +#define PROTOCOL_VERSION "2.3.0" #define DEFAULT_CHUNK_SIZE (10 * 1024 * 1024) Config* config_create(void); diff --git a/src/shared/file_receive.c b/src/shared/file_receive.c index a610d67..1c3800a 100644 --- a/src/shared/file_receive.c +++ b/src/shared/file_receive.c @@ -314,12 +314,19 @@ File* receive_incremental_check(int fd, const Config* config, bool* skipped) { unsigned long long check_size; long long check_mtime; + long long check_mtime_nsec; uint64_t check_checksum = 0; if (!receive_n_data(fd, &check_size, sizeof(check_size)) || !receive_n_data(fd, &check_mtime, sizeof(check_mtime))) { free(check_path); return NULL; } + if (!receive_n_data(fd, &check_mtime_nsec, sizeof(check_mtime_nsec)) || check_mtime_nsec < 0 || + check_mtime_nsec >= 1000000000LL) { + free(check_path); + send_status(fd, STATUS_ERROR); + return NULL; + } if (config->checksum && !receive_n_data(fd, &check_checksum, sizeof(check_checksum))) { free(check_path); return NULL; @@ -384,7 +391,12 @@ File* receive_incremental_check(int fd, const Config* config, bool* skipped) { free(old_data); old_data = NULL; } else if (match) { - match = (long long)st.st_mtime == check_mtime; + long long old_mtime_nsec = 0; +#ifdef __linux__ + old_mtime_nsec = st.st_mtim.tv_nsec; +#endif + match = metadata_mtime_matches(st.st_mtime, old_mtime_nsec, (time_t)check_mtime, + (long)check_mtime_nsec, config->modify_window); } if (match) { diff --git a/src/shared/metadata.c b/src/shared/metadata.c index 0c37309..ed1f9eb 100644 --- a/src/shared/metadata.c +++ b/src/shared/metadata.c @@ -24,6 +24,23 @@ typedef char static_assert_mode_t_fits[(sizeof(mode_t) <= sizeof(int32_t)) ? 1 : typedef char static_assert_uid_t_fits[(sizeof(uid_t) <= sizeof(int32_t)) ? 1 : -1]; typedef char static_assert_gid_t_fits[(sizeof(gid_t) <= sizeof(int32_t)) ? 1 : -1]; +bool metadata_mtime_matches(time_t left_sec, long left_nsec, time_t right_sec, long right_nsec, + int modify_window) { + int64_t left = (int64_t)left_sec; + int64_t right = (int64_t)right_sec; + int64_t seconds; + int64_t nanoseconds; + + if (left > right || (left == right && left_nsec >= right_nsec)) { + seconds = left - right; + nanoseconds = (int64_t)left_nsec - (int64_t)right_nsec; + } else { + seconds = right - left; + nanoseconds = (int64_t)right_nsec - (int64_t)left_nsec; + } + return seconds < modify_window || (seconds == modify_window && nanoseconds == 0); +} + void metadata_to_buf(char** buf, const FileMetadata* m) { int32_t present = (m != NULL) ? 1 : 0; memcpy(*buf, &present, sizeof(present)); diff --git a/src/shared/metadata.h b/src/shared/metadata.h index b8d7ba9..8353724 100644 --- a/src/shared/metadata.h +++ b/src/shared/metadata.h @@ -5,6 +5,7 @@ #include #include #include +#include /* * Wire format (introduced in protocol version 2.0.0): @@ -32,4 +33,8 @@ FileMetadata* metadata_receive(int file_descriptor, int* ok); void file_restore_metadata(const char* path, const FileMetadata* metadata); bool file_restore_metadata_fd(int fd, const FileMetadata* metadata); +/* Compare timestamps using rsync's whole-second modification window. */ +bool metadata_mtime_matches(time_t left_sec, long left_nsec, time_t right_sec, long right_nsec, + int modify_window); + #endif diff --git a/tests/integration/test_features.py b/tests/integration/test_features.py index 312add4..e10beee 100644 --- a/tests/integration/test_features.py +++ b/tests/integration/test_features.py @@ -213,6 +213,27 @@ class TestIncremental: with open(received_file, "rb") as f: assert f.read() == b"hello world\n" + def test_modify_window_allows_subsecond_mtime_difference(self, shared_server): + clean_dir(DEST_DIR) + result, _ = run_client(SOURCE_DIR, DEST_DIR, flags=["-M"], port=shared_server.port) + assert result.returncode == 0 + + received = get_dest_received_dir(DEST_DIR, SOURCE_DIR) + source_file = os.path.join(SOURCE_DIR, "small.txt") + received_file = os.path.join(received, "small.txt") + source_stat = os.stat(source_file) + with open(received_file, "wb") as f: + f.write(b"modified!!!\n") + os.utime(received_file, ns=(source_stat.st_atime_ns, + source_stat.st_mtime_ns - 1500000000)) + + result, _ = run_client(SOURCE_DIR, DEST_DIR, + flags=["-M", "--incremental", "--modify-window=2"], + port=shared_server.port) + assert result.returncode == 0, f"Modify-window sync failed: {result.stderr[:200]}" + with open(received_file, "rb") as f: + assert f.read() == b"modified!!!\n" + class TestDelete: def test_delete_removes_extra_files(self, shared_server): diff --git a/tests/test_client_cli.c b/tests/test_client_cli.c index 656fd40..3337658 100644 --- a/tests/test_client_cli.c +++ b/tests/test_client_cli.c @@ -247,6 +247,43 @@ static void test_parse_args_valid_compression_level() { config_delete(cfg); } +static void test_parse_args_modify_window() { + Config* cfg = config_create(); + char* argv[] = {"fastsync", "--modify-window=3", "/src", "/dst"}; + int positional_args[2]; + int positional_count = 0; + + EXPECT_EQ_INT(parse_args(cfg, 4, argv, positional_args, &positional_count), 0); + EXPECT_EQ_INT(cfg->modify_window, 3); + config_delete(cfg); + + cfg = config_create(); + char* short_argv[] = {"fastsync", "-@", "7", "/src", "/dst"}; + positional_count = 0; + EXPECT_EQ_INT(parse_args(cfg, 5, short_argv, positional_args, &positional_count), 0); + EXPECT_EQ_INT(cfg->modify_window, 7); + config_delete(cfg); + + cfg = config_create(); + char* attached_argv[] = {"fastsync", "-@11", "/src", "/dst"}; + positional_count = 0; + EXPECT_EQ_INT(parse_args(cfg, 4, attached_argv, positional_args, &positional_count), 0); + EXPECT_EQ_INT(cfg->modify_window, 11); + config_delete(cfg); +} + +static void test_parse_args_rejects_invalid_modify_window() { + const char* values[] = {"-1", "not-a-number", ""}; + for (size_t i = 0; i < sizeof(values) / sizeof(values[0]); i++) { + Config* cfg = config_create(); + char* argv[] = {"fastsync", "--modify-window", (char*)values[i], "/src", "/dst"}; + int positional_args[2]; + int positional_count = 0; + EXPECT_EQ_INT(parse_args(cfg, 5, argv, positional_args, &positional_count), -1); + config_delete(cfg); + } +} + /* Test parse_args unknown option returns error */ static void test_parse_args_unknown_option() { Config* cfg = config_create(); @@ -357,6 +394,8 @@ void test_client_cli() { test_parse_args_invalid_server_port(); test_parse_args_invalid_compression_level(); test_parse_args_valid_compression_level(); + test_parse_args_modify_window(); + test_parse_args_rejects_invalid_modify_window(); test_parse_args_unknown_option(); test_parse_args_rejects_unimplemented_options(); test_parse_args_archive(); diff --git a/tests/test_config.c b/tests/test_config.c index 4f124db..bce637a 100644 --- a/tests/test_config.c +++ b/tests/test_config.c @@ -126,6 +126,7 @@ static void test_config_send_receive() { send_cfg->use_metadata = true; send_cfg->compression_level = 5; send_cfg->chunk_size = 1024; + send_cfg->modify_window = 4; /* Use socketpair for bidirectional communication */ int p[2]; @@ -160,6 +161,8 @@ static void test_config_send_receive() { ok = false; if (recv_cfg->chunk_size != 1024) ok = false; + if (recv_cfg->modify_window != 4) + ok = false; } config_delete(recv_cfg); close(p[0]); diff --git a/tests/test_metadata.c b/tests/test_metadata.c index 6db0d31..02fee7f 100644 --- a/tests/test_metadata.c +++ b/tests/test_metadata.c @@ -122,6 +122,14 @@ static void test_metadata_rejects_invalid_values() { close(p[1]); } +static void test_metadata_mtime_window() { + EXPECT_TRUE(metadata_mtime_matches(100, 100000000, 101, 600000000, 2)); + EXPECT_FALSE(metadata_mtime_matches(100, 100000000, 102, 600000000, 2)); + EXPECT_TRUE(metadata_mtime_matches(100, 100000000, 102, 100000000, 2)); + EXPECT_FALSE(metadata_mtime_matches(100, 100000000, 100, 100000001, 0)); + EXPECT_TRUE(metadata_mtime_matches(100, 100000000, 100, 100000000, 0)); +} + static void test_file_restore_metadata() { const char* path = "temp_meta_restore_test.txt"; const char* content = "test content"; @@ -151,5 +159,6 @@ void test_metadata() { test_metadata_send_receive_roundtrip(); test_metadata_send_null(); test_metadata_rejects_invalid_values(); + test_metadata_mtime_window(); test_file_restore_metadata(); }