fix: bind client job events to command leases
This commit is contained in:
+15
-12
@@ -45,8 +45,8 @@ seconds. A connection stable for 60 seconds resets the backoff.
|
||||
|
||||
## Version and feature negotiation
|
||||
|
||||
The envelope and registration both state the sender version. Major `1` accepts
|
||||
only major `1`; a different major is rejected. The negotiated minor is the
|
||||
The envelope and registration both state the sender version. Major `2` accepts
|
||||
only major `2`; a different major is rejected. The negotiated minor is the
|
||||
highest mutually supported minor no greater than either endpoint's advertised
|
||||
minor.
|
||||
|
||||
@@ -96,19 +96,22 @@ client state snapshot before a route can become ready.
|
||||
|
||||
Every mutating command includes an expected job revision or immutable job
|
||||
definition plus the expected last global event sequence. A stale revision or
|
||||
sequence is rejected without side effects. The control daemon sends only one
|
||||
active step command for a job and waits for its persisted completion event
|
||||
before commanding the next participant. This single-writer lease lets whichever
|
||||
client owns the current step allocate the next global per-job event sequence;
|
||||
the following command starts from the sequence control has durably accepted.
|
||||
sequence is rejected without side effects. Assignment is represented solely by
|
||||
the durable assignment `CommandAck`; it produces no `JobEvent`. The control
|
||||
daemon sends only one active step command for a job and waits for its persisted
|
||||
completion event before commanding the next participant.
|
||||
|
||||
## Events, progress, and reconciliation
|
||||
|
||||
Each client-originated job event has a unique ID, monotonically increasing
|
||||
global per-job `sequence`, and resulting job revision. Duplicate IDs/sequences
|
||||
are idempotent only when their complete content matches. Control never grants
|
||||
concurrent event-writer leases for one job. A gap or conflicting duplicate
|
||||
pauses destructive orchestration and requests snapshots.
|
||||
global per-job `sequence`, resulting job revision, and the `command_id` of its
|
||||
owning `ExecuteStepCommand` or `CancelJobCommand`. Duplicate IDs/sequences are
|
||||
idempotent only when their complete content matches. Control verifies the
|
||||
client, acknowledged command lease, expected cursor, and active step before
|
||||
accepting it. A gap, conflict, or stale lease triggers a durable authoritative
|
||||
`ReconcileJobCommand`: the client retires only the named stale command leases
|
||||
and restores that job cursor, without deleting resource data or unrelated job
|
||||
state.
|
||||
|
||||
`fraction_complete` is current-step progress and
|
||||
`overall_fraction_complete` is the weighted five- or three-step job progress;
|
||||
@@ -120,7 +123,7 @@ On registration, active-job cursors provide the client's revision, last event
|
||||
sequence, state, and commit flag. Reconciliation applies these rules:
|
||||
|
||||
1. Equal cursors resume normal delivery.
|
||||
2. A client behind receives safe replay/snapshot commands.
|
||||
2. A client behind receives an authoritative reconciliation command.
|
||||
3. Control behind requests and validates the client's full job snapshot.
|
||||
4. Conflicting commit evidence reserves the resource and requires manual
|
||||
reconciliation; neither side performs cleanup.
|
||||
|
||||
+1
-1
@@ -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"]
|
||||
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
"""Archive Control data-node daemon."""
|
||||
|
||||
PROTO_COMMIT = "4ec852014dad74606d4078b3ae1aa208c814b033"
|
||||
PROTO_COMMIT = "03b6751d420e8e4e4cc9aeef4f76b13e28a7265f"
|
||||
|
||||
|
||||
@@ -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
@@ -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:
|
||||
|
||||
@@ -18,7 +18,7 @@ class ProtocolError(ValueError):
|
||||
|
||||
def new_envelope() -> envelope_pb2.Envelope:
|
||||
envelope = envelope_pb2.Envelope()
|
||||
envelope.protocol_version.major = 1
|
||||
envelope.protocol_version.major = 2
|
||||
envelope.message_id = str(uuid.uuid4())
|
||||
envelope.sent_at.FromDatetime(datetime.now(timezone.utc))
|
||||
return envelope
|
||||
@@ -59,7 +59,7 @@ def decode(data: str | bytes, max_bytes: int = 1024 * 1024) -> envelope_pb2.Enve
|
||||
json_format.ParseError,
|
||||
) as exc:
|
||||
raise ProtocolError("invalid control envelope") from exc
|
||||
if envelope.protocol_version.major != 1:
|
||||
if envelope.protocol_version.major != 2:
|
||||
raise ProtocolError("unsupported protocol major version")
|
||||
_canonical_uuid(envelope.message_id, "message_id")
|
||||
if envelope.correlation_id:
|
||||
|
||||
@@ -38,6 +38,7 @@ class FileOperationConflict(RuntimeError):
|
||||
class CommandAcceptance:
|
||||
duplicate: bool
|
||||
acknowledgement_json: str
|
||||
state: str
|
||||
|
||||
|
||||
class ClientStore:
|
||||
@@ -85,7 +86,7 @@ class ClientStore:
|
||||
connection.execute("BEGIN IMMEDIATE")
|
||||
existing = connection.execute(
|
||||
"""
|
||||
SELECT payload_sha256, command_json, acknowledgement_json
|
||||
SELECT payload_sha256, command_json, acknowledgement_json, state
|
||||
FROM commands WHERE command_id = ?
|
||||
""",
|
||||
(command_id,),
|
||||
@@ -93,17 +94,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,2 +1,3 @@
|
||||
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||
"""Generated archive_control.v1 bindings."""
|
||||
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
||||
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||
# -*- coding: utf-8 -*-
|
||||
# Generated by the protocol buffer compiler. DO NOT EDIT!
|
||||
# NO CHECKED-IN PROTOBUF GENCODE
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
||||
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||
import datetime
|
||||
|
||||
from archive_control.v1 import common_pb2 as _common_pb2
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
||||
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||
# -*- coding: utf-8 -*-
|
||||
# Generated by the protocol buffer compiler. DO NOT EDIT!
|
||||
# NO CHECKED-IN PROTOBUF GENCODE
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
||||
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||
import datetime
|
||||
|
||||
from google.protobuf import timestamp_pb2 as _timestamp_pb2
|
||||
|
||||
File diff suppressed because one or more lines are too long
@@ -1,4 +1,4 @@
|
||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
||||
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||
import datetime
|
||||
|
||||
from archive_control.v1 import common_pb2 as _common_pb2
|
||||
@@ -108,8 +108,20 @@ class RequestJobSnapshotCommand(_message.Message):
|
||||
job_ids: _containers.RepeatedScalarFieldContainer[str]
|
||||
def __init__(self, job_ids: _Optional[_Iterable[str]] = ...) -> None: ...
|
||||
|
||||
class ReconcileJobCommand(_message.Message):
|
||||
__slots__ = ("authoritative_job", "authoritative_last_event_sequence", "superseded_command_ids", "reason")
|
||||
AUTHORITATIVE_JOB_FIELD_NUMBER: _ClassVar[int]
|
||||
AUTHORITATIVE_LAST_EVENT_SEQUENCE_FIELD_NUMBER: _ClassVar[int]
|
||||
SUPERSEDED_COMMAND_IDS_FIELD_NUMBER: _ClassVar[int]
|
||||
REASON_FIELD_NUMBER: _ClassVar[int]
|
||||
authoritative_job: _job_pb2.JobRecord
|
||||
authoritative_last_event_sequence: int
|
||||
superseded_command_ids: _containers.RepeatedScalarFieldContainer[str]
|
||||
reason: str
|
||||
def __init__(self, authoritative_job: _Optional[_Union[_job_pb2.JobRecord, _Mapping]] = ..., authoritative_last_event_sequence: _Optional[int] = ..., superseded_command_ids: _Optional[_Iterable[str]] = ..., reason: _Optional[str] = ...) -> None: ...
|
||||
|
||||
class Command(_message.Message):
|
||||
__slots__ = ("command_id", "created_at", "assign_job", "execute_step", "cancel_job", "ensure_route", "inventory_query", "request_job_snapshot")
|
||||
__slots__ = ("command_id", "created_at", "assign_job", "execute_step", "cancel_job", "ensure_route", "inventory_query", "request_job_snapshot", "reconcile_job")
|
||||
COMMAND_ID_FIELD_NUMBER: _ClassVar[int]
|
||||
CREATED_AT_FIELD_NUMBER: _ClassVar[int]
|
||||
ASSIGN_JOB_FIELD_NUMBER: _ClassVar[int]
|
||||
@@ -118,6 +130,7 @@ class Command(_message.Message):
|
||||
ENSURE_ROUTE_FIELD_NUMBER: _ClassVar[int]
|
||||
INVENTORY_QUERY_FIELD_NUMBER: _ClassVar[int]
|
||||
REQUEST_JOB_SNAPSHOT_FIELD_NUMBER: _ClassVar[int]
|
||||
RECONCILE_JOB_FIELD_NUMBER: _ClassVar[int]
|
||||
command_id: str
|
||||
created_at: _timestamp_pb2.Timestamp
|
||||
assign_job: AssignJobCommand
|
||||
@@ -126,7 +139,8 @@ class Command(_message.Message):
|
||||
ensure_route: EnsureRouteCommand
|
||||
inventory_query: _inventory_pb2.InventoryQuery
|
||||
request_job_snapshot: RequestJobSnapshotCommand
|
||||
def __init__(self, command_id: _Optional[str] = ..., created_at: _Optional[_Union[datetime.datetime, _timestamp_pb2.Timestamp, _Mapping]] = ..., assign_job: _Optional[_Union[AssignJobCommand, _Mapping]] = ..., execute_step: _Optional[_Union[ExecuteStepCommand, _Mapping]] = ..., cancel_job: _Optional[_Union[CancelJobCommand, _Mapping]] = ..., ensure_route: _Optional[_Union[EnsureRouteCommand, _Mapping]] = ..., inventory_query: _Optional[_Union[_inventory_pb2.InventoryQuery, _Mapping]] = ..., request_job_snapshot: _Optional[_Union[RequestJobSnapshotCommand, _Mapping]] = ...) -> None: ...
|
||||
reconcile_job: ReconcileJobCommand
|
||||
def __init__(self, command_id: _Optional[str] = ..., created_at: _Optional[_Union[datetime.datetime, _timestamp_pb2.Timestamp, _Mapping]] = ..., assign_job: _Optional[_Union[AssignJobCommand, _Mapping]] = ..., execute_step: _Optional[_Union[ExecuteStepCommand, _Mapping]] = ..., cancel_job: _Optional[_Union[CancelJobCommand, _Mapping]] = ..., ensure_route: _Optional[_Union[EnsureRouteCommand, _Mapping]] = ..., inventory_query: _Optional[_Union[_inventory_pb2.InventoryQuery, _Mapping]] = ..., request_job_snapshot: _Optional[_Union[RequestJobSnapshotCommand, _Mapping]] = ..., reconcile_job: _Optional[_Union[ReconcileJobCommand, _Mapping]] = ...) -> None: ...
|
||||
|
||||
class CommandAck(_message.Message):
|
||||
__slots__ = ("command_id", "status", "error")
|
||||
@@ -139,7 +153,7 @@ class CommandAck(_message.Message):
|
||||
def __init__(self, command_id: _Optional[str] = ..., status: _Optional[_Union[CommandAckStatus, str]] = ..., error: _Optional[_Union[_common_pb2.Error, _Mapping]] = ...) -> None: ...
|
||||
|
||||
class JobEvent(_message.Message):
|
||||
__slots__ = ("event_id", "job_id", "sequence", "job_revision", "type", "state", "committed", "progress", "error", "observed_resource", "observed_placement", "occurred_at")
|
||||
__slots__ = ("event_id", "job_id", "sequence", "job_revision", "type", "state", "committed", "progress", "error", "observed_resource", "observed_placement", "occurred_at", "command_id")
|
||||
EVENT_ID_FIELD_NUMBER: _ClassVar[int]
|
||||
JOB_ID_FIELD_NUMBER: _ClassVar[int]
|
||||
SEQUENCE_FIELD_NUMBER: _ClassVar[int]
|
||||
@@ -152,6 +166,7 @@ class JobEvent(_message.Message):
|
||||
OBSERVED_RESOURCE_FIELD_NUMBER: _ClassVar[int]
|
||||
OBSERVED_PLACEMENT_FIELD_NUMBER: _ClassVar[int]
|
||||
OCCURRED_AT_FIELD_NUMBER: _ClassVar[int]
|
||||
COMMAND_ID_FIELD_NUMBER: _ClassVar[int]
|
||||
event_id: str
|
||||
job_id: str
|
||||
sequence: int
|
||||
@@ -164,7 +179,8 @@ class JobEvent(_message.Message):
|
||||
observed_resource: _resource_pb2.ResourceStateFingerprint
|
||||
observed_placement: _resource_pb2.Placement
|
||||
occurred_at: _timestamp_pb2.Timestamp
|
||||
def __init__(self, event_id: _Optional[str] = ..., job_id: _Optional[str] = ..., sequence: _Optional[int] = ..., job_revision: _Optional[int] = ..., type: _Optional[_Union[JobEventType, str]] = ..., state: _Optional[_Union[_job_pb2.JobState, str]] = ..., committed: _Optional[bool] = ..., progress: _Optional[_Union[_job_pb2.JobProgress, _Mapping]] = ..., error: _Optional[_Union[_common_pb2.Error, _Mapping]] = ..., observed_resource: _Optional[_Union[_resource_pb2.ResourceStateFingerprint, _Mapping]] = ..., observed_placement: _Optional[_Union[_resource_pb2.Placement, _Mapping]] = ..., occurred_at: _Optional[_Union[datetime.datetime, _timestamp_pb2.Timestamp, _Mapping]] = ...) -> None: ...
|
||||
command_id: str
|
||||
def __init__(self, event_id: _Optional[str] = ..., job_id: _Optional[str] = ..., sequence: _Optional[int] = ..., job_revision: _Optional[int] = ..., type: _Optional[_Union[JobEventType, str]] = ..., state: _Optional[_Union[_job_pb2.JobState, str]] = ..., committed: _Optional[bool] = ..., progress: _Optional[_Union[_job_pb2.JobProgress, _Mapping]] = ..., error: _Optional[_Union[_common_pb2.Error, _Mapping]] = ..., observed_resource: _Optional[_Union[_resource_pb2.ResourceStateFingerprint, _Mapping]] = ..., observed_placement: _Optional[_Union[_resource_pb2.Placement, _Mapping]] = ..., occurred_at: _Optional[_Union[datetime.datetime, _timestamp_pb2.Timestamp, _Mapping]] = ..., command_id: _Optional[str] = ...) -> None: ...
|
||||
|
||||
class JobSnapshot(_message.Message):
|
||||
__slots__ = ("job", "last_event_sequence", "in_flight_command_ids")
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
||||
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||
# -*- coding: utf-8 -*-
|
||||
# Generated by the protocol buffer compiler. DO NOT EDIT!
|
||||
# NO CHECKED-IN PROTOBUF GENCODE
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
||||
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||
import datetime
|
||||
|
||||
from archive_control.v1 import client_pb2 as _client_pb2
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
||||
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||
# -*- coding: utf-8 -*-
|
||||
# Generated by the protocol buffer compiler. DO NOT EDIT!
|
||||
# NO CHECKED-IN PROTOBUF GENCODE
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
||||
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||
from archive_control.v1 import common_pb2 as _common_pb2
|
||||
from archive_control.v1 import resource_pb2 as _resource_pb2
|
||||
from google.protobuf.internal import containers as _containers
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
||||
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||
# -*- coding: utf-8 -*-
|
||||
# Generated by the protocol buffer compiler. DO NOT EDIT!
|
||||
# NO CHECKED-IN PROTOBUF GENCODE
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
||||
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||
import datetime
|
||||
|
||||
from archive_control.v1 import common_pb2 as _common_pb2
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
||||
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||
# -*- coding: utf-8 -*-
|
||||
# Generated by the protocol buffer compiler. DO NOT EDIT!
|
||||
# NO CHECKED-IN PROTOBUF GENCODE
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
||||
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||
import datetime
|
||||
|
||||
from google.protobuf import timestamp_pb2 as _timestamp_pb2
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
||||
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||
# -*- coding: utf-8 -*-
|
||||
# Generated by the protocol buffer compiler. DO NOT EDIT!
|
||||
# NO CHECKED-IN PROTOBUF GENCODE
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
||||
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||
import datetime
|
||||
|
||||
from google.protobuf import timestamp_pb2 as _timestamp_pb2
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
||||
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||
# -*- coding: utf-8 -*-
|
||||
# Generated by the protocol buffer compiler. DO NOT EDIT!
|
||||
# NO CHECKED-IN PROTOBUF GENCODE
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
# archive-control-proto commit: 4ec852014dad74606d4078b3ae1aa208c814b033
|
||||
# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f
|
||||
import datetime
|
||||
|
||||
from archive_control.v1 import resource_pb2 as _resource_pb2
|
||||
|
||||
@@ -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
@@ -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),
|
||||
|
||||
@@ -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()
|
||||
|
||||
Reference in New Issue
Block a user