Compare commits

...
17 Commits
Author SHA1 Message Date
cabbage 1009defc46 release: archive clients v0.1.14 2026-07-27 03:21:15 +00:00
cabbage 8ed8e75447 feat: allow unrestricted independent job concurrency 2026-07-27 03:19:47 +00:00
cabbage 1a06da3984 fix: reconnect after silent control heartbeat loss 2026-07-27 02:45:21 +00:00
cabbage 4ba5e92248 fix: map x2 route override by folder path 2026-07-26 04:34:08 +00:00
cabbage a6da224f28 docs: require hardlink-safe route mount topology 2026-07-26 04:26:23 +00:00
cabbage 90925fd321 fix: use job-specific sparse transfer progress 2026-07-25 03:01:34 +00:00
cabbage 9df2e6991e fix: publish initial sparse transfer progress 2026-07-25 02:50:22 +00:00
cabbage d2fa69a1d6 fix: report partial Syncthing transfer progress 2026-07-25 02:37:30 +00:00
cabbage 0f94388f49 fix: advance target steps past peer events 2026-07-24 23:27:11 +00:00
cabbage 94a3233ce2 fix: clean interrupted staging temporaries 2026-07-24 17:00:10 +00:00
cabbage 053b0135b3 fix: defer slow Syncthing rescans 2026-07-24 16:13:39 +00:00
cabbage 57acc90363 fix: honor hardlinks across client mount topology 2026-07-24 16:05:20 +00:00
cabbage a73bcf870f config: reduce zero-copy space reserve default 2026-07-24 15:26:33 +00:00
cabbage 748ca49837 fix: avoid duplicate capacity reservation for hardlinks 2026-07-24 13:48:24 +00:00
cabbage 589fd4a41d fix: normalize existing Syncthing route paths 2026-07-24 12:52:30 +00:00
cabbage d951e2e4fc fix: require qBittorrent metainfo export support 2026-07-24 12:37:40 +00:00
cabbage 21ea431c56 fix: expand configured service usernames 2026-07-23 16:34:37 +00:00
31 changed files with 987 additions and 114 deletions
+31 -4
View File
@@ -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:
+1 -1
View File
@@ -19,7 +19,7 @@ reconnect_jitter = true
stall_after = "30m" # Status warning only; it does not fail the job. stall_after = "30m" # Status warning only; it does not fail the job.
verification_timeout = "30m" verification_timeout = "30m"
poll_interval = "1s" 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] [backup]
interval = "6h" interval = "6h"
+1 -1
View File
@@ -2,7 +2,7 @@ name: archive-control-archive
services: services:
archive-client: archive-client:
image: sodium/archive-clients:v0.1.0 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"]
+4 -1
View File
@@ -19,7 +19,7 @@ reconnect_jitter = true
stall_after = "30m" # Status warning only; it does not fail the job. stall_after = "30m" # Status warning only; it does not fail the job.
verification_timeout = "30m" verification_timeout = "30m"
poll_interval = "1s" 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] [backup]
interval = "6h" interval = "6h"
@@ -39,4 +39,7 @@ 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"
# 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" }
advertised_addresses = ["dynamic"] advertised_addresses = ["dynamic"]
+1 -2
View File
@@ -2,7 +2,7 @@ name: archive-control-cache
services: services:
archive-client: archive-client:
image: sodium/archive-clients:v0.1.0 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/DownloadsSync
+5 -1
View File
@@ -19,7 +19,7 @@ reconnect_jitter = true
stall_after = "30m" # Status warning only; it does not fail the job. stall_after = "30m" # Status warning only; it does not fail the job.
verification_timeout = "30m" verification_timeout = "30m"
poll_interval = "1s" 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] [backup]
interval = "6h" interval = "6h"
@@ -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"]
+1 -2
View File
@@ -2,7 +2,7 @@ name: archive-control-cache
services: services:
archive-client: archive-client:
image: sodium/archive-clients:v0.1.0 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
View File
@@ -35,8 +35,11 @@ expand these rules but must not contradict them.
manifests, and observed external state. manifests, and observed external state.
- SQLite backups use the online backup API. Defaults are every six hours, 12 - SQLite backups use the online backup API. Defaults are every six hours, 12
recent, 14 daily, and eight weekly copies, plus pre/post-migration backups. recent, 14 daily, and eight weekly copies, plus pre/post-migration backups.
- Default concurrency is one active data-moving job per client and per route. - Default concurrency is unrestricted: every eligible queued job is admitted
Disjoint node pairs may run concurrently. Queueing is durable FIFO. in one scheduler pass and independent jobs execute concurrently on a client.
Physical qBittorrent, Syncthing, disk, and network capacity are therefore
the natural limit. `enforce_concurrency_limits=true` restores the optional
legacy per-client/per-route gates. Commands for one job remain serialized.
- A queued job owns a per-resource reservation. Cancelling it removes only the - A queued job owns a per-resource reservation. Cancelling it removes only the
queued record/reservation and never sends cleanup commands. queued record/reservation and never sends cleanup commands.
- Connectivity or transfer stalls wait indefinitely. A configurable 30-minute - Connectivity or transfer stalls wait indefinitely. A configurable 30-minute
+15 -2
View File
@@ -44,6 +44,9 @@ command_max_attempts = 3 # Total sends, including the initial attempt.
stall_after = "30m" # Warning state only; jobs continue waiting. 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
@@ -102,7 +105,7 @@ stall_after = "30m" # Warning only; no automatic job failure.
verification_timeout = "30m" # Stopped qB full-recheck deadline. verification_timeout = "30m" # Stopped qB full-recheck deadline.
poll_interval = "1s" # Active qB/Syncthing observation interval. poll_interval = "1s" # Active qB/Syncthing observation interval.
# Worst-case copy fallback must leave this many bytes free. # 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] [backup]
interval = "6h" interval = "6h"
@@ -125,6 +128,16 @@ local_root = "/data/sync"
advertised_addresses = ["dynamic"] 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 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
override them per node. override them per node.
@@ -213,7 +226,7 @@ cache/archive routes according to policy.
```yaml ```yaml
services: services:
archive-client: archive-client:
image: sodium/archive-clients:v0.1.0 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"]
+2 -1
View File
@@ -105,7 +105,8 @@ follow-up task.
### Actions ### Actions
1. Implement qBittorrent cookie authentication and capability adapters for 1. Implement qBittorrent cookie authentication and capability adapters for
supported 4.44.6 and 5.x APIs. supported 4.54.6 and 5.x APIs; reject older releases because they cannot
export exact torrent metainfo through the public API.
2. Implement on-demand torrent summaries, scoped hash lookup, lazy complete 2. Implement on-demand torrent summaries, scoped hash lookup, lazy complete
content trees, v1/v2 identity validation, torrent export, stopped add, content trees, v1/v2 identity validation, torrent export, stopped add,
selection application, recheck monitoring, and entry-only deletion. selection application, recheck monitoring, and entry-only deletion.
+11 -5
View File
@@ -6,7 +6,7 @@ shapes.
## qBittorrent Web API ## qBittorrent Web API
The implementation targets supported qBittorrent 4.44.6 and 5.x releases and The implementation targets qBittorrent 4.54.6 and 5.x releases and
detects the application and Web API versions at startup. The authoritative detects the application and Web API versions at startup. The authoritative
references are the official [5.0 WebUI API](https://github.com/qbittorrent/qBittorrent/wiki/WebUI-API-%28qBittorrent-5.0%29) references are the official [5.0 WebUI API](https://github.com/qbittorrent/qBittorrent/wiki/WebUI-API-%28qBittorrent-5.0%29)
and [4.1-compatible WebUI API](https://github.com/qbittorrent/qBittorrent/wiki/WebUI-API-%28qBittorrent-4.1%29). and [4.1-compatible WebUI API](https://github.com/qbittorrent/qBittorrent/wiki/WebUI-API-%28qBittorrent-4.1%29).
@@ -16,7 +16,7 @@ and [4.1-compatible WebUI API](https://github.com/qbittorrent/qBittorrent/wiki/W
| Need | Web API family | Archive Control rule | | Need | Web API family | Archive Control rule |
| --- | --- | --- | | --- | --- | --- |
| Authenticate | `auth/login`, cookie session | Log no credentials/cookies; reauthenticate once on expiry | | Authenticate | `auth/login`, cookie session | Log no credentials/cookies; reauthenticate once on expiry |
| Detect compatibility | `app/version`, `app/webapiVersion`, build info | Advertise exact versions and select adapter | | Detect compatibility | `app/version`, `app/webapiVersion`, build info | Require qBittorrent 4.5.0+ because exact metainfo export is mandatory |
| List/lookup | `torrents/info`, properties | On-demand and hash-scoped where possible | | List/lookup | `torrents/info`, properties | On-demand and hash-scoped where possible |
| Read file state | `torrents/files` | Normalize indices, paths, selected/skipped, size, progress | | Read file state | `torrents/files` | Normalize indices, paths, selected/skipped, size, progress |
| Export metainfo | `torrents/export` | Required before staging; fail if exact metainfo unavailable | | Export metainfo | `torrents/export` | Required before staging; fail if exact metainfo unavailable |
@@ -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, stopped states. It does not mutate qBittorrent preferences, categories, tags,
limits, queueing defaults, or global save-path behavior. 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 ### Inventory normalization
For each torrent the client derives canonical v1/v2 identity from reliable API 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 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
View File
@@ -133,10 +133,13 @@ a resource may be active or queued at a time. This prevents incompatible
baselines even when jobs would use different nodes or observe different sides baselines even when jobs would use different nodes or observe different sides
of a hybrid identity. of a hybrid identity.
Defaults allow one active data-moving job per client and one per route. A job By default every eligible queued job is claimed in the same scheduler pass and
must acquire its source client, target client, route, and resource reservation independent jobs run concurrently on each client; physical service and storage
atomically. Disjoint node pairs may run concurrently. Route setup is a capacity are the limit. Set `enforce_concurrency_limits=true` on control to
preflight activity and does not permit a data step to bypass these leases. enable the optional one-per-client/one-per-route gates. A job always retains
its resource reservation, and commands for the same job remain serialized.
Route setup is a preflight activity and does not permit a data step to bypass
the resource reservation.
Offline nodes do not prevent unrelated jobs from being listed or run. A job Offline nodes do not prevent unrelated jobs from being listed or run. A job
requiring an offline node remains waiting indefinitely; it does not consume an requiring an offline node remains waiting indefinitely; it does not consume an
@@ -154,6 +157,12 @@ advances to the next step. On reconnect:
4. Control observes relevant qBittorrent, Syncthing, staging, and manifest 4. Control observes relevant qBittorrent, Syncthing, staging, and manifest
state before selecting retry, resume, compensation, cleanup, or manual state before selecting retry, resume, compensation, cleanup, or manual
intervention. intervention.
If a connection disappears while a synchronous job operation is in progress,
the client journals events locally, reconnects indefinitely, and serializes a
replayed command behind the in-flight operation for that job. The replay reads
the journal rather than repeating the data operation. This applies to every
transfer step, including post-commit staging cleanup.
5. A command is reissued with its original ID when the acceptance result is 5. A command is reissued with its original ID when the acceptance result is
uncertain. uncertain.
+4 -2
View File
@@ -25,8 +25,10 @@ What do you want to do?
[ Evict Cache ] [ Job Status ] [ Evict Cache ] [ Job Status ]
``` ```
All subsequent pages edit this message. `Cancel` closes the active selection All subsequent pages edit this message. Leaving a resource/tree selection
flow. `Back` returns one level while retaining validated filters/selections. returns to the Archive Control operation menu without closing the shared
conversation. Confirmation `Cancel` returns to its immediate prior selection;
the main menu alone can close Archive Control.
Callback payloads contain opaque session/action IDs, not resource names or Callback payloads contain opaque session/action IDs, not resource names or
paths, and are validated against persisted session revision and expiry. paths, and are validated against persisted session revision and expiry.
+3 -3
View File
@@ -89,8 +89,8 @@ overlapping selections, reordered/duplicated events, and partial journals.
### Adapter contract tests ### Adapter contract tests
Pinned qBittorrent and Syncthing container versions cover every supported API Pinned qBittorrent 4.54.6 and 5.x and Syncthing container versions cover every
family. Assertions include stopped add, exact selection, export, full recheck, supported API family. Assertions include stopped add, exact selection, export, full recheck,
download-attempt detection, entry-only delete, folder/device idempotency, download-attempt detection, entry-only delete, folder/device idempotency,
events fallback, need/completion proof, and redacted errors. Partfile adapters events fallback, need/completion proof, and redacted errors. Partfile adapters
are tested only against explicitly supported qBittorrent/libtorrent fixtures. are tested only against explicitly supported qBittorrent/libtorrent fixtures.
@@ -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 |
+4
View File
@@ -102,6 +102,10 @@ and placement policy differ.
external state rather than being downloaded or silently deselected. external state rather than being downloaded or silently deselected.
- Confirm free space/reserve, permissions, sparse capability, supported - Confirm free space/reserve, permissions, sparse capability, supported
partfile format, route health, and negotiated protocol features. 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. - Export the source `.torrent` metainfo.
- Persist the exact source, target baseline, requested selection, and delta. - Persist the exact source, target baseline, requested selection, and delta.
+1 -1
View File
@@ -21,7 +21,7 @@ verification_timeout = "30m"
# Poll qBittorrent and Syncthing at this interval while a step is active. # Poll qBittorrent and Syncthing at this interval while a step is active.
poll_interval = "1s" poll_interval = "1s"
# Space that must remain free after a worst-case copy fallback. # Space that must remain free after a worst-case copy fallback.
free_space_reserve_bytes = 1073741824 free_space_reserve_bytes = 33554432
[backup] [backup]
interval = "6h" interval = "6h"
+1 -1
View File
@@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta"
[project] [project]
name = "archive-clients" name = "archive-clients"
version = "0.1.0" version = "0.1.14"
requires-python = ">=3.11" requires-python = ">=3.11"
dependencies = ["protobuf==7.35.1", "websockets==16.0"] dependencies = ["protobuf==7.35.1", "websockets==16.0"]
+64 -11
View File
@@ -27,17 +27,28 @@ _ENV = re.compile(r"\$\{([A-Za-z_][A-Za-z0-9_]*)\}")
class RootMapping: class RootMapping:
api_root: PurePosixPath api_root: PurePosixPath
local_root: Path local_root: Path
local_path_overrides: tuple[tuple[PurePosixPath, Path], ...] = ()
def api_to_local(self, api_path: str) -> Path: def api_to_local(self, api_path: str) -> Path:
candidate = PurePosixPath(api_path) candidate = PurePosixPath(api_path)
try: for api_root, local_root in self.local_path_overrides:
relative = candidate.relative_to(self.api_root) relative = _safe_relative(candidate, api_root)
except ValueError as exc: if relative is not None:
raise ConfigError("API path is outside its configured root") from exc return local_root.joinpath(*relative.parts)
if any(part in {"", ".", ".."} for part in relative.parts): relative = _safe_relative(candidate, self.api_root)
raise ConfigError("API path contains an unsafe component") if relative is None:
raise ConfigError("API path is outside its configured root")
return self.local_root.joinpath(*relative.parts) 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) @dataclass(frozen=True)
class ServiceConfig: class ServiceConfig:
@@ -48,10 +59,13 @@ class ServiceConfig:
password_file: Path | None = None password_file: Path | None = None
api_key_file: Path | None = None api_key_file: Path | None = None
advertised_addresses: tuple[str, ...] = () advertised_addresses: tuple[str, ...] = ()
local_path_overrides: tuple[tuple[PurePosixPath, Path], ...] = ()
@property @property
def roots(self) -> RootMapping: 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: def read_password(self) -> str | None:
return ( return (
@@ -82,7 +96,7 @@ class JobsConfig:
stall_after: float = 30 * 60 stall_after: float = 30 * 60
verification_timeout: float = 30 * 60 verification_timeout: float = 30 * 60
poll_interval: float = 1 poll_interval: float = 1
free_space_reserve_bytes: int = 1024 * 1024 * 1024 free_space_reserve_bytes: int = 32 * 1024 * 1024
@dataclass(frozen=True) @dataclass(frozen=True)
@@ -160,7 +174,7 @@ def _service(value: Any, name: str) -> ServiceConfig:
raise ConfigError(f"{name} must be a table") raise ConfigError(f"{name} must be a table")
allowed = { allowed = {
"endpoint", "api_root", "local_root", "username", "password_file", "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) _keys(value, allowed, name)
api_root = PurePosixPath(_string(value, "api_root")) api_root = PurePosixPath(_string(value, "api_root"))
@@ -169,6 +183,8 @@ def _service(value: Any, name: str) -> ServiceConfig:
username = value.get("username") username = value.get("username")
if username is not None and (not isinstance(username, str) or not username): if username is not None and (not isinstance(username, str) or not username):
raise ConfigError(f"{name}.username must be a non-empty string") raise ConfigError(f"{name}.username must be a non-empty string")
if username is not None:
username = _expand(username)
addresses = value.get("advertised_addresses", []) addresses = value.get("advertised_addresses", [])
if not isinstance(addresses, list) or any( if not isinstance(addresses, list) or any(
not isinstance(address, str) or not address for address in addresses not isinstance(address, str) or not address for address in addresses
@@ -180,6 +196,7 @@ def _service(value: Any, name: str) -> ServiceConfig:
raise ConfigError("qbittorrent username and password_file are required") raise ConfigError("qbittorrent username and password_file are required")
if name == "syncthing" and "api_key_file" not in value: if name == "syncthing" and "api_key_file" not in value:
raise ConfigError("syncthing api_key_file is required") raise ConfigError("syncthing api_key_file is required")
overrides = _local_path_overrides(value, api_root, name)
return ServiceConfig( return ServiceConfig(
_endpoint(value, "endpoint", {"http", "https"}), api_root, _endpoint(value, "endpoint", {"http", "https"}), api_root,
_absolute_path(value, "local_root"), username, _absolute_path(value, "local_root"), username,
@@ -187,10 +204,46 @@ def _service(value: Any, name: str) -> ServiceConfig:
if "password_file" in value else None, if "password_file" in value else None,
_absolute_path(value, "api_key_file") _absolute_path(value, "api_key_file")
if "api_key_file" in value else None, 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: def _connection(value: Any) -> ConnectionConfig:
if not isinstance(value, dict): if not isinstance(value, dict):
raise ConfigError("connection must be a table") raise ConfigError("connection must be a table")
@@ -237,7 +290,7 @@ def _jobs(value: Any) -> JobsConfig:
), ),
poll_interval=_duration(value.get("poll_interval", "1s")), poll_interval=_duration(value.get("poll_interval", "1s")),
free_space_reserve_bytes=_positive_int( 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", "jobs.free_space_reserve_bytes",
), ),
) )
+38 -3
View File
@@ -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
) )
+196 -35
View File
@@ -4,12 +4,15 @@ from __future__ import annotations
import hashlib import hashlib
import json import json
import logging
import os
import shutil import shutil
import stat
import threading import threading
import time import time
import uuid import uuid
from pathlib import Path, PurePosixPath 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.protocol import decode_message, encode_message
from archive_clients.eviction import ( from archive_clients.eviction import (
@@ -25,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,
@@ -51,6 +55,9 @@ class JobCancelled(JobExecutionError):
pass pass
logger = logging.getLogger(__name__)
class ClientJobExecutor: class ClientJobExecutor:
def __init__( def __init__(
self, self,
@@ -65,7 +72,7 @@ class ClientJobExecutor:
sparse_supported: bool, sparse_supported: bool,
poll_interval: float = 1, poll_interval: float = 1,
verification_timeout: float = 30 * 60, 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.client_id = client_id
self.qbittorrent = qbittorrent self.qbittorrent = qbittorrent
@@ -79,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()
@@ -108,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,
@@ -125,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,
@@ -160,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
@@ -541,15 +593,29 @@ class ClientJobExecutor:
sha256_hex=hashlib.sha256(resource.metainfo_bytes).hexdigest(), sha256_hex=hashlib.sha256(resource.metainfo_bytes).hexdigest(),
) )
del artifact 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._require_space(
self.route_path(definition.transfer.route_id), route_root,
definition.transfer.transfer_delta_logical_bytes 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), + len(resource.metainfo_bytes),
) )
stage_transfer( stage_transfer(
manifest, manifest,
source_root=self.qb_root, source_root=self.qb_root,
sync_root=self.route_path(definition.transfer.route_id), sync_root=route_root,
store=self.store, store=self.store,
artifact_sources={"metainfo/source.torrent": metainfo_path}, artifact_sources={"metainfo/source.torrent": metainfo_path},
sparse_supported=self.sparse_supported, sparse_supported=self.sparse_supported,
@@ -563,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,
@@ -607,9 +683,18 @@ class ClientJobExecutor:
"target materialization was sent to the wrong client" "target materialization was sent to the wrong client"
) )
published = load_published_transfer(self._job_directory(definition)) 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._require_space(
self.qb_root, 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) info_hash = _info_hash(definition)
resource = self.qbittorrent.get_resource(info_hash) resource = self.qbittorrent.get_resource(info_hash)
@@ -794,6 +879,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():
@@ -843,6 +933,37 @@ class ClientJobExecutor:
f"{required} bytes required including reserve" 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( def _observer(
self, definition: job_pb2.JobDefinition self, definition: job_pb2.JobDefinition
) -> SyncthingTransferObserver: ) -> SyncthingTransferObserver:
@@ -1089,6 +1210,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):
@@ -1096,3 +1219,41 @@ def _job_error_code(error: Exception) -> int:
if isinstance(error, EvictionError): if isinstance(error, EvictionError):
return common_pb2.ERROR_CODE_PRECONDITION_FAILED return common_pb2.ERROR_CODE_PRECONDITION_FAILED
return common_pb2.ERROR_CODE_INTERNAL 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 -1
View File
@@ -37,7 +37,9 @@ def discover_routes(
relative = normalized_api_path.relative_to(roots.api_root) relative = normalized_api_path.relative_to(roots.api_root)
except (ConfigError, ValueError): except (ConfigError, ValueError):
continue 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 continue
if relative == PurePosixPath("."): if relative == PurePosixPath("."):
continue continue
+3 -1
View File
@@ -159,7 +159,9 @@ def _supported_qb_version(version: str) -> bool:
if match is None: if match is None:
return False return False
major, minor = int(match.group(1)), int(match.group(2)) major, minor = int(match.group(1)), int(match.group(2))
return major == 5 or (major == 4 and minor in {4, 5, 6}) # torrents/export, required to preserve exact metainfo before staging,
# was introduced in qBittorrent 4.5.0.
return major == 5 or (major == 4 and minor in {5, 6})
def _text_get( def _text_get(
+114 -19
View File
@@ -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),
@@ -356,7 +384,9 @@ class SyncthingRouteManager:
local_path = self.config.roots.api_to_local(api_path) local_path = self.config.roots.api_to_local(api_path)
except ConfigError as exc: except ConfigError as exc:
raise RoutePathConflict("route path is outside the sync root") from 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: try:
local_path.resolve(strict=False).relative_to(resolved_root) local_path.resolve(strict=False).relative_to(resolved_root)
except ValueError as exc: except ValueError as exc:
@@ -369,8 +399,8 @@ class SyncthingRouteManager:
raise RoutePathConflict("new route path contains unrelated data") raise RoutePathConflict("new route path contains unrelated data")
path.mkdir(parents=True, exist_ok=True) path.mkdir(parents=True, exist_ok=True)
@staticmethod
def _validate_folder( def _validate_folder(
self,
folder: dict[str, Any], folder: dict[str, Any],
api_path: str, api_path: str,
local_device_id: str, local_device_id: str,
@@ -382,7 +412,15 @@ class SyncthingRouteManager:
for item in raw_devices for item in raw_devices
if isinstance(item, dict) and isinstance(item.get("deviceID"), str) if isinstance(item, dict) and isinstance(item.get("deviceID"), str)
} if isinstance(raw_devices, list) else set() } 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") raise RoutePathConflict("existing route ID uses a different path")
if folder.get("type") != "sendreceive": if folder.get("type") != "sendreceive":
raise RouteSetupError("existing route folder is not sendreceive") raise RouteSetupError("existing route folder is not sendreceive")
@@ -391,6 +429,27 @@ class SyncthingRouteManager:
if folder.get("paused") is True: if folder.get("paused") is True:
raise RouteSetupError("existing route folder is paused") 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: def _wait(self, deadline: float) -> None:
if time.monotonic() >= deadline: if time.monotonic() >= deadline:
raise RouteSetupTimeout("route setup timed out") raise RouteSetupTimeout("route setup timed out")
@@ -441,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"):
+30
View File
@@ -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,
+25 -3
View File
@@ -1,7 +1,8 @@
import os import os
import tempfile import tempfile
import unittest import unittest
from pathlib import Path from pathlib import Path, PurePosixPath
from unittest.mock import patch
from archive_clients.config import ClientConfig, ConfigError, RootMapping from archive_clients.config import ClientConfig, ConfigError, RootMapping
@@ -18,16 +19,22 @@ class ConfigTests(unittest.TestCase):
(root / "sync").mkdir() (root / "sync").mkdir()
config_path = root / "client.toml" config_path = root / "client.toml"
config_path.write_text( config_path.write_text(
_config(root, role="cache"), encoding="utf-8" _config(root, role="cache").replace(
'username = "admin"', 'username = "${QB_USER}"'
),
encoding="utf-8",
) )
config = ClientConfig.load(config_path, "archive") with patch.dict(os.environ, {"QB_USER": "admin"}):
config = ClientConfig.load(config_path, "archive")
self.assertEqual(config.role, "archive") self.assertEqual(config.role, "archive")
self.assertEqual(config.qbittorrent.username, "admin")
self.assertEqual(config.read_shared_token(), "token") self.assertEqual(config.read_shared_token(), "token")
self.assertEqual( self.assertEqual(
config.qbittorrent.roots.api_to_local("/downloads/a/b"), config.qbittorrent.roots.api_to_local("/downloads/a/b"),
root / "qb/a/b", root / "qb/a/b",
) )
self.assertEqual(config.jobs.stall_after, 30 * 60) self.assertEqual(config.jobs.stall_after, 30 * 60)
self.assertEqual(config.jobs.free_space_reserve_bytes, 32 * 1024 * 1024)
self.assertEqual( self.assertEqual(
config.syncthing.advertised_addresses, ("dynamic",) config.syncthing.advertised_addresses, ("dynamic",)
) )
@@ -56,6 +63,21 @@ class ConfigTests(unittest.TestCase):
with self.assertRaises(ConfigError): with self.assertRaises(ConfigError):
mapping.api_to_local("/elsewhere/file") 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): def test_endpoint_scheme_and_job_keys_are_strict(self):
with tempfile.TemporaryDirectory() as directory: with tempfile.TemporaryDirectory() as directory:
root = Path(directory) root = Path(directory)
+88
View File
@@ -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)
+248 -5
View File
@@ -1,5 +1,6 @@
import hashlib import hashlib
import tempfile import tempfile
import threading
import unittest import unittest
from pathlib import Path from pathlib import Path
from unittest.mock import Mock, patch from unittest.mock import Mock, patch
@@ -11,6 +12,7 @@ from archive_clients.jobs import (
JobExecutionError, JobExecutionError,
_resource_fingerprint, _resource_fingerprint,
) )
from archive_clients.syncthing import RouteSetupError
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,88 @@ 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_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)
@@ -69,14 +152,166 @@ class ClientJobHappyPathTests(unittest.TestCase):
with self.subTest(operation=operation), tempfile.TemporaryDirectory() as directory: with self.subTest(operation=operation), tempfile.TemporaryDirectory() as directory:
self._run_transfer(Path(directory), operation) 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" source_root = root / "source"
target_root = root / "target" target_root = root / "target"
route_root = root / "route" route_root = root / "route"
source_root.mkdir() source_root.mkdir()
target_root.mkdir() target_root.mkdir()
route_root.mkdir() route_root.mkdir()
content = b"archive-control-happy-path"
(source_root / "fixture.bin").write_bytes(content) (source_root / "fixture.bin").write_bytes(content)
info = { info = {
b"length": len(content), b"length": len(content),
@@ -137,7 +372,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,
@@ -185,13 +420,21 @@ class ClientJobHappyPathTests(unittest.TestCase):
) )
final = None final = None
for executor, step in pipeline: for executor, step in pipeline:
events = executor.execute(control_pb2.ExecuteStepCommand( command = control_pb2.ExecuteStepCommand(
job_id=definition.job_id, job_id=definition.job_id,
expected_job_revision=cursor_revision, expected_job_revision=cursor_revision,
expected_last_event_sequence=cursor_sequence, expected_last_event_sequence=cursor_sequence,
step=step, step=step,
attempt=1, 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( self.assertEqual(
[event.sequence for event in events], [event.sequence for event in events],
list(range( list(range(
+26
View File
@@ -62,6 +62,32 @@ class RouteDiscoveryTests(unittest.TestCase):
self.assertEqual(routes[0].local_relative_path, "DownloadsSync") self.assertEqual(routes[0].local_relative_path, "DownloadsSync")
self.assertEqual(routes[0].state, route_pb2.ROUTE_STATE_DISCOVERED) 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__": if __name__ == "__main__":
unittest.main() unittest.main()
+3 -3
View File
@@ -43,11 +43,11 @@ class _Opener:
class ServiceProbeTests(unittest.TestCase): class ServiceProbeTests(unittest.TestCase):
def test_supported_qbittorrent_versions_include_live_4_4_api(self): def test_supported_qbittorrent_versions_require_export_capability(self):
for version in ("v4.4.5", "v4.5.5", "v4.6.7", "v5.2.3"): for version in ("v4.5.0", "v4.6.7", "v5.2.3"):
with self.subTest(version=version): with self.subTest(version=version):
self.assertTrue(_supported_qb_version(version)) self.assertTrue(_supported_qb_version(version))
for version in ("v4.3.9", "v6.0.0", "invalid"): for version in ("v4.3.9", "v4.4.5", "v6.0.0", "invalid"):
with self.subTest(version=version): with self.subTest(version=version):
self.assertFalse(_supported_qb_version(version)) self.assertFalse(_supported_qb_version(version))
+18
View File
@@ -125,6 +125,24 @@ class SyncthingRouteManagerTests(unittest.TestCase):
self.manager.configure(self.spec, time.monotonic() + 1) self.manager.configure(self.spec, time.monotonic() + 1)
self.assertEqual(self.transport.puts, []) 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): def test_bidirectional_nonce_and_ack_are_required(self):
configured = self.manager.configure(self.spec, time.monotonic() + 1) configured = self.manager.configure(self.spec, time.monotonic() + 1)
peer_nonce = configured.local_path / ( peer_nonce = configured.local_path / (
+28
View File
@@ -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)