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 b1d51ae243 feat(computing-unit-managing-service): add admin endpoint
to list all computing units (#6854)
b1d51ae243 is described below
commit b1d51ae2438d7d5e78ce5d3a226eec8c87c21e9d
Author: Kunwoo (Chris) <[email protected]>
AuthorDate: Wed Aug 5 13:29:57 2026 -0400
feat(computing-unit-managing-service): add admin endpoint to list all
computing units (#6854)
### What changes were proposed in this PR?
Adds an ADMIN-only endpoint that returns every non-terminated computing
unit across all users. Needed for the admin Computing Units dashboard.
The endpoint queries all units via `WorkflowComputingUnitDao`, excludes
rows with a non-null `terminate_time`, joins `UserDao` for each owner's
name and avatar, and reuses the shared `ComputingUnitHelpers`
(introduced in the preceding refactor PR) to reconcile vanished pods and
resolve status/metrics.
The new resource is registered in `ComputingUnitManagingService`, and
the startup access-control coverage check is extended to include it.
`accessPrivilege` is reported as `WRITE` for admin rows (an admin can
manage any unit it can see), and `isOwner` reflects whether the
requesting admin happens to own the unit. Note that the mutating
endpoints gate on ownership alone and have no ADMIN bypass yet, so a
client must not present these rows as writable until they do — recorded
in the resource's scaladoc.
> 📌 **#6853 has landed and this is rebased onto it.** The diff is now
just the admin resource + registration + specs, in a single commit.
### Any related issues, documentation, discussions?
Closes #6477. Part of #6476 (Admin Computing Units Dashboard).
Originates from discussion #6322.
### How was this PR tested?
Added `AdminComputingUnitResourceSpec` (drives `listAllComputingUnits`
end-to-end over the embedded DB: terminated rows excluded, `WRITE`
access, `isOwner`, an empty table, and a caller owning none of the
listed units) and extended `ComputingUnitManagingServiceRunSpec` to
verify the admin resource registers along with the auth features that
make `@RolesAllowed` enforceable.
Verified on the rebased commit by CI's `platform
(computing-unit-managing-service)` job, which runs
`ComputingUnitManagingService/test` and the scalafmt checks: 79 tests, 0
failed, 0 canceled (Postgres is reachable there, so the DB-gated `run()`
registration test executes rather than cancelling).
### Was this PR authored or co-authored using generative AI tooling?
Generated-by: Claude Code, Claude Opus 4.8
---
.../service/ComputingUnitManagingService.scala | 2 +
.../resource/AdminComputingUnitResource.scala | 111 +++++++++++++++
.../ComputingUnitManagingServiceRunSpec.scala | 65 +++++++++
.../resource/AdminComputingUnitResourceSpec.scala | 158 +++++++++++++++++++++
4 files changed, 336 insertions(+)
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 db63bbf2eb..f0dffc89e1 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
@@ -27,6 +27,7 @@ import org.apache.texera.common.config.StorageConfig
import org.apache.texera.auth.{AuthFeatures, RequestLoggingFilter,
RoleAnnotationEnforcer}
import org.apache.texera.dao.SqlServer
import org.apache.texera.service.resource.{
+ AdminComputingUnitResource,
ComputingUnitAccessResource,
ComputingUnitManagingResource,
HealthCheckResource
@@ -66,6 +67,7 @@ class ComputingUnitManagingService extends
Application[ComputingUnitManagingServ
environment.jersey().register(new ComputingUnitManagingResource)
environment.jersey().register(new ComputingUnitAccessResource)
+ environment.jersey().register(new AdminComputingUnitResource)
RoleAnnotationEnforcer.enforce(
environment.jersey.getResourceConfig,
diff --git
a/computing-unit-managing-service/src/main/scala/org/apache/texera/service/resource/AdminComputingUnitResource.scala
b/computing-unit-managing-service/src/main/scala/org/apache/texera/service/resource/AdminComputingUnitResource.scala
new file mode 100644
index 0000000000..45d137a4b4
--- /dev/null
+++
b/computing-unit-managing-service/src/main/scala/org/apache/texera/service/resource/AdminComputingUnitResource.scala
@@ -0,0 +1,111 @@
+/*
+ * 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 io.dropwizard.auth.Auth
+import jakarta.annotation.security.RolesAllowed
+import jakarta.ws.rs.{GET, Path, Produces}
+import jakarta.ws.rs.core.MediaType
+import org.apache.texera.auth.SessionUser
+import org.apache.texera.dao.SqlServer
+import org.apache.texera.dao.jooq.generated.Tables.WORKFLOW_COMPUTING_UNIT
+import org.apache.texera.dao.jooq.generated.enums.PrivilegeEnum
+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.DashboardWorkflowComputingUnit
+import org.apache.texera.service.util.ComputingUnitHelpers
+import org.jooq.DSLContext
+
+import scala.jdk.CollectionConverters.CollectionHasAsScala
+
+object AdminComputingUnitResource {
+ private def context: DSLContext =
+ SqlServer
+ .getInstance()
+ .createDSLContext()
+}
+
+// Path is /computing-unit/admin, not the /admin/<domain> shape other admin
resources use: the
+// gateway only routes /api/computing-unit here, so /api/admin/... never
reaches this service.
+// Jersey matches the literal /admin/list ahead of
ComputingUnitManagingResource's @Path("/{cuid}"),
+// so don't add another two-segment literal path under /computing-unit without
re-checking this.
+@Produces(Array(MediaType.APPLICATION_JSON))
+@Path("/computing-unit/admin")
+@RolesAllowed(Array("ADMIN"))
+class AdminComputingUnitResource {
+
+ import AdminComputingUnitResource._
+
+ /**
+ * List every non-terminated computing unit across all users (ADMIN-only).
+ *
+ * TODO: rows report WRITE, but the mutating endpoints on
[[ComputingUnitManagingResource]]
+ * still gate on ownership with no ADMIN bypass, so an admin acting on a
unit it does not own is
+ * rejected. Add that bypass; until then a client must not present these
rows as writable.
+ */
+ @GET
+ @Path("/list")
+ def listAllComputingUnits(
+ @Auth user: SessionUser
+ ): List[DashboardWorkflowComputingUnit] = {
+ val ctx = context
+
+ // Filter active units in SQL. Explicit `IS NULL`, not the DAO's
fetchByTerminateTime(null),
+ // which renders `terminate_time IN (null)` and matches nothing.
+ val activeUnits =
+ ctx
+ .selectFrom(WORKFLOW_COMPUTING_UNIT)
+ .where(WORKFLOW_COMPUTING_UNIT.TERMINATE_TIME.isNull)
+ .fetchInto(classOf[WorkflowComputingUnit])
+ .asScala
+ .toList
+
+ val podPhases = ComputingUnitHelpers.podPhasesFor(activeUnits)
+
+ // Wrap only the reconcile write, not the Kubernetes round trips above:
keeps a pooled
+ // connection from being held open across them, while the transaction
makes the batchUpdate
+ // all-or-nothing (autocommit would retire only some vanished units on a
mid-batch failure).
+
+ val liveUnits = SqlServer.withTransaction(ctx) { txCtx =>
+ ComputingUnitHelpers.reconcileVanishedKubernetesUnits(
+ new WorkflowComputingUnitDao(txCtx.configuration()),
+ activeUnits,
+ podPhases
+ )
+ }
+
+ // Metrics only for survivors, so fetch after reconciliation.
+ val podMetrics = ComputingUnitHelpers.podMetricsFor(liveUnits)
+
+ val userDao = new UserDao(ctx.configuration())
+ val ownerInfo = ComputingUnitHelpers.resolveOwnerInfo(userDao,
liveUnits.map(_.getUid).distinct)
+
+ liveUnits.map { unit =>
+ ComputingUnitHelpers.buildDashboardUnit(
+ unit,
+ isOwner = unit.getUid.equals(user.getUid),
+ accessPrivilege = PrivilegeEnum.WRITE,
+ ownerInfo = ownerInfo,
+ podPhases = podPhases,
+ podMetrics = podMetrics
+ )
+ }
+ }
+}
diff --git
a/computing-unit-managing-service/src/test/scala/org/apache/texera/service/ComputingUnitManagingServiceRunSpec.scala
b/computing-unit-managing-service/src/test/scala/org/apache/texera/service/ComputingUnitManagingServiceRunSpec.scala
index d2162d48c7..e4694e2d20 100644
---
a/computing-unit-managing-service/src/test/scala/org/apache/texera/service/ComputingUnitManagingServiceRunSpec.scala
+++
b/computing-unit-managing-service/src/test/scala/org/apache/texera/service/ComputingUnitManagingServiceRunSpec.scala
@@ -19,25 +19,90 @@
package org.apache.texera.service
+import io.dropwizard.auth.AuthDynamicFeature
+import io.dropwizard.core.setup.Environment
+import io.dropwizard.jersey.DropwizardResourceConfig
+import io.dropwizard.jersey.setup.JerseyEnvironment
+import io.dropwizard.jetty.MutableServletContextHandler
import org.apache.texera.auth.RoleAnnotationEnforcer
+import org.apache.texera.common.config.StorageConfig
+import org.apache.texera.dao.SqlServer
import org.apache.texera.service.resource.{
+ AdminComputingUnitResource,
ComputingUnitAccessResource,
ComputingUnitManagingResource,
HealthCheckResource
}
+import org.glassfish.jersey.server.filter.RolesAllowedDynamicFeature
+import org.mockito.ArgumentMatchers.isA
+import org.mockito.Mockito.{mock, verify, when}
import org.scalatest.flatspec.AnyFlatSpec
import org.scalatest.matchers.should.Matchers
+import java.sql.DriverManager
+
class ComputingUnitManagingServiceRunSpec extends AnyFlatSpec with Matchers {
+ // Booting the service opens a real connection pool, so this path is only
exercisable where
+ // Postgres is provisioned at the configured JDBC URL (as in CI).
+ private def databaseReachable: Boolean =
+ try {
+ DriverManager
+ .getConnection(
+ StorageConfig.jdbcUrl,
+ StorageConfig.jdbcUsername,
+ StorageConfig.jdbcPassword
+ )
+ .close()
+ true
+ } catch {
+ case _: Throwable => false
+ }
+
// Every endpoint this service registers declares
@RolesAllowed/@PermitAll/@DenyAll.
"ComputingUnitManagingService's registered resources" should "all declare
access control" in {
RoleAnnotationEnforcer.findUnannotatedEndpoints(
Seq(
classOf[ComputingUnitManagingResource],
classOf[ComputingUnitAccessResource],
+ classOf[AdminComputingUnitResource],
classOf[HealthCheckResource]
)
) shouldBe empty
}
+
+ "ComputingUnitManagingService.run" should "register the admin resource on
the Jersey environment" in {
+ assume(
+ databaseReachable,
+ "run() requires a reachable Postgres at the configured JDBC URL
(provided in CI)"
+ )
+
+ val jersey = mock(classOf[JerseyEnvironment])
+ val context = mock(classOf[MutableServletContextHandler])
+ val env = mock(classOf[Environment])
+ when(env.jersey).thenReturn(jersey)
+ when(env.getApplicationContext).thenReturn(context)
+
when(jersey.getResourceConfig).thenReturn(DropwizardResourceConfig.forTesting())
+
+ try {
+ new ComputingUnitManagingService()
+ .run(mock(classOf[ComputingUnitManagingServiceConfiguration]), env)
+
+ verify(jersey).register(isA(classOf[ComputingUnitManagingResource]))
+ verify(jersey).register(isA(classOf[ComputingUnitAccessResource]))
+ verify(jersey).register(isA(classOf[AdminComputingUnitResource]))
+ verify(jersey).setUrlPattern("/api/*")
+ // Without these two, Jersey never enforces AdminComputingUnitResource's
+ // @RolesAllowed(ADMIN) and the cross-user listing is readable by any
authenticated user.
+ verify(jersey).register(isA(classOf[AuthDynamicFeature]))
+ verify(jersey).register(classOf[RolesAllowedDynamicFeature])
+ } finally {
+ // run() calls SqlServer.initConnection, which opens a real HikariCP
pool against the
+ // configured JDBC URL. Close it so the pool's threads/connections don't
outlive this suite
+ // in the shared forked JVM. Swallowed because run() may have failed
before initConnection,
+ // and a throw here would mask that failure.
+ try SqlServer.getInstance().close()
+ catch { case _: Throwable => () }
+ }
+ }
}
diff --git
a/computing-unit-managing-service/src/test/scala/org/apache/texera/service/resource/AdminComputingUnitResourceSpec.scala
b/computing-unit-managing-service/src/test/scala/org/apache/texera/service/resource/AdminComputingUnitResourceSpec.scala
new file mode 100644
index 0000000000..a721986b69
--- /dev/null
+++
b/computing-unit-managing-service/src/test/scala/org/apache/texera/service/resource/AdminComputingUnitResourceSpec.scala
@@ -0,0 +1,158 @@
+/*
+ * 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 jakarta.annotation.security.RolesAllowed
+import org.apache.texera.auth.SessionUser
+import org.apache.texera.dao.MockTexeraDB
+import org.apache.texera.dao.jooq.generated.Tables.WORKFLOW_COMPUTING_UNIT
+import org.apache.texera.dao.jooq.generated.enums.{
+ PrivilegeEnum,
+ UserRoleEnum,
+ WorkflowComputingUnitTypeEnum
+}
+import org.apache.texera.dao.jooq.generated.tables.daos.{UserDao,
WorkflowComputingUnitDao}
+import org.apache.texera.dao.jooq.generated.tables.pojos.{User,
WorkflowComputingUnit}
+import org.scalatest.flatspec.AnyFlatSpec
+import org.scalatest.matchers.should.Matchers
+import org.scalatest.{BeforeAndAfterAll, BeforeAndAfterEach}
+
+import java.sql.Timestamp
+
+class AdminComputingUnitResourceSpec
+ extends AnyFlatSpec
+ with Matchers
+ with BeforeAndAfterAll
+ with BeforeAndAfterEach
+ with MockTexeraDB {
+
+ override protected def beforeAll(): Unit = {
+ super.beforeAll()
+ initializeDBAndReplaceDSLContext()
+ }
+
+ // MockTexeraDB does not truncate between tests, and this endpoint lists
every unit in the
+ // database, so each case must start from a known-empty table — otherwise
the exact-match
+ // assertions below depend on the order ScalaTest happens to run the cases
in. Safe to wipe
+ // unqualified: MockTexeraDB provisions a fresh database per suite, so these
are the only rows.
+ override protected def beforeEach(): Unit = {
+ super.beforeEach()
+ getDSLContext.deleteFrom(WORKFLOW_COMPUTING_UNIT).execute()
+ }
+
+ override protected def afterAll(): Unit =
+ try shutdownDB()
+ finally super.afterAll()
+
+ private def makeUser(uid: Int, name: String): User = {
+ val u = new User()
+ u.setUid(uid)
+ u.setName(name)
+ u.setEmail(s"[email protected]")
+ u.setRole(UserRoleEnum.ADMIN)
+ u.setPassword("password")
+ u.setGoogleAvatar(s"avatar-$uid")
+ u
+ }
+
+ // Only local units are used here: a kubernetes unit would make the listing
reach for the real
+ // KubernetesClient singleton, and its stub seam is private[util].
+ private def localUnit(cuid: Int, uid: Int, name: String):
WorkflowComputingUnit = {
+ val u = new WorkflowComputingUnit()
+ u.setCuid(cuid)
+ u.setUid(uid)
+ u.setName(name)
+ u.setType(WorkflowComputingUnitTypeEnum.local)
+ u
+ }
+
+ // The class-level @RolesAllowed(ADMIN) is what makes Jersey reject
non-ADMIN callers; that this
+ // spec calls the method directly is exactly why the annotation has to be
asserted here. The
+ // enforcement half (Jersey actually honouring it) is pinned by
+ // ComputingUnitManagingServiceRunSpec, which verifies
RolesAllowedDynamicFeature is registered.
+ "AdminComputingUnitResource" should "declare @RolesAllowed(ADMIN) at the
class level" in {
+ val annotation =
classOf[AdminComputingUnitResource].getAnnotation(classOf[RolesAllowed])
+ annotation should not be null
+ annotation.value.toSeq shouldBe Seq("ADMIN")
+ }
+
+ "listAllComputingUnits" should "return every non-terminated unit across
users, marked WRITE" in {
+ val userDao = new UserDao(getDSLContext.configuration())
+ val unitDao = new WorkflowComputingUnitDao(getDSLContext.configuration())
+ val admin = makeUser(700, "admin")
+ userDao.insert(admin)
+ userDao.insert(makeUser(701, "other"))
+ unitDao.insert(localUnit(cuid = 700, uid = 700, name = "admin-cu"))
+ unitDao.insert(localUnit(cuid = 701, uid = 701, name = "other-cu"))
+ // A terminated unit must be excluded by the SQL filter.
+ val terminated = localUnit(cuid = 702, uid = 701, name = "terminated-cu")
+ terminated.setTerminateTime(new Timestamp(0L))
+ unitDao.insert(terminated)
+
+ val result = new AdminComputingUnitResource().listAllComputingUnits(new
SessionUser(admin))
+
+ result.map(_.computingUnit.getCuid.intValue()) should contain
theSameElementsAs Seq(700, 701)
+ all(result.map(_.accessPrivilege)) shouldBe PrivilegeEnum.WRITE
+ all(result.map(_.status)) shouldBe "Running" // local units
+ val byCuid = result.map(r => r.computingUnit.getCuid.intValue() -> r).toMap
+ // isOwner tracks the caller; owner name/avatar are joined from the user
table.
+ byCuid(700).isOwner shouldBe true
+ byCuid(700).ownerName shouldBe "admin"
+ byCuid(700).ownerGoogleAvatar shouldBe "avatar-700"
+ byCuid(701).isOwner shouldBe false
+ byCuid(701).ownerName shouldBe "other"
+ }
+
+ it should "return an empty list when every unit is terminated" in {
+ val userDao = new UserDao(getDSLContext.configuration())
+ val unitDao = new WorkflowComputingUnitDao(getDSLContext.configuration())
+ val admin = makeUser(710, "lonely-admin")
+ userDao.insert(admin)
+ val terminated = localUnit(cuid = 710, uid = 710, name =
"only-terminated-cu")
+ terminated.setTerminateTime(new Timestamp(0L))
+ unitDao.insert(terminated)
+
+ new AdminComputingUnitResource().listAllComputingUnits(new
SessionUser(admin)) shouldBe empty
+ }
+
+ it should "return an empty list when no unit exists at all" in {
+ val admin = makeUser(720, "empty-admin")
+ new UserDao(getDSLContext.configuration()).insert(admin)
+
+ new AdminComputingUnitResource().listAllComputingUnits(new
SessionUser(admin)) shouldBe empty
+ }
+
+ // A caller whose own user row was deleted (or who is an admin from another
realm) still sees
+ // every unit, but owns none of them — isOwner must not accidentally default
to true.
+ it should "mark no row as owned when the caller owns none of the units" in {
+ val userDao = new UserDao(getDSLContext.configuration())
+ val unitDao = new WorkflowComputingUnitDao(getDSLContext.configuration())
+ val admin = makeUser(730, "outsider-admin")
+ userDao.insert(admin)
+ userDao.insert(makeUser(731, "owner"))
+ unitDao.insert(localUnit(cuid = 731, uid = 731, name = "someone-elses-cu"))
+
+ val result = new AdminComputingUnitResource().listAllComputingUnits(new
SessionUser(admin))
+
+ result should have size 1
+ result.head.isOwner shouldBe false
+ result.head.ownerName shouldBe "owner"
+ }
+}