Compare commits

..
4 Commits
9 changed files with 206 additions and 13 deletions
+6 -2
View File
@@ -11,8 +11,12 @@ CMD ["python", "-m", "unittest", "discover", "-s", "tests", "-v"]
FROM python:3.11-slim FROM python:3.11-slim
ENV PYTHONDONTWRITEBYTECODE=1 PYTHONUNBUFFERED=1 ENV PYTHONDONTWRITEBYTECODE=1 PYTHONUNBUFFERED=1
RUN groupadd --gid 1001 archive-control \ RUN if ! getent group 1001 >/dev/null; then \
&& useradd --uid 1001 --gid 1001 --no-create-home archive-control groupadd --gid 1001 archive-control; \
fi \
&& if ! getent passwd 1001 >/dev/null; then \
useradd --uid 1001 --gid 1001 --no-create-home archive-control; \
fi
COPY --from=builder /wheels /wheels COPY --from=builder /wheels /wheels
RUN pip install --no-cache-dir /wheels/*.whl && rm -rf /wheels RUN pip install --no-cache-dir /wheels/*.whl && rm -rf /wheels
USER 1001:1001 USER 1001:1001
+22
View File
@@ -259,6 +259,28 @@ Generated protobuf Python bindings are committed in each consumer with the
exact `archive-control-proto` tag/commit recorded. Release order is proto, exact `archive-control-proto` tag/commit recorded. Release order is proto,
control consumer, client consumer, E2E, then image publication. control consumer, client consumer, E2E, then image publication.
### Reproducible multi-platform Buildx lifecycle
The named Buildx builders are local acceleration/cache only; they are not a
deployment dependency and can be removed after publication. To create a fresh
builder, verify its platforms, publish a release, and remove it afterwards:
```bash
docker buildx create --name archive-control-release --driver docker-container --use
docker buildx inspect --bootstrap
docker buildx build --platform linux/amd64,linux/arm64 \
--tag sodium/archive-clients:vX.Y.Z --push .
docker buildx rm archive-control-release
```
`docker buildx inspect` must list both `linux/amd64` and `linux/arm64` before
publishing. If the host has no arm64 emulation, install/configure it according
to the host Docker distribution before the build; do not publish a partial
single-platform tag. Retain the pushed manifest digest in the release notes
and deploy the immutable tag or digest. The optional `archive-control-qemu`
builder follows the same lifecycle when it is used for an emulation smoke
build.
## Operator usage ## Operator usage
1. Prepare local qB, sync, state, backup, config, and secret mounts with the 1. Prepare local qB, sync, state, backup, config, and secret mounts with the
+1 -1
View File
@@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta"
[project] [project]
name = "archive-clients" name = "archive-clients"
version = "0.1.14" version = "0.1.17"
requires-python = ">=3.11" requires-python = ">=3.11"
dependencies = ["protobuf==7.35.1", "websockets==16.0"] dependencies = ["protobuf==7.35.1", "websockets==16.0"]
+9
View File
@@ -671,6 +671,15 @@ class ClientJobExecutor:
# Syncthing can expose the per-job directory before the # Syncthing can expose the per-job directory before the
# ready marker and manifest have arrived atomically as a set. # ready marker and manifest have arrived atomically as a set.
pass pass
except RouteSetupError as exc:
# A network interruption can temporarily make the local
# Syncthing REST API unavailable. The staged payload is
# durable and the job must remain recoverable; keep polling
# instead of turning a transient outage into a failed job.
logger.warning("syncthing_status_deferred", extra={
"job_id": definition.job_id,
"error_type": type(exc).__name__,
})
time.sleep(self.poll_interval) time.sleep(self.poll_interval)
def _target_materialize( def _target_materialize(
+14 -1
View File
@@ -3,6 +3,7 @@
from __future__ import annotations from __future__ import annotations
import json import json
import logging
import threading import threading
import time import time
import uuid import uuid
@@ -16,6 +17,8 @@ from urllib import error, parse, request
from archive_clients.config import ServiceConfig from archive_clients.config import ServiceConfig
from archive_clients.resources import NormalizedResource, normalize_resource from archive_clients.resources import NormalizedResource, normalize_resource
logger = logging.getLogger(__name__)
class QBittorrentError(RuntimeError): class QBittorrentError(RuntimeError):
pass pass
@@ -72,7 +75,17 @@ class QBittorrentReader:
torrent_hash = torrent.get("hash") torrent_hash = torrent.get("hash")
if not isinstance(torrent_hash, str): if not isinstance(torrent_hash, str):
raise QBittorrentError("qBittorrent torrent hash is invalid") raise QBittorrentError("qBittorrent torrent hash is invalid")
result.append(self._normalize(torrent, torrent_hash)) try:
result.append(self._normalize(torrent, torrent_hash))
except ValueError as exc:
# A stale or malformed qBittorrent entry must not hide every
# otherwise valid resource from Archive/Evict inventory.
# Keep it ineligible and leave an operator-visible diagnosis.
logger.warning(
"Skipping qBittorrent resource with inconsistent "
"metadata: hash=%s name=%r reason=%s",
torrent_hash, torrent.get("name"), exc,
)
return result return result
def get_resource(self, torrent_hash: str) -> NormalizedResource | None: def get_resource(self, torrent_hash: str) -> NormalizedResource | None:
+32 -8
View File
@@ -7,6 +7,8 @@ import json
import os import os
import sqlite3 import sqlite3
import stat import stat
import threading
from contextlib import contextmanager
from dataclasses import dataclass from dataclasses import dataclass
from pathlib import Path from pathlib import Path
@@ -14,6 +16,12 @@ from pathlib import Path
SCHEMA_VERSION = 3 SCHEMA_VERSION = 3
# A client daemon executes unrelated jobs concurrently. SQLite still permits
# only one writer, so serialize this process's short state transactions rather
# than allowing an otherwise healthy job to fail after its busy timeout.
_DATABASE_LOCK = threading.RLock()
class CommandConflict(RuntimeError): class CommandConflict(RuntimeError):
pass pass
@@ -557,14 +565,30 @@ class ClientStore:
).fetchall() ).fetchall()
return [dict(row) for row in rows] return [dict(row) for row in rows]
def _connect(self) -> sqlite3.Connection: @contextmanager
connection = sqlite3.connect(self.database, isolation_level=None, timeout=5) def _connect(self):
connection.row_factory = sqlite3.Row """Yield one connection while serializing local SQLite writers.
connection.execute("PRAGMA foreign_keys = ON")
connection.execute("PRAGMA journal_mode = WAL") The longer SQLite timeout also covers a short lock held by a separate
connection.execute("PRAGMA synchronous = FULL") maintenance process such as a backup. The lock is deliberately held
connection.execute("PRAGMA busy_timeout = 5000") for the full transaction, including ``BEGIN IMMEDIATE``.
return connection """
with _DATABASE_LOCK:
connection = sqlite3.connect(
self.database, isolation_level=None, timeout=30
)
try:
connection.row_factory = sqlite3.Row
connection.execute("PRAGMA foreign_keys = ON")
connection.execute("PRAGMA journal_mode = WAL")
connection.execute("PRAGMA synchronous = FULL")
connection.execute("PRAGMA busy_timeout = 30000")
# Preserve the original ``with connection`` commit/rollback
# behavior used by every store operation.
with connection:
yield connection
finally:
connection.close()
def _canonical(value: object) -> str: def _canonical(value: object) -> str:
+39 -1
View File
@@ -12,7 +12,7 @@ from archive_clients.jobs import (
JobExecutionError, JobExecutionError,
_resource_fingerprint, _resource_fingerprint,
) )
from archive_clients.syncthing import RouteSetupError from archive_clients.syncthing import RouteSetupError, SyncthingTransferStatus
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
@@ -39,6 +39,44 @@ class SlowRescanSyncthing(CompleteSyncthing):
class ClientJobHappyPathTests(unittest.TestCase): class ClientJobHappyPathTests(unittest.TestCase):
def test_syncthing_api_outage_during_transfer_is_retried(self):
"""A transient local REST outage must not terminally fail the job."""
with tempfile.TemporaryDirectory() as directory:
root = Path(directory)
store = ClientStore(root / "client.db")
store.initialize()
definition = job_pb2.JobDefinition(
job_id=str(uuid4()),
transfer={
"source_client_id": "cache-1",
"target_client_id": "archive-1",
"route_id": "route-1",
},
)
observer = Mock()
observer.status.side_effect = [
RouteSetupError("Syncthing API is unavailable"),
SyncthingTransferStatus(1, 42, 42, True, 0),
]
executor = ClientJobExecutor(
client_id="archive-1",
qbittorrent=Mock(),
store=store,
qb_root=root,
qb_api_root=Path("/downloads"),
route_path=lambda _: root,
syncthing_transport=Mock(),
sparse_supported=True,
poll_interval=0,
)
executor._observer = Mock(return_value=observer)
progress = Mock()
executor._wait_for_syncthing(definition, progress)
self.assertEqual(observer.status.call_count, 2)
progress.assert_called_once_with(1, 42, 42, "0 Syncthing items still needed")
def test_reconnect_replay_never_duplicates_any_transfer_step(self): def test_reconnect_replay_never_duplicates_any_transfer_step(self):
for step in ( for step in (
job_pb2.JOB_STEP_KIND_SOURCE_STAGE, job_pb2.JOB_STEP_KIND_SOURCE_STAGE,
+48
View File
@@ -205,6 +205,54 @@ class QBittorrentReaderTests(unittest.TestCase):
self.assertEqual(resources, []) self.assertEqual(resources, [])
self.assertEqual(len(opener.calls), 2) self.assertEqual(len(opener.calls), 2)
def test_listing_skips_malformed_torrent_and_keeps_valid_resources(self):
def torrent(name, content):
info = {
b"length": len(content), b"name": name.encode(),
b"piece length": 16384, b"pieces": b"x" * 20,
}
return encode({b"info": info}), hashlib.sha1(encode(info)).hexdigest()
first_bytes, first_hash = torrent("first.txt", b"one")
malformed_bytes, malformed_hash = torrent("broken.txt", b"bad")
second_bytes, second_hash = torrent("second.txt", b"two")
torrents = [
{"hash": first_hash, "name": "first.txt", "state": "uploading"},
{"hash": malformed_hash, "name": "broken.txt", "state": "stalledUP"},
{"hash": second_hash, "name": "second.txt", "state": "uploading"},
]
responses = [
b"Ok.", json.dumps(torrents).encode(),
b'[{"index":0,"name":"first.txt","size":3,"progress":1,"priority":1}]',
first_bytes,
b'[{"index":0,"name":"broken.txt","size":3,"progress":1,"priority":1},'
b'{"index":1,"name":"unexpected.txt","size":1,"progress":1,"priority":1}]',
malformed_bytes,
b'[{"index":0,"name":"second.txt","size":3,"progress":1,"priority":1}]',
second_bytes,
]
with tempfile.TemporaryDirectory() as directory:
root = Path(directory)
password = root / "password"
password.write_text("secret", encoding="utf-8")
os.chmod(password, 0o600)
config = ServiceConfig(
"http://qb", PurePosixPath("/downloads"), root,
username="admin", password_file=password,
)
opener = _Opener(responses)
with patch(
"archive_clients.qbittorrent.request.build_opener",
return_value=opener,
), self.assertLogs("archive_clients.qbittorrent", "WARNING") as logs:
resources = QBittorrentReader(config).list_resources()
self.assertEqual(
[resource.summary.display_name for resource in resources],
["first.txt", "second.txt"],
)
self.assertIn(malformed_hash, "\n".join(logs.output))
def test_stopped_add_selection_recheck_and_entry_only_delete(self): def test_stopped_add_selection_recheck_and_entry_only_delete(self):
torrent_hash = "a" * 40 torrent_hash = "a" * 40
responses = [ responses = [
+35
View File
@@ -1,6 +1,7 @@
import tempfile import tempfile
import unittest import unittest
import os import os
import threading
from pathlib import Path from pathlib import Path
from uuid import uuid4 from uuid import uuid4
@@ -13,6 +14,40 @@ from archive_clients.state import (
class ClientStoreTests(unittest.TestCase): class ClientStoreTests(unittest.TestCase):
def test_database_connections_serialize_local_writers(self):
"""Concurrent jobs share one daemon DB without SQLite lock failures."""
with tempfile.TemporaryDirectory() as directory:
database = Path(directory) / "state.db"
first = ClientStore(database)
second = ClientStore(database)
first.initialize()
barrier = threading.Barrier(2)
failures: list[Exception] = []
def write(store, command_id):
try:
barrier.wait()
store.accept_command(command_id, '{"kind":"job"}', '{"ok":true}')
except Exception as exc: # pragma: no cover - assertion below
failures.append(exc)
left = threading.Thread(target=write, args=(first, "command-1"))
right = threading.Thread(target=write, args=(second, "command-2"))
left.start()
right.start()
left.join(1)
right.join(1)
self.assertFalse(left.is_alive())
self.assertFalse(right.is_alive())
self.assertEqual(failures, [])
self.assertEqual(len(first.list_accepted_commands()), 2)
with first._connect() as connection:
self.assertEqual(
connection.execute("PRAGMA busy_timeout").fetchone()[0],
30000,
)
def test_command_acceptance_is_durable_and_content_addressed(self): def test_command_acceptance_is_durable_and_content_addressed(self):
with tempfile.TemporaryDirectory() as directory: with tempfile.TemporaryDirectory() as directory:
database = Path(directory) / "state.db" database = Path(directory) / "state.db"