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-7800-54fba318c15b38d128e7d28969bd18c72f8fd206 in repository https://gitbox.apache.org/repos/asf/texera.git
commit 70bd3b3caef58febd250765e107811718aa5a5db Author: Xinyuan Lin <[email protected]> AuthorDate: Wed Sep 23 03:58:24 2026 +0000 fix(test): stop the file-handle leak aborting the operator suite, seed the interval spec (#7800) ### What changes were proposed in this PR? Four test-only fixes across 3 spec files. **89 insertions, 25 deletions; `git diff -- '*/src/main/*'` is empty.** ### 1. A leaked file handle was aborting the operator suite and killing the module's coverage run The exact mechanism, established by instrumenting the run rather than by inspection: in `"create LargeBinary when reading file with LARGE_BINARY type"`, `createTuplesFromFile` throws mid-`map` with `IllegalStateException: LargeBinaryManager.create() requires a base URI, but none was set on the current thread`. That abandons the `AutoClosingIterator` before exhaustion, so its close-on-exhaustion hook never fires. `afterAll` then fails with `test_large_binary.txt: The process cannot access the file because it is being used by another process`, aborting the suite. The blast radius is the whole module: | | before | after | |---|---|---| | Suites | 289 completed, **1 aborted** | **290 completed, 0 aborted** | | sbt exit | 1 | **0** | | jacoco report | **no report directory emitted** | `jacoco.xml`, 1.8 MB | | Tests | 2301 succeeded | 2301 succeeded | So `WorkflowOperator/jacoco` could not produce coverage for *any* file in the module on Windows. POSIX `unlink` masks this on Linux CI. **The fix is in the test**: the read is wrapped so the iterator is drained in a `finally`, firing the close hook on both the success and failure paths. `AutoClosingIterator` and `createTuplesFromFile` are untouched — the leak is a test that stops early, not a broken production contract. Also worth recording: the `.zip` fixture was stranded too. It only *looked* fine because `.gitignore:10` is `*.zip`; `afterAll` never reached it, because the `.txt` delete threw first. ### 2. `IntervalOpExecSpec` is now deterministic It imported `scala.util.Random.{nextInt, nextLong}` and used them at four sites — input ordering twice, a 1k-row dataset, and the interval constant — so `IntervalJoinOpExec`'s coverage footprint drifted between runs of identical source (CI has reported 2 missed + 14 partial where a local run gave 0 + 13). Now a fixed `Seed`, with a fresh `new Random(Seed)` created at each use site so determinism does not depend on test execution order either. No assertion changed. Verified rather than assumed: instrumented to print every generated input, two runs produced byte-identical output (22 lines, empty `diff`), and **two full `WorkflowOperator/jacoco` runs now give byte-identical `IntervalJoinOpExec` counters and all 110 per-line entries.** One observation left alone as out of scope: the 1k test deterministically yields 0 matches, because random 64-bit longs essentially never fall inside a sub-1000-wide window, so `outputTuples.size == bruteForceResult` is `0 == 0`. That was equally true before this change. ### 3. `PveResourceSpec`'s traversal assertion said something it did not test `getPythonBin(testCuid, "..") shouldBe None` appeared to pin the guard at `PveManager.scala:91`, but `".."` **matches** the name regex on line 88 (`^[A-Za-z0-9._-]+$` admits dots), and `<VenvRoot>/<cuid>/../pve` normalizes to `<VenvRoot>/pve` — still under root. The `None` came from the `Files.exists` check on line 92. **Proven, not argued:** deleting line 91 produced results identical to baseline (43 run, 37 succeeded, 6 failed, 1 canceled) and the traversal test still **passed**. Nothing in the spec detects that guard's removal. The production file was reverted afterwards. Split into two honest cases: - `"reject pveNames containing a path separator"` keeps `"../../../etc"` and `"foo/bar"` — both carry a `/`, so they genuinely pin the **name regex**. - `"return None for a dot-only pveName, which has no venv"` keeps the `".."` case and pins what it actually pins: a name with no venv on disk yields `None`. A comment records that line 91 is **unreachable by construction** — the regex forbids `/` and `cuid` is an `Int` — so it is defensive code rather than untested code. The guard is left in place; removing it is a production decision, not a test one. ### 4. The same spec's LARGE_BINARY test could not fail Listed below as a deferred defect in the first revision of this PR; fixed here instead, since it is the same test and the leak fix above is what makes its failure path reachable at all. The test wrapped its whole body in `catch { case e: Exception => info(...) }`. ScalaTest's `assert` throws `TestFailedException`, which extends `Exception`, so every assertion was swallowed and reported only through `info(...)` — which this build's `-u`-only reporter (`build.sbt:37`) does not surface. **The test could not fail for any reason.** What it was really tolerating is the environment, and in two places rather than one: | Environment gap | Where it throws | | --- | --- | | No large-binary base URI bound to the thread | `LargeBinaryManager.create()` → `IllegalStateException` | | No reachable S3 endpoint | `LargeBinaryOutputStream.close()` → `IOException` | The second one matters more than it looks: `WorkflowOperator` is wired as a plain `dependsOn(WorkflowCore)` with no testcontainers in its test scope, so **there is no configuration of this module in which those assertions could ever have run** — not on a developer machine, not in CI. Both are now handled explicitly instead of by catching everything: - **Seed the base URI** the way the coordinator does in production, and the way `LargeBinaryManagerSpec` does in test (`setCurrentBaseUri(baseUriForExecution(eid))`), cleared in a `finally` so it cannot leak into another suite reusing the thread. That gap stops being an environment gap. - **Probe S3 up front** and `cancel` when it is absent, which leaves the scan itself **uncaught**. A catch around the scan — even narrowed to `IOException` — would still report a genuine operator failure as a skip. `directoryExists` was rejected as the probe: it throws `NoSuchBucketException` on a fresh bucket, which would skip the test on an environment where S3 *is* present. - **Assertions moved outside the try**, so a `TestFailedException` propagates. The drain from §1 stays, and now also drops a throw of its own, so it cannot displace an in-flight failure via the `finally`. **Proven by mutation, not argued:** changing `startsWith("s3://")` to `startsWith("gs://")` now produces `1 TEST FAILED`, sbt exit 1. Before this change the identical mutation still passed. Worth recording for anyone relying on `cancel` for visibility: **`-u` erases it.** A cancelled test is written to the JUnit XML as a plain `<testcase>` with no `<skipped/>` child and `skipped="0"` on the suite — byte-identical to a pass, so Codecov cannot tell "skipped, no S3" from "green". The console summary's `canceled 1` names no test. Hence the one `Console.err.println` on the cancel path: it is the only channel that reports *which* test skipped and *why*. ### Verification - Full `WorkflowOperator/jacoco`: **290 suites, 0 aborted, sbt exit 0**, report emitted. 2300 succeeded + 1 cancelled; the cancelled one is §4's test on a machine with no S3 endpoint, which previously counted as a vacuous pass. - §4's assertions were exercised against a working destination, since Docker was unavailable locally: a temporary subclass of the spec pointed `StorageConfig.s3Endpoint` at a JDK `HttpServer` stub, so the **real** assertions ran rather than copies. **7 succeeded, 0 cancelled**, with the stub observing the genuine `HEAD` → `POST ?uploads` → `PUT ?partNumber=1` → `POST ?uploadId` sequence carrying the file's bytes. Harness deleted; it is not in this diff. - `PveResourceSpec`: 44 run, 38 succeeded — both reshaped tests pass. The 6 failures are pre-existing and environmental (the mock fabricates `bin/python` while `PveManager` looks for `Scripts/python.exe` on Windows); baseline shows the identical 6. - `git diff -- '*/src/main/*'` is empty; no stranded fixtures (checked with `git status --ignored`). - Lint green: `WorkflowOperator/Test/scalafmtCheck`, `WorkflowOperator/scalafixAll --check`, `WorkflowExecutionService/Test/scalafmtCheck`, `WorkflowExecutionService/scalafixAll --check`. ### Any related issues, documentation, discussions? Closes #7799 ### How was this PR tested? ``` sbt "WorkflowOperator/jacoco" ``` ``` [info] Suites: completed 290, aborted 0 [info] Tests: succeeded 2300, failed 0, canceled 1, ignored 0, pending 2 ``` ### Was this PR authored or co-authored using generative AI tooling? Generated-by: Claude Code (Opus 5) --- .../pythonvirtualenvironment/PveResourceSpec.scala | 18 +++++- .../operator/intervalJoin/IntervalOpExecSpec.scala | 21 ++++-- .../scan/file/FileScanSourceOpExecSpec.scala | 75 +++++++++++++++++----- 3 files changed, 89 insertions(+), 25 deletions(-) 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 dd59386e64..4e004f7b13 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 @@ -407,12 +407,26 @@ class PveResourceSpec PveManager.getPythonBin(testCuid, "no-such-env") shouldBe None } - it should "reject pveNames containing path-traversal segments" in { - PveManager.getPythonBin(testCuid, "..") shouldBe None + it should "reject pveNames containing a path separator" in { + // Both names carry a '/', so SafePveName rejects them before any path is built. PveManager.getPythonBin(testCuid, "../../../etc") shouldBe None PveManager.getPythonBin(testCuid, "foo/bar") shouldBe None } + it should "return None for a dot-only pveName, which has no venv" in { + // This does NOT exercise the `!resolved.startsWith(root)` guard in + // getPythonBin. SafePveName's character class admits '.', so ".." is accepted, + // and <VenvRoot>/<cuid>/../pve normalises to <VenvRoot>/pve — still under the + // root. The None below therefore comes from the Files.exists check, exactly as + // it does for any other name with no venv on disk. + // + // That guard is unreachable by construction rather than untested: SafePveName + // forbids '/', so no accepted name can leave the root, and cuid is an Int. + // Deleting the guard leaves this spec entirely green, so no assertion here can + // pin it; making it reachable (or dropping it) is a production decision. + PveManager.getPythonBin(testCuid, "..") shouldBe None + } + it should "reject pveNames with disallowed characters" in { PveManager.getPythonBin(testCuid, "") shouldBe None PveManager.getPythonBin(testCuid, "name with spaces") shouldBe None diff --git a/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/intervalJoin/IntervalOpExecSpec.scala b/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/intervalJoin/IntervalOpExecSpec.scala index cfc3f360e6..9063105e0a 100644 --- a/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/intervalJoin/IntervalOpExecSpec.scala +++ b/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/intervalJoin/IntervalOpExecSpec.scala @@ -28,11 +28,18 @@ import org.scalatest.flatspec.AnyFlatSpec import java.sql.Timestamp import scala.collection.mutable.ArrayBuffer -import scala.util.Random.{nextInt, nextLong} +import scala.util.Random class IntervalOpExecSpec extends AnyFlatSpec with BeforeAndAfter { val left: Int = 0 val right: Int = 1 + // The generated inputs, the order they are fed in, and the interval constant all + // used to come from the global unseeded `scala.util.Random`, so every run drove + // the operator down a different set of branches. Each generator below is seeded + // from this constant and created fresh at its use site, so the inputs are + // identical on every run and independent of test execution order. + private val Seed: Long = 20260819L + var opDesc: IntervalJoinOpDesc = _ var counter: Int = 0 @@ -240,8 +247,9 @@ class IntervalOpExecSpec extends AnyFlatSpec with BeforeAndAfter { counter = 0 var leftIndex: Int = 0 var rightIndex: Int = 0 - val leftOrder = LazyList.continually(nextInt(10)).take(leftInput.length).toList - val rightOrder = LazyList.continually(nextInt(10)).take(rightInput.length).toList + val orderRandom = new Random(Seed) + val leftOrder = LazyList.continually(orderRandom.nextInt(10)).take(leftInput.length).toList + val rightOrder = LazyList.continually(orderRandom.nextInt(10)).take(rightInput.length).toList val outputTuples: ArrayBuffer[Tuple] = new ArrayBuffer[Tuple] while (leftIndex < leftOrder.size || rightIndex < rightOrder.size) { @@ -483,8 +491,9 @@ class IntervalOpExecSpec extends AnyFlatSpec with BeforeAndAfter { } it should "test larger dataset(1k)" in { - val pointList: Array[Long] = LazyList.continually(nextLong()).take(1000).toArray - val rangeList: Array[Long] = LazyList.continually(nextLong()).take(1000).toArray + val dataRandom = new Random(Seed) + val pointList: Array[Long] = LazyList.continually(dataRandom.nextLong()).take(1000).toArray + val rangeList: Array[Long] = LazyList.continually(dataRandom.nextLong()).take(1000).toArray testJoin[Long]( "point", "range", @@ -492,7 +501,7 @@ class IntervalOpExecSpec extends AnyFlatSpec with BeforeAndAfter { includeRightBound = true, AttributeType.LONG, TimeIntervalType.DAY, - nextInt(1000).toLong, + dataRandom.nextInt(1000).toLong, pointList, rangeList ) diff --git a/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/source/scan/file/FileScanSourceOpExecSpec.scala b/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/source/scan/file/FileScanSourceOpExecSpec.scala index 206cc33f32..36c8ffb4e9 100644 --- a/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/source/scan/file/FileScanSourceOpExecSpec.scala +++ b/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/source/scan/file/FileScanSourceOpExecSpec.scala @@ -22,8 +22,10 @@ package org.apache.texera.amber.operator.source.scan.file import org.apache.texera.amber.core.tuple.{AttributeType, LargeBinary, Schema, SchemaEnforceable} import org.apache.texera.amber.operator.source.scan.{FileAttributeType, FileDecodingMethod} import org.apache.texera.amber.util.JSONUtils.objectMapper +import org.apache.texera.service.util.{LargeBinaryManager, S3StorageClient} import org.scalatest.BeforeAndAfterAll import org.scalatest.flatspec.AnyFlatSpec +import software.amazon.awssdk.core.exception.SdkException import java.io.{BufferedOutputStream, FileOutputStream} import java.net.URI @@ -36,6 +38,9 @@ import java.util.zip.{ZipEntry, ZipOutputStream} */ class FileScanSourceOpExecSpec extends AnyFlatSpec with BeforeAndAfterAll { + /** Execution id whose object prefix the large binaries created here are written under. */ + private val testExecutionId: Long = 8888L + private val testDir = Path .of(sys.env.getOrElse("TEXERA_HOME", ".")) .resolve("common/workflow-operator/src/test/resources") @@ -122,28 +127,64 @@ class FileScanSourceOpExecSpec extends AnyFlatSpec with BeforeAndAfterAll { // Execution Tests it should "create LargeBinary when reading file with LARGE_BINARY type" in { + // FileScanUtils streams every LARGE_BINARY field to S3 through LargeBinaryOutputStream, so + // this test can only run where an S3 endpoint is reachable, and WorkflowOperator's test + // scope provides none (LargeBinaryManagerSpec gets one from testcontainers). Probe for the + // endpoint up front and cancel when it is absent, leaving the scan below uncaught: a catch + // around the scan would report a genuine operator failure as a skip, which is how this + // test used to be unable to fail at all. + try S3StorageClient.createBucketIfNotExist(LargeBinaryManager.DEFAULT_BUCKET) + catch { + case e: SdkException => + // On stderr because `-u` is the only ScalaTest reporter this build configures, so + // info() and alert() go nowhere -- and the JUnit XML records a cancelled test as a + // plain passing <testcase>, leaving the console's "canceled 1" as the sole other clue. + Console.err.println(s"[FileScanSourceOpExecSpec] skipping LARGE_BINARY read: $e") + cancel("no reachable S3 endpoint in this test scope", e) + } + val desc = createDescriptor() desc.setResolvedFileName(URI.create(testFile.toUri.toString)) val executor = new FileScanSourceOpExec(objectMapper.writeValueAsString(desc)) - try { - executor.open() - val tuples = executor.produceTuple().toSeq - executor.close() - - assert(tuples.size == 1) - val field = tuples.head - .asInstanceOf[SchemaEnforceable] - .enforceSchema(desc.sourceSchema()) - .getField[Any]("line") - - assert(field.isInstanceOf[LargeBinary]) - assert(field.asInstanceOf[LargeBinary].getUri.startsWith("s3://")) - } catch { - case e: Exception => - info(s"S3 not configured: ${e.getMessage}") - } + // FileScanUtils mints each LARGE_BINARY field through LargeBinaryManager.create(), which + // reads a per-execution base URI off the calling thread. Seed it the way the coordinator + // does in production (and LargeBinaryManagerSpec does in test) so an unset base URI can + // never be mistaken for the operator misbehaving. + LargeBinaryManager.setCurrentBaseUri(LargeBinaryManager.baseUriForExecution(testExecutionId)) + val tuples = + try { + executor.open() + val rows = executor.produceTuple() + try { + rows.toSeq + } finally { + // `produceTuple` hands back an AutoClosingIterator, which releases the underlying + // file handle only once `hasNext` turns false, so a scan that throws part way + // abandons the iterator with the handle still open. Windows then refuses to delete + // `testFile` in `afterAll`, which aborts the whole suite and buries the very failure + // that caused it; POSIX `unlink` hides the same leak on Linux CI. Draining fires the + // close hook on both paths, and a throw from the drain itself is dropped so it + // cannot displace an in-flight failure. + try while (rows.hasNext) rows.next() + catch { case _: Exception => } + executor.close() + } + } finally { + // An empty value clears the thread-local, so the seeded base URI cannot leak into + // another suite that happens to reuse this thread. + LargeBinaryManager.setCurrentBaseUri("") + } + + assert(tuples.size == 1) + val field = tuples.head + .asInstanceOf[SchemaEnforceable] + .enforceSchema(desc.sourceSchema()) + .getField[Any]("line") + + assert(field.isInstanceOf[LargeBinary]) + assert(field.asInstanceOf[LargeBinary].getUri.startsWith("s3://")) } // LargeBinary Tests
