Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
a73bcf870f | ||
|
|
748ca49837 | ||
|
|
589fd4a41d | ||
|
|
d951e2e4fc |
@@ -19,7 +19,7 @@ reconnect_jitter = true
|
||||
stall_after = "30m" # Status warning only; it does not fail the job.
|
||||
verification_timeout = "30m"
|
||||
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]
|
||||
interval = "6h"
|
||||
|
||||
@@ -2,7 +2,7 @@ name: archive-control-archive
|
||||
|
||||
services:
|
||||
archive-client:
|
||||
image: sodium/archive-clients:v0.1.1
|
||||
image: sodium/archive-clients:v0.1.5
|
||||
user: "1000:1000"
|
||||
restart: unless-stopped
|
||||
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.
|
||||
verification_timeout = "30m"
|
||||
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]
|
||||
interval = "6h"
|
||||
|
||||
@@ -2,7 +2,7 @@ name: archive-control-cache
|
||||
|
||||
services:
|
||||
archive-client:
|
||||
image: sodium/archive-clients:v0.1.1
|
||||
image: sodium/archive-clients:v0.1.5
|
||||
user: "1001:1001"
|
||||
restart: unless-stopped
|
||||
network_mode: host
|
||||
|
||||
@@ -19,7 +19,7 @@ reconnect_jitter = true
|
||||
stall_after = "30m" # Status warning only; it does not fail the job.
|
||||
verification_timeout = "30m"
|
||||
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]
|
||||
interval = "6h"
|
||||
|
||||
@@ -2,7 +2,7 @@ name: archive-control-cache
|
||||
|
||||
services:
|
||||
archive-client:
|
||||
image: sodium/archive-clients:v0.1.1
|
||||
image: sodium/archive-clients:v0.1.5
|
||||
user: "1001:1001"
|
||||
restart: unless-stopped
|
||||
network_mode: host
|
||||
|
||||
@@ -102,7 +102,7 @@ stall_after = "30m" # Warning only; no automatic job failure.
|
||||
verification_timeout = "30m" # Stopped qB full-recheck deadline.
|
||||
poll_interval = "1s" # Active qB/Syncthing observation interval.
|
||||
# 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]
|
||||
interval = "6h"
|
||||
@@ -213,7 +213,7 @@ cache/archive routes according to policy.
|
||||
```yaml
|
||||
services:
|
||||
archive-client:
|
||||
image: sodium/archive-clients:v0.1.1
|
||||
image: sodium/archive-clients:v0.1.5
|
||||
user: "1001:1001"
|
||||
restart: unless-stopped
|
||||
command: ["archive-client", "--config", "/etc/archive-control/client.toml"]
|
||||
|
||||
@@ -105,7 +105,8 @@ follow-up task.
|
||||
### Actions
|
||||
|
||||
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
|
||||
content trees, v1/v2 identity validation, torrent export, stopped add,
|
||||
selection application, recheck monitoring, and entry-only deletion.
|
||||
|
||||
@@ -6,7 +6,7 @@ shapes.
|
||||
|
||||
## 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
|
||||
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).
|
||||
@@ -16,7 +16,7 @@ and [4.1-compatible WebUI API](https://github.com/qbittorrent/qBittorrent/wiki/W
|
||||
| Need | Web API family | Archive Control rule |
|
||||
| --- | --- | --- |
|
||||
| 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 |
|
||||
| Read file state | `torrents/files` | Normalize indices, paths, selected/skipped, size, progress |
|
||||
| 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,
|
||||
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
|
||||
|
||||
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
|
||||
|
||||
Pinned qBittorrent and Syncthing container versions cover every supported API
|
||||
family. Assertions include stopped add, exact selection, export, full recheck,
|
||||
Pinned qBittorrent 4.5–4.6 and 5.x and Syncthing container versions cover every
|
||||
supported API family. Assertions include stopped add, exact selection, export, full recheck,
|
||||
download-attempt detection, entry-only delete, folder/device idempotency,
|
||||
events fallback, need/completion proof, and redacted errors. Partfile adapters
|
||||
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.
|
||||
- Confirm free space/reserve, permissions, sparse capability, supported
|
||||
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.
|
||||
- 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_interval = "1s"
|
||||
# Space that must remain free after a worst-case copy fallback.
|
||||
free_space_reserve_bytes = 1073741824
|
||||
free_space_reserve_bytes = 33554432
|
||||
|
||||
[backup]
|
||||
interval = "6h"
|
||||
|
||||
+1
-1
@@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta"
|
||||
|
||||
[project]
|
||||
name = "archive-clients"
|
||||
version = "0.1.1"
|
||||
version = "0.1.5"
|
||||
requires-python = ">=3.11"
|
||||
dependencies = ["protobuf==7.35.1", "websockets==16.0"]
|
||||
|
||||
|
||||
@@ -82,7 +82,7 @@ class JobsConfig:
|
||||
stall_after: float = 30 * 60
|
||||
verification_timeout: float = 30 * 60
|
||||
poll_interval: float = 1
|
||||
free_space_reserve_bytes: int = 1024 * 1024 * 1024
|
||||
free_space_reserve_bytes: int = 32 * 1024 * 1024
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
@@ -239,7 +239,7 @@ def _jobs(value: Any) -> JobsConfig:
|
||||
),
|
||||
poll_interval=_duration(value.get("poll_interval", "1s")),
|
||||
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",
|
||||
),
|
||||
)
|
||||
|
||||
@@ -5,11 +5,12 @@ from __future__ import annotations
|
||||
import hashlib
|
||||
import json
|
||||
import shutil
|
||||
import stat
|
||||
import threading
|
||||
import time
|
||||
import uuid
|
||||
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.eviction import (
|
||||
@@ -65,7 +66,7 @@ class ClientJobExecutor:
|
||||
sparse_supported: bool,
|
||||
poll_interval: float = 1,
|
||||
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.qbittorrent = qbittorrent
|
||||
@@ -541,15 +542,29 @@ class ClientJobExecutor:
|
||||
sha256_hex=hashlib.sha256(resource.metainfo_bytes).hexdigest(),
|
||||
)
|
||||
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.route_path(definition.transfer.route_id),
|
||||
definition.transfer.transfer_delta_logical_bytes
|
||||
route_root,
|
||||
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),
|
||||
)
|
||||
stage_transfer(
|
||||
manifest,
|
||||
source_root=self.qb_root,
|
||||
sync_root=self.route_path(definition.transfer.route_id),
|
||||
sync_root=route_root,
|
||||
store=self.store,
|
||||
artifact_sources={"metainfo/source.torrent": metainfo_path},
|
||||
sparse_supported=self.sparse_supported,
|
||||
@@ -607,9 +622,18 @@ class ClientJobExecutor:
|
||||
"target materialization was sent to the wrong client"
|
||||
)
|
||||
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.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)
|
||||
resource = self.qbittorrent.get_resource(info_hash)
|
||||
@@ -843,6 +867,36 @@ class ClientJobExecutor:
|
||||
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
|
||||
):
|
||||
required += logical_bytes
|
||||
return required
|
||||
|
||||
def _observer(
|
||||
self, definition: job_pb2.JobDefinition
|
||||
) -> SyncthingTransferObserver:
|
||||
|
||||
@@ -159,7 +159,9 @@ def _supported_qb_version(version: str) -> bool:
|
||||
if match is None:
|
||||
return False
|
||||
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(
|
||||
|
||||
@@ -369,8 +369,8 @@ class SyncthingRouteManager:
|
||||
raise RoutePathConflict("new route path contains unrelated data")
|
||||
path.mkdir(parents=True, exist_ok=True)
|
||||
|
||||
@staticmethod
|
||||
def _validate_folder(
|
||||
self,
|
||||
folder: dict[str, Any],
|
||||
api_path: str,
|
||||
local_device_id: str,
|
||||
@@ -382,7 +382,15 @@ class SyncthingRouteManager:
|
||||
for item in raw_devices
|
||||
if isinstance(item, dict) and isinstance(item.get("deviceID"), str)
|
||||
} 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")
|
||||
if folder.get("type") != "sendreceive":
|
||||
raise RouteSetupError("existing route folder is not sendreceive")
|
||||
@@ -391,6 +399,27 @@ class SyncthingRouteManager:
|
||||
if folder.get("paused") is True:
|
||||
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:
|
||||
if time.monotonic() >= deadline:
|
||||
raise RouteSetupTimeout("route setup timed out")
|
||||
|
||||
@@ -34,6 +34,7 @@ class ConfigTests(unittest.TestCase):
|
||||
root / "qb/a/b",
|
||||
)
|
||||
self.assertEqual(config.jobs.stall_after, 30 * 60)
|
||||
self.assertEqual(config.jobs.free_space_reserve_bytes, 32 * 1024 * 1024)
|
||||
self.assertEqual(
|
||||
config.syncthing.advertised_addresses, ("dynamic",)
|
||||
)
|
||||
|
||||
+27
-4
@@ -69,14 +69,29 @@ class ClientJobHappyPathTests(unittest.TestCase):
|
||||
with self.subTest(operation=operation), tempfile.TemporaryDirectory() as directory:
|
||||
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 _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"
|
||||
target_root = root / "target"
|
||||
route_root = root / "route"
|
||||
source_root.mkdir()
|
||||
target_root.mkdir()
|
||||
route_root.mkdir()
|
||||
content = b"archive-control-happy-path"
|
||||
(source_root / "fixture.bin").write_bytes(content)
|
||||
info = {
|
||||
b"length": len(content),
|
||||
@@ -185,13 +200,21 @@ class ClientJobHappyPathTests(unittest.TestCase):
|
||||
)
|
||||
final = None
|
||||
for executor, step in pipeline:
|
||||
events = executor.execute(control_pb2.ExecuteStepCommand(
|
||||
command = control_pb2.ExecuteStepCommand(
|
||||
job_id=definition.job_id,
|
||||
expected_job_revision=cursor_revision,
|
||||
expected_last_event_sequence=cursor_sequence,
|
||||
step=step,
|
||||
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(
|
||||
[event.sequence for event in events],
|
||||
list(range(
|
||||
|
||||
@@ -43,11 +43,11 @@ class _Opener:
|
||||
|
||||
|
||||
class ServiceProbeTests(unittest.TestCase):
|
||||
def test_supported_qbittorrent_versions_include_live_4_4_api(self):
|
||||
for version in ("v4.4.5", "v4.5.5", "v4.6.7", "v5.2.3"):
|
||||
def test_supported_qbittorrent_versions_require_export_capability(self):
|
||||
for version in ("v4.5.0", "v4.6.7", "v5.2.3"):
|
||||
with self.subTest(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):
|
||||
self.assertFalse(_supported_qb_version(version))
|
||||
|
||||
|
||||
@@ -125,6 +125,24 @@ class SyncthingRouteManagerTests(unittest.TestCase):
|
||||
self.manager.configure(self.spec, time.monotonic() + 1)
|
||||
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):
|
||||
configured = self.manager.configure(self.spec, time.monotonic() + 1)
|
||||
peer_nonce = configured.local_path / (
|
||||
|
||||
Reference in New Issue
Block a user