Codeseys/composer-replication-framework
0
1"""EKSExecutor — production Amazon EKS / Kubernetes-backed serverless executor.2 3This is the v0-finished k8s sibling of `ModalSpawnExecutor`. It implements4the `ServerlessExecutor` Protocol against the Kubernetes ``BatchV1Api`` using5the **single Indexed Job** topology recommended for gang-scheduled DiLoCo6replicas.7 8Topology (the load-bearing design choice)9------------------------------------------10There are two ways to map N replicas onto k8s:11 12 (A) ONE Indexed Job — ``completions=N, parallelism=N,13 completionMode='Indexed'``. The control plane assigns each pod a14 ``JOB_COMPLETION_INDEX`` 0..N-1 which IS the rank, all pods share one15 rendezvous URI, scheduling is atomic, and a single delete cancels the16 whole cohort.17 (B) N separate non-indexed Jobs, one per rank.18 19`EKSExecutor` uses **(A)** because it is the better fit for DiLoCo: rank20assignment is free, scheduling is gang-atomic, and one delete tears down the21cohort — which matches ``ObjectStoreAllReduce``'s all-or-nothing barrier. The22reconciliation with the per-replica ``ReplicaHandle`` contract: ``launch_replicas``23creates ONE Indexed Job but still returns N ``ReplicaHandle`` objects24(``handles[i].rank == i``) whose ``metadata`` stores the SHARED25``job_name``/``namespace`` plus that rank.26 27This is materially different from ``ModalSpawnExecutor`` where each handle is28an independent ``FunctionCall``:29 30 * ``poll(handle)`` reads the single Job status and checks whether31 ``handle.rank`` is in the run-length-compressed ``completed_indexes`` /32 ``failed_indexes`` strings.33 * ``cancel(handle)`` on ANY handle deletes the WHOLE Job (intentional gang34 semantics — cancelling one rank tears down the whole replica cohort).35 36Rank plumbing37-------------38The repo's ``replica_entrypoint`` reads ``REPLICA_RANK``. We bridge the k8s39completion-index to that env var via the downward API rather than relying on40the auto-injected ``JOB_COMPLETION_INDEX``::41 42 V1EnvVar(43 name="REPLICA_RANK",44 value_from=V1EnvVarSource(field_ref=V1ObjectFieldSelector(45 field_path="metadata.annotations['batch.kubernetes.io/job-completion-index']")),46 )47 48so the unchanged entrypoint's ``REPLICA_RANK`` read just works. ``WORLD_SIZE``49is set as a literal env var.50 51S3 rendezvous via IRSA / Pod Identity52-------------------------------------53``EKSExecutor`` accepts ``service_account_name`` and references it on the54PodSpec. The EKS Pod Identity / IRSA mutating webhook then injects55``AWS_ROLE_ARN`` + ``AWS_WEB_IDENTITY_TOKEN_FILE`` (and a projected token56volume) into the pod, so ``boto3``/``s3fs``/``fsspec`` pick up credentials via57the web-identity provider with ZERO code change inside the replica — the58``s3://`` rendezvous works out of the box. ``EKSExecutor`` only REFERENCES a59pre-annotated ServiceAccount; it never creates IAM/OIDC resources.60 61Sandboxing (advanced, optional)62-------------------------------63``runtime_class_name`` references a pre-existing cluster-scoped ``RuntimeClass``64(``runsc`` for gVisor, ``kata`` for Kata). It defaults to ``None``.65 66.. warning::67 Combining ``gpu`` with ``runtime_class_name`` is advanced and unverified.68 gVisor (runsc) needs ``nvproxy`` enabled and only supports a fixed allowlist69 of NVIDIA driver families; Kata runs a microVM that caps CPU/mem and needs70 GPU passthrough (PCIe/IOMMU + NVIDIA Kata Manager + CDI). Do not silently71 combine the two without operator validation. ``EKSExecutor`` cannot create72 the RuntimeClass — it only references one that already exists.73 74References75----------76- k8s Indexed Jobs: https://kubernetes.io/docs/tasks/job/indexed-parallel-processing-static/77- kubernetes-client/python job_crud example + V1JobSpec / V1JobStatus docs78- EKS IRSA: https://docs.aws.amazon.com/eks/latest/userguide/iam-roles-for-service-accounts.html79- ADR-005 (executor protocol design)80"""81from __future__ import annotations82 83import time84import uuid85from collections.abc import Callable, Mapping86from typing import Any87 88from composer_replication.diloco.serverless.executor import (89 ReplicaHandle,90)91 92# Logical GPU spec ("A100"/"H100") -> (gpu_count_string, node_selector merge).93# The Protocol's `gpu` arg is a logical name; map it to a concrete EKS node94# class + GPU count rather than passing the opaque string straight through.95_GPU_SPEC_TABLE: dict[str, tuple[str, dict[str, str]]] = {96 "A100": ("1", {"node.kubernetes.io/instance-type": "p4d.24xlarge"}),97 "H100": ("1", {"node.kubernetes.io/instance-type": "p5.48xlarge"}),98 "A10G": ("1", {"node.kubernetes.io/instance-type": "g5.xlarge"}),99 "T4": ("1", {"node.kubernetes.io/instance-type": "g4dn.xlarge"}),100}101 102 103def _expand_indexes(spec: str | None) -> set[int]:104 """Expand a run-length-compressed completion-index string to a set.105 106 The k8s ``V1JobStatus.completed_indexes`` / ``failed_indexes`` fields are107 strings like ``"1,3-5,7"`` (comma-separated singletons and ``a-b`` ranges).108 ``_expand_indexes("1,3-5,7") == {1, 3, 4, 5, 7}``. Empty/None -> empty set.109 """110 out: set[int] = set()111 if not spec:112 return out113 for token in spec.split(","):114 token = token.strip()115 if not token:116 continue117 if "-" in token:118 lo_s, _, hi_s = token.partition("-")119 try:120 lo, hi = int(lo_s), int(hi_s)121 except ValueError:122 continue123 if hi < lo:124 lo, hi = hi, lo125 out.update(range(lo, hi + 1))126 else:127 try:128 out.add(int(token))129 except ValueError:130 continue131 return out132 133 134class EKSExecutor:135 """Run N DiLoCo replicas as a single Kubernetes Indexed Job on EKS.136 137 Implements the `ServerlessExecutor` Protocol. ``launch_replicas`` creates138 ONE Indexed Job (``completions == parallelism == n_replicas``,139 ``completionMode='Indexed'``) and returns N ``ReplicaHandle`` objects that140 all share the same ``job_name``/``namespace`` (gang semantics).141 142 Args:143 image: container image that has ``composer_replication`` installed and144 runs the replica entrypoint.145 namespace: k8s namespace for the Job. Default ``"default"``.146 service_account_name: ServiceAccount to attach to the PodSpec for IRSA /147 EKS Pod Identity S3 access. ``EKSExecutor`` references it; it does148 NOT create it or any IAM/OIDC resources.149 node_selector: extra node selector merged into the GPU node selector.150 tolerations: PodSpec tolerations. If GPU is requested and the caller did151 not supply tolerations, the standard ``nvidia.com/gpu`` NoSchedule152 toleration is added automatically.153 runtime_class_name: optional pre-existing RuntimeClass (e.g. ``"gvisor"``154 / ``"kata"``). Default ``None``. See the module-level warning before155 combining with ``gpu``.156 command: container command. Defaults to the repo replica entrypoint157 module ``["python", "-m",158 "composer_replication.diloco.serverless.replica_entrypoint"]``.159 cpu_request / memory_request: PodSpec resource requests.160 ttl_seconds_after_finished: auto-delete the finished Job (and its pods,161 cascadingly) after this many seconds. Default 3600.162 backoff_limit: Job retry budget. Default 0 (fail-fast — RL gangs usually163 do NOT want the k8s default of 6 retries).164 gpu_resource_key: the GPU resource key. Default ``"nvidia.com/gpu"``.165 run_id: optional run id baked into the generated Job name.166 batch_api / core_api: dependency-injected ``BatchV1Api`` / ``CoreV1Api``167 instances. When ``None`` (the default), they are built lazily on168 first use via in-cluster or kube-config loading. Tests inject mocks.169 170 Raises:171 RuntimeError: if the ``kubernetes`` client is not installed AND no api172 was injected (the import is needed to construct V1 model objects).173 """174 175 backend_name = "eks"176 # Pods are network-isolated by default; rendezvous is S3 (ObjectStoreAllReduce).177 supports_inter_replica_network = False178 179 def __init__(180 self,181 image: str,182 *,183 namespace: str = "default",184 service_account_name: str | None = None,185 node_selector: dict[str, str] | None = None,186 tolerations: list[Any] | None = None,187 runtime_class_name: str | None = None,188 command: list[str] | None = None,189 cpu_request: str = "4",190 memory_request: str = "16Gi",191 ttl_seconds_after_finished: int = 3600,192 backoff_limit: int = 0,193 gpu_resource_key: str = "nvidia.com/gpu",194 run_id: str | None = None,195 batch_api: Any = None,196 core_api: Any = None,197 ) -> None:198 # `kubernetes` is only strictly required when we have to BUILD V1 model199 # objects ourselves (launch_replicas) or load cluster config (when no200 # api is injected). We surface a clear error here only if we definitely201 # need it and it is absent — i.e. when no api was injected. When apis202 # ARE injected (tests, or callers that pre-built clients), we tolerate a203 # missing top-level package and lazy-import `client` per call.204 if batch_api is None or core_api is None:205 try:206 import kubernetes # noqa: F401207 except ImportError as e:208 raise RuntimeError(209 'EKSExecutor requires the kubernetes client: '210 'pip install "kubernetes>=29" (or '211 "`pip install -e .[serverless]`). Got: " + repr(e)212 ) from e213 214 self.image = image215 self.namespace = namespace216 self.service_account_name = service_account_name217 self.node_selector = dict(node_selector) if node_selector else None218 self.tolerations = list(tolerations) if tolerations else None219 self.runtime_class_name = runtime_class_name220 self.command = command or [221 "python",222 "-m",223 "composer_replication.diloco.serverless.replica_entrypoint",224 ]225 self.cpu_request = cpu_request226 self.memory_request = memory_request227 self.ttl_seconds_after_finished = ttl_seconds_after_finished228 self.backoff_limit = backoff_limit229 self.gpu_resource_key = gpu_resource_key230 self.run_id = run_id or "diloco"231 232 self._batch_api = batch_api233 self._core_api = core_api234 # rank -> {"job_name", "namespace", "result"}; lets poll/collect cache.235 self._handles: dict[int, dict[str, Any]] = {}236 237 # -----------------------------------------------------------------238 # Lazy client construction (config loading only when not injected)239 # -----------------------------------------------------------------240 241 def _load_config(self) -> None:242 """Load k8s config once: in-cluster first, then ~/.kube/config."""243 from kubernetes import config244 245 try:246 config.load_incluster_config()247 except config.ConfigException:248 config.load_kube_config()249 250 def _batch(self) -> Any:251 if self._batch_api is None:252 from kubernetes import client253 254 self._load_config()255 self._batch_api = client.BatchV1Api()256 return self._batch_api257 258 def _core(self) -> Any:259 if self._core_api is None:260 from kubernetes import client261 262 self._load_config()263 self._core_api = client.CoreV1Api()264 return self._core_api265 266 # -----------------------------------------------------------------267 # Job-spec construction268 # -----------------------------------------------------------------269 270 def _build_env(271 self, world_size: int, entrypoint_args: Mapping[str, Any]272 ) -> list[Any]:273 """Build the container env list, including the downward-API rank var."""274 from kubernetes import client275 276 env: list[Any] = [277 # REPLICA_RANK from the per-pod completion-index annotation via the278 # downward API — bridges k8s indexing to the repo entrypoint's279 # REPLICA_RANK read with no entrypoint change.280 client.V1EnvVar(281 name="REPLICA_RANK",282 value_from=client.V1EnvVarSource(283 field_ref=client.V1ObjectFieldSelector(284 field_path=(285 "metadata.annotations["286 "'batch.kubernetes.io/job-completion-index']"287 )288 )289 ),290 ),291 client.V1EnvVar(name="WORLD_SIZE", value=str(world_size)),292 ]293 # rendezvous_uri (and any other scalar kwargs) passed as literal env so294 # the entrypoint / user code can read them. `rank_env` is the295 # LocalProcessExecutor convention — drop it (same as ModalSpawnExecutor).296 for key, value in entrypoint_args.items():297 if key == "rank_env":298 continue299 if isinstance(value, (str, int, float, bool)):300 env.append(301 client.V1EnvVar(name=key.upper(), value=str(value))302 )303 return env304 305 def _build_resources(self, gpu: str | None) -> tuple[Any, dict[str, str], list[Any]]:306 """Build V1ResourceRequirements + (node_selector, tolerations) for GPU.307 308 Returns (resources, node_selector, tolerations). The GPU count is309 ALWAYS a STRING ('1', not int 1) — the OpenAPI type for the limits map310 is dict[str, str] and an int can serialize wrong or raise.311 """312 from kubernetes import client313 314 requests = {"cpu": self.cpu_request, "memory": self.memory_request}315 limits: dict[str, str] = {}316 node_selector: dict[str, str] = dict(self.node_selector or {})317 tolerations: list[Any] = list(self.tolerations or [])318 319 if gpu is not None:320 gpu_count, gpu_node_selector = _GPU_SPEC_TABLE.get(321 gpu, ("1", {})322 )323 # STRING, always.324 limits[self.gpu_resource_key] = str(gpu_count)325 # Merge the mapped node selector under any caller-supplied one326 # (caller wins on key conflicts).327 for k, v in gpu_node_selector.items():328 node_selector.setdefault(k, v)329 # Auto-add the GPU NoSchedule toleration unless the caller overrode330 # tolerations explicitly.331 if not self.tolerations:332 tolerations.append(333 client.V1Toleration(334 key=self.gpu_resource_key,335 operator="Exists",336 effect="NoSchedule",337 )338 )339 340 resources = client.V1ResourceRequirements(341 requests=requests,342 limits=limits or None,343 )344 return resources, node_selector, tolerations345 346 def _build_job(347 self,348 *,349 job_name: str,350 n_replicas: int,351 gpu: str | None,352 timeout: int,353 entrypoint_args: Mapping[str, Any],354 ) -> Any:355 """Assemble the full V1Job (Indexed) bottom-up."""356 from kubernetes import client357 358 env = self._build_env(n_replicas, entrypoint_args)359 resources, node_selector, tolerations = self._build_resources(gpu)360 361 container = client.V1Container(362 name="replica",363 image=self.image,364 command=list(self.command),365 env=env,366 resources=resources,367 )368 369 pod_spec = client.V1PodSpec(370 restart_policy="Never", # required for Indexed jobs / fail-fast RL371 containers=[container],372 service_account_name=self.service_account_name,373 node_selector=node_selector or None,374 tolerations=tolerations or None,375 runtime_class_name=self.runtime_class_name,376 )377 378 labels = {"app": "composer-diloco", "job-name": job_name}379 pod_template = client.V1PodTemplateSpec(380 metadata=client.V1ObjectMeta(labels=labels),381 spec=pod_spec,382 )383 384 job_spec = client.V1JobSpec(385 template=pod_template,386 completions=n_replicas,387 parallelism=n_replicas,388 completion_mode="Indexed",389 backoff_limit=self.backoff_limit,390 ttl_seconds_after_finished=self.ttl_seconds_after_finished,391 active_deadline_seconds=timeout,392 )393 394 return client.V1Job(395 api_version="batch/v1",396 kind="Job",397 metadata=client.V1ObjectMeta(name=job_name, labels=labels),398 spec=job_spec,399 )400 401 # -----------------------------------------------------------------402 # ServerlessExecutor Protocol403 # -----------------------------------------------------------------404 405 def launch_replicas(406 self,407 n_replicas: int,408 entrypoint: str | Callable[..., Any],409 entrypoint_args: Mapping[str, Any],410 *,411 gpu: str | None = None,412 timeout: int = 3600,413 ) -> list[ReplicaHandle]:414 """Create ONE Indexed Job of N pods and return N rank-ordered handles.415 416 ``entrypoint`` is ignored when it names a Callable (k8s runs a container417 command, not an in-process callable); the container command is fixed at418 construction (``command`` ctor arg). The repo entrypoint module is the419 default. ``entrypoint_args`` scalar kwargs are passed as upper-cased env420 vars so ``replica_entrypoint`` / user code can read them. ``gpu`` maps to421 a ``nvidia.com/gpu`` limit + node selector; ``timeout`` becomes the Job's422 ``active_deadline_seconds`` hard wall-clock kill.423 """424 del entrypoint # k8s runs a container command, not an in-process fn425 426 if n_replicas < 1:427 raise ValueError(f"n_replicas must be >= 1, got {n_replicas}")428 429 job_name = f"{self.run_id}-{uuid.uuid4().hex[:8]}"430 job = self._build_job(431 job_name=job_name,432 n_replicas=n_replicas,433 gpu=gpu,434 timeout=timeout,435 entrypoint_args=entrypoint_args,436 )437 438 self._batch().create_namespaced_job(namespace=self.namespace, body=job)439 440 handles: list[ReplicaHandle] = []441 for rank in range(n_replicas):442 handles.append(443 ReplicaHandle(444 rank=rank,445 backend_name=self.backend_name,446 metadata={447 "job_name": job_name,448 "namespace": self.namespace,449 "rank": rank,450 },451 )452 )453 self._handles[rank] = {454 "job_name": job_name,455 "namespace": self.namespace,456 "result": None,457 }458 return handles459 460 def poll(self, handle: ReplicaHandle) -> str:461 """Poll this rank's status off the shared Indexed Job.462 463 Reads ``read_namespaced_job_status`` once, then maps the whole-job464 status to this rank: ``rank in completed_indexes`` -> ``succeeded``;465 ``rank in failed_indexes`` -> ``failed``; ``active > 0`` -> ``running``;466 else ``pending``. A 404 (Job deleted/cancelled) -> ``cancelled``.467 468 Returns one of: ``pending`` | ``running`` | ``succeeded`` | ``failed`` |469 ``cancelled``.470 """471 from kubernetes.client.exceptions import ApiException472 473 job_name = handle.metadata["job_name"]474 namespace = handle.metadata["namespace"]475 rank = handle.metadata.get("rank", handle.rank)476 477 try:478 status = self._batch().read_namespaced_job_status(479 name=job_name, namespace=namespace480 ).status481 except ApiException as e:482 if getattr(e, "status", None) == 404:483 return "cancelled"484 raise485 486 completed = _expand_indexes(getattr(status, "completed_indexes", None))487 if rank in completed:488 return "succeeded"489 490 failed = _expand_indexes(getattr(status, "failed_indexes", None))491 if rank in failed:492 return "failed"493 494 # Whole-job terminal Failed (e.g. DeadlineExceeded / backoff) with no495 # per-index attribution -> treat this rank as failed.496 for cond in (getattr(status, "conditions", None) or []):497 if (498 getattr(cond, "type", None) == "Failed"499 and getattr(cond, "status", None) == "True"500 ):501 return "failed"502 503 active = getattr(status, "active", None) or 0504 if active > 0:505 return "running"506 return "pending"507 508 def stream_logs(self, handle: ReplicaHandle, *, n_lines: int = 200) -> str:509 """Read recent logs for this rank's pod.510 511 Finds the pod whose ``batch.kubernetes.io/job-completion-index``512 annotation (or label) equals the rank, then reads its log tail. Returns513 a placeholder string (rather than raising) when the pod has not started514 or the Job is gone — mirrors ``LocalProcessExecutor``.515 """516 from kubernetes.client.exceptions import ApiException517 518 job_name = handle.metadata["job_name"]519 namespace = handle.metadata["namespace"]520 rank = handle.metadata.get("rank", handle.rank)521 idx_key = "batch.kubernetes.io/job-completion-index"522 523 try:524 pods = self._core().list_namespaced_pod(525 namespace=namespace, label_selector=f"job-name={job_name}"526 )527 except ApiException:528 return f"<rank {rank}: job not found / no pods yet>"529 530 pod_name = None531 for pod in getattr(pods, "items", None) or []:532 meta = getattr(pod, "metadata", None)533 annotations = getattr(meta, "annotations", None) or {}534 labels = getattr(meta, "labels", None) or {}535 if annotations.get(idx_key) == str(rank) or labels.get(idx_key) == str(rank):536 pod_name = getattr(meta, "name", None)537 break538 539 if pod_name is None:540 # Fall back to the deterministic name prefix on k8s >= 1.28.541 prefix = f"{job_name}-{rank}-"542 for pod in getattr(pods, "items", None) or []:543 name = getattr(getattr(pod, "metadata", None), "name", "") or ""544 if name.startswith(prefix):545 pod_name = name546 break547 548 if pod_name is None:549 return f"<rank {rank}: pod not started / no logs yet>"550 551 try:552 return self._core().read_namespaced_pod_log(553 name=pod_name,554 namespace=namespace,555 container="replica",556 tail_lines=n_lines,557 )558 except ApiException as e:559 if getattr(e, "status", None) in (400, 404):560 return f"<rank {rank}: pod not started / no logs yet>"561 raise562 563 def cancel(self, handle: ReplicaHandle) -> None:564 """Delete the WHOLE shared Indexed Job (gang teardown).565 566 Because ``EKSExecutor`` uses one shared Indexed Job, cancelling ANY rank567 tears down the entire replica cohort — intentional gang semantics for568 the DiLoCo all-reduce barrier (a single straggler being cancelled should569 not leave the rest spinning and burning GPU).570 571 Uses ``propagation_policy='Background'`` so the pods are cascadingly572 deleted (the k8s default ORPHANS pods, which would keep burning GPU —573 the exact failure mode for RL). Idempotent: a 404 (already deleted) is574 swallowed, and an unknown handle never raises, honoring the Protocol's575 "no exception if already terminated" contract.576 """577 from kubernetes import client578 from kubernetes.client.exceptions import ApiException579 580 job_name = handle.metadata.get("job_name")581 namespace = handle.metadata.get("namespace", self.namespace)582 if not job_name:583 return # unknown handle — no-op584 585 try:586 self._batch().delete_namespaced_job(587 name=job_name,588 namespace=namespace,589 body=client.V1DeleteOptions(590 propagation_policy="Background",591 grace_period_seconds=0,592 ),593 )594 except ApiException as e:595 # R5: swallow ONLY already-terminated signals (404 Not Found, 409596 # Conflict on a job mid-deletion). A genuinely unexpected API error597 # (403 forbidden, 500, malformed request) must NOT be reported as a598 # successful cancel — re-raise so a real teardown failure (leaking599 # GPU-burning pods) is visible rather than silently swallowed.600 if getattr(e, "status", None) in (404, 409):601 return # already deleted / mid-deletion — idempotent no-op602 raise603 604 def collect(605 self,606 handles: list[ReplicaHandle],607 *,608 timeout: int | None = None,609 ) -> list[dict[str, Any]]:610 """Poll until every rank reaches a terminal state or the deadline.611 612 Sleeps between polls (Job status is eventually consistent — do not613 hammer the API server). Returns per-rank result dicts in handles order::614 615 {"rank", "status", "exit_code", "error", "job_name"}616 617 ``exit_code`` is 0 for succeeded, 1 for failed, ``None`` for618 running/pending/cancelled — matching the Protocol's documented shape.619 """620 deadline = time.time() + (timeout if timeout is not None else 86400)621 poll_interval = float(self._collect_poll_interval())622 terminal = {"succeeded", "failed", "cancelled"}623 results_by_rank: dict[int, dict[str, Any]] = {}624 625 pending = list(handles)626 while pending and time.time() < deadline:627 still_pending: list[ReplicaHandle] = []628 for h in pending:629 state = self.poll(h)630 if state in terminal:631 results_by_rank[h.rank] = self._result_dict(h, state)632 else:633 still_pending.append(h)634 pending = still_pending635 if not pending:636 break637 remaining = deadline - time.time()638 if remaining <= 0:639 break640 time.sleep(min(poll_interval, max(0.0, remaining)))641 642 # Any rank still non-terminal at the deadline -> report its last state.643 for h in pending:644 state = self.poll(h)645 results_by_rank[h.rank] = self._result_dict(h, state)646 647 return [results_by_rank[h.rank] for h in handles]648 649 # -----------------------------------------------------------------650 # Internals651 # -----------------------------------------------------------------652 653 def _collect_poll_interval(self) -> float:654 """Seconds between collect() polls. Overridable in tests."""655 return 5.0656 657 @staticmethod658 def _result_dict(handle: ReplicaHandle, state: str) -> dict[str, Any]:659 exit_code = {"succeeded": 0, "failed": 1}.get(state, None)660 error = None661 if state == "failed":662 error = f"rank {handle.rank} reported failed by Job status"663 elif state == "cancelled":664 error = f"rank {handle.rank} Job no longer exists (cancelled)"665 elif state in ("running", "pending"):666 error = f"rank {handle.rank} not terminal at deadline (state={state})"667 return {668 "rank": handle.rank,669 "status": state,670 "exit_code": exit_code,671 "error": error,672 "job_name": handle.metadata.get("job_name"),673 # R6: cross-backend uniformity with Local/Modal/SageMaker collect()674 # shapes. EKS replicas write their real output to the S3 rendezvous675 # (ObjectStoreAllReduce), not back through the k8s API, so the Job676 # status carries no in-band payload — the value is the rendezvous677 # URI when known (callers read the artifact from S3), else None.678 "result": handle.metadata.get("rendezvous_uri"),679 }680 681 682__all__ = ["EKSExecutor"]683 