use crate::error::{OxCacheError, OxCacheResult};
use async_trait::async_trait;
use std::time::Duration;
#[async_trait]
pub trait LockProvider: Send + Sync {
async fn try_lock(&mut self) -> OxCacheResult<bool>;
async fn lock(&mut self) -> OxCacheResult<()>;
async fn unlock(&mut self) -> OxCacheResult<()>;
async fn is_held(&self) -> OxCacheResult<bool>;
fn fencing_token(&self) -> Option<u64> {
None
}
}
use super::lock::DistributedLock;
#[async_trait]
impl LockProvider for DistributedLock {
async fn try_lock(&mut self) -> OxCacheResult<bool> {
match self.acquire().await {
Ok(acquired) => Ok(acquired),
Err(OxCacheError::Operation(msg)) if msg.contains("already held by another owner") => {
Ok(false)
}
Err(e) => Err(e),
}
}
async fn lock(&mut self) -> OxCacheResult<()> {
let mut delay = Duration::from_millis(50);
let max_delay = Duration::from_secs(2);
loop {
match self.try_lock().await? {
true => return Ok(()),
false => {
tokio::time::sleep(delay).await;
delay = (delay * 2).min(max_delay);
}
}
}
}
async fn unlock(&mut self) -> OxCacheResult<()> {
self.release().await
}
async fn is_held(&self) -> OxCacheResult<bool> {
DistributedLock::is_held(self).await
}
fn fencing_token(&self) -> Option<u64> {
let token = DistributedLock::token(self);
if token == 0 { None } else { Some(token) }
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn lock_provider_is_object_safe() {
fn _assert_dyn_safe(_: &dyn LockProvider) {}
}
#[test]
fn default_lock_provider_alias() {
fn _assert_same<T: LockProvider>() {}
fn _check() {
_assert_same::<super::super::DefaultLockProvider>();
}
}
}