fix: defer slow Syncthing rescans
This commit is contained in:
@@ -2,7 +2,7 @@ name: archive-control-archive
|
|||||||
|
|
||||||
services:
|
services:
|
||||||
archive-client:
|
archive-client:
|
||||||
image: sodium/archive-clients:v0.1.6
|
image: sodium/archive-clients:v0.1.7
|
||||||
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"]
|
||||||
|
|||||||
@@ -2,7 +2,7 @@ name: archive-control-cache
|
|||||||
|
|
||||||
services:
|
services:
|
||||||
archive-client:
|
archive-client:
|
||||||
image: sodium/archive-clients:v0.1.6
|
image: sodium/archive-clients:v0.1.7
|
||||||
user: "1001:1001"
|
user: "1001:1001"
|
||||||
restart: unless-stopped
|
restart: unless-stopped
|
||||||
network_mode: host
|
network_mode: host
|
||||||
|
|||||||
@@ -2,7 +2,7 @@ name: archive-control-cache
|
|||||||
|
|
||||||
services:
|
services:
|
||||||
archive-client:
|
archive-client:
|
||||||
image: sodium/archive-clients:v0.1.6
|
image: sodium/archive-clients:v0.1.7
|
||||||
user: "1001:1001"
|
user: "1001:1001"
|
||||||
restart: unless-stopped
|
restart: unless-stopped
|
||||||
network_mode: host
|
network_mode: host
|
||||||
|
|||||||
@@ -218,7 +218,7 @@ cache/archive routes according to policy.
|
|||||||
```yaml
|
```yaml
|
||||||
services:
|
services:
|
||||||
archive-client:
|
archive-client:
|
||||||
image: sodium/archive-clients:v0.1.6
|
image: sodium/archive-clients:v0.1.7
|
||||||
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"]
|
||||||
|
|||||||
+1
-1
@@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta"
|
|||||||
|
|
||||||
[project]
|
[project]
|
||||||
name = "archive-clients"
|
name = "archive-clients"
|
||||||
version = "0.1.6"
|
version = "0.1.7"
|
||||||
requires-python = ">=3.11"
|
requires-python = ">=3.11"
|
||||||
dependencies = ["protobuf==7.35.1", "websockets==16.0"]
|
dependencies = ["protobuf==7.35.1", "websockets==16.0"]
|
||||||
|
|
||||||
|
|||||||
@@ -4,6 +4,7 @@ from __future__ import annotations
|
|||||||
|
|
||||||
import hashlib
|
import hashlib
|
||||||
import json
|
import json
|
||||||
|
import logging
|
||||||
import os
|
import os
|
||||||
import shutil
|
import shutil
|
||||||
import stat
|
import stat
|
||||||
@@ -27,7 +28,7 @@ from archive_clients.qbittorrent import (
|
|||||||
)
|
)
|
||||||
from archive_clients.resources import NormalizedResource
|
from archive_clients.resources import NormalizedResource
|
||||||
from archive_clients.state import ClientStore
|
from archive_clients.state import ClientStore
|
||||||
from archive_clients.syncthing import SyncthingTransferObserver
|
from archive_clients.syncthing import RouteSetupError, SyncthingTransferObserver
|
||||||
from archive_clients.transfer import (
|
from archive_clients.transfer import (
|
||||||
TransferError,
|
TransferError,
|
||||||
TransferIntegrityError,
|
TransferIntegrityError,
|
||||||
@@ -53,6 +54,9 @@ class JobCancelled(JobExecutionError):
|
|||||||
pass
|
pass
|
||||||
|
|
||||||
|
|
||||||
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
|
||||||
class ClientJobExecutor:
|
class ClientJobExecutor:
|
||||||
def __init__(
|
def __init__(
|
||||||
self,
|
self,
|
||||||
@@ -579,7 +583,17 @@ class ClientJobExecutor:
|
|||||||
f"staged {completed} of {total} bytes",
|
f"staged {completed} of {total} bytes",
|
||||||
),
|
),
|
||||||
)
|
)
|
||||||
|
try:
|
||||||
self._observer(definition).rescan()
|
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(
|
def _wait_for_syncthing(
|
||||||
self,
|
self,
|
||||||
@@ -1145,6 +1159,8 @@ def _job_error_code(error: Exception) -> int:
|
|||||||
return common_pb2.ERROR_CODE_INTEGRITY_CHECK_FAILED
|
return common_pb2.ERROR_CODE_INTEGRITY_CHECK_FAILED
|
||||||
if isinstance(error, PermissionError):
|
if isinstance(error, PermissionError):
|
||||||
return common_pb2.ERROR_CODE_PERMISSION_DENIED
|
return common_pb2.ERROR_CODE_PERMISSION_DENIED
|
||||||
|
if isinstance(error, RouteSetupError):
|
||||||
|
return common_pb2.ERROR_CODE_UNAVAILABLE
|
||||||
if isinstance(error, JobExecutionError):
|
if isinstance(error, JobExecutionError):
|
||||||
return common_pb2.ERROR_CODE_PRECONDITION_FAILED
|
return common_pb2.ERROR_CODE_PRECONDITION_FAILED
|
||||||
if isinstance(error, TransferError):
|
if isinstance(error, TransferError):
|
||||||
|
|||||||
+16
-1
@@ -11,6 +11,7 @@ from archive_clients.jobs import (
|
|||||||
JobExecutionError,
|
JobExecutionError,
|
||||||
_resource_fingerprint,
|
_resource_fingerprint,
|
||||||
)
|
)
|
||||||
|
from archive_clients.syncthing import RouteSetupError
|
||||||
from archive_clients.resources import normalize_resource
|
from archive_clients.resources import normalize_resource
|
||||||
from archive_clients.state import ClientStore
|
from archive_clients.state import ClientStore
|
||||||
from archive_control.v1 import control_pb2, job_pb2
|
from archive_control.v1 import control_pb2, job_pb2
|
||||||
@@ -31,6 +32,11 @@ class CompleteSyncthing:
|
|||||||
self.posts.append(path)
|
self.posts.append(path)
|
||||||
|
|
||||||
|
|
||||||
|
class SlowRescanSyncthing(CompleteSyncthing):
|
||||||
|
def post(self, path):
|
||||||
|
raise RouteSetupError("Syncthing API is unavailable")
|
||||||
|
|
||||||
|
|
||||||
class ClientJobHappyPathTests(unittest.TestCase):
|
class ClientJobHappyPathTests(unittest.TestCase):
|
||||||
def test_capacity_guard_fails_before_data_movement(self):
|
def test_capacity_guard_fails_before_data_movement(self):
|
||||||
with tempfile.TemporaryDirectory() as directory:
|
with tempfile.TemporaryDirectory() as directory:
|
||||||
@@ -97,6 +103,14 @@ class ClientJobHappyPathTests(unittest.TestCase):
|
|||||||
)
|
)
|
||||||
self.assertEqual(required, 7)
|
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(
|
def _run_transfer(
|
||||||
self,
|
self,
|
||||||
root: Path,
|
root: Path,
|
||||||
@@ -104,6 +118,7 @@ class ClientJobHappyPathTests(unittest.TestCase):
|
|||||||
*,
|
*,
|
||||||
source_stage_free_bytes: int | None = None,
|
source_stage_free_bytes: int | None = None,
|
||||||
content: bytes = b"archive-control-happy-path",
|
content: bytes = b"archive-control-happy-path",
|
||||||
|
syncthing: CompleteSyncthing | None = None,
|
||||||
):
|
):
|
||||||
source_root = root / "source"
|
source_root = root / "source"
|
||||||
target_root = root / "target"
|
target_root = root / "target"
|
||||||
@@ -171,7 +186,7 @@ class ClientJobHappyPathTests(unittest.TestCase):
|
|||||||
target_qb.get_resource.side_effect = [
|
target_qb.get_resource.side_effect = [
|
||||||
None, None, resource, resource,
|
None, None, resource, resource,
|
||||||
]
|
]
|
||||||
syncthing = CompleteSyncthing()
|
syncthing = syncthing or CompleteSyncthing()
|
||||||
source = ClientJobExecutor(
|
source = ClientJobExecutor(
|
||||||
client_id=source_id,
|
client_id=source_id,
|
||||||
qbittorrent=source_qb,
|
qbittorrent=source_qb,
|
||||||
|
|||||||
Reference in New Issue
Block a user