diff --git a/fluxer_media_proxy/src/coalescer.rs b/fluxer_media_proxy/src/coalescer.rs index 088bb0af5..f08f19b04 100644 --- a/fluxer_media_proxy/src/coalescer.rs +++ b/fluxer_media_proxy/src/coalescer.rs @@ -25,6 +25,24 @@ pub struct ByteCoalescer { in_flight: Mutex>>, } +struct CleanupGuard<'a> { + in_flight: &'a Mutex>>, + key: &'a str, + slot: &'a Arc, + completed: bool, +} + +impl Drop for CleanupGuard<'_> { + fn drop(&mut self) { + if !self.completed { + self.in_flight.lock().remove(self.key); + let mut state = self.slot.state.lock(); + *state = Some(Err(CoalescerError::WorkFailed)); + self.slot.notify.notify_waiters(); + } + } +} + impl ByteCoalescer { pub fn new() -> Self { Self::default() @@ -71,10 +89,21 @@ impl ByteCoalescer { crate::metrics::GLOBAL .coalescer_leader .fetch_add(1, std::sync::atomic::Ordering::Relaxed); + + let mut guard = CleanupGuard { + in_flight: &self.in_flight, + key: key.as_str(), + slot: &slot, + completed: false, + }; + let result = work().await.map(Bytes::from).map_err(coalesced_work_error); *slot.state.lock() = Some(result.clone()); + guard.completed = true; + drop(guard); slot.notify.notify_waiters(); self.in_flight.lock().remove(&key); + result } else { drop(work);