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 14de4583fd fix(kubernetes): terminate idle computing units (#6046)
14de4583fd is described below
commit 14de4583fd8078a8639f458f0eb8e34a3490a9a7
Author: Yichen Ren <[email protected]>
AuthorDate: Fri Sep 25 22:41:33 2026 +0000
fix(kubernetes): terminate idle computing units (#6046)
### What changes were proposed in this PR?
Following discussion https://github.com/apache/texera/discussions/6264,
this PR adds backend-side cleanup for idle Kubernetes computing units.
The main change is a scheduled cleanup task in the computing unit
managing service that periodically scans active Kubernetes computing
units and terminates units that have been inactive longer than a
configurable timeout.
The implementation includes the following changes:
- Added new Kubernetes configuration entries for:
- computing unit idle timeout
- computing unit idle check interval
- Exposed both settings through environment-variable-based configuration
so deployment-side overrides can be applied without code changes.
- Added a scheduled background task in `ComputingUnitManagingService`
that runs the idle cleanup logic at a fixed interval.
- Added idle Kubernetes computing unit termination logic in
`ComputingUnitManagingResource`:
- only considers Kubernetes computing units that are not already
terminated
- checks whether the computing unit has any active workflow executions
- computes the latest execution activity timestamp from existing
execution metadata
- terminates the Kubernetes pod when the computing unit is considered
idle past the configured timeout
- updates the computing unit termination time in the database after
cleanup
The timeout and check interval are configurable through environment
variables, so the behavior can be tuned for different deployment or
testing needs without modifying the code.
### Any related issues, documentation, discussions?
Fixes #5362
### How was this PR tested?
Tested locally on the Kubernetes deployment flow.
https://github.com/user-attachments/assets/98e30808-49f7-4397-a6ae-3cc536e8e583
### Was this PR authored or co-authored using generative AI tooling?
Generated-by: OpenAI Codex GPT-5
Co-authored-by: Claude Opus 5 <[email protected]>
---
bin/k8s/values-development.yaml | 8 +
bin/k8s/values.yaml | 8 +
common/config/src/main/resources/kubernetes.conf | 13 +
.../texera/common/config/KubernetesConfig.scala | 9 +
.../common/config/KubernetesConfigSpec.scala | 13 +
.../service/ComputingUnitManagingService.scala | 53 +-
.../resource/ComputingUnitManagingResource.scala | 257 +++++++++-
.../texera/service/util/ComputingUnitHelpers.scala | 10 +-
.../service/util/IdleComputingUnitCleanupJob.scala | 107 ++++
.../ComputingUnitIdleCleanupSchedulerSpec.scala | 107 ++++
.../resource/ComputingUnitIdleCleanupSpec.scala | 544 +++++++++++++++++++++
.../service/util/ComputingUnitHelpersSpec.scala | 3 +
.../util/IdleComputingUnitCleanupJobSpec.scala | 109 +++++
sql/changelog.xml | 5 +
sql/texera_ddl.sql | 6 +
sql/updates/51.sql | 45 ++
16 files changed, 1287 insertions(+), 10 deletions(-)
diff --git a/bin/k8s/values-development.yaml b/bin/k8s/values-development.yaml
index 0216a44239..5cbda8ee0d 100644
--- a/bin/k8s/values-development.yaml
+++ b/bin/k8s/values-development.yaml
@@ -360,6 +360,14 @@ texeraEnvVars:
value: "false"
- name: KUBERNETES_COMPUTING_UNIT_ENABLED
value: "true"
+ # Periodic termination of idle Kubernetes computing units. Off by default;
set the timeout to
+ # match how long a deployment's users expect an untouched unit to stay alive
before enabling.
+ - name: KUBERNETES_COMPUTING_UNIT_IDLE_CLEANUP_ENABLED
+ value: "false"
+ - name: KUBERNETES_COMPUTING_UNIT_IDLE_TIMEOUT_MINUTES
+ value: "1440"
+ - name: KUBERNETES_COMPUTING_UNIT_IDLE_CHECK_INTERVAL_MINUTES
+ value: "60"
- name: KUBERNETES_IMAGE_PULL_POLICY
value: "IfNotPresent"
- name: GUI_WORKFLOW_WORKSPACE_PYTHON_LANGUAGE_SERVER_PORT
diff --git a/bin/k8s/values.yaml b/bin/k8s/values.yaml
index a0ac347118..c72120a2e2 100644
--- a/bin/k8s/values.yaml
+++ b/bin/k8s/values.yaml
@@ -461,6 +461,14 @@ texeraEnvVars:
value: "false"
- name: KUBERNETES_COMPUTING_UNIT_ENABLED
value: "true"
+ # Periodic termination of idle Kubernetes computing units. Off by default;
set the timeout to
+ # match how long a deployment's users expect an untouched unit to stay alive
before enabling.
+ - name: KUBERNETES_COMPUTING_UNIT_IDLE_CLEANUP_ENABLED
+ value: "false"
+ - name: KUBERNETES_COMPUTING_UNIT_IDLE_TIMEOUT_MINUTES
+ value: "1440"
+ - name: KUBERNETES_COMPUTING_UNIT_IDLE_CHECK_INTERVAL_MINUTES
+ value: "60"
- name: KUBERNETES_IMAGE_PULL_POLICY
value: "IfNotPresent"
- name: GUI_WORKFLOW_WORKSPACE_PYTHON_LANGUAGE_SERVER_PORT
diff --git a/common/config/src/main/resources/kubernetes.conf
b/common/config/src/main/resources/kubernetes.conf
index 6278679a7c..e3977ae91f 100644
--- a/common/config/src/main/resources/kubernetes.conf
+++ b/common/config/src/main/resources/kubernetes.conf
@@ -48,6 +48,19 @@ kubernetes {
max-num-of-running-computing-units-per-user = 10
max-num-of-running-computing-units-per-user =
${?MAX_NUM_OF_RUNNING_COMPUTING_UNITS_PER_USER}
+ # Periodically terminate Kubernetes CUs that have been idle past the timeout
below. Off by
+ # default: the sweep deletes pods on a timer, so a deployment opts in only
once its idle timeout
+ # has been validated against how its users actually work.
+ computing-unit-idle-cleanup-enabled = false
+ computing-unit-idle-cleanup-enabled =
${?KUBERNETES_COMPUTING_UNIT_IDLE_CLEANUP_ENABLED}
+
+ # Terminate Kubernetes CUs whose latest workflow execution is older than
this.
+ computing-unit-idle-timeout-minutes = 1440
+ computing-unit-idle-timeout-minutes =
${?KUBERNETES_COMPUTING_UNIT_IDLE_TIMEOUT_MINUTES}
+
+ computing-unit-idle-check-interval-minutes = 60
+ computing-unit-idle-check-interval-minutes =
${?KUBERNETES_COMPUTING_UNIT_IDLE_CHECK_INTERVAL_MINUTES}
+
computing-unit-cpu-limit-options = "1,2,4"
computing-unit-cpu-limit-options =
${?KUBERNETES_COMPUTING_UNIT_CPU_LIMIT_OPTIONS}
diff --git
a/common/config/src/main/scala/org/apache/texera/common/config/KubernetesConfig.scala
b/common/config/src/main/scala/org/apache/texera/common/config/KubernetesConfig.scala
index 7b56c41864..e537d3f19a 100644
---
a/common/config/src/main/scala/org/apache/texera/common/config/KubernetesConfig.scala
+++
b/common/config/src/main/scala/org/apache/texera/common/config/KubernetesConfig.scala
@@ -39,6 +39,15 @@ object KubernetesConfig {
val maxNumOfRunningComputingUnitsPerUser: Int =
conf.getInt("kubernetes.max-num-of-running-computing-units-per-user")
+ val computingUnitIdleCleanupEnabled: Boolean =
+ conf.getBoolean("kubernetes.computing-unit-idle-cleanup-enabled")
+
+ val computingUnitIdleTimeoutMinutes: Long =
+ conf.getLong("kubernetes.computing-unit-idle-timeout-minutes")
+
+ val computingUnitIdleCheckIntervalMinutes: Long =
+ conf.getLong("kubernetes.computing-unit-idle-check-interval-minutes")
+
val cpuLimitOptions: List[String] =
conf
.getString("kubernetes.computing-unit-cpu-limit-options")
diff --git
a/common/config/src/test/scala/org/apache/texera/common/config/KubernetesConfigSpec.scala
b/common/config/src/test/scala/org/apache/texera/common/config/KubernetesConfigSpec.scala
index 7dbeedd07c..a9d413c0c5 100644
---
a/common/config/src/test/scala/org/apache/texera/common/config/KubernetesConfigSpec.scala
+++
b/common/config/src/test/scala/org/apache/texera/common/config/KubernetesConfigSpec.scala
@@ -71,6 +71,19 @@ class KubernetesConfigSpec extends AnyFlatSpec with Matchers
{
KubernetesConfig.maxNumOfRunningComputingUnitsPerUser should be >= 0
}
+ "KubernetesConfig idle computing unit cleanup settings" should "resolve to
their kubernetes.conf defaults" in {
+ // The sweep deletes pods on a timer, so it stays off until a deployment
opts in.
+ ifUnset("KUBERNETES_COMPUTING_UNIT_IDLE_CLEANUP_ENABLED")(
+ KubernetesConfig.computingUnitIdleCleanupEnabled shouldBe false
+ )
+ ifUnset("KUBERNETES_COMPUTING_UNIT_IDLE_TIMEOUT_MINUTES")(
+ KubernetesConfig.computingUnitIdleTimeoutMinutes shouldBe 1440
+ )
+ ifUnset("KUBERNETES_COMPUTING_UNIT_IDLE_CHECK_INTERVAL_MINUTES")(
+ KubernetesConfig.computingUnitIdleCheckIntervalMinutes shouldBe 60
+ )
+ }
+
"KubernetesConfig jupyter settings" should "resolve to their kubernetes.conf
defaults" in {
KubernetesConfig.jupyterPortNumber shouldBe 8888
// Off by default and keyed separately from kubernetes.enabled, so
enabling computing
diff --git
a/computing-unit-managing-service/src/main/scala/org/apache/texera/service/ComputingUnitManagingService.scala
b/computing-unit-managing-service/src/main/scala/org/apache/texera/service/ComputingUnitManagingService.scala
index 97c9159184..785c719fe3 100644
---
a/computing-unit-managing-service/src/main/scala/org/apache/texera/service/ComputingUnitManagingService.scala
+++
b/computing-unit-managing-service/src/main/scala/org/apache/texera/service/ComputingUnitManagingService.scala
@@ -23,7 +23,7 @@ import com.fasterxml.jackson.module.scala.DefaultScalaModule
import io.dropwizard.configuration.{EnvironmentVariableSubstitutor,
SubstitutingSourceProvider}
import io.dropwizard.core.Application
import io.dropwizard.core.setup.{Bootstrap, Environment}
-import org.apache.texera.common.config.StorageConfig
+import org.apache.texera.common.config.{KubernetesConfig, StorageConfig}
import org.apache.texera.auth.{AuthFeatures, RequestLoggingFilter,
RoleAnnotationEnforcer}
import org.apache.texera.dao.SqlServer
import org.apache.texera.service.resource.{
@@ -33,9 +33,45 @@ import org.apache.texera.service.resource.{
CuratedImageResource,
HealthCheckResource
}
+import org.apache.texera.service.util.IdleComputingUnitCleanupJob
+import org.slf4j.LoggerFactory
import java.nio.file.Path
class ComputingUnitManagingService extends
Application[ComputingUnitManagingServiceConfiguration] {
+ private val logger =
LoggerFactory.getLogger(classOf[ComputingUnitManagingService])
+
+ private def initSqlServer(): Unit =
+ SqlServer.initConnection(
+ StorageConfig.jdbcUrl,
+ StorageConfig.jdbcUsername,
+ StorageConfig.jdbcPassword
+ )
+
+ /**
+ * Registers the periodic idle computing unit cleanup job on the
application lifecycle when
+ * enabled. Extracted from `run` (and kept free of any global config reads)
so the conditional
+ * wiring can be unit-tested with a standalone `Environment`.
+ */
+ private[service] def registerIdleComputingUnitCleanup(
+ environment: Environment,
+ enabled: Boolean,
+ idleTimeoutMinutes: Long,
+ intervalMinutes: Long
+ ): Unit =
+ if (enabled) {
+ // The job's scheduler rejects a non-positive delay, and a misconfigured
sweep should leave
+ // the rest of the service usable, so log it and skip rather than abort
startup.
+ if (idleTimeoutMinutes <= 0 || intervalMinutes <= 0) {
+ logger.warn(
+ s"Idle Kubernetes computing unit cleanup is disabled: timeout and
check interval must " +
+ s"both be positive but are $idleTimeoutMinutes and
$intervalMinutes minute(s)"
+ )
+ } else {
+ environment
+ .lifecycle()
+ .manage(new IdleComputingUnitCleanupJob(idleTimeoutMinutes,
intervalMinutes))
+ }
+ }
override def initialize(
bootstrap: Bootstrap[ComputingUnitManagingServiceConfiguration]
@@ -60,11 +96,7 @@ class ComputingUnitManagingService extends
Application[ComputingUnitManagingServ
AuthFeatures.register(environment)
- SqlServer.initConnection(
- StorageConfig.jdbcUrl,
- StorageConfig.jdbcUsername,
- StorageConfig.jdbcPassword
- )
+ initSqlServer()
environment.jersey().register(new ComputingUnitManagingResource)
environment.jersey().register(new ComputingUnitAccessResource)
@@ -76,6 +108,15 @@ class ComputingUnitManagingService extends
Application[ComputingUnitManagingServ
"ComputingUnitManagingService"
)
+ // Periodically terminate Kubernetes computing units their owners have
stopped using
+ registerIdleComputingUnitCleanup(
+ environment,
+ KubernetesConfig.kubernetesComputingUnitEnabled &&
+ KubernetesConfig.computingUnitIdleCleanupEnabled,
+ KubernetesConfig.computingUnitIdleTimeoutMinutes,
+ KubernetesConfig.computingUnitIdleCheckIntervalMinutes
+ )
+
// Route request logs through SLF4J, controlled by TEXERA_SERVICE_LOG_LEVEL
RequestLoggingFilter.register(environment.getApplicationContext)
}
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 a281a0768e..2b6b55a667 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
@@ -42,9 +42,15 @@ import org.apache.texera.common.config.{
}
import org.apache.texera.dao.SqlServer
import org.apache.texera.dao.SqlServer.withTransaction
+import org.apache.texera.dao.jooq.generated.Tables.{
+ USER,
+ WORKFLOW_COMPUTING_UNIT,
+ WORKFLOW_EXECUTIONS
+}
import org.apache.texera.dao.jooq.generated.enums.{
PrivilegeEnum,
UserRoleEnum,
+ WorkflowComputingUnitTerminationReasonEnum,
WorkflowComputingUnitTypeEnum
}
import org.apache.texera.dao.jooq.generated.tables.daos.{
@@ -60,19 +66,250 @@ import org.apache.texera.service.util.{
InsufficientComputingUnitQuota,
KubernetesClient
}
-import org.jooq.{DSLContext, EnumType}
+import org.jooq.{Condition, DSLContext, EnumType}
+import org.jooq.impl.DSL.{boolOr, exists, max, selectOne}
+import org.slf4j.LoggerFactory
import play.api.libs.json._
import java.sql.Timestamp
import scala.annotation.unused
import scala.jdk.CollectionConverters.CollectionHasAsScala
+import scala.util.control.NonFatal
object ComputingUnitManagingResource {
+ private[resource] val logger =
LoggerFactory.getLogger(classOf[ComputingUnitManagingResource])
+
private def context: DSLContext =
SqlServer
.getInstance()
.createDSLContext()
+ private[resource] case class IdleComputingUnitCandidate(
+ unit: WorkflowComputingUnit,
+ username: Option[String]
+ )
+
+ /**
+ * The codes persisted in `workflow_executions.status`. These are the
collapsed codes produced
+ * by amber's `Utils.maptoStatusCode`, NOT the ordinals of
`WorkflowAggregatedState`, and this
+ * service cannot depend on the amber module to reuse either. Keep this in
sync with
+ * `Utils.maptoStatusCode`.
+ *
+ * Only the terminal codes are listed, and the sweep treats every other
code as an execution
+ * still in flight. That direction matters: `maptoStatusCode` collapses
PAUSING, RESUMING,
+ * UNKNOWN and TERMINATED alike to -1, so enumerating the non-terminal
codes instead would
+ * leave a paused-or-resuming execution looking idle and get its computing
unit deleted out
+ * from under its owner.
+ *
+ * Known limitation: a non-terminal code is trusted without a time bound,
so a row that never
+ * reaches 3/4/5 keeps its computing unit off this sweep indefinitely --
not for one more
+ * sweep, but permanently. Two ways in:
+ * - a row left stuck at RUNNING -- a coordinator that died in place, an
OOM inside the
+ * container -- because nothing ever rewrites it;
+ * - a row whose *final* status is -1, because `maptoStatusCode` gives an
execution that
+ * ended the same code as one that is merely paused. Today that is
UNKNOWN, the fallback
+ * `WorkflowExecution.getState` returns for a worker mix it cannot
name. TERMINATED maps
+ * to -1 as well, but is a worker-level state:
`ExecutionUtils.aggregateStates` reports a
+ * COMPLETED-or-TERMINATED worker set as COMPLETED, so it does not
reach this column
+ * today. If amber ever persists it, it lands in this same bucket.
+ * Neither is separable here -- a -1 carries nothing that distinguishes an
ended execution
+ * from a live one, so the fix belongs in what amber persists. Nothing else
reclaims such a
+ * unit either: `ComputingUnitHelpers.reconcileVanishedKubernetesUnits`
only runs when someone
+ * calls a listing endpoint, and it keys off a vanished pod rather than a
stale execution row.
+ * Tracked as a follow-up in apache/texera#8618.
+ */
+ private[resource] object TerminalWorkflowExecutionStatus extends Enumeration
{
+ val Completed: Value = Value(3)
+ val Failed: Value = Value(4)
+ val Killed: Value = Value(5)
+
+ def dbStatuses: Seq[java.lang.Short] =
+ values.toSeq.map(status => Short.box(status.id.toShort))
+ }
+
+ private[resource] def lastComputingUnitActivityTime(
+ unit: WorkflowComputingUnit,
+ latestUpdateTime: Option[Timestamp],
+ latestStartTime: Option[Timestamp]
+ ): Timestamp =
+ Seq(
+ latestUpdateTime,
+ latestStartTime,
+ Option(unit.getCreationTime)
+ ).flatten.maxBy(_.getTime)
+
+ private[resource] def shouldTerminateIdleComputingUnit(
+ hasActiveExecution: Boolean,
+ lastExecutionTime: Timestamp,
+ cutoff: Timestamp
+ ): Boolean =
+ !hasActiveExecution && lastExecutionTime.before(cutoff)
+
+ /**
+ * Terminates every Kubernetes computing unit whose last execution activity
is older than
+ * `idleTimeoutMinutes`, returning the units terminated so the caller can
log their owners.
+ */
+ def terminateIdleKubernetesComputingUnits(
+ idleTimeoutMinutes: Long
+ ): List[TerminatedComputingUnitInfo] =
+ terminateIdleKubernetesComputingUnits(
+ idleTimeoutMinutes,
+ new Timestamp(System.currentTimeMillis()),
+ KubernetesClient
+ )
+
+ /**
+ * The client is a by-name parameter -- not the global singleton -- so the
sweep is unit-testable
+ * with a stub and a sweep that terminates nothing never forces the
singleton; the public
+ * overload binds the production [[KubernetesClient]]. Same seam as
+ * [[org.apache.texera.service.util.ComputingUnitHelpers.singleUnitStatus]].
+ */
+ private[resource] def terminateIdleKubernetesComputingUnits(
+ idleTimeoutMinutes: Long,
+ now: Timestamp,
+ k8s: => KubernetesClient
+ ): List[TerminatedComputingUnitInfo] = {
+ val cutoff = new Timestamp(now.getTime - idleTimeoutMinutes * 60 * 1000)
+ idleKubernetesComputingUnitCandidates(cutoff).flatMap(candidate =>
+ terminateIdleKubernetesComputingUnitCandidate(candidate, cutoff, now,
k8s)
+ )
+ }
+
+ private[resource] def idleKubernetesComputingUnitCandidates(
+ cutoff: Timestamp
+ ): List[IdleComputingUnitCandidate] = {
+ // All three questions asked per computing unit -- is any execution still
active, when did an
+ // execution last report progress, when did one last start -- are
aggregates over the same rows
+ // grouped by the same key, so one grouped query answers them for every
unit at once. The left
+ // joins keep units that have no executions (both max() are NULL) and
units whose owner row is
+ // gone (name is NULL), matching what a per-unit scan would produce.
+ val latestUpdateTime = max(WORKFLOW_EXECUTIONS.LAST_UPDATE_TIME)
+ val latestStartTime = max(WORKFLOW_EXECUTIONS.STARTING_TIME)
+ val hasActiveExecution =
+
boolOr(WORKFLOW_EXECUTIONS.STATUS.notIn(TerminalWorkflowExecutionStatus.dbStatuses:
_*))
+
+ withTransaction(context) { ctx =>
+ ctx
+ .select(
+ WORKFLOW_COMPUTING_UNIT.asterisk(),
+ USER.NAME,
+ latestUpdateTime,
+ latestStartTime,
+ hasActiveExecution
+ )
+ .from(WORKFLOW_COMPUTING_UNIT)
+ .leftJoin(WORKFLOW_EXECUTIONS)
+ .on(WORKFLOW_EXECUTIONS.CUID.eq(WORKFLOW_COMPUTING_UNIT.CUID))
+ .leftJoin(USER)
+ .on(USER.UID.eq(WORKFLOW_COMPUTING_UNIT.UID))
+ .where(
+ WORKFLOW_COMPUTING_UNIT.TYPE
+ .eq(WorkflowComputingUnitTypeEnum.kubernetes)
+ .and(WORKFLOW_COMPUTING_UNIT.TERMINATE_TIME.isNull)
+ )
+ .groupBy(WORKFLOW_COMPUTING_UNIT.CUID, USER.NAME)
+ .fetch()
+ .asScala
+ .flatMap { record =>
+ val unit =
record.into(WORKFLOW_COMPUTING_UNIT).into(classOf[WorkflowComputingUnit])
+ val lastExecutionTime = lastComputingUnitActivityTime(
+ unit,
+ Option(record.get(latestUpdateTime)),
+ Option(record.get(latestStartTime))
+ )
+
+ // bool_or over zero matching executions yields NULL, which means
"no active execution"
+ val active =
Option(record.get(hasActiveExecution)).exists(_.booleanValue())
+ if (shouldTerminateIdleComputingUnit(active, lastExecutionTime,
cutoff)) {
+ Some(
+ IdleComputingUnitCandidate(
+ unit,
+ Option(record.get(USER.NAME)).filter(_.nonEmpty)
+ )
+ )
+ } else {
+ None
+ }
+ }
+ .toList
+ }
+ }
+
+ /**
+ * Every execution row that would have kept `cuid` out of the scan's
result: one that is not in
+ * a terminal state, or one whose activity lands at or after `cutoff`.
Re-asserted inside the
+ * terminating UPDATE so a run started between the scan and the update
takes the unit off the
+ * table -- the scan reads every candidate before terminating any of them,
so that window is as
+ * wide as the whole sweep, not an instant.
+ */
+ private[resource] def liveExecutionExists(cuid: Integer, cutoff: Timestamp):
Condition =
+ exists(
+ selectOne()
+ .from(WORKFLOW_EXECUTIONS)
+ .where(
+ WORKFLOW_EXECUTIONS.CUID
+ .eq(cuid)
+ .and(
+ WORKFLOW_EXECUTIONS.STATUS
+ .notIn(TerminalWorkflowExecutionStatus.dbStatuses: _*)
+ .or(WORKFLOW_EXECUTIONS.STARTING_TIME.ge(cutoff))
+ .or(WORKFLOW_EXECUTIONS.LAST_UPDATE_TIME.ge(cutoff))
+ )
+ )
+ )
+
+ private[resource] def terminateIdleKubernetesComputingUnitCandidate(
+ candidate: IdleComputingUnitCandidate,
+ cutoff: Timestamp,
+ terminationTime: Timestamp,
+ k8s: => KubernetesClient
+ ): Option[TerminatedComputingUnitInfo] = {
+ val unit = candidate.unit
+ val cuid = unit.getCuid
+ val reason = WorkflowComputingUnitTerminationReasonEnum.GARBAGE_COLLECTED
+ try {
+ withTransaction(context) { ctx =>
+ // Stamp the row first and delete the pod second, within one
transaction per unit. The
+ // guards the scan applied are repeated in the WHERE clause, so a unit
a user terminated
+ // or started a run on in the meantime updates zero rows and is left
alone. Deleting an
+ // absent pod is a no-op and deleting a live one is idempotent, so
letting a delete
+ // failure roll the stamp back only costs a retry next round --
whereas stamping after a
+ // failed delete would leave a live pod behind a row that says
terminated.
+ val marked = ctx
+ .update(WORKFLOW_COMPUTING_UNIT)
+ .set(WORKFLOW_COMPUTING_UNIT.TERMINATE_TIME, terminationTime)
+ .set(WORKFLOW_COMPUTING_UNIT.TERMINATION_REASON, reason)
+ .where(
+ WORKFLOW_COMPUTING_UNIT.CUID
+ .eq(cuid)
+ .and(WORKFLOW_COMPUTING_UNIT.TERMINATE_TIME.isNull)
+
.and(WORKFLOW_COMPUTING_UNIT.TYPE.eq(WorkflowComputingUnitTypeEnum.kubernetes))
+ .andNot(liveExecutionExists(cuid, cutoff))
+ )
+ .execute() == 1
+
+ if (!marked) {
+ None
+ } else {
+ k8s.deletePod(cuid)
+ Some(
+ TerminatedComputingUnitInfo(
+ cuid = cuid,
+ name = unit.getName,
+ uid = unit.getUid,
+ username = candidate.username,
+ reason = reason
+ )
+ )
+ }
+ }
+ } catch {
+ case NonFatal(t) =>
+ logger.warn(s"Failed to terminate idle Kubernetes computing unit
cuid=$cuid", t)
+ None
+ }
+ }
+
private def icebergEnvironmentVariables: Map[String, Any] = {
val base = Map[String, Any](
EnvironmentalVariable.ENV_ICEBERG_CATALOG_TYPE ->
StorageConfig.icebergCatalogType
@@ -167,6 +404,14 @@ object ComputingUnitManagingResource {
) ++ requiredComputingUnitEnv(EnvironmentalVariable.get) ++
optionalComputingUnitEnv(EnvironmentalVariable.get)
+ case class TerminatedComputingUnitInfo(
+ cuid: Integer,
+ name: String,
+ uid: Integer,
+ username: Option[String],
+ reason: WorkflowComputingUnitTerminationReasonEnum
+ )
+
case class WorkflowComputingUnitCreationParams(
name: String,
unitType: String,
@@ -215,7 +460,6 @@ object ComputingUnitManagingResource {
@Produces(Array(MediaType.APPLICATION_JSON))
@Path("/computing-unit")
class ComputingUnitManagingResource {
-
private def getComputingUnitByCuid(ctx: DSLContext, cuid: Int):
WorkflowComputingUnit = {
val wcDao = new WorkflowComputingUnitDao(ctx.configuration())
val unit = wcDao.fetchOneByCuid(cuid)
@@ -678,8 +922,17 @@ class ComputingUnitManagingResource {
KubernetesClient.deletePod(cuid)
}
+ val terminationReason =
WorkflowComputingUnitTerminationReasonEnum.USER_REQUESTED
unit.setTerminateTime(new Timestamp(System.currentTimeMillis()))
+ unit.setTerminationReason(terminationReason)
cuDao.update(unit)
+ // owner_* describes the unit, terminated_by_* the caller: an ADMIN may
terminate a unit
+ // they do not own, so a single `uid`/`username` pair would mix the two
identities.
+ logger.info(
+ s"Terminated computing unit: cuid=${unit.getCuid},
name=${unit.getName}, " +
+ s"owner_uid=${unit.getUid}, terminated_by_uid=${user.getUid}, " +
+ s"terminated_by=${user.getName},
reason=${terminationReason.getLiteral}"
+ )
}
Response.ok().build()
}
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 509603c243..f2e6781cc8 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
@@ -18,7 +18,10 @@
package org.apache.texera.service.util
-import org.apache.texera.dao.jooq.generated.enums.WorkflowComputingUnitTypeEnum
+import org.apache.texera.dao.jooq.generated.enums.{
+ WorkflowComputingUnitTerminationReasonEnum,
+ WorkflowComputingUnitTypeEnum
+}
import org.apache.texera.dao.jooq.generated.tables.daos.{UserDao,
WorkflowComputingUnitDao}
import org.apache.texera.dao.jooq.generated.tables.pojos.WorkflowComputingUnit
import org.apache.texera.service.resource.ComputingUnitManagingResource.{
@@ -200,7 +203,10 @@ object ComputingUnitHelpers {
val vanished = partitioned._2
if (vanished.nonEmpty) {
val now = new Timestamp(System.currentTimeMillis())
- vanished.foreach(_.setTerminateTime(now))
+ vanished.foreach { unit =>
+ unit.setTerminateTime(now)
+
unit.setTerminationReason(WorkflowComputingUnitTerminationReasonEnum.GARBAGE_COLLECTED)
+ }
dao.update(vanished.asJava)
}
partitioned._1
diff --git
a/computing-unit-managing-service/src/main/scala/org/apache/texera/service/util/IdleComputingUnitCleanupJob.scala
b/computing-unit-managing-service/src/main/scala/org/apache/texera/service/util/IdleComputingUnitCleanupJob.scala
new file mode 100644
index 0000000000..07fd5dc5ab
--- /dev/null
+++
b/computing-unit-managing-service/src/main/scala/org/apache/texera/service/util/IdleComputingUnitCleanupJob.scala
@@ -0,0 +1,107 @@
+/*
+ * 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 com.typesafe.scalalogging.LazyLogging
+import io.dropwizard.lifecycle.Managed
+import org.apache.texera.service.resource.ComputingUnitManagingResource
+import
org.apache.texera.service.resource.ComputingUnitManagingResource.TerminatedComputingUnitInfo
+
+import java.util.concurrent.{Executors, ScheduledExecutorService, TimeUnit}
+
+/**
+ * Periodically terminates Kubernetes computing units whose last workflow
execution activity is
+ * older than the idle timeout, reclaiming pods their owners have stopped
using.
+ *
+ * @param idleTimeoutMinutes Idle time (in minutes) after which a unit is
terminated.
+ * @param intervalMinutes Delay (in minutes) between cleanup rounds.
+ */
+class IdleComputingUnitCleanupJob(
+ idleTimeoutMinutes: Long,
+ intervalMinutes: Long,
+ terminateIdleComputingUnits: Long => List[TerminatedComputingUnitInfo] =
+ ComputingUnitManagingResource.terminateIdleKubernetesComputingUnits
+) extends Managed
+ with LazyLogging {
+
+ require(idleTimeoutMinutes > 0, s"idleTimeoutMinutes must be > 0 (got
$idleTimeoutMinutes)")
+ require(intervalMinutes > 0, s"intervalMinutes must be > 0 (got
$intervalMinutes)")
+
+ private var executor: ScheduledExecutorService = _
+
+ override def start(): Unit = {
+ executor = Executors.newSingleThreadScheduledExecutor((runnable: Runnable)
=> {
+ val thread = new Thread(runnable, "idle-computing-unit-cleanup")
+ thread.setDaemon(true)
+ thread
+ })
+ executor.scheduleWithFixedDelay(
+ () => runScheduledTick(),
+ // Small fixed initial delay so a restart doesn't postpone the backlog
of already-idle units
+ // by up to a full interval.
+ 1L,
+ intervalMinutes,
+ TimeUnit.MINUTES
+ )
+ }
+
+ /**
+ * Runs one cleanup round for the scheduler. Visible for testing. Catches
every Throwable
+ * because an exception escaping the scheduled task would cancel the
fixed-delay schedule and
+ * silently stop all future cleanup rounds.
+ */
+ private[util] def runScheduledTick(): Unit =
+ try {
+ runCleanupOnce()
+ } catch {
+ case t: Throwable => logger.error("Idle computing unit cleanup round
failed", t)
+ }
+
+ /**
+ * Runs a single cleanup round. Logs each terminated unit with its owner so
a user who lost a
+ * computing unit can be traced in the logs. Idempotent: units already
terminated are not
+ * revisited, and failures are retried on the next round.
+ *
+ * @return The units terminated in this round.
+ */
+ private[util] def runCleanupOnce(): List[TerminatedComputingUnitInfo] = {
+ val terminated = terminateIdleComputingUnits(idleTimeoutMinutes)
+ if (terminated.nonEmpty) {
+ val terminatedDetails = terminated
+ .map(unit =>
+ // Same owner_* keys the manual termination path logs, so both are
grepped alike; the
+ // sweep has no acting user, which is what reason=GARBAGE_COLLECTED
already says.
+ s"cuid=${unit.cuid}, name=${unit.name}, owner_uid=${unit.uid}, " +
+ s"owner=${unit.username.getOrElse("unknown")},
reason=${unit.reason.getLiteral}"
+ )
+ .mkString("; ")
+ logger.info(
+ s"Terminated ${terminated.size} idle Kubernetes computing unit(s):
$terminatedDetails"
+ )
+ }
+ terminated
+ }
+
+ override def stop(): Unit = {
+ if (executor != null) {
+ executor.shutdown()
+ }
+ }
+}
diff --git
a/computing-unit-managing-service/src/test/scala/org/apache/texera/service/ComputingUnitIdleCleanupSchedulerSpec.scala
b/computing-unit-managing-service/src/test/scala/org/apache/texera/service/ComputingUnitIdleCleanupSchedulerSpec.scala
new file mode 100644
index 0000000000..94fa929e35
--- /dev/null
+++
b/computing-unit-managing-service/src/test/scala/org/apache/texera/service/ComputingUnitIdleCleanupSchedulerSpec.scala
@@ -0,0 +1,107 @@
+/*
+ * 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
+
+import io.dropwizard.core.setup.Environment
+import io.dropwizard.lifecycle.JettyManaged
+import org.apache.texera.service.util.IdleComputingUnitCleanupJob
+import org.scalatest.flatspec.AnyFlatSpec
+import org.scalatest.matchers.should.Matchers
+
+import scala.jdk.CollectionConverters._
+
+/**
+ * Spec for the conditional wiring of the idle computing unit cleanup job.
The sweep itself is
+ * covered by `ComputingUnitIdleCleanupSpec`, and the job's scheduling by
+ * `IdleComputingUnitCleanupJobSpec`; what is pinned here is only which
configurations put a job
+ * on the lifecycle.
+ */
+class ComputingUnitIdleCleanupSchedulerSpec extends AnyFlatSpec with Matchers {
+
+ private val service = new ComputingUnitManagingService()
+
+ /** The IdleComputingUnitCleanupJob instances registered on an environment's
lifecycle. */
+ private def registeredCleanupJobs(
+ environment: Environment
+ ): Seq[IdleComputingUnitCleanupJob] =
+ environment
+ .lifecycle()
+ .getManagedObjects
+ .asScala
+ .collect {
+ case managed: JettyManaged
+ if managed.getManaged.isInstanceOf[IdleComputingUnitCleanupJob] =>
+ managed.getManaged.asInstanceOf[IdleComputingUnitCleanupJob]
+ }
+ .toSeq
+
+ "registerIdleComputingUnitCleanup" should "manage an
IdleComputingUnitCleanupJob on the lifecycle when enabled" in {
+ val environment = new Environment("test-computing-unit-managing-service")
+ service.registerIdleComputingUnitCleanup(
+ environment,
+ enabled = true,
+ idleTimeoutMinutes = 1440,
+ intervalMinutes = 60
+ )
+ registeredCleanupJobs(environment) should have size 1
+ }
+
+ it should "register nothing when disabled" in {
+ val environment = new Environment("test-computing-unit-managing-service")
+ service.registerIdleComputingUnitCleanup(
+ environment,
+ enabled = false,
+ idleTimeoutMinutes = 1440,
+ intervalMinutes = 60
+ )
+ registeredCleanupJobs(environment) shouldBe empty
+ }
+
+ it should "not construct the job (so not throw) when disabled even with
invalid config" in {
+ // idleTimeoutMinutes/intervalMinutes are invalid (0), but because enabled
= false the job is
+ // never constructed, so IdleComputingUnitCleanupJob's require(...) is
never evaluated and
+ // nothing throws. This pins that the enabled check guards construction,
not just registration.
+ val environment = new Environment("test-computing-unit-managing-service")
+ service.registerIdleComputingUnitCleanup(
+ environment,
+ enabled = false,
+ idleTimeoutMinutes = 0,
+ intervalMinutes = 0
+ )
+ registeredCleanupJobs(environment) shouldBe empty
+ }
+
+ it should "skip the sweep rather than abort startup when enabled with a
non-positive timeout or interval" in {
+ // A misconfigured sweep leaves the rest of the service perfectly usable,
and the scheduler
+ // would reject a non-positive delay outright, so the wiring logs and
skips instead of letting
+ // the job's require(...) take the whole service down at startup.
+ Seq((0L, 60L), (-1L, 60L), (1440L, 0L), (1440L, -1L)).foreach {
+ case (idleTimeoutMinutes, intervalMinutes) =>
+ val environment = new
Environment("test-computing-unit-managing-service")
+ noException should be thrownBy
service.registerIdleComputingUnitCleanup(
+ environment,
+ enabled = true,
+ idleTimeoutMinutes = idleTimeoutMinutes,
+ intervalMinutes = intervalMinutes
+ )
+ registeredCleanupJobs(environment) shouldBe empty
+ }
+ }
+}
diff --git
a/computing-unit-managing-service/src/test/scala/org/apache/texera/service/resource/ComputingUnitIdleCleanupSpec.scala
b/computing-unit-managing-service/src/test/scala/org/apache/texera/service/resource/ComputingUnitIdleCleanupSpec.scala
new file mode 100644
index 0000000000..3a736550c8
--- /dev/null
+++
b/computing-unit-managing-service/src/test/scala/org/apache/texera/service/resource/ComputingUnitIdleCleanupSpec.scala
@@ -0,0 +1,544 @@
+/*
+ * 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.resource
+
+import org.apache.texera.auth.SessionUser
+import org.apache.texera.dao.MockTexeraDB
+import org.apache.texera.dao.jooq.generated.Tables.{
+ USER => USER_TABLE,
+ WORKFLOW,
+ WORKFLOW_COMPUTING_UNIT,
+ WORKFLOW_EXECUTIONS,
+ WORKFLOW_VERSION
+}
+import org.apache.texera.dao.jooq.generated.enums.{
+ WorkflowComputingUnitTerminationReasonEnum,
+ WorkflowComputingUnitTypeEnum
+}
+import org.apache.texera.dao.jooq.generated.tables.daos.{
+ UserDao,
+ WorkflowComputingUnitDao,
+ WorkflowDao,
+ WorkflowExecutionsDao,
+ WorkflowVersionDao
+}
+import org.apache.texera.dao.jooq.generated.tables.pojos.{
+ User,
+ Workflow,
+ WorkflowComputingUnit,
+ WorkflowExecutions,
+ WorkflowVersion
+}
+import
org.apache.texera.service.resource.ComputingUnitManagingResource.TerminatedComputingUnitInfo
+import org.apache.texera.service.util.KubernetesClient
+import org.mockito.ArgumentMatchers.anyInt
+import org.mockito.Mockito.{doThrow, mock, never, verify, verifyNoInteractions}
+import org.scalatest.flatspec.AnyFlatSpec
+import org.scalatest.matchers.should.Matchers
+import org.scalatest.{BeforeAndAfterAll, BeforeAndAfterEach}
+
+import java.sql.Timestamp
+import java.util.UUID
+import java.util.concurrent.TimeUnit
+
+class ComputingUnitIdleCleanupSpec
+ extends AnyFlatSpec
+ with Matchers
+ with BeforeAndAfterAll
+ with BeforeAndAfterEach
+ with MockTexeraDB {
+
+ private val testUserId = 810000 + scala.util.Random.nextInt(10000)
+ private val testWorkflowId = 820000 + scala.util.Random.nextInt(10000)
+ private val now = new Timestamp(TimeUnit.DAYS.toMillis(20))
+ private val idleTimeoutMinutes = 60L
+
+ private var userDao: UserDao = _
+ private var workflowDao: WorkflowDao = _
+ private var workflowVersionDao: WorkflowVersionDao = _
+ private var workflowComputingUnitDao: WorkflowComputingUnitDao = _
+ private var workflowExecutionsDao: WorkflowExecutionsDao = _
+ private var testVersion: WorkflowVersion = _
+
+ override protected def beforeAll(): Unit =
+ initializeDBAndReplaceDSLContext()
+
+ override protected def beforeEach(): Unit = {
+ userDao = new UserDao(getDSLContext.configuration())
+ workflowDao = new WorkflowDao(getDSLContext.configuration())
+ workflowVersionDao = new WorkflowVersionDao(getDSLContext.configuration())
+ workflowComputingUnitDao = new
WorkflowComputingUnitDao(getDSLContext.configuration())
+ workflowExecutionsDao = new
WorkflowExecutionsDao(getDSLContext.configuration())
+
+ cleanupTestData()
+
+ val user = new User
+ user.setUid(testUserId)
+ user.setName("idle-cu-owner")
+ user.setEmail(s"idle-cu-${UUID.randomUUID()}@example.com")
+ userDao.insert(user)
+
+ val workflow = new Workflow
+ workflow.setWid(testWorkflowId)
+ workflow.setName("idle-cu-workflow")
+ workflow.setContent("{}")
+ workflow.setCreationTime(new Timestamp(now.getTime -
TimeUnit.DAYS.toMillis(2)))
+ workflow.setLastModifiedTime(new Timestamp(now.getTime -
TimeUnit.DAYS.toMillis(2)))
+ workflowDao.insert(workflow)
+
+ testVersion = new WorkflowVersion
+ testVersion.setWid(testWorkflowId)
+ testVersion.setContent("{}")
+ testVersion.setCreationTime(new Timestamp(now.getTime -
TimeUnit.DAYS.toMillis(2)))
+ workflowVersionDao.insert(testVersion)
+ }
+
+ override protected def afterEach(): Unit =
+ cleanupTestData()
+
+ override protected def afterAll(): Unit =
+ shutdownDB()
+
+ private def cleanupTestData(): Unit = {
+ getDSLContext
+ .deleteFrom(WORKFLOW_EXECUTIONS)
+ .where(WORKFLOW_EXECUTIONS.UID.eq(testUserId))
+ .execute()
+ getDSLContext
+ .deleteFrom(WORKFLOW_COMPUTING_UNIT)
+ .where(WORKFLOW_COMPUTING_UNIT.UID.eq(testUserId))
+ .execute()
+ getDSLContext
+ .deleteFrom(WORKFLOW_VERSION)
+ .where(WORKFLOW_VERSION.WID.eq(testWorkflowId))
+ .execute()
+
getDSLContext.deleteFrom(WORKFLOW).where(WORKFLOW.WID.eq(testWorkflowId)).execute()
+
getDSLContext.deleteFrom(USER_TABLE).where(USER_TABLE.UID.eq(testUserId)).execute()
+ }
+
+ private def timestampMinutesBefore(minutes: Long): Timestamp =
+ new Timestamp(now.getTime - TimeUnit.MINUTES.toMillis(minutes))
+
+ private def insertComputingUnit(
+ name: String,
+ unitType: WorkflowComputingUnitTypeEnum =
WorkflowComputingUnitTypeEnum.kubernetes,
+ creationMinutesBefore: Long = 120,
+ terminated: Boolean = false
+ ): WorkflowComputingUnit = {
+ val unit = new WorkflowComputingUnit
+ unit.setUid(testUserId)
+ unit.setName(name)
+ unit.setCreationTime(timestampMinutesBefore(creationMinutesBefore))
+ unit.setType(unitType)
+ unit.setUri("kubernetes://test")
+ unit.setResource("{}")
+ if (terminated) {
+ unit.setTerminateTime(timestampMinutesBefore(10))
+
unit.setTerminationReason(WorkflowComputingUnitTerminationReasonEnum.USER_REQUESTED)
+ }
+ workflowComputingUnitDao.insert(unit)
+ unit
+ }
+
+ private def insertExecution(
+ unit: WorkflowComputingUnit,
+ status: Short,
+ startingMinutesBefore: Long,
+ lastUpdateMinutesBefore: Option[Long] = None
+ ): Unit = {
+ val execution = new WorkflowExecutions
+ execution.setVid(testVersion.getVid)
+ execution.setUid(testUserId)
+ execution.setCuid(unit.getCuid)
+ execution.setStatus(status)
+ execution.setStartingTime(timestampMinutesBefore(startingMinutesBefore))
+ lastUpdateMinutesBefore.foreach(minutes =>
+ execution.setLastUpdateTime(timestampMinutesBefore(minutes))
+ )
+ execution.setBookmarked(false)
+ execution.setName("execution-" + UUID.randomUUID().toString.substring(0,
8))
+ execution.setEnvironmentVersion("test-env")
+ workflowExecutionsDao.insert(execution)
+ }
+
+ private def sessionUser(
+ uid: Integer = testUserId,
+ name: String = "idle-cu-owner"
+ ): SessionUser = {
+ val user = new User
+ user.setUid(uid)
+ user.setName(name)
+ new SessionUser(user)
+ }
+
+ // The sweep drives the Kubernetes client through the same by-name seam
ComputingUnitHelpers
+ // uses, so a stub stands in for the production singleton.
+ private def stubKubernetesClient(): KubernetesClient =
mock(classOf[KubernetesClient])
+
+ "TerminatedComputingUnitInfo" should "carry the terminated unit and its
owner" in {
+ val terminated = TerminatedComputingUnitInfo(
+ cuid = 1,
+ name = "plain-unit-info",
+ uid = testUserId,
+ username = Some("idle-cu-owner"),
+ reason = WorkflowComputingUnitTerminationReasonEnum.GARBAGE_COLLECTED
+ )
+ terminated.cuid shouldBe 1
+ terminated.username shouldBe Some("idle-cu-owner")
+ terminated.reason shouldBe
WorkflowComputingUnitTerminationReasonEnum.GARBAGE_COLLECTED
+ }
+
+ "terminateIdleKubernetesComputingUnits" should "garbage collect only
inactive Kubernetes computing units past the timeout" in {
+ val stale = insertComputingUnit("stale")
+ val active = insertComputingUnit("active")
+ val recent = insertComputingUnit("recent")
+ val local = insertComputingUnit("local",
WorkflowComputingUnitTypeEnum.local)
+ val alreadyTerminated = insertComputingUnit("already-terminated",
terminated = true)
+
+ insertExecution(active, status = 1, startingMinutesBefore = 180)
+ insertExecution(
+ recent,
+ status = 3,
+ startingMinutesBefore = 180,
+ lastUpdateMinutesBefore = Some(5)
+ )
+ insertExecution(local, status = 3, startingMinutesBefore = 180)
+
+ val k8s = stubKubernetesClient()
+ val terminated =
ComputingUnitManagingResource.terminateIdleKubernetesComputingUnits(
+ idleTimeoutMinutes,
+ now,
+ k8s
+ )
+
+ terminated.map(_.cuid) shouldBe List(stale.getCuid)
+ terminated.head.username shouldBe Some("idle-cu-owner")
+ terminated.head.reason shouldBe
WorkflowComputingUnitTerminationReasonEnum.GARBAGE_COLLECTED
+ verify(k8s).deletePod(stale.getCuid)
+ verify(k8s, never()).deletePod(active.getCuid)
+ verify(k8s, never()).deletePod(recent.getCuid)
+ verify(k8s, never()).deletePod(local.getCuid)
+
+ val staleAfterCleanup =
workflowComputingUnitDao.fetchOneByCuid(stale.getCuid)
+ staleAfterCleanup.getTerminateTime shouldBe now
+ staleAfterCleanup.getTerminationReason shouldBe
WorkflowComputingUnitTerminationReasonEnum.GARBAGE_COLLECTED
+ workflowComputingUnitDao.fetchOneByCuid(active.getCuid).getTerminateTime
shouldBe null
+ workflowComputingUnitDao.fetchOneByCuid(recent.getCuid).getTerminateTime
shouldBe null
+ workflowComputingUnitDao.fetchOneByCuid(local.getCuid).getTerminateTime
shouldBe null
+
workflowComputingUnitDao.fetchOneByCuid(alreadyTerminated.getCuid).getTerminationReason
shouldBe
+ WorkflowComputingUnitTerminationReasonEnum.USER_REQUESTED
+ }
+
+ it should "never touch the cluster when no unit is idle" in {
+ val active = insertComputingUnit("active-only")
+ insertExecution(active, status = 1, startingMinutesBefore = 180)
+ val k8s = stubKubernetesClient()
+
+ ComputingUnitManagingResource.terminateIdleKubernetesComputingUnits(
+ idleTimeoutMinutes,
+ now,
+ k8s
+ ) shouldBe empty
+
+ verifyNoInteractions(k8s)
+ }
+
+ it should "roll the termination back and keep collecting other units when
one pod deletion fails" in {
+ val failing = insertComputingUnit("stale-delete-fails")
+ val successful = insertComputingUnit("stale-delete-succeeds")
+ val k8s = stubKubernetesClient()
+ doThrow(new RuntimeException("pod deletion
failed")).when(k8s).deletePod(failing.getCuid)
+
+ val terminated =
ComputingUnitManagingResource.terminateIdleKubernetesComputingUnits(
+ idleTimeoutMinutes,
+ now,
+ k8s
+ )
+
+ terminated.map(_.cuid) shouldBe List(successful.getCuid)
+ // The stamp is written before the pod is deleted, so the failed delete
rolls it back and the
+ // unit is retried on the next sweep rather than being left as a live pod
marked terminated.
+ workflowComputingUnitDao.fetchOneByCuid(failing.getCuid).getTerminateTime
shouldBe null
+
workflowComputingUnitDao.fetchOneByCuid(failing.getCuid).getTerminationReason
shouldBe null
+ workflowComputingUnitDao
+ .fetchOneByCuid(successful.getCuid)
+ .getTerminationReason shouldBe
+ WorkflowComputingUnitTerminationReasonEnum.GARBAGE_COLLECTED
+ }
+
+ it should "mark an idle unit terminated even when its pod is already gone"
in {
+ // Deleting an absent pod is a no-op in the Kubernetes API, so the sweep
issues the delete
+ // unconditionally rather than paying a second round trip to ask whether
the pod exists.
+ val stale = insertComputingUnit("stale-missing-pod")
+ val k8s = stubKubernetesClient()
+
+ val terminated =
ComputingUnitManagingResource.terminateIdleKubernetesComputingUnits(
+ idleTimeoutMinutes,
+ now,
+ k8s
+ )
+
+ terminated.map(_.cuid) shouldBe List(stale.getCuid)
+ verify(k8s).deletePod(stale.getCuid)
+
workflowComputingUnitDao.fetchOneByCuid(stale.getCuid).getTerminationReason
shouldBe
+ WorkflowComputingUnitTerminationReasonEnum.GARBAGE_COLLECTED
+ }
+
+ it should "keep units whose latest execution is in any non-terminal state
running" in {
+ // 0/1/2 are UNINITIALIZED-or-READY, RUNNING and PAUSED; -1 is what
maptoStatusCode collapses
+ // PAUSING, RESUMING, UNKNOWN and TERMINATED to. None of them may be read
as idle.
+ val units = Seq[Short](0, 1, 2, -1).map { status =>
+ val unit = insertComputingUnit(s"active-status-$status")
+ insertExecution(unit, status = status, startingMinutesBefore = 180)
+ unit
+ }
+ val k8s = stubKubernetesClient()
+
+ ComputingUnitManagingResource.terminateIdleKubernetesComputingUnits(
+ idleTimeoutMinutes,
+ now,
+ k8s
+ ) shouldBe empty
+
+ verify(k8s, never()).deletePod(anyInt())
+ units.foreach(unit =>
+ workflowComputingUnitDao.fetchOneByCuid(unit.getCuid).getTerminateTime
shouldBe null
+ )
+ }
+
+ it should "garbage collect a unit whose latest execution finished past the
timeout" in {
+ // 3/4/5 are COMPLETED, FAILED and KILLED -- the terminal codes, so the
unit counts as idle.
+ Seq[Short](3, 4, 5).foreach { status =>
+ val stale = insertComputingUnit(s"stale-status-$status")
+ insertExecution(
+ stale,
+ status = status,
+ startingMinutesBefore = 180,
+ lastUpdateMinutesBefore = Some(90)
+ )
+ val k8s = stubKubernetesClient()
+
+ val terminated =
ComputingUnitManagingResource.terminateIdleKubernetesComputingUnits(
+ idleTimeoutMinutes,
+ now,
+ k8s
+ )
+
+ terminated.map(_.cuid) shouldBe List(stale.getCuid)
+ verify(k8s).deletePod(stale.getCuid)
+ workflowComputingUnitDao
+ .fetchOneByCuid(stale.getCuid)
+ .getTerminationReason shouldBe
+ WorkflowComputingUnitTerminationReasonEnum.GARBAGE_COLLECTED
+ }
+ }
+
+ it should "omit an empty owner name from terminated unit info" in {
+ val user = userDao.fetchOneByUid(testUserId)
+ user.setName("")
+ userDao.update(user)
+ insertComputingUnit("stale-empty-owner")
+
+ val terminated =
ComputingUnitManagingResource.terminateIdleKubernetesComputingUnits(
+ idleTimeoutMinutes,
+ now,
+ stubKubernetesClient()
+ )
+
+ terminated should have size 1
+ terminated.head.username shouldBe None
+ }
+
+ // The sweep reads every candidate before terminating any of them, so the
window between the scan
+ // and a given unit's UPDATE is as wide as the whole sweep. These drive the
two steps separately
+ // and mutate the DB in between, which is what that window looks like from
the unit's side.
+ private def scanCandidate(unit: WorkflowComputingUnit) = {
+ val cutoff = timestampMinutesBefore(idleTimeoutMinutes)
+ val candidate = ComputingUnitManagingResource
+ .idleKubernetesComputingUnitCandidates(cutoff)
+ .find(_.unit.getCuid == unit.getCuid)
+ candidate should not be empty
+ (candidate.get, cutoff)
+ }
+
+ "the terminating update" should "leave a unit alone when a run starts after
the scan" in {
+ val unit = insertComputingUnit("raced-run-started")
+ val (candidate, cutoff) = scanCandidate(unit)
+
+ // The user starts a workflow on the unit after the scan listed it as idle.
+ insertExecution(unit, status = 1, startingMinutesBefore = 0)
+
+ val k8s = stubKubernetesClient()
+
ComputingUnitManagingResource.terminateIdleKubernetesComputingUnitCandidate(
+ candidate,
+ cutoff,
+ now,
+ k8s
+ ) shouldBe None
+
+ verifyNoInteractions(k8s)
+ workflowComputingUnitDao.fetchOneByCuid(unit.getCuid).getTerminateTime
shouldBe null
+ }
+
+ it should "leave a unit alone when an execution finished after the scan" in {
+ // Terminal status, so the status half of the guard does not fire -- only
the activity
+ // timestamps place the execution at or after the cutoff.
+ val unit = insertComputingUnit("raced-run-finished")
+ val (candidate, cutoff) = scanCandidate(unit)
+
+ insertExecution(
+ unit,
+ status = 3,
+ startingMinutesBefore = 1,
+ lastUpdateMinutesBefore = Some(1)
+ )
+
+ val k8s = stubKubernetesClient()
+
ComputingUnitManagingResource.terminateIdleKubernetesComputingUnitCandidate(
+ candidate,
+ cutoff,
+ now,
+ k8s
+ ) shouldBe None
+
+ verifyNoInteractions(k8s)
+ workflowComputingUnitDao.fetchOneByCuid(unit.getCuid).getTerminateTime
shouldBe null
+ }
+
+ it should "leave a unit alone when it was terminated after the scan" in {
+ val unit = insertComputingUnit("raced-user-terminated")
+ val (candidate, cutoff) = scanCandidate(unit)
+
+ // Stamped directly rather than through terminateComputingUnit, whose
kubernetes branch talks
+ // to the real cluster singleton; what the guard reads is the row, and
this is the row a user
+ // termination leaves behind.
+ val terminatedByUser =
workflowComputingUnitDao.fetchOneByCuid(unit.getCuid)
+ terminatedByUser.setTerminateTime(timestampMinutesBefore(1))
+ terminatedByUser.setTerminationReason(
+ WorkflowComputingUnitTerminationReasonEnum.USER_REQUESTED
+ )
+ workflowComputingUnitDao.update(terminatedByUser)
+
+ val k8s = stubKubernetesClient()
+
ComputingUnitManagingResource.terminateIdleKubernetesComputingUnitCandidate(
+ candidate,
+ cutoff,
+ now,
+ k8s
+ ) shouldBe None
+
+ verifyNoInteractions(k8s)
+ workflowComputingUnitDao.fetchOneByCuid(unit.getCuid).getTerminationReason
shouldBe
+ WorkflowComputingUnitTerminationReasonEnum.USER_REQUESTED
+ }
+
+ it should "terminate the unit when nothing changed since the scan" in {
+ val unit = insertComputingUnit("unraced")
+ val (candidate, cutoff) = scanCandidate(unit)
+ // A stale execution, still entirely before the cutoff, must not trip the
guard.
+ insertExecution(
+ unit,
+ status = 3,
+ startingMinutesBefore = 180,
+ lastUpdateMinutesBefore = Some(90)
+ )
+
+ val k8s = stubKubernetesClient()
+ ComputingUnitManagingResource
+ .terminateIdleKubernetesComputingUnitCandidate(candidate, cutoff, now,
k8s)
+ .map(_.cuid) shouldBe Some(unit.getCuid)
+
+ verify(k8s).deletePod(unit.getCuid)
+ workflowComputingUnitDao.fetchOneByCuid(unit.getCuid).getTerminationReason
shouldBe
+ WorkflowComputingUnitTerminationReasonEnum.GARBAGE_COLLECTED
+ }
+
+ "terminateComputingUnit" should "mark manual termination as user requested"
in {
+ val local = insertComputingUnit("manual-local",
WorkflowComputingUnitTypeEnum.local)
+
+ val response =
+ new
ComputingUnitManagingResource().terminateComputingUnit(local.getCuid,
sessionUser())
+
+ response.getStatus shouldBe 200
+ val terminated = workflowComputingUnitDao.fetchOneByCuid(local.getCuid)
+ terminated.getTerminateTime should not be null
+ terminated.getTerminationReason shouldBe
WorkflowComputingUnitTerminationReasonEnum.USER_REQUESTED
+ }
+
+ it should "reject manual termination from a non-owner" in {
+ val local = insertComputingUnit("manual-local-non-owner",
WorkflowComputingUnitTypeEnum.local)
+
+ val response = new ComputingUnitManagingResource().terminateComputingUnit(
+ local.getCuid,
+ sessionUser(uid = testUserId + 1)
+ )
+
+ response.getStatus shouldBe 400
+ workflowComputingUnitDao.fetchOneByCuid(local.getCuid).getTerminateTime
shouldBe null
+ }
+
+ "lastComputingUnitActivityTime" should "prefer the latest execution
timestamp over creation time" in {
+ val unit = new WorkflowComputingUnit
+ unit.setCreationTime(timestampMinutesBefore(120))
+
+ ComputingUnitManagingResource.lastComputingUnitActivityTime(
+ unit,
+ latestUpdateTime = Some(timestampMinutesBefore(10)),
+ latestStartTime = Some(timestampMinutesBefore(30))
+ ) shouldBe timestampMinutesBefore(10)
+ }
+
+ it should "fall back to start time and then creation time" in {
+ val unit = new WorkflowComputingUnit
+ unit.setCreationTime(timestampMinutesBefore(120))
+
+ ComputingUnitManagingResource.lastComputingUnitActivityTime(
+ unit,
+ latestUpdateTime = None,
+ latestStartTime = Some(timestampMinutesBefore(30))
+ ) shouldBe timestampMinutesBefore(30)
+
+ ComputingUnitManagingResource.lastComputingUnitActivityTime(
+ unit,
+ latestUpdateTime = None,
+ latestStartTime = None
+ ) shouldBe timestampMinutesBefore(120)
+ }
+
+ "shouldTerminateIdleComputingUnit" should "require both no active execution
and activity before cutoff" in {
+ val cutoff = timestampMinutesBefore(60)
+
+ ComputingUnitManagingResource.shouldTerminateIdleComputingUnit(
+ hasActiveExecution = false,
+ lastExecutionTime = timestampMinutesBefore(61),
+ cutoff = cutoff
+ ) shouldBe true
+ ComputingUnitManagingResource.shouldTerminateIdleComputingUnit(
+ hasActiveExecution = true,
+ lastExecutionTime = timestampMinutesBefore(61),
+ cutoff = cutoff
+ ) shouldBe false
+ ComputingUnitManagingResource.shouldTerminateIdleComputingUnit(
+ hasActiveExecution = false,
+ lastExecutionTime = timestampMinutesBefore(60),
+ cutoff = cutoff
+ ) shouldBe false
+ }
+}
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 5cf6fae214..18554cd1fe 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
@@ -24,6 +24,7 @@ import org.apache.texera.dao.MockTexeraDB
import org.apache.texera.dao.jooq.generated.enums.{
PrivilegeEnum,
UserRoleEnum,
+ WorkflowComputingUnitTerminationReasonEnum,
WorkflowComputingUnitTypeEnum
}
import org.apache.texera.dao.jooq.generated.tables.daos.{UserDao,
WorkflowComputingUnitDao}
@@ -317,6 +318,8 @@ class ComputingUnitHelpersSpec
live.map(_.getCuid) should contain theSameElementsAs Seq(600, 602)
computingUnitDao.fetchOneByCuid(601).getTerminateTime should not be null
+ computingUnitDao.fetchOneByCuid(601).getTerminationReason shouldBe
+ WorkflowComputingUnitTerminationReasonEnum.GARBAGE_COLLECTED
computingUnitDao.fetchOneByCuid(600).getTerminateTime shouldBe null
computingUnitDao.fetchOneByCuid(602).getTerminateTime shouldBe null
}
diff --git
a/computing-unit-managing-service/src/test/scala/org/apache/texera/service/util/IdleComputingUnitCleanupJobSpec.scala
b/computing-unit-managing-service/src/test/scala/org/apache/texera/service/util/IdleComputingUnitCleanupJobSpec.scala
new file mode 100644
index 0000000000..d9ec45f3e0
--- /dev/null
+++
b/computing-unit-managing-service/src/test/scala/org/apache/texera/service/util/IdleComputingUnitCleanupJobSpec.scala
@@ -0,0 +1,109 @@
+/*
+ * 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
org.apache.texera.dao.jooq.generated.enums.WorkflowComputingUnitTerminationReasonEnum
+import
org.apache.texera.service.resource.ComputingUnitManagingResource.TerminatedComputingUnitInfo
+import org.scalatest.flatspec.AnyFlatSpec
+import org.scalatest.matchers.should.Matchers
+
+/**
+ * Spec for [[IdleComputingUnitCleanupJob]]. The sweep it runs is covered by
+ * `ComputingUnitIdleCleanupSpec`; what is pinned here is the job's own
contract — it forwards
+ * the configured timeout, reports what it terminated, and never lets a
failing round cancel the
+ * fixed-delay schedule.
+ */
+class IdleComputingUnitCleanupJobSpec extends AnyFlatSpec with Matchers {
+
+ private def terminatedUnit(cuid: Int, username: Option[String]) =
+ TerminatedComputingUnitInfo(
+ cuid = cuid,
+ name = s"unit-$cuid",
+ uid = 7,
+ username = username,
+ reason = WorkflowComputingUnitTerminationReasonEnum.GARBAGE_COLLECTED
+ )
+
+ "IdleComputingUnitCleanupJob" should "reject a non-positive timeout or
interval at construction" in {
+ assertThrows[IllegalArgumentException](new IdleComputingUnitCleanupJob(0,
60))
+ assertThrows[IllegalArgumentException](new IdleComputingUnitCleanupJob(-1,
60))
+ assertThrows[IllegalArgumentException](new
IdleComputingUnitCleanupJob(1440, 0))
+ assertThrows[IllegalArgumentException](new
IdleComputingUnitCleanupJob(1440, -1))
+ }
+
+ "runCleanupOnce" should "pass the configured idle timeout to the sweep and
return what it terminated" in {
+ var seenTimeouts = List.empty[Long]
+ val job = new IdleComputingUnitCleanupJob(
+ idleTimeoutMinutes = 1440,
+ intervalMinutes = 60,
+ terminateIdleComputingUnits = timeout => {
+ seenTimeouts = seenTimeouts :+ timeout
+ List(terminatedUnit(1, Some("owner")), terminatedUnit(2, None))
+ }
+ )
+
+ job.runCleanupOnce().map(_.cuid) shouldBe List(1, 2)
+ seenTimeouts shouldBe List(1440L)
+ }
+
+ it should "return empty when the sweep terminates nothing" in {
+ val job = new IdleComputingUnitCleanupJob(1440, 60, _ => List.empty)
+ job.runCleanupOnce() shouldBe empty
+ }
+
+ "runScheduledTick" should "swallow a failing round so the fixed-delay
schedule survives" in {
+ // An exception escaping the scheduled task cancels the schedule outright,
silently stopping
+ // every future round, so the tick has to absorb it and let the next round
retry.
+ val job = new IdleComputingUnitCleanupJob(
+ 1440,
+ 60,
+ _ => throw new RuntimeException("sweep failed")
+ )
+ noException should be thrownBy job.runScheduledTick()
+ }
+
+ it should "run the sweep on a successful round" in {
+ var invocations = 0
+ val job = new IdleComputingUnitCleanupJob(
+ 1440,
+ 60,
+ _ => {
+ invocations += 1
+ List.empty
+ }
+ )
+ job.runScheduledTick()
+ invocations shouldBe 1
+ }
+
+ "the job lifecycle" should "allow stop() before start() without throwing" in
{
+ noException should be thrownBy new IdleComputingUnitCleanupJob(1440, 60, _
=> List.empty).stop()
+ }
+
+ it should "start() then stop() without throwing" in {
+ val job = new IdleComputingUnitCleanupJob(1440, 60, _ => List.empty)
+ try {
+ noException should be thrownBy job.start()
+ } finally {
+ // Always stop so a started daemon executor never leaks between tests.
+ job.stop()
+ }
+ }
+}
diff --git a/sql/changelog.xml b/sql/changelog.xml
index b16670669a..97724697c6 100644
--- a/sql/changelog.xml
+++ b/sql/changelog.xml
@@ -164,6 +164,11 @@
<sqlFile path="sql/updates/50.sql"/>
</changeSet>
+ <!-- Record why a computing unit was terminated and index
workflow_executions.cuid -->
+ <changeSet id="51" author="yrenat">
+ <sqlFile path="sql/updates/51.sql"/>
+ </changeSet>
+
<!-- example changeSet
<changeSet id="1" author="author">
<sqlFile path="sql/updates/1.sql"/>
diff --git a/sql/texera_ddl.sql b/sql/texera_ddl.sql
index 4a7e483db0..dde189976a 100644
--- a/sql/texera_ddl.sql
+++ b/sql/texera_ddl.sql
@@ -98,6 +98,7 @@ CREATE TYPE user_role_enum AS ENUM ('INACTIVE', 'RESTRICTED',
'REGULAR', 'ADMIN'
CREATE TYPE action_enum AS ENUM ('like', 'unlike', 'view', 'clone');
CREATE TYPE privilege_enum AS ENUM ('NONE', 'READ', 'WRITE');
CREATE TYPE workflow_computing_unit_type_enum AS ENUM ('local', 'kubernetes');
+CREATE TYPE workflow_computing_unit_termination_reason_enum AS ENUM
('USER_REQUESTED', 'GARBAGE_COLLECTED');
CREATE TYPE provider_type_enum AS ENUM ('LOCAL', 'GOOGLE', 'ORCID', 'APPLE');
CREATE TYPE user_warehouse_flavor_enum AS ENUM ('local', 'aws');
CREATE TYPE default_view_enum AS ENUM ('CANVAS', 'FORM');
@@ -246,6 +247,7 @@ CREATE TABLE IF NOT EXISTS workflow_computing_unit
cuid SERIAL PRIMARY KEY,
creation_time TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP,
terminate_time TIMESTAMP DEFAULT NULL,
+ termination_reason workflow_computing_unit_termination_reason_enum DEFAULT
NULL,
type workflow_computing_unit_type_enum,
uri TEXT NOT NULL DEFAULT '',
resource TEXT DEFAULT '',
@@ -333,6 +335,10 @@ CREATE TABLE IF NOT EXISTS workflow_executions
FOREIGN KEY (whid) REFERENCES user_warehouse(whid) ON DELETE SET NULL
);
+-- Postgres indexes only the referenced side of a foreign key, so cuid needs
its own index for the
+-- idle computing unit sweep and every other per-computing-unit lookup on this
table.
+CREATE INDEX idx_workflow_executions_cuid ON workflow_executions (cuid);
+
-- dataset
CREATE TABLE IF NOT EXISTS dataset
(
diff --git a/sql/updates/51.sql b/sql/updates/51.sql
new file mode 100644
index 0000000000..bbf5dc6b18
--- /dev/null
+++ b/sql/updates/51.sql
@@ -0,0 +1,45 @@
+/*
+ * 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.
+ */
+
+\c texera_db
+
+SET search_path TO texera_db;
+
+BEGIN;
+
+DO $$
+BEGIN
+ IF NOT EXISTS (
+ SELECT 1 FROM pg_type WHERE typname =
'workflow_computing_unit_termination_reason_enum'
+ ) THEN
+ CREATE TYPE workflow_computing_unit_termination_reason_enum AS ENUM (
+ 'USER_REQUESTED',
+ 'GARBAGE_COLLECTED'
+ );
+ END IF;
+END $$;
+
+ALTER TABLE workflow_computing_unit
+ ADD COLUMN IF NOT EXISTS termination_reason
workflow_computing_unit_termination_reason_enum DEFAULT NULL;
+
+-- Postgres indexes only the referenced side of a foreign key, so cuid needs
its own index for the
+-- idle computing unit sweep and every other per-computing-unit lookup on this
table.
+CREATE INDEX IF NOT EXISTS idx_workflow_executions_cuid ON workflow_executions
(cuid);
+
+COMMIT;