Coverage for gco / services / api_routes / jobs.py: 100.00%

349 statements  

« prev     ^ index     » next       coverage.py v7.13.5, created at 2026-09-14 22:07 +0000

1"""Job listing, details, logs, events, metrics, delete, and retry endpoints.""" 

2 

3from __future__ import annotations 

4 

5import logging 

6import re 

7from datetime import UTC, datetime, timedelta 

8from typing import Any, cast 

9 

10from fastapi import APIRouter, HTTPException, Query 

11from fastapi.responses import JSONResponse, Response 

12from kubernetes import client as kubernetes_client 

13 

14from gco.models import ManifestSubmissionRequest 

15from gco.services.api_shared import ( 

16 BulkDeleteRequest, 

17 _check_namespace, 

18 _check_processor, 

19 _collect_pod_scheduling, 

20 _empty_scheduling_info, 

21 _parse_event_to_dict, 

22 _parse_job_to_dict, 

23 _parse_pod_to_dict, 

24 internal_server_error, 

25) 

26from gco.services.structured_logging import sanitize_log_value 

27 

28router = APIRouter(prefix="/api/v1/jobs", tags=["Jobs"]) 

29logger = logging.getLogger(__name__) 

30 

31_LABEL_NAME_RE = re.compile(r"^[A-Za-z0-9](?:[-_.A-Za-z0-9]{0,61}[A-Za-z0-9])?$") 

32_DNS_LABEL_RE = re.compile(r"^[a-z0-9](?:[-a-z0-9]{0,61}[a-z0-9])?$") 

33 

34 

35def _decode_pod_log_response(response: Any) -> str: 

36 """Decode the raw Kubernetes log body without converting bytes to their repr.""" 

37 release_conn = getattr(response, "release_conn", None) 

38 try: 

39 payload = response.data if hasattr(response, "data") else response 

40 if isinstance(payload, str): 

41 return payload 

42 if isinstance(payload, (bytes, bytearray)): 

43 return bytes(payload).decode("utf-8", errors="replace") 

44 raise TypeError(f"Kubernetes returned an unsupported pod log payload: {type(payload)!r}") 

45 finally: 

46 if callable(release_conn): 

47 release_conn() 

48 

49 

50def _parse_exact_label_selector(selector: str | None) -> list[tuple[str, str]]: 

51 """Parse the API's deliberately narrow, fail-closed selector subset.""" 

52 if selector is None: 

53 return [] 

54 

55 requirements: list[tuple[str, str]] = [] 

56 for raw_clause in selector.split(","): 

57 clause = raw_clause.strip() 

58 if not clause or clause.count("=") != 1: 

59 raise ValueError( 

60 "Label selectors must be comma-separated exact matches in key=value form" 

61 ) 

62 

63 key, value = (part.strip() for part in clause.split("=", 1)) 

64 if "/" in key: 

65 if key.count("/") != 1: 

66 raise ValueError(f"Invalid label key in selector: {key!r}") 

67 prefix, name = key.split("/", 1) 

68 valid_prefix = len(prefix) <= 253 and all( 

69 _DNS_LABEL_RE.fullmatch(part) for part in prefix.split(".") 

70 ) 

71 else: 

72 name = key 

73 valid_prefix = True 

74 

75 if not valid_prefix or not _LABEL_NAME_RE.fullmatch(name): 

76 raise ValueError(f"Invalid label key in selector: {key!r}") 

77 if value and not _LABEL_NAME_RE.fullmatch(value): 

78 raise ValueError(f"Invalid label value in selector for {key!r}") 

79 requirements.append((key, value)) 

80 

81 return requirements 

82 

83 

84def _labels_match(labels: Any, requirements: list[tuple[str, str]]) -> bool: 

85 """Return whether all parsed exact-match requirements are satisfied.""" 

86 if not isinstance(labels, dict): 

87 return False 

88 return all(labels.get(key) == value for key, value in requirements) 

89 

90 

91@router.get("") 

92async def list_jobs( 

93 namespace: str | None = Query(None, description="Filter by namespace"), 

94 status: str | None = Query(None, description="Filter by status"), 

95 limit: int = Query(50, ge=1, le=1000, description="Maximum number of jobs to return"), 

96 offset: int = Query(0, ge=0, description="Number of jobs to skip"), 

97 sort: str = Query("createdAt:desc", description="Sort field and order (field:asc|desc)"), 

98 label_selector: str | None = Query( 

99 None, 

100 max_length=1024, 

101 description="Comma-separated exact-match label filters (key=value only)", 

102 ), 

103) -> Response: 

104 """List Kubernetes Jobs with pagination and filtering.""" 

105 processor = _check_processor() 

106 

107 try: 

108 selector_requirements = _parse_exact_label_selector(label_selector) 

109 all_jobs = await processor.list_jobs(namespace=namespace, status_filter=status) 

110 

111 if selector_requirements: 

112 all_jobs = [ 

113 job 

114 for job in all_jobs 

115 if _labels_match(job.get("metadata", {}).get("labels", {}), selector_requirements) 

116 ] 

117 

118 sort_field, sort_order = "createdAt", "desc" 

119 if ":" in sort: 

120 sort_field, sort_order = sort.split(":", 1) 

121 

122 def get_sort_key(job: dict[str, Any]) -> Any: 

123 if sort_field == "createdAt": 

124 return job.get("metadata", {}).get("creationTimestamp", "") 

125 if sort_field == "name": 

126 return job.get("metadata", {}).get("name", "") 

127 if sort_field == "status": 

128 return job.get("status", {}).get("active", 0) 

129 return "" 

130 

131 all_jobs.sort(key=get_sort_key, reverse=(sort_order == "desc")) 

132 

133 total = len(all_jobs) 

134 paginated_jobs = all_jobs[offset : offset + limit] 

135 

136 response = { 

137 "cluster_id": processor.cluster_id, 

138 "region": processor.region, 

139 "timestamp": datetime.now(UTC).isoformat(), 

140 "total": total, 

141 "limit": limit, 

142 "offset": offset, 

143 "has_more": (offset + limit) < total, 

144 "count": len(paginated_jobs), 

145 "jobs": paginated_jobs, 

146 } 

147 

148 return JSONResponse(status_code=200, content=response) 

149 

150 except ValueError as e: 

151 raise HTTPException(status_code=400, detail=str(e)) from e 

152 except Exception as e: 

153 raise internal_server_error("listing jobs", e) from e 

154 

155 

156def _job_scheduling(processor: Any, namespace: str, name: str) -> dict[str, Any]: 

157 """Report the nodes a Job's pods landed on, and each node's hardware. 

158 

159 Placement is supplementary to the Job read: a failure here (pods already 

160 garbage-collected, Node read refused) must never turn a successful job 

161 lookup into an error, so the failure is recorded in the payload instead. 

162 """ 

163 try: 

164 pods = processor.core_v1.list_namespaced_pod( 

165 namespace=namespace, label_selector=f"job-name={name}" 

166 ) 

167 return _collect_pod_scheduling(processor.core_v1, pods.items) 

168 except Exception as exc: 

169 # namespace/name arrive from the request path and the Kubernetes error 

170 # echoes them back, so all three are sanitized before logging (CWE-117). 

171 logger.warning( 

172 "Could not resolve node placement for %s/%s: %s", 

173 sanitize_log_value(namespace), 

174 sanitize_log_value(name), 

175 sanitize_log_value(exc), 

176 ) 

177 info = _empty_scheduling_info() 

178 info["node_lookup_error"] = str(exc) 

179 return info 

180 

181 

182@router.get("/{namespace}/{name}") 

183async def get_job(namespace: str, name: str) -> Response: 

184 """Get details of a specific Job, including where its pods were scheduled.""" 

185 processor = _check_processor() 

186 _check_namespace(namespace, processor) 

187 

188 try: 

189 job = processor.batch_v1.read_namespaced_job(name=name, namespace=namespace) 

190 job_info = _parse_job_to_dict(job) 

191 

192 response = { 

193 "cluster_id": processor.cluster_id, 

194 "region": processor.region, 

195 "timestamp": datetime.now(UTC).isoformat(), 

196 **job_info, 

197 "scheduling": _job_scheduling(processor, namespace, name), 

198 } 

199 

200 return JSONResponse(status_code=200, content=response) 

201 

202 except Exception as e: 

203 if "NotFound" in str(e) or "404" in str(e): 

204 raise HTTPException( 

205 status_code=404, detail=f"Job '{name}' not found in namespace '{namespace}'" 

206 ) from e 

207 raise internal_server_error("getting job", e) from e 

208 

209 

210@router.get("/{namespace}/{name}/logs") 

211async def get_job_logs( 

212 namespace: str, 

213 name: str, 

214 container: str | None = Query(None, description="Container name (for multi-container pods)"), 

215 tail: int = Query(100, ge=1, le=10000, description="Number of lines from the end"), 

216 previous: bool = Query(False, description="Get logs from previous terminated container"), 

217 since_seconds: int | None = Query( 

218 None, ge=1, description="Only return logs newer than N seconds" 

219 ), 

220 timestamps: bool = Query(False, description="Include timestamps in log lines"), 

221) -> Response: 

222 """Get logs from a Job's pods.""" 

223 from kubernetes.client.rest import ApiException as K8sApiException 

224 

225 processor = _check_processor() 

226 _check_namespace(namespace, processor) 

227 

228 try: 

229 try: 

230 processor.batch_v1.read_namespaced_job(name=name, namespace=namespace) 

231 except K8sApiException as e: 

232 if e.status == 404: 

233 raise HTTPException( 

234 status_code=404, 

235 detail=f"Job '{name}' not found in namespace '{namespace}'", 

236 ) from e 

237 raise 

238 

239 pods = processor.core_v1.list_namespaced_pod( 

240 namespace=namespace, label_selector=f"job-name={name}" 

241 ) 

242 

243 if not pods.items: 

244 raise HTTPException( 

245 status_code=404, 

246 detail=( 

247 f"No pods found for job '{name}'. " 

248 "The job may have completed and pods were cleaned up " 

249 "(ttlSecondsAfterFinished). Use 'gco jobs get' to check job status." 

250 ), 

251 ) 

252 

253 sorted_pods = sorted( 

254 pods.items, 

255 key=lambda p: p.metadata.creation_timestamp or datetime.min.replace(tzinfo=UTC), 

256 reverse=True, 

257 ) 

258 pod = sorted_pods[0] 

259 

260 pod_phase = pod.status.phase if pod.status else "Unknown" 

261 if pod_phase == "Pending": 

262 raise HTTPException( 

263 status_code=409, 

264 detail=( 

265 f"Pod '{pod.metadata.name}' is still Pending — logs are not yet available. " 

266 "The node may still be provisioning. Use 'gco jobs events' to check." 

267 ), 

268 ) 

269 

270 log_kwargs: dict[str, Any] = { 

271 "name": pod.metadata.name, 

272 "namespace": namespace, 

273 "tail_lines": tail, 

274 "previous": previous, 

275 "timestamps": timestamps, 

276 } 

277 if container: 

278 log_kwargs["container"] = container 

279 if since_seconds: 

280 log_kwargs["since_seconds"] = since_seconds 

281 

282 try: 

283 log_response = processor.core_v1.read_namespaced_pod_log( 

284 **log_kwargs, _preload_content=False 

285 ) 

286 logs = _decode_pod_log_response(log_response) 

287 except K8sApiException as e: 

288 if e.status == 400: 

289 error_body = str(e.body) if e.body else str(e.reason) 

290 if "waiting" in error_body.lower() or "not found" in error_body.lower(): 

291 available = [c.name for c in pod.spec.containers] 

292 raise HTTPException( 

293 status_code=400, 

294 detail=( 

295 f"Logs not available: {error_body}. " 

296 f"Pod phase: {pod_phase}. " 

297 f"Available containers: {available}" 

298 ), 

299 ) from e 

300 raise HTTPException(status_code=400, detail=f"Bad request: {error_body}") from e 

301 raise 

302 

303 available_containers = [c.name for c in pod.spec.containers] 

304 init_containers = [c.name for c in (pod.spec.init_containers or [])] 

305 

306 response = { 

307 "cluster_id": processor.cluster_id, 

308 "region": processor.region, 

309 "timestamp": datetime.now(UTC).isoformat(), 

310 "job_name": name, 

311 "namespace": namespace, 

312 "pod_name": pod.metadata.name, 

313 "container": container or (available_containers[0] if available_containers else None), 

314 "available_containers": available_containers, 

315 "init_containers": init_containers, 

316 "previous": previous, 

317 "tail_lines": tail, 

318 "logs": logs, 

319 } 

320 

321 return JSONResponse(status_code=200, content=response) 

322 

323 except HTTPException: 

324 raise 

325 except K8sApiException as e: 

326 logger.error(f"Kubernetes API error getting job logs: {e.status} {e.reason}") 

327 raise HTTPException( 

328 status_code=502, detail=f"Kubernetes API error: {e.status} {e.reason}" 

329 ) from e 

330 except Exception as e: 

331 raise internal_server_error("getting job logs", e) from e 

332 

333 

334@router.get("/{namespace}/{name}/events") 

335async def get_job_events(namespace: str, name: str) -> Response: 

336 """Get events related to a Job.""" 

337 processor = _check_processor() 

338 _check_namespace(namespace, processor) 

339 

340 try: 

341 field_selector = f"involvedObject.name={name},involvedObject.kind=Job" 

342 job_events = processor.core_v1.list_namespaced_event( 

343 namespace=namespace, field_selector=field_selector 

344 ) 

345 

346 pods = processor.core_v1.list_namespaced_pod( 

347 namespace=namespace, label_selector=f"job-name={name}" 

348 ) 

349 

350 pod_events = [] 

351 for pod in pods.items: 

352 field_selector = f"involvedObject.name={pod.metadata.name},involvedObject.kind=Pod" 

353 events = processor.core_v1.list_namespaced_event( 

354 namespace=namespace, field_selector=field_selector 

355 ) 

356 pod_events.extend(events.items) 

357 

358 all_events = [_parse_event_to_dict(e) for e in job_events.items] 

359 all_events.extend([_parse_event_to_dict(e) for e in pod_events]) 

360 all_events.sort( 

361 key=lambda e: e.get("lastTimestamp") or e.get("firstTimestamp") or "", reverse=True 

362 ) 

363 

364 response = { 

365 "cluster_id": processor.cluster_id, 

366 "region": processor.region, 

367 "timestamp": datetime.now(UTC).isoformat(), 

368 "job_name": name, 

369 "namespace": namespace, 

370 "count": len(all_events), 

371 "events": all_events, 

372 } 

373 

374 return JSONResponse(status_code=200, content=response) 

375 

376 except Exception as e: 

377 raise internal_server_error("getting job events", e) from e 

378 

379 

380@router.get("/{namespace}/{name}/pods") 

381async def get_job_pods(namespace: str, name: str) -> Response: 

382 """Get pods belonging to a Job, with the hardware each one landed on.""" 

383 processor = _check_processor() 

384 _check_namespace(namespace, processor) 

385 

386 try: 

387 pods = processor.core_v1.list_namespaced_pod( 

388 namespace=namespace, label_selector=f"job-name={name}" 

389 ) 

390 pod_list = [_parse_pod_to_dict(pod) for pod in pods.items] 

391 

392 try: 

393 scheduling = _collect_pod_scheduling(processor.core_v1, pods.items) 

394 except Exception as exc: # pragma: no cover - collector swallows its own 

395 logger.warning( 

396 "Could not resolve node placement for %s/%s: %s", 

397 sanitize_log_value(namespace), 

398 sanitize_log_value(name), 

399 sanitize_log_value(exc), 

400 ) 

401 scheduling = _empty_scheduling_info() 

402 scheduling["node_lookup_error"] = str(exc) 

403 

404 # Denormalized onto each pod so a caller iterating pods does not have to 

405 # join against the scheduling block to answer "what did this pod run on". 

406 node_index = {node["name"]: node for node in scheduling["nodes"]} 

407 for pod_entry in pod_list: 

408 node = node_index.get(pod_entry.get("spec", {}).get("nodeName")) 

409 pod_entry["node"] = ( 

410 { 

411 "name": node["name"], 

412 "instance_type": node["instance_type"], 

413 "capacity_type": node["capacity_type"], 

414 "labels": node["labels"], 

415 } 

416 if node 

417 else None 

418 ) 

419 

420 response = { 

421 "cluster_id": processor.cluster_id, 

422 "region": processor.region, 

423 "timestamp": datetime.now(UTC).isoformat(), 

424 "job_name": name, 

425 "namespace": namespace, 

426 "count": len(pod_list), 

427 "pods": pod_list, 

428 "scheduling": scheduling, 

429 } 

430 

431 return JSONResponse(status_code=200, content=response) 

432 

433 except Exception as e: 

434 raise internal_server_error("getting job pods", e) from e 

435 

436 

437@router.get("/{namespace}/{name}/pods/{pod_name}/logs") 

438async def get_pod_logs( 

439 namespace: str, 

440 name: str, 

441 pod_name: str, 

442 container: str | None = Query(None, description="Container name"), 

443 tail: int = Query(100, ge=1, le=10000, description="Number of lines from the end"), 

444 previous: bool = Query(False, description="Get logs from previous terminated container"), 

445) -> Response: 

446 """Get logs from a specific pod belonging to a Job.""" 

447 processor = _check_processor() 

448 _check_namespace(namespace, processor) 

449 

450 try: 

451 pod = processor.core_v1.read_namespaced_pod(name=pod_name, namespace=namespace) 

452 job_name_label = pod.metadata.labels.get("job-name") 

453 if job_name_label != name: 

454 raise HTTPException( 

455 status_code=400, detail=f"Pod '{pod_name}' does not belong to job '{name}'" 

456 ) 

457 

458 log_kwargs: dict[str, Any] = { 

459 "name": pod_name, 

460 "namespace": namespace, 

461 "tail_lines": tail, 

462 "previous": previous, 

463 } 

464 if container: 

465 log_kwargs["container"] = container 

466 

467 log_response = processor.core_v1.read_namespaced_pod_log( 

468 **log_kwargs, _preload_content=False 

469 ) 

470 logs = _decode_pod_log_response(log_response) 

471 

472 response = { 

473 "cluster_id": processor.cluster_id, 

474 "region": processor.region, 

475 "timestamp": datetime.now(UTC).isoformat(), 

476 "job_name": name, 

477 "namespace": namespace, 

478 "pod_name": pod_name, 

479 "container": container, 

480 "logs": logs, 

481 } 

482 

483 return JSONResponse(status_code=200, content=response) 

484 

485 except HTTPException: 

486 raise 

487 except Exception as e: 

488 if "NotFound" in str(e) or "404" in str(e): 

489 raise HTTPException(status_code=404, detail=f"Pod '{pod_name}' not found") from e 

490 raise internal_server_error("getting pod logs", e) from e 

491 

492 

493@router.get("/{namespace}/{name}/metrics") 

494async def get_job_metrics(namespace: str, name: str) -> Response: 

495 """Get resource usage metrics for a Job's pods.""" 

496 processor = _check_processor() 

497 _check_namespace(namespace, processor) 

498 

499 try: 

500 pods = processor.core_v1.list_namespaced_pod( 

501 namespace=namespace, label_selector=f"job-name={name}" 

502 ) 

503 

504 if not pods.items: 

505 raise HTTPException(status_code=404, detail=f"No pods found for job '{name}'") 

506 

507 pod_metrics = [] 

508 total_cpu_millicores = 0 

509 total_memory_bytes = 0 

510 

511 try: 

512 for pod in pods.items: 

513 try: 

514 metrics = processor.custom_objects.get_namespaced_custom_object( 

515 group="metrics.k8s.io", 

516 version="v1beta1", 

517 namespace=namespace, 

518 plural="pods", 

519 name=pod.metadata.name, 

520 ) 

521 

522 containers_metrics = [] 

523 for container in metrics.get("containers", []): 

524 cpu_str = container.get("usage", {}).get("cpu", "0") 

525 memory_str = container.get("usage", {}).get("memory", "0") 

526 

527 cpu_millicores = 0 

528 if cpu_str.endswith("n"): 

529 cpu_millicores = int(cpu_str[:-1]) // 1000000 

530 elif cpu_str.endswith("m"): 

531 cpu_millicores = int(cpu_str[:-1]) 

532 else: 

533 cpu_millicores = int(cpu_str) * 1000 

534 

535 memory_bytes = 0 

536 if memory_str.endswith("Ki"): 

537 memory_bytes = int(memory_str[:-2]) * 1024 

538 elif memory_str.endswith("Mi"): 

539 memory_bytes = int(memory_str[:-2]) * 1024 * 1024 

540 elif memory_str.endswith("Gi"): 

541 memory_bytes = int(memory_str[:-2]) * 1024 * 1024 * 1024 

542 else: 

543 memory_bytes = int(memory_str) 

544 

545 total_cpu_millicores += cpu_millicores 

546 total_memory_bytes += memory_bytes 

547 

548 containers_metrics.append( 

549 { 

550 "name": container.get("name"), 

551 "cpu_millicores": cpu_millicores, 

552 "memory_bytes": memory_bytes, 

553 "memory_mib": round(memory_bytes / (1024 * 1024), 2), 

554 } 

555 ) 

556 

557 pod_metrics.append( 

558 {"pod_name": pod.metadata.name, "containers": containers_metrics} 

559 ) 

560 

561 except Exception as e: 

562 logger.warning(f"Could not get metrics for pod {pod.metadata.name}: {e}") 

563 pod_metrics.append( 

564 {"pod_name": pod.metadata.name, "error": "Metrics not available"} 

565 ) 

566 

567 except Exception as e: 

568 logger.warning(f"Metrics API not available: {e}") 

569 

570 response = { 

571 "cluster_id": processor.cluster_id, 

572 "region": processor.region, 

573 "timestamp": datetime.now(UTC).isoformat(), 

574 "job_name": name, 

575 "namespace": namespace, 

576 "summary": { 

577 "total_cpu_millicores": total_cpu_millicores, 

578 "total_memory_bytes": total_memory_bytes, 

579 "total_memory_mib": round(total_memory_bytes / (1024 * 1024), 2), 

580 "pod_count": len(pods.items), 

581 }, 

582 "pods": pod_metrics, 

583 } 

584 

585 return JSONResponse(status_code=200, content=response) 

586 

587 except HTTPException: 

588 raise 

589 except Exception as e: 

590 raise internal_server_error("getting job metrics", e) from e 

591 

592 

593@router.delete("/{namespace}/{name}") 

594async def delete_job( 

595 namespace: str, 

596 name: str, 

597 expected_uid: str | None = Query( 

598 None, 

599 min_length=1, 

600 max_length=128, 

601 description="Optional immutable Kubernetes UID deletion precondition", 

602 ), 

603) -> Response: 

604 """Delete a Job, optionally requiring its immutable Kubernetes UID.""" 

605 processor = _check_processor() 

606 _check_namespace(namespace, processor) 

607 

608 try: 

609 delete_kwargs: dict[str, Any] = { 

610 "name": name, 

611 "namespace": namespace, 

612 } 

613 deleted_uid: str | None = None 

614 if expected_uid is not None: 

615 current = processor.batch_v1.read_namespaced_job(name=name, namespace=namespace) 

616 deleted_uid = str(getattr(current.metadata, "uid", "") or "") 

617 if deleted_uid != expected_uid: 

618 raise HTTPException( 

619 status_code=409, 

620 detail=( 

621 f"Job '{name}' UID changed; expected {expected_uid!r}, " 

622 f"found {deleted_uid or 'unknown'!r}" 

623 ), 

624 ) 

625 delete_kwargs["body"] = kubernetes_client.V1DeleteOptions( 

626 propagation_policy="Background", 

627 preconditions=kubernetes_client.V1Preconditions(uid=expected_uid), 

628 ) 

629 else: 

630 delete_kwargs["propagation_policy"] = "Background" 

631 processor.batch_v1.delete_namespaced_job(**delete_kwargs) 

632 

633 response = { 

634 "cluster_id": processor.cluster_id, 

635 "region": processor.region, 

636 "timestamp": datetime.now(UTC).isoformat(), 

637 "job_name": name, 

638 "namespace": namespace, 

639 "uid": deleted_uid, 

640 "status": "deleted", 

641 "message": "Job deleted successfully", 

642 } 

643 

644 return JSONResponse(status_code=200, content=response) 

645 

646 except HTTPException: 

647 raise 

648 except Exception as e: 

649 status = getattr(e, "status", None) 

650 if status == 409: 

651 raise HTTPException( 

652 status_code=409, 

653 detail=f"Job '{name}' changed before its UID-preconditioned delete", 

654 ) from e 

655 if status == 404 or "NotFound" in str(e) or "404" in str(e): 

656 raise HTTPException( 

657 status_code=404, detail=f"Job '{name}' not found in namespace '{namespace}'" 

658 ) from e 

659 raise internal_server_error("deleting job", e) from e 

660 

661 

662@router.delete("") 

663async def bulk_delete_jobs(request: BulkDeleteRequest) -> Response: 

664 """Bulk delete jobs based on filters.""" 

665 

666 processor = _check_processor() 

667 

668 try: 

669 selector_requirements = _parse_exact_label_selector(request.label_selector) 

670 status_filter = request.status.value if request.status else None 

671 all_jobs = await processor.list_jobs( 

672 namespace=request.namespace, status_filter=status_filter 

673 ) 

674 

675 jobs_to_delete = [] 

676 cutoff_time = None 

677 if request.older_than_days: 

678 cutoff_time = datetime.now(UTC) - timedelta(days=request.older_than_days) 

679 

680 for job in all_jobs: 

681 if cutoff_time: 

682 created_str = job.get("metadata", {}).get("creationTimestamp") 

683 if created_str: 

684 created = datetime.fromisoformat(created_str.replace("Z", "+00:00")) 

685 # Kubernetes normally returns an aware RFC3339 timestamp, 

686 # but older fixtures/clients may provide a naive value. 

687 # Normalize both forms to aware UTC before comparison. 

688 if created.tzinfo is None: 

689 created = created.replace(tzinfo=UTC) 

690 else: 

691 created = created.astimezone(UTC) 

692 if created > cutoff_time: 

693 continue 

694 

695 if selector_requirements and not _labels_match( 

696 job.get("metadata", {}).get("labels", {}), selector_requirements 

697 ): 

698 continue 

699 

700 jobs_to_delete.append(job) 

701 

702 deleted_jobs = [] 

703 failed_jobs = [] 

704 

705 if not request.dry_run: 

706 for job in jobs_to_delete: 

707 job_name = job.get("metadata", {}).get("name") 

708 job_namespace = job.get("metadata", {}).get("namespace") 

709 try: 

710 processor.batch_v1.delete_namespaced_job( 

711 name=job_name, namespace=job_namespace, propagation_policy="Background" 

712 ) 

713 deleted_jobs.append({"name": job_name, "namespace": job_namespace}) 

714 except Exception as e: 

715 failed_jobs.append( 

716 {"name": job_name, "namespace": job_namespace, "error": str(e)} 

717 ) 

718 

719 response: dict[str, Any] = { 

720 "cluster_id": processor.cluster_id, 

721 "region": processor.region, 

722 "timestamp": datetime.now(UTC).isoformat(), 

723 "dry_run": request.dry_run, 

724 "total_matched": len(jobs_to_delete), 

725 "deleted_count": len(deleted_jobs), 

726 "failed_count": len(failed_jobs), 

727 "jobs": ( 

728 [ 

729 { 

730 "name": j.get("metadata", {}).get("name"), 

731 "namespace": j.get("metadata", {}).get("namespace"), 

732 } 

733 for j in jobs_to_delete 

734 ] 

735 if request.dry_run 

736 else deleted_jobs 

737 ), 

738 "failed": failed_jobs if failed_jobs else None, 

739 } 

740 

741 return JSONResponse(status_code=200, content=response) 

742 

743 except ValueError as e: 

744 raise HTTPException(status_code=400, detail=str(e)) from e 

745 except Exception as e: 

746 raise internal_server_error("bulk deleting jobs", e) from e 

747 

748 

749@router.post("/{namespace}/{name}/retry") 

750async def retry_job(namespace: str, name: str) -> Response: 

751 """Retry a failed job by creating a new job from its spec.""" 

752 processor = _check_processor() 

753 _check_namespace(namespace, processor) 

754 

755 try: 

756 try: 

757 original_job = processor.batch_v1.read_namespaced_job(name=name, namespace=namespace) 

758 except Exception as e: 

759 if "NotFound" in str(e) or "404" in str(e): 

760 raise HTTPException( 

761 status_code=404, detail=f"Job '{name}' not found in namespace '{namespace}'" 

762 ) from e 

763 raise 

764 

765 new_name = f"{name}-retry-{datetime.now(UTC).strftime('%Y%m%d%H%M%S')}" 

766 

767 new_job_manifest = { 

768 "apiVersion": "batch/v1", 

769 "kind": "Job", 

770 "metadata": { 

771 "name": new_name, 

772 "namespace": namespace, 

773 "labels": { 

774 **(original_job.metadata.labels or {}), 

775 "gco.io/retry-of": name, 

776 }, 

777 "annotations": { 

778 **(original_job.metadata.annotations or {}), 

779 "gco.io/original-job": name, 

780 }, 

781 }, 

782 "spec": { 

783 "parallelism": original_job.spec.parallelism, 

784 "completions": original_job.spec.completions, 

785 "backoffLimit": original_job.spec.backoff_limit, 

786 "template": original_job.spec.template.to_dict(), 

787 }, 

788 } 

789 

790 spec_dict = cast(dict[str, Any], new_job_manifest["spec"]) 

791 template_dict = spec_dict.get("template", {}) 

792 if isinstance(template_dict, dict) and "status" in template_dict: 

793 del template_dict["status"] 

794 

795 submission_request = ManifestSubmissionRequest( 

796 manifests=[new_job_manifest], namespace=namespace, dry_run=False, validate=True 

797 ) 

798 

799 result = await processor.process_manifest_submission(submission_request) 

800 

801 response = { 

802 "cluster_id": processor.cluster_id, 

803 "region": processor.region, 

804 "timestamp": datetime.now(UTC).isoformat(), 

805 "original_job": name, 

806 "new_job": new_name, 

807 "namespace": namespace, 

808 "success": result.success, 

809 "message": ( 

810 "Job retry created successfully" if result.success else "Failed to create retry job" 

811 ), 

812 "errors": result.errors, 

813 } 

814 

815 status_code = 201 if result.success else 400 

816 return JSONResponse(status_code=status_code, content=response) 

817 

818 except HTTPException: 

819 raise 

820 except Exception as e: 

821 raise internal_server_error("retrying job", e) from e