Compare commits
5
Commits
| 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
|
||||
RUN pip install --no-cache-dir /wheels/*.whl
|
||||
COPY tests ./tests
|
||||
COPY scripts ./scripts
|
||||
CMD ["python", "-m", "unittest", "discover", "-s", "tests", "-v"]
|
||||
|
||||
FROM python:3.11-slim
|
||||
|
||||
@@ -95,6 +95,7 @@ backup_dir = "/var/backups/archive-control"
|
||||
registration_timeout = "10s" # First response must arrive within this window.
|
||||
heartbeat_interval = "15s" # Server may negotiate a different effective value.
|
||||
offline_timeout = "45s"
|
||||
outbound_enqueue_timeout = "15s" # Reconnect rather than wedging if sends stop draining.
|
||||
reconnect_initial = "1s"
|
||||
reconnect_max = "60s" # Retry forever, never wait longer than this.
|
||||
reconnect_reset_after = "60s"
|
||||
@@ -262,24 +263,24 @@ control consumer, client consumer, E2E, then image publication.
|
||||
### Reproducible multi-platform Buildx lifecycle
|
||||
|
||||
The named Buildx builders are local acceleration/cache only; they are not a
|
||||
deployment dependency and can be removed after publication. To create a fresh
|
||||
builder, verify its platforms, publish a release, and remove it afterwards:
|
||||
deployment dependency and can be removed after publication. Use the checked-in
|
||||
publisher: it installs only the non-native binfmt handler, creates a temporary
|
||||
rootless BuildKit builder, verifies both platforms, publishes the index,
|
||||
displays its digest, then removes both temporary resources.
|
||||
|
||||
```bash
|
||||
docker buildx create --name archive-control-release --driver docker-container --use
|
||||
docker buildx inspect --bootstrap
|
||||
docker buildx build --platform linux/amd64,linux/arm64 \
|
||||
--tag sodium/archive-clients:vX.Y.Z --push .
|
||||
docker buildx rm archive-control-release
|
||||
scripts/publish-image.sh vX.Y.Z
|
||||
scripts/publish-image.sh vX.Y.Z --also-latest
|
||||
```
|
||||
|
||||
`docker buildx inspect` must list both `linux/amd64` and `linux/arm64` before
|
||||
publishing. If the host has no arm64 emulation, install/configure it according
|
||||
to the host Docker distribution before the build; do not publish a partial
|
||||
single-platform tag. Retain the pushed manifest digest in the release notes
|
||||
and deploy the immutable tag or digest. The optional `archive-control-qemu`
|
||||
builder follows the same lifecycle when it is used for an emulation smoke
|
||||
build.
|
||||
The publisher deliberately uses `moby/buildkit:rootless` with
|
||||
`--oci-worker-no-process-sandbox`. On nested Docker hosts, the default OCI
|
||||
sandbox can fail while masking `/proc/acpi` for an emulated build; rootless
|
||||
BuildKit confines that compatibility setting to the disposable builder. It
|
||||
refuses to publish unless `docker buildx inspect` reports both `linux/amd64`
|
||||
and `linux/arm64`, and removes the builder and binfmt handler on success,
|
||||
failure, or interruption. Retain the displayed manifest digest in release
|
||||
notes and deploy the immutable tag or digest.
|
||||
|
||||
## Operator usage
|
||||
|
||||
|
||||
+15
-12
@@ -45,8 +45,8 @@ seconds. A connection stable for 60 seconds resets the backoff.
|
||||
|
||||
## Version and feature negotiation
|
||||
|
||||
The envelope and registration both state the sender version. Major `1` accepts
|
||||
only major `1`; a different major is rejected. The negotiated minor is the
|
||||
The envelope and registration both state the sender version. Major `2` accepts
|
||||
only major `2`; a different major is rejected. The negotiated minor is the
|
||||
highest mutually supported minor no greater than either endpoint's advertised
|
||||
minor.
|
||||
|
||||
@@ -96,19 +96,22 @@ client state snapshot before a route can become ready.
|
||||
|
||||
Every mutating command includes an expected job revision or immutable job
|
||||
definition plus the expected last global event sequence. A stale revision or
|
||||
sequence is rejected without side effects. The control daemon sends only one
|
||||
active step command for a job and waits for its persisted completion event
|
||||
before commanding the next participant. This single-writer lease lets whichever
|
||||
client owns the current step allocate the next global per-job event sequence;
|
||||
the following command starts from the sequence control has durably accepted.
|
||||
sequence is rejected without side effects. Assignment is represented solely by
|
||||
the durable assignment `CommandAck`; it produces no `JobEvent`. The control
|
||||
daemon sends only one active step command for a job and waits for its persisted
|
||||
completion event before commanding the next participant.
|
||||
|
||||
## Events, progress, and reconciliation
|
||||
|
||||
Each client-originated job event has a unique ID, monotonically increasing
|
||||
global per-job `sequence`, and resulting job revision. Duplicate IDs/sequences
|
||||
are idempotent only when their complete content matches. Control never grants
|
||||
concurrent event-writer leases for one job. A gap or conflicting duplicate
|
||||
pauses destructive orchestration and requests snapshots.
|
||||
global per-job `sequence`, resulting job revision, and the `command_id` of its
|
||||
owning `ExecuteStepCommand` or `CancelJobCommand`. Duplicate IDs/sequences are
|
||||
idempotent only when their complete content matches. Control verifies the
|
||||
client, acknowledged command lease, expected cursor, and active step before
|
||||
accepting it. A gap, conflict, or stale lease triggers a durable authoritative
|
||||
`ReconcileJobCommand`: the client retires only the named stale command leases
|
||||
and restores that job cursor, without deleting resource data or unrelated job
|
||||
state.
|
||||
|
||||
`fraction_complete` is current-step progress and
|
||||
`overall_fraction_complete` is the weighted five- or three-step job progress;
|
||||
@@ -120,7 +123,7 @@ On registration, active-job cursors provide the client's revision, last event
|
||||
sequence, state, and commit flag. Reconciliation applies these rules:
|
||||
|
||||
1. Equal cursors resume normal delivery.
|
||||
2. A client behind receives safe replay/snapshot commands.
|
||||
2. A client behind receives an authoritative reconciliation command.
|
||||
3. Control behind requests and validates the client's full job snapshot.
|
||||
4. Conflicting commit evidence reserves the resource and requires manual
|
||||
reconciliation; neither side performs cleanup.
|
||||
|
||||
+4
-1
@@ -125,7 +125,10 @@ and placement policy differ.
|
||||
and explicit save path, or apply a transactional delta to an existing
|
||||
torrent. Set selected/skipped state, run a full recheck of the target union,
|
||||
and fail immediately if qBittorrent attempts to download. Successful recheck
|
||||
atomically advances the placement generation and is the commit point.
|
||||
atomically advances the placement generation and is the commit point. Before
|
||||
qBittorrent is resumed, the client applies `a+rx` to the verified resource
|
||||
directories and `a+r` to its verified files, preserving ownership, write
|
||||
bits, and special mode bits so other local applications can read the data.
|
||||
5. **Staging Cleanup.** Remove only job-owned staging artifacts from both
|
||||
endpoints. A post-commit cleanup failure produces `CLEANUP_REQUIRED`; it
|
||||
never rolls back or deletes the committed placement.
|
||||
|
||||
+1
-1
@@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta"
|
||||
|
||||
[project]
|
||||
name = "archive-clients"
|
||||
version = "0.1.19"
|
||||
version = "0.1.23"
|
||||
requires-python = ">=3.11"
|
||||
dependencies = ["protobuf==7.35.1", "websockets==16.0"]
|
||||
|
||||
|
||||
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."""
|
||||
|
||||
PROTO_COMMIT = "4ec852014dad74606d4078b3ae1aa208c814b033"
|
||||
PROTO_COMMIT = "03b6751d420e8e4e4cc9aeef4f76b13e28a7265f"
|
||||
|
||||
|
||||
@@ -85,6 +85,7 @@ class ConnectionConfig:
|
||||
registration_timeout: float = 10
|
||||
heartbeat_interval: float = 15
|
||||
offline_timeout: float = 45
|
||||
outbound_enqueue_timeout: float = 15
|
||||
reconnect_initial: float = 1
|
||||
reconnect_max: float = 60
|
||||
reconnect_reset_after: float = 60
|
||||
@@ -249,6 +250,7 @@ def _connection(value: Any) -> ConnectionConfig:
|
||||
raise ConfigError("connection must be a table")
|
||||
_keys(value, {
|
||||
"registration_timeout", "heartbeat_interval", "offline_timeout",
|
||||
"outbound_enqueue_timeout",
|
||||
"reconnect_initial", "reconnect_max", "reconnect_reset_after",
|
||||
"reconnect_jitter",
|
||||
}, "connection")
|
||||
@@ -256,6 +258,9 @@ def _connection(value: Any) -> ConnectionConfig:
|
||||
registration_timeout=_duration(value.get("registration_timeout", "10s")),
|
||||
heartbeat_interval=_duration(value.get("heartbeat_interval", "15s")),
|
||||
offline_timeout=_duration(value.get("offline_timeout", "45s")),
|
||||
outbound_enqueue_timeout=_duration(
|
||||
value.get("outbound_enqueue_timeout", "15s")
|
||||
),
|
||||
reconnect_initial=_duration(value.get("reconnect_initial", "1s")),
|
||||
reconnect_max=_duration(value.get("reconnect_max", "60s")),
|
||||
reconnect_reset_after=_duration(value.get("reconnect_reset_after", "60s")),
|
||||
@@ -265,6 +270,8 @@ def _connection(value: Any) -> ConnectionConfig:
|
||||
raise ConfigError("reconnect_initial cannot exceed reconnect_max")
|
||||
if result.offline_timeout <= result.heartbeat_interval:
|
||||
raise ConfigError("offline_timeout must exceed heartbeat_interval")
|
||||
if result.outbound_enqueue_timeout <= 0:
|
||||
raise ConfigError("outbound_enqueue_timeout must be positive")
|
||||
if not isinstance(result.reconnect_jitter, bool):
|
||||
raise ConfigError("reconnect_jitter must be boolean")
|
||||
return result
|
||||
|
||||
+164
-63
@@ -214,7 +214,7 @@ class ArchiveClientDaemon:
|
||||
raise RuntimeError("registration response correlation mismatch")
|
||||
if response.register_response.status != client_pb2.REGISTRATION_STATUS_ACCEPTED:
|
||||
raise RuntimeError("control rejected registration")
|
||||
if response.register_response.negotiated_version.major != 1:
|
||||
if response.register_response.negotiated_version.major != 2:
|
||||
raise RuntimeError("control negotiated an unsupported protocol version")
|
||||
logger.info("control_connection_registered")
|
||||
outbound: asyncio.Queue[str] = asyncio.Queue(maxsize=100)
|
||||
@@ -229,16 +229,11 @@ class ArchiveClientDaemon:
|
||||
# detects the missing application heartbeats and reconnects.
|
||||
while True:
|
||||
try:
|
||||
frame = await asyncio.wait_for(
|
||||
websocket.recv(),
|
||||
self.config.connection.offline_timeout,
|
||||
frame = await self._next_frame(
|
||||
websocket, writer
|
||||
)
|
||||
except ConnectionClosedOK:
|
||||
return
|
||||
except asyncio.TimeoutError as exc:
|
||||
raise RuntimeError(
|
||||
"control heartbeat timed out"
|
||||
) from exc
|
||||
await self._handle(decode(frame), outbound, command_tasks)
|
||||
finally:
|
||||
writer.cancel()
|
||||
@@ -323,6 +318,63 @@ class ArchiveClientDaemon:
|
||||
while True:
|
||||
await websocket.send(await outbound.get())
|
||||
|
||||
async def _next_frame(self, websocket: Any, writer: asyncio.Task[None]):
|
||||
"""Receive one frame while supervising the connection writer.
|
||||
|
||||
The previous receive-only wait let a failed writer go unnoticed until
|
||||
the remote side happened to close or a heartbeat timeout elapsed.
|
||||
Waiting for both directions makes a failed send an immediate,
|
||||
reconnectable connection failure.
|
||||
"""
|
||||
|
||||
receiver = asyncio.create_task(websocket.recv())
|
||||
try:
|
||||
done, _ = await asyncio.wait(
|
||||
{receiver, writer},
|
||||
timeout=self.config.connection.offline_timeout,
|
||||
return_when=asyncio.FIRST_COMPLETED,
|
||||
)
|
||||
if writer in done:
|
||||
if not receiver.done():
|
||||
receiver.cancel()
|
||||
await asyncio.gather(receiver, return_exceptions=True)
|
||||
if writer.cancelled():
|
||||
raise RuntimeError("control writer was cancelled")
|
||||
error = writer.exception()
|
||||
if error is not None:
|
||||
raise RuntimeError("control writer failed") from error
|
||||
raise RuntimeError("control writer stopped")
|
||||
if receiver not in done:
|
||||
receiver.cancel()
|
||||
await asyncio.gather(receiver, return_exceptions=True)
|
||||
raise RuntimeError("control heartbeat timed out")
|
||||
return receiver.result()
|
||||
except BaseException:
|
||||
if not receiver.done():
|
||||
receiver.cancel()
|
||||
await asyncio.gather(receiver, return_exceptions=True)
|
||||
raise
|
||||
|
||||
async def _enqueue(
|
||||
self, outbound: asyncio.Queue[str], message: str
|
||||
) -> None:
|
||||
"""Bound connection-local backpressure so a dead writer cannot wedge I/O.
|
||||
|
||||
A job event is durable before it is sent, so dropping this particular
|
||||
connection after a bounded wait is safe: reconciliation/redelivery on
|
||||
the next session will recover it. In contrast, indefinitely waiting
|
||||
for a full queue can prevent heartbeat acknowledgements from being
|
||||
read or sent, which prevents the reconnect supervisor from running.
|
||||
"""
|
||||
|
||||
try:
|
||||
await asyncio.wait_for(
|
||||
outbound.put(message),
|
||||
timeout=self.config.connection.outbound_enqueue_timeout,
|
||||
)
|
||||
except asyncio.TimeoutError as exc:
|
||||
raise RuntimeError("control outbound queue is blocked") from exc
|
||||
|
||||
async def _handle(
|
||||
self,
|
||||
envelope: Any,
|
||||
@@ -334,13 +386,20 @@ class ArchiveClientDaemon:
|
||||
response = new_envelope()
|
||||
response.correlation_id = envelope.message_id
|
||||
response.heartbeat_ack.sequence = envelope.heartbeat.sequence
|
||||
await outbound.put(encode(response))
|
||||
await self._enqueue(outbound, encode(response))
|
||||
elif payload == "command":
|
||||
await self._accept_command(envelope, outbound, command_tasks)
|
||||
elif payload == "protocol_error":
|
||||
logger.warning(
|
||||
"control_reported_protocol_error",
|
||||
extra={"error_code": envelope.protocol_error.error.code},
|
||||
extra={
|
||||
"error_code": envelope.protocol_error.error.code,
|
||||
"error_detail": envelope.protocol_error.error.message,
|
||||
"offending_message_id": (
|
||||
envelope.protocol_error.offending_message_id
|
||||
),
|
||||
"retryable": envelope.protocol_error.error.retryable,
|
||||
},
|
||||
)
|
||||
|
||||
async def _accept_command(
|
||||
@@ -368,12 +427,36 @@ class ArchiveClientDaemon:
|
||||
accepted = None
|
||||
accepted_for_execution = False
|
||||
try:
|
||||
accepted = await asyncio.to_thread(
|
||||
self.store.accept_command,
|
||||
command.command_id,
|
||||
encode_message(command),
|
||||
encode_message(acknowledgement),
|
||||
)
|
||||
if command.WhichOneof("payload") == "reconcile_job":
|
||||
reconciliation = command.reconcile_job
|
||||
accepted = await asyncio.to_thread(
|
||||
self.store.accept_reconcile_command,
|
||||
command.command_id,
|
||||
encode_message(command),
|
||||
encode_message(acknowledgement),
|
||||
job_id=reconciliation.authoritative_job.definition.job_id,
|
||||
definition_json=encode_message(
|
||||
reconciliation.authoritative_job.definition
|
||||
),
|
||||
state=job_pb2.JobState.Name(
|
||||
reconciliation.authoritative_job.state
|
||||
),
|
||||
revision=reconciliation.authoritative_job.revision,
|
||||
last_event_sequence=(
|
||||
reconciliation.authoritative_last_event_sequence
|
||||
),
|
||||
committed=reconciliation.authoritative_job.committed,
|
||||
superseded_command_ids=list(
|
||||
reconciliation.superseded_command_ids
|
||||
),
|
||||
)
|
||||
else:
|
||||
accepted = await asyncio.to_thread(
|
||||
self.store.accept_command,
|
||||
command.command_id,
|
||||
encode_message(command),
|
||||
encode_message(acknowledgement),
|
||||
)
|
||||
if accepted.duplicate:
|
||||
acknowledgement = decode_message(
|
||||
accepted.acknowledgement_json, control_pb2.CommandAck()
|
||||
@@ -381,6 +464,7 @@ class ArchiveClientDaemon:
|
||||
if (
|
||||
acknowledgement.status
|
||||
== control_pb2.COMMAND_ACK_STATUS_ACCEPTED
|
||||
and accepted.state == "accepted"
|
||||
):
|
||||
accepted_for_execution = True
|
||||
acknowledgement.status = (
|
||||
@@ -398,7 +482,7 @@ class ArchiveClientDaemon:
|
||||
response = new_envelope()
|
||||
response.correlation_id = envelope.message_id
|
||||
response.command_ack.CopyFrom(acknowledgement)
|
||||
await outbound.put(encode(response))
|
||||
await self._enqueue(outbound, encode(response))
|
||||
|
||||
if (
|
||||
accepted is not None
|
||||
@@ -419,7 +503,7 @@ class ArchiveClientDaemon:
|
||||
job_snapshot.last_event_sequence = int(
|
||||
row["last_event_sequence"]
|
||||
)
|
||||
await outbound.put(encode(snapshot))
|
||||
await self._enqueue(outbound, encode(snapshot))
|
||||
elif (
|
||||
accepted is not None
|
||||
and accepted_for_execution
|
||||
@@ -472,7 +556,7 @@ class ArchiveClientDaemon:
|
||||
outbound: asyncio.Queue[str],
|
||||
command_tasks: set[asyncio.Task[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):
|
||||
acknowledgement = decode_message(
|
||||
str(row["acknowledgement_json"]), control_pb2.CommandAck()
|
||||
@@ -483,28 +567,20 @@ class ArchiveClientDaemon:
|
||||
str(row["command_json"]), control_pb2.Command()
|
||||
)
|
||||
if command.WhichOneof("payload") == "ensure_route":
|
||||
self._schedule_route_command(
|
||||
command, "", outbound, command_tasks
|
||||
)
|
||||
elif (
|
||||
command.WhichOneof("payload") in {
|
||||
"assign_job", "execute_step", "cancel_job"
|
||||
}
|
||||
and self.jobs is not None
|
||||
):
|
||||
if command.WhichOneof("payload") == "cancel_job":
|
||||
self.jobs.request_cancel(command.cancel_job.job_id)
|
||||
job_commands.append(command)
|
||||
if job_commands:
|
||||
task = asyncio.create_task(
|
||||
self._resume_job_commands(job_commands, outbound),
|
||||
name="resume-job-commands",
|
||||
)
|
||||
command_tasks.add(task)
|
||||
task.add_done_callback(
|
||||
lambda completed: self._command_finished(
|
||||
completed, command_tasks
|
||||
)
|
||||
# A newer ensure command for the same route is sufficient to
|
||||
# restore the process-local route path and replay its own
|
||||
# updates. Replaying every historical ready command causes
|
||||
# redundant Syncthing configuration after a restart.
|
||||
latest_route_commands[command.ensure_route.route.route_id] = command
|
||||
# Job commands are intentionally not replayed here. A command
|
||||
# acknowledgement is durable on both sides; the control daemon
|
||||
# redelivers an unacknowledged command, while registration
|
||||
# reconciliation requests snapshots for active jobs. Replaying
|
||||
# every historical command re-emits overlapping event ranges and
|
||||
# can flood the single connection's bounded outbound queue.
|
||||
for command in latest_route_commands.values():
|
||||
self._schedule_route_command(
|
||||
command, "", outbound, command_tasks, restore_ready=True
|
||||
)
|
||||
|
||||
async def _resume_route_commands(
|
||||
@@ -600,34 +676,46 @@ class ArchiveClientDaemon:
|
||||
loop = asyncio.get_running_loop()
|
||||
|
||||
def emit(event):
|
||||
if not self.store.is_command_active(command.command_id):
|
||||
# Reconciliation retired this lease while its blocking
|
||||
# worker was still unwinding. Its durable local record is
|
||||
# not allowed to re-enter the control event stream.
|
||||
return
|
||||
response = new_envelope()
|
||||
response.correlation_id = correlation_id
|
||||
response.job_event.CopyFrom(event)
|
||||
future = asyncio.run_coroutine_threadsafe(
|
||||
outbound.put(encode(response)), loop
|
||||
self._enqueue(outbound, encode(response)), loop
|
||||
)
|
||||
try:
|
||||
future.result(timeout=30)
|
||||
future.result(
|
||||
timeout=self.config.connection.outbound_enqueue_timeout
|
||||
+ 1
|
||||
)
|
||||
except Exception:
|
||||
# The event is already durable in the client DB and will
|
||||
# be replayed after reconnect.
|
||||
pass
|
||||
|
||||
events = await asyncio.to_thread(
|
||||
self.jobs.execute, command.execute_step, emit
|
||||
self.jobs.execute, command.execute_step, emit, command.command_id
|
||||
)
|
||||
streamed = True
|
||||
elif payload == "cancel_job":
|
||||
events = await asyncio.to_thread(
|
||||
self.jobs.cancel, command.cancel_job
|
||||
self.jobs.cancel, command.cancel_job, command.command_id
|
||||
)
|
||||
else:
|
||||
raise JobExecutionError("job command payload is unsupported")
|
||||
for event in (() if streamed else events):
|
||||
if not await asyncio.to_thread(
|
||||
self.store.is_command_active, command.command_id
|
||||
):
|
||||
return
|
||||
response = new_envelope()
|
||||
response.correlation_id = correlation_id
|
||||
response.job_event.CopyFrom(event)
|
||||
await outbound.put(encode(response))
|
||||
await self._enqueue(outbound, encode(response))
|
||||
|
||||
def _schedule_route_command(
|
||||
self,
|
||||
@@ -635,12 +723,15 @@ class ArchiveClientDaemon:
|
||||
correlation_id: str,
|
||||
outbound: asyncio.Queue[str],
|
||||
command_tasks: set[asyncio.Task[None]],
|
||||
restore_ready: bool = False,
|
||||
) -> None:
|
||||
if command.command_id in self._active_route_commands:
|
||||
return
|
||||
self._active_route_commands.add(command.command_id)
|
||||
task = asyncio.create_task(
|
||||
self._execute_route(command, correlation_id, outbound),
|
||||
self._execute_route(
|
||||
command, correlation_id, outbound, restore_ready=restore_ready
|
||||
),
|
||||
name=f"route-{command.ensure_route.route.route_id}",
|
||||
)
|
||||
command_tasks.add(task)
|
||||
@@ -671,13 +762,14 @@ class ArchiveClientDaemon:
|
||||
response = new_envelope()
|
||||
response.correlation_id = correlation_id
|
||||
response.inventory_chunk.CopyFrom(chunk)
|
||||
await outbound.put(encode(response))
|
||||
await self._enqueue(outbound, encode(response))
|
||||
|
||||
async def _execute_route(
|
||||
self,
|
||||
command: Any,
|
||||
correlation_id: str,
|
||||
outbound: asyncio.Queue[str],
|
||||
restore_ready: bool = False,
|
||||
) -> None:
|
||||
assert self.routes is not None
|
||||
spec = command.ensure_route.route
|
||||
@@ -696,21 +788,19 @@ class ArchiveClientDaemon:
|
||||
response.route_update.CopyFrom(
|
||||
decode_message(str(row["update_json"]), control_pb2.RouteUpdate())
|
||||
)
|
||||
await outbound.put(encode(response))
|
||||
await self._enqueue(outbound, encode(response))
|
||||
if attempt["state"] == "ready":
|
||||
# Route readiness is durable, but this process-local lookup is
|
||||
# not. Reconfigure (which validates the existing Syncthing folder
|
||||
# and recreates its local directory if necessary) before replaying
|
||||
# the stored READY updates. This is essential after a daemon or
|
||||
# mount restart: a previously ready route may otherwise point at a
|
||||
# path that is no longer present, and a later source-stage command
|
||||
# would fail with a bare ENOENT.
|
||||
configured = await asyncio.to_thread(
|
||||
self.routes.configure,
|
||||
spec,
|
||||
time.monotonic() + spec.setup_timeout_seconds,
|
||||
)
|
||||
self._known_route_paths[spec.route_id] = configured.local_path
|
||||
if restore_ready:
|
||||
# Route readiness is durable, but this process-local lookup is
|
||||
# not. Reconfigure (which validates the existing Syncthing
|
||||
# folder and recreates its local directory if necessary) after
|
||||
# a daemon restart, before replaying stored READY updates.
|
||||
configured = await asyncio.to_thread(
|
||||
self.routes.configure,
|
||||
spec,
|
||||
time.monotonic() + spec.setup_timeout_seconds,
|
||||
)
|
||||
self._known_route_paths[spec.route_id] = configured.local_path
|
||||
return
|
||||
if attempt["state"] == "failed":
|
||||
return
|
||||
@@ -863,7 +953,7 @@ class ArchiveClientDaemon:
|
||||
response = new_envelope()
|
||||
response.correlation_id = correlation_id
|
||||
response.route_update.CopyFrom(update)
|
||||
await outbound.put(encode(response))
|
||||
await self._enqueue(outbound, encode(response))
|
||||
|
||||
def _local_route(
|
||||
self, spec: route_pb2.EnsureRouteSpec, state: int
|
||||
@@ -990,6 +1080,17 @@ class ArchiveClientDaemon:
|
||||
acknowledgement.error.message = "job assignment is invalid"
|
||||
else:
|
||||
acknowledgement.status = control_pb2.COMMAND_ACK_STATUS_ACCEPTED
|
||||
elif command.WhichOneof("payload") == "reconcile_job":
|
||||
record = command.reconcile_job.authoritative_job
|
||||
if (
|
||||
not record.definition.job_id
|
||||
or command.reconcile_job.authoritative_last_event_sequence < 0
|
||||
):
|
||||
acknowledgement.status = control_pb2.COMMAND_ACK_STATUS_REJECTED
|
||||
acknowledgement.error.code = common_pb2.ERROR_CODE_INVALID_ARGUMENT
|
||||
acknowledgement.error.message = "job reconciliation is invalid"
|
||||
else:
|
||||
acknowledgement.status = control_pb2.COMMAND_ACK_STATUS_ACCEPTED
|
||||
elif command.WhichOneof("payload") == "execute_step":
|
||||
step = command.execute_step
|
||||
if self.jobs is None:
|
||||
|
||||
+92
-18
@@ -97,33 +97,25 @@ class ClientJobExecutor:
|
||||
) -> list[control_pb2.JobEvent]:
|
||||
definition = command.job
|
||||
self._validate_definition(definition)
|
||||
replay = self._replay(
|
||||
definition.job_id, command.expected_last_event_sequence
|
||||
self.store.ensure_job_definition(
|
||||
definition.job_id, encode_message(definition)
|
||||
)
|
||||
if replay:
|
||||
return replay
|
||||
event = self._event(
|
||||
definition,
|
||||
sequence=command.expected_last_event_sequence + 1,
|
||||
revision=command.expected_job_revision,
|
||||
event_type=control_pb2.JOB_EVENT_TYPE_ASSIGNED,
|
||||
state=job_pb2.JOB_STATE_PREPARING,
|
||||
committed=False,
|
||||
)
|
||||
self._record(definition, event)
|
||||
return [event]
|
||||
# Assignment succeeds through CommandAck. It must not let a client
|
||||
# allocate a globally ordered JobEvent cursor.
|
||||
return []
|
||||
|
||||
def execute(
|
||||
self,
|
||||
command: control_pb2.ExecuteStepCommand,
|
||||
event_callback: Callable[[control_pb2.JobEvent], None] | None = None,
|
||||
command_id: str = "",
|
||||
) -> list[control_pb2.JobEvent]:
|
||||
# asyncio cancellation of a connection-bound task cannot stop the
|
||||
# synchronous filesystem/qB operation already running in its worker
|
||||
# thread. A replay after reconnect therefore waits for that operation
|
||||
# and then reads its durable event journal instead of executing twice.
|
||||
with self._execution_lock(command.job_id):
|
||||
return self._execute_locked(command, event_callback)
|
||||
return self._execute_locked(command, event_callback, command_id)
|
||||
|
||||
def _execution_lock(self, job_id: str) -> threading.Lock:
|
||||
with self._execution_locks_guard:
|
||||
@@ -133,6 +125,7 @@ class ClientJobExecutor:
|
||||
self,
|
||||
command: control_pb2.ExecuteStepCommand,
|
||||
event_callback: Callable[[control_pb2.JobEvent], None] | None = None,
|
||||
command_id: str = "",
|
||||
) -> list[control_pb2.JobEvent]:
|
||||
definition = self._definition(command.job_id)
|
||||
replay = self._replay(
|
||||
@@ -186,6 +179,7 @@ class ClientJobExecutor:
|
||||
),
|
||||
step=command.step,
|
||||
step_state=job_pb2.STEP_STATE_RUNNING,
|
||||
command_id=command_id,
|
||||
)
|
||||
self._record(definition, started)
|
||||
cursor = started
|
||||
@@ -228,6 +222,7 @@ class ClientJobExecutor:
|
||||
committed=cursor.committed,
|
||||
step=command.step,
|
||||
step_state=job_pb2.STEP_STATE_RUNNING,
|
||||
command_id=command_id,
|
||||
)
|
||||
event.progress.fraction_complete = max(
|
||||
0.0, min(float(fraction), 1.0)
|
||||
@@ -258,6 +253,7 @@ class ClientJobExecutor:
|
||||
committed=started.committed,
|
||||
step=command.step,
|
||||
step_state=job_pb2.STEP_STATE_CANCELLED,
|
||||
command_id=command_id,
|
||||
)
|
||||
cancelling.error.code = common_pb2.ERROR_CODE_CANCELLED
|
||||
cancelling.error.message = str(error)
|
||||
@@ -315,6 +311,7 @@ class ClientJobExecutor:
|
||||
committed=started.committed,
|
||||
step=command.step,
|
||||
step_state=job_pb2.STEP_STATE_FAILED,
|
||||
command_id=command_id,
|
||||
)
|
||||
failed.error.code = _job_error_code(error)
|
||||
failed.error.message = str(error) or type(error).__name__
|
||||
@@ -357,6 +354,7 @@ class ClientJobExecutor:
|
||||
committed=committed,
|
||||
step=command.step,
|
||||
step_state=job_pb2.STEP_STATE_SUCCEEDED,
|
||||
command_id=command_id,
|
||||
)
|
||||
if result is not None:
|
||||
succeeded.observed_placement.CopyFrom(result)
|
||||
@@ -367,7 +365,7 @@ class ClientJobExecutor:
|
||||
return emitted
|
||||
|
||||
def cancel(
|
||||
self, command: control_pb2.CancelJobCommand
|
||||
self, command: control_pb2.CancelJobCommand, command_id: str = ""
|
||||
) -> list[control_pb2.JobEvent]:
|
||||
definition = self._definition(command.job_id)
|
||||
self.request_cancel(command.job_id)
|
||||
@@ -418,6 +416,7 @@ class ClientJobExecutor:
|
||||
committed=committed,
|
||||
step=job_pb2.JOB_STEP_KIND_ROLLBACK,
|
||||
step_state=job_pb2.STEP_STATE_SUCCEEDED,
|
||||
command_id=command_id,
|
||||
)
|
||||
if placement is not None:
|
||||
event.observed_placement.CopyFrom(placement)
|
||||
@@ -809,13 +808,14 @@ class ClientJobExecutor:
|
||||
fraction, 0, 0, "qBittorrent stopped recheck"
|
||||
),
|
||||
)
|
||||
if should_start:
|
||||
self.qbittorrent.start(qb_torrent_id)
|
||||
verified = self.qbittorrent.get_resource(info_hash)
|
||||
if verified is None:
|
||||
raise JobExecutionError(
|
||||
"verified qBittorrent resource disappeared"
|
||||
)
|
||||
_normalize_verified_resource_permissions(self.qb_root, verified)
|
||||
if should_start:
|
||||
self.qbittorrent.start(qb_torrent_id)
|
||||
placement = resource_pb2.Placement(
|
||||
client_id=self.client_id,
|
||||
state=resource_pb2.PLACEMENT_STATE_PRESENT,
|
||||
@@ -1050,6 +1050,7 @@ class ClientJobExecutor:
|
||||
committed: bool,
|
||||
step: int = job_pb2.JOB_STEP_KIND_UNSPECIFIED,
|
||||
step_state: int = job_pb2.STEP_STATE_UNSPECIFIED,
|
||||
command_id: str = "",
|
||||
) -> control_pb2.JobEvent:
|
||||
event = control_pb2.JobEvent(
|
||||
event_id=str(uuid.uuid4()),
|
||||
@@ -1059,6 +1060,7 @@ class ClientJobExecutor:
|
||||
type=event_type,
|
||||
state=state,
|
||||
committed=committed,
|
||||
command_id=command_id,
|
||||
)
|
||||
event.occurred_at.GetCurrentTime()
|
||||
if step != job_pb2.JOB_STEP_KIND_UNSPECIFIED:
|
||||
@@ -1190,6 +1192,78 @@ def _info_hash(definition: job_pb2.JobDefinition) -> str:
|
||||
return value
|
||||
|
||||
|
||||
def _normalize_verified_resource_permissions(
|
||||
qb_root: Path, resource: NormalizedResource
|
||||
) -> None:
|
||||
"""Apply ``a+rx``/``a+r`` to exactly qB-verified completed content.
|
||||
|
||||
qBittorrent may finish a stopped recheck with restrictive modes inherited
|
||||
from source materialization. Normalize only selected, fully completed
|
||||
regular files and their real parent directories; never traverse a symlink
|
||||
or broaden permissions on the configured qB root itself.
|
||||
"""
|
||||
root_metadata = qb_root.lstat()
|
||||
if not stat.S_ISDIR(root_metadata.st_mode):
|
||||
raise JobExecutionError("configured qB root is not a real directory")
|
||||
directories: set[Path] = set()
|
||||
files: set[Path] = set()
|
||||
for item in resource.files:
|
||||
if not item.selected or item.completed_bytes != item.logical_bytes:
|
||||
continue
|
||||
relative = _resource_relative_path(item.canonical_path)
|
||||
current = qb_root
|
||||
for component in relative.parts[:-1]:
|
||||
current = current / component
|
||||
metadata = current.lstat()
|
||||
if not stat.S_ISDIR(metadata.st_mode):
|
||||
raise JobExecutionError(
|
||||
"verified resource parent is not a real directory"
|
||||
)
|
||||
directories.add(current)
|
||||
candidate = current / relative.name
|
||||
metadata = candidate.lstat()
|
||||
if not stat.S_ISREG(metadata.st_mode):
|
||||
raise JobExecutionError("verified resource file is not regular")
|
||||
files.add(candidate)
|
||||
for directory in sorted(directories, key=lambda value: len(value.parts)):
|
||||
metadata = directory.lstat()
|
||||
if not stat.S_ISDIR(metadata.st_mode):
|
||||
raise JobExecutionError(
|
||||
"verified resource parent changed during permission update"
|
||||
)
|
||||
os.chmod(
|
||||
directory,
|
||||
stat.S_IMODE(metadata.st_mode) | 0o555,
|
||||
follow_symlinks=False,
|
||||
)
|
||||
for path in files:
|
||||
metadata = path.lstat()
|
||||
if not stat.S_ISREG(metadata.st_mode):
|
||||
raise JobExecutionError(
|
||||
"verified resource file changed during permission update"
|
||||
)
|
||||
os.chmod(
|
||||
path,
|
||||
stat.S_IMODE(metadata.st_mode) | 0o444,
|
||||
follow_symlinks=False,
|
||||
)
|
||||
|
||||
|
||||
def _resource_relative_path(value: str) -> PurePosixPath:
|
||||
if (
|
||||
not value
|
||||
or "\x00" in value
|
||||
or "\\" in value
|
||||
or value.startswith("/")
|
||||
or any(part in {"", ".", ".."} for part in value.split("/"))
|
||||
):
|
||||
raise JobExecutionError("verified resource path is unsafe")
|
||||
path = PurePosixPath(value)
|
||||
if path.is_absolute():
|
||||
raise JobExecutionError("verified resource path is unsafe")
|
||||
return path
|
||||
|
||||
|
||||
def _fingerprint_matches(
|
||||
current: resource_pb2.ResourceStateFingerprint,
|
||||
expected: resource_pb2.ResourceStateFingerprint,
|
||||
|
||||
@@ -18,7 +18,7 @@ class ProtocolError(ValueError):
|
||||
|
||||
def new_envelope() -> envelope_pb2.Envelope:
|
||||
envelope = envelope_pb2.Envelope()
|
||||
envelope.protocol_version.major = 1
|
||||
envelope.protocol_version.major = 2
|
||||
envelope.message_id = str(uuid.uuid4())
|
||||
envelope.sent_at.FromDatetime(datetime.now(timezone.utc))
|
||||
return envelope
|
||||
@@ -59,7 +59,7 @@ def decode(data: str | bytes, max_bytes: int = 1024 * 1024) -> envelope_pb2.Enve
|
||||
json_format.ParseError,
|
||||
) as exc:
|
||||
raise ProtocolError("invalid control envelope") from exc
|
||||
if envelope.protocol_version.major != 1:
|
||||
if envelope.protocol_version.major != 2:
|
||||
raise ProtocolError("unsupported protocol major version")
|
||||
_canonical_uuid(envelope.message_id, "message_id")
|
||||
if envelope.correlation_id:
|
||||
|
||||
@@ -38,6 +38,7 @@ class FileOperationConflict(RuntimeError):
|
||||
class CommandAcceptance:
|
||||
duplicate: bool
|
||||
acknowledgement_json: str
|
||||
state: str
|
||||
|
||||
|
||||
class ClientStore:
|
||||
@@ -85,7 +86,7 @@ class ClientStore:
|
||||
connection.execute("BEGIN IMMEDIATE")
|
||||
existing = connection.execute(
|
||||
"""
|
||||
SELECT payload_sha256, command_json, acknowledgement_json
|
||||
SELECT payload_sha256, command_json, acknowledgement_json, state
|
||||
FROM commands WHERE command_id = ?
|
||||
""",
|
||||
(command_id,),
|
||||
@@ -93,17 +94,95 @@ class ClientStore:
|
||||
if existing:
|
||||
if existing["payload_sha256"] != digest or existing["command_json"] != payload:
|
||||
raise CommandConflict("command ID was reused with different content")
|
||||
return CommandAcceptance(True, existing["acknowledgement_json"])
|
||||
return CommandAcceptance(
|
||||
True, existing["acknowledgement_json"], existing["state"]
|
||||
)
|
||||
status = json.loads(acknowledgement).get("status")
|
||||
state = (
|
||||
"rejected"
|
||||
if status is not None
|
||||
and status != "COMMAND_ACK_STATUS_ACCEPTED"
|
||||
else "accepted"
|
||||
)
|
||||
connection.execute(
|
||||
"""
|
||||
INSERT INTO commands (
|
||||
command_id, payload_sha256, command_json,
|
||||
acknowledgement_json, state
|
||||
) VALUES (?, ?, ?, ?, 'received')
|
||||
) VALUES (?, ?, ?, ?, ?)
|
||||
""",
|
||||
(command_id, digest, payload, acknowledgement),
|
||||
(command_id, digest, payload, acknowledgement, state),
|
||||
)
|
||||
return CommandAcceptance(False, acknowledgement)
|
||||
return CommandAcceptance(False, acknowledgement, state)
|
||||
|
||||
def accept_reconcile_command(
|
||||
self,
|
||||
command_id: str,
|
||||
command_json: str,
|
||||
acknowledgement_json: str,
|
||||
*,
|
||||
job_id: str,
|
||||
definition_json: str,
|
||||
state: str,
|
||||
revision: int,
|
||||
last_event_sequence: int,
|
||||
committed: bool,
|
||||
superseded_command_ids: list[str],
|
||||
) -> CommandAcceptance:
|
||||
"""Atomically durably accept and apply a reconciliation command.
|
||||
|
||||
A reconciliation acknowledgement is meaningful only after its cursor
|
||||
and retired leases have reached SQLite. Keeping both operations in
|
||||
one transaction makes a reconnect either redeliver the command or
|
||||
observe its completed effect; it cannot observe a bare acknowledgement.
|
||||
"""
|
||||
payload = _canonical(json.loads(command_json))
|
||||
acknowledgement = _canonical(json.loads(acknowledgement_json))
|
||||
definition = _canonical(json.loads(definition_json))
|
||||
digest = hashlib.sha256(payload.encode()).hexdigest()
|
||||
with self._connect() as connection:
|
||||
connection.execute("BEGIN IMMEDIATE")
|
||||
existing = connection.execute(
|
||||
"""
|
||||
SELECT payload_sha256, command_json, acknowledgement_json, state
|
||||
FROM commands WHERE command_id = ?
|
||||
""",
|
||||
(command_id,),
|
||||
).fetchone()
|
||||
if existing:
|
||||
if existing["payload_sha256"] != digest or existing["command_json"] != payload:
|
||||
raise CommandConflict("command ID was reused with different content")
|
||||
return CommandAcceptance(
|
||||
True, existing["acknowledgement_json"], existing["state"]
|
||||
)
|
||||
status = json.loads(acknowledgement).get("status")
|
||||
command_state = (
|
||||
"rejected"
|
||||
if status is not None
|
||||
and status != "COMMAND_ACK_STATUS_ACCEPTED"
|
||||
else "accepted"
|
||||
)
|
||||
connection.execute(
|
||||
"""
|
||||
INSERT INTO commands (
|
||||
command_id, payload_sha256, command_json,
|
||||
acknowledgement_json, state
|
||||
) VALUES (?, ?, ?, ?, ?)
|
||||
""",
|
||||
(command_id, digest, payload, acknowledgement, command_state),
|
||||
)
|
||||
if command_state == "accepted":
|
||||
self._reconcile_job_connection(
|
||||
connection,
|
||||
job_id=job_id,
|
||||
definition=definition,
|
||||
state=state,
|
||||
revision=revision,
|
||||
last_event_sequence=last_event_sequence,
|
||||
committed=committed,
|
||||
superseded_command_ids=superseded_command_ids,
|
||||
)
|
||||
return CommandAcceptance(False, acknowledgement, command_state)
|
||||
|
||||
def list_active_job_cursors(self) -> list[dict[str, object]]:
|
||||
with self._connect() as connection:
|
||||
@@ -122,7 +201,7 @@ class ClientStore:
|
||||
rows = connection.execute(
|
||||
"""
|
||||
SELECT command_id, command_json, acknowledgement_json
|
||||
FROM commands ORDER BY rowid
|
||||
FROM commands WHERE state = 'accepted' ORDER BY rowid
|
||||
"""
|
||||
).fetchall()
|
||||
return [dict(row) for row in rows]
|
||||
@@ -199,6 +278,114 @@ class ClientStore:
|
||||
),
|
||||
)
|
||||
|
||||
def ensure_job_definition(self, job_id: str, definition_json: str) -> None:
|
||||
"""Durably record an assignment without inventing a global event."""
|
||||
definition = _canonical(json.loads(definition_json))
|
||||
with self._connect() as connection:
|
||||
connection.execute("BEGIN IMMEDIATE")
|
||||
existing = connection.execute(
|
||||
"SELECT definition_json FROM jobs WHERE job_id = ?", (job_id,)
|
||||
).fetchone()
|
||||
if existing is not None:
|
||||
if existing["definition_json"] != definition:
|
||||
raise JobConflict("job definition is immutable")
|
||||
return
|
||||
connection.execute(
|
||||
"""
|
||||
INSERT INTO jobs (job_id, definition_json, state, revision,
|
||||
last_event_sequence, committed)
|
||||
VALUES (?, ?, 'JOB_STATE_QUEUED', 0, 0, 0)
|
||||
""",
|
||||
(job_id, definition),
|
||||
)
|
||||
|
||||
def is_command_active(self, command_id: str) -> bool:
|
||||
with self._connect() as connection:
|
||||
row = connection.execute(
|
||||
"SELECT state FROM commands WHERE command_id = ?", (command_id,)
|
||||
).fetchone()
|
||||
return row is not None and row["state"] == "accepted"
|
||||
|
||||
def reconcile_job(
|
||||
self,
|
||||
*,
|
||||
job_id: str,
|
||||
definition_json: str,
|
||||
state: str,
|
||||
revision: int,
|
||||
last_event_sequence: int,
|
||||
committed: bool,
|
||||
superseded_command_ids: list[str],
|
||||
) -> None:
|
||||
"""Apply control's cursor, retiring only stale command leases.
|
||||
|
||||
This drops unacknowledged local journal rows beyond control's cursor;
|
||||
resource files and operation artifacts are intentionally retained.
|
||||
"""
|
||||
definition = _canonical(json.loads(definition_json))
|
||||
with self._connect() as connection:
|
||||
connection.execute("BEGIN IMMEDIATE")
|
||||
self._reconcile_job_connection(
|
||||
connection,
|
||||
job_id=job_id,
|
||||
definition=definition,
|
||||
state=state,
|
||||
revision=revision,
|
||||
last_event_sequence=last_event_sequence,
|
||||
committed=committed,
|
||||
superseded_command_ids=superseded_command_ids,
|
||||
)
|
||||
|
||||
@staticmethod
|
||||
def _reconcile_job_connection(
|
||||
connection: sqlite3.Connection,
|
||||
*,
|
||||
job_id: str,
|
||||
definition: str,
|
||||
state: str,
|
||||
revision: int,
|
||||
last_event_sequence: int,
|
||||
committed: bool,
|
||||
superseded_command_ids: list[str],
|
||||
) -> None:
|
||||
existing = connection.execute(
|
||||
"SELECT definition_json FROM jobs WHERE job_id = ?", (job_id,)
|
||||
).fetchone()
|
||||
if existing is not None and existing["definition_json"] != definition:
|
||||
raise JobConflict("reconciliation has a different job definition")
|
||||
connection.execute(
|
||||
"DELETE FROM events WHERE job_id = ? AND sequence > ?",
|
||||
(job_id, last_event_sequence),
|
||||
)
|
||||
connection.execute(
|
||||
"""
|
||||
INSERT INTO jobs (job_id, definition_json, state, revision,
|
||||
last_event_sequence, committed)
|
||||
VALUES (?, ?, ?, ?, ?, ?)
|
||||
ON CONFLICT(job_id) DO UPDATE SET
|
||||
state = excluded.state, revision = excluded.revision,
|
||||
last_event_sequence = excluded.last_event_sequence,
|
||||
committed = excluded.committed
|
||||
""",
|
||||
(job_id, definition, state, revision, last_event_sequence, int(committed)),
|
||||
)
|
||||
if superseded_command_ids:
|
||||
placeholders = ",".join("?" for _ in superseded_command_ids)
|
||||
rows = connection.execute(
|
||||
f"SELECT command_id, command_json FROM commands WHERE command_id IN ({placeholders})",
|
||||
superseded_command_ids,
|
||||
).fetchall()
|
||||
owned_ids = [
|
||||
row["command_id"] for row in rows
|
||||
if _command_job_id(row["command_json"]) == job_id
|
||||
]
|
||||
if owned_ids:
|
||||
owned_placeholders = ",".join("?" for _ in owned_ids)
|
||||
connection.execute(
|
||||
f"UPDATE commands SET state = 'superseded' WHERE command_id IN ({owned_placeholders})",
|
||||
owned_ids,
|
||||
)
|
||||
|
||||
def begin_file_operation(
|
||||
self,
|
||||
operation_id: str,
|
||||
@@ -591,6 +778,18 @@ class ClientStore:
|
||||
connection.close()
|
||||
|
||||
|
||||
def _command_job_id(command_json: str) -> str:
|
||||
"""Return the job target from canonical protobuf JSON, if it has one."""
|
||||
command = json.loads(command_json)
|
||||
if "assignJob" in command:
|
||||
return str(command["assignJob"].get("job", {}).get("jobId", ""))
|
||||
if "executeStep" in command:
|
||||
return str(command["executeStep"].get("jobId", ""))
|
||||
if "cancelJob" in command:
|
||||
return str(command["cancelJob"].get("jobId", ""))
|
||||
return ""
|
||||
|
||||
|
||||
def _canonical(value: object) -> str:
|
||||
return json.dumps(value, sort_keys=True, separators=(",", ":"), ensure_ascii=False)
|
||||
|
||||
|
||||
@@ -1,2 +1,3 @@
|
||||
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||
"""Generated archive_control.v1 bindings."""
|
||||
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
||||
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||
# -*- coding: utf-8 -*-
|
||||
# Generated by the protocol buffer compiler. DO NOT EDIT!
|
||||
# NO CHECKED-IN PROTOBUF GENCODE
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
||||
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||
import datetime
|
||||
|
||||
from archive_control.v1 import common_pb2 as _common_pb2
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
||||
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||
# -*- coding: utf-8 -*-
|
||||
# Generated by the protocol buffer compiler. DO NOT EDIT!
|
||||
# NO CHECKED-IN PROTOBUF GENCODE
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
||||
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||
import datetime
|
||||
|
||||
from google.protobuf import timestamp_pb2 as _timestamp_pb2
|
||||
|
||||
File diff suppressed because one or more lines are too long
@@ -1,4 +1,4 @@
|
||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
||||
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||
import datetime
|
||||
|
||||
from archive_control.v1 import common_pb2 as _common_pb2
|
||||
@@ -108,8 +108,20 @@ class RequestJobSnapshotCommand(_message.Message):
|
||||
job_ids: _containers.RepeatedScalarFieldContainer[str]
|
||||
def __init__(self, job_ids: _Optional[_Iterable[str]] = ...) -> None: ...
|
||||
|
||||
class ReconcileJobCommand(_message.Message):
|
||||
__slots__ = ("authoritative_job", "authoritative_last_event_sequence", "superseded_command_ids", "reason")
|
||||
AUTHORITATIVE_JOB_FIELD_NUMBER: _ClassVar[int]
|
||||
AUTHORITATIVE_LAST_EVENT_SEQUENCE_FIELD_NUMBER: _ClassVar[int]
|
||||
SUPERSEDED_COMMAND_IDS_FIELD_NUMBER: _ClassVar[int]
|
||||
REASON_FIELD_NUMBER: _ClassVar[int]
|
||||
authoritative_job: _job_pb2.JobRecord
|
||||
authoritative_last_event_sequence: int
|
||||
superseded_command_ids: _containers.RepeatedScalarFieldContainer[str]
|
||||
reason: str
|
||||
def __init__(self, authoritative_job: _Optional[_Union[_job_pb2.JobRecord, _Mapping]] = ..., authoritative_last_event_sequence: _Optional[int] = ..., superseded_command_ids: _Optional[_Iterable[str]] = ..., reason: _Optional[str] = ...) -> None: ...
|
||||
|
||||
class Command(_message.Message):
|
||||
__slots__ = ("command_id", "created_at", "assign_job", "execute_step", "cancel_job", "ensure_route", "inventory_query", "request_job_snapshot")
|
||||
__slots__ = ("command_id", "created_at", "assign_job", "execute_step", "cancel_job", "ensure_route", "inventory_query", "request_job_snapshot", "reconcile_job")
|
||||
COMMAND_ID_FIELD_NUMBER: _ClassVar[int]
|
||||
CREATED_AT_FIELD_NUMBER: _ClassVar[int]
|
||||
ASSIGN_JOB_FIELD_NUMBER: _ClassVar[int]
|
||||
@@ -118,6 +130,7 @@ class Command(_message.Message):
|
||||
ENSURE_ROUTE_FIELD_NUMBER: _ClassVar[int]
|
||||
INVENTORY_QUERY_FIELD_NUMBER: _ClassVar[int]
|
||||
REQUEST_JOB_SNAPSHOT_FIELD_NUMBER: _ClassVar[int]
|
||||
RECONCILE_JOB_FIELD_NUMBER: _ClassVar[int]
|
||||
command_id: str
|
||||
created_at: _timestamp_pb2.Timestamp
|
||||
assign_job: AssignJobCommand
|
||||
@@ -126,7 +139,8 @@ class Command(_message.Message):
|
||||
ensure_route: EnsureRouteCommand
|
||||
inventory_query: _inventory_pb2.InventoryQuery
|
||||
request_job_snapshot: RequestJobSnapshotCommand
|
||||
def __init__(self, command_id: _Optional[str] = ..., created_at: _Optional[_Union[datetime.datetime, _timestamp_pb2.Timestamp, _Mapping]] = ..., assign_job: _Optional[_Union[AssignJobCommand, _Mapping]] = ..., execute_step: _Optional[_Union[ExecuteStepCommand, _Mapping]] = ..., cancel_job: _Optional[_Union[CancelJobCommand, _Mapping]] = ..., ensure_route: _Optional[_Union[EnsureRouteCommand, _Mapping]] = ..., inventory_query: _Optional[_Union[_inventory_pb2.InventoryQuery, _Mapping]] = ..., request_job_snapshot: _Optional[_Union[RequestJobSnapshotCommand, _Mapping]] = ...) -> None: ...
|
||||
reconcile_job: ReconcileJobCommand
|
||||
def __init__(self, command_id: _Optional[str] = ..., created_at: _Optional[_Union[datetime.datetime, _timestamp_pb2.Timestamp, _Mapping]] = ..., assign_job: _Optional[_Union[AssignJobCommand, _Mapping]] = ..., execute_step: _Optional[_Union[ExecuteStepCommand, _Mapping]] = ..., cancel_job: _Optional[_Union[CancelJobCommand, _Mapping]] = ..., ensure_route: _Optional[_Union[EnsureRouteCommand, _Mapping]] = ..., inventory_query: _Optional[_Union[_inventory_pb2.InventoryQuery, _Mapping]] = ..., request_job_snapshot: _Optional[_Union[RequestJobSnapshotCommand, _Mapping]] = ..., reconcile_job: _Optional[_Union[ReconcileJobCommand, _Mapping]] = ...) -> None: ...
|
||||
|
||||
class CommandAck(_message.Message):
|
||||
__slots__ = ("command_id", "status", "error")
|
||||
@@ -139,7 +153,7 @@ class CommandAck(_message.Message):
|
||||
def __init__(self, command_id: _Optional[str] = ..., status: _Optional[_Union[CommandAckStatus, str]] = ..., error: _Optional[_Union[_common_pb2.Error, _Mapping]] = ...) -> None: ...
|
||||
|
||||
class JobEvent(_message.Message):
|
||||
__slots__ = ("event_id", "job_id", "sequence", "job_revision", "type", "state", "committed", "progress", "error", "observed_resource", "observed_placement", "occurred_at")
|
||||
__slots__ = ("event_id", "job_id", "sequence", "job_revision", "type", "state", "committed", "progress", "error", "observed_resource", "observed_placement", "occurred_at", "command_id")
|
||||
EVENT_ID_FIELD_NUMBER: _ClassVar[int]
|
||||
JOB_ID_FIELD_NUMBER: _ClassVar[int]
|
||||
SEQUENCE_FIELD_NUMBER: _ClassVar[int]
|
||||
@@ -152,6 +166,7 @@ class JobEvent(_message.Message):
|
||||
OBSERVED_RESOURCE_FIELD_NUMBER: _ClassVar[int]
|
||||
OBSERVED_PLACEMENT_FIELD_NUMBER: _ClassVar[int]
|
||||
OCCURRED_AT_FIELD_NUMBER: _ClassVar[int]
|
||||
COMMAND_ID_FIELD_NUMBER: _ClassVar[int]
|
||||
event_id: str
|
||||
job_id: str
|
||||
sequence: int
|
||||
@@ -164,7 +179,8 @@ class JobEvent(_message.Message):
|
||||
observed_resource: _resource_pb2.ResourceStateFingerprint
|
||||
observed_placement: _resource_pb2.Placement
|
||||
occurred_at: _timestamp_pb2.Timestamp
|
||||
def __init__(self, event_id: _Optional[str] = ..., job_id: _Optional[str] = ..., sequence: _Optional[int] = ..., job_revision: _Optional[int] = ..., type: _Optional[_Union[JobEventType, str]] = ..., state: _Optional[_Union[_job_pb2.JobState, str]] = ..., committed: _Optional[bool] = ..., progress: _Optional[_Union[_job_pb2.JobProgress, _Mapping]] = ..., error: _Optional[_Union[_common_pb2.Error, _Mapping]] = ..., observed_resource: _Optional[_Union[_resource_pb2.ResourceStateFingerprint, _Mapping]] = ..., observed_placement: _Optional[_Union[_resource_pb2.Placement, _Mapping]] = ..., occurred_at: _Optional[_Union[datetime.datetime, _timestamp_pb2.Timestamp, _Mapping]] = ...) -> None: ...
|
||||
command_id: str
|
||||
def __init__(self, event_id: _Optional[str] = ..., job_id: _Optional[str] = ..., sequence: _Optional[int] = ..., job_revision: _Optional[int] = ..., type: _Optional[_Union[JobEventType, str]] = ..., state: _Optional[_Union[_job_pb2.JobState, str]] = ..., committed: _Optional[bool] = ..., progress: _Optional[_Union[_job_pb2.JobProgress, _Mapping]] = ..., error: _Optional[_Union[_common_pb2.Error, _Mapping]] = ..., observed_resource: _Optional[_Union[_resource_pb2.ResourceStateFingerprint, _Mapping]] = ..., observed_placement: _Optional[_Union[_resource_pb2.Placement, _Mapping]] = ..., occurred_at: _Optional[_Union[datetime.datetime, _timestamp_pb2.Timestamp, _Mapping]] = ..., command_id: _Optional[str] = ...) -> None: ...
|
||||
|
||||
class JobSnapshot(_message.Message):
|
||||
__slots__ = ("job", "last_event_sequence", "in_flight_command_ids")
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
||||
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||
# -*- coding: utf-8 -*-
|
||||
# Generated by the protocol buffer compiler. DO NOT EDIT!
|
||||
# NO CHECKED-IN PROTOBUF GENCODE
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
||||
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||
import datetime
|
||||
|
||||
from archive_control.v1 import client_pb2 as _client_pb2
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
||||
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||
# -*- coding: utf-8 -*-
|
||||
# Generated by the protocol buffer compiler. DO NOT EDIT!
|
||||
# NO CHECKED-IN PROTOBUF GENCODE
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
||||
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||
from archive_control.v1 import common_pb2 as _common_pb2
|
||||
from archive_control.v1 import resource_pb2 as _resource_pb2
|
||||
from google.protobuf.internal import containers as _containers
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
||||
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||
# -*- coding: utf-8 -*-
|
||||
# Generated by the protocol buffer compiler. DO NOT EDIT!
|
||||
# NO CHECKED-IN PROTOBUF GENCODE
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
||||
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||
import datetime
|
||||
|
||||
from archive_control.v1 import common_pb2 as _common_pb2
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
||||
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||
# -*- coding: utf-8 -*-
|
||||
# Generated by the protocol buffer compiler. DO NOT EDIT!
|
||||
# NO CHECKED-IN PROTOBUF GENCODE
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
||||
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||
import datetime
|
||||
|
||||
from google.protobuf import timestamp_pb2 as _timestamp_pb2
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
||||
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||
# -*- coding: utf-8 -*-
|
||||
# Generated by the protocol buffer compiler. DO NOT EDIT!
|
||||
# NO CHECKED-IN PROTOBUF GENCODE
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
||||
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||
import datetime
|
||||
|
||||
from google.protobuf import timestamp_pb2 as _timestamp_pb2
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
||||
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||
# -*- coding: utf-8 -*-
|
||||
# Generated by the protocol buffer compiler. DO NOT EDIT!
|
||||
# NO CHECKED-IN PROTOBUF GENCODE
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
||||
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||
import datetime
|
||||
|
||||
from archive_control.v1 import resource_pb2 as _resource_pb2
|
||||
|
||||
+119
-2
@@ -74,7 +74,7 @@ class DaemonTransportTests(unittest.IsolatedAsyncioTestCase):
|
||||
response.register_response.status = (
|
||||
client_pb2.REGISTRATION_STATUS_ACCEPTED
|
||||
)
|
||||
response.register_response.negotiated_version.major = 1
|
||||
response.register_response.negotiated_version.major = 2
|
||||
await websocket.send(encode(response))
|
||||
# Deliberately keep TCP/WebSocket open but send no application
|
||||
# heartbeats. This models a stale proxy/server-side session.
|
||||
@@ -110,6 +110,123 @@ class DaemonTransportTests(unittest.IsolatedAsyncioTestCase):
|
||||
):
|
||||
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):
|
||||
with tempfile.TemporaryDirectory() as directory:
|
||||
root = Path(directory)
|
||||
@@ -368,7 +485,7 @@ class DaemonTransportTests(unittest.IsolatedAsyncioTestCase):
|
||||
response = new_envelope()
|
||||
response.correlation_id = registration.message_id
|
||||
response.register_response.status = client_pb2.REGISTRATION_STATUS_ACCEPTED
|
||||
response.register_response.negotiated_version.major = 1
|
||||
response.register_response.negotiated_version.major = 2
|
||||
await websocket.send(encode(response))
|
||||
heartbeat = new_envelope()
|
||||
heartbeat.heartbeat.sequence = 7
|
||||
|
||||
+56
-9
@@ -1,4 +1,6 @@
|
||||
import hashlib
|
||||
import os
|
||||
import stat
|
||||
import tempfile
|
||||
import threading
|
||||
import unittest
|
||||
@@ -10,12 +12,13 @@ from archive_clients.bencode import encode
|
||||
from archive_clients.jobs import (
|
||||
ClientJobExecutor,
|
||||
JobExecutionError,
|
||||
_normalize_verified_resource_permissions,
|
||||
_resource_fingerprint,
|
||||
)
|
||||
from archive_clients.syncthing import RouteSetupError, SyncthingTransferStatus
|
||||
from archive_clients.resources import normalize_resource
|
||||
from archive_clients.resources import NormalizedResource, normalize_resource
|
||||
from archive_clients.state import ClientStore
|
||||
from archive_control.v1 import control_pb2, job_pb2
|
||||
from archive_control.v1 import control_pb2, job_pb2, resource_pb2
|
||||
|
||||
|
||||
class CompleteSyncthing:
|
||||
@@ -39,6 +42,50 @@ class SlowRescanSyncthing(CompleteSyncthing):
|
||||
|
||||
|
||||
class ClientJobHappyPathTests(unittest.TestCase):
|
||||
def test_verified_resource_permissions_are_readable_by_other_apps(self):
|
||||
with tempfile.TemporaryDirectory() as directory:
|
||||
root = Path(directory) / "qb"
|
||||
resource_directory = root / "resource"
|
||||
resource_directory.mkdir(parents=True)
|
||||
completed = resource_directory / "complete.bin"
|
||||
incomplete = resource_directory / "incomplete.bin"
|
||||
completed.write_bytes(b"complete")
|
||||
incomplete.write_bytes(b"incomplete")
|
||||
os.chmod(resource_directory, 0o300)
|
||||
os.chmod(completed, 0o200)
|
||||
os.chmod(incomplete, 0o200)
|
||||
resource = NormalizedResource(
|
||||
resource_pb2.ResourceSummary(),
|
||||
(
|
||||
resource_pb2.TorrentFile(
|
||||
file_index=0,
|
||||
canonical_path="resource/complete.bin",
|
||||
logical_bytes=len(b"complete"),
|
||||
completed_bytes=len(b"complete"),
|
||||
selected=True,
|
||||
),
|
||||
resource_pb2.TorrentFile(
|
||||
file_index=1,
|
||||
canonical_path="resource/incomplete.bin",
|
||||
logical_bytes=len(b"incomplete"),
|
||||
completed_bytes=0,
|
||||
selected=True,
|
||||
),
|
||||
),
|
||||
Mock(),
|
||||
)
|
||||
_normalize_verified_resource_permissions(root, resource)
|
||||
self.assertEqual(
|
||||
stat.S_IMODE(resource_directory.stat().st_mode) & 0o555,
|
||||
0o555,
|
||||
)
|
||||
self.assertEqual(
|
||||
stat.S_IMODE(completed.stat().st_mode) & 0o444, 0o444
|
||||
)
|
||||
self.assertEqual(
|
||||
stat.S_IMODE(incomplete.stat().st_mode), 0o200
|
||||
)
|
||||
|
||||
def test_syncthing_api_outage_during_transfer_is_retried(self):
|
||||
"""A transient local REST outage must not terminally fail the job."""
|
||||
with tempfile.TemporaryDirectory() as directory:
|
||||
@@ -436,19 +483,19 @@ class ClientJobHappyPathTests(unittest.TestCase):
|
||||
|
||||
source_assigned = source.assign(control_pb2.AssignJobCommand(
|
||||
job=definition,
|
||||
expected_job_revision=1,
|
||||
expected_job_revision=0,
|
||||
expected_last_event_sequence=0,
|
||||
))
|
||||
target_assigned = target.assign(control_pb2.AssignJobCommand(
|
||||
job=definition,
|
||||
expected_job_revision=1,
|
||||
expected_last_event_sequence=1,
|
||||
expected_job_revision=0,
|
||||
expected_last_event_sequence=0,
|
||||
))
|
||||
self.assertEqual(source_assigned[0].sequence, 1)
|
||||
self.assertEqual(target_assigned[0].sequence, 2)
|
||||
self.assertEqual(source_assigned, [])
|
||||
self.assertEqual(target_assigned, [])
|
||||
|
||||
cursor_revision = 1
|
||||
cursor_sequence = 2
|
||||
cursor_revision = 0
|
||||
cursor_sequence = 0
|
||||
pipeline = (
|
||||
(source, job_pb2.JOB_STEP_KIND_SOURCE_STAGE),
|
||||
(target, job_pb2.JOB_STEP_KIND_SYNCTHING_TRANSFER),
|
||||
|
||||
@@ -177,6 +177,67 @@ class ClientStoreTests(unittest.TestCase):
|
||||
"job-1", "baseline", {"selected": [2]}
|
||||
)
|
||||
|
||||
def test_reconciliation_retires_only_stale_job_leases(self):
|
||||
with tempfile.TemporaryDirectory() as directory:
|
||||
store = ClientStore(Path(directory) / "state.db")
|
||||
store.initialize()
|
||||
stale = store.accept_command(
|
||||
"stale", '{"executeStep":{"jobId":"job-1"}}',
|
||||
'{"status":"COMMAND_ACK_STATUS_ACCEPTED"}',
|
||||
)
|
||||
store.accept_command(
|
||||
"other", '{"executeStep":{"jobId":"job-2"}}',
|
||||
'{"status":"COMMAND_ACK_STATUS_ACCEPTED"}',
|
||||
)
|
||||
self.assertEqual(stale.state, "accepted")
|
||||
store.save_job(
|
||||
"job-1", '{"jobId":"job-1"}', "JOB_STATE_RUNNING", 3, 3, False,
|
||||
)
|
||||
store.reconcile_job(
|
||||
job_id="job-1", definition_json='{"jobId":"job-1"}',
|
||||
state="JOB_STATE_QUEUED", revision=0,
|
||||
last_event_sequence=0, committed=False,
|
||||
superseded_command_ids=["stale"],
|
||||
)
|
||||
self.assertFalse(store.is_command_active("stale"))
|
||||
self.assertTrue(store.is_command_active("other"))
|
||||
row = store.job_snapshot_rows(["job-1"])[0]
|
||||
self.assertEqual(row["last_event_sequence"], 0)
|
||||
self.assertEqual(row["state"], "JOB_STATE_QUEUED")
|
||||
|
||||
def test_reconciliation_acknowledgement_and_state_are_atomic(self):
|
||||
with tempfile.TemporaryDirectory() as directory:
|
||||
store = ClientStore(Path(directory) / "state.db")
|
||||
store.initialize()
|
||||
store.accept_command(
|
||||
"stale", '{"executeStep":{"jobId":"job-1"}}',
|
||||
'{"status":"COMMAND_ACK_STATUS_ACCEPTED"}',
|
||||
)
|
||||
accepted = store.accept_reconcile_command(
|
||||
"reconcile-1", '{"reconcileJob":{"authoritativeJob":{}}}',
|
||||
'{"status":"COMMAND_ACK_STATUS_ACCEPTED"}',
|
||||
job_id="job-1", definition_json='{"jobId":"job-1"}',
|
||||
state="JOB_STATE_QUEUED", revision=0,
|
||||
last_event_sequence=0, committed=False,
|
||||
superseded_command_ids=["stale"],
|
||||
)
|
||||
self.assertFalse(accepted.duplicate)
|
||||
self.assertEqual(accepted.state, "accepted")
|
||||
self.assertTrue(store.is_command_active("reconcile-1"))
|
||||
self.assertFalse(store.is_command_active("stale"))
|
||||
self.assertEqual(
|
||||
store.job_snapshot_rows(["job-1"])[0]["last_event_sequence"], 0
|
||||
)
|
||||
duplicate = store.accept_reconcile_command(
|
||||
"reconcile-1", '{"reconcileJob":{"authoritativeJob":{}}}',
|
||||
'{"status":"COMMAND_ACK_STATUS_ACCEPTED"}',
|
||||
job_id="job-1", definition_json='{"jobId":"job-1"}',
|
||||
state="JOB_STATE_QUEUED", revision=0,
|
||||
last_event_sequence=0, committed=False,
|
||||
superseded_command_ids=["stale"],
|
||||
)
|
||||
self.assertTrue(duplicate.duplicate)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
|
||||
Reference in New Issue
Block a user