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 bdc6d2a90e fix(WorkflowExecutionService): shutdown console writer
thread on unsubscribe (#7914)
bdc6d2a90e is described below
commit bdc6d2a90eb1415014a4f705b7bb7cca31bb9688
Author: Martin Vu <[email protected]>
AuthorDate: Fri Aug 28 03:38:22 2026 +0000
fix(WorkflowExecutionService): shutdown console writer thread on
unsubscribe (#7914)
### What changes were proposed in this PR?
Fixes a console writer thread leak in `ExecutionConsoleService`.
The console writer executor was not being shut down when an execution
service was unsubscribed. This could leave `texera-console-writer`
threads alive after workflow execution finished.
Before:
<img width="1240" height="60" alt="Screenshot 2026-08-23 at 11 40 04 PM"
src="https://github.com/user-attachments/assets/d199f3c0-a463-4645-addf-a5dc2ee456d8"
/>
After:
<img width="1348" height="48" alt="Screenshot 2026-08-23 at 6 40 25 PM"
src="https://github.com/user-attachments/assets/a2e0c9bc-607f-451f-bccb-ec4ea4315a96"
/>
<img width="1348" height="60" alt="Screenshot 2026-08-23 at 6 41 53 PM"
src="https://github.com/user-attachments/assets/3d7742ef-db7a-4790-a43d-d8b5bfa69d58"
/>
<img width="1348" height="60" alt="Screenshot 2026-08-23 at 6 42 18 PM"
src="https://github.com/user-attachments/assets/c2e7a6d6-e516-4dd1-8ed9-7d6cccf22f1f"
/>
<img width="1348" height="60" alt="Screenshot 2026-08-23 at 6 42 42 PM"
src="https://github.com/user-attachments/assets/084bed9e-1338-4e27-9563-ad22708d75aa"
/>
<img width="1348" height="60" alt="Screenshot 2026-08-23 at 6 42 55 PM"
src="https://github.com/user-attachments/assets/ee11bcca-35fa-453e-91e2-8f3e06065c1d"
/>
This PR:
- Shuts down the console writer executor during `unsubscribeAll()`.
- Waits for termination and falls back to `shutdownNow()` if necessary.
- Closes active console message writers and clears the writer map.
- Adds a test verifying the console writer executor is shut down and
terminated.
### Any related issues, documentation, discussions?
Fixes #7455
### How was this PR tested?
Ran:
```bash
sbt "project WorkflowExecutionService" "testOnly
*ExecutionConsoleServiceSpec -- -z unsubscribeAll"
```
Manual testing:
Ran workflows multiple times and checked the console writer threads
with:
```bash
jcmd 57436 Thread.print | grep "texera-console-writer"
```
```bash
jcmd 67548 Thread.print | grep "texera-console-writer"
```
Verified that the console writer threads are terminated after workflow
execution completes and unsubscribeAll() is called.
### Was this PR authored or co-authored using generative AI tooling?
Generated-by: ChatGPT (5.5 mini)
---
.../web/service/ExecutionConsoleService.scala | 25 +++++++++++++++++++++
.../web/service/ExecutionConsoleServiceSpec.scala | 26 ++++++++++++++++++++++
2 files changed, 51 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 55f72c35d8..fd16c654f1 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 b1d647d035..ce88a46e6e 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,14 @@ 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
import org.scalatest.matchers.should.Matchers
+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 +445,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
+ }
+ }
+ }
}