From f20a1c11a7cd995025a744964544cfe7bc1c9de4 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Mon, 15 Sep 2025 11:03:43 +0000 Subject: [PATCH 1/2] Refactor RetriedStorage to use a generic RetryState struct Co-authored-by: miles.frankel --- src/storage.rs | 218 ++++++++++++++++++++++++------------------------- 1 file changed, 106 insertions(+), 112 deletions(-) diff --git a/src/storage.rs b/src/storage.rs index 6c5da46..e7f1569 100644 --- a/src/storage.rs +++ b/src/storage.rs @@ -795,6 +795,11 @@ fn build_lease_entry_message( } // Retry wrapper that composes busy-retry behavior around a `Storage` instance. +struct RetryState { + attempt: u32, + owned: Option, +} + pub struct RetriedStorage { inner: std::sync::Arc, } @@ -812,30 +817,66 @@ impl RetriedStorage { } impl RetriedStorage { + #[inline] + fn run_blocking_maybe_inline(f: impl FnOnce() -> Result) -> Result { + match tokio::runtime::Handle::try_current() { + Ok(handle) + if handle.runtime_flavor() == tokio::runtime::RuntimeFlavor::CurrentThread => + { + f() + } + _ => tokio::task::block_in_place(f), + } + } + + async fn run_blocking_on_pool( + f: impl FnOnce() -> Result + Send + 'static, + ) -> Result { + let res: Result = tokio::task::spawn_blocking(f) + .await + .map_err(Into::::into)?; + res + } + + async fn busy_retry_with_state( + name: &str, + state: std::sync::Arc>>, + f: F, + ) -> Result + where + O: Send, + T: Send, + Fut: std::future::Future> + Send, + F: Fn(std::sync::Arc>>) -> Fut + Send + Sync + Clone, + { + let state_ref = std::sync::Arc::clone(&state); + let f_outer = f.clone(); + with_rocksdb_busy_retry_async(name, move || { + let state_inner = std::sync::Arc::clone(&state_ref); + let f_cloned = f_outer.clone(); + async move { f_cloned(state_inner).await } + }) + .await + } + pub async fn add_available_item_from_parts( &self, id: &[u8], contents: &[u8], visibility_timeout_secs: u64, ) -> Result<()> { + type Owned = (std::sync::Arc<[u8]>, std::sync::Arc<[u8]>); let inner = std::sync::Arc::clone(&self.inner); - - type OwnedBuffers = (std::sync::Arc<[u8]>, std::sync::Arc<[u8]>); - struct State { - attempt: u32, - owned: Option, - } - let state = std::sync::Arc::new(std::sync::Mutex::new(State { + let state = std::sync::Arc::new(std::sync::Mutex::new(RetryState:: { attempt: 0, owned: None, })); - with_rocksdb_busy_retry_async("add_available_item_from_parts", || { + Self::busy_retry_with_state("add_available_item_from_parts", state.clone(), move |st| { let inner_ref = std::sync::Arc::clone(&inner); - let state_ref = std::sync::Arc::clone(&state); async move { { - let mut guard = state_ref.lock().unwrap(); + let mut guard = st.lock().unwrap(); if guard.attempt > 0 && guard.owned.is_none() { guard.owned = Some(( std::sync::Arc::<[u8]>::from(id.to_vec()), @@ -843,48 +884,32 @@ impl RetriedStorage { )); } } - if let Some((id_arc, contents_arc)) = { - let guard = state_ref.lock().unwrap(); - guard.owned.as_ref().cloned() - } { - tokio::task::spawn_blocking(move || { + + let res = if let Some(owned) = { st.lock().unwrap().owned.as_ref().cloned() } { + Self::run_blocking_on_pool(move || { inner_ref.add_available_item_from_parts( - &id_arc, - &contents_arc, + &owned.0, + &owned.1, visibility_timeout_secs, ) }) .await - .map_err(Into::::into)??; - let mut guard = state_ref.lock().unwrap(); - guard.attempt = guard.attempt.saturating_add(1); - Ok(()) } else { - // First attempt: zero-copy direct or block_in_place depending on runtime - let res = match tokio::runtime::Handle::try_current() { - Ok(handle) - if handle.runtime_flavor() - == tokio::runtime::RuntimeFlavor::CurrentThread => - { - inner_ref.add_available_item_from_parts( - id, - contents, - visibility_timeout_secs, - ) - } - _ => tokio::task::block_in_place(move || { - inner_ref.add_available_item_from_parts( - id, - contents, - visibility_timeout_secs, - ) - }), - }; - // Increment attempt after finishing - let mut guard = state_ref.lock().unwrap(); - guard.attempt = guard.attempt.saturating_add(1); - res.map(|_| ()) + Self::run_blocking_maybe_inline(|| { + inner_ref.add_available_item_from_parts( + id, + contents, + visibility_timeout_secs, + ) + }) + .map(|_| ()) + }; + + { + let mut g = st.lock().unwrap(); + g.attempt = g.attempt.saturating_add(1); } + res } }) .await @@ -901,22 +926,18 @@ impl RetriedStorage { let inner = std::sync::Arc::clone(&self.inner); let borrowed_arc = std::sync::Arc::new(borrowed); - struct State { - attempt: u32, - owned: Option>>, - } - let state = std::sync::Arc::new(std::sync::Mutex::new(State { + type Owned = std::sync::Arc>; + let state = std::sync::Arc::new(std::sync::Mutex::new(RetryState:: { attempt: 0, owned: None, })); - with_rocksdb_busy_retry_async("add_available_items_from_parts", || { + Self::busy_retry_with_state("add_available_items_from_parts", state.clone(), move |st| { let inner_ref = std::sync::Arc::clone(&inner); - let state_ref = std::sync::Arc::clone(&state); let borrowed_ref = std::sync::Arc::clone(&borrowed_arc); async move { { - let mut guard = state_ref.lock().unwrap(); + let mut guard = st.lock().unwrap(); if guard.attempt > 0 && guard.owned.is_none() { let owned_vec: Vec = borrowed_ref .iter() @@ -931,40 +952,27 @@ impl RetriedStorage { } } - if let Some(owned_ref) = { state_ref.lock().unwrap().owned.as_ref().cloned() } { - tokio::task::spawn_blocking(move || { + let res = if let Some(owned_ref) = { st.lock().unwrap().owned.as_ref().cloned() } { + Self::run_blocking_on_pool(move || { let iter = owned_ref .iter() .map(|(id, (c, v))| (id.as_ref(), (c.as_ref(), *v))); inner_ref.add_available_items_from_parts(iter) }) .await - .map_err(Into::::into)??; - let mut guard = state_ref.lock().unwrap(); - guard.attempt = guard.attempt.saturating_add(1); - Ok(()) } else { - let res = match tokio::runtime::Handle::try_current() { - Ok(handle) - if handle.runtime_flavor() - == tokio::runtime::RuntimeFlavor::CurrentThread => - { - let iter = borrowed_ref.iter().map(|(id, (c, v))| ((*id), ((*c), *v))); - inner_ref.add_available_items_from_parts(iter) - } - _ => { - let borrowed_ref = std::sync::Arc::clone(&borrowed_ref); - tokio::task::block_in_place(move || { - let iter = - borrowed_ref.iter().map(|(id, (c, v))| ((*id), ((*c), *v))); - inner_ref.add_available_items_from_parts(iter) - }) - } - }; - let mut guard = state_ref.lock().unwrap(); - guard.attempt = guard.attempt.saturating_add(1); - res.map(|_| ()) + Self::run_blocking_maybe_inline(|| { + let iter = borrowed_ref.iter().map(|(id, (c, v))| ((*id), ((*c), *v))); + inner_ref.add_available_items_from_parts(iter) + }) + .map(|_| ()) + }; + + { + let mut g = st.lock().unwrap(); + g.attempt = g.attempt.saturating_add(1); } + res } }) .await @@ -1024,54 +1032,40 @@ impl RetriedStorage { } pub async fn remove_in_progress_item(&self, id: &[u8], lease: &Lease) -> Result { + type Owned = std::sync::Arc<[u8]>; let lease_copy = *lease; let inner = std::sync::Arc::clone(&self.inner); - - struct State { - attempt: u32, - owned_id: Option>, - } - let state = std::sync::Arc::new(std::sync::Mutex::new(State { + let state = std::sync::Arc::new(std::sync::Mutex::new(RetryState:: { attempt: 0, - owned_id: None, + owned: None, })); - with_rocksdb_busy_retry_async("remove_in_progress_item", || { + Self::busy_retry_with_state("remove_in_progress_item", state.clone(), move |st| { let inner_ref = std::sync::Arc::clone(&inner); - let state_ref = std::sync::Arc::clone(&state); async move { { - let mut guard = state_ref.lock().unwrap(); - if guard.attempt > 0 && guard.owned_id.is_none() { - guard.owned_id = Some(std::sync::Arc::<[u8]>::from(id.to_vec())); + let mut guard = st.lock().unwrap(); + if guard.attempt > 0 && guard.owned.is_none() { + guard.owned = Some(std::sync::Arc::<[u8]>::from(id.to_vec())); } } - if let Some(id_arc) = { state_ref.lock().unwrap().owned_id.as_ref().cloned() } { - let out = tokio::task::spawn_blocking(move || { + let out = if let Some(id_arc) = { st.lock().unwrap().owned.as_ref().cloned() } { + Self::run_blocking_on_pool(move || { inner_ref.remove_in_progress_item(&id_arc, &lease_copy) }) - .await - .map_err(Into::::into)??; - let mut guard = state_ref.lock().unwrap(); - guard.attempt = guard.attempt.saturating_add(1); - Ok(out) + .await? } else { - let out = match tokio::runtime::Handle::try_current() { - Ok(handle) - if handle.runtime_flavor() - == tokio::runtime::RuntimeFlavor::CurrentThread => - { - inner_ref.remove_in_progress_item(id, &lease_copy) - } - _ => tokio::task::block_in_place(move || { - inner_ref.remove_in_progress_item(id, &lease_copy) - }), - }?; - let mut guard = state_ref.lock().unwrap(); - guard.attempt = guard.attempt.saturating_add(1); - Ok(out) + Self::run_blocking_maybe_inline(|| { + inner_ref.remove_in_progress_item(id, &lease_copy) + })? + }; + + { + let mut g = st.lock().unwrap(); + g.attempt = g.attempt.saturating_add(1); } + Ok(out) } }) .await From 27cc44785075a475c0ed025a478149931347b3a1 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Mon, 15 Sep 2025 11:27:37 +0000 Subject: [PATCH 2/2] Remove unused run_blocking_maybe_inline helper Co-authored-by: miles.frankel --- src/storage.rs | 36 ++++++++---------------------------- 1 file changed, 8 insertions(+), 28 deletions(-) diff --git a/src/storage.rs b/src/storage.rs index e7f1569..4d37156 100644 --- a/src/storage.rs +++ b/src/storage.rs @@ -817,18 +817,6 @@ impl RetriedStorage { } impl RetriedStorage { - #[inline] - fn run_blocking_maybe_inline(f: impl FnOnce() -> Result) -> Result { - match tokio::runtime::Handle::try_current() { - Ok(handle) - if handle.runtime_flavor() == tokio::runtime::RuntimeFlavor::CurrentThread => - { - f() - } - _ => tokio::task::block_in_place(f), - } - } - async fn run_blocking_on_pool( f: impl FnOnce() -> Result + Send + 'static, ) -> Result { @@ -895,14 +883,9 @@ impl RetriedStorage { }) .await } else { - Self::run_blocking_maybe_inline(|| { - inner_ref.add_available_item_from_parts( - id, - contents, - visibility_timeout_secs, - ) - }) - .map(|_| ()) + inner_ref + .add_available_item_from_parts(id, contents, visibility_timeout_secs) + .map(|_| ()) }; { @@ -961,11 +944,10 @@ impl RetriedStorage { }) .await } else { - Self::run_blocking_maybe_inline(|| { - let iter = borrowed_ref.iter().map(|(id, (c, v))| ((*id), ((*c), *v))); - inner_ref.add_available_items_from_parts(iter) - }) - .map(|_| ()) + let iter = borrowed_ref.iter().map(|(id, (c, v))| ((*id), ((*c), *v))); + inner_ref + .add_available_items_from_parts(iter) + .map(|_| ()) }; { @@ -1056,9 +1038,7 @@ impl RetriedStorage { }) .await? } else { - Self::run_blocking_maybe_inline(|| { - inner_ref.remove_in_progress_item(id, &lease_copy) - })? + inner_ref.remove_in_progress_item(id, &lease_copy)? }; {