Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
b67ca18403 | ||
|
|
6446846ee8 | ||
|
|
c717aea394 | ||
|
|
ea35ba2758 | ||
|
|
5e5e2f9095 | ||
|
|
0b11e2a3a2 | ||
|
|
1009defc46 | ||
|
|
8ed8e75447 | ||
|
|
1a06da3984 | ||
|
|
4ba5e92248 | ||
|
|
a6da224f28 | ||
|
|
90925fd321 | ||
|
|
9df2e6991e | ||
|
|
d2fa69a1d6 | ||
|
|
0f94388f49 | ||
|
|
94a3233ce2 | ||
|
|
053b0135b3 |
+6
-2
@@ -11,8 +11,12 @@ CMD ["python", "-m", "unittest", "discover", "-s", "tests", "-v"]
|
|||||||
|
|
||||||
FROM python:3.11-slim
|
FROM python:3.11-slim
|
||||||
ENV PYTHONDONTWRITEBYTECODE=1 PYTHONUNBUFFERED=1
|
ENV PYTHONDONTWRITEBYTECODE=1 PYTHONUNBUFFERED=1
|
||||||
RUN groupadd --gid 1001 archive-control \
|
RUN if ! getent group 1001 >/dev/null; then \
|
||||||
&& useradd --uid 1001 --gid 1001 --no-create-home archive-control
|
groupadd --gid 1001 archive-control; \
|
||||||
|
fi \
|
||||||
|
&& if ! getent passwd 1001 >/dev/null; then \
|
||||||
|
useradd --uid 1001 --gid 1001 --no-create-home archive-control; \
|
||||||
|
fi
|
||||||
COPY --from=builder /wheels /wheels
|
COPY --from=builder /wheels /wheels
|
||||||
RUN pip install --no-cache-dir /wheels/*.whl && rm -rf /wheels
|
RUN pip install --no-cache-dir /wheels/*.whl && rm -rf /wheels
|
||||||
USER 1001:1001
|
USER 1001:1001
|
||||||
|
|||||||
@@ -13,10 +13,37 @@ initial x1/x2/lithium topology.
|
|||||||
`archive_control_token`, `qb_password`, and `syncthing_api_key`, each a
|
`archive_control_token`, `qb_password`, and `syncthing_api_key`, each a
|
||||||
regular non-empty file with mode `0600`.
|
regular non-empty file with mode `0600`.
|
||||||
|
|
||||||
The Syncthing mounts intentionally reproduce each instance's `/var/syncthing`
|
## Hardlink-safe bind-mount topology
|
||||||
layout, including nested data binds. This lets route discovery and route
|
|
||||||
provisioning use one safe API-to-local path mapping without altering an
|
For an archive source, the qB content path and every Syncthing route used for
|
||||||
existing Syncthing configuration.
|
staging must resolve through the **same container mount**. Matching host
|
||||||
|
filesystem device IDs alone is insufficient: two separate Docker bind mounts
|
||||||
|
have different mount IDs and `link(2)` may return `EXDEV` across them. The
|
||||||
|
client deliberately treats that case as copy-only and performs a full payload
|
||||||
|
free-space check.
|
||||||
|
|
||||||
|
When a Syncthing route is physically nested below the qB root, mount the qB
|
||||||
|
root once and map the exact Syncthing API folder through it:
|
||||||
|
|
||||||
|
```yaml
|
||||||
|
volumes:
|
||||||
|
- /srv/downloads:/data/qb
|
||||||
|
- /srv/syncthing-config:/data/sync
|
||||||
|
```
|
||||||
|
|
||||||
|
```toml
|
||||||
|
[syncthing]
|
||||||
|
api_root = "/var/syncthing"
|
||||||
|
local_root = "/data/sync"
|
||||||
|
local_path_overrides = { "/var/syncthing/Downloads/Sync" = "/data/qb/Sync" }
|
||||||
|
```
|
||||||
|
|
||||||
|
Do **not** additionally mount `/srv/downloads/Sync` at a path beneath
|
||||||
|
`/data/sync`. The override is the authoritative mapping for that folder and
|
||||||
|
keeps qB source files and staging destinations in one mount namespace. Use
|
||||||
|
the folder's normalized API-visible **path** as the override key (for example,
|
||||||
|
Syncthing `~/Downloads/Sync` becomes `/var/syncthing/Downloads/Sync`); do not
|
||||||
|
use the folder ID.
|
||||||
|
|
||||||
Before starting a stack, validate it with:
|
Before starting a stack, validate it with:
|
||||||
|
|
||||||
|
|||||||
@@ -2,7 +2,7 @@ name: archive-control-archive
|
|||||||
|
|
||||||
services:
|
services:
|
||||||
archive-client:
|
archive-client:
|
||||||
image: sodium/archive-clients:v0.1.6
|
image: sodium/archive-clients:v0.1.14
|
||||||
user: "1000:1000"
|
user: "1000:1000"
|
||||||
restart: unless-stopped
|
restart: unless-stopped
|
||||||
command: ["--config", "/etc/archive-control/client.toml"]
|
command: ["--config", "/etc/archive-control/client.toml"]
|
||||||
|
|||||||
@@ -2,7 +2,7 @@ name: archive-control-cache
|
|||||||
|
|
||||||
services:
|
services:
|
||||||
archive-client:
|
archive-client:
|
||||||
image: sodium/archive-clients:v0.1.6
|
image: sodium/archive-clients:v0.1.14
|
||||||
user: "1001:1001"
|
user: "1001:1001"
|
||||||
restart: unless-stopped
|
restart: unless-stopped
|
||||||
network_mode: host
|
network_mode: host
|
||||||
|
|||||||
@@ -39,4 +39,8 @@ endpoint = "http://127.0.0.1:8384"
|
|||||||
api_key_file = "/run/secrets/syncthing_api_key"
|
api_key_file = "/run/secrets/syncthing_api_key"
|
||||||
api_root = "/var/syncthing"
|
api_root = "/var/syncthing"
|
||||||
local_root = "/data/sync"
|
local_root = "/data/sync"
|
||||||
|
# `DownloadsSync-X2` is ~/Downloads/Sync on the host, nested below the qB
|
||||||
|
# root. Resolve it through /data/qb rather than a second nested bind mount so
|
||||||
|
# source staging can hardlink it.
|
||||||
|
local_path_overrides = { "/var/syncthing/Downloads/Sync" = "/data/qb/Sync" }
|
||||||
advertised_addresses = ["dynamic"]
|
advertised_addresses = ["dynamic"]
|
||||||
|
|||||||
@@ -2,7 +2,7 @@ name: archive-control-cache
|
|||||||
|
|
||||||
services:
|
services:
|
||||||
archive-client:
|
archive-client:
|
||||||
image: sodium/archive-clients:v0.1.6
|
image: sodium/archive-clients:v0.1.14
|
||||||
user: "1001:1001"
|
user: "1001:1001"
|
||||||
restart: unless-stopped
|
restart: unless-stopped
|
||||||
network_mode: host
|
network_mode: host
|
||||||
@@ -16,4 +16,3 @@ services:
|
|||||||
- ./backups:/var/backups/archive-control
|
- ./backups:/var/backups/archive-control
|
||||||
- /home/ubuntu/Downloads:/data/qb
|
- /home/ubuntu/Downloads:/data/qb
|
||||||
- /home/ubuntu/compose/syncthing/st_home:/data/sync
|
- /home/ubuntu/compose/syncthing/st_home:/data/sync
|
||||||
- /home/ubuntu/Downloads/Sync:/data/sync/Downloads/Sync
|
|
||||||
|
|||||||
+5
-2
@@ -35,8 +35,11 @@ expand these rules but must not contradict them.
|
|||||||
manifests, and observed external state.
|
manifests, and observed external state.
|
||||||
- SQLite backups use the online backup API. Defaults are every six hours, 12
|
- SQLite backups use the online backup API. Defaults are every six hours, 12
|
||||||
recent, 14 daily, and eight weekly copies, plus pre/post-migration backups.
|
recent, 14 daily, and eight weekly copies, plus pre/post-migration backups.
|
||||||
- Default concurrency is one active data-moving job per client and per route.
|
- Default concurrency is unrestricted: every eligible queued job is admitted
|
||||||
Disjoint node pairs may run concurrently. Queueing is durable FIFO.
|
in one scheduler pass and independent jobs execute concurrently on a client.
|
||||||
|
Physical qBittorrent, Syncthing, disk, and network capacity are therefore
|
||||||
|
the natural limit. `enforce_concurrency_limits=true` restores the optional
|
||||||
|
legacy per-client/per-route gates. Commands for one job remain serialized.
|
||||||
- A queued job owns a per-resource reservation. Cancelling it removes only the
|
- A queued job owns a per-resource reservation. Cancelling it removes only the
|
||||||
queued record/reservation and never sends cleanup commands.
|
queued record/reservation and never sends cleanup commands.
|
||||||
- Connectivity or transfer stalls wait indefinitely. A configurable 30-minute
|
- Connectivity or transfer stalls wait indefinitely. A configurable 30-minute
|
||||||
|
|||||||
@@ -44,6 +44,9 @@ command_max_attempts = 3 # Total sends, including the initial attempt.
|
|||||||
stall_after = "30m" # Warning state only; jobs continue waiting.
|
stall_after = "30m" # Warning state only; jobs continue waiting.
|
||||||
route_policy = "on_demand" # Alternative: eager_mesh.
|
route_policy = "on_demand" # Alternative: eager_mesh.
|
||||||
route_setup_timeout = "30m"
|
route_setup_timeout = "30m"
|
||||||
|
# Default: allow all eligible jobs; disk/network/service capacity is the limit.
|
||||||
|
enforce_concurrency_limits = false
|
||||||
|
# Used only when enforce_concurrency_limits is true.
|
||||||
max_active_per_client = 1
|
max_active_per_client = 1
|
||||||
max_active_per_route = 1
|
max_active_per_route = 1
|
||||||
max_envelope_bytes = 1048576
|
max_envelope_bytes = 1048576
|
||||||
@@ -128,7 +131,12 @@ advertised_addresses = ["dynamic"]
|
|||||||
When an existing Syncthing folder is physically nested in the qB data root,
|
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
|
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
|
the same client bind mount. This enables hardlinks without creating two Docker
|
||||||
mount boundaries for the same host files.
|
mount boundaries for the same host files. Do not add a second bind mount for
|
||||||
|
the nested folder: Linux treats it as a distinct mount even when it has the
|
||||||
|
same `st_dev`, and the client correctly falls back to copy-only capacity
|
||||||
|
accounting. The override key is the normalized Syncthing folder path beneath
|
||||||
|
`api_root`, not its folder ID. See the production deployment README for the
|
||||||
|
required compose and override pattern.
|
||||||
|
|
||||||
The remaining node examples omit optional `[connection]`, `[jobs]`, and
|
The remaining node examples omit optional `[connection]`, `[jobs]`, and
|
||||||
`[backup]` tables and therefore use these same defaults; deployments may
|
`[backup]` tables and therefore use these same defaults; deployments may
|
||||||
@@ -218,7 +226,7 @@ cache/archive routes according to policy.
|
|||||||
```yaml
|
```yaml
|
||||||
services:
|
services:
|
||||||
archive-client:
|
archive-client:
|
||||||
image: sodium/archive-clients:v0.1.6
|
image: sodium/archive-clients:v0.1.14
|
||||||
user: "1001:1001"
|
user: "1001:1001"
|
||||||
restart: unless-stopped
|
restart: unless-stopped
|
||||||
command: ["archive-client", "--config", "/etc/archive-control/client.toml"]
|
command: ["archive-client", "--config", "/etc/archive-control/client.toml"]
|
||||||
@@ -251,6 +259,28 @@ Generated protobuf Python bindings are committed in each consumer with the
|
|||||||
exact `archive-control-proto` tag/commit recorded. Release order is proto,
|
exact `archive-control-proto` tag/commit recorded. Release order is proto,
|
||||||
control consumer, client consumer, E2E, then image publication.
|
control consumer, client consumer, E2E, then image publication.
|
||||||
|
|
||||||
|
### Reproducible multi-platform Buildx lifecycle
|
||||||
|
|
||||||
|
The named Buildx builders are local acceleration/cache only; they are not a
|
||||||
|
deployment dependency and can be removed after publication. To create a fresh
|
||||||
|
builder, verify its platforms, publish a release, and remove it afterwards:
|
||||||
|
|
||||||
|
```bash
|
||||||
|
docker buildx create --name archive-control-release --driver docker-container --use
|
||||||
|
docker buildx inspect --bootstrap
|
||||||
|
docker buildx build --platform linux/amd64,linux/arm64 \
|
||||||
|
--tag sodium/archive-clients:vX.Y.Z --push .
|
||||||
|
docker buildx rm archive-control-release
|
||||||
|
```
|
||||||
|
|
||||||
|
`docker buildx inspect` must list both `linux/amd64` and `linux/arm64` before
|
||||||
|
publishing. If the host has no arm64 emulation, install/configure it according
|
||||||
|
to the host Docker distribution before the build; do not publish a partial
|
||||||
|
single-platform tag. Retain the pushed manifest digest in the release notes
|
||||||
|
and deploy the immutable tag or digest. The optional `archive-control-qemu`
|
||||||
|
builder follows the same lifecycle when it is used for an emulation smoke
|
||||||
|
build.
|
||||||
|
|
||||||
## Operator usage
|
## Operator usage
|
||||||
|
|
||||||
1. Prepare local qB, sync, state, backup, config, and secret mounts with the
|
1. Prepare local qB, sync, state, backup, config, and secret mounts with the
|
||||||
|
|||||||
@@ -89,9 +89,11 @@ left unchanged.
|
|||||||
|
|
||||||
Syncthing may serialize a folder path relative to its home as `~/...`. The
|
Syncthing may serialize a folder path relative to its home as `~/...`. The
|
||||||
client normalizes that notation beneath the configured API-visible sync root
|
client normalizes that notation beneath the configured API-visible sync root
|
||||||
before applying the API-to-local root mapping. Deployments must mirror
|
before applying the API-to-local root mapping. When that folder is nested
|
||||||
Syncthing's nested bind mounts into the client so the normalized API path and
|
below qB's content root, an exact `local_path_overrides` entry must map it
|
||||||
the client filesystem path refer to the same bytes.
|
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
|
Provisioning uses idempotent device and folder configuration updates. Each
|
||||||
client receives its peer device ID and optional advertised addresses
|
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
|
baselines even when jobs would use different nodes or observe different sides
|
||||||
of a hybrid identity.
|
of a hybrid identity.
|
||||||
|
|
||||||
Defaults allow one active data-moving job per client and one per route. A job
|
By default every eligible queued job is claimed in the same scheduler pass and
|
||||||
must acquire its source client, target client, route, and resource reservation
|
independent jobs run concurrently on each client; physical service and storage
|
||||||
atomically. Disjoint node pairs may run concurrently. Route setup is a
|
capacity are the limit. Set `enforce_concurrency_limits=true` on control to
|
||||||
preflight activity and does not permit a data step to bypass these leases.
|
enable the optional one-per-client/one-per-route gates. A job always retains
|
||||||
|
its resource reservation, and commands for the same job remain serialized.
|
||||||
|
Route setup is a preflight activity and does not permit a data step to bypass
|
||||||
|
the resource reservation.
|
||||||
|
|
||||||
Offline nodes do not prevent unrelated jobs from being listed or run. A job
|
Offline nodes do not prevent unrelated jobs from being listed or run. A job
|
||||||
requiring an offline node remains waiting indefinitely; it does not consume an
|
requiring an offline node remains waiting indefinitely; it does not consume an
|
||||||
@@ -154,6 +157,12 @@ advances to the next step. On reconnect:
|
|||||||
4. Control observes relevant qBittorrent, Syncthing, staging, and manifest
|
4. Control observes relevant qBittorrent, Syncthing, staging, and manifest
|
||||||
state before selecting retry, resume, compensation, cleanup, or manual
|
state before selecting retry, resume, compensation, cleanup, or manual
|
||||||
intervention.
|
intervention.
|
||||||
|
|
||||||
|
If a connection disappears while a synchronous job operation is in progress,
|
||||||
|
the client journals events locally, reconnects indefinitely, and serializes a
|
||||||
|
replayed command behind the in-flight operation for that job. The replay reads
|
||||||
|
the journal rather than repeating the data operation. This applies to every
|
||||||
|
transfer step, including post-commit staging cleanup.
|
||||||
5. A command is reissued with its original ID when the acceptance result is
|
5. A command is reissued with its original ID when the acceptance result is
|
||||||
uncertain.
|
uncertain.
|
||||||
|
|
||||||
|
|||||||
+4
-2
@@ -25,8 +25,10 @@ What do you want to do?
|
|||||||
[ Evict Cache ] [ Job Status ]
|
[ Evict Cache ] [ Job Status ]
|
||||||
```
|
```
|
||||||
|
|
||||||
All subsequent pages edit this message. `Cancel` closes the active selection
|
All subsequent pages edit this message. Leaving a resource/tree selection
|
||||||
flow. `Back` returns one level while retaining validated filters/selections.
|
returns to the Archive Control operation menu without closing the shared
|
||||||
|
conversation. Confirmation `Cancel` returns to its immediate prior selection;
|
||||||
|
the main menu alone can close Archive Control.
|
||||||
Callback payloads contain opaque session/action IDs, not resource names or
|
Callback payloads contain opaque session/action IDs, not resource names or
|
||||||
paths, and are validated against persisted session revision and expiry.
|
paths, and are validated against persisted session revision and expiry.
|
||||||
|
|
||||||
|
|||||||
+1
-1
@@ -225,7 +225,7 @@ Fixtures include small deterministic v1, v2, and hybrid torrents with:
|
|||||||
| Recovery | restart every step on source/target/control, lost DB with proof, ambiguous loss fails closed |
|
| Recovery | restart every step on source/target/control, lost DB with proof, ambiguous loss fails closed |
|
||||||
| Messaging | duplicate command/event, lost ack, sequence gap, stale revision, duplicate client ID |
|
| Messaging | duplicate command/event, lost ack, sequence gap, stale revision, duplicate client ID |
|
||||||
| Safety | attempted download, corrupt same-size file, symlink/path escape, special file, partfile mismatch |
|
| Safety | attempted download, corrupt same-size file, symlink/path escape, special file, partfile mismatch |
|
||||||
| Operations | stall retains percent, offline wait, cancellation each phase, postcommit cleanup retry |
|
| Operations | stall retains percent, offline wait, cancellation each phase, postcommit cleanup retry, disconnect/reconnect during every transfer step with exactly-once replay |
|
||||||
| Database | online backup under load, pre/post migration, retention, corrupt backup rejection, offline restore |
|
| Database | online backup under load, pre/post migration, retention, corrupt backup rejection, offline restore |
|
||||||
| UI | pagination/filter chain, stale/double callback, persisted session, clear vs evict separation |
|
| UI | pagination/filter chain, stale/double callback, persisted session, clear vs evict separation |
|
||||||
|
|
||||||
|
|||||||
+1
-1
@@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta"
|
|||||||
|
|
||||||
[project]
|
[project]
|
||||||
name = "archive-clients"
|
name = "archive-clients"
|
||||||
version = "0.1.6"
|
version = "0.1.18"
|
||||||
requires-python = ">=3.11"
|
requires-python = ">=3.11"
|
||||||
dependencies = ["protobuf==7.35.1", "websockets==16.0"]
|
dependencies = ["protobuf==7.35.1", "websockets==16.0"]
|
||||||
|
|
||||||
|
|||||||
@@ -11,6 +11,7 @@ from pathlib import Path, PurePosixPath
|
|||||||
from typing import Any
|
from typing import Any
|
||||||
|
|
||||||
from websockets.asyncio.client import connect
|
from websockets.asyncio.client import connect
|
||||||
|
from websockets.exceptions import ConnectionClosedOK
|
||||||
|
|
||||||
from archive_clients.backup import SQLiteBackupManager
|
from archive_clients.backup import SQLiteBackupManager
|
||||||
from archive_clients.config import ClientConfig
|
from archive_clients.config import ClientConfig
|
||||||
@@ -42,6 +43,19 @@ from archive_control.v1 import (
|
|||||||
logger = logging.getLogger(__name__)
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
|
||||||
|
def _job_id_for_command(command: control_pb2.Command) -> str:
|
||||||
|
"""Return the durable job key for commands whose execution is serialized."""
|
||||||
|
|
||||||
|
payload = command.WhichOneof("payload")
|
||||||
|
if payload == "assign_job":
|
||||||
|
return command.assign_job.job.job_id
|
||||||
|
if payload == "execute_step":
|
||||||
|
return command.execute_step.job_id
|
||||||
|
if payload == "cancel_job":
|
||||||
|
return command.cancel_job.job_id
|
||||||
|
raise ValueError(f"command {command.command_id} does not execute a job")
|
||||||
|
|
||||||
|
|
||||||
class ArchiveClientDaemon:
|
class ArchiveClientDaemon:
|
||||||
def __init__(
|
def __init__(
|
||||||
self,
|
self,
|
||||||
@@ -94,7 +108,10 @@ class ArchiveClientDaemon:
|
|||||||
self._lease = DatabaseLease(config.state_db)
|
self._lease = DatabaseLease(config.state_db)
|
||||||
self._active_route_commands: set[str] = set()
|
self._active_route_commands: set[str] = set()
|
||||||
self._active_job_commands: set[str] = set()
|
self._active_job_commands: set[str] = set()
|
||||||
self._job_execution_lock = asyncio.Lock()
|
# Commands for one job remain ordered locally, while unrelated jobs
|
||||||
|
# may use the node's available qB/Syncthing/filesystem capacity in
|
||||||
|
# parallel. The control daemon owns admission policy.
|
||||||
|
self._job_execution_locks: dict[str, asyncio.Lock] = {}
|
||||||
self.jobs = (
|
self.jobs = (
|
||||||
ClientJobExecutor(
|
ClientJobExecutor(
|
||||||
client_id=config.client_id,
|
client_id=config.client_id,
|
||||||
@@ -205,7 +222,23 @@ class ArchiveClientDaemon:
|
|||||||
command_tasks: set[asyncio.Task[None]] = set()
|
command_tasks: set[asyncio.Task[None]] = set()
|
||||||
try:
|
try:
|
||||||
await self._resume_commands(outbound, command_tasks)
|
await self._resume_commands(outbound, command_tasks)
|
||||||
async for frame in websocket:
|
# A proxy can leave the TCP/WebSocket socket apparently open
|
||||||
|
# after the control server has discarded its session. The
|
||||||
|
# server then cannot deliver durable commands and its pending
|
||||||
|
# outbox remains stranded unless the client independently
|
||||||
|
# detects the missing application heartbeats and reconnects.
|
||||||
|
while True:
|
||||||
|
try:
|
||||||
|
frame = await asyncio.wait_for(
|
||||||
|
websocket.recv(),
|
||||||
|
self.config.connection.offline_timeout,
|
||||||
|
)
|
||||||
|
except ConnectionClosedOK:
|
||||||
|
return
|
||||||
|
except asyncio.TimeoutError as exc:
|
||||||
|
raise RuntimeError(
|
||||||
|
"control heartbeat timed out"
|
||||||
|
) from exc
|
||||||
await self._handle(decode(frame), outbound, command_tasks)
|
await self._handle(decode(frame), outbound, command_tasks)
|
||||||
finally:
|
finally:
|
||||||
writer.cancel()
|
writer.cancel()
|
||||||
@@ -543,7 +576,9 @@ class ArchiveClientDaemon:
|
|||||||
correlation_id: str,
|
correlation_id: str,
|
||||||
outbound: asyncio.Queue[str],
|
outbound: asyncio.Queue[str],
|
||||||
) -> None:
|
) -> None:
|
||||||
async with self._job_execution_lock:
|
job_id = _job_id_for_command(command)
|
||||||
|
lock = self._job_execution_locks.setdefault(job_id, asyncio.Lock())
|
||||||
|
async with lock:
|
||||||
await self._execute_job_command_locked(
|
await self._execute_job_command_locked(
|
||||||
command, correlation_id, outbound
|
command, correlation_id, outbound
|
||||||
)
|
)
|
||||||
|
|||||||
+105
-29
@@ -4,6 +4,7 @@ from __future__ import annotations
|
|||||||
|
|
||||||
import hashlib
|
import hashlib
|
||||||
import json
|
import json
|
||||||
|
import logging
|
||||||
import os
|
import os
|
||||||
import shutil
|
import shutil
|
||||||
import stat
|
import stat
|
||||||
@@ -27,10 +28,11 @@ from archive_clients.qbittorrent import (
|
|||||||
)
|
)
|
||||||
from archive_clients.resources import NormalizedResource
|
from archive_clients.resources import NormalizedResource
|
||||||
from archive_clients.state import ClientStore
|
from archive_clients.state import ClientStore
|
||||||
from archive_clients.syncthing import SyncthingTransferObserver
|
from archive_clients.syncthing import RouteSetupError, SyncthingTransferObserver
|
||||||
from archive_clients.transfer import (
|
from archive_clients.transfer import (
|
||||||
TransferError,
|
TransferError,
|
||||||
TransferIntegrityError,
|
TransferIntegrityError,
|
||||||
|
cleanup_partial_transfer,
|
||||||
cleanup_transfer,
|
cleanup_transfer,
|
||||||
load_published_transfer,
|
load_published_transfer,
|
||||||
materialize_transfer,
|
materialize_transfer,
|
||||||
@@ -53,6 +55,9 @@ class JobCancelled(JobExecutionError):
|
|||||||
pass
|
pass
|
||||||
|
|
||||||
|
|
||||||
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
|
||||||
class ClientJobExecutor:
|
class ClientJobExecutor:
|
||||||
def __init__(
|
def __init__(
|
||||||
self,
|
self,
|
||||||
@@ -81,6 +86,8 @@ class ClientJobExecutor:
|
|||||||
self.verification_timeout = verification_timeout
|
self.verification_timeout = verification_timeout
|
||||||
self.free_space_reserve_bytes = free_space_reserve_bytes
|
self.free_space_reserve_bytes = free_space_reserve_bytes
|
||||||
self._cancel_events: dict[str, threading.Event] = {}
|
self._cancel_events: dict[str, threading.Event] = {}
|
||||||
|
self._execution_locks: dict[str, threading.Lock] = {}
|
||||||
|
self._execution_locks_guard = threading.Lock()
|
||||||
|
|
||||||
def request_cancel(self, job_id: str) -> None:
|
def request_cancel(self, job_id: str) -> None:
|
||||||
self._cancel_events.setdefault(job_id, threading.Event()).set()
|
self._cancel_events.setdefault(job_id, threading.Event()).set()
|
||||||
@@ -110,12 +117,32 @@ class ClientJobExecutor:
|
|||||||
self,
|
self,
|
||||||
command: control_pb2.ExecuteStepCommand,
|
command: control_pb2.ExecuteStepCommand,
|
||||||
event_callback: Callable[[control_pb2.JobEvent], None] | None = None,
|
event_callback: Callable[[control_pb2.JobEvent], None] | None = None,
|
||||||
|
) -> list[control_pb2.JobEvent]:
|
||||||
|
# asyncio cancellation of a connection-bound task cannot stop the
|
||||||
|
# synchronous filesystem/qB operation already running in its worker
|
||||||
|
# thread. A replay after reconnect therefore waits for that operation
|
||||||
|
# and then reads its durable event journal instead of executing twice.
|
||||||
|
with self._execution_lock(command.job_id):
|
||||||
|
return self._execute_locked(command, event_callback)
|
||||||
|
|
||||||
|
def _execution_lock(self, job_id: str) -> threading.Lock:
|
||||||
|
with self._execution_locks_guard:
|
||||||
|
return self._execution_locks.setdefault(job_id, threading.Lock())
|
||||||
|
|
||||||
|
def _execute_locked(
|
||||||
|
self,
|
||||||
|
command: control_pb2.ExecuteStepCommand,
|
||||||
|
event_callback: Callable[[control_pb2.JobEvent], None] | None = None,
|
||||||
) -> list[control_pb2.JobEvent]:
|
) -> list[control_pb2.JobEvent]:
|
||||||
definition = self._definition(command.job_id)
|
definition = self._definition(command.job_id)
|
||||||
replay = self._replay(
|
replay = self._replay(
|
||||||
command.job_id, command.expected_last_event_sequence
|
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_STEP_SUCCEEDED,
|
||||||
control_pb2.JOB_EVENT_TYPE_COMMITTED,
|
control_pb2.JOB_EVENT_TYPE_COMMITTED,
|
||||||
control_pb2.JOB_EVENT_TYPE_SUCCEEDED,
|
control_pb2.JOB_EVENT_TYPE_SUCCEEDED,
|
||||||
@@ -127,32 +154,51 @@ class ClientJobExecutor:
|
|||||||
for event in replay:
|
for event in replay:
|
||||||
event_callback(event)
|
event_callback(event)
|
||||||
return replay
|
return replay
|
||||||
started = replay[0] if replay else self._event(
|
if step_replay:
|
||||||
definition,
|
started = step_replay[0]
|
||||||
sequence=command.expected_last_event_sequence + 1,
|
cursor = step_replay[-1]
|
||||||
revision=command.expected_job_revision + 1,
|
emitted = list(step_replay)
|
||||||
event_type=control_pb2.JOB_EVENT_TYPE_STEP_STARTED,
|
else:
|
||||||
state=job_pb2.JOB_STATE_RUNNING,
|
previous = replay[-1] if replay else None
|
||||||
committed=(
|
started = self._event(
|
||||||
self._committed(command.job_id)
|
definition,
|
||||||
or command.step == job_pb2.JOB_STEP_KIND_STAGING_CLEANUP
|
sequence=(
|
||||||
or (
|
previous.sequence + 1
|
||||||
command.step == job_pb2.JOB_STEP_KIND_SAFE_FILE_UNLINK
|
if previous is not None
|
||||||
and self.store.get_job_artifact(
|
else command.expected_last_event_sequence + 1
|
||||||
command.job_id, "qb-entry-removed"
|
),
|
||||||
) is not None
|
revision=(
|
||||||
)
|
previous.job_revision + 1
|
||||||
),
|
if previous is not None
|
||||||
step=command.step,
|
else command.expected_job_revision + 1
|
||||||
step_state=job_pb2.STEP_STATE_RUNNING,
|
),
|
||||||
)
|
event_type=control_pb2.JOB_EVENT_TYPE_STEP_STARTED,
|
||||||
if not replay:
|
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)
|
self._record(definition, started)
|
||||||
|
cursor = started
|
||||||
|
emitted = [started]
|
||||||
if event_callback is not None:
|
if event_callback is not None:
|
||||||
event_callback(started)
|
for event in emitted:
|
||||||
emitted = [started]
|
event_callback(event)
|
||||||
cursor = started
|
# A newly-started step must publish its first observation. In
|
||||||
speed_sample = [time.monotonic(), 0]
|
# 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(
|
def progress(
|
||||||
fraction: float,
|
fraction: float,
|
||||||
@@ -162,8 +208,12 @@ class ClientJobExecutor:
|
|||||||
) -> None:
|
) -> None:
|
||||||
nonlocal cursor
|
nonlocal cursor
|
||||||
now = time.monotonic()
|
now = time.monotonic()
|
||||||
elapsed = now - float(speed_sample[0])
|
previous_time = speed_sample[0]
|
||||||
if fraction < 1 and elapsed < 1:
|
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
|
return
|
||||||
speed = (
|
speed = (
|
||||||
max(0, bytes_complete - int(speed_sample[1])) / elapsed
|
max(0, bytes_complete - int(speed_sample[1])) / elapsed
|
||||||
@@ -579,7 +629,17 @@ class ClientJobExecutor:
|
|||||||
f"staged {completed} of {total} bytes",
|
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(
|
def _wait_for_syncthing(
|
||||||
self,
|
self,
|
||||||
@@ -611,6 +671,15 @@ class ClientJobExecutor:
|
|||||||
# Syncthing can expose the per-job directory before the
|
# Syncthing can expose the per-job directory before the
|
||||||
# ready marker and manifest have arrived atomically as a set.
|
# ready marker and manifest have arrived atomically as a set.
|
||||||
pass
|
pass
|
||||||
|
except RouteSetupError as exc:
|
||||||
|
# A network interruption can temporarily make the local
|
||||||
|
# Syncthing REST API unavailable. The staged payload is
|
||||||
|
# durable and the job must remain recoverable; keep polling
|
||||||
|
# instead of turning a transient outage into a failed job.
|
||||||
|
logger.warning("syncthing_status_deferred", extra={
|
||||||
|
"job_id": definition.job_id,
|
||||||
|
"error_type": type(exc).__name__,
|
||||||
|
})
|
||||||
time.sleep(self.poll_interval)
|
time.sleep(self.poll_interval)
|
||||||
|
|
||||||
def _target_materialize(
|
def _target_materialize(
|
||||||
@@ -819,6 +888,11 @@ class ClientJobExecutor:
|
|||||||
qb_root=job_directory,
|
qb_root=job_directory,
|
||||||
store=self.store,
|
store=self.store,
|
||||||
)
|
)
|
||||||
|
cleanup_partial_transfer(
|
||||||
|
job_directory,
|
||||||
|
job_id=definition.job_id,
|
||||||
|
store=self.store,
|
||||||
|
)
|
||||||
for name in ("ready.json", "manifest.json"):
|
for name in ("ready.json", "manifest.json"):
|
||||||
candidate = job_directory / name
|
candidate = job_directory / name
|
||||||
if candidate.is_file() and not candidate.is_symlink():
|
if candidate.is_file() and not candidate.is_symlink():
|
||||||
@@ -1145,6 +1219,8 @@ def _job_error_code(error: Exception) -> int:
|
|||||||
return common_pb2.ERROR_CODE_INTEGRITY_CHECK_FAILED
|
return common_pb2.ERROR_CODE_INTEGRITY_CHECK_FAILED
|
||||||
if isinstance(error, PermissionError):
|
if isinstance(error, PermissionError):
|
||||||
return common_pb2.ERROR_CODE_PERMISSION_DENIED
|
return common_pb2.ERROR_CODE_PERMISSION_DENIED
|
||||||
|
if isinstance(error, RouteSetupError):
|
||||||
|
return common_pb2.ERROR_CODE_UNAVAILABLE
|
||||||
if isinstance(error, JobExecutionError):
|
if isinstance(error, JobExecutionError):
|
||||||
return common_pb2.ERROR_CODE_PRECONDITION_FAILED
|
return common_pb2.ERROR_CODE_PRECONDITION_FAILED
|
||||||
if isinstance(error, TransferError):
|
if isinstance(error, TransferError):
|
||||||
|
|||||||
@@ -3,6 +3,7 @@
|
|||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
import json
|
import json
|
||||||
|
import logging
|
||||||
import threading
|
import threading
|
||||||
import time
|
import time
|
||||||
import uuid
|
import uuid
|
||||||
@@ -16,6 +17,8 @@ from urllib import error, parse, request
|
|||||||
from archive_clients.config import ServiceConfig
|
from archive_clients.config import ServiceConfig
|
||||||
from archive_clients.resources import NormalizedResource, normalize_resource
|
from archive_clients.resources import NormalizedResource, normalize_resource
|
||||||
|
|
||||||
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
|
||||||
class QBittorrentError(RuntimeError):
|
class QBittorrentError(RuntimeError):
|
||||||
pass
|
pass
|
||||||
@@ -72,7 +75,17 @@ class QBittorrentReader:
|
|||||||
torrent_hash = torrent.get("hash")
|
torrent_hash = torrent.get("hash")
|
||||||
if not isinstance(torrent_hash, str):
|
if not isinstance(torrent_hash, str):
|
||||||
raise QBittorrentError("qBittorrent torrent hash is invalid")
|
raise QBittorrentError("qBittorrent torrent hash is invalid")
|
||||||
result.append(self._normalize(torrent, torrent_hash))
|
try:
|
||||||
|
result.append(self._normalize(torrent, torrent_hash))
|
||||||
|
except ValueError as exc:
|
||||||
|
# A stale or malformed qBittorrent entry must not hide every
|
||||||
|
# otherwise valid resource from Archive/Evict inventory.
|
||||||
|
# Keep it ineligible and leave an operator-visible diagnosis.
|
||||||
|
logger.warning(
|
||||||
|
"Skipping qBittorrent resource with inconsistent "
|
||||||
|
"metadata: hash=%s name=%r reason=%s",
|
||||||
|
torrent_hash, torrent.get("name"), exc,
|
||||||
|
)
|
||||||
return result
|
return result
|
||||||
|
|
||||||
def get_resource(self, torrent_hash: str) -> NormalizedResource | None:
|
def get_resource(self, torrent_hash: str) -> NormalizedResource | None:
|
||||||
|
|||||||
@@ -91,12 +91,21 @@ def normalize_resource(
|
|||||||
observed_at: datetime | None = None,
|
observed_at: datetime | None = None,
|
||||||
) -> NormalizedResource:
|
) -> NormalizedResource:
|
||||||
metainfo = decode_metainfo(metainfo_bytes)
|
metainfo = decode_metainfo(metainfo_bytes)
|
||||||
if len(raw_files) != len(metainfo.files):
|
# qBittorrent/libtorrent may omit torrent padding files from
|
||||||
|
# /torrents/files while keeping them in the exported metainfo.
|
||||||
|
visible_metainfo_files = tuple(
|
||||||
|
item for item in metainfo.files if not item.padding
|
||||||
|
)
|
||||||
|
if len(raw_files) == len(metainfo.files):
|
||||||
|
matched_metainfo_files = metainfo.files
|
||||||
|
elif len(raw_files) == len(visible_metainfo_files):
|
||||||
|
matched_metainfo_files = visible_metainfo_files
|
||||||
|
else:
|
||||||
raise ResourceError("qBittorrent and metainfo file counts differ")
|
raise ResourceError("qBittorrent and metainfo file counts differ")
|
||||||
files = []
|
files = []
|
||||||
canonical = True
|
canonical = True
|
||||||
for expected_index, (raw, meta_file) in enumerate(
|
for expected_index, (raw, meta_file) in enumerate(
|
||||||
zip(raw_files, metainfo.files, strict=True)
|
zip(raw_files, matched_metainfo_files, strict=True)
|
||||||
):
|
):
|
||||||
index = _integer(raw.get("index"), "file index")
|
index = _integer(raw.get("index"), "file index")
|
||||||
if index != expected_index:
|
if index != expected_index:
|
||||||
|
|||||||
@@ -7,6 +7,8 @@ import json
|
|||||||
import os
|
import os
|
||||||
import sqlite3
|
import sqlite3
|
||||||
import stat
|
import stat
|
||||||
|
import threading
|
||||||
|
from contextlib import contextmanager
|
||||||
from dataclasses import dataclass
|
from dataclasses import dataclass
|
||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
|
|
||||||
@@ -14,6 +16,12 @@ from pathlib import Path
|
|||||||
SCHEMA_VERSION = 3
|
SCHEMA_VERSION = 3
|
||||||
|
|
||||||
|
|
||||||
|
# A client daemon executes unrelated jobs concurrently. SQLite still permits
|
||||||
|
# only one writer, so serialize this process's short state transactions rather
|
||||||
|
# than allowing an otherwise healthy job to fail after its busy timeout.
|
||||||
|
_DATABASE_LOCK = threading.RLock()
|
||||||
|
|
||||||
|
|
||||||
class CommandConflict(RuntimeError):
|
class CommandConflict(RuntimeError):
|
||||||
pass
|
pass
|
||||||
|
|
||||||
@@ -557,14 +565,30 @@ class ClientStore:
|
|||||||
).fetchall()
|
).fetchall()
|
||||||
return [dict(row) for row in rows]
|
return [dict(row) for row in rows]
|
||||||
|
|
||||||
def _connect(self) -> sqlite3.Connection:
|
@contextmanager
|
||||||
connection = sqlite3.connect(self.database, isolation_level=None, timeout=5)
|
def _connect(self):
|
||||||
connection.row_factory = sqlite3.Row
|
"""Yield one connection while serializing local SQLite writers.
|
||||||
connection.execute("PRAGMA foreign_keys = ON")
|
|
||||||
connection.execute("PRAGMA journal_mode = WAL")
|
The longer SQLite timeout also covers a short lock held by a separate
|
||||||
connection.execute("PRAGMA synchronous = FULL")
|
maintenance process such as a backup. The lock is deliberately held
|
||||||
connection.execute("PRAGMA busy_timeout = 5000")
|
for the full transaction, including ``BEGIN IMMEDIATE``.
|
||||||
return connection
|
"""
|
||||||
|
with _DATABASE_LOCK:
|
||||||
|
connection = sqlite3.connect(
|
||||||
|
self.database, isolation_level=None, timeout=30
|
||||||
|
)
|
||||||
|
try:
|
||||||
|
connection.row_factory = sqlite3.Row
|
||||||
|
connection.execute("PRAGMA foreign_keys = ON")
|
||||||
|
connection.execute("PRAGMA journal_mode = WAL")
|
||||||
|
connection.execute("PRAGMA synchronous = FULL")
|
||||||
|
connection.execute("PRAGMA busy_timeout = 30000")
|
||||||
|
# Preserve the original ``with connection`` commit/rollback
|
||||||
|
# behavior used by every store operation.
|
||||||
|
with connection:
|
||||||
|
yield connection
|
||||||
|
finally:
|
||||||
|
connection.close()
|
||||||
|
|
||||||
|
|
||||||
def _canonical(value: object) -> str:
|
def _canonical(value: object) -> str:
|
||||||
|
|||||||
@@ -139,19 +139,14 @@ class SyncthingTransferObserver:
|
|||||||
|
|
||||||
def status(self) -> SyncthingTransferStatus:
|
def status(self) -> SyncthingTransferStatus:
|
||||||
published = load_published_transfer(self.local_job_directory)
|
published = load_published_transfer(self.local_job_directory)
|
||||||
total = 0
|
declared_files = [
|
||||||
for entry in published.manifest.files:
|
(entry.payload_relative_path, entry.logical_bytes)
|
||||||
total += _verified_job_file_size(
|
for entry in published.manifest.files
|
||||||
self.local_job_directory,
|
] + [
|
||||||
entry.payload_relative_path,
|
(artifact.payload_relative_path, artifact.logical_bytes)
|
||||||
entry.logical_bytes,
|
for artifact in published.manifest.artifacts
|
||||||
)
|
]
|
||||||
for artifact in published.manifest.artifacts:
|
total = sum(size for _, size in declared_files)
|
||||||
total += _verified_job_file_size(
|
|
||||||
self.local_job_directory,
|
|
||||||
artifact.payload_relative_path,
|
|
||||||
artifact.logical_bytes,
|
|
||||||
)
|
|
||||||
|
|
||||||
completion = self.transport.get_json(
|
completion = self.transport.get_json(
|
||||||
"/rest/db/completion?"
|
"/rest/db/completion?"
|
||||||
@@ -178,11 +173,44 @@ class SyncthingTransferObserver:
|
|||||||
name for name in needed_names
|
name for name in needed_names
|
||||||
if name == self.job_relative_path or name.startswith(prefix)
|
if name == self.job_relative_path or name.startswith(prefix)
|
||||||
}
|
}
|
||||||
fraction = float(raw_fraction) / 100
|
complete = not relevant
|
||||||
complete = fraction == 1 and 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(
|
return SyncthingTransferStatus(
|
||||||
fraction,
|
fraction,
|
||||||
total if complete else int(total * fraction),
|
completed_bytes,
|
||||||
total,
|
total,
|
||||||
complete,
|
complete,
|
||||||
len(relevant),
|
len(relevant),
|
||||||
@@ -472,6 +500,42 @@ def _verified_job_file_size(
|
|||||||
return metadata.st_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]:
|
def _needed_names(value: dict[str, Any]) -> set[str]:
|
||||||
result: set[str] = set()
|
result: set[str] = set()
|
||||||
for key in ("progress", "queued", "rest"):
|
for key in ("progress", "queued", "rest"):
|
||||||
|
|||||||
@@ -456,6 +456,36 @@ def cleanup_transfer(job_directory: Path) -> bool:
|
|||||||
return True
|
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:
|
def canonical_message_json(message: object) -> bytes:
|
||||||
value = json_format.MessageToDict(
|
value = json_format.MessageToDict(
|
||||||
message,
|
message,
|
||||||
|
|||||||
@@ -22,6 +22,94 @@ from archive_control.v1 import (
|
|||||||
|
|
||||||
|
|
||||||
class DaemonTransportTests(unittest.IsolatedAsyncioTestCase):
|
class DaemonTransportTests(unittest.IsolatedAsyncioTestCase):
|
||||||
|
async def test_unrelated_job_commands_execute_concurrently(self):
|
||||||
|
with tempfile.TemporaryDirectory() as directory:
|
||||||
|
root = Path(directory)
|
||||||
|
token = root / "token"
|
||||||
|
token.write_text("shared-secret", encoding="utf-8")
|
||||||
|
os.chmod(token, 0o600)
|
||||||
|
service = ServiceConfig("http://local", PurePosixPath("/api"), root)
|
||||||
|
config = ClientConfig(
|
||||||
|
"cache-1", "Cache 1", "cache", "ws://control", token,
|
||||||
|
root / "state.db", root / "backups", service, service,
|
||||||
|
)
|
||||||
|
probe = FilesystemProbe(root, True, True, True, True, True)
|
||||||
|
daemon = ArchiveClientDaemon(config, [probe, probe], [])
|
||||||
|
started: set[str] = set()
|
||||||
|
release = asyncio.Event()
|
||||||
|
|
||||||
|
async def execute(command, correlation_id, outbound):
|
||||||
|
started.add(command.assign_job.job.job_id)
|
||||||
|
await release.wait()
|
||||||
|
|
||||||
|
daemon._execute_job_command_locked = execute
|
||||||
|
commands = []
|
||||||
|
for _ in range(2):
|
||||||
|
command = control_pb2.Command(command_id=str(uuid4()))
|
||||||
|
command.assign_job.job.job_id = str(uuid4())
|
||||||
|
commands.append(command)
|
||||||
|
outbound = asyncio.Queue()
|
||||||
|
tasks = [
|
||||||
|
asyncio.create_task(
|
||||||
|
daemon._execute_job_command(command, "", outbound)
|
||||||
|
)
|
||||||
|
for command in commands
|
||||||
|
]
|
||||||
|
for _ in range(100):
|
||||||
|
if len(started) == 2:
|
||||||
|
break
|
||||||
|
await asyncio.sleep(0.01)
|
||||||
|
self.assertEqual(len(started), 2)
|
||||||
|
release.set()
|
||||||
|
await asyncio.gather(*tasks)
|
||||||
|
|
||||||
|
async def test_silent_control_connection_ends_for_reconnect(self):
|
||||||
|
"""A lost server heartbeat must not leave durable commands stranded."""
|
||||||
|
|
||||||
|
async def control(websocket):
|
||||||
|
registration = decode(await websocket.recv())
|
||||||
|
self.assertEqual(registration.WhichOneof("payload"), "register_request")
|
||||||
|
response = new_envelope()
|
||||||
|
response.correlation_id = registration.message_id
|
||||||
|
response.register_response.status = (
|
||||||
|
client_pb2.REGISTRATION_STATUS_ACCEPTED
|
||||||
|
)
|
||||||
|
response.register_response.negotiated_version.major = 1
|
||||||
|
await websocket.send(encode(response))
|
||||||
|
# Deliberately keep TCP/WebSocket open but send no application
|
||||||
|
# heartbeats. This models a stale proxy/server-side session.
|
||||||
|
await websocket.wait_closed()
|
||||||
|
|
||||||
|
with tempfile.TemporaryDirectory() as directory:
|
||||||
|
root = Path(directory)
|
||||||
|
token = root / "token"
|
||||||
|
token.write_text("shared-secret", encoding="utf-8")
|
||||||
|
os.chmod(token, 0o600)
|
||||||
|
async with serve(control, "127.0.0.1", 0, ping_interval=None) as server:
|
||||||
|
port = server.sockets[0].getsockname()[1]
|
||||||
|
service = ServiceConfig(
|
||||||
|
"http://local", PurePosixPath("/api"), root,
|
||||||
|
)
|
||||||
|
config = ClientConfig(
|
||||||
|
"cache-1", "Cache 1", "cache",
|
||||||
|
f"ws://127.0.0.1:{port}", token,
|
||||||
|
root / "state.db", root / "backups", service, service,
|
||||||
|
ConnectionConfig(
|
||||||
|
registration_timeout=1,
|
||||||
|
offline_timeout=0.05,
|
||||||
|
reconnect_initial=0.01,
|
||||||
|
reconnect_max=0.01,
|
||||||
|
reconnect_jitter=False,
|
||||||
|
),
|
||||||
|
)
|
||||||
|
probe = FilesystemProbe(root, True, True, True, True, True)
|
||||||
|
daemon = ArchiveClientDaemon(config, [probe, probe], [])
|
||||||
|
await asyncio.to_thread(daemon.store.initialize)
|
||||||
|
with self.assertRaisesRegex(
|
||||||
|
RuntimeError, "control heartbeat timed out"
|
||||||
|
):
|
||||||
|
await asyncio.wait_for(daemon._connection(), 1)
|
||||||
|
|
||||||
async def test_eviction_assignment_and_steps_are_admitted(self):
|
async def test_eviction_assignment_and_steps_are_admitted(self):
|
||||||
with tempfile.TemporaryDirectory() as directory:
|
with tempfile.TemporaryDirectory() as directory:
|
||||||
root = Path(directory)
|
root = Path(directory)
|
||||||
|
|||||||
+240
-1
@@ -1,5 +1,6 @@
|
|||||||
import hashlib
|
import hashlib
|
||||||
import tempfile
|
import tempfile
|
||||||
|
import threading
|
||||||
import unittest
|
import unittest
|
||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
from unittest.mock import Mock, patch
|
from unittest.mock import Mock, patch
|
||||||
@@ -11,6 +12,7 @@ from archive_clients.jobs import (
|
|||||||
JobExecutionError,
|
JobExecutionError,
|
||||||
_resource_fingerprint,
|
_resource_fingerprint,
|
||||||
)
|
)
|
||||||
|
from archive_clients.syncthing import RouteSetupError, SyncthingTransferStatus
|
||||||
from archive_clients.resources import normalize_resource
|
from archive_clients.resources import normalize_resource
|
||||||
from archive_clients.state import ClientStore
|
from archive_clients.state import ClientStore
|
||||||
from archive_control.v1 import control_pb2, job_pb2
|
from archive_control.v1 import control_pb2, job_pb2
|
||||||
@@ -31,7 +33,126 @@ class CompleteSyncthing:
|
|||||||
self.posts.append(path)
|
self.posts.append(path)
|
||||||
|
|
||||||
|
|
||||||
|
class SlowRescanSyncthing(CompleteSyncthing):
|
||||||
|
def post(self, path):
|
||||||
|
raise RouteSetupError("Syncthing API is unavailable")
|
||||||
|
|
||||||
|
|
||||||
class ClientJobHappyPathTests(unittest.TestCase):
|
class ClientJobHappyPathTests(unittest.TestCase):
|
||||||
|
def test_syncthing_api_outage_during_transfer_is_retried(self):
|
||||||
|
"""A transient local REST outage must not terminally fail the job."""
|
||||||
|
with tempfile.TemporaryDirectory() as directory:
|
||||||
|
root = Path(directory)
|
||||||
|
store = ClientStore(root / "client.db")
|
||||||
|
store.initialize()
|
||||||
|
definition = job_pb2.JobDefinition(
|
||||||
|
job_id=str(uuid4()),
|
||||||
|
transfer={
|
||||||
|
"source_client_id": "cache-1",
|
||||||
|
"target_client_id": "archive-1",
|
||||||
|
"route_id": "route-1",
|
||||||
|
},
|
||||||
|
)
|
||||||
|
observer = Mock()
|
||||||
|
observer.status.side_effect = [
|
||||||
|
RouteSetupError("Syncthing API is unavailable"),
|
||||||
|
SyncthingTransferStatus(1, 42, 42, True, 0),
|
||||||
|
]
|
||||||
|
executor = ClientJobExecutor(
|
||||||
|
client_id="archive-1",
|
||||||
|
qbittorrent=Mock(),
|
||||||
|
store=store,
|
||||||
|
qb_root=root,
|
||||||
|
qb_api_root=Path("/downloads"),
|
||||||
|
route_path=lambda _: root,
|
||||||
|
syncthing_transport=Mock(),
|
||||||
|
sparse_supported=True,
|
||||||
|
poll_interval=0,
|
||||||
|
)
|
||||||
|
executor._observer = Mock(return_value=observer)
|
||||||
|
progress = Mock()
|
||||||
|
|
||||||
|
executor._wait_for_syncthing(definition, progress)
|
||||||
|
|
||||||
|
self.assertEqual(observer.status.call_count, 2)
|
||||||
|
progress.assert_called_once_with(1, 42, 42, "0 Syncthing items still needed")
|
||||||
|
|
||||||
|
def test_reconnect_replay_never_duplicates_any_transfer_step(self):
|
||||||
|
for step in (
|
||||||
|
job_pb2.JOB_STEP_KIND_SOURCE_STAGE,
|
||||||
|
job_pb2.JOB_STEP_KIND_SYNCTHING_TRANSFER,
|
||||||
|
job_pb2.JOB_STEP_KIND_TARGET_MATERIALIZE,
|
||||||
|
job_pb2.JOB_STEP_KIND_QB_VERIFY,
|
||||||
|
job_pb2.JOB_STEP_KIND_STAGING_CLEANUP,
|
||||||
|
):
|
||||||
|
with self.subTest(step=step), tempfile.TemporaryDirectory() as directory:
|
||||||
|
root = Path(directory)
|
||||||
|
store = ClientStore(root / "client.db")
|
||||||
|
store.initialize()
|
||||||
|
definition = job_pb2.JobDefinition(
|
||||||
|
job_id=str(uuid4()),
|
||||||
|
idempotency_key=str(uuid4()),
|
||||||
|
operation=job_pb2.JOB_OPERATION_ARCHIVE,
|
||||||
|
resource_display_name="reconnect fixture",
|
||||||
|
transfer={
|
||||||
|
"source_client_id": "cache-1",
|
||||||
|
"target_client_id": "archive-1",
|
||||||
|
"route_id": "route-1",
|
||||||
|
},
|
||||||
|
)
|
||||||
|
definition.resource_id.info_hash_v1_hex = "a" * 40
|
||||||
|
definition.created_at.GetCurrentTime()
|
||||||
|
executor = ClientJobExecutor(
|
||||||
|
client_id="cache-1",
|
||||||
|
qbittorrent=Mock(),
|
||||||
|
store=store,
|
||||||
|
qb_root=root,
|
||||||
|
qb_api_root=Path("/downloads"),
|
||||||
|
route_path=lambda _: root,
|
||||||
|
syncthing_transport=Mock(),
|
||||||
|
sparse_supported=True,
|
||||||
|
)
|
||||||
|
executor.assign(control_pb2.AssignJobCommand(
|
||||||
|
job=definition,
|
||||||
|
expected_job_revision=1,
|
||||||
|
expected_last_event_sequence=0,
|
||||||
|
))
|
||||||
|
entered = threading.Event()
|
||||||
|
release = threading.Event()
|
||||||
|
|
||||||
|
def execute_step(*_args):
|
||||||
|
entered.set()
|
||||||
|
release.wait(1)
|
||||||
|
return None
|
||||||
|
|
||||||
|
executor._execute_step = Mock(side_effect=execute_step)
|
||||||
|
command = control_pb2.ExecuteStepCommand(
|
||||||
|
job_id=definition.job_id,
|
||||||
|
expected_job_revision=1,
|
||||||
|
expected_last_event_sequence=1,
|
||||||
|
step=step,
|
||||||
|
attempt=1,
|
||||||
|
)
|
||||||
|
results: list[list[control_pb2.JobEvent]] = []
|
||||||
|
first = threading.Thread(
|
||||||
|
target=lambda: results.append(executor.execute(command))
|
||||||
|
)
|
||||||
|
second = threading.Thread(
|
||||||
|
target=lambda: results.append(executor.execute(command))
|
||||||
|
)
|
||||||
|
first.start()
|
||||||
|
self.assertTrue(entered.wait(1))
|
||||||
|
second.start()
|
||||||
|
release.set()
|
||||||
|
first.join(1)
|
||||||
|
second.join(1)
|
||||||
|
|
||||||
|
self.assertFalse(first.is_alive())
|
||||||
|
self.assertFalse(second.is_alive())
|
||||||
|
self.assertEqual(executor._execute_step.call_count, 1)
|
||||||
|
self.assertEqual(len(results), 2)
|
||||||
|
self.assertEqual(results[0][-1].event_id, results[1][-1].event_id)
|
||||||
|
|
||||||
def test_capacity_guard_fails_before_data_movement(self):
|
def test_capacity_guard_fails_before_data_movement(self):
|
||||||
with tempfile.TemporaryDirectory() as directory:
|
with tempfile.TemporaryDirectory() as directory:
|
||||||
root = Path(directory)
|
root = Path(directory)
|
||||||
@@ -97,6 +218,123 @@ class ClientJobHappyPathTests(unittest.TestCase):
|
|||||||
)
|
)
|
||||||
self.assertEqual(required, 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(
|
def _run_transfer(
|
||||||
self,
|
self,
|
||||||
root: Path,
|
root: Path,
|
||||||
@@ -104,6 +342,7 @@ class ClientJobHappyPathTests(unittest.TestCase):
|
|||||||
*,
|
*,
|
||||||
source_stage_free_bytes: int | None = None,
|
source_stage_free_bytes: int | None = None,
|
||||||
content: bytes = b"archive-control-happy-path",
|
content: bytes = b"archive-control-happy-path",
|
||||||
|
syncthing: CompleteSyncthing | None = None,
|
||||||
):
|
):
|
||||||
source_root = root / "source"
|
source_root = root / "source"
|
||||||
target_root = root / "target"
|
target_root = root / "target"
|
||||||
@@ -171,7 +410,7 @@ class ClientJobHappyPathTests(unittest.TestCase):
|
|||||||
target_qb.get_resource.side_effect = [
|
target_qb.get_resource.side_effect = [
|
||||||
None, None, resource, resource,
|
None, None, resource, resource,
|
||||||
]
|
]
|
||||||
syncthing = CompleteSyncthing()
|
syncthing = syncthing or CompleteSyncthing()
|
||||||
source = ClientJobExecutor(
|
source = ClientJobExecutor(
|
||||||
client_id=source_id,
|
client_id=source_id,
|
||||||
qbittorrent=source_qb,
|
qbittorrent=source_qb,
|
||||||
|
|||||||
@@ -205,6 +205,54 @@ class QBittorrentReaderTests(unittest.TestCase):
|
|||||||
self.assertEqual(resources, [])
|
self.assertEqual(resources, [])
|
||||||
self.assertEqual(len(opener.calls), 2)
|
self.assertEqual(len(opener.calls), 2)
|
||||||
|
|
||||||
|
def test_listing_skips_malformed_torrent_and_keeps_valid_resources(self):
|
||||||
|
def torrent(name, content):
|
||||||
|
info = {
|
||||||
|
b"length": len(content), b"name": name.encode(),
|
||||||
|
b"piece length": 16384, b"pieces": b"x" * 20,
|
||||||
|
}
|
||||||
|
return encode({b"info": info}), hashlib.sha1(encode(info)).hexdigest()
|
||||||
|
|
||||||
|
first_bytes, first_hash = torrent("first.txt", b"one")
|
||||||
|
malformed_bytes, malformed_hash = torrent("broken.txt", b"bad")
|
||||||
|
second_bytes, second_hash = torrent("second.txt", b"two")
|
||||||
|
torrents = [
|
||||||
|
{"hash": first_hash, "name": "first.txt", "state": "uploading"},
|
||||||
|
{"hash": malformed_hash, "name": "broken.txt", "state": "stalledUP"},
|
||||||
|
{"hash": second_hash, "name": "second.txt", "state": "uploading"},
|
||||||
|
]
|
||||||
|
responses = [
|
||||||
|
b"Ok.", json.dumps(torrents).encode(),
|
||||||
|
b'[{"index":0,"name":"first.txt","size":3,"progress":1,"priority":1}]',
|
||||||
|
first_bytes,
|
||||||
|
b'[{"index":0,"name":"broken.txt","size":3,"progress":1,"priority":1},'
|
||||||
|
b'{"index":1,"name":"unexpected.txt","size":1,"progress":1,"priority":1}]',
|
||||||
|
malformed_bytes,
|
||||||
|
b'[{"index":0,"name":"second.txt","size":3,"progress":1,"priority":1}]',
|
||||||
|
second_bytes,
|
||||||
|
]
|
||||||
|
with tempfile.TemporaryDirectory() as directory:
|
||||||
|
root = Path(directory)
|
||||||
|
password = root / "password"
|
||||||
|
password.write_text("secret", encoding="utf-8")
|
||||||
|
os.chmod(password, 0o600)
|
||||||
|
config = ServiceConfig(
|
||||||
|
"http://qb", PurePosixPath("/downloads"), root,
|
||||||
|
username="admin", password_file=password,
|
||||||
|
)
|
||||||
|
opener = _Opener(responses)
|
||||||
|
with patch(
|
||||||
|
"archive_clients.qbittorrent.request.build_opener",
|
||||||
|
return_value=opener,
|
||||||
|
), self.assertLogs("archive_clients.qbittorrent", "WARNING") as logs:
|
||||||
|
resources = QBittorrentReader(config).list_resources()
|
||||||
|
|
||||||
|
self.assertEqual(
|
||||||
|
[resource.summary.display_name for resource in resources],
|
||||||
|
["first.txt", "second.txt"],
|
||||||
|
)
|
||||||
|
self.assertIn(malformed_hash, "\n".join(logs.output))
|
||||||
|
|
||||||
def test_stopped_add_selection_recheck_and_entry_only_delete(self):
|
def test_stopped_add_selection_recheck_and_entry_only_delete(self):
|
||||||
torrent_hash = "a" * 40
|
torrent_hash = "a" * 40
|
||||||
responses = [
|
responses = [
|
||||||
|
|||||||
@@ -118,6 +118,47 @@ class ResourceTests(unittest.TestCase):
|
|||||||
self.assertEqual(root.available_file_count, 1)
|
self.assertEqual(root.available_file_count, 1)
|
||||||
self.assertEqual(root.available_logical_bytes, 3)
|
self.assertEqual(root.available_logical_bytes, 3)
|
||||||
|
|
||||||
|
def test_qbittorrent_hidden_padding_files_are_normalized(self):
|
||||||
|
info = {
|
||||||
|
b"files": [
|
||||||
|
{b"length": 3, b"path": [b"first.bin"]},
|
||||||
|
{
|
||||||
|
b"attr": b"p", b"length": 5,
|
||||||
|
b"path": [b"_____padding_file_5"],
|
||||||
|
},
|
||||||
|
{b"length": 7, b"path": [b"last.bin"]},
|
||||||
|
],
|
||||||
|
b"name": b"with-padding",
|
||||||
|
b"piece length": 16384,
|
||||||
|
b"pieces": b"x" * 20,
|
||||||
|
}
|
||||||
|
metainfo = encode({b"info": info})
|
||||||
|
torrent_hash = hashlib.sha1(encode(info)).hexdigest()
|
||||||
|
normalized = normalize_resource(
|
||||||
|
{
|
||||||
|
"hash": torrent_hash,
|
||||||
|
"name": "with-padding",
|
||||||
|
"state": "stalledUP",
|
||||||
|
},
|
||||||
|
[
|
||||||
|
{
|
||||||
|
"index": 0, "name": "with-padding/first.bin",
|
||||||
|
"size": 3, "completed": 3, "priority": 1,
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"index": 1, "name": "with-padding/last.bin",
|
||||||
|
"size": 7, "completed": 7, "priority": 1,
|
||||||
|
},
|
||||||
|
],
|
||||||
|
metainfo,
|
||||||
|
)
|
||||||
|
self.assertEqual(
|
||||||
|
[item.canonical_path for item in normalized.files],
|
||||||
|
["with-padding/first.bin", "with-padding/last.bin"],
|
||||||
|
)
|
||||||
|
self.assertEqual(normalized.summary.total_file_count, 2)
|
||||||
|
self.assertEqual(normalized.summary.selected_complete_bytes, 10)
|
||||||
|
|
||||||
def test_noncanonical_and_unsafe_paths_are_distinct(self):
|
def test_noncanonical_and_unsafe_paths_are_distinct(self):
|
||||||
info = {
|
info = {
|
||||||
b"length": 3,
|
b"length": 3,
|
||||||
|
|||||||
@@ -1,6 +1,7 @@
|
|||||||
import tempfile
|
import tempfile
|
||||||
import unittest
|
import unittest
|
||||||
import os
|
import os
|
||||||
|
import threading
|
||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
from uuid import uuid4
|
from uuid import uuid4
|
||||||
|
|
||||||
@@ -13,6 +14,40 @@ from archive_clients.state import (
|
|||||||
|
|
||||||
|
|
||||||
class ClientStoreTests(unittest.TestCase):
|
class ClientStoreTests(unittest.TestCase):
|
||||||
|
def test_database_connections_serialize_local_writers(self):
|
||||||
|
"""Concurrent jobs share one daemon DB without SQLite lock failures."""
|
||||||
|
with tempfile.TemporaryDirectory() as directory:
|
||||||
|
database = Path(directory) / "state.db"
|
||||||
|
first = ClientStore(database)
|
||||||
|
second = ClientStore(database)
|
||||||
|
first.initialize()
|
||||||
|
barrier = threading.Barrier(2)
|
||||||
|
failures: list[Exception] = []
|
||||||
|
|
||||||
|
def write(store, command_id):
|
||||||
|
try:
|
||||||
|
barrier.wait()
|
||||||
|
store.accept_command(command_id, '{"kind":"job"}', '{"ok":true}')
|
||||||
|
except Exception as exc: # pragma: no cover - assertion below
|
||||||
|
failures.append(exc)
|
||||||
|
|
||||||
|
left = threading.Thread(target=write, args=(first, "command-1"))
|
||||||
|
right = threading.Thread(target=write, args=(second, "command-2"))
|
||||||
|
left.start()
|
||||||
|
right.start()
|
||||||
|
left.join(1)
|
||||||
|
right.join(1)
|
||||||
|
|
||||||
|
self.assertFalse(left.is_alive())
|
||||||
|
self.assertFalse(right.is_alive())
|
||||||
|
self.assertEqual(failures, [])
|
||||||
|
self.assertEqual(len(first.list_accepted_commands()), 2)
|
||||||
|
with first._connect() as connection:
|
||||||
|
self.assertEqual(
|
||||||
|
connection.execute("PRAGMA busy_timeout").fetchone()[0],
|
||||||
|
30000,
|
||||||
|
)
|
||||||
|
|
||||||
def test_command_acceptance_is_durable_and_content_addressed(self):
|
def test_command_acceptance_is_durable_and_content_addressed(self):
|
||||||
with tempfile.TemporaryDirectory() as directory:
|
with tempfile.TemporaryDirectory() as directory:
|
||||||
database = Path(directory) / "state.db"
|
database = Path(directory) / "state.db"
|
||||||
|
|||||||
@@ -11,6 +11,7 @@ from archive_clients.transfer import (
|
|||||||
FileMaterializer,
|
FileMaterializer,
|
||||||
TransferIntegrityError,
|
TransferIntegrityError,
|
||||||
canonical_message_json,
|
canonical_message_json,
|
||||||
|
cleanup_partial_transfer,
|
||||||
load_published_transfer,
|
load_published_transfer,
|
||||||
materialize_transfer,
|
materialize_transfer,
|
||||||
stage_transfer,
|
stage_transfer,
|
||||||
@@ -206,6 +207,33 @@ class TransferHappyPathTests(unittest.TestCase):
|
|||||||
os.stat(copied).st_blocks * 512, os.stat(copied).st_size
|
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):
|
def test_canonical_manifest_json_is_stable(self):
|
||||||
manifest = self._manifest()
|
manifest = self._manifest()
|
||||||
first = canonical_message_json(manifest)
|
first = canonical_message_json(manifest)
|
||||||
|
|||||||
Reference in New Issue
Block a user