Compare commits

...
4 Commits
21 changed files with 169 additions and 33 deletions
+1 -1
View File
@@ -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"
+1 -1
View File
@@ -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"]
+1 -1
View File
@@ -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"
+1 -1
View File
@@ -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
+1 -1
View File
@@ -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"
+1 -1
View File
@@ -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
+2 -2
View File
@@ -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"]
+2 -1
View File
@@ -105,7 +105,8 @@ follow-up task.
### Actions
1. Implement qBittorrent cookie authentication and capability adapters for
supported 4.44.6 and 5.x APIs.
supported 4.54.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 -2
View File
@@ -6,7 +6,7 @@ shapes.
## qBittorrent Web API
The implementation targets supported qBittorrent 4.44.6 and 5.x releases and
The implementation targets qBittorrent 4.54.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
View File
@@ -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.54.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.
+4
View File
@@ -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.
+1 -1
View File
@@ -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
View File
@@ -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"]
+2 -2
View File
@@ -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",
),
)
+60 -6
View File
@@ -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:
+3 -1
View File
@@ -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(
+31 -2
View File
@@ -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")
+1
View File
@@ -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
View File
@@ -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(
+3 -3
View File
@@ -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))
+18
View File
@@ -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 / (