Compare commits
2
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
cd3a9b3c0d | ||
|
|
0311f6f084 |
@@ -263,24 +263,24 @@ control consumer, client consumer, E2E, then image publication.
|
|||||||
### Reproducible multi-platform Buildx lifecycle
|
### Reproducible multi-platform Buildx lifecycle
|
||||||
|
|
||||||
The named Buildx builders are local acceleration/cache only; they are not a
|
The named Buildx builders are local acceleration/cache only; they are not a
|
||||||
deployment dependency and can be removed after publication. To create a fresh
|
deployment dependency and can be removed after publication. Use the checked-in
|
||||||
builder, verify its platforms, publish a release, and remove it afterwards:
|
publisher: it installs only the non-native binfmt handler, creates a temporary
|
||||||
|
rootless BuildKit builder, verifies both platforms, publishes the index,
|
||||||
|
displays its digest, then removes both temporary resources.
|
||||||
|
|
||||||
```bash
|
```bash
|
||||||
docker buildx create --name archive-control-release --driver docker-container --use
|
scripts/publish-image.sh vX.Y.Z
|
||||||
docker buildx inspect --bootstrap
|
scripts/publish-image.sh vX.Y.Z --also-latest
|
||||||
docker buildx build --platform linux/amd64,linux/arm64 \
|
|
||||||
--tag sodium/archive-clients:vX.Y.Z --push .
|
|
||||||
docker buildx rm archive-control-release
|
|
||||||
```
|
```
|
||||||
|
|
||||||
`docker buildx inspect` must list both `linux/amd64` and `linux/arm64` before
|
The publisher deliberately uses `moby/buildkit:rootless` with
|
||||||
publishing. If the host has no arm64 emulation, install/configure it according
|
`--oci-worker-no-process-sandbox`. On nested Docker hosts, the default OCI
|
||||||
to the host Docker distribution before the build; do not publish a partial
|
sandbox can fail while masking `/proc/acpi` for an emulated build; rootless
|
||||||
single-platform tag. Retain the pushed manifest digest in the release notes
|
BuildKit confines that compatibility setting to the disposable builder. It
|
||||||
and deploy the immutable tag or digest. The optional `archive-control-qemu`
|
refuses to publish unless `docker buildx inspect` reports both `linux/amd64`
|
||||||
builder follows the same lifecycle when it is used for an emulation smoke
|
and `linux/arm64`, and removes the builder and binfmt handler on success,
|
||||||
build.
|
failure, or interruption. Retain the displayed manifest digest in release
|
||||||
|
notes and deploy the immutable tag or digest.
|
||||||
|
|
||||||
## Operator usage
|
## Operator usage
|
||||||
|
|
||||||
|
|||||||
+15
-12
@@ -45,8 +45,8 @@ seconds. A connection stable for 60 seconds resets the backoff.
|
|||||||
|
|
||||||
## Version and feature negotiation
|
## Version and feature negotiation
|
||||||
|
|
||||||
The envelope and registration both state the sender version. Major `1` accepts
|
The envelope and registration both state the sender version. Major `2` accepts
|
||||||
only major `1`; a different major is rejected. The negotiated minor is the
|
only major `2`; a different major is rejected. The negotiated minor is the
|
||||||
highest mutually supported minor no greater than either endpoint's advertised
|
highest mutually supported minor no greater than either endpoint's advertised
|
||||||
minor.
|
minor.
|
||||||
|
|
||||||
@@ -96,19 +96,22 @@ client state snapshot before a route can become ready.
|
|||||||
|
|
||||||
Every mutating command includes an expected job revision or immutable job
|
Every mutating command includes an expected job revision or immutable job
|
||||||
definition plus the expected last global event sequence. A stale revision or
|
definition plus the expected last global event sequence. A stale revision or
|
||||||
sequence is rejected without side effects. The control daemon sends only one
|
sequence is rejected without side effects. Assignment is represented solely by
|
||||||
active step command for a job and waits for its persisted completion event
|
the durable assignment `CommandAck`; it produces no `JobEvent`. The control
|
||||||
before commanding the next participant. This single-writer lease lets whichever
|
daemon sends only one active step command for a job and waits for its persisted
|
||||||
client owns the current step allocate the next global per-job event sequence;
|
completion event before commanding the next participant.
|
||||||
the following command starts from the sequence control has durably accepted.
|
|
||||||
|
|
||||||
## Events, progress, and reconciliation
|
## Events, progress, and reconciliation
|
||||||
|
|
||||||
Each client-originated job event has a unique ID, monotonically increasing
|
Each client-originated job event has a unique ID, monotonically increasing
|
||||||
global per-job `sequence`, and resulting job revision. Duplicate IDs/sequences
|
global per-job `sequence`, resulting job revision, and the `command_id` of its
|
||||||
are idempotent only when their complete content matches. Control never grants
|
owning `ExecuteStepCommand` or `CancelJobCommand`. Duplicate IDs/sequences are
|
||||||
concurrent event-writer leases for one job. A gap or conflicting duplicate
|
idempotent only when their complete content matches. Control verifies the
|
||||||
pauses destructive orchestration and requests snapshots.
|
client, acknowledged command lease, expected cursor, and active step before
|
||||||
|
accepting it. A gap, conflict, or stale lease triggers a durable authoritative
|
||||||
|
`ReconcileJobCommand`: the client retires only the named stale command leases
|
||||||
|
and restores that job cursor, without deleting resource data or unrelated job
|
||||||
|
state.
|
||||||
|
|
||||||
`fraction_complete` is current-step progress and
|
`fraction_complete` is current-step progress and
|
||||||
`overall_fraction_complete` is the weighted five- or three-step job progress;
|
`overall_fraction_complete` is the weighted five- or three-step job progress;
|
||||||
@@ -120,7 +123,7 @@ On registration, active-job cursors provide the client's revision, last event
|
|||||||
sequence, state, and commit flag. Reconciliation applies these rules:
|
sequence, state, and commit flag. Reconciliation applies these rules:
|
||||||
|
|
||||||
1. Equal cursors resume normal delivery.
|
1. Equal cursors resume normal delivery.
|
||||||
2. A client behind receives safe replay/snapshot commands.
|
2. A client behind receives an authoritative reconciliation command.
|
||||||
3. Control behind requests and validates the client's full job snapshot.
|
3. Control behind requests and validates the client's full job snapshot.
|
||||||
4. Conflicting commit evidence reserves the resource and requires manual
|
4. Conflicting commit evidence reserves the resource and requires manual
|
||||||
reconciliation; neither side performs cleanup.
|
reconciliation; neither side performs cleanup.
|
||||||
|
|||||||
+1
-1
@@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta"
|
|||||||
|
|
||||||
[project]
|
[project]
|
||||||
name = "archive-clients"
|
name = "archive-clients"
|
||||||
version = "0.1.19"
|
version = "0.1.21"
|
||||||
requires-python = ">=3.11"
|
requires-python = ">=3.11"
|
||||||
dependencies = ["protobuf==7.35.1", "websockets==16.0"]
|
dependencies = ["protobuf==7.35.1", "websockets==16.0"]
|
||||||
|
|
||||||
|
|||||||
Executable
+86
@@ -0,0 +1,86 @@
|
|||||||
|
#!/usr/bin/env bash
|
||||||
|
# Publish one complete Archive Control client OCI image index.
|
||||||
|
set -euo pipefail
|
||||||
|
|
||||||
|
usage() {
|
||||||
|
cat <<'EOF'
|
||||||
|
Usage: scripts/publish-image.sh <semver> [--also-latest]
|
||||||
|
|
||||||
|
Builds and pushes sodium/archive-clients for linux/amd64 and linux/arm64.
|
||||||
|
Docker must already be authenticated to the target registry.
|
||||||
|
EOF
|
||||||
|
}
|
||||||
|
|
||||||
|
if [[ $# -lt 1 || ${1:-} == '-h' || ${1:-} == '--help' ]]; then
|
||||||
|
usage
|
||||||
|
exit 2
|
||||||
|
fi
|
||||||
|
|
||||||
|
tag=$1
|
||||||
|
shift
|
||||||
|
also_latest=false
|
||||||
|
if [[ ${1:-} == '--also-latest' ]]; then
|
||||||
|
also_latest=true
|
||||||
|
shift
|
||||||
|
fi
|
||||||
|
if [[ $# -ne 0 ]]; then
|
||||||
|
usage
|
||||||
|
exit 2
|
||||||
|
fi
|
||||||
|
|
||||||
|
repository=${IMAGE_REPOSITORY:-sodium/archive-clients}
|
||||||
|
builder=${BUILDER:-archive-control-release}
|
||||||
|
buildkit_image=${BUILDKIT_IMAGE:-moby/buildkit:rootless}
|
||||||
|
# Nested Docker hosts can reject default OCI /proc mount masking. Rootless
|
||||||
|
# BuildKit confines this compatibility flag to the disposable release builder.
|
||||||
|
buildkitd_flags=${BUILDKITD_FLAGS:---oci-worker-no-process-sandbox}
|
||||||
|
binfmt_image=${BINFMT_IMAGE:-tonistiigi/binfmt}
|
||||||
|
host_arch=$(docker version --format '{{.Server.Arch}}')
|
||||||
|
case ${host_arch} in
|
||||||
|
arm64|aarch64) emulated_arch=amd64 ;;
|
||||||
|
amd64|x86_64) emulated_arch=arm64 ;;
|
||||||
|
*) echo "Unsupported Docker server architecture: ${host_arch}" >&2; exit 1 ;;
|
||||||
|
esac
|
||||||
|
|
||||||
|
builder_created=false
|
||||||
|
binfmt_installed=false
|
||||||
|
cleanup() {
|
||||||
|
local status=$?
|
||||||
|
trap - EXIT
|
||||||
|
if [[ ${builder_created} == true ]]; then
|
||||||
|
docker buildx rm "${builder}" >/dev/null 2>&1 || true
|
||||||
|
fi
|
||||||
|
if [[ ${binfmt_installed} == true ]]; then
|
||||||
|
docker run --privileged --rm "${binfmt_image}" --uninstall "${emulated_arch}" \
|
||||||
|
>/dev/null 2>&1 || true
|
||||||
|
fi
|
||||||
|
exit "${status}"
|
||||||
|
}
|
||||||
|
trap cleanup EXIT
|
||||||
|
|
||||||
|
docker run --privileged --rm "${binfmt_image}" --install "${emulated_arch}" \
|
||||||
|
>/dev/null
|
||||||
|
binfmt_installed=true
|
||||||
|
if docker buildx inspect "${builder}" >/dev/null 2>&1; then
|
||||||
|
docker buildx rm "${builder}" >/dev/null
|
||||||
|
fi
|
||||||
|
docker buildx create --name "${builder}" --driver docker-container \
|
||||||
|
--driver-opt "image=${buildkit_image}" \
|
||||||
|
--buildkitd-flags "${buildkitd_flags}" --use >/dev/null
|
||||||
|
builder_created=true
|
||||||
|
|
||||||
|
platforms=$(docker buildx inspect "${builder}" --bootstrap 2>&1)
|
||||||
|
for platform in linux/amd64 linux/arm64; do
|
||||||
|
if ! grep -Fq "${platform}" <<<"${platforms}"; then
|
||||||
|
echo "Builder ${builder} does not support ${platform}; refusing partial release." >&2
|
||||||
|
exit 1
|
||||||
|
fi
|
||||||
|
done
|
||||||
|
|
||||||
|
tags=(--tag "${repository}:${tag}")
|
||||||
|
if [[ ${also_latest} == true ]]; then
|
||||||
|
tags+=(--tag "${repository}:latest")
|
||||||
|
fi
|
||||||
|
docker buildx build --pull --platform linux/amd64,linux/arm64 \
|
||||||
|
--push "${tags[@]}" .
|
||||||
|
docker buildx imagetools inspect "${repository}:${tag}"
|
||||||
@@ -1,4 +1,4 @@
|
|||||||
"""Archive Control data-node daemon."""
|
"""Archive Control data-node daemon."""
|
||||||
|
|
||||||
PROTO_COMMIT = "4ec852014dad74606d4078b3ae1aa208c814b033"
|
PROTO_COMMIT = "03b6751d420e8e4e4cc9aeef4f76b13e28a7265f"
|
||||||
|
|
||||||
|
|||||||
@@ -214,7 +214,7 @@ class ArchiveClientDaemon:
|
|||||||
raise RuntimeError("registration response correlation mismatch")
|
raise RuntimeError("registration response correlation mismatch")
|
||||||
if response.register_response.status != client_pb2.REGISTRATION_STATUS_ACCEPTED:
|
if response.register_response.status != client_pb2.REGISTRATION_STATUS_ACCEPTED:
|
||||||
raise RuntimeError("control rejected registration")
|
raise RuntimeError("control rejected registration")
|
||||||
if response.register_response.negotiated_version.major != 1:
|
if response.register_response.negotiated_version.major != 2:
|
||||||
raise RuntimeError("control negotiated an unsupported protocol version")
|
raise RuntimeError("control negotiated an unsupported protocol version")
|
||||||
logger.info("control_connection_registered")
|
logger.info("control_connection_registered")
|
||||||
outbound: asyncio.Queue[str] = asyncio.Queue(maxsize=100)
|
outbound: asyncio.Queue[str] = asyncio.Queue(maxsize=100)
|
||||||
@@ -440,6 +440,7 @@ class ArchiveClientDaemon:
|
|||||||
if (
|
if (
|
||||||
acknowledgement.status
|
acknowledgement.status
|
||||||
== control_pb2.COMMAND_ACK_STATUS_ACCEPTED
|
== control_pb2.COMMAND_ACK_STATUS_ACCEPTED
|
||||||
|
and accepted.state == "accepted"
|
||||||
):
|
):
|
||||||
accepted_for_execution = True
|
accepted_for_execution = True
|
||||||
acknowledgement.status = (
|
acknowledgement.status = (
|
||||||
@@ -525,6 +526,22 @@ class ArchiveClientDaemon:
|
|||||||
outbound,
|
outbound,
|
||||||
command_tasks,
|
command_tasks,
|
||||||
)
|
)
|
||||||
|
elif (
|
||||||
|
accepted is not None
|
||||||
|
and accepted_for_execution
|
||||||
|
and command.WhichOneof("payload") == "reconcile_job"
|
||||||
|
):
|
||||||
|
reconciliation = command.reconcile_job
|
||||||
|
await asyncio.to_thread(
|
||||||
|
self.store.reconcile_job,
|
||||||
|
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),
|
||||||
|
)
|
||||||
|
|
||||||
async def _resume_commands(
|
async def _resume_commands(
|
||||||
self,
|
self,
|
||||||
@@ -644,6 +661,11 @@ class ArchiveClientDaemon:
|
|||||||
loop = asyncio.get_running_loop()
|
loop = asyncio.get_running_loop()
|
||||||
|
|
||||||
def emit(event):
|
def emit(event):
|
||||||
|
if not self.store.is_command_active(command.command_id):
|
||||||
|
# Reconciliation retired this lease while its blocking
|
||||||
|
# worker was still unwinding. Its durable local record is
|
||||||
|
# not allowed to re-enter the control event stream.
|
||||||
|
return
|
||||||
response = new_envelope()
|
response = new_envelope()
|
||||||
response.correlation_id = correlation_id
|
response.correlation_id = correlation_id
|
||||||
response.job_event.CopyFrom(event)
|
response.job_event.CopyFrom(event)
|
||||||
@@ -661,16 +683,20 @@ class ArchiveClientDaemon:
|
|||||||
pass
|
pass
|
||||||
|
|
||||||
events = await asyncio.to_thread(
|
events = await asyncio.to_thread(
|
||||||
self.jobs.execute, command.execute_step, emit
|
self.jobs.execute, command.execute_step, emit, command.command_id
|
||||||
)
|
)
|
||||||
streamed = True
|
streamed = True
|
||||||
elif payload == "cancel_job":
|
elif payload == "cancel_job":
|
||||||
events = await asyncio.to_thread(
|
events = await asyncio.to_thread(
|
||||||
self.jobs.cancel, command.cancel_job
|
self.jobs.cancel, command.cancel_job, command.command_id
|
||||||
)
|
)
|
||||||
else:
|
else:
|
||||||
raise JobExecutionError("job command payload is unsupported")
|
raise JobExecutionError("job command payload is unsupported")
|
||||||
for event in (() if streamed else events):
|
for event in (() if streamed else events):
|
||||||
|
if not await asyncio.to_thread(
|
||||||
|
self.store.is_command_active, command.command_id
|
||||||
|
):
|
||||||
|
return
|
||||||
response = new_envelope()
|
response = new_envelope()
|
||||||
response.correlation_id = correlation_id
|
response.correlation_id = correlation_id
|
||||||
response.job_event.CopyFrom(event)
|
response.job_event.CopyFrom(event)
|
||||||
@@ -1039,6 +1065,17 @@ class ArchiveClientDaemon:
|
|||||||
acknowledgement.error.message = "job assignment is invalid"
|
acknowledgement.error.message = "job assignment is invalid"
|
||||||
else:
|
else:
|
||||||
acknowledgement.status = control_pb2.COMMAND_ACK_STATUS_ACCEPTED
|
acknowledgement.status = control_pb2.COMMAND_ACK_STATUS_ACCEPTED
|
||||||
|
elif command.WhichOneof("payload") == "reconcile_job":
|
||||||
|
record = command.reconcile_job.authoritative_job
|
||||||
|
if (
|
||||||
|
not record.definition.job_id
|
||||||
|
or command.reconcile_job.authoritative_last_event_sequence < 0
|
||||||
|
):
|
||||||
|
acknowledgement.status = control_pb2.COMMAND_ACK_STATUS_REJECTED
|
||||||
|
acknowledgement.error.code = common_pb2.ERROR_CODE_INVALID_ARGUMENT
|
||||||
|
acknowledgement.error.message = "job reconciliation is invalid"
|
||||||
|
else:
|
||||||
|
acknowledgement.status = control_pb2.COMMAND_ACK_STATUS_ACCEPTED
|
||||||
elif command.WhichOneof("payload") == "execute_step":
|
elif command.WhichOneof("payload") == "execute_step":
|
||||||
step = command.execute_step
|
step = command.execute_step
|
||||||
if self.jobs is None:
|
if self.jobs is None:
|
||||||
|
|||||||
+17
-16
@@ -97,33 +97,25 @@ class ClientJobExecutor:
|
|||||||
) -> list[control_pb2.JobEvent]:
|
) -> list[control_pb2.JobEvent]:
|
||||||
definition = command.job
|
definition = command.job
|
||||||
self._validate_definition(definition)
|
self._validate_definition(definition)
|
||||||
replay = self._replay(
|
self.store.ensure_job_definition(
|
||||||
definition.job_id, command.expected_last_event_sequence
|
definition.job_id, encode_message(definition)
|
||||||
)
|
)
|
||||||
if replay:
|
# Assignment succeeds through CommandAck. It must not let a client
|
||||||
return replay
|
# allocate a globally ordered JobEvent cursor.
|
||||||
event = self._event(
|
return []
|
||||||
definition,
|
|
||||||
sequence=command.expected_last_event_sequence + 1,
|
|
||||||
revision=command.expected_job_revision,
|
|
||||||
event_type=control_pb2.JOB_EVENT_TYPE_ASSIGNED,
|
|
||||||
state=job_pb2.JOB_STATE_PREPARING,
|
|
||||||
committed=False,
|
|
||||||
)
|
|
||||||
self._record(definition, event)
|
|
||||||
return [event]
|
|
||||||
|
|
||||||
def execute(
|
def execute(
|
||||||
self,
|
self,
|
||||||
command: control_pb2.ExecuteStepCommand,
|
command: control_pb2.ExecuteStepCommand,
|
||||||
event_callback: Callable[[control_pb2.JobEvent], None] | None = None,
|
event_callback: Callable[[control_pb2.JobEvent], None] | None = None,
|
||||||
|
command_id: str = "",
|
||||||
) -> list[control_pb2.JobEvent]:
|
) -> list[control_pb2.JobEvent]:
|
||||||
# asyncio cancellation of a connection-bound task cannot stop the
|
# asyncio cancellation of a connection-bound task cannot stop the
|
||||||
# synchronous filesystem/qB operation already running in its worker
|
# synchronous filesystem/qB operation already running in its worker
|
||||||
# thread. A replay after reconnect therefore waits for that operation
|
# thread. A replay after reconnect therefore waits for that operation
|
||||||
# and then reads its durable event journal instead of executing twice.
|
# and then reads its durable event journal instead of executing twice.
|
||||||
with self._execution_lock(command.job_id):
|
with self._execution_lock(command.job_id):
|
||||||
return self._execute_locked(command, event_callback)
|
return self._execute_locked(command, event_callback, command_id)
|
||||||
|
|
||||||
def _execution_lock(self, job_id: str) -> threading.Lock:
|
def _execution_lock(self, job_id: str) -> threading.Lock:
|
||||||
with self._execution_locks_guard:
|
with self._execution_locks_guard:
|
||||||
@@ -133,6 +125,7 @@ class ClientJobExecutor:
|
|||||||
self,
|
self,
|
||||||
command: control_pb2.ExecuteStepCommand,
|
command: control_pb2.ExecuteStepCommand,
|
||||||
event_callback: Callable[[control_pb2.JobEvent], None] | None = None,
|
event_callback: Callable[[control_pb2.JobEvent], None] | None = None,
|
||||||
|
command_id: str = "",
|
||||||
) -> list[control_pb2.JobEvent]:
|
) -> list[control_pb2.JobEvent]:
|
||||||
definition = self._definition(command.job_id)
|
definition = self._definition(command.job_id)
|
||||||
replay = self._replay(
|
replay = self._replay(
|
||||||
@@ -186,6 +179,7 @@ class ClientJobExecutor:
|
|||||||
),
|
),
|
||||||
step=command.step,
|
step=command.step,
|
||||||
step_state=job_pb2.STEP_STATE_RUNNING,
|
step_state=job_pb2.STEP_STATE_RUNNING,
|
||||||
|
command_id=command_id,
|
||||||
)
|
)
|
||||||
self._record(definition, started)
|
self._record(definition, started)
|
||||||
cursor = started
|
cursor = started
|
||||||
@@ -228,6 +222,7 @@ class ClientJobExecutor:
|
|||||||
committed=cursor.committed,
|
committed=cursor.committed,
|
||||||
step=command.step,
|
step=command.step,
|
||||||
step_state=job_pb2.STEP_STATE_RUNNING,
|
step_state=job_pb2.STEP_STATE_RUNNING,
|
||||||
|
command_id=command_id,
|
||||||
)
|
)
|
||||||
event.progress.fraction_complete = max(
|
event.progress.fraction_complete = max(
|
||||||
0.0, min(float(fraction), 1.0)
|
0.0, min(float(fraction), 1.0)
|
||||||
@@ -258,6 +253,7 @@ class ClientJobExecutor:
|
|||||||
committed=started.committed,
|
committed=started.committed,
|
||||||
step=command.step,
|
step=command.step,
|
||||||
step_state=job_pb2.STEP_STATE_CANCELLED,
|
step_state=job_pb2.STEP_STATE_CANCELLED,
|
||||||
|
command_id=command_id,
|
||||||
)
|
)
|
||||||
cancelling.error.code = common_pb2.ERROR_CODE_CANCELLED
|
cancelling.error.code = common_pb2.ERROR_CODE_CANCELLED
|
||||||
cancelling.error.message = str(error)
|
cancelling.error.message = str(error)
|
||||||
@@ -315,6 +311,7 @@ class ClientJobExecutor:
|
|||||||
committed=started.committed,
|
committed=started.committed,
|
||||||
step=command.step,
|
step=command.step,
|
||||||
step_state=job_pb2.STEP_STATE_FAILED,
|
step_state=job_pb2.STEP_STATE_FAILED,
|
||||||
|
command_id=command_id,
|
||||||
)
|
)
|
||||||
failed.error.code = _job_error_code(error)
|
failed.error.code = _job_error_code(error)
|
||||||
failed.error.message = str(error) or type(error).__name__
|
failed.error.message = str(error) or type(error).__name__
|
||||||
@@ -357,6 +354,7 @@ class ClientJobExecutor:
|
|||||||
committed=committed,
|
committed=committed,
|
||||||
step=command.step,
|
step=command.step,
|
||||||
step_state=job_pb2.STEP_STATE_SUCCEEDED,
|
step_state=job_pb2.STEP_STATE_SUCCEEDED,
|
||||||
|
command_id=command_id,
|
||||||
)
|
)
|
||||||
if result is not None:
|
if result is not None:
|
||||||
succeeded.observed_placement.CopyFrom(result)
|
succeeded.observed_placement.CopyFrom(result)
|
||||||
@@ -367,7 +365,7 @@ class ClientJobExecutor:
|
|||||||
return emitted
|
return emitted
|
||||||
|
|
||||||
def cancel(
|
def cancel(
|
||||||
self, command: control_pb2.CancelJobCommand
|
self, command: control_pb2.CancelJobCommand, command_id: str = ""
|
||||||
) -> list[control_pb2.JobEvent]:
|
) -> list[control_pb2.JobEvent]:
|
||||||
definition = self._definition(command.job_id)
|
definition = self._definition(command.job_id)
|
||||||
self.request_cancel(command.job_id)
|
self.request_cancel(command.job_id)
|
||||||
@@ -418,6 +416,7 @@ class ClientJobExecutor:
|
|||||||
committed=committed,
|
committed=committed,
|
||||||
step=job_pb2.JOB_STEP_KIND_ROLLBACK,
|
step=job_pb2.JOB_STEP_KIND_ROLLBACK,
|
||||||
step_state=job_pb2.STEP_STATE_SUCCEEDED,
|
step_state=job_pb2.STEP_STATE_SUCCEEDED,
|
||||||
|
command_id=command_id,
|
||||||
)
|
)
|
||||||
if placement is not None:
|
if placement is not None:
|
||||||
event.observed_placement.CopyFrom(placement)
|
event.observed_placement.CopyFrom(placement)
|
||||||
@@ -1050,6 +1049,7 @@ class ClientJobExecutor:
|
|||||||
committed: bool,
|
committed: bool,
|
||||||
step: int = job_pb2.JOB_STEP_KIND_UNSPECIFIED,
|
step: int = job_pb2.JOB_STEP_KIND_UNSPECIFIED,
|
||||||
step_state: int = job_pb2.STEP_STATE_UNSPECIFIED,
|
step_state: int = job_pb2.STEP_STATE_UNSPECIFIED,
|
||||||
|
command_id: str = "",
|
||||||
) -> control_pb2.JobEvent:
|
) -> control_pb2.JobEvent:
|
||||||
event = control_pb2.JobEvent(
|
event = control_pb2.JobEvent(
|
||||||
event_id=str(uuid.uuid4()),
|
event_id=str(uuid.uuid4()),
|
||||||
@@ -1059,6 +1059,7 @@ class ClientJobExecutor:
|
|||||||
type=event_type,
|
type=event_type,
|
||||||
state=state,
|
state=state,
|
||||||
committed=committed,
|
committed=committed,
|
||||||
|
command_id=command_id,
|
||||||
)
|
)
|
||||||
event.occurred_at.GetCurrentTime()
|
event.occurred_at.GetCurrentTime()
|
||||||
if step != job_pb2.JOB_STEP_KIND_UNSPECIFIED:
|
if step != job_pb2.JOB_STEP_KIND_UNSPECIFIED:
|
||||||
|
|||||||
@@ -18,7 +18,7 @@ class ProtocolError(ValueError):
|
|||||||
|
|
||||||
def new_envelope() -> envelope_pb2.Envelope:
|
def new_envelope() -> envelope_pb2.Envelope:
|
||||||
envelope = envelope_pb2.Envelope()
|
envelope = envelope_pb2.Envelope()
|
||||||
envelope.protocol_version.major = 1
|
envelope.protocol_version.major = 2
|
||||||
envelope.message_id = str(uuid.uuid4())
|
envelope.message_id = str(uuid.uuid4())
|
||||||
envelope.sent_at.FromDatetime(datetime.now(timezone.utc))
|
envelope.sent_at.FromDatetime(datetime.now(timezone.utc))
|
||||||
return envelope
|
return envelope
|
||||||
@@ -59,7 +59,7 @@ def decode(data: str | bytes, max_bytes: int = 1024 * 1024) -> envelope_pb2.Enve
|
|||||||
json_format.ParseError,
|
json_format.ParseError,
|
||||||
) as exc:
|
) as exc:
|
||||||
raise ProtocolError("invalid control envelope") from exc
|
raise ProtocolError("invalid control envelope") from exc
|
||||||
if envelope.protocol_version.major != 1:
|
if envelope.protocol_version.major != 2:
|
||||||
raise ProtocolError("unsupported protocol major version")
|
raise ProtocolError("unsupported protocol major version")
|
||||||
_canonical_uuid(envelope.message_id, "message_id")
|
_canonical_uuid(envelope.message_id, "message_id")
|
||||||
if envelope.correlation_id:
|
if envelope.correlation_id:
|
||||||
|
|||||||
@@ -38,6 +38,7 @@ class FileOperationConflict(RuntimeError):
|
|||||||
class CommandAcceptance:
|
class CommandAcceptance:
|
||||||
duplicate: bool
|
duplicate: bool
|
||||||
acknowledgement_json: str
|
acknowledgement_json: str
|
||||||
|
state: str
|
||||||
|
|
||||||
|
|
||||||
class ClientStore:
|
class ClientStore:
|
||||||
@@ -85,7 +86,7 @@ class ClientStore:
|
|||||||
connection.execute("BEGIN IMMEDIATE")
|
connection.execute("BEGIN IMMEDIATE")
|
||||||
existing = connection.execute(
|
existing = connection.execute(
|
||||||
"""
|
"""
|
||||||
SELECT payload_sha256, command_json, acknowledgement_json
|
SELECT payload_sha256, command_json, acknowledgement_json, state
|
||||||
FROM commands WHERE command_id = ?
|
FROM commands WHERE command_id = ?
|
||||||
""",
|
""",
|
||||||
(command_id,),
|
(command_id,),
|
||||||
@@ -93,17 +94,26 @@ class ClientStore:
|
|||||||
if existing:
|
if existing:
|
||||||
if existing["payload_sha256"] != digest or existing["command_json"] != payload:
|
if existing["payload_sha256"] != digest or existing["command_json"] != payload:
|
||||||
raise CommandConflict("command ID was reused with different content")
|
raise CommandConflict("command ID was reused with different content")
|
||||||
return CommandAcceptance(True, existing["acknowledgement_json"])
|
return CommandAcceptance(
|
||||||
|
True, existing["acknowledgement_json"], existing["state"]
|
||||||
|
)
|
||||||
|
status = json.loads(acknowledgement).get("status")
|
||||||
|
state = (
|
||||||
|
"rejected"
|
||||||
|
if status is not None
|
||||||
|
and status != "COMMAND_ACK_STATUS_ACCEPTED"
|
||||||
|
else "accepted"
|
||||||
|
)
|
||||||
connection.execute(
|
connection.execute(
|
||||||
"""
|
"""
|
||||||
INSERT INTO commands (
|
INSERT INTO commands (
|
||||||
command_id, payload_sha256, command_json,
|
command_id, payload_sha256, command_json,
|
||||||
acknowledgement_json, state
|
acknowledgement_json, state
|
||||||
) VALUES (?, ?, ?, ?, 'received')
|
) VALUES (?, ?, ?, ?, ?)
|
||||||
""",
|
""",
|
||||||
(command_id, digest, payload, acknowledgement),
|
(command_id, digest, payload, acknowledgement, state),
|
||||||
)
|
)
|
||||||
return CommandAcceptance(False, acknowledgement)
|
return CommandAcceptance(False, acknowledgement, state)
|
||||||
|
|
||||||
def list_active_job_cursors(self) -> list[dict[str, object]]:
|
def list_active_job_cursors(self) -> list[dict[str, object]]:
|
||||||
with self._connect() as connection:
|
with self._connect() as connection:
|
||||||
@@ -122,7 +132,7 @@ class ClientStore:
|
|||||||
rows = connection.execute(
|
rows = connection.execute(
|
||||||
"""
|
"""
|
||||||
SELECT command_id, command_json, acknowledgement_json
|
SELECT command_id, command_json, acknowledgement_json
|
||||||
FROM commands ORDER BY rowid
|
FROM commands WHERE state = 'accepted' ORDER BY rowid
|
||||||
"""
|
"""
|
||||||
).fetchall()
|
).fetchall()
|
||||||
return [dict(row) for row in rows]
|
return [dict(row) for row in rows]
|
||||||
@@ -199,6 +209,91 @@ 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")
|
||||||
|
existing = connection.execute(
|
||||||
|
"SELECT definition_json FROM jobs WHERE job_id = ?", (job_id,)
|
||||||
|
).fetchone()
|
||||||
|
if existing is not None and existing["definition_json"] != definition:
|
||||||
|
raise JobConflict("reconciliation has a different job definition")
|
||||||
|
connection.execute(
|
||||||
|
"DELETE FROM events WHERE job_id = ? AND sequence > ?",
|
||||||
|
(job_id, last_event_sequence),
|
||||||
|
)
|
||||||
|
connection.execute(
|
||||||
|
"""
|
||||||
|
INSERT INTO jobs (job_id, definition_json, state, revision,
|
||||||
|
last_event_sequence, committed)
|
||||||
|
VALUES (?, ?, ?, ?, ?, ?)
|
||||||
|
ON CONFLICT(job_id) DO UPDATE SET
|
||||||
|
state = excluded.state, revision = excluded.revision,
|
||||||
|
last_event_sequence = excluded.last_event_sequence,
|
||||||
|
committed = excluded.committed
|
||||||
|
""",
|
||||||
|
(job_id, definition, state, revision, last_event_sequence, int(committed)),
|
||||||
|
)
|
||||||
|
if superseded_command_ids:
|
||||||
|
placeholders = ",".join("?" for _ in superseded_command_ids)
|
||||||
|
rows = connection.execute(
|
||||||
|
f"SELECT command_id, command_json FROM commands WHERE command_id IN ({placeholders})",
|
||||||
|
superseded_command_ids,
|
||||||
|
).fetchall()
|
||||||
|
owned_ids = [
|
||||||
|
row["command_id"] for row in rows
|
||||||
|
if _command_job_id(row["command_json"]) == job_id
|
||||||
|
]
|
||||||
|
if owned_ids:
|
||||||
|
owned_placeholders = ",".join("?" for _ in owned_ids)
|
||||||
|
connection.execute(
|
||||||
|
f"UPDATE commands SET state = 'superseded' WHERE command_id IN ({owned_placeholders})",
|
||||||
|
owned_ids,
|
||||||
|
)
|
||||||
|
|
||||||
def begin_file_operation(
|
def begin_file_operation(
|
||||||
self,
|
self,
|
||||||
operation_id: str,
|
operation_id: str,
|
||||||
@@ -591,6 +686,18 @@ class ClientStore:
|
|||||||
connection.close()
|
connection.close()
|
||||||
|
|
||||||
|
|
||||||
|
def _command_job_id(command_json: str) -> str:
|
||||||
|
"""Return the job target from canonical protobuf JSON, if it has one."""
|
||||||
|
command = json.loads(command_json)
|
||||||
|
if "assignJob" in command:
|
||||||
|
return str(command["assignJob"].get("job", {}).get("jobId", ""))
|
||||||
|
if "executeStep" in command:
|
||||||
|
return str(command["executeStep"].get("jobId", ""))
|
||||||
|
if "cancelJob" in command:
|
||||||
|
return str(command["cancelJob"].get("jobId", ""))
|
||||||
|
return ""
|
||||||
|
|
||||||
|
|
||||||
def _canonical(value: object) -> str:
|
def _canonical(value: object) -> str:
|
||||||
return json.dumps(value, sort_keys=True, separators=(",", ":"), ensure_ascii=False)
|
return json.dumps(value, sort_keys=True, separators=(",", ":"), ensure_ascii=False)
|
||||||
|
|
||||||
|
|||||||
@@ -1,2 +1,3 @@
|
|||||||
|
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||||
"""Generated archive_control.v1 bindings."""
|
"""Generated archive_control.v1 bindings."""
|
||||||
|
|
||||||
|
|||||||
@@ -1,4 +1,4 @@
|
|||||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||||
# -*- coding: utf-8 -*-
|
# -*- coding: utf-8 -*-
|
||||||
# Generated by the protocol buffer compiler. DO NOT EDIT!
|
# Generated by the protocol buffer compiler. DO NOT EDIT!
|
||||||
# NO CHECKED-IN PROTOBUF GENCODE
|
# NO CHECKED-IN PROTOBUF GENCODE
|
||||||
|
|||||||
@@ -1,4 +1,4 @@
|
|||||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||||
import datetime
|
import datetime
|
||||||
|
|
||||||
from archive_control.v1 import common_pb2 as _common_pb2
|
from archive_control.v1 import common_pb2 as _common_pb2
|
||||||
|
|||||||
@@ -1,4 +1,4 @@
|
|||||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||||
# -*- coding: utf-8 -*-
|
# -*- coding: utf-8 -*-
|
||||||
# Generated by the protocol buffer compiler. DO NOT EDIT!
|
# Generated by the protocol buffer compiler. DO NOT EDIT!
|
||||||
# NO CHECKED-IN PROTOBUF GENCODE
|
# NO CHECKED-IN PROTOBUF GENCODE
|
||||||
|
|||||||
@@ -1,4 +1,4 @@
|
|||||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||||
import datetime
|
import datetime
|
||||||
|
|
||||||
from google.protobuf import timestamp_pb2 as _timestamp_pb2
|
from google.protobuf import timestamp_pb2 as _timestamp_pb2
|
||||||
|
|||||||
File diff suppressed because one or more lines are too long
@@ -1,4 +1,4 @@
|
|||||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||||
import datetime
|
import datetime
|
||||||
|
|
||||||
from archive_control.v1 import common_pb2 as _common_pb2
|
from archive_control.v1 import common_pb2 as _common_pb2
|
||||||
@@ -108,8 +108,20 @@ class RequestJobSnapshotCommand(_message.Message):
|
|||||||
job_ids: _containers.RepeatedScalarFieldContainer[str]
|
job_ids: _containers.RepeatedScalarFieldContainer[str]
|
||||||
def __init__(self, job_ids: _Optional[_Iterable[str]] = ...) -> None: ...
|
def __init__(self, job_ids: _Optional[_Iterable[str]] = ...) -> None: ...
|
||||||
|
|
||||||
|
class ReconcileJobCommand(_message.Message):
|
||||||
|
__slots__ = ("authoritative_job", "authoritative_last_event_sequence", "superseded_command_ids", "reason")
|
||||||
|
AUTHORITATIVE_JOB_FIELD_NUMBER: _ClassVar[int]
|
||||||
|
AUTHORITATIVE_LAST_EVENT_SEQUENCE_FIELD_NUMBER: _ClassVar[int]
|
||||||
|
SUPERSEDED_COMMAND_IDS_FIELD_NUMBER: _ClassVar[int]
|
||||||
|
REASON_FIELD_NUMBER: _ClassVar[int]
|
||||||
|
authoritative_job: _job_pb2.JobRecord
|
||||||
|
authoritative_last_event_sequence: int
|
||||||
|
superseded_command_ids: _containers.RepeatedScalarFieldContainer[str]
|
||||||
|
reason: str
|
||||||
|
def __init__(self, authoritative_job: _Optional[_Union[_job_pb2.JobRecord, _Mapping]] = ..., authoritative_last_event_sequence: _Optional[int] = ..., superseded_command_ids: _Optional[_Iterable[str]] = ..., reason: _Optional[str] = ...) -> None: ...
|
||||||
|
|
||||||
class Command(_message.Message):
|
class Command(_message.Message):
|
||||||
__slots__ = ("command_id", "created_at", "assign_job", "execute_step", "cancel_job", "ensure_route", "inventory_query", "request_job_snapshot")
|
__slots__ = ("command_id", "created_at", "assign_job", "execute_step", "cancel_job", "ensure_route", "inventory_query", "request_job_snapshot", "reconcile_job")
|
||||||
COMMAND_ID_FIELD_NUMBER: _ClassVar[int]
|
COMMAND_ID_FIELD_NUMBER: _ClassVar[int]
|
||||||
CREATED_AT_FIELD_NUMBER: _ClassVar[int]
|
CREATED_AT_FIELD_NUMBER: _ClassVar[int]
|
||||||
ASSIGN_JOB_FIELD_NUMBER: _ClassVar[int]
|
ASSIGN_JOB_FIELD_NUMBER: _ClassVar[int]
|
||||||
@@ -118,6 +130,7 @@ class Command(_message.Message):
|
|||||||
ENSURE_ROUTE_FIELD_NUMBER: _ClassVar[int]
|
ENSURE_ROUTE_FIELD_NUMBER: _ClassVar[int]
|
||||||
INVENTORY_QUERY_FIELD_NUMBER: _ClassVar[int]
|
INVENTORY_QUERY_FIELD_NUMBER: _ClassVar[int]
|
||||||
REQUEST_JOB_SNAPSHOT_FIELD_NUMBER: _ClassVar[int]
|
REQUEST_JOB_SNAPSHOT_FIELD_NUMBER: _ClassVar[int]
|
||||||
|
RECONCILE_JOB_FIELD_NUMBER: _ClassVar[int]
|
||||||
command_id: str
|
command_id: str
|
||||||
created_at: _timestamp_pb2.Timestamp
|
created_at: _timestamp_pb2.Timestamp
|
||||||
assign_job: AssignJobCommand
|
assign_job: AssignJobCommand
|
||||||
@@ -126,7 +139,8 @@ class Command(_message.Message):
|
|||||||
ensure_route: EnsureRouteCommand
|
ensure_route: EnsureRouteCommand
|
||||||
inventory_query: _inventory_pb2.InventoryQuery
|
inventory_query: _inventory_pb2.InventoryQuery
|
||||||
request_job_snapshot: RequestJobSnapshotCommand
|
request_job_snapshot: RequestJobSnapshotCommand
|
||||||
def __init__(self, command_id: _Optional[str] = ..., created_at: _Optional[_Union[datetime.datetime, _timestamp_pb2.Timestamp, _Mapping]] = ..., assign_job: _Optional[_Union[AssignJobCommand, _Mapping]] = ..., execute_step: _Optional[_Union[ExecuteStepCommand, _Mapping]] = ..., cancel_job: _Optional[_Union[CancelJobCommand, _Mapping]] = ..., ensure_route: _Optional[_Union[EnsureRouteCommand, _Mapping]] = ..., inventory_query: _Optional[_Union[_inventory_pb2.InventoryQuery, _Mapping]] = ..., request_job_snapshot: _Optional[_Union[RequestJobSnapshotCommand, _Mapping]] = ...) -> None: ...
|
reconcile_job: ReconcileJobCommand
|
||||||
|
def __init__(self, command_id: _Optional[str] = ..., created_at: _Optional[_Union[datetime.datetime, _timestamp_pb2.Timestamp, _Mapping]] = ..., assign_job: _Optional[_Union[AssignJobCommand, _Mapping]] = ..., execute_step: _Optional[_Union[ExecuteStepCommand, _Mapping]] = ..., cancel_job: _Optional[_Union[CancelJobCommand, _Mapping]] = ..., ensure_route: _Optional[_Union[EnsureRouteCommand, _Mapping]] = ..., inventory_query: _Optional[_Union[_inventory_pb2.InventoryQuery, _Mapping]] = ..., request_job_snapshot: _Optional[_Union[RequestJobSnapshotCommand, _Mapping]] = ..., reconcile_job: _Optional[_Union[ReconcileJobCommand, _Mapping]] = ...) -> None: ...
|
||||||
|
|
||||||
class CommandAck(_message.Message):
|
class CommandAck(_message.Message):
|
||||||
__slots__ = ("command_id", "status", "error")
|
__slots__ = ("command_id", "status", "error")
|
||||||
@@ -139,7 +153,7 @@ class CommandAck(_message.Message):
|
|||||||
def __init__(self, command_id: _Optional[str] = ..., status: _Optional[_Union[CommandAckStatus, str]] = ..., error: _Optional[_Union[_common_pb2.Error, _Mapping]] = ...) -> None: ...
|
def __init__(self, command_id: _Optional[str] = ..., status: _Optional[_Union[CommandAckStatus, str]] = ..., error: _Optional[_Union[_common_pb2.Error, _Mapping]] = ...) -> None: ...
|
||||||
|
|
||||||
class JobEvent(_message.Message):
|
class JobEvent(_message.Message):
|
||||||
__slots__ = ("event_id", "job_id", "sequence", "job_revision", "type", "state", "committed", "progress", "error", "observed_resource", "observed_placement", "occurred_at")
|
__slots__ = ("event_id", "job_id", "sequence", "job_revision", "type", "state", "committed", "progress", "error", "observed_resource", "observed_placement", "occurred_at", "command_id")
|
||||||
EVENT_ID_FIELD_NUMBER: _ClassVar[int]
|
EVENT_ID_FIELD_NUMBER: _ClassVar[int]
|
||||||
JOB_ID_FIELD_NUMBER: _ClassVar[int]
|
JOB_ID_FIELD_NUMBER: _ClassVar[int]
|
||||||
SEQUENCE_FIELD_NUMBER: _ClassVar[int]
|
SEQUENCE_FIELD_NUMBER: _ClassVar[int]
|
||||||
@@ -152,6 +166,7 @@ class JobEvent(_message.Message):
|
|||||||
OBSERVED_RESOURCE_FIELD_NUMBER: _ClassVar[int]
|
OBSERVED_RESOURCE_FIELD_NUMBER: _ClassVar[int]
|
||||||
OBSERVED_PLACEMENT_FIELD_NUMBER: _ClassVar[int]
|
OBSERVED_PLACEMENT_FIELD_NUMBER: _ClassVar[int]
|
||||||
OCCURRED_AT_FIELD_NUMBER: _ClassVar[int]
|
OCCURRED_AT_FIELD_NUMBER: _ClassVar[int]
|
||||||
|
COMMAND_ID_FIELD_NUMBER: _ClassVar[int]
|
||||||
event_id: str
|
event_id: str
|
||||||
job_id: str
|
job_id: str
|
||||||
sequence: int
|
sequence: int
|
||||||
@@ -164,7 +179,8 @@ class JobEvent(_message.Message):
|
|||||||
observed_resource: _resource_pb2.ResourceStateFingerprint
|
observed_resource: _resource_pb2.ResourceStateFingerprint
|
||||||
observed_placement: _resource_pb2.Placement
|
observed_placement: _resource_pb2.Placement
|
||||||
occurred_at: _timestamp_pb2.Timestamp
|
occurred_at: _timestamp_pb2.Timestamp
|
||||||
def __init__(self, event_id: _Optional[str] = ..., job_id: _Optional[str] = ..., sequence: _Optional[int] = ..., job_revision: _Optional[int] = ..., type: _Optional[_Union[JobEventType, str]] = ..., state: _Optional[_Union[_job_pb2.JobState, str]] = ..., committed: _Optional[bool] = ..., progress: _Optional[_Union[_job_pb2.JobProgress, _Mapping]] = ..., error: _Optional[_Union[_common_pb2.Error, _Mapping]] = ..., observed_resource: _Optional[_Union[_resource_pb2.ResourceStateFingerprint, _Mapping]] = ..., observed_placement: _Optional[_Union[_resource_pb2.Placement, _Mapping]] = ..., occurred_at: _Optional[_Union[datetime.datetime, _timestamp_pb2.Timestamp, _Mapping]] = ...) -> None: ...
|
command_id: str
|
||||||
|
def __init__(self, event_id: _Optional[str] = ..., job_id: _Optional[str] = ..., sequence: _Optional[int] = ..., job_revision: _Optional[int] = ..., type: _Optional[_Union[JobEventType, str]] = ..., state: _Optional[_Union[_job_pb2.JobState, str]] = ..., committed: _Optional[bool] = ..., progress: _Optional[_Union[_job_pb2.JobProgress, _Mapping]] = ..., error: _Optional[_Union[_common_pb2.Error, _Mapping]] = ..., observed_resource: _Optional[_Union[_resource_pb2.ResourceStateFingerprint, _Mapping]] = ..., observed_placement: _Optional[_Union[_resource_pb2.Placement, _Mapping]] = ..., occurred_at: _Optional[_Union[datetime.datetime, _timestamp_pb2.Timestamp, _Mapping]] = ..., command_id: _Optional[str] = ...) -> None: ...
|
||||||
|
|
||||||
class JobSnapshot(_message.Message):
|
class JobSnapshot(_message.Message):
|
||||||
__slots__ = ("job", "last_event_sequence", "in_flight_command_ids")
|
__slots__ = ("job", "last_event_sequence", "in_flight_command_ids")
|
||||||
|
|||||||
@@ -1,4 +1,4 @@
|
|||||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||||
# -*- coding: utf-8 -*-
|
# -*- coding: utf-8 -*-
|
||||||
# Generated by the protocol buffer compiler. DO NOT EDIT!
|
# Generated by the protocol buffer compiler. DO NOT EDIT!
|
||||||
# NO CHECKED-IN PROTOBUF GENCODE
|
# NO CHECKED-IN PROTOBUF GENCODE
|
||||||
|
|||||||
@@ -1,4 +1,4 @@
|
|||||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||||
import datetime
|
import datetime
|
||||||
|
|
||||||
from archive_control.v1 import client_pb2 as _client_pb2
|
from archive_control.v1 import client_pb2 as _client_pb2
|
||||||
|
|||||||
@@ -1,4 +1,4 @@
|
|||||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||||
# -*- coding: utf-8 -*-
|
# -*- coding: utf-8 -*-
|
||||||
# Generated by the protocol buffer compiler. DO NOT EDIT!
|
# Generated by the protocol buffer compiler. DO NOT EDIT!
|
||||||
# NO CHECKED-IN PROTOBUF GENCODE
|
# NO CHECKED-IN PROTOBUF GENCODE
|
||||||
|
|||||||
@@ -1,4 +1,4 @@
|
|||||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||||
from archive_control.v1 import common_pb2 as _common_pb2
|
from archive_control.v1 import common_pb2 as _common_pb2
|
||||||
from archive_control.v1 import resource_pb2 as _resource_pb2
|
from archive_control.v1 import resource_pb2 as _resource_pb2
|
||||||
from google.protobuf.internal import containers as _containers
|
from google.protobuf.internal import containers as _containers
|
||||||
|
|||||||
@@ -1,4 +1,4 @@
|
|||||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||||
# -*- coding: utf-8 -*-
|
# -*- coding: utf-8 -*-
|
||||||
# Generated by the protocol buffer compiler. DO NOT EDIT!
|
# Generated by the protocol buffer compiler. DO NOT EDIT!
|
||||||
# NO CHECKED-IN PROTOBUF GENCODE
|
# NO CHECKED-IN PROTOBUF GENCODE
|
||||||
|
|||||||
@@ -1,4 +1,4 @@
|
|||||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||||
import datetime
|
import datetime
|
||||||
|
|
||||||
from archive_control.v1 import common_pb2 as _common_pb2
|
from archive_control.v1 import common_pb2 as _common_pb2
|
||||||
|
|||||||
@@ -1,4 +1,4 @@
|
|||||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||||
# -*- coding: utf-8 -*-
|
# -*- coding: utf-8 -*-
|
||||||
# Generated by the protocol buffer compiler. DO NOT EDIT!
|
# Generated by the protocol buffer compiler. DO NOT EDIT!
|
||||||
# NO CHECKED-IN PROTOBUF GENCODE
|
# NO CHECKED-IN PROTOBUF GENCODE
|
||||||
|
|||||||
@@ -1,4 +1,4 @@
|
|||||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||||
import datetime
|
import datetime
|
||||||
|
|
||||||
from google.protobuf import timestamp_pb2 as _timestamp_pb2
|
from google.protobuf import timestamp_pb2 as _timestamp_pb2
|
||||||
|
|||||||
@@ -1,4 +1,4 @@
|
|||||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||||
# -*- coding: utf-8 -*-
|
# -*- coding: utf-8 -*-
|
||||||
# Generated by the protocol buffer compiler. DO NOT EDIT!
|
# Generated by the protocol buffer compiler. DO NOT EDIT!
|
||||||
# NO CHECKED-IN PROTOBUF GENCODE
|
# NO CHECKED-IN PROTOBUF GENCODE
|
||||||
|
|||||||
@@ -1,4 +1,4 @@
|
|||||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||||
import datetime
|
import datetime
|
||||||
|
|
||||||
from google.protobuf import timestamp_pb2 as _timestamp_pb2
|
from google.protobuf import timestamp_pb2 as _timestamp_pb2
|
||||||
|
|||||||
@@ -1,4 +1,4 @@
|
|||||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||||
# -*- coding: utf-8 -*-
|
# -*- coding: utf-8 -*-
|
||||||
# Generated by the protocol buffer compiler. DO NOT EDIT!
|
# Generated by the protocol buffer compiler. DO NOT EDIT!
|
||||||
# NO CHECKED-IN PROTOBUF GENCODE
|
# NO CHECKED-IN PROTOBUF GENCODE
|
||||||
|
|||||||
@@ -1,4 +1,4 @@
|
|||||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||||
import datetime
|
import datetime
|
||||||
|
|
||||||
from archive_control.v1 import resource_pb2 as _resource_pb2
|
from archive_control.v1 import resource_pb2 as _resource_pb2
|
||||||
|
|||||||
@@ -74,7 +74,7 @@ class DaemonTransportTests(unittest.IsolatedAsyncioTestCase):
|
|||||||
response.register_response.status = (
|
response.register_response.status = (
|
||||||
client_pb2.REGISTRATION_STATUS_ACCEPTED
|
client_pb2.REGISTRATION_STATUS_ACCEPTED
|
||||||
)
|
)
|
||||||
response.register_response.negotiated_version.major = 1
|
response.register_response.negotiated_version.major = 2
|
||||||
await websocket.send(encode(response))
|
await websocket.send(encode(response))
|
||||||
# Deliberately keep TCP/WebSocket open but send no application
|
# Deliberately keep TCP/WebSocket open but send no application
|
||||||
# heartbeats. This models a stale proxy/server-side session.
|
# heartbeats. This models a stale proxy/server-side session.
|
||||||
@@ -120,7 +120,7 @@ class DaemonTransportTests(unittest.IsolatedAsyncioTestCase):
|
|||||||
response.register_response.status = (
|
response.register_response.status = (
|
||||||
client_pb2.REGISTRATION_STATUS_ACCEPTED
|
client_pb2.REGISTRATION_STATUS_ACCEPTED
|
||||||
)
|
)
|
||||||
response.register_response.negotiated_version.major = 1
|
response.register_response.negotiated_version.major = 2
|
||||||
await websocket.send(encode(response))
|
await websocket.send(encode(response))
|
||||||
await websocket.wait_closed()
|
await websocket.wait_closed()
|
||||||
|
|
||||||
@@ -158,7 +158,7 @@ class DaemonTransportTests(unittest.IsolatedAsyncioTestCase):
|
|||||||
response.register_response.status = (
|
response.register_response.status = (
|
||||||
client_pb2.REGISTRATION_STATUS_ACCEPTED
|
client_pb2.REGISTRATION_STATUS_ACCEPTED
|
||||||
)
|
)
|
||||||
response.register_response.negotiated_version.major = 1
|
response.register_response.negotiated_version.major = 2
|
||||||
await websocket.send(encode(response))
|
await websocket.send(encode(response))
|
||||||
for sequence in range(1, 102):
|
for sequence in range(1, 102):
|
||||||
heartbeat = new_envelope()
|
heartbeat = new_envelope()
|
||||||
@@ -485,7 +485,7 @@ class DaemonTransportTests(unittest.IsolatedAsyncioTestCase):
|
|||||||
response = new_envelope()
|
response = new_envelope()
|
||||||
response.correlation_id = registration.message_id
|
response.correlation_id = registration.message_id
|
||||||
response.register_response.status = client_pb2.REGISTRATION_STATUS_ACCEPTED
|
response.register_response.status = client_pb2.REGISTRATION_STATUS_ACCEPTED
|
||||||
response.register_response.negotiated_version.major = 1
|
response.register_response.negotiated_version.major = 2
|
||||||
await websocket.send(encode(response))
|
await websocket.send(encode(response))
|
||||||
heartbeat = new_envelope()
|
heartbeat = new_envelope()
|
||||||
heartbeat.heartbeat.sequence = 7
|
heartbeat.heartbeat.sequence = 7
|
||||||
|
|||||||
+7
-7
@@ -436,19 +436,19 @@ class ClientJobHappyPathTests(unittest.TestCase):
|
|||||||
|
|
||||||
source_assigned = source.assign(control_pb2.AssignJobCommand(
|
source_assigned = source.assign(control_pb2.AssignJobCommand(
|
||||||
job=definition,
|
job=definition,
|
||||||
expected_job_revision=1,
|
expected_job_revision=0,
|
||||||
expected_last_event_sequence=0,
|
expected_last_event_sequence=0,
|
||||||
))
|
))
|
||||||
target_assigned = target.assign(control_pb2.AssignJobCommand(
|
target_assigned = target.assign(control_pb2.AssignJobCommand(
|
||||||
job=definition,
|
job=definition,
|
||||||
expected_job_revision=1,
|
expected_job_revision=0,
|
||||||
expected_last_event_sequence=1,
|
expected_last_event_sequence=0,
|
||||||
))
|
))
|
||||||
self.assertEqual(source_assigned[0].sequence, 1)
|
self.assertEqual(source_assigned, [])
|
||||||
self.assertEqual(target_assigned[0].sequence, 2)
|
self.assertEqual(target_assigned, [])
|
||||||
|
|
||||||
cursor_revision = 1
|
cursor_revision = 0
|
||||||
cursor_sequence = 2
|
cursor_sequence = 0
|
||||||
pipeline = (
|
pipeline = (
|
||||||
(source, job_pb2.JOB_STEP_KIND_SOURCE_STAGE),
|
(source, job_pb2.JOB_STEP_KIND_SOURCE_STAGE),
|
||||||
(target, job_pb2.JOB_STEP_KIND_SYNCTHING_TRANSFER),
|
(target, job_pb2.JOB_STEP_KIND_SYNCTHING_TRANSFER),
|
||||||
|
|||||||
@@ -177,6 +177,34 @@ class ClientStoreTests(unittest.TestCase):
|
|||||||
"job-1", "baseline", {"selected": [2]}
|
"job-1", "baseline", {"selected": [2]}
|
||||||
)
|
)
|
||||||
|
|
||||||
|
def test_reconciliation_retires_only_stale_job_leases(self):
|
||||||
|
with tempfile.TemporaryDirectory() as directory:
|
||||||
|
store = ClientStore(Path(directory) / "state.db")
|
||||||
|
store.initialize()
|
||||||
|
stale = store.accept_command(
|
||||||
|
"stale", '{"executeStep":{"jobId":"job-1"}}',
|
||||||
|
'{"status":"COMMAND_ACK_STATUS_ACCEPTED"}',
|
||||||
|
)
|
||||||
|
store.accept_command(
|
||||||
|
"other", '{"executeStep":{"jobId":"job-2"}}',
|
||||||
|
'{"status":"COMMAND_ACK_STATUS_ACCEPTED"}',
|
||||||
|
)
|
||||||
|
self.assertEqual(stale.state, "accepted")
|
||||||
|
store.save_job(
|
||||||
|
"job-1", '{"jobId":"job-1"}', "JOB_STATE_RUNNING", 3, 3, False,
|
||||||
|
)
|
||||||
|
store.reconcile_job(
|
||||||
|
job_id="job-1", definition_json='{"jobId":"job-1"}',
|
||||||
|
state="JOB_STATE_QUEUED", revision=0,
|
||||||
|
last_event_sequence=0, committed=False,
|
||||||
|
superseded_command_ids=["stale"],
|
||||||
|
)
|
||||||
|
self.assertFalse(store.is_command_active("stale"))
|
||||||
|
self.assertTrue(store.is_command_active("other"))
|
||||||
|
row = store.job_snapshot_rows(["job-1"])[0]
|
||||||
|
self.assertEqual(row["last_event_sequence"], 0)
|
||||||
|
self.assertEqual(row["state"], "JOB_STATE_QUEUED")
|
||||||
|
|
||||||
|
|
||||||
if __name__ == "__main__":
|
if __name__ == "__main__":
|
||||||
unittest.main()
|
unittest.main()
|
||||||
|
|||||||
Reference in New Issue
Block a user