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 ### 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
View File
@@ -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
View File
@@ -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"]
+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.""" """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") 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
View File
@@ -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:
+2 -2
View File
@@ -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:
+113 -6
View File
@@ -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
View File
@@ -1,2 +1,3 @@
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
"""Generated archive_control.v1 bindings.""" """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 -*- # -*- 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 -1
View File
@@ -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 -1
View File
@@ -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 -1
View File
@@ -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
+21 -5
View File
@@ -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 -1
View File
@@ -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 -1
View File
@@ -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 -1
View File
@@ -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 -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 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 -1
View File
@@ -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 -1
View File
@@ -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 -1
View File
@@ -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 -1
View File
@@ -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 -1
View File
@@ -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 -1
View File
@@ -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 -1
View File
@@ -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 -1
View File
@@ -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
+4 -4
View File
@@ -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
View File
@@ -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),
+28
View File
@@ -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()