Yicong-Huang commented on code in PR #7207:
URL: https://github.com/apache/texera/pull/7207#discussion_r3741403525


##########
common/workflow-operator/src/test/scala/org/apache/texera/amber/util/PythonCodeRawInvalidTextSpec.scala:
##########
@@ -49,8 +56,81 @@ final class PythonCodeRawInvalidTextSpec extends AnyFunSuite 
{
   private val MaxDepth: Int = 3
   private val AcceptPackages: Seq[String] = 
Seq("org.apache.texera.amber.operator")
 
+  /** Budget for one whole fanned-out pass over every descriptor. Deliberately 
far
+    * above a real run (well under a second), so it only ever fires on a hang.
+    */
+  private val PassTimeout: FiniteDuration = 10.minutes
+
+  /** Runs the given work concurrently and returns the results in submission
+    * order, rethrowing the first failure so it fails the test. Sized to the 
pool
+    * so the threads match the workers available to serve them.
+    *
+    * Daemon threads: a task parked on a subprocess pipe answers no interrupt, 
so
+    * `shutdownNow` need not end it, and a non-daemon one left there would 
hold the
+    * JVM — and the build — open after [[PassTimeout]] has already failed the 
test.
+    */
+  private def awaitAll[T](work: Seq[() => T]): Seq[T] = {
+    val threads = Executors.newFixedThreadPool(
+      PythonWorkerPool.maxWorkers,
+      (r: Runnable) => {
+        val t = new Thread(r, "py-compile-check")
+        t.setDaemon(true)
+        t
+      }
+    )
+    try {
+      implicit val ec: ExecutionContext = 
ExecutionContext.fromExecutorService(threads)
+      Await.result(Future.sequence(work.map(w => Future(w()))), PassTimeout)
+    } finally threads.shutdownNow()
+  }
+
+  /** Syntax-checks one generated module, through a pooled worker when 
available.
+    *
+    * The worker is launched with the same `-I -S` isolation the one-shot path
+    * uses, so what the check accepts is unchanged; it just stops paying an
+    * interpreter boot — the whole cost of a check whose real work is under a
+    * millisecond — once per descriptor. A hard worker crash falls back to the
+    * spawn, so behavior is never worse than before the pool.
+    */
+  private def syntaxCheck(
+      pythonExecutable: String,
+      pythonSource: String,
+      descriptorName: String
+  ): Either[String, Unit] = {
+    def viaPool: Either[String, Unit] = {
+      val request = objectMapper.createObjectNode()
+      request.put("source", pythonSource)
+      request.put("name", s"$descriptorName.py")
+      val outcome = PythonWorkerPool.run(
+        resourcePath = "/python/py_compile_worker.py",
+        launchArgs = Seq.empty,
+        pythonExe = pythonExecutable,
+        request = request,
+        interpreterArgs = Seq("-I", "-S")
+      )
+      if (outcome.exit == 0) Right(())
+      else {
+        val output = if (outcome.stderr.trim.nonEmpty) outcome.stderr.trim 
else "(no output)"
+        Left(
+          s"py_compile failed (exit=${outcome.exit})\nOutput:\n" +
+            truncateBlock(output, maxLines = 40, maxChars = 8000)
+        )
+      }
+    }
+
+    if (PythonWorkerPool.enabled) {
+      try viaPool
+      catch {
+        case _: PythonWorkerPool.WorkerDiedException =>
+          pyCompile(pythonExecutable, pythonSource)
+      }

Review Comment:
   Pool startup can fail without a `WorkerDiedException`: `pb.start()` throws a 
bare `IOException` (`PythonWorkerPool.scala:250`), and `borrow` rethrows it 
verbatim. That escapes this fallback and aborts the suite. The pre-pool code 
wrapped the same call in `Try` and returned a named `Left` (:160-163), so the 
"never worse than before" claim at :93 does not hold here. `error=11` under 
this PR's 4-way concurrency is a live way to reach it.
   
   Also: the fallback is silent, so a pool whose workers all die still reports 
green.
   
   ```suggestion
         catch {
           case _: PythonWorkerPool.WorkerDiedException | _: 
java.io.IOException =>
             pyCompile(pythonExecutable, pythonSource)
         }
   ```



##########
common/workflow-operator/build.sbt:
##########
@@ -35,6 +35,18 @@ ThisBuild / conflictManager := ConflictManager.latestRevision
 // Restrict parallel execution of tests to avoid conflicts
 Global / concurrentRestrictions += Tags.limit(Tags.Test, 1)
 
+// A test needing more than a bare Python interpreter is tagged, so the amber 
job
+// excludes it and amber-integration, which installs operator-requirements.txt,
+// runs it. The amber job already sets this env var on the step that invokes
+// WorkflowOperator/jacoco, so no workflow change is needed for the exclusion.
+// -P4 bounds ScalaTest's ParallelTestExecution pool on the integration side, 
whose
+// tagged tests drive Python subprocesses. It matches PythonWorkerPool's 
default
+// worker cap so the two bounds agree rather than multiply.
+Test / testOptions ++= TestFilters.integrationSplit(
+  envVar = "AMBER_TEST_FILTER",
+  tag = "org.apache.texera.amber.operator.tags.IntegrationTest",
+  integrationOnlyExtra = Seq("-P4")

Review Comment:
   `-P4` is inert here. No suite in the module mixes in `ParallelTestExecution` 
— grep finds the string only in the comment above — and the single tagged test 
is the pandas/plotly probe, which spawns its own subprocess and never touches 
the pool.
   
   The rationale is also wrong on its own terms. `maxWorkers` caps each 
*sub-pool*, keyed by script/args/env, so parallel suites on distinct keys 
multiply rather than agree. It is env-overridable too, while `-P4` is a 
literal. I would drop the parameter and add it back with the first pooled 
tagged test.



##########
common/workflow-operator/src/test/scala/org/apache/texera/amber/util/python/PythonWorkerPool.scala:
##########
@@ -0,0 +1,445 @@
+/*
+ * 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.amber.util.python
+
+import com.fasterxml.jackson.databind.node.ObjectNode
+import com.typesafe.scalalogging.LazyLogging
+import org.apache.texera.amber.util.JSONUtils.objectMapper
+
+import java.io.{BufferedReader, BufferedWriter, InputStreamReader, 
OutputStreamWriter}
+import java.nio.charset.StandardCharsets
+import java.nio.file.{Files, Path, StandardCopyOption}
+import java.util.concurrent.{
+  Callable,
+  ConcurrentHashMap,
+  ExecutionException,
+  ExecutorService,
+  Executors,
+  LinkedBlockingQueue,
+  TimeUnit,
+  TimeoutException
+}
+import java.util.concurrent.atomic.AtomicInteger
+import scala.annotation.tailrec
+import scala.collection.mutable
+import scala.jdk.CollectionConverters._
+import scala.util.control.NonFatal
+
+/**
+  * Pools of persistent Python "worker" processes that eliminate the per-call
+  * interpreter-boot + import cost a test otherwise pays on every subprocess
+  * spawn. Testing operators one at a time does not scale when each one costs a
+  * spawn: a bare `-I -S` interpreter boots in ~25 ms, and once pandas and 
plotly
+  * are imported a spawn costs ~260-310 ms — ~96% of a job whose real work is
+  * ~4 ms. A worker pays that once at startup, then serves many jobs over its
+  * lifetime, so N spawns become one.
+  *
+  * Lives in test scope here, rather than beside a single caller, because 
tests in
+  * several modules run generated operator code and would otherwise each
+  * hand-roll a driver, a stdout protocol and a timeout. Other modules reach it
+  * through a `test->test` dependency on this one.
+  *
+  * Generic over the worker script — the pool never interprets the payload — so
+  * one implementation serves a syntax check, template execution and DataFrame
+  * comparison alike. Each distinct (resource, interpreterArgs, launchArgs,
+  * python, env) combination gets its own sub-pool.
+  *
+  * Protocol (line-delimited JSON, shared by all worker scripts):
+  *   startup   worker -> pool:  {"ready": true}
+  *   request   pool -> worker:  <caller-supplied JSON object>\n
+  *   response  worker -> pool:  {"exit": <int>, "stdout": "...", "stderr": 
"..."}\n
+  *
+  * Concurrency: callers submit from several threads at once — a spec run with
+  * ScalaTest's `-P4`, or one test fanning its cases out — so each sub-pool 
holds
+  * up to [[maxWorkers]] workers, each serving one job at a time (borrow -> run
+  * -> return). A worker script may chdir per job, so a worker must never run 
two
+  * jobs at once — the borrow/return discipline guarantees that.
+  *
+  * Robustness: an ordinary *job* failure comes back as an [[Outcome]] with
+  * `exit != 0` (worker keeps running). A hard interpreter crash ends a worker;
+  * the pool detects the EOF / broken pipe, discards it, and throws
+  * [[WorkerDiedException]] so the caller can fall back to a one-shot 
subprocess
+  * — behavior is then never worse than the pre-pool path. A worker that stays
+  * alive but stops answering ends the same way, on the [[Timeouts]] below.
+  */
+object PythonWorkerPool extends LazyLogging {
+
+  /** Worker response: process-like exit code plus captured streams. */
+  final case class Outcome(exit: Int, stdout: String, stderr: String)
+
+  /** Thrown when a worker dies mid-job (hard crash / broken pipe). Callers
+    * catch this and fall back to a one-shot subprocess.
+    */
+  final class WorkerDiedException(message: String, cause: Throwable = null)
+      extends RuntimeException(message, cause)
+
+  /** Feature toggle. `TEXERA_TEST_PYTHON_WORKER=0` (or `false`/`off`) forces 
the
+    * one-subprocess-per-call paths everywhere — an escape hatch for debugging 
a
+    * suspected isolation leak. Default on.
+    */
+  val enabled: Boolean =
+    !sys.env
+      .get("TEXERA_TEST_PYTHON_WORKER")
+      .map(_.trim.toLowerCase)
+      .exists(Set("0", "false", "off"))
+
+  /** Max live workers per sub-pool. Defaults to 4 to match ScalaTest's `-P4`, 
so
+    * the two concurrency bounds agree on how many interpreters may be live;
+    * override via `TEXERA_TEST_PYTHON_WORKERS`. Public so a caller fanning out
+    * jobs within one test can size that fan-out to the workers it will get.
+    */
+  val maxWorkers: Int =
+    sys.env
+      .get("TEXERA_TEST_PYTHON_WORKERS")
+      .flatMap(s => scala.util.Try(s.trim.toInt).toOption)
+      .filter(_ > 0)
+      .getOrElse(4)
+
+  /** How long a caller waits on a worker before the pool kills and discards 
it.
+    * A read on a process pipe cannot be interrupted — a suite or executor 
timeout
+    * leaves the reading thread stuck on it — so a worker that stays alive 
without
+    * answering has to be bounded here. `response` keeps the 30 seconds the
+    * one-shot spawn this pool replaced allowed a job; `startup` is longer 
because
+    * a worker imports its libraries before it reports ready, and a loaded CI
+    * machine makes that slow. Override in seconds via
+    * `TEXERA_TEST_PYTHON_WORKER_TIMEOUT` / 
`TEXERA_TEST_PYTHON_WORKER_STARTUP_TIMEOUT`.
+    */
+  final case class Timeouts(responseMillis: Long, startupMillis: Long)
+
+  object Timeouts {
+    private def envSeconds(name: String, default: Long): Long =
+      sys.env
+        .get(name)
+        .flatMap(s => scala.util.Try(s.trim.toLong).toOption)
+        .filter(_ > 0)
+        .getOrElse(default) * 1000
+
+    val Default: Timeouts = Timeouts(
+      responseMillis = envSeconds("TEXERA_TEST_PYTHON_WORKER_TIMEOUT", 30),
+      startupMillis = envSeconds("TEXERA_TEST_PYTHON_WORKER_STARTUP_TIMEOUT", 
60)
+    )
+  }
+
+  /**
+    * Run one job through a pooled worker for `resourcePath`, launched as
+    * `pythonExe <interpreterArgs> <script> <launchArgs>` with extra 
environment
+    * `env`. `request` is the worker-specific JSON payload (the pool does not
+    * interpret it). Throws [[WorkerDiedException]] on a hard worker crash.
+    *
+    * `interpreterArgs` are the flags that must precede the script — a syntax
+    * checker wants `-I -S` so it validates under the same isolation a one-shot
+    * `python -I -S -m py_compile` gave it. `launchArgs` are the script's own
+    * (e.g. `--serve`), and `env` carries what a flag cannot (e.g. PYTHONPATH).
+    * All three are part of a worker's identity: one started differently is not
+    * interchangeable, so it gets its own sub-pool. `timeouts` is not — it 
bounds
+    * this call, so a caller whose jobs are slower than most can raise it 
without
+    * splitting the pool.
+    */
+  def run(
+      resourcePath: String,
+      launchArgs: Seq[String],
+      pythonExe: String,
+      request: ObjectNode,
+      env: Map[String, String] = Map.empty,
+      interpreterArgs: Seq[String] = Seq.empty,
+      timeouts: Timeouts = Timeouts.Default
+  ): Outcome = {
+    val pool = pools.computeIfAbsent(
+      Key(resourcePath, pythonExe, interpreterArgs.toList, launchArgs.toList, 
env.toList.sorted),
+      _ => new Pool(resourcePath, launchArgs, pythonExe, env, interpreterArgs)
+    )
+    pool.run(request, timeouts)
+  }
+
+  /** What makes two launches the same worker. Compared field by field rather 
than
+    * as one joined string, so a value carrying whatever the separator was — a
+    * python path with a space in it, a PYTHONPATH — cannot make two different
+    * launches share a pool.
+    */
+  private final case class Key(
+      resourcePath: String,
+      pythonExe: String,
+      interpreterArgs: List[String],
+      launchArgs: List[String],
+      env: List[(String, String)]
+  )
+
+  /** How long a caller at the worker cap waits before re-examining it. Not a
+    * deadline — see [[Pool.borrow]].
+    */
+  private val CapRecheckMillis: Long = 250
+
+  private val pools = new ConcurrentHashMap[Key, Pool]()
+
+  Runtime.getRuntime.addShutdownHook(new Thread(() => shutdownAll()))
+
+  private def shutdownAll(): Unit =
+    pools.values().forEach(_.shutdown())
+
+  // A single sub-pool: up to `maxWorkers` live workers for one worker script.
+  private final class Pool(
+      resourcePath: String,
+      launchArgs: Seq[String],
+      pythonExe: String,
+      env: Map[String, String],
+      interpreterArgs: Seq[String]
+  ) {
+    private val idle = new LinkedBlockingQueue[Worker]()
+    private val liveCount = new AtomicInteger(0)
+    private val all = mutable.Set.empty[Worker] // guarded by `all`
+    @volatile private var script: Path = _
+
+    def run(request: ObjectNode, timeouts: Timeouts): Outcome = {
+      val w = borrow(timeouts)
+      try {
+        val outcome = w.run(request, timeouts.responseMillis)
+        idle.offer(w) // healthy — return to pool
+        outcome
+      } catch {
+        case e: WorkerDiedException =>
+          discard(w)
+          throw e
+      }

Review Comment:
   `Worker.run` converts `NonFatal` into `WorkerDiedException`. But 
`scala.util.control.NonFatal` excludes `InterruptedException`, and both 
blocking calls on that path throw exactly that (`lines.poll` :415, `write.get` 
:392).
   
   Such a worker stays in `all`, never returns to `idle`, and `liveCount` is 
never decremented — the slot is gone for the life of the JVM. After 
`maxWorkers` losses `borrow` spins its 250 ms cap loop with no log until the 
10-minute `PassTimeout`. The interrupt comes from `awaitAll`'s own `finally 
threads.shutdownNow()`.
   
   ```suggestion
         } catch {
           case e: Throwable =>
             discard(w)
             throw e
         }
   ```



##########
common/workflow-operator/src/test/scala/org/apache/texera/amber/util/python/PythonWorkerPoolSpec.scala:
##########
@@ -0,0 +1,172 @@
+/*
+ * 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.amber.util.python
+
+import org.apache.texera.amber.util.JSONUtils.objectMapper
+import org.scalatest.funsuite.AnyFunSuite
+
+import java.util.concurrent.{Executors, TimeUnit}
+import scala.concurrent.duration._
+import scala.concurrent.{Await, ExecutionContext, 
ExecutionContextExecutorService, Future}
+import scala.util.Try
+
+/**
+  * What the pool owes a caller when a worker misbehaves. An ordinary job 
failure
+  * is the other suites' business; this one is about a worker that stays alive 
and
+  * stops taking part, which is the case that does not end by itself: a crash
+  * closes the pipe and the pending read returns, while silence would hold the
+  * caller forever — neither a read nor a write on a process pipe answers an
+  * interrupt or a deadline, so no suite-level timeout can release one.
+  *
+  * Every wait here is bounded and runs on daemon threads, so a regression 
fails
+  * these tests instead of wedging the run: a lost non-daemon thread parked on 
a
+  * pipe would keep the JVM, and the build, alive.
+  *
+  * The fixture worker is stdlib-only and runs under `-I -S`, so this needs an
+  * interpreter but none of the operator packages.
+  */
+final class PythonWorkerPoolSpec extends AnyFunSuite {
+
+  private val HangingWorker = "/python/hanging_worker.py"
+  private val CompileWorker = "/python/py_compile_worker.py"
+
+  /** Short enough to keep the suite quick, far enough above process startup to
+    * not be mistaken for one: the fixture never answers, so it cannot race.
+    */
+  private val ShortTimeouts: PythonWorkerPool.Timeouts =
+    PythonWorkerPool.Timeouts(responseMillis = 1500, startupMillis = 1500)
+
+  /** Ceiling on a whole case, well above the deadlines under test. Reaching it
+    * means something never gave up.
+    */
+  private val Bound: FiniteDuration = 25.seconds
+
+  /** Any interpreter serves — the fixture imports only `json` and `time` — so
+    * this deliberately skips the configured `python.path` the suites that need
+    * pandas resolve. A machine without one cancels rather than fails.
+    */
+  private def python(): String = {
+    def isRunnable(exe: String): Boolean =
+      Try(new ProcessBuilder(exe, 
"--version").redirectErrorStream(true).start()).toOption
+        .exists { p =>
+          if (p.waitFor(5, TimeUnit.SECONDS)) p.exitValue() == 0 else { 
p.destroyForcibly(); false }
+        }
+
+    List("python3", "python", "py").find(isRunnable).getOrElse(cancel("no 
runnable python"))
+  }
+
+  private def onDaemonThreads[T](threads: Int)(body: ExecutionContext => T): T 
= {
+    val pool = Executors.newFixedThreadPool(
+      threads,
+      (r: Runnable) => {
+        val t = new Thread(r, "pool-spec-caller")
+        t.setDaemon(true)
+        t
+      }
+    )
+    val ec: ExecutionContextExecutorService = 
ExecutionContext.fromExecutorService(pool)
+    try body(ec)
+    finally pool.shutdownNow()
+  }
+
+  private def hangingCall(
+      launchArgs: Seq[String],
+      request: com.fasterxml.jackson.databind.node.ObjectNode = 
objectMapper.createObjectNode()
+  ): PythonWorkerPool.Outcome =
+    PythonWorkerPool.run(
+      resourcePath = HangingWorker,
+      launchArgs = launchArgs,
+      pythonExe = python(),
+      request = request,
+      interpreterArgs = Seq("-I", "-S"),
+      timeouts = ShortTimeouts
+    )
+
+  /** The call, on a daemon thread and under [[Bound]], expected to give up. */
+  private def interceptBounded(call: => Any): 
PythonWorkerPool.WorkerDiedException =
+    intercept[PythonWorkerPool.WorkerDiedException] {
+      onDaemonThreads(1)(ec => Await.result(Future(call)(ec), Bound))
+    }
+
+  test("a worker that takes the job and stops answering is killed and 
reported") {
+    val startedAt = System.nanoTime()
+    val thrown = interceptBounded(hangingCall(Seq.empty))
+    val elapsedMillis = (System.nanoTime() - startedAt) / 1000000
+
+    assert(thrown.getMessage.contains("did not answer"))
+    assert(thrown.getMessage.contains("killed it"))
+    // Well under the default response budget: what fired is the timeout 
passed in,
+    // not a wait that happened to end.
+    assert(elapsedMillis < PythonWorkerPool.Timeouts.Default.responseMillis / 
2)
+  }
+
+  test("a worker that never signals ready is killed and reported") {
+    val startedAt = System.nanoTime()
+    val thrown = interceptBounded(hangingCall(Seq("--hang-before-ready")))
+    val elapsedMillis = (System.nanoTime() - startedAt) / 1000000
+
+    assert(thrown.getMessage.contains("did not signal ready"))
+    assert(elapsedMillis < PythonWorkerPool.Timeouts.Default.startupMillis / 2)
+  }
+
+  test("a worker that never reads its request is killed and reported") {
+    val request = objectMapper.createObjectNode()
+    // Past any pipe buffer, so the write cannot simply be handed to the 
kernel and
+    // left there: it is the blocked write itself that has to be given up on.
+    request.put("source", "x" * (4 * 1024 * 1024))
+
+    val thrown = interceptBounded(hangingCall(Seq("--deaf"), request))
+
+    assert(thrown.getMessage.contains("did not read its request"))
+    assert(thrown.getMessage.contains("killed it"))
+  }
+
+  test("a caller waiting at the worker cap is not stranded by a discarded 
worker") {
+    // One caller more than there are workers, all onto a worker that goes 
quiet:
+    // those holding a worker time out and are discarded, which frees a slot
+    // without handing anything back, and the caller waiting at the cap has to
+    // notice that rather than wait for a hand-back that never comes.
+    val callers = PythonWorkerPool.maxWorkers + 1
+
+    val outcomes = onDaemonThreads(callers) { implicit ec =>
+      
Await.result(Future.sequence(Seq.fill(callers)(Future(Try(hangingCall(Seq.empty))))),
 Bound)
+    }
+
+    assert(outcomes.length == callers)
+    assert(outcomes.forall(_.isFailure))
+  }
+
+  test("the pool still serves jobs after it has discarded a timed-out worker") 
{

Review Comment:
   `interceptBounded(hangingCall(...))` discards a worker in the 
`HangingWorker` sub-pool, but the follow-up job passes `resourcePath = 
CompileWorker` (:163). `Key` includes `resourcePath` 
(`PythonWorkerPool.scala:176-182`), so `computeIfAbsent` builds a fresh `Pool` 
with an empty `idle`.
   
   The assertion passes whether or not the discard left the pool it acted on 
usable. Pointing the follow-up at the same sub-pool would make it test its 
name. The cap case at :142-155 does cover replacement-after-discard, so the 
mechanism is not wholly untested.



##########
project/TestFilters.scala:
##########
@@ -0,0 +1,46 @@
+/*
+ * 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.
+ */
+
+import sbt._
+
+/**
+ * Selects a module's tagged tests for the fast-unit job or the integration 
job:
+ * skip-integration excludes them, integration-only runs only them, unset runs
+ * everything. Shared because the mapping is identical in every module, while 
the
+ * env var and the tag are not — the tag annotation has to live somewhere the
+ * module's own Test config can see.
+ */
+object TestFilters {
+
+  /** @param integrationOnlyExtra further ScalaTest args for the integration 
side,

Review Comment:
   This uses the `@param` form but documents only `integrationOnlyExtra`. 
`envVar` and `tag` are the two a caller has to get right: the env var must 
match what the workflow sets, and the tag must match the annotation. Worth a 
line each here.



##########
common/workflow-operator/src/test/scala/org/apache/texera/amber/util/PythonCodeRawInvalidTextSpec.scala:
##########
@@ -225,21 +305,23 @@ final class PythonCodeRawInvalidTextSpec extends 
AnyFunSuite {
     }
 
     val total = descriptorCandidates.size
-    var ok = 0
-    var checked = 0
-
-    val allFindings = descriptorCandidates.flatMap { descriptorClass =>
-      checked += 1
+    val ok = new AtomicInteger(0)
+    val checked = new AtomicInteger(0)
 
+    // Checked concurrently: the fan-out is what turns the pool's workers into
+    // parallel interpreters rather than a queue in front of one. The pool 
bounds
+    // it — a submission past the cap blocks until a worker is returned.

Review Comment:
   Both halves of the second sentence are off. The fan-out is bounded by 
`awaitAll`'s executor, which is sized to `PythonWorkerPool.maxWorkers` (:73-74) 
— nothing is submitted past the cap. And `borrow` at the cap deliberately does 
*not* block for a returned worker; its own comment 
(`PythonWorkerPool.scala:235-239`) says waiting would strand the caller.
   
   ```suggestion
       // Checked concurrently: the fan-out is what turns the pool's workers 
into
       // parallel interpreters rather than a queue in front of one. The 
executor is
       // sized to maxWorkers, so nothing is submitted past the cap.
   ```



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to