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 1a584332af test(amber): make PveResourceSpec hermetic via
process-runner seam (#6024)
1a584332af is described below
commit 1a584332af50eb21e715df91a108a3d7c81736cc
Author: Yicong Huang <[email protected]>
AuthorDate: Mon Jun 29 21:15:16 2026 -0700
test(amber): make PveResourceSpec hermetic via process-runner seam (#6024)
### What changes were proposed in this PR?
`PveResourceSpec` runs in the `amber` unit CI job
(`AMBER_TEST_FILTER=skip-integration`) but is not tagged
`@IntegrationTest`, so every run shelled out to **real `pip` over the
network**:
- each `createNewPve` ran `pip install -r requirements.txt` (the full
amber dependency set), and
- the lazy `resolveSystemPackages()` installed `requirements.txt` into a
throwaway venv, then `pip freeze`
— roughly 6 full installs per spec run. That made a PR-blocking unit
spec **network-dependent (flaky)** and **slow**.
This PR funnels every child process in `PveManager` (venv creation, pip
install / uninstall / freeze) through a single injectable seam:
```scala
private[pythonvirtualenvironment] var runProcess: ProcessRunner =
(command, env, logger) => Process(command, None, env: _*).!(logger)
```
Production wiring is unchanged (the default runner executes the command
for real). `PveResourceSpec` swaps in a **ScalaMock `mockFunction`**
whose handler fabricates the `<venv>/bin/{python,pip}` layout, emits the
resolved system set on `freeze`, and returns a configurable exit code —
so the spec is **fully hermetic: no venv, no pip, no network**.
`PveManager` still owns the metadata files (`user-packages.txt`) and
queue messages, so that logic stays under test.
Because failures are now cheap to trigger, this also adds negative
coverage that was previously impractical: venv-create failure,
requirements-install failure, user-package install failure, and
system-package rejection (using `pyarrow`, a hard amber dependency).
| Before | After |
| --- | --- |
| ~6 real `pip install` per run, hits PyPI | 0 network calls |
| flaky on network hiccups, ~minutes | deterministic, ~8s |
| only happy-path assertions | + 4 negative/failure cases |
### Any related issues, documentation, discussions?
Closes #6023
### How was this PR tested?
`PveManager` is an `object`, so the test points the shared `runProcess`
at the ScalaMock `mockFunction` in `beforeAll` and restores the real
runner in `afterAll`. ScalaMock expectations are per-test, so
`expectProcessCalls()` (an `anyNumberOfTimes` handler) is invoked at the
top of each test that exercises a process. Run with JDK 17:
```bash
STORAGE_JDBC_USERNAME=texera STORAGE_JDBC_PASSWORD=password \
sbt "WorkflowExecutionService/testOnly
org.apache.texera.web.resource.pythonvirtualenvironment.PveResourceSpec"
# -> Tests: succeeded 25, failed 0, in ~8s
```
Also green:
```bash
sbt scalafmtCheckAll
sbt "WorkflowExecutionService/scalafixAll --check"
```
### Was this PR authored or co-authored using generative AI tooling?
Generated-by: Claude Code (Claude Opus 4.8)
---
.../pythonvirtualenvironment/PveManager.scala | 50 +++++----
.../pythonvirtualenvironment/PveResourceSpec.scala | 124 ++++++++++++++++++++-
2 files changed, 154 insertions(+), 20 deletions(-)
diff --git
a/amber/src/main/scala/org/apache/texera/web/resource/pythonvirtualenvironment/PveManager.scala
b/amber/src/main/scala/org/apache/texera/web/resource/pythonvirtualenvironment/PveManager.scala
index 423a5b4a58..bb3910cbc3 100644
---
a/amber/src/main/scala/org/apache/texera/web/resource/pythonvirtualenvironment/PveManager.scala
+++
b/amber/src/main/scala/org/apache/texera/web/resource/pythonvirtualenvironment/PveManager.scala
@@ -100,6 +100,16 @@ object PveManager extends LazyLogging {
"PIP_NO_INPUT" -> "1"
)
+ // Test seam: every child process (venv creation, pip
install/uninstall/freeze)
+ // funnels through this so unit tests can run hermetically — no real venv, no
+ // pip, no network. Production wiring runs the command for real;
PveResourceSpec
+ // swaps in a fake that fabricates the venv layout and emits canned output.
+ private[pythonvirtualenvironment] type ProcessRunner =
+ (Seq[String], Seq[(String, String)], ProcessLogger) => Int
+
+ private[pythonvirtualenvironment] var runProcess: ProcessRunner =
+ (command, env, logger) => Process(command, None, env: _*).!(logger)
+
private def readPackageFile(path: Path): Seq[String] = {
if (Files.exists(path)) {
Files
@@ -131,13 +141,17 @@ object PveManager extends LazyLogging {
try {
val python = venvPython(tempVenv).toString
val createCode =
- Process(Seq(PythonUtils.getPythonExecutable, "-m", "venv",
tempVenv.toString)).!
+ runProcess(
+ Seq(PythonUtils.getPythonExecutable, "-m", "venv",
tempVenv.toString),
+ Nil,
+ ProcessLogger(_ => (), _ => ())
+ )
if (createCode != 0) {
logger.error(s"failed to create temp venv for system-package
resolution (exit=$createCode)")
return Seq.empty
}
- val installCode = Process(
+ val installCode = runProcess(
Seq(
python,
"-u",
@@ -150,16 +164,18 @@ object PveManager extends LazyLogging {
"-r",
requirementsPath.toString
),
- None,
- pipEnv.toSeq: _*
- ).!
+ pipEnv.toSeq,
+ ProcessLogger(_ => (), _ => ())
+ )
if (installCode != 0) {
logger.error(s"failed to install requirements into temp venv
(exit=$installCode)")
return Seq.empty
}
val collected = scala.collection.mutable.ListBuffer[String]()
- val freezeCode = Process(Seq(python, "-m", "pip", "freeze")).!(
+ val freezeCode = runProcess(
+ Seq(python, "-m", "pip", "freeze"),
+ Nil,
ProcessLogger(line => collected += line, _ => ())
)
if (freezeCode != 0) {
@@ -207,7 +223,7 @@ object PveManager extends LazyLogging {
args: Seq[String],
queue: BlockingQueue[String]
): Int = {
- Process(
+ runProcess(
Seq(
python,
"-u",
@@ -218,9 +234,7 @@ object PveManager extends LazyLogging {
"off",
"--no-input"
) ++ args,
- None,
- pipEnv.toSeq: _*
- ).!(
+ pipEnv.toSeq,
ProcessLogger(
out => queue.put(s"[pip] $out"),
err => queue.put(s"[pip][ERR] $err")
@@ -259,7 +273,9 @@ object PveManager extends LazyLogging {
Files.createDirectories(venvDirPath.getParent)
- val createCode = Process(Seq(createVenvPython, "-m", "venv",
venvDirPath.toString)).!(
+ val createCode = runProcess(
+ Seq(createVenvPython, "-m", "venv", venvDirPath.toString),
+ Nil,
ProcessLogger(
out => queue.put(s"[pve] $out"),
err => queue.put(s"[pve][ERR] $err")
@@ -521,7 +537,9 @@ object PveManager extends LazyLogging {
}
try {
- val command = Process(
+ val output = scala.collection.mutable.ListBuffer[String]()
+
+ val exitCode = runProcess(
Seq(
python,
"-u",
@@ -531,13 +549,7 @@ object PveManager extends LazyLogging {
"-y",
trimmedPackageName
),
- None,
- pipEnv.toSeq: _*
- )
-
- val output = scala.collection.mutable.ListBuffer[String]()
-
- val exitCode = command.!(
+ pipEnv.toSeq,
ProcessLogger(
out => {
logger.info(s"[pip] $out")
diff --git
a/amber/src/test/scala/org/apache/texera/web/resource/pythonvirtualenvironment/PveResourceSpec.scala
b/amber/src/test/scala/org/apache/texera/web/resource/pythonvirtualenvironment/PveResourceSpec.scala
index 69c993282d..416bd07a74 100644
---
a/amber/src/test/scala/org/apache/texera/web/resource/pythonvirtualenvironment/PveResourceSpec.scala
+++
b/amber/src/test/scala/org/apache/texera/web/resource/pythonvirtualenvironment/PveResourceSpec.scala
@@ -25,6 +25,7 @@ import
org.apache.texera.dao.jooq.generated.Tables.VIRTUAL_ENVIRONMENTS
import org.apache.texera.dao.jooq.generated.tables.daos.UserDao
import org.apache.texera.dao.jooq.generated.tables.pojos.User
import
org.apache.texera.web.resource.pythonvirtualenvironment.PveResource.SavePvePayload
+import org.scalamock.scalatest.MockFactory
import org.scalatest.{BeforeAndAfterAll, BeforeAndAfterEach}
import org.scalatest.flatspec.AnyFlatSpec
import org.scalatest.matchers.should.Matchers
@@ -34,10 +35,12 @@ import java.util.UUID
import java.util.concurrent.LinkedBlockingQueue
import javax.ws.rs.core.Response
import scala.jdk.CollectionConverters._
+import scala.sys.process.ProcessLogger
class PveResourceSpec
extends AnyFlatSpec
with Matchers
+ with MockFactory
with BeforeAndAfterAll
with BeforeAndAfterEach
with MockTexeraDB {
@@ -48,7 +51,59 @@ class PveResourceSpec
private var testRoot: Path = _
private var queue: LinkedBlockingQueue[String] = _
+ // Exit codes the mock returns for the next venv / pip invocation. Reset to
+ // success in beforeEach; individual tests flip one to force a failure.
+ private var venvExit = 0
+ private var installExit = 0
+ private var uninstallExit = 0
+
+ // What the mocked `pip freeze` reports as the resolved system set. pyarrow
is
+ // always a hard dependency in amber/requirements.txt, so it stands in for "a
+ // system package the user may neither install nor delete".
+ private val systemFreeze = Seq("pyarrow==23.0.1")
+
+ private val realRunner = PveManager.runProcess
+
+ // Mocks every child process PveManager spawns (venv creation, pip
+ // install/uninstall/freeze) so the spec is hermetic — no real venv, no pip,
+ // no network. ScalaMock expectations are per-test, so expectProcessCalls()
is
+ // called at the top of each test that exercises a process. The single
handler
+ // dispatches on the command: a venv create fabricates
<dir>/bin/{python,pip},
+ // freeze emits the system set, install/uninstall just return the configured
+ // exit code. PveManager still owns the metadata files and queue messages.
+ private val runProcessMock =
+ mockFunction[Seq[String], Seq[(String, String)], ProcessLogger, Int]
+
+ private def expectProcessCalls(): Unit =
+ runProcessMock
+ .expects(*, *, *)
+ .onCall { (command: Seq[String], _: Seq[(String, String)], logger:
ProcessLogger) =>
+ if (command.contains("venv")) {
+ if (venvExit == 0) {
+ val bin = Paths.get(command.last).resolve("bin")
+ Files.createDirectories(bin)
+ Seq("python", "pip").foreach { exe =>
+ val f = bin.resolve(exe)
+ Files.write(f, Array.emptyByteArray)
+ f.toFile.setExecutable(true)
+ }
+ }
+ venvExit
+ } else if (command.contains("freeze")) {
+ systemFreeze.foreach(line => logger.out(line))
+ 0
+ } else if (command.contains("uninstall")) {
+ logger.out("mock uninstall")
+ uninstallExit
+ } else if (command.contains("install")) {
+ logger.out("mock install")
+ installExit
+ } else 0
+ }
+ .anyNumberOfTimes()
+
override protected def beforeAll(): Unit = {
+ PveManager.runProcess = runProcessMock
initializeDBAndReplaceDSLContext()
val userDao = new UserDao(getDSLContext.configuration())
val user = new User
@@ -59,9 +114,15 @@ class PveResourceSpec
userDao.insert(user)
}
- override protected def afterAll(): Unit = shutdownDB()
+ override protected def afterAll(): Unit = {
+ PveManager.runProcess = realRunner
+ shutdownDB()
+ }
override protected def beforeEach(): Unit = {
+ venvExit = 0
+ installExit = 0
+ uninstallExit = 0
testPveName = s"testenv${System.currentTimeMillis()}"
testRoot = Paths.get("/tmp/texera-pve/venvs").resolve(testCuid.toString)
queue = new LinkedBlockingQueue[String]()
@@ -80,6 +141,7 @@ class PveResourceSpec
}
"PveManager" should "create a new PVE and list it" in {
+ expectProcessCalls()
PveManager.createNewPve(testCuid, queue, testPveName)
val logs = queueText()
@@ -99,6 +161,7 @@ class PveResourceSpec
}
"PveManager" should "install a user package and list it for the PVE" in {
+ expectProcessCalls()
PveManager.createNewPve(testCuid, queue, testPveName)
val packageName = "colorama"
@@ -129,6 +192,7 @@ class PveResourceSpec
}
"PveManager" should "delete a user package and remove it from the PVE
package list" in {
+ expectProcessCalls()
PveManager.createNewPve(testCuid, queue, testPveName)
val packageName = "colorama"
@@ -167,7 +231,64 @@ class PveResourceSpec
pve.get.userPackages should not contain packageSpec
}
+ "PveManager" should "report an error when venv creation fails" in {
+ expectProcessCalls()
+ venvExit = 1
+
+ PveManager.createNewPve(testCuid, queue, testPveName)
+
+ val logs = queueText()
+ logs should include("[PVE][ERR] Failed to create venv")
+ Files.exists(testRoot.resolve(testPveName).resolve("pve")) shouldBe false
+ }
+
+ it should "report an error when the system requirements install fails" in {
+ expectProcessCalls()
+ installExit = 1
+
+ PveManager.createNewPve(testCuid, queue, testPveName)
+
+ val logs = queueText()
+ logs should include("[PVE][ERR] Failed to install requirements files")
+ }
+
+ it should "refuse to install a package that is part of the system set" in {
+ expectProcessCalls()
+ PveManager.createNewPve(testCuid, queue, testPveName)
+ queue.clear()
+
+ PveManager.installUserPackages(List("pyarrow==23.0.1"), testCuid, queue,
testPveName)
+
+ val logs = queueText()
+ logs should include("[PVE][ERR] pyarrow==23.0.1 is a system package")
+
+ PveManager
+ .getEnvironments(testCuid)
+ .find(_.pveName == testPveName)
+ .get
+ .userPackages should not contain "pyarrow==23.0.1"
+ }
+
+ it should "report an error when a user package install fails" in {
+ expectProcessCalls()
+ PveManager.createNewPve(testCuid, queue, testPveName)
+ installExit = 1
+ queue.clear()
+
+ PveManager.installUserPackages(List("colorama==0.4.6"), testCuid, queue,
testPveName)
+
+ val logs = queueText()
+ logs should include("[PVE][ERR] Failed to install package:
colorama==0.4.6")
+
+ PveManager
+ .getEnvironments(testCuid)
+ .find(_.pveName == testPveName)
+ .get
+ .userPackages should not contain "colorama==0.4.6"
+ }
+
"PveManager" should "delete all PVEs for a computing unit" in {
+ expectProcessCalls()
PveManager.createNewPve(testCuid, queue, testPveName)
Files.exists(testRoot.resolve(testPveName)) shouldBe true
@@ -179,6 +300,7 @@ class PveResourceSpec
}
"PveManager.getPythonBin" should "return Some for an existing venv" in {
+ expectProcessCalls()
PveManager.createNewPve(testCuid, queue, testPveName)
val result = PveManager.getPythonBin(testCuid, testPveName)