Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
94a3233ce2 | ||
|
|
053b0135b3 | ||
|
|
57acc90363 | ||
|
|
a73bcf870f | ||
|
|
748ca49837 | ||
|
|
589fd4a41d | ||
|
|
d951e2e4fc |
@@ -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"
|
||||||
|
|||||||
@@ -2,7 +2,7 @@ name: archive-control-archive
|
|||||||
|
|
||||||
services:
|
services:
|
||||||
archive-client:
|
archive-client:
|
||||||
image: sodium/archive-clients:v0.1.1
|
image: sodium/archive-clients:v0.1.8
|
||||||
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"]
|
||||||
|
|||||||
@@ -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"]
|
||||||
|
|||||||
@@ -2,7 +2,7 @@ name: archive-control-cache
|
|||||||
|
|
||||||
services:
|
services:
|
||||||
archive-client:
|
archive-client:
|
||||||
image: sodium/archive-clients:v0.1.1
|
image: sodium/archive-clients:v0.1.8
|
||||||
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
|
|
||||||
|
|||||||
@@ -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"
|
||||||
|
|||||||
@@ -2,7 +2,7 @@ name: archive-control-cache
|
|||||||
|
|
||||||
services:
|
services:
|
||||||
archive-client:
|
archive-client:
|
||||||
image: sodium/archive-clients:v0.1.1
|
image: sodium/archive-clients:v0.1.8
|
||||||
user: "1001:1001"
|
user: "1001:1001"
|
||||||
restart: unless-stopped
|
restart: unless-stopped
|
||||||
network_mode: host
|
network_mode: host
|
||||||
|
|||||||
@@ -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.1
|
image: sodium/archive-clients:v0.1.8
|
||||||
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"]
|
||||||
|
|||||||
@@ -105,7 +105,8 @@ follow-up task.
|
|||||||
### Actions
|
### Actions
|
||||||
|
|
||||||
1. Implement qBittorrent cookie authentication and capability adapters for
|
1. Implement qBittorrent cookie authentication and capability adapters for
|
||||||
supported 4.4–4.6 and 5.x APIs.
|
supported 4.5–4.6 and 5.x APIs; reject older releases because they cannot
|
||||||
|
export exact torrent metainfo through the public API.
|
||||||
2. Implement on-demand torrent summaries, scoped hash lookup, lazy complete
|
2. Implement on-demand torrent summaries, scoped hash lookup, lazy complete
|
||||||
content trees, v1/v2 identity validation, torrent export, stopped add,
|
content trees, v1/v2 identity validation, torrent export, stopped add,
|
||||||
selection application, recheck monitoring, and entry-only deletion.
|
selection application, recheck monitoring, and entry-only deletion.
|
||||||
|
|||||||
@@ -6,7 +6,7 @@ shapes.
|
|||||||
|
|
||||||
## qBittorrent Web API
|
## qBittorrent Web API
|
||||||
|
|
||||||
The implementation targets supported qBittorrent 4.4–4.6 and 5.x releases and
|
The implementation targets qBittorrent 4.5–4.6 and 5.x releases and
|
||||||
detects the application and Web API versions at startup. The authoritative
|
detects the application and Web API versions at startup. The authoritative
|
||||||
references are the official [5.0 WebUI API](https://github.com/qbittorrent/qBittorrent/wiki/WebUI-API-%28qBittorrent-5.0%29)
|
references are the official [5.0 WebUI API](https://github.com/qbittorrent/qBittorrent/wiki/WebUI-API-%28qBittorrent-5.0%29)
|
||||||
and [4.1-compatible WebUI API](https://github.com/qbittorrent/qBittorrent/wiki/WebUI-API-%28qBittorrent-4.1%29).
|
and [4.1-compatible WebUI API](https://github.com/qbittorrent/qBittorrent/wiki/WebUI-API-%28qBittorrent-4.1%29).
|
||||||
@@ -16,7 +16,7 @@ and [4.1-compatible WebUI API](https://github.com/qbittorrent/qBittorrent/wiki/W
|
|||||||
| Need | Web API family | Archive Control rule |
|
| Need | Web API family | Archive Control rule |
|
||||||
| --- | --- | --- |
|
| --- | --- | --- |
|
||||||
| Authenticate | `auth/login`, cookie session | Log no credentials/cookies; reauthenticate once on expiry |
|
| Authenticate | `auth/login`, cookie session | Log no credentials/cookies; reauthenticate once on expiry |
|
||||||
| Detect compatibility | `app/version`, `app/webapiVersion`, build info | Advertise exact versions and select adapter |
|
| Detect compatibility | `app/version`, `app/webapiVersion`, build info | Require qBittorrent 4.5.0+ because exact metainfo export is mandatory |
|
||||||
| List/lookup | `torrents/info`, properties | On-demand and hash-scoped where possible |
|
| List/lookup | `torrents/info`, properties | On-demand and hash-scoped where possible |
|
||||||
| Read file state | `torrents/files` | Normalize indices, paths, selected/skipped, size, progress |
|
| Read file state | `torrents/files` | Normalize indices, paths, selected/skipped, size, progress |
|
||||||
| Export metainfo | `torrents/export` | Required before staging; fail if exact metainfo unavailable |
|
| Export metainfo | `torrents/export` | Required before staging; fail if exact metainfo unavailable |
|
||||||
@@ -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
|
||||||
|
|||||||
+2
-2
@@ -89,8 +89,8 @@ overlapping selections, reordered/duplicated events, and partial journals.
|
|||||||
|
|
||||||
### Adapter contract tests
|
### Adapter contract tests
|
||||||
|
|
||||||
Pinned qBittorrent and Syncthing container versions cover every supported API
|
Pinned qBittorrent 4.5–4.6 and 5.x and Syncthing container versions cover every
|
||||||
family. Assertions include stopped add, exact selection, export, full recheck,
|
supported API family. Assertions include stopped add, exact selection, export, full recheck,
|
||||||
download-attempt detection, entry-only delete, folder/device idempotency,
|
download-attempt detection, entry-only delete, folder/device idempotency,
|
||||||
events fallback, need/completion proof, and redacted errors. Partfile adapters
|
events fallback, need/completion proof, and redacted errors. Partfile adapters
|
||||||
are tested only against explicitly supported qBittorrent/libtorrent fixtures.
|
are tested only against explicitly supported qBittorrent/libtorrent fixtures.
|
||||||
|
|||||||
@@ -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.
|
||||||
|
|
||||||
|
|||||||
@@ -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
@@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta"
|
|||||||
|
|
||||||
[project]
|
[project]
|
||||||
name = "archive-clients"
|
name = "archive-clients"
|
||||||
version = "0.1.1"
|
version = "0.1.8"
|
||||||
requires-python = ">=3.11"
|
requires-python = ">=3.11"
|
||||||
dependencies = ["protobuf==7.35.1", "websockets==16.0"]
|
dependencies = ["protobuf==7.35.1", "websockets==16.0"]
|
||||||
|
|
||||||
|
|||||||
@@ -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",
|
||||||
),
|
),
|
||||||
)
|
)
|
||||||
|
|||||||
+124
-8
@@ -4,12 +4,15 @@ from __future__ import annotations
|
|||||||
|
|
||||||
import hashlib
|
import hashlib
|
||||||
import json
|
import json
|
||||||
|
import logging
|
||||||
|
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 (
|
||||||
@@ -25,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,
|
||||||
@@ -51,6 +55,9 @@ class JobCancelled(JobExecutionError):
|
|||||||
pass
|
pass
|
||||||
|
|
||||||
|
|
||||||
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
|
||||||
class ClientJobExecutor:
|
class ClientJobExecutor:
|
||||||
def __init__(
|
def __init__(
|
||||||
self,
|
self,
|
||||||
@@ -65,7 +72,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 +548,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,
|
||||||
@@ -563,7 +584,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,
|
||||||
@@ -607,9 +638,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)
|
||||||
@@ -794,6 +834,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():
|
||||||
@@ -843,6 +888,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:
|
||||||
@@ -1089,6 +1165,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):
|
||||||
@@ -1096,3 +1174,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", "\\")
|
||||||
|
)
|
||||||
|
|||||||
@@ -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
|
||||||
|
|||||||
@@ -159,7 +159,9 @@ def _supported_qb_version(version: str) -> bool:
|
|||||||
if match is None:
|
if match is None:
|
||||||
return False
|
return False
|
||||||
major, minor = int(match.group(1)), int(match.group(2))
|
major, minor = int(match.group(1)), int(match.group(2))
|
||||||
return major == 5 or (major == 4 and minor in {4, 5, 6})
|
# torrents/export, required to preserve exact metainfo before staging,
|
||||||
|
# was introduced in qBittorrent 4.5.0.
|
||||||
|
return major == 5 or (major == 4 and minor in {5, 6})
|
||||||
|
|
||||||
|
|
||||||
def _text_get(
|
def _text_get(
|
||||||
|
|||||||
@@ -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")
|
||||||
|
|||||||
@@ -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,
|
||||||
|
|||||||
+17
-1
@@ -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)
|
||||||
|
|||||||
+62
-5
@@ -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:
|
||||||
@@ -69,14 +75,57 @@ 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 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 _run_transfer(
|
||||||
|
self,
|
||||||
|
root: Path,
|
||||||
|
operation: int,
|
||||||
|
*,
|
||||||
|
source_stage_free_bytes: int | None = None,
|
||||||
|
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"
|
||||||
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),
|
||||||
@@ -137,7 +186,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,
|
||||||
@@ -185,13 +234,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(
|
||||||
|
|||||||
@@ -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()
|
||||||
|
|||||||
@@ -43,11 +43,11 @@ class _Opener:
|
|||||||
|
|
||||||
|
|
||||||
class ServiceProbeTests(unittest.TestCase):
|
class ServiceProbeTests(unittest.TestCase):
|
||||||
def test_supported_qbittorrent_versions_include_live_4_4_api(self):
|
def test_supported_qbittorrent_versions_require_export_capability(self):
|
||||||
for version in ("v4.4.5", "v4.5.5", "v4.6.7", "v5.2.3"):
|
for version in ("v4.5.0", "v4.6.7", "v5.2.3"):
|
||||||
with self.subTest(version=version):
|
with self.subTest(version=version):
|
||||||
self.assertTrue(_supported_qb_version(version))
|
self.assertTrue(_supported_qb_version(version))
|
||||||
for version in ("v4.3.9", "v6.0.0", "invalid"):
|
for version in ("v4.3.9", "v4.4.5", "v6.0.0", "invalid"):
|
||||||
with self.subTest(version=version):
|
with self.subTest(version=version):
|
||||||
self.assertFalse(_supported_qb_version(version))
|
self.assertFalse(_supported_qb_version(version))
|
||||||
|
|
||||||
|
|||||||
@@ -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 / (
|
||||||
|
|||||||
@@ -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)
|
||||||
|
|||||||
Reference in New Issue
Block a user