Coverage for scripts / live_release_validation / actions / central_queue.py: 100.00%
82 statements
« prev ^ index » next coverage.py v7.13.5, created at 2026-09-14 22:07 +0000
« prev ^ index » next coverage.py v7.13.5, created at 2026-09-14 22:07 +0000
1"""central-queue: the idempotent DynamoDB-backed queue lifecycle."""
3from __future__ import annotations
5import copy
6from typing import Any
8from ..checks.central_queue import (
9 _central_manifest,
10 _central_queue_job_id,
11 _get_central_queue_job,
12 _read_central_job_item,
13 _reconcile_central_workload_identity,
14 _register_central_job,
15 _validate_central_job_identity,
16 _wait_for_central_queue_appearance,
17 _wait_for_central_queue_terminal,
18)
19from ..checks.jobs import (
20 _complete_job_lifecycle,
21 _effective_job_identity,
22 _job_appearance_timeout,
23 _register_job,
24 _response_json,
25 _run_token,
26)
27from ..constants import (
28 _RUN_JOB_LABEL,
29)
30from ..models import RunContext
33def action_central_queue_lifecycle(ctx: RunContext) -> dict[str, Any]:
34 """Exercise the idempotent DynamoDB queue and require terminal persistence."""
35 manifest, name, namespace, marker = _central_manifest(ctx)
36 target_region = ctx.deployment_regions[0]
37 record = _register_job(
38 ctx,
39 name=name,
40 namespace=namespace,
41 execution_region=target_region,
42 path="dynamodb",
43 reactivate_deleted=False,
44 )
45 transport_region = record.get("transport_region")
46 idempotency_key = f"gco-live-validation:{ctx.settings.run_id}:central"
47 job_id = _central_queue_job_id(idempotency_key)
48 body = {
49 "manifest": manifest,
50 "target_region": target_region,
51 "namespace": namespace,
52 "priority": 100,
53 "labels": {_RUN_JOB_LABEL: _run_token(ctx.settings.run_id)},
54 }
56 envelope = {
57 "transport": "central-queue",
58 "body": body,
59 "idempotency_key": idempotency_key,
60 "job_id": job_id,
61 "transport_region": transport_region,
62 }
63 ctx.prepare_job_submission(record, envelope=envelope, resumable=True)
65 central_record = _register_central_job(
66 ctx,
67 job_id=job_id,
68 idempotency_key=idempotency_key,
69 record=record,
70 marker=marker,
71 body=body,
72 )
73 initial_state = str(record.get("submission_state") or "")
74 queue_job = _get_central_queue_job(ctx, central_record)
75 submission: dict[str, Any]
76 if queue_job is None and initial_state in {"prepared", "submitting", "submitted"}:
77 if initial_state in {"prepared", "submitting"}:
78 ctx.begin_job_submission(
79 record,
80 reconciliation_timeout_seconds=_job_appearance_timeout(ctx),
81 )
82 persisted_envelope = record.get("submission_envelope")
83 if not isinstance(persisted_envelope, dict) or persisted_envelope != envelope:
84 raise RuntimeError("Central queue replay envelope changed")
85 response = ctx.aws_client.make_authenticated_request(
86 method="POST",
87 path="/api/v1/queue/jobs",
88 body=copy.deepcopy(persisted_envelope["body"]),
89 headers={"Idempotency-Key": str(persisted_envelope["idempotency_key"])},
90 target_region=persisted_envelope.get("transport_region"),
91 )
92 if response.status_code == 409:
93 raise RuntimeError(
94 "Central queue rejected the exact idempotent replay because request drift was detected"
95 )
96 if response.status_code not in {200, 201}:
97 raise RuntimeError(
98 f"Central queue submission failed: {response.status_code} {response.text}"
99 )
100 submission = _response_json(response, "Central queue submission")
101 queued_job = submission.get("job")
102 if not isinstance(queued_job, dict):
103 raise RuntimeError("Central queue response omitted job")
104 _validate_central_job_identity(central_record, queued_job)
105 ctx.finish_job_submission(
106 record,
107 submission,
108 appearance_timeout_seconds=_job_appearance_timeout(ctx),
109 )
110 central_record["submission"] = submission
111 central_record["submission_state"] = "submitted"
112 central_record["appearance_deadline"] = record["appearance_deadline"]
113 ctx.persist()
114 elif queue_job is not None:
115 submission = (
116 record.get("submission")
117 or central_record.get("submission")
118 or {"reconciled_existing_job": True, "job": queue_job}
119 )
120 if not isinstance(submission, dict):
121 raise RuntimeError("Checkpointed central queue submission is malformed")
122 submitted_job = submission.get("job")
123 if isinstance(submitted_job, dict):
124 _validate_central_job_identity(central_record, submitted_job)
125 if initial_state in {"submitting", "submitted", "appeared"}:
126 ctx.finish_job_submission(
127 record,
128 submission,
129 appearance_timeout_seconds=_job_appearance_timeout(ctx),
130 )
131 central_record["submission"] = submission
132 central_record["submission_state"] = "reconciled"
133 central_record["appearance_deadline"] = record.get("appearance_deadline")
134 ctx.persist()
135 else:
136 raise RuntimeError(
137 f"Central queue record is absent and state {initial_state!r} is not replayable"
138 )
140 _wait_for_central_queue_appearance(ctx, central_record)
141 final_job, history = _wait_for_central_queue_terminal(ctx, central_record)
142 central_record["status"] = str(final_job.get("status") or "unknown")
143 central_record["status_history"] = history
144 ctx.persist()
145 if final_job.get("status") != "succeeded":
146 raise RuntimeError(
147 f"Central queue job {job_id} finished as {final_job.get('status')}: "
148 f"{final_job.get('error_message') or 'no error message'}"
149 )
151 item = _read_central_job_item(ctx, job_id)
152 if item.get("status") != "succeeded":
153 raise RuntimeError(f"DynamoDB record {job_id} is {item.get('status')}, expected succeeded")
154 record = _reconcile_central_workload_identity(
155 ctx,
156 central_record,
157 item,
158 workload_record=record,
159 )
160 if record.get("deleted"):
161 evidence = record.get("validation_evidence")
162 if not isinstance(evidence, dict) or evidence.get("marker") != marker:
163 raise RuntimeError(
164 "Central Job was deleted without checkpointed live-validation evidence; "
165 "refusing an idempotency replay"
166 )
167 actual_name, actual_namespace = _effective_job_identity(record)
168 expected_evidence = {
169 "name": actual_name,
170 "namespace": actual_namespace,
171 "uid": record.get("k8s_job_uid"),
172 "central_queue_job_id": job_id,
173 }
174 for key, expected in expected_evidence.items():
175 if evidence.get(key) != expected:
176 raise RuntimeError(
177 f"Central Job deletion evidence does not match actual identity: {key}"
178 )
179 workload_lifecycle = {
180 **copy.deepcopy(evidence),
181 "deletion": {"reconciled_checkpointed_deletion": True},
182 }
183 else:
184 workload_lifecycle = _complete_job_lifecycle(ctx, record=record, marker=marker)
186 central_record["status"] = "succeeded"
187 central_record["cleanup_complete"] = True
188 central_record["workload_lifecycle"] = workload_lifecycle
189 ctx.persist()
190 return {
191 "submission": submission,
192 "job_id": job_id,
193 "job_name": name,
194 "namespace": namespace,
195 "k8s_job_name": record.get("k8s_job_name"),
196 "k8s_job_namespace": record.get("k8s_job_namespace"),
197 "k8s_job_uid": record.get("k8s_job_uid"),
198 "target_region": target_region,
199 "status_history": history,
200 "final_job": final_job,
201 "dynamodb_item": item,
202 "workload_lifecycle": workload_lifecycle,
203 }