Compare commits

...
28 Commits
Author SHA1 Message Date
cabbage 62a0f24437 fix: recover client control connection safely 2026-08-13 06:25:51 +00:00
cabbage eabbda3c3a fix: recover durable route paths after restart 2026-08-03 11:19:45 +00:00
cabbage a7bc2bf087 test: probe hardlink staging in deployment preflight 2026-08-03 11:13:10 +00:00
cabbage 25d40daea6 Run deployment preflight through client image 2026-08-03 11:08:42 +00:00
cabbage f3517bd21b Route automatic staging through qB data mounts 2026-08-03 10:55:13 +00:00
cabbage 300210c17d Document deployment preflight workflow 2026-08-03 10:22:38 +00:00
cabbage 1f97faeab5 Add deployment preflight and fix helium Syncthing root 2026-08-03 10:16:32 +00:00
cabbage 1c2590c3c8 Add helium archive node deployment manifests 2026-08-03 09:42:53 +00:00
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
cabbage 053b0135b3 fix: defer slow Syncthing rescans 2026-07-24 16:13:39 +00:00
cabbage 57acc90363 fix: honor hardlinks across client mount topology 2026-07-24 16:05:20 +00:00
cabbage a73bcf870f config: reduce zero-copy space reserve default 2026-07-24 15:26:33 +00:00
45 changed files with 2035 additions and 145 deletions
+7 -2
View File
@@ -7,12 +7,17 @@ RUN pip wheel --no-cache-dir --wheel-dir /wheels .
FROM builder AS test
RUN pip install --no-cache-dir /wheels/*.whl
COPY tests ./tests
COPY scripts ./scripts
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
+7
View File
@@ -77,6 +77,13 @@ four eager-mesh routes, transfer/eviction/unarchive workflows, and the
adversarial matrix without live infrastructure or Telegram. See
`e2e/README.md`.
## Deployment Preflight Tests
Check deployment configuration immediately before bringing a new node online.
The reusable read-only preflight validates the deployed image's filesystem and
API view against its Docker bind mounts. See details in
[docs/deployment-preflight.md](docs/deployment-preflight.md).
The complete cross-project design, protocol, workflow, safety, deployment,
Telegram UX, and testing documentation is published under [`docs/`](docs/).
+50 -6
View File
@@ -1,30 +1,74 @@
# Production deployment
These files are the non-secret, host-specific deployment manifests for the
initial x1/x2/lithium topology.
initial x1/x2/lithium topology and the later helium archive node.
- Install the x1 and x2 files as
`~/compose/ArchiveControl-cache/{compose.yaml,client.toml}`.
- Install the lithium files as
`~/compose/ArchiveControl-archive/{compose.yaml,client.toml}`.
- The helium directory contains three isolated compose-project manifests. Install
them under `~/Repositories/compose/{qbittorrent-helium,syncthing-helium,ArchiveControl-archive}`
as described in `helium/README.md`.
- Create sibling `state`, `backups`, and `secrets` directories owned by the
configured container UID/GID.
- Secret files are never committed. Each `secrets` directory contains
`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
Before starting a stack, validate it with:
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.
Automatic routes always use the Syncthing API path `routes/<route-id>`. Bind
that path from a dedicated directory beneath the qB data root, then map the
same path through the qB client mount:
```yaml
# Syncthing compose project
volumes:
- /srv/syncthing-config:/var/syncthing
- /srv/downloads/.archive-control-routes:/var/syncthing/routes
# Archive Control client compose project
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/routes" = "/data/qb/.archive-control-routes" }
```
Do **not** mount the route directory separately into the client. The override
is the authoritative mapping and keeps qB source files and automatic route
folders in one mount namespace. Existing manually configured Syncthing folders
below the qB tree may retain their own exact overrides. Use API-visible paths,
not Syncthing folder IDs, as override keys.
Before starting a stack, validate its config and then run the generic host
preflight with that machine's own paths and container names:
```sh
docker compose config
docker compose run --rm archive-client --check-config
python3 scripts/preflight-deployment.py --help
```
The check is fail-fast and performs local permission, filesystem, sparse-file,
hard-link, and reflink probes. Normal startup additionally probes the local
qBittorrent and Syncthing APIs before registration.
`scripts/preflight-deployment.py` is deliberately topology-agnostic: pass the
host's `client.toml` and its client, qBittorrent, and Syncthing container names.
It resolves Docker mounts rather than assuming project names or host paths. It
also catches a common fatal error: a Syncthing `api_root` that names its config
volume rather than the mounted shared-data tree.
+34
View File
@@ -0,0 +1,34 @@
# Helium archive-node deployment
This directory contains the non-secret manifests for the dedicated helium
archive node. Install the files into these separate compose projects:
- `~/Repositories/compose/qbittorrent-helium/compose.yaml`
- `~/Repositories/compose/syncthing-helium/compose.yaml`
- `~/Repositories/compose/ArchiveControl-archive/{compose.yaml,client.toml}`
The qBittorrent API and Syncthing GUI/API bind only to loopback. Syncthing
transport/discovery ports are published normally. qBittorrent data and
Syncthing route folders both live under `/media/Data2`; the Archive Control
container therefore mounts `/media/Data2` once at `/data/storage` so archive
placements can hardlink rather than make a full copy.
`/media/Data2/Downloading` already contains unrelated files. The fresh qB
instance must not import or manage them; only Archive Control-created torrents
are managed.
Before the first client start, run the repository preflight on the target host
(replace container names if that host uses different compose project names):
```sh
python3 scripts/preflight-deployment.py \
--client-config ~/Repositories/compose/ArchiveControl-archive/client.toml \
--client-container archive-control-helium-archive-client-1 \
--syncthing-container helium-syncthing-syncthing-1 \
--qbittorrent-container helium-qbittorrent-qbittorrent-1
```
It is read-only: it checks the configured API roots resolve to the same host
paths as the client roots, that both data roots share one client bind mount for
hardlinks, secret-file permissions, filesystem capabilities, and the local qB
and Syncthing APIs. It does not send the control token or modify any container.
@@ -0,0 +1,19 @@
name: archive-control-helium
services:
archive-client:
image: sodium/archive-clients:v0.1.19
user: "1000:1000"
restart: unless-stopped
network_mode: host
command: ["--config", "/etc/archive-control/client.toml"]
environment:
QB_USER: admin
volumes:
- ./client.toml:/etc/archive-control/client.toml:ro
- ./secrets:/run/secrets:ro
- ./state:/var/lib/archive-control
- ./backups:/var/backups/archive-control
# Do not split this into separate qB and Sync bind mounts: hardlinks
# require both trees to resolve within this one mount namespace.
- /media/Data2:/data/storage
+42
View File
@@ -0,0 +1,42 @@
client_id = "helium-archive"
display_name = "helium archive"
role = "archive"
control_endpoint = "ws://bot.everdream.xyz:8766/archive_control"
shared_token_file = "/run/secrets/archive_control_token"
state_db = "/var/lib/archive-control/client.db"
backup_dir = "/var/backups/archive-control"
[connection]
registration_timeout = "10s"
heartbeat_interval = "15s"
offline_timeout = "45s"
reconnect_initial = "1s"
reconnect_max = "60s"
reconnect_reset_after = "60s"
reconnect_jitter = true
[jobs]
stall_after = "30m"
verification_timeout = "30m"
poll_interval = "1s"
free_space_reserve_bytes = 33554432
[backup]
interval = "6h"
recent = 12
daily = 14
weekly = 8
[qbittorrent]
endpoint = "http://127.0.0.1:8081"
username = "${QB_USER}"
password_file = "/run/secrets/qb_password"
api_root = "/downloads/Downloading"
local_root = "/data/storage/Downloading"
[syncthing]
endpoint = "http://127.0.0.1:8384"
api_key_file = "/run/secrets/syncthing_api_key"
api_root = "/var/syncthing/Sync"
local_root = "/data/storage/Sync"
advertised_addresses = ["dynamic"]
@@ -0,0 +1,24 @@
name: helium-qbittorrent
services:
qbittorrent:
image: sodium/qbittorrent-nox:5.2.3-lt2-1-proxy-listen-amd64
environment:
PUID: "1000"
PGID: "1000"
TZ: Etc/UTC
WEBUI_PORT: "8081"
TORRENTING_PORT: "6881"
command: ["--webui-port=8081", "--confirm-legal-notice"]
volumes:
- ./config:/config
# This is deliberately one parent mount. Archive Control maps both the
# qB and Syncthing subtrees through the same mount namespace.
- /media/Data2:/downloads
ports:
# Keep the published and internal WebUI ports identical: qBittorrent
# rejects a forwarded request whose Host header carries a different port.
- "127.0.0.1:8081:8081"
- "6881:6881/tcp"
- "6881:6881/udp"
restart: unless-stopped
+18
View File
@@ -0,0 +1,18 @@
#!/bin/sh
set -eu
home=/var/syncthing
api_key=$(cat /run/secrets/syncthing_api_key)
if [ ! -f "$home/config.xml" ]; then
syncthing generate --home="$home" --no-port-probing
fi
exec syncthing serve \
--home="$home" \
--gui-address=http://0.0.0.0:8384 \
--gui-apikey="$api_key" \
--no-browser \
--no-port-probing \
--no-restart \
--no-upgrade
@@ -0,0 +1,19 @@
name: helium-syncthing
services:
syncthing:
image: syncthing/syncthing:2.1.2@sha256:4464f4161dd0251e20d46bb3aec83363db75d80cef1abdd5d5fd4054b04a004d
user: "1000:1000"
hostname: helium
entrypoint: ["/bootstrap/start-syncthing.sh"]
volumes:
- ./start-syncthing.sh:/bootstrap/start-syncthing.sh:ro
- ./secrets:/run/secrets:ro
- ./config:/var/syncthing
- /media/Data2/Sync:/var/syncthing/Sync
ports:
- "127.0.0.1:8384:8384"
- "22000:22000/tcp"
- "22000:22000/udp"
- "21027:21027/udp"
restart: unless-stopped
+2 -1
View File
@@ -19,7 +19,7 @@ reconnect_jitter = true
stall_after = "30m" # Status warning only; it does not fail the job.
verification_timeout = "30m"
poll_interval = "1s"
free_space_reserve_bytes = 1073741824 # Rechecked immediately before work.
free_space_reserve_bytes = 33554432 # 32 MiB; payload bytes are added for copy paths.
[backup]
interval = "6h"
@@ -39,4 +39,5 @@ endpoint = "http://syncthing:8384"
api_key_file = "/run/secrets/syncthing_api_key"
api_root = "/var/syncthing"
local_root = "/data/sync"
local_path_overrides = { "/var/syncthing/routes" = "/data/qb/.archive-control-routes" }
advertised_addresses = ["dynamic"]
+1 -1
View File
@@ -2,7 +2,7 @@ name: archive-control-archive
services:
archive-client:
image: sodium/archive-clients:v0.1.4
image: sodium/archive-clients:v0.1.14
user: "1000:1000"
restart: unless-stopped
command: ["--config", "/etc/archive-control/client.toml"]
+4 -1
View File
@@ -19,7 +19,7 @@ reconnect_jitter = true
stall_after = "30m" # Status warning only; it does not fail the job.
verification_timeout = "30m"
poll_interval = "1s"
free_space_reserve_bytes = 1073741824 # Rechecked immediately before work.
free_space_reserve_bytes = 33554432 # 32 MiB; payload bytes are added for copy paths.
[backup]
interval = "6h"
@@ -39,4 +39,7 @@ endpoint = "http://127.0.0.1:8384"
api_key_file = "/run/secrets/syncthing_api_key"
api_root = "/var/syncthing"
local_root = "/data/sync"
# This existing folder is physically inside the qB data tree. Map it through
# that same client bind mount so source staging can use hardlinks.
local_path_overrides = { "/var/syncthing/DownloadsSync" = "/data/qb/Sync", "/var/syncthing/routes" = "/data/qb/.archive-control-routes" }
advertised_addresses = ["dynamic"]
+1 -2
View File
@@ -2,7 +2,7 @@ name: archive-control-cache
services:
archive-client:
image: sodium/archive-clients:v0.1.4
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/DownloadsSync
+5 -1
View File
@@ -19,7 +19,7 @@ reconnect_jitter = true
stall_after = "30m" # Status warning only; it does not fail the job.
verification_timeout = "30m"
poll_interval = "1s"
free_space_reserve_bytes = 1073741824 # Rechecked immediately before work.
free_space_reserve_bytes = 33554432 # 32 MiB; payload bytes are added for copy paths.
[backup]
interval = "6h"
@@ -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", "/var/syncthing/routes" = "/data/qb/.archive-control-routes" }
advertised_addresses = ["dynamic"]
+1 -2
View File
@@ -2,7 +2,7 @@ name: archive-control-cache
services:
archive-client:
image: sodium/archive-clients:v0.1.4
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
+3 -2
View File
@@ -16,8 +16,9 @@ credentials and host-specific details.
6. [Storage safety and recovery](storage-safety.md)
7. [qBittorrent and Syncthing integration](service-apis.md)
8. [Configuration, deployment, and usage](deployment-and-usage.md)
9. [Telegram UX](telegram-ux.md)
10. [Testing strategy](testing.md)
9. [Deployment preflight tests](deployment-preflight.md)
10. [Telegram UX](telegram-ux.md)
11. [Testing strategy](testing.md)
The independent protobuf source of truth lives in the
`cabbage/archive-control-proto` repository. Generated bindings in this
+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
+38 -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
@@ -92,6 +95,7 @@ backup_dir = "/var/backups/archive-control"
registration_timeout = "10s" # First response must arrive within this window.
heartbeat_interval = "15s" # Server may negotiate a different effective value.
offline_timeout = "45s"
outbound_enqueue_timeout = "15s" # Reconnect rather than wedging if sends stop draining.
reconnect_initial = "1s"
reconnect_max = "60s" # Retry forever, never wait longer than this.
reconnect_reset_after = "60s"
@@ -102,7 +106,7 @@ stall_after = "30m" # Warning only; no automatic job failure.
verification_timeout = "30m" # Stopped qB full-recheck deadline.
poll_interval = "1s" # Active qB/Syncthing observation interval.
# Worst-case copy fallback must leave this many bytes free.
free_space_reserve_bytes = 1073741824
free_space_reserve_bytes = 33554432 # 32 MiB; payload bytes are added for copy paths.
[backup]
interval = "6h"
@@ -125,6 +129,16 @@ local_root = "/data/sync"
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. 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
override them per node.
@@ -213,7 +227,7 @@ cache/archive routes according to policy.
```yaml
services:
archive-client:
image: sodium/archive-clients:v0.1.4
image: sodium/archive-clients:v0.1.14
user: "1001:1001"
restart: unless-stopped
command: ["archive-client", "--config", "/etc/archive-control/client.toml"]
@@ -246,6 +260,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
+121
View File
@@ -0,0 +1,121 @@
# Deployment preflight tests
Run the deployment preflight on every new node after its qBittorrent and
Syncthing stacks are running, but before starting the Archive Control daemon
for normal operation or creating any routes/jobs. It does not contact the
control daemon, alter Syncthing/qBittorrent state, or print secret contents.
It does create and remove two unique, zero-byte probe files while proving that
the live client mount topology permits hard-link staging; no resource data or
configuration is modified.
The script is [`scripts/preflight-deployment.py`](../scripts/preflight-deployment.py).
It deliberately takes paths and container names as arguments rather than
assuming hostnames, compose project names, mount paths, credentials, or roles.
It can therefore be used unchanged on cache and archive nodes.
## Requirements
- Run it on the Docker host being checked.
- Use Python 3.8 or newer. The script reads the mounted configuration through
the client image, so it does not depend on the host Python TOML library.
- Have Docker CLI access to the daemon.
- Have the checkout containing the script available on that host, or copy only
this script into the node's compose directory.
- Bring up qBittorrent and Syncthing first. Do not start the normal Archive
Control daemon yet.
The disposable client container below uses the exact `archive-client` service
image, mounts, environment, and network specified by that node's compose file.
Consequently its qB/Syncthing API probes validate the same runtime environment
the daemon will use, not the host shell's network namespace.
## New-node procedure
The example uses generic names. Replace paths, compose project directories,
service names, and real container names for the node being deployed.
```sh
# 1. Bring up only the local dependencies.
cd /srv/compose/syncthing-node && docker compose up -d
cd /srv/compose/qbittorrent-node && docker compose up -d
# 2. Create a temporary client from the exact production image/config, without
# running the daemon or registering it with control.
cd /srv/compose/ArchiveControl-archive
docker compose run -d --no-deps --name archive-control-preflight \
--entrypoint sleep archive-client infinity
# 3. Discover the real dependency container names if needed.
docker ps --format '{{.Names}}'
# 4. Run the deployment checks. They include a disposable hard-link probe.
python3 /path/to/archive-clients/scripts/preflight-deployment.py \
--client-config /srv/compose/ArchiveControl-archive/client.toml \
--client-container archive-control-preflight \
--syncthing-container syncthing-node-syncthing-1 \
--qbittorrent-container qbittorrent-node-qbittorrent-1
# 5. Remove the temporary container, then start the real daemon only after a
# successful result.
docker rm -f archive-control-preflight
docker compose up -d archive-client
```
If the real client configuration is mounted somewhere other than the standard
`/etc/archive-control/client.toml`, pass that in-container location explicitly:
```sh
python3 scripts/preflight-deployment.py ... \
--container-config /custom/path/client.toml
```
## Required arguments
| Argument | Meaning |
| --- | --- |
| `--client-config` | Host path to the exact `client.toml` that will be mounted into the daemon. |
| `--client-container` | Running disposable or already-running client container to inspect and probe. |
| `--syncthing-container` | Running Syncthing container for this node. |
| `--qbittorrent-container` | Running qBittorrent container for this node. |
| `--container-config` | Optional `client.toml` path inside the client container; defaults to `/etc/archive-control/client.toml`. |
## What it checks
- The future automatic `routes/<route-id>` path resolves through Docker mounts
to the same host path as its client-side mapping.
- That future route path and qBittorrent content root use one client bind
mount, so hard-link staging remains possible rather than silently falling
back to a space-consuming copy.
- A real `link(2)` operation between a unique zero-byte file in the qB root
and one in the future automatic-route root. The probe verifies that both
names refer to the same inode and removes them unconditionally.
- Every existing `archive-control:*` Syncthing folder beneath `routes/` still
has its local directory. This catches a route-root bind-mount migration that
would otherwise hide an already configured folder and later fail a job with
`No such file or directory`.
- The token, qB password, and Syncthing API-key files are non-empty regular
files with no group/world permissions.
- The client image can read its configuration and reports usable permissions,
sparse-file support, and filesystem capabilities for both roots.
- qBittorrent authentication/version compatibility and Syncthing
authentication/device identity are healthy from the client container.
In particular, if `syncthing.api_root` is the Syncthing configuration root,
bind its `routes` subdirectory from a dedicated directory beneath qB's data
root and add a `local_path_overrides` mapping for that exact API path. This
prevents automatic route folders being created on a small configuration
filesystem while the client expects to hardlink from the qB data mount.
When converting an existing node, stop its client and Syncthing containers,
move each existing `routes/<route-id>` directory from the old Syncthing config
tree into the new qB-backed route-root directory, then recreate Syncthing and
run this preflight. Do not discard existing route directories: they can contain
route handshake state or an in-progress transfer namespace.
## Failure handling
Treat a nonzero exit status as a deployment blocker. Correct the compose
mounts, `client.toml`, secret permissions, service credentials, or service
image/version that the message identifies, then rerun the same command. Do not
work around a mapping failure by creating route folders manually: the route
provisioner relies on those API/local path mappings being exact.
+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
@@ -21,7 +21,7 @@ verification_timeout = "30m"
# Poll qBittorrent and Syncthing at this interval while a step is active.
poll_interval = "1s"
# Space that must remain free after a worst-case copy fallback.
free_space_reserve_bytes = 1073741824
free_space_reserve_bytes = 33554432
[backup]
interval = "6h"
+1 -1
View File
@@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta"
[project]
name = "archive-clients"
version = "0.1.4"
version = "0.1.19"
requires-python = ">=3.11"
dependencies = ["protobuf==7.35.1", "websockets==16.0"]
+327
View File
@@ -0,0 +1,327 @@
#!/usr/bin/env python3
"""Validate a host's Archive Control Docker deployment before it is used.
This intentionally relies only on the config file and the named containers on
the host where it runs. It never reads secret contents or calls a remote
control daemon. The hard-link probe creates uniquely named, empty files in
the configured qB and automatic-route directories and removes them afterward.
"""
from __future__ import annotations
import argparse
import json
import os
import stat
import subprocess
import sys
from dataclasses import dataclass
from pathlib import Path, PurePosixPath
from typing import Any
class CheckFailure(RuntimeError):
pass
@dataclass(frozen=True)
class Mapping:
source: Path
destination: PurePosixPath
def docker_inspect(container: str) -> dict[str, Any]:
result = subprocess.run(
["docker", "inspect", container], check=False, text=True,
capture_output=True,
)
if result.returncode:
raise CheckFailure(f"cannot inspect container {container!r}")
try:
value = json.loads(result.stdout)
return value[0]
except (json.JSONDecodeError, IndexError, TypeError) as exc:
raise CheckFailure(f"invalid Docker inspection for {container!r}") from exc
def map_path(container: dict[str, Any], path: str) -> Mapping:
candidate = PurePosixPath(path)
matches: list[tuple[int, Mapping]] = []
for mount in container.get("Mounts", []):
if mount.get("Type") not in {"bind", "volume"}:
continue
source, destination = mount.get("Source"), mount.get("Destination")
if not isinstance(source, str) or not isinstance(destination, str):
continue
destination_path = PurePosixPath(destination)
try:
relative = candidate.relative_to(destination_path)
except ValueError:
continue
matches.append((len(destination_path.parts), Mapping(
Path(source).joinpath(*relative.parts), destination_path
)))
if not matches:
raise CheckFailure(f"{path} is not backed by a Docker bind/volume mount")
return max(matches, key=lambda item: item[0])[1]
def container_user(container: dict[str, Any]) -> str | None:
value = container.get("Config", {}).get("User")
return value if isinstance(value, str) and value else None
def require_same_path(label: str, left: Mapping, right: Mapping) -> None:
if left.source.resolve(strict=False) != right.source.resolve(strict=False):
raise CheckFailure(
f"{label} host paths differ: {left.source} != {right.source}"
)
def require_regular_secret(path: Path) -> None:
try:
mode = path.stat().st_mode
except OSError as exc:
raise CheckFailure(f"secret is unavailable: {path}") from exc
if not stat.S_ISREG(mode) or stat.S_IMODE(mode) & 0o077 or path.stat().st_size == 0:
raise CheckFailure(f"secret must be a non-empty mode-0600 regular file: {path}")
def route_local_path(syncthing: dict[str, Any]) -> str:
"""Resolve the fixed, future automatic ``routes/`` path in the client."""
api_root = PurePosixPath(syncthing["api_root"])
candidate = api_root / "routes"
overrides = syncthing.get("local_path_overrides", {})
if not isinstance(overrides, dict):
raise CheckFailure("syncthing.local_path_overrides must be a table")
matches: list[tuple[int, PurePosixPath, str]] = []
for raw_api, raw_local in overrides.items():
if not isinstance(raw_api, str) or not isinstance(raw_local, str):
raise CheckFailure("syncthing.local_path_overrides entries are invalid")
root = PurePosixPath(raw_api)
try:
relative = candidate.relative_to(root)
except ValueError:
continue
matches.append((len(root.parts), relative, raw_local))
if matches:
_, relative, local_root = max(matches, key=lambda item: item[0])
return str(Path(local_root).joinpath(*relative.parts))
try:
relative = candidate.relative_to(api_root)
except ValueError as exc:
raise CheckFailure("future route path is outside syncthing.api_root") from exc
return str(Path(syncthing["local_root"]).joinpath(*relative.parts))
def run_client_check(container: str, config: str, flag: str) -> None:
result = subprocess.run(
["docker", "exec", container, "archive-client", "--config", config, flag],
check=False, text=True, capture_output=True,
)
if result.returncode:
detail = result.stderr.strip() or result.stdout.strip() or "failed"
raise CheckFailure(f"client {flag} failed: {detail}")
print(result.stdout.strip())
def mounted_config(container: str, config: str) -> dict[str, Any]:
"""Read normalized non-secret config through the client image itself."""
program = '''import json,sys
from pathlib import Path
from archive_clients.config import ClientConfig
config=ClientConfig.load(Path(sys.argv[1]))
value={"shared_token_file":str(config.shared_token_file)}
value["qbittorrent"]={"api_root":str(config.qbittorrent.api_root),"local_root":str(config.qbittorrent.local_root),"password_file":str(config.qbittorrent.password_file)}
value["syncthing"]={"api_root":str(config.syncthing.api_root),"local_root":str(config.syncthing.local_root),"api_key_file":str(config.syncthing.api_key_file),"local_path_overrides":{str(api):str(local) for api,local in config.syncthing.local_path_overrides}}
print(json.dumps(value,sort_keys=True))'''
result = subprocess.run(
["docker", "exec", container, "python", "-c", program, config],
check=False, text=True, capture_output=True,
)
if result.returncode:
detail = result.stderr.strip() or result.stdout.strip() or "failed"
raise CheckFailure(f"cannot load mounted client config: {detail}")
try:
value = json.loads(result.stdout)
except json.JSONDecodeError as exc:
raise CheckFailure("mounted client config output is invalid") from exc
if not isinstance(value, dict):
raise CheckFailure("mounted client config is invalid")
return value
def run_service_check(container: str, config: str) -> None:
program = '''import json,sys
from pathlib import Path
from archive_clients.config import ClientConfig
from archive_clients.probes import probe_root
from archive_clients.services import probe_qbittorrent,probe_syncthing
config=ClientConfig.load(Path(sys.argv[1]))
sync_probe=probe_root(config.syncthing.local_root)
probes=[probe_qbittorrent(config.qbittorrent),probe_syncthing(config.syncthing,sparse_supported=sync_probe.sparse_files)]
print(json.dumps({"services":[{"service":p.service,"state":p.state,"detail":p.detail,"version":p.version,"device_id":p.device_id} for p in probes]},sort_keys=True))
raise SystemExit(0 if all(p.state == 1 for p in probes) else 1)'''
result = subprocess.run(
["docker", "exec", container, "python", "-c", program, config],
check=False, text=True, capture_output=True,
)
if result.returncode:
detail = result.stderr.strip() or result.stdout.strip() or "failed"
raise CheckFailure(f"client local-service probe failed: {detail}")
print(result.stdout.strip())
def run_hardlink_probe(
container: str, qb_root: str, route_root: str, user: str | None = None,
) -> None:
"""Prove that the daemon can hard-link from qB data into automatic routes."""
program = r'''import os,sys,tempfile
source_root,route_root=sys.argv[1:]
source_path=None
destination_path=None
try:
source_fd,source_path=tempfile.mkstemp(prefix=".archive-control-preflight-",dir=source_root)
os.close(source_fd)
destination_fd,destination_path=tempfile.mkstemp(prefix=".archive-control-preflight-",dir=route_root)
os.close(destination_fd)
os.unlink(destination_path)
os.link(source_path,destination_path)
source_stat=os.stat(source_path)
destination_stat=os.stat(destination_path)
if source_stat.st_dev != destination_stat.st_dev or source_stat.st_ino != destination_stat.st_ino:
raise RuntimeError("link(2) did not produce the same filesystem inode")
print("hard-link staging probe passed")
finally:
for path in (destination_path,source_path):
if path:
try:
os.unlink(path)
except FileNotFoundError:
pass
'''
command = ["docker", "exec"]
if user:
command.extend(["--user", user])
command.extend([container, "python", "-c", program, qb_root, route_root])
result = subprocess.run(
command,
check=False, text=True, capture_output=True,
)
if result.returncode:
detail = result.stderr.strip() or result.stdout.strip() or "failed"
raise CheckFailure(f"hard-link staging probe failed: {detail}")
print(result.stdout.strip())
def run_existing_route_directory_check(container: str, config: str) -> None:
"""Reject a mount migration that hides an already configured route folder."""
program = r'''import json,os,sys
from pathlib import Path,PurePosixPath
from archive_clients.config import ClientConfig
from archive_clients.syncthing import SyncthingHttp
config=ClientConfig.load(Path(sys.argv[1]))
transport=SyncthingHttp(config.syncthing)
route_root=config.syncthing.api_root / "routes"
missing=[]
checked=[]
for folder in transport.get_json("/rest/config").get("folders",[]):
if not isinstance(folder,dict) or not str(folder.get("label","")).startswith("archive-control:"):
continue
raw=folder.get("path")
if not isinstance(raw,str) or not raw:
continue
api=PurePosixPath(raw)
if api.parts and api.parts[0] == "~":
api=config.syncthing.api_root.joinpath(*api.parts[1:])
try:
api.relative_to(route_root)
except ValueError:
continue
local=config.syncthing.roots.api_to_local(api.as_posix())
checked.append(str(local))
if not local.is_dir():
missing.append({"folder_id":folder.get("id",""),"path":str(local)})
print(json.dumps({"checked":checked,"missing":missing},sort_keys=True))
raise SystemExit(1 if missing else 0)
'''
result = subprocess.run(
["docker", "exec", container, "python", "-c", program, config],
check=False, text=True, capture_output=True,
)
if result.returncode:
detail = result.stderr.strip() or result.stdout.strip() or "failed"
raise CheckFailure(
"an existing Archive Control Syncthing route directory is missing; "
f"complete the route-directory migration before starting jobs: {detail}"
)
print(result.stdout.strip())
def main(argv: list[str] | None = None) -> int:
parser = argparse.ArgumentParser(
description="Read-only Archive Control Docker deployment preflight"
)
parser.add_argument("--client-config", type=Path, required=True,
help="host path to client.toml")
parser.add_argument("--client-container", required=True)
parser.add_argument("--syncthing-container", required=True)
parser.add_argument("--qbittorrent-container", required=True)
parser.add_argument("--container-config", default="/etc/archive-control/client.toml",
help="client.toml path inside the client container")
args = parser.parse_args(argv)
try:
client = docker_inspect(args.client_container)
syncthing = docker_inspect(args.syncthing_container)
qbittorrent = docker_inspect(args.qbittorrent_container)
configured_host_path = map_path(client, args.container_config).source
if configured_host_path.resolve(strict=False) != args.client_config.resolve(strict=False):
raise CheckFailure(
"--client-config does not match the file mounted into the client"
)
config = mounted_config(args.client_container, args.container_config)
qb = config["qbittorrent"]
sync = config["syncthing"]
client_qb = map_path(client, qb["local_root"])
client_route = map_path(client, route_local_path(sync))
syncthing_route = map_path(
syncthing, str(PurePosixPath(sync["api_root"]) / "routes")
)
qb_api = map_path(qbittorrent, qb["api_root"])
require_same_path(
"future Syncthing routes/local route mapping",
syncthing_route, client_route,
)
require_same_path("qBittorrent api_root/local_root", qb_api, client_qb)
if client_qb.destination != client_route.destination:
raise CheckFailure(
"qBittorrent and future route roots use separate client bind "
"mounts; hard-link staging would be unavailable"
)
run_hardlink_probe(
args.client_container, qb["local_root"], route_local_path(sync),
container_user(client),
)
run_existing_route_directory_check(
args.client_container, args.container_config
)
for key in ("shared_token_file",):
require_regular_secret(map_path(client, config[key]).source)
for service, key in ((qb, "password_file"), (sync, "api_key_file")):
require_regular_secret(map_path(client, service[key]).source)
run_client_check(args.client_container, args.container_config, "--check-config")
run_service_check(args.client_container, args.container_config)
except (CheckFailure, KeyError, OSError) as exc:
print(f"preflight failed: {exc}", file=sys.stderr)
return 1
print("preflight passed: bind mappings, hard-link staging, secrets, filesystem capabilities, and local APIs are healthy")
return 0
if __name__ == "__main__":
raise SystemExit(main())
+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_")
+69 -11
View File
@@ -27,17 +27,28 @@ _ENV = re.compile(r"\$\{([A-Za-z_][A-Za-z0-9_]*)\}")
class RootMapping:
api_root: PurePosixPath
local_root: Path
local_path_overrides: tuple[tuple[PurePosixPath, Path], ...] = ()
def api_to_local(self, api_path: str) -> Path:
candidate = PurePosixPath(api_path)
try:
relative = candidate.relative_to(self.api_root)
except ValueError as exc:
raise ConfigError("API path is outside its configured root") from exc
if any(part in {"", ".", ".."} for part in relative.parts):
raise ConfigError("API path contains an unsafe component")
for api_root, local_root in self.local_path_overrides:
relative = _safe_relative(candidate, api_root)
if relative is not None:
return local_root.joinpath(*relative.parts)
relative = _safe_relative(candidate, self.api_root)
if relative is None:
raise ConfigError("API path is outside its configured root")
return self.local_root.joinpath(*relative.parts)
def local_root_for_api(self, api_path: str) -> Path:
candidate = PurePosixPath(api_path)
for api_root, local_root in self.local_path_overrides:
if _safe_relative(candidate, api_root) is not None:
return local_root
if _safe_relative(candidate, self.api_root) is None:
raise ConfigError("API path is outside its configured root")
return self.local_root
@dataclass(frozen=True)
class ServiceConfig:
@@ -48,10 +59,13 @@ class ServiceConfig:
password_file: Path | None = None
api_key_file: Path | None = None
advertised_addresses: tuple[str, ...] = ()
local_path_overrides: tuple[tuple[PurePosixPath, Path], ...] = ()
@property
def roots(self) -> RootMapping:
return RootMapping(self.api_root, self.local_root)
return RootMapping(
self.api_root, self.local_root, self.local_path_overrides
)
def read_password(self) -> str | None:
return (
@@ -71,6 +85,7 @@ class ConnectionConfig:
registration_timeout: float = 10
heartbeat_interval: float = 15
offline_timeout: float = 45
outbound_enqueue_timeout: float = 15
reconnect_initial: float = 1
reconnect_max: float = 60
reconnect_reset_after: float = 60
@@ -82,7 +97,7 @@ class JobsConfig:
stall_after: float = 30 * 60
verification_timeout: float = 30 * 60
poll_interval: float = 1
free_space_reserve_bytes: int = 1024 * 1024 * 1024
free_space_reserve_bytes: int = 32 * 1024 * 1024
@dataclass(frozen=True)
@@ -160,7 +175,7 @@ def _service(value: Any, name: str) -> ServiceConfig:
raise ConfigError(f"{name} must be a table")
allowed = {
"endpoint", "api_root", "local_root", "username", "password_file",
"api_key_file", "advertised_addresses",
"api_key_file", "advertised_addresses", "local_path_overrides",
}
_keys(value, allowed, name)
api_root = PurePosixPath(_string(value, "api_root"))
@@ -182,6 +197,7 @@ def _service(value: Any, name: str) -> ServiceConfig:
raise ConfigError("qbittorrent username and password_file are required")
if name == "syncthing" and "api_key_file" not in value:
raise ConfigError("syncthing api_key_file is required")
overrides = _local_path_overrides(value, api_root, name)
return ServiceConfig(
_endpoint(value, "endpoint", {"http", "https"}), api_root,
_absolute_path(value, "local_root"), username,
@@ -189,15 +205,52 @@ def _service(value: Any, name: str) -> ServiceConfig:
if "password_file" in value else None,
_absolute_path(value, "api_key_file")
if "api_key_file" in value else None,
tuple(addresses),
tuple(addresses), overrides,
)
def _local_path_overrides(
value: dict[str, Any], api_root: PurePosixPath, name: str
) -> tuple[tuple[PurePosixPath, Path], ...]:
raw = value.get("local_path_overrides", {})
if name != "syncthing" and raw:
raise ConfigError(f"{name}.local_path_overrides is unsupported")
if not isinstance(raw, dict):
raise ConfigError(f"{name}.local_path_overrides must be a table")
parsed: list[tuple[PurePosixPath, Path]] = []
for raw_api_path, raw_local_path in raw.items():
if not isinstance(raw_api_path, str) or not isinstance(raw_local_path, str):
raise ConfigError(f"{name}.local_path_overrides entries must be strings")
candidate = PurePosixPath(raw_api_path)
if not candidate.is_absolute() or ".." in candidate.parts:
raise ConfigError(f"{name}.local_path_overrides API path is invalid")
if _safe_relative(candidate, api_root) is None:
raise ConfigError(f"{name}.local_path_overrides API path is outside root")
local = Path(raw_local_path)
if not local.is_absolute():
raise ConfigError(f"{name}.local_path_overrides local path is invalid")
parsed.append((candidate, local))
return tuple(sorted(parsed, key=lambda item: len(item[0].parts), reverse=True))
def _safe_relative(
candidate: PurePosixPath, root: PurePosixPath
) -> PurePosixPath | None:
try:
relative = candidate.relative_to(root)
except ValueError:
return None
if any(part in {"", ".", ".."} for part in relative.parts):
raise ConfigError("API path contains an unsafe component")
return relative
def _connection(value: Any) -> ConnectionConfig:
if not isinstance(value, dict):
raise ConfigError("connection must be a table")
_keys(value, {
"registration_timeout", "heartbeat_interval", "offline_timeout",
"outbound_enqueue_timeout",
"reconnect_initial", "reconnect_max", "reconnect_reset_after",
"reconnect_jitter",
}, "connection")
@@ -205,6 +258,9 @@ def _connection(value: Any) -> ConnectionConfig:
registration_timeout=_duration(value.get("registration_timeout", "10s")),
heartbeat_interval=_duration(value.get("heartbeat_interval", "15s")),
offline_timeout=_duration(value.get("offline_timeout", "45s")),
outbound_enqueue_timeout=_duration(
value.get("outbound_enqueue_timeout", "15s")
),
reconnect_initial=_duration(value.get("reconnect_initial", "1s")),
reconnect_max=_duration(value.get("reconnect_max", "60s")),
reconnect_reset_after=_duration(value.get("reconnect_reset_after", "60s")),
@@ -214,6 +270,8 @@ def _connection(value: Any) -> ConnectionConfig:
raise ConfigError("reconnect_initial cannot exceed reconnect_max")
if result.offline_timeout <= result.heartbeat_interval:
raise ConfigError("offline_timeout must exceed heartbeat_interval")
if result.outbound_enqueue_timeout <= 0:
raise ConfigError("outbound_enqueue_timeout must be positive")
if not isinstance(result.reconnect_jitter, bool):
raise ConfigError("reconnect_jitter must be boolean")
return result
@@ -239,7 +297,7 @@ def _jobs(value: Any) -> JobsConfig:
),
poll_interval=_duration(value.get("poll_interval", "1s")),
free_space_reserve_bytes=_positive_int(
value.get("free_space_reserve_bytes", 1024 * 1024 * 1024),
value.get("free_space_reserve_bytes", 32 * 1024 * 1024),
"jobs.free_space_reserve_bytes",
),
)
+136 -37
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,18 @@ 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 self._next_frame(
websocket, writer
)
except ConnectionClosedOK:
return
await self._handle(decode(frame), outbound, command_tasks)
finally:
writer.cancel()
@@ -290,6 +318,63 @@ class ArchiveClientDaemon:
while True:
await websocket.send(await outbound.get())
async def _next_frame(self, websocket: Any, writer: asyncio.Task[None]):
"""Receive one frame while supervising the connection writer.
The previous receive-only wait let a failed writer go unnoticed until
the remote side happened to close or a heartbeat timeout elapsed.
Waiting for both directions makes a failed send an immediate,
reconnectable connection failure.
"""
receiver = asyncio.create_task(websocket.recv())
try:
done, _ = await asyncio.wait(
{receiver, writer},
timeout=self.config.connection.offline_timeout,
return_when=asyncio.FIRST_COMPLETED,
)
if writer in done:
if not receiver.done():
receiver.cancel()
await asyncio.gather(receiver, return_exceptions=True)
if writer.cancelled():
raise RuntimeError("control writer was cancelled")
error = writer.exception()
if error is not None:
raise RuntimeError("control writer failed") from error
raise RuntimeError("control writer stopped")
if receiver not in done:
receiver.cancel()
await asyncio.gather(receiver, return_exceptions=True)
raise RuntimeError("control heartbeat timed out")
return receiver.result()
except BaseException:
if not receiver.done():
receiver.cancel()
await asyncio.gather(receiver, return_exceptions=True)
raise
async def _enqueue(
self, outbound: asyncio.Queue[str], message: str
) -> None:
"""Bound connection-local backpressure so a dead writer cannot wedge I/O.
A job event is durable before it is sent, so dropping this particular
connection after a bounded wait is safe: reconciliation/redelivery on
the next session will recover it. In contrast, indefinitely waiting
for a full queue can prevent heartbeat acknowledgements from being
read or sent, which prevents the reconnect supervisor from running.
"""
try:
await asyncio.wait_for(
outbound.put(message),
timeout=self.config.connection.outbound_enqueue_timeout,
)
except asyncio.TimeoutError as exc:
raise RuntimeError("control outbound queue is blocked") from exc
async def _handle(
self,
envelope: Any,
@@ -301,13 +386,20 @@ class ArchiveClientDaemon:
response = new_envelope()
response.correlation_id = envelope.message_id
response.heartbeat_ack.sequence = envelope.heartbeat.sequence
await outbound.put(encode(response))
await self._enqueue(outbound, encode(response))
elif payload == "command":
await self._accept_command(envelope, outbound, command_tasks)
elif payload == "protocol_error":
logger.warning(
"control_reported_protocol_error",
extra={"error_code": envelope.protocol_error.error.code},
extra={
"error_code": envelope.protocol_error.error.code,
"error_detail": envelope.protocol_error.error.message,
"offending_message_id": (
envelope.protocol_error.offending_message_id
),
"retryable": envelope.protocol_error.error.retryable,
},
)
async def _accept_command(
@@ -365,7 +457,7 @@ class ArchiveClientDaemon:
response = new_envelope()
response.correlation_id = envelope.message_id
response.command_ack.CopyFrom(acknowledgement)
await outbound.put(encode(response))
await self._enqueue(outbound, encode(response))
if (
accepted is not None
@@ -386,7 +478,7 @@ class ArchiveClientDaemon:
job_snapshot.last_event_sequence = int(
row["last_event_sequence"]
)
await outbound.put(encode(snapshot))
await self._enqueue(outbound, encode(snapshot))
elif (
accepted is not None
and accepted_for_execution
@@ -439,7 +531,6 @@ class ArchiveClientDaemon:
outbound: asyncio.Queue[str],
command_tasks: set[asyncio.Task[None]],
) -> None:
job_commands: list[control_pb2.Command] = []
for row in await asyncio.to_thread(self.store.list_accepted_commands):
acknowledgement = decode_message(
str(row["acknowledgement_json"]), control_pb2.CommandAck()
@@ -451,28 +542,14 @@ class ArchiveClientDaemon:
)
if command.WhichOneof("payload") == "ensure_route":
self._schedule_route_command(
command, "", outbound, command_tasks
command, "", outbound, command_tasks, restore_ready=True
)
elif (
command.WhichOneof("payload") in {
"assign_job", "execute_step", "cancel_job"
}
and self.jobs is not None
):
if command.WhichOneof("payload") == "cancel_job":
self.jobs.request_cancel(command.cancel_job.job_id)
job_commands.append(command)
if job_commands:
task = asyncio.create_task(
self._resume_job_commands(job_commands, outbound),
name="resume-job-commands",
)
command_tasks.add(task)
task.add_done_callback(
lambda completed: self._command_finished(
completed, command_tasks
)
)
# Job commands are intentionally not replayed here. A command
# acknowledgement is durable on both sides; the control daemon
# redelivers an unacknowledged command, while registration
# reconciliation requests snapshots for active jobs. Replaying
# every historical command re-emits overlapping event ranges and
# can flood the single connection's bounded outbound queue.
async def _resume_route_commands(
self,
@@ -543,7 +620,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
)
@@ -569,10 +648,13 @@ class ArchiveClientDaemon:
response.correlation_id = correlation_id
response.job_event.CopyFrom(event)
future = asyncio.run_coroutine_threadsafe(
outbound.put(encode(response)), loop
self._enqueue(outbound, encode(response)), loop
)
try:
future.result(timeout=30)
future.result(
timeout=self.config.connection.outbound_enqueue_timeout
+ 1
)
except Exception:
# The event is already durable in the client DB and will
# be replayed after reconnect.
@@ -592,7 +674,7 @@ class ArchiveClientDaemon:
response = new_envelope()
response.correlation_id = correlation_id
response.job_event.CopyFrom(event)
await outbound.put(encode(response))
await self._enqueue(outbound, encode(response))
def _schedule_route_command(
self,
@@ -600,12 +682,15 @@ class ArchiveClientDaemon:
correlation_id: str,
outbound: asyncio.Queue[str],
command_tasks: set[asyncio.Task[None]],
restore_ready: bool = False,
) -> None:
if command.command_id in self._active_route_commands:
return
self._active_route_commands.add(command.command_id)
task = asyncio.create_task(
self._execute_route(command, correlation_id, outbound),
self._execute_route(
command, correlation_id, outbound, restore_ready=restore_ready
),
name=f"route-{command.ensure_route.route.route_id}",
)
command_tasks.add(task)
@@ -636,13 +721,14 @@ class ArchiveClientDaemon:
response = new_envelope()
response.correlation_id = correlation_id
response.inventory_chunk.CopyFrom(chunk)
await outbound.put(encode(response))
await self._enqueue(outbound, encode(response))
async def _execute_route(
self,
command: Any,
correlation_id: str,
outbound: asyncio.Queue[str],
restore_ready: bool = False,
) -> None:
assert self.routes is not None
spec = command.ensure_route.route
@@ -661,8 +747,21 @@ class ArchiveClientDaemon:
response.route_update.CopyFrom(
decode_message(str(row["update_json"]), control_pb2.RouteUpdate())
)
await outbound.put(encode(response))
if attempt["state"] in {"ready", "failed"}:
await self._enqueue(outbound, encode(response))
if attempt["state"] == "ready":
if restore_ready:
# Route readiness is durable, but this process-local lookup is
# not. Reconfigure (which validates the existing Syncthing
# folder and recreates its local directory if necessary) after
# a daemon restart, before replaying stored READY updates.
configured = await asyncio.to_thread(
self.routes.configure,
spec,
time.monotonic() + spec.setup_timeout_seconds,
)
self._known_route_paths[spec.route_id] = configured.local_path
return
if attempt["state"] == "failed":
return
deadline = time.monotonic() + spec.setup_timeout_seconds
@@ -813,7 +912,7 @@ class ArchiveClientDaemon:
response = new_envelope()
response.correlation_id = correlation_id
response.route_update.CopyFrom(update)
await outbound.put(encode(response))
await self._enqueue(outbound, encode(response))
def _local_route(
self, spec: route_pb2.EnsureRouteSpec, state: int
+146 -30
View File
@@ -4,6 +4,8 @@ from __future__ import annotations
import hashlib
import json
import logging
import os
import shutil
import stat
import threading
@@ -26,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,
@@ -52,6 +55,9 @@ class JobCancelled(JobExecutionError):
pass
logger = logging.getLogger(__name__)
class ClientJobExecutor:
def __init__(
self,
@@ -66,7 +72,7 @@ class ClientJobExecutor:
sparse_supported: bool,
poll_interval: float = 1,
verification_timeout: float = 30 * 60,
free_space_reserve_bytes: int = 1024 * 1024 * 1024,
free_space_reserve_bytes: int = 32 * 1024 * 1024,
):
self.client_id = client_id
self.qbittorrent = qbittorrent
@@ -80,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()
@@ -109,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,
@@ -126,32 +154,51 @@ class ClientJobExecutor:
for event in replay:
event_callback(event)
return replay
started = replay[0] if replay else self._event(
definition,
sequence=command.expected_last_event_sequence + 1,
revision=command.expected_job_revision + 1,
event_type=control_pb2.JOB_EVENT_TYPE_STEP_STARTED,
state=job_pb2.JOB_STATE_RUNNING,
committed=(
self._committed(command.job_id)
or command.step == job_pb2.JOB_STEP_KIND_STAGING_CLEANUP
or (
command.step == job_pb2.JOB_STEP_KIND_SAFE_FILE_UNLINK
and self.store.get_job_artifact(
command.job_id, "qb-entry-removed"
) is not None
)
),
step=command.step,
step_state=job_pb2.STEP_STATE_RUNNING,
)
if not replay:
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=(
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=(
self._committed(command.job_id)
or command.step == job_pb2.JOB_STEP_KIND_STAGING_CLEANUP
or (
command.step == job_pb2.JOB_STEP_KIND_SAFE_FILE_UNLINK
and self.store.get_job_artifact(
command.job_id, "qb-entry-removed"
) is not None
)
),
step=command.step,
step_state=job_pb2.STEP_STATE_RUNNING,
)
self._record(definition, started)
cursor = started
emitted = [started]
if event_callback is not None:
event_callback(started)
emitted = [started]
cursor = started
speed_sample = [time.monotonic(), 0]
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,
@@ -161,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
@@ -578,7 +629,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,
@@ -610,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(
@@ -818,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():
@@ -893,6 +968,7 @@ class ClientJobExecutor:
if (
not stat.S_ISREG(metadata.st_mode)
or metadata.st_dev != destination_device
or _mount_id(source) != _mount_id(destination_root)
):
required += logical_bytes
return required
@@ -1143,6 +1219,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):
@@ -1150,3 +1228,41 @@ def _job_error_code(error: Exception) -> int:
if isinstance(error, EvictionError):
return common_pb2.ERROR_CODE_PRECONDITION_FAILED
return common_pb2.ERROR_CODE_INTERNAL
def _mount_id(path: Path) -> str | None:
"""Return Linux's effective mount ID for a path when procfs is available.
Bind mounts can share ``st_dev`` while still rejecting ``link(2)`` with
``EXDEV``. Mount IDs distinguish that case without creating probe files
inside a Syncthing folder.
"""
try:
target = os.path.realpath(path)
best: tuple[int, str] | None = None
with open("/proc/self/mountinfo", encoding="utf-8") as source:
for line in source:
fields = line.rstrip("\n").split(" ")
if len(fields) < 5:
continue
mountpoint = _unescape_mount_path(fields[4])
if target != mountpoint and not target.startswith(
mountpoint.rstrip("/") + "/"
):
continue
candidate = (len(mountpoint), fields[0])
if best is None or candidate[0] > best[0]:
best = candidate
return None if best is None else best[1]
except OSError:
return None
def _unescape_mount_path(value: str) -> str:
return (
value.replace("\\040", " ")
.replace("\\011", "\t")
.replace("\\012", "\n")
.replace("\\134", "\\")
)
+14 -1
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")
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
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:
+3 -1
View File
@@ -37,7 +37,9 @@ def discover_routes(
relative = normalized_api_path.relative_to(roots.api_root)
except (ConfigError, ValueError):
continue
if not _lexically_within(local_path, roots.local_root):
if not _lexically_within(
local_path, roots.local_root_for_api(normalized_api_path.as_posix())
):
continue
if relative == PurePosixPath("."):
continue
+32 -8
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)
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
@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 = 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:
+83 -17
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),
@@ -356,7 +384,9 @@ class SyncthingRouteManager:
local_path = self.config.roots.api_to_local(api_path)
except ConfigError as exc:
raise RoutePathConflict("route path is outside the sync root") from exc
resolved_root = self.config.local_root.resolve(strict=False)
resolved_root = self.config.roots.local_root_for_api(api_path).resolve(
strict=False
)
try:
local_path.resolve(strict=False).relative_to(resolved_root)
except ValueError as exc:
@@ -470,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,
+17 -1
View File
@@ -1,7 +1,7 @@
import os
import tempfile
import unittest
from pathlib import Path
from pathlib import Path, PurePosixPath
from unittest.mock import patch
from archive_clients.config import ClientConfig, ConfigError, RootMapping
@@ -34,6 +34,7 @@ class ConfigTests(unittest.TestCase):
root / "qb/a/b",
)
self.assertEqual(config.jobs.stall_after, 30 * 60)
self.assertEqual(config.jobs.free_space_reserve_bytes, 32 * 1024 * 1024)
self.assertEqual(
config.syncthing.advertised_addresses, ("dynamic",)
)
@@ -62,6 +63,21 @@ class ConfigTests(unittest.TestCase):
with self.assertRaises(ConfigError):
mapping.api_to_local("/elsewhere/file")
def test_mapping_can_override_one_syncthing_folder_locally(self):
mapping = RootMapping(
PurePosixPath("/sync"),
Path("/local/sync"),
((PurePosixPath("/sync/DownloadsSync"), Path("/local/qb/Sync")),),
)
self.assertEqual(
mapping.api_to_local("/sync/DownloadsSync/job/ready.json"),
Path("/local/qb/Sync/job/ready.json"),
)
self.assertEqual(
mapping.local_root_for_api("/sync/DownloadsSync"),
Path("/local/qb/Sync"),
)
def test_endpoint_scheme_and_job_keys_are_strict(self):
with tempfile.TemporaryDirectory() as directory:
root = Path(directory)
+209
View File
@@ -22,6 +22,211 @@ 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_failed_writer_ends_connection_without_waiting_for_heartbeat(self):
"""A send failure must immediately reach the reconnect supervisor."""
async def control(websocket):
registration = decode(await websocket.recv())
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))
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,
)
probe = FilesystemProbe(root, True, True, True, True, True)
daemon = ArchiveClientDaemon(config, [probe, probe], [])
await asyncio.to_thread(daemon.store.initialize)
async def failed_writer(websocket, outbound):
raise OSError("simulated broken socket")
daemon._writer = failed_writer
with self.assertRaisesRegex(RuntimeError, "writer failed"):
await asyncio.wait_for(daemon._connection(), 1)
async def test_full_outbound_queue_aborts_connection_instead_of_blocking_heartbeats(self):
"""Bulk output cannot indefinitely block the receive/heartbeat loop."""
async def control(websocket):
registration = decode(await websocket.recv())
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))
for sequence in range(1, 102):
heartbeat = new_envelope()
heartbeat.heartbeat.sequence = sequence
await websocket.send(encode(heartbeat))
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(outbound_enqueue_timeout=0.01),
)
probe = FilesystemProbe(root, True, True, True, True, True)
daemon = ArchiveClientDaemon(config, [probe, probe], [])
await asyncio.to_thread(daemon.store.initialize)
async def stopped_writer(websocket, outbound):
await asyncio.Event().wait()
daemon._writer = stopped_writer
with self.assertRaisesRegex(RuntimeError, "outbound queue is blocked"):
await asyncio.wait_for(daemon._connection(), 2)
async def test_reconnect_resume_skips_historical_job_commands(self):
"""Registration reconciliation, not command replay, recovers job state."""
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], [])
await asyncio.to_thread(daemon.store.initialize)
command = control_pb2.Command(command_id=str(uuid4()))
command.execute_step.job_id = str(uuid4())
command.execute_step.expected_last_event_sequence = 1
acknowledgement = control_pb2.CommandAck(
command_id=command.command_id,
status=control_pb2.COMMAND_ACK_STATUS_ACCEPTED,
)
await asyncio.to_thread(
daemon.store.accept_command,
command.command_id,
encode_message(command),
encode_message(acknowledgement),
)
daemon.jobs = Mock()
outbound = asyncio.Queue()
tasks = set()
await daemon._resume_commands(outbound, tasks)
self.assertEqual(tasks, set())
self.assertTrue(outbound.empty())
async def test_eviction_assignment_and_steps_are_admitted(self):
with tempfile.TemporaryDirectory() as directory:
root = Path(directory)
@@ -192,6 +397,10 @@ class DaemonTransportTests(unittest.IsolatedAsyncioTestCase):
for _ in range(3)
]
self.assertEqual([item.sequence for item in resumed], [1, 2, 3])
self.assertEqual(manager.configure.call_count, 2)
self.assertEqual(
restarted._route_path("route-1"), root / "routes/route-1"
)
async def test_slow_inventory_does_not_block_heartbeat(self):
with tempfile.TemporaryDirectory() as directory:
+83
View File
@@ -0,0 +1,83 @@
import importlib.util
import sys
import unittest
from unittest import mock
from pathlib import Path
_SCRIPT = Path(__file__).parents[1] / "scripts" / "preflight-deployment.py"
_SPEC = importlib.util.spec_from_file_location("deployment_preflight", _SCRIPT)
assert _SPEC is not None and _SPEC.loader is not None
preflight = importlib.util.module_from_spec(_SPEC)
sys.modules[_SPEC.name] = preflight
_SPEC.loader.exec_module(preflight)
class DeploymentPreflightTests(unittest.TestCase):
def test_maps_nested_container_path_to_host_source(self):
container = {
"Mounts": [
{"Type": "bind", "Source": "/srv/data", "Destination": "/data"},
{"Type": "bind", "Source": "/srv/sync", "Destination": "/data/Sync"},
]
}
mapping = preflight.map_path(container, "/data/Sync/routes/route-1")
self.assertEqual(mapping.source, Path("/srv/sync/routes/route-1"))
self.assertEqual(str(mapping.destination), "/data/Sync")
def test_rejects_api_and_client_roots_from_different_host_paths(self):
with self.assertRaisesRegex(preflight.CheckFailure, "host paths differ"):
preflight.require_same_path(
"Syncthing api_root/local_root",
preflight.Mapping(Path("/srv/syncthing-config/routes"), Path("/var/syncthing")),
preflight.Mapping(Path("/srv/sync/routes"), Path("/data/storage")),
)
def test_rejects_unmounted_container_path(self):
with self.assertRaisesRegex(preflight.CheckFailure, "not backed"):
preflight.map_path({"Mounts": []}, "/missing")
def test_maps_future_route_through_syncthing_override(self):
config = {
"api_root": "/var/syncthing",
"local_root": "/data/sync",
"local_path_overrides": {
"/var/syncthing/routes": "/data/qb/.archive-control-routes",
},
}
self.assertEqual(
preflight.route_local_path(config),
"/data/qb/.archive-control-routes",
)
def test_runs_hardlink_probe_with_qb_and_route_roots(self):
completed = __import__("subprocess").CompletedProcess(
args=[], returncode=0, stdout="hard-link staging probe passed\n", stderr=""
)
with mock.patch.object(preflight.subprocess, "run", return_value=completed) as run:
preflight.run_hardlink_probe(
"client", "/data/qb", "/data/qb/routes", "1001:1001"
)
args = run.call_args.args[0]
self.assertEqual(args[:6], ["docker", "exec", "--user", "1001:1001", "client", "python"])
self.assertEqual(args[-2:], ["/data/qb", "/data/qb/routes"])
def test_hardlink_probe_failure_is_actionable(self):
completed = __import__("subprocess").CompletedProcess(
args=[], returncode=1, stdout="", stderr="[Errno 18] Invalid cross-device link"
)
with mock.patch.object(preflight.subprocess, "run", return_value=completed):
with self.assertRaisesRegex(preflight.CheckFailure, "hard-link staging probe failed"):
preflight.run_hardlink_probe("client", "/data/qb", "/data/qb/routes")
def test_existing_route_directory_failure_is_actionable(self):
completed = __import__("subprocess").CompletedProcess(
args=[], returncode=1, stdout='{"missing":[{"folder_id":"route-1"}]}', stderr=""
)
with mock.patch.object(preflight.subprocess, "run", return_value=completed):
with self.assertRaisesRegex(preflight.CheckFailure, "route directory is missing"):
preflight.run_existing_route_directory_check("client", "/config.toml")
if __name__ == "__main__":
unittest.main()
+260 -2
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,6 +12,7 @@ from archive_clients.jobs import (
JobExecutionError,
_resource_fingerprint,
)
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
@@ -31,7 +33,126 @@ 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_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)
@@ -74,10 +195,146 @@ class ClientJobHappyPathTests(unittest.TestCase):
self._run_transfer(
Path(directory),
job_pb2.JOB_OPERATION_ARCHIVE,
source_stage_free_bytes=1024 * 1024 * 1024 + 1024,
source_stage_free_bytes=32 * 1024 * 1024 + 1024,
content=b"x" * 4096,
)
def test_mount_boundary_requires_copy_space_even_with_same_device(self):
with tempfile.TemporaryDirectory() as directory:
root = Path(directory)
source_root = root / "source"
destination_root = root / "destination"
source_root.mkdir()
destination_root.mkdir()
(source_root / "fixture.bin").write_bytes(b"fixture")
with patch(
"archive_clients.jobs._mount_id",
side_effect=("source-mount", "destination-mount"),
):
required = ClientJobExecutor._copy_required_bytes(
source_root,
destination_root,
(("fixture.bin", 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 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,
@@ -85,6 +342,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"
@@ -152,7 +410,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,
+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,
+26
View File
@@ -62,6 +62,32 @@ class RouteDiscoveryTests(unittest.TestCase):
self.assertEqual(routes[0].local_relative_path, "DownloadsSync")
self.assertEqual(routes[0].state, route_pb2.ROUTE_STATE_DISCOVERED)
def test_discovery_accepts_a_safe_local_folder_override(self):
with tempfile.TemporaryDirectory() as directory:
root = Path(directory)
override = root / "qb/Sync"
override.mkdir(parents=True)
roots = RootMapping(
PurePosixPath("/var/syncthing"),
root / "sync",
((PurePosixPath("/var/syncthing/DownloadsSync"), override),),
)
routes = discover_routes(
{"folders": [{
"id": "DownloadsSync",
"path": "~/DownloadsSync",
"type": "sendreceive",
"devices": [
{"deviceID": "LOCAL"},
{"deviceID": "ARCHIVE"},
],
}]},
"LOCAL",
roots,
True,
)
self.assertEqual([route.route_id for route in routes], ["DownloadsSync"])
if __name__ == "__main__":
unittest.main()
+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)