This is an automated email from the ASF dual-hosted git repository.
github-merge-queue[bot] pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/texera.git
The following commit(s) were added to refs/heads/main by this push:
new 33bd07bf27 feat(computing-unit): surface failed/unhealthy computing
units instead of an endless "Connecting" (#7944)
33bd07bf27 is described below
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