Skip to content
Draft
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
310 changes: 74 additions & 236 deletions src/storage.rs
Original file line number Diff line number Diff line change
Expand Up @@ -771,188 +771,77 @@ impl RetriedStorage<Storage> {
visibility_timeout_secs: u64,
) -> Result<()> {
let inner = std::sync::Arc::clone(&self.inner);

type OwnedBuffers = (std::sync::Arc<[u8]>, std::sync::Arc<[u8]>);
struct State {
attempt: u32,
owned: Option<OwnedBuffers>,
}
let state = std::sync::Arc::new(std::sync::Mutex::new(State {
attempt: 0,
owned: None,
}));

with_rocksdb_busy_retry_async("add_available_item_from_parts", || {
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.is_none() {
guard.owned = Some((
std::sync::Arc::<[u8]>::from(id.to_vec()),
std::sync::Arc::<[u8]>::from(contents.to_vec()),
));
}
}
if let Some((id_arc, contents_arc)) = {
let guard = state_ref.lock().unwrap();
guard.owned.as_ref().cloned()
} {
tokio::task::spawn_blocking(move || {
inner_ref.add_available_item_from_parts(
&id_arc,
&contents_arc,
visibility_timeout_secs,
)
})
.await
.map_err(Into::<Error>::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(|_| ())
}
}
// Own inputs once to avoid per-retry allocations and allow a single blocking task
let id_owned: std::sync::Arc<[u8]> = std::sync::Arc::from(id.to_vec());
let contents_owned: std::sync::Arc<[u8]> = std::sync::Arc::from(contents.to_vec());

tokio::task::spawn_blocking(move || {
with_rocksdb_busy_retry_blocking("add_available_item_from_parts", || {
inner.add_available_item_from_parts(
id_owned.as_ref(),
contents_owned.as_ref(),
visibility_timeout_secs,
)
})
})
.await
.map_err(Into::<Error>::into)??;
Ok(())
}

pub async fn add_available_items_from_parts<'a, I>(&self, items: I) -> Result<()>
where
I: IntoIterator<Item = (&'a [u8], (&'a [u8], u64))>,
{
type BorrowedItem<'b> = (&'b [u8], (&'b [u8], u64));
type OwnedItem = (std::sync::Arc<[u8]>, (std::sync::Arc<[u8]>, u64));

let borrowed: Vec<BorrowedItem<'a>> = items.into_iter().collect();
let owned: Vec<OwnedItem> = items
.into_iter()
.map(|(id, (c, v))| {
(
std::sync::Arc::<[u8]>::from(id),
(std::sync::Arc::<[u8]>::from(c), v),
)
})
.collect();
let inner = std::sync::Arc::clone(&self.inner);
let borrowed_arc = std::sync::Arc::new(borrowed);

struct State {
attempt: u32,
owned: Option<std::sync::Arc<Vec<OwnedItem>>>,
}
let state = std::sync::Arc::new(std::sync::Mutex::new(State {
attempt: 0,
owned: None,
}));

with_rocksdb_busy_retry_async("add_available_items_from_parts", || {
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();
if guard.attempt > 0 && guard.owned.is_none() {
let owned_vec: Vec<OwnedItem> = borrowed_ref
.iter()
.map(|(id, (c, v))| {
(
std::sync::Arc::<[u8]>::from(*id),
(std::sync::Arc::<[u8]>::from(*c), *v),
)
})
.collect();
guard.owned = Some(std::sync::Arc::new(owned_vec));
}
}

if let Some(owned_ref) = { state_ref.lock().unwrap().owned.as_ref().cloned() } {
tokio::task::spawn_blocking(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::<Error>::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(|_| ())
}
}
tokio::task::spawn_blocking(move || {
with_rocksdb_busy_retry_blocking("add_available_items_from_parts", || {
let iter = owned
.iter()
.map(|(id, (c, v))| (id.as_ref(), (c.as_ref(), *v)));
inner.add_available_items_from_parts(iter)
})
})
.await
.map_err(Into::<Error>::into)??;
Ok(())
}

pub async fn peek_next_visibility_ts_secs(&self) -> Result<Option<u64>> {
let inner = std::sync::Arc::clone(&self.inner);
with_rocksdb_busy_retry_async("peek_next_visibility_ts_secs", || {
let inner_ref = std::sync::Arc::clone(&inner);
async move {
let out =
tokio::task::spawn_blocking(move || inner_ref.peek_next_visibility_ts_secs())
.await
.map_err(Into::<Error>::into)??;
Ok(out)
}
let out = tokio::task::spawn_blocking(move || {
with_rocksdb_busy_retry_blocking("peek_next_visibility_ts_secs", || {
inner.peek_next_visibility_ts_secs()
})
})
.await
.map_err(Into::<Error>::into)??;
Ok(out)
}

pub async fn get_next_available_entries(
&self,
n: usize,
) -> Result<(Lease, Vec<PolledItemOwnedReader>)> {
let inner = std::sync::Arc::clone(&self.inner);
with_rocksdb_busy_retry_async("get_next_available_entries", || {
let inner_ref = std::sync::Arc::clone(&inner);
async move {
let out =
tokio::task::spawn_blocking(move || inner_ref.get_next_available_entries(n))
.await
.map_err(Into::<Error>::into)??;
Ok(out)
}
let out = tokio::task::spawn_blocking(move || {
with_rocksdb_busy_retry_blocking("get_next_available_entries", || {
inner.get_next_available_entries(n)
})
})
.await
.map_err(Into::<Error>::into)??;
Ok(out)
}

pub async fn get_next_available_entries_with_lease(
Expand All @@ -961,118 +850,67 @@ impl RetriedStorage<Storage> {
lease_validity_secs: u64,
) -> Result<(Lease, Vec<PolledItemOwnedReader>)> {
let inner = std::sync::Arc::clone(&self.inner);
with_rocksdb_busy_retry_async("get_next_available_entries_with_lease", || {
let inner_ref = std::sync::Arc::clone(&inner);
async move {
let out = tokio::task::spawn_blocking(move || {
inner_ref.get_next_available_entries_with_lease(n, lease_validity_secs)
})
.await
.map_err(Into::<Error>::into)??;
Ok(out)
}
let out = tokio::task::spawn_blocking(move || {
with_rocksdb_busy_retry_blocking("get_next_available_entries_with_lease", || {
inner.get_next_available_entries_with_lease(n, lease_validity_secs)
})
})
.await
.map_err(Into::<Error>::into)??;
Ok(out)
}

pub async fn remove_in_progress_item(&self, id: &[u8], lease: &Lease) -> Result<bool> {
let lease_copy = *lease;
let inner = std::sync::Arc::clone(&self.inner);

struct State {
attempt: u32,
owned_id: Option<std::sync::Arc<[u8]>>,
}
let state = std::sync::Arc::new(std::sync::Mutex::new(State {
attempt: 0,
owned_id: None,
}));

with_rocksdb_busy_retry_async("remove_in_progress_item", || {
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()));
}
}

if let Some(id_arc) = { state_ref.lock().unwrap().owned_id.as_ref().cloned() } {
let out = tokio::task::spawn_blocking(move || {
inner_ref.remove_in_progress_item(&id_arc, &lease_copy)
})
.await
.map_err(Into::<Error>::into)??;
let mut guard = state_ref.lock().unwrap();
guard.attempt = guard.attempt.saturating_add(1);
Ok(out)
} 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)
}
}
let id_owned: std::sync::Arc<[u8]> = std::sync::Arc::from(id.to_vec());
let out = tokio::task::spawn_blocking(move || {
with_rocksdb_busy_retry_blocking("remove_in_progress_item", || {
inner.remove_in_progress_item(id_owned.as_ref(), &lease_copy)
})
})
.await
.map_err(Into::<Error>::into)??;
Ok(out)
}

pub async fn expire_due_leases(&self) -> Result<usize> {
let inner = std::sync::Arc::clone(&self.inner);
with_rocksdb_busy_retry_async("expire_due_leases", || {
let inner_ref = std::sync::Arc::clone(&inner);
async move {
let out = tokio::task::spawn_blocking(move || inner_ref.expire_due_leases())
.await
.map_err(Into::<Error>::into)??;
Ok(out)
}
let out = tokio::task::spawn_blocking(move || {
with_rocksdb_busy_retry_blocking("expire_due_leases", || inner.expire_due_leases())
})
.await
.map_err(Into::<Error>::into)??;
Ok(out)
}

pub async fn extend_lease(&self, lease: &Lease, lease_validity_secs: u64) -> Result<bool> {
let lease = *lease;
let inner = std::sync::Arc::clone(&self.inner);
with_rocksdb_busy_retry_async("extend_lease", || {
let inner_ref = std::sync::Arc::clone(&inner);
async move {
let out = tokio::task::spawn_blocking(move || {
inner_ref.extend_lease(&lease, lease_validity_secs)
})
.await
.map_err(Into::<Error>::into)??;
Ok(out)
}
let out = tokio::task::spawn_blocking(move || {
with_rocksdb_busy_retry_blocking("extend_lease", || {
inner.extend_lease(&lease, lease_validity_secs)
})
})
.await
.map_err(Into::<Error>::into)??;
Ok(out)
}
}

async fn with_rocksdb_busy_retry_async<T, Fut, F>(name: &str, mut f: F) -> Result<T>
// removed unused async retry helper (consolidated to blocking variant)

fn with_rocksdb_busy_retry_blocking<T, F>(name: &str, mut f: F) -> Result<T>
where
Fut: std::future::Future<Output = Result<T>>,
F: FnMut() -> Fut,
F: FnMut() -> Result<T>,
{
const MAX_RETRIES: u32 = 15;
const BASE_DELAY_MS: u64 = 10;
const MAX_DELAY_MS: u64 = 5000;

let mut attempt: u32 = 0;
loop {
match f().await {
match f() {
Ok(val) => return Ok(val),
Err(e) => {
let is_busy = matches!(
Expand All @@ -1093,8 +931,8 @@ where
} else {
rand::thread_rng().gen_range(0..=ceiling)
};
tracing::debug!(operation = %name, attempt = attempt + 1, delay_ms = jitter_ms, "RocksDB Busy: backing off with jitter");
tokio::time::sleep(std::time::Duration::from_millis(jitter_ms)).await;
tracing::debug!(operation = %name, attempt = attempt + 1, delay_ms = jitter_ms, "RocksDB Busy: backing off with jitter (blocking)");
std::thread::sleep(std::time::Duration::from_millis(jitter_ms));
attempt = attempt.saturating_add(1);
}
}
Expand Down