Compare commits

..
4 Commits
19 changed files with 332 additions and 37 deletions
+1 -1
View File
@@ -19,7 +19,7 @@ reconnect_jitter = true
stall_after = "30m" # Status warning only; it does not fail the job. stall_after = "30m" # Status warning only; it does not fail the job.
verification_timeout = "30m" verification_timeout = "30m"
poll_interval = "1s" poll_interval = "1s"
free_space_reserve_bytes = 1073741824 # Rechecked immediately before work. free_space_reserve_bytes = 33554432 # 32 MiB; payload bytes are added for copy paths.
[backup] [backup]
interval = "6h" interval = "6h"
+1 -1
View File
@@ -2,7 +2,7 @@ name: archive-control-archive
services: services:
archive-client: archive-client:
image: sodium/archive-clients:v0.1.2 image: sodium/archive-clients:v0.1.6
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"]
+4 -1
View File
@@ -19,7 +19,7 @@ reconnect_jitter = true
stall_after = "30m" # Status warning only; it does not fail the job. stall_after = "30m" # Status warning only; it does not fail the job.
verification_timeout = "30m" verification_timeout = "30m"
poll_interval = "1s" poll_interval = "1s"
free_space_reserve_bytes = 1073741824 # Rechecked immediately before work. free_space_reserve_bytes = 33554432 # 32 MiB; payload bytes are added for copy paths.
[backup] [backup]
interval = "6h" interval = "6h"
@@ -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.2 image: sodium/archive-clients:v0.1.6
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
@@ -19,7 +19,7 @@ reconnect_jitter = true
stall_after = "30m" # Status warning only; it does not fail the job. stall_after = "30m" # Status warning only; it does not fail the job.
verification_timeout = "30m" verification_timeout = "30m"
poll_interval = "1s" poll_interval = "1s"
free_space_reserve_bytes = 1073741824 # Rechecked immediately before work. free_space_reserve_bytes = 33554432 # 32 MiB; payload bytes are added for copy paths.
[backup] [backup]
interval = "6h" interval = "6h"
+1 -1
View File
@@ -2,7 +2,7 @@ name: archive-control-cache
services: services:
archive-client: archive-client:
image: sodium/archive-clients:v0.1.2 image: sodium/archive-clients:v0.1.6
user: "1001:1001" user: "1001:1001"
restart: unless-stopped restart: unless-stopped
network_mode: host network_mode: host
+7 -2
View File
@@ -102,7 +102,7 @@ stall_after = "30m" # Warning only; no automatic job failure.
verification_timeout = "30m" # Stopped qB full-recheck deadline. verification_timeout = "30m" # Stopped qB full-recheck deadline.
poll_interval = "1s" # Active qB/Syncthing observation interval. poll_interval = "1s" # Active qB/Syncthing observation interval.
# Worst-case copy fallback must leave this many bytes free. # Worst-case copy fallback must leave this many bytes free.
free_space_reserve_bytes = 1073741824 free_space_reserve_bytes = 33554432 # 32 MiB; payload bytes are added for copy paths.
[backup] [backup]
interval = "6h" interval = "6h"
@@ -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.2 image: sodium/archive-clients:v0.1.6
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"]
+4
View File
@@ -29,6 +29,10 @@ The adapter accounts for terminology/behavior changes such as paused versus
stopped states. It does not mutate qBittorrent preferences, categories, tags, stopped states. It does not mutate qBittorrent preferences, categories, tags,
limits, queueing defaults, or global save-path behavior. limits, queueing defaults, or global save-path behavior.
Syncthing route validation normalizes both absolute and `~/` folder-path
spellings under the configured sync root before comparing an existing folder;
an equivalent pre-existing pair is adopted without rewriting it.
### Inventory normalization ### Inventory normalization
For each torrent the client derives canonical v1/v2 identity from reliable API For each torrent the client derives canonical v1/v2 identity from reliable API
+4
View File
@@ -102,6 +102,10 @@ and placement policy differ.
external state rather than being downloaded or silently deselected. external state rather than being downloaded or silently deselected.
- Confirm free space/reserve, permissions, sparse capability, supported - Confirm free space/reserve, permissions, sparse capability, supported
partfile format, route health, and negotiated protocol features. partfile format, route health, and negotiated protocol features.
- Charge payload bytes against free space only when the specific source and
destination files cannot be hardlinked. Same-filesystem hardlink stages and
merges retain the configured reserve plus small metadata artifacts, rather
than reserving a duplicate logical payload.
- Export the source `.torrent` metainfo. - Export the source `.torrent` metainfo.
- Persist the exact source, target baseline, requested selection, and delta. - Persist the exact source, target baseline, requested selection, and delta.
+1 -1
View File
@@ -21,7 +21,7 @@ verification_timeout = "30m"
# Poll qBittorrent and Syncthing at this interval while a step is active. # Poll qBittorrent and Syncthing at this interval while a step is active.
poll_interval = "1s" poll_interval = "1s"
# Space that must remain free after a worst-case copy fallback. # Space that must remain free after a worst-case copy fallback.
free_space_reserve_bytes = 1073741824 free_space_reserve_bytes = 33554432
[backup] [backup]
interval = "6h" interval = "6h"
+1 -1
View File
@@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta"
[project] [project]
name = "archive-clients" name = "archive-clients"
version = "0.1.2" version = "0.1.6"
requires-python = ">=3.11" requires-python = ">=3.11"
dependencies = ["protobuf==7.35.1", "websockets==16.0"] dependencies = ["protobuf==7.35.1", "websockets==16.0"]
+62 -11
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 (
@@ -82,7 +96,7 @@ class JobsConfig:
stall_after: float = 30 * 60 stall_after: float = 30 * 60
verification_timeout: float = 30 * 60 verification_timeout: float = 30 * 60
poll_interval: float = 1 poll_interval: float = 1
free_space_reserve_bytes: int = 1024 * 1024 * 1024 free_space_reserve_bytes: int = 32 * 1024 * 1024
@dataclass(frozen=True) @dataclass(frozen=True)
@@ -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")
@@ -239,7 +290,7 @@ def _jobs(value: Any) -> JobsConfig:
), ),
poll_interval=_duration(value.get("poll_interval", "1s")), poll_interval=_duration(value.get("poll_interval", "1s")),
free_space_reserve_bytes=_positive_int( free_space_reserve_bytes=_positive_int(
value.get("free_space_reserve_bytes", 1024 * 1024 * 1024), value.get("free_space_reserve_bytes", 32 * 1024 * 1024),
"jobs.free_space_reserve_bytes", "jobs.free_space_reserve_bytes",
), ),
) )
+100 -6
View File
@@ -4,12 +4,14 @@ from __future__ import annotations
import hashlib import hashlib
import json import json
import os
import shutil import shutil
import stat
import threading import threading
import time import time
import uuid import uuid
from pathlib import Path, PurePosixPath from pathlib import Path, PurePosixPath
from typing import Callable from typing import Callable, Iterable
from archive_clients.protocol import decode_message, encode_message from archive_clients.protocol import decode_message, encode_message
from archive_clients.eviction import ( from archive_clients.eviction import (
@@ -65,7 +67,7 @@ class ClientJobExecutor:
sparse_supported: bool, sparse_supported: bool,
poll_interval: float = 1, poll_interval: float = 1,
verification_timeout: float = 30 * 60, verification_timeout: float = 30 * 60,
free_space_reserve_bytes: int = 1024 * 1024 * 1024, free_space_reserve_bytes: int = 32 * 1024 * 1024,
): ):
self.client_id = client_id self.client_id = client_id
self.qbittorrent = qbittorrent self.qbittorrent = qbittorrent
@@ -541,15 +543,29 @@ class ClientJobExecutor:
sha256_hex=hashlib.sha256(resource.metainfo_bytes).hexdigest(), sha256_hex=hashlib.sha256(resource.metainfo_bytes).hexdigest(),
) )
del artifact del artifact
route_root = self.route_path(definition.transfer.route_id)
# Staging is normally zero-copy when qB's data root and the paired
# Syncthing route share a filesystem. Do not reserve the complete
# logical payload in that case: FileMaterializer will use link(2),
# which consumes only directory/inode metadata. Retain the metainfo
# allowance and reserve, and account for any source files that really
# must fall back to a data-copy path.
self._require_space( self._require_space(
self.route_path(definition.transfer.route_id), route_root,
definition.transfer.transfer_delta_logical_bytes self._copy_required_bytes(
self.qb_root,
route_root,
(
(entry.target_canonical_path, entry.logical_bytes)
for entry in manifest.files
),
)
+ len(resource.metainfo_bytes), + len(resource.metainfo_bytes),
) )
stage_transfer( stage_transfer(
manifest, manifest,
source_root=self.qb_root, source_root=self.qb_root,
sync_root=self.route_path(definition.transfer.route_id), sync_root=route_root,
store=self.store, store=self.store,
artifact_sources={"metainfo/source.torrent": metainfo_path}, artifact_sources={"metainfo/source.torrent": metainfo_path},
sparse_supported=self.sparse_supported, sparse_supported=self.sparse_supported,
@@ -607,9 +623,18 @@ class ClientJobExecutor:
"target materialization was sent to the wrong client" "target materialization was sent to the wrong client"
) )
published = load_published_transfer(self._job_directory(definition)) published = load_published_transfer(self._job_directory(definition))
# The target can likewise hardlink an arrived Syncthing payload into
# qB's content root when those directories share a filesystem.
self._require_space( self._require_space(
self.qb_root, self.qb_root,
definition.transfer.transfer_delta_logical_bytes, self._copy_required_bytes(
published.job_directory,
self.qb_root,
(
(entry.payload_relative_path, entry.logical_bytes)
for entry in published.manifest.files
),
),
) )
info_hash = _info_hash(definition) info_hash = _info_hash(definition)
resource = self.qbittorrent.get_resource(info_hash) resource = self.qbittorrent.get_resource(info_hash)
@@ -843,6 +868,37 @@ class ClientJobExecutor:
f"{required} bytes required including reserve" f"{required} bytes required including reserve"
) )
@staticmethod
def _copy_required_bytes(
source_root: Path,
destination_root: Path,
files: Iterable[tuple[str, int]],
) -> int:
"""Return logical bytes that cannot be materialized by hardlink.
A hardlink is possible only for regular files on the destination
filesystem. Conservatively charge a file when it cannot be inspected;
the normal materializer will then provide the precise integrity error.
"""
destination_device = destination_root.stat().st_dev
required = 0
for relative_path, logical_bytes in files:
relative = PurePosixPath(relative_path)
source = source_root.joinpath(*relative.parts)
try:
metadata = source.stat(follow_symlinks=False)
except OSError:
required += logical_bytes
continue
if (
not stat.S_ISREG(metadata.st_mode)
or metadata.st_dev != destination_device
or _mount_id(source) != _mount_id(destination_root)
):
required += logical_bytes
return required
def _observer( def _observer(
self, definition: job_pb2.JobDefinition self, definition: job_pb2.JobDefinition
) -> SyncthingTransferObserver: ) -> SyncthingTransferObserver:
@@ -1096,3 +1152,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
+34 -3
View File
@@ -356,7 +356,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:
@@ -369,8 +371,8 @@ class SyncthingRouteManager:
raise RoutePathConflict("new route path contains unrelated data") raise RoutePathConflict("new route path contains unrelated data")
path.mkdir(parents=True, exist_ok=True) path.mkdir(parents=True, exist_ok=True)
@staticmethod
def _validate_folder( def _validate_folder(
self,
folder: dict[str, Any], folder: dict[str, Any],
api_path: str, api_path: str,
local_device_id: str, local_device_id: str,
@@ -382,7 +384,15 @@ class SyncthingRouteManager:
for item in raw_devices for item in raw_devices
if isinstance(item, dict) and isinstance(item.get("deviceID"), str) if isinstance(item, dict) and isinstance(item.get("deviceID"), str)
} if isinstance(raw_devices, list) else set() } if isinstance(raw_devices, list) else set()
if folder.get("path") != api_path: try:
configured_api_path = self._normalized_folder_api_path(
folder.get("path")
)
except RoutePathConflict as exc:
raise RoutePathConflict(
"existing route ID uses a different path"
) from exc
if configured_api_path != api_path:
raise RoutePathConflict("existing route ID uses a different path") raise RoutePathConflict("existing route ID uses a different path")
if folder.get("type") != "sendreceive": if folder.get("type") != "sendreceive":
raise RouteSetupError("existing route folder is not sendreceive") raise RouteSetupError("existing route folder is not sendreceive")
@@ -391,6 +401,27 @@ class SyncthingRouteManager:
if folder.get("paused") is True: if folder.get("paused") is True:
raise RouteSetupError("existing route folder is paused") raise RouteSetupError("existing route folder is paused")
def _normalized_folder_api_path(self, value: Any) -> str:
"""Normalize Syncthing's absolute and home-relative path spellings."""
if not isinstance(value, str) or not value:
raise RoutePathConflict("existing route folder path is invalid")
candidate = PurePosixPath(value)
if candidate.parts and candidate.parts[0] == "~":
candidate = self.config.api_root.joinpath(*candidate.parts[1:])
if not candidate.is_absolute() or any(
part in {"", ".", ".."} for part in candidate.parts
):
raise RoutePathConflict("existing route folder path is unsafe")
normalized = candidate.as_posix()
try:
self.config.roots.api_to_local(normalized)
except ConfigError as exc:
raise RoutePathConflict(
"existing route folder path is outside the sync root"
) from exc
return normalized
def _wait(self, deadline: float) -> None: def _wait(self, deadline: float) -> None:
if time.monotonic() >= deadline: if time.monotonic() >= deadline:
raise RouteSetupTimeout("route setup timed out") raise RouteSetupTimeout("route setup timed out")
+17 -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
@@ -34,6 +34,7 @@ class ConfigTests(unittest.TestCase):
root / "qb/a/b", root / "qb/a/b",
) )
self.assertEqual(config.jobs.stall_after, 30 * 60) self.assertEqual(config.jobs.stall_after, 30 * 60)
self.assertEqual(config.jobs.free_space_reserve_bytes, 32 * 1024 * 1024)
self.assertEqual( self.assertEqual(
config.syncthing.advertised_addresses, ("dynamic",) config.syncthing.advertised_addresses, ("dynamic",)
) )
@@ -62,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)
+46 -4
View File
@@ -69,14 +69,48 @@ class ClientJobHappyPathTests(unittest.TestCase):
with self.subTest(operation=operation), tempfile.TemporaryDirectory() as directory: with self.subTest(operation=operation), tempfile.TemporaryDirectory() as directory:
self._run_transfer(Path(directory), operation) self._run_transfer(Path(directory), operation)
def _run_transfer(self, root: Path, operation: int): def test_same_filesystem_source_stage_does_not_require_payload_space(self):
with tempfile.TemporaryDirectory() as directory:
self._run_transfer(
Path(directory),
job_pb2.JOB_OPERATION_ARCHIVE,
source_stage_free_bytes=32 * 1024 * 1024 + 1024,
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 _run_transfer(
self,
root: Path,
operation: int,
*,
source_stage_free_bytes: int | None = None,
content: bytes = b"archive-control-happy-path",
):
source_root = root / "source" source_root = root / "source"
target_root = root / "target" target_root = root / "target"
route_root = root / "route" route_root = root / "route"
source_root.mkdir() source_root.mkdir()
target_root.mkdir() target_root.mkdir()
route_root.mkdir() route_root.mkdir()
content = b"archive-control-happy-path"
(source_root / "fixture.bin").write_bytes(content) (source_root / "fixture.bin").write_bytes(content)
info = { info = {
b"length": len(content), b"length": len(content),
@@ -185,13 +219,21 @@ class ClientJobHappyPathTests(unittest.TestCase):
) )
final = None final = None
for executor, step in pipeline: for executor, step in pipeline:
events = executor.execute(control_pb2.ExecuteStepCommand( command = control_pb2.ExecuteStepCommand(
job_id=definition.job_id, job_id=definition.job_id,
expected_job_revision=cursor_revision, expected_job_revision=cursor_revision,
expected_last_event_sequence=cursor_sequence, expected_last_event_sequence=cursor_sequence,
step=step, step=step,
attempt=1, attempt=1,
)) )
if executor is source and step == job_pb2.JOB_STEP_KIND_SOURCE_STAGE and source_stage_free_bytes is not None:
with patch(
"archive_clients.jobs.shutil.disk_usage",
return_value=Mock(free=source_stage_free_bytes),
):
events = executor.execute(command)
else:
events = executor.execute(command)
self.assertEqual( self.assertEqual(
[event.sequence for event in events], [event.sequence for event in events],
list(range( list(range(
+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()
+18
View File
@@ -125,6 +125,24 @@ class SyncthingRouteManagerTests(unittest.TestCase):
self.manager.configure(self.spec, time.monotonic() + 1) self.manager.configure(self.spec, time.monotonic() + 1)
self.assertEqual(self.transport.puts, []) self.assertEqual(self.transport.puts, [])
def test_existing_home_relative_folder_path_is_accepted(self):
self.transport.config["devices"].append({"deviceID": "PEER"})
self.transport.config["folders"].append(
{
"id": "route-1",
"path": "~/routes/route-1",
"type": "sendreceive",
"devices": [{"deviceID": "LOCAL"}, {"deviceID": "PEER"}],
}
)
configured = self.manager.configure(
self.spec, time.monotonic() + 1
)
self.assertEqual(self.transport.puts, [])
self.assertFalse(configured.local_route.archive_control_created)
def test_bidirectional_nonce_and_ack_are_required(self): def test_bidirectional_nonce_and_ack_are_required(self):
configured = self.manager.configure(self.spec, time.monotonic() + 1) configured = self.manager.configure(self.spec, time.monotonic() + 1)
peer_nonce = configured.local_path / ( peer_nonce = configured.local_path / (