Compare commits

..
17 Commits
Author SHA1 Message Date
cabbage 515599f2d4 Recognize BitComet padding file names 2026-08-03 07:10:55 +00:00
cabbage b67ca18403 release: archive clients v0.1.18 2026-08-03 07:02:24 +00:00
cabbage 6446846ee8 Accept qBittorrent hidden padding files 2026-08-03 07:01:10 +00:00
cabbage c717aea394 Make release image UID setup base compatible 2026-08-03 06:49:19 +00:00
cabbage ea35ba2758 Skip malformed qBittorrent inventory entries 2026-08-03 06:46:44 +00:00
cabbage 5e5e2f9095 Document Buildx builder lifecycle 2026-07-30 05:36:15 +00:00
cabbage 0b11e2a3a2 Recover transient Syncthing and SQLite job failures 2026-07-28 14:03:30 +00:00
cabbage 1009defc46 release: archive clients v0.1.14 2026-07-27 03:21:15 +00:00
cabbage 8ed8e75447 feat: allow unrestricted independent job concurrency 2026-07-27 03:19:47 +00:00
cabbage 1a06da3984 fix: reconnect after silent control heartbeat loss 2026-07-27 02:45:21 +00:00
cabbage 4ba5e92248 fix: map x2 route override by folder path 2026-07-26 04:34:08 +00:00
cabbage a6da224f28 docs: require hardlink-safe route mount topology 2026-07-26 04:26:23 +00:00
cabbage 90925fd321 fix: use job-specific sparse transfer progress 2026-07-25 03:01:34 +00:00
cabbage 9df2e6991e fix: publish initial sparse transfer progress 2026-07-25 02:50:22 +00:00
cabbage d2fa69a1d6 fix: report partial Syncthing transfer progress 2026-07-25 02:37:30 +00:00
cabbage 0f94388f49 fix: advance target steps past peer events 2026-07-24 23:27:11 +00:00
cabbage 94a3233ce2 fix: clean interrupted staging temporaries 2026-07-24 17:00:10 +00:00
27 changed files with 873 additions and 84 deletions
+6 -2
View File
@@ -11,8 +11,12 @@ CMD ["python", "-m", "unittest", "discover", "-s", "tests", "-v"]
FROM python:3.11-slim
ENV PYTHONDONTWRITEBYTECODE=1 PYTHONUNBUFFERED=1
RUN groupadd --gid 1001 archive-control \
&& useradd --uid 1001 --gid 1001 --no-create-home archive-control
RUN if ! getent group 1001 >/dev/null; then \
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
RUN pip install --no-cache-dir /wheels/*.whl && rm -rf /wheels
USER 1001:1001
+31 -4
View File
@@ -13,10 +13,37 @@ initial x1/x2/lithium topology.
`archive_control_token`, `qb_password`, and `syncthing_api_key`, each a
regular non-empty file with mode `0600`.
The Syncthing mounts intentionally reproduce each instance's `/var/syncthing`
layout, including nested data binds. This lets route discovery and route
provisioning use one safe API-to-local path mapping without altering an
existing Syncthing configuration.
## Hardlink-safe bind-mount topology
For an archive source, the qB content path and every Syncthing route used for
staging must resolve through the **same container mount**. Matching host
filesystem device IDs alone is insufficient: two separate Docker bind mounts
have different mount IDs and `link(2)` may return `EXDEV` across them. The
client deliberately treats that case as copy-only and performs a full payload
free-space check.
When a Syncthing route is physically nested below the qB root, mount the qB
root once and map the exact Syncthing API folder through it:
```yaml
volumes:
- /srv/downloads:/data/qb
- /srv/syncthing-config:/data/sync
```
```toml
[syncthing]
api_root = "/var/syncthing"
local_root = "/data/sync"
local_path_overrides = { "/var/syncthing/Downloads/Sync" = "/data/qb/Sync" }
```
Do **not** additionally mount `/srv/downloads/Sync` at a path beneath
`/data/sync`. The override is the authoritative mapping for that folder and
keeps qB source files and staging destinations in one mount namespace. Use
the folder's normalized API-visible **path** as the override key (for example,
Syncthing `~/Downloads/Sync` becomes `/var/syncthing/Downloads/Sync`); do not
use the folder ID.
Before starting a stack, validate it with:
+1 -1
View File
@@ -2,7 +2,7 @@ name: archive-control-archive
services:
archive-client:
image: sodium/archive-clients:v0.1.7
image: sodium/archive-clients:v0.1.14
user: "1000:1000"
restart: unless-stopped
command: ["--config", "/etc/archive-control/client.toml"]
+1 -1
View File
@@ -2,7 +2,7 @@ name: archive-control-cache
services:
archive-client:
image: sodium/archive-clients:v0.1.7
image: sodium/archive-clients:v0.1.14
user: "1001:1001"
restart: unless-stopped
network_mode: host
+4
View File
@@ -39,4 +39,8 @@ endpoint = "http://127.0.0.1:8384"
api_key_file = "/run/secrets/syncthing_api_key"
api_root = "/var/syncthing"
local_root = "/data/sync"
# `DownloadsSync-X2` is ~/Downloads/Sync on the host, nested below the qB
# root. Resolve it through /data/qb rather than a second nested bind mount so
# source staging can hardlink it.
local_path_overrides = { "/var/syncthing/Downloads/Sync" = "/data/qb/Sync" }
advertised_addresses = ["dynamic"]
+1 -2
View File
@@ -2,7 +2,7 @@ name: archive-control-cache
services:
archive-client:
image: sodium/archive-clients:v0.1.7
image: sodium/archive-clients:v0.1.14
user: "1001:1001"
restart: unless-stopped
network_mode: host
@@ -16,4 +16,3 @@ services:
- ./backups:/var/backups/archive-control
- /home/ubuntu/Downloads:/data/qb
- /home/ubuntu/compose/syncthing/st_home:/data/sync
- /home/ubuntu/Downloads/Sync:/data/sync/Downloads/Sync
+5 -2
View File
@@ -35,8 +35,11 @@ expand these rules but must not contradict them.
manifests, and observed external state.
- SQLite backups use the online backup API. Defaults are every six hours, 12
recent, 14 daily, and eight weekly copies, plus pre/post-migration backups.
- Default concurrency is one active data-moving job per client and per route.
Disjoint node pairs may run concurrently. Queueing is durable FIFO.
- Default concurrency is unrestricted: every eligible queued job is admitted
in one scheduler pass and independent jobs execute concurrently on a client.
Physical qBittorrent, Syncthing, disk, and network capacity are therefore
the natural limit. `enforce_concurrency_limits=true` restores the optional
legacy per-client/per-route gates. Commands for one job remain serialized.
- A queued job owns a per-resource reservation. Cancelling it removes only the
queued record/reservation and never sends cleanup commands.
- Connectivity or transfer stalls wait indefinitely. A configurable 30-minute
+32 -2
View File
@@ -44,6 +44,9 @@ command_max_attempts = 3 # Total sends, including the initial attempt.
stall_after = "30m" # Warning state only; jobs continue waiting.
route_policy = "on_demand" # Alternative: eager_mesh.
route_setup_timeout = "30m"
# Default: allow all eligible jobs; disk/network/service capacity is the limit.
enforce_concurrency_limits = false
# Used only when enforce_concurrency_limits is true.
max_active_per_client = 1
max_active_per_route = 1
max_envelope_bytes = 1048576
@@ -128,7 +131,12 @@ advertised_addresses = ["dynamic"]
When an existing Syncthing folder is physically nested in the qB data root,
use a `local_path_overrides` entry to map that exact Syncthing API path through
the same client bind mount. This enables hardlinks without creating two Docker
mount boundaries for the same host files.
mount boundaries for the same host files. Do not add a second bind mount for
the nested folder: Linux treats it as a distinct mount even when it has the
same `st_dev`, and the client correctly falls back to copy-only capacity
accounting. The override key is the normalized Syncthing folder path beneath
`api_root`, not its folder ID. See the production deployment README for the
required compose and override pattern.
The remaining node examples omit optional `[connection]`, `[jobs]`, and
`[backup]` tables and therefore use these same defaults; deployments may
@@ -218,7 +226,7 @@ cache/archive routes according to policy.
```yaml
services:
archive-client:
image: sodium/archive-clients:v0.1.7
image: sodium/archive-clients:v0.1.14
user: "1001:1001"
restart: unless-stopped
command: ["archive-client", "--config", "/etc/archive-control/client.toml"]
@@ -251,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
+5 -3
View File
@@ -89,9 +89,11 @@ left unchanged.
Syncthing may serialize a folder path relative to its home as `~/...`. The
client normalizes that notation beneath the configured API-visible sync root
before applying the API-to-local root mapping. Deployments must mirror
Syncthing's nested bind mounts into the client so the normalized API path and
the client filesystem path refer to the same bytes.
before applying the API-to-local root mapping. When that folder is nested
below qB's content root, an exact `local_path_overrides` entry must map it
through the qB bind mount. Do not mirror it as a second nested client bind
mount: it denotes the same host bytes but a distinct mount namespace boundary,
which prevents hardlink staging.
Provisioning uses idempotent device and folder configuration updates. Each
client receives its peer device ID and optional advertised addresses
+13 -4
View File
@@ -133,10 +133,13 @@ a resource may be active or queued at a time. This prevents incompatible
baselines even when jobs would use different nodes or observe different sides
of a hybrid identity.
Defaults allow one active data-moving job per client and one per route. A job
must acquire its source client, target client, route, and resource reservation
atomically. Disjoint node pairs may run concurrently. Route setup is a
preflight activity and does not permit a data step to bypass these leases.
By default every eligible queued job is claimed in the same scheduler pass and
independent jobs run concurrently on each client; physical service and storage
capacity are the limit. Set `enforce_concurrency_limits=true` on control to
enable the optional one-per-client/one-per-route gates. A job always retains
its resource reservation, and commands for the same job remain serialized.
Route setup is a preflight activity and does not permit a data step to bypass
the resource reservation.
Offline nodes do not prevent unrelated jobs from being listed or run. A job
requiring an offline node remains waiting indefinitely; it does not consume an
@@ -154,6 +157,12 @@ advances to the next step. On reconnect:
4. Control observes relevant qBittorrent, Syncthing, staging, and manifest
state before selecting retry, resume, compensation, cleanup, or manual
intervention.
If a connection disappears while a synchronous job operation is in progress,
the client journals events locally, reconnects indefinitely, and serializes a
replayed command behind the in-flight operation for that job. The replay reads
the journal rather than repeating the data operation. This applies to every
transfer step, including post-commit staging cleanup.
5. A command is reissued with its original ID when the acceptance result is
uncertain.
+4 -2
View File
@@ -25,8 +25,10 @@ What do you want to do?
[ Evict Cache ] [ Job Status ]
```
All subsequent pages edit this message. `Cancel` closes the active selection
flow. `Back` returns one level while retaining validated filters/selections.
All subsequent pages edit this message. Leaving a resource/tree selection
returns to the Archive Control operation menu without closing the shared
conversation. Confirmation `Cancel` returns to its immediate prior selection;
the main menu alone can close Archive Control.
Callback payloads contain opaque session/action IDs, not resource names or
paths, and are validated against persisted session revision and expiry.
+1 -1
View File
@@ -225,7 +225,7 @@ Fixtures include small deterministic v1, v2, and hybrid torrents with:
| Recovery | restart every step on source/target/control, lost DB with proof, ambiguous loss fails closed |
| Messaging | duplicate command/event, lost ack, sequence gap, stale revision, duplicate client ID |
| Safety | attempted download, corrupt same-size file, symlink/path escape, special file, partfile mismatch |
| Operations | stall retains percent, offline wait, cancellation each phase, postcommit cleanup retry |
| Operations | stall retains percent, offline wait, cancellation each phase, postcommit cleanup retry, disconnect/reconnect during every transfer step with exactly-once replay |
| Database | online backup under load, pre/post migration, retention, corrupt backup rejection, offline restore |
| UI | pagination/filter chain, stale/double callback, persisted session, clear vs evict separation |
+1 -1
View File
@@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta"
[project]
name = "archive-clients"
version = "0.1.7"
version = "0.1.19"
requires-python = ">=3.11"
dependencies = ["protobuf==7.35.1", "websockets==16.0"]
+11 -1
View File
@@ -168,7 +168,7 @@ def _v1_files(info: dict[bytes, Any]) -> list[MetaFile]:
components = [name] + [_component(part) for part in raw_path]
result.append(MetaFile(
"/".join(components), _length(item.get(b"length")),
_padding(item),
_padding(item) or _bitcomet_padding_name(components[-1]),
))
return result
@@ -223,3 +223,13 @@ def _length(value: Any) -> int:
def _padding(value: dict[bytes, Any]) -> bool:
attributes = value.get(b"attr", b"")
return isinstance(attributes, bytes) and b"p" in attributes
def _bitcomet_padding_name(component: str) -> bool:
"""Recognize BitComet's legacy padding-file convention.
Such v1 torrents often omit the standard ``attr=p`` flag, but
qBittorrent/libtorrent still hides these synthetic entries from its file
API. The exact reserved prefix is the interoperable marker.
"""
return component.startswith("_____padding_file_")
+38 -3
View File
@@ -11,6 +11,7 @@ from pathlib import Path, PurePosixPath
from typing import Any
from websockets.asyncio.client import connect
from websockets.exceptions import ConnectionClosedOK
from archive_clients.backup import SQLiteBackupManager
from archive_clients.config import ClientConfig
@@ -42,6 +43,19 @@ from archive_control.v1 import (
logger = logging.getLogger(__name__)
def _job_id_for_command(command: control_pb2.Command) -> str:
"""Return the durable job key for commands whose execution is serialized."""
payload = command.WhichOneof("payload")
if payload == "assign_job":
return command.assign_job.job.job_id
if payload == "execute_step":
return command.execute_step.job_id
if payload == "cancel_job":
return command.cancel_job.job_id
raise ValueError(f"command {command.command_id} does not execute a job")
class ArchiveClientDaemon:
def __init__(
self,
@@ -94,7 +108,10 @@ class ArchiveClientDaemon:
self._lease = DatabaseLease(config.state_db)
self._active_route_commands: set[str] = set()
self._active_job_commands: set[str] = set()
self._job_execution_lock = asyncio.Lock()
# Commands for one job remain ordered locally, while unrelated jobs
# may use the node's available qB/Syncthing/filesystem capacity in
# parallel. The control daemon owns admission policy.
self._job_execution_locks: dict[str, asyncio.Lock] = {}
self.jobs = (
ClientJobExecutor(
client_id=config.client_id,
@@ -205,7 +222,23 @@ class ArchiveClientDaemon:
command_tasks: set[asyncio.Task[None]] = set()
try:
await self._resume_commands(outbound, command_tasks)
async for frame in websocket:
# A proxy can leave the TCP/WebSocket socket apparently open
# after the control server has discarded its session. The
# server then cannot deliver durable commands and its pending
# outbox remains stranded unless the client independently
# detects the missing application heartbeats and reconnects.
while True:
try:
frame = await asyncio.wait_for(
websocket.recv(),
self.config.connection.offline_timeout,
)
except ConnectionClosedOK:
return
except asyncio.TimeoutError as exc:
raise RuntimeError(
"control heartbeat timed out"
) from exc
await self._handle(decode(frame), outbound, command_tasks)
finally:
writer.cancel()
@@ -543,7 +576,9 @@ class ArchiveClientDaemon:
correlation_id: str,
outbound: asyncio.Queue[str],
) -> None:
async with self._job_execution_lock:
job_id = _job_id_for_command(command)
lock = self._job_execution_locks.setdefault(job_id, asyncio.Lock())
async with lock:
await self._execute_job_command_locked(
command, correlation_id, outbound
)
+71 -11
View File
@@ -32,6 +32,7 @@ from archive_clients.syncthing import RouteSetupError, SyncthingTransferObserver
from archive_clients.transfer import (
TransferError,
TransferIntegrityError,
cleanup_partial_transfer,
cleanup_transfer,
load_published_transfer,
materialize_transfer,
@@ -85,6 +86,8 @@ class ClientJobExecutor:
self.verification_timeout = verification_timeout
self.free_space_reserve_bytes = free_space_reserve_bytes
self._cancel_events: dict[str, threading.Event] = {}
self._execution_locks: dict[str, threading.Lock] = {}
self._execution_locks_guard = threading.Lock()
def request_cancel(self, job_id: str) -> None:
self._cancel_events.setdefault(job_id, threading.Event()).set()
@@ -114,12 +117,32 @@ class ClientJobExecutor:
self,
command: control_pb2.ExecuteStepCommand,
event_callback: Callable[[control_pb2.JobEvent], None] | None = None,
) -> list[control_pb2.JobEvent]:
# asyncio cancellation of a connection-bound task cannot stop the
# synchronous filesystem/qB operation already running in its worker
# thread. A replay after reconnect therefore waits for that operation
# and then reads its durable event journal instead of executing twice.
with self._execution_lock(command.job_id):
return self._execute_locked(command, event_callback)
def _execution_lock(self, job_id: str) -> threading.Lock:
with self._execution_locks_guard:
return self._execution_locks.setdefault(job_id, threading.Lock())
def _execute_locked(
self,
command: control_pb2.ExecuteStepCommand,
event_callback: Callable[[control_pb2.JobEvent], None] | None = None,
) -> list[control_pb2.JobEvent]:
definition = self._definition(command.job_id)
replay = self._replay(
command.job_id, command.expected_last_event_sequence
)
if replay and replay[-1].type in {
step_replay = [
event for event in replay
if event.progress.step == command.step
]
if step_replay and step_replay[-1].type in {
control_pb2.JOB_EVENT_TYPE_STEP_SUCCEEDED,
control_pb2.JOB_EVENT_TYPE_COMMITTED,
control_pb2.JOB_EVENT_TYPE_SUCCEEDED,
@@ -131,10 +154,24 @@ class ClientJobExecutor:
for event in replay:
event_callback(event)
return replay
started = replay[0] if replay else self._event(
if step_replay:
started = step_replay[0]
cursor = step_replay[-1]
emitted = list(step_replay)
else:
previous = replay[-1] if replay else None
started = self._event(
definition,
sequence=command.expected_last_event_sequence + 1,
revision=command.expected_job_revision + 1,
sequence=(
previous.sequence + 1
if previous is not None
else command.expected_last_event_sequence + 1
),
revision=(
previous.job_revision + 1
if previous is not None
else command.expected_job_revision + 1
),
event_type=control_pb2.JOB_EVENT_TYPE_STEP_STARTED,
state=job_pb2.JOB_STATE_RUNNING,
committed=(
@@ -150,13 +187,18 @@ class ClientJobExecutor:
step=command.step,
step_state=job_pb2.STEP_STATE_RUNNING,
)
if not replay:
self._record(definition, started)
if event_callback is not None:
event_callback(started)
emitted = [started]
cursor = started
speed_sample = [time.monotonic(), 0]
emitted = [started]
if event_callback is not None:
for event in emitted:
event_callback(event)
# A newly-started step must publish its first observation. In
# particular, a sparse Syncthing temporary file may retain the same
# allocated-block count for a long time while data is still needed;
# suppressing that first observation made a healthy transfer look
# permanently stalled to the control daemon.
speed_sample: list[float | int | None] = [None, 0]
def progress(
fraction: float,
@@ -166,8 +208,12 @@ class ClientJobExecutor:
) -> None:
nonlocal cursor
now = time.monotonic()
elapsed = now - float(speed_sample[0])
if fraction < 1 and elapsed < 1:
previous_time = speed_sample[0]
elapsed = (
now - float(previous_time)
if previous_time is not None else 0.0
)
if previous_time is not None and fraction < 1 and elapsed < 1:
return
speed = (
max(0, bytes_complete - int(speed_sample[1])) / elapsed
@@ -625,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(
@@ -833,6 +888,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():
+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:
+11 -2
View File
@@ -91,12 +91,21 @@ def normalize_resource(
observed_at: datetime | None = None,
) -> NormalizedResource:
metainfo = decode_metainfo(metainfo_bytes)
if len(raw_files) != len(metainfo.files):
# qBittorrent/libtorrent may omit torrent padding files from
# /torrents/files while keeping them in the exported metainfo.
visible_metainfo_files = tuple(
item for item in metainfo.files if not item.padding
)
if len(raw_files) == len(metainfo.files):
matched_metainfo_files = metainfo.files
elif len(raw_files) == len(visible_metainfo_files):
matched_metainfo_files = visible_metainfo_files
else:
raise ResourceError("qBittorrent and metainfo file counts differ")
files = []
canonical = True
for expected_index, (raw, meta_file) in enumerate(
zip(raw_files, metainfo.files, strict=True)
zip(raw_files, matched_metainfo_files, strict=True)
):
index = _integer(raw.get("index"), "file index")
if index != expected_index:
+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:
+80 -16
View File
@@ -139,19 +139,14 @@ class SyncthingTransferObserver:
def status(self) -> SyncthingTransferStatus:
published = load_published_transfer(self.local_job_directory)
total = 0
for entry in published.manifest.files:
total += _verified_job_file_size(
self.local_job_directory,
entry.payload_relative_path,
entry.logical_bytes,
)
for artifact in published.manifest.artifacts:
total += _verified_job_file_size(
self.local_job_directory,
artifact.payload_relative_path,
artifact.logical_bytes,
)
declared_files = [
(entry.payload_relative_path, entry.logical_bytes)
for entry in published.manifest.files
] + [
(artifact.payload_relative_path, artifact.logical_bytes)
for artifact in published.manifest.artifacts
]
total = sum(size for _, size in declared_files)
completion = self.transport.get_json(
"/rest/db/completion?"
@@ -178,11 +173,44 @@ class SyncthingTransferObserver:
name for name in needed_names
if name == self.job_relative_path or name.startswith(prefix)
}
fraction = float(raw_fraction) / 100
complete = fraction == 1 and not relevant
complete = not relevant
if complete:
for relative_path, expected_bytes in declared_files:
_verified_job_file_size(
self.local_job_directory,
relative_path,
expected_bytes,
)
completed_bytes = total
else:
observed_bytes = sum(
_received_job_file_bytes(
self.local_job_directory,
relative_path,
expected_bytes,
)
for relative_path, expected_bytes in declared_files
)
observed_fraction = observed_bytes / total if total else 1.0
# ``/db/completion`` is folder-wide. It can be near zero when a
# route contains a freshly-created item even though the tracked
# job's temporary payload has already received many blocks. For
# an incomplete temporary file, its allocated blocks are the only
# job-specific signal, so do not cap them with that unrelated
# folder aggregate. Once every declared file is atomically
# present, retain the completion value as a conservative guard
# until the need queue has caught up.
fraction = (
observed_fraction
if observed_fraction < 1.0
else min(1.0, float(raw_fraction) / 100)
)
completed_bytes = int(total * fraction)
if complete:
fraction = 1.0
return SyncthingTransferStatus(
fraction,
total if complete else int(total * fraction),
completed_bytes,
total,
complete,
len(relevant),
@@ -472,6 +500,42 @@ def _verified_job_file_size(
return metadata.st_size
def _received_job_file_bytes(
job_directory: Path,
relative_path: str,
expected_bytes: int,
) -> int:
"""Return a conservative receive estimate for one job-owned file.
Syncthing writes incomplete files as ``.syncthing.<name>.tmp`` and may
pre-size that sparse temporary to its final logical length. Allocated
blocks, rather than ``st_size``, therefore provide the useful progress
signal until the final atomic rename occurs.
"""
relative = PurePosixPath(relative_path)
current = job_directory
for component in relative.parts[:-1]:
current = current / component
final = current / relative.name
try:
metadata = final.lstat()
except FileNotFoundError:
metadata = None
if metadata is not None:
if stat.S_ISREG(metadata.st_mode) and metadata.st_size == expected_bytes:
return expected_bytes
return 0
temporary = current / f".syncthing.{relative.name}.tmp"
try:
temporary_metadata = temporary.lstat()
except FileNotFoundError:
return 0
if not stat.S_ISREG(temporary_metadata.st_mode):
return 0
return min(expected_bytes, temporary_metadata.st_blocks * 512)
def _needed_names(value: dict[str, Any]) -> set[str]:
result: set[str] = set()
for key in ("progress", "queued", "rest"):
+30
View File
@@ -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,
+88
View File
@@ -22,6 +22,94 @@ from archive_control.v1 import (
class DaemonTransportTests(unittest.IsolatedAsyncioTestCase):
async def test_unrelated_job_commands_execute_concurrently(self):
with tempfile.TemporaryDirectory() as directory:
root = Path(directory)
token = root / "token"
token.write_text("shared-secret", encoding="utf-8")
os.chmod(token, 0o600)
service = ServiceConfig("http://local", PurePosixPath("/api"), root)
config = ClientConfig(
"cache-1", "Cache 1", "cache", "ws://control", token,
root / "state.db", root / "backups", service, service,
)
probe = FilesystemProbe(root, True, True, True, True, True)
daemon = ArchiveClientDaemon(config, [probe, probe], [])
started: set[str] = set()
release = asyncio.Event()
async def execute(command, correlation_id, outbound):
started.add(command.assign_job.job.job_id)
await release.wait()
daemon._execute_job_command_locked = execute
commands = []
for _ in range(2):
command = control_pb2.Command(command_id=str(uuid4()))
command.assign_job.job.job_id = str(uuid4())
commands.append(command)
outbound = asyncio.Queue()
tasks = [
asyncio.create_task(
daemon._execute_job_command(command, "", outbound)
)
for command in commands
]
for _ in range(100):
if len(started) == 2:
break
await asyncio.sleep(0.01)
self.assertEqual(len(started), 2)
release.set()
await asyncio.gather(*tasks)
async def test_silent_control_connection_ends_for_reconnect(self):
"""A lost server heartbeat must not leave durable commands stranded."""
async def control(websocket):
registration = decode(await websocket.recv())
self.assertEqual(registration.WhichOneof("payload"), "register_request")
response = new_envelope()
response.correlation_id = registration.message_id
response.register_response.status = (
client_pb2.REGISTRATION_STATUS_ACCEPTED
)
response.register_response.negotiated_version.major = 1
await websocket.send(encode(response))
# Deliberately keep TCP/WebSocket open but send no application
# heartbeats. This models a stale proxy/server-side session.
await websocket.wait_closed()
with tempfile.TemporaryDirectory() as directory:
root = Path(directory)
token = root / "token"
token.write_text("shared-secret", encoding="utf-8")
os.chmod(token, 0o600)
async with serve(control, "127.0.0.1", 0, ping_interval=None) as server:
port = server.sockets[0].getsockname()[1]
service = ServiceConfig(
"http://local", PurePosixPath("/api"), root,
)
config = ClientConfig(
"cache-1", "Cache 1", "cache",
f"ws://127.0.0.1:{port}", token,
root / "state.db", root / "backups", service, service,
ConnectionConfig(
registration_timeout=1,
offline_timeout=0.05,
reconnect_initial=0.01,
reconnect_max=0.01,
reconnect_jitter=False,
),
)
probe = FilesystemProbe(root, True, True, True, True, True)
daemon = ArchiveClientDaemon(config, [probe, probe], [])
await asyncio.to_thread(daemon.store.initialize)
with self.assertRaisesRegex(
RuntimeError, "control heartbeat timed out"
):
await asyncio.wait_for(daemon._connection(), 1)
async def test_eviction_assignment_and_steps_are_admitted(self):
with tempfile.TemporaryDirectory() as directory:
root = Path(directory)
+225 -1
View File
@@ -1,5 +1,6 @@
import hashlib
import tempfile
import threading
import unittest
from pathlib import Path
from unittest.mock import Mock, patch
@@ -11,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
@@ -38,6 +39,120 @@ 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,
job_pb2.JOB_STEP_KIND_SYNCTHING_TRANSFER,
job_pb2.JOB_STEP_KIND_TARGET_MATERIALIZE,
job_pb2.JOB_STEP_KIND_QB_VERIFY,
job_pb2.JOB_STEP_KIND_STAGING_CLEANUP,
):
with self.subTest(step=step), tempfile.TemporaryDirectory() as directory:
root = Path(directory)
store = ClientStore(root / "client.db")
store.initialize()
definition = job_pb2.JobDefinition(
job_id=str(uuid4()),
idempotency_key=str(uuid4()),
operation=job_pb2.JOB_OPERATION_ARCHIVE,
resource_display_name="reconnect fixture",
transfer={
"source_client_id": "cache-1",
"target_client_id": "archive-1",
"route_id": "route-1",
},
)
definition.resource_id.info_hash_v1_hex = "a" * 40
definition.created_at.GetCurrentTime()
executor = ClientJobExecutor(
client_id="cache-1",
qbittorrent=Mock(),
store=store,
qb_root=root,
qb_api_root=Path("/downloads"),
route_path=lambda _: root,
syncthing_transport=Mock(),
sparse_supported=True,
)
executor.assign(control_pb2.AssignJobCommand(
job=definition,
expected_job_revision=1,
expected_last_event_sequence=0,
))
entered = threading.Event()
release = threading.Event()
def execute_step(*_args):
entered.set()
release.wait(1)
return None
executor._execute_step = Mock(side_effect=execute_step)
command = control_pb2.ExecuteStepCommand(
job_id=definition.job_id,
expected_job_revision=1,
expected_last_event_sequence=1,
step=step,
attempt=1,
)
results: list[list[control_pb2.JobEvent]] = []
first = threading.Thread(
target=lambda: results.append(executor.execute(command))
)
second = threading.Thread(
target=lambda: results.append(executor.execute(command))
)
first.start()
self.assertTrue(entered.wait(1))
second.start()
release.set()
first.join(1)
second.join(1)
self.assertFalse(first.is_alive())
self.assertFalse(second.is_alive())
self.assertEqual(executor._execute_step.call_count, 1)
self.assertEqual(len(results), 2)
self.assertEqual(results[0][-1].event_id, results[1][-1].event_id)
def test_capacity_guard_fails_before_data_movement(self):
with tempfile.TemporaryDirectory() as directory:
root = Path(directory)
@@ -111,6 +226,115 @@ class ClientJobHappyPathTests(unittest.TestCase):
syncthing=SlowRescanSyncthing(),
)
def test_target_step_advances_past_replayed_source_completion(self):
with tempfile.TemporaryDirectory() as directory:
root = Path(directory)
store = ClientStore(root / "client.db")
store.initialize()
definition = job_pb2.JobDefinition(
job_id=str(uuid4()),
idempotency_key=str(uuid4()),
operation=job_pb2.JOB_OPERATION_ARCHIVE,
resource_display_name="fixture",
transfer={
"source_client_id": "cache-1",
"target_client_id": "archive-1",
"route_id": "route-1",
},
)
definition.resource_id.info_hash_v1_hex = "a" * 40
definition.created_at.GetCurrentTime()
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.assign(control_pb2.AssignJobCommand(
job=definition,
expected_job_revision=1,
expected_last_event_sequence=1,
))
source_complete = executor._event(
definition,
sequence=5,
revision=4,
event_type=control_pb2.JOB_EVENT_TYPE_STEP_SUCCEEDED,
state=job_pb2.JOB_STATE_RUNNING,
committed=False,
step=job_pb2.JOB_STEP_KIND_SOURCE_STAGE,
step_state=job_pb2.STEP_STATE_SUCCEEDED,
)
executor._record(definition, source_complete)
with patch.object(executor, "_wait_for_syncthing"):
events = executor.execute(control_pb2.ExecuteStepCommand(
job_id=definition.job_id,
expected_job_revision=4,
expected_last_event_sequence=5,
step=job_pb2.JOB_STEP_KIND_SYNCTHING_TRANSFER,
attempt=1,
))
self.assertEqual(events[0].sequence, 6)
self.assertEqual(
events[0].progress.step,
job_pb2.JOB_STEP_KIND_SYNCTHING_TRANSFER,
)
self.assertEqual(events[-1].sequence, 7)
self.assertEqual(
events[-1].type,
control_pb2.JOB_EVENT_TYPE_STEP_SUCCEEDED,
)
def test_first_partial_progress_is_durable(self):
"""A sparse transfer must not lose its only initial observation."""
with tempfile.TemporaryDirectory() as directory:
root = Path(directory)
store = ClientStore(root / "client.db")
store.initialize()
definition = job_pb2.JobDefinition(
job_id=str(uuid4()),
idempotency_key=str(uuid4()),
operation=job_pb2.JOB_OPERATION_ARCHIVE,
transfer={
"source_client_id": "cache-1",
"target_client_id": "archive-1",
"route_id": "route-1",
},
)
definition.resource_id.info_hash_v1_hex = "a" * 40
definition.created_at.GetCurrentTime()
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.assign(control_pb2.AssignJobCommand(
job=definition, expected_job_revision=1,
expected_last_event_sequence=0,
))
def partial_progress(_definition, _step, progress):
progress(0.001, 10, 10_000, "still receiving")
with patch.object(executor, "_execute_step", partial_progress):
events = executor.execute(control_pb2.ExecuteStepCommand(
job_id=definition.job_id, expected_job_revision=1,
expected_last_event_sequence=1,
step=job_pb2.JOB_STEP_KIND_SYNCTHING_TRANSFER, attempt=1,
))
self.assertEqual(len(events), 3)
self.assertEqual(events[1].type, control_pb2.JOB_EVENT_TYPE_PROGRESS)
self.assertEqual(events[1].progress.bytes_complete, 10)
def _run_transfer(
self,
root: Path,
+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 = [
+41
View File
@@ -118,6 +118,47 @@ class ResourceTests(unittest.TestCase):
self.assertEqual(root.available_file_count, 1)
self.assertEqual(root.available_logical_bytes, 3)
def test_qbittorrent_hidden_padding_files_are_normalized(self):
info = {
b"files": [
{b"length": 3, b"path": [b"first.bin"]},
{
b"length": 5,
b"path": [b"_____padding_file_5"],
},
{b"length": 7, b"path": [b"last.bin"]},
],
b"name": b"with-padding",
b"piece length": 16384,
b"pieces": b"x" * 20,
}
metainfo = encode({b"info": info})
torrent_hash = hashlib.sha1(encode(info)).hexdigest()
normalized = normalize_resource(
{
"hash": torrent_hash,
"name": "with-padding",
"state": "stalledUP",
},
[
{
"index": 0, "name": "with-padding/first.bin",
"size": 3, "completed": 3, "priority": 1,
},
{
"index": 1, "name": "with-padding/last.bin",
"size": 7, "completed": 7, "priority": 1,
},
],
metainfo,
)
self.assertEqual(
[item.canonical_path for item in normalized.files],
["with-padding/first.bin", "with-padding/last.bin"],
)
self.assertEqual(normalized.summary.total_file_count, 2)
self.assertEqual(normalized.summary.selected_complete_bytes, 10)
def test_noncanonical_and_unsafe_paths_are_distinct(self):
info = {
b"length": 3,
+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"
+28
View File
@@ -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)