This is an automated email from the ASF dual-hosted git repository. github-merge-queue[bot] pushed a commit to branch gh-readonly-queue/main/pr-6046-36fd40ee04029064a66955e51062e4d8e470298d in repository https://gitbox.apache.org/repos/asf/texera.git
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;
