From 748ca49837e8881c9d837604c29a78ffded4dbc2 Mon Sep 17 00:00:00 2001 From: Cabbagec Date: Fri, 24 Jul 2026 13:48:24 +0000 Subject: [PATCH] fix: avoid duplicate capacity reservation for hardlinks --- deploy/production/lithium/compose.yaml | 2 +- deploy/production/x1/compose.yaml | 2 +- deploy/production/x2/compose.yaml | 2 +- docs/deployment-and-usage.md | 2 +- docs/workflows.md | 4 ++ pyproject.toml | 2 +- src/archive_clients/jobs.py | 64 ++++++++++++++++++++++++-- tests/test_jobs.py | 31 +++++++++++-- 8 files changed, 95 insertions(+), 14 deletions(-) diff --git a/deploy/production/lithium/compose.yaml b/deploy/production/lithium/compose.yaml index df7d537..0b90a04 100644 --- a/deploy/production/lithium/compose.yaml +++ b/deploy/production/lithium/compose.yaml @@ -2,7 +2,7 @@ name: archive-control-archive services: archive-client: - image: sodium/archive-clients:v0.1.3 + image: sodium/archive-clients:v0.1.4 user: "1000:1000" restart: unless-stopped command: ["--config", "/etc/archive-control/client.toml"] diff --git a/deploy/production/x1/compose.yaml b/deploy/production/x1/compose.yaml index 40e4438..5eae5aa 100644 --- a/deploy/production/x1/compose.yaml +++ b/deploy/production/x1/compose.yaml @@ -2,7 +2,7 @@ name: archive-control-cache services: archive-client: - image: sodium/archive-clients:v0.1.3 + image: sodium/archive-clients:v0.1.4 user: "1001:1001" restart: unless-stopped network_mode: host diff --git a/deploy/production/x2/compose.yaml b/deploy/production/x2/compose.yaml index 14cdedd..41263c8 100644 --- a/deploy/production/x2/compose.yaml +++ b/deploy/production/x2/compose.yaml @@ -2,7 +2,7 @@ name: archive-control-cache services: archive-client: - image: sodium/archive-clients:v0.1.3 + image: sodium/archive-clients:v0.1.4 user: "1001:1001" restart: unless-stopped network_mode: host diff --git a/docs/deployment-and-usage.md b/docs/deployment-and-usage.md index c2608d7..600e5e4 100644 --- a/docs/deployment-and-usage.md +++ b/docs/deployment-and-usage.md @@ -213,7 +213,7 @@ cache/archive routes according to policy. ```yaml services: archive-client: - image: sodium/archive-clients:v0.1.3 + image: sodium/archive-clients:v0.1.4 user: "1001:1001" restart: unless-stopped command: ["archive-client", "--config", "/etc/archive-control/client.toml"] diff --git a/docs/workflows.md b/docs/workflows.md index 375c3c6..67ad77d 100644 --- a/docs/workflows.md +++ b/docs/workflows.md @@ -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. diff --git a/pyproject.toml b/pyproject.toml index c97ce5b..c574c23 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta" [project] name = "archive-clients" -version = "0.1.3" +version = "0.1.4" requires-python = ">=3.11" dependencies = ["protobuf==7.35.1", "websockets==16.0"] diff --git a/src/archive_clients/jobs.py b/src/archive_clients/jobs.py index e5450d5..1674a0d 100644 --- a/src/archive_clients/jobs.py +++ b/src/archive_clients/jobs.py @@ -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 ( @@ -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: diff --git a/tests/test_jobs.py b/tests/test_jobs.py index dbe385e..34a385b 100644 --- a/tests/test_jobs.py +++ b/tests/test_jobs.py @@ -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=1024 * 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(