Compare commits
7
Commits
62a0f24437
...
v0.1.24
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
87cca59d3a | ||
|
|
618bb3d206 | ||
|
|
a1bf3e4315 | ||
|
|
ae7175719b | ||
|
|
1c8128fe55 | ||
|
|
cd3a9b3c0d | ||
|
|
0311f6f084 |
@@ -120,6 +120,8 @@ username = "${QB_USER}"
|
||||
password_file = "/run/secrets/qb_password"
|
||||
api_root = "/downloads"
|
||||
local_root = "/data/qb"
|
||||
# Optional only when a nested qB path uses a different client-visible mount.
|
||||
# local_path_overrides = { "/downloads/fast" = "/data/qb-fast" }
|
||||
|
||||
[syncthing]
|
||||
endpoint = "http://syncthing:8384"
|
||||
@@ -139,6 +141,14 @@ 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.
|
||||
|
||||
qBittorrent content is resolved differently: every existing torrent keeps its
|
||||
qB-reported `save_path`. The client maps that path beneath `qbittorrent.api_root`
|
||||
to its own mount before a source, existing-target, permission, or eviction
|
||||
operation. Thus `/downloads/Downloading` naturally maps below `/data/qb`;
|
||||
there is no migration or per-resource configuration. Add a qB
|
||||
`local_path_overrides` entry only when that nested API prefix is a separate
|
||||
client mount.
|
||||
|
||||
The remaining node examples omit optional `[connection]`, `[jobs]`, and
|
||||
`[backup]` tables and therefore use these same defaults; deployments may
|
||||
override them per node.
|
||||
@@ -263,24 +273,25 @@ 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:
|
||||
deployment dependency and can be removed after publication. Use the checked-in
|
||||
publisher: it installs only the non-native binfmt handler, creates a temporary
|
||||
rootless BuildKit builder, verifies both platforms, publishes the index,
|
||||
displays its digest, then removes both temporary resources.
|
||||
|
||||
```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
|
||||
scripts/publish-image.sh vX.Y.Z
|
||||
scripts/publish-image.sh vX.Y.Z --also-latest
|
||||
```
|
||||
|
||||
`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.
|
||||
The publisher deliberately uses `moby/buildkit:rootless` with
|
||||
`--oci-worker-no-process-sandbox`. On nested Docker hosts, the default OCI
|
||||
sandbox can fail while masking `/proc/acpi` for an emulated build; rootless
|
||||
BuildKit confines that compatibility setting to the disposable builder. It
|
||||
waits briefly for the new worker to observe binfmt, then refuses to publish
|
||||
unless `docker buildx inspect` reports both `linux/amd64` and `linux/arm64`.
|
||||
On capability failure it prints that inspection output and removes the builder
|
||||
and binfmt handler on success, failure, or interruption. Retain the displayed
|
||||
manifest digest in release notes and deploy the immutable tag or digest.
|
||||
|
||||
## Operator usage
|
||||
|
||||
|
||||
@@ -86,6 +86,9 @@ python3 scripts/preflight-deployment.py ... \
|
||||
- That future route path and qBittorrent content root use one client bind
|
||||
mount, so hard-link staging remains possible rather than silently falling
|
||||
back to a space-consuming copy.
|
||||
- Every explicit `qbittorrent.local_path_overrides` entry maps the same host
|
||||
path in the qBittorrent and client containers. This protects nested qB save
|
||||
paths that use a dedicated bind mount.
|
||||
- A real `link(2)` operation between a unique zero-byte file in the qB root
|
||||
and one in the future automatic-route root. The probe verifies that both
|
||||
names refer to the same inode and removes them unconditionally.
|
||||
@@ -96,7 +99,8 @@ python3 scripts/preflight-deployment.py ... \
|
||||
- The token, qB password, and Syncthing API-key files are non-empty regular
|
||||
files with no group/world permissions.
|
||||
- The client image can read its configuration and reports usable permissions,
|
||||
sparse-file support, and filesystem capabilities for both roots.
|
||||
sparse-file support, and filesystem capabilities for both primary roots and
|
||||
every configured qBittorrent local-path override.
|
||||
- qBittorrent authentication/version compatibility and Syncthing
|
||||
authentication/device identity are healthy from the client container.
|
||||
|
||||
@@ -106,6 +110,15 @@ root and add a `local_path_overrides` mapping for that exact API path. This
|
||||
prevents automatic route folders being created on a small configuration
|
||||
filesystem while the client expects to hardlink from the qB data mount.
|
||||
|
||||
qBittorrent resources may use any `save_path` below `qbittorrent.api_root`.
|
||||
At job preflight the client maps that qB API path to its local mount and uses
|
||||
it as the resource root; data is never moved to fit Archive Control. A resource
|
||||
outside the configured qB API root, or whose resolved directory is unavailable
|
||||
or not a real directory in the client container, fails that job before staging
|
||||
or eviction. Use `qbittorrent.local_path_overrides` only for a nested qB API
|
||||
prefix backed by a distinct client mount; this preflight verifies the Docker
|
||||
bind topology for each such override.
|
||||
|
||||
When converting an existing node, stop its client and Syncthing containers,
|
||||
move each existing `routes/<route-id>` directory from the old Syncthing config
|
||||
tree into the new qB-backed route-root directory, then recreate Syncthing and
|
||||
|
||||
+15
-12
@@ -45,8 +45,8 @@ seconds. A connection stable for 60 seconds resets the backoff.
|
||||
|
||||
## Version and feature negotiation
|
||||
|
||||
The envelope and registration both state the sender version. Major `1` accepts
|
||||
only major `1`; a different major is rejected. The negotiated minor is the
|
||||
The envelope and registration both state the sender version. Major `2` accepts
|
||||
only major `2`; a different major is rejected. The negotiated minor is the
|
||||
highest mutually supported minor no greater than either endpoint's advertised
|
||||
minor.
|
||||
|
||||
@@ -96,19 +96,22 @@ client state snapshot before a route can become ready.
|
||||
|
||||
Every mutating command includes an expected job revision or immutable job
|
||||
definition plus the expected last global event sequence. A stale revision or
|
||||
sequence is rejected without side effects. The control daemon sends only one
|
||||
active step command for a job and waits for its persisted completion event
|
||||
before commanding the next participant. This single-writer lease lets whichever
|
||||
client owns the current step allocate the next global per-job event sequence;
|
||||
the following command starts from the sequence control has durably accepted.
|
||||
sequence is rejected without side effects. Assignment is represented solely by
|
||||
the durable assignment `CommandAck`; it produces no `JobEvent`. The control
|
||||
daemon sends only one active step command for a job and waits for its persisted
|
||||
completion event before commanding the next participant.
|
||||
|
||||
## Events, progress, and reconciliation
|
||||
|
||||
Each client-originated job event has a unique ID, monotonically increasing
|
||||
global per-job `sequence`, and resulting job revision. Duplicate IDs/sequences
|
||||
are idempotent only when their complete content matches. Control never grants
|
||||
concurrent event-writer leases for one job. A gap or conflicting duplicate
|
||||
pauses destructive orchestration and requests snapshots.
|
||||
global per-job `sequence`, resulting job revision, and the `command_id` of its
|
||||
owning `ExecuteStepCommand` or `CancelJobCommand`. Duplicate IDs/sequences are
|
||||
idempotent only when their complete content matches. Control verifies the
|
||||
client, acknowledged command lease, expected cursor, and active step before
|
||||
accepting it. A gap, conflict, or stale lease triggers a durable authoritative
|
||||
`ReconcileJobCommand`: the client retires only the named stale command leases
|
||||
and restores that job cursor, without deleting resource data or unrelated job
|
||||
state.
|
||||
|
||||
`fraction_complete` is current-step progress and
|
||||
`overall_fraction_complete` is the weighted five- or three-step job progress;
|
||||
@@ -120,7 +123,7 @@ On registration, active-job cursors provide the client's revision, last event
|
||||
sequence, state, and commit flag. Reconciliation applies these rules:
|
||||
|
||||
1. Equal cursors resume normal delivery.
|
||||
2. A client behind receives safe replay/snapshot commands.
|
||||
2. A client behind receives an authoritative reconciliation command.
|
||||
3. Control behind requests and validates the client's full job snapshot.
|
||||
4. Conflicting commit evidence reserves the resource and requires manual
|
||||
reconciliation; neither side performs cleanup.
|
||||
|
||||
+19
-2
@@ -42,8 +42,8 @@ and exported-metainfo data, keeps the qB-local hash separately, and calculates:
|
||||
- normalized selected and selected-complete file-index sets;
|
||||
- torrent runtime state;
|
||||
- canonical path flag and content revision;
|
||||
- a save-path fingerprint based on validated path mapping, never a leaked
|
||||
absolute host path.
|
||||
- a validated qBittorrent `save_path`, retained only in the client's local
|
||||
normalized observation.
|
||||
|
||||
qBittorrent 5 may expose a pure-v2 or hybrid torrent under the first 20 bytes
|
||||
of its v2 hash while separately advertising full `infohash_v1` and
|
||||
@@ -56,6 +56,23 @@ All entries retain qBittorrent's stable torrent file index. Renamed/noncanonical
|
||||
content paths are rejected in v1 because they cannot be transported and merged
|
||||
without ambiguity.
|
||||
|
||||
### Per-torrent local content roots
|
||||
|
||||
`save_path` is not inventory, placement, or protocol data. It is qBittorrent
|
||||
metadata used only by the daemon that queried qBittorrent. Before staging a
|
||||
source, merging into an existing target, applying post-recheck permissions, or
|
||||
evicting files, that daemon maps the torrent's API-visible `save_path` through
|
||||
`qbittorrent.api_root`/`local_root`. The most-specific
|
||||
`qbittorrent.local_path_overrides` mapping wins when a nested path is exposed
|
||||
through a distinct client container mount.
|
||||
|
||||
An ordinary nested qBittorrent path such as `/media/Data/Downloading` needs no
|
||||
per-resource configuration: it resolves beneath the configured root. A path
|
||||
outside that root, an unmapped distinct mount, or a mapped local path that is
|
||||
not a visible real directory fails the affected job before filesystem mutation.
|
||||
This check is intentionally per resource; an unrelated malformed qBittorrent
|
||||
entry cannot prevent normal resources from being staged or transferred.
|
||||
|
||||
### Verification guard
|
||||
|
||||
The client captures transfer counters and state before recheck, issues recheck,
|
||||
|
||||
@@ -14,6 +14,13 @@ client:
|
||||
a regular file or an explicitly created directory;
|
||||
5. verifies every operation remains beneath the local configured root.
|
||||
|
||||
For qBittorrent content, the root is resolved per resource: qB's authoritative
|
||||
API-visible `save_path` is mapped beneath `qbittorrent.api_root` (or a more
|
||||
specific configured local override) into the client namespace. This permits
|
||||
existing nested save paths without moving data, but does not permit paths
|
||||
outside the configured boundary. The resolved path remains local client state
|
||||
and is never sent to the control daemon.
|
||||
|
||||
Sockets, devices, FIFOs, symlinks, and other special entries fail preflight.
|
||||
Permission or ownership mismatch is fail-fast. Archive Control never changes
|
||||
source ownership or mode to make a job pass.
|
||||
@@ -150,4 +157,3 @@ Offline tooling provides list, verify, and restore. Restore requires stopped
|
||||
daemon access, verifies the chosen backup, preserves the suspect database under
|
||||
a timestamped name, installs the replacement atomically, and runs integrity and
|
||||
schema checks before normal startup. Backups contain no configured secrets.
|
||||
|
||||
|
||||
+4
-1
@@ -125,7 +125,10 @@ and placement policy differ.
|
||||
and explicit save path, or apply a transactional delta to an existing
|
||||
torrent. Set selected/skipped state, run a full recheck of the target union,
|
||||
and fail immediately if qBittorrent attempts to download. Successful recheck
|
||||
atomically advances the placement generation and is the commit point.
|
||||
atomically advances the placement generation and is the commit point. Before
|
||||
qBittorrent is resumed, the client applies `a+rx` to the verified resource
|
||||
directories and `a+r` to its verified files, preserving ownership, write
|
||||
bits, and special mode bits so other local applications can read the data.
|
||||
5. **Staging Cleanup.** Remove only job-owned staging artifacts from both
|
||||
endpoints. A post-commit cleanup failure produces `CLEANUP_REQUIRED`; it
|
||||
never rolls back or deletes the committed placement.
|
||||
|
||||
+1
-1
@@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta"
|
||||
|
||||
[project]
|
||||
name = "archive-clients"
|
||||
version = "0.1.19"
|
||||
version = "0.1.23"
|
||||
requires-python = ">=3.11"
|
||||
dependencies = ["protobuf==7.35.1", "websockets==16.0"]
|
||||
|
||||
|
||||
@@ -134,7 +134,7 @@ from pathlib import Path
|
||||
from archive_clients.config import ClientConfig
|
||||
config=ClientConfig.load(Path(sys.argv[1]))
|
||||
value={"shared_token_file":str(config.shared_token_file)}
|
||||
value["qbittorrent"]={"api_root":str(config.qbittorrent.api_root),"local_root":str(config.qbittorrent.local_root),"password_file":str(config.qbittorrent.password_file)}
|
||||
value["qbittorrent"]={"api_root":str(config.qbittorrent.api_root),"local_root":str(config.qbittorrent.local_root),"password_file":str(config.qbittorrent.password_file),"local_path_overrides":{str(api):str(local) for api,local in config.qbittorrent.local_path_overrides}}
|
||||
value["syncthing"]={"api_root":str(config.syncthing.api_root),"local_root":str(config.syncthing.local_root),"api_key_file":str(config.syncthing.api_key_file),"local_path_overrides":{str(api):str(local) for api,local in config.syncthing.local_path_overrides}}
|
||||
print(json.dumps(value,sort_keys=True))'''
|
||||
result = subprocess.run(
|
||||
@@ -174,6 +174,26 @@ raise SystemExit(0 if all(p.state == 1 for p in probes) else 1)'''
|
||||
print(result.stdout.strip())
|
||||
|
||||
|
||||
def require_qb_override_mappings(
|
||||
client: dict[str, Any], qbittorrent: dict[str, Any], qb: dict[str, Any],
|
||||
) -> None:
|
||||
"""Prove every explicit qB API/local override sees the same host bytes."""
|
||||
|
||||
overrides = qb.get("local_path_overrides", {})
|
||||
if not isinstance(overrides, dict):
|
||||
raise CheckFailure("qbittorrent.local_path_overrides must be a table")
|
||||
for api_path, local_path in overrides.items():
|
||||
if not isinstance(api_path, str) or not isinstance(local_path, str):
|
||||
raise CheckFailure(
|
||||
"qbittorrent.local_path_overrides entries are invalid"
|
||||
)
|
||||
require_same_path(
|
||||
f"qBittorrent override {api_path}/local mapping",
|
||||
map_path(qbittorrent, api_path),
|
||||
map_path(client, local_path),
|
||||
)
|
||||
|
||||
|
||||
def run_hardlink_probe(
|
||||
container: str, qb_root: str, route_root: str, user: str | None = None,
|
||||
) -> None:
|
||||
@@ -298,6 +318,7 @@ def main(argv: list[str] | None = None) -> int:
|
||||
syncthing_route, client_route,
|
||||
)
|
||||
require_same_path("qBittorrent api_root/local_root", qb_api, client_qb)
|
||||
require_qb_override_mappings(client, qbittorrent, qb)
|
||||
if client_qb.destination != client_route.destination:
|
||||
raise CheckFailure(
|
||||
"qBittorrent and future route roots use separate client bind "
|
||||
|
||||
Executable
+105
@@ -0,0 +1,105 @@
|
||||
#!/usr/bin/env bash
|
||||
# Publish one complete Archive Control client OCI image index.
|
||||
set -euo pipefail
|
||||
|
||||
usage() {
|
||||
cat <<'EOF'
|
||||
Usage: scripts/publish-image.sh <semver> [--also-latest]
|
||||
|
||||
Builds and pushes sodium/archive-clients for linux/amd64 and linux/arm64.
|
||||
Docker must already be authenticated to the target registry.
|
||||
EOF
|
||||
}
|
||||
|
||||
if [[ $# -lt 1 || ${1:-} == '-h' || ${1:-} == '--help' ]]; then
|
||||
usage
|
||||
exit 2
|
||||
fi
|
||||
|
||||
tag=$1
|
||||
shift
|
||||
also_latest=false
|
||||
if [[ ${1:-} == '--also-latest' ]]; then
|
||||
also_latest=true
|
||||
shift
|
||||
fi
|
||||
if [[ $# -ne 0 ]]; then
|
||||
usage
|
||||
exit 2
|
||||
fi
|
||||
|
||||
repository=${IMAGE_REPOSITORY:-sodium/archive-clients}
|
||||
builder=${BUILDER:-archive-control-release}
|
||||
buildkit_image=${BUILDKIT_IMAGE:-moby/buildkit:rootless}
|
||||
# Nested Docker hosts can reject default OCI /proc mount masking. Rootless
|
||||
# BuildKit confines this compatibility flag to the disposable release builder.
|
||||
buildkitd_flags=${BUILDKITD_FLAGS:---oci-worker-no-process-sandbox}
|
||||
binfmt_image=${BINFMT_IMAGE:-tonistiigi/binfmt}
|
||||
host_arch=$(docker version --format '{{.Server.Arch}}')
|
||||
case ${host_arch} in
|
||||
arm64|aarch64) emulated_arch=amd64 ;;
|
||||
amd64|x86_64) emulated_arch=arm64 ;;
|
||||
*) echo "Unsupported Docker server architecture: ${host_arch}" >&2; exit 1 ;;
|
||||
esac
|
||||
|
||||
builder_created=false
|
||||
binfmt_installed=false
|
||||
cleanup() {
|
||||
local status=$?
|
||||
trap - EXIT
|
||||
if [[ ${builder_created} == true ]]; then
|
||||
docker buildx rm "${builder}" >/dev/null 2>&1 || true
|
||||
fi
|
||||
if [[ ${binfmt_installed} == true ]]; then
|
||||
docker run --privileged --rm "${binfmt_image}" --uninstall "${emulated_arch}" \
|
||||
>/dev/null 2>&1 || true
|
||||
fi
|
||||
exit "${status}"
|
||||
}
|
||||
trap cleanup EXIT
|
||||
|
||||
docker run --privileged --rm "${binfmt_image}" --install "${emulated_arch}" \
|
||||
>/dev/null
|
||||
binfmt_installed=true
|
||||
if docker buildx inspect "${builder}" >/dev/null 2>&1; then
|
||||
docker buildx rm "${builder}" >/dev/null
|
||||
fi
|
||||
docker buildx create --name "${builder}" --driver docker-container \
|
||||
--driver-opt "image=${buildkit_image}" \
|
||||
--buildkitd-flags "${buildkitd_flags}" --use >/dev/null
|
||||
builder_created=true
|
||||
|
||||
# A newly-created rootless worker can publish its native platform before it has
|
||||
# observed the just-registered binfmt handler. Do not mistake that brief
|
||||
# startup state for a partial-release-capable builder.
|
||||
platforms=''
|
||||
supports_all=false
|
||||
for attempt in {1..10}; do
|
||||
platforms=$(docker buildx inspect "${builder}" --bootstrap 2>&1)
|
||||
supports_all=true
|
||||
for platform in linux/amd64 linux/arm64; do
|
||||
if ! grep -Fq "${platform}" <<<"${platforms}"; then
|
||||
supports_all=false
|
||||
break
|
||||
fi
|
||||
done
|
||||
if [[ ${supports_all} == true ]]; then
|
||||
break
|
||||
fi
|
||||
if [[ ${attempt} -lt 10 ]]; then
|
||||
sleep 1
|
||||
fi
|
||||
done
|
||||
if [[ ${supports_all} != true ]]; then
|
||||
echo "Builder ${builder} does not support both required platforms; refusing partial release." >&2
|
||||
printf '%s\n' "${platforms}" >&2
|
||||
exit 1
|
||||
fi
|
||||
|
||||
tags=(--tag "${repository}:${tag}")
|
||||
if [[ ${also_latest} == true ]]; then
|
||||
tags+=(--tag "${repository}:latest")
|
||||
fi
|
||||
docker buildx build --pull --platform linux/amd64,linux/arm64 \
|
||||
--push "${tags[@]}" .
|
||||
docker buildx imagetools inspect "${repository}:${tag}"
|
||||
@@ -1,4 +1,4 @@
|
||||
"""Archive Control data-node daemon."""
|
||||
|
||||
PROTO_COMMIT = "4ec852014dad74606d4078b3ae1aa208c814b033"
|
||||
PROTO_COMMIT = "03b6751d420e8e4e4cc9aeef4f76b13e28a7265f"
|
||||
|
||||
|
||||
@@ -27,10 +27,16 @@ def main(argv: Sequence[str] | None = None) -> int:
|
||||
parser.add_argument("--check-config", action="store_true")
|
||||
arguments = parser.parse_args(argv)
|
||||
config = ClientConfig.load(arguments.config, arguments.mode)
|
||||
qb_extra_roots = tuple(sorted({
|
||||
root
|
||||
for _, root in config.qbittorrent.local_path_overrides
|
||||
if root != config.qbittorrent.local_root
|
||||
}))
|
||||
probes = [
|
||||
probe_root(config.qbittorrent.local_root),
|
||||
probe_root(config.syncthing.local_root),
|
||||
]
|
||||
qb_extra_probes = tuple(probe_root(root) for root in qb_extra_roots)
|
||||
probe_writable_directory(config.state_db.parent)
|
||||
probe_writable_directory(config.backup_dir)
|
||||
shared_token = config.read_shared_token()
|
||||
@@ -40,7 +46,10 @@ def main(argv: Sequence[str] | None = None) -> int:
|
||||
print(json.dumps({
|
||||
"client_id": config.client_id,
|
||||
"role": config.role,
|
||||
"filesystems": [probe.__dict__ | {"root": str(probe.root)} for probe in probes],
|
||||
"filesystems": [
|
||||
probe.__dict__ | {"root": str(probe.root)}
|
||||
for probe in (*probes[:1], *qb_extra_probes, *probes[1:])
|
||||
],
|
||||
}, sort_keys=True))
|
||||
return 0
|
||||
configure_logging(
|
||||
@@ -68,7 +77,7 @@ def main(argv: Sequence[str] | None = None) -> int:
|
||||
},
|
||||
)
|
||||
asyncio.run(ArchiveClientDaemon(
|
||||
config, probes, service_probes,
|
||||
config, probes, service_probes, qb_extra_probes=qb_extra_probes,
|
||||
resource_reader=QBittorrentReader(config.qbittorrent),
|
||||
).run())
|
||||
except KeyboardInterrupt:
|
||||
|
||||
@@ -213,8 +213,6 @@ 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]] = []
|
||||
|
||||
@@ -64,6 +64,7 @@ class ArchiveClientDaemon:
|
||||
service_probes: list[ServiceProbe],
|
||||
resource_reader: QBittorrentReader | None = None,
|
||||
route_manager: SyncthingRouteManager | None = None,
|
||||
qb_extra_probes: tuple[FilesystemProbe, ...] = (),
|
||||
):
|
||||
if len(probes) != 2:
|
||||
raise ValueError(
|
||||
@@ -71,6 +72,7 @@ class ArchiveClientDaemon:
|
||||
)
|
||||
self.config = config
|
||||
self.probes = probes
|
||||
self.qb_probes = (probes[0], *qb_extra_probes)
|
||||
self.service_probes = service_probes
|
||||
self.inventory = (
|
||||
InventoryService(resource_reader, config.client_id)
|
||||
@@ -119,9 +121,13 @@ class ArchiveClientDaemon:
|
||||
store=self.store,
|
||||
qb_root=config.qbittorrent.local_root,
|
||||
qb_api_root=config.qbittorrent.api_root,
|
||||
qb_roots=config.qbittorrent.roots,
|
||||
route_path=self._route_path,
|
||||
syncthing_transport=self.routes.transport,
|
||||
sparse_supported=all(probe.sparse_files for probe in probes),
|
||||
sparse_supported=all(
|
||||
probe.sparse_files
|
||||
for probe in (*self.qb_probes, probes[1])
|
||||
),
|
||||
poll_interval=config.jobs.poll_interval,
|
||||
verification_timeout=config.jobs.verification_timeout,
|
||||
free_space_reserve_bytes=(
|
||||
@@ -214,7 +220,7 @@ class ArchiveClientDaemon:
|
||||
raise RuntimeError("registration response correlation mismatch")
|
||||
if response.register_response.status != client_pb2.REGISTRATION_STATUS_ACCEPTED:
|
||||
raise RuntimeError("control rejected registration")
|
||||
if response.register_response.negotiated_version.major != 1:
|
||||
if response.register_response.negotiated_version.major != 2:
|
||||
raise RuntimeError("control negotiated an unsupported protocol version")
|
||||
logger.info("control_connection_registered")
|
||||
outbound: asyncio.Queue[str] = asyncio.Queue(maxsize=100)
|
||||
@@ -427,12 +433,36 @@ class ArchiveClientDaemon:
|
||||
accepted = None
|
||||
accepted_for_execution = False
|
||||
try:
|
||||
accepted = await asyncio.to_thread(
|
||||
self.store.accept_command,
|
||||
command.command_id,
|
||||
encode_message(command),
|
||||
encode_message(acknowledgement),
|
||||
)
|
||||
if command.WhichOneof("payload") == "reconcile_job":
|
||||
reconciliation = command.reconcile_job
|
||||
accepted = await asyncio.to_thread(
|
||||
self.store.accept_reconcile_command,
|
||||
command.command_id,
|
||||
encode_message(command),
|
||||
encode_message(acknowledgement),
|
||||
job_id=reconciliation.authoritative_job.definition.job_id,
|
||||
definition_json=encode_message(
|
||||
reconciliation.authoritative_job.definition
|
||||
),
|
||||
state=job_pb2.JobState.Name(
|
||||
reconciliation.authoritative_job.state
|
||||
),
|
||||
revision=reconciliation.authoritative_job.revision,
|
||||
last_event_sequence=(
|
||||
reconciliation.authoritative_last_event_sequence
|
||||
),
|
||||
committed=reconciliation.authoritative_job.committed,
|
||||
superseded_command_ids=list(
|
||||
reconciliation.superseded_command_ids
|
||||
),
|
||||
)
|
||||
else:
|
||||
accepted = await asyncio.to_thread(
|
||||
self.store.accept_command,
|
||||
command.command_id,
|
||||
encode_message(command),
|
||||
encode_message(acknowledgement),
|
||||
)
|
||||
if accepted.duplicate:
|
||||
acknowledgement = decode_message(
|
||||
accepted.acknowledgement_json, control_pb2.CommandAck()
|
||||
@@ -440,6 +470,7 @@ class ArchiveClientDaemon:
|
||||
if (
|
||||
acknowledgement.status
|
||||
== control_pb2.COMMAND_ACK_STATUS_ACCEPTED
|
||||
and accepted.state == "accepted"
|
||||
):
|
||||
accepted_for_execution = True
|
||||
acknowledgement.status = (
|
||||
@@ -531,6 +562,7 @@ class ArchiveClientDaemon:
|
||||
outbound: asyncio.Queue[str],
|
||||
command_tasks: set[asyncio.Task[None]],
|
||||
) -> None:
|
||||
latest_route_commands: dict[str, control_pb2.Command] = {}
|
||||
for row in await asyncio.to_thread(self.store.list_accepted_commands):
|
||||
acknowledgement = decode_message(
|
||||
str(row["acknowledgement_json"]), control_pb2.CommandAck()
|
||||
@@ -541,15 +573,21 @@ class ArchiveClientDaemon:
|
||||
str(row["command_json"]), control_pb2.Command()
|
||||
)
|
||||
if command.WhichOneof("payload") == "ensure_route":
|
||||
self._schedule_route_command(
|
||||
command, "", outbound, command_tasks, restore_ready=True
|
||||
)
|
||||
# A newer ensure command for the same route is sufficient to
|
||||
# restore the process-local route path and replay its own
|
||||
# updates. Replaying every historical ready command causes
|
||||
# redundant Syncthing configuration after a restart.
|
||||
latest_route_commands[command.ensure_route.route.route_id] = command
|
||||
# Job commands are intentionally not replayed here. A command
|
||||
# acknowledgement is durable on both sides; the control daemon
|
||||
# redelivers an unacknowledged command, while registration
|
||||
# reconciliation requests snapshots for active jobs. Replaying
|
||||
# every historical command re-emits overlapping event ranges and
|
||||
# can flood the single connection's bounded outbound queue.
|
||||
for command in latest_route_commands.values():
|
||||
self._schedule_route_command(
|
||||
command, "", outbound, command_tasks, restore_ready=True
|
||||
)
|
||||
|
||||
async def _resume_route_commands(
|
||||
self,
|
||||
@@ -644,6 +682,11 @@ class ArchiveClientDaemon:
|
||||
loop = asyncio.get_running_loop()
|
||||
|
||||
def emit(event):
|
||||
if not self.store.is_command_active(command.command_id):
|
||||
# Reconciliation retired this lease while its blocking
|
||||
# worker was still unwinding. Its durable local record is
|
||||
# not allowed to re-enter the control event stream.
|
||||
return
|
||||
response = new_envelope()
|
||||
response.correlation_id = correlation_id
|
||||
response.job_event.CopyFrom(event)
|
||||
@@ -661,16 +704,20 @@ class ArchiveClientDaemon:
|
||||
pass
|
||||
|
||||
events = await asyncio.to_thread(
|
||||
self.jobs.execute, command.execute_step, emit
|
||||
self.jobs.execute, command.execute_step, emit, command.command_id
|
||||
)
|
||||
streamed = True
|
||||
elif payload == "cancel_job":
|
||||
events = await asyncio.to_thread(
|
||||
self.jobs.cancel, command.cancel_job
|
||||
self.jobs.cancel, command.cancel_job, command.command_id
|
||||
)
|
||||
else:
|
||||
raise JobExecutionError("job command payload is unsupported")
|
||||
for event in (() if streamed else events):
|
||||
if not await asyncio.to_thread(
|
||||
self.store.is_command_active, command.command_id
|
||||
):
|
||||
return
|
||||
response = new_envelope()
|
||||
response.correlation_id = correlation_id
|
||||
response.job_event.CopyFrom(event)
|
||||
@@ -1039,6 +1086,17 @@ class ArchiveClientDaemon:
|
||||
acknowledgement.error.message = "job assignment is invalid"
|
||||
else:
|
||||
acknowledgement.status = control_pb2.COMMAND_ACK_STATUS_ACCEPTED
|
||||
elif command.WhichOneof("payload") == "reconcile_job":
|
||||
record = command.reconcile_job.authoritative_job
|
||||
if (
|
||||
not record.definition.job_id
|
||||
or command.reconcile_job.authoritative_last_event_sequence < 0
|
||||
):
|
||||
acknowledgement.status = control_pb2.COMMAND_ACK_STATUS_REJECTED
|
||||
acknowledgement.error.code = common_pb2.ERROR_CODE_INVALID_ARGUMENT
|
||||
acknowledgement.error.message = "job reconciliation is invalid"
|
||||
else:
|
||||
acknowledgement.status = control_pb2.COMMAND_ACK_STATUS_ACCEPTED
|
||||
elif command.WhichOneof("payload") == "execute_step":
|
||||
step = command.execute_step
|
||||
if self.jobs is None:
|
||||
|
||||
@@ -7,7 +7,7 @@ import hashlib
|
||||
import os
|
||||
import stat
|
||||
from pathlib import Path, PurePosixPath
|
||||
from typing import Iterable
|
||||
from typing import Callable, Iterable
|
||||
|
||||
from archive_clients.qbittorrent import QBittorrentReader
|
||||
from archive_clients.resources import NormalizedResource
|
||||
@@ -23,7 +23,7 @@ def verify_and_snapshot(
|
||||
job_id: str,
|
||||
resource: NormalizedResource,
|
||||
selected_indices: Iterable[int],
|
||||
qb_root: Path,
|
||||
content_root: Path,
|
||||
store: ClientStore,
|
||||
) -> dict[str, object]:
|
||||
selected = set(selected_indices)
|
||||
@@ -38,7 +38,7 @@ def verify_and_snapshot(
|
||||
f"cache file {index} is not selected and complete"
|
||||
)
|
||||
relative = _relative(item.canonical_path)
|
||||
path = qb_root.joinpath(*relative.parts)
|
||||
path = content_root.joinpath(*relative.parts)
|
||||
try:
|
||||
metadata = path.lstat()
|
||||
except FileNotFoundError as exc:
|
||||
@@ -65,7 +65,7 @@ def verify_and_snapshot(
|
||||
"sha256": _sha256(path),
|
||||
}
|
||||
)
|
||||
snapshot = {"files": files}
|
||||
snapshot = {"content_root": str(content_root), "files": files}
|
||||
return store.put_job_artifact(job_id, "eviction-snapshot", snapshot)["value"]
|
||||
|
||||
|
||||
@@ -90,9 +90,9 @@ def remove_qb_entry(
|
||||
def safe_unlink(
|
||||
*,
|
||||
job_id: str,
|
||||
qb_root: Path,
|
||||
qbittorrent: QBittorrentReader,
|
||||
store: ClientStore,
|
||||
resource_root: Callable[[NormalizedResource], Path],
|
||||
) -> dict[str, object]:
|
||||
completed = store.get_job_artifact(job_id, "eviction-unlinked")
|
||||
if completed is not None:
|
||||
@@ -101,11 +101,19 @@ def safe_unlink(
|
||||
if snapshot_row is None:
|
||||
raise EvictionError("eviction snapshot is missing")
|
||||
snapshot = snapshot_row["value"]
|
||||
if not isinstance(snapshot, dict):
|
||||
raise EvictionError("eviction snapshot is invalid")
|
||||
content_root_value = snapshot.get("content_root")
|
||||
if not isinstance(content_root_value, str):
|
||||
raise EvictionError("eviction snapshot content root is missing")
|
||||
content_root = Path(content_root_value)
|
||||
if not content_root.is_absolute() or content_root.is_symlink():
|
||||
raise EvictionError("eviction snapshot content root is unsafe")
|
||||
files = snapshot.get("files")
|
||||
if not isinstance(files, list):
|
||||
raise EvictionError("eviction snapshot is invalid")
|
||||
|
||||
referenced = _remaining_paths(qbittorrent.list_resources())
|
||||
referenced = _remaining_paths(qbittorrent.list_resources(), resource_root)
|
||||
removed: list[str] = []
|
||||
retained: list[dict[str, str]] = []
|
||||
directories: set[Path] = set()
|
||||
@@ -113,10 +121,10 @@ def safe_unlink(
|
||||
if not isinstance(record, dict) or not isinstance(record.get("path"), str):
|
||||
raise EvictionError("eviction file record is invalid")
|
||||
relative = _relative(record["path"])
|
||||
if relative.as_posix() in referenced:
|
||||
path = content_root.joinpath(*relative.parts)
|
||||
if path in referenced:
|
||||
retained.append({"path": relative.as_posix(), "reason": "shared"})
|
||||
continue
|
||||
path = qb_root.joinpath(*relative.parts)
|
||||
try:
|
||||
metadata = path.lstat()
|
||||
except FileNotFoundError:
|
||||
@@ -142,7 +150,7 @@ def safe_unlink(
|
||||
path.unlink()
|
||||
removed.append(relative.as_posix())
|
||||
parent = path.parent
|
||||
while parent != qb_root:
|
||||
while parent != content_root:
|
||||
directories.add(parent)
|
||||
parent = parent.parent
|
||||
|
||||
@@ -153,7 +161,7 @@ def safe_unlink(
|
||||
try:
|
||||
directory.rmdir()
|
||||
removed_directories.append(
|
||||
directory.relative_to(qb_root).as_posix()
|
||||
directory.relative_to(content_root).as_posix()
|
||||
)
|
||||
except OSError as exc:
|
||||
if exc.errno not in {errno.ENOTEMPTY, errno.ENOENT}:
|
||||
@@ -214,12 +222,24 @@ def compensate_materialized_files(
|
||||
return removed
|
||||
|
||||
|
||||
def _remaining_paths(resources: Iterable[NormalizedResource]) -> set[str]:
|
||||
return {
|
||||
item.canonical_path
|
||||
for resource in resources
|
||||
for item in resource.files
|
||||
}
|
||||
def _remaining_paths(
|
||||
resources: Iterable[NormalizedResource],
|
||||
resource_root: Callable[[NormalizedResource], Path],
|
||||
) -> set[Path]:
|
||||
result: set[Path] = set()
|
||||
for resource in resources:
|
||||
try:
|
||||
root = resource_root(resource)
|
||||
except (OSError, RuntimeError, ValueError):
|
||||
# A malformed unrelated qB entry must not stop an otherwise
|
||||
# safe eviction. Its unresolvable path cannot be considered a
|
||||
# shared path under the verified eviction root.
|
||||
continue
|
||||
result.update(
|
||||
root.joinpath(*_relative(item.canonical_path).parts)
|
||||
for item in resource.files
|
||||
)
|
||||
return result
|
||||
|
||||
|
||||
def _relative(value: str) -> PurePosixPath:
|
||||
|
||||
+169
-33
@@ -14,6 +14,7 @@ import uuid
|
||||
from pathlib import Path, PurePosixPath
|
||||
from typing import Callable, Iterable
|
||||
|
||||
from archive_clients.config import ConfigError, RootMapping
|
||||
from archive_clients.protocol import decode_message, encode_message
|
||||
from archive_clients.eviction import (
|
||||
EvictionError,
|
||||
@@ -67,6 +68,7 @@ class ClientJobExecutor:
|
||||
store: ClientStore,
|
||||
qb_root: Path,
|
||||
qb_api_root: PurePosixPath,
|
||||
qb_roots: RootMapping | None = None,
|
||||
route_path: Callable[[str], Path],
|
||||
syncthing_transport: object,
|
||||
sparse_supported: bool,
|
||||
@@ -78,7 +80,10 @@ class ClientJobExecutor:
|
||||
self.qbittorrent = qbittorrent
|
||||
self.store = store
|
||||
self.qb_root = qb_root
|
||||
self.qb_api_root = qb_api_root
|
||||
self.qb_api_root = PurePosixPath(qb_api_root)
|
||||
self.qb_roots = qb_roots or RootMapping(
|
||||
self.qb_api_root, self.qb_root
|
||||
)
|
||||
self.route_path = route_path
|
||||
self.syncthing_transport = syncthing_transport
|
||||
self.sparse_supported = sparse_supported
|
||||
@@ -97,33 +102,25 @@ class ClientJobExecutor:
|
||||
) -> list[control_pb2.JobEvent]:
|
||||
definition = command.job
|
||||
self._validate_definition(definition)
|
||||
replay = self._replay(
|
||||
definition.job_id, command.expected_last_event_sequence
|
||||
self.store.ensure_job_definition(
|
||||
definition.job_id, encode_message(definition)
|
||||
)
|
||||
if replay:
|
||||
return replay
|
||||
event = self._event(
|
||||
definition,
|
||||
sequence=command.expected_last_event_sequence + 1,
|
||||
revision=command.expected_job_revision,
|
||||
event_type=control_pb2.JOB_EVENT_TYPE_ASSIGNED,
|
||||
state=job_pb2.JOB_STATE_PREPARING,
|
||||
committed=False,
|
||||
)
|
||||
self._record(definition, event)
|
||||
return [event]
|
||||
# Assignment succeeds through CommandAck. It must not let a client
|
||||
# allocate a globally ordered JobEvent cursor.
|
||||
return []
|
||||
|
||||
def execute(
|
||||
self,
|
||||
command: control_pb2.ExecuteStepCommand,
|
||||
event_callback: Callable[[control_pb2.JobEvent], None] | None = None,
|
||||
command_id: str = "",
|
||||
) -> 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)
|
||||
return self._execute_locked(command, event_callback, command_id)
|
||||
|
||||
def _execution_lock(self, job_id: str) -> threading.Lock:
|
||||
with self._execution_locks_guard:
|
||||
@@ -133,6 +130,7 @@ class ClientJobExecutor:
|
||||
self,
|
||||
command: control_pb2.ExecuteStepCommand,
|
||||
event_callback: Callable[[control_pb2.JobEvent], None] | None = None,
|
||||
command_id: str = "",
|
||||
) -> list[control_pb2.JobEvent]:
|
||||
definition = self._definition(command.job_id)
|
||||
replay = self._replay(
|
||||
@@ -186,6 +184,7 @@ class ClientJobExecutor:
|
||||
),
|
||||
step=command.step,
|
||||
step_state=job_pb2.STEP_STATE_RUNNING,
|
||||
command_id=command_id,
|
||||
)
|
||||
self._record(definition, started)
|
||||
cursor = started
|
||||
@@ -228,6 +227,7 @@ class ClientJobExecutor:
|
||||
committed=cursor.committed,
|
||||
step=command.step,
|
||||
step_state=job_pb2.STEP_STATE_RUNNING,
|
||||
command_id=command_id,
|
||||
)
|
||||
event.progress.fraction_complete = max(
|
||||
0.0, min(float(fraction), 1.0)
|
||||
@@ -258,6 +258,7 @@ class ClientJobExecutor:
|
||||
committed=started.committed,
|
||||
step=command.step,
|
||||
step_state=job_pb2.STEP_STATE_CANCELLED,
|
||||
command_id=command_id,
|
||||
)
|
||||
cancelling.error.code = common_pb2.ERROR_CODE_CANCELLED
|
||||
cancelling.error.message = str(error)
|
||||
@@ -315,6 +316,7 @@ class ClientJobExecutor:
|
||||
committed=started.committed,
|
||||
step=command.step,
|
||||
step_state=job_pb2.STEP_STATE_FAILED,
|
||||
command_id=command_id,
|
||||
)
|
||||
failed.error.code = _job_error_code(error)
|
||||
failed.error.message = str(error) or type(error).__name__
|
||||
@@ -357,6 +359,7 @@ class ClientJobExecutor:
|
||||
committed=committed,
|
||||
step=command.step,
|
||||
step_state=job_pb2.STEP_STATE_SUCCEEDED,
|
||||
command_id=command_id,
|
||||
)
|
||||
if result is not None:
|
||||
succeeded.observed_placement.CopyFrom(result)
|
||||
@@ -367,7 +370,7 @@ class ClientJobExecutor:
|
||||
return emitted
|
||||
|
||||
def cancel(
|
||||
self, command: control_pb2.CancelJobCommand
|
||||
self, command: control_pb2.CancelJobCommand, command_id: str = ""
|
||||
) -> list[control_pb2.JobEvent]:
|
||||
definition = self._definition(command.job_id)
|
||||
self.request_cancel(command.job_id)
|
||||
@@ -418,6 +421,7 @@ class ClientJobExecutor:
|
||||
committed=committed,
|
||||
step=job_pb2.JOB_STEP_KIND_ROLLBACK,
|
||||
step_state=job_pb2.STEP_STATE_SUCCEEDED,
|
||||
command_id=command_id,
|
||||
)
|
||||
if placement is not None:
|
||||
event.observed_placement.CopyFrom(placement)
|
||||
@@ -480,6 +484,7 @@ class ClientJobExecutor:
|
||||
"archive coverage no longer covers the cache selection"
|
||||
)
|
||||
resource = self._resource(definition)
|
||||
resource_root = self._resource_root(resource)
|
||||
current = _resource_fingerprint(resource, self.client_id)
|
||||
if not _fingerprint_matches(
|
||||
current, definition.eviction.cache_fingerprint
|
||||
@@ -491,7 +496,7 @@ class ClientJobExecutor:
|
||||
job_id=definition.job_id,
|
||||
resource=resource,
|
||||
selected_indices=requested,
|
||||
qb_root=self.qb_root,
|
||||
content_root=resource_root,
|
||||
store=self.store,
|
||||
)
|
||||
return None
|
||||
@@ -506,9 +511,9 @@ class ClientJobExecutor:
|
||||
if step == job_pb2.JOB_STEP_KIND_SAFE_FILE_UNLINK:
|
||||
safe_unlink(
|
||||
job_id=definition.job_id,
|
||||
qb_root=self.qb_root,
|
||||
qbittorrent=self.qbittorrent,
|
||||
store=self.store,
|
||||
resource_root=self._resource_root,
|
||||
)
|
||||
placement = resource_pb2.Placement(
|
||||
client_id=self.client_id,
|
||||
@@ -539,7 +544,8 @@ class ClientJobExecutor:
|
||||
raise JobExecutionError(
|
||||
"source resource changed after job confirmation"
|
||||
)
|
||||
self._reject_unsafe_partfile(definition, resource)
|
||||
source_root = self._resource_root(resource)
|
||||
self._reject_unsafe_partfile(definition, resource, source_root)
|
||||
indices = _selection_indices(definition.transfer.transfer_delta_files)
|
||||
by_index = {item.file_index: item for item in resource.files}
|
||||
if not indices or any(index not in by_index for index in indices):
|
||||
@@ -603,7 +609,7 @@ class ClientJobExecutor:
|
||||
self._require_space(
|
||||
route_root,
|
||||
self._copy_required_bytes(
|
||||
self.qb_root,
|
||||
source_root,
|
||||
route_root,
|
||||
(
|
||||
(entry.target_canonical_path, entry.logical_bytes)
|
||||
@@ -614,7 +620,7 @@ class ClientJobExecutor:
|
||||
)
|
||||
stage_transfer(
|
||||
manifest,
|
||||
source_root=self.qb_root,
|
||||
source_root=source_root,
|
||||
sync_root=route_root,
|
||||
store=self.store,
|
||||
artifact_sources={"metainfo/source.torrent": metainfo_path},
|
||||
@@ -692,23 +698,27 @@ class ClientJobExecutor:
|
||||
"target materialization was sent to the wrong client"
|
||||
)
|
||||
published = load_published_transfer(self._job_directory(definition))
|
||||
info_hash = _info_hash(definition)
|
||||
resource = self.qbittorrent.get_resource(info_hash)
|
||||
target_root = (
|
||||
self._resource_root(resource)
|
||||
if resource is not None else self.qb_root
|
||||
)
|
||||
# The target can likewise hardlink an arrived Syncthing payload into
|
||||
# qB's content root when those directories share a filesystem.
|
||||
self._require_space(
|
||||
self.qb_root,
|
||||
target_root,
|
||||
self._copy_required_bytes(
|
||||
published.job_directory,
|
||||
self.qb_root,
|
||||
target_root,
|
||||
(
|
||||
(entry.payload_relative_path, entry.logical_bytes)
|
||||
for entry in published.manifest.files
|
||||
),
|
||||
),
|
||||
)
|
||||
info_hash = _info_hash(definition)
|
||||
resource = self.qbittorrent.get_resource(info_hash)
|
||||
if resource is not None:
|
||||
self._reject_unsafe_partfile(definition, resource)
|
||||
self._reject_unsafe_partfile(definition, resource, target_root)
|
||||
if definition.transfer.HasField("target_baseline_fingerprint"):
|
||||
if resource is None or not _fingerprint_matches(
|
||||
_resource_fingerprint(resource, self.client_id),
|
||||
@@ -733,13 +743,14 @@ class ClientJobExecutor:
|
||||
== resource_pb2.TORRENT_RUNTIME_STATE_STOPPED
|
||||
),
|
||||
"total_file_count": len(published.manifest.files),
|
||||
"content_root": str(target_root),
|
||||
}
|
||||
self.store.put_job_artifact(
|
||||
definition.job_id, "target-baseline", baseline
|
||||
)
|
||||
materialize_transfer(
|
||||
published,
|
||||
target_root=self.qb_root,
|
||||
target_root=target_root,
|
||||
store=self.store,
|
||||
sparse_supported=self.sparse_supported,
|
||||
cancel_check=lambda: self._raise_if_cancelled(
|
||||
@@ -809,13 +820,16 @@ class ClientJobExecutor:
|
||||
fraction, 0, 0, "qBittorrent stopped recheck"
|
||||
),
|
||||
)
|
||||
if should_start:
|
||||
self.qbittorrent.start(qb_torrent_id)
|
||||
verified = self.qbittorrent.get_resource(info_hash)
|
||||
if verified is None:
|
||||
raise JobExecutionError(
|
||||
"verified qBittorrent resource disappeared"
|
||||
)
|
||||
_normalize_verified_resource_permissions(
|
||||
self._resource_root(verified), verified
|
||||
)
|
||||
if should_start:
|
||||
self.qbittorrent.start(qb_torrent_id)
|
||||
placement = resource_pb2.Placement(
|
||||
client_id=self.client_id,
|
||||
state=resource_pb2.PLACEMENT_STATE_PRESENT,
|
||||
@@ -852,6 +866,7 @@ class ClientJobExecutor:
|
||||
)
|
||||
return
|
||||
baseline = baseline_row["value"]
|
||||
content_root = self._artifact_content_root(baseline)
|
||||
info_hash = _info_hash(definition)
|
||||
present = self.qbittorrent.get_resource(info_hash)
|
||||
if present is not None:
|
||||
@@ -870,7 +885,7 @@ class ClientJobExecutor:
|
||||
self.qbittorrent.delete_entry(qb_torrent_id)
|
||||
compensate_materialized_files(
|
||||
job_id=definition.job_id,
|
||||
qb_root=self.qb_root,
|
||||
qb_root=content_root,
|
||||
store=self.store,
|
||||
)
|
||||
|
||||
@@ -906,10 +921,11 @@ class ClientJobExecutor:
|
||||
self,
|
||||
definition: job_pb2.JobDefinition,
|
||||
resource: NormalizedResource,
|
||||
content_root: Path,
|
||||
) -> None:
|
||||
roots = {self.qb_root}
|
||||
roots = {content_root}
|
||||
for item in resource.files:
|
||||
candidate = self.qb_root / PurePosixPath(item.canonical_path).parts[0]
|
||||
candidate = content_root / PurePosixPath(item.canonical_path).parts[0]
|
||||
roots.add(candidate if candidate.is_dir() else candidate.parent)
|
||||
hashes = {
|
||||
value.lower() for value in (
|
||||
@@ -928,6 +944,52 @@ class ClientJobExecutor:
|
||||
"safely transferred by this client"
|
||||
)
|
||||
|
||||
def _resource_root(self, resource: NormalizedResource) -> Path:
|
||||
"""Resolve a qB resource's save path within this client's mounts.
|
||||
|
||||
qB returns paths in its own API/container namespace. The configured
|
||||
root mapping translates those paths into the client namespace and
|
||||
supports nested save paths and explicit prefix overrides. A missing
|
||||
save path is retained only for older in-process test fixtures; every
|
||||
normalized qB API resource has one.
|
||||
"""
|
||||
|
||||
if resource.save_path is None:
|
||||
return self.qb_root
|
||||
try:
|
||||
root = self.qb_roots.api_to_local(resource.save_path.as_posix())
|
||||
except ConfigError as exc:
|
||||
raise JobExecutionError(
|
||||
"qBittorrent resource save path is outside the configured "
|
||||
"API root or has no local path mapping"
|
||||
) from exc
|
||||
try:
|
||||
metadata = root.lstat()
|
||||
except FileNotFoundError as exc:
|
||||
raise JobExecutionError(
|
||||
"qBittorrent resource save path is not visible in the "
|
||||
"client container"
|
||||
) from exc
|
||||
if not stat.S_ISDIR(metadata.st_mode):
|
||||
raise JobExecutionError(
|
||||
"qBittorrent resource save path is not a real directory in "
|
||||
"the client container"
|
||||
)
|
||||
return root
|
||||
|
||||
def _artifact_content_root(self, baseline: object) -> Path:
|
||||
if not isinstance(baseline, dict):
|
||||
raise JobExecutionError("target baseline is invalid")
|
||||
value = baseline.get("content_root")
|
||||
if not isinstance(value, str):
|
||||
# A pre-existing durable baseline predates per-resource roots.
|
||||
# It can only have materialized to the configured target root.
|
||||
return self.qb_root
|
||||
root = Path(value)
|
||||
if not root.is_absolute() or root.is_symlink():
|
||||
raise JobExecutionError("target baseline content root is unsafe")
|
||||
return root
|
||||
|
||||
def _raise_if_cancelled(self, job_id: str) -> None:
|
||||
event = self._cancel_events.get(job_id)
|
||||
if event is not None and event.is_set():
|
||||
@@ -1050,6 +1112,7 @@ class ClientJobExecutor:
|
||||
committed: bool,
|
||||
step: int = job_pb2.JOB_STEP_KIND_UNSPECIFIED,
|
||||
step_state: int = job_pb2.STEP_STATE_UNSPECIFIED,
|
||||
command_id: str = "",
|
||||
) -> control_pb2.JobEvent:
|
||||
event = control_pb2.JobEvent(
|
||||
event_id=str(uuid.uuid4()),
|
||||
@@ -1059,6 +1122,7 @@ class ClientJobExecutor:
|
||||
type=event_type,
|
||||
state=state,
|
||||
committed=committed,
|
||||
command_id=command_id,
|
||||
)
|
||||
event.occurred_at.GetCurrentTime()
|
||||
if step != job_pb2.JOB_STEP_KIND_UNSPECIFIED:
|
||||
@@ -1190,6 +1254,78 @@ def _info_hash(definition: job_pb2.JobDefinition) -> str:
|
||||
return value
|
||||
|
||||
|
||||
def _normalize_verified_resource_permissions(
|
||||
qb_root: Path, resource: NormalizedResource
|
||||
) -> None:
|
||||
"""Apply ``a+rx``/``a+r`` to exactly qB-verified completed content.
|
||||
|
||||
qBittorrent may finish a stopped recheck with restrictive modes inherited
|
||||
from source materialization. Normalize only selected, fully completed
|
||||
regular files and their real parent directories; never traverse a symlink
|
||||
or broaden permissions on the configured qB root itself.
|
||||
"""
|
||||
root_metadata = qb_root.lstat()
|
||||
if not stat.S_ISDIR(root_metadata.st_mode):
|
||||
raise JobExecutionError("configured qB root is not a real directory")
|
||||
directories: set[Path] = set()
|
||||
files: set[Path] = set()
|
||||
for item in resource.files:
|
||||
if not item.selected or item.completed_bytes != item.logical_bytes:
|
||||
continue
|
||||
relative = _resource_relative_path(item.canonical_path)
|
||||
current = qb_root
|
||||
for component in relative.parts[:-1]:
|
||||
current = current / component
|
||||
metadata = current.lstat()
|
||||
if not stat.S_ISDIR(metadata.st_mode):
|
||||
raise JobExecutionError(
|
||||
"verified resource parent is not a real directory"
|
||||
)
|
||||
directories.add(current)
|
||||
candidate = current / relative.name
|
||||
metadata = candidate.lstat()
|
||||
if not stat.S_ISREG(metadata.st_mode):
|
||||
raise JobExecutionError("verified resource file is not regular")
|
||||
files.add(candidate)
|
||||
for directory in sorted(directories, key=lambda value: len(value.parts)):
|
||||
metadata = directory.lstat()
|
||||
if not stat.S_ISDIR(metadata.st_mode):
|
||||
raise JobExecutionError(
|
||||
"verified resource parent changed during permission update"
|
||||
)
|
||||
os.chmod(
|
||||
directory,
|
||||
stat.S_IMODE(metadata.st_mode) | 0o555,
|
||||
follow_symlinks=False,
|
||||
)
|
||||
for path in files:
|
||||
metadata = path.lstat()
|
||||
if not stat.S_ISREG(metadata.st_mode):
|
||||
raise JobExecutionError(
|
||||
"verified resource file changed during permission update"
|
||||
)
|
||||
os.chmod(
|
||||
path,
|
||||
stat.S_IMODE(metadata.st_mode) | 0o444,
|
||||
follow_symlinks=False,
|
||||
)
|
||||
|
||||
|
||||
def _resource_relative_path(value: str) -> PurePosixPath:
|
||||
if (
|
||||
not value
|
||||
or "\x00" in value
|
||||
or "\\" in value
|
||||
or value.startswith("/")
|
||||
or any(part in {"", ".", ".."} for part in value.split("/"))
|
||||
):
|
||||
raise JobExecutionError("verified resource path is unsafe")
|
||||
path = PurePosixPath(value)
|
||||
if path.is_absolute():
|
||||
raise JobExecutionError("verified resource path is unsafe")
|
||||
return path
|
||||
|
||||
|
||||
def _fingerprint_matches(
|
||||
current: resource_pb2.ResourceStateFingerprint,
|
||||
expected: resource_pb2.ResourceStateFingerprint,
|
||||
|
||||
@@ -18,7 +18,7 @@ class ProtocolError(ValueError):
|
||||
|
||||
def new_envelope() -> envelope_pb2.Envelope:
|
||||
envelope = envelope_pb2.Envelope()
|
||||
envelope.protocol_version.major = 1
|
||||
envelope.protocol_version.major = 2
|
||||
envelope.message_id = str(uuid.uuid4())
|
||||
envelope.sent_at.FromDatetime(datetime.now(timezone.utc))
|
||||
return envelope
|
||||
@@ -59,7 +59,7 @@ def decode(data: str | bytes, max_bytes: int = 1024 * 1024) -> envelope_pb2.Enve
|
||||
json_format.ParseError,
|
||||
) as exc:
|
||||
raise ProtocolError("invalid control envelope") from exc
|
||||
if envelope.protocol_version.major != 1:
|
||||
if envelope.protocol_version.major != 2:
|
||||
raise ProtocolError("unsupported protocol major version")
|
||||
_canonical_uuid(envelope.message_id, "message_id")
|
||||
if envelope.correlation_id:
|
||||
|
||||
@@ -19,10 +19,21 @@ class ResourceError(ValueError):
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class NormalizedResource:
|
||||
"""A normalized qB observation for protocol data and local file work.
|
||||
|
||||
``save_path`` is qBittorrent's API-visible per-torrent content root. It
|
||||
is deliberately local-only: clients resolve it through their qB path
|
||||
mapping immediately before filesystem work, and it is never serialized in
|
||||
inventory, placements, or control protocol messages.
|
||||
"""
|
||||
|
||||
summary: resource_pb2.ResourceSummary
|
||||
files: tuple[resource_pb2.TorrentFile, ...]
|
||||
metainfo: Metainfo
|
||||
metainfo_bytes: bytes = b""
|
||||
# qBittorrent's API-visible save path is intentionally local-only. It
|
||||
# must never become part of inventory or placement protocol messages.
|
||||
save_path: PurePosixPath | None = None
|
||||
|
||||
|
||||
def build_content_tree(
|
||||
@@ -137,6 +148,7 @@ def normalize_resource(
|
||||
character not in "0123456789abcdef" for character in qb_torrent_id
|
||||
):
|
||||
raise ResourceError("torrent hash is invalid")
|
||||
save_path = _save_path(torrent.get("save_path"))
|
||||
summary = resource_pb2.ResourceSummary(
|
||||
qb_torrent_id=qb_torrent_id,
|
||||
display_name=_string(torrent.get("name"), "torrent name"),
|
||||
@@ -172,7 +184,9 @@ def normalize_resource(
|
||||
revision_data, sort_keys=True, separators=(",", ":"),
|
||||
).encode("utf-8")).hexdigest()
|
||||
summary.observed_at.FromDatetime(observed_at)
|
||||
return NormalizedResource(summary, tuple(files), metainfo, metainfo_bytes)
|
||||
return NormalizedResource(
|
||||
summary, tuple(files), metainfo, metainfo_bytes, save_path
|
||||
)
|
||||
|
||||
|
||||
def _set_selection(target: Any, indices: list[int]) -> None:
|
||||
@@ -217,6 +231,20 @@ def _path(value: Any) -> str:
|
||||
return candidate.as_posix()
|
||||
|
||||
|
||||
def _save_path(value: Any) -> PurePosixPath:
|
||||
"""Validate qBittorrent's API-visible per-torrent content root."""
|
||||
|
||||
path = _string(value, "torrent save path")
|
||||
candidate = PurePosixPath(path)
|
||||
if (
|
||||
not candidate.is_absolute()
|
||||
or ".." in candidate.parts
|
||||
or "." in candidate.parts
|
||||
):
|
||||
raise ResourceError("qBittorrent torrent save path is unsafe")
|
||||
return candidate
|
||||
|
||||
|
||||
def _integer(value: Any, name: str) -> int:
|
||||
if isinstance(value, bool) or not isinstance(value, int) or value < 0:
|
||||
raise ResourceError(f"{name} is invalid")
|
||||
|
||||
@@ -38,6 +38,7 @@ class FileOperationConflict(RuntimeError):
|
||||
class CommandAcceptance:
|
||||
duplicate: bool
|
||||
acknowledgement_json: str
|
||||
state: str
|
||||
|
||||
|
||||
class ClientStore:
|
||||
@@ -85,7 +86,7 @@ class ClientStore:
|
||||
connection.execute("BEGIN IMMEDIATE")
|
||||
existing = connection.execute(
|
||||
"""
|
||||
SELECT payload_sha256, command_json, acknowledgement_json
|
||||
SELECT payload_sha256, command_json, acknowledgement_json, state
|
||||
FROM commands WHERE command_id = ?
|
||||
""",
|
||||
(command_id,),
|
||||
@@ -93,17 +94,95 @@ class ClientStore:
|
||||
if existing:
|
||||
if existing["payload_sha256"] != digest or existing["command_json"] != payload:
|
||||
raise CommandConflict("command ID was reused with different content")
|
||||
return CommandAcceptance(True, existing["acknowledgement_json"])
|
||||
return CommandAcceptance(
|
||||
True, existing["acknowledgement_json"], existing["state"]
|
||||
)
|
||||
status = json.loads(acknowledgement).get("status")
|
||||
state = (
|
||||
"rejected"
|
||||
if status is not None
|
||||
and status != "COMMAND_ACK_STATUS_ACCEPTED"
|
||||
else "accepted"
|
||||
)
|
||||
connection.execute(
|
||||
"""
|
||||
INSERT INTO commands (
|
||||
command_id, payload_sha256, command_json,
|
||||
acknowledgement_json, state
|
||||
) VALUES (?, ?, ?, ?, 'received')
|
||||
) VALUES (?, ?, ?, ?, ?)
|
||||
""",
|
||||
(command_id, digest, payload, acknowledgement),
|
||||
(command_id, digest, payload, acknowledgement, state),
|
||||
)
|
||||
return CommandAcceptance(False, acknowledgement)
|
||||
return CommandAcceptance(False, acknowledgement, state)
|
||||
|
||||
def accept_reconcile_command(
|
||||
self,
|
||||
command_id: str,
|
||||
command_json: str,
|
||||
acknowledgement_json: str,
|
||||
*,
|
||||
job_id: str,
|
||||
definition_json: str,
|
||||
state: str,
|
||||
revision: int,
|
||||
last_event_sequence: int,
|
||||
committed: bool,
|
||||
superseded_command_ids: list[str],
|
||||
) -> CommandAcceptance:
|
||||
"""Atomically durably accept and apply a reconciliation command.
|
||||
|
||||
A reconciliation acknowledgement is meaningful only after its cursor
|
||||
and retired leases have reached SQLite. Keeping both operations in
|
||||
one transaction makes a reconnect either redeliver the command or
|
||||
observe its completed effect; it cannot observe a bare acknowledgement.
|
||||
"""
|
||||
payload = _canonical(json.loads(command_json))
|
||||
acknowledgement = _canonical(json.loads(acknowledgement_json))
|
||||
definition = _canonical(json.loads(definition_json))
|
||||
digest = hashlib.sha256(payload.encode()).hexdigest()
|
||||
with self._connect() as connection:
|
||||
connection.execute("BEGIN IMMEDIATE")
|
||||
existing = connection.execute(
|
||||
"""
|
||||
SELECT payload_sha256, command_json, acknowledgement_json, state
|
||||
FROM commands WHERE command_id = ?
|
||||
""",
|
||||
(command_id,),
|
||||
).fetchone()
|
||||
if existing:
|
||||
if existing["payload_sha256"] != digest or existing["command_json"] != payload:
|
||||
raise CommandConflict("command ID was reused with different content")
|
||||
return CommandAcceptance(
|
||||
True, existing["acknowledgement_json"], existing["state"]
|
||||
)
|
||||
status = json.loads(acknowledgement).get("status")
|
||||
command_state = (
|
||||
"rejected"
|
||||
if status is not None
|
||||
and status != "COMMAND_ACK_STATUS_ACCEPTED"
|
||||
else "accepted"
|
||||
)
|
||||
connection.execute(
|
||||
"""
|
||||
INSERT INTO commands (
|
||||
command_id, payload_sha256, command_json,
|
||||
acknowledgement_json, state
|
||||
) VALUES (?, ?, ?, ?, ?)
|
||||
""",
|
||||
(command_id, digest, payload, acknowledgement, command_state),
|
||||
)
|
||||
if command_state == "accepted":
|
||||
self._reconcile_job_connection(
|
||||
connection,
|
||||
job_id=job_id,
|
||||
definition=definition,
|
||||
state=state,
|
||||
revision=revision,
|
||||
last_event_sequence=last_event_sequence,
|
||||
committed=committed,
|
||||
superseded_command_ids=superseded_command_ids,
|
||||
)
|
||||
return CommandAcceptance(False, acknowledgement, command_state)
|
||||
|
||||
def list_active_job_cursors(self) -> list[dict[str, object]]:
|
||||
with self._connect() as connection:
|
||||
@@ -122,7 +201,7 @@ class ClientStore:
|
||||
rows = connection.execute(
|
||||
"""
|
||||
SELECT command_id, command_json, acknowledgement_json
|
||||
FROM commands ORDER BY rowid
|
||||
FROM commands WHERE state = 'accepted' ORDER BY rowid
|
||||
"""
|
||||
).fetchall()
|
||||
return [dict(row) for row in rows]
|
||||
@@ -199,6 +278,114 @@ class ClientStore:
|
||||
),
|
||||
)
|
||||
|
||||
def ensure_job_definition(self, job_id: str, definition_json: str) -> None:
|
||||
"""Durably record an assignment without inventing a global event."""
|
||||
definition = _canonical(json.loads(definition_json))
|
||||
with self._connect() as connection:
|
||||
connection.execute("BEGIN IMMEDIATE")
|
||||
existing = connection.execute(
|
||||
"SELECT definition_json FROM jobs WHERE job_id = ?", (job_id,)
|
||||
).fetchone()
|
||||
if existing is not None:
|
||||
if existing["definition_json"] != definition:
|
||||
raise JobConflict("job definition is immutable")
|
||||
return
|
||||
connection.execute(
|
||||
"""
|
||||
INSERT INTO jobs (job_id, definition_json, state, revision,
|
||||
last_event_sequence, committed)
|
||||
VALUES (?, ?, 'JOB_STATE_QUEUED', 0, 0, 0)
|
||||
""",
|
||||
(job_id, definition),
|
||||
)
|
||||
|
||||
def is_command_active(self, command_id: str) -> bool:
|
||||
with self._connect() as connection:
|
||||
row = connection.execute(
|
||||
"SELECT state FROM commands WHERE command_id = ?", (command_id,)
|
||||
).fetchone()
|
||||
return row is not None and row["state"] == "accepted"
|
||||
|
||||
def reconcile_job(
|
||||
self,
|
||||
*,
|
||||
job_id: str,
|
||||
definition_json: str,
|
||||
state: str,
|
||||
revision: int,
|
||||
last_event_sequence: int,
|
||||
committed: bool,
|
||||
superseded_command_ids: list[str],
|
||||
) -> None:
|
||||
"""Apply control's cursor, retiring only stale command leases.
|
||||
|
||||
This drops unacknowledged local journal rows beyond control's cursor;
|
||||
resource files and operation artifacts are intentionally retained.
|
||||
"""
|
||||
definition = _canonical(json.loads(definition_json))
|
||||
with self._connect() as connection:
|
||||
connection.execute("BEGIN IMMEDIATE")
|
||||
self._reconcile_job_connection(
|
||||
connection,
|
||||
job_id=job_id,
|
||||
definition=definition,
|
||||
state=state,
|
||||
revision=revision,
|
||||
last_event_sequence=last_event_sequence,
|
||||
committed=committed,
|
||||
superseded_command_ids=superseded_command_ids,
|
||||
)
|
||||
|
||||
@staticmethod
|
||||
def _reconcile_job_connection(
|
||||
connection: sqlite3.Connection,
|
||||
*,
|
||||
job_id: str,
|
||||
definition: str,
|
||||
state: str,
|
||||
revision: int,
|
||||
last_event_sequence: int,
|
||||
committed: bool,
|
||||
superseded_command_ids: list[str],
|
||||
) -> None:
|
||||
existing = connection.execute(
|
||||
"SELECT definition_json FROM jobs WHERE job_id = ?", (job_id,)
|
||||
).fetchone()
|
||||
if existing is not None and existing["definition_json"] != definition:
|
||||
raise JobConflict("reconciliation has a different job definition")
|
||||
connection.execute(
|
||||
"DELETE FROM events WHERE job_id = ? AND sequence > ?",
|
||||
(job_id, last_event_sequence),
|
||||
)
|
||||
connection.execute(
|
||||
"""
|
||||
INSERT INTO jobs (job_id, definition_json, state, revision,
|
||||
last_event_sequence, committed)
|
||||
VALUES (?, ?, ?, ?, ?, ?)
|
||||
ON CONFLICT(job_id) DO UPDATE SET
|
||||
state = excluded.state, revision = excluded.revision,
|
||||
last_event_sequence = excluded.last_event_sequence,
|
||||
committed = excluded.committed
|
||||
""",
|
||||
(job_id, definition, state, revision, last_event_sequence, int(committed)),
|
||||
)
|
||||
if superseded_command_ids:
|
||||
placeholders = ",".join("?" for _ in superseded_command_ids)
|
||||
rows = connection.execute(
|
||||
f"SELECT command_id, command_json FROM commands WHERE command_id IN ({placeholders})",
|
||||
superseded_command_ids,
|
||||
).fetchall()
|
||||
owned_ids = [
|
||||
row["command_id"] for row in rows
|
||||
if _command_job_id(row["command_json"]) == job_id
|
||||
]
|
||||
if owned_ids:
|
||||
owned_placeholders = ",".join("?" for _ in owned_ids)
|
||||
connection.execute(
|
||||
f"UPDATE commands SET state = 'superseded' WHERE command_id IN ({owned_placeholders})",
|
||||
owned_ids,
|
||||
)
|
||||
|
||||
def begin_file_operation(
|
||||
self,
|
||||
operation_id: str,
|
||||
@@ -591,6 +778,18 @@ class ClientStore:
|
||||
connection.close()
|
||||
|
||||
|
||||
def _command_job_id(command_json: str) -> str:
|
||||
"""Return the job target from canonical protobuf JSON, if it has one."""
|
||||
command = json.loads(command_json)
|
||||
if "assignJob" in command:
|
||||
return str(command["assignJob"].get("job", {}).get("jobId", ""))
|
||||
if "executeStep" in command:
|
||||
return str(command["executeStep"].get("jobId", ""))
|
||||
if "cancelJob" in command:
|
||||
return str(command["cancelJob"].get("jobId", ""))
|
||||
return ""
|
||||
|
||||
|
||||
def _canonical(value: object) -> str:
|
||||
return json.dumps(value, sort_keys=True, separators=(",", ":"), ensure_ascii=False)
|
||||
|
||||
|
||||
@@ -1,2 +1,3 @@
|
||||
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||
"""Generated archive_control.v1 bindings."""
|
||||
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
||||
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||
# -*- coding: utf-8 -*-
|
||||
# Generated by the protocol buffer compiler. DO NOT EDIT!
|
||||
# NO CHECKED-IN PROTOBUF GENCODE
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
||||
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||
import datetime
|
||||
|
||||
from archive_control.v1 import common_pb2 as _common_pb2
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
||||
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||
# -*- coding: utf-8 -*-
|
||||
# Generated by the protocol buffer compiler. DO NOT EDIT!
|
||||
# NO CHECKED-IN PROTOBUF GENCODE
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
||||
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||
import datetime
|
||||
|
||||
from google.protobuf import timestamp_pb2 as _timestamp_pb2
|
||||
|
||||
File diff suppressed because one or more lines are too long
@@ -1,4 +1,4 @@
|
||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
||||
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||
import datetime
|
||||
|
||||
from archive_control.v1 import common_pb2 as _common_pb2
|
||||
@@ -108,8 +108,20 @@ class RequestJobSnapshotCommand(_message.Message):
|
||||
job_ids: _containers.RepeatedScalarFieldContainer[str]
|
||||
def __init__(self, job_ids: _Optional[_Iterable[str]] = ...) -> None: ...
|
||||
|
||||
class ReconcileJobCommand(_message.Message):
|
||||
__slots__ = ("authoritative_job", "authoritative_last_event_sequence", "superseded_command_ids", "reason")
|
||||
AUTHORITATIVE_JOB_FIELD_NUMBER: _ClassVar[int]
|
||||
AUTHORITATIVE_LAST_EVENT_SEQUENCE_FIELD_NUMBER: _ClassVar[int]
|
||||
SUPERSEDED_COMMAND_IDS_FIELD_NUMBER: _ClassVar[int]
|
||||
REASON_FIELD_NUMBER: _ClassVar[int]
|
||||
authoritative_job: _job_pb2.JobRecord
|
||||
authoritative_last_event_sequence: int
|
||||
superseded_command_ids: _containers.RepeatedScalarFieldContainer[str]
|
||||
reason: str
|
||||
def __init__(self, authoritative_job: _Optional[_Union[_job_pb2.JobRecord, _Mapping]] = ..., authoritative_last_event_sequence: _Optional[int] = ..., superseded_command_ids: _Optional[_Iterable[str]] = ..., reason: _Optional[str] = ...) -> None: ...
|
||||
|
||||
class Command(_message.Message):
|
||||
__slots__ = ("command_id", "created_at", "assign_job", "execute_step", "cancel_job", "ensure_route", "inventory_query", "request_job_snapshot")
|
||||
__slots__ = ("command_id", "created_at", "assign_job", "execute_step", "cancel_job", "ensure_route", "inventory_query", "request_job_snapshot", "reconcile_job")
|
||||
COMMAND_ID_FIELD_NUMBER: _ClassVar[int]
|
||||
CREATED_AT_FIELD_NUMBER: _ClassVar[int]
|
||||
ASSIGN_JOB_FIELD_NUMBER: _ClassVar[int]
|
||||
@@ -118,6 +130,7 @@ class Command(_message.Message):
|
||||
ENSURE_ROUTE_FIELD_NUMBER: _ClassVar[int]
|
||||
INVENTORY_QUERY_FIELD_NUMBER: _ClassVar[int]
|
||||
REQUEST_JOB_SNAPSHOT_FIELD_NUMBER: _ClassVar[int]
|
||||
RECONCILE_JOB_FIELD_NUMBER: _ClassVar[int]
|
||||
command_id: str
|
||||
created_at: _timestamp_pb2.Timestamp
|
||||
assign_job: AssignJobCommand
|
||||
@@ -126,7 +139,8 @@ class Command(_message.Message):
|
||||
ensure_route: EnsureRouteCommand
|
||||
inventory_query: _inventory_pb2.InventoryQuery
|
||||
request_job_snapshot: RequestJobSnapshotCommand
|
||||
def __init__(self, command_id: _Optional[str] = ..., created_at: _Optional[_Union[datetime.datetime, _timestamp_pb2.Timestamp, _Mapping]] = ..., assign_job: _Optional[_Union[AssignJobCommand, _Mapping]] = ..., execute_step: _Optional[_Union[ExecuteStepCommand, _Mapping]] = ..., cancel_job: _Optional[_Union[CancelJobCommand, _Mapping]] = ..., ensure_route: _Optional[_Union[EnsureRouteCommand, _Mapping]] = ..., inventory_query: _Optional[_Union[_inventory_pb2.InventoryQuery, _Mapping]] = ..., request_job_snapshot: _Optional[_Union[RequestJobSnapshotCommand, _Mapping]] = ...) -> None: ...
|
||||
reconcile_job: ReconcileJobCommand
|
||||
def __init__(self, command_id: _Optional[str] = ..., created_at: _Optional[_Union[datetime.datetime, _timestamp_pb2.Timestamp, _Mapping]] = ..., assign_job: _Optional[_Union[AssignJobCommand, _Mapping]] = ..., execute_step: _Optional[_Union[ExecuteStepCommand, _Mapping]] = ..., cancel_job: _Optional[_Union[CancelJobCommand, _Mapping]] = ..., ensure_route: _Optional[_Union[EnsureRouteCommand, _Mapping]] = ..., inventory_query: _Optional[_Union[_inventory_pb2.InventoryQuery, _Mapping]] = ..., request_job_snapshot: _Optional[_Union[RequestJobSnapshotCommand, _Mapping]] = ..., reconcile_job: _Optional[_Union[ReconcileJobCommand, _Mapping]] = ...) -> None: ...
|
||||
|
||||
class CommandAck(_message.Message):
|
||||
__slots__ = ("command_id", "status", "error")
|
||||
@@ -139,7 +153,7 @@ class CommandAck(_message.Message):
|
||||
def __init__(self, command_id: _Optional[str] = ..., status: _Optional[_Union[CommandAckStatus, str]] = ..., error: _Optional[_Union[_common_pb2.Error, _Mapping]] = ...) -> None: ...
|
||||
|
||||
class JobEvent(_message.Message):
|
||||
__slots__ = ("event_id", "job_id", "sequence", "job_revision", "type", "state", "committed", "progress", "error", "observed_resource", "observed_placement", "occurred_at")
|
||||
__slots__ = ("event_id", "job_id", "sequence", "job_revision", "type", "state", "committed", "progress", "error", "observed_resource", "observed_placement", "occurred_at", "command_id")
|
||||
EVENT_ID_FIELD_NUMBER: _ClassVar[int]
|
||||
JOB_ID_FIELD_NUMBER: _ClassVar[int]
|
||||
SEQUENCE_FIELD_NUMBER: _ClassVar[int]
|
||||
@@ -152,6 +166,7 @@ class JobEvent(_message.Message):
|
||||
OBSERVED_RESOURCE_FIELD_NUMBER: _ClassVar[int]
|
||||
OBSERVED_PLACEMENT_FIELD_NUMBER: _ClassVar[int]
|
||||
OCCURRED_AT_FIELD_NUMBER: _ClassVar[int]
|
||||
COMMAND_ID_FIELD_NUMBER: _ClassVar[int]
|
||||
event_id: str
|
||||
job_id: str
|
||||
sequence: int
|
||||
@@ -164,7 +179,8 @@ class JobEvent(_message.Message):
|
||||
observed_resource: _resource_pb2.ResourceStateFingerprint
|
||||
observed_placement: _resource_pb2.Placement
|
||||
occurred_at: _timestamp_pb2.Timestamp
|
||||
def __init__(self, event_id: _Optional[str] = ..., job_id: _Optional[str] = ..., sequence: _Optional[int] = ..., job_revision: _Optional[int] = ..., type: _Optional[_Union[JobEventType, str]] = ..., state: _Optional[_Union[_job_pb2.JobState, str]] = ..., committed: _Optional[bool] = ..., progress: _Optional[_Union[_job_pb2.JobProgress, _Mapping]] = ..., error: _Optional[_Union[_common_pb2.Error, _Mapping]] = ..., observed_resource: _Optional[_Union[_resource_pb2.ResourceStateFingerprint, _Mapping]] = ..., observed_placement: _Optional[_Union[_resource_pb2.Placement, _Mapping]] = ..., occurred_at: _Optional[_Union[datetime.datetime, _timestamp_pb2.Timestamp, _Mapping]] = ...) -> None: ...
|
||||
command_id: str
|
||||
def __init__(self, event_id: _Optional[str] = ..., job_id: _Optional[str] = ..., sequence: _Optional[int] = ..., job_revision: _Optional[int] = ..., type: _Optional[_Union[JobEventType, str]] = ..., state: _Optional[_Union[_job_pb2.JobState, str]] = ..., committed: _Optional[bool] = ..., progress: _Optional[_Union[_job_pb2.JobProgress, _Mapping]] = ..., error: _Optional[_Union[_common_pb2.Error, _Mapping]] = ..., observed_resource: _Optional[_Union[_resource_pb2.ResourceStateFingerprint, _Mapping]] = ..., observed_placement: _Optional[_Union[_resource_pb2.Placement, _Mapping]] = ..., occurred_at: _Optional[_Union[datetime.datetime, _timestamp_pb2.Timestamp, _Mapping]] = ..., command_id: _Optional[str] = ...) -> None: ...
|
||||
|
||||
class JobSnapshot(_message.Message):
|
||||
__slots__ = ("job", "last_event_sequence", "in_flight_command_ids")
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
||||
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||
# -*- coding: utf-8 -*-
|
||||
# Generated by the protocol buffer compiler. DO NOT EDIT!
|
||||
# NO CHECKED-IN PROTOBUF GENCODE
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
||||
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||
import datetime
|
||||
|
||||
from archive_control.v1 import client_pb2 as _client_pb2
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
||||
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||
# -*- coding: utf-8 -*-
|
||||
# Generated by the protocol buffer compiler. DO NOT EDIT!
|
||||
# NO CHECKED-IN PROTOBUF GENCODE
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
||||
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||
from archive_control.v1 import common_pb2 as _common_pb2
|
||||
from archive_control.v1 import resource_pb2 as _resource_pb2
|
||||
from google.protobuf.internal import containers as _containers
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
||||
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||
# -*- coding: utf-8 -*-
|
||||
# Generated by the protocol buffer compiler. DO NOT EDIT!
|
||||
# NO CHECKED-IN PROTOBUF GENCODE
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
||||
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||
import datetime
|
||||
|
||||
from archive_control.v1 import common_pb2 as _common_pb2
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
||||
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||
# -*- coding: utf-8 -*-
|
||||
# Generated by the protocol buffer compiler. DO NOT EDIT!
|
||||
# NO CHECKED-IN PROTOBUF GENCODE
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
||||
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||
import datetime
|
||||
|
||||
from google.protobuf import timestamp_pb2 as _timestamp_pb2
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
||||
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||
# -*- coding: utf-8 -*-
|
||||
# Generated by the protocol buffer compiler. DO NOT EDIT!
|
||||
# NO CHECKED-IN PROTOBUF GENCODE
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
||||
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||
import datetime
|
||||
|
||||
from google.protobuf import timestamp_pb2 as _timestamp_pb2
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
||||
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||
# -*- coding: utf-8 -*-
|
||||
# Generated by the protocol buffer compiler. DO NOT EDIT!
|
||||
# NO CHECKED-IN PROTOBUF GENCODE
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
||||
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||
import datetime
|
||||
|
||||
from archive_control.v1 import resource_pb2 as _resource_pb2
|
||||
|
||||
@@ -0,0 +1,59 @@
|
||||
import io
|
||||
import json
|
||||
import os
|
||||
import tempfile
|
||||
import unittest
|
||||
from contextlib import redirect_stdout
|
||||
from pathlib import Path
|
||||
|
||||
from archive_clients.cli import main
|
||||
|
||||
|
||||
class ClientCliTests(unittest.TestCase):
|
||||
def test_check_config_probes_every_qb_override_root(self):
|
||||
with tempfile.TemporaryDirectory() as directory:
|
||||
root = Path(directory)
|
||||
for name in ("token", "qb-password", "syncthing-key"):
|
||||
path = root / name
|
||||
path.write_text(name, encoding="utf-8")
|
||||
os.chmod(path, 0o600)
|
||||
for name in ("qb", "qb-fast", "sync", "backups"):
|
||||
(root / name).mkdir()
|
||||
config = root / "client.toml"
|
||||
config.write_text(
|
||||
f'''client_id = "cache-1"
|
||||
display_name = "Cache 1"
|
||||
role = "cache"
|
||||
control_endpoint = "ws://control/archive_control"
|
||||
shared_token_file = "{root / "token"}"
|
||||
state_db = "{root / "state.db"}"
|
||||
backup_dir = "{root / "backups"}"
|
||||
|
||||
[qbittorrent]
|
||||
endpoint = "http://qb"
|
||||
username = "admin"
|
||||
password_file = "{root / "qb-password"}"
|
||||
api_root = "/downloads"
|
||||
local_root = "{root / "qb"}"
|
||||
local_path_overrides = {{ "/downloads/fast" = "{root / "qb-fast"}" }}
|
||||
|
||||
[syncthing]
|
||||
endpoint = "http://syncthing"
|
||||
api_key_file = "{root / "syncthing-key"}"
|
||||
api_root = "/sync"
|
||||
local_root = "{root / "sync"}"
|
||||
''',
|
||||
encoding="utf-8",
|
||||
)
|
||||
output = io.StringIO()
|
||||
with redirect_stdout(output):
|
||||
self.assertEqual(main(["--config", str(config), "--check-config"]), 0)
|
||||
reported = json.loads(output.getvalue())
|
||||
self.assertEqual(
|
||||
[item["root"] for item in reported["filesystems"]],
|
||||
[str(root / "qb"), str(root / "qb-fast"), str(root / "sync")],
|
||||
)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
@@ -78,6 +78,35 @@ class ConfigTests(unittest.TestCase):
|
||||
Path("/local/qb/Sync"),
|
||||
)
|
||||
|
||||
def test_qbittorrent_can_override_a_nested_save_path(self):
|
||||
with tempfile.TemporaryDirectory() as directory:
|
||||
root = Path(directory)
|
||||
for name in ("token", "qb-password", "syncthing-key"):
|
||||
path = root / name
|
||||
path.write_text(name, encoding="utf-8")
|
||||
os.chmod(path, 0o600)
|
||||
(root / "qb").mkdir()
|
||||
(root / "fast").mkdir()
|
||||
(root / "sync").mkdir()
|
||||
config_path = root / "client.toml"
|
||||
config_path.write_text(
|
||||
_config(root).replace(
|
||||
f'local_root = "{root / "qb"}"',
|
||||
f'local_root = "{root / "qb"}"\n'
|
||||
"local_path_overrides = { \"/downloads/fast\" = "
|
||||
f'"{root / "fast"}" }}',
|
||||
1,
|
||||
),
|
||||
encoding="utf-8",
|
||||
)
|
||||
config = ClientConfig.load(config_path)
|
||||
self.assertEqual(
|
||||
config.qbittorrent.roots.api_to_local(
|
||||
"/downloads/fast/resource/file.bin"
|
||||
),
|
||||
root / "fast/resource/file.bin",
|
||||
)
|
||||
|
||||
def test_endpoint_scheme_and_job_keys_are_strict(self):
|
||||
with tempfile.TemporaryDirectory() as directory:
|
||||
root = Path(directory)
|
||||
|
||||
@@ -74,7 +74,7 @@ class DaemonTransportTests(unittest.IsolatedAsyncioTestCase):
|
||||
response.register_response.status = (
|
||||
client_pb2.REGISTRATION_STATUS_ACCEPTED
|
||||
)
|
||||
response.register_response.negotiated_version.major = 1
|
||||
response.register_response.negotiated_version.major = 2
|
||||
await websocket.send(encode(response))
|
||||
# Deliberately keep TCP/WebSocket open but send no application
|
||||
# heartbeats. This models a stale proxy/server-side session.
|
||||
@@ -120,7 +120,7 @@ class DaemonTransportTests(unittest.IsolatedAsyncioTestCase):
|
||||
response.register_response.status = (
|
||||
client_pb2.REGISTRATION_STATUS_ACCEPTED
|
||||
)
|
||||
response.register_response.negotiated_version.major = 1
|
||||
response.register_response.negotiated_version.major = 2
|
||||
await websocket.send(encode(response))
|
||||
await websocket.wait_closed()
|
||||
|
||||
@@ -158,7 +158,7 @@ class DaemonTransportTests(unittest.IsolatedAsyncioTestCase):
|
||||
response.register_response.status = (
|
||||
client_pb2.REGISTRATION_STATUS_ACCEPTED
|
||||
)
|
||||
response.register_response.negotiated_version.major = 1
|
||||
response.register_response.negotiated_version.major = 2
|
||||
await websocket.send(encode(response))
|
||||
for sequence in range(1, 102):
|
||||
heartbeat = new_envelope()
|
||||
@@ -485,7 +485,7 @@ class DaemonTransportTests(unittest.IsolatedAsyncioTestCase):
|
||||
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
|
||||
response.register_response.negotiated_version.major = 2
|
||||
await websocket.send(encode(response))
|
||||
heartbeat = new_envelope()
|
||||
heartbeat.heartbeat.sequence = 7
|
||||
|
||||
@@ -50,6 +50,26 @@ class DeploymentPreflightTests(unittest.TestCase):
|
||||
"/data/qb/.archive-control-routes",
|
||||
)
|
||||
|
||||
def test_qb_override_must_map_to_the_same_host_path(self):
|
||||
client = {"Mounts": [
|
||||
{"Type": "bind", "Source": "/srv/fast", "Destination": "/data/fast"},
|
||||
]}
|
||||
qbittorrent = {"Mounts": [
|
||||
{"Type": "bind", "Source": "/srv/fast", "Destination": "/downloads/fast"},
|
||||
]}
|
||||
preflight.require_qb_override_mappings(
|
||||
client, qbittorrent,
|
||||
{"local_path_overrides": {"/downloads/fast": "/data/fast"}},
|
||||
)
|
||||
with self.assertRaisesRegex(preflight.CheckFailure, "host paths differ"):
|
||||
preflight.require_qb_override_mappings(
|
||||
client,
|
||||
{"Mounts": [{
|
||||
"Type": "bind", "Source": "/other", "Destination": "/downloads/fast",
|
||||
}]},
|
||||
{"local_path_overrides": {"/downloads/fast": "/data/fast"}},
|
||||
)
|
||||
|
||||
def test_runs_hardlink_probe_with_qb_and_route_roots(self):
|
||||
completed = __import__("subprocess").CompletedProcess(
|
||||
args=[], returncode=0, stdout="hard-link staging probe passed\n", stderr=""
|
||||
|
||||
+47
-7
@@ -1,5 +1,7 @@
|
||||
import tempfile
|
||||
import unittest
|
||||
from dataclasses import replace
|
||||
from pathlib import PurePosixPath
|
||||
from pathlib import Path
|
||||
|
||||
from archive_clients.eviction import (
|
||||
@@ -86,7 +88,7 @@ class EvictionTests(unittest.TestCase):
|
||||
job_id="job-1",
|
||||
resource=evicted,
|
||||
selected_indices=[0, 1],
|
||||
qb_root=self.root,
|
||||
content_root=self.root,
|
||||
store=self.store,
|
||||
)
|
||||
remove_qb_entry(
|
||||
@@ -97,9 +99,9 @@ class EvictionTests(unittest.TestCase):
|
||||
)
|
||||
result = safe_unlink(
|
||||
job_id="job-1",
|
||||
qb_root=self.root,
|
||||
qbittorrent=qb,
|
||||
store=self.store,
|
||||
resource_root=lambda _: self.root,
|
||||
)
|
||||
self.assertTrue((self.root / "tree/shared.bin").exists())
|
||||
self.assertFalse((self.root / "tree/owned.bin").exists())
|
||||
@@ -121,7 +123,7 @@ class EvictionTests(unittest.TestCase):
|
||||
job_id="job-1",
|
||||
resource=evicted,
|
||||
selected_indices=[0],
|
||||
qb_root=self.root,
|
||||
content_root=self.root,
|
||||
store=self.store,
|
||||
)
|
||||
path.unlink()
|
||||
@@ -134,9 +136,9 @@ class EvictionTests(unittest.TestCase):
|
||||
)
|
||||
result = safe_unlink(
|
||||
job_id="job-1",
|
||||
qb_root=self.root,
|
||||
qbittorrent=qb,
|
||||
store=self.store,
|
||||
resource_root=lambda _: self.root,
|
||||
)
|
||||
self.assertTrue(path.exists())
|
||||
self.assertEqual(
|
||||
@@ -152,7 +154,7 @@ class EvictionTests(unittest.TestCase):
|
||||
job_id="job-1",
|
||||
resource=evicted,
|
||||
selected_indices=[0],
|
||||
qb_root=self.root,
|
||||
content_root=self.root,
|
||||
store=self.store,
|
||||
)
|
||||
remove_qb_entry(
|
||||
@@ -163,19 +165,57 @@ class EvictionTests(unittest.TestCase):
|
||||
)
|
||||
first = safe_unlink(
|
||||
job_id="job-1",
|
||||
qb_root=self.root,
|
||||
qbittorrent=qb,
|
||||
store=self.store,
|
||||
resource_root=lambda _: self.root,
|
||||
)
|
||||
second = safe_unlink(
|
||||
job_id="job-1",
|
||||
qb_root=self.root,
|
||||
qbittorrent=qb,
|
||||
store=self.store,
|
||||
resource_root=lambda _: self.root,
|
||||
)
|
||||
self.assertEqual(first, second)
|
||||
self.assertEqual(qb.deleted, ["a" * 40])
|
||||
|
||||
def test_nested_save_path_uses_its_own_root_and_not_a_same_named_peer(self):
|
||||
nested = self.root / "Downloading"
|
||||
nested.mkdir()
|
||||
evicted = replace(
|
||||
normalized("a" * 40, ["resource/file.bin"]),
|
||||
save_path=PurePosixPath("/downloads/Downloading"),
|
||||
)
|
||||
peer = replace(
|
||||
normalized("b" * 40, ["resource/file.bin"]),
|
||||
save_path=PurePosixPath("/downloads"),
|
||||
)
|
||||
nested_file = nested / "resource/file.bin"
|
||||
nested_file.parent.mkdir()
|
||||
nested_file.write_bytes(b"x" * evicted.files[0].logical_bytes)
|
||||
root_file = self.root / "resource/file.bin"
|
||||
root_file.parent.mkdir()
|
||||
root_file.write_bytes(b"x" * peer.files[0].logical_bytes)
|
||||
qb = FakeQB(evicted, [peer])
|
||||
|
||||
verify_and_snapshot(
|
||||
job_id="job-1", resource=evicted, selected_indices=[0],
|
||||
content_root=nested, store=self.store,
|
||||
)
|
||||
remove_qb_entry(
|
||||
job_id="job-1", torrent_hash="a" * 40,
|
||||
qbittorrent=qb, store=self.store,
|
||||
)
|
||||
safe_unlink(
|
||||
job_id="job-1", qbittorrent=qb, store=self.store,
|
||||
resource_root=lambda item: (
|
||||
nested
|
||||
if item.save_path == PurePosixPath("/downloads/Downloading")
|
||||
else self.root
|
||||
),
|
||||
)
|
||||
self.assertFalse(nested_file.exists())
|
||||
self.assertTrue(root_file.exists())
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
|
||||
+120
-12
@@ -1,21 +1,26 @@
|
||||
import hashlib
|
||||
import os
|
||||
import stat
|
||||
import tempfile
|
||||
import threading
|
||||
import unittest
|
||||
from pathlib import Path
|
||||
from dataclasses import replace
|
||||
from pathlib import Path, PurePosixPath
|
||||
from unittest.mock import Mock, patch
|
||||
from uuid import uuid4
|
||||
|
||||
from archive_clients.bencode import encode
|
||||
from archive_clients.config import RootMapping
|
||||
from archive_clients.jobs import (
|
||||
ClientJobExecutor,
|
||||
JobExecutionError,
|
||||
_normalize_verified_resource_permissions,
|
||||
_resource_fingerprint,
|
||||
)
|
||||
from archive_clients.syncthing import RouteSetupError, SyncthingTransferStatus
|
||||
from archive_clients.resources import normalize_resource
|
||||
from archive_clients.resources import NormalizedResource, normalize_resource
|
||||
from archive_clients.state import ClientStore
|
||||
from archive_control.v1 import control_pb2, job_pb2
|
||||
from archive_control.v1 import control_pb2, job_pb2, resource_pb2
|
||||
|
||||
|
||||
class CompleteSyncthing:
|
||||
@@ -39,6 +44,50 @@ class SlowRescanSyncthing(CompleteSyncthing):
|
||||
|
||||
|
||||
class ClientJobHappyPathTests(unittest.TestCase):
|
||||
def test_verified_resource_permissions_are_readable_by_other_apps(self):
|
||||
with tempfile.TemporaryDirectory() as directory:
|
||||
root = Path(directory) / "qb"
|
||||
resource_directory = root / "resource"
|
||||
resource_directory.mkdir(parents=True)
|
||||
completed = resource_directory / "complete.bin"
|
||||
incomplete = resource_directory / "incomplete.bin"
|
||||
completed.write_bytes(b"complete")
|
||||
incomplete.write_bytes(b"incomplete")
|
||||
os.chmod(resource_directory, 0o300)
|
||||
os.chmod(completed, 0o200)
|
||||
os.chmod(incomplete, 0o200)
|
||||
resource = NormalizedResource(
|
||||
resource_pb2.ResourceSummary(),
|
||||
(
|
||||
resource_pb2.TorrentFile(
|
||||
file_index=0,
|
||||
canonical_path="resource/complete.bin",
|
||||
logical_bytes=len(b"complete"),
|
||||
completed_bytes=len(b"complete"),
|
||||
selected=True,
|
||||
),
|
||||
resource_pb2.TorrentFile(
|
||||
file_index=1,
|
||||
canonical_path="resource/incomplete.bin",
|
||||
logical_bytes=len(b"incomplete"),
|
||||
completed_bytes=0,
|
||||
selected=True,
|
||||
),
|
||||
),
|
||||
Mock(),
|
||||
)
|
||||
_normalize_verified_resource_permissions(root, resource)
|
||||
self.assertEqual(
|
||||
stat.S_IMODE(resource_directory.stat().st_mode) & 0o555,
|
||||
0o555,
|
||||
)
|
||||
self.assertEqual(
|
||||
stat.S_IMODE(completed.stat().st_mode) & 0o444, 0o444
|
||||
)
|
||||
self.assertEqual(
|
||||
stat.S_IMODE(incomplete.stat().st_mode), 0o200
|
||||
)
|
||||
|
||||
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:
|
||||
@@ -199,6 +248,55 @@ class ClientJobHappyPathTests(unittest.TestCase):
|
||||
content=b"x" * 4096,
|
||||
)
|
||||
|
||||
def test_nested_qb_save_path_stages_from_mapped_subdirectory(self):
|
||||
with tempfile.TemporaryDirectory() as directory:
|
||||
self._run_transfer(
|
||||
Path(directory),
|
||||
job_pb2.JOB_OPERATION_UNARCHIVE,
|
||||
source_save_path="/downloads/Downloading",
|
||||
)
|
||||
|
||||
def test_qb_save_path_preflight_rejects_unmapped_path(self):
|
||||
with tempfile.TemporaryDirectory() as directory:
|
||||
root = Path(directory)
|
||||
executor = ClientJobExecutor(
|
||||
client_id="cache-1", qbittorrent=Mock(),
|
||||
store=ClientStore(root / "state.db"), qb_root=root,
|
||||
qb_api_root=PurePosixPath("/downloads"),
|
||||
route_path=lambda _: root, syncthing_transport=Mock(),
|
||||
sparse_supported=True,
|
||||
)
|
||||
resource = NormalizedResource(
|
||||
resource_pb2.ResourceSummary(), (), Mock(),
|
||||
save_path=PurePosixPath("/outside"),
|
||||
)
|
||||
with self.assertRaisesRegex(JobExecutionError, "outside"):
|
||||
executor._resource_root(resource)
|
||||
|
||||
def test_qb_save_path_override_uses_its_dedicated_local_mount(self):
|
||||
with tempfile.TemporaryDirectory() as directory:
|
||||
root = Path(directory)
|
||||
primary = root / "primary"
|
||||
override = root / "override"
|
||||
primary.mkdir()
|
||||
override.mkdir()
|
||||
executor = ClientJobExecutor(
|
||||
client_id="cache-1", qbittorrent=Mock(),
|
||||
store=ClientStore(root / "state.db"), qb_root=primary,
|
||||
qb_api_root=PurePosixPath("/downloads"),
|
||||
qb_roots=RootMapping(
|
||||
PurePosixPath("/downloads"), primary,
|
||||
((PurePosixPath("/downloads/slow"), override),),
|
||||
),
|
||||
route_path=lambda _: root, syncthing_transport=Mock(),
|
||||
sparse_supported=True,
|
||||
)
|
||||
resource = NormalizedResource(
|
||||
resource_pb2.ResourceSummary(), (), Mock(),
|
||||
save_path=PurePosixPath("/downloads/slow"),
|
||||
)
|
||||
self.assertEqual(executor._resource_root(resource), override)
|
||||
|
||||
def test_mount_boundary_requires_copy_space_even_with_same_device(self):
|
||||
with tempfile.TemporaryDirectory() as directory:
|
||||
root = Path(directory)
|
||||
@@ -343,6 +441,7 @@ class ClientJobHappyPathTests(unittest.TestCase):
|
||||
source_stage_free_bytes: int | None = None,
|
||||
content: bytes = b"archive-control-happy-path",
|
||||
syncthing: CompleteSyncthing | None = None,
|
||||
source_save_path: str = "/downloads",
|
||||
):
|
||||
source_root = root / "source"
|
||||
target_root = root / "target"
|
||||
@@ -350,7 +449,12 @@ class ClientJobHappyPathTests(unittest.TestCase):
|
||||
source_root.mkdir()
|
||||
target_root.mkdir()
|
||||
route_root.mkdir()
|
||||
(source_root / "fixture.bin").write_bytes(content)
|
||||
save_relative = PurePosixPath(source_save_path).relative_to(
|
||||
PurePosixPath("/downloads")
|
||||
)
|
||||
source_content_root = source_root.joinpath(*save_relative.parts)
|
||||
source_content_root.mkdir(parents=True, exist_ok=True)
|
||||
(source_content_root / "fixture.bin").write_bytes(content)
|
||||
info = {
|
||||
b"length": len(content),
|
||||
b"name": b"fixture.bin",
|
||||
@@ -364,6 +468,7 @@ class ClientJobHappyPathTests(unittest.TestCase):
|
||||
"hash": torrent_hash,
|
||||
"name": "fixture.bin",
|
||||
"state": "uploading",
|
||||
"save_path": source_save_path,
|
||||
},
|
||||
[{
|
||||
"index": 0,
|
||||
@@ -407,8 +512,11 @@ class ClientJobHappyPathTests(unittest.TestCase):
|
||||
source_qb = Mock()
|
||||
source_qb.get_resource.return_value = resource
|
||||
target_qb = Mock()
|
||||
target_resource = replace(
|
||||
resource, save_path=PurePosixPath("/downloads")
|
||||
)
|
||||
target_qb.get_resource.side_effect = [
|
||||
None, None, resource, resource,
|
||||
None, None, target_resource, target_resource,
|
||||
]
|
||||
syncthing = syncthing or CompleteSyncthing()
|
||||
source = ClientJobExecutor(
|
||||
@@ -436,19 +544,19 @@ class ClientJobHappyPathTests(unittest.TestCase):
|
||||
|
||||
source_assigned = source.assign(control_pb2.AssignJobCommand(
|
||||
job=definition,
|
||||
expected_job_revision=1,
|
||||
expected_job_revision=0,
|
||||
expected_last_event_sequence=0,
|
||||
))
|
||||
target_assigned = target.assign(control_pb2.AssignJobCommand(
|
||||
job=definition,
|
||||
expected_job_revision=1,
|
||||
expected_last_event_sequence=1,
|
||||
expected_job_revision=0,
|
||||
expected_last_event_sequence=0,
|
||||
))
|
||||
self.assertEqual(source_assigned[0].sequence, 1)
|
||||
self.assertEqual(target_assigned[0].sequence, 2)
|
||||
self.assertEqual(source_assigned, [])
|
||||
self.assertEqual(target_assigned, [])
|
||||
|
||||
cursor_revision = 1
|
||||
cursor_sequence = 2
|
||||
cursor_revision = 0
|
||||
cursor_sequence = 0
|
||||
pipeline = (
|
||||
(source, job_pb2.JOB_STEP_KIND_SOURCE_STAGE),
|
||||
(target, job_pb2.JOB_STEP_KIND_SYNCTHING_TRANSFER),
|
||||
|
||||
@@ -55,6 +55,7 @@ class QBittorrentReaderTests(unittest.TestCase):
|
||||
b"Ok.",
|
||||
json.dumps([{
|
||||
"hash": torrent_hash, "name": "a.txt", "state": "uploading",
|
||||
"save_path": "/downloads",
|
||||
}]).encode(),
|
||||
b'[{"index":0,"name":"a.txt","size":3,"progress":1,"priority":1}]',
|
||||
torrent_bytes,
|
||||
@@ -102,6 +103,7 @@ class QBittorrentReaderTests(unittest.TestCase):
|
||||
b"Ok.",
|
||||
json.dumps([{
|
||||
"hash": qb_hash, "name": "a.txt", "state": "uploading",
|
||||
"save_path": "/downloads",
|
||||
}]).encode(),
|
||||
b'[{"index":0,"name":"a.txt","size":3,"progress":1,"priority":1}]',
|
||||
torrent_bytes,
|
||||
@@ -157,6 +159,7 @@ class QBittorrentReaderTests(unittest.TestCase):
|
||||
"infohash_v2": v2_hash,
|
||||
"name": "a.txt",
|
||||
"state": "uploading",
|
||||
"save_path": "/downloads",
|
||||
}]).encode(),
|
||||
b'[{"index":0,"name":"a.txt","size":3,"progress":1,"priority":1}]',
|
||||
torrent_bytes,
|
||||
@@ -217,9 +220,9 @@ class QBittorrentReaderTests(unittest.TestCase):
|
||||
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"},
|
||||
{"hash": first_hash, "name": "first.txt", "state": "uploading", "save_path": "/downloads"},
|
||||
{"hash": malformed_hash, "name": "broken.txt", "state": "stalledUP", "save_path": "/downloads"},
|
||||
{"hash": second_hash, "name": "second.txt", "state": "uploading", "save_path": "/downloads"},
|
||||
]
|
||||
responses = [
|
||||
b"Ok.", json.dumps(torrents).encode(),
|
||||
|
||||
+22
-1
@@ -40,6 +40,7 @@ class ResourceTests(unittest.TestCase):
|
||||
"hash": decoded.info_hash_v2_hex,
|
||||
"name": "v2.bin",
|
||||
"state": "stoppedUP",
|
||||
"save_path": "/downloads",
|
||||
},
|
||||
[{
|
||||
"index": 0,
|
||||
@@ -91,6 +92,7 @@ class ResourceTests(unittest.TestCase):
|
||||
"hash": decoded.info_hash_v1_hex,
|
||||
"name": "resource",
|
||||
"state": "stoppedUP",
|
||||
"save_path": "/downloads/nested",
|
||||
},
|
||||
[
|
||||
{
|
||||
@@ -106,6 +108,7 @@ class ResourceTests(unittest.TestCase):
|
||||
observed,
|
||||
)
|
||||
summary = normalized.summary
|
||||
self.assertEqual(normalized.save_path.as_posix(), "/downloads/nested")
|
||||
self.assertEqual(
|
||||
summary.runtime_state, resource_pb2.TORRENT_RUNTIME_STATE_STOPPED
|
||||
)
|
||||
@@ -139,6 +142,7 @@ class ResourceTests(unittest.TestCase):
|
||||
"hash": torrent_hash,
|
||||
"name": "with-padding",
|
||||
"state": "stalledUP",
|
||||
"save_path": "/downloads",
|
||||
},
|
||||
[
|
||||
{
|
||||
@@ -169,7 +173,7 @@ class ResourceTests(unittest.TestCase):
|
||||
metainfo = encode({b"info": info})
|
||||
torrent = {
|
||||
"hash": hashlib.sha1(encode(info)).hexdigest(),
|
||||
"name": "renamed", "state": "uploading",
|
||||
"name": "renamed", "state": "uploading", "save_path": "/downloads",
|
||||
}
|
||||
renamed = normalize_resource(torrent, [{
|
||||
"index": 0, "name": "renamed.txt", "size": 3,
|
||||
@@ -182,6 +186,23 @@ class ResourceTests(unittest.TestCase):
|
||||
"progress": 1.0, "priority": 1,
|
||||
}], metainfo)
|
||||
|
||||
def test_missing_or_out_of_shape_save_path_is_rejected(self):
|
||||
info = {
|
||||
b"length": 3, b"name": b"a.txt", b"piece length": 16384,
|
||||
b"pieces": b"x" * 20,
|
||||
}
|
||||
metainfo = encode({b"info": info})
|
||||
torrent = {
|
||||
"hash": hashlib.sha1(encode(info)).hexdigest(),
|
||||
"name": "a.txt", "state": "uploading", "save_path": "relative",
|
||||
}
|
||||
with self.assertRaisesRegex(ResourceError, "save path is unsafe"):
|
||||
normalize_resource(torrent, [{
|
||||
"index": 0, "name": "a.txt", "size": 3,
|
||||
"completed": 3, "priority": 1,
|
||||
}], metainfo)
|
||||
|
||||
|
||||
def test_noncanonical_bencode_is_rejected(self):
|
||||
with self.assertRaisesRegex(BencodeError, "unsorted"):
|
||||
decode_metainfo(b"d4:infod1:b1:x1:a1:yee")
|
||||
|
||||
@@ -177,6 +177,67 @@ class ClientStoreTests(unittest.TestCase):
|
||||
"job-1", "baseline", {"selected": [2]}
|
||||
)
|
||||
|
||||
def test_reconciliation_retires_only_stale_job_leases(self):
|
||||
with tempfile.TemporaryDirectory() as directory:
|
||||
store = ClientStore(Path(directory) / "state.db")
|
||||
store.initialize()
|
||||
stale = store.accept_command(
|
||||
"stale", '{"executeStep":{"jobId":"job-1"}}',
|
||||
'{"status":"COMMAND_ACK_STATUS_ACCEPTED"}',
|
||||
)
|
||||
store.accept_command(
|
||||
"other", '{"executeStep":{"jobId":"job-2"}}',
|
||||
'{"status":"COMMAND_ACK_STATUS_ACCEPTED"}',
|
||||
)
|
||||
self.assertEqual(stale.state, "accepted")
|
||||
store.save_job(
|
||||
"job-1", '{"jobId":"job-1"}', "JOB_STATE_RUNNING", 3, 3, False,
|
||||
)
|
||||
store.reconcile_job(
|
||||
job_id="job-1", definition_json='{"jobId":"job-1"}',
|
||||
state="JOB_STATE_QUEUED", revision=0,
|
||||
last_event_sequence=0, committed=False,
|
||||
superseded_command_ids=["stale"],
|
||||
)
|
||||
self.assertFalse(store.is_command_active("stale"))
|
||||
self.assertTrue(store.is_command_active("other"))
|
||||
row = store.job_snapshot_rows(["job-1"])[0]
|
||||
self.assertEqual(row["last_event_sequence"], 0)
|
||||
self.assertEqual(row["state"], "JOB_STATE_QUEUED")
|
||||
|
||||
def test_reconciliation_acknowledgement_and_state_are_atomic(self):
|
||||
with tempfile.TemporaryDirectory() as directory:
|
||||
store = ClientStore(Path(directory) / "state.db")
|
||||
store.initialize()
|
||||
store.accept_command(
|
||||
"stale", '{"executeStep":{"jobId":"job-1"}}',
|
||||
'{"status":"COMMAND_ACK_STATUS_ACCEPTED"}',
|
||||
)
|
||||
accepted = store.accept_reconcile_command(
|
||||
"reconcile-1", '{"reconcileJob":{"authoritativeJob":{}}}',
|
||||
'{"status":"COMMAND_ACK_STATUS_ACCEPTED"}',
|
||||
job_id="job-1", definition_json='{"jobId":"job-1"}',
|
||||
state="JOB_STATE_QUEUED", revision=0,
|
||||
last_event_sequence=0, committed=False,
|
||||
superseded_command_ids=["stale"],
|
||||
)
|
||||
self.assertFalse(accepted.duplicate)
|
||||
self.assertEqual(accepted.state, "accepted")
|
||||
self.assertTrue(store.is_command_active("reconcile-1"))
|
||||
self.assertFalse(store.is_command_active("stale"))
|
||||
self.assertEqual(
|
||||
store.job_snapshot_rows(["job-1"])[0]["last_event_sequence"], 0
|
||||
)
|
||||
duplicate = store.accept_reconcile_command(
|
||||
"reconcile-1", '{"reconcileJob":{"authoritativeJob":{}}}',
|
||||
'{"status":"COMMAND_ACK_STATUS_ACCEPTED"}',
|
||||
job_id="job-1", definition_json='{"jobId":"job-1"}',
|
||||
state="JOB_STATE_QUEUED", revision=0,
|
||||
last_event_sequence=0, committed=False,
|
||||
superseded_command_ids=["stale"],
|
||||
)
|
||||
self.assertTrue(duplicate.duplicate)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
|
||||
Reference in New Issue
Block a user