Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
94a3233ce2 | ||
|
|
053b0135b3 |
@@ -2,7 +2,7 @@ name: archive-control-archive
|
||||
|
||||
services:
|
||||
archive-client:
|
||||
image: sodium/archive-clients:v0.1.6
|
||||
image: sodium/archive-clients:v0.1.8
|
||||
user: "1000:1000"
|
||||
restart: unless-stopped
|
||||
command: ["--config", "/etc/archive-control/client.toml"]
|
||||
|
||||
@@ -2,7 +2,7 @@ name: archive-control-cache
|
||||
|
||||
services:
|
||||
archive-client:
|
||||
image: sodium/archive-clients:v0.1.6
|
||||
image: sodium/archive-clients:v0.1.8
|
||||
user: "1001:1001"
|
||||
restart: unless-stopped
|
||||
network_mode: host
|
||||
|
||||
@@ -2,7 +2,7 @@ name: archive-control-cache
|
||||
|
||||
services:
|
||||
archive-client:
|
||||
image: sodium/archive-clients:v0.1.6
|
||||
image: sodium/archive-clients:v0.1.8
|
||||
user: "1001:1001"
|
||||
restart: unless-stopped
|
||||
network_mode: host
|
||||
|
||||
@@ -218,7 +218,7 @@ cache/archive routes according to policy.
|
||||
```yaml
|
||||
services:
|
||||
archive-client:
|
||||
image: sodium/archive-clients:v0.1.6
|
||||
image: sodium/archive-clients:v0.1.8
|
||||
user: "1001:1001"
|
||||
restart: unless-stopped
|
||||
command: ["archive-client", "--config", "/etc/archive-control/client.toml"]
|
||||
|
||||
+1
-1
@@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta"
|
||||
|
||||
[project]
|
||||
name = "archive-clients"
|
||||
version = "0.1.6"
|
||||
version = "0.1.8"
|
||||
requires-python = ">=3.11"
|
||||
dependencies = ["protobuf==7.35.1", "websockets==16.0"]
|
||||
|
||||
|
||||
@@ -4,6 +4,7 @@ from __future__ import annotations
|
||||
|
||||
import hashlib
|
||||
import json
|
||||
import logging
|
||||
import os
|
||||
import shutil
|
||||
import stat
|
||||
@@ -27,10 +28,11 @@ from archive_clients.qbittorrent import (
|
||||
)
|
||||
from archive_clients.resources import NormalizedResource
|
||||
from archive_clients.state import ClientStore
|
||||
from archive_clients.syncthing import SyncthingTransferObserver
|
||||
from archive_clients.syncthing import RouteSetupError, SyncthingTransferObserver
|
||||
from archive_clients.transfer import (
|
||||
TransferError,
|
||||
TransferIntegrityError,
|
||||
cleanup_partial_transfer,
|
||||
cleanup_transfer,
|
||||
load_published_transfer,
|
||||
materialize_transfer,
|
||||
@@ -53,6 +55,9 @@ class JobCancelled(JobExecutionError):
|
||||
pass
|
||||
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
class ClientJobExecutor:
|
||||
def __init__(
|
||||
self,
|
||||
@@ -579,7 +584,17 @@ class ClientJobExecutor:
|
||||
f"staged {completed} of {total} bytes",
|
||||
),
|
||||
)
|
||||
self._observer(definition).rescan()
|
||||
try:
|
||||
self._observer(definition).rescan()
|
||||
except RouteSetupError as exc:
|
||||
# A manual Syncthing scan may take longer than the bounded REST
|
||||
# request timeout for a large newly linked file. The folder watch
|
||||
# and the following transfer observation still converge, so do
|
||||
# not roll back an otherwise durable staged transfer.
|
||||
logger.warning("syncthing_rescan_deferred", extra={
|
||||
"job_id": definition.job_id,
|
||||
"error_type": type(exc).__name__,
|
||||
})
|
||||
|
||||
def _wait_for_syncthing(
|
||||
self,
|
||||
@@ -819,6 +834,11 @@ class ClientJobExecutor:
|
||||
qb_root=job_directory,
|
||||
store=self.store,
|
||||
)
|
||||
cleanup_partial_transfer(
|
||||
job_directory,
|
||||
job_id=definition.job_id,
|
||||
store=self.store,
|
||||
)
|
||||
for name in ("ready.json", "manifest.json"):
|
||||
candidate = job_directory / name
|
||||
if candidate.is_file() and not candidate.is_symlink():
|
||||
@@ -1145,6 +1165,8 @@ def _job_error_code(error: Exception) -> int:
|
||||
return common_pb2.ERROR_CODE_INTEGRITY_CHECK_FAILED
|
||||
if isinstance(error, PermissionError):
|
||||
return common_pb2.ERROR_CODE_PERMISSION_DENIED
|
||||
if isinstance(error, RouteSetupError):
|
||||
return common_pb2.ERROR_CODE_UNAVAILABLE
|
||||
if isinstance(error, JobExecutionError):
|
||||
return common_pb2.ERROR_CODE_PRECONDITION_FAILED
|
||||
if isinstance(error, TransferError):
|
||||
|
||||
@@ -456,6 +456,36 @@ def cleanup_transfer(job_directory: Path) -> bool:
|
||||
return True
|
||||
|
||||
|
||||
def cleanup_partial_transfer(
|
||||
job_directory: Path,
|
||||
*,
|
||||
job_id: str,
|
||||
store: ClientStore,
|
||||
) -> None:
|
||||
"""Remove copy/reflink temporaries left before a transfer is published.
|
||||
|
||||
A source-stage failure before ``manifest.json``/``ready.json`` exists
|
||||
cannot use :func:`cleanup_transfer`. The file-operation journal is the
|
||||
authoritative list of owned destinations, so derive each temporary path
|
||||
from it rather than recursively removing arbitrary content from a shared
|
||||
Syncthing folder.
|
||||
"""
|
||||
|
||||
for row in store.file_operation_rows(job_id):
|
||||
intent = json.loads(str(row["intent_json"]))
|
||||
destination = job_directory / _relative_path(str(intent["destination"]))
|
||||
temporary = _temporary_path(destination, str(row["operation_id"]))
|
||||
try:
|
||||
metadata = temporary.lstat()
|
||||
except FileNotFoundError:
|
||||
continue
|
||||
if not stat.S_ISREG(metadata.st_mode):
|
||||
raise TransferIntegrityError(
|
||||
"job-owned temporary cleanup path is not a regular file"
|
||||
)
|
||||
temporary.unlink()
|
||||
|
||||
|
||||
def canonical_message_json(message: object) -> bytes:
|
||||
value = json_format.MessageToDict(
|
||||
message,
|
||||
|
||||
+16
-1
@@ -11,6 +11,7 @@ from archive_clients.jobs import (
|
||||
JobExecutionError,
|
||||
_resource_fingerprint,
|
||||
)
|
||||
from archive_clients.syncthing import RouteSetupError
|
||||
from archive_clients.resources import normalize_resource
|
||||
from archive_clients.state import ClientStore
|
||||
from archive_control.v1 import control_pb2, job_pb2
|
||||
@@ -31,6 +32,11 @@ class CompleteSyncthing:
|
||||
self.posts.append(path)
|
||||
|
||||
|
||||
class SlowRescanSyncthing(CompleteSyncthing):
|
||||
def post(self, path):
|
||||
raise RouteSetupError("Syncthing API is unavailable")
|
||||
|
||||
|
||||
class ClientJobHappyPathTests(unittest.TestCase):
|
||||
def test_capacity_guard_fails_before_data_movement(self):
|
||||
with tempfile.TemporaryDirectory() as directory:
|
||||
@@ -97,6 +103,14 @@ class ClientJobHappyPathTests(unittest.TestCase):
|
||||
)
|
||||
self.assertEqual(required, 7)
|
||||
|
||||
def test_source_stage_survives_a_timed_out_syncthing_rescan(self):
|
||||
with tempfile.TemporaryDirectory() as directory:
|
||||
self._run_transfer(
|
||||
Path(directory),
|
||||
job_pb2.JOB_OPERATION_ARCHIVE,
|
||||
syncthing=SlowRescanSyncthing(),
|
||||
)
|
||||
|
||||
def _run_transfer(
|
||||
self,
|
||||
root: Path,
|
||||
@@ -104,6 +118,7 @@ class ClientJobHappyPathTests(unittest.TestCase):
|
||||
*,
|
||||
source_stage_free_bytes: int | None = None,
|
||||
content: bytes = b"archive-control-happy-path",
|
||||
syncthing: CompleteSyncthing | None = None,
|
||||
):
|
||||
source_root = root / "source"
|
||||
target_root = root / "target"
|
||||
@@ -171,7 +186,7 @@ class ClientJobHappyPathTests(unittest.TestCase):
|
||||
target_qb.get_resource.side_effect = [
|
||||
None, None, resource, resource,
|
||||
]
|
||||
syncthing = CompleteSyncthing()
|
||||
syncthing = syncthing or CompleteSyncthing()
|
||||
source = ClientJobExecutor(
|
||||
client_id=source_id,
|
||||
qbittorrent=source_qb,
|
||||
|
||||
@@ -11,6 +11,7 @@ from archive_clients.transfer import (
|
||||
FileMaterializer,
|
||||
TransferIntegrityError,
|
||||
canonical_message_json,
|
||||
cleanup_partial_transfer,
|
||||
load_published_transfer,
|
||||
materialize_transfer,
|
||||
stage_transfer,
|
||||
@@ -206,6 +207,33 @@ class TransferHappyPathTests(unittest.TestCase):
|
||||
os.stat(copied).st_blocks * 512, os.stat(copied).st_size
|
||||
)
|
||||
|
||||
def test_partial_cleanup_removes_only_journalled_temporary(self):
|
||||
job_id = str(uuid4())
|
||||
job_root = self.sync / ".archive-control/jobs" / job_id
|
||||
destination = job_root / "payload/album/one.bin"
|
||||
destination.parent.mkdir(parents=True)
|
||||
operation_id = "interrupted-copy"
|
||||
self.store.begin_file_operation(
|
||||
operation_id,
|
||||
job_id,
|
||||
'{"destination":"payload/album/one.bin"}',
|
||||
)
|
||||
temporary = destination.with_name(
|
||||
".one.bin.archive-control-"
|
||||
+ hashlib.sha256(operation_id.encode()).hexdigest()[:16]
|
||||
+ ".tmp"
|
||||
)
|
||||
temporary.write_bytes(b"partial")
|
||||
unrelated = destination.parent / "keep-me"
|
||||
unrelated.write_bytes(b"unrelated")
|
||||
|
||||
cleanup_partial_transfer(
|
||||
job_root, job_id=job_id, store=self.store
|
||||
)
|
||||
|
||||
self.assertFalse(temporary.exists())
|
||||
self.assertTrue(unrelated.is_file())
|
||||
|
||||
def test_canonical_manifest_json_is_stable(self):
|
||||
manifest = self._manifest()
|
||||
first = canonical_message_json(manifest)
|
||||
|
||||
Reference in New Issue
Block a user