dwsmith1983 commented on code in PR #5613:
URL: https://github.com/apache/datafusion-comet/pull/5613#discussion_r4133984221


##########
native/core/src/execution/memory_pools/fair_pool.rs:
##########
@@ -90,9 +144,110 @@ impl CometFairMemoryPool {
             state: Mutex::new(CometFairPoolState {
                 used: 0,
                 consumers: HashMap::new(),
+                anchor_held: false,
             }),
         }
     }
+
+    /// Whether an anchor request came back covered. A declined anchor is a 
zero grant, so
+    /// there is nothing to hand back either way.
+    fn anchor_granted(acquired: i64) -> bool {
+        let granted = usize::try_from(acquired).unwrap_or(0);
+        if granted > ANCHOR_BYTES {
+            warn!("Requested {ANCHOR_BYTES} bytes from the JVM but it reports 
{granted} granted");
+        }
+        granted >= ANCHOR_BYTES
+    }
+
+    /// Takes the anchor on the first grow and retries it while Spark declines 
it, as a
+    /// request of its own that never rides on a real grow. Spark declines it 
only while the
+    /// task sits at its share, so the extra JNI call is paid on that path 
alone and never
+    /// once the anchor is held. A `try_grow` rolls back its reservation if 
this fails.
+    fn take_missing_anchor(&self) -> CometResult<()> {
+        if self.state.lock().anchor_held {
+            return Ok(());
+        }
+        // The lock is not held across the call.
+        if 
!Self::anchor_granted(self.spark.manager().acquire_anchor(ANCHOR_BYTES)?) {
+            return Ok(());
+        }
+        {
+            let mut state = self.state.lock();
+            if !state.anchor_held {
+                state.anchor_held = true;
+                return Ok(());
+            }
+        }
+        // A grow or a release on another thread took the anchor meanwhile. 
This byte was never
+        // booked, so a failed return only leaves Spark holding it until the 
task ends, as on
+        // drop.
+        if let Err(e) = self.spark.manager().release_anchor(ANCHOR_BYTES) {
+            warn!("Failed to return a duplicate memory pool anchor byte: 
{e:?}");
+        }
+        Ok(())
+    }
+
+    /// Hands `bytes` that Spark granted this pool back to it. While the 
anchor is missing, one
+    /// of them is kept as the anchor instead: this pool holds at least those 
bytes from Spark,
+    /// so the release cannot take the task's balance to zero under an acquire 
parked there.
+    /// Claiming the anchor under the lock means one release keeps it, and an 
anchor retry that
+    /// lands afterwards hands its byte back as a duplicate. The JVM releases 
to Spark before it
+    /// moves its counters, so a failed call has moved nothing: the claim is 
rolled back and the
+    /// pool is unanchored again, with Spark still holding the bytes, and 
later grows retry the
+    /// anchor as usual. A retry that returned its byte while the claim stood 
holds nothing.
+    fn release_to_spark(&self, bytes: usize) -> CometResult<()> {
+        let keep_anchor = !std::mem::replace(&mut 
self.state.lock().anchor_held, true);
+        if !keep_anchor {
+            return self.spark.manager().release(bytes);
+        }
+        let released = self.spark.manager().release_keeping_anchor(bytes);
+        if released.is_err() {
+            self.state.lock().anchor_held = false;
+        }
+        released
+    }
+
+    /// Settles a release the JVM has accepted: the bytes come off the pool's 
total and the
+    /// consumer's share only now, so a grow is never admitted on bytes Spark 
still holds.
+    fn settle_release(&self, consumer: usize, bytes: usize) {
+        self.state.lock().settle(consumer, bytes, 0);
+    }
+
+    /// Settles a finished JVM acquire. `charged` bytes went on the pool's 
total and the
+    /// consumer's share before the call and `held` is what Spark granted and 
still holds, which
+    /// stays charged until it is handed back. The difference is rolled back.
+    fn settle_acquire(&self, consumer: usize, charged: usize, held: usize) {
+        self.state.lock().settle(consumer, charged, held);
+    }
+}
+
+impl Drop for CometFairMemoryPool {
+    /// The last plan of the task letting go of the pool runs this, on 
whatever thread that
+    /// happens on; `with_env` attaches the thread and the JVM object handle 
outlives the pool.
+    /// Nothing may panic out of a drop, so failures are only logged, and 
Spark frees the task's
+    /// whole balance when the task ends anyway.
+    fn drop(&mut self) {
+        let state = self.state.get_mut();
+        if state.used != 0 {
+            warn!(
+                "Task {} dropped CometFairMemoryPool with {} bytes still 
reserved ({} bytes overcommitted)",
+                self.spark.task_attempt_id(),
+                state.used,
+                self.spark.overcommit()
+            );
+        }
+        if !state.anchor_held {
+            return;
+        }
+        let released = 
std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
+            self.spark.manager().release_anchor(ANCHOR_BYTES)

Review Comment:
   > Could the anchor be shared across pool generations, or could replacement 
acquisition be coordinated with completed teardown without holding the global 
registry lock across JNI?
   
   Coordinated with completed teardown, in ac0a691. The registry entry used to 
go away before the fair pool's drop handed the anchor back, so a replacement 
could be created, park its first acquire in Spark and wake to the removed task 
entry, as your test showed. Now the entry stays, with an expired `Weak`, until 
the old pool has finished dropping, anchor release included, and 
`acquire_task_shared_pool` waits on a condvar while such an entry exists. A 
guard field declared after the pool removes the entry and notifies, so the JNI 
release still runs with no Rust lock held and the registry lock is only taken 
for the map update.
   
   I did not share the anchor across generations. The byte is charged to the 
`CometTaskMemoryManager` that acquired it, and a replacement created between 
"no successor, release" and the release landing would hit the same window, so 
it needs this ordering anyway.
   
   Two regressions, both failing on the previous head: 
`a_replacement_waits_for_the_previous_pool_to_finish_dropping` holds the old 
pool inside its drop and checks no replacement is handed out until it finishes, 
and `the_next_plan_of_a_task_waits_for_the_dying_pool_to_hand_back_its_anchor` 
is your scenario on the Spark stub (100 bytes, other tasks holding 74 and 25, 
the dying pool holding only its anchor).
   
   The cost is that `createPlan` can wait for one release call made by another 
thread.
   
   The branch also picked up main's #6261, which runs Spark acquires in 
`block_in_place`. The anchor acquire now goes through the same call (cd19f9a), 
so a parked anchor request no longer holds its Tokio worker.
   



-- 
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]

Reply via email to