Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
a73bcf870f | ||
|
|
748ca49837 |
@@ -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.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"]
|
||||||
|
|||||||
@@ -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.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
|
||||||
|
|||||||
@@ -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.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
|
||||||
|
|||||||
@@ -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"]
|
||||||
|
|||||||
@@ -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.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"]
|
||||||
|
|
||||||
|
|||||||
@@ -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",
|
||||||
),
|
),
|
||||||
)
|
)
|
||||||
|
|||||||
@@ -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:
|
||||||
|
|||||||
@@ -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
@@ -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(
|
||||||
|
|||||||
Reference in New Issue
Block a user