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-8472-ded7ba1b49d0dc6644bc34b765b727607ad9d654 in repository https://gitbox.apache.org/repos/asf/texera.git
commit 7dbdb728807ce53c6266b5b49b278d04ad2df390 Author: Tanishq Gandhi <[email protected]> AuthorDate: Thu Sep 10 18:59:42 2026 +0000 feat(computing-unit): name the configuration a unit cannot start without (#8472) ### What changes were proposed in this PR? Six environment variables were read with `Option.get`, so a deployment missing any of them failed with `None.get`. Dropwizard shows that as "There was an error processing your request" — no variable named, no cause, and six candidates to check. Every missing variable is now listed at once, as a 503. All at once because when these are missing they are usually all missing: the helm chart supplies them, so they are present together or absent together. Reviewing what each one is for shrank the list. A variable belongs in it only if the unit has no usable default without it: | Variable | Without it, the unit | Verdict | | --- | --- | --- | | `FILE_SERVICE_GET_DATASET_PRESIGNED_URL_ENDPOINT` | falls back to `localhost:9092` | required | | `FILE_SERVICE_UPLOAD_ONE_FILE_TO_DATASET_ENDPOINT` | falls back to `localhost:9092` | required | | `AUTH_JWT_SECRET` | falls back to the published literal in `auth.conf` | required | | `MAX_WORKFLOW_WEBSOCKET_REQUEST_PAYLOAD_SIZE_KB` | uses `application.conf`'s 1024 | optional, forwarded when set | | `USER_SYS_ENABLED` | ignores it — conf key removed in #3831 | dropped | | `SCHEDULE_GENERATOR_ENABLE_COST_BASED_SCHEDULE_GENERATOR` | ignores it — conf key removed in #3542 | dropped | Required values are forwarded raw: `LakeFSFileDocument` and `ResultExportService` trim the endpoints themselves, and the secret has to stay byte-identical to what `AuthConfig` read or the unit and this service verify the token against different keys. The optional override is trimmed, because HOCON reads `" 1024"` as a string and refuses it as an int — the unit then dies at startup naming nothing. **No change for any working deployment.** A chart deployment sets all of these, unpadded, so every forwarded value is byte-identical to today's; the only difference in the pod's environment is the two dropped variables, which nothing there read. A deployment that omits the payload size is no longer refused. ### Any related issues, documentation, discussions? Closes #8467 Part of #8466 ### How was this PR tested? Ten new tests for `requiredComputingUnitEnv`, and the existing default-image assertion updated. | Case | What it pins | | --- | --- | | all set | exactly the three keys come back, values untouched | | not needed | none of the three defaulted/dead variables is required, nor forwarded | | one missing | that one is named — **and the ones that are set are not** | | blank | a whitespace-only value counts as missing | | padded | forwarded raw, so unit and manager verify against the same key | | wording | "Unset or blank environment variable(s)", not "missing" | | all missing | every name in one message | | override set | the payload size is forwarded, trimmed | | override unset/blank | nothing forwarded, so `application.conf`'s 1024 stands | The negative half of the one-missing case was verified by mutation: computing `missing` as every name whenever one is absent leaves the suite green without it, and fails only that test with it. ``` sbt "ComputingUnitManagingService/testOnly org.apache.texera.service.resource.ComputingUnitManagingResourceSpec" \ "Config/testOnly org.apache.texera.common.config.KubernetesConfigSpec" ComputingUnitManagingResourceSpec 41 passed, 0 failed KubernetesConfigSpec 6 passed, 0 failed ``` ### Was this PR authored or co-authored using generative AI tooling? Generated-by: Claude Code (Claude Opus 5) --- common/config/src/main/resources/kubernetes.conf | 4 +- .../common/config/KubernetesConfigSpec.scala | 3 +- .../resource/ComputingUnitManagingResource.scala | 80 +++++++++++----- .../ComputingUnitManagingResourceSpec.scala | 105 ++++++++++++++++++++- 4 files changed, 167 insertions(+), 25 deletions(-) diff --git a/common/config/src/main/resources/kubernetes.conf b/common/config/src/main/resources/kubernetes.conf index da72b930d3..e3edf31aff 100644 --- a/common/config/src/main/resources/kubernetes.conf +++ b/common/config/src/main/resources/kubernetes.conf @@ -34,7 +34,9 @@ kubernetes { compute-unit-pod-name-prefix = "computing-unit" compute-unit-pod-name-prefix = ${?KUBERNETES_COMPUTE_UNIT_POD_NAME_PREFIX} - image-name = "bobbai/texera-workflow-computing-unit:dev" + # Nightly image, tracking main. The chart always sets KUBERNETES_IMAGE_NAME to its own + # registry and tag, so only a service run outside the chart reaches this default. + image-name = "ghcr.io/apache/texera-workflow-execution-coordinator:latest" image-name = ${?KUBERNETES_IMAGE_NAME} image-pull-policy = "Always" 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 3b3c194b63..7dbeedd07c 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 @@ -49,7 +49,8 @@ class KubernetesConfigSpec extends AnyFlatSpec with Matchers { KubernetesConfig.computeUnitPoolNamespace shouldBe "texera-workflow-computing-unit-pool" ) ifUnset("KUBERNETES_IMAGE_NAME")( - KubernetesConfig.computeUnitImageName shouldBe "bobbai/texera-workflow-computing-unit:dev" + KubernetesConfig.computeUnitImageName shouldBe + "ghcr.io/apache/texera-workflow-execution-coordinator:latest" ) ifUnset("KUBERNETES_IMAGE_PULL_POLICY")( KubernetesConfig.computingUnitImagePullPolicy shouldBe "Always" 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 3a249d296e..c5df9078b0 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 @@ -93,6 +93,61 @@ object ComputingUnitManagingResource { } } + // Required: the endpoints default to localhost:9092 (LakeFSFileDocument, + // ResultExportService) and the secret to a published literal (auth.conf), none of which + // suits a real deployment. Forwarded raw -- the endpoints are trimmed by their own readers, + // and trimming the secret would leave the unit and this service verifying the token against + // different keys, since AuthConfig does not trim. + private val requiredComputingUnitEnvNames: Seq[String] = Seq( + EnvironmentalVariable.ENV_FILE_SERVICE_GET_DATASET_PRESIGNED_URL_ENDPOINT, + EnvironmentalVariable.ENV_FILE_SERVICE_UPLOAD_ONE_FILE_TO_DATASET_ENDPOINT, + EnvironmentalVariable.ENV_AUTH_JWT_SECRET + ) + + // Overrides, forwarded only when set: application.conf defaults the payload size to 1024, + // so its absence is not an error. USER_SYS_ENABLED and + // SCHEDULE_GENERATOR_ENABLE_COST_BASED_SCHEDULE_GENERATOR are absent from both lists -- + // their conf keys went away with #3831 and #3542, so nothing reads them. + // TODO: use AmberConfig here; it is only accessible in workflow-executing-service + private val optionalComputingUnitEnvNames: Seq[String] = Seq( + EnvironmentalVariable.ENV_MAX_WORKFLOW_WEBSOCKET_REQUEST_PAYLOAD_SIZE_KB + ) + + /** + * Returns the variables, or fails with a 503 listing every one that is unset or blank. + * + * A WebApplicationException so the message survives: dropwizard replaces a plain 500's + * with generic text. 503 because the deployment is not ready, not the request wrong. + */ + private[resource] def requiredComputingUnitEnv( + lookup: String => Option[String] + ): Map[String, String] = { + // Blank counts as missing: the chart renders every value as "{{ .value }}", so an unset + // one arrives as "" rather than absent. + val looked = + requiredComputingUnitEnvNames.map(name => name -> lookup(name).filter(_.trim.nonEmpty)) + val missing = looked.collect { case (name, None) => name } + if (missing.nonEmpty) { + throw new ServiceUnavailableException( + "This deployment cannot create a computing unit. Unset or blank environment " + + s"variable(s): ${missing.mkString(", ")}." + ) + } + looked.collect { case (name, Some(value)) => name -> value }.toMap + } + + /** + * The overrides that are set, trimmed. A blank one is dropped rather than forwarded, and + * a padded one is trimmed, because HOCON reads " 1024" as a string and refuses it as an + * int -- the unit then dies at startup naming nothing. + */ + private[resource] def optionalComputingUnitEnv( + lookup: String => Option[String] + ): Map[String, String] = + optionalComputingUnitEnvNames.flatMap { name => + lookup(name).map(_.trim).filter(_.nonEmpty).map(name -> _) + }.toMap + // Environment variables passed to the created computing unit(pod) private lazy val computingUnitEnvironmentVariables: Map[String, Any] = icebergEnvironmentVariables ++ Map( @@ -108,28 +163,9 @@ object ComputingUnitManagingResource { EnvironmentalVariable.ENV_S3_ENDPOINT -> StorageConfig.s3Endpoint, EnvironmentalVariable.ENV_S3_REGION -> StorageConfig.s3Region, EnvironmentalVariable.ENV_S3_AUTH_USERNAME -> StorageConfig.s3Username, - EnvironmentalVariable.ENV_S3_AUTH_PASSWORD -> StorageConfig.s3Password, - EnvironmentalVariable.ENV_FILE_SERVICE_GET_DATASET_PRESIGNED_URL_ENDPOINT -> EnvironmentalVariable - .get(EnvironmentalVariable.ENV_FILE_SERVICE_GET_DATASET_PRESIGNED_URL_ENDPOINT) - .get, - EnvironmentalVariable.ENV_FILE_SERVICE_UPLOAD_ONE_FILE_TO_DATASET_ENDPOINT -> EnvironmentalVariable - .get(EnvironmentalVariable.ENV_FILE_SERVICE_UPLOAD_ONE_FILE_TO_DATASET_ENDPOINT) - .get, - // Variables for amber setting - // TODO: use AmberConfig for the following items. Currently AmberConfig is only accessible in workflow-executing-service - EnvironmentalVariable.ENV_SCHEDULE_GENERATOR_ENABLE_COST_BASED_SCHEDULE_GENERATOR -> EnvironmentalVariable - .get(EnvironmentalVariable.ENV_SCHEDULE_GENERATOR_ENABLE_COST_BASED_SCHEDULE_GENERATOR) - .get, - EnvironmentalVariable.ENV_USER_SYS_ENABLED -> EnvironmentalVariable - .get(EnvironmentalVariable.ENV_USER_SYS_ENABLED) - .get, - EnvironmentalVariable.ENV_MAX_WORKFLOW_WEBSOCKET_REQUEST_PAYLOAD_SIZE_KB -> EnvironmentalVariable - .get(EnvironmentalVariable.ENV_MAX_WORKFLOW_WEBSOCKET_REQUEST_PAYLOAD_SIZE_KB) - .get, - EnvironmentalVariable.ENV_AUTH_JWT_SECRET -> EnvironmentalVariable - .get(EnvironmentalVariable.ENV_AUTH_JWT_SECRET) - .get - ) + EnvironmentalVariable.ENV_S3_AUTH_PASSWORD -> StorageConfig.s3Password + ) ++ requiredComputingUnitEnv(EnvironmentalVariable.get) ++ + optionalComputingUnitEnv(EnvironmentalVariable.get) case class WorkflowComputingUnitCreationParams( name: String, diff --git a/computing-unit-managing-service/src/test/scala/org/apache/texera/service/resource/ComputingUnitManagingResourceSpec.scala b/computing-unit-managing-service/src/test/scala/org/apache/texera/service/resource/ComputingUnitManagingResourceSpec.scala index 4cbb0e0781..34b3c996cb 100644 --- a/computing-unit-managing-service/src/test/scala/org/apache/texera/service/resource/ComputingUnitManagingResourceSpec.scala +++ b/computing-unit-managing-service/src/test/scala/org/apache/texera/service/resource/ComputingUnitManagingResourceSpec.scala @@ -19,7 +19,13 @@ package org.apache.texera.service.resource -import jakarta.ws.rs.{BadRequestException, ForbiddenException, NotFoundException} +import jakarta.ws.rs.{ + BadRequestException, + ForbiddenException, + NotFoundException, + ServiceUnavailableException +} +import org.apache.texera.common.config.EnvironmentalVariable import org.apache.texera.auth.SessionUser import org.apache.texera.common.config.KubernetesConfig.maxNumOfRunningComputingUnitsPerUser import org.apache.texera.dao.MockTexeraDB @@ -460,4 +466,101 @@ class ComputingUnitManagingResourceSpec a[NotFoundException] should be thrownBy resource.getComputingUnitResourceLimit("99999", user) } + + // Mirrors the private production list, so a name added or dropped there fails the all-set + // case. + private val requiredEnvNames = Seq( + EnvironmentalVariable.ENV_FILE_SERVICE_GET_DATASET_PRESIGNED_URL_ENDPOINT, + EnvironmentalVariable.ENV_FILE_SERVICE_UPLOAD_ONE_FILE_TO_DATASET_ENDPOINT, + EnvironmentalVariable.ENV_AUTH_JWT_SECRET + ) + + private val payloadSize = EnvironmentalVariable.ENV_MAX_WORKFLOW_WEBSOCKET_REQUEST_PAYLOAD_SIZE_KB + + "requiredComputingUnitEnv" should "return every variable when all are set" in { + val env = ComputingUnitManagingResource.requiredComputingUnitEnv(name => Some(s"value-$name")) + env.keySet shouldBe requiredEnvNames.toSet + env.values.foreach(_ should startWith("value-")) + } + + // USER_SYS_ENABLED and SCHEDULE_GENERATOR_ENABLE_COST_BASED_SCHEDULE_GENERATOR lost their + // conf keys (#3831, #3542); the payload size defaults to 1024 in application.conf. None of + // the three stops a unit from starting, so none may refuse to create one. + it should "not require a variable the unit does not need" in { + val notNeeded = Seq( + EnvironmentalVariable.ENV_USER_SYS_ENABLED, + EnvironmentalVariable.ENV_SCHEDULE_GENERATOR_ENABLE_COST_BASED_SCHEDULE_GENERATOR, + payloadSize + ) + val env = ComputingUnitManagingResource.requiredComputingUnitEnv(name => + if (notNeeded.contains(name)) None else Some("set") + ) + notNeeded.foreach(name => env.keySet should not contain name) + } + + it should "name the missing variable and leave the ones that are set out of it" in { + val absent = EnvironmentalVariable.ENV_AUTH_JWT_SECRET + val thrown = intercept[ServiceUnavailableException] { + ComputingUnitManagingResource.requiredComputingUnitEnv(name => + if (name == absent) None else Some("set") + ) + } + thrown.getMessage should include(absent) + requiredEnvNames + .filterNot(_ == absent) + .foreach(name => thrown.getMessage should not include name) + } + + // The chart renders every value as "{{ .value }}", so an unset one arrives as "". + it should "treat a blank variable as missing" in { + val blank = EnvironmentalVariable.ENV_AUTH_JWT_SECRET + val thrown = intercept[ServiceUnavailableException] { + ComputingUnitManagingResource.requiredComputingUnitEnv(name => + if (name == blank) Some(" ") else Some("set") + ) + } + thrown.getMessage should include(blank) + } + + // AuthConfig does not trim, so a trimmed copy would verify against a different key. The + // endpoints are trimmed by their own readers, so they need nothing here either. + it should "hand on the value untrimmed" in { + val env = + ComputingUnitManagingResource.requiredComputingUnitEnv(_ => Some(" s3cret ")) + env(EnvironmentalVariable.ENV_AUTH_JWT_SECRET) shouldBe " s3cret " + } + + // The variable is there in the pod's env, so calling it "missing" would read as wrong. + it should "say unset or blank rather than missing" in { + val thrown = intercept[ServiceUnavailableException] { + ComputingUnitManagingResource.requiredComputingUnitEnv(_ => Some(" ")) + } + thrown.getMessage should include("Unset or blank environment variable(s)") + } + + it should "name every missing variable at once" in { + val thrown = intercept[ServiceUnavailableException] { + ComputingUnitManagingResource.requiredComputingUnitEnv(_ => None) + } + requiredEnvNames.foreach(name => thrown.getMessage should include(name)) + } + + "optionalComputingUnitEnv" should "forward an override that is set" in { + ComputingUnitManagingResource.optionalComputingUnitEnv(_ => Some("2048")) shouldBe + Map(payloadSize -> "2048") + } + + // HOCON refuses " 2048" as an int, and the unit then dies at startup naming nothing. + it should "trim what it forwards" in { + ComputingUnitManagingResource.optionalComputingUnitEnv(_ => Some(" 2048\n")) shouldBe + Map(payloadSize -> "2048") + } + + // application.conf already defaults this to 1024; forwarding "" would override the default + // with a value HOCON cannot read as an int. + it should "forward nothing when unset or blank" in { + ComputingUnitManagingResource.optionalComputingUnitEnv(_ => None) shouldBe empty + ComputingUnitManagingResource.optionalComputingUnitEnv(_ => Some(" ")) shouldBe empty + } + }
