Compare commits

..
2 Commits
31 changed files with 388 additions and 107 deletions
+14 -14
View File
@@ -263,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.
+1 -1
View File
@@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta"
[project]
name = "archive-clients"
version = "0.1.19"
version = "0.1.21"
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"
+40 -3
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)
@@ -440,6 +440,7 @@ class ArchiveClientDaemon:
if (
acknowledgement.status
== control_pb2.COMMAND_ACK_STATUS_ACCEPTED
and accepted.state == "accepted"
):
accepted_for_execution = True
acknowledgement.status = (
@@ -525,6 +526,22 @@ class ArchiveClientDaemon:
outbound,
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(
self,
@@ -644,6 +661,11 @@ class ArchiveClientDaemon:
loop = asyncio.get_running_loop()
def emit(event):
if not self.store.is_command_active(command.command_id):
# Reconciliation retired this lease while its blocking
# worker was still unwinding. Its durable local record is
# not allowed to re-enter the control event stream.
return
response = new_envelope()
response.correlation_id = correlation_id
response.job_event.CopyFrom(event)
@@ -661,16 +683,20 @@ class ArchiveClientDaemon:
pass
events = await asyncio.to_thread(
self.jobs.execute, command.execute_step, emit
self.jobs.execute, command.execute_step, emit, command.command_id
)
streamed = True
elif payload == "cancel_job":
events = await asyncio.to_thread(
self.jobs.cancel, command.cancel_job
self.jobs.cancel, command.cancel_job, command.command_id
)
else:
raise JobExecutionError("job command payload is unsupported")
for event in (() if streamed else events):
if not await asyncio.to_thread(
self.store.is_command_active, command.command_id
):
return
response = new_envelope()
response.correlation_id = correlation_id
response.job_event.CopyFrom(event)
@@ -1039,6 +1065,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:
+17 -16
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)
@@ -1050,6 +1049,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 +1059,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:
+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:
+113 -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,26 @@ 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 list_active_job_cursors(self) -> list[dict[str, object]]:
with self._connect() as connection:
@@ -122,7 +132,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 +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(
self,
operation_id: str,
@@ -591,6 +686,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
+4 -4
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.
@@ -120,7 +120,7 @@ class DaemonTransportTests(unittest.IsolatedAsyncioTestCase):
response.register_response.status = (
client_pb2.REGISTRATION_STATUS_ACCEPTED
)
response.register_response.negotiated_version.major = 1
response.register_response.negotiated_version.major = 2
await websocket.send(encode(response))
await websocket.wait_closed()
@@ -158,7 +158,7 @@ class DaemonTransportTests(unittest.IsolatedAsyncioTestCase):
response.register_response.status = (
client_pb2.REGISTRATION_STATUS_ACCEPTED
)
response.register_response.negotiated_version.major = 1
response.register_response.negotiated_version.major = 2
await websocket.send(encode(response))
for sequence in range(1, 102):
heartbeat = new_envelope()
@@ -485,7 +485,7 @@ class DaemonTransportTests(unittest.IsolatedAsyncioTestCase):
response = new_envelope()
response.correlation_id = registration.message_id
response.register_response.status = client_pb2.REGISTRATION_STATUS_ACCEPTED
response.register_response.negotiated_version.major = 1
response.register_response.negotiated_version.major = 2
await websocket.send(encode(response))
heartbeat = new_envelope()
heartbeat.heartbeat.sequence = 7
+7 -7
View File
@@ -436,19 +436,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),
+28
View File
@@ -177,6 +177,34 @@ 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")
if __name__ == "__main__":
unittest.main()