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
« 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."""
3from __future__ import annotations
5import logging
6import re
7from datetime import UTC, datetime, timedelta
8from typing import Any, cast
10from fastapi import APIRouter, HTTPException, Query
11from fastapi.responses import JSONResponse, Response
12from kubernetes import client as kubernetes_client
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
28router = APIRouter(prefix="/api/v1/jobs", tags=["Jobs"])
29logger = logging.getLogger(__name__)
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])?$")
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()
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 []
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 )
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
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))
81 return requirements
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)
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()
107 try:
108 selector_requirements = _parse_exact_label_selector(label_selector)
109 all_jobs = await processor.list_jobs(namespace=namespace, status_filter=status)
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 ]
118 sort_field, sort_order = "createdAt", "desc"
119 if ":" in sort:
120 sort_field, sort_order = sort.split(":", 1)
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 ""
131 all_jobs.sort(key=get_sort_key, reverse=(sort_order == "desc"))
133 total = len(all_jobs)
134 paginated_jobs = all_jobs[offset : offset + limit]
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 }
148 return JSONResponse(status_code=200, content=response)
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
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.
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
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)
188 try:
189 job = processor.batch_v1.read_namespaced_job(name=name, namespace=namespace)
190 job_info = _parse_job_to_dict(job)
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 }
200 return JSONResponse(status_code=200, content=response)
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
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
225 processor = _check_processor()
226 _check_namespace(namespace, processor)
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
239 pods = processor.core_v1.list_namespaced_pod(
240 namespace=namespace, label_selector=f"job-name={name}"
241 )
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 )
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]
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 )
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
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
303 available_containers = [c.name for c in pod.spec.containers]
304 init_containers = [c.name for c in (pod.spec.init_containers or [])]
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 }
321 return JSONResponse(status_code=200, content=response)
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
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)
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 )
346 pods = processor.core_v1.list_namespaced_pod(
347 namespace=namespace, label_selector=f"job-name={name}"
348 )
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)
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 )
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 }
374 return JSONResponse(status_code=200, content=response)
376 except Exception as e:
377 raise internal_server_error("getting job events", e) from e
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)
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]
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)
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 )
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 }
431 return JSONResponse(status_code=200, content=response)
433 except Exception as e:
434 raise internal_server_error("getting job pods", e) from e
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)
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 )
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
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)
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 }
483 return JSONResponse(status_code=200, content=response)
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
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)
499 try:
500 pods = processor.core_v1.list_namespaced_pod(
501 namespace=namespace, label_selector=f"job-name={name}"
502 )
504 if not pods.items:
505 raise HTTPException(status_code=404, detail=f"No pods found for job '{name}'")
507 pod_metrics = []
508 total_cpu_millicores = 0
509 total_memory_bytes = 0
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 )
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")
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
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)
545 total_cpu_millicores += cpu_millicores
546 total_memory_bytes += memory_bytes
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 )
557 pod_metrics.append(
558 {"pod_name": pod.metadata.name, "containers": containers_metrics}
559 )
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 )
567 except Exception as e:
568 logger.warning(f"Metrics API not available: {e}")
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 }
585 return JSONResponse(status_code=200, content=response)
587 except HTTPException:
588 raise
589 except Exception as e:
590 raise internal_server_error("getting job metrics", e) from e
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)
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)
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 }
644 return JSONResponse(status_code=200, content=response)
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
662@router.delete("")
663async def bulk_delete_jobs(request: BulkDeleteRequest) -> Response:
664 """Bulk delete jobs based on filters."""
666 processor = _check_processor()
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 )
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)
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
695 if selector_requirements and not _labels_match(
696 job.get("metadata", {}).get("labels", {}), selector_requirements
697 ):
698 continue
700 jobs_to_delete.append(job)
702 deleted_jobs = []
703 failed_jobs = []
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 )
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 }
741 return JSONResponse(status_code=200, content=response)
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
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)
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
765 new_name = f"{name}-retry-{datetime.now(UTC).strftime('%Y%m%d%H%M%S')}"
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 }
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"]
795 submission_request = ManifestSubmissionRequest(
796 manifests=[new_job_manifest], namespace=namespace, dry_run=False, validate=True
797 )
799 result = await processor.process_manifest_submission(submission_request)
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 }
815 status_code = 201 if result.success else 400
816 return JSONResponse(status_code=status_code, content=response)
818 except HTTPException:
819 raise
820 except Exception as e:
821 raise internal_server_error("retrying job", e) from e