Compare commits

...
2 Commits
14 changed files with 104 additions and 22 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.3 image: sodium/archive-clients:v0.1.5
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"]
+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.3 image: sodium/archive-clients:v0.1.5
user: "1001:1001" user: "1001:1001"
restart: unless-stopped restart: unless-stopped
network_mode: host 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. 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.3 image: sodium/archive-clients:v0.1.5
user: "1001:1001" user: "1001:1001"
restart: unless-stopped restart: unless-stopped
network_mode: host 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. 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"
@@ -213,7 +213,7 @@ cache/archive routes according to policy.
```yaml ```yaml
services: services:
archive-client: archive-client:
image: sodium/archive-clients:v0.1.3 image: sodium/archive-clients:v0.1.5
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
@@ -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.3" version = "0.1.5"
requires-python = ">=3.11" requires-python = ">=3.11"
dependencies = ["protobuf==7.35.1", "websockets==16.0"] dependencies = ["protobuf==7.35.1", "websockets==16.0"]
+2 -2
View File
@@ -82,7 +82,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)
@@ -239,7 +239,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",
), ),
) )
+60 -6
View File
@@ -5,11 +5,12 @@ from __future__ import annotations
import hashlib import hashlib
import json import json
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 +66,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 +542,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 +622,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 +867,36 @@ 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
):
required += logical_bytes
return required
def _observer( def _observer(
self, definition: job_pb2.JobDefinition self, definition: job_pb2.JobDefinition
) -> SyncthingTransferObserver: ) -> SyncthingTransferObserver:
+1
View File
@@ -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",)
) )
+27 -4
View File
@@ -69,14 +69,29 @@ 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 _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 +200,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(