Compare commits

..
7 Commits
15 changed files with 521 additions and 64 deletions
+1 -1
View File
@@ -2,7 +2,7 @@ name: archive-control-archive
services: services:
archive-client: archive-client:
image: sodium/archive-clients:v0.1.5 image: sodium/archive-clients:v0.1.12
user: "1000:1000" user: "1000:1000"
restart: unless-stopped restart: unless-stopped
command: ["--config", "/etc/archive-control/client.toml"] command: ["--config", "/etc/archive-control/client.toml"]
+3
View File
@@ -39,4 +39,7 @@ endpoint = "http://127.0.0.1:8384"
api_key_file = "/run/secrets/syncthing_api_key" api_key_file = "/run/secrets/syncthing_api_key"
api_root = "/var/syncthing" api_root = "/var/syncthing"
local_root = "/data/sync" local_root = "/data/sync"
# This existing folder is physically inside the qB data tree. Map it through
# that same client bind mount so source staging can use hardlinks.
local_path_overrides = { "/var/syncthing/DownloadsSync" = "/data/qb/Sync" }
advertised_addresses = ["dynamic"] advertised_addresses = ["dynamic"]
+1 -2
View File
@@ -2,7 +2,7 @@ name: archive-control-cache
services: services:
archive-client: archive-client:
image: sodium/archive-clients:v0.1.5 image: sodium/archive-clients:v0.1.12
user: "1001:1001" user: "1001:1001"
restart: unless-stopped restart: unless-stopped
network_mode: host network_mode: host
@@ -16,4 +16,3 @@ services:
- ./backups:/var/backups/archive-control - ./backups:/var/backups/archive-control
- /home/ubuntu/Downloads:/data/qb - /home/ubuntu/Downloads:/data/qb
- /home/ubuntu/compose/syncthing/st_home:/data/sync - /home/ubuntu/compose/syncthing/st_home:/data/sync
- /home/ubuntu/Downloads/Sync:/data/sync/DownloadsSync
+1 -1
View File
@@ -2,7 +2,7 @@ name: archive-control-cache
services: services:
archive-client: archive-client:
image: sodium/archive-clients:v0.1.5 image: sodium/archive-clients:v0.1.12
user: "1001:1001" user: "1001:1001"
restart: unless-stopped restart: unless-stopped
network_mode: host network_mode: host
+6 -1
View File
@@ -125,6 +125,11 @@ local_root = "/data/sync"
advertised_addresses = ["dynamic"] advertised_addresses = ["dynamic"]
``` ```
When an existing Syncthing folder is physically nested in the qB data root,
use a `local_path_overrides` entry to map that exact Syncthing API path through
the same client bind mount. This enables hardlinks without creating two Docker
mount boundaries for the same host files.
The remaining node examples omit optional `[connection]`, `[jobs]`, and The remaining node examples omit optional `[connection]`, `[jobs]`, and
`[backup]` tables and therefore use these same defaults; deployments may `[backup]` tables and therefore use these same defaults; deployments may
override them per node. override them per node.
@@ -213,7 +218,7 @@ cache/archive routes according to policy.
```yaml ```yaml
services: services:
archive-client: archive-client:
image: sodium/archive-clients:v0.1.5 image: sodium/archive-clients:v0.1.12
user: "1001:1001" user: "1001:1001"
restart: unless-stopped restart: unless-stopped
command: ["archive-client", "--config", "/etc/archive-control/client.toml"] command: ["archive-client", "--config", "/etc/archive-control/client.toml"]
+1 -1
View File
@@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta"
[project] [project]
name = "archive-clients" name = "archive-clients"
version = "0.1.5" version = "0.1.12"
requires-python = ">=3.11" requires-python = ">=3.11"
dependencies = ["protobuf==7.35.1", "websockets==16.0"] dependencies = ["protobuf==7.35.1", "websockets==16.0"]
+60 -9
View File
@@ -27,17 +27,28 @@ _ENV = re.compile(r"\$\{([A-Za-z_][A-Za-z0-9_]*)\}")
class RootMapping: class RootMapping:
api_root: PurePosixPath api_root: PurePosixPath
local_root: Path local_root: Path
local_path_overrides: tuple[tuple[PurePosixPath, Path], ...] = ()
def api_to_local(self, api_path: str) -> Path: def api_to_local(self, api_path: str) -> Path:
candidate = PurePosixPath(api_path) candidate = PurePosixPath(api_path)
try: for api_root, local_root in self.local_path_overrides:
relative = candidate.relative_to(self.api_root) relative = _safe_relative(candidate, api_root)
except ValueError as exc: if relative is not None:
raise ConfigError("API path is outside its configured root") from exc return local_root.joinpath(*relative.parts)
if any(part in {"", ".", ".."} for part in relative.parts): relative = _safe_relative(candidate, self.api_root)
raise ConfigError("API path contains an unsafe component") if relative is None:
raise ConfigError("API path is outside its configured root")
return self.local_root.joinpath(*relative.parts) return self.local_root.joinpath(*relative.parts)
def local_root_for_api(self, api_path: str) -> Path:
candidate = PurePosixPath(api_path)
for api_root, local_root in self.local_path_overrides:
if _safe_relative(candidate, api_root) is not None:
return local_root
if _safe_relative(candidate, self.api_root) is None:
raise ConfigError("API path is outside its configured root")
return self.local_root
@dataclass(frozen=True) @dataclass(frozen=True)
class ServiceConfig: class ServiceConfig:
@@ -48,10 +59,13 @@ class ServiceConfig:
password_file: Path | None = None password_file: Path | None = None
api_key_file: Path | None = None api_key_file: Path | None = None
advertised_addresses: tuple[str, ...] = () advertised_addresses: tuple[str, ...] = ()
local_path_overrides: tuple[tuple[PurePosixPath, Path], ...] = ()
@property @property
def roots(self) -> RootMapping: def roots(self) -> RootMapping:
return RootMapping(self.api_root, self.local_root) return RootMapping(
self.api_root, self.local_root, self.local_path_overrides
)
def read_password(self) -> str | None: def read_password(self) -> str | None:
return ( return (
@@ -160,7 +174,7 @@ def _service(value: Any, name: str) -> ServiceConfig:
raise ConfigError(f"{name} must be a table") raise ConfigError(f"{name} must be a table")
allowed = { allowed = {
"endpoint", "api_root", "local_root", "username", "password_file", "endpoint", "api_root", "local_root", "username", "password_file",
"api_key_file", "advertised_addresses", "api_key_file", "advertised_addresses", "local_path_overrides",
} }
_keys(value, allowed, name) _keys(value, allowed, name)
api_root = PurePosixPath(_string(value, "api_root")) api_root = PurePosixPath(_string(value, "api_root"))
@@ -182,6 +196,7 @@ def _service(value: Any, name: str) -> ServiceConfig:
raise ConfigError("qbittorrent username and password_file are required") raise ConfigError("qbittorrent username and password_file are required")
if name == "syncthing" and "api_key_file" not in value: if name == "syncthing" and "api_key_file" not in value:
raise ConfigError("syncthing api_key_file is required") raise ConfigError("syncthing api_key_file is required")
overrides = _local_path_overrides(value, api_root, name)
return ServiceConfig( return ServiceConfig(
_endpoint(value, "endpoint", {"http", "https"}), api_root, _endpoint(value, "endpoint", {"http", "https"}), api_root,
_absolute_path(value, "local_root"), username, _absolute_path(value, "local_root"), username,
@@ -189,10 +204,46 @@ def _service(value: Any, name: str) -> ServiceConfig:
if "password_file" in value else None, if "password_file" in value else None,
_absolute_path(value, "api_key_file") _absolute_path(value, "api_key_file")
if "api_key_file" in value else None, if "api_key_file" in value else None,
tuple(addresses), tuple(addresses), overrides,
) )
def _local_path_overrides(
value: dict[str, Any], api_root: PurePosixPath, name: str
) -> tuple[tuple[PurePosixPath, Path], ...]:
raw = value.get("local_path_overrides", {})
if name != "syncthing" and raw:
raise ConfigError(f"{name}.local_path_overrides is unsupported")
if not isinstance(raw, dict):
raise ConfigError(f"{name}.local_path_overrides must be a table")
parsed: list[tuple[PurePosixPath, Path]] = []
for raw_api_path, raw_local_path in raw.items():
if not isinstance(raw_api_path, str) or not isinstance(raw_local_path, str):
raise ConfigError(f"{name}.local_path_overrides entries must be strings")
candidate = PurePosixPath(raw_api_path)
if not candidate.is_absolute() or ".." in candidate.parts:
raise ConfigError(f"{name}.local_path_overrides API path is invalid")
if _safe_relative(candidate, api_root) is None:
raise ConfigError(f"{name}.local_path_overrides API path is outside root")
local = Path(raw_local_path)
if not local.is_absolute():
raise ConfigError(f"{name}.local_path_overrides local path is invalid")
parsed.append((candidate, local))
return tuple(sorted(parsed, key=lambda item: len(item[0].parts), reverse=True))
def _safe_relative(
candidate: PurePosixPath, root: PurePosixPath
) -> PurePosixPath | None:
try:
relative = candidate.relative_to(root)
except ValueError:
return None
if any(part in {"", ".", ".."} for part in relative.parts):
raise ConfigError("API path contains an unsafe component")
return relative
def _connection(value: Any) -> ConnectionConfig: def _connection(value: Any) -> ConnectionConfig:
if not isinstance(value, dict): if not isinstance(value, dict):
raise ConfigError("connection must be a table") raise ConfigError("connection must be a table")
+118 -29
View File
@@ -4,6 +4,8 @@ from __future__ import annotations
import hashlib import hashlib
import json import json
import logging
import os
import shutil import shutil
import stat import stat
import threading import threading
@@ -26,10 +28,11 @@ from archive_clients.qbittorrent import (
) )
from archive_clients.resources import NormalizedResource from archive_clients.resources import NormalizedResource
from archive_clients.state import ClientStore from archive_clients.state import ClientStore
from archive_clients.syncthing import SyncthingTransferObserver from archive_clients.syncthing import RouteSetupError, SyncthingTransferObserver
from archive_clients.transfer import ( from archive_clients.transfer import (
TransferError, TransferError,
TransferIntegrityError, TransferIntegrityError,
cleanup_partial_transfer,
cleanup_transfer, cleanup_transfer,
load_published_transfer, load_published_transfer,
materialize_transfer, materialize_transfer,
@@ -52,6 +55,9 @@ class JobCancelled(JobExecutionError):
pass pass
logger = logging.getLogger(__name__)
class ClientJobExecutor: class ClientJobExecutor:
def __init__( def __init__(
self, self,
@@ -114,7 +120,11 @@ class ClientJobExecutor:
replay = self._replay( replay = self._replay(
command.job_id, command.expected_last_event_sequence command.job_id, command.expected_last_event_sequence
) )
if replay and replay[-1].type in { step_replay = [
event for event in replay
if event.progress.step == command.step
]
if step_replay and step_replay[-1].type in {
control_pb2.JOB_EVENT_TYPE_STEP_SUCCEEDED, control_pb2.JOB_EVENT_TYPE_STEP_SUCCEEDED,
control_pb2.JOB_EVENT_TYPE_COMMITTED, control_pb2.JOB_EVENT_TYPE_COMMITTED,
control_pb2.JOB_EVENT_TYPE_SUCCEEDED, control_pb2.JOB_EVENT_TYPE_SUCCEEDED,
@@ -126,32 +136,51 @@ class ClientJobExecutor:
for event in replay: for event in replay:
event_callback(event) event_callback(event)
return replay return replay
started = replay[0] if replay else self._event( if step_replay:
definition, started = step_replay[0]
sequence=command.expected_last_event_sequence + 1, cursor = step_replay[-1]
revision=command.expected_job_revision + 1, emitted = list(step_replay)
event_type=control_pb2.JOB_EVENT_TYPE_STEP_STARTED, else:
state=job_pb2.JOB_STATE_RUNNING, previous = replay[-1] if replay else None
committed=( started = self._event(
self._committed(command.job_id) definition,
or command.step == job_pb2.JOB_STEP_KIND_STAGING_CLEANUP sequence=(
or ( previous.sequence + 1
command.step == job_pb2.JOB_STEP_KIND_SAFE_FILE_UNLINK if previous is not None
and self.store.get_job_artifact( else command.expected_last_event_sequence + 1
command.job_id, "qb-entry-removed" ),
) is not None revision=(
) previous.job_revision + 1
), if previous is not None
step=command.step, else command.expected_job_revision + 1
step_state=job_pb2.STEP_STATE_RUNNING, ),
) event_type=control_pb2.JOB_EVENT_TYPE_STEP_STARTED,
if not replay: state=job_pb2.JOB_STATE_RUNNING,
committed=(
self._committed(command.job_id)
or command.step == job_pb2.JOB_STEP_KIND_STAGING_CLEANUP
or (
command.step == job_pb2.JOB_STEP_KIND_SAFE_FILE_UNLINK
and self.store.get_job_artifact(
command.job_id, "qb-entry-removed"
) is not None
)
),
step=command.step,
step_state=job_pb2.STEP_STATE_RUNNING,
)
self._record(definition, started) self._record(definition, started)
cursor = started
emitted = [started]
if event_callback is not None: if event_callback is not None:
event_callback(started) for event in emitted:
emitted = [started] event_callback(event)
cursor = started # A newly-started step must publish its first observation. In
speed_sample = [time.monotonic(), 0] # particular, a sparse Syncthing temporary file may retain the same
# allocated-block count for a long time while data is still needed;
# suppressing that first observation made a healthy transfer look
# permanently stalled to the control daemon.
speed_sample: list[float | int | None] = [None, 0]
def progress( def progress(
fraction: float, fraction: float,
@@ -161,8 +190,12 @@ class ClientJobExecutor:
) -> None: ) -> None:
nonlocal cursor nonlocal cursor
now = time.monotonic() now = time.monotonic()
elapsed = now - float(speed_sample[0]) previous_time = speed_sample[0]
if fraction < 1 and elapsed < 1: elapsed = (
now - float(previous_time)
if previous_time is not None else 0.0
)
if previous_time is not None and fraction < 1 and elapsed < 1:
return return
speed = ( speed = (
max(0, bytes_complete - int(speed_sample[1])) / elapsed max(0, bytes_complete - int(speed_sample[1])) / elapsed
@@ -578,7 +611,17 @@ class ClientJobExecutor:
f"staged {completed} of {total} bytes", f"staged {completed} of {total} bytes",
), ),
) )
self._observer(definition).rescan() try:
self._observer(definition).rescan()
except RouteSetupError as exc:
# A manual Syncthing scan may take longer than the bounded REST
# request timeout for a large newly linked file. The folder watch
# and the following transfer observation still converge, so do
# not roll back an otherwise durable staged transfer.
logger.warning("syncthing_rescan_deferred", extra={
"job_id": definition.job_id,
"error_type": type(exc).__name__,
})
def _wait_for_syncthing( def _wait_for_syncthing(
self, self,
@@ -818,6 +861,11 @@ class ClientJobExecutor:
qb_root=job_directory, qb_root=job_directory,
store=self.store, store=self.store,
) )
cleanup_partial_transfer(
job_directory,
job_id=definition.job_id,
store=self.store,
)
for name in ("ready.json", "manifest.json"): for name in ("ready.json", "manifest.json"):
candidate = job_directory / name candidate = job_directory / name
if candidate.is_file() and not candidate.is_symlink(): if candidate.is_file() and not candidate.is_symlink():
@@ -893,6 +941,7 @@ class ClientJobExecutor:
if ( if (
not stat.S_ISREG(metadata.st_mode) not stat.S_ISREG(metadata.st_mode)
or metadata.st_dev != destination_device or metadata.st_dev != destination_device
or _mount_id(source) != _mount_id(destination_root)
): ):
required += logical_bytes required += logical_bytes
return required return required
@@ -1143,6 +1192,8 @@ def _job_error_code(error: Exception) -> int:
return common_pb2.ERROR_CODE_INTEGRITY_CHECK_FAILED return common_pb2.ERROR_CODE_INTEGRITY_CHECK_FAILED
if isinstance(error, PermissionError): if isinstance(error, PermissionError):
return common_pb2.ERROR_CODE_PERMISSION_DENIED return common_pb2.ERROR_CODE_PERMISSION_DENIED
if isinstance(error, RouteSetupError):
return common_pb2.ERROR_CODE_UNAVAILABLE
if isinstance(error, JobExecutionError): if isinstance(error, JobExecutionError):
return common_pb2.ERROR_CODE_PRECONDITION_FAILED return common_pb2.ERROR_CODE_PRECONDITION_FAILED
if isinstance(error, TransferError): if isinstance(error, TransferError):
@@ -1150,3 +1201,41 @@ def _job_error_code(error: Exception) -> int:
if isinstance(error, EvictionError): if isinstance(error, EvictionError):
return common_pb2.ERROR_CODE_PRECONDITION_FAILED return common_pb2.ERROR_CODE_PRECONDITION_FAILED
return common_pb2.ERROR_CODE_INTERNAL return common_pb2.ERROR_CODE_INTERNAL
def _mount_id(path: Path) -> str | None:
"""Return Linux's effective mount ID for a path when procfs is available.
Bind mounts can share ``st_dev`` while still rejecting ``link(2)`` with
``EXDEV``. Mount IDs distinguish that case without creating probe files
inside a Syncthing folder.
"""
try:
target = os.path.realpath(path)
best: tuple[int, str] | None = None
with open("/proc/self/mountinfo", encoding="utf-8") as source:
for line in source:
fields = line.rstrip("\n").split(" ")
if len(fields) < 5:
continue
mountpoint = _unescape_mount_path(fields[4])
if target != mountpoint and not target.startswith(
mountpoint.rstrip("/") + "/"
):
continue
candidate = (len(mountpoint), fields[0])
if best is None or candidate[0] > best[0]:
best = candidate
return None if best is None else best[1]
except OSError:
return None
def _unescape_mount_path(value: str) -> str:
return (
value.replace("\\040", " ")
.replace("\\011", "\t")
.replace("\\012", "\n")
.replace("\\134", "\\")
)
+3 -1
View File
@@ -37,7 +37,9 @@ def discover_routes(
relative = normalized_api_path.relative_to(roots.api_root) relative = normalized_api_path.relative_to(roots.api_root)
except (ConfigError, ValueError): except (ConfigError, ValueError):
continue continue
if not _lexically_within(local_path, roots.local_root): if not _lexically_within(
local_path, roots.local_root_for_api(normalized_api_path.as_posix())
):
continue continue
if relative == PurePosixPath("."): if relative == PurePosixPath("."):
continue continue
+83 -17
View File
@@ -139,19 +139,14 @@ class SyncthingTransferObserver:
def status(self) -> SyncthingTransferStatus: def status(self) -> SyncthingTransferStatus:
published = load_published_transfer(self.local_job_directory) published = load_published_transfer(self.local_job_directory)
total = 0 declared_files = [
for entry in published.manifest.files: (entry.payload_relative_path, entry.logical_bytes)
total += _verified_job_file_size( for entry in published.manifest.files
self.local_job_directory, ] + [
entry.payload_relative_path, (artifact.payload_relative_path, artifact.logical_bytes)
entry.logical_bytes, for artifact in published.manifest.artifacts
) ]
for artifact in published.manifest.artifacts: total = sum(size for _, size in declared_files)
total += _verified_job_file_size(
self.local_job_directory,
artifact.payload_relative_path,
artifact.logical_bytes,
)
completion = self.transport.get_json( completion = self.transport.get_json(
"/rest/db/completion?" "/rest/db/completion?"
@@ -178,11 +173,44 @@ class SyncthingTransferObserver:
name for name in needed_names name for name in needed_names
if name == self.job_relative_path or name.startswith(prefix) if name == self.job_relative_path or name.startswith(prefix)
} }
fraction = float(raw_fraction) / 100 complete = not relevant
complete = fraction == 1 and not relevant if complete:
for relative_path, expected_bytes in declared_files:
_verified_job_file_size(
self.local_job_directory,
relative_path,
expected_bytes,
)
completed_bytes = total
else:
observed_bytes = sum(
_received_job_file_bytes(
self.local_job_directory,
relative_path,
expected_bytes,
)
for relative_path, expected_bytes in declared_files
)
observed_fraction = observed_bytes / total if total else 1.0
# ``/db/completion`` is folder-wide. It can be near zero when a
# route contains a freshly-created item even though the tracked
# job's temporary payload has already received many blocks. For
# an incomplete temporary file, its allocated blocks are the only
# job-specific signal, so do not cap them with that unrelated
# folder aggregate. Once every declared file is atomically
# present, retain the completion value as a conservative guard
# until the need queue has caught up.
fraction = (
observed_fraction
if observed_fraction < 1.0
else min(1.0, float(raw_fraction) / 100)
)
completed_bytes = int(total * fraction)
if complete:
fraction = 1.0
return SyncthingTransferStatus( return SyncthingTransferStatus(
fraction, fraction,
total if complete else int(total * fraction), completed_bytes,
total, total,
complete, complete,
len(relevant), len(relevant),
@@ -356,7 +384,9 @@ class SyncthingRouteManager:
local_path = self.config.roots.api_to_local(api_path) local_path = self.config.roots.api_to_local(api_path)
except ConfigError as exc: except ConfigError as exc:
raise RoutePathConflict("route path is outside the sync root") from exc raise RoutePathConflict("route path is outside the sync root") from exc
resolved_root = self.config.local_root.resolve(strict=False) resolved_root = self.config.roots.local_root_for_api(api_path).resolve(
strict=False
)
try: try:
local_path.resolve(strict=False).relative_to(resolved_root) local_path.resolve(strict=False).relative_to(resolved_root)
except ValueError as exc: except ValueError as exc:
@@ -470,6 +500,42 @@ def _verified_job_file_size(
return metadata.st_size return metadata.st_size
def _received_job_file_bytes(
job_directory: Path,
relative_path: str,
expected_bytes: int,
) -> int:
"""Return a conservative receive estimate for one job-owned file.
Syncthing writes incomplete files as ``.syncthing.<name>.tmp`` and may
pre-size that sparse temporary to its final logical length. Allocated
blocks, rather than ``st_size``, therefore provide the useful progress
signal until the final atomic rename occurs.
"""
relative = PurePosixPath(relative_path)
current = job_directory
for component in relative.parts[:-1]:
current = current / component
final = current / relative.name
try:
metadata = final.lstat()
except FileNotFoundError:
metadata = None
if metadata is not None:
if stat.S_ISREG(metadata.st_mode) and metadata.st_size == expected_bytes:
return expected_bytes
return 0
temporary = current / f".syncthing.{relative.name}.tmp"
try:
temporary_metadata = temporary.lstat()
except FileNotFoundError:
return 0
if not stat.S_ISREG(temporary_metadata.st_mode):
return 0
return min(expected_bytes, temporary_metadata.st_blocks * 512)
def _needed_names(value: dict[str, Any]) -> set[str]: def _needed_names(value: dict[str, Any]) -> set[str]:
result: set[str] = set() result: set[str] = set()
for key in ("progress", "queued", "rest"): for key in ("progress", "queued", "rest"):
+30
View File
@@ -456,6 +456,36 @@ def cleanup_transfer(job_directory: Path) -> bool:
return True return True
def cleanup_partial_transfer(
job_directory: Path,
*,
job_id: str,
store: ClientStore,
) -> None:
"""Remove copy/reflink temporaries left before a transfer is published.
A source-stage failure before ``manifest.json``/``ready.json`` exists
cannot use :func:`cleanup_transfer`. The file-operation journal is the
authoritative list of owned destinations, so derive each temporary path
from it rather than recursively removing arbitrary content from a shared
Syncthing folder.
"""
for row in store.file_operation_rows(job_id):
intent = json.loads(str(row["intent_json"]))
destination = job_directory / _relative_path(str(intent["destination"]))
temporary = _temporary_path(destination, str(row["operation_id"]))
try:
metadata = temporary.lstat()
except FileNotFoundError:
continue
if not stat.S_ISREG(metadata.st_mode):
raise TransferIntegrityError(
"job-owned temporary cleanup path is not a regular file"
)
temporary.unlink()
def canonical_message_json(message: object) -> bytes: def canonical_message_json(message: object) -> bytes:
value = json_format.MessageToDict( value = json_format.MessageToDict(
message, message,
+16 -1
View File
@@ -1,7 +1,7 @@
import os import os
import tempfile import tempfile
import unittest import unittest
from pathlib import Path from pathlib import Path, PurePosixPath
from unittest.mock import patch from unittest.mock import patch
from archive_clients.config import ClientConfig, ConfigError, RootMapping from archive_clients.config import ClientConfig, ConfigError, RootMapping
@@ -63,6 +63,21 @@ class ConfigTests(unittest.TestCase):
with self.assertRaises(ConfigError): with self.assertRaises(ConfigError):
mapping.api_to_local("/elsewhere/file") mapping.api_to_local("/elsewhere/file")
def test_mapping_can_override_one_syncthing_folder_locally(self):
mapping = RootMapping(
PurePosixPath("/sync"),
Path("/local/sync"),
((PurePosixPath("/sync/DownloadsSync"), Path("/local/qb/Sync")),),
)
self.assertEqual(
mapping.api_to_local("/sync/DownloadsSync/job/ready.json"),
Path("/local/qb/Sync/job/ready.json"),
)
self.assertEqual(
mapping.local_root_for_api("/sync/DownloadsSync"),
Path("/local/qb/Sync"),
)
def test_endpoint_scheme_and_job_keys_are_strict(self): def test_endpoint_scheme_and_job_keys_are_strict(self):
with tempfile.TemporaryDirectory() as directory: with tempfile.TemporaryDirectory() as directory:
root = Path(directory) root = Path(directory)
+144 -1
View File
@@ -11,6 +11,7 @@ from archive_clients.jobs import (
JobExecutionError, JobExecutionError,
_resource_fingerprint, _resource_fingerprint,
) )
from archive_clients.syncthing import RouteSetupError
from archive_clients.resources import normalize_resource from archive_clients.resources import normalize_resource
from archive_clients.state import ClientStore from archive_clients.state import ClientStore
from archive_control.v1 import control_pb2, job_pb2 from archive_control.v1 import control_pb2, job_pb2
@@ -31,6 +32,11 @@ class CompleteSyncthing:
self.posts.append(path) self.posts.append(path)
class SlowRescanSyncthing(CompleteSyncthing):
def post(self, path):
raise RouteSetupError("Syncthing API is unavailable")
class ClientJobHappyPathTests(unittest.TestCase): class ClientJobHappyPathTests(unittest.TestCase):
def test_capacity_guard_fails_before_data_movement(self): def test_capacity_guard_fails_before_data_movement(self):
with tempfile.TemporaryDirectory() as directory: with tempfile.TemporaryDirectory() as directory:
@@ -78,6 +84,142 @@ class ClientJobHappyPathTests(unittest.TestCase):
content=b"x" * 4096, content=b"x" * 4096,
) )
def test_mount_boundary_requires_copy_space_even_with_same_device(self):
with tempfile.TemporaryDirectory() as directory:
root = Path(directory)
source_root = root / "source"
destination_root = root / "destination"
source_root.mkdir()
destination_root.mkdir()
(source_root / "fixture.bin").write_bytes(b"fixture")
with patch(
"archive_clients.jobs._mount_id",
side_effect=("source-mount", "destination-mount"),
):
required = ClientJobExecutor._copy_required_bytes(
source_root,
destination_root,
(("fixture.bin", 7),),
)
self.assertEqual(required, 7)
def test_source_stage_survives_a_timed_out_syncthing_rescan(self):
with tempfile.TemporaryDirectory() as directory:
self._run_transfer(
Path(directory),
job_pb2.JOB_OPERATION_ARCHIVE,
syncthing=SlowRescanSyncthing(),
)
def test_target_step_advances_past_replayed_source_completion(self):
with tempfile.TemporaryDirectory() as directory:
root = Path(directory)
store = ClientStore(root / "client.db")
store.initialize()
definition = job_pb2.JobDefinition(
job_id=str(uuid4()),
idempotency_key=str(uuid4()),
operation=job_pb2.JOB_OPERATION_ARCHIVE,
resource_display_name="fixture",
transfer={
"source_client_id": "cache-1",
"target_client_id": "archive-1",
"route_id": "route-1",
},
)
definition.resource_id.info_hash_v1_hex = "a" * 40
definition.created_at.GetCurrentTime()
executor = ClientJobExecutor(
client_id="archive-1",
qbittorrent=Mock(),
store=store,
qb_root=root,
qb_api_root=Path("/downloads"),
route_path=lambda _: root,
syncthing_transport=Mock(),
sparse_supported=True,
poll_interval=0,
)
executor.assign(control_pb2.AssignJobCommand(
job=definition,
expected_job_revision=1,
expected_last_event_sequence=1,
))
source_complete = executor._event(
definition,
sequence=5,
revision=4,
event_type=control_pb2.JOB_EVENT_TYPE_STEP_SUCCEEDED,
state=job_pb2.JOB_STATE_RUNNING,
committed=False,
step=job_pb2.JOB_STEP_KIND_SOURCE_STAGE,
step_state=job_pb2.STEP_STATE_SUCCEEDED,
)
executor._record(definition, source_complete)
with patch.object(executor, "_wait_for_syncthing"):
events = executor.execute(control_pb2.ExecuteStepCommand(
job_id=definition.job_id,
expected_job_revision=4,
expected_last_event_sequence=5,
step=job_pb2.JOB_STEP_KIND_SYNCTHING_TRANSFER,
attempt=1,
))
self.assertEqual(events[0].sequence, 6)
self.assertEqual(
events[0].progress.step,
job_pb2.JOB_STEP_KIND_SYNCTHING_TRANSFER,
)
self.assertEqual(events[-1].sequence, 7)
self.assertEqual(
events[-1].type,
control_pb2.JOB_EVENT_TYPE_STEP_SUCCEEDED,
)
def test_first_partial_progress_is_durable(self):
"""A sparse transfer must not lose its only initial observation."""
with tempfile.TemporaryDirectory() as directory:
root = Path(directory)
store = ClientStore(root / "client.db")
store.initialize()
definition = job_pb2.JobDefinition(
job_id=str(uuid4()),
idempotency_key=str(uuid4()),
operation=job_pb2.JOB_OPERATION_ARCHIVE,
transfer={
"source_client_id": "cache-1",
"target_client_id": "archive-1",
"route_id": "route-1",
},
)
definition.resource_id.info_hash_v1_hex = "a" * 40
definition.created_at.GetCurrentTime()
executor = ClientJobExecutor(
client_id="archive-1", qbittorrent=Mock(), store=store,
qb_root=root, qb_api_root=Path("/downloads"),
route_path=lambda _: root, syncthing_transport=Mock(),
sparse_supported=True, poll_interval=0,
)
executor.assign(control_pb2.AssignJobCommand(
job=definition, expected_job_revision=1,
expected_last_event_sequence=0,
))
def partial_progress(_definition, _step, progress):
progress(0.001, 10, 10_000, "still receiving")
with patch.object(executor, "_execute_step", partial_progress):
events = executor.execute(control_pb2.ExecuteStepCommand(
job_id=definition.job_id, expected_job_revision=1,
expected_last_event_sequence=1,
step=job_pb2.JOB_STEP_KIND_SYNCTHING_TRANSFER, attempt=1,
))
self.assertEqual(len(events), 3)
self.assertEqual(events[1].type, control_pb2.JOB_EVENT_TYPE_PROGRESS)
self.assertEqual(events[1].progress.bytes_complete, 10)
def _run_transfer( def _run_transfer(
self, self,
root: Path, root: Path,
@@ -85,6 +227,7 @@ class ClientJobHappyPathTests(unittest.TestCase):
*, *,
source_stage_free_bytes: int | None = None, source_stage_free_bytes: int | None = None,
content: bytes = b"archive-control-happy-path", content: bytes = b"archive-control-happy-path",
syncthing: CompleteSyncthing | None = None,
): ):
source_root = root / "source" source_root = root / "source"
target_root = root / "target" target_root = root / "target"
@@ -152,7 +295,7 @@ class ClientJobHappyPathTests(unittest.TestCase):
target_qb.get_resource.side_effect = [ target_qb.get_resource.side_effect = [
None, None, resource, resource, None, None, resource, resource,
] ]
syncthing = CompleteSyncthing() syncthing = syncthing or CompleteSyncthing()
source = ClientJobExecutor( source = ClientJobExecutor(
client_id=source_id, client_id=source_id,
qbittorrent=source_qb, qbittorrent=source_qb,
+26
View File
@@ -62,6 +62,32 @@ class RouteDiscoveryTests(unittest.TestCase):
self.assertEqual(routes[0].local_relative_path, "DownloadsSync") self.assertEqual(routes[0].local_relative_path, "DownloadsSync")
self.assertEqual(routes[0].state, route_pb2.ROUTE_STATE_DISCOVERED) self.assertEqual(routes[0].state, route_pb2.ROUTE_STATE_DISCOVERED)
def test_discovery_accepts_a_safe_local_folder_override(self):
with tempfile.TemporaryDirectory() as directory:
root = Path(directory)
override = root / "qb/Sync"
override.mkdir(parents=True)
roots = RootMapping(
PurePosixPath("/var/syncthing"),
root / "sync",
((PurePosixPath("/var/syncthing/DownloadsSync"), override),),
)
routes = discover_routes(
{"folders": [{
"id": "DownloadsSync",
"path": "~/DownloadsSync",
"type": "sendreceive",
"devices": [
{"deviceID": "LOCAL"},
{"deviceID": "ARCHIVE"},
],
}]},
"LOCAL",
roots,
True,
)
self.assertEqual([route.route_id for route in routes], ["DownloadsSync"])
if __name__ == "__main__": if __name__ == "__main__":
unittest.main() unittest.main()
+28
View File
@@ -11,6 +11,7 @@ from archive_clients.transfer import (
FileMaterializer, FileMaterializer,
TransferIntegrityError, TransferIntegrityError,
canonical_message_json, canonical_message_json,
cleanup_partial_transfer,
load_published_transfer, load_published_transfer,
materialize_transfer, materialize_transfer,
stage_transfer, stage_transfer,
@@ -206,6 +207,33 @@ class TransferHappyPathTests(unittest.TestCase):
os.stat(copied).st_blocks * 512, os.stat(copied).st_size os.stat(copied).st_blocks * 512, os.stat(copied).st_size
) )
def test_partial_cleanup_removes_only_journalled_temporary(self):
job_id = str(uuid4())
job_root = self.sync / ".archive-control/jobs" / job_id
destination = job_root / "payload/album/one.bin"
destination.parent.mkdir(parents=True)
operation_id = "interrupted-copy"
self.store.begin_file_operation(
operation_id,
job_id,
'{"destination":"payload/album/one.bin"}',
)
temporary = destination.with_name(
".one.bin.archive-control-"
+ hashlib.sha256(operation_id.encode()).hexdigest()[:16]
+ ".tmp"
)
temporary.write_bytes(b"partial")
unrelated = destination.parent / "keep-me"
unrelated.write_bytes(b"unrelated")
cleanup_partial_transfer(
job_root, job_id=job_id, store=self.store
)
self.assertFalse(temporary.exists())
self.assertTrue(unrelated.is_file())
def test_canonical_manifest_json_is_stable(self): def test_canonical_manifest_json_is_stable(self):
manifest = self._manifest() manifest = self._manifest()
first = canonical_message_json(manifest) first = canonical_message_json(manifest)