Compare commits

..
3 Commits
8 changed files with 200 additions and 11 deletions
+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,
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
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]
name = "archive-clients"
version = "0.1.14"
version = "0.1.16"
requires-python = ">=3.11"
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
# ready marker and manifest have arrived atomically as a set.
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)
def _target_materialize(
+13
View File
@@ -3,6 +3,7 @@
from __future__ import annotations
import json
import logging
import threading
import time
import uuid
@@ -16,6 +17,8 @@ from urllib import error, parse, request
from archive_clients.config import ServiceConfig
from archive_clients.resources import NormalizedResource, normalize_resource
logger = logging.getLogger(__name__)
class QBittorrentError(RuntimeError):
pass
@@ -72,7 +75,17 @@ class QBittorrentReader:
torrent_hash = torrent.get("hash")
if not isinstance(torrent_hash, str):
raise QBittorrentError("qBittorrent torrent hash is invalid")
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
def get_resource(self, torrent_hash: str) -> NormalizedResource | None:
+28 -4
View File
@@ -7,6 +7,8 @@ import json
import os
import sqlite3
import stat
import threading
from contextlib import contextmanager
from dataclasses import dataclass
from pathlib import Path
@@ -14,6 +16,12 @@ from pathlib import Path
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):
pass
@@ -557,14 +565,30 @@ class ClientStore:
).fetchall()
return [dict(row) for row in rows]
def _connect(self) -> sqlite3.Connection:
connection = sqlite3.connect(self.database, isolation_level=None, timeout=5)
@contextmanager
def _connect(self):
"""Yield one connection while serializing local SQLite writers.
The longer SQLite timeout also covers a short lock held by a separate
maintenance process such as a backup. The lock is deliberately held
for the full transaction, including ``BEGIN IMMEDIATE``.
"""
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 = 5000")
return connection
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:
+39 -1
View File
@@ -12,7 +12,7 @@ from archive_clients.jobs import (
JobExecutionError,
_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.state import ClientStore
from archive_control.v1 import control_pb2, job_pb2
@@ -39,6 +39,44 @@ class SlowRescanSyncthing(CompleteSyncthing):
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):
for step in (
job_pb2.JOB_STEP_KIND_SOURCE_STAGE,
+48
View File
@@ -205,6 +205,54 @@ class QBittorrentReaderTests(unittest.TestCase):
self.assertEqual(resources, [])
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):
torrent_hash = "a" * 40
responses = [
+35
View File
@@ -1,6 +1,7 @@
import tempfile
import unittest
import os
import threading
from pathlib import Path
from uuid import uuid4
@@ -13,6 +14,40 @@ from archive_clients.state import (
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):
with tempfile.TemporaryDirectory() as directory:
database = Path(directory) / "state.db"