diff --git a/.github/workflows/tests.yml b/.github/workflows/tests.yml new file mode 100644 index 0000000..457dfab --- /dev/null +++ b/.github/workflows/tests.yml @@ -0,0 +1,24 @@ +name: Tests + +on: + push: + pull_request: + +permissions: + contents: read + +jobs: + test: + runs-on: ubuntu-latest + strategy: + matrix: + python-version: ['3.10', '3.14'] + steps: + - uses: actions/checkout@v4 + - uses: actions/setup-python@v5 + with: + python-version: ${{ matrix.python-version }} + - name: Install transfer tools + run: sudo apt-get update && sudo apt-get install -y rsync rclone fpart + - name: Run regression and local integration tests + run: python -m unittest discover -s tests -v diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..8eef5cb --- /dev/null +++ b/.gitignore @@ -0,0 +1,4 @@ +__pycache__/ +*.py[cod] +.venv/ +.ruff_cache/ diff --git a/README.md b/README.md index 813ce4d..3b73e42 100644 --- a/README.md +++ b/README.md @@ -1,155 +1,213 @@ # dsync -dsync is a python script developed to move data, fast. It utilizes fpart(https://github.com/martymac/fpart) to break apart the source directory so that a copy can be run accross multiple instances of rsync(https://rsync.samba.org/) or rclone(https://rclone.org) on a single, or multiple hosts. -## Getting Started +dsync partitions a local directory and copies its contents using parallel +[rsync](https://rsync.samba.org/) or [rclone](https://rclone.org/) processes. +[fpart](https://www.fpart.org/) balances chunks by size; basic chunking is also +available without fpart. Transfers can run locally or across SSH worker hosts +that share access to the source filesystem. -These are the requirements and necessary steps in order to get you up and running with dsync. +## Requirements -### Requirements and prerequisites +- Linux and Python 3.10 or later. No Python packages are required. +- rsync for filesystem transfers, or rclone for `--cloud` transfers. +- fpart unless using `--no-fpart` or reusing existing chunks. +- SSH with noninteractive authentication when using worker or destination hosts. +```sh +git clone https://github.com/daltschu22/dsync.git +cd dsync +python3 dsync.py --help ``` -Python 3+ -fpart (https://github.com/martymac/fpart) -rsync(https://rsync.samba.org/) -rclone(https://rclone.org) -``` - -## Installing and Configuring -Clone the repo or download the necessary files. -``` -git clone https://github.com/daltschu11/dsync.git +Install the external tools through your distribution, for example: +```sh +sudo apt install rsync fpart rclone ``` -Install fpart, rsync, and rclone if using one or all. - -Rsync should already be installed on most linux systems. -Fpart instructions can be found in the README: https://github.com/martymac/fpart/blob/master/README -rclone installation instructions can be found here: https://rclone.org/downloads/ - -### Configuring rclone +## Usage -You need to have an endpoint configured if you plan on running dsync with rclone. +Copy the **contents** of a directory with up to four concurrent processes: -The documentation for that is here: https://rclone.org/commands/rclone_config/ -But, its fairly simple. - -And example of a section of the config file for a Google Cloud Storage configuration is here: -``` -[google-cloud-bucket-1] -type = google cloud storage -client_id = -client_secret = -project_number = -service_account_file = /path/to/my/json-file.json -object_acl = -bucket_acl = -location = us -storage_class = COLDLINE +```sh +python3 dsync.py /mnt/source /mnt/destination -n 4 ``` -Or you can run through the config interactively and select n for new config: +Use basic chunking, including hidden files and directories: +```sh +python3 dsync.py /mnt/source /mnt/destination -n 4 --no-fpart ``` -$ rclone config -Current remotes: - -Name Type -==== ==== +With rsync, basic chunks contain top-level entries and each transfer recurses +into its assigned directories. With rclone, basic chunking walks the source +recursively and distributes file paths across chunks. It streams the listing +to disk instead of retaining the whole tree in memory. For uneven directory +sizes, fpart generally provides better load balancing. Both modes exclude +`.zfs` and names beginning with `.snapshot` at every level. -e) Edit existing remote -n) New remote -d) Delete remote -r) Rename remote -c) Copy remote -s) Set configuration password -q) Quit config +Upload using an endpoint created with [`rclone config`](https://rclone.org/commands/rclone_config/): -e/n/d/r/c/s/q> n +```sh +python3 dsync.py /mnt/source remote:bucket/path -n 4 --cloud \ + --rclone-config ~/.config/rclone/rclone.conf ``` -## Running dsync +`--cloud` uses `rclone copy`, with two file transfers per process. Local rclone +destinations are also supported. Rsync uses archive mode. Neither backend +deletes destination files to mirror the source. -dsync is run with a defined number of threads, a source, and a destination. These are the required arguments. - If defaults are used, the working directory that dsync uses to store working files is ~/dsync_working/. - If defaults are used, the log directory is stored inside the working directory. - dsync will default to using rsync (For local NFS transfers). - dsync will also default to using fpart to chunk out the source directory. +Preview a transfer with `--dry-run`. It still creates local chunks and logs, +but does not create the destination or perform cloud write/delete probes: -Default flags for rsync is `-av` which means rsync runs in `archive` mode. -Default flags for rclone is `-v`. rclone will default to using the `copy` command. Which means source data will only be copied, not moved or deleted. -Fpart will default to ignoring `.zfs` and `.snapshot*` directories. +```sh +python3 dsync.py /mnt/source remote:bucket/path -n 4 --cloud --dry-run +``` -If dsync is run again with the same source directory, it will rerun the chunking process. - Dont forget to add the `--reuse` flag in order to reuse the stored chunk files if you dont want to rerun the chunking. +The command waits for all transfers. Any partitioning or transfer failure +returns a nonzero exit status and identifies the relevant error log. Ctrl+C +and SIGTERM stop active local partition/transfer process groups before exiting +with status 130 and 143 respectively. Cleanup waits up to five seconds before +killing surviving group members, even if their parent has exited. Further +interruptions during cleanup do not abandon the remaining processes. Remote +worker cleanup depends on SSH and the remote tool's disconnect behavior. + +## Working files, logs, and reuse + +The default working directory is `~/dsync_working`. Logs default to its `logs` +subdirectory; override these with `--working-dir` and `--log-output`. +Working/log directories must not overlap the source or a local destination. +Local source and destination trees must also be disjoint. + +The working directory reserves `chunks/`, `manifest.json`, and `.dsync.lock`. +Only one run can use it at a time. Use separate working **and log** directories +for independent simultaneous runs. Logs are overwritten on subsequent runs. +Chunk generation is staged so a partitioning failure preserves the previous +completed chunk set. Regeneration replaces only unchanged chunk files recorded +in a valid manifest. Unrecognized files, modified chunks, and symlinked chunk +directories cause an error and are preserved, including during dry runs. Use +a new working directory or move those files aside after reviewing them. An +interrupted replacement may also require a new working directory. + +Reuse a successful partition without rescanning the source: + +```sh +python3 dsync.py /mnt/source /mnt/destination -n 2 --reuse +``` -dsync can be pointed towards a cloud location by using the `--cloud` flag. - Your destination will need to follow the rclone convention of `configured-endpoint:bucket/path/path` - This will launch dsync using rclone as the transfer tool. Please define your cloud endpoints in the rclone config file before running. - You can also feed rclone a config file using `--rclone-config` otherwise it will default to `~/.config/rclone/rclone.conf`. +Reuse verifies the source path and filesystem identity, backend, chunking +mode, and chunk checksums. Repeat `--cloud` and/or `--no-fpart` if used during +generation. `-n` still caps concurrent processes, even if more chunks exist. +Legacy chunks without a manifest must be regenerated once. -fpart can be skipped by using --no-fpart. This will perform a rudimentary chunking of the top level directories of the source path. +Reuse does **not** discover new files within the source tree. Regenerate chunks +when files are added or removed. Before each cloud chunk is copied, a Python +check on the transfer host verifies that all listed source entries exist and +spools the list to a temporary file. Missing entries fail that chunk instead +of being silently skipped by rclone. This adds a metadata lookup per entry +and temporary disk space for the list, and cannot prevent changes after the +check. Use a stable source snapshot when consistency is required. -You can run rsync or clone in dry run mode using the flag `--dry-run` +## Multiple hosts +Host files contain one hostname or `user@hostname` per line. Blank lines and +lines beginning with `#` are ignored. Hosts are assigned round-robin. -If you run dsync with the `-h` flag you will get the usage: -``` -uusage: dsync.py [-h] -n NUMBER [--no-fpart] [--source-hosts SOURCE_HOSTS] - [--destination-hosts DESTINATION_HOSTS] [--reuse] [--cloud] - [--dry-run] [--rclone-config RCLONE_CONFIG] - [--working-dir /working/dir/] [--log-output /log/dir/] - /source/path/ /destination/path/ OR - cloud-prefix:bucket-name/path/in/bucket/ - -Uses fpart to bag up filesystems into defined chunks, then transfers them -using rsync/rclone either on a single host - -positional arguments: - /source/path/ Source path -- use absolute paths! (dsync always - behaves as if you used a trailing slash in rsync!) - /destination/path/ OR cloud-prefix:bucket-name/path/in/bucket/ - Destination path -- use absolute paths! - -optional arguments: - -h, --help show this help message and exit - -n NUMBER, --number NUMBER - Pack files into chunks and kickoff - transfers - --no-fpart Run without fpart in basic mode (Chunks consist of top - level files/dirs) WARNING: BROKEN WITH RCLONE, ONLY - TRANSFERS FILES! - --source-hosts SOURCE_HOSTS - Provide a file with a list of hosts you want to run - the transfers to run on (will evenly balance out the # - of transfers with the number of hosts) - --destination-hosts DESTINATION_HOSTS - Provide a file with a list of hosts you want the - transfers to run against (For example if you have a - number of remote hosts with an NFS storage mount) - --reuse Reuse existing chunk files from same source, and same - working directory - --cloud Upload data to a cloud provider using rclone instead - of local rsync - --dry-run Run rclone or rsync in dry run mode (Wont actually - copy anything) - --rclone-config RCLONE_CONFIG - Path to config file for rclone (If not defined, will - default to ~/.config/rclone/rclone.conf) - --working-dir /working/dir/ - Directory in which temp files will be stored while - running (default is your home dir ~/dsync_working/) - --log-output /log/dir/ - location for the log files (Default is in the working - directory ~/dsync_working/logs/) +```sh +python3 dsync.py /mnt/source /mnt/destination -n 8 \ + --source-hosts workers.txt --destination-hosts storage.txt ``` -## License +- Source workers must see the same absolute source path and have the transfer + tool on their `PATH`. Chunk contents are sent over SSH stdin, so workers do + not need access to the working directory. +- Without `--destination-hosts`, the destination is local to each source worker. +- With `--destination-hosts`, each destination host must expose the same + destination storage. The option applies only to rsync. +- Rclone workers must have the configuration at the same absolute + `--rclone-config` path, `python3` on their `PATH`, and writable temporary + storage for their chunk list. Encrypted rclone configurations must be + unlocked noninteractively, for example through `RCLONE_CONFIG_PASS` on the + transfer host. SSH workers need noninteractive access to any destination + hosts they use. + +Filesystem destinations on SSH hosts must be absolute paths; dsync passes them +unchanged so the destination host resolves any symlinks or `..` components. +Source and configuration paths sent to workers retain their symlink spelling +(relative paths and `~` are expanded on the controller). Destination overlap +checks apply to local transfers; the controller cannot validate a remote +host's filesystem layout. + +Both SSH hops use batch authentication, preserve stdin without a terminal, +allow 15 seconds to connect, and send keepalives every 15 seconds with three +missed replies allowed. Host keys must already be trusted. This detects a dead +SSH connection; it does not impose a total transfer deadline or resolve a +stalled NFS mount. Each host must mount the intended shared source/destination +storage: matching paths or filenames alone do not establish filesystem identity. + +Host files currently accept DNS names, IPv4 addresses, SSH aliases, and optional +usernames; IPv6 literals are not supported. Rsync remote destinations should +be supplied through `--destination-hosts`, rather than a `host:path` positional +argument. + +## Filename and metadata handling + +Paths with spaces, quotes, or shell metacharacters are passed safely as +arguments. Rsync chunk lists are NUL-delimited, including support for newline +filenames. Rclone uses raw file lists so leading/trailing spaces and names +starting with `#` or `;` are preserved. For compatibility with older rclone +versions, filenames containing newlines or carriage returns cause a clear +error before cloud transfers start. + +Rsync preserves symlinks and empty directories. Rclone uses its standard copy +semantics: empty directories and symlinks are not uploaded. Parallel rsync +jobs can update common parent directories, so final directory modification +times are not guaranteed to match the source. This tool does not provide a +filesystem snapshot or cross-chunk hard-link preservation. + +## Failures and large transfers + +A failed run can leave completed chunks at the destination. Fix the connection +or permissions and rerun, using `--reuse` only while its source listing remains +valid. Rsync and rclone compare existing destination files on the next run; +byte-level upload resumption and incomplete-object cleanup depend on the backend. +Rclone's own retry and network timeout settings remain in effect. A disconnected +worker may continue running until SSH/the remote tool notices the disconnect. + +Basic chunk generation keeps at most 32 chunk files open. Its cloud traversal +streams filenames and holds pending directory paths, so a directory containing +many files does not require a complete filename list in Python memory. Fpart +partitioning and the transfer tools still have their own memory requirements. +Cloud preflight stores one file list per active process in that host's temporary +directory; allow scratch space proportional to those lists. + +`-n` limits transfer processes, not total resource usage. Each rclone process +has two file transfers plus its own checkers, buffers, and destination listings. +Increase concurrency gradually while watching memory, NFS metadata load, cloud +rate limits, and log/scratch disk space. Rclone's `RCLONE_NO_TRAVERSE=true` can +reduce repeated destination listings for small chunks against a large remote, +but can be slower for large unchanged file sets; set it on the transfer hosts +only after comparing that workload. See [rclone's tuning guidance](https://rclone.org/docs/#no-traverse). + +## Development + +The script is organized around `Fpart`, `Rsync`, `Rclone`, and `FilesystemOps`. +`run()` walks through path checks, tool selection, chunk preparation, and +transfers. The tools share process scheduling and cancellation helpers; +`Popen` provides the process handles needed for parallel transfers and cleanup. + +```sh +python3 -m unittest discover -s tests -v +``` -This project is licensed under the MIT License - see the [LICENSE](LICENSE) file for details +The suite covers actual rsync, rclone, and fpart transfers in temporary local +directories, plus command quoting, exit status, dry runs, chunk reuse, +concurrency limits, locking, interruption, missing cloud source files, and a +connection drop against a temporary loopback HTTP endpoint. A low-file-limit +test verifies chunk generation with more chunks than available file descriptors. +External-tool integration tests are skipped when the corresponding binaries +are missing. No cloud credentials or external network destinations are used. -## Acknowledgments +## License ----- \ No newline at end of file +[MIT](LICENSE). diff --git a/checkpyversion.py b/checkpyversion.py deleted file mode 100644 index dc5a60f..0000000 --- a/checkpyversion.py +++ /dev/null @@ -1,5 +0,0 @@ -import sys -def check_py_version(): - if sys.version_info <= (3, 0): - sys.stdout.write("Sorry, requires Python 3.x, not Python 2.x\n") - sys.exit(1) \ No newline at end of file diff --git a/dsync.py b/dsync.py index 048a8d0..b71ea77 100755 --- a/dsync.py +++ b/dsync.py @@ -1,514 +1,622 @@ -#!/usr/bin/env python +#!/usr/bin/env python3 +"""Partition a local filesystem and copy it with bounded parallel transfers.""" # Written by daltschu22 -- https://github.com/daltschu22 -import sys -import os import argparse -from shutil import which -import glob -import subprocess -import checkpyversion +from collections import OrderedDict +from contextlib import contextmanager +import fcntl +import hashlib import itertools +import json +import os +from pathlib import Path +import re +import shlex +import shutil +import signal +import subprocess +import sys +import tempfile +import time -def check_linux(): - if sys.platform != "linux" and sys.platform != "linux2": - print("ERROR: Must run this on a linux machine") - print(sys.platform) - sys.exit() - -def parse_arguments(): - parser = argparse.ArgumentParser(description="Uses fpart to bag up filesystems into defined chunks, \ - then transfers them using rsync/rclone either on a single host") # or clustered") - parser.add_argument('source', metavar='/source/path/', help='Source path -- use absolute paths! \ - (dsync always behaves as if you used a trailing slash in rsync!)') - parser.add_argument('dest', metavar='/destination/path/ OR cloud-prefix:bucket-name/path/in/bucket/', help='Destination path -- use absolute paths!') - parser.add_argument('-n', '--number', type=int, action='store', required=True, help='Pack files into chunks and kickoff transfers') - parser.add_argument('--no-fpart', action='store_true', required=False, help='Run without fpart in basic mode (Chunks consist of top level files/dirs) \ - WARNING: BROKEN WITH RCLONE, ONLY TRANSFERS FILES!') - # parser.add_argument('--fpart-options', action='store', required=False, help='Override the default fpart options (list those here)') - # parser.add_argument('--rsync-options', action='store', required=False, help='Override the default rsync options (list those here)') - parser.add_argument('--source-hosts', default=None, action='store', required=False, help='Provide a file with a list of hosts you want to run the transfers to run on \ - (will evenly balance out the # of transfers with the number of hosts)') - parser.add_argument('--destination-hosts', default=None, action='store', required=False, help='Provide a file with a list of hosts you want the transfers to run against \ - (For example if you have a number of remote hosts with an NFS storage mount)') - parser.add_argument('--reuse', action='store_true', required=False, help='Reuse existing chunk files from same source, and same working directory') - parser.add_argument('--cloud', action='store_true', required=False, help='Upload data to a cloud provider using rclone instead of local rsync') - parser.add_argument('--dry-run', action='store_true', required=False, help='Run rclone or rsync in dry run mode (Wont actually copy anything)') - parser.add_argument( - '--rclone-config', - type=str, - action='store', - required=False, - default=os.path.expanduser('~/.config/rclone/rclone.conf'), - help='Path to config file for rclone (If not defined, will default to ~/.config/rclone/rclone.conf)' - ) - parser.add_argument( - '--working-dir', - action='store', - required=False, - default=os.path.expanduser('~/dsync_working/'), - metavar='/working/dir/', - help='Directory in which temp files will be stored while running (default is your home dir ~/dsync_working/)' - ) - parser.add_argument( - '--log-output', - action='store', - required=False, - metavar='/log/dir/', - default=os.path.expanduser('~/dsync_working/logs/'), - help='location for the log files (Default is in the working directory ~/dsync_working/logs/)' - ) - - if len(sys.argv[2:]) == 0: - parser.print_help() - parser.exit() - - return parser.parse_args() -class Fpart: - # Class for running fpart +class SyncError(Exception): + """An actionable transfer or configuration error.""" - def __init__(self): - self.check_fpart() - def check_fpart(self): # Check if fpart binary exists. - self.fpart_bin = which('fpart') - if self.fpart_bin is None: - print("ERROR: fpart not installed!") - sys.exit() +class SyncInterrupted(KeyboardInterrupt): + def __init__(self, signum): + self.signum = signum - def run_fpart(self, fpart_command, source, log_dir): # Run fpart against the given path. - log_stdout_path = str(log_dir + 'fpart.out') - log_stderr_path = str(log_dir + 'fpart.err') - try: - with open(log_stdout_path, 'w') as out, open(log_stderr_path, 'w') as err: - process = subprocess.Popen(fpart_command, cwd=source, shell=True, stderr=err, stdout=out) - process.wait() - except subprocess.CalledProcessError: - print("ERROR: Something went wrong when running fpart!") - sys.exit() - except Exception as e: - print(e) - sys.exit() - - def generate_chunks(self, file_ops, working_dir, thread_num, source, log_dir): - file_ops.delete_chunks(working_dir) # Delete existing chunks - # Assemble fpart arguments and run fpart - chunk_path = working_dir + 'chunk' - fpart_command = ' '.join([ - self.fpart_bin, - '-Z', - '-x .zfs -x .snapshot*', - '-n %s' % (str(thread_num)), - '-o', - chunk_path, - '.' - ]) - - self.run_fpart(fpart_command, source, log_dir) # Run fpart to create chunk files. - - chunk_pattern = 'chunk*' - chunks = file_ops.list_files_byname(working_dir, chunk_pattern) - chunk_count = len(chunks) - - return chunk_count, chunks -class Rsync: - # Class for running rsync +class Cancellation: + """Defer interruption while a newly spawned process is being registered.""" def __init__(self): - self.check_rsync() - - def check_rsync(self): # Check if rsync binary exists. - self.rsync_bin = which('rsync') - if self.rsync_bin is None: - print("ERROR: rsync not installed!") - sys.exit() - - def run_rsync(self, rsync_command, log_dir, log_stdout_path, log_stderr_path): # Run rsync with given paths. + self.signum = None + self.deferred = 0 + self.cleaning_up = False + + def __call__(self, signum, frame): + if not self.cleaning_up: + self.signum = self.signum or signum + if not self.deferred: + self.raise_if_pending() + + def raise_if_pending(self): + if self.signum is not None and not self.cleaning_up: + self.cleaning_up = True + raise SyncInterrupted(self.signum) + + +@contextmanager +def cancellation_handlers(): + handler = Cancellation() + previous = {sig: signal.getsignal(sig) for sig in (signal.SIGINT, signal.SIGTERM)} + try: + for sig in previous: + signal.signal(sig, handler) + yield + finally: + for sig, old_handler in previous.items(): + signal.signal(sig, old_handler) + + +@contextmanager +def defer_cancellation(): + handler = signal.getsignal(signal.SIGTERM) + if isinstance(handler, Cancellation): + handler.deferred += 1 + try: + yield + finally: + if isinstance(handler, Cancellation): + handler.deferred -= 1 + if not handler.deferred: + handler.raise_if_pending() + + +def start_process(command, processes, **kwargs): + with defer_cancellation(): + process = subprocess.Popen(command, start_new_session=True, **kwargs) + processes.append(process) + return process + + +def positive_int(value): + number = int(value) + if number < 1: + raise argparse.ArgumentTypeError('must be greater than zero') + return number + + +def parse_arguments(argv=None): + parser = argparse.ArgumentParser( + description='Partition a local source and copy its contents with parallel rsync/rclone processes.') + parser.add_argument('source', help='Local source directory (contents are copied)') + parser.add_argument('dest', help='Destination directory or rclone remote:path with --cloud') + parser.add_argument('-n', '--number', type=positive_int, required=True, + help='Number of chunks and maximum concurrent transfer processes') + parser.add_argument('--no-fpart', action='store_true', + help='Use basic chunking (top-level entries for rsync; recursive files for rclone)') + parser.add_argument('--source-hosts', help='File containing SSH hosts on which to run transfers') + parser.add_argument('--destination-hosts', help='File containing rsync destination SSH hosts') + parser.add_argument('--reuse', action='store_true', help='Reuse validated chunks from the same source and mode') + parser.add_argument('--cloud', action='store_true', help='Copy using rclone') + parser.add_argument('--dry-run', action='store_true', help='Preview transfers without writing to the destination') + parser.add_argument('--rclone-config', default='~/.config/rclone/rclone.conf', help='Rclone config file') + parser.add_argument('--working-dir', default='~/dsync_working/', help='Directory for chunks and run state') + parser.add_argument('--log-output', help='Log directory (default: WORKING_DIR/logs)') + return parser.parse_args(argv) + + +def executable(name): + path = shutil.which(name) + if path is None: + raise SyncError('{} is not installed or not on PATH'.format(name)) + return path + + +def local_path(value): + return Path(value).expanduser().resolve() + + +def absolute_path(value): + """Expand a controller-relative path without dereferencing its symlinks.""" + return Path(value).expanduser().absolute() + + +def ssh_options(): + # Use the same noninteractive, bounded connection on both SSH hops. + return ['-T', '-o', 'StdinNull=no', '-o', 'BatchMode=yes', '-o', 'ConnectTimeout=15', + '-o', 'ServerAliveInterval=15', '-o', 'ServerAliveCountMax=3'] + + +def read_hosts(filename): + if not filename: + return [] + hosts = [] + for line in local_path(filename).read_text().splitlines(): + host = line.strip() + if not host or host.startswith('#'): + continue + # Hosts become SSH operands and rsync host:path prefixes, never shell code. + if not re.fullmatch(r'(?:[A-Za-z0-9_][A-Za-z0-9_.-]*@)?[A-Za-z0-9_][A-Za-z0-9_.-]*', host): + raise SyncError('Invalid SSH host {!r} in {}'.format(host, filename)) + hosts.append(host) + if not hosts: + raise SyncError('Host file is empty: {}'.format(filename)) + return hosts + + +def inside(path, parent): + return path == parent or parent in path.parents + + +@contextmanager +def working_lock(working): + working.mkdir(parents=True, exist_ok=True) + with (working / '.dsync.lock').open('a') as lock: try: - with open(log_stdout_path, 'w') as out, open(log_stderr_path, 'w') as err: - subprocess.Popen(rsync_command, shell=True, stdout=out, stderr=err) - except subprocess.CalledProcessError: - print("ERROR: Something went wrong when running rsync!") - sys.exit() - except Exception as e: - print(e) - sys.exit() - - def sync_chunks(self, chunks, source, dest, log_dir, rsync_optional_args): # Run through chunks from fpart and pass to run_rsync() to be ran. - x = 0 - - if 'list_of_source_hosts' in rsync_optional_args: # Adds ssh formatting to rsync command string - list_of_source_hosts = rsync_optional_args.get('list_of_source_hosts') - round_robin_source_hosts = itertools.cycle(list_of_source_hosts) - - if 'list_of_dest_hosts' in rsync_optional_args: - list_of_dest_hosts = rsync_optional_args.get('list_of_dest_hosts') - round_robin_dest_hosts = itertools.cycle(list_of_dest_hosts) - - for chunk in chunks: - log_stdout_path = str(log_dir + 'rsync.out.' + str(x)) - log_stderr_path = str(log_dir + 'rsync.err.' + str(x)) - - rsync_bin = self.rsync_bin - rsync_flags = '-av' - rsync_recursive = '--recursive' - rsync_files_from = '--files-from {}'.format(chunk) - rsync_source = source - - if 'dry_run_yesno' in rsync_optional_args: - rsync_dry_run = '--dry-run' - else: - rsync_dry_run = '' + fcntl.flock(lock, fcntl.LOCK_EX | fcntl.LOCK_NB) + except BlockingIOError: + raise SyncError('Another dsync run is using {}'.format(working)) from None + yield + + +def excluded(name): + return name == '.zfs' or name.startswith('.snapshot') + + +def basic_entries(source, cloud): + if not cloud: + with os.scandir(source) as entries: + for entry in entries: + if not excluded(entry.name): + # Rsync treats leading # and ; as comments, even with --from0. + yield b'./' + os.fsencode(entry.name) + return + + # Walk without accumulating every filename in a wide directory. Keep one + # scandir handle open and only queue directories; rclone skips symlinks. + directories = [str(source)] + while directories: + directory = directories.pop() + with os.scandir(directory) as entries: + for entry in entries: + if excluded(entry.name) or entry.is_symlink(): + continue + if entry.is_dir(follow_symlinks=False): + directories.append(entry.path) + else: + yield os.fsencode(os.path.relpath(entry.path, source)) + + +def cloud_entry(entry): + if b'\n' in entry or b'\r' in entry: + raise SyncError('Rclone chunk lists cannot represent newline/carriage-return filenames: {!r}'.format(entry)) + return entry + b'\n' + + +def nul_entries(path): + pending = b'' + with path.open('rb') as handle: + while True: + block = handle.read(65536) + if not block: + break + entries = (pending + block).split(b'\0') + pending = entries.pop() + yield from entries + if pending: + raise SyncError('Incomplete NUL-delimited chunk: {}'.format(path)) + + +def signal_group(process, signum): + # Reap the leader if possible, but its exit says nothing about descendants. + process.poll() + try: + os.killpg(process.pid, signum) + return True + except ProcessLookupError: + return False + + +def stop_processes(processes, grace_seconds=5): + with defer_cancellation(): + groups = [process for process in processes if signal_group(process, signal.SIGTERM)] + deadline = time.monotonic() + grace_seconds + while groups and time.monotonic() < deadline: + groups = [process for process in groups if signal_group(process, 0)] + if groups: + time.sleep(min(0.05, max(0, deadline - time.monotonic()))) + for process in groups: + signal_group(process, signal.SIGKILL) + for process in processes: + process.wait() - if 'list_of_source_hosts' in rsync_optional_args: # Adds ssh formatting to rsync command string - source_host = next(round_robin_source_hosts) # Round robins the list of hosts - source_host_ssh_head = "ssh {} '".format(source_host) - source_host_ssh_tail = "'" - else: - source_host_ssh_head = '' - source_host_ssh_tail = '' - - if 'list_of_dest_hosts' in rsync_optional_args: - rsync_dest_host = next(round_robin_dest_hosts) # Round robins the list of hosts - rsync_dest = ''.join([ # Join the dest host with the dest path formatted for rsync/rclone - rsync_dest_host, - ':', - dest - ]) - else: - rsync_dest = dest - - rsync_command = ' '.join([ - source_host_ssh_head, - rsync_bin, - rsync_flags, - rsync_recursive, - rsync_files_from, - rsync_dry_run, - rsync_source, - rsync_dest, - source_host_ssh_tail - ]) - - # print('-- ' + rsync_command) # Used for testing if you dont want to actually run the commands - self.run_rsync(rsync_command, log_dir, log_stdout_path, log_stderr_path) - x += 1 - - # subprocess.call(['ps -ef | grep /usr/bin/rsync | grep chunk | grep -v grep'], shell=True) -class Rclone: - # Class for running rclone +class Fpart: + """Run fpart and prepare its file lists for the selected transfer tool.""" def __init__(self): - self.check_rclone() - self.threads = 2 + self.fpart_bin = executable('fpart') - def check_rclone(self): # Check if rclone binary exists. - self.rclone_bin = which('rclone') - if self.rclone_bin is None: - print("ERROR: rclone not installed!") - sys.exit() - - def run_rclone(self, rclone_command, log_dir, log_stdout_path, log_stderr_path): + def run_fpart(self, command, source, out, err): + processes = [] try: - with open(log_stdout_path, 'w') as out, open(log_stderr_path, 'w') as err: - subprocess.Popen(rclone_command, shell=True, stderr=err, stdout=out) - except subprocess.CalledProcessError: - print("ERROR: Something went wrong when running rclone!") - sys.exit() - except Exception as e: - print(e) - sys.exit() - - def sync_chunks(self, chunks, source, dest, log_dir, rclone_optional_args): - x = 0 - - if 'list_of_source_hosts' in rclone_optional_args: # Adds ssh formatting to rsync command string - list_of_source_hosts = rclone_optional_args.get('list_of_source_hosts') - round_robin_source_hosts = itertools.cycle(list_of_source_hosts) - - for chunk in chunks: - log_stdout_path = str(log_dir + 'rclone.out.' + str(x)) - log_stderr_path = str(log_dir + 'rclone.err.' + str(x)) - - rclone_bin = self.rclone_bin - rclone_copy = 'copy' - rclone_flags = '-v' - rclone_transfers = '--transfers {}'.format(self.threads) - rclone_files_from = '--files-from {}'.format(chunk) - rclone_source = source - rclone_dest = dest - - if 'dry_run_yesno' in rclone_optional_args: - rclone_dry_run = '--dry-run' - else: - rclone_dry_run = '' - - if 'list_of_source_hosts' in rclone_optional_args: # Adds ssh formatting to rsync command string - source_host = next(round_robin_source_hosts) # Round robins the list of hosts - source_host_ssh_head = "ssh {} '".format(source_host) - source_host_ssh_tail = "'" - else: - source_host_ssh_head = '' - source_host_ssh_tail = '' - - rclone_command = ' '.join([ - source_host_ssh_head, - rclone_bin, - rclone_copy, - rclone_flags, - rclone_transfers, - rclone_files_from, - rclone_dry_run, - rclone_source, - rclone_dest, - source_host_ssh_tail - ]) - - # print('-- ' + rclone_command) #Used for testing if you dont want to actually run the commands - self.run_rclone(rclone_command, log_dir, log_stdout_path, log_stderr_path) - x += 1 - - def test_write_perms(self, dest, log_dir): # Touches a file to remote to check permissions so you can fail before kicking off all the rclones - log_stdout_path = str(log_dir + 'check_write_perms.out') - log_stderr_path = str(log_dir + 'check_write_perms.err') - test_file_path = str(dest + 'testfile.dsync') - rclone_command = ' '.join([self.rclone_bin, 'touch', '-v', test_file_path]) - - self.run_rclone(rclone_command, log_dir, log_stdout_path, log_stderr_path) - - # Cleanup test file - self.cleanup_write_perms_test(dest, log_dir, test_file_path) - - def cleanup_write_perms_test(self, dest, log_dir, test_file_path): # Cleanup test touch file - print(" + Removing test file") - log_stdout_path = str(log_dir + 'cleanup_write_perms.out') - log_stderr_path = str(log_dir + 'cleanup_write_perms.err') - rclone_command = ' '.join([self.rclone_bin, 'deletefile', '-v', test_file_path]) - - self.run_rclone(rclone_command, log_dir, log_stdout_path, log_stderr_path) - - def clean_fpart_chunks(self, chunks): - for chunk in chunks: - sed_command = ' '.join(['sed -i \'s|^./||\'', chunk]) - process = subprocess.Popen(sed_command, shell=True) - process.wait() - -class Filesystem_Ops(): - # Class for running various filesystem operations - - def make_path(self, make_dir): # Check if the path exists and if not, create it. + process = start_process(command, processes, cwd=source, stdout=out, stderr=err) + return process.wait() + finally: + stop_processes(processes) + + def generate_chunks(self, directory, source, number, cloud, logs): + command = [self.fpart_bin, '-0', '-x', '.zfs', '-x', '.snapshot*', + '-n', str(number), '-o', str(directory / 'chunk')] + if not cloud: + command.append('-z') # Preserve empty directories without recursively copying chunks twice. + command.append('.') + with (logs / 'fpart.out').open('wb') as out, (logs / 'fpart.err').open('wb') as err: + status = self.run_fpart(command, source, out, err) + if status: + raise SyncError('fpart failed (exit {}); see {}'.format(status, logs / 'fpart.err')) + # Fpart can return zero after filesystem traversal errors. Without verbose + # flags its only normal stderr output is partition statistics. Fail closed + # on other diagnostics instead of reporting an incomplete copy as success. + with (logs / 'fpart.err').open('rb') as errors: + for line in errors: + if line.strip() and not re.fullmatch(rb'Part #\d+: size = \d+, files = \d+', line.strip()): + raise SyncError('fpart reported a diagnostic; see {}'.format(logs / 'fpart.err')) + chunks = sorted(path for path in directory.iterdir() + if re.fullmatch(r'chunk\.\d+', path.name) and path.stat().st_size) + if cloud: + for path in chunks: + converted = path.with_suffix(path.suffix + '.tmp') + with converted.open('wb') as out: + for entry in nul_entries(path): + if entry.startswith(b'./'): + entry = entry[2:] + out.write(cloud_entry(entry)) + converted.replace(path) + return chunks + + +def digest(path): + result = hashlib.sha256() + with path.open('rb') as handle: + for block in iter(lambda: handle.read(65536), b''): + result.update(block) + return result.hexdigest() + + +def chunk_identity(source, args): + stat = source.stat() + return {'version': 1, 'source': str(source), 'device': stat.st_dev, 'inode': stat.st_ino, + 'cloud': args.cloud, 'no_fpart': args.no_fpart} + + +def read_manifest(path): + if path.is_symlink(): + raise ValueError('manifest must not be a symlink') + manifest = json.loads(path.read_text()) + identity = manifest['identity'] + if identity['version'] != 1 or not isinstance(identity['source'], str): + raise ValueError('unrecognized manifest identity') + for key in ('device', 'inode'): + if type(identity[key]) is not int: + raise ValueError('invalid source identity') + for key in ('cloud', 'no_fpart'): + if type(identity[key]) is not bool: + raise ValueError('invalid chunking mode') + for name, checksum in manifest['chunks'].items(): + if not re.fullmatch(r'chunk\.\d+', name) or not re.fullmatch(r'[0-9a-f]{64}', checksum): + raise ValueError('invalid chunk record') + return manifest + + +class FilesystemOps: + """Manage basic chunking, saved chunks, and their ownership checks.""" + + def __init__(self, source, working_dir, log_dir): + self.source = source + self.working_dir = working_dir + self.log_dir = log_dir + + def no_fpart_chunk_gen(self, directory, number, cloud): + chunks = [] + handles = OrderedDict() + try: + for index, entry in enumerate(basic_entries(self.source, cloud)): + slot = index % number + if slot == len(chunks): + path = directory / 'chunk.{}'.format(slot) + chunks.append(path) + handle = handles.pop(slot, None) + if handle is None: + if len(handles) >= 32: + _, oldest = handles.popitem(last=False) + oldest.close() + handle = chunks[slot].open('ab') + handles[slot] = handle + handle.write(cloud_entry(entry) if cloud else entry + b'\0') + finally: + for handle in handles.values(): + handle.close() + return chunks + + def owned_chunks(self): + """Return only verified files that a previous dsync run created.""" + working = self.working_dir + target = working / 'chunks' + manifest_path = working / 'manifest.json' + try: + if target.is_symlink() or (target.exists() and not target.is_dir()): + raise ValueError('chunks must be a directory, not a file or symlink') + entries = set(target.iterdir()) if target.exists() else set() + if manifest_path.exists() or manifest_path.is_symlink(): + manifest = read_manifest(manifest_path) + chunks = [target / name for name in manifest['chunks']] + if entries != set(chunks): + raise ValueError('chunk directory contains unrecognized or missing files') + for path in chunks: + if path.is_symlink() or not path.is_file() or digest(path) != manifest['chunks'][path.name]: + raise ValueError('chunk is not an unchanged dsync file: {}'.format(path.name)) + return chunks + if entries: + raise ValueError('existing chunks have no ownership manifest') + return [] + except (OSError, ValueError, KeyError, TypeError, AttributeError) as error: + raise SyncError('Refusing to replace working files: {}. Use a new working directory or move ' + 'the existing files aside after reviewing them.'.format(error)) from error + + def check_existing_chunks(self, args): + working = self.working_dir + manifest_path = working / 'manifest.json' + identity = chunk_identity(self.source, args) try: - print(' + ' + make_dir + " does not exist. making...") - os.makedirs(make_dir) - except OSError as e: - if "Permission denied" in str(e): - print("ERROR: Cannot create path {} due to permissions".format(make_dir)) + manifest = read_manifest(manifest_path) + if manifest['identity'] != identity: + raise ValueError('source or chunking mode differs') + chunks = [] + for name, checksum in manifest['chunks'].items(): + path = working / 'chunks' / name + if (working / 'chunks').is_symlink() or path.is_symlink() or digest(path) != checksum: + raise ValueError('chunk changed: {}'.format(name)) + chunks.append(path) + return chunks + except (OSError, ValueError, KeyError, TypeError, AttributeError) as error: + raise SyncError('Cannot reuse chunks: {}. Regenerate without --reuse; if working files ' + 'were modified, use a new working directory.'.format(error)) from error + + def prepare_chunks(self, args): + source = self.source + working = self.working_dir + logs = self.log_dir + manifest_path = working / 'manifest.json' + identity = chunk_identity(source, args) + if args.reuse: + return self.check_existing_chunks(args) + + self.owned_chunks() # Refuse unrelated data before doing any partition work. + # Build in isolation; a failed partition must not replace the previous good set. + with tempfile.TemporaryDirectory(prefix='.partition-', dir=working) as temp: + directory = Path(temp) + if args.no_fpart: + chunks = self.no_fpart_chunk_gen(directory, args.number, args.cloud) else: - print(e) - sys.exit() + chunks = Fpart().generate_chunks(directory, source, args.number, args.cloud, logs) + manifest = {'identity': identity, 'chunks': {p.name: digest(p) for p in chunks}} + new_manifest = directory / 'manifest.json' + new_manifest.write_text(json.dumps(manifest, indent=2) + '\n') + target = working / 'chunks' + previous_chunks = self.owned_chunks() # Recheck after the source scan. + # Invalidate before replacing so interruption cannot leave reusable stale state. + manifest_path.unlink(missing_ok=True) + for path in previous_chunks: + path.unlink() + target.mkdir(exist_ok=True) + for path in chunks: + # Exclusive creation refuses any file that appeared since validation. + os.link(path, target / path.name) + new_manifest.replace(manifest_path) + return [target / path.name for path in chunks] - def check_path(self, dir_path): # Check if the path exists, if it does return True, if it doesnt return False. - if os.path.exists(dir_path): - return True - else: - return False - def check_tilde(self, dir_path): # Check if path starts with '~'. If so, expand it. - if dir_path: - if dir_path.startswith('~'): - dir_path = os.path.expanduser(dir_path) - - return dir_path - - def trailing_slash(self, dir_path): # Check if path has trailing slash, if it doesnt, add one. - if dir_path: - dir_path = os.path.join(dir_path, '') - - return dir_path - - def list_files_byname(self, dir_path, pattern): # List all files in a directory that match pattern. - glob_name = dir_path + pattern - file_list = glob.glob(glob_name) +class Rsync: + """Build archive transfers from rsync's NUL-delimited chunk lists.""" - return file_list + name = 'rsync' - def check_read_perms(self, path): - access = os.access(path, os.R_OK) + def __init__(self, binary): + self.rsync_bin = binary - return access + def build_command(self, args, source, dest, source_host, dest_host): + rsync_bin = Path(self.rsync_bin).name if source_host else self.rsync_bin + command = [rsync_bin, '-av', '--protect-args', '--from0', '--files-from=-'] + if args.no_fpart: + command += ['--recursive', '--exclude=.zfs', '--exclude=.snapshot*'] + if dest_host: + # This SSH client runs on the worker when --source-hosts is used. + command += ['--rsh', shlex.join(['ssh', *ssh_options()])] + dest = '{}:{}'.format(dest_host, dest) + return finish_transfer_command(command, args, source, dest, source_host) - def check_write_perms(self, path): - access = os.access(path, os.W_OK) - return access +class Rclone: + """Build rclone copy transfers with the requested configuration.""" + + name = 'rclone' + + # Executed on the host reading the source. Rclone silently ignores missing + # --files-from entries, so validate first and give it a rewindable file list. + # The temporary file is unlinked automatically; exec keeps only its stdin fd. + source_check = '''import os, sys, tempfile +try: + source = os.fsencode(sys.argv[1]) + with tempfile.TemporaryFile() as files: + for line in sys.stdin.buffer: + entry = line[:-1] if line.endswith(b'\\n') else line + if entry: + os.lstat(os.path.join(source, entry)) + files.write(line) + files.seek(0) + os.dup2(files.fileno(), 0) + os.execvp(sys.argv[2], sys.argv[2:]) +except OSError as error: + print('ERROR: rclone source check or launch failed: {}'.format(error), file=sys.stderr) + sys.exit(1) +''' + + def __init__(self, binary): + self.rclone_bin = binary + self.threads = 2 - def read_file_into_list(self, file): - try: - with open(file, 'r') as f: - lines = f.read().splitlines() - except IOError as err: - print("ERROR: Cannot read or file doesnt exist! {0}: {1}".format(file, err)) - sys.exit() - - return lines - - def delete_chunks(self, working_dir): - chunk_pattern = 'chunk*' - chunks = self.list_files_byname(working_dir, chunk_pattern) - print(" + Removing previous chunk files...") - for chunk in chunks: - os.remove(chunk) - - def no_fpart_chunk_gen(self, working_dir, source, thread_num): - depth = '*' # Setting a default of 1 levels deep just for now - path = str(source + depth) - file_list = [os.path.basename(x) for x in glob.glob(path)] - - chunks = [file_list[i::thread_num] for i in range(thread_num)] - - x = 0 - chunk_name_list = [] - for chunk in chunks: - chunk_name = working_dir + 'chunk.' + str(x) - chunk_name_list.append(chunk_name) - - with open(chunk_name, 'w') as f: - for line in chunk: - f.write("%s\n" % line) - x += 1 - - return chunk_name_list - - def check_existing_chunks(self, working_dir, source): # Check if chunks already exist. - chunk_pattern = 'chunk*' - chunks = self.list_files_byname(working_dir, chunk_pattern) - if chunks: # Check that there arent 0 chunks. - # first_chunk_file = open(chunks[0], 'r')#BROKEN? - # if source in first_chunk_file.read(): #Check if the source directory is present inside the chunk file. - chunk_count = len(chunks) - true_false = True - else: - chunk_count = 0 - true_false = False - - return true_false, chunk_count, chunks - -def main(): - args = parse_arguments() # Parse arguments - - checkpyversion.check_py_version() # Check that you are running python3 - check_linux() # Check that you are running on linux - - file_ops = Filesystem_Ops() # Initialize filesystem ops class - - # Set variables from cmd line arguments - source = file_ops.trailing_slash(file_ops.check_tilde(args.source)) - dest = file_ops.trailing_slash(file_ops.check_tilde(args.dest)) - thread_num = args.number - no_fpart = args.no_fpart - reuse = args.reuse - working_dir = file_ops.trailing_slash(file_ops.check_tilde(args.working_dir)) - log_dir = file_ops.trailing_slash(file_ops.check_tilde(args.log_output)) - to_cloud = args.cloud - rclone_config = args.rclone_config - dry_run_yesno = args.dry_run - source_hosts = file_ops.check_tilde(args.source_hosts) - dest_hosts = file_ops.check_tilde(args.destination_hosts) - - # Initialize fpart class. - if not no_fpart: - fpart_class = Fpart() - - # If to_cloud is true, use rclone. If not using cloud, use rsync. - if to_cloud: - rclone_class = Rclone() - elif not to_cloud: - rsync_class = Rsync() - - print("-- Checking Paths...") - # Check if source exists, if not exit. - if not file_ops.check_path(source): - print("ERROR: Source path does not exist!") - sys.exit() - - source_read_access = file_ops.check_read_perms(source) # Check if you have read perms on source - if not source_read_access: - print("WARNING: You seem to not have read permissions on the source!") - - # Check if working paths exist, if not make them - if not file_ops.check_path(working_dir): - file_ops.make_path(working_dir) - if not file_ops.check_path(log_dir): - file_ops.make_path(log_dir) - - if not to_cloud or not dest_hosts: # Only run the destination check/creation if using rsync. - if not file_ops.check_path(dest): - file_ops.make_path(dest) - - # Check if not using fpart or reusing existing chunks - if reuse: # If reuse is true, use existing chunk files - reuse_true_false, chunk_count, chunks = file_ops.check_existing_chunks(working_dir, source) # If chunks exist, use those instead of generating. - if reuse_true_false: - print("-- Reusing {0} existing chunk files... (Thread count will be changed to {0})".format(str(chunk_count))) - elif not reuse_true_false: - print("ERROR: The existing chunks dont match the source directory or dont exist!") - sys.exit() + def build_command(self, args, source, dest, source_host, dest_host): + rclone_bin = Path(self.rclone_bin).name if source_host else self.rclone_bin + config = absolute_path(args.rclone_config) if source_host else local_path(args.rclone_config) + command = [rclone_bin, 'copy', '-v', '--transfers', str(self.threads), + '--ask-password=false', '--config', str(config), '--files-from-raw', '-'] + return finish_transfer_command(command, args, source, dest, source_host, + source_check=self.source_check) + + +def finish_transfer_command(command, args, source, dest, source_host, source_check=None): + if args.dry_run: + command.append('--dry-run') + command += ['--', os.path.join(str(source), ''), dest] + if source_check: + python = 'python3' if source_host else sys.executable + command = [python, '-c', source_check, str(source), *command] + if source_host: + command = [executable('ssh'), *ssh_options(), '--', source_host, + ' '.join(shlex.quote(arg) for arg in command)] + return command + + +def run_transfers(args, chunks, transfer, source, dest, source_hosts, dest_hosts, logs): + source_cycle = itertools.cycle(source_hosts or [None]) + dest_cycle = itertools.cycle(dest_hosts or [None]) + jobs = iter(enumerate(chunks)) + active = [] + processes = [] + failures = [] + exhausted = False + tool = transfer.name + try: + while active or not exhausted: + while len(active) < args.number and not exhausted: + try: + index, chunk = next(jobs) + except StopIteration: + exhausted = True + break + command = transfer.build_command(args, source, dest, + next(source_cycle), next(dest_cycle)) + error_path = logs / '{}.err.{}'.format(tool, index) + with chunk.open('rb') as files, (logs / '{}.out.{}'.format(tool, index)).open('wb') as out, error_path.open('wb') as err: + process = start_process(command, processes, stdin=files, stdout=out, stderr=err) + active.append((process, error_path)) + pending = [] + for process, error_path in active: + status = process.poll() + if status is None: + pending.append((process, error_path)) + elif status: + failures.append('exit {}: {}'.format(status, error_path)) + # Retain groups with surviving descendants even after their leader exits. + processes = [process for process in processes + if process.poll() is None or signal_group(process, 0)] + active = pending + if active: + time.sleep(0.05) + finally: + stop_processes(processes) + if failures: + raise SyncError('{} transfer(s) failed; {}'.format(len(failures), '; '.join(failures))) + + +def run(args): + # Check paths and host lists before creating working files. + source = local_path(args.source) + working = local_path(args.working_dir) + logs = local_path(args.log_output) if args.log_output else working / 'logs' + if inside(logs, working / 'chunks'): + raise SyncError('Log directory must not be inside the reserved working chunks directory') + if not source.is_dir(): + raise SyncError('Source is not a directory: {}'.format(source)) + source_hosts = read_hosts(args.source_hosts) + dest_hosts = read_hosts(args.destination_hosts) + if args.cloud and dest_hosts: + raise SyncError('--destination-hosts applies only to rsync, not --cloud') + if any(inside(path, source) or inside(source, path) for path in (working, logs)): + raise SyncError('Working/log directories and the source tree must not overlap') + cloud_remote = args.cloud and ':' in args.dest and not os.path.isabs(args.dest) + remote_filesystem = bool(source_hosts or dest_hosts) and not cloud_remote + if cloud_remote: + dest = args.dest + elif remote_filesystem: + if not os.path.isabs(args.dest): + raise SyncError('Remote filesystem destinations must use an absolute path') + dest = args.dest # Only the destination host can resolve its symlinks. + else: + if not os.path.isabs(os.path.expanduser(args.dest)) and ':' in args.dest: + raise SyncError('Use --destination-hosts for remote rsync destinations') + dest = str(local_path(args.dest)) + if not cloud_remote and not remote_filesystem: + destination = local_path(dest) + if inside(destination, source) or inside(source, destination): + raise SyncError('Source and destination directories must not overlap') + if any(inside(path, destination) or inside(destination, path) for path in (working, logs)): + raise SyncError('Working/log directories and the destination tree must not overlap') + # Choose the transfer tool; source workers use their own PATH. + tool = 'rclone' if args.cloud else 'rsync' + binary = tool if source_hosts else executable(tool) + if args.cloud: + transfer = Rclone(binary) else: - if no_fpart: # Run without fpart (list files/dirs 2 dirs deep) - print("-- Breaking source directory into chunks...") - chunks = file_ops.no_fpart_chunk_gen(working_dir, source, thread_num) - if not no_fpart: - print("-- Running fpart to break source directory into chunks...") - chunk_count, chunks = fpart_class.generate_chunks(file_ops, working_dir, thread_num, source, log_dir) # Generate chunks. - - if source_hosts: # Get a list of the source hosts - print("-- Using a list of source hosts...") - list_of_source_hosts = file_ops.read_file_into_list(source_hosts) - - if dest_hosts: - print("-- Using a list of destination hosts...") - list_of_dest_hosts = file_ops.read_file_into_list(dest_hosts) - - if dry_run_yesno: # Warn the user that no files will be transferred - print("WARNING: --dry-run used (NO FILES WILL ACTUALLY BE TRANSFERRED!)") - - if to_cloud: # If you are running a cloud transfer - print("-- Using rclone...") - rclone_class.rclone_config_file = rclone_config # Add the rclone config file to the class - - print(" + Testing rclone write permissions to bucket") - rclone_class.test_write_perms(dest, log_dir) - - if not no_fpart: - print(" + Cleaning up fpart chunks (remove './')") - rclone_class.clean_fpart_chunks(chunks) - - print(" + Running rclone's...") - rclone_optional_args = {} # dictionary for any optional stuff - if dry_run_yesno: - rclone_optional_args.update({"dry_run_yesno": dry_run_yesno}) # Add dry run flag to dict - if source_hosts: - rclone_optional_args.update({"list_of_source_hosts": list_of_source_hosts}) # Add list of source hosts to dict - - rclone_class.sync_chunks(chunks, source, dest, log_dir, rclone_optional_args) - - else: # If you are running a local transfer - print("-- Using rsync...") - print(" + Running rsync's...") - rsync_optional_args = {} # dictionary for any optional stuff - if dry_run_yesno: - rsync_optional_args.update({"dry_run_yesno": dry_run_yesno}) # Add dry run flag to dict - if source_hosts: - rsync_optional_args.update({"list_of_source_hosts": list_of_source_hosts}) # Add list of source hosts to dict - if dest_hosts: - rsync_optional_args.update({"list_of_dest_hosts": list_of_dest_hosts}) # Add list of dest hosts to dict - - rsync_class.sync_chunks(chunks, source, dest, log_dir, rsync_optional_args) # Run rsync - - -if __name__ == "__main__": - main() + transfer = Rsync(binary) + if source_hosts or dest_hosts: + executable('ssh') + + file_ops = FilesystemOps(source, working, logs) + with working_lock(working): + logs.mkdir(parents=True, exist_ok=True) + # Generate chunks, or validate the saved file lists. + chunks = file_ops.prepare_chunks(args) + print('-- {} {} chunk(s), up to {} concurrent transfers'.format( + 'Reusing' if args.reuse else 'Prepared', len(chunks), args.number), flush=True) + if args.dry_run: + print('-- Dry run: destination will not be changed', flush=True) + if not chunks: + print('-- No entries to transfer') + return + if not args.cloud and not dest_hosts and not source_hosts and not args.dry_run: + Path(dest).mkdir(parents=True, exist_ok=True) + # Copy the chunks and wait for every transfer to finish. + transfer_source = absolute_path(args.source) if source_hosts else source + run_transfers(args, chunks, transfer, transfer_source, dest, source_hosts, dest_hosts, logs) + print('-- {} completed successfully; logs: {}'.format('Dry run' if args.dry_run else 'Transfer', logs)) + + +def main(argv=None): + args = parse_arguments(argv) + try: + with cancellation_handlers(): + run(args) + except SyncInterrupted as error: + print('ERROR: Interrupted by {}; active local process groups stopped'.format( + signal.Signals(error.signum).name), file=sys.stderr) + return 128 + error.signum + except KeyboardInterrupt: + print('ERROR: Interrupted; active local transfer processes stopped', file=sys.stderr) + return 130 + except (SyncError, OSError) as error: + print('ERROR: {}'.format(error), file=sys.stderr) + return 1 + return 0 + + +if __name__ == '__main__': + sys.exit(main()) diff --git a/tests/test_dsync.py b/tests/test_dsync.py new file mode 100644 index 0000000..ea38757 --- /dev/null +++ b/tests/test_dsync.py @@ -0,0 +1,685 @@ +"""Regression tests; integration tests use only temporary local directories.""" + +import os +from contextlib import nullcontext +from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer +from pathlib import Path +import shlex +import shutil +import signal +import socket +import subprocess +import sys +import tempfile +import textwrap +import threading +import time +import unittest +from unittest import mock + +import dsync + + +SCRIPT = Path(dsync.__file__).resolve() + + +class SyncTests(unittest.TestCase): + def setUp(self): + self.temp = tempfile.TemporaryDirectory(prefix='dsync test ') + self.addCleanup(self.temp.cleanup) + self.root = Path(self.temp.name) + self.source = self.root / "source ' $(touch INJECTED)" + self.source.mkdir() + self.dest = self.root / 'destination space' + self.working = self.root / "work ' space" + self.config = self.root / 'custom config.conf' + self.config.write_text('[testremote]\ntype = alias\nremote = {}\n'.format(self.dest)) + + def cli(self, *options, success=True): + result = subprocess.run( + [sys.executable, str(SCRIPT), str(self.source), str(self.dest), + '-n', '3', '--working-dir', str(self.working), *map(str, options)], + capture_output=True, text=True, timeout=30) + if success: + self.assertEqual(result.returncode, 0, result.stdout + result.stderr) + else: + self.assertNotEqual(result.returncode, 0, result.stdout + result.stderr) + return result + + def fixture(self): + files = ['plain', '.hidden', '#hash', ';semicolon', ' leading and trailing ', + "quote'$(touch INJECTED)", 'nested/deep/file', 'nested/.hidden'] + for index, name in enumerate(files): + path = self.source / name + path.parent.mkdir(parents=True, exist_ok=True) + path.write_text('content {}\n'.format(index)) + (self.source / 'empty').mkdir() + for name in ['.zfs', '.snapshot-old', 'nested/.snapshot']: + folder = self.source / name + folder.mkdir() + (folder / 'excluded').write_text('excluded') + return files + + def assert_copied(self, files): + for name in files: + self.assertEqual((self.dest / name).read_bytes(), (self.source / name).read_bytes()) + self.assertFalse((self.dest / '.zfs').exists()) + self.assertFalse((self.dest / '.snapshot-old').exists()) + self.assertFalse((self.dest / 'nested/.snapshot').exists()) + self.assertFalse((self.source / 'INJECTED').exists()) + + @unittest.skipUnless(shutil.which('rsync'), 'rsync required') + def test_rsync_basic_preserves_hidden_special_names_empty_dirs_and_symlinks(self): + files = self.fixture() + (self.source / 'line\nbreak').write_text('newline') + files.append('line\nbreak') + (self.source / 'link').symlink_to('plain') + self.cli('--no-fpart') + self.assert_copied(files) + self.assertTrue((self.dest / 'empty').is_dir()) + self.assertEqual(os.readlink(self.dest / 'link'), 'plain') + + @unittest.skipUnless(shutil.which('rsync') and shutil.which('fpart'), 'rsync and fpart required') + def test_rsync_fpart_preserves_files_empty_dirs_and_symlinks(self): + files = self.fixture() + (self.source / 'line\nbreak').write_text('newline') + files.append('line\nbreak') + (self.source / 'link').symlink_to('plain') + self.cli() + self.assert_copied(files) + self.assertTrue((self.dest / 'empty').is_dir()) + self.assertEqual(os.readlink(self.dest / 'link'), 'plain') + + @unittest.skipUnless(shutil.which('rclone'), 'rclone required') + def test_rclone_basic_recursive_and_raw_names(self): + files = self.fixture() + self.cli('--no-fpart', '--cloud', '--rclone-config', self.config) + self.assert_copied(files) + + @unittest.skipUnless(shutil.which('rclone') and shutil.which('fpart'), 'rclone and fpart required') + def test_rclone_fpart(self): + files = self.fixture() + self.cli('--cloud', '--rclone-config', self.config) + self.assert_copied(files) + + @unittest.skipUnless(shutil.which('rclone'), 'rclone required') + def test_rclone_honors_config_and_preserves_existing_testfile(self): + (self.source / 'file').write_text('data') + self.dest.mkdir() + (self.dest / 'testfile.dsync').write_text('keep me') + result = subprocess.run( + [sys.executable, str(SCRIPT), str(self.source), 'testremote:', '-n', '2', + '--working-dir', str(self.working), '--no-fpart', '--cloud', + '--rclone-config', str(self.config)], capture_output=True, text=True, timeout=30, + cwd=self.root) + self.assertEqual(result.returncode, 0, result.stderr) + self.assertEqual((self.dest / 'file').read_text(), 'data') + self.assertEqual((self.dest / 'testfile.dsync').read_text(), 'keep me') + self.assertFalse((self.root / 'testremote:').exists()) + + def test_dry_runs_do_not_create_destination(self): + (self.source / 'file').write_text('data') + for tool in ['rsync', 'rclone']: + if not shutil.which(tool): + continue + with self.subTest(tool=tool): + options = ['--no-fpart', '--dry-run'] + if tool == 'rclone': + options += ['--cloud', '--rclone-config', self.config] + self.cli(*options) + self.assertFalse(self.dest.exists()) + self.assertTrue((self.working / 'logs').is_dir()) + + def test_nonpositive_counts_and_missing_arguments_fail(self): + for count in ['0', '-1']: + self.cli('--no-fpart', '-n', count, success=False) + result = subprocess.run([sys.executable, str(SCRIPT)], capture_output=True, timeout=5) + self.assertNotEqual(result.returncode, 0) + + def test_missing_source_and_invalid_host_files_fail(self): + for content in ['', '# only a comment\n', '-oProxyCommand=evil', 'host;touch bad']: + hosts = self.root / 'hosts' + hosts.write_text(content) + self.cli('--no-fpart', '--source-hosts', hosts, success=False) + self.source.rmdir() + self.cli('--no-fpart', success=False) + + def test_host_comments_and_whitespace(self): + hosts = self.root / 'hosts' + hosts.write_text(' # comment\n\n host-one \nuser@host.two\n') + self.assertEqual(dsync.read_hosts(str(hosts)), ['host-one', 'user@host.two']) + + def test_overlap_rejected_before_creating_work(self): + self.cli('--no-fpart', '--working-dir', self.source / 'work', success=False) + self.assertFalse((self.source / 'work').exists()) + self.dest = self.source / 'dest' + self.cli('--no-fpart', success=False) + self.assertFalse(self.dest.exists()) + + def test_work_directory_cannot_contain_source_or_destination(self): + (self.source / 'keep').write_text('data') + self.cli('--no-fpart', '--working-dir', self.root, success=False) + self.assertEqual((self.source / 'keep').read_text(), 'data') + self.dest = self.working / 'chunks' + self.cli('--no-fpart', success=False) + self.assertFalse(self.working.exists()) + + def test_logs_cannot_be_deleted_by_chunk_replacement(self): + self.cli('--no-fpart', '--log-output', self.working / 'chunks' / 'logs', success=False) + self.assertFalse(self.working.exists()) + + def test_unowned_chunks_survive_dry_run(self): + chunks = self.working / 'chunks' + chunks.mkdir(parents=True) + (chunks / 'unrelated-backup').write_text('keep backup') + (chunks / 'chunk.0').write_text('keep even a matching filename') + (self.source / 'file').write_text('data') + result = self.cli('--no-fpart', '--dry-run', success=False) + self.assertIn('Refusing to replace working files', result.stderr) + self.assertEqual((chunks / 'unrelated-backup').read_text(), 'keep backup') + self.assertEqual((chunks / 'chunk.0').read_text(), 'keep even a matching filename') + self.assertFalse(self.dest.exists()) + + @unittest.skipUnless(shutil.which('rsync'), 'rsync required') + def test_unrecognized_file_in_owned_chunks_preserves_previous_generation(self): + (self.source / 'file').write_text('data') + self.cli('--no-fpart') + chunks = self.working / 'chunks' + manifest = (self.working / 'manifest.json').read_bytes() + previous = {path.name: path.read_bytes() for path in chunks.iterdir()} + (chunks / 'unrelated-backup').write_text('keep backup') + self.cli('--no-fpart', success=False) + self.assertEqual((chunks / 'unrelated-backup').read_text(), 'keep backup') + self.assertEqual((self.working / 'manifest.json').read_bytes(), manifest) + for name, content in previous.items(): + self.assertEqual((chunks / name).read_bytes(), content) + + def test_chunk_symlink_is_refused_without_touching_target(self): + self.working.mkdir() + external = self.root / 'external' + external.mkdir() + (external / 'backup').write_text('keep') + (self.working / 'chunks').symlink_to(external, target_is_directory=True) + (self.source / 'file').write_text('data') + self.cli('--no-fpart', '--dry-run', success=False) + self.assertEqual((external / 'backup').read_text(), 'keep') + self.assertTrue((self.working / 'chunks').is_symlink()) + + @unittest.skipUnless(shutil.which('rsync'), 'rsync required') + def test_changed_owned_chunk_is_not_deleted_on_regeneration(self): + (self.source / 'file').write_text('data') + self.cli('--no-fpart') + chunk = next((self.working / 'chunks').iterdir()) + chunk.write_text('user replacement') + self.cli('--no-fpart', success=False) + self.assertEqual(chunk.read_text(), 'user replacement') + + @unittest.skipUnless(shutil.which('rsync'), 'rsync required') + def test_owned_chunks_can_be_regenerated_with_fewer_jobs(self): + for index in range(4): + (self.source / str(index)).write_text('data') + self.cli('--no-fpart') + self.assertEqual(len(list((self.working / 'chunks').iterdir())), 3) + self.cli('--no-fpart', '-n', '1') + self.assertEqual(len(list((self.working / 'chunks').iterdir())), 1) + self.cli('--no-fpart', '--reuse') + + def test_remote_dest_does_not_create_local_directory(self): + hosts = self.root / 'hosts' + hosts.write_text('destination-host\n') + (self.source / 'file').write_text('data') + args = dsync.parse_arguments([str(self.source), str(self.dest), '-n', '1', '--no-fpart', + '--working-dir', str(self.working), '--destination-hosts', str(hosts)]) + with mock.patch.object(dsync, 'run_transfers') as transfer, mock.patch.object(dsync, 'executable', side_effect=lambda name: name): + dsync.run(args) + transfer.assert_called_once() + self.assertFalse(self.dest.exists()) + + def test_remote_dest_preserves_controller_symlinks_and_parent_components(self): + hosts = self.root / 'hosts' + hosts.write_text('test-host\n') + controller_only = self.root / 'controller-only' + controller_only.mkdir() + self.dest.symlink_to(controller_only, target_is_directory=True) + (self.source / 'file').write_text('data') + for host_option in ['--source-hosts', '--destination-hosts']: + for suffix in ['', '/../backup']: + with self.subTest(host_option=host_option, suffix=suffix): + destination = str(self.dest) + suffix + args = dsync.parse_arguments([str(self.source), destination, '-n', '1', '--no-fpart', + '--working-dir', str(self.working), host_option, str(hosts)]) + with mock.patch.object(dsync, 'run_transfers') as transfer: + dsync.run(args) + self.assertEqual(transfer.call_args.args[4], destination) + + def test_worker_source_and_config_preserve_controller_symlinks(self): + hosts = self.root / 'hosts' + hosts.write_text('test-host\n') + source_alias = self.root / 'source-alias' + source_alias.symlink_to(self.source, target_is_directory=True) + config_alias = self.root / 'config-alias' + config_alias.symlink_to(self.config) + (self.source / 'file').write_text('data') + args = dsync.parse_arguments([str(source_alias), 'remote:bucket', '-n', '1', '--no-fpart', '--cloud', + '--rclone-config', str(config_alias), '--working-dir', str(self.working), + '--source-hosts', str(hosts)]) + with mock.patch.object(dsync, 'run_transfers') as transfer: + dsync.run(args) + self.assertEqual(transfer.call_args.args[3], source_alias) + command = dsync.Rclone('rclone').build_command(args, source_alias, 'remote:bucket', 'test-host', None) + remote = shlex.split(command[-1]) + self.assertEqual(remote[remote.index('--config') + 1], str(config_alias)) + + def test_remote_dest_rejects_controller_relative_paths(self): + hosts = self.root / 'hosts' + hosts.write_text('test-host\n') + self.dest = Path('relative-destination') + result = self.cli('--no-fpart', '--destination-hosts', hosts, success=False) + self.assertIn('must use an absolute path', result.stderr) + self.assertFalse(self.working.exists()) + + def test_basic_chunking_avoids_empty_jobs(self): + (self.source / '.hidden').write_text('data') + directory = self.root / 'parts' + directory.mkdir() + file_ops = dsync.FilesystemOps(self.source, self.working, self.working / 'logs') + chunks = file_ops.no_fpart_chunk_gen(directory, 100, False) + self.assertEqual(len(chunks), 1) + self.assertEqual(chunks[0].read_bytes(), b'./.hidden\0') + + @unittest.skipUnless(shutil.which('rsync'), 'rsync required') + def test_empty_source_has_no_jobs_and_is_reusable(self): + self.cli('--no-fpart') + self.cli('--no-fpart', '--reuse') + self.assertFalse(self.dest.exists()) + + @unittest.skipUnless(shutil.which('rsync'), 'rsync required') + def test_reuse_validates_source_mode_and_chunk_integrity(self): + (self.source / 'file').write_text('data') + self.cli('--no-fpart') + self.cli('--no-fpart', '--reuse', '-n', '1') + self.cli('--reuse', success=False) # Different chunk mode. + original = self.source + self.source = self.root / 'other-source' + self.source.mkdir() + self.cli('--no-fpart', '--reuse', success=False) + self.source = original + chunk = next((self.working / 'chunks').iterdir()) + chunk.write_bytes(b'tampered\0') + self.cli('--no-fpart', '--reuse', success=False) + + def test_reuse_without_manifest_fails(self): + self.cli('--no-fpart', '--reuse', success=False) + + def test_lock_rejects_competing_run(self): + with dsync.working_lock(self.working): + self.cli('--no-fpart', success=False) + + def test_cloud_rejects_newlines_and_retains_previous_manifest(self): + args = dsync.parse_arguments([str(self.source), 'remote:', '-n', '2', '--cloud', '--no-fpart']) + self.working.mkdir() + logs = self.working / 'logs' + logs.mkdir() + (self.source / 'valid').write_text('data') + file_ops = dsync.FilesystemOps(self.source, self.working, logs) + previous = file_ops.prepare_chunks(args) + manifest = (self.working / 'manifest.json').read_bytes() + (self.source / 'bad\nname').write_text('data') + with self.assertRaises(dsync.SyncError): + file_ops.prepare_chunks(args) + self.assertEqual((self.working / 'manifest.json').read_bytes(), manifest) + self.assertTrue(previous[0].exists()) + + def test_ssh_command_quotes_and_streams_chunk_on_stdin(self): + args = dsync.parse_arguments([str(self.source), str(self.dest), '-n', '2', '--no-fpart']) + with mock.patch.object(dsync, 'executable', return_value='/usr/bin/ssh'): + command = dsync.Rsync('/custom/bin/rsync').build_command( + args, self.source, str(self.dest), 'user@worker', 'storage') + separator = command.index('--') + self.assertEqual(command[:separator], ['/usr/bin/ssh', *dsync.ssh_options()]) + self.assertEqual(command[separator + 1], 'user@worker') + remote = shlex.split(command[separator + 2]) + self.assertEqual(remote[0], 'rsync') + self.assertIn('--files-from=-', remote) + self.assertEqual(remote[-2], str(self.source) + '/') + self.assertEqual(remote[-1], 'storage:' + str(self.dest)) + destination_ssh = shlex.split(remote[remote.index('--rsh') + 1]) + self.assertEqual(destination_ssh, ['ssh', *dsync.ssh_options()]) + + @unittest.skipUnless(shutil.which('rclone'), 'rclone required') + def test_cloud_reuse_fails_when_a_listed_file_has_disappeared(self): + file = self.source / 'file' + file.write_text('data') + self.cli('--cloud', '--no-fpart', '--dry-run', '--rclone-config', self.config) + file.unlink() + self.cli('--cloud', '--no-fpart', '--reuse', '--rclone-config', self.config, success=False) + self.assertIn('source check', (self.working / 'logs/rclone.err.0').read_text()) + self.assertFalse(self.dest.exists()) + + @unittest.skipUnless(shutil.which('rclone'), 'rclone required') + def test_cloud_worker_validates_its_own_source_mount(self): + (self.source / 'file').write_text('data') + worker_source = self.root / 'empty-worker-mount' + worker_source.mkdir() + hosts = self.root / 'hosts' + hosts.write_text('test-worker\n') + self.fake_tool('ssh', ''' + import os, shlex, subprocess, sys + source, worker = os.environ['DSYNC_TEST_SOURCE'], os.environ['DSYNC_TEST_WORKER'] + command = shlex.split(sys.argv[-1]) + command = [worker + arg[len(source):] if arg in (source, source + '/') else arg + for arg in command] + sys.exit(subprocess.call(command)) + ''') + env = {'PATH': str(self.root / 'bin') + os.pathsep + os.environ['PATH'], + 'DSYNC_TEST_SOURCE': str(self.source), 'DSYNC_TEST_WORKER': str(worker_source)} + with mock.patch.dict(os.environ, env): + self.cli('--cloud', '--no-fpart', '--source-hosts', hosts, + '--rclone-config', self.config, success=False) + self.assertIn(str(worker_source), (self.working / 'logs/rclone.err.0').read_text()) + self.assertFalse(self.dest.exists()) + + @unittest.skipUnless(shutil.which('rclone'), 'rclone required') + def test_cloud_connection_drop_is_reported_as_failure(self): + requests = [] + + class DisconnectingServer(BaseHTTPRequestHandler): + def do_PROPFIND(self): + requests.append(self.path) + self.close_connection = True + self.connection.shutdown(socket.SHUT_RDWR) + self.connection.close() + + def log_message(self, *args): + pass + + with ThreadingHTTPServer(('127.0.0.1', 0), DisconnectingServer) as server: + thread = threading.Thread(target=server.serve_forever, daemon=True) + thread.start() + try: + self.config.write_text('[testremote]\ntype = webdav\nurl = http://127.0.0.1:{}\n'.format( + server.server_port)) + (self.source / 'file').write_text('data') + self.dest = 'testremote:backup' + with mock.patch.dict(os.environ, {'RCLONE_RETRIES': '1', 'RCLONE_LOW_LEVEL_RETRIES': '1', + 'RCLONE_TIMEOUT': '1s', 'RCLONE_CONTIMEOUT': '1s'}): + result = self.cli('--cloud', '--no-fpart', '--rclone-config', self.config, success=False) + self.assertTrue(requests, 'rclone did not reach the test endpoint') + self.assertIn('transfer(s) failed', result.stderr) + self.assertIn('EOF', (self.working / 'logs/rclone.err.0').read_text()) + finally: + server.shutdown() + thread.join(timeout=5) + + def test_many_chunks_respect_a_small_descriptor_limit(self): + for index in range(260): + (self.source / str(index)).touch() + directory = self.root / 'chunks' + directory.mkdir() + script = ''' +import resource, sys +from pathlib import Path +from dsync import FilesystemOps +resource.setrlimit(resource.RLIMIT_NOFILE, (64, resource.getrlimit(resource.RLIMIT_NOFILE)[1])) +source, chunks = map(Path, sys.argv[1:]) +FilesystemOps(source, chunks, chunks).no_fpart_chunk_gen(chunks, 128, False) +''' + result = subprocess.run([sys.executable, '-c', script, str(self.source), str(directory)], + cwd=SCRIPT.parent, capture_output=True, text=True, timeout=10) + self.assertEqual(result.returncode, 0, result.stderr) + actual = [entry for chunk in directory.iterdir() for entry in dsync.nul_entries(chunk)] + self.assertEqual(sorted(actual), sorted(b'./' + str(index).encode() for index in range(260))) + self.assertEqual(len(list(directory.iterdir())), 128) + + def test_cloud_listing_yields_before_reading_a_whole_directory(self): + entry = mock.Mock(name='file-entry') + entry.name = 'first' + entry.path = str(self.source / 'first') + entry.is_symlink.return_value = False + entry.is_dir.return_value = False + + def entries(): + yield entry + raise AssertionError('directory listing was consumed eagerly') + + with mock.patch.object(dsync.os, 'scandir', return_value=nullcontext(entries())): + files = dsync.basic_entries(self.source, True) + try: + self.assertEqual(next(files), b'first') + finally: + files.close() + + @unittest.skipUnless(shutil.which('rsync') and shutil.which('rclone'), 'transfer tools required') + def test_source_host_execution_with_local_ssh_shim(self): + files = self.fixture() + hosts = self.root / 'hosts' + hosts.write_text('test-worker\n') + self.fake_tool('ssh', ''' + import subprocess, sys + sys.exit(subprocess.call(['/bin/sh', '-c', sys.argv[-1]])) + ''') + with mock.patch.dict(os.environ, {'PATH': str(self.root / 'bin') + os.pathsep + os.environ['PATH']}): + for cloud in [False, True]: + with self.subTest(cloud=cloud): + options = ['--no-fpart', '--source-hosts', hosts] + if cloud: + options += ['--cloud', '--rclone-config', self.config] + self.cli(*options) + self.assert_copied(files) + shutil.rmtree(self.dest) + + def fake_tool(self, name, code): + bindir = self.root / 'bin' + bindir.mkdir(exist_ok=True) + tool = bindir / name + tool.write_text('#!{}\n'.format(sys.executable) + textwrap.dedent(code)) + tool.chmod(0o755) + return tool + + def test_transfer_and_partition_failures_are_nonzero(self): + self.fake_tool('rsync', 'import sys\nsys.exit(23)\n') + self.fake_tool('fpart', 'import sys\nsys.exit(4)\n') + (self.source / 'file').write_text('data') + with mock.patch.dict(os.environ, {'PATH': str(self.root / 'bin') + os.pathsep + os.environ['PATH']}): + result = self.cli('--no-fpart', success=False) + self.assertIn('exit 23', result.stderr) + result = self.cli(success=False) + self.assertIn('fpart failed', result.stderr) + + def test_fpart_diagnostic_with_zero_exit_is_a_failure(self): + self.fake_tool('fpart', "import sys\nprint('./unreadable: Permission denied', file=sys.stderr)\n") + self.fake_tool('rsync', "raise AssertionError('transfer should not start')\n") + (self.source / 'file').write_text('data') + with mock.patch.dict(os.environ, {'PATH': str(self.root / 'bin') + os.pathsep + os.environ['PATH']}): + result = self.cli(success=False) + self.assertIn('fpart reported a diagnostic', result.stderr) + self.assertFalse(self.dest.exists()) + self.assertFalse((self.working / 'manifest.json').exists()) + + @unittest.skipUnless(shutil.which('rsync') and shutil.which('fpart') and os.geteuid() != 0, + 'rsync, fpart and a non-root user required') + def test_fpart_unreadable_directory_cannot_silently_succeed(self): + unreadable = self.source / 'unreadable' + unreadable.mkdir() + (unreadable / 'file').write_text('data') + unreadable.chmod(0) + try: + self.cli(success=False) + self.assertFalse(self.dest.exists()) + finally: + unreadable.chmod(0o700) + + def test_bounded_concurrency_waits_for_completion(self): + events = self.root / 'events' + tool = self.fake_tool('rsync', ''' + import os, sys, time + from pathlib import Path + event = Path(os.environ['DSYNC_TEST_EVENTS']) + with event.open('a') as log: + log.write('start %s\\n' % os.getpid()) + sys.stdin.buffer.read() + time.sleep(.15) + with event.open('a') as log: + log.write('end %s\\n' % os.getpid()) + ''') + chunks = [] + for index in range(5): + chunk = self.root / 'chunk.{}'.format(index) + chunk.write_bytes(b'file\0') + chunks.append(chunk) + logs = self.root / 'logs' + logs.mkdir() + args = dsync.parse_arguments([str(self.source), str(self.dest), '-n', '2', '--no-fpart']) + with mock.patch.dict(os.environ, {'DSYNC_TEST_EVENTS': str(events)}): + dsync.run_transfers(args, chunks, dsync.Rsync(str(tool)), self.source, str(self.dest), [], [], logs) + active = peak = started = 0 + for line in events.read_text().splitlines(): + if line.startswith('start'): + active += 1 + started += 1 + else: + active -= 1 + peak = max(peak, active) + self.assertEqual(started, 5) + self.assertEqual(peak, 2) + self.assertEqual(active, 0) + + def test_interrupt_stops_transfer_process(self): + self.assert_signal_stops_process(signal.SIGINT, 'rsync') + + def test_sigterm_stops_transfer_process(self): + self.assert_signal_stops_process(signal.SIGTERM, 'rsync') + + def test_sigterm_stops_fpart_process(self): + self.assert_signal_stops_process(signal.SIGTERM, 'fpart') + + def assert_signal_stops_process(self, signum, tool): + pidfile = self.root / 'pid' + self.fake_tool(tool, ''' + import os, time + from pathlib import Path + Path(os.environ['DSYNC_TEST_PID']).write_text(str(os.getpid())) + time.sleep(30) + ''') + (self.source / 'file').write_text('data') + env = dict(os.environ, PATH=str(self.root / 'bin') + os.pathsep + os.environ['PATH'], + DSYNC_TEST_PID=str(pidfile)) + options = ['--no-fpart'] if tool == 'rsync' else [] + process = subprocess.Popen([sys.executable, str(SCRIPT), str(self.source), str(self.dest), + '-n', '1', *options, '--working-dir', str(self.working)], + env=env, stdout=subprocess.PIPE, stderr=subprocess.PIPE) + try: + deadline = time.monotonic() + 5 + while not pidfile.exists() and time.monotonic() < deadline: + time.sleep(.02) + self.assertTrue(pidfile.exists()) + process.send_signal(signum) + out, err = process.communicate(timeout=10) + self.assertEqual(process.returncode, 128 + signum, out + err) + with self.assertRaises(ProcessLookupError): + os.kill(int(pidfile.read_text()), 0) + finally: + if process.poll() is None: + process.kill() + process.communicate() + if pidfile.exists(): + try: + os.killpg(int(pidfile.read_text()), signal.SIGKILL) + except ProcessLookupError: + pass + + def test_cancellation_waits_until_spawned_process_is_registered(self): + processes = [] + fake_process = object() + + def spawn(*args, **kwargs): + os.kill(os.getpid(), signal.SIGTERM) + return fake_process + + previous = signal.getsignal(signal.SIGTERM) + with self.assertRaises(dsync.SyncInterrupted): + with dsync.cancellation_handlers(), mock.patch.object(dsync.subprocess, 'Popen', side_effect=spawn): + dsync.start_process(['rsync'], processes) + self.assertEqual(processes, [fake_process]) + self.assertIs(signal.getsignal(signal.SIGTERM), previous) + + def test_signal_during_cleanup_does_not_abandon_process(self): + ready = self.root / 'ready' + child_code = ''' +import signal, sys, time +from pathlib import Path +signal.signal(signal.SIGTERM, signal.SIG_IGN) +Path(sys.argv[1]).touch() +time.sleep(30) +''' + process = subprocess.Popen([sys.executable, '-c', child_code, str(ready)], start_new_session=True) + signal_group = dsync.signal_group + + def interrupt_cleanup(process, signum): + # Deliver a real signal at a deterministic point during cleanup. + os.kill(os.getpid(), signal.SIGTERM) + return signal_group(process, signum) + + try: + deadline = time.monotonic() + 5 + while not ready.exists() and time.monotonic() < deadline: + time.sleep(.02) + self.assertTrue(ready.exists()) + with self.assertRaises(dsync.SyncInterrupted): + with dsync.cancellation_handlers(), mock.patch.object( + dsync, 'signal_group', side_effect=interrupt_cleanup): + dsync.stop_processes([process], grace_seconds=.2) + self.assertEqual(process.poll(), -signal.SIGKILL) + finally: + if process.poll() is None: + os.killpg(process.pid, signal.SIGKILL) + process.wait() + + def test_cleanup_kills_descendant_after_leader_exits(self): + pidfile = self.root / 'grandchild.pid' + child_code = ''' +import os, signal, sys, time +from pathlib import Path +signal.signal(signal.SIGTERM, signal.SIG_IGN) +Path(sys.argv[1]).write_text(str(os.getpid())) +time.sleep(30) +''' + parent_code = ''' +import subprocess, sys, time +subprocess.Popen([sys.executable, '-c', sys.argv[2], sys.argv[1]]) +time.sleep(30) +''' + process = subprocess.Popen([sys.executable, '-c', parent_code, str(pidfile), child_code], + start_new_session=True) + try: + deadline = time.monotonic() + 5 + while not pidfile.exists() and time.monotonic() < deadline: + time.sleep(.02) + self.assertTrue(pidfile.exists()) + child_pid = int(pidfile.read_text()) + # The leader is already gone when cleanup begins. Its descendant + # still owns the process group and ignores graceful termination. + process.terminate() + process.wait(timeout=5) + dsync.stop_processes([process], grace_seconds=.1) + deadline = time.monotonic() + 5 + while time.monotonic() < deadline: + try: + state = Path('/proc/{}/stat'.format(child_pid)).read_text().rsplit(')', 1)[1].split()[0] + except FileNotFoundError: + break + if state == 'Z': # An orphaned zombie is stopped; PID 1 must reap it. + break + time.sleep(.02) + else: + self.fail('descendant survived process-group cleanup') + finally: + try: + os.killpg(process.pid, signal.SIGKILL) + except ProcessLookupError: + pass + process.wait() + + +if __name__ == '__main__': + unittest.main()