Compare commits
5
Commits
v0.1.20
...
ae7175719b
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
ae7175719b | ||
|
|
1c8128fe55 | ||
|
|
cd3a9b3c0d | ||
|
|
0311f6f084 | ||
|
|
62a0f24437 |
@@ -7,6 +7,7 @@ RUN pip wheel --no-cache-dir --wheel-dir /wheels .
|
|||||||
FROM builder AS test
|
FROM builder AS test
|
||||||
RUN pip install --no-cache-dir /wheels/*.whl
|
RUN pip install --no-cache-dir /wheels/*.whl
|
||||||
COPY tests ./tests
|
COPY tests ./tests
|
||||||
|
COPY scripts ./scripts
|
||||||
CMD ["python", "-m", "unittest", "discover", "-s", "tests", "-v"]
|
CMD ["python", "-m", "unittest", "discover", "-s", "tests", "-v"]
|
||||||
|
|
||||||
FROM python:3.11-slim
|
FROM python:3.11-slim
|
||||||
|
|||||||
@@ -95,6 +95,7 @@ backup_dir = "/var/backups/archive-control"
|
|||||||
registration_timeout = "10s" # First response must arrive within this window.
|
registration_timeout = "10s" # First response must arrive within this window.
|
||||||
heartbeat_interval = "15s" # Server may negotiate a different effective value.
|
heartbeat_interval = "15s" # Server may negotiate a different effective value.
|
||||||
offline_timeout = "45s"
|
offline_timeout = "45s"
|
||||||
|
outbound_enqueue_timeout = "15s" # Reconnect rather than wedging if sends stop draining.
|
||||||
reconnect_initial = "1s"
|
reconnect_initial = "1s"
|
||||||
reconnect_max = "60s" # Retry forever, never wait longer than this.
|
reconnect_max = "60s" # Retry forever, never wait longer than this.
|
||||||
reconnect_reset_after = "60s"
|
reconnect_reset_after = "60s"
|
||||||
@@ -262,24 +263,24 @@ control consumer, client consumer, E2E, then image publication.
|
|||||||
### Reproducible multi-platform Buildx lifecycle
|
### Reproducible multi-platform Buildx lifecycle
|
||||||
|
|
||||||
The named Buildx builders are local acceleration/cache only; they are not a
|
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
|
deployment dependency and can be removed after publication. Use the checked-in
|
||||||
builder, verify its platforms, publish a release, and remove it afterwards:
|
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
|
```bash
|
||||||
docker buildx create --name archive-control-release --driver docker-container --use
|
scripts/publish-image.sh vX.Y.Z
|
||||||
docker buildx inspect --bootstrap
|
scripts/publish-image.sh vX.Y.Z --also-latest
|
||||||
docker buildx build --platform linux/amd64,linux/arm64 \
|
|
||||||
--tag sodium/archive-clients:vX.Y.Z --push .
|
|
||||||
docker buildx rm archive-control-release
|
|
||||||
```
|
```
|
||||||
|
|
||||||
`docker buildx inspect` must list both `linux/amd64` and `linux/arm64` before
|
The publisher deliberately uses `moby/buildkit:rootless` with
|
||||||
publishing. If the host has no arm64 emulation, install/configure it according
|
`--oci-worker-no-process-sandbox`. On nested Docker hosts, the default OCI
|
||||||
to the host Docker distribution before the build; do not publish a partial
|
sandbox can fail while masking `/proc/acpi` for an emulated build; rootless
|
||||||
single-platform tag. Retain the pushed manifest digest in the release notes
|
BuildKit confines that compatibility setting to the disposable builder. It
|
||||||
and deploy the immutable tag or digest. The optional `archive-control-qemu`
|
refuses to publish unless `docker buildx inspect` reports both `linux/amd64`
|
||||||
builder follows the same lifecycle when it is used for an emulation smoke
|
and `linux/arm64`, and removes the builder and binfmt handler on success,
|
||||||
build.
|
failure, or interruption. Retain the displayed manifest digest in release
|
||||||
|
notes and deploy the immutable tag or digest.
|
||||||
|
|
||||||
## Operator usage
|
## Operator usage
|
||||||
|
|
||||||
|
|||||||
+15
-12
@@ -45,8 +45,8 @@ seconds. A connection stable for 60 seconds resets the backoff.
|
|||||||
|
|
||||||
## Version and feature negotiation
|
## Version and feature negotiation
|
||||||
|
|
||||||
The envelope and registration both state the sender version. Major `1` accepts
|
The envelope and registration both state the sender version. Major `2` accepts
|
||||||
only major `1`; a different major is rejected. The negotiated minor is the
|
only major `2`; a different major is rejected. The negotiated minor is the
|
||||||
highest mutually supported minor no greater than either endpoint's advertised
|
highest mutually supported minor no greater than either endpoint's advertised
|
||||||
minor.
|
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
|
Every mutating command includes an expected job revision or immutable job
|
||||||
definition plus the expected last global event sequence. A stale revision or
|
definition plus the expected last global event sequence. A stale revision or
|
||||||
sequence is rejected without side effects. The control daemon sends only one
|
sequence is rejected without side effects. Assignment is represented solely by
|
||||||
active step command for a job and waits for its persisted completion event
|
the durable assignment `CommandAck`; it produces no `JobEvent`. The control
|
||||||
before commanding the next participant. This single-writer lease lets whichever
|
daemon sends only one active step command for a job and waits for its persisted
|
||||||
client owns the current step allocate the next global per-job event sequence;
|
completion event before commanding the next participant.
|
||||||
the following command starts from the sequence control has durably accepted.
|
|
||||||
|
|
||||||
## Events, progress, and reconciliation
|
## Events, progress, and reconciliation
|
||||||
|
|
||||||
Each client-originated job event has a unique ID, monotonically increasing
|
Each client-originated job event has a unique ID, monotonically increasing
|
||||||
global per-job `sequence`, and resulting job revision. Duplicate IDs/sequences
|
global per-job `sequence`, resulting job revision, and the `command_id` of its
|
||||||
are idempotent only when their complete content matches. Control never grants
|
owning `ExecuteStepCommand` or `CancelJobCommand`. Duplicate IDs/sequences are
|
||||||
concurrent event-writer leases for one job. A gap or conflicting duplicate
|
idempotent only when their complete content matches. Control verifies the
|
||||||
pauses destructive orchestration and requests snapshots.
|
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
|
`fraction_complete` is current-step progress and
|
||||||
`overall_fraction_complete` is the weighted five- or three-step job progress;
|
`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:
|
sequence, state, and commit flag. Reconciliation applies these rules:
|
||||||
|
|
||||||
1. Equal cursors resume normal delivery.
|
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.
|
3. Control behind requests and validates the client's full job snapshot.
|
||||||
4. Conflicting commit evidence reserves the resource and requires manual
|
4. Conflicting commit evidence reserves the resource and requires manual
|
||||||
reconciliation; neither side performs cleanup.
|
reconciliation; neither side performs cleanup.
|
||||||
|
|||||||
+4
-1
@@ -125,7 +125,10 @@ and placement policy differ.
|
|||||||
and explicit save path, or apply a transactional delta to an existing
|
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,
|
torrent. Set selected/skipped state, run a full recheck of the target union,
|
||||||
and fail immediately if qBittorrent attempts to download. Successful recheck
|
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
|
5. **Staging Cleanup.** Remove only job-owned staging artifacts from both
|
||||||
endpoints. A post-commit cleanup failure produces `CLEANUP_REQUIRED`; it
|
endpoints. A post-commit cleanup failure produces `CLEANUP_REQUIRED`; it
|
||||||
never rolls back or deletes the committed placement.
|
never rolls back or deletes the committed placement.
|
||||||
|
|||||||
+1
-1
@@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta"
|
|||||||
|
|
||||||
[project]
|
[project]
|
||||||
name = "archive-clients"
|
name = "archive-clients"
|
||||||
version = "0.1.19"
|
version = "0.1.23"
|
||||||
requires-python = ">=3.11"
|
requires-python = ">=3.11"
|
||||||
dependencies = ["protobuf==7.35.1", "websockets==16.0"]
|
dependencies = ["protobuf==7.35.1", "websockets==16.0"]
|
||||||
|
|
||||||
|
|||||||
Executable
+86
@@ -0,0 +1,86 @@
|
|||||||
|
#!/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
|
||||||
|
|
||||||
|
platforms=$(docker buildx inspect "${builder}" --bootstrap 2>&1)
|
||||||
|
for platform in linux/amd64 linux/arm64; do
|
||||||
|
if ! grep -Fq "${platform}" <<<"${platforms}"; then
|
||||||
|
echo "Builder ${builder} does not support ${platform}; refusing partial release." >&2
|
||||||
|
exit 1
|
||||||
|
fi
|
||||||
|
done
|
||||||
|
|
||||||
|
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."""
|
"""Archive Control data-node daemon."""
|
||||||
|
|
||||||
PROTO_COMMIT = "4ec852014dad74606d4078b3ae1aa208c814b033"
|
PROTO_COMMIT = "03b6751d420e8e4e4cc9aeef4f76b13e28a7265f"
|
||||||
|
|
||||||
|
|||||||
@@ -85,6 +85,7 @@ class ConnectionConfig:
|
|||||||
registration_timeout: float = 10
|
registration_timeout: float = 10
|
||||||
heartbeat_interval: float = 15
|
heartbeat_interval: float = 15
|
||||||
offline_timeout: float = 45
|
offline_timeout: float = 45
|
||||||
|
outbound_enqueue_timeout: float = 15
|
||||||
reconnect_initial: float = 1
|
reconnect_initial: float = 1
|
||||||
reconnect_max: float = 60
|
reconnect_max: float = 60
|
||||||
reconnect_reset_after: float = 60
|
reconnect_reset_after: float = 60
|
||||||
@@ -249,6 +250,7 @@ def _connection(value: Any) -> ConnectionConfig:
|
|||||||
raise ConfigError("connection must be a table")
|
raise ConfigError("connection must be a table")
|
||||||
_keys(value, {
|
_keys(value, {
|
||||||
"registration_timeout", "heartbeat_interval", "offline_timeout",
|
"registration_timeout", "heartbeat_interval", "offline_timeout",
|
||||||
|
"outbound_enqueue_timeout",
|
||||||
"reconnect_initial", "reconnect_max", "reconnect_reset_after",
|
"reconnect_initial", "reconnect_max", "reconnect_reset_after",
|
||||||
"reconnect_jitter",
|
"reconnect_jitter",
|
||||||
}, "connection")
|
}, "connection")
|
||||||
@@ -256,6 +258,9 @@ def _connection(value: Any) -> ConnectionConfig:
|
|||||||
registration_timeout=_duration(value.get("registration_timeout", "10s")),
|
registration_timeout=_duration(value.get("registration_timeout", "10s")),
|
||||||
heartbeat_interval=_duration(value.get("heartbeat_interval", "15s")),
|
heartbeat_interval=_duration(value.get("heartbeat_interval", "15s")),
|
||||||
offline_timeout=_duration(value.get("offline_timeout", "45s")),
|
offline_timeout=_duration(value.get("offline_timeout", "45s")),
|
||||||
|
outbound_enqueue_timeout=_duration(
|
||||||
|
value.get("outbound_enqueue_timeout", "15s")
|
||||||
|
),
|
||||||
reconnect_initial=_duration(value.get("reconnect_initial", "1s")),
|
reconnect_initial=_duration(value.get("reconnect_initial", "1s")),
|
||||||
reconnect_max=_duration(value.get("reconnect_max", "60s")),
|
reconnect_max=_duration(value.get("reconnect_max", "60s")),
|
||||||
reconnect_reset_after=_duration(value.get("reconnect_reset_after", "60s")),
|
reconnect_reset_after=_duration(value.get("reconnect_reset_after", "60s")),
|
||||||
@@ -265,6 +270,8 @@ def _connection(value: Any) -> ConnectionConfig:
|
|||||||
raise ConfigError("reconnect_initial cannot exceed reconnect_max")
|
raise ConfigError("reconnect_initial cannot exceed reconnect_max")
|
||||||
if result.offline_timeout <= result.heartbeat_interval:
|
if result.offline_timeout <= result.heartbeat_interval:
|
||||||
raise ConfigError("offline_timeout must exceed heartbeat_interval")
|
raise ConfigError("offline_timeout must exceed heartbeat_interval")
|
||||||
|
if result.outbound_enqueue_timeout <= 0:
|
||||||
|
raise ConfigError("outbound_enqueue_timeout must be positive")
|
||||||
if not isinstance(result.reconnect_jitter, bool):
|
if not isinstance(result.reconnect_jitter, bool):
|
||||||
raise ConfigError("reconnect_jitter must be boolean")
|
raise ConfigError("reconnect_jitter must be boolean")
|
||||||
return result
|
return result
|
||||||
|
|||||||
+164
-63
@@ -214,7 +214,7 @@ class ArchiveClientDaemon:
|
|||||||
raise RuntimeError("registration response correlation mismatch")
|
raise RuntimeError("registration response correlation mismatch")
|
||||||
if response.register_response.status != client_pb2.REGISTRATION_STATUS_ACCEPTED:
|
if response.register_response.status != client_pb2.REGISTRATION_STATUS_ACCEPTED:
|
||||||
raise RuntimeError("control rejected registration")
|
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")
|
raise RuntimeError("control negotiated an unsupported protocol version")
|
||||||
logger.info("control_connection_registered")
|
logger.info("control_connection_registered")
|
||||||
outbound: asyncio.Queue[str] = asyncio.Queue(maxsize=100)
|
outbound: asyncio.Queue[str] = asyncio.Queue(maxsize=100)
|
||||||
@@ -229,16 +229,11 @@ class ArchiveClientDaemon:
|
|||||||
# detects the missing application heartbeats and reconnects.
|
# detects the missing application heartbeats and reconnects.
|
||||||
while True:
|
while True:
|
||||||
try:
|
try:
|
||||||
frame = await asyncio.wait_for(
|
frame = await self._next_frame(
|
||||||
websocket.recv(),
|
websocket, writer
|
||||||
self.config.connection.offline_timeout,
|
|
||||||
)
|
)
|
||||||
except ConnectionClosedOK:
|
except ConnectionClosedOK:
|
||||||
return
|
return
|
||||||
except asyncio.TimeoutError as exc:
|
|
||||||
raise RuntimeError(
|
|
||||||
"control heartbeat timed out"
|
|
||||||
) from exc
|
|
||||||
await self._handle(decode(frame), outbound, command_tasks)
|
await self._handle(decode(frame), outbound, command_tasks)
|
||||||
finally:
|
finally:
|
||||||
writer.cancel()
|
writer.cancel()
|
||||||
@@ -323,6 +318,63 @@ class ArchiveClientDaemon:
|
|||||||
while True:
|
while True:
|
||||||
await websocket.send(await outbound.get())
|
await websocket.send(await outbound.get())
|
||||||
|
|
||||||
|
async def _next_frame(self, websocket: Any, writer: asyncio.Task[None]):
|
||||||
|
"""Receive one frame while supervising the connection writer.
|
||||||
|
|
||||||
|
The previous receive-only wait let a failed writer go unnoticed until
|
||||||
|
the remote side happened to close or a heartbeat timeout elapsed.
|
||||||
|
Waiting for both directions makes a failed send an immediate,
|
||||||
|
reconnectable connection failure.
|
||||||
|
"""
|
||||||
|
|
||||||
|
receiver = asyncio.create_task(websocket.recv())
|
||||||
|
try:
|
||||||
|
done, _ = await asyncio.wait(
|
||||||
|
{receiver, writer},
|
||||||
|
timeout=self.config.connection.offline_timeout,
|
||||||
|
return_when=asyncio.FIRST_COMPLETED,
|
||||||
|
)
|
||||||
|
if writer in done:
|
||||||
|
if not receiver.done():
|
||||||
|
receiver.cancel()
|
||||||
|
await asyncio.gather(receiver, return_exceptions=True)
|
||||||
|
if writer.cancelled():
|
||||||
|
raise RuntimeError("control writer was cancelled")
|
||||||
|
error = writer.exception()
|
||||||
|
if error is not None:
|
||||||
|
raise RuntimeError("control writer failed") from error
|
||||||
|
raise RuntimeError("control writer stopped")
|
||||||
|
if receiver not in done:
|
||||||
|
receiver.cancel()
|
||||||
|
await asyncio.gather(receiver, return_exceptions=True)
|
||||||
|
raise RuntimeError("control heartbeat timed out")
|
||||||
|
return receiver.result()
|
||||||
|
except BaseException:
|
||||||
|
if not receiver.done():
|
||||||
|
receiver.cancel()
|
||||||
|
await asyncio.gather(receiver, return_exceptions=True)
|
||||||
|
raise
|
||||||
|
|
||||||
|
async def _enqueue(
|
||||||
|
self, outbound: asyncio.Queue[str], message: str
|
||||||
|
) -> None:
|
||||||
|
"""Bound connection-local backpressure so a dead writer cannot wedge I/O.
|
||||||
|
|
||||||
|
A job event is durable before it is sent, so dropping this particular
|
||||||
|
connection after a bounded wait is safe: reconciliation/redelivery on
|
||||||
|
the next session will recover it. In contrast, indefinitely waiting
|
||||||
|
for a full queue can prevent heartbeat acknowledgements from being
|
||||||
|
read or sent, which prevents the reconnect supervisor from running.
|
||||||
|
"""
|
||||||
|
|
||||||
|
try:
|
||||||
|
await asyncio.wait_for(
|
||||||
|
outbound.put(message),
|
||||||
|
timeout=self.config.connection.outbound_enqueue_timeout,
|
||||||
|
)
|
||||||
|
except asyncio.TimeoutError as exc:
|
||||||
|
raise RuntimeError("control outbound queue is blocked") from exc
|
||||||
|
|
||||||
async def _handle(
|
async def _handle(
|
||||||
self,
|
self,
|
||||||
envelope: Any,
|
envelope: Any,
|
||||||
@@ -334,13 +386,20 @@ class ArchiveClientDaemon:
|
|||||||
response = new_envelope()
|
response = new_envelope()
|
||||||
response.correlation_id = envelope.message_id
|
response.correlation_id = envelope.message_id
|
||||||
response.heartbeat_ack.sequence = envelope.heartbeat.sequence
|
response.heartbeat_ack.sequence = envelope.heartbeat.sequence
|
||||||
await outbound.put(encode(response))
|
await self._enqueue(outbound, encode(response))
|
||||||
elif payload == "command":
|
elif payload == "command":
|
||||||
await self._accept_command(envelope, outbound, command_tasks)
|
await self._accept_command(envelope, outbound, command_tasks)
|
||||||
elif payload == "protocol_error":
|
elif payload == "protocol_error":
|
||||||
logger.warning(
|
logger.warning(
|
||||||
"control_reported_protocol_error",
|
"control_reported_protocol_error",
|
||||||
extra={"error_code": envelope.protocol_error.error.code},
|
extra={
|
||||||
|
"error_code": envelope.protocol_error.error.code,
|
||||||
|
"error_detail": envelope.protocol_error.error.message,
|
||||||
|
"offending_message_id": (
|
||||||
|
envelope.protocol_error.offending_message_id
|
||||||
|
),
|
||||||
|
"retryable": envelope.protocol_error.error.retryable,
|
||||||
|
},
|
||||||
)
|
)
|
||||||
|
|
||||||
async def _accept_command(
|
async def _accept_command(
|
||||||
@@ -368,12 +427,36 @@ class ArchiveClientDaemon:
|
|||||||
accepted = None
|
accepted = None
|
||||||
accepted_for_execution = False
|
accepted_for_execution = False
|
||||||
try:
|
try:
|
||||||
accepted = await asyncio.to_thread(
|
if command.WhichOneof("payload") == "reconcile_job":
|
||||||
self.store.accept_command,
|
reconciliation = command.reconcile_job
|
||||||
command.command_id,
|
accepted = await asyncio.to_thread(
|
||||||
encode_message(command),
|
self.store.accept_reconcile_command,
|
||||||
encode_message(acknowledgement),
|
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:
|
if accepted.duplicate:
|
||||||
acknowledgement = decode_message(
|
acknowledgement = decode_message(
|
||||||
accepted.acknowledgement_json, control_pb2.CommandAck()
|
accepted.acknowledgement_json, control_pb2.CommandAck()
|
||||||
@@ -381,6 +464,7 @@ class ArchiveClientDaemon:
|
|||||||
if (
|
if (
|
||||||
acknowledgement.status
|
acknowledgement.status
|
||||||
== control_pb2.COMMAND_ACK_STATUS_ACCEPTED
|
== control_pb2.COMMAND_ACK_STATUS_ACCEPTED
|
||||||
|
and accepted.state == "accepted"
|
||||||
):
|
):
|
||||||
accepted_for_execution = True
|
accepted_for_execution = True
|
||||||
acknowledgement.status = (
|
acknowledgement.status = (
|
||||||
@@ -398,7 +482,7 @@ class ArchiveClientDaemon:
|
|||||||
response = new_envelope()
|
response = new_envelope()
|
||||||
response.correlation_id = envelope.message_id
|
response.correlation_id = envelope.message_id
|
||||||
response.command_ack.CopyFrom(acknowledgement)
|
response.command_ack.CopyFrom(acknowledgement)
|
||||||
await outbound.put(encode(response))
|
await self._enqueue(outbound, encode(response))
|
||||||
|
|
||||||
if (
|
if (
|
||||||
accepted is not None
|
accepted is not None
|
||||||
@@ -419,7 +503,7 @@ class ArchiveClientDaemon:
|
|||||||
job_snapshot.last_event_sequence = int(
|
job_snapshot.last_event_sequence = int(
|
||||||
row["last_event_sequence"]
|
row["last_event_sequence"]
|
||||||
)
|
)
|
||||||
await outbound.put(encode(snapshot))
|
await self._enqueue(outbound, encode(snapshot))
|
||||||
elif (
|
elif (
|
||||||
accepted is not None
|
accepted is not None
|
||||||
and accepted_for_execution
|
and accepted_for_execution
|
||||||
@@ -472,7 +556,7 @@ class ArchiveClientDaemon:
|
|||||||
outbound: asyncio.Queue[str],
|
outbound: asyncio.Queue[str],
|
||||||
command_tasks: set[asyncio.Task[None]],
|
command_tasks: set[asyncio.Task[None]],
|
||||||
) -> None:
|
) -> None:
|
||||||
job_commands: list[control_pb2.Command] = []
|
latest_route_commands: dict[str, control_pb2.Command] = {}
|
||||||
for row in await asyncio.to_thread(self.store.list_accepted_commands):
|
for row in await asyncio.to_thread(self.store.list_accepted_commands):
|
||||||
acknowledgement = decode_message(
|
acknowledgement = decode_message(
|
||||||
str(row["acknowledgement_json"]), control_pb2.CommandAck()
|
str(row["acknowledgement_json"]), control_pb2.CommandAck()
|
||||||
@@ -483,28 +567,20 @@ class ArchiveClientDaemon:
|
|||||||
str(row["command_json"]), control_pb2.Command()
|
str(row["command_json"]), control_pb2.Command()
|
||||||
)
|
)
|
||||||
if command.WhichOneof("payload") == "ensure_route":
|
if command.WhichOneof("payload") == "ensure_route":
|
||||||
self._schedule_route_command(
|
# A newer ensure command for the same route is sufficient to
|
||||||
command, "", outbound, command_tasks
|
# restore the process-local route path and replay its own
|
||||||
)
|
# updates. Replaying every historical ready command causes
|
||||||
elif (
|
# redundant Syncthing configuration after a restart.
|
||||||
command.WhichOneof("payload") in {
|
latest_route_commands[command.ensure_route.route.route_id] = command
|
||||||
"assign_job", "execute_step", "cancel_job"
|
# Job commands are intentionally not replayed here. A command
|
||||||
}
|
# acknowledgement is durable on both sides; the control daemon
|
||||||
and self.jobs is not None
|
# redelivers an unacknowledged command, while registration
|
||||||
):
|
# reconciliation requests snapshots for active jobs. Replaying
|
||||||
if command.WhichOneof("payload") == "cancel_job":
|
# every historical command re-emits overlapping event ranges and
|
||||||
self.jobs.request_cancel(command.cancel_job.job_id)
|
# can flood the single connection's bounded outbound queue.
|
||||||
job_commands.append(command)
|
for command in latest_route_commands.values():
|
||||||
if job_commands:
|
self._schedule_route_command(
|
||||||
task = asyncio.create_task(
|
command, "", outbound, command_tasks, restore_ready=True
|
||||||
self._resume_job_commands(job_commands, outbound),
|
|
||||||
name="resume-job-commands",
|
|
||||||
)
|
|
||||||
command_tasks.add(task)
|
|
||||||
task.add_done_callback(
|
|
||||||
lambda completed: self._command_finished(
|
|
||||||
completed, command_tasks
|
|
||||||
)
|
|
||||||
)
|
)
|
||||||
|
|
||||||
async def _resume_route_commands(
|
async def _resume_route_commands(
|
||||||
@@ -600,34 +676,46 @@ class ArchiveClientDaemon:
|
|||||||
loop = asyncio.get_running_loop()
|
loop = asyncio.get_running_loop()
|
||||||
|
|
||||||
def emit(event):
|
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 = new_envelope()
|
||||||
response.correlation_id = correlation_id
|
response.correlation_id = correlation_id
|
||||||
response.job_event.CopyFrom(event)
|
response.job_event.CopyFrom(event)
|
||||||
future = asyncio.run_coroutine_threadsafe(
|
future = asyncio.run_coroutine_threadsafe(
|
||||||
outbound.put(encode(response)), loop
|
self._enqueue(outbound, encode(response)), loop
|
||||||
)
|
)
|
||||||
try:
|
try:
|
||||||
future.result(timeout=30)
|
future.result(
|
||||||
|
timeout=self.config.connection.outbound_enqueue_timeout
|
||||||
|
+ 1
|
||||||
|
)
|
||||||
except Exception:
|
except Exception:
|
||||||
# The event is already durable in the client DB and will
|
# The event is already durable in the client DB and will
|
||||||
# be replayed after reconnect.
|
# be replayed after reconnect.
|
||||||
pass
|
pass
|
||||||
|
|
||||||
events = await asyncio.to_thread(
|
events = await asyncio.to_thread(
|
||||||
self.jobs.execute, command.execute_step, emit
|
self.jobs.execute, command.execute_step, emit, command.command_id
|
||||||
)
|
)
|
||||||
streamed = True
|
streamed = True
|
||||||
elif payload == "cancel_job":
|
elif payload == "cancel_job":
|
||||||
events = await asyncio.to_thread(
|
events = await asyncio.to_thread(
|
||||||
self.jobs.cancel, command.cancel_job
|
self.jobs.cancel, command.cancel_job, command.command_id
|
||||||
)
|
)
|
||||||
else:
|
else:
|
||||||
raise JobExecutionError("job command payload is unsupported")
|
raise JobExecutionError("job command payload is unsupported")
|
||||||
for event in (() if streamed else events):
|
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 = new_envelope()
|
||||||
response.correlation_id = correlation_id
|
response.correlation_id = correlation_id
|
||||||
response.job_event.CopyFrom(event)
|
response.job_event.CopyFrom(event)
|
||||||
await outbound.put(encode(response))
|
await self._enqueue(outbound, encode(response))
|
||||||
|
|
||||||
def _schedule_route_command(
|
def _schedule_route_command(
|
||||||
self,
|
self,
|
||||||
@@ -635,12 +723,15 @@ class ArchiveClientDaemon:
|
|||||||
correlation_id: str,
|
correlation_id: str,
|
||||||
outbound: asyncio.Queue[str],
|
outbound: asyncio.Queue[str],
|
||||||
command_tasks: set[asyncio.Task[None]],
|
command_tasks: set[asyncio.Task[None]],
|
||||||
|
restore_ready: bool = False,
|
||||||
) -> None:
|
) -> None:
|
||||||
if command.command_id in self._active_route_commands:
|
if command.command_id in self._active_route_commands:
|
||||||
return
|
return
|
||||||
self._active_route_commands.add(command.command_id)
|
self._active_route_commands.add(command.command_id)
|
||||||
task = asyncio.create_task(
|
task = asyncio.create_task(
|
||||||
self._execute_route(command, correlation_id, outbound),
|
self._execute_route(
|
||||||
|
command, correlation_id, outbound, restore_ready=restore_ready
|
||||||
|
),
|
||||||
name=f"route-{command.ensure_route.route.route_id}",
|
name=f"route-{command.ensure_route.route.route_id}",
|
||||||
)
|
)
|
||||||
command_tasks.add(task)
|
command_tasks.add(task)
|
||||||
@@ -671,13 +762,14 @@ class ArchiveClientDaemon:
|
|||||||
response = new_envelope()
|
response = new_envelope()
|
||||||
response.correlation_id = correlation_id
|
response.correlation_id = correlation_id
|
||||||
response.inventory_chunk.CopyFrom(chunk)
|
response.inventory_chunk.CopyFrom(chunk)
|
||||||
await outbound.put(encode(response))
|
await self._enqueue(outbound, encode(response))
|
||||||
|
|
||||||
async def _execute_route(
|
async def _execute_route(
|
||||||
self,
|
self,
|
||||||
command: Any,
|
command: Any,
|
||||||
correlation_id: str,
|
correlation_id: str,
|
||||||
outbound: asyncio.Queue[str],
|
outbound: asyncio.Queue[str],
|
||||||
|
restore_ready: bool = False,
|
||||||
) -> None:
|
) -> None:
|
||||||
assert self.routes is not None
|
assert self.routes is not None
|
||||||
spec = command.ensure_route.route
|
spec = command.ensure_route.route
|
||||||
@@ -696,21 +788,19 @@ class ArchiveClientDaemon:
|
|||||||
response.route_update.CopyFrom(
|
response.route_update.CopyFrom(
|
||||||
decode_message(str(row["update_json"]), control_pb2.RouteUpdate())
|
decode_message(str(row["update_json"]), control_pb2.RouteUpdate())
|
||||||
)
|
)
|
||||||
await outbound.put(encode(response))
|
await self._enqueue(outbound, encode(response))
|
||||||
if attempt["state"] == "ready":
|
if attempt["state"] == "ready":
|
||||||
# Route readiness is durable, but this process-local lookup is
|
if restore_ready:
|
||||||
# not. Reconfigure (which validates the existing Syncthing folder
|
# Route readiness is durable, but this process-local lookup is
|
||||||
# and recreates its local directory if necessary) before replaying
|
# not. Reconfigure (which validates the existing Syncthing
|
||||||
# the stored READY updates. This is essential after a daemon or
|
# folder and recreates its local directory if necessary) after
|
||||||
# mount restart: a previously ready route may otherwise point at a
|
# a daemon restart, before replaying stored READY updates.
|
||||||
# path that is no longer present, and a later source-stage command
|
configured = await asyncio.to_thread(
|
||||||
# would fail with a bare ENOENT.
|
self.routes.configure,
|
||||||
configured = await asyncio.to_thread(
|
spec,
|
||||||
self.routes.configure,
|
time.monotonic() + spec.setup_timeout_seconds,
|
||||||
spec,
|
)
|
||||||
time.monotonic() + spec.setup_timeout_seconds,
|
self._known_route_paths[spec.route_id] = configured.local_path
|
||||||
)
|
|
||||||
self._known_route_paths[spec.route_id] = configured.local_path
|
|
||||||
return
|
return
|
||||||
if attempt["state"] == "failed":
|
if attempt["state"] == "failed":
|
||||||
return
|
return
|
||||||
@@ -863,7 +953,7 @@ class ArchiveClientDaemon:
|
|||||||
response = new_envelope()
|
response = new_envelope()
|
||||||
response.correlation_id = correlation_id
|
response.correlation_id = correlation_id
|
||||||
response.route_update.CopyFrom(update)
|
response.route_update.CopyFrom(update)
|
||||||
await outbound.put(encode(response))
|
await self._enqueue(outbound, encode(response))
|
||||||
|
|
||||||
def _local_route(
|
def _local_route(
|
||||||
self, spec: route_pb2.EnsureRouteSpec, state: int
|
self, spec: route_pb2.EnsureRouteSpec, state: int
|
||||||
@@ -990,6 +1080,17 @@ class ArchiveClientDaemon:
|
|||||||
acknowledgement.error.message = "job assignment is invalid"
|
acknowledgement.error.message = "job assignment is invalid"
|
||||||
else:
|
else:
|
||||||
acknowledgement.status = control_pb2.COMMAND_ACK_STATUS_ACCEPTED
|
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":
|
elif command.WhichOneof("payload") == "execute_step":
|
||||||
step = command.execute_step
|
step = command.execute_step
|
||||||
if self.jobs is None:
|
if self.jobs is None:
|
||||||
|
|||||||
+92
-18
@@ -97,33 +97,25 @@ class ClientJobExecutor:
|
|||||||
) -> list[control_pb2.JobEvent]:
|
) -> list[control_pb2.JobEvent]:
|
||||||
definition = command.job
|
definition = command.job
|
||||||
self._validate_definition(definition)
|
self._validate_definition(definition)
|
||||||
replay = self._replay(
|
self.store.ensure_job_definition(
|
||||||
definition.job_id, command.expected_last_event_sequence
|
definition.job_id, encode_message(definition)
|
||||||
)
|
)
|
||||||
if replay:
|
# Assignment succeeds through CommandAck. It must not let a client
|
||||||
return replay
|
# allocate a globally ordered JobEvent cursor.
|
||||||
event = self._event(
|
return []
|
||||||
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]
|
|
||||||
|
|
||||||
def execute(
|
def execute(
|
||||||
self,
|
self,
|
||||||
command: control_pb2.ExecuteStepCommand,
|
command: control_pb2.ExecuteStepCommand,
|
||||||
event_callback: Callable[[control_pb2.JobEvent], None] | None = None,
|
event_callback: Callable[[control_pb2.JobEvent], None] | None = None,
|
||||||
|
command_id: str = "",
|
||||||
) -> list[control_pb2.JobEvent]:
|
) -> list[control_pb2.JobEvent]:
|
||||||
# asyncio cancellation of a connection-bound task cannot stop the
|
# asyncio cancellation of a connection-bound task cannot stop the
|
||||||
# synchronous filesystem/qB operation already running in its worker
|
# synchronous filesystem/qB operation already running in its worker
|
||||||
# thread. A replay after reconnect therefore waits for that operation
|
# thread. A replay after reconnect therefore waits for that operation
|
||||||
# and then reads its durable event journal instead of executing twice.
|
# and then reads its durable event journal instead of executing twice.
|
||||||
with self._execution_lock(command.job_id):
|
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:
|
def _execution_lock(self, job_id: str) -> threading.Lock:
|
||||||
with self._execution_locks_guard:
|
with self._execution_locks_guard:
|
||||||
@@ -133,6 +125,7 @@ class ClientJobExecutor:
|
|||||||
self,
|
self,
|
||||||
command: control_pb2.ExecuteStepCommand,
|
command: control_pb2.ExecuteStepCommand,
|
||||||
event_callback: Callable[[control_pb2.JobEvent], None] | None = None,
|
event_callback: Callable[[control_pb2.JobEvent], None] | None = None,
|
||||||
|
command_id: str = "",
|
||||||
) -> list[control_pb2.JobEvent]:
|
) -> list[control_pb2.JobEvent]:
|
||||||
definition = self._definition(command.job_id)
|
definition = self._definition(command.job_id)
|
||||||
replay = self._replay(
|
replay = self._replay(
|
||||||
@@ -186,6 +179,7 @@ class ClientJobExecutor:
|
|||||||
),
|
),
|
||||||
step=command.step,
|
step=command.step,
|
||||||
step_state=job_pb2.STEP_STATE_RUNNING,
|
step_state=job_pb2.STEP_STATE_RUNNING,
|
||||||
|
command_id=command_id,
|
||||||
)
|
)
|
||||||
self._record(definition, started)
|
self._record(definition, started)
|
||||||
cursor = started
|
cursor = started
|
||||||
@@ -228,6 +222,7 @@ class ClientJobExecutor:
|
|||||||
committed=cursor.committed,
|
committed=cursor.committed,
|
||||||
step=command.step,
|
step=command.step,
|
||||||
step_state=job_pb2.STEP_STATE_RUNNING,
|
step_state=job_pb2.STEP_STATE_RUNNING,
|
||||||
|
command_id=command_id,
|
||||||
)
|
)
|
||||||
event.progress.fraction_complete = max(
|
event.progress.fraction_complete = max(
|
||||||
0.0, min(float(fraction), 1.0)
|
0.0, min(float(fraction), 1.0)
|
||||||
@@ -258,6 +253,7 @@ class ClientJobExecutor:
|
|||||||
committed=started.committed,
|
committed=started.committed,
|
||||||
step=command.step,
|
step=command.step,
|
||||||
step_state=job_pb2.STEP_STATE_CANCELLED,
|
step_state=job_pb2.STEP_STATE_CANCELLED,
|
||||||
|
command_id=command_id,
|
||||||
)
|
)
|
||||||
cancelling.error.code = common_pb2.ERROR_CODE_CANCELLED
|
cancelling.error.code = common_pb2.ERROR_CODE_CANCELLED
|
||||||
cancelling.error.message = str(error)
|
cancelling.error.message = str(error)
|
||||||
@@ -315,6 +311,7 @@ class ClientJobExecutor:
|
|||||||
committed=started.committed,
|
committed=started.committed,
|
||||||
step=command.step,
|
step=command.step,
|
||||||
step_state=job_pb2.STEP_STATE_FAILED,
|
step_state=job_pb2.STEP_STATE_FAILED,
|
||||||
|
command_id=command_id,
|
||||||
)
|
)
|
||||||
failed.error.code = _job_error_code(error)
|
failed.error.code = _job_error_code(error)
|
||||||
failed.error.message = str(error) or type(error).__name__
|
failed.error.message = str(error) or type(error).__name__
|
||||||
@@ -357,6 +354,7 @@ class ClientJobExecutor:
|
|||||||
committed=committed,
|
committed=committed,
|
||||||
step=command.step,
|
step=command.step,
|
||||||
step_state=job_pb2.STEP_STATE_SUCCEEDED,
|
step_state=job_pb2.STEP_STATE_SUCCEEDED,
|
||||||
|
command_id=command_id,
|
||||||
)
|
)
|
||||||
if result is not None:
|
if result is not None:
|
||||||
succeeded.observed_placement.CopyFrom(result)
|
succeeded.observed_placement.CopyFrom(result)
|
||||||
@@ -367,7 +365,7 @@ class ClientJobExecutor:
|
|||||||
return emitted
|
return emitted
|
||||||
|
|
||||||
def cancel(
|
def cancel(
|
||||||
self, command: control_pb2.CancelJobCommand
|
self, command: control_pb2.CancelJobCommand, command_id: str = ""
|
||||||
) -> list[control_pb2.JobEvent]:
|
) -> list[control_pb2.JobEvent]:
|
||||||
definition = self._definition(command.job_id)
|
definition = self._definition(command.job_id)
|
||||||
self.request_cancel(command.job_id)
|
self.request_cancel(command.job_id)
|
||||||
@@ -418,6 +416,7 @@ class ClientJobExecutor:
|
|||||||
committed=committed,
|
committed=committed,
|
||||||
step=job_pb2.JOB_STEP_KIND_ROLLBACK,
|
step=job_pb2.JOB_STEP_KIND_ROLLBACK,
|
||||||
step_state=job_pb2.STEP_STATE_SUCCEEDED,
|
step_state=job_pb2.STEP_STATE_SUCCEEDED,
|
||||||
|
command_id=command_id,
|
||||||
)
|
)
|
||||||
if placement is not None:
|
if placement is not None:
|
||||||
event.observed_placement.CopyFrom(placement)
|
event.observed_placement.CopyFrom(placement)
|
||||||
@@ -809,13 +808,14 @@ class ClientJobExecutor:
|
|||||||
fraction, 0, 0, "qBittorrent stopped recheck"
|
fraction, 0, 0, "qBittorrent stopped recheck"
|
||||||
),
|
),
|
||||||
)
|
)
|
||||||
if should_start:
|
|
||||||
self.qbittorrent.start(qb_torrent_id)
|
|
||||||
verified = self.qbittorrent.get_resource(info_hash)
|
verified = self.qbittorrent.get_resource(info_hash)
|
||||||
if verified is None:
|
if verified is None:
|
||||||
raise JobExecutionError(
|
raise JobExecutionError(
|
||||||
"verified qBittorrent resource disappeared"
|
"verified qBittorrent resource disappeared"
|
||||||
)
|
)
|
||||||
|
_normalize_verified_resource_permissions(self.qb_root, verified)
|
||||||
|
if should_start:
|
||||||
|
self.qbittorrent.start(qb_torrent_id)
|
||||||
placement = resource_pb2.Placement(
|
placement = resource_pb2.Placement(
|
||||||
client_id=self.client_id,
|
client_id=self.client_id,
|
||||||
state=resource_pb2.PLACEMENT_STATE_PRESENT,
|
state=resource_pb2.PLACEMENT_STATE_PRESENT,
|
||||||
@@ -1050,6 +1050,7 @@ class ClientJobExecutor:
|
|||||||
committed: bool,
|
committed: bool,
|
||||||
step: int = job_pb2.JOB_STEP_KIND_UNSPECIFIED,
|
step: int = job_pb2.JOB_STEP_KIND_UNSPECIFIED,
|
||||||
step_state: int = job_pb2.STEP_STATE_UNSPECIFIED,
|
step_state: int = job_pb2.STEP_STATE_UNSPECIFIED,
|
||||||
|
command_id: str = "",
|
||||||
) -> control_pb2.JobEvent:
|
) -> control_pb2.JobEvent:
|
||||||
event = control_pb2.JobEvent(
|
event = control_pb2.JobEvent(
|
||||||
event_id=str(uuid.uuid4()),
|
event_id=str(uuid.uuid4()),
|
||||||
@@ -1059,6 +1060,7 @@ class ClientJobExecutor:
|
|||||||
type=event_type,
|
type=event_type,
|
||||||
state=state,
|
state=state,
|
||||||
committed=committed,
|
committed=committed,
|
||||||
|
command_id=command_id,
|
||||||
)
|
)
|
||||||
event.occurred_at.GetCurrentTime()
|
event.occurred_at.GetCurrentTime()
|
||||||
if step != job_pb2.JOB_STEP_KIND_UNSPECIFIED:
|
if step != job_pb2.JOB_STEP_KIND_UNSPECIFIED:
|
||||||
@@ -1190,6 +1192,78 @@ def _info_hash(definition: job_pb2.JobDefinition) -> str:
|
|||||||
return value
|
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(
|
def _fingerprint_matches(
|
||||||
current: resource_pb2.ResourceStateFingerprint,
|
current: resource_pb2.ResourceStateFingerprint,
|
||||||
expected: resource_pb2.ResourceStateFingerprint,
|
expected: resource_pb2.ResourceStateFingerprint,
|
||||||
|
|||||||
@@ -18,7 +18,7 @@ class ProtocolError(ValueError):
|
|||||||
|
|
||||||
def new_envelope() -> envelope_pb2.Envelope:
|
def new_envelope() -> envelope_pb2.Envelope:
|
||||||
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.message_id = str(uuid.uuid4())
|
||||||
envelope.sent_at.FromDatetime(datetime.now(timezone.utc))
|
envelope.sent_at.FromDatetime(datetime.now(timezone.utc))
|
||||||
return envelope
|
return envelope
|
||||||
@@ -59,7 +59,7 @@ def decode(data: str | bytes, max_bytes: int = 1024 * 1024) -> envelope_pb2.Enve
|
|||||||
json_format.ParseError,
|
json_format.ParseError,
|
||||||
) as exc:
|
) as exc:
|
||||||
raise ProtocolError("invalid control envelope") from 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")
|
raise ProtocolError("unsupported protocol major version")
|
||||||
_canonical_uuid(envelope.message_id, "message_id")
|
_canonical_uuid(envelope.message_id, "message_id")
|
||||||
if envelope.correlation_id:
|
if envelope.correlation_id:
|
||||||
|
|||||||
@@ -38,6 +38,7 @@ class FileOperationConflict(RuntimeError):
|
|||||||
class CommandAcceptance:
|
class CommandAcceptance:
|
||||||
duplicate: bool
|
duplicate: bool
|
||||||
acknowledgement_json: str
|
acknowledgement_json: str
|
||||||
|
state: str
|
||||||
|
|
||||||
|
|
||||||
class ClientStore:
|
class ClientStore:
|
||||||
@@ -85,7 +86,7 @@ class ClientStore:
|
|||||||
connection.execute("BEGIN IMMEDIATE")
|
connection.execute("BEGIN IMMEDIATE")
|
||||||
existing = connection.execute(
|
existing = connection.execute(
|
||||||
"""
|
"""
|
||||||
SELECT payload_sha256, command_json, acknowledgement_json
|
SELECT payload_sha256, command_json, acknowledgement_json, state
|
||||||
FROM commands WHERE command_id = ?
|
FROM commands WHERE command_id = ?
|
||||||
""",
|
""",
|
||||||
(command_id,),
|
(command_id,),
|
||||||
@@ -93,17 +94,95 @@ class ClientStore:
|
|||||||
if existing:
|
if existing:
|
||||||
if existing["payload_sha256"] != digest or existing["command_json"] != payload:
|
if existing["payload_sha256"] != digest or existing["command_json"] != payload:
|
||||||
raise CommandConflict("command ID was reused with different content")
|
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(
|
connection.execute(
|
||||||
"""
|
"""
|
||||||
INSERT INTO commands (
|
INSERT INTO commands (
|
||||||
command_id, payload_sha256, command_json,
|
command_id, payload_sha256, command_json,
|
||||||
acknowledgement_json, state
|
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]]:
|
def list_active_job_cursors(self) -> list[dict[str, object]]:
|
||||||
with self._connect() as connection:
|
with self._connect() as connection:
|
||||||
@@ -122,7 +201,7 @@ class ClientStore:
|
|||||||
rows = connection.execute(
|
rows = connection.execute(
|
||||||
"""
|
"""
|
||||||
SELECT command_id, command_json, acknowledgement_json
|
SELECT command_id, command_json, acknowledgement_json
|
||||||
FROM commands ORDER BY rowid
|
FROM commands WHERE state = 'accepted' ORDER BY rowid
|
||||||
"""
|
"""
|
||||||
).fetchall()
|
).fetchall()
|
||||||
return [dict(row) for row in rows]
|
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(
|
def begin_file_operation(
|
||||||
self,
|
self,
|
||||||
operation_id: str,
|
operation_id: str,
|
||||||
@@ -591,6 +778,18 @@ class ClientStore:
|
|||||||
connection.close()
|
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:
|
def _canonical(value: object) -> str:
|
||||||
return json.dumps(value, sort_keys=True, separators=(",", ":"), ensure_ascii=False)
|
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."""
|
"""Generated archive_control.v1 bindings."""
|
||||||
|
|
||||||
|
|||||||
@@ -1,4 +1,4 @@
|
|||||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||||
# -*- coding: utf-8 -*-
|
# -*- coding: utf-8 -*-
|
||||||
# Generated by the protocol buffer compiler. DO NOT EDIT!
|
# Generated by the protocol buffer compiler. DO NOT EDIT!
|
||||||
# NO CHECKED-IN PROTOBUF GENCODE
|
# NO CHECKED-IN PROTOBUF GENCODE
|
||||||
|
|||||||
@@ -1,4 +1,4 @@
|
|||||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||||
import datetime
|
import datetime
|
||||||
|
|
||||||
from archive_control.v1 import common_pb2 as _common_pb2
|
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 -*-
|
# -*- coding: utf-8 -*-
|
||||||
# Generated by the protocol buffer compiler. DO NOT EDIT!
|
# Generated by the protocol buffer compiler. DO NOT EDIT!
|
||||||
# NO CHECKED-IN PROTOBUF GENCODE
|
# NO CHECKED-IN PROTOBUF GENCODE
|
||||||
|
|||||||
@@ -1,4 +1,4 @@
|
|||||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||||
import datetime
|
import datetime
|
||||||
|
|
||||||
from google.protobuf import timestamp_pb2 as _timestamp_pb2
|
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
|
import datetime
|
||||||
|
|
||||||
from archive_control.v1 import common_pb2 as _common_pb2
|
from archive_control.v1 import common_pb2 as _common_pb2
|
||||||
@@ -108,8 +108,20 @@ class RequestJobSnapshotCommand(_message.Message):
|
|||||||
job_ids: _containers.RepeatedScalarFieldContainer[str]
|
job_ids: _containers.RepeatedScalarFieldContainer[str]
|
||||||
def __init__(self, job_ids: _Optional[_Iterable[str]] = ...) -> None: ...
|
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):
|
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]
|
COMMAND_ID_FIELD_NUMBER: _ClassVar[int]
|
||||||
CREATED_AT_FIELD_NUMBER: _ClassVar[int]
|
CREATED_AT_FIELD_NUMBER: _ClassVar[int]
|
||||||
ASSIGN_JOB_FIELD_NUMBER: _ClassVar[int]
|
ASSIGN_JOB_FIELD_NUMBER: _ClassVar[int]
|
||||||
@@ -118,6 +130,7 @@ class Command(_message.Message):
|
|||||||
ENSURE_ROUTE_FIELD_NUMBER: _ClassVar[int]
|
ENSURE_ROUTE_FIELD_NUMBER: _ClassVar[int]
|
||||||
INVENTORY_QUERY_FIELD_NUMBER: _ClassVar[int]
|
INVENTORY_QUERY_FIELD_NUMBER: _ClassVar[int]
|
||||||
REQUEST_JOB_SNAPSHOT_FIELD_NUMBER: _ClassVar[int]
|
REQUEST_JOB_SNAPSHOT_FIELD_NUMBER: _ClassVar[int]
|
||||||
|
RECONCILE_JOB_FIELD_NUMBER: _ClassVar[int]
|
||||||
command_id: str
|
command_id: str
|
||||||
created_at: _timestamp_pb2.Timestamp
|
created_at: _timestamp_pb2.Timestamp
|
||||||
assign_job: AssignJobCommand
|
assign_job: AssignJobCommand
|
||||||
@@ -126,7 +139,8 @@ class Command(_message.Message):
|
|||||||
ensure_route: EnsureRouteCommand
|
ensure_route: EnsureRouteCommand
|
||||||
inventory_query: _inventory_pb2.InventoryQuery
|
inventory_query: _inventory_pb2.InventoryQuery
|
||||||
request_job_snapshot: RequestJobSnapshotCommand
|
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):
|
class CommandAck(_message.Message):
|
||||||
__slots__ = ("command_id", "status", "error")
|
__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: ...
|
def __init__(self, command_id: _Optional[str] = ..., status: _Optional[_Union[CommandAckStatus, str]] = ..., error: _Optional[_Union[_common_pb2.Error, _Mapping]] = ...) -> None: ...
|
||||||
|
|
||||||
class JobEvent(_message.Message):
|
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]
|
EVENT_ID_FIELD_NUMBER: _ClassVar[int]
|
||||||
JOB_ID_FIELD_NUMBER: _ClassVar[int]
|
JOB_ID_FIELD_NUMBER: _ClassVar[int]
|
||||||
SEQUENCE_FIELD_NUMBER: _ClassVar[int]
|
SEQUENCE_FIELD_NUMBER: _ClassVar[int]
|
||||||
@@ -152,6 +166,7 @@ class JobEvent(_message.Message):
|
|||||||
OBSERVED_RESOURCE_FIELD_NUMBER: _ClassVar[int]
|
OBSERVED_RESOURCE_FIELD_NUMBER: _ClassVar[int]
|
||||||
OBSERVED_PLACEMENT_FIELD_NUMBER: _ClassVar[int]
|
OBSERVED_PLACEMENT_FIELD_NUMBER: _ClassVar[int]
|
||||||
OCCURRED_AT_FIELD_NUMBER: _ClassVar[int]
|
OCCURRED_AT_FIELD_NUMBER: _ClassVar[int]
|
||||||
|
COMMAND_ID_FIELD_NUMBER: _ClassVar[int]
|
||||||
event_id: str
|
event_id: str
|
||||||
job_id: str
|
job_id: str
|
||||||
sequence: int
|
sequence: int
|
||||||
@@ -164,7 +179,8 @@ class JobEvent(_message.Message):
|
|||||||
observed_resource: _resource_pb2.ResourceStateFingerprint
|
observed_resource: _resource_pb2.ResourceStateFingerprint
|
||||||
observed_placement: _resource_pb2.Placement
|
observed_placement: _resource_pb2.Placement
|
||||||
occurred_at: _timestamp_pb2.Timestamp
|
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):
|
class JobSnapshot(_message.Message):
|
||||||
__slots__ = ("job", "last_event_sequence", "in_flight_command_ids")
|
__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 -*-
|
# -*- coding: utf-8 -*-
|
||||||
# Generated by the protocol buffer compiler. DO NOT EDIT!
|
# Generated by the protocol buffer compiler. DO NOT EDIT!
|
||||||
# NO CHECKED-IN PROTOBUF GENCODE
|
# NO CHECKED-IN PROTOBUF GENCODE
|
||||||
|
|||||||
@@ -1,4 +1,4 @@
|
|||||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||||
import datetime
|
import datetime
|
||||||
|
|
||||||
from archive_control.v1 import client_pb2 as _client_pb2
|
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 -*-
|
# -*- coding: utf-8 -*-
|
||||||
# Generated by the protocol buffer compiler. DO NOT EDIT!
|
# Generated by the protocol buffer compiler. DO NOT EDIT!
|
||||||
# NO CHECKED-IN PROTOBUF GENCODE
|
# 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 common_pb2 as _common_pb2
|
||||||
from archive_control.v1 import resource_pb2 as _resource_pb2
|
from archive_control.v1 import resource_pb2 as _resource_pb2
|
||||||
from google.protobuf.internal import containers as _containers
|
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 -*-
|
# -*- coding: utf-8 -*-
|
||||||
# Generated by the protocol buffer compiler. DO NOT EDIT!
|
# Generated by the protocol buffer compiler. DO NOT EDIT!
|
||||||
# NO CHECKED-IN PROTOBUF GENCODE
|
# NO CHECKED-IN PROTOBUF GENCODE
|
||||||
|
|||||||
@@ -1,4 +1,4 @@
|
|||||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||||
import datetime
|
import datetime
|
||||||
|
|
||||||
from archive_control.v1 import common_pb2 as _common_pb2
|
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 -*-
|
# -*- coding: utf-8 -*-
|
||||||
# Generated by the protocol buffer compiler. DO NOT EDIT!
|
# Generated by the protocol buffer compiler. DO NOT EDIT!
|
||||||
# NO CHECKED-IN PROTOBUF GENCODE
|
# NO CHECKED-IN PROTOBUF GENCODE
|
||||||
|
|||||||
@@ -1,4 +1,4 @@
|
|||||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||||
import datetime
|
import datetime
|
||||||
|
|
||||||
from google.protobuf import timestamp_pb2 as _timestamp_pb2
|
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 -*-
|
# -*- coding: utf-8 -*-
|
||||||
# Generated by the protocol buffer compiler. DO NOT EDIT!
|
# Generated by the protocol buffer compiler. DO NOT EDIT!
|
||||||
# NO CHECKED-IN PROTOBUF GENCODE
|
# NO CHECKED-IN PROTOBUF GENCODE
|
||||||
|
|||||||
@@ -1,4 +1,4 @@
|
|||||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||||
import datetime
|
import datetime
|
||||||
|
|
||||||
from google.protobuf import timestamp_pb2 as _timestamp_pb2
|
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 -*-
|
# -*- coding: utf-8 -*-
|
||||||
# Generated by the protocol buffer compiler. DO NOT EDIT!
|
# Generated by the protocol buffer compiler. DO NOT EDIT!
|
||||||
# NO CHECKED-IN PROTOBUF GENCODE
|
# NO CHECKED-IN PROTOBUF GENCODE
|
||||||
|
|||||||
@@ -1,4 +1,4 @@
|
|||||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||||
import datetime
|
import datetime
|
||||||
|
|
||||||
from archive_control.v1 import resource_pb2 as _resource_pb2
|
from archive_control.v1 import resource_pb2 as _resource_pb2
|
||||||
|
|||||||
+119
-2
@@ -74,7 +74,7 @@ class DaemonTransportTests(unittest.IsolatedAsyncioTestCase):
|
|||||||
response.register_response.status = (
|
response.register_response.status = (
|
||||||
client_pb2.REGISTRATION_STATUS_ACCEPTED
|
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.send(encode(response))
|
||||||
# Deliberately keep TCP/WebSocket open but send no application
|
# Deliberately keep TCP/WebSocket open but send no application
|
||||||
# heartbeats. This models a stale proxy/server-side session.
|
# heartbeats. This models a stale proxy/server-side session.
|
||||||
@@ -110,6 +110,123 @@ class DaemonTransportTests(unittest.IsolatedAsyncioTestCase):
|
|||||||
):
|
):
|
||||||
await asyncio.wait_for(daemon._connection(), 1)
|
await asyncio.wait_for(daemon._connection(), 1)
|
||||||
|
|
||||||
|
async def test_failed_writer_ends_connection_without_waiting_for_heartbeat(self):
|
||||||
|
"""A send failure must immediately reach the reconnect supervisor."""
|
||||||
|
|
||||||
|
async def control(websocket):
|
||||||
|
registration = decode(await websocket.recv())
|
||||||
|
response = new_envelope()
|
||||||
|
response.correlation_id = registration.message_id
|
||||||
|
response.register_response.status = (
|
||||||
|
client_pb2.REGISTRATION_STATUS_ACCEPTED
|
||||||
|
)
|
||||||
|
response.register_response.negotiated_version.major = 2
|
||||||
|
await websocket.send(encode(response))
|
||||||
|
await websocket.wait_closed()
|
||||||
|
|
||||||
|
with tempfile.TemporaryDirectory() as directory:
|
||||||
|
root = Path(directory)
|
||||||
|
token = root / "token"
|
||||||
|
token.write_text("shared-secret", encoding="utf-8")
|
||||||
|
os.chmod(token, 0o600)
|
||||||
|
async with serve(control, "127.0.0.1", 0, ping_interval=None) as server:
|
||||||
|
port = server.sockets[0].getsockname()[1]
|
||||||
|
service = ServiceConfig("http://local", PurePosixPath("/api"), root)
|
||||||
|
config = ClientConfig(
|
||||||
|
"cache-1", "Cache 1", "cache",
|
||||||
|
f"ws://127.0.0.1:{port}", token,
|
||||||
|
root / "state.db", root / "backups", service, service,
|
||||||
|
)
|
||||||
|
probe = FilesystemProbe(root, True, True, True, True, True)
|
||||||
|
daemon = ArchiveClientDaemon(config, [probe, probe], [])
|
||||||
|
await asyncio.to_thread(daemon.store.initialize)
|
||||||
|
|
||||||
|
async def failed_writer(websocket, outbound):
|
||||||
|
raise OSError("simulated broken socket")
|
||||||
|
|
||||||
|
daemon._writer = failed_writer
|
||||||
|
with self.assertRaisesRegex(RuntimeError, "writer failed"):
|
||||||
|
await asyncio.wait_for(daemon._connection(), 1)
|
||||||
|
|
||||||
|
async def test_full_outbound_queue_aborts_connection_instead_of_blocking_heartbeats(self):
|
||||||
|
"""Bulk output cannot indefinitely block the receive/heartbeat loop."""
|
||||||
|
|
||||||
|
async def control(websocket):
|
||||||
|
registration = decode(await websocket.recv())
|
||||||
|
response = new_envelope()
|
||||||
|
response.correlation_id = registration.message_id
|
||||||
|
response.register_response.status = (
|
||||||
|
client_pb2.REGISTRATION_STATUS_ACCEPTED
|
||||||
|
)
|
||||||
|
response.register_response.negotiated_version.major = 2
|
||||||
|
await websocket.send(encode(response))
|
||||||
|
for sequence in range(1, 102):
|
||||||
|
heartbeat = new_envelope()
|
||||||
|
heartbeat.heartbeat.sequence = sequence
|
||||||
|
await websocket.send(encode(heartbeat))
|
||||||
|
await websocket.wait_closed()
|
||||||
|
|
||||||
|
with tempfile.TemporaryDirectory() as directory:
|
||||||
|
root = Path(directory)
|
||||||
|
token = root / "token"
|
||||||
|
token.write_text("shared-secret", encoding="utf-8")
|
||||||
|
os.chmod(token, 0o600)
|
||||||
|
async with serve(control, "127.0.0.1", 0, ping_interval=None) as server:
|
||||||
|
port = server.sockets[0].getsockname()[1]
|
||||||
|
service = ServiceConfig("http://local", PurePosixPath("/api"), root)
|
||||||
|
config = ClientConfig(
|
||||||
|
"cache-1", "Cache 1", "cache",
|
||||||
|
f"ws://127.0.0.1:{port}", token,
|
||||||
|
root / "state.db", root / "backups", service, service,
|
||||||
|
ConnectionConfig(outbound_enqueue_timeout=0.01),
|
||||||
|
)
|
||||||
|
probe = FilesystemProbe(root, True, True, True, True, True)
|
||||||
|
daemon = ArchiveClientDaemon(config, [probe, probe], [])
|
||||||
|
await asyncio.to_thread(daemon.store.initialize)
|
||||||
|
|
||||||
|
async def stopped_writer(websocket, outbound):
|
||||||
|
await asyncio.Event().wait()
|
||||||
|
|
||||||
|
daemon._writer = stopped_writer
|
||||||
|
with self.assertRaisesRegex(RuntimeError, "outbound queue is blocked"):
|
||||||
|
await asyncio.wait_for(daemon._connection(), 2)
|
||||||
|
|
||||||
|
async def test_reconnect_resume_skips_historical_job_commands(self):
|
||||||
|
"""Registration reconciliation, not command replay, recovers job state."""
|
||||||
|
|
||||||
|
with tempfile.TemporaryDirectory() as directory:
|
||||||
|
root = Path(directory)
|
||||||
|
token = root / "token"
|
||||||
|
token.write_text("shared-secret", encoding="utf-8")
|
||||||
|
os.chmod(token, 0o600)
|
||||||
|
service = ServiceConfig("http://local", PurePosixPath("/api"), root)
|
||||||
|
config = ClientConfig(
|
||||||
|
"cache-1", "Cache 1", "cache", "ws://control", token,
|
||||||
|
root / "state.db", root / "backups", service, service,
|
||||||
|
)
|
||||||
|
probe = FilesystemProbe(root, True, True, True, True, True)
|
||||||
|
daemon = ArchiveClientDaemon(config, [probe, probe], [])
|
||||||
|
await asyncio.to_thread(daemon.store.initialize)
|
||||||
|
command = control_pb2.Command(command_id=str(uuid4()))
|
||||||
|
command.execute_step.job_id = str(uuid4())
|
||||||
|
command.execute_step.expected_last_event_sequence = 1
|
||||||
|
acknowledgement = control_pb2.CommandAck(
|
||||||
|
command_id=command.command_id,
|
||||||
|
status=control_pb2.COMMAND_ACK_STATUS_ACCEPTED,
|
||||||
|
)
|
||||||
|
await asyncio.to_thread(
|
||||||
|
daemon.store.accept_command,
|
||||||
|
command.command_id,
|
||||||
|
encode_message(command),
|
||||||
|
encode_message(acknowledgement),
|
||||||
|
)
|
||||||
|
daemon.jobs = Mock()
|
||||||
|
outbound = asyncio.Queue()
|
||||||
|
tasks = set()
|
||||||
|
await daemon._resume_commands(outbound, tasks)
|
||||||
|
self.assertEqual(tasks, set())
|
||||||
|
self.assertTrue(outbound.empty())
|
||||||
|
|
||||||
async def test_eviction_assignment_and_steps_are_admitted(self):
|
async def test_eviction_assignment_and_steps_are_admitted(self):
|
||||||
with tempfile.TemporaryDirectory() as directory:
|
with tempfile.TemporaryDirectory() as directory:
|
||||||
root = Path(directory)
|
root = Path(directory)
|
||||||
@@ -368,7 +485,7 @@ class DaemonTransportTests(unittest.IsolatedAsyncioTestCase):
|
|||||||
response = new_envelope()
|
response = new_envelope()
|
||||||
response.correlation_id = registration.message_id
|
response.correlation_id = registration.message_id
|
||||||
response.register_response.status = client_pb2.REGISTRATION_STATUS_ACCEPTED
|
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.send(encode(response))
|
||||||
heartbeat = new_envelope()
|
heartbeat = new_envelope()
|
||||||
heartbeat.heartbeat.sequence = 7
|
heartbeat.heartbeat.sequence = 7
|
||||||
|
|||||||
+56
-9
@@ -1,4 +1,6 @@
|
|||||||
import hashlib
|
import hashlib
|
||||||
|
import os
|
||||||
|
import stat
|
||||||
import tempfile
|
import tempfile
|
||||||
import threading
|
import threading
|
||||||
import unittest
|
import unittest
|
||||||
@@ -10,12 +12,13 @@ from archive_clients.bencode import encode
|
|||||||
from archive_clients.jobs import (
|
from archive_clients.jobs import (
|
||||||
ClientJobExecutor,
|
ClientJobExecutor,
|
||||||
JobExecutionError,
|
JobExecutionError,
|
||||||
|
_normalize_verified_resource_permissions,
|
||||||
_resource_fingerprint,
|
_resource_fingerprint,
|
||||||
)
|
)
|
||||||
from archive_clients.syncthing import RouteSetupError, SyncthingTransferStatus
|
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_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:
|
class CompleteSyncthing:
|
||||||
@@ -39,6 +42,50 @@ class SlowRescanSyncthing(CompleteSyncthing):
|
|||||||
|
|
||||||
|
|
||||||
class ClientJobHappyPathTests(unittest.TestCase):
|
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):
|
def test_syncthing_api_outage_during_transfer_is_retried(self):
|
||||||
"""A transient local REST outage must not terminally fail the job."""
|
"""A transient local REST outage must not terminally fail the job."""
|
||||||
with tempfile.TemporaryDirectory() as directory:
|
with tempfile.TemporaryDirectory() as directory:
|
||||||
@@ -436,19 +483,19 @@ class ClientJobHappyPathTests(unittest.TestCase):
|
|||||||
|
|
||||||
source_assigned = source.assign(control_pb2.AssignJobCommand(
|
source_assigned = source.assign(control_pb2.AssignJobCommand(
|
||||||
job=definition,
|
job=definition,
|
||||||
expected_job_revision=1,
|
expected_job_revision=0,
|
||||||
expected_last_event_sequence=0,
|
expected_last_event_sequence=0,
|
||||||
))
|
))
|
||||||
target_assigned = target.assign(control_pb2.AssignJobCommand(
|
target_assigned = target.assign(control_pb2.AssignJobCommand(
|
||||||
job=definition,
|
job=definition,
|
||||||
expected_job_revision=1,
|
expected_job_revision=0,
|
||||||
expected_last_event_sequence=1,
|
expected_last_event_sequence=0,
|
||||||
))
|
))
|
||||||
self.assertEqual(source_assigned[0].sequence, 1)
|
self.assertEqual(source_assigned, [])
|
||||||
self.assertEqual(target_assigned[0].sequence, 2)
|
self.assertEqual(target_assigned, [])
|
||||||
|
|
||||||
cursor_revision = 1
|
cursor_revision = 0
|
||||||
cursor_sequence = 2
|
cursor_sequence = 0
|
||||||
pipeline = (
|
pipeline = (
|
||||||
(source, job_pb2.JOB_STEP_KIND_SOURCE_STAGE),
|
(source, job_pb2.JOB_STEP_KIND_SOURCE_STAGE),
|
||||||
(target, job_pb2.JOB_STEP_KIND_SYNCTHING_TRANSFER),
|
(target, job_pb2.JOB_STEP_KIND_SYNCTHING_TRANSFER),
|
||||||
|
|||||||
@@ -177,6 +177,67 @@ class ClientStoreTests(unittest.TestCase):
|
|||||||
"job-1", "baseline", {"selected": [2]}
|
"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__":
|
if __name__ == "__main__":
|
||||||
unittest.main()
|
unittest.main()
|
||||||
|
|||||||
Reference in New Issue
Block a user