Compare commits

...
5 Commits
34 changed files with 889 additions and 170 deletions
+1
View File
@@ -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
+15 -14
View File
@@ -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
View File
@@ -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
View File
@@ -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
View File
@@ -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"]
+86
View File
@@ -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 -1
View File
@@ -1,4 +1,4 @@
"""Archive Control data-node daemon."""
PROTO_COMMIT = "4ec852014dad74606d4078b3ae1aa208c814b033"
PROTO_COMMIT = "03b6751d420e8e4e4cc9aeef4f76b13e28a7265f"
+7
View File
@@ -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
View File
@@ -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
View File
@@ -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,
+2 -2
View File
@@ -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:
+205 -6
View File
@@ -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
View File
@@ -1,2 +1,3 @@
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
"""Generated archive_control.v1 bindings."""
+1 -1
View File
@@ -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 -1
View File
@@ -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 -1
View File
@@ -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 -1
View File
@@ -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
+21 -5
View File
@@ -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 -1
View File
@@ -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 -1
View File
@@ -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 -1
View File
@@ -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 -1
View File
@@ -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 -1
View File
@@ -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 -1
View File
@@ -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 -1
View File
@@ -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 -1
View File
@@ -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 -1
View File
@@ -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 -1
View File
@@ -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 -1
View File
@@ -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 -1
View File
@@ -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
View File
@@ -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
View File
@@ -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),
+61
View File
@@ -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()