diff --git a/src/client/scanner.c b/src/client/scanner.c index 6d83c8e..b768001 100644 --- a/src/client/scanner.c +++ b/src/client/scanner.c @@ -205,6 +205,17 @@ typedef struct { bool referent_error; } ScannerEntry; +/* One inspected directory entry buffered so the sequential scanner can emit the + stream in rsync's flist order. `name` is the raw dirent name (owned here); + `entry` is the scanner_inspect_entry() result whose path/link_target are owned + when `inspection == 1`; `inspection` is that call's return code (1 keep, + 0 skip, <0 fatal). */ +typedef struct { + char* name; + ScannerEntry entry; + int inspection; +} SortedEntry; + /* --one-file-system (-x) decision. Only directories can carry a different * device than their parent (mount points), so this is checked when a child * directory is about to be descended into. */ @@ -687,6 +698,114 @@ skip: return 0; } +static void sorted_entry_destroy(void* item) { + SortedEntry* se = (SortedEntry*)item; + if (!se) + return; + free(se->name); + free(se->entry.path); + free(se->entry.link_target); +} + +/* rsync flist order within one directory: non-directories first, then + directories, each group by ascending name. strcmp() compares as unsigned + char, matching rsync's f_name_cmp(). */ +static int sorted_entry_cmp(const void* a, const void* b) { + const SortedEntry* x = (const SortedEntry*)a; + const SortedEntry* y = (const SortedEntry*)b; + bool x_dir = x->inspection > 0 && x->entry.is_directory; + bool y_dir = y->inspection > 0 && y->entry.is_directory; + if (x_dir != y_dir) + return x_dir ? 1 : -1; + return strcmp(x->name, y->name); +} + +static void scanner_free_sorted(DirectoryScanner* scanner) { + SortedEntry* entries = (SortedEntry*)scanner->sorted_entries; + for (size_t i = 0; i < scanner->sorted_count; i++) + sorted_entry_destroy(&entries[i]); + free(entries); + scanner->sorted_entries = NULL; + scanner->sorted_count = 0; + scanner->sorted_index = 0; +} + +/* Read every entry of the open directory, inspect it once and store it sorted in + rsync's flist order. Returns 0 on success, -1 on a fatal error (the caller + aborts the scan). */ +static int scanner_buffer_current_directory(DirectoryScanner* scanner) { + size_t capacity = 64; + size_t count = 0; + SortedEntry* entries = malloc(capacity * sizeof(*entries)); + if (!entries) { + scanner->failed = true; + return -1; + } + const struct dirent* dirent; + while ((dirent = readdir(scanner->current_dir)) != NULL) { + if (strcmp(dirent->d_name, ".") == 0 || strcmp(dirent->d_name, "..") == 0) + continue; + if (count == capacity) { + size_t next = capacity * 2; + SortedEntry* grown = realloc(entries, next * sizeof(*entries)); + if (!grown) { + scanner->failed = true; + break; + } + entries = grown; + capacity = next; + } + char* name = str_dup(dirent->d_name); + if (!name) { + scanner->failed = true; + break; + } + char* link_rel = child_rel_path(scanner->current_rel, dirent->d_name); + if (!link_rel) { + free(name); + scanner->failed = true; + break; + } + int inspection = scanner_inspect_entry(&scanner->options, scanner->current_path, link_rel, + dirent->d_name, &entries[count].entry); + free(link_rel); + if (inspection < 0) { + free(name); + scanner->failed = true; + break; + } + entries[count].name = name; + entries[count].inspection = inspection; + count++; + } + if (scanner->failed) { + for (size_t i = 0; i < count; i++) + sorted_entry_destroy(&entries[i]); + free(entries); + return -1; + } + qsort(entries, count, sizeof(*entries), sorted_entry_cmp); + scanner->sorted_entries = entries; + scanner->sorted_count = count; + scanner->sorted_index = 0; + return 0; +} + +/* Push this directory's collected child directories onto the LIFO stack in + reverse so the first (ascending) child is popped first (depth-first). */ +static void scanner_push_pending_dirs(DirectoryScanner* scanner) { + ArrayList* pending = (ArrayList*)scanner->pending_dirs; + if (!pending) + return; + for (int i = pending->size - 1; i >= 0; i--) { + if (!queue_push(scanner->directories, pending->items[i])) { + dir_entry_destroy(pending->items[i]); + scanner->failed = true; + } + } + pending->size = 0; +} + DirectoryScanner* directory_scanner_create_with_options(const char* root_directory, const ScannerOptions* options) { if (!root_directory || !options) @@ -704,12 +823,22 @@ DirectoryScanner* directory_scanner_create_with_options(const char* root_directo free(scanner); return NULL; } + scanner->pending_dirs = array_list_create(NULL); + if (!scanner->pending_dirs) { + queue_destroy(scanner->directories); + free(scanner); + return NULL; + } scanner->current_dir = NULL; scanner->current_path = NULL; scanner->current_depth = 0; scanner->failed = false; + scanner->sorted_entries = NULL; + scanner->sorted_count = 0; + scanner->sorted_index = 0; scanner->root_path = str_dup(root_directory); if (!scanner->root_path) { + array_list_delete(scanner->pending_dirs); queue_destroy(scanner->directories); free(scanner); return NULL; @@ -729,6 +858,7 @@ DirectoryScanner* directory_scanner_create_with_options(const char* root_directo scanner->filter_nodes = array_list_create(filter_node_destroy); if (!scanner->filter_nodes) { free(scanner->root_path); + array_list_delete(scanner->pending_dirs); queue_destroy(scanner->directories); free(scanner); return NULL; @@ -739,6 +869,7 @@ DirectoryScanner* directory_scanner_create_with_options(const char* root_directo if (stat(root_directory, &root_stats) != 0) { log_perror("Could not stat source directory"); free(scanner->root_path); + array_list_delete(scanner->pending_dirs); queue_destroy(scanner->directories); array_list_delete(scanner->filter_nodes); free(scanner); @@ -757,6 +888,7 @@ DirectoryScanner* directory_scanner_create_with_options(const char* root_directo if (!queue_enqueue(scanner->directories, root)) { dir_entry_destroy(root); free(scanner->root_path); + array_list_delete(scanner->pending_dirs); queue_destroy(scanner->directories); array_list_delete(scanner->filter_nodes); free(scanner); @@ -808,6 +940,13 @@ void directory_scanner_destroy(DirectoryScanner* scanner) { free(scanner->current_path); free(scanner->current_rel); free(scanner->root_path); + scanner_free_sorted(scanner); + ArrayList* pending = (ArrayList*)scanner->pending_dirs; + if (pending) { + for (int i = 0; i < pending->size; i++) + dir_entry_destroy(pending->items[i]); + array_list_delete(pending); + } array_list_delete(scanner->filter_nodes); array_list_delete(scanner->dirs_batch); queue_destroy(scanner->directories); @@ -975,7 +1114,7 @@ static int open_next_directory(DirectoryScanner* scanner) { scanner->current_path = NULL; while (!queue_is_empty(scanner->directories)) { - DirEntry* de = (DirEntry*)queue_dequeue(scanner->directories); + DirEntry* de = (DirEntry*)queue_pop(scanner->directories); scanner->current_path = de->path; scanner->current_depth = de->depth; /* The seed directory inherits the scanner's configured context (the root @@ -1060,6 +1199,14 @@ static int open_next_directory(DirectoryScanner* scanner) { scanner->failed = true; return -1; } + /* Buffer and sort this directory's entries in rsync's flist order. */ + if (scanner_buffer_current_directory(scanner) != 0) { + closedir(scanner->current_dir); + scanner->current_dir = NULL; + free(scanner->current_path); + scanner->current_path = NULL; + return -1; + } return 1; } return 0; @@ -1415,8 +1562,7 @@ Chunk* directory_scanner_next(DirectoryScanner* scanner) { break; } - const struct dirent* entry = readdir(scanner->current_dir); - if (entry == NULL) { + if (scanner->sorted_index >= scanner->sorted_count) { /* The directory is exhausted: if nothing was transferred or descended from it, recreate it at the destination as an explicit entry. */ if (scanner->options.emit_empty_dirs && !scanner->current_dir_produced && @@ -1425,10 +1571,12 @@ Chunk* directory_scanner_next(DirectoryScanner* scanner) { if (!scanner_emit_empty_dir(scanner, chunk_data)) scanner->failed = true; } + scanner_push_pending_dirs(scanner); closedir(scanner->current_dir); scanner->current_dir = NULL; free(scanner->current_path); scanner->current_path = NULL; + scanner_free_sorted(scanner); if (scanner->failed) { array_list_delete(chunk_data); return NULL; @@ -1436,26 +1584,15 @@ Chunk* directory_scanner_next(DirectoryScanner* scanner) { continue; } - if (strcmp(entry->d_name, ".") == 0 || strcmp(entry->d_name, "..") == 0) - continue; + SortedEntry* sorted = &((SortedEntry*)scanner->sorted_entries)[scanner->sorted_index++]; + const char* name = sorted->name; + ScannerEntry* inspected = &sorted->entry; + int inspection = sorted->inspection; - ScannerEntry inspected; - char* link_rel = child_rel_path(scanner->current_rel, entry->d_name); - if (!link_rel) { - scanner->failed = true; - break; - } - int inspection = scanner_inspect_entry(&scanner->options, scanner->current_path, link_rel, - entry->d_name, &inspected); - free(link_rel); - if (inspection < 0) { - scanner->failed = true; - break; - } if (inspection == 0) { /* A dereferenced symlink with no referent is a partial-transfer error (rsync exit 23): record it as a non-fatal scan I/O error. */ - if (inspected.referent_error) + if (inspected->referent_error) scanner->io_error = true; /* A user-selection exclude protects its destination mirror from --delete unless --delete-excluded; a size prune is always protected. Other @@ -1463,23 +1600,23 @@ Chunk* directory_scanner_next(DirectoryScanner* scanner) { --files-from the protected prefix must be the entry's bare relative wire path, not its source path (which would not match the destination layout and would leave the mirror deletable). */ - if (inspected.excluded) { + if (inspected->excluded) { char* protected_path; if (scanner->relative_mode) { - protected_path = child_rel_path(scanner->current_rel, entry->d_name); + protected_path = child_rel_path(scanner->current_rel, name); } else if (scanner->options.relative_prefix) { - char* relc = child_rel_path(scanner->current_rel, entry->d_name); + char* relc = child_rel_path(scanner->current_rel, name); protected_path = relc ? scanner_prefix_send_path(scanner->options.relative_prefix, relc) : NULL; free(relc); } else { - protected_path = path_cat(scanner->current_path, entry->d_name); + protected_path = path_cat(scanner->current_path, name); } if (!protected_path) { scanner->failed = true; break; } - if (inspected.size_excluded) + if (inspected->size_excluded) scanner_record_size_skipped(scanner, protected_path); else scanner_record_excluded(scanner, protected_path); @@ -1487,23 +1624,22 @@ Chunk* directory_scanner_next(DirectoryScanner* scanner) { } continue; } - char* cur_path = inspected.path; - struct stat stats = inspected.stats; + char* cur_path = inspected->path; + struct stat stats = inspected->stats; /* --files-from allow-set and the filter layer apply to files and to * directories (an excluded directory is not descended into). */ - bool is_dir = inspected.is_directory; - char* rel = child_rel_path(scanner->current_rel, entry->d_name); + bool is_dir = inspected->is_directory; + char* rel = child_rel_path(scanner->current_rel, name); if (!rel) { - free(cur_path); scanner->failed = true; break; } bool protect = false; bool passes_selection = entry_passes_selection( - scanner->options.file_list, scanner->options.base_filters, scanner->current_node, rel, - entry->d_name, is_dir, scanner->options.per_dir_filters, - scanner->options.exclude_per_dir_filter_files, &protect); + scanner->options.file_list, scanner->options.base_filters, scanner->current_node, rel, name, + is_dir, scanner->options.per_dir_filters, scanner->options.exclude_per_dir_filter_files, + &protect); /* A sender-side hide leaves the entry out of the transfer; an independent receiver-side protect rule keeps a transferred entry's destination mirror from being deleted. Both are recorded in the same protection set. */ @@ -1525,7 +1661,6 @@ Chunk* directory_scanner_next(DirectoryScanner* scanner) { char* wrel = scanner_prefix_send_path(scanner->options.relative_prefix, rel); if (!wrel) { free(rel); - free(cur_path); scanner->failed = true; break; } @@ -1542,13 +1677,11 @@ Chunk* directory_scanner_next(DirectoryScanner* scanner) { char* rel_copy = needs_rel ? str_dup(rel) : NULL; free(rel); if (rel_copy == NULL && needs_rel) { - free(cur_path); scanner->failed = true; break; } if (!passes_selection) { free(rel_copy); - free(cur_path); continue; } @@ -1563,12 +1696,10 @@ Chunk* directory_scanner_next(DirectoryScanner* scanner) { File* mount = scanner_build_dir_file(cur_path, &stats, &scanner->options); if (mount == NULL || !array_list_add(chunk_data, mount)) { file_destroy(mount); - free(cur_path); scanner->failed = true; break; } scanner->current_dir_produced = true; - free(cur_path); continue; } /* --list-only: list directory entries too (rsync prints them), even @@ -1577,7 +1708,6 @@ Chunk* directory_scanner_next(DirectoryScanner* scanner) { File* dir = scanner_build_dir_file(cur_path, &stats, &scanner->options); if (dir == NULL || !array_list_add(chunk_data, dir)) { file_destroy(dir); - free(cur_path); scanner->failed = true; break; } @@ -1586,32 +1716,29 @@ Chunk* directory_scanner_next(DirectoryScanner* scanner) { int next_depth = scanner->current_depth + 1; if (scanner->options.max_depth <= 0 || next_depth < scanner->options.max_depth) { DirEntry* de = dir_entry_create(cur_path, next_depth, scanner->current_node); - if (!de || !queue_enqueue(scanner->directories, de)) { + if (!de || !array_list_add((ArrayList*)scanner->pending_dirs, de)) { dir_entry_destroy(de); scanner->failed = true; } } - free(cur_path); } else { if (scanner->options.max_depth > 0 && scanner->current_depth + 1 > scanner->options.max_depth) { free(rel_copy); - free(cur_path); continue; } File* file = file_create(cur_path); - free(cur_path); if (file == NULL) { free(rel_copy); - free(inspected.link_target); - inspected.link_target = NULL; + free(inspected->link_target); + inspected->link_target = NULL; scanner->failed = true; continue; } - if (inspected.is_symlink) { + if (inspected->is_symlink) { file->is_symlink = true; - file->symlink_target = inspected.link_target; - inspected.link_target = NULL; + file->symlink_target = inspected->link_target; + inspected->link_target = NULL; } else { file->data->size = stats.st_size; } diff --git a/src/client/scanner.h b/src/client/scanner.h index 18f61a6..55fc3bb 100644 --- a/src/client/scanner.h +++ b/src/client/scanner.h @@ -10,6 +10,7 @@ #include "stop_condition.h" #include #include +#include #include #include #include @@ -188,6 +189,17 @@ typedef struct { int current_depth; dev_t root_dev; bool failed; + /* rsync-order traversal: each opened directory's entries are inspected once + and buffered (an internal SortedEntry[] owned here) sorted as rsync's flist + orders them -- non-directories ascending, then directories ascending. The + entries are walked in order and child directories are collected in + `pending_dirs` (an ArrayList of DirEntry*, owned here) and pushed onto the + LIFO `directories` stack in reverse at directory exhaustion, so the emitted + stream is depth-first like rsync. `sorted_*` are reset per directory. */ + void* sorted_entries; + size_t sorted_count; + size_t sorted_index; + void* pending_dirs; /* Recursive scan: whether the open directory yielded any transferred or descended entry. When it did not, closing it emits a directory entry so the empty source directory is recreated at the destination (rsync diff --git a/src/shared/queue.c b/src/shared/queue.c index e4e5f09..3ebac08 100644 --- a/src/shared/queue.c +++ b/src/shared/queue.c @@ -139,6 +139,23 @@ void* queue_dequeue(Queue* queue) { return item; } +bool queue_push(Queue* queue, void* item) { + return queue_enqueue(queue, item); +} + +void* queue_pop(Queue* queue) { + if (queue == NULL || queue_is_empty(queue)) { + log_perror("ERROR: Could not pop from null or empty queue."); + return NULL; + } + + queue->rear = (queue->rear - 1 + queue->capacity) % queue->capacity; + void* item = queue->items[queue->rear]; + queue->items[queue->rear] = NULL; + queue->size--; + return item; +} + void* queue_dequeue_multithreaded(Queue* queue, mtx_t* mutex, cnd_t* condition_not_empty, cnd_t* condition_not_full, const bool* other_thread_done) { mtx_lock(mutex); diff --git a/src/shared/queue.h b/src/shared/queue.h index 7a8bd98..3fddd9d 100644 --- a/src/shared/queue.h +++ b/src/shared/queue.h @@ -28,4 +28,11 @@ void* queue_dequeue(Queue* queue); void* queue_dequeue_multithreaded(Queue* queue, mtx_t* mutex, cnd_t* condition_not_empty, cnd_t* condition_not_full, const bool* other_thread_done); +/* LIFO stack operations over the same ring buffer. queue_push() is the enqueue + primitive; queue_pop() removes from the rear, so a sequence of pushes is + returned in reverse order. Used by the sequential scanner's depth-first + traversal. */ +bool queue_push(Queue* queue, void* item); +void* queue_pop(Queue* queue); + #endif diff --git a/tests/integration/test_parity_order.py b/tests/integration/test_parity_order.py new file mode 100644 index 0000000..44dc16b --- /dev/null +++ b/tests/integration/test_parity_order.py @@ -0,0 +1,135 @@ +"""Differential rsync-parity coverage for FastSync's transfer/delete ORDER. + +rsync walks a source tree in its sorted flist order: within each directory the +non-directories come first (ascending name), then the subdirectories (ascending +name), each subdirectory immediately followed by its own subtree (depth-first). +The sequential scanner now reproduces that order, which makes both the +``--info=name`` stream and the ``--delete-during`` deletion sequence match real +``rsync 3.4.1`` exactly. ``--threads`` has no rsync analogue and is unordered. + +Every test skips cleanly when rsync is absent. +""" +import os +import shutil +import subprocess +import sys + +import pytest + +sys.path.insert(0, os.path.dirname(__file__)) +from common import ( # noqa: E402 + TEST_DATA_DIR, + ServerManager, + clean_dir, + get_dest_received_dir, + run_client, +) + +RSYNC = shutil.which("rsync") +requires_rsync = pytest.mark.skipif(RSYNC is None, reason="rsync 3.4.1 not installed") + +MTIME = 1_500_000_000 + +_TREE = { + "a.txt": b"a\n", + "b.txt": b"b\n", + "z.txt": b"z\n", + "a_dir/f.txt": b"f\n", + "a_dir/deep/g.txt": b"g\n", + "m_dir/h.txt": b"h\n", + "Z_dir/i.txt": b"i\n", +} + + +def _write(path, data): + os.makedirs(os.path.dirname(path), exist_ok=True) + with open(path, "wb") as fh: + fh.write(data) + os.utime(path, (MTIME, MTIME)) + + +def _rsync(args): + env = dict(os.environ, LC_ALL="C") + return subprocess.run([RSYNC] + args, capture_output=True, text=True, env=env, timeout=120) + + +def _deleting(text): + out = [] + for line in text.splitlines(): + stripped = line.strip() + if stripped.startswith("*deleting") or stripped.startswith("deleting"): + out.append(stripped.split()[-1]) + return out + + +class TestTransferOrderParity: + @requires_rsync + def test_info_name_file_order_matches_rsync(self, shared_server): + source = os.path.join(TEST_DATA_DIR, "order_name_src") + clean_dir(source) + for rel, data in _TREE.items(): + _write(os.path.join(source, rel), data) + + rdst = os.path.join(TEST_DATA_DIR, "order_name_rdst") + clean_dir(rdst) + rs = _rsync(["-a", "--info=name", source + "/", rdst + "/"]) + assert rs.returncode == 0, rs.stderr + # rsync also names the directories (trailing '/'); FastSync names the + # transferred entries. Compare the file/symlink sequence, which is what + # the traversal order determines. + rsync_files = [l for l in rs.stdout.splitlines() if l.strip() and not l.endswith("/")] + + fdst = os.path.join(TEST_DATA_DIR, "order_name_fdst") + clean_dir(fdst) + result, _ = run_client(source, fdst, flags=["-a", "--info=name"], + port=shared_server.port) + assert result.returncode == 0, (result.stderr or result.stdout)[:300] + fsync_files = [ + l for l in result.stdout.splitlines() + if l.strip() and l.strip() != "./" and not l.startswith("sending") + ] + assert fsync_files == rsync_files, ( + f"transfer order differs\nrsync={rsync_files}\nfastsync={fsync_files}") + + +class TestDeleteOrderParity: + _EXTRA = { + "a_extra.txt": b"a\n", + "z_extra.txt": b"z\n", + "a_extra_dir/f": b"f\n", + "z_extra_dir/f": b"f\n", + "a_extra_dir/sub/g": b"g\n", + } + + @requires_rsync + def test_delete_during_deletion_order_matches_rsync(self): + source = os.path.join(TEST_DATA_DIR, "order_del_src") + clean_dir(source) + _write(os.path.join(source, "keep.txt"), b"k\n") + _write(os.path.join(source, "keepdir", "x.txt"), b"x\n") + _write(os.path.join(source, "keep2", "y.txt"), b"y\n") + + rdst = os.path.join(TEST_DATA_DIR, "order_del_rdst") + clean_dir(rdst) + for rel, data in self._EXTRA.items(): + _write(os.path.join(rdst, rel), data) + rs = _rsync(["-a", "--delete-during", "--info=del", source + "/", rdst + "/"]) + assert rs.returncode == 0, rs.stderr + + fdst = os.path.join(TEST_DATA_DIR, "order_del_fdst") + clean_dir(fdst) + received = get_dest_received_dir(fdst, source) + for rel, data in self._EXTRA.items(): + _write(os.path.join(received, rel), data) + with ServerManager() as server: + server.start(extra_args=["--allow-delete"]) + result, _ = run_client(source, fdst, flags=["-a", "--delete-during", "--info=del"], + port=server.port) + assert result.returncode == 0, (result.stderr or result.stdout)[:300] + + rsync_order = _deleting(rs.stdout) + fsync_order = _deleting(result.stdout) + assert sorted(fsync_order) == sorted(rsync_order), ( + f"deleted set differs\nrsync={rsync_order}\nfastsync={fsync_order}") + assert fsync_order == rsync_order, ( + f"deletion order differs\nrsync={rsync_order}\nfastsync={fsync_order}")