This is an automated email from the ASF dual-hosted git repository. github-merge-queue[bot] pushed a commit to branch gh-readonly-queue/main/pr-7944-14de4583fd8078a8639f458f0eb8e34a3490a9a7 in repository https://gitbox.apache.org/repos/asf/texera.git
commit 33bd07bf27648870ecaf54da0d24c85763186cda Author: Eugene Gu <[email protected]> AuthorDate: Sat Sep 26 05:25:24 2026 +0000 feat(computing-unit): surface failed/unhealthy computing units instead of an endless "Connecting" (#7944) ### What changes were proposed in this PR? When a computing unit becomes unhealthy (the pod is OOM-killed into a crash loop, evicted for disk pressure, stuck on an image pull, or its node becomes unreachable), the UI used to show **"Connecting"** with a "starting up" tooltip forever because status resolution only looked at `pod.status.phase` and the frontend did not represent the complete backend status vocabulary. This PR follows the decisions settled in #7670: mirror Kubernetes for the status vocabulary, expose actionable failure reasons only to authorized viewers, and perform no automatic recovery—users delete and recreate unhealthy units, using the existing termination flow that already handles dead pods. Because a status the UI refuses to act on is only worth as much as the refusal, the change also closes the paths that ignored it: every run entry point now declines to start work on a unit in a terminal state, and a dropdown row the list already greys out can no longer be selected by clicking it. #### Before → after ```text Before: unhealthy unit → Pending/Connecting forever → no actionable explanation After: unhealthy unit → Failed/Unknown + Unit Unavailable → authorized reason in the row tooltip Before: terminating unit → frontend Pending fallback → Connecting After: terminating unit → explicit Terminating state → Shutting Down Before: run on a dead unit → results cleared, request dropped on a closed socket, no feedback After: run on a dead unit → refused with an actionable toast, previous results left on screen Before: click a greyed-out row → unit selected anyway, and remembered for next time After: click a greyed-out row → ignored, as the row's own "Cannot select." tooltip says ``` #### Screenshots <!-- Attach the files from ~/Downloads/pr-7944-screenshots/ in place of each FILE: line below. --> Every frame below is the running application with the computing-unit status rewritten in flight, the automated equivalent of the DevTools response-override recipe. The workflow held one operator, so no "Invalid" or "Empty" state could mask the unit state. **Before.** An unhealthy pod could only be reported as `Pending`, so the run button said **Connecting** for ever and the tooltip claimed the unit was starting up. This one frame is the counterpart to both "after" frames, because every unhealthy or terminating unit used to collapse into this same state. <img width="1362" height="368" alt="01-before-connecting" src="https://github.com/user-attachments/assets/7c3382f0-f6c5-4fd7-8ac7-8b70b696c6a8" /> **After, an unreachable node.** Red badge, a disabled **Unit Unavailable** run button, and the authorized owner-facing reason on the row tooltip. <img width="1362" height="368" alt="02-after-unit-unavailable" src="https://github.com/user-attachments/assets/f8027a33-54b2-4814-a266-3e0d9d415c35" /> **After, a unit being deleted.** Gold transient-state badge and a disabled **Shutting Down**, instead of falling through to `Pending`. <img width="1362" height="368" alt="03-after-shutting-down" src="https://github.com/user-attachments/assets/aad7b5c6-1647-4d04-8d0e-9682ef19dd2c" /> **After, a crash loop caused by an OOM kill.** The reason names what happened and what to do about it. <img width="1362" height="388" alt="04-after-failed-crashloop-oom" src="https://github.com/user-attachments/assets/56c4294c-70ec-4590-b8a8-34a6ba88c193" /> **After, the Form View run button with the socket still connected.** This is the second entry point, which has its own run-button logic. In this exact state it previously fell through to an enabled **Run**; it is now a disabled **Unavailable**. <img width="1494" height="176" alt="05-after-form-view-unavailable" src="https://github.com/user-attachments/assets/12b0edcd-355d-4693-903c-b49ed4602753" /> A recovered OOM-killed unit stays green and runnable, but the tooltip carries the out-of-memory warning with the restart count: <img width="550" height="241" alt="Screenshot 2026-08-23 at 1 05 13 AM" src="https://github.com/user-attachments/assets/555d440c-411a-4b6f-a970-bef3a7ac2555" /> #### Status vocabulary and reasons `ComputingUnitState` now explicitly represents `Running`, `Pending`, `Failed`, `Unknown`, and `Terminating`; truly unrecognized future status strings retain the existing `Pending` fallback. The full mapping from observed pod state to status and authorized `statusReason` is: | Observed pod state | Status | Authorized `statusReason` | |---|---|---| | `deletionTimestamp` set | `Terminating` | — | | phase `Failed` + reason `Evicted`, message mentions ephemeral/disk | `Failed` | "The computing unit was evicted because it ran out of local disk storage. Consider storing less data on the unit's local file system, or recreate it with more storage." | | phase `Failed` + reason `Evicted`, other | `Failed` | "The computing unit was evicted by the cluster (\<first sentence of the cluster message, capped at 120 chars\>). Consider recreating it." | | container waiting `ImagePullBackOff` / `ErrImagePull` / `InvalidImageName` | `Failed` | "The computing unit's image could not be pulled. Please recreate the unit or contact an administrator." | | container waiting `CrashLoopBackOff`, last termination `OOMKilled` | `Failed` | "The computing unit keeps crashing because it runs out of memory. Please terminate it and recreate it with a higher memory limit." | | container waiting `CrashLoopBackOff`, other | `Failed` | "The computing unit is repeatedly crashing (restarted N times). Please terminate and recreate it, or contact an administrator." | | phase `Failed`, not evicted | `Failed` | "The computing unit stopped unexpectedly. Please terminate and recreate it, or contact an administrator." | | phase `Unknown` | `Unknown` | "The state of the computing unit cannot be determined (its node may be unreachable)." | | phase `Pending` + `PodScheduled=False/Unschedulable` condition | `Pending` | "The computing unit is waiting for cluster resources to become available." | | phase `Running` + a container's last termination `OOMKilled` | `Running` | "The last run was terminated because the computing unit ran out of memory (restarted N times). Consider recreating the unit with a higher memory limit before running the same workload." | | phase `Running`, healthy | `Running` | — | | anything else / pod absent / `local` unit | `Pending` / `Running` (unchanged) | — | Precedence is top-to-bottom; in particular, a `CrashLoopBackOff` failure wins over the recovered-OOM warning, and `Terminating` wins over everything. `statusReason` visibility is independent of ownership: owners receive it through the regular list and direct lookup, administrators receive it for every row from the dedicated admin-list endpoint without those rows being marked as owned, and ordinary shared users receive `null` and therefore see generic frontend text. Frontend fallbacks when `statusReason` is absent are `Running` → "Ready to use", `Pending` → "Computing unit is starting up", `Failed`/`Unknown` → "This computing unit is unavailable.", `Terminating` → "Computing unit is shutting down", and any other value → the status word itself. #### Backend - The existing single namespace-level pod listing now yields a `PodStatusSnapshot` per pod (phase, deletion timestamp, pod reason/message, `Unschedulable` condition, and each container's waiting reason, last termination reason, and restart count) via a pure, unit-testable transform, with **no additional Kubernetes round trips**. - With `restartPolicy: Always`, an OOM-killed container restarts in place and the pod phase never leaves `Running`, so OOM kills are detected through `containerStatuses[].lastState.terminated.reason`; a recovered unit stays `Running` with a warning, while a unit that cannot recover and enters `CrashLoopBackOff` is reported as `Failed`. - Vanished-pod reconciliation (#6853/#6854) remains keyed only on pod presence, so a present but `Failed` pod is not treated as vanished. - `local` units remain always `Running`; local liveness remains outside this PR's scope, as settled in #7670. - Dashboard-row construction now receives `canViewStatusReason` separately from `isOwner`, which lets the admin endpoint expose actionable reasons while preserving truthful ownership and ordinary shared-user suppression. - The unused `getComputingUnitStatus` wrapper was removed, and the two equivalent `Running` branches were collapsed without changing status precedence or OOM-warning behavior. #### Frontend - The dashboard DTO status union and the internal `ComputingUnitState` enum both contain all five states, and the status service maps `Terminating` explicitly while preserving `Pending` as the fallback for genuinely unrecognized future values. - One shared `unavailableComputingUnitReason` utility answers whether a unit can never accept work, distinguishing `Terminating` from `Failed`/`Unknown`. Both run buttons and the execution service read the terminal states from it, so the three cannot drift apart. - A selected `Terminating` unit shows a disabled **Shutting Down** run-button state with a loading icon before WebSocket connectivity is considered; invalid and empty workflow messages retain precedence, `Pending` still shows **Connecting**, and `Failed`/`Unknown` still show **Unit Unavailable**. - The Form View run button, which is a second entry point with its own state logic, gained the same terminal states in the same relative position, shortened to **Shutting Down** and **Unavailable** to match its siblings. Its "connecting" condition now excludes terminal units, which would otherwise have spun there for ever, and its Stop branch is offered only while the socket can actually deliver the kill. - `Failed`/`Unknown` use a red dropdown badge, while `Pending`/`Terminating` use gold because both are transient rather than failed states. - No state adds a text suffix next to the unit name; status is conveyed consistently through badge color, row tooltip, and run-button behavior without truncating labels at normal dropdown widths. - Each dropdown row has one status-tooltip surface: a `Running` row shows its status or recovered-OOM warning, while every non-running row appends "Cannot select." without duplicating a trailing period. - Row-tooltip composition is now the pure `getComputingUnitRowTooltip` utility instead of three component methods, and its punctuation/status matrix is tested in the utility spec. - The unreferenced `.unit-disconnected` and `.unit-terminating` selectors were removed; the live `.unit-connecting` selector remains bound to `Pending` rows. - `statusReason` is typed `string | null` on the DTO, which is what `Option[String]` through `DefaultScalaModule` actually puts on the wire; the previous optional-only type could not represent it. #### Refusing work a unit cannot do Naming a terminal state is not enough on its own, because three of the four run entry points never looked at the unit status: run-up-to from the operator menu, the Time Travel replay, and the Form View run button all call `ExecuteWorkflowService` directly, and only the canvas button was gated, by its disabled state alone. On a terminal unit the old flow reset the execution state, clearing the results already on screen, and then handed the request to a WebSocket that has no delivery check and no error path, so the user silently lost their previous results and no run started. - `ExecuteWorkflowService` gained a guard modelled on the existing warehouse refusal: it is checked at both public entry points, before the warehouse guard and before `resetExecutionState()`, so a refused click leaves the screen as it was. The unit is checked first because picking a warehouse cannot rescue a dead unit. - The guard reads the same shared utility the two run buttons use, so it never refuses a click that either button allows. - In the dropdown, a row the list marks unselectable is no longer selectable by clicking it. ng-zorro greys such a row out and calls `preventDefault`/`stopPropagation`, but neither cancels the `(click)` binding on the same `<li>`, so a click still selected the unit. Since the per-workflow selection is now remembered, that stray selection was also written to storage and restored on every later load, because the recall path only checks that the remembered `cuid` still exists, never its status. The guard lives in a row-only handler rather than in the shared selection method, which the creation flow also calls with a freshly created and therefore `Pending` unit. ### Any related issues, documentation, discussions? Closes #7669. The status vocabulary, reason visibility, and no-auto-recovery policy were discussed in #7670; administrator reason visibility and the review cleanups address feedback on this PR. The run-entry-point guard and the Form View run button address the Copilot review on this PR, which observed that the new states reached only `MenuComponent`. ### How was this PR tested? - TDD covers every backend decision-table row with `PodBuilder`-built pods through the pure snapshot and mapping functions, including both eviction wordings, all three image-pull reasons, crash loops with and without OOM history, crash-loop precedence over recovered OOM, multi-container pods, authorization gating, and the absent-pod path used by creation polling; fabric8 null guards for `status`, `conditions`, `containerStatuses`, `state`, `lastState`, and `terminated` have dedicated coverage. - Additional edge cases cover Terminating precedence over eviction and status-less pods, eviction-message truncation/whitespace/case-insensitivity, mixed multi-container pods, stale/malformed `Unschedulable` conditions, image-pull precedence over crash loops, empty `statusReason` fallback behavior, and the pinned `Succeeded` → `Pending` behavior. - The refusal guard is pinned per entry point: a terminating, `Failed` or `Unknown` unit is refused, nothing reaches the WebSocket, and the execution state is **not** reset, so a test fails if the guard is ever moved after the reset. `Running` and `Pending` units still run, and a workflow with no unit selected at all still runs and keeps the existing warning. - The Form View button is pinned for both terminal labels with the socket up and down, for invalid and empty workflows keeping precedence over a terminal unit, for a terminal unit being named ahead of a missing warehouse, for `Pending` still saying Connecting, and for Stop being withheld when it could not be delivered. - The dropdown row is pinned for a refused click on a `Failed`, `Unknown`, `Terminating` and `Pending` row, and for a `Running` row still being selected and remembered; creating a unit, which arrives `Pending`, is pinned as still being selected. - Mutation checks confirmed that swapping precedence, mapping Evicted to Pending, removing wording branches, hardcoding restart counts, bypassing reason authorization, mapping absent pods to Unknown, or removing run-button terminal-state branches causes tests to fail. A further twenty mutations on the shared utility, the refusal guard, the Form View branch, the Stop condition and the row handler were each killed by at least one test, with no survivors. - Backend: `ComputingUnitManagingService / Test` passes 238/238. - Frontend: the five affected specs pass 626/626 with `NODE_OPTIONS="--no-experimental-webstorage" npx ng test --watch=false --include src/app/common/util/computing-unit.util.spec.ts --include src/app/workspace/service/execute-workflow/execute-workflow.service.spec.ts --include src/app/workspace/component/workflow-form/workflow-form.component.spec.ts --include src/app/workspace/component/menu/menu.component.spec.ts --include src/app/workspace/component/power-button/computing-unit-selection.component.spec.ts`. The whole frontend suite is green under the same flag; without it, local Node 25 shadows jsdom's `localStorage` and the suite is both noisy and non-deterministic, which CI does not hit because it pins Node 24. - Type checking passed with `tsc --noEmit`, Scala and frontend formatting checks passed, scoped frontend ESLint passed, and `git diff --check` passed. - Manual verification in Chrome with response overrides on both computing-unit polling endpoints: switching a unit's status between `Running`, `Pending`, `Terminating`, `Failed` and `Unknown` produces the badge colour, row tooltip and run-button state described above. Click path: open a workflow in the workspace, open the computing-unit picker in the toolbar, override the array response for the dropdown rows and the single-unit response for the run button, then observe both within one polling interval. `Option[String] → null` serialization relies on `DefaultScalaModule` registered on the service's Dropwizard `ObjectMapper`; there is no end-to-end JSON serialization test for the field, though the DTO type now admits `null` and both consumers treat it as absent. A bare pod in phase `Succeeded` still maps to `Pending`, which is pinned by a test and practically unreachable under `restartPolicy: Always`. Follow-ups intentionally outside this PR's scope remain a frontend timeout for the "pod Running but engine unreachable" zombie case, UX for silent vanish reconciliation, and `local` computing-unit liveness. ### Was this PR authored or co-authored using generative AI tooling? Co-authored by: Claude Code (Claude Fable 5) --- .../resource/AdminComputingUnitResource.scala | 7 +- .../resource/ComputingUnitManagingResource.scala | 31 +- .../service/resource/ComputingUnitState.scala | 7 +- .../texera/service/util/ComputingUnitHelpers.scala | 200 +++++-- .../texera/service/util/KubernetesClient.scala | 16 +- .../texera/service/util/PodStatusSnapshot.scala | 105 ++++ .../service/util/ComputingUnitHelpersSpec.scala | 628 +++++++++++++++++++-- .../texera/service/util/KubernetesClientSpec.scala | 143 ++++- .../computing-unit-status.service.spec.ts | 35 +- .../computing-unit-status.service.ts | 6 + .../type/computing-unit-connection.interface.ts | 3 + .../src/app/common/type/workflow-computing-unit.ts | 4 +- .../app/common/util/computing-unit.util.spec.ts | 125 ++++ .../src/app/common/util/computing-unit.util.ts | 46 ++ .../component/menu/menu.component.spec.ts | 64 ++- .../app/workspace/component/menu/menu.component.ts | 21 + .../computing-unit-selection.component.html | 18 +- .../computing-unit-selection.component.scss | 23 +- .../computing-unit-selection.component.spec.ts | 53 +- .../computing-unit-selection.component.ts | 25 +- .../workflow-form/workflow-form.component.spec.ts | 109 ++++ .../workflow-form/workflow-form.component.ts | 24 +- .../execute-workflow.service.spec.ts | 104 ++++ .../execute-workflow/execute-workflow.service.ts | 30 + 24 files changed, 1640 insertions(+), 187 deletions(-) diff --git a/computing-unit-managing-service/src/main/scala/org/apache/texera/service/resource/AdminComputingUnitResource.scala b/computing-unit-managing-service/src/main/scala/org/apache/texera/service/resource/AdminComputingUnitResource.scala index 28c95999b2..799c58b270 100644 --- a/computing-unit-managing-service/src/main/scala/org/apache/texera/service/resource/AdminComputingUnitResource.scala +++ b/computing-unit-managing-service/src/main/scala/org/apache/texera/service/resource/AdminComputingUnitResource.scala @@ -73,7 +73,7 @@ class AdminComputingUnitResource { .asScala .toList - val podPhases = ComputingUnitHelpers.podPhasesFor(activeUnits) + val podSnapshots = ComputingUnitHelpers.podSnapshotsFor(activeUnits) // Wrap only the reconcile write, not the Kubernetes round trips above: keeps a pooled // connection from being held open across them, while the transaction makes the batchUpdate @@ -83,7 +83,7 @@ class AdminComputingUnitResource { ComputingUnitHelpers.reconcileVanishedKubernetesUnits( new WorkflowComputingUnitDao(txCtx.configuration()), activeUnits, - podPhases + podSnapshots ) } @@ -97,9 +97,10 @@ class AdminComputingUnitResource { ComputingUnitHelpers.buildDashboardUnit( unit, isOwner = unit.getUid.equals(user.getUid), + canViewStatusReason = true, accessPrivilege = PrivilegeEnum.WRITE, ownerInfo = ownerInfo, - podPhases = podPhases, + podSnapshots = podSnapshots, podMetrics = podMetrics ) } diff --git a/computing-unit-managing-service/src/main/scala/org/apache/texera/service/resource/ComputingUnitManagingResource.scala b/computing-unit-managing-service/src/main/scala/org/apache/texera/service/resource/ComputingUnitManagingResource.scala index 2b6b55a667..a9326e14a4 100644 --- a/computing-unit-managing-service/src/main/scala/org/apache/texera/service/resource/ComputingUnitManagingResource.scala +++ b/computing-unit-managing-service/src/main/scala/org/apache/texera/service/resource/ComputingUnitManagingResource.scala @@ -439,6 +439,9 @@ object ComputingUnitManagingResource { case class DashboardWorkflowComputingUnit( computingUnit: WorkflowComputingUnit, status: String, + // User-friendly explanation of a failing/degraded status; serialized as null when there is + // nothing to explain or the endpoint does not authorize the caller to view it. + statusReason: Option[String], metrics: WorkflowComputingUnitMetrics, isOwner: Boolean, accessPrivilege: EnumType, @@ -744,9 +747,13 @@ class ComputingUnitManagingResource { } } + // The creator is always the owner, so the status reason is never withheld here. + val (status, statusReason) = + ComputingUnitHelpers.getComputingUnitStatusWithReason(insertedUnit) DashboardWorkflowComputingUnit( insertedUnit, - ComputingUnitHelpers.getComputingUnitStatus(insertedUnit).toString, + status.toString, + statusReason, ComputingUnitHelpers.getComputingUnitMetrics(insertedUnit), isOwner = true, accessPrivilege = PrivilegeEnum.WRITE, @@ -812,14 +819,14 @@ class ComputingUnitManagingResource { }.toMap val candidateUnits = unitsWithPrivilege.map { case (unit, _) => unit } - // Pod phases decide which Kubernetes units are still alive. - val podPhases = ComputingUnitHelpers.podPhasesFor(candidateUnits) + // Pod snapshots decide which Kubernetes units are still alive (by pod-name presence). + val podSnapshots = ComputingUnitHelpers.podSnapshotsFor(candidateUnits) val liveUnits = ComputingUnitHelpers.reconcileVanishedKubernetesUnits( computingUnitDao, candidateUnits, - podPhases + podSnapshots ) // Metrics only for survivors, so fetch after reconciliation. @@ -829,12 +836,14 @@ class ComputingUnitManagingResource { ComputingUnitHelpers.resolveOwnerInfo(userDao, liveUnits.map(_.getUid).distinct) liveUnits.map { unit => + val isOwner = unit.getUid.equals(uid) ComputingUnitHelpers.buildDashboardUnit( unit, - isOwner = unit.getUid.equals(uid), + isOwner = isOwner, + canViewStatusReason = isOwner, accessPrivilege = privilegeByCuid(unit.getCuid), ownerInfo = ownerInfoMap, - podPhases = podPhases, + podSnapshots = podSnapshots, podMetrics = podMetrics ) } @@ -864,11 +873,17 @@ class ComputingUnitManagingResource { val ownerUsername: String = ownerUser.flatMap(u => Option(u.getName).filter(_.nonEmpty)).orNull + val isOwner = unit.getUid.equals(user.getUid) + val (status, statusReason) = ComputingUnitHelpers.getComputingUnitStatusWithReason(unit) + DashboardWorkflowComputingUnit( computingUnit = unit, - status = ComputingUnitHelpers.getComputingUnitStatus(unit).toString, + status = status.toString, + // The direct and regular-user listing endpoints remain owner-gated; the separate admin + // listing authorizes administrators to view reasons without marking them as owners. + statusReason = if (isOwner) statusReason else None, metrics = ComputingUnitHelpers.getComputingUnitMetrics(unit), - isOwner = unit.getUid.equals(user.getUid), + isOwner = isOwner, accessPrivilege = { val cuAccessDao = new ComputingUnitUserAccessDao(context.configuration()) val access = cuAccessDao diff --git a/computing-unit-managing-service/src/main/scala/org/apache/texera/service/resource/ComputingUnitState.scala b/computing-unit-managing-service/src/main/scala/org/apache/texera/service/resource/ComputingUnitState.scala index e6ff217cfb..4ea3b2e10a 100644 --- a/computing-unit-managing-service/src/main/scala/org/apache/texera/service/resource/ComputingUnitState.scala +++ b/computing-unit-managing-service/src/main/scala/org/apache/texera/service/resource/ComputingUnitState.scala @@ -18,7 +18,12 @@ package org.apache.texera.service.resource +/** + * Status vocabulary of a computing unit, mirroring the Kubernetes pod lifecycle: `Failed`, + * `Unknown` and `Terminating` map from the pod's phase / deletion timestamp so a dead unit is + * reported as dead instead of forever `Pending`. + */ object ComputingUnitState extends Enumeration { type ComputingUnitState = Value - val Running, Pending = Value + val Running, Pending, Failed, Unknown, Terminating = Value } diff --git a/computing-unit-managing-service/src/main/scala/org/apache/texera/service/util/ComputingUnitHelpers.scala b/computing-unit-managing-service/src/main/scala/org/apache/texera/service/util/ComputingUnitHelpers.scala index f2e6781cc8..e0506b52fb 100644 --- a/computing-unit-managing-service/src/main/scala/org/apache/texera/service/util/ComputingUnitHelpers.scala +++ b/computing-unit-managing-service/src/main/scala/org/apache/texera/service/util/ComputingUnitHelpers.scala @@ -28,7 +28,14 @@ import org.apache.texera.service.resource.ComputingUnitManagingResource.{ DashboardWorkflowComputingUnit, WorkflowComputingUnitMetrics } -import org.apache.texera.service.resource.ComputingUnitState.{ComputingUnitState, Pending, Running} +import org.apache.texera.service.resource.ComputingUnitState.{ + ComputingUnitState, + Failed, + Pending, + Running, + Terminating, + Unknown +} import org.jooq.EnumType import java.sql.Timestamp @@ -57,43 +64,136 @@ object ComputingUnitHelpers { .toMap } - def getComputingUnitStatus(unit: WorkflowComputingUnit): ComputingUnitState = - singleUnitStatus(unit, KubernetesClient) + /** Single-unit status plus its user-facing reason (see [[kubernetesStatusAndReason]]). */ + def getComputingUnitStatusWithReason( + unit: WorkflowComputingUnit + ): (ComputingUnitState, Option[String]) = + singleUnitStatusAndReason(unit, KubernetesClient) /** * Single-unit status via a per-unit pod lookup (a targeted GET, cheaper than listing the whole * namespace). The client is a by-name parameter — not the global singleton — so the kubernetes * branch is unit-testable with a stub and the local/unknown branches never force the singleton; - * the public overload binds the production [[KubernetesClient]]. (Metrics has no analogous seam: + * the public overloads bind the production [[KubernetesClient]]. (Metrics has no analogous seam: * its per-unit lookup already fans out to the whole namespace and the bulk (unit, podMetrics) * overload already covers the cpu/memory resolution, so nothing there is worth pinning.) */ - private[util] def singleUnitStatus( + private[util] def singleUnitStatusAndReason( unit: WorkflowComputingUnit, k8s: => KubernetesClient - ): ComputingUnitState = { + ): (ComputingUnitState, Option[String]) = { unit.getType match { // Local CUs are always “running” case WorkflowComputingUnitTypeEnum.local => - Running + (Running, None) - // Kubernetes CUs – only explicit “Running” counts as running + // Kubernetes CUs – resolved from the pod's status snapshot case WorkflowComputingUnitTypeEnum.kubernetes => - // Guard the pod status the same way the bulk getAllPodPhases does: a pod with no - // status yet has a null getStatus, so map through Option to avoid an NPE. val client = k8s - val phaseOpt = client - .getPodByName(client.generatePodName(unit.getCuid)) - .flatMap(pod => Option(pod.getStatus).map(_.getPhase)) - - if (phaseOpt.contains("Running")) Running else Pending + kubernetesStatusAndReason( + client + .getPodByName(client.generatePodName(unit.getCuid)) + .map(PodStatusSnapshot.fromPod) + ) // Any other (unknown) type is treated as pending case _ => - Pending + (Pending, None) + } + } + + // User-facing wording for each failure mode. Deliberately actionable prose, never a raw + // Kubernetes dump; buildDashboardUnit withholds these unless the caller permits access. + private val ImagePullWaitingReasons = Set("ImagePullBackOff", "ErrImagePull", "InvalidImageName") + private val EvictedDiskReason = + "The computing unit was evicted because it ran out of local disk storage. Consider " + + "storing less data on the unit's local file system, or recreate it with more storage." + private val ImagePullReason = + "The computing unit's image could not be pulled. Please recreate the unit or contact " + + "an administrator." + private val CrashLoopOomReason = + "The computing unit keeps crashing because it runs out of memory. Please terminate it " + + "and recreate it with a higher memory limit." + private val GenericFailedReason = + "The computing unit stopped unexpectedly. Please terminate and recreate it, or contact " + + "an administrator." + private val UnknownStateReason = + "The state of the computing unit cannot be determined (its node may be unreachable)." + private val UnschedulableReason = + "The computing unit is waiting for cluster resources to become available." + + private def evictedReason(podMessage: Option[String]): String = { + val mentionsDisk = + podMessage.exists { message => + val lower = message.toLowerCase + lower.contains("ephemeral") || lower.contains("disk") + } + if (mentionsDisk) EvictedDiskReason + else { + // First sentence of the cluster's message, capped so the tooltip stays readable. + val shortReason = podMessage + .map(_.takeWhile(_ != '.').trim) + .filter(_.nonEmpty) + .map(sentence => if (sentence.length > 120) sentence.take(120).trim + "..." else sentence) + .getOrElse("Evicted") + s"The computing unit was evicted by the cluster ($shortReason). Consider recreating it." } } + private def crashLoopReason(restartCount: Int): String = + s"The computing unit is repeatedly crashing (restarted $restartCount times). Please " + + "terminate and recreate it, or contact an administrator." + + private def recoveredOomWarning(restartCount: Int): String = + s"The last run was terminated because the computing unit ran out of memory (restarted " + + s"$restartCount times). Consider recreating the unit with a higher memory limit before " + + "running the same workload." + + /** + * Pure (snapshot -> state, reason) mapping, mirroring the Kubernetes pod lifecycle. An absent + * pod stays Pending — exactly today's behavior — because the vanish reconciliation, not this + * mapping, is what retires units whose pods are gone. + * + * Note the restartPolicy-Always subtlety: an OOM-killed container restarts in place with the + * pod phase still "Running", so OOM kills and crash loops are read from the container-level + * fields, and a waiting-state failure takes precedence over the recovered-OOM warning. + */ + private[util] def kubernetesStatusAndReason( + snapshotOpt: Option[PodStatusSnapshot] + ): (ComputingUnitState, Option[String]) = + snapshotOpt match { + case None => (Pending, None) + case Some(snapshot) => + val phase = snapshot.phase.getOrElse("") + val imagePullFailed = + snapshot.containers.exists(_.waitingReason.exists(ImagePullWaitingReasons.contains)) + val crashLooping = snapshot.containers.find(_.waitingReason.contains("CrashLoopBackOff")) + val oomKilled = snapshot.containers.find(_.lastTerminatedReason.contains("OOMKilled")) + + if (snapshot.terminating) + (Terminating, None) + else if (phase == "Failed" && snapshot.podReason.contains("Evicted")) + (Failed, Some(evictedReason(snapshot.podMessage))) + else if (imagePullFailed) + (Failed, Some(ImagePullReason)) + else if (crashLooping.isDefined) { + val container = crashLooping.get + if (container.lastTerminatedReason.contains("OOMKilled")) + (Failed, Some(CrashLoopOomReason)) + else + (Failed, Some(crashLoopReason(container.restartCount))) + } else if (phase == "Failed") + (Failed, Some(GenericFailedReason)) + else if (phase == "Unknown") + (Unknown, Some(UnknownStateReason)) + else if (phase == "Pending" && snapshot.unschedulable) + (Pending, Some(UnschedulableReason)) + else if (phase == "Running") + (Running, oomKilled.map(container => recoveredOomWarning(container.restartCount))) + else + (Pending, None) + } + def getComputingUnitMetrics(unit: WorkflowComputingUnit): WorkflowComputingUnitMetrics = { unit.getType match { case WorkflowComputingUnitTypeEnum.local => @@ -110,23 +210,20 @@ object ComputingUnitHelpers { } /** - * Resolves status from a pre-fetched pod-phase map instead of a per-unit cluster call, so a - * listing costs O(1) round trips rather than one per unit. + * Resolves status (and its user-facing reason) from a pre-fetched pod-snapshot map instead + * of a per-unit cluster call, so a listing costs O(1) round trips rather than one per unit. */ - def getComputingUnitStatus( + def getComputingUnitStatusAndReason( unit: WorkflowComputingUnit, - podPhases: Map[String, String] - ): ComputingUnitState = { + podSnapshots: Map[String, PodStatusSnapshot] + ): (ComputingUnitState, Option[String]) = { unit.getType match { case WorkflowComputingUnitTypeEnum.local => - Running + (Running, None) case WorkflowComputingUnitTypeEnum.kubernetes => - // A missing entry or null phase both count as not-Running. - if (podPhases.get(KubernetesClient.generatePodName(unit.getCuid)).contains("Running")) - Running - else Pending + kubernetesStatusAndReason(podSnapshots.get(KubernetesClient.generatePodName(unit.getCuid))) case _ => - Pending + (Pending, None) } } @@ -156,19 +253,19 @@ object ComputingUnitHelpers { case _ => false } - // Pod phases/metrics for the namespace; skipped (empty) when no Kubernetes unit is present, so a - // cluster-free listing issues no round trip. Same seam as singleUnitStatus: the public overload - // binds the production singleton, passing it by-name to the private[util] overload, which forces - // it only inside the guard's true branch — so the empty path never touches the client, and tests - // drive the private overload with a stub. - def podPhasesFor(units: Seq[WorkflowComputingUnit]): Map[String, String] = - podPhasesFor(units, KubernetesClient) + // Pod snapshots/metrics for the namespace; skipped (empty) when no Kubernetes unit is present, + // so a cluster-free listing issues no round trip. Same seam as singleUnitStatusAndReason: the + // public overload binds the production singleton, passing it by-name to the private[util] + // overload, which forces it only inside the guard's true branch — so the empty path never + // touches the client, and tests drive the private overload with a stub. + def podSnapshotsFor(units: Seq[WorkflowComputingUnit]): Map[String, PodStatusSnapshot] = + podSnapshotsFor(units, KubernetesClient) - private[util] def podPhasesFor( + private[util] def podSnapshotsFor( units: Seq[WorkflowComputingUnit], k8s: => KubernetesClient - ): Map[String, String] = - if (units.exists(isKubernetes)) k8s.getAllPodPhases else Map.empty + ): Map[String, PodStatusSnapshot] = + if (units.exists(isKubernetes)) k8s.getAllPodStatusSnapshots else Map.empty def podMetricsFor(units: Seq[WorkflowComputingUnit]): Map[String, Map[String, String]] = podMetricsFor(units, KubernetesClient) @@ -179,16 +276,19 @@ object ComputingUnitHelpers { ): Map[String, Map[String, String]] = if (units.exists(isKubernetes)) k8s.getAllPodMetrics else Map.empty - /** A Kubernetes unit whose pod is absent from `podPhases` (deleted or TTL GC-ed). */ - private def isVanished(unit: WorkflowComputingUnit, podPhases: Map[String, String]): Boolean = - isKubernetes(unit) && !podPhases.contains(KubernetesClient.generatePodName(unit.getCuid)) + /** A Kubernetes unit whose pod is absent from `podSnapshots` (deleted or TTL GC-ed). */ + private def isVanished( + unit: WorkflowComputingUnit, + podSnapshots: Map[String, PodStatusSnapshot] + ): Boolean = + isKubernetes(unit) && !podSnapshots.contains(KubernetesClient.generatePodName(unit.getCuid)) - /** Partition into `(live, vanished)` by `podPhases`. Pure (no I/O), so it is unit-testable. */ + /** Partition into `(live, vanished)` by `podSnapshots`. Pure (no I/O), so it is unit-testable. */ def partitionLiveUnits( units: List[WorkflowComputingUnit], - podPhases: Map[String, String] + podSnapshots: Map[String, PodStatusSnapshot] ): (List[WorkflowComputingUnit], List[WorkflowComputingUnit]) = - units.partition(unit => !isVanished(unit, podPhases)) + units.partition(unit => !isVanished(unit, podSnapshots)) /** * Stamp `terminateTime` on vanished Kubernetes units (one batched update) and return the live @@ -197,9 +297,9 @@ object ComputingUnitHelpers { def reconcileVanishedKubernetesUnits( dao: WorkflowComputingUnitDao, units: List[WorkflowComputingUnit], - podPhases: Map[String, String] + podSnapshots: Map[String, PodStatusSnapshot] ): List[WorkflowComputingUnit] = { - val partitioned = partitionLiveUnits(units, podPhases) + val partitioned = partitionLiveUnits(units, podSnapshots) val vanished = partitioned._2 if (vanished.nonEmpty) { val now = new Timestamp(System.currentTimeMillis()) @@ -215,19 +315,25 @@ object ComputingUnitHelpers { /** * Build one dashboard row; status/metrics come from the pre-fetched maps (no per-unit K8s * call). Shared by both listing endpoints so row shape and owner-info fallback stay identical. + * Status-reason visibility is separate from ownership: regular-user callers permit owners, + * while the admin listing permits every row. Shared non-admin users get a bare status and the + * frontend shows a generic "unavailable" instead. */ def buildDashboardUnit( unit: WorkflowComputingUnit, isOwner: Boolean, + canViewStatusReason: Boolean, accessPrivilege: EnumType, ownerInfo: Map[Integer, (String, String)], - podPhases: Map[String, String], + podSnapshots: Map[String, PodStatusSnapshot], podMetrics: Map[String, Map[String, String]] ): DashboardWorkflowComputingUnit = { val owner = ownerInfo.getOrElse(unit.getUid, (null, null)) + val (status, statusReason) = getComputingUnitStatusAndReason(unit, podSnapshots) DashboardWorkflowComputingUnit( computingUnit = unit, - status = getComputingUnitStatus(unit, podPhases).toString, + status = status.toString, + statusReason = if (canViewStatusReason) statusReason else None, metrics = getComputingUnitMetrics(unit, podMetrics), isOwner = isOwner, accessPrivilege = accessPrivilege, diff --git a/computing-unit-managing-service/src/main/scala/org/apache/texera/service/util/KubernetesClient.scala b/computing-unit-managing-service/src/main/scala/org/apache/texera/service/util/KubernetesClient.scala index 76f2abc7bd..704513159d 100644 --- a/computing-unit-managing-service/src/main/scala/org/apache/texera/service/util/KubernetesClient.scala +++ b/computing-unit-managing-service/src/main/scala/org/apache/texera/service/util/KubernetesClient.scala @@ -62,17 +62,17 @@ class KubernetesClient( } /** - * Phase of every pod in the namespace, keyed by pod name, in one call — so a bulk listing - * avoids a per-unit lookup. Unfiltered so callers can test a unit's presence by its pod-name - * key; a pod with no status yet maps to a `null` phase but still appears. + * Status snapshot of every pod in the namespace, keyed by pod name, in one call — so a bulk + * listing avoids a per-unit lookup. Unfiltered so callers can test a unit's presence by its + * pod-name key; a pod with no status yet maps to an empty snapshot but still appears. */ - def getAllPodPhases: Map[String, String] = - phasesByPodName(client.pods().inNamespace(namespace).list().getItems.asScala) + def getAllPodStatusSnapshots: Map[String, PodStatusSnapshot] = + snapshotsByPodName(client.pods().inNamespace(namespace).list().getItems.asScala) - /** Pure fabric8 -> map transform: a pod with no status yet maps to a `null` phase. */ - private[util] def phasesByPodName(pods: Iterable[Pod]): Map[String, String] = + /** Pure fabric8 -> map transform; the per-pod extraction is [[PodStatusSnapshot.fromPod]]. */ + private[util] def snapshotsByPodName(pods: Iterable[Pod]): Map[String, PodStatusSnapshot] = pods - .map(pod => pod.getMetadata.getName -> Option(pod.getStatus).map(_.getPhase).orNull) + .map(pod => pod.getMetadata.getName -> PodStatusSnapshot.fromPod(pod)) .toMap // Flatten a pod's per-container resource usage into a single metric -> value map. diff --git a/computing-unit-managing-service/src/main/scala/org/apache/texera/service/util/PodStatusSnapshot.scala b/computing-unit-managing-service/src/main/scala/org/apache/texera/service/util/PodStatusSnapshot.scala new file mode 100644 index 0000000000..1d4e3066be --- /dev/null +++ b/computing-unit-managing-service/src/main/scala/org/apache/texera/service/util/PodStatusSnapshot.scala @@ -0,0 +1,105 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package org.apache.texera.service.util + +import io.fabric8.kubernetes.api.model.Pod + +import scala.jdk.CollectionConverters._ + +/** The slice of one container's status that the state mapping needs. */ +case class ContainerStatusSnapshot( + waitingReason: Option[String], + lastTerminatedReason: Option[String], + restartCount: Int +) + +/** + * The slice of a pod's status that computing-unit state resolution needs, extracted once per + * pod so the bulk listing keeps its single namespace-wide call. Kept as a plain case class + * (not the fabric8 Pod) so the (snapshot -> state) mapping in [[ComputingUnitHelpers]] is a + * pure, builder-testable function and so vanish reconciliation can keep testing key presence + * on the snapshot map exactly as it did on the phase map. + * + * Container-level fields matter because our pods run with restartPolicy Always: an OOM-killed + * container restarts in place and the pod phase stays "Running", so OOM kills and crash loops + * are only visible through `containerStatuses[].lastState/state`. + */ +case class PodStatusSnapshot( + phase: Option[String], + terminating: Boolean, + podReason: Option[String], + podMessage: Option[String], + unschedulable: Boolean, + containers: Seq[ContainerStatusSnapshot] +) + +object PodStatusSnapshot { + + /** The snapshot of a pod whose status has not been populated yet (getStatus == null). */ + val empty: PodStatusSnapshot = + PodStatusSnapshot( + phase = None, + terminating = false, + podReason = None, + podMessage = None, + unschedulable = false, + containers = Seq.empty + ) + + /** Pure fabric8 -> snapshot extraction; every nested status object is null-guarded. */ + def fromPod(pod: Pod): PodStatusSnapshot = { + val terminating = + Option(pod.getMetadata).flatMap(m => Option(m.getDeletionTimestamp)).isDefined + Option(pod.getStatus) match { + case None => empty.copy(terminating = terminating) + case Some(status) => + val unschedulable = Option(status.getConditions) + .map(_.asScala) + .getOrElse(Seq.empty) + .exists(condition => + condition.getType == "PodScheduled" && + condition.getStatus == "False" && + condition.getReason == "Unschedulable" + ) + val containers = Option(status.getContainerStatuses) + .map(_.asScala.toSeq) + .getOrElse(Seq.empty) + .map(containerStatus => + ContainerStatusSnapshot( + waitingReason = Option(containerStatus.getState) + .flatMap(state => Option(state.getWaiting)) + .flatMap(waiting => Option(waiting.getReason)), + lastTerminatedReason = Option(containerStatus.getLastState) + .flatMap(state => Option(state.getTerminated)) + .flatMap(terminated => Option(terminated.getReason)), + restartCount = Option(containerStatus.getRestartCount).map(_.intValue()).getOrElse(0) + ) + ) + PodStatusSnapshot( + phase = Option(status.getPhase), + terminating = terminating, + podReason = Option(status.getReason), + podMessage = Option(status.getMessage), + unschedulable = unschedulable, + containers = containers + ) + } + } +} diff --git a/computing-unit-managing-service/src/test/scala/org/apache/texera/service/util/ComputingUnitHelpersSpec.scala b/computing-unit-managing-service/src/test/scala/org/apache/texera/service/util/ComputingUnitHelpersSpec.scala index 18554cd1fe..2b7eb00957 100644 --- a/computing-unit-managing-service/src/test/scala/org/apache/texera/service/util/ComputingUnitHelpersSpec.scala +++ b/computing-unit-managing-service/src/test/scala/org/apache/texera/service/util/ComputingUnitHelpersSpec.scala @@ -30,7 +30,13 @@ import org.apache.texera.dao.jooq.generated.enums.{ import org.apache.texera.dao.jooq.generated.tables.daos.{UserDao, WorkflowComputingUnitDao} import org.apache.texera.dao.jooq.generated.tables.pojos.{User, WorkflowComputingUnit} import org.apache.texera.service.resource.ComputingUnitManagingResource.WorkflowComputingUnitMetrics -import org.apache.texera.service.resource.ComputingUnitState.{Pending, Running} +import org.apache.texera.service.resource.ComputingUnitState.{ + Failed, + Pending, + Running, + Terminating, + Unknown +} import org.mockito.Mockito.{mock, when} import org.scalatest.BeforeAndAfterAll import org.scalatest.flatspec.AnyFlatSpec @@ -91,35 +97,133 @@ class ComputingUnitHelpersSpec // A pod whose status has not been populated yet (getStatus == null). private def statuslessPod(): Pod = new PodBuilder().build() - "getComputingUnitStatus" should "return Running for a local unit" in { - ComputingUnitHelpers.getComputingUnitStatus(localUnit()) shouldBe Running - } - - it should "return Pending for an unknown (untyped) unit" in { - ComputingUnitHelpers.getComputingUnitStatus(untypedUnit()) shouldBe Pending + // ── PodBuilder recipes for the failure-state decision table ───────── + // Each builds the minimal fabric8 Pod that a real cluster would report for + // the scenario, then goes through the same pure fromPod transform production + // uses, so these tests cover the snapshot extraction and the state mapping + // together. + + private def terminatingPod(): Pod = + new PodBuilder() + .withNewMetadata() + .withDeletionTimestamp("2026-08-20T00:00:00Z") + .endMetadata() + .withNewStatus() + .withPhase("Running") + .endStatus() + .build() + + private def evictedPod(message: String): Pod = + new PodBuilder() + .withNewStatus() + .withPhase("Failed") + .withReason("Evicted") + .withMessage(message) + .endStatus() + .build() + + /** A pod whose (single) container is stuck in the given waiting reason. */ + private def waitingContainerPod( + phase: String, + waitingReason: String, + lastTerminatedReason: Option[String] = None, + restartCount: Int = 0 + ): Pod = { + val builder = new PodBuilder() + .withNewStatus() + .withPhase(phase) + .addNewContainerStatus() + .withRestartCount(restartCount) + .withNewState() + .withNewWaiting() + .withReason(waitingReason) + .endWaiting() + .endState() + lastTerminatedReason + .fold(builder)(reason => + builder + .withNewLastState() + .withNewTerminated() + .withReason(reason) + .endTerminated() + .endLastState() + ) + .endContainerStatus() + .endStatus() + .build() } - // The kubernetes branch does a per-unit pod lookup, so singleUnitStatus is driven through a - // stubbed client (the public getComputingUnitStatus binds the production singleton). - "singleUnitStatus" should "return Running for a kubernetes unit whose pod phase is Running" in { + /** A running pod whose container previously terminated with the given reason. */ + private def restartedPod(lastTerminatedReason: String, restartCount: Int): Pod = + new PodBuilder() + .withNewStatus() + .withPhase("Running") + .addNewContainerStatus() + .withRestartCount(restartCount) + .withNewLastState() + .withNewTerminated() + .withReason(lastTerminatedReason) + .endTerminated() + .endLastState() + .endContainerStatus() + .endStatus() + .build() + + private def unschedulablePod(): Pod = + new PodBuilder() + .withNewStatus() + .withPhase("Pending") + .addNewCondition() + .withType("PodScheduled") + .withStatus("False") + .withReason("Unschedulable") + .endCondition() + .endStatus() + .build() + + private def snapshotOf(pod: Pod): PodStatusSnapshot = PodStatusSnapshot.fromPod(pod) + + private def statusAndReasonOf( + pod: Pod + ): (org.apache.texera.service.resource.ComputingUnitState.ComputingUnitState, Option[String]) = + ComputingUnitHelpers.kubernetesStatusAndReason(Some(snapshotOf(pod))) + + // The kubernetes branch does a per-unit pod lookup, so singleUnitStatusAndReason is driven + // through a stubbed client (the public overloads bind the production singleton). + "singleUnitStatusAndReason" should "return Running for a kubernetes unit whose pod phase is Running" in { val k8s = mock(classOf[KubernetesClient]) when(k8s.generatePodName(40)).thenReturn("computing-unit-40") when(k8s.getPodByName("computing-unit-40")).thenReturn(Some(podWithPhase("Running"))) - ComputingUnitHelpers.singleUnitStatus(kubernetesUnit(40), k8s) shouldBe Running + ComputingUnitHelpers.singleUnitStatusAndReason(kubernetesUnit(40), k8s) shouldBe + ((Running, None)) } it should "return Pending for a kubernetes unit whose pod has no status yet" in { val k8s = mock(classOf[KubernetesClient]) when(k8s.generatePodName(41)).thenReturn("computing-unit-41") when(k8s.getPodByName("computing-unit-41")).thenReturn(Some(statuslessPod())) - ComputingUnitHelpers.singleUnitStatus(kubernetesUnit(41), k8s) shouldBe Pending + ComputingUnitHelpers.singleUnitStatusAndReason(kubernetesUnit(41), k8s) shouldBe + ((Pending, None)) } it should "return Pending for a kubernetes unit whose pod is absent" in { val k8s = mock(classOf[KubernetesClient]) when(k8s.generatePodName(42)).thenReturn("computing-unit-42") when(k8s.getPodByName("computing-unit-42")).thenReturn(None) - ComputingUnitHelpers.singleUnitStatus(kubernetesUnit(42), k8s) shouldBe Pending + ComputingUnitHelpers.singleUnitStatusAndReason(kubernetesUnit(42), k8s) shouldBe + ((Pending, None)) + } + + it should "surface the failure state and reason of a dead pod" in { + val k8s = mock(classOf[KubernetesClient]) + when(k8s.generatePodName(43)).thenReturn("computing-unit-43") + when(k8s.getPodByName("computing-unit-43")) + .thenReturn(Some(waitingContainerPod("Pending", "ImagePullBackOff"))) + val (state, reason) = ComputingUnitHelpers.singleUnitStatusAndReason(kubernetesUnit(43), k8s) + state shouldBe Failed + reason shouldBe Some( + "The computing unit's image could not be pulled. Please recreate the unit or contact an administrator." + ) } "getComputingUnitMetrics" should "return NaN metrics for a local unit" in { @@ -134,29 +238,402 @@ class ComputingUnitHelpersSpec // ── Bulk variants resolving from pre-fetched pod maps ──────────────── - "getComputingUnitStatus(unit, podPhases)" should "return Running for a local unit" in { - ComputingUnitHelpers.getComputingUnitStatus(localUnit(), Map.empty) shouldBe Running + "getComputingUnitStatusAndReason(unit, podSnapshots)" should "return Running for a local unit" in { + ComputingUnitHelpers.getComputingUnitStatusAndReason(localUnit(), Map.empty) shouldBe + ((Running, None)) } it should "return Running for a kubernetes unit whose pod phase is Running" in { val unit = kubernetesUnit(7) - val podPhases = Map(KubernetesClient.generatePodName(7) -> "Running") - ComputingUnitHelpers.getComputingUnitStatus(unit, podPhases) shouldBe Running + val podSnapshots = + Map(KubernetesClient.generatePodName(7) -> snapshotOf(podWithPhase("Running"))) + ComputingUnitHelpers.getComputingUnitStatusAndReason(unit, podSnapshots) shouldBe + ((Running, None)) } it should "return Pending for a kubernetes unit whose pod is absent or not Running" in { val unit = kubernetesUnit(8) - ComputingUnitHelpers.getComputingUnitStatus(unit, Map.empty) shouldBe Pending - ComputingUnitHelpers.getComputingUnitStatus( + ComputingUnitHelpers.getComputingUnitStatusAndReason(unit, Map.empty) shouldBe ((Pending, None)) + ComputingUnitHelpers.getComputingUnitStatusAndReason( unit, - Map(KubernetesClient.generatePodName(8) -> "Pending") - ) shouldBe Pending + Map(KubernetesClient.generatePodName(8) -> snapshotOf(podWithPhase("Pending"))) + ) shouldBe ((Pending, None)) } it should "treat a null phase as not Running" in { val unit = kubernetesUnit(9) - val podPhases = Map(KubernetesClient.generatePodName(9) -> (null: String)) - ComputingUnitHelpers.getComputingUnitStatus(unit, podPhases) shouldBe Pending + val podSnapshots = Map(KubernetesClient.generatePodName(9) -> snapshotOf(statuslessPod())) + ComputingUnitHelpers.getComputingUnitStatusAndReason(unit, podSnapshots) shouldBe + ((Pending, None)) + } + + it should "return Pending for an unknown (untyped) unit" in { + ComputingUnitHelpers.getComputingUnitStatusAndReason(untypedUnit(), Map.empty) shouldBe + ((Pending, None)) + } + + // ── kubernetesStatusAndReason: one test per decision-table row ─────── + + "kubernetesStatusAndReason" should "map an absent pod to Pending with no reason" in { + ComputingUnitHelpers.kubernetesStatusAndReason(None) shouldBe ((Pending, None)) + } + + it should "map a pod with a deletion timestamp to Terminating regardless of phase" in { + statusAndReasonOf(terminatingPod()) shouldBe ((Terminating, None)) + } + + it should "map a status-less pod that is being deleted to Terminating, not Pending" in { + // A pod can carry a deletion timestamp before its status is ever populated; the + // empty-status snapshot must still honor the terminating flag. + val pod = new PodBuilder() + .withNewMetadata() + .withDeletionTimestamp("2026-08-20T00:00:00Z") + .endMetadata() + .build() + statusAndReasonOf(pod) shouldBe ((Terminating, None)) + } + + it should "let Terminating win over an eviction" in { + // Deleting an already-evicted pod is the owner acting on the failure: the row should + // show the deletion in progress, not keep explaining the eviction. + val pod = new PodBuilder() + .withNewMetadata() + .withDeletionTimestamp("2026-08-20T00:00:00Z") + .endMetadata() + .withNewStatus() + .withPhase("Failed") + .withReason("Evicted") + .withMessage("Pod ephemeral local storage usage exceeds the total limit of containers 1Gi.") + .endStatus() + .build() + statusAndReasonOf(pod) shouldBe ((Terminating, None)) + } + + it should "map a disk-pressure eviction to Failed with the local-disk reason" in { + val (state, reason) = statusAndReasonOf( + evictedPod("Pod ephemeral local storage usage exceeds the total limit of containers 1Gi.") + ) + state shouldBe Failed + reason shouldBe Some( + "The computing unit was evicted because it ran out of local disk storage. Consider " + + "storing less data on the unit's local file system, or recreate it with more storage." + ) + } + + it should "map any other eviction to Failed with a short version of the cluster's reason" in { + val (state, reason) = statusAndReasonOf( + evictedPod("The node was low on resource: memory. Container was using more than its request.") + ) + state shouldBe Failed + reason shouldBe Some( + "The computing unit was evicted by the cluster (The node was low on resource: memory). " + + "Consider recreating it." + ) + } + + it should "fall back to the pod's Evicted reason when the eviction carries no message" in { + val (state, reason) = statusAndReasonOf(evictedPod(null)) + state shouldBe Failed + reason shouldBe Some( + "The computing unit was evicted by the cluster (Evicted). Consider recreating it." + ) + } + + it should "fall back to the pod's Evicted reason when the eviction message is only whitespace" in { + val (state, reason) = statusAndReasonOf(evictedPod(" ")) + state shouldBe Failed + reason shouldBe Some( + "The computing unit was evicted by the cluster (Evicted). Consider recreating it." + ) + } + + it should "use the whole eviction message when it contains no period" in { + val (state, reason) = statusAndReasonOf(evictedPod("The node was under memory pressure")) + state shouldBe Failed + reason shouldBe Some( + "The computing unit was evicted by the cluster (The node was under memory pressure). " + + "Consider recreating it." + ) + } + + it should "truncate an eviction sentence longer than 120 characters" in { + val sentence = "The node reported " + ("x" * 120) // 138 chars, no period + val (state, reason) = statusAndReasonOf(evictedPod(sentence)) + state shouldBe Failed + reason shouldBe Some( + s"The computing unit was evicted by the cluster (${sentence.take(120)}...). " + + "Consider recreating it." + ) + } + + it should "detect the disk-pressure wording case-insensitively" in { + Seq( + "EPHEMERAL storage limit exceeded", + "The node was low on Disk space" + ).foreach { message => + val (state, reason) = statusAndReasonOf(evictedPod(message)) + state shouldBe Failed + reason shouldBe Some( + "The computing unit was evicted because it ran out of local disk storage. Consider " + + "storing less data on the unit's local file system, or recreate it with more storage." + ) + } + } + + it should "let the eviction wording win over a crash-looping container" in { + // An evicted pod may still report a crash-looping container from before the eviction; + // the eviction is the root cause, so its wording must take precedence. + val pod = new PodBuilder() + .withNewStatus() + .withPhase("Failed") + .withReason("Evicted") + .withMessage("Pod ephemeral local storage usage exceeds the total limit of containers 1Gi.") + .addNewContainerStatus() + .withRestartCount(5) + .withNewState() + .withNewWaiting() + .withReason("CrashLoopBackOff") + .endWaiting() + .endState() + .endContainerStatus() + .endStatus() + .build() + + val (state, reason) = statusAndReasonOf(pod) + state shouldBe Failed + reason shouldBe Some( + "The computing unit was evicted because it ran out of local disk storage. Consider " + + "storing less data on the unit's local file system, or recreate it with more storage." + ) + } + + it should "map every image-pull waiting reason to Failed with the image-pull reason" in { + Seq("ImagePullBackOff", "ErrImagePull", "InvalidImageName").foreach { waitingReason => + val (state, reason) = statusAndReasonOf(waitingContainerPod("Pending", waitingReason)) + state shouldBe Failed + reason shouldBe Some( + "The computing unit's image could not be pulled. Please recreate the unit or contact " + + "an administrator." + ) + } + } + + it should "map a crash loop after an OOM kill to Failed with the out-of-memory reason" in { + // restartPolicy Always keeps the pod phase "Running" through an OOM crash loop, so the + // waiting/lastState container fields are the only signal. + val (state, reason) = statusAndReasonOf( + waitingContainerPod( + "Running", + "CrashLoopBackOff", + lastTerminatedReason = Some("OOMKilled"), + restartCount = 4 + ) + ) + state shouldBe Failed + reason shouldBe Some( + "The computing unit keeps crashing because it runs out of memory. Please terminate it " + + "and recreate it with a higher memory limit." + ) + } + + it should "map a plain crash loop to Failed with the restart count" in { + val (state, reason) = statusAndReasonOf( + waitingContainerPod( + "Running", + "CrashLoopBackOff", + lastTerminatedReason = Some("Error"), + restartCount = 7 + ) + ) + state shouldBe Failed + reason shouldBe Some( + "The computing unit is repeatedly crashing (restarted 7 times). Please terminate and " + + "recreate it, or contact an administrator." + ) + } + + it should "fail the unit when any container crash-loops, not just the first one" in { + // Multi-container pod (e.g. a sidecar): the healthy first container must not mask the + // crash-looping second one. + val pod = new PodBuilder() + .withNewStatus() + .withPhase("Running") + .addNewContainerStatus() + .withRestartCount(0) + .endContainerStatus() + .addNewContainerStatus() + .withRestartCount(3) + .withNewState() + .withNewWaiting() + .withReason("CrashLoopBackOff") + .endWaiting() + .endState() + .endContainerStatus() + .endStatus() + .build() + + val (state, reason) = statusAndReasonOf(pod) + state shouldBe Failed + reason shouldBe Some( + "The computing unit is repeatedly crashing (restarted 3 times). Please terminate and " + + "recreate it, or contact an administrator." + ) + } + + it should "word a crash loop by the crash-looping container's own history, not a sibling's" in { + // Container A recovered from an OOM kill and is running; container B crash-loops for a + // non-OOM reason. The wording (and restart count) must come from B's own lastState, so + // A's OOM history must not upgrade the message to the out-of-memory variant. + val pod = new PodBuilder() + .withNewStatus() + .withPhase("Running") + .addNewContainerStatus() + .withRestartCount(2) + .withNewLastState() + .withNewTerminated() + .withReason("OOMKilled") + .endTerminated() + .endLastState() + .endContainerStatus() + .addNewContainerStatus() + .withRestartCount(6) + .withNewState() + .withNewWaiting() + .withReason("CrashLoopBackOff") + .endWaiting() + .endState() + .withNewLastState() + .withNewTerminated() + .withReason("Error") + .endTerminated() + .endLastState() + .endContainerStatus() + .endStatus() + .build() + + val (state, reason) = statusAndReasonOf(pod) + state shouldBe Failed + reason shouldBe Some( + "The computing unit is repeatedly crashing (restarted 6 times). Please terminate and " + + "recreate it, or contact an administrator." + ) + } + + it should "let an image-pull failure win over a crash-looping sibling container" in { + // Precedence pin: the image-pull check runs before the crash-loop one, because a unit + // whose image cannot be pulled can only be fixed by recreating it. + val pod = new PodBuilder() + .withNewStatus() + .withPhase("Pending") + .addNewContainerStatus() + .withRestartCount(0) + .withNewState() + .withNewWaiting() + .withReason("ErrImagePull") + .endWaiting() + .endState() + .endContainerStatus() + .addNewContainerStatus() + .withRestartCount(3) + .withNewState() + .withNewWaiting() + .withReason("CrashLoopBackOff") + .endWaiting() + .endState() + .endContainerStatus() + .endStatus() + .build() + + val (state, reason) = statusAndReasonOf(pod) + state shouldBe Failed + reason shouldBe Some( + "The computing unit's image could not be pulled. Please recreate the unit or contact " + + "an administrator." + ) + } + + it should "map a non-evicted Failed phase to Failed with the generic reason" in { + val (state, reason) = statusAndReasonOf(podWithPhase("Failed")) + state shouldBe Failed + reason shouldBe Some( + "The computing unit stopped unexpectedly. Please terminate and recreate it, or contact " + + "an administrator." + ) + } + + it should "map an Unknown phase to Unknown with the unreachable-node reason" in { + val (state, reason) = statusAndReasonOf(podWithPhase("Unknown")) + state shouldBe Unknown + reason shouldBe Some( + "The state of the computing unit cannot be determined (its node may be unreachable)." + ) + } + + it should "keep an unschedulable pod Pending but explain the wait" in { + val (state, reason) = statusAndReasonOf(unschedulablePod()) + state shouldBe Pending + reason shouldBe Some( + "The computing unit is waiting for cluster resources to become available." + ) + } + + it should "ignore a stale Unschedulable condition once the pod is Running" in { + // The unschedulable explanation is gated on phase Pending, so a leftover + // PodScheduled=False condition on a pod that has since started must not resurface it. + val pod = new PodBuilder() + .withNewStatus() + .withPhase("Running") + .addNewCondition() + .withType("PodScheduled") + .withStatus("False") + .withReason("Unschedulable") + .endCondition() + .endStatus() + .build() + statusAndReasonOf(pod) shouldBe ((Running, None)) + } + + it should "keep a recovered OOM-killed pod Running but attach the memory warning" in { + // The container restarted in place after the OOM kill (restartPolicy Always), so the unit + // is usable again — the reason is a warning, not a failure. + val (state, reason) = statusAndReasonOf(restartedPod("OOMKilled", restartCount = 2)) + state shouldBe Running + reason shouldBe Some( + "The last run was terminated because the computing unit ran out of memory (restarted 2 " + + "times). Consider recreating the unit with a higher memory limit before running the " + + "same workload." + ) + } + + it should "not attach the memory warning to a container that restarted for another reason" in { + statusAndReasonOf(restartedPod("Error", restartCount = 1)) shouldBe ((Running, None)) + } + + it should "let a crash-loop failure win over the recovered-OOM warning" in { + // Both signals present: the waiting CrashLoopBackOff means the unit is NOT usable, so the + // failure must take precedence over the phase-Running OOM warning. + val (state, _) = statusAndReasonOf( + waitingContainerPod( + "Running", + "CrashLoopBackOff", + lastTerminatedReason = Some("OOMKilled"), + restartCount = 3 + ) + ) + state shouldBe Failed + } + + it should "map a plain Running phase to Running and anything else to Pending, without reasons" in { + statusAndReasonOf(podWithPhase("Running")) shouldBe ((Running, None)) + statusAndReasonOf(podWithPhase("Pending")) shouldBe ((Pending, None)) + statusAndReasonOf(podWithPhase("SomethingNew")) shouldBe ((Pending, None)) + } + + it should "map a Succeeded phase to Pending (current behavior, pinned deliberately)" in { + // "Succeeded" has no dedicated branch and falls through to Pending. With restartPolicy + // Always a computing-unit pod essentially never completes, but if one ever did, its unit + // would show as connecting indefinitely — pinned here so any future change is conscious. + statusAndReasonOf(podWithPhase("Succeeded")) shouldBe ((Pending, None)) } "getComputingUnitMetrics(unit, podMetrics)" should "return NaN metrics for a local unit" in { @@ -178,6 +655,11 @@ class ComputingUnitHelpersSpec WorkflowComputingUnitMetrics("", "") } + it should "return NaN for an unknown (untyped) unit" in { + ComputingUnitHelpers.getComputingUnitMetrics(untypedUnit(), Map.empty) shouldBe + WorkflowComputingUnitMetrics("NaN", "NaN") + } + // ── partitionLiveUnits ─────────────────────────────────────────────── "partitionLiveUnits" should "treat local units as always live" in { @@ -190,14 +672,29 @@ class ComputingUnitHelpersSpec it should "classify a kubernetes unit as live iff its pod is present in the map" in { val present = kubernetesUnit(20) val gone = kubernetesUnit(21) - val podPhases = Map(KubernetesClient.generatePodName(20) -> "Running") + val podSnapshots = + Map(KubernetesClient.generatePodName(20) -> snapshotOf(podWithPhase("Running"))) - val (live, vanished) = ComputingUnitHelpers.partitionLiveUnits(List(present, gone), podPhases) + val (live, vanished) = + ComputingUnitHelpers.partitionLiveUnits(List(present, gone), podSnapshots) live.map(_.getCuid) shouldBe List(20) vanished.map(_.getCuid) shouldBe List(21) } + it should "treat a unit whose pod is present but failed as live (dead, not vanished)" in { + // Failure states must NOT change the vanish semantics: a pod that still exists in the + // namespace — however broken — keeps its unit un-terminated so the owner can see why. + val failed = kubernetesUnit(22) + val podSnapshots = + Map(KubernetesClient.generatePodName(22) -> snapshotOf(podWithPhase("Failed"))) + + val (live, vanished) = ComputingUnitHelpers.partitionLiveUnits(List(failed), podSnapshots) + + live.map(_.getCuid) shouldBe List(22) + vanished shouldBe empty + } + it should "treat an untyped (null-type) unit as live (never kubernetes)" in { val (live, vanished) = ComputingUnitHelpers.partitionLiveUnits(List(untypedUnit()), Map.empty) live should have size 1 @@ -213,9 +710,10 @@ class ComputingUnitHelpersSpec val row = ComputingUnitHelpers.buildDashboardUnit( unit, isOwner = true, + canViewStatusReason = true, accessPrivilege = PrivilegeEnum.READ, ownerInfo = Map((100: Integer) -> ("avatar", "owner")), - podPhases = Map(podName -> "Running"), + podSnapshots = Map(podName -> snapshotOf(podWithPhase("Running"))), podMetrics = Map(podName -> Map("cpu" -> "100m", "memory" -> "64Mi")) ) @@ -223,6 +721,7 @@ class ComputingUnitHelpersSpec row.isOwner shouldBe true row.accessPrivilege shouldBe PrivilegeEnum.READ row.status shouldBe "Running" + row.statusReason shouldBe None row.metrics shouldBe WorkflowComputingUnitMetrics("100m", "64Mi") row.ownerAvatar shouldBe "avatar" row.ownerName shouldBe "owner" @@ -232,9 +731,10 @@ class ComputingUnitHelpersSpec val row = ComputingUnitHelpers.buildDashboardUnit( localUnit(cuid = 31, uid = 200), isOwner = false, + canViewStatusReason = false, accessPrivilege = PrivilegeEnum.WRITE, ownerInfo = Map.empty, - podPhases = Map.empty, + podSnapshots = Map.empty, podMetrics = Map.empty ) @@ -244,28 +744,49 @@ class ComputingUnitHelpersSpec row.metrics shouldBe WorkflowComputingUnitMetrics("NaN", "NaN") } - // ── Bulk variants: unknown (untyped) branch ────────────────────────── + it should "include the status reason only when the caller permits it" in { + val unit = kubernetesUnit(cuid = 32, uid = 100) + val podName = KubernetesClient.generatePodName(32) + val podSnapshots = Map(podName -> snapshotOf(podWithPhase("Failed"))) + + def build(isOwner: Boolean, canViewStatusReason: Boolean) = + ComputingUnitHelpers.buildDashboardUnit( + unit, + isOwner = isOwner, + canViewStatusReason = canViewStatusReason, + accessPrivilege = PrivilegeEnum.READ, + ownerInfo = Map.empty, + podSnapshots = podSnapshots, + podMetrics = Map.empty + ) - "getComputingUnitStatus(unit, podPhases)" should "return Pending for an unknown (untyped) unit" in { - ComputingUnitHelpers.getComputingUnitStatus(untypedUnit(), Map.empty) shouldBe Pending - } + val ownerRow = build(isOwner = true, canViewStatusReason = true) + ownerRow.status shouldBe "Failed" + ownerRow.statusReason shouldBe Some( + "The computing unit stopped unexpectedly. Please terminate and recreate it, or contact " + + "an administrator." + ) - "getComputingUnitMetrics(unit, podMetrics)" should "return NaN for an unknown (untyped) unit" in { - ComputingUnitHelpers.getComputingUnitMetrics(untypedUnit(), Map.empty) shouldBe - WorkflowComputingUnitMetrics("NaN", "NaN") + val sharedRow = build(isOwner = false, canViewStatusReason = false) + sharedRow.status shouldBe "Failed" + sharedRow.statusReason shouldBe None + + val adminRow = build(isOwner = false, canViewStatusReason = true) + adminRow.isOwner shouldBe false + adminRow.statusReason shouldBe ownerRow.statusReason } - // ── podPhasesFor / podMetricsFor guards ────────────────────────────── + // ── podSnapshotsFor / podMetricsFor guards ─────────────────────────── - "podPhasesFor" should "return empty (issuing no cluster call) when no kubernetes unit is present" in { - ComputingUnitHelpers.podPhasesFor(List(localUnit(), untypedUnit())) shouldBe empty + "podSnapshotsFor" should "return empty (issuing no cluster call) when no kubernetes unit is present" in { + ComputingUnitHelpers.podSnapshotsFor(List(localUnit(), untypedUnit())) shouldBe empty } - it should "fetch all pod phases once when a kubernetes unit is present" in { + it should "fetch all pod snapshots once when a kubernetes unit is present" in { val k8s = mock(classOf[KubernetesClient]) - val phases = Map("computing-unit-50" -> "Running") - when(k8s.getAllPodPhases).thenReturn(phases) - ComputingUnitHelpers.podPhasesFor(List(kubernetesUnit(50)), k8s) shouldBe phases + val snapshots = Map("computing-unit-50" -> snapshotOf(podWithPhase("Running"))) + when(k8s.getAllPodStatusSnapshots).thenReturn(snapshots) + ComputingUnitHelpers.podSnapshotsFor(List(kubernetesUnit(50)), k8s) shouldBe snapshots } "podMetricsFor" should "return empty (issuing no cluster call) when no kubernetes unit is present" in { @@ -308,12 +829,13 @@ class ComputingUnitHelpersSpec Seq(present, gone, local).foreach(computingUnitDao.insert(_)) // Only the pod for cuid 600 exists; cuid 601's pod has vanished. - val podPhases = Map(KubernetesClient.generatePodName(600) -> "Running") + val podSnapshots = + Map(KubernetesClient.generatePodName(600) -> snapshotOf(podWithPhase("Running"))) val live = ComputingUnitHelpers.reconcileVanishedKubernetesUnits( computingUnitDao, List(present, gone, local), - podPhases + podSnapshots ) live.map(_.getCuid) should contain theSameElementsAs Seq(600, 602) @@ -323,4 +845,24 @@ class ComputingUnitHelpersSpec computingUnitDao.fetchOneByCuid(600).getTerminateTime shouldBe null computingUnitDao.fetchOneByCuid(602).getTerminateTime shouldBe null } + + it should "not terminate a unit whose pod is present but in a failure state" in { + userDao.insert(makeUser(610, "dave", "[email protected]", null)) + + val failed = kubernetesUnit(611, 610) + failed.setName("failed") + computingUnitDao.insert(failed) + + val podSnapshots = + Map(KubernetesClient.generatePodName(611) -> snapshotOf(podWithPhase("Failed"))) + val live = + ComputingUnitHelpers.reconcileVanishedKubernetesUnits( + computingUnitDao, + List(failed), + podSnapshots + ) + + live.map(_.getCuid) shouldBe List(611) + computingUnitDao.fetchOneByCuid(611).getTerminateTime shouldBe null + } } diff --git a/computing-unit-managing-service/src/test/scala/org/apache/texera/service/util/KubernetesClientSpec.scala b/computing-unit-managing-service/src/test/scala/org/apache/texera/service/util/KubernetesClientSpec.scala index 8931cb8b51..2723c35a53 100644 --- a/computing-unit-managing-service/src/test/scala/org/apache/texera/service/util/KubernetesClientSpec.scala +++ b/computing-unit-managing-service/src/test/scala/org/apache/texera/service/util/KubernetesClientSpec.scala @@ -56,13 +56,13 @@ import org.scalatest.matchers.should.Matchers import scala.jdk.CollectionConverters._ // Two layers are exercised here: -// * the pure fabric8 -> map transforms (phasesByPodName / metricsByPodName) with +// * the pure fabric8 -> map transforms (snapshotsByPodName / metricsByPodName) with // builder-constructed model objects, so the transform logic needs no client, and -// * the thin namespace-wide wrappers (getAllPodPhases / getAllPodMetrics / getPodMetrics), -// which are driven through a freshly constructed KubernetesClient whose fabric8 client is a -// Mockito stub — no live cluster and no mutable global. -// The status/metrics *decision* logic that consumes these maps (Running vs Pending, cpu/memory -// resolution) is covered by ComputingUnitHelpersSpec. +// * the thin namespace-wide wrappers (getAllPodStatusSnapshots / getAllPodMetrics / +// getPodMetrics), which are driven through a freshly constructed KubernetesClient whose +// fabric8 client is a Mockito stub — no live cluster and no mutable global. +// The status/metrics *decision* logic that consumes these maps (Running vs Failed vs Pending, +// cpu/memory resolution) is covered by ComputingUnitHelpersSpec. class KubernetesClientSpec extends AnyFlatSpec with Matchers { private val testAccessControlServiceUrl = @@ -142,16 +142,123 @@ class KubernetesClientSpec extends AnyFlatSpec with Matchers { KubernetesClient.generatePodName(0) shouldBe "computing-unit-0" } - "phasesByPodName" should "map every pod name to its phase" in { - val phases = KubernetesClient.phasesByPodName(Seq(pod(1, "Running"), pod(2, "Pending"))) - phases(KubernetesClient.generatePodName(1)) shouldBe "Running" - phases(KubernetesClient.generatePodName(2)) shouldBe "Pending" + "snapshotsByPodName" should "map every pod name to a snapshot of its phase" in { + val snapshots = KubernetesClient.snapshotsByPodName(Seq(pod(1, "Running"), pod(2, "Pending"))) + snapshots(KubernetesClient.generatePodName(1)).phase shouldBe Some("Running") + snapshots(KubernetesClient.generatePodName(2)).phase shouldBe Some("Pending") } - it should "map a pod with no status to a null phase but still include it" in { - val phases = KubernetesClient.phasesByPodName(Seq(statuslessPod(3))) - phases should contain key KubernetesClient.generatePodName(3) - phases(KubernetesClient.generatePodName(3)) shouldBe null + it should "map a pod with no status to an empty snapshot but still include it" in { + val snapshots = KubernetesClient.snapshotsByPodName(Seq(statuslessPod(3))) + snapshots should contain key KubernetesClient.generatePodName(3) + snapshots(KubernetesClient.generatePodName(3)) shouldBe PodStatusSnapshot.empty + } + + // ── PodStatusSnapshot.fromPod: the pure fabric8 -> snapshot extraction ── + // The decision table consuming these fields lives in ComputingUnitHelpersSpec; here each + // fabric8 field is pinned to its snapshot slot, including the null-guard paths. + + "PodStatusSnapshot.fromPod" should "capture the phase, pod-level reason and message" in { + val evicted = new PodBuilder() + .withNewStatus() + .withPhase("Failed") + .withReason("Evicted") + .withMessage("The node was low on resource: ephemeral-storage.") + .endStatus() + .build() + + val snapshot = PodStatusSnapshot.fromPod(evicted) + snapshot.phase shouldBe Some("Failed") + snapshot.podReason shouldBe Some("Evicted") + snapshot.podMessage shouldBe Some("The node was low on resource: ephemeral-storage.") + snapshot.terminating shouldBe false + snapshot.unschedulable shouldBe false + snapshot.containers shouldBe empty + } + + it should "flag a pod carrying a deletion timestamp as terminating" in { + val deleted = new PodBuilder() + .withNewMetadata() + .withName(KubernetesClient.generatePodName(9)) + .withDeletionTimestamp("2026-08-20T00:00:00Z") + .endMetadata() + .withNewStatus() + .withPhase("Running") + .endStatus() + .build() + + PodStatusSnapshot.fromPod(deleted).terminating shouldBe true + } + + it should "flag an Unschedulable PodScheduled=False condition and only that condition" in { + def podWithCondition(condType: String, status: String, reason: String) = + new PodBuilder() + .withNewStatus() + .withPhase("Pending") + .addNewCondition() + .withType(condType) + .withStatus(status) + .withReason(reason) + .endCondition() + .endStatus() + .build() + + PodStatusSnapshot + .fromPod(podWithCondition("PodScheduled", "False", "Unschedulable")) + .unschedulable shouldBe true + // A scheduled pod, or a False condition of another type/reason, must not count. + PodStatusSnapshot + .fromPod(podWithCondition("PodScheduled", "True", "")) + .unschedulable shouldBe false + PodStatusSnapshot + .fromPod(podWithCondition("Ready", "False", "Unschedulable")) + .unschedulable shouldBe false + // A malformed condition carrying the Unschedulable reason with status True must not + // count either: all three fields have to line up. + PodStatusSnapshot + .fromPod(podWithCondition("PodScheduled", "True", "Unschedulable")) + .unschedulable shouldBe false + } + + it should "capture each container's waiting reason, last termination reason and restart count" in { + val crashLooping = new PodBuilder() + .withNewStatus() + .withPhase("Running") + .addNewContainerStatus() + .withRestartCount(5) + .withNewState() + .withNewWaiting() + .withReason("CrashLoopBackOff") + .endWaiting() + .endState() + .withNewLastState() + .withNewTerminated() + .withReason("OOMKilled") + .endTerminated() + .endLastState() + .endContainerStatus() + .endStatus() + .build() + + val container = PodStatusSnapshot.fromPod(crashLooping).containers.head + container.waitingReason shouldBe Some("CrashLoopBackOff") + container.lastTerminatedReason shouldBe Some("OOMKilled") + container.restartCount shouldBe 5 + } + + it should "map a container status with no state details to an empty container snapshot" in { + val bare = new PodBuilder() + .withNewStatus() + .withPhase("Running") + .addNewContainerStatus() + .endContainerStatus() + .endStatus() + .build() + + val container = PodStatusSnapshot.fromPod(bare).containers.head + container.waitingReason shouldBe None + container.lastTerminatedReason shouldBe None + container.restartCount shouldBe 0 } "metricsByPodName" should "flatten each pod's container usage into a cpu/memory map" in { @@ -162,13 +269,13 @@ class KubernetesClientSpec extends AnyFlatSpec with Matchers { // ── namespace-wide wrappers, driven through a stubbed fabric8 client ── // These pin the fabric8 fluent-chain plumbing (list() / top().pods().metrics()); the value - // transform they delegate to is already pinned by the phasesByPodName / metricsByPodName tests, - // so they only assert that the namespace items flow through keyed by pod name. + // transform they delegate to is already pinned by the snapshotsByPodName / metricsByPodName + // tests, so they only assert that the namespace items flow through keyed by pod name. - "getAllPodPhases" should "list the namespace pods and key them by pod name" in { + "getAllPodStatusSnapshots" should "list the namespace pods and key them by pod name" in { val k8s = new KubernetesClient(stubbedClient(Seq(pod(1, "Running"), pod(2, "Pending")), Seq.empty)) - k8s.getAllPodPhases.keySet shouldBe + k8s.getAllPodStatusSnapshots.keySet shouldBe Set(KubernetesClient.generatePodName(1), KubernetesClient.generatePodName(2)) } diff --git a/frontend/src/app/common/service/computing-unit/computing-unit-status/computing-unit-status.service.spec.ts b/frontend/src/app/common/service/computing-unit/computing-unit-status/computing-unit-status.service.spec.ts index d1c18ab450..7052cb3619 100644 --- a/frontend/src/app/common/service/computing-unit/computing-unit-status/computing-unit-status.service.spec.ts +++ b/frontend/src/app/common/service/computing-unit/computing-unit-status/computing-unit-status.service.spec.ts @@ -199,12 +199,45 @@ describe("ComputingUnitStatusService", () => { expect(status).toBe(ComputingUnitState.Pending); }); - it("getStatus() maps an unrecognized status to Pending (default branch)", () => { + it("getStatus() maps a Failed unit to ComputingUnitState.Failed", () => { + (service as any).selectedUnitSubject.next({ + computingUnit: { cuid: 1 }, + status: "Failed", + } as unknown as DashboardWorkflowComputingUnit); + + let status: ComputingUnitState | undefined; + service.getStatus().subscribe(s => (status = s)); + expect(status).toBe(ComputingUnitState.Failed); + }); + + it("getStatus() maps an Unknown unit to ComputingUnitState.Unknown", () => { + (service as any).selectedUnitSubject.next({ + computingUnit: { cuid: 1 }, + status: "Unknown", + } as unknown as DashboardWorkflowComputingUnit); + + let status: ComputingUnitState | undefined; + service.getStatus().subscribe(s => (status = s)); + expect(status).toBe(ComputingUnitState.Unknown); + }); + + it("getStatus() maps a Terminating unit to ComputingUnitState.Terminating", () => { (service as any).selectedUnitSubject.next({ computingUnit: { cuid: 1 }, status: "Terminating", } as unknown as DashboardWorkflowComputingUnit); + let status: ComputingUnitState | undefined; + service.getStatus().subscribe(s => (status = s)); + expect(status).toBe(ComputingUnitState.Terminating); + }); + + it("getStatus() keeps an unrecognized future status on the Pending fallback", () => { + (service as any).selectedUnitSubject.next({ + computingUnit: { cuid: 1 }, + status: "FutureStatus", + } as unknown as DashboardWorkflowComputingUnit); + let status: ComputingUnitState | undefined; service.getStatus().subscribe(s => (status = s)); expect(status).toBe(ComputingUnitState.Pending); diff --git a/frontend/src/app/common/service/computing-unit/computing-unit-status/computing-unit-status.service.ts b/frontend/src/app/common/service/computing-unit/computing-unit-status/computing-unit-status.service.ts index e5b6c44388..93dd009335 100644 --- a/frontend/src/app/common/service/computing-unit/computing-unit-status/computing-unit-status.service.ts +++ b/frontend/src/app/common/service/computing-unit/computing-unit-status/computing-unit-status.service.ts @@ -229,6 +229,12 @@ export class ComputingUnitStatusService implements OnDestroy { return ComputingUnitState.Running; case "Pending": return ComputingUnitState.Pending; + case "Failed": + return ComputingUnitState.Failed; + case "Unknown": + return ComputingUnitState.Unknown; + case "Terminating": + return ComputingUnitState.Terminating; default: return ComputingUnitState.Pending; } diff --git a/frontend/src/app/common/type/computing-unit-connection.interface.ts b/frontend/src/app/common/type/computing-unit-connection.interface.ts index e0cf04a270..eed75fa91d 100644 --- a/frontend/src/app/common/type/computing-unit-connection.interface.ts +++ b/frontend/src/app/common/type/computing-unit-connection.interface.ts @@ -24,5 +24,8 @@ export enum ComputingUnitState { Running = "Running", Pending = "Pending", + Failed = "Failed", + Unknown = "Unknown", + Terminating = "Terminating", NoComputingUnit = "No Computing Unit", } diff --git a/frontend/src/app/common/type/workflow-computing-unit.ts b/frontend/src/app/common/type/workflow-computing-unit.ts index 8515aa32e8..a883b39200 100644 --- a/frontend/src/app/common/type/workflow-computing-unit.ts +++ b/frontend/src/app/common/type/workflow-computing-unit.ts @@ -46,7 +46,9 @@ export interface WorkflowComputingUnitMetrics { export interface DashboardWorkflowComputingUnit { computingUnit: WorkflowComputingUnit; - status: "Running" | "Pending"; + status: "Running" | "Pending" | "Failed" | "Unknown" | "Terminating"; + // Explanation of a failing or degraded status. Sent only to authorized viewers; null otherwise. + statusReason?: string | null; metrics: WorkflowComputingUnitMetrics; isOwner: boolean; accessPrivilege: "READ" | "WRITE" | "NONE"; diff --git a/frontend/src/app/common/util/computing-unit.util.spec.ts b/frontend/src/app/common/util/computing-unit.util.spec.ts index ab6d3b2a4f..88ff3e514c 100644 --- a/frontend/src/app/common/util/computing-unit.util.spec.ts +++ b/frontend/src/app/common/util/computing-unit.util.spec.ts @@ -20,8 +20,10 @@ import { TestBed, ComponentFixture } from "@angular/core/testing"; import { NZ_MODAL_DATA } from "ng-zorro-antd/modal"; import type { DashboardWorkflowComputingUnit } from "../type/workflow-computing-unit"; +import { ComputingUnitState } from "../type/computing-unit-connection.interface"; import { ComputingUnitMetadataComponent, + unavailableComputingUnitReason, parseResourceUnit, parseResourceNumber, cpuResourceConversion, @@ -30,6 +32,7 @@ import { memoryPercentage, validateName, getComputingUnitBadgeColor, + getComputingUnitRowTooltip, getComputingUnitStatusTooltip, getComputingUnitCpuStatus, getComputingUnitMemoryStatus, @@ -188,11 +191,48 @@ describe("getComputingUnitBadgeColor", () => { expect(getComputingUnitBadgeColor("Pending")).toBe("gold"); }); + it("should keep a terminating unit gold (transient, not broken)", () => { + expect(getComputingUnitBadgeColor("Terminating")).toBe("gold"); + }); + + it("should mark failure states red", () => { + expect(getComputingUnitBadgeColor("Failed")).toBe("red"); + expect(getComputingUnitBadgeColor("Unknown")).toBe("red"); + }); + it("should default to red for unknown statuses", () => { expect(getComputingUnitBadgeColor("Terminated")).toBe("red"); }); }); +describe("unavailableComputingUnitReason", () => { + it("names a shutting-down unit separately from a dead one", () => { + expect(unavailableComputingUnitReason(ComputingUnitState.Terminating)).toBe("terminating"); + }); + + it("reports a unit that can never come back as unavailable", () => { + expect(unavailableComputingUnitReason(ComputingUnitState.Failed)).toBe("unavailable"); + expect(unavailableComputingUnitReason(ComputingUnitState.Unknown)).toBe("unavailable"); + }); + + it("reports no reason for a unit that can still accept work, or for no unit at all", () => { + expect(unavailableComputingUnitReason(ComputingUnitState.Running)).toBeUndefined(); + // Pending is still starting up, so it must not be blocked. + expect(unavailableComputingUnitReason(ComputingUnitState.Pending)).toBeUndefined(); + expect(unavailableComputingUnitReason(ComputingUnitState.NoComputingUnit)).toBeUndefined(); + expect(unavailableComputingUnitReason(undefined)).toBeUndefined(); + }); + + it("reads a dashboard unit's own status field, not only the connection enum", () => { + // The enum and the DTO string share the same values, so one helper serves both. + expect(unavailableComputingUnitReason(makeUnit({ status: "Terminating" }).status)).toBe("terminating"); + expect(unavailableComputingUnitReason(makeUnit({ status: "Failed" }).status)).toBe("unavailable"); + expect(unavailableComputingUnitReason(makeUnit({ status: "Unknown" }).status)).toBe("unavailable"); + expect(unavailableComputingUnitReason(makeUnit({ status: "Running" }).status)).toBeUndefined(); + expect(unavailableComputingUnitReason(makeUnit({ status: "Pending" }).status)).toBeUndefined(); + }); +}); + describe("getComputingUnitStatusTooltip", () => { it("should map statuses to tooltip text", () => { expect(getComputingUnitStatusTooltip({ status: "Running" } as unknown as DashboardWorkflowComputingUnit)).toBe( @@ -203,10 +243,95 @@ describe("getComputingUnitStatusTooltip", () => { ); }); + it("should show the backend's status reason whenever one is present", () => { + // An authorized viewer's reason is already user-friendly, so it wins over every canned text: + // a failure explanation, a Pending unschedulable note, or a Running recovered-OOM warning. + expect( + getComputingUnitStatusTooltip({ + status: "Failed", + statusReason: "The computing unit keeps crashing because it runs out of memory.", + } as unknown as DashboardWorkflowComputingUnit) + ).toBe("The computing unit keeps crashing because it runs out of memory."); + expect( + getComputingUnitStatusTooltip({ + status: "Running", + statusReason: "The last run was terminated because the computing unit ran out of memory.", + } as unknown as DashboardWorkflowComputingUnit) + ).toBe("The last run was terminated because the computing unit ran out of memory."); + expect( + getComputingUnitStatusTooltip({ + status: "Pending", + statusReason: "The computing unit is waiting for cluster resources to become available.", + } as unknown as DashboardWorkflowComputingUnit) + ).toBe("The computing unit is waiting for cluster resources to become available."); + }); + + it("should show a generic unavailable message for reason-less failure states", () => { + expect(getComputingUnitStatusTooltip({ status: "Failed" } as unknown as DashboardWorkflowComputingUnit)).toBe( + "This computing unit is unavailable." + ); + expect(getComputingUnitStatusTooltip({ status: "Unknown" } as unknown as DashboardWorkflowComputingUnit)).toBe( + "This computing unit is unavailable." + ); + }); + + it("should not let an empty-string statusReason shadow the canned fallback", () => { + // The backend omits the reason (undefined) for ordinary shared users, but an empty string + // must behave the same way: it is falsy, so the status switch still decides the text. + expect( + getComputingUnitStatusTooltip({ + status: "Failed", + statusReason: "", + } as unknown as DashboardWorkflowComputingUnit) + ).toBe("This computing unit is unavailable."); + }); + + it("should fall back to the canned text when the backend sends a null status reason", () => { + // The backend sends a missing reason as null, which must fall back just like undefined. + expect(getComputingUnitStatusTooltip(makeUnit({ status: "Failed", statusReason: null }))).toBe( + "This computing unit is unavailable." + ); + expect(getComputingUnitStatusTooltip(makeUnit({ status: "Pending", statusReason: null }))).toBe( + "Computing unit is starting up" + ); + }); + it("should fall back to the raw status for unknown statuses", () => { expect(getComputingUnitStatusTooltip({ status: "Terminated" } as unknown as DashboardWorkflowComputingUnit)).toBe( "Terminated" ); + // Terminating is the one transient state with its own canned text, so the raw + // enum word never reaches the user. + expect(getComputingUnitStatusTooltip({ status: "Terminating" } as unknown as DashboardWorkflowComputingUnit)).toBe( + "Computing unit is shutting down" + ); + }); +}); + +describe("getComputingUnitRowTooltip", () => { + it("should show only the status tooltip for a selectable running unit", () => { + expect(getComputingUnitRowTooltip(makeUnit({ status: "Running" }))).toBe("Ready to use"); + expect(getComputingUnitRowTooltip(makeUnit({ status: "Running", statusReason: "OOM warning." }))).toBe( + "OOM warning." + ); + }); + + it("should explain why a non-running unit cannot be selected without doubling punctuation", () => { + expect(getComputingUnitRowTooltip(makeUnit({ status: "Pending" }))).toBe( + "Computing unit is starting up. Cannot select." + ); + expect(getComputingUnitRowTooltip(makeUnit({ status: "Failed" }))).toBe( + "This computing unit is unavailable. Cannot select." + ); + expect(getComputingUnitRowTooltip(makeUnit({ status: "Failed", statusReason: "Crash loop detected." }))).toBe( + "Crash loop detected. Cannot select." + ); + }); + + it("should fall back to the canned text when the backend sends a null status reason", () => { + expect(getComputingUnitRowTooltip(makeUnit({ status: "Failed", statusReason: null }))).toBe( + "This computing unit is unavailable. Cannot select." + ); }); }); diff --git a/frontend/src/app/common/util/computing-unit.util.ts b/frontend/src/app/common/util/computing-unit.util.ts index aab9fec2a3..4a1f25f1b2 100644 --- a/frontend/src/app/common/util/computing-unit.util.ts +++ b/frontend/src/app/common/util/computing-unit.util.ts @@ -20,6 +20,7 @@ import { Component, inject } from "@angular/core"; import { NZ_MODAL_DATA } from "ng-zorro-antd/modal"; import { DashboardWorkflowComputingUnit } from "../type/workflow-computing-unit"; +import { ComputingUnitState } from "../type/computing-unit-connection.interface"; @Component({ template: ` @@ -193,28 +194,73 @@ export function validateName(trimmedName: string): string | null { return null; } +/** Reason returned by unavailableComputingUnitReason. */ +export type UnavailableComputingUnitReason = "terminating" | "unavailable"; + +/** + * Returns the reason a computing unit cannot accept work: "terminating" for Terminating, + * "unavailable" for Failed or Unknown, and undefined for any other status or no status. + * + * Accepts either a ComputingUnitState or the status string of a DashboardWorkflowComputingUnit. + */ +export function unavailableComputingUnitReason( + status: ComputingUnitState | DashboardWorkflowComputingUnit["status"] | undefined +): UnavailableComputingUnitReason | undefined { + switch (status) { + case ComputingUnitState.Terminating: + return "terminating"; + case ComputingUnitState.Failed: + case ComputingUnitState.Unknown: + return "unavailable"; + default: + return undefined; + } +} + export function getComputingUnitBadgeColor(status: string): string { switch (status) { case "Running": return "green"; + // Pending and Terminating are transient, not broken, so both stay gold. case "Pending": + case "Terminating": return "gold"; + // Failed / Unknown and anything unrecognized. default: return "red"; } } export function getComputingUnitStatusTooltip(entry: DashboardWorkflowComputingUnit): string { + // A provided statusReason is already user-friendly and more specific than any canned + // text (failure cause, unschedulable wait, or a recovered-OOM warning on a Running unit). + if (entry.statusReason) { + return entry.statusReason; + } switch (entry.status) { case "Running": return "Ready to use"; case "Pending": return "Computing unit is starting up"; + // Ordinary shared users get no reason from the backend, only the generic text. + case "Failed": + case "Unknown": + return "This computing unit is unavailable."; + case "Terminating": + return "Computing unit is shutting down"; default: return entry.status; } } +export function getComputingUnitRowTooltip(entry: DashboardWorkflowComputingUnit): string { + const statusTooltip = getComputingUnitStatusTooltip(entry); + if (entry.status === "Running") { + return statusTooltip; + } + return `${statusTooltip.replace(/\.$/, "")}. Cannot select.`; +} + export function getComputingUnitCpuStatus(percentage: number): "success" | "exception" | "active" | "normal" { if (percentage > 90) return "exception"; if (percentage > 50) return "normal"; diff --git a/frontend/src/app/workspace/component/menu/menu.component.spec.ts b/frontend/src/app/workspace/component/menu/menu.component.spec.ts index e9ab4f7e90..83ceb5127e 100644 --- a/frontend/src/app/workspace/component/menu/menu.component.spec.ts +++ b/frontend/src/app/workspace/component/menu/menu.component.spec.ts @@ -389,10 +389,10 @@ describe("MenuComponent", () => { expect(resumeSpy).toHaveBeenCalled(); }); - it("returns 'Connecting' when a unit exists but the websocket is not connected", () => { + it("returns 'Connecting' for a Pending unit when the websocket is not connected", () => { component.isWorkflowValid = true; component.isWorkflowEmpty = false; - component.computingUnitStatus = ComputingUnitState.Running; + component.computingUnitStatus = ComputingUnitState.Pending; Object.defineProperty(component.workflowWebsocketService, "isConnected", { get: () => false, configurable: true, @@ -404,6 +404,66 @@ describe("MenuComponent", () => { expect(behavior.disable).toBe(true); }); + it.each([false, true])( + "returns 'Shutting Down' for a Terminating unit when websocket connected is %s", + isConnected => { + component.isWorkflowValid = true; + component.isWorkflowEmpty = false; + component.computingUnitStatus = ComputingUnitState.Terminating; + component.executionState = ExecutionState.Uninitialized; + Object.defineProperty(component.workflowWebsocketService, "isConnected", { + get: () => isConnected, + configurable: true, + }); + const runSpy = vi.spyOn(component, "runWorkflow"); + + const behavior = component.getRunButtonBehavior(); + behavior.onClick(); + + expect(behavior.text).toBe("Shutting Down"); + expect(behavior.icon).toBe("loading"); + expect(behavior.disable).toBe(true); + expect(runSpy).not.toHaveBeenCalled(); + } + ); + + it.each([ + { valid: false, empty: false, text: "Invalid Workflow", icon: "warning" }, + { valid: true, empty: true, text: "Empty Workflow", icon: "info-circle" }, + ])("keeps '$text' precedence while the computing unit is Terminating", ({ valid, empty, text, icon }) => { + component.isWorkflowValid = valid; + component.isWorkflowEmpty = empty; + component.computingUnitStatus = ComputingUnitState.Terminating; + + const behavior = component.getRunButtonBehavior(); + + expect(behavior.text).toBe(text); + expect(behavior.icon).toBe(icon); + expect(behavior.disable).toBe(true); + }); + + it("returns 'Unit Unavailable' instead of an endless 'Connecting' spinner for a dead unit", () => { + // A Failed/Unknown unit can never become connected, so the dead-unit check must come + // before the websocket-disconnected branch that would otherwise spin forever. + component.isWorkflowValid = true; + component.isWorkflowEmpty = false; + + Object.defineProperty(component.workflowWebsocketService, "isConnected", { + get: () => false, + configurable: true, + }); + + for (const status of [ComputingUnitState.Failed, ComputingUnitState.Unknown]) { + component.computingUnitStatus = status; + + const behavior = component.getRunButtonBehavior(); + + expect(behavior.text).toBe("Unit Unavailable"); + expect(behavior.icon).toBe("warning"); + expect(behavior.disable).toBe(true); + } + }); + /** Puts the component into the valid, connected state the state switch needs. */ function connected(): void { component.isWorkflowValid = true; diff --git a/frontend/src/app/workspace/component/menu/menu.component.ts b/frontend/src/app/workspace/component/menu/menu.component.ts index e320fc484e..54f5bde599 100644 --- a/frontend/src/app/workspace/component/menu/menu.component.ts +++ b/frontend/src/app/workspace/component/menu/menu.component.ts @@ -50,6 +50,7 @@ import { USER_WORKFLOW, workspaceFormUrl } from "../../../app-routing.constant"; import { ComputingUnitStatusService } from "../../../common/service/computing-unit/computing-unit-status/computing-unit-status.service"; import { WarehouseService } from "../../../common/service/warehouse/warehouse.service"; import { ComputingUnitState } from "../../../common/type/computing-unit-connection.interface"; +import { unavailableComputingUnitReason } from "../../../common/util/computing-unit.util"; import { ComputingUnitSelectionComponent } from "../power-button/computing-unit-selection.component"; import { GuiConfigService } from "../../../common/service/gui-config.service"; import { DashboardWorkflowComputingUnit } from "../../../common/type/workflow-computing-unit"; @@ -413,6 +414,26 @@ export class MenuComponent implements OnInit, OnDestroy { }; } + // Checked before the "Connecting" branch below, which would otherwise spin forever: + // these units are not coming back. + const unavailableReason = unavailableComputingUnitReason(this.computingUnitStatus); + if (unavailableReason === "terminating") { + return { + text: "Shutting Down", + icon: "loading", + disable: true, + onClick: () => {}, + }; + } + if (unavailableReason === "unavailable") { + return { + text: "Unit Unavailable", + icon: "warning", + disable: true, + onClick: () => {}, + }; + } + // This handles the case where a unit exists but we're not connected to it if (this.computingUnitStatus !== ComputingUnitState.NoComputingUnit && !this.workflowWebsocketService.isConnected) { return { diff --git a/frontend/src/app/workspace/component/power-button/computing-unit-selection.component.html b/frontend/src/app/workspace/component/power-button/computing-unit-selection.component.html index 7511f52bc5..0702ab491b 100644 --- a/frontend/src/app/workspace/component/power-button/computing-unit-selection.component.html +++ b/frontend/src/app/workspace/component/power-button/computing-unit-selection.component.html @@ -198,8 +198,8 @@ 'unit-selected': isSelectedUnit(unit), 'unit-connecting': unit.status === 'Pending', }" - [nz-tooltip]="cannotSelectUnit(unit) ? getUnitStatusTooltip(unit) + '. Cannot select.' : ''" - (click)="onPickComputingUnit(unit)"> + [nz-tooltip]="getRowTooltip(unit)" + (click)="onClickComputingUnitRow(unit)"> <div class="computing-unit-row"> <texera-user-avatar [avatar]="unit.ownerAvatar" @@ -209,19 +209,9 @@ [style.opacity]="0.7"> </texera-user-avatar> <div class="computing-unit-name"> - <nz-badge - [nzColor]="getBadgeColor(unit.status)" - [nz-tooltip]="getUnitStatusTooltip(unit)"></nz-badge> - <span - *ngIf="editingNameOfUnit !== unit.computingUnit.cuid; else editableUnitName" - nz-tooltip - [nzTooltipTitle]="unit.computingUnit.uri"> + <nz-badge [nzColor]="getBadgeColor(unit.status)"></nz-badge> + <span *ngIf="editingNameOfUnit !== unit.computingUnit.cuid; else editableUnitName"> {{ unit.computingUnit.name }} - <span - *ngIf="unit.status === 'Pending'" - class="unit-status-indicator" - >(Connecting)</span - > </span> <ng-template #editableUnitName> <input diff --git a/frontend/src/app/workspace/component/power-button/computing-unit-selection.component.scss b/frontend/src/app/workspace/component/power-button/computing-unit-selection.component.scss index 3913fec422..226751bca1 100644 --- a/frontend/src/app/workspace/component/power-button/computing-unit-selection.component.scss +++ b/frontend/src/app/workspace/component/power-button/computing-unit-selection.component.scss @@ -62,23 +62,10 @@ color: #1890ff; } - &.unit-disconnected { - color: #ff4d4f; - } - - &.unit-terminating { - color: #faad14; - } - &[disabled] { opacity: 0.6; cursor: not-allowed !important; - .unit-status-indicator { - opacity: 1; - font-weight: 500; - } - .computing-unit-terminate-icon { opacity: 1; cursor: pointer !important; @@ -114,20 +101,12 @@ white-space: nowrap; } -.computing-unit-name span, -.unit-status-indicator { +.computing-unit-name span { overflow: hidden; text-overflow: ellipsis; white-space: nowrap; } -.unit-status-indicator { - margin-left: 4px; - font-size: 0.85em; - font-style: italic; - opacity: 0.8; -} - .computing-unit-terminate-icon { margin-left: auto; flex-shrink: 0; diff --git a/frontend/src/app/workspace/component/power-button/computing-unit-selection.component.spec.ts b/frontend/src/app/workspace/component/power-button/computing-unit-selection.component.spec.ts index 432d97d379..705524b42f 100644 --- a/frontend/src/app/workspace/component/power-button/computing-unit-selection.component.spec.ts +++ b/frontend/src/app/workspace/component/power-button/computing-unit-selection.component.spec.ts @@ -74,6 +74,7 @@ function makeComputingUnit( uri: string; type: WorkflowComputingUnitType; status: string; + statusReason: string; isOwner: boolean; }> = {} ): DashboardWorkflowComputingUnit { @@ -83,6 +84,7 @@ function makeComputingUnit( uri = `uri-${cuid}`, type = "kubernetes", status = "Running", + statusReason = undefined, isOwner = true, } = overrides; return { @@ -104,6 +106,7 @@ function makeComputingUnit( }, }, status: status as DashboardWorkflowComputingUnit["status"], + statusReason, metrics: { cpuUsage: "N/A", memoryUsage: "N/A" }, isOwner, accessPrivilege: "WRITE", @@ -215,7 +218,12 @@ describe("PowerButtonComponent", () => { component.workflowId = 7; const selectSpy = vi.spyOn(component, "selectComputingUnit").mockImplementation(() => {}); const modal = fixture.debugElement.query(By.directive(ComputingUnitCreateModalComponent)).componentInstance; - modal.unitCreated.emit({ computingUnit: { cuid: 42 } } as unknown as DashboardWorkflowComputingUnit); + // A new unit is Pending. Pins that the row guard stays out of onPickComputingUnit, + // which would otherwise stop a new unit from being selected. + modal.unitCreated.emit({ + computingUnit: { cuid: 42 }, + status: "Pending", + } as unknown as DashboardWorkflowComputingUnit); expect(selectSpy).toHaveBeenCalledWith(7, 42); }); }); @@ -2063,15 +2071,9 @@ describe("PowerButtonComponent", () => { it("maps a status to a badge color", () => { expect(component.getBadgeColor("Running")).toBe("green"); expect(component.getBadgeColor("Pending")).toBe("gold"); + expect(component.getBadgeColor("Terminating")).toBe("gold"); expect(component.getBadgeColor("Failed")).toBe("red"); - }); - - it("describes a unit's status as a tooltip", () => { - expect(component.getUnitStatusTooltip(makeComputingUnit({ status: "Running" }))).toBe("Ready to use"); - expect(component.getUnitStatusTooltip(makeComputingUnit({ status: "Pending" }))).toBe( - "Computing unit is starting up" - ); - expect(component.getUnitStatusTooltip(makeComputingUnit({ status: "Failed" }))).toBe("Failed"); + expect(component.getBadgeColor("Unknown")).toBe("red"); }); }); @@ -2216,6 +2218,39 @@ describe("PowerButtonComponent", () => { expect(selectSpy).toHaveBeenCalledWith(5, 2); }); + it.each(["Failed", "Unknown", "Terminating", "Pending"] as const)( + "ignores a click on a %s row, which nz-menu only greys out", + async status => { + // nzDisabled does not stop the (click) on the same <li>, so our own guard must. + component.allComputingUnits = [ + makeComputingUnit({ cuid: 1, name: "Alpha" }), + makeComputingUnit({ cuid: 2, name: "Beta", status }), + ]; + fixture.detectChanges(); + // Spy on the only writer of the key, so the check holds even where Storage is unusable. + const rememberSpy = vi.spyOn(component as any, "rememberComputingUnit"); + const rows = await openDropdown(); + + click(rows[1]); + + expect(component.selectedComputingUnit).toBeNull(); + expect(selectSpy).not.toHaveBeenCalled(); + expect(rememberSpy).not.toHaveBeenCalled(); + expect(Object.keys(localStorage)).not.toContain("computing-unit-of-workflow-5"); + } + ); + + it("still selects a Running row, and remembers it", async () => { + const rememberSpy = vi.spyOn(component as any, "rememberComputingUnit"); + const rows = await openDropdown(); + + click(rows[0]); + + expect(component.selectedComputingUnit?.computingUnit.cuid).toBe(1); + expect(selectSpy).toHaveBeenCalledWith(5, 1); + expect(rememberSpy).toHaveBeenCalledWith(5, 1); + }); + it("commits a rename when the inline editor loses focus", async () => { const renameSpy = vi .spyOn(TestBed.inject(WorkflowComputingUnitManagingService), "renameComputingUnit") diff --git a/frontend/src/app/workspace/component/power-button/computing-unit-selection.component.ts b/frontend/src/app/workspace/component/power-button/computing-unit-selection.component.ts index da869a6b62..d196e6d762 100644 --- a/frontend/src/app/workspace/component/power-button/computing-unit-selection.component.ts +++ b/frontend/src/app/workspace/component/power-button/computing-unit-selection.component.ts @@ -50,7 +50,7 @@ import { memoryPercentage, validateName, getComputingUnitBadgeColor, - getComputingUnitStatusTooltip, + getComputingUnitRowTooltip, getComputingUnitCpuStatus, getComputingUnitMemoryStatus, getComputingUnitCpuLimitUnit, @@ -138,6 +138,8 @@ type PveDraft = { ], }) export class ComputingUnitSelectionComponent implements OnInit { + readonly getRowTooltip = getComputingUnitRowTooltip; + // variables for creating a virtual environment pves: PveDraft[] = []; systemPackages: { name: string; version: string }[] = []; @@ -456,6 +458,20 @@ export class ComputingUnitSelectionComponent implements OnInit { this.rememberComputingUnit(this.workflowId, cuid); } + /** + * Click handler for a dropdown row. nzDisabled only greys the row out; it does not stop this + * (click) on the same <li>, so without this check a disabled unit could still be selected. + * + * Not inside onPickComputingUnit on purpose: creating a unit also calls that, and a new unit is + * Pending, so a guard there would stop new units from being selected. + */ + public onClickComputingUnitRow(unit: DashboardWorkflowComputingUnit): void { + if (this.cannotSelectUnit(unit)) { + return; + } + this.onPickComputingUnit(unit); + } + /** * The live selection lives only in ComputingUnitStatusService, re-derived on load from the * last execution -- but that only exists once the workflow has run (pick a unit, reload @@ -878,13 +894,6 @@ export class ComputingUnitSelectionComponent implements OnInit { return getComputingUnitMemoryStatus(this.getMemoryPercentage()); } - /** - * Returns a descriptive tooltip for a specific unit's status - */ - getUnitStatusTooltip(unit: DashboardWorkflowComputingUnit): string { - return getComputingUnitStatusTooltip(unit); - } - public async onClickOpenShareAccess(cuid: number): Promise<void> { this.computingUnitActionsService.openShareAccessModal(cuid, true); } diff --git a/frontend/src/app/workspace/component/workflow-form/workflow-form.component.spec.ts b/frontend/src/app/workspace/component/workflow-form/workflow-form.component.spec.ts index 752816414f..ca5694bdef 100644 --- a/frontend/src/app/workspace/component/workflow-form/workflow-form.component.spec.ts +++ b/frontend/src/app/workspace/component/workflow-form/workflow-form.component.spec.ts @@ -1765,6 +1765,115 @@ describe("WorkflowFormComponent", () => { expect(component.runButtonState).toEqual({ label: "Connecting", icon: "loading", disabled: true }); }); + it("asks for a unit instead of a dead Stop when the selected unit vanishes mid-run", () => { + build(formViewWorkflow).ngOnInit(); + makeReady(); // a valid workflow runs on a unit the reader can write to + h.executionStateStream.next({ current: { state: ExecutionState.Running } }); // a run is in flight + h.statusStream.next(ComputingUnitState.NoComputingUnit); // the selected unit left the list + h.workflowWebsocketService.isConnected = false; // and its socket is gone + + // The run never looks finished, so the button must name the problem instead of a dead Stop. + expect(component.isRunning).toBe(true); + expect(component.runButtonState).toEqual({ label: "Computing Unit", icon: "plus-circle", disabled: true }); + }); + + it("still offers a deliverable Stop when the selected unit vanishes mid-run but its socket is up", () => { + build(formViewWorkflow).ngOnInit(); + h.statusStream.next(ComputingUnitState.NoComputingUnit); + h.executionStateStream.next({ current: { state: ExecutionState.Running } }); + h.workflowWebsocketService.isConnected = true; + + expect(component.runButtonState).toEqual({ label: "Stop", icon: "stop", disabled: false }); + }); + + it.each([ + [ComputingUnitState.Terminating, { label: "Shutting Down", icon: "loading", disabled: true }], + [ComputingUnitState.Failed, { label: "Unavailable", icon: "warning", disabled: true }], + [ComputingUnitState.Unknown, { label: "Unavailable", icon: "warning", disabled: true }], + ])("names a %s unit instead of a dead Stop when it dies mid-run", (state, expected) => { + build(formViewWorkflow).ngOnInit(); + makeReady(); // a valid workflow runs on a unit the reader can write to + h.executionStateStream.next({ current: { state: ExecutionState.Running } }); // a run is in flight + h.statusStream.next(state); // the unit dies + h.workflowWebsocketService.isConnected = false; // and its socket goes with it + + // The run never looks finished, so the button must name the problem instead of a dead Stop. + expect(component.isRunning).toBe(true); + expect(component.runButtonState).toEqual(expected); + }); + + it.each([ComputingUnitState.Terminating, ComputingUnitState.Failed, ComputingUnitState.Unknown])( + "still offers a deliverable Stop when a run is in flight on a %s unit whose socket is up", + state => { + build(formViewWorkflow).ngOnInit(); + h.statusStream.next(state); + h.executionStateStream.next({ current: { state: ExecutionState.Running } }); + h.workflowWebsocketService.isConnected = true; + + // The socket is still up, so Stop can be delivered and must stay enabled. + expect(component.runButtonState).toEqual({ label: "Stop", icon: "stop", disabled: false }); + } + ); + + it.each([false, true])( + "disables and says Shutting Down for a terminating unit, with the socket connected %s", + isConnected => { + build(formViewWorkflow).ngOnInit(); + makeReady(); + h.workflowWebsocketService.isConnected = isConnected; + h.statusStream.next(ComputingUnitState.Terminating); + + // A terminating unit is not coming back, so it must not show "Connecting". + expect(component.isConnecting).toBe(false); + expect(component.runButtonState).toEqual({ label: "Shutting Down", icon: "loading", disabled: true }); + } + ); + + it.each([ComputingUnitState.Failed, ComputingUnitState.Unknown])( + "disables and says Unavailable for a %s unit whose socket never comes up", + status => { + build(formViewWorkflow).ngOnInit(); + makeReady(); + h.workflowWebsocketService.isConnected = false; + h.statusStream.next(status); + + expect(component.isConnecting).toBe(false); + expect(component.runButtonState).toEqual({ label: "Unavailable", icon: "warning", disabled: true }); + } + ); + + it.each([ + { errors: { op: {} }, empty: false, label: "Invalid", icon: "warning" }, + { errors: {}, empty: true, label: "Empty", icon: "info-circle" }, + ])("keeps '$label' ahead of a terminal computing unit, as the canvas does", ({ errors, empty, label, icon }) => { + build(formViewWorkflow).ngOnInit(); + makeReady(); + h.statusStream.next(ComputingUnitState.Failed); + h.validationStream.next({ errors, workflowEmpty: empty }); + + expect(component.runButtonState).toEqual({ label, icon, disabled: true }); + }); + + it("names the terminal unit before the missing warehouse, which cannot rescue a dead unit", () => { + build(formViewWorkflow).ngOnInit(); + makeReady(); + h.config.env.warehouseEnabled = true; + h.warehouseService.selectWarehouse(undefined); + h.statusStream.next(ComputingUnitState.Failed); + + expect(component.runButtonState).toEqual({ label: "Unavailable", icon: "warning", disabled: true }); + }); + + it("still says Connecting for a Pending unit, which is starting up rather than dead", () => { + build(formViewWorkflow).ngOnInit(); + makeReady(); + h.workflowWebsocketService.isConnected = false; + h.statusStream.next(ComputingUnitState.Pending); + + expect(component.isConnecting).toBe(true); + expect(component.runButtonState).toEqual({ label: "Connecting", icon: "loading", disabled: true }); + }); + it("repaints when the websocket connection status changes", () => { build(formViewWorkflow).ngOnInit(); h.cdr.markForCheck.mockClear(); diff --git a/frontend/src/app/workspace/component/workflow-form/workflow-form.component.ts b/frontend/src/app/workspace/component/workflow-form/workflow-form.component.ts index 9b06c05273..4fbd753aad 100644 --- a/frontend/src/app/workspace/component/workflow-form/workflow-form.component.ts +++ b/frontend/src/app/workspace/component/workflow-form/workflow-form.component.ts @@ -40,6 +40,7 @@ import { EditableLabelWrapperComponent } from "../../../common/formly/editable-l import { FormFieldBinding, Workflow, WorkflowContent } from "../../../common/type/workflow"; import { ComputingUnitStatusService } from "../../../common/service/computing-unit/computing-unit-status/computing-unit-status.service"; import { ComputingUnitState } from "../../../common/type/computing-unit-connection.interface"; +import { unavailableComputingUnitReason } from "../../../common/util/computing-unit.util"; import { DashboardWorkflowComputingUnit } from "../../../common/type/workflow-computing-unit"; import { WorkflowPersistService } from "../../../common/service/workflow-persist/workflow-persist.service"; import { NotificationService } from "../../../common/service/notification/notification.service"; @@ -1562,10 +1563,14 @@ export class WorkflowFormComponent implements OnInit, OnDestroy { * A unit is picked but its socket is still coming up -- the same window the operator canvas shows * "Connecting" and disables its run button. Read from the exact condition the canvas uses * (menu.component's getRunButtonBehavior), so the two stay in step. + * + * Terminal units are excluded: they are not coming back, so runButtonState names them instead. */ public get isConnecting(): boolean { return ( - this.computingUnitStatus !== ComputingUnitState.NoComputingUnit && !this.workflowWebsocketService.isConnected + this.computingUnitStatus !== ComputingUnitState.NoComputingUnit && + unavailableComputingUnitReason(this.computingUnitStatus) === undefined && + !this.workflowWebsocketService.isConnected ); } @@ -1596,6 +1601,11 @@ export class WorkflowFormComponent implements OnInit, OnDestroy { * to Run and Stop, with no pause/resume: while a run is in flight the button stops (kills) it, * otherwise it runs. (The canvas offers Pause/Resume/Submitting and a clickable Connect; a form * reader does not, and picks the unit in the embedded selector instead.) + * + * A unit that cannot accept work also disables it, named as on the canvas via + * unavailableComputingUnitReason but with shorter labels. One difference: mid-run the canvas shows + * "Shutting Down", but here Stop wins while the socket can still deliver it, because Stop is this + * button's only run control. */ public get runButtonState(): { label: string; icon: string; disabled: boolean } { // Connecting is checked before Stop on purpose: if the socket drops mid-run, a "Stop" would @@ -1604,7 +1614,9 @@ export class WorkflowFormComponent implements OnInit, OnDestroy { if (this.isConnecting) { return { label: "Connecting", icon: "loading", disabled: true }; } - if (this.isRunning) { + // Stop is shown only while the socket can deliver the kill. A run on a unit that died or + // vanished mid-run falls through, so a later branch names the problem. + if (this.isRunning && this.workflowWebsocketService.isConnected) { return { label: "Stop", icon: "stop", disabled: false }; } if (!this.isWorkflowValid) { @@ -1613,6 +1625,14 @@ export class WorkflowFormComponent implements OnInit, OnDestroy { if (this.isWorkflowEmpty) { return { label: "Empty", icon: "info-circle", disabled: true }; } + // Before the access and warehouse checks: neither can fix a dead unit. + const unavailableReason = unavailableComputingUnitReason(this.computingUnitStatus); + if (unavailableReason === "terminating") { + return { label: "Shutting Down", icon: "loading", disabled: true }; + } + if (unavailableReason === "unavailable") { + return { label: "Unavailable", icon: "warning", disabled: true }; + } if (this.hasNoComputingUnit) { return { label: "Computing Unit", icon: "plus-circle", disabled: true }; } diff --git a/frontend/src/app/workspace/service/execute-workflow/execute-workflow.service.spec.ts b/frontend/src/app/workspace/service/execute-workflow/execute-workflow.service.spec.ts index b3bdbc1ec7..650a2d32a6 100644 --- a/frontend/src/app/workspace/service/execute-workflow/execute-workflow.service.spec.ts +++ b/frontend/src/app/workspace/service/execute-workflow/execute-workflow.service.spec.ts @@ -42,6 +42,7 @@ import { WorkflowUtilService } from "../workflow-graph/util/workflow-util.servic import { WorkflowSettings } from "src/app/common/type/workflow"; import { ComputingUnitStatusService } from "../../../common/service/computing-unit/computing-unit-status/computing-unit-status.service"; +import { DashboardWorkflowComputingUnit } from "../../../common/type/workflow-computing-unit"; import { WarehouseService } from "../../../common/service/warehouse/warehouse.service"; import { AuthService } from "src/app/common/service/user/auth.service"; import { StubAuthService } from "src/app/common/service/user/stub-auth.service"; @@ -565,6 +566,109 @@ describe("ExecuteWorkflowService", () => { } })); + // ---- refusing to run on a unit that cannot accept work ------------------------------------- + + describe("refusing to run on a computing unit that cannot accept work", () => { + // Run-up-to and Time Travel replay bypass the run buttons and call this service directly, + // so the service itself must refuse, before it resets the results on screen. + const selectUnitWithStatus = (status: DashboardWorkflowComputingUnit["status"]): void => { + vi.spyOn(service["computingUnitStatusService"], "getSelectedComputingUnitValue").mockReturnValue({ + computingUnit: { cuid: 7 }, + status, + } as unknown as DashboardWorkflowComputingUnit); + }; + + const replayInfo: ReplayExecutionInfo = { eid: 42, interaction: "step-3" }; + + const entryPoints: { name: string; run: () => void }[] = [ + { + name: "executeWorkflowWithEmailNotification", + run: () => service.executeWorkflowWithEmailNotification("e", false), + }, + { name: "executeWorkflowWithReplay", run: () => service.executeWorkflowWithReplay(replayInfo) }, + ]; + + for (const entryPoint of entryPoints) { + it(`${entryPoint.name} refuses a terminating unit and keeps the previous results on screen`, fakeAsync(() => { + selectUnitWithStatus("Terminating"); + const wsSendSpy = vi.spyOn(service["workflowWebsocketService"], "send"); + const resetSpy = vi.spyOn(service, "resetExecutionState"); + const statusResetSpy = vi.spyOn(TestBed.inject(WorkflowStatusService), "resetStatus"); + const errorSpy = vi.spyOn(TestBed.inject(NotificationService), "error").mockReturnValue(undefined as never); + + entryPoint.run(); + tick(FORM_DEBOUNCE_TIME_MS + 1); + flush(); + + expect(wsSendSpy).not.toHaveBeenCalledWith("WorkflowExecuteRequest", expect.anything()); + expect(resetSpy).not.toHaveBeenCalled(); + expect(statusResetSpy).not.toHaveBeenCalled(); + expect(errorSpy).toHaveBeenCalledWith( + "The selected computing unit is shutting down. Wait for it to finish, then select or create another one." + ); + })); + + it.each(["Failed", "Unknown"] as const)( + `${entryPoint.name} refuses a %s unit and keeps the previous results on screen`, + status => { + selectUnitWithStatus(status); + const wsSendSpy = vi.spyOn(service["workflowWebsocketService"], "send"); + const resetSpy = vi.spyOn(service, "resetExecutionState"); + const statusResetSpy = vi.spyOn(TestBed.inject(WorkflowStatusService), "resetStatus"); + const errorSpy = vi.spyOn(TestBed.inject(NotificationService), "error").mockReturnValue(undefined as never); + + entryPoint.run(); + + expect(wsSendSpy).not.toHaveBeenCalledWith("WorkflowExecuteRequest", expect.anything()); + expect(resetSpy).not.toHaveBeenCalled(); + expect(statusResetSpy).not.toHaveBeenCalled(); + expect(errorSpy).toHaveBeenCalledWith( + "The selected computing unit is unavailable. Select a running unit or create a new one." + ); + } + ); + + it.each(["Running", "Pending"] as const)(`${entryPoint.name} still runs on a %s unit`, status => { + // Pending is still starting up, so it must not be blocked. + selectUnitWithStatus(status); + const sendExecutionRequestSpy = vi.spyOn(service, "sendExecutionRequest").mockImplementation(() => {}); + const errorSpy = vi.spyOn(TestBed.inject(NotificationService), "error").mockReturnValue(undefined as never); + + entryPoint.run(); + + expect(sendExecutionRequestSpy).toHaveBeenCalled(); + expect(errorSpy).not.toHaveBeenCalled(); + }); + } + + it("refuses on the unit before the warehouse, since a warehouse cannot rescue a dead unit", () => { + TestBed.inject(GuiConfigService).env.warehouseEnabled = true; + try { + TestBed.inject(WarehouseService).selectWarehouse(undefined); + selectUnitWithStatus("Failed"); + const errorSpy = vi.spyOn(TestBed.inject(NotificationService), "error").mockReturnValue(undefined as never); + + service.executeWorkflowWithEmailNotification("e", false); + + expect(errorSpy).toHaveBeenCalledTimes(1); + expect(errorSpy).toHaveBeenCalledWith( + "The selected computing unit is unavailable. Select a running unit or create a new one." + ); + } finally { + TestBed.inject(GuiConfigService).env.warehouseEnabled = false; + } + }); + + it("runs normally when no unit is selected at all, leaving that to the existing warning", () => { + // No unit selected is not a terminal state; sendExecutionRequest already warns about it. + const sendExecutionRequestSpy = vi.spyOn(service, "sendExecutionRequest").mockImplementation(() => {}); + + service.executeWorkflowWithEmailNotification("e", false); + + expect(sendExecutionRequestSpy).toHaveBeenCalled(); + }); + }); + it("sendExecutionRequest carries the picked warehouse id, and none when unset (#7817)", fakeAsync(() => { const warehouseService = TestBed.inject(WarehouseService); const wsSendSpy = vi.spyOn(service["workflowWebsocketService"], "send"); diff --git a/frontend/src/app/workspace/service/execute-workflow/execute-workflow.service.ts b/frontend/src/app/workspace/service/execute-workflow/execute-workflow.service.ts index c3633721d7..a321381681 100644 --- a/frontend/src/app/workspace/service/execute-workflow/execute-workflow.service.ts +++ b/frontend/src/app/workspace/service/execute-workflow/execute-workflow.service.ts @@ -48,6 +48,7 @@ import { intersection } from "../../../common/util/set"; import { WorkflowSettings } from "../../../common/type/workflow"; import { ComputingUnitStatusService } from "../../../common/service/computing-unit/computing-unit-status/computing-unit-status.service"; +import { unavailableComputingUnitReason } from "../../../common/util/computing-unit.util"; import { WarehouseService } from "../../../common/service/warehouse/warehouse.service"; import { GuiConfigService } from "../../../common/service/gui-config.service"; @@ -272,6 +273,9 @@ export class ExecuteWorkflowService { targetOperatorId ); const settings = this.workflowActionService.getWorkflowSettings(); + if (this.refuseToRunOnUnavailableUnit()) { + return; + } if (this.refuseToRunWithoutWarehouse()) { return; } @@ -287,6 +291,9 @@ export class ExecuteWorkflowService { public executeWorkflowWithReplay(replayExecutionInfo: ReplayExecutionInfo): void { const logicalPlan = ExecuteWorkflowService.getLogicalPlanRequest(this.workflowActionService.getTexeraGraph()); const settings = this.workflowActionService.getWorkflowSettings(); + if (this.refuseToRunOnUnavailableUnit()) { + return; + } if (this.refuseToRunWithoutWarehouse()) { return; } @@ -301,6 +308,29 @@ export class ExecuteWorkflowService { ); } + /** + * Refuses to run on a unit that cannot accept work: shows a toast and returns true. The run + * buttons already disable themselves, but run-up-to and Time Travel replay call this service + * directly. + * + * Checked before resetExecutionState(), so a refused click keeps the results on screen, and + * before the warehouse check, since a warehouse cannot fix a dead unit. + */ + private refuseToRunOnUnavailableUnit(): boolean { + const reason = unavailableComputingUnitReason( + this.computingUnitStatusService.getSelectedComputingUnitValue()?.status + ); + if (reason === undefined) { + return false; + } + this.notificationService.error( + reason === "terminating" + ? "The selected computing unit is shutting down. Wait for it to finish, then select or create another one." + : "The selected computing unit is unavailable. Select a running unit or create a new one." + ); + return true; + } + /** * While the deployment requires a warehouse (#7817) and none is picked, * refuses with a toast and returns true. Checked at every public entry
