Copilot commented on code in PR #13166:
URL: https://github.com/apache/gluten/pull/13166#discussion_r4178357321
##########
gluten-core/src/main/scala/org/apache/spark/task/TaskResources.scala:
##########
@@ -339,6 +358,7 @@ class TaskResourceRegistry extends Logging {
}
priorityToResourcesMapping.clear()
resources.clear()
+ released = true
Review Comment:
`released` is set only after every callback finishes, but `releaseAll()`
deliberately lets fatal and control-flow throwables escape immediately. The
completion listener still removes the map entry in its `finally`, so a caller
that fetched this registry before removal can acquire its lock after such a
failure and successfully register a resource into an unreachable registry. Mark
the registry released on every exit (for example, before invoking callbacks or
in an unconditional `finally`).
##########
gluten-core/src/test/scala/org/apache/gluten/task/TaskResourceSuite.scala:
##########
@@ -119,4 +122,95 @@ class TaskResourceSuite extends AnyFunSuite with SQLHelper
{
// The registry was cleared on the way out, so a fresh task runs cleanly.
TaskResources.runUnsafe(assert(TaskResources.inSparkTask()))
}
+
+ test("Run unsafe - release callbacks run outside the global registry lock") {
+ // On a real executor the release callbacks include JNI native teardown
that
+ // can take milliseconds. Running them under the JVM-global registry lock
+ // serialized every concurrent task's completion and registration on the
+ // slowest callback; a concurrent task start must not block on it.
+ var concurrentTaskStarted = false
+ // Captured on the starter thread; read on this thread only after the join
+ // below confirms the thread terminated (the !isAlive assert gates it),
which
+ // establishes the happens-before to read these plain locals safely.
+ var starterFailure: Option[Throwable] = None
+ TaskResources.runUnsafe {
+ TaskResources.addResource(
+ UUID.randomUUID().toString,
+ new TaskResource {
+ override def release(): Unit = {
+ val starter = new Thread(
+ () => {
+ try {
+ TaskResources.runUnsafe {
+ TaskResources.addResource(
+ UUID.randomUUID().toString,
+ new TaskResource {
+ override def release(): Unit = {}
+ override def resourceName(): String = "concurrent task
resource"
+ }
+ )
+ }
+ concurrentTaskStarted = true
+ } catch {
+ case NonFatal(t) => starterFailure = Some(t)
+ }
+ })
+ starter.start()
+ starter.join(30000)
+ assert(
+ !starter.isAlive,
+ "concurrent task start is still blocked after the release
callback returned")
+ }
+ override def resourceName(): String = "slow releaser"
+ }
+ )
+ }
+ // Surface a genuine unrelated failure with its real stack instead of the
+ // misleading "blocked" message below.
+ starterFailure.foreach(t => throw t)
+ assert(concurrentTaskStarted, "concurrent task start was blocked by the
release pass")
+ }
+
+ test("Run unsafe - registering into a released registry fails") {
+ // A thread sharing the task's TaskContext can fetch the registry while
the release
+ // pass runs outside the global lock. Its registration must fail loudly
instead of
+ // landing in the cleared registry, where nothing would ever release it.
+ var lateFailure: Option[Throwable] = None
+ var lateReleased = false
+ var child: Thread = null
+ TaskResources.runUnsafe {
+ val tc = TaskContext.get()
+ TaskResources.addResource(
+ UUID.randomUUID().toString,
+ new TaskResource {
+ override def release(): Unit = {
+ child = new Thread(
+ () => {
+ SparkTaskUtil.setTaskContext(tc)
+ try {
+ TaskResources.addResource(
+ UUID.randomUUID().toString,
+ new TaskResource {
+ override def release(): Unit = lateReleased = true
+ override def resourceName(): String = "late resource"
+ })
+ } catch {
+ case NonFatal(t) => lateFailure = Some(t)
+ } finally {
+ SparkTaskUtil.unsetTaskContext()
+ }
+ })
+ child.start()
+ // Simulate a slow native teardown.
+ Thread.sleep(500)
Review Comment:
The sleep does not establish that the child fetched the registry before the
completion listener removes it. If the child is not scheduled within 500 ms,
`getTaskResourceRegistry()` throws the same `IllegalStateException` for a
missing map entry, so all assertions pass even if the new `released` guard is
absent. Use deterministic coordination or test `TaskResourceRegistry` directly
so registration is known to be waiting on the registry lock before release
completes.
--
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]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]