Compare commits
30
Commits
v0.1.2
...
62a0f24437
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
62a0f24437 | ||
|
|
eabbda3c3a | ||
|
|
a7bc2bf087 | ||
|
|
25d40daea6 | ||
|
|
f3517bd21b | ||
|
|
300210c17d | ||
|
|
1f97faeab5 | ||
|
|
1c2590c3c8 | ||
|
|
515599f2d4 | ||
|
|
b67ca18403 | ||
|
|
6446846ee8 | ||
|
|
c717aea394 | ||
|
|
ea35ba2758 | ||
|
|
5e5e2f9095 | ||
|
|
0b11e2a3a2 | ||
|
|
1009defc46 | ||
|
|
8ed8e75447 | ||
|
|
1a06da3984 | ||
|
|
4ba5e92248 | ||
|
|
a6da224f28 | ||
|
|
90925fd321 | ||
|
|
9df2e6991e | ||
|
|
d2fa69a1d6 | ||
|
|
0f94388f49 | ||
|
|
94a3233ce2 | ||
|
|
053b0135b3 | ||
|
|
57acc90363 | ||
|
|
a73bcf870f | ||
|
|
748ca49837 | ||
|
|
589fd4a41d |
+7
-2
@@ -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
|
||||
|
||||
@@ -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/).
|
||||
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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
|
||||
@@ -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
|
||||
Executable
+18
@@ -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
|
||||
@@ -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"]
|
||||
|
||||
@@ -2,7 +2,7 @@ name: archive-control-archive
|
||||
|
||||
services:
|
||||
archive-client:
|
||||
image: sodium/archive-clients:v0.1.2
|
||||
image: sodium/archive-clients:v0.1.14
|
||||
user: "1000:1000"
|
||||
restart: unless-stopped
|
||||
command: ["--config", "/etc/archive-control/client.toml"]
|
||||
|
||||
@@ -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"]
|
||||
|
||||
@@ -2,7 +2,7 @@ name: archive-control-cache
|
||||
|
||||
services:
|
||||
archive-client:
|
||||
image: sodium/archive-clients:v0.1.2
|
||||
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
|
||||
|
||||
@@ -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"]
|
||||
|
||||
@@ -2,7 +2,7 @@ name: archive-control-cache
|
||||
|
||||
services:
|
||||
archive-client:
|
||||
image: sodium/archive-clients:v0.1.2
|
||||
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
@@ -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
@@ -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
|
||||
|
||||
@@ -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.2
|
||||
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
|
||||
|
||||
@@ -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.
|
||||
@@ -29,6 +29,10 @@ The adapter accounts for terminology/behavior changes such as paused versus
|
||||
stopped states. It does not mutate qBittorrent preferences, categories, tags,
|
||||
limits, queueing defaults, or global save-path behavior.
|
||||
|
||||
Syncthing route validation normalizes both absolute and `~/` folder-path
|
||||
spellings under the configured sync root before comparing an existing folder;
|
||||
an equivalent pre-existing pair is adopted without rewriting it.
|
||||
|
||||
### Inventory normalization
|
||||
|
||||
For each torrent the client derives canonical v1/v2 identity from reliable API
|
||||
@@ -85,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
@@ -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
@@ -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
@@ -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 |
|
||||
|
||||
|
||||
@@ -102,6 +102,10 @@ and placement policy differ.
|
||||
external state rather than being downloaded or silently deselected.
|
||||
- Confirm free space/reserve, permissions, sparse capability, supported
|
||||
partfile format, route health, and negotiated protocol features.
|
||||
- Charge payload bytes against free space only when the specific source and
|
||||
destination files cannot be hardlinked. Same-filesystem hardlink stages and
|
||||
merges retain the configured reserve plus small metadata artifacts, rather
|
||||
than reserving a duplicate logical payload.
|
||||
- Export the source `.torrent` metainfo.
|
||||
- Persist the exact source, target baseline, requested selection, and delta.
|
||||
|
||||
|
||||
@@ -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
@@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta"
|
||||
|
||||
[project]
|
||||
name = "archive-clients"
|
||||
version = "0.1.2"
|
||||
version = "0.1.19"
|
||||
requires-python = ">=3.11"
|
||||
dependencies = ["protobuf==7.35.1", "websockets==16.0"]
|
||||
|
||||
|
||||
Executable
+327
@@ -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())
|
||||
@@ -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_")
|
||||
|
||||
@@ -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
@@ -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
|
||||
|
||||
+205
-35
@@ -4,12 +4,15 @@ from __future__ import annotations
|
||||
|
||||
import hashlib
|
||||
import json
|
||||
import logging
|
||||
import os
|
||||
import shutil
|
||||
import stat
|
||||
import threading
|
||||
import time
|
||||
import uuid
|
||||
from pathlib import Path, PurePosixPath
|
||||
from typing import Callable
|
||||
from typing import Callable, Iterable
|
||||
|
||||
from archive_clients.protocol import decode_message, encode_message
|
||||
from archive_clients.eviction import (
|
||||
@@ -25,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,
|
||||
@@ -51,6 +55,9 @@ class JobCancelled(JobExecutionError):
|
||||
pass
|
||||
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
class ClientJobExecutor:
|
||||
def __init__(
|
||||
self,
|
||||
@@ -65,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
|
||||
@@ -79,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()
|
||||
@@ -108,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,
|
||||
@@ -125,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,
|
||||
@@ -160,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
|
||||
@@ -541,15 +593,29 @@ class ClientJobExecutor:
|
||||
sha256_hex=hashlib.sha256(resource.metainfo_bytes).hexdigest(),
|
||||
)
|
||||
del artifact
|
||||
route_root = self.route_path(definition.transfer.route_id)
|
||||
# Staging is normally zero-copy when qB's data root and the paired
|
||||
# Syncthing route share a filesystem. Do not reserve the complete
|
||||
# logical payload in that case: FileMaterializer will use link(2),
|
||||
# which consumes only directory/inode metadata. Retain the metainfo
|
||||
# allowance and reserve, and account for any source files that really
|
||||
# must fall back to a data-copy path.
|
||||
self._require_space(
|
||||
self.route_path(definition.transfer.route_id),
|
||||
definition.transfer.transfer_delta_logical_bytes
|
||||
route_root,
|
||||
self._copy_required_bytes(
|
||||
self.qb_root,
|
||||
route_root,
|
||||
(
|
||||
(entry.target_canonical_path, entry.logical_bytes)
|
||||
for entry in manifest.files
|
||||
),
|
||||
)
|
||||
+ len(resource.metainfo_bytes),
|
||||
)
|
||||
stage_transfer(
|
||||
manifest,
|
||||
source_root=self.qb_root,
|
||||
sync_root=self.route_path(definition.transfer.route_id),
|
||||
sync_root=route_root,
|
||||
store=self.store,
|
||||
artifact_sources={"metainfo/source.torrent": metainfo_path},
|
||||
sparse_supported=self.sparse_supported,
|
||||
@@ -563,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,
|
||||
@@ -595,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(
|
||||
@@ -607,9 +692,18 @@ class ClientJobExecutor:
|
||||
"target materialization was sent to the wrong client"
|
||||
)
|
||||
published = load_published_transfer(self._job_directory(definition))
|
||||
# The target can likewise hardlink an arrived Syncthing payload into
|
||||
# qB's content root when those directories share a filesystem.
|
||||
self._require_space(
|
||||
self.qb_root,
|
||||
definition.transfer.transfer_delta_logical_bytes,
|
||||
self._copy_required_bytes(
|
||||
published.job_directory,
|
||||
self.qb_root,
|
||||
(
|
||||
(entry.payload_relative_path, entry.logical_bytes)
|
||||
for entry in published.manifest.files
|
||||
),
|
||||
),
|
||||
)
|
||||
info_hash = _info_hash(definition)
|
||||
resource = self.qbittorrent.get_resource(info_hash)
|
||||
@@ -794,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():
|
||||
@@ -843,6 +942,37 @@ class ClientJobExecutor:
|
||||
f"{required} bytes required including reserve"
|
||||
)
|
||||
|
||||
@staticmethod
|
||||
def _copy_required_bytes(
|
||||
source_root: Path,
|
||||
destination_root: Path,
|
||||
files: Iterable[tuple[str, int]],
|
||||
) -> int:
|
||||
"""Return logical bytes that cannot be materialized by hardlink.
|
||||
|
||||
A hardlink is possible only for regular files on the destination
|
||||
filesystem. Conservatively charge a file when it cannot be inspected;
|
||||
the normal materializer will then provide the precise integrity error.
|
||||
"""
|
||||
|
||||
destination_device = destination_root.stat().st_dev
|
||||
required = 0
|
||||
for relative_path, logical_bytes in files:
|
||||
relative = PurePosixPath(relative_path)
|
||||
source = source_root.joinpath(*relative.parts)
|
||||
try:
|
||||
metadata = source.stat(follow_symlinks=False)
|
||||
except OSError:
|
||||
required += logical_bytes
|
||||
continue
|
||||
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
|
||||
|
||||
def _observer(
|
||||
self, definition: job_pb2.JobDefinition
|
||||
) -> SyncthingTransferObserver:
|
||||
@@ -1089,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):
|
||||
@@ -1096,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", "\\")
|
||||
)
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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:
|
||||
@@ -369,8 +399,8 @@ class SyncthingRouteManager:
|
||||
raise RoutePathConflict("new route path contains unrelated data")
|
||||
path.mkdir(parents=True, exist_ok=True)
|
||||
|
||||
@staticmethod
|
||||
def _validate_folder(
|
||||
self,
|
||||
folder: dict[str, Any],
|
||||
api_path: str,
|
||||
local_device_id: str,
|
||||
@@ -382,7 +412,15 @@ class SyncthingRouteManager:
|
||||
for item in raw_devices
|
||||
if isinstance(item, dict) and isinstance(item.get("deviceID"), str)
|
||||
} if isinstance(raw_devices, list) else set()
|
||||
if folder.get("path") != api_path:
|
||||
try:
|
||||
configured_api_path = self._normalized_folder_api_path(
|
||||
folder.get("path")
|
||||
)
|
||||
except RoutePathConflict as exc:
|
||||
raise RoutePathConflict(
|
||||
"existing route ID uses a different path"
|
||||
) from exc
|
||||
if configured_api_path != api_path:
|
||||
raise RoutePathConflict("existing route ID uses a different path")
|
||||
if folder.get("type") != "sendreceive":
|
||||
raise RouteSetupError("existing route folder is not sendreceive")
|
||||
@@ -391,6 +429,27 @@ class SyncthingRouteManager:
|
||||
if folder.get("paused") is True:
|
||||
raise RouteSetupError("existing route folder is paused")
|
||||
|
||||
def _normalized_folder_api_path(self, value: Any) -> str:
|
||||
"""Normalize Syncthing's absolute and home-relative path spellings."""
|
||||
|
||||
if not isinstance(value, str) or not value:
|
||||
raise RoutePathConflict("existing route folder path is invalid")
|
||||
candidate = PurePosixPath(value)
|
||||
if candidate.parts and candidate.parts[0] == "~":
|
||||
candidate = self.config.api_root.joinpath(*candidate.parts[1:])
|
||||
if not candidate.is_absolute() or any(
|
||||
part in {"", ".", ".."} for part in candidate.parts
|
||||
):
|
||||
raise RoutePathConflict("existing route folder path is unsafe")
|
||||
normalized = candidate.as_posix()
|
||||
try:
|
||||
self.config.roots.api_to_local(normalized)
|
||||
except ConfigError as exc:
|
||||
raise RoutePathConflict(
|
||||
"existing route folder path is outside the sync root"
|
||||
) from exc
|
||||
return normalized
|
||||
|
||||
def _wait(self, deadline: float) -> None:
|
||||
if time.monotonic() >= deadline:
|
||||
raise RouteSetupTimeout("route setup timed out")
|
||||
@@ -441,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"):
|
||||
|
||||
@@ -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
@@ -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)
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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()
|
||||
+286
-5
@@ -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)
|
||||
@@ -69,14 +190,166 @@ class ClientJobHappyPathTests(unittest.TestCase):
|
||||
with self.subTest(operation=operation), tempfile.TemporaryDirectory() as directory:
|
||||
self._run_transfer(Path(directory), operation)
|
||||
|
||||
def _run_transfer(self, root: Path, operation: int):
|
||||
def test_same_filesystem_source_stage_does_not_require_payload_space(self):
|
||||
with tempfile.TemporaryDirectory() as directory:
|
||||
self._run_transfer(
|
||||
Path(directory),
|
||||
job_pb2.JOB_OPERATION_ARCHIVE,
|
||||
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,
|
||||
operation: int,
|
||||
*,
|
||||
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"
|
||||
route_root = root / "route"
|
||||
source_root.mkdir()
|
||||
target_root.mkdir()
|
||||
route_root.mkdir()
|
||||
content = b"archive-control-happy-path"
|
||||
(source_root / "fixture.bin").write_bytes(content)
|
||||
info = {
|
||||
b"length": len(content),
|
||||
@@ -137,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,
|
||||
@@ -185,13 +458,21 @@ class ClientJobHappyPathTests(unittest.TestCase):
|
||||
)
|
||||
final = None
|
||||
for executor, step in pipeline:
|
||||
events = executor.execute(control_pb2.ExecuteStepCommand(
|
||||
command = control_pb2.ExecuteStepCommand(
|
||||
job_id=definition.job_id,
|
||||
expected_job_revision=cursor_revision,
|
||||
expected_last_event_sequence=cursor_sequence,
|
||||
step=step,
|
||||
attempt=1,
|
||||
))
|
||||
)
|
||||
if executor is source and step == job_pb2.JOB_STEP_KIND_SOURCE_STAGE and source_stage_free_bytes is not None:
|
||||
with patch(
|
||||
"archive_clients.jobs.shutil.disk_usage",
|
||||
return_value=Mock(free=source_stage_free_bytes),
|
||||
):
|
||||
events = executor.execute(command)
|
||||
else:
|
||||
events = executor.execute(command)
|
||||
self.assertEqual(
|
||||
[event.sequence for event in events],
|
||||
list(range(
|
||||
|
||||
@@ -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 = [
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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()
|
||||
|
||||
@@ -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"
|
||||
|
||||
@@ -125,6 +125,24 @@ class SyncthingRouteManagerTests(unittest.TestCase):
|
||||
self.manager.configure(self.spec, time.monotonic() + 1)
|
||||
self.assertEqual(self.transport.puts, [])
|
||||
|
||||
def test_existing_home_relative_folder_path_is_accepted(self):
|
||||
self.transport.config["devices"].append({"deviceID": "PEER"})
|
||||
self.transport.config["folders"].append(
|
||||
{
|
||||
"id": "route-1",
|
||||
"path": "~/routes/route-1",
|
||||
"type": "sendreceive",
|
||||
"devices": [{"deviceID": "LOCAL"}, {"deviceID": "PEER"}],
|
||||
}
|
||||
)
|
||||
|
||||
configured = self.manager.configure(
|
||||
self.spec, time.monotonic() + 1
|
||||
)
|
||||
|
||||
self.assertEqual(self.transport.puts, [])
|
||||
self.assertFalse(configured.local_route.archive_control_created)
|
||||
|
||||
def test_bidirectional_nonce_and_ack_are_required(self):
|
||||
configured = self.manager.configure(self.spec, time.monotonic() + 1)
|
||||
peer_nonce = configured.local_path / (
|
||||
|
||||
@@ -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)
|
||||
|
||||
Reference in New Issue
Block a user