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

1"""central-queue: the idempotent DynamoDB-backed queue lifecycle.""" 

2 

3from __future__ import annotations 

4 

5import copy 

6from typing import Any 

7 

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 

31 

32 

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 } 

55 

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) 

64 

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 ) 

139 

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 ) 

150 

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) 

185 

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 }