From e7a974e0e734657b3d929af0df1704ce0f3293c6 Mon Sep 17 00:00:00 2001 From: Martin Vu <22mvu7@gmail.com> Date: Mon, 24 Aug 2026 01:06:56 -0700 Subject: [PATCH 1/3] Fix console writer thread leak on unsubscribe --- .../web/service/ExecutionConsoleService.scala | 25 ++++++++++++++++ .../service/ExecutionConsoleServiceSpec.scala | 29 +++++++++++++++++++ 2 files changed, 54 insertions(+) diff --git a/amber/src/main/scala/org/apache/texera/web/service/ExecutionConsoleService.scala b/amber/src/main/scala/org/apache/texera/web/service/ExecutionConsoleService.scala index 55f72c35d81..fd16c654f10 100644 --- a/amber/src/main/scala/org/apache/texera/web/service/ExecutionConsoleService.scala +++ b/amber/src/main/scala/org/apache/texera/web/service/ExecutionConsoleService.scala @@ -221,6 +221,31 @@ class ExecutionConsoleService( } ) + override def unsubscribeAll(): Unit = { + consoleMessageOpIdToWriterMap.values.foreach { writer => + try { + writer.close() + } catch { + case e: Exception => + logger.error("Failed to close console message writer during unsubscribeAll", e) + } + } + consoleMessageOpIdToWriterMap.clear() + + super.unsubscribeAll() + + consoleWriterThread.shutdown() + try { + if (!consoleWriterThread.awaitTermination(5, java.util.concurrent.TimeUnit.SECONDS)) { + consoleWriterThread.shutdownNow() + } + } catch { + case _: InterruptedException => + consoleWriterThread.shutdownNow() + Thread.currentThread().interrupt() + } + } + /** * Processes a console message for display, performing truncation if needed. * This method uses the shared implementation in ConsoleMessageProcessor. diff --git a/amber/src/test/scala/org/apache/texera/web/service/ExecutionConsoleServiceSpec.scala b/amber/src/test/scala/org/apache/texera/web/service/ExecutionConsoleServiceSpec.scala index b1d647d035c..a6b7a0933ca 100644 --- a/amber/src/test/scala/org/apache/texera/web/service/ExecutionConsoleServiceSpec.scala +++ b/amber/src/test/scala/org/apache/texera/web/service/ExecutionConsoleServiceSpec.scala @@ -42,10 +42,17 @@ import org.apache.texera.web.model.websocket.request.python.DebugCommandRequest import org.apache.texera.web.storage.ExecutionStateStore import org.scalamock.scalatest.MockFactory import org.scalatest.BeforeAndAfterAll + import org.scalatest.flatspec.AnyFlatSpecLike import org.scalatest.matchers.should.Matchers + +import org.scalatest.concurrent.Eventually.eventually +import org.scalatest.concurrent.PatienceConfiguration.{Interval, Timeout} +import org.scalatest.time.{Millis, Span} + import java.time.Instant +import java.util.concurrent.ExecutorService import scala.collection.mutable.ListBuffer import scala.reflect.ClassTag @@ -441,4 +448,26 @@ class ExecutionConsoleServiceSpec keys should not contain "Worker:WF1-udf1-main-0" } } + + "unsubscribeAll" should "shutdown consoleWriterThread" in { + withFixture { f => + f.client.consoleCallback(message(title = "test")) + + val threadField = classOf[ExecutionConsoleService].getDeclaredField("consoleWriterThread") + threadField.setAccessible(true) + val executor = threadField.get(f.service).asInstanceOf[ExecutorService] + + // Verify it is initially active + executor.isShutdown shouldBe false + + // Trigger the teardown + f.service.unsubscribeAll() + + // Use Eventually to wait for async termination without blocking arbitrarily + eventually(Timeout(Span(2000, Millis)), Interval(Span(50, Millis))) { + executor.isShutdown shouldBe true + executor.isTerminated shouldBe true + } + } + } } From ded85e0d3aee004bacc313e0c357aaade1632bb1 Mon Sep 17 00:00:00 2001 From: Martin Vu <22mvu7@gmail.com> Date: Mon, 24 Aug 2026 01:23:05 -0700 Subject: [PATCH 2/3] Fix impor format --- .../texera/web/service/ExecutionConsoleServiceSpec.scala | 6 ++---- 1 file changed, 2 insertions(+), 4 deletions(-) diff --git a/amber/src/test/scala/org/apache/texera/web/service/ExecutionConsoleServiceSpec.scala b/amber/src/test/scala/org/apache/texera/web/service/ExecutionConsoleServiceSpec.scala index a6b7a0933ca..34cec01bc78 100644 --- a/amber/src/test/scala/org/apache/texera/web/service/ExecutionConsoleServiceSpec.scala +++ b/amber/src/test/scala/org/apache/texera/web/service/ExecutionConsoleServiceSpec.scala @@ -43,12 +43,10 @@ import org.apache.texera.web.storage.ExecutionStateStore import org.scalamock.scalatest.MockFactory import org.scalatest.BeforeAndAfterAll -import org.scalatest.flatspec.AnyFlatSpecLike -import org.scalatest.matchers.should.Matchers - - import org.scalatest.concurrent.Eventually.eventually import org.scalatest.concurrent.PatienceConfiguration.{Interval, Timeout} +import org.scalatest.flatspec.AnyFlatSpecLike +import org.scalatest.matchers.should.Matchers import org.scalatest.time.{Millis, Span} import java.time.Instant From e49b30df5d76b5dc845796f56a2ba0587ac922b7 Mon Sep 17 00:00:00 2001 From: Martin Vu <22mvu7@gmail.com> Date: Mon, 24 Aug 2026 01:26:05 -0700 Subject: [PATCH 3/3] remove space --- .../apache/texera/web/service/ExecutionConsoleServiceSpec.scala | 1 - 1 file changed, 1 deletion(-) diff --git a/amber/src/test/scala/org/apache/texera/web/service/ExecutionConsoleServiceSpec.scala b/amber/src/test/scala/org/apache/texera/web/service/ExecutionConsoleServiceSpec.scala index 34cec01bc78..ce88a46e6ec 100644 --- a/amber/src/test/scala/org/apache/texera/web/service/ExecutionConsoleServiceSpec.scala +++ b/amber/src/test/scala/org/apache/texera/web/service/ExecutionConsoleServiceSpec.scala @@ -42,7 +42,6 @@ import org.apache.texera.web.model.websocket.request.python.DebugCommandRequest import org.apache.texera.web.storage.ExecutionStateStore import org.scalamock.scalatest.MockFactory import org.scalatest.BeforeAndAfterAll - import org.scalatest.concurrent.Eventually.eventually import org.scalatest.concurrent.PatienceConfiguration.{Interval, Timeout} import org.scalatest.flatspec.AnyFlatSpecLike