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-7914-e54ecd2970ac8d2099c1e8f061cf053d32dd5e85
in repository https://gitbox.apache.org/repos/asf/texera.git

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
+      }
+    }
+  }
 }

Reply via email to