From 0311f6f084fc7ddf7e7540c988335d0483fdf91a Mon Sep 17 00:00:00 2001 From: Cabbagec Date: Thu, 13 Aug 2026 07:09:03 +0000 Subject: [PATCH] fix: bind client job events to command leases --- docs/protocol.md | 27 ++--- pyproject.toml | 2 +- src/archive_clients/__init__.py | 2 +- src/archive_clients/daemon.py | 43 +++++++- src/archive_clients/jobs.py | 33 ++++--- src/archive_clients/protocol.py | 4 +- src/archive_clients/state.py | 119 +++++++++++++++++++++-- src/archive_control/v1/__init__.py | 1 + src/archive_control/v1/client_pb2.py | 2 +- src/archive_control/v1/client_pb2.pyi | 2 +- src/archive_control/v1/common_pb2.py | 2 +- src/archive_control/v1/common_pb2.pyi | 2 +- src/archive_control/v1/control_pb2.py | 42 ++++---- src/archive_control/v1/control_pb2.pyi | 26 ++++- src/archive_control/v1/envelope_pb2.py | 2 +- src/archive_control/v1/envelope_pb2.pyi | 2 +- src/archive_control/v1/inventory_pb2.py | 2 +- src/archive_control/v1/inventory_pb2.pyi | 2 +- src/archive_control/v1/job_pb2.py | 2 +- src/archive_control/v1/job_pb2.pyi | 2 +- src/archive_control/v1/resource_pb2.py | 2 +- src/archive_control/v1/resource_pb2.pyi | 2 +- src/archive_control/v1/route_pb2.py | 2 +- src/archive_control/v1/route_pb2.pyi | 2 +- src/archive_control/v1/transfer_pb2.py | 2 +- src/archive_control/v1/transfer_pb2.pyi | 2 +- tests/test_daemon.py | 8 +- tests/test_jobs.py | 14 +-- tests/test_state.py | 28 ++++++ 29 files changed, 288 insertions(+), 93 deletions(-) diff --git a/docs/protocol.md b/docs/protocol.md index 69163d1..fb9b139 100644 --- a/docs/protocol.md +++ b/docs/protocol.md @@ -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. diff --git a/pyproject.toml b/pyproject.toml index e07c548..6f2ae30 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -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"] diff --git a/src/archive_clients/__init__.py b/src/archive_clients/__init__.py index 0458cb7..da9e69f 100644 --- a/src/archive_clients/__init__.py +++ b/src/archive_clients/__init__.py @@ -1,4 +1,4 @@ """Archive Control data-node daemon.""" -PROTO_COMMIT = "4ec852014dad74606d4078b3ae1aa208c814b033" +PROTO_COMMIT = "03b6751d420e8e4e4cc9aeef4f76b13e28a7265f" diff --git a/src/archive_clients/daemon.py b/src/archive_clients/daemon.py index c2dec05..565bf62 100644 --- a/src/archive_clients/daemon.py +++ b/src/archive_clients/daemon.py @@ -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: diff --git a/src/archive_clients/jobs.py b/src/archive_clients/jobs.py index 3391a78..7a91021 100644 --- a/src/archive_clients/jobs.py +++ b/src/archive_clients/jobs.py @@ -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: diff --git a/src/archive_clients/protocol.py b/src/archive_clients/protocol.py index 0b2619b..d87ac5b 100644 --- a/src/archive_clients/protocol.py +++ b/src/archive_clients/protocol.py @@ -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: diff --git a/src/archive_clients/state.py b/src/archive_clients/state.py index 2504a19..242ffbd 100644 --- a/src/archive_clients/state.py +++ b/src/archive_clients/state.py @@ -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) diff --git a/src/archive_control/v1/__init__.py b/src/archive_control/v1/__init__.py index 6c01800..b1c4902 100644 --- a/src/archive_control/v1/__init__.py +++ b/src/archive_control/v1/__init__.py @@ -1,2 +1,3 @@ +# archive-control-proto commit: 03b6751d420e8e4e4cc9aeef4f76b13e28a7265f """Generated archive_control.v1 bindings.""" diff --git a/src/archive_control/v1/client_pb2.py b/src/archive_control/v1/client_pb2.py index f421b6e..c3f3058 100644 --- a/src/archive_control/v1/client_pb2.py +++ b/src/archive_control/v1/client_pb2.py @@ -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 diff --git a/src/archive_control/v1/client_pb2.pyi b/src/archive_control/v1/client_pb2.pyi index 2f1b138..fc270b6 100644 --- a/src/archive_control/v1/client_pb2.pyi +++ b/src/archive_control/v1/client_pb2.pyi @@ -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 diff --git a/src/archive_control/v1/common_pb2.py b/src/archive_control/v1/common_pb2.py index 563dd81..6dc23b3 100644 --- a/src/archive_control/v1/common_pb2.py +++ b/src/archive_control/v1/common_pb2.py @@ -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 diff --git a/src/archive_control/v1/common_pb2.pyi b/src/archive_control/v1/common_pb2.pyi index 9f2d749..2fccb9d 100644 --- a/src/archive_control/v1/common_pb2.pyi +++ b/src/archive_control/v1/common_pb2.pyi @@ -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 diff --git a/src/archive_control/v1/control_pb2.py b/src/archive_control/v1/control_pb2.py index 2eac574..62256b7 100644 --- a/src/archive_control/v1/control_pb2.py +++ b/src/archive_control/v1/control_pb2.py @@ -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 @@ -31,17 +31,17 @@ from archive_control.v1 import route_pb2 as archive__control_dot_v1_dot_route__p from google.protobuf import timestamp_pb2 as google_dot_protobuf_dot_timestamp__pb2 -DESCRIPTOR = _descriptor_pool.Default().AddSerializedFile(b'\n archive_control/v1/control.proto\x12\x12\x61rchive_control.v1\x1a\x1f\x61rchive_control/v1/common.proto\x1a\"archive_control/v1/inventory.proto\x1a\x1c\x61rchive_control/v1/job.proto\x1a!archive_control/v1/resource.proto\x1a\x1e\x61rchive_control/v1/route.proto\x1a\x1fgoogle/protobuf/timestamp.proto\"\xbc\x01\n\x10\x41ssignJobCommand\x12\x33\n\x03job\x18\x01 \x01(\x0b\x32!.archive_control.v1.JobDefinitionR\x03job\x12\x32\n\x15\x65xpected_job_revision\x18\x02 \x01(\x04R\x13\x65xpectedJobRevision\x12?\n\x1c\x65xpected_last_event_sequence\x18\x03 \x01(\x04R\x19\x65xpectedLastEventSequence\"\xef\x01\n\x12\x45xecuteStepCommand\x12\x15\n\x06job_id\x18\x01 \x01(\tR\x05jobId\x12\x32\n\x15\x65xpected_job_revision\x18\x02 \x01(\x04R\x13\x65xpectedJobRevision\x12\x33\n\x04step\x18\x03 \x01(\x0e\x32\x1f.archive_control.v1.JobStepKindR\x04step\x12\x18\n\x07\x61ttempt\x18\x04 \x01(\rR\x07\x61ttempt\x12?\n\x1c\x65xpected_last_event_sequence\x18\x05 \x01(\x04R\x19\x65xpectedLastEventSequence\"\xb6\x01\n\x10\x43\x61ncelJobCommand\x12\x15\n\x06job_id\x18\x01 \x01(\tR\x05jobId\x12\x32\n\x15\x65xpected_job_revision\x18\x02 \x01(\x04R\x13\x65xpectedJobRevision\x12\x16\n\x06reason\x18\x03 \x01(\tR\x06reason\x12?\n\x1c\x65xpected_last_event_sequence\x18\x04 \x01(\x04R\x19\x65xpectedLastEventSequence\"O\n\x12\x45nsureRouteCommand\x12\x39\n\x05route\x18\x01 \x01(\x0b\x32#.archive_control.v1.EnsureRouteSpecR\x05route\"4\n\x19RequestJobSnapshotCommand\x12\x17\n\x07job_ids\x18\x01 \x03(\tR\x06jobIds\"\xc8\x04\n\x07\x43ommand\x12\x1d\n\ncommand_id\x18\x01 \x01(\tR\tcommandId\x12\x39\n\ncreated_at\x18\x02 \x01(\x0b\x32\x1a.google.protobuf.TimestampR\tcreatedAt\x12\x45\n\nassign_job\x18\n \x01(\x0b\x32$.archive_control.v1.AssignJobCommandH\x00R\tassignJob\x12K\n\x0c\x65xecute_step\x18\x0b \x01(\x0b\x32&.archive_control.v1.ExecuteStepCommandH\x00R\x0b\x65xecuteStep\x12\x45\n\ncancel_job\x18\x0c \x01(\x0b\x32$.archive_control.v1.CancelJobCommandH\x00R\tcancelJob\x12K\n\x0c\x65nsure_route\x18\r \x01(\x0b\x32&.archive_control.v1.EnsureRouteCommandH\x00R\x0b\x65nsureRoute\x12M\n\x0finventory_query\x18\x0e \x01(\x0b\x32\".archive_control.v1.InventoryQueryH\x00R\x0einventoryQuery\x12\x61\n\x14request_job_snapshot\x18\x0f \x01(\x0b\x32-.archive_control.v1.RequestJobSnapshotCommandH\x00R\x12requestJobSnapshotB\t\n\x07payload\"\x9a\x01\n\nCommandAck\x12\x1d\n\ncommand_id\x18\x01 \x01(\tR\tcommandId\x12<\n\x06status\x18\x02 \x01(\x0e\x32$.archive_control.v1.CommandAckStatusR\x06status\x12/\n\x05\x65rror\x18\x03 \x01(\x0b\x32\x19.archive_control.v1.ErrorR\x05\x65rror\"\xd7\x04\n\x08JobEvent\x12\x19\n\x08\x65vent_id\x18\x01 \x01(\tR\x07\x65ventId\x12\x15\n\x06job_id\x18\x02 \x01(\tR\x05jobId\x12\x1a\n\x08sequence\x18\x03 \x01(\x04R\x08sequence\x12!\n\x0cjob_revision\x18\x04 \x01(\x04R\x0bjobRevision\x12\x34\n\x04type\x18\x05 \x01(\x0e\x32 .archive_control.v1.JobEventTypeR\x04type\x12\x32\n\x05state\x18\x06 \x01(\x0e\x32\x1c.archive_control.v1.JobStateR\x05state\x12\x1c\n\tcommitted\x18\x07 \x01(\x08R\tcommitted\x12;\n\x08progress\x18\x08 \x01(\x0b\x32\x1f.archive_control.v1.JobProgressR\x08progress\x12/\n\x05\x65rror\x18\t \x01(\x0b\x32\x19.archive_control.v1.ErrorR\x05\x65rror\x12Y\n\x11observed_resource\x18\n \x01(\x0b\x32,.archive_control.v1.ResourceStateFingerprintR\x10observedResource\x12L\n\x12observed_placement\x18\x0b \x01(\x0b\x32\x1d.archive_control.v1.PlacementR\x11observedPlacement\x12;\n\x0boccurred_at\x18\x0c \x01(\x0b\x32\x1a.google.protobuf.TimestampR\noccurredAt\"\xa1\x01\n\x0bJobSnapshot\x12/\n\x03job\x18\x01 \x01(\x0b\x32\x1d.archive_control.v1.JobRecordR\x03job\x12.\n\x13last_event_sequence\x18\x02 \x01(\x04R\x11lastEventSequence\x12\x31\n\x15in_flight_command_ids\x18\x03 \x03(\tR\x12inFlightCommandIds\"\xac\x02\n\x13\x43lientStateSnapshot\x12\x1f\n\x0bsnapshot_id\x18\x01 \x01(\tR\nsnapshotId\x12@\n\x0b\x61\x63tive_jobs\x18\x02 \x03(\x0b\x32\x1f.archive_control.v1.JobSnapshotR\nactiveJobs\x12\x36\n\x06routes\x18\x03 \x03(\x0b\x32\x1e.archive_control.v1.LocalRouteR\x06routes\x12=\n\x08services\x18\x04 \x03(\x0b\x32!.archive_control.v1.ServiceHealthR\x08services\x12;\n\x0bobserved_at\x18\x05 \x01(\x0b\x32\x1a.google.protobuf.TimestampR\nobservedAt\"\x97\x02\n\x0bRouteUpdate\x12\x1d\n\ncommand_id\x18\x01 \x01(\tR\tcommandId\x12\x34\n\x05route\x18\x02 \x01(\x0b\x32\x1e.archive_control.v1.LocalRouteR\x05route\x12I\n\x0cverification\x18\x03 \x01(\x0b\x32%.archive_control.v1.RouteVerificationR\x0cverification\x12/\n\x05\x65rror\x18\x04 \x01(\x0b\x32\x19.archive_control.v1.ErrorR\x05\x65rror\x12\x1b\n\tupdate_id\x18\x05 \x01(\tR\x08updateId\x12\x1a\n\x08sequence\x18\x06 \x01(\x04R\x08sequence\"r\n\rProtocolError\x12/\n\x05\x65rror\x18\x01 \x01(\x0b\x32\x19.archive_control.v1.ErrorR\x05\x65rror\x12\x30\n\x14offending_message_id\x18\x02 \x01(\tR\x12offendingMessageId*\x9a\x01\n\x10\x43ommandAckStatus\x12\"\n\x1e\x43OMMAND_ACK_STATUS_UNSPECIFIED\x10\x00\x12\x1f\n\x1b\x43OMMAND_ACK_STATUS_ACCEPTED\x10\x01\x12 \n\x1c\x43OMMAND_ACK_STATUS_DUPLICATE\x10\x02\x12\x1f\n\x1b\x43OMMAND_ACK_STATUS_REJECTED\x10\x03*\xe9\x03\n\x0cJobEventType\x12\x1e\n\x1aJOB_EVENT_TYPE_UNSPECIFIED\x10\x00\x12\x1b\n\x17JOB_EVENT_TYPE_ASSIGNED\x10\x01\x12\x1f\n\x1bJOB_EVENT_TYPE_STEP_STARTED\x10\x02\x12\x1b\n\x17JOB_EVENT_TYPE_PROGRESS\x10\x03\x12\x1a\n\x16JOB_EVENT_TYPE_WAITING\x10\x04\x12\x1a\n\x16JOB_EVENT_TYPE_STALLED\x10\x05\x12!\n\x1dJOB_EVENT_TYPE_STEP_SUCCEEDED\x10\x06\x12\x1c\n\x18JOB_EVENT_TYPE_COMMITTED\x10\x07\x12\x1d\n\x19JOB_EVENT_TYPE_CANCELLING\x10\x08\x12#\n\x1fJOB_EVENT_TYPE_ROLLBACK_STARTED\x10\t\x12%\n!JOB_EVENT_TYPE_ROLLBACK_SUCCEEDED\x10\n\x12#\n\x1fJOB_EVENT_TYPE_CLEANUP_REQUIRED\x10\x0b\x12\x1c\n\x18JOB_EVENT_TYPE_SUCCEEDED\x10\x0c\x12\x19\n\x15JOB_EVENT_TYPE_FAILED\x10\r\x12\x1c\n\x18JOB_EVENT_TYPE_CANCELLED\x10\x0e\x62\x06proto3') +DESCRIPTOR = _descriptor_pool.Default().AddSerializedFile(b'\n archive_control/v1/control.proto\x12\x12\x61rchive_control.v1\x1a\x1f\x61rchive_control/v1/common.proto\x1a\"archive_control/v1/inventory.proto\x1a\x1c\x61rchive_control/v1/job.proto\x1a!archive_control/v1/resource.proto\x1a\x1e\x61rchive_control/v1/route.proto\x1a\x1fgoogle/protobuf/timestamp.proto\"\xbc\x01\n\x10\x41ssignJobCommand\x12\x33\n\x03job\x18\x01 \x01(\x0b\x32!.archive_control.v1.JobDefinitionR\x03job\x12\x32\n\x15\x65xpected_job_revision\x18\x02 \x01(\x04R\x13\x65xpectedJobRevision\x12?\n\x1c\x65xpected_last_event_sequence\x18\x03 \x01(\x04R\x19\x65xpectedLastEventSequence\"\xef\x01\n\x12\x45xecuteStepCommand\x12\x15\n\x06job_id\x18\x01 \x01(\tR\x05jobId\x12\x32\n\x15\x65xpected_job_revision\x18\x02 \x01(\x04R\x13\x65xpectedJobRevision\x12\x33\n\x04step\x18\x03 \x01(\x0e\x32\x1f.archive_control.v1.JobStepKindR\x04step\x12\x18\n\x07\x61ttempt\x18\x04 \x01(\rR\x07\x61ttempt\x12?\n\x1c\x65xpected_last_event_sequence\x18\x05 \x01(\x04R\x19\x65xpectedLastEventSequence\"\xb6\x01\n\x10\x43\x61ncelJobCommand\x12\x15\n\x06job_id\x18\x01 \x01(\tR\x05jobId\x12\x32\n\x15\x65xpected_job_revision\x18\x02 \x01(\x04R\x13\x65xpectedJobRevision\x12\x16\n\x06reason\x18\x03 \x01(\tR\x06reason\x12?\n\x1c\x65xpected_last_event_sequence\x18\x04 \x01(\x04R\x19\x65xpectedLastEventSequence\"O\n\x12\x45nsureRouteCommand\x12\x39\n\x05route\x18\x01 \x01(\x0b\x32#.archive_control.v1.EnsureRouteSpecR\x05route\"4\n\x19RequestJobSnapshotCommand\x12\x17\n\x07job_ids\x18\x01 \x03(\tR\x06jobIds\"\xfa\x01\n\x13ReconcileJobCommand\x12J\n\x11\x61uthoritative_job\x18\x01 \x01(\x0b\x32\x1d.archive_control.v1.JobRecordR\x10\x61uthoritativeJob\x12I\n!authoritative_last_event_sequence\x18\x02 \x01(\x04R\x1e\x61uthoritativeLastEventSequence\x12\x34\n\x16superseded_command_ids\x18\x03 \x03(\tR\x14supersededCommandIds\x12\x16\n\x06reason\x18\x04 \x01(\tR\x06reason\"\x98\x05\n\x07\x43ommand\x12\x1d\n\ncommand_id\x18\x01 \x01(\tR\tcommandId\x12\x39\n\ncreated_at\x18\x02 \x01(\x0b\x32\x1a.google.protobuf.TimestampR\tcreatedAt\x12\x45\n\nassign_job\x18\n \x01(\x0b\x32$.archive_control.v1.AssignJobCommandH\x00R\tassignJob\x12K\n\x0c\x65xecute_step\x18\x0b \x01(\x0b\x32&.archive_control.v1.ExecuteStepCommandH\x00R\x0b\x65xecuteStep\x12\x45\n\ncancel_job\x18\x0c \x01(\x0b\x32$.archive_control.v1.CancelJobCommandH\x00R\tcancelJob\x12K\n\x0c\x65nsure_route\x18\r \x01(\x0b\x32&.archive_control.v1.EnsureRouteCommandH\x00R\x0b\x65nsureRoute\x12M\n\x0finventory_query\x18\x0e \x01(\x0b\x32\".archive_control.v1.InventoryQueryH\x00R\x0einventoryQuery\x12\x61\n\x14request_job_snapshot\x18\x0f \x01(\x0b\x32-.archive_control.v1.RequestJobSnapshotCommandH\x00R\x12requestJobSnapshot\x12N\n\rreconcile_job\x18\x10 \x01(\x0b\x32\'.archive_control.v1.ReconcileJobCommandH\x00R\x0creconcileJobB\t\n\x07payload\"\x9a\x01\n\nCommandAck\x12\x1d\n\ncommand_id\x18\x01 \x01(\tR\tcommandId\x12<\n\x06status\x18\x02 \x01(\x0e\x32$.archive_control.v1.CommandAckStatusR\x06status\x12/\n\x05\x65rror\x18\x03 \x01(\x0b\x32\x19.archive_control.v1.ErrorR\x05\x65rror\"\xf6\x04\n\x08JobEvent\x12\x19\n\x08\x65vent_id\x18\x01 \x01(\tR\x07\x65ventId\x12\x15\n\x06job_id\x18\x02 \x01(\tR\x05jobId\x12\x1a\n\x08sequence\x18\x03 \x01(\x04R\x08sequence\x12!\n\x0cjob_revision\x18\x04 \x01(\x04R\x0bjobRevision\x12\x34\n\x04type\x18\x05 \x01(\x0e\x32 .archive_control.v1.JobEventTypeR\x04type\x12\x32\n\x05state\x18\x06 \x01(\x0e\x32\x1c.archive_control.v1.JobStateR\x05state\x12\x1c\n\tcommitted\x18\x07 \x01(\x08R\tcommitted\x12;\n\x08progress\x18\x08 \x01(\x0b\x32\x1f.archive_control.v1.JobProgressR\x08progress\x12/\n\x05\x65rror\x18\t \x01(\x0b\x32\x19.archive_control.v1.ErrorR\x05\x65rror\x12Y\n\x11observed_resource\x18\n \x01(\x0b\x32,.archive_control.v1.ResourceStateFingerprintR\x10observedResource\x12L\n\x12observed_placement\x18\x0b \x01(\x0b\x32\x1d.archive_control.v1.PlacementR\x11observedPlacement\x12;\n\x0boccurred_at\x18\x0c \x01(\x0b\x32\x1a.google.protobuf.TimestampR\noccurredAt\x12\x1d\n\ncommand_id\x18\r \x01(\tR\tcommandId\"\xa1\x01\n\x0bJobSnapshot\x12/\n\x03job\x18\x01 \x01(\x0b\x32\x1d.archive_control.v1.JobRecordR\x03job\x12.\n\x13last_event_sequence\x18\x02 \x01(\x04R\x11lastEventSequence\x12\x31\n\x15in_flight_command_ids\x18\x03 \x03(\tR\x12inFlightCommandIds\"\xac\x02\n\x13\x43lientStateSnapshot\x12\x1f\n\x0bsnapshot_id\x18\x01 \x01(\tR\nsnapshotId\x12@\n\x0b\x61\x63tive_jobs\x18\x02 \x03(\x0b\x32\x1f.archive_control.v1.JobSnapshotR\nactiveJobs\x12\x36\n\x06routes\x18\x03 \x03(\x0b\x32\x1e.archive_control.v1.LocalRouteR\x06routes\x12=\n\x08services\x18\x04 \x03(\x0b\x32!.archive_control.v1.ServiceHealthR\x08services\x12;\n\x0bobserved_at\x18\x05 \x01(\x0b\x32\x1a.google.protobuf.TimestampR\nobservedAt\"\x97\x02\n\x0bRouteUpdate\x12\x1d\n\ncommand_id\x18\x01 \x01(\tR\tcommandId\x12\x34\n\x05route\x18\x02 \x01(\x0b\x32\x1e.archive_control.v1.LocalRouteR\x05route\x12I\n\x0cverification\x18\x03 \x01(\x0b\x32%.archive_control.v1.RouteVerificationR\x0cverification\x12/\n\x05\x65rror\x18\x04 \x01(\x0b\x32\x19.archive_control.v1.ErrorR\x05\x65rror\x12\x1b\n\tupdate_id\x18\x05 \x01(\tR\x08updateId\x12\x1a\n\x08sequence\x18\x06 \x01(\x04R\x08sequence\"r\n\rProtocolError\x12/\n\x05\x65rror\x18\x01 \x01(\x0b\x32\x19.archive_control.v1.ErrorR\x05\x65rror\x12\x30\n\x14offending_message_id\x18\x02 \x01(\tR\x12offendingMessageId*\x9a\x01\n\x10\x43ommandAckStatus\x12\"\n\x1e\x43OMMAND_ACK_STATUS_UNSPECIFIED\x10\x00\x12\x1f\n\x1b\x43OMMAND_ACK_STATUS_ACCEPTED\x10\x01\x12 \n\x1c\x43OMMAND_ACK_STATUS_DUPLICATE\x10\x02\x12\x1f\n\x1b\x43OMMAND_ACK_STATUS_REJECTED\x10\x03*\xe9\x03\n\x0cJobEventType\x12\x1e\n\x1aJOB_EVENT_TYPE_UNSPECIFIED\x10\x00\x12\x1b\n\x17JOB_EVENT_TYPE_ASSIGNED\x10\x01\x12\x1f\n\x1bJOB_EVENT_TYPE_STEP_STARTED\x10\x02\x12\x1b\n\x17JOB_EVENT_TYPE_PROGRESS\x10\x03\x12\x1a\n\x16JOB_EVENT_TYPE_WAITING\x10\x04\x12\x1a\n\x16JOB_EVENT_TYPE_STALLED\x10\x05\x12!\n\x1dJOB_EVENT_TYPE_STEP_SUCCEEDED\x10\x06\x12\x1c\n\x18JOB_EVENT_TYPE_COMMITTED\x10\x07\x12\x1d\n\x19JOB_EVENT_TYPE_CANCELLING\x10\x08\x12#\n\x1fJOB_EVENT_TYPE_ROLLBACK_STARTED\x10\t\x12%\n!JOB_EVENT_TYPE_ROLLBACK_SUCCEEDED\x10\n\x12#\n\x1fJOB_EVENT_TYPE_CLEANUP_REQUIRED\x10\x0b\x12\x1c\n\x18JOB_EVENT_TYPE_SUCCEEDED\x10\x0c\x12\x19\n\x15JOB_EVENT_TYPE_FAILED\x10\r\x12\x1c\n\x18JOB_EVENT_TYPE_CANCELLED\x10\x0e\x62\x06proto3') _globals = globals() _builder.BuildMessageAndEnumDescriptors(DESCRIPTOR, _globals) _builder.BuildTopDescriptorsAndMessages(DESCRIPTOR, 'archive_control.v1.control_pb2', _globals) if not _descriptor._USE_C_DESCRIPTORS: DESCRIPTOR._loaded_options = None - _globals['_COMMANDACKSTATUS']._serialized_start=3220 - _globals['_COMMANDACKSTATUS']._serialized_end=3374 - _globals['_JOBEVENTTYPE']._serialized_start=3377 - _globals['_JOBEVENTTYPE']._serialized_end=3866 + _globals['_COMMANDACKSTATUS']._serialized_start=3584 + _globals['_COMMANDACKSTATUS']._serialized_end=3738 + _globals['_JOBEVENTTYPE']._serialized_start=3741 + _globals['_JOBEVENTTYPE']._serialized_end=4230 _globals['_ASSIGNJOBCOMMAND']._serialized_start=256 _globals['_ASSIGNJOBCOMMAND']._serialized_end=444 _globals['_EXECUTESTEPCOMMAND']._serialized_start=447 @@ -52,18 +52,20 @@ if not _descriptor._USE_C_DESCRIPTORS: _globals['_ENSUREROUTECOMMAND']._serialized_end=952 _globals['_REQUESTJOBSNAPSHOTCOMMAND']._serialized_start=954 _globals['_REQUESTJOBSNAPSHOTCOMMAND']._serialized_end=1006 - _globals['_COMMAND']._serialized_start=1009 - _globals['_COMMAND']._serialized_end=1593 - _globals['_COMMANDACK']._serialized_start=1596 - _globals['_COMMANDACK']._serialized_end=1750 - _globals['_JOBEVENT']._serialized_start=1753 - _globals['_JOBEVENT']._serialized_end=2352 - _globals['_JOBSNAPSHOT']._serialized_start=2355 - _globals['_JOBSNAPSHOT']._serialized_end=2516 - _globals['_CLIENTSTATESNAPSHOT']._serialized_start=2519 - _globals['_CLIENTSTATESNAPSHOT']._serialized_end=2819 - _globals['_ROUTEUPDATE']._serialized_start=2822 - _globals['_ROUTEUPDATE']._serialized_end=3101 - _globals['_PROTOCOLERROR']._serialized_start=3103 - _globals['_PROTOCOLERROR']._serialized_end=3217 + _globals['_RECONCILEJOBCOMMAND']._serialized_start=1009 + _globals['_RECONCILEJOBCOMMAND']._serialized_end=1259 + _globals['_COMMAND']._serialized_start=1262 + _globals['_COMMAND']._serialized_end=1926 + _globals['_COMMANDACK']._serialized_start=1929 + _globals['_COMMANDACK']._serialized_end=2083 + _globals['_JOBEVENT']._serialized_start=2086 + _globals['_JOBEVENT']._serialized_end=2716 + _globals['_JOBSNAPSHOT']._serialized_start=2719 + _globals['_JOBSNAPSHOT']._serialized_end=2880 + _globals['_CLIENTSTATESNAPSHOT']._serialized_start=2883 + _globals['_CLIENTSTATESNAPSHOT']._serialized_end=3183 + _globals['_ROUTEUPDATE']._serialized_start=3186 + _globals['_ROUTEUPDATE']._serialized_end=3465 + _globals['_PROTOCOLERROR']._serialized_start=3467 + _globals['_PROTOCOLERROR']._serialized_end=3581 # @@protoc_insertion_point(module_scope) diff --git a/src/archive_control/v1/control_pb2.pyi b/src/archive_control/v1/control_pb2.pyi index 21e4da8..73cc187 100644 --- a/src/archive_control/v1/control_pb2.pyi +++ b/src/archive_control/v1/control_pb2.pyi @@ -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") diff --git a/src/archive_control/v1/envelope_pb2.py b/src/archive_control/v1/envelope_pb2.py index 54efcca..9d57951 100644 --- a/src/archive_control/v1/envelope_pb2.py +++ b/src/archive_control/v1/envelope_pb2.py @@ -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 diff --git a/src/archive_control/v1/envelope_pb2.pyi b/src/archive_control/v1/envelope_pb2.pyi index 2c400c1..1df4126 100644 --- a/src/archive_control/v1/envelope_pb2.pyi +++ b/src/archive_control/v1/envelope_pb2.pyi @@ -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 diff --git a/src/archive_control/v1/inventory_pb2.py b/src/archive_control/v1/inventory_pb2.py index 165af54..6801b65 100644 --- a/src/archive_control/v1/inventory_pb2.py +++ b/src/archive_control/v1/inventory_pb2.py @@ -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 diff --git a/src/archive_control/v1/inventory_pb2.pyi b/src/archive_control/v1/inventory_pb2.pyi index b25fba9..b505f29 100644 --- a/src/archive_control/v1/inventory_pb2.pyi +++ b/src/archive_control/v1/inventory_pb2.pyi @@ -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 diff --git a/src/archive_control/v1/job_pb2.py b/src/archive_control/v1/job_pb2.py index d3b595f..4c8814e 100644 --- a/src/archive_control/v1/job_pb2.py +++ b/src/archive_control/v1/job_pb2.py @@ -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 diff --git a/src/archive_control/v1/job_pb2.pyi b/src/archive_control/v1/job_pb2.pyi index d5b13a0..a1ece65 100644 --- a/src/archive_control/v1/job_pb2.pyi +++ b/src/archive_control/v1/job_pb2.pyi @@ -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 diff --git a/src/archive_control/v1/resource_pb2.py b/src/archive_control/v1/resource_pb2.py index b42874d..3b886d0 100644 --- a/src/archive_control/v1/resource_pb2.py +++ b/src/archive_control/v1/resource_pb2.py @@ -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 diff --git a/src/archive_control/v1/resource_pb2.pyi b/src/archive_control/v1/resource_pb2.pyi index 2474c21..59795bc 100644 --- a/src/archive_control/v1/resource_pb2.pyi +++ b/src/archive_control/v1/resource_pb2.pyi @@ -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 diff --git a/src/archive_control/v1/route_pb2.py b/src/archive_control/v1/route_pb2.py index df7957b..0886b4f 100644 --- a/src/archive_control/v1/route_pb2.py +++ b/src/archive_control/v1/route_pb2.py @@ -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 diff --git a/src/archive_control/v1/route_pb2.pyi b/src/archive_control/v1/route_pb2.pyi index abe25c7..06e59cd 100644 --- a/src/archive_control/v1/route_pb2.pyi +++ b/src/archive_control/v1/route_pb2.pyi @@ -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 diff --git a/src/archive_control/v1/transfer_pb2.py b/src/archive_control/v1/transfer_pb2.py index e31a5ad..9f18c7e 100644 --- a/src/archive_control/v1/transfer_pb2.py +++ b/src/archive_control/v1/transfer_pb2.py @@ -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 diff --git a/src/archive_control/v1/transfer_pb2.pyi b/src/archive_control/v1/transfer_pb2.pyi index 0f65722..6ee0f47 100644 --- a/src/archive_control/v1/transfer_pb2.pyi +++ b/src/archive_control/v1/transfer_pb2.pyi @@ -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 diff --git a/tests/test_daemon.py b/tests/test_daemon.py index 1157e69..ad30a79 100644 --- a/tests/test_daemon.py +++ b/tests/test_daemon.py @@ -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 diff --git a/tests/test_jobs.py b/tests/test_jobs.py index 317cb3e..dabbed4 100644 --- a/tests/test_jobs.py +++ b/tests/test_jobs.py @@ -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), diff --git a/tests/test_state.py b/tests/test_state.py index da40aa0..e34b185 100644 --- a/tests/test_state.py +++ b/tests/test_state.py @@ -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()