Team Ai
Modelpublic

Codeseys/composer-replication-framework

sourceHugging Facemitupdated 4mo agoView on Hugging Face
0likes
eks.py683 linesDownload Raw Back to serverless
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