Compare commits

...
16 Commits
36 changed files with 1203 additions and 47 deletions
+6 -2
View File
@@ -11,8 +11,12 @@ CMD ["python", "-m", "unittest", "discover", "-s", "tests", "-v"]
FROM python:3.11-slim FROM python:3.11-slim
ENV PYTHONDONTWRITEBYTECODE=1 PYTHONUNBUFFERED=1 ENV PYTHONDONTWRITEBYTECODE=1 PYTHONUNBUFFERED=1
RUN groupadd --gid 1001 archive-control \ RUN if ! getent group 1001 >/dev/null; then \
&& useradd --uid 1001 --gid 1001 --no-create-home archive-control groupadd --gid 1001 archive-control; \
fi \
&& if ! getent passwd 1001 >/dev/null; then \
useradd --uid 1001 --gid 1001 --no-create-home archive-control; \
fi
COPY --from=builder /wheels /wheels COPY --from=builder /wheels /wheels
RUN pip install --no-cache-dir /wheels/*.whl && rm -rf /wheels RUN pip install --no-cache-dir /wheels/*.whl && rm -rf /wheels
USER 1001:1001 USER 1001:1001
+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 adversarial matrix without live infrastructure or Telegram. See
`e2e/README.md`. `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, The complete cross-project design, protocol, workflow, safety, deployment,
Telegram UX, and testing documentation is published under [`docs/`](docs/). Telegram UX, and testing documentation is published under [`docs/`](docs/).
+28 -11
View File
@@ -1,12 +1,15 @@
# Production deployment # Production deployment
These files are the non-secret, host-specific deployment manifests for the 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 - Install the x1 and x2 files as
`~/compose/ArchiveControl-cache/{compose.yaml,client.toml}`. `~/compose/ArchiveControl-cache/{compose.yaml,client.toml}`.
- Install the lithium files as - Install the lithium files as
`~/compose/ArchiveControl-archive/{compose.yaml,client.toml}`. `~/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 - Create sibling `state`, `backups`, and `secrets` directories owned by the
configured container UID/GID. configured container UID/GID.
- Secret files are never committed. Each `secrets` directory contains - Secret files are never committed. Each `secrets` directory contains
@@ -22,10 +25,17 @@ 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 client deliberately treats that case as copy-only and performs a full payload
free-space check. free-space check.
When a Syncthing route is physically nested below the qB root, mount the qB Automatic routes always use the Syncthing API path `routes/<route-id>`. Bind
root once and map the exact Syncthing API folder through it: that path from a dedicated directory beneath the qB data root, then map the
same path through the qB client mount:
```yaml ```yaml
# Syncthing compose project
volumes:
- /srv/syncthing-config:/var/syncthing
- /srv/downloads/.archive-control-routes:/var/syncthing/routes
# Archive Control client compose project
volumes: volumes:
- /srv/downloads:/data/qb - /srv/downloads:/data/qb
- /srv/syncthing-config:/data/sync - /srv/syncthing-config:/data/sync
@@ -35,23 +45,30 @@ volumes:
[syncthing] [syncthing]
api_root = "/var/syncthing" api_root = "/var/syncthing"
local_root = "/data/sync" local_root = "/data/sync"
local_path_overrides = { "/var/syncthing/Downloads/Sync" = "/data/qb/Sync" } local_path_overrides = { "/var/syncthing/routes" = "/data/qb/.archive-control-routes" }
``` ```
Do **not** additionally mount `/srv/downloads/Sync` at a path beneath Do **not** mount the route directory separately into the client. The override
`/data/sync`. The override is the authoritative mapping for that folder and is the authoritative mapping and keeps qB source files and automatic route
keeps qB source files and staging destinations in one mount namespace. Use folders in one mount namespace. Existing manually configured Syncthing folders
the folder's normalized API-visible **path** as the override key (for example, below the qB tree may retain their own exact overrides. Use API-visible paths,
Syncthing `~/Downloads/Sync` becomes `/var/syncthing/Downloads/Sync`); do not not Syncthing folder IDs, as override keys.
use the folder ID.
Before starting a stack, validate it with: Before starting a stack, validate its config and then run the generic host
preflight with that machine's own paths and container names:
```sh ```sh
docker compose config docker compose config
docker compose run --rm archive-client --check-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, The check is fail-fast and performs local permission, filesystem, sparse-file,
hard-link, and reflink probes. Normal startup additionally probes the local hard-link, and reflink probes. Normal startup additionally probes the local
qBittorrent and Syncthing APIs before registration. 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
+1
View File
@@ -39,4 +39,5 @@ endpoint = "http://syncthing:8384"
api_key_file = "/run/secrets/syncthing_api_key" api_key_file = "/run/secrets/syncthing_api_key"
api_root = "/var/syncthing" api_root = "/var/syncthing"
local_root = "/data/sync" local_root = "/data/sync"
local_path_overrides = { "/var/syncthing/routes" = "/data/qb/.archive-control-routes" }
advertised_addresses = ["dynamic"] advertised_addresses = ["dynamic"]
+1 -1
View File
@@ -2,7 +2,7 @@ name: archive-control-archive
services: services:
archive-client: archive-client:
image: sodium/archive-clients:v0.1.13 image: sodium/archive-clients:v0.1.14
user: "1000:1000" user: "1000:1000"
restart: unless-stopped restart: unless-stopped
command: ["--config", "/etc/archive-control/client.toml"] command: ["--config", "/etc/archive-control/client.toml"]
+1 -1
View File
@@ -41,5 +41,5 @@ api_root = "/var/syncthing"
local_root = "/data/sync" local_root = "/data/sync"
# This existing folder is physically inside the qB data tree. Map it through # This existing folder is physically inside the qB data tree. Map it through
# that same client bind mount so source staging can use hardlinks. # that same client bind mount so source staging can use hardlinks.
local_path_overrides = { "/var/syncthing/DownloadsSync" = "/data/qb/Sync" } local_path_overrides = { "/var/syncthing/DownloadsSync" = "/data/qb/Sync", "/var/syncthing/routes" = "/data/qb/.archive-control-routes" }
advertised_addresses = ["dynamic"] advertised_addresses = ["dynamic"]
+1 -1
View File
@@ -2,7 +2,7 @@ name: archive-control-cache
services: services:
archive-client: archive-client:
image: sodium/archive-clients:v0.1.13 image: sodium/archive-clients:v0.1.14
user: "1001:1001" user: "1001:1001"
restart: unless-stopped restart: unless-stopped
network_mode: host network_mode: host
+1 -1
View File
@@ -42,5 +42,5 @@ local_root = "/data/sync"
# `DownloadsSync-X2` is ~/Downloads/Sync on the host, nested below the qB # `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 # root. Resolve it through /data/qb rather than a second nested bind mount so
# source staging can hardlink it. # source staging can hardlink it.
local_path_overrides = { "/var/syncthing/Downloads/Sync" = "/data/qb/Sync" } local_path_overrides = { "/var/syncthing/Downloads/Sync" = "/data/qb/Sync", "/var/syncthing/routes" = "/data/qb/.archive-control-routes" }
advertised_addresses = ["dynamic"] advertised_addresses = ["dynamic"]
+1 -1
View File
@@ -2,7 +2,7 @@ name: archive-control-cache
services: services:
archive-client: archive-client:
image: sodium/archive-clients:v0.1.13 image: sodium/archive-clients:v0.1.14
user: "1001:1001" user: "1001:1001"
restart: unless-stopped restart: unless-stopped
network_mode: host network_mode: host
+3 -2
View File
@@ -16,8 +16,9 @@ credentials and host-specific details.
6. [Storage safety and recovery](storage-safety.md) 6. [Storage safety and recovery](storage-safety.md)
7. [qBittorrent and Syncthing integration](service-apis.md) 7. [qBittorrent and Syncthing integration](service-apis.md)
8. [Configuration, deployment, and usage](deployment-and-usage.md) 8. [Configuration, deployment, and usage](deployment-and-usage.md)
9. [Telegram UX](telegram-ux.md) 9. [Deployment preflight tests](deployment-preflight.md)
10. [Testing strategy](testing.md) 10. [Telegram UX](telegram-ux.md)
11. [Testing strategy](testing.md)
The independent protobuf source of truth lives in the The independent protobuf source of truth lives in the
`cabbage/archive-control-proto` repository. Generated bindings in this `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. manifests, and observed external state.
- SQLite backups use the online backup API. Defaults are every six hours, 12 - 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. 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. - Default concurrency is unrestricted: every eligible queued job is admitted
Disjoint node pairs may run concurrently. Queueing is durable FIFO. 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 - A queued job owns a per-resource reservation. Cancelling it removes only the
queued record/reservation and never sends cleanup commands. queued record/reservation and never sends cleanup commands.
- Connectivity or transfer stalls wait indefinitely. A configurable 30-minute - Connectivity or transfer stalls wait indefinitely. A configurable 30-minute
+26 -1
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. stall_after = "30m" # Warning state only; jobs continue waiting.
route_policy = "on_demand" # Alternative: eager_mesh. route_policy = "on_demand" # Alternative: eager_mesh.
route_setup_timeout = "30m" 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_client = 1
max_active_per_route = 1 max_active_per_route = 1
max_envelope_bytes = 1048576 max_envelope_bytes = 1048576
@@ -223,7 +226,7 @@ cache/archive routes according to policy.
```yaml ```yaml
services: services:
archive-client: archive-client:
image: sodium/archive-clients:v0.1.13 image: sodium/archive-clients:v0.1.14
user: "1001:1001" user: "1001:1001"
restart: unless-stopped restart: unless-stopped
command: ["archive-client", "--config", "/etc/archive-control/client.toml"] command: ["archive-client", "--config", "/etc/archive-control/client.toml"]
@@ -256,6 +259,28 @@ Generated protobuf Python bindings are committed in each consumer with the
exact `archive-control-proto` tag/commit recorded. Release order is proto, exact `archive-control-proto` tag/commit recorded. Release order is proto,
control consumer, client consumer, E2E, then image publication. control consumer, client consumer, E2E, then image publication.
### Reproducible multi-platform Buildx lifecycle
The named Buildx builders are local acceleration/cache only; they are not a
deployment dependency and can be removed after publication. To create a fresh
builder, verify its platforms, publish a release, and remove it afterwards:
```bash
docker buildx create --name archive-control-release --driver docker-container --use
docker buildx inspect --bootstrap
docker buildx build --platform linux/amd64,linux/arm64 \
--tag sodium/archive-clients:vX.Y.Z --push .
docker buildx rm archive-control-release
```
`docker buildx inspect` must list both `linux/amd64` and `linux/arm64` before
publishing. If the host has no arm64 emulation, install/configure it according
to the host Docker distribution before the build; do not publish a partial
single-platform tag. Retain the pushed manifest digest in the release notes
and deploy the immutable tag or digest. The optional `archive-control-qemu`
builder follows the same lifecycle when it is used for an emulation smoke
build.
## Operator usage ## Operator usage
1. Prepare local qB, sync, state, backup, config, and secret mounts with the 1. Prepare local qB, sync, state, backup, config, and secret mounts with the
+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.
+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 baselines even when jobs would use different nodes or observe different sides
of a hybrid identity. of a hybrid identity.
Defaults allow one active data-moving job per client and one per route. A job By default every eligible queued job is claimed in the same scheduler pass and
must acquire its source client, target client, route, and resource reservation independent jobs run concurrently on each client; physical service and storage
atomically. Disjoint node pairs may run concurrently. Route setup is a capacity are the limit. Set `enforce_concurrency_limits=true` on control to
preflight activity and does not permit a data step to bypass these leases. 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 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 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 4. Control observes relevant qBittorrent, Syncthing, staging, and manifest
state before selecting retry, resume, compensation, cleanup, or manual state before selecting retry, resume, compensation, cleanup, or manual
intervention. 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 5. A command is reissued with its original ID when the acceptance result is
uncertain. uncertain.
+4 -2
View File
@@ -25,8 +25,10 @@ What do you want to do?
[ Evict Cache ] [ Job Status ] [ Evict Cache ] [ Job Status ]
``` ```
All subsequent pages edit this message. `Cancel` closes the active selection All subsequent pages edit this message. Leaving a resource/tree selection
flow. `Back` returns one level while retaining validated filters/selections. 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 Callback payloads contain opaque session/action IDs, not resource names or
paths, and are validated against persisted session revision and expiry. 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 | | 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 | | 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 | | 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 | | 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 | | UI | pagination/filter chain, stale/double callback, persisted session, clear vs evict separation |
+1 -1
View File
@@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta"
[project] [project]
name = "archive-clients" name = "archive-clients"
version = "0.1.13" version = "0.1.19"
requires-python = ">=3.11" requires-python = ">=3.11"
dependencies = ["protobuf==7.35.1", "websockets==16.0"] 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] components = [name] + [_component(part) for part in raw_path]
result.append(MetaFile( result.append(MetaFile(
"/".join(components), _length(item.get(b"length")), "/".join(components), _length(item.get(b"length")),
_padding(item), _padding(item) or _bitcomet_padding_name(components[-1]),
)) ))
return result return result
@@ -223,3 +223,13 @@ def _length(value: Any) -> int:
def _padding(value: dict[bytes, Any]) -> bool: def _padding(value: dict[bytes, Any]) -> bool:
attributes = value.get(b"attr", b"") attributes = value.get(b"attr", b"")
return isinstance(attributes, bytes) and b"p" in attributes 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_")
+36 -3
View File
@@ -43,6 +43,19 @@ from archive_control.v1 import (
logger = logging.getLogger(__name__) 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: class ArchiveClientDaemon:
def __init__( def __init__(
self, self,
@@ -95,7 +108,10 @@ class ArchiveClientDaemon:
self._lease = DatabaseLease(config.state_db) self._lease = DatabaseLease(config.state_db)
self._active_route_commands: set[str] = set() self._active_route_commands: set[str] = set()
self._active_job_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 = ( self.jobs = (
ClientJobExecutor( ClientJobExecutor(
client_id=config.client_id, client_id=config.client_id,
@@ -560,7 +576,9 @@ class ArchiveClientDaemon:
correlation_id: str, correlation_id: str,
outbound: asyncio.Queue[str], outbound: asyncio.Queue[str],
) -> None: ) -> 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( await self._execute_job_command_locked(
command, correlation_id, outbound command, correlation_id, outbound
) )
@@ -679,7 +697,22 @@ class ArchiveClientDaemon:
decode_message(str(row["update_json"]), control_pb2.RouteUpdate()) decode_message(str(row["update_json"]), control_pb2.RouteUpdate())
) )
await outbound.put(encode(response)) await outbound.put(encode(response))
if attempt["state"] in {"ready", "failed"}: if attempt["state"] == "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) before replaying
# the stored READY updates. This is essential after a daemon or
# mount restart: a previously ready route may otherwise point at a
# path that is no longer present, and a later source-stage command
# would fail with a bare ENOENT.
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 return
deadline = time.monotonic() + spec.setup_timeout_seconds deadline = time.monotonic() + spec.setup_timeout_seconds
+27
View File
@@ -86,6 +86,8 @@ class ClientJobExecutor:
self.verification_timeout = verification_timeout self.verification_timeout = verification_timeout
self.free_space_reserve_bytes = free_space_reserve_bytes self.free_space_reserve_bytes = free_space_reserve_bytes
self._cancel_events: dict[str, threading.Event] = {} 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: def request_cancel(self, job_id: str) -> None:
self._cancel_events.setdefault(job_id, threading.Event()).set() self._cancel_events.setdefault(job_id, threading.Event()).set()
@@ -115,6 +117,22 @@ class ClientJobExecutor:
self, self,
command: control_pb2.ExecuteStepCommand, command: control_pb2.ExecuteStepCommand,
event_callback: Callable[[control_pb2.JobEvent], None] | None = None, 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]: ) -> list[control_pb2.JobEvent]:
definition = self._definition(command.job_id) definition = self._definition(command.job_id)
replay = self._replay( replay = self._replay(
@@ -653,6 +671,15 @@ class ClientJobExecutor:
# Syncthing can expose the per-job directory before the # Syncthing can expose the per-job directory before the
# ready marker and manifest have arrived atomically as a set. # ready marker and manifest have arrived atomically as a set.
pass pass
except RouteSetupError as exc:
# A network interruption can temporarily make the local
# Syncthing REST API unavailable. The staged payload is
# durable and the job must remain recoverable; keep polling
# instead of turning a transient outage into a failed job.
logger.warning("syncthing_status_deferred", extra={
"job_id": definition.job_id,
"error_type": type(exc).__name__,
})
time.sleep(self.poll_interval) time.sleep(self.poll_interval)
def _target_materialize( def _target_materialize(
+14 -1
View File
@@ -3,6 +3,7 @@
from __future__ import annotations from __future__ import annotations
import json import json
import logging
import threading import threading
import time import time
import uuid import uuid
@@ -16,6 +17,8 @@ from urllib import error, parse, request
from archive_clients.config import ServiceConfig from archive_clients.config import ServiceConfig
from archive_clients.resources import NormalizedResource, normalize_resource from archive_clients.resources import NormalizedResource, normalize_resource
logger = logging.getLogger(__name__)
class QBittorrentError(RuntimeError): class QBittorrentError(RuntimeError):
pass pass
@@ -72,7 +75,17 @@ class QBittorrentReader:
torrent_hash = torrent.get("hash") torrent_hash = torrent.get("hash")
if not isinstance(torrent_hash, str): if not isinstance(torrent_hash, str):
raise QBittorrentError("qBittorrent torrent hash is invalid") raise QBittorrentError("qBittorrent torrent hash is invalid")
result.append(self._normalize(torrent, torrent_hash)) try:
result.append(self._normalize(torrent, torrent_hash))
except ValueError as exc:
# A stale or malformed qBittorrent entry must not hide every
# otherwise valid resource from Archive/Evict inventory.
# Keep it ineligible and leave an operator-visible diagnosis.
logger.warning(
"Skipping qBittorrent resource with inconsistent "
"metadata: hash=%s name=%r reason=%s",
torrent_hash, torrent.get("name"), exc,
)
return result return result
def get_resource(self, torrent_hash: str) -> NormalizedResource | None: def get_resource(self, torrent_hash: str) -> NormalizedResource | None:
+11 -2
View File
@@ -91,12 +91,21 @@ def normalize_resource(
observed_at: datetime | None = None, observed_at: datetime | None = None,
) -> NormalizedResource: ) -> NormalizedResource:
metainfo = decode_metainfo(metainfo_bytes) 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") raise ResourceError("qBittorrent and metainfo file counts differ")
files = [] files = []
canonical = True canonical = True
for expected_index, (raw, meta_file) in enumerate( 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") index = _integer(raw.get("index"), "file index")
if index != expected_index: if index != expected_index:
+32 -8
View File
@@ -7,6 +7,8 @@ import json
import os import os
import sqlite3 import sqlite3
import stat import stat
import threading
from contextlib import contextmanager
from dataclasses import dataclass from dataclasses import dataclass
from pathlib import Path from pathlib import Path
@@ -14,6 +16,12 @@ from pathlib import Path
SCHEMA_VERSION = 3 SCHEMA_VERSION = 3
# A client daemon executes unrelated jobs concurrently. SQLite still permits
# only one writer, so serialize this process's short state transactions rather
# than allowing an otherwise healthy job to fail after its busy timeout.
_DATABASE_LOCK = threading.RLock()
class CommandConflict(RuntimeError): class CommandConflict(RuntimeError):
pass pass
@@ -557,14 +565,30 @@ class ClientStore:
).fetchall() ).fetchall()
return [dict(row) for row in rows] return [dict(row) for row in rows]
def _connect(self) -> sqlite3.Connection: @contextmanager
connection = sqlite3.connect(self.database, isolation_level=None, timeout=5) def _connect(self):
connection.row_factory = sqlite3.Row """Yield one connection while serializing local SQLite writers.
connection.execute("PRAGMA foreign_keys = ON")
connection.execute("PRAGMA journal_mode = WAL") The longer SQLite timeout also covers a short lock held by a separate
connection.execute("PRAGMA synchronous = FULL") maintenance process such as a backup. The lock is deliberately held
connection.execute("PRAGMA busy_timeout = 5000") for the full transaction, including ``BEGIN IMMEDIATE``.
return connection """
with _DATABASE_LOCK:
connection = sqlite3.connect(
self.database, isolation_level=None, timeout=30
)
try:
connection.row_factory = sqlite3.Row
connection.execute("PRAGMA foreign_keys = ON")
connection.execute("PRAGMA journal_mode = WAL")
connection.execute("PRAGMA synchronous = FULL")
connection.execute("PRAGMA busy_timeout = 30000")
# Preserve the original ``with connection`` commit/rollback
# behavior used by every store operation.
with connection:
yield connection
finally:
connection.close()
def _canonical(value: object) -> str: def _canonical(value: object) -> str:
+45
View File
@@ -22,6 +22,47 @@ from archive_control.v1 import (
class DaemonTransportTests(unittest.IsolatedAsyncioTestCase): 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): async def test_silent_control_connection_ends_for_reconnect(self):
"""A lost server heartbeat must not leave durable commands stranded.""" """A lost server heartbeat must not leave durable commands stranded."""
@@ -239,6 +280,10 @@ class DaemonTransportTests(unittest.IsolatedAsyncioTestCase):
for _ in range(3) for _ in range(3)
] ]
self.assertEqual([item.sequence for item in resumed], [1, 2, 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): async def test_slow_inventory_does_not_block_heartbeat(self):
with tempfile.TemporaryDirectory() as directory: 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()
+116 -1
View File
@@ -1,5 +1,6 @@
import hashlib import hashlib
import tempfile import tempfile
import threading
import unittest import unittest
from pathlib import Path from pathlib import Path
from unittest.mock import Mock, patch from unittest.mock import Mock, patch
@@ -11,7 +12,7 @@ from archive_clients.jobs import (
JobExecutionError, JobExecutionError,
_resource_fingerprint, _resource_fingerprint,
) )
from archive_clients.syncthing import RouteSetupError from archive_clients.syncthing import RouteSetupError, SyncthingTransferStatus
from archive_clients.resources import normalize_resource from archive_clients.resources import normalize_resource
from archive_clients.state import ClientStore from archive_clients.state import ClientStore
from archive_control.v1 import control_pb2, job_pb2 from archive_control.v1 import control_pb2, job_pb2
@@ -38,6 +39,120 @@ class SlowRescanSyncthing(CompleteSyncthing):
class ClientJobHappyPathTests(unittest.TestCase): class ClientJobHappyPathTests(unittest.TestCase):
def test_syncthing_api_outage_during_transfer_is_retried(self):
"""A transient local REST outage must not terminally fail the job."""
with tempfile.TemporaryDirectory() as directory:
root = Path(directory)
store = ClientStore(root / "client.db")
store.initialize()
definition = job_pb2.JobDefinition(
job_id=str(uuid4()),
transfer={
"source_client_id": "cache-1",
"target_client_id": "archive-1",
"route_id": "route-1",
},
)
observer = Mock()
observer.status.side_effect = [
RouteSetupError("Syncthing API is unavailable"),
SyncthingTransferStatus(1, 42, 42, True, 0),
]
executor = ClientJobExecutor(
client_id="archive-1",
qbittorrent=Mock(),
store=store,
qb_root=root,
qb_api_root=Path("/downloads"),
route_path=lambda _: root,
syncthing_transport=Mock(),
sparse_supported=True,
poll_interval=0,
)
executor._observer = Mock(return_value=observer)
progress = Mock()
executor._wait_for_syncthing(definition, progress)
self.assertEqual(observer.status.call_count, 2)
progress.assert_called_once_with(1, 42, 42, "0 Syncthing items still needed")
def test_reconnect_replay_never_duplicates_any_transfer_step(self):
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): def test_capacity_guard_fails_before_data_movement(self):
with tempfile.TemporaryDirectory() as directory: with tempfile.TemporaryDirectory() as directory:
root = Path(directory) root = Path(directory)
+48
View File
@@ -205,6 +205,54 @@ class QBittorrentReaderTests(unittest.TestCase):
self.assertEqual(resources, []) self.assertEqual(resources, [])
self.assertEqual(len(opener.calls), 2) self.assertEqual(len(opener.calls), 2)
def test_listing_skips_malformed_torrent_and_keeps_valid_resources(self):
def torrent(name, content):
info = {
b"length": len(content), b"name": name.encode(),
b"piece length": 16384, b"pieces": b"x" * 20,
}
return encode({b"info": info}), hashlib.sha1(encode(info)).hexdigest()
first_bytes, first_hash = torrent("first.txt", b"one")
malformed_bytes, malformed_hash = torrent("broken.txt", b"bad")
second_bytes, second_hash = torrent("second.txt", b"two")
torrents = [
{"hash": first_hash, "name": "first.txt", "state": "uploading"},
{"hash": malformed_hash, "name": "broken.txt", "state": "stalledUP"},
{"hash": second_hash, "name": "second.txt", "state": "uploading"},
]
responses = [
b"Ok.", json.dumps(torrents).encode(),
b'[{"index":0,"name":"first.txt","size":3,"progress":1,"priority":1}]',
first_bytes,
b'[{"index":0,"name":"broken.txt","size":3,"progress":1,"priority":1},'
b'{"index":1,"name":"unexpected.txt","size":1,"progress":1,"priority":1}]',
malformed_bytes,
b'[{"index":0,"name":"second.txt","size":3,"progress":1,"priority":1}]',
second_bytes,
]
with tempfile.TemporaryDirectory() as directory:
root = Path(directory)
password = root / "password"
password.write_text("secret", encoding="utf-8")
os.chmod(password, 0o600)
config = ServiceConfig(
"http://qb", PurePosixPath("/downloads"), root,
username="admin", password_file=password,
)
opener = _Opener(responses)
with patch(
"archive_clients.qbittorrent.request.build_opener",
return_value=opener,
), self.assertLogs("archive_clients.qbittorrent", "WARNING") as logs:
resources = QBittorrentReader(config).list_resources()
self.assertEqual(
[resource.summary.display_name for resource in resources],
["first.txt", "second.txt"],
)
self.assertIn(malformed_hash, "\n".join(logs.output))
def test_stopped_add_selection_recheck_and_entry_only_delete(self): def test_stopped_add_selection_recheck_and_entry_only_delete(self):
torrent_hash = "a" * 40 torrent_hash = "a" * 40
responses = [ responses = [
+41
View File
@@ -118,6 +118,47 @@ class ResourceTests(unittest.TestCase):
self.assertEqual(root.available_file_count, 1) self.assertEqual(root.available_file_count, 1)
self.assertEqual(root.available_logical_bytes, 3) 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): def test_noncanonical_and_unsafe_paths_are_distinct(self):
info = { info = {
b"length": 3, b"length": 3,
+35
View File
@@ -1,6 +1,7 @@
import tempfile import tempfile
import unittest import unittest
import os import os
import threading
from pathlib import Path from pathlib import Path
from uuid import uuid4 from uuid import uuid4
@@ -13,6 +14,40 @@ from archive_clients.state import (
class ClientStoreTests(unittest.TestCase): class ClientStoreTests(unittest.TestCase):
def test_database_connections_serialize_local_writers(self):
"""Concurrent jobs share one daemon DB without SQLite lock failures."""
with tempfile.TemporaryDirectory() as directory:
database = Path(directory) / "state.db"
first = ClientStore(database)
second = ClientStore(database)
first.initialize()
barrier = threading.Barrier(2)
failures: list[Exception] = []
def write(store, command_id):
try:
barrier.wait()
store.accept_command(command_id, '{"kind":"job"}', '{"ok":true}')
except Exception as exc: # pragma: no cover - assertion below
failures.append(exc)
left = threading.Thread(target=write, args=(first, "command-1"))
right = threading.Thread(target=write, args=(second, "command-2"))
left.start()
right.start()
left.join(1)
right.join(1)
self.assertFalse(left.is_alive())
self.assertFalse(right.is_alive())
self.assertEqual(failures, [])
self.assertEqual(len(first.list_accepted_commands()), 2)
with first._connect() as connection:
self.assertEqual(
connection.execute("PRAGMA busy_timeout").fetchone()[0],
30000,
)
def test_command_acceptance_is_durable_and_content_addressed(self): def test_command_acceptance_is_durable_and_content_addressed(self):
with tempfile.TemporaryDirectory() as directory: with tempfile.TemporaryDirectory() as directory:
database = Path(directory) / "state.db" database = Path(directory) / "state.db"