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-8038-97e3585a6ac7aadbea629e7350a96d1a84f7b51c in repository https://gitbox.apache.org/repos/asf/texera.git
commit daeefc4dd1ddc38fe9791fe086eea988581d2d36 Author: Xinyuan Lin <[email protected]> AuthorDate: Thu Aug 27 08:20:57 2026 +0000 test(amber): cover ComputingUnitMaster runtime construction and cleanup (#8038) ### What changes were proposed in this PR? `ComputingUnitMasterSpec` goes from 34 tests to 40, reaching `createAmberRuntime` and the two cleanup `catch` blocks. | Metric | Before | After | |---|---|---| | **Codecov (fully-covered lines)** | 67/93 = 72.0% | **75/93 = 80.6%** | | JaCoCo line-hit | 73/93 | 79/93 | | Branch arms | 24/30 | 26/30 | **+8 fully-covered lines and +2 branch arms.** Attribution is exact: lines 75–80 (`createAmberRuntime`'s body) give +6, and lines 277 and 291 (the two `catch` blocks) give +1 line and +1 arm each. Nothing else moved. **A strict reviewer should count this as 6, not 8, and I would rather say so.** The two lines at 277/291 close a scala-logging `isWarnEnabled` macro guard, and no killing mutation exists for them on author-written code — the untaken arm belongs to the macro, not to anything in this file. The six `createAmberRuntime` lines are real and are mutation-proven. The existing spec was already good; this is an extension, not a rescue. ### What the reviewers found Two adversarial reviewers returned four findings against the first draft. The substance of the repair was **mutation strength, not coverage** — and that is worth being explicit about, because the numbers above did not move: **The two tests added by the repair round contribute zero new lines.** The forwarding test re-executes lines 75–80 that the first `createAmberRuntime` test already covered, and the `URISyntaxException` test executes line 285 (already covered) and by construction never reaches 286–291. They exist to kill mutants the original tests could not, and the coverage delta is identical to the builder's. Two pre-existing tests were also strengthened in place with cause assertions. Ten mutations were run in the final pass, **all killed**, each credited to the single test that failed — every run reported exactly `39 succeeded, 1 failed`, with the failing name read out of `target/test-reports/TEST-*.xml` rather than the sbt log, which never names it. The protocol is auditable: production was snapshotted to a scratch dir first (sha256 `e084b184…`), every revert came from that copy rather than from git, and before each apply a script asserted both that the file was byte-identical to the snapshot and that the anchor occurred exactly once, refusing to write otherwise. After each revert the hash was re-checked and the production diff confirmed empty, so no mutant was ever live during another's compile. **Three mutations are reported as unverified rather than as kills.** Deleting either `try`/`catch`, narrowing `case e: Throwable` to `IOException` in `dropCollections`, and reordering `dropCollections` against `deleteReplayLog` inside `cleanExecutions` are all believed to be killed by the pre-existing tests, but the final pass did not re-run them. They are listed as credited-on-trust, not as verified kills. ### Verification Measured with two scoped sbt-jacoco runs, scoping via a throwaway `Tests.Filter` on the **suite name** (never a `-z` phrase filter), byte-identical on both sides, one fresh sbt JVM per measurement with the jacoco directory wiped between them, counters read per-line out of `jacoco.xml`. The before-state ran the `HEAD` spec restored into place (34/34 green); the after-state ran the final spec (40/40 green). The throwaway scope file was deleted. Cross-suite contamination was checked explicitly, because this file touches the process-wide actor system and amber's suites share one strictly-serial JVM: `ComputingUnitMasterSpec` + `SessionStateSpec` + `WorkflowWebsocketResourceSpec` + `CoordinatorSpec` + `AmberClientSpec` in a single JVM — 70 tests, 5 suites, all green. ### Deliberately not included Lines 102–114 and 185–188 are refused rather than merely hard. `main()` would `System.exit` the shared test JVM, and the cleanup branch sits behind a private `static final boolean` that cannot be flipped from a test. Reaching either needs a production seam. Lines 91, 93, 120 and 181 remain partial and are left that way — they are on the do-not-chase list (scalac-generated and logging-guard arms). No production file is touched. ### Any related issues, documentation, discussions? Closes #8037 ### How was this PR tested? ``` sbt "WorkflowExecutionService/testOnly org.apache.texera.web.ComputingUnitMasterSpec" ``` ``` [info] Total number of tests run: 40 [info] Tests: succeeded 40, failed 0, canceled 0, ignored 0, pending 0 [info] All tests passed. ``` `WorkflowExecutionService/Test/scalafmtCheck` and `WorkflowExecutionService/Test/scalafix --check` both pass. Re-run after rebasing onto current `main`. ### Was this PR authored or co-authored using generative AI tooling? Generated-by: Claude Code (Opus 5) --------- Signed-off-by: Xinyuan Lin <[email protected]> Co-authored-by: Copilot Autofix powered by AI <[email protected]> --- .../texera/web/ComputingUnitMasterSpec.scala | 271 ++++++++++++++++++++- 1 file changed, 258 insertions(+), 13 deletions(-) diff --git a/amber/src/test/scala/org/apache/texera/web/ComputingUnitMasterSpec.scala b/amber/src/test/scala/org/apache/texera/web/ComputingUnitMasterSpec.scala index 5607777821..bfd9a02c5f 100644 --- a/amber/src/test/scala/org/apache/texera/web/ComputingUnitMasterSpec.scala +++ b/amber/src/test/scala/org/apache/texera/web/ComputingUnitMasterSpec.scala @@ -31,8 +31,17 @@ import io.dropwizard.jersey.validation.Validators import io.dropwizard.setup.{Bootstrap, Environment} import io.dropwizard.websockets.WebsocketBundle import org.apache.commons.jcs3.access.exception.InvalidArgumentException -import org.apache.pekko.actor.ActorSystem +import org.apache.pekko.actor.{ActorRef, ActorSystem} +import org.apache.texera.amber.core.virtualidentity.{ExecutionIdentity, WorkflowIdentity} +import org.apache.texera.amber.core.workflow.{PhysicalPlan, WorkflowContext} +import org.apache.texera.amber.engine.architecture.coordinator.CoordinatorConfig +import org.apache.texera.amber.engine.architecture.worker.WorkflowWorker.FaultToleranceConfig import org.apache.texera.amber.engine.common.AmberRuntime +import org.apache.texera.amber.engine.common.client.AmberClient +import org.apache.texera.amber.engine.common.virtualidentity.util.COORDINATOR +import org.apache.texera.amber.engine.e2e.TestUtils.buildWorkflow +import org.apache.texera.amber.operator.TestOperators +import org.apache.texera.amber.util.VirtualIdentityUtils import org.apache.texera.common.config.ApplicationConfig import org.apache.texera.dao.{MockTexeraDB, SqlServer} import org.apache.texera.dao.jooq.generated.Tables.{ @@ -54,23 +63,39 @@ import org.eclipse.jetty.servlet.FilterHolder import org.eclipse.jetty.websocket.server.WebSocketUpgradeFilter import org.glassfish.jersey.server.filter.RolesAllowedDynamicFeature import org.scalatest.BeforeAndAfterAll +import org.scalatest.concurrent.Eventually import org.scalatest.flatspec.AnyFlatSpec import org.scalatest.matchers.should.Matchers +import org.scalatest.time.{Millis, Seconds, Span} import java.lang.reflect.Proxy +import java.net.URISyntaxException import java.nio.file.{Files, Path} import java.sql.Timestamp +import java.util.concurrent.ConcurrentLinkedQueue import javax.servlet.http.{HttpServletRequest, HttpServletResponse} import javax.servlet.{DispatcherType, FilterChain, ServletContext, ServletContextEvent} import javax.websocket.server.ServerContainer +import javax.websocket.{RemoteEndpoint, Session} +import scala.concurrent.Await +import scala.concurrent.duration.DurationInt import scala.jdk.CollectionConverters._ class ComputingUnitMasterSpec extends AnyFlatSpec with Matchers with BeforeAndAfterAll + with Eventually with MockTexeraDB { + /** + * The coordinator the `createAmberRuntime` tests build does its scheduling on a pekko + * dispatcher thread, so what it was configured with only becomes observable a moment + * after the factory returns. Every wait below is bounded by this. + */ + implicit override val patienceConfig: PatienceConfig = + PatienceConfig(timeout = Span(20, Seconds), interval = Span(100, Millis)) + // ComputingUnitMaster.run schedules the result-cleanup task through AmberRuntime's // process-wide actor system, and only touches `.scheduler`, so a plain ActorSystem is // enough here; AmberRuntime.startActorMaster would additionally bind pekko artery on @@ -174,26 +199,40 @@ class ComputingUnitMasterSpec else if (returnType == java.lang.Long.TYPE) java.lang.Long.valueOf(0L) else null - /** Collects the events a named logger emits at `level` or above while `body` runs. */ - private def eventsLoggedAt[T](loggerName: String, level: Level)( - body: => T - ): Seq[ILoggingEvent] = { + /** + * Runs `body` with a logback ListAppender attached to `loggerName` forced to `level`, and + * hands it a live view of what has been collected so far. [[eventsLoggedAt]] covers the + * usual case where the code under test logs on the calling thread; this variant exists for + * the coordinator, which logs from a pekko dispatcher thread after the call returns, so the + * test has to be able to poll. The copy is taken under the appender's own lock because + * `AppenderBase.doAppend` is synchronized on the appender and the writer is another thread. + */ + private def whileCapturing[T](loggerName: String, level: Level)( + body: (() => Seq[ILoggingEvent]) => T + ): T = { val logger = org.slf4j.LoggerFactory.getLogger(loggerName).asInstanceOf[LogbackLogger] val appender = new ListAppender[ILoggingEvent] val previousLevel = logger.getLevel appender.start() logger.addAppender(appender) logger.setLevel(level) - try { - body - appender.list.asScala.toSeq - } finally { + try body(() => appender.synchronized(appender.list.asScala.toList)) + finally { logger.setLevel(previousLevel) logger.detachAppender(appender) appender.stop() } } + /** Collects the events a named logger emits at `level` or above while `body` runs. */ + private def eventsLoggedAt[T](loggerName: String, level: Level)( + body: => T + ): Seq[ILoggingEvent] = + whileCapturing(loggerName, level) { collected => + body + collected() + } + private def loggedAt[T](loggerName: String, level: Level)(body: => T): Seq[String] = eventsLoggedAt(loggerName, level)(body).map(_.getFormattedMessage) @@ -318,6 +357,170 @@ class ComputingUnitMasterSpec } } + /** + * The context every client below is built for. Deliberately NOT `new WorkflowContext()`: + * the default ids are exactly what a factory that dropped its `workflowContext` argument + * would fall back to, and a fixture that hands production the mutant's own literal cannot + * tell the two apart. + */ + private def specWorkflowContext: WorkflowContext = + new WorkflowContext( + workflowId = WorkflowIdentity(90210L), + executionId = ExecutionIdentity(90211L) + ) + + /** + * Builds a client the way `WorkflowExecutionService` does. The empty plan plus an + * all-None config is the recipe AmberClientSpec and three `web.service` specs already + * use: the AmberClient constructor blocks on an InitializeRequest that spawns a real + * Coordinator child, and an empty plan lets that finish with no engine behind it. The + * client is always shut down afterwards -- amber's suites share one serially-run JVM, so + * a leaked client would leave a live actor tree behind for every later suite. + */ + private def withAmberRuntime( + workflowContext: WorkflowContext = specWorkflowContext, + physicalPlan: PhysicalPlan = PhysicalPlan(Set.empty, Set.empty), + conf: CoordinatorConfig = CoordinatorConfig(None, None, None, None), + errorHandler: Throwable => Unit = _ => () + )(body: AmberClient => Unit): Unit = { + val client = ComputingUnitMaster.createAmberRuntime( + workflowContext, + physicalPlan, + conf, + errorHandler + ) + try body(client) + finally client.shutdown() + } + + /** + * Registers one SessionState against a stubbed websocket for the duration of `body` and + * hands it the JSON payloads pushed to that session. `SessionState` keeps a JVM-global + * registry, so the entry is always removed again -- a leaked one would receive events + * from every later suite that builds a coordinator. + */ + private def withRegisteredSession(body: (() => Seq[String]) => Unit): Unit = { + val pushed = new ConcurrentLinkedQueue[String]() + val remote = stub(classOf[RemoteEndpoint.Async]) { + case ("sendText", args) => pushed.add(args.head.asInstanceOf[String]); null + } + val session = stub(classOf[Session]) { + case ("getAsyncRemote", _) => remote.asInstanceOf[AnyRef] + } + val sessionId = "computing-unit-master-spec-" + java.util.UUID.randomUUID() + SessionState.setState(sessionId, new SessionState(session)) + try body(() => pushed.asScala.toList) + finally SessionState.removeState(sessionId) + } + + /** File names directly inside `folder`. */ + private def entriesIn(folder: Path): Seq[String] = { + val listing = Files.list(folder) + try listing.iterator().asScala.map(_.getFileName.toString).toList + finally listing.close() + } + + /** + * Reads the only field of `fieldType` off an AmberClient. Every field there is + * class-private, and the client actor's is additionally name-mangled, so they are looked + * up by TYPE rather than by a brittle spelling of the name. + */ + private def amberClientField[T](client: AmberClient, fieldType: Class[T]): T = { + val field = classOf[AmberClient].getDeclaredFields + .find(candidate => fieldType.isAssignableFrom(candidate.getType)) + .getOrElse(fail(s"AmberClient no longer holds a ${fieldType.getSimpleName} field")) + field.setAccessible(true) + field.get(client).asInstanceOf[T] + } + + "createAmberRuntime" should "build the client on the process-wide actor system" in { + withAmberRuntime() { client => + val clientActor = amberClientField(client, classOf[ActorRef]) + + // The factory has to hand AmberClient the process-wide AmberRuntime.actorSystem and + // not one of its own: in production that is the system startActorMaster bound to + // artery and joined to the cluster, and a workflow run on a private system would be + // unreachable from the rest of the runtime. Identity is asserted as well as the name + // -- resolving the path inside this suite's own system fails if the actor was created + // anywhere else, even in a second system that happened to carry the same name. + clientActor.path.address.system shouldBe testSystem.name + val resolved = Await.result( + testSystem + .actorSelection("/" + clientActor.path.elements.mkString("/")) + .resolveOne(10.seconds), + 10.seconds + ) + resolved shouldBe clientActor + } + } + + it should "hand the client the caller's error handler" in { + val errorHandler: Throwable => Unit = _ => () + + withAmberRuntime(errorHandler = errorHandler) { client => + // AmberClient routes everything a registered callback throws to this function and + // nowhere else, so a factory that quietly substituted a handler of its own would + // swallow callback failures at the single place production builds a client. + amberClientField(client, classOf[Function1[_, _]]) should be theSameInstanceAs errorHandler + } + } + + it should "forward the workflow context, physical plan and coordinator config it is handed" in { + // The two assertions above pin arguments 1 and 5. Arguments 2, 3 and 4 are consumed + // inside the AmberClient constructor and retained on no field of it (`javap -p` keeps + // only errorHandler, clientActor, timeout, registeredObservables, isActive and + // coordinatorInterface), so they can only be observed through what the coordinator the + // constructor spawns then does with them. Production passes `workflow.context`, + // `workflow.physicalPlan` and the service's own config at WorkflowExecutionService:124; + // a factory that substituted defaults for any of the three would run every workflow with + // an empty plan, no workflow/execution identity and fault tolerance disabled. + val scanOp = TestOperators.headerlessSmallCsvScanOpDesc() + val context = specWorkflowContext + val physicalPlan = buildWorkflow(List(scanOp), List.empty, context).physicalPlan + val logFolder = Files.createTempDirectory("computing-unit-master-spec-fault-tolerance") + val conf = + CoordinatorConfig(None, None, None, Some(FaultToleranceConfig(writeTo = logFolder.toUri))) + // s"${toShorterString(COORDINATOR)}] [<simple name>" is AmberLogging's naming scheme. + val scheduleGeneratorLogger = + s"${VirtualIdentityUtils.toShorterString(COORDINATOR)}] [${classOf[org.apache.texera.amber.engine.architecture.scheduling.CostBasedScheduleGenerator].getSimpleName}" + try { + withRegisteredSession { pushedToSession => + whileCapturing(scheduleGeneratorLogger, Level.INFO) { generatorEvents => + withAmberRuntime(context, physicalPlan, conf) { _ => + // The coordinator schedules on a dispatcher thread, so all three signals appear + // shortly after the factory returns rather than inside the call. + eventually { + // physicalPlan: the coordinator partitions the plan it was given into regions + // and pushes their logical operator ids to every open session. The empty plan + // yields no regions at all, so this id can only come from the plan passed in. + pushedToSession().mkString("\n") should include(scanOp.operatorIdentifier.id) + + // workflowContext: nothing downstream retains it, so the schedule generator's + // own diagnostic line is the only in-process evidence of which workflow and + // execution the runtime was built for. `new WorkflowContext()` would say 1/1. + generatorEvents().map(_.getFormattedMessage).mkString("\n") should ( + include("WID: 90210") and include("EID: 90211") + ) + + // conf: a fault-tolerance URI makes WorkflowActor open a real VFS record + // storage and the replay log manager create the coordinator's log file inside + // it. The all-None config yields an EmptyRecordStorage, which never writes. + entriesIn(logFolder) should not be empty + } + } + } + } + } finally { + // Best effort: the replay writer may still hold the log file open on Windows. + entriesIn(logFolder).foreach(name => + try Files.deleteIfExists(logFolder.resolve(name)) + catch { case _: Throwable => () } + ) + try Files.deleteIfExists(logFolder) + catch { case _: Throwable => () } + } + } + "initialize" should "wrap the configuration source provider in the environment substitutor" in { val bootstrap = new Bootstrap[Configuration](master) @@ -512,9 +715,14 @@ class ComputingUnitMasterSpec } it should "warn instead of propagating when the result pointer is unparsable" in { - val messages = loggedAt(masterLoggerName, Level.WARN)(dropCollections("{not json")) + val events = eventsLoggedAt(masterLoggerName, Level.WARN)(dropCollections("{not json")) - messages shouldBe Seq("result collection cleanup failed.") + events.map(_.getFormattedMessage) shouldBe Seq("result collection cleanup failed.") + // The cause is asserted, not just the text. "result collection cleanup failed." names + // neither the execution nor the pointer nor the failure, so the attached stack trace is + // the entire diagnostic content of this line -- and dropping the second argument (the + // one-arg Logger.warn overload) leaves the message identical and the operator blind. + events.head.getThrowableProxy should not be null } it should "warn instead of propagating on an unsupported storage type" in { @@ -529,6 +737,18 @@ class ComputingUnitMasterSpec messages shouldBe Seq("result collection cleanup failed.") } + it should "stay silent about an unparsable pointer once warn logging is switched off" in { + // A failed cleanup is advisory, so it has to sit at warn and vanish completely when + // TEXERA_SERVICE_LOG_LEVEL is raised to error; the warn-level tests above pass just as + // well if the catch block is promoted to logger.error, which would keep shouting at a + // deployment that asked for quiet. The swallow still has to hold with logging off. + val events = eventsLoggedAt(masterLoggerName, Level.ERROR) { + noException should be thrownBy dropCollections("{not json") + } + + events shouldBe empty + } + "deleteReplayLog" should "skip an execution with no log location" in { // Without the early return the empty location becomes a scheme-less URI and the // storage lookup would blow up into a delete-failed warning. @@ -549,9 +769,34 @@ class ComputingUnitMasterSpec it should "warn instead of propagating when the log location is unusable" in { // A scheme-less location cannot be resolved to any record storage. - val messages = loggedAt(masterLoggerName, Level.WARN)(deleteReplayLog("no-scheme/logs")) + val events = eventsLoggedAt(masterLoggerName, Level.WARN)(deleteReplayLog("no-scheme/logs")) + + events.map(_.getFormattedMessage) shouldBe Seq("failed to delete log at no-scheme/logs") + // As in dropCollections, the cause is asserted and not merely the text: the message says + // which location failed but nothing about why, so dropping the throwable argument would + // leave an advisory line that cannot be acted on. + events.head.getThrowableProxy should not be null + } + + it should "propagate a syntactically invalid log location instead of warning it" in { + // Documented, NOT endorsed. `new URI(logLocation)` sits outside the try, so this is the + // one failure mode deleteReplayLog does not swallow: the checked URISyntaxException + // travels out through cleanExecutions' foreach and aborts the cleanup of every remaining + // execution in the batch, which reads as an oversight next to the advisory warn every + // other failure here gets. Pinned because it is invisible otherwise -- the four other + // inputs this suite uses are all syntactically valid -- and because a fix should be a + // deliberate change to this test, not a silent one. `new URI` rejects the space at index 3. + an[URISyntaxException] should be thrownBy deleteReplayLog("has space/logs") + } + + it should "stay silent about an unusable location once warn logging is switched off" in { + // Same contract as the dropCollections case: the delete failure is advisory, so raising + // the service log level to error must silence it outright rather than merely reword it. + val events = eventsLoggedAt(masterLoggerName, Level.ERROR) { + noException should be thrownBy deleteReplayLog("no-scheme/logs") + } - messages shouldBe Seq("failed to delete log at no-scheme/logs") + events shouldBe empty } "cleanExecutions" should "clean every execution it is handed" in {
