use std::sync::atomic::Ordering;
use super::reference_resource_counter::ReferenceResourceBase;
pub trait ReferenceResource: Send + Sync {
fn base(&self) -> &ReferenceResourceBase;
#[inline]
fn hold(&self) -> bool {
if !self.base().available.load(Ordering::Relaxed) {
return false;
}
let _guard = self.base().hold_lock.lock();
if self.base().available.load(Ordering::Acquire) {
let prev_count = self.base().ref_count.fetch_add(1, Ordering::Relaxed);
if prev_count > 0 {
return true;
} else {
self.base().ref_count.fetch_sub(1, Ordering::Relaxed);
}
}
false
}
#[inline]
fn is_available(&self) -> bool {
self.base().available.load(Ordering::Relaxed)
}
fn shutdown(&self, interval_forcibly: u64) {
use rocketmq_common::TimeUtils::current_millis;
if self.base().available.load(Ordering::Acquire) {
self.base().available.store(false, Ordering::Release);
self.base()
.first_shutdown_timestamp
.store(current_millis(), Ordering::Release);
self.release();
} else if self.get_ref_count() > 0 {
let elapsed = current_millis().saturating_sub(self.base().first_shutdown_timestamp.load(Ordering::Acquire));
if elapsed >= interval_forcibly {
let current_count = self.get_ref_count();
self.base().ref_count.store(-1000 - current_count, Ordering::Release);
self.release();
}
}
}
#[inline]
fn release(&self) {
let value = self.base().ref_count.fetch_sub(1, Ordering::Release) - 1;
if value > 0 {
return;
}
std::sync::atomic::fence(Ordering::Acquire);
let _guard = self.base().release_lock.lock();
if !self.base().cleanup_over.load(Ordering::Relaxed) {
let cleanup_result = self.cleanup(value);
self.base().cleanup_over.store(cleanup_result, Ordering::Release);
}
}
#[inline]
fn get_ref_count(&self) -> i64 {
self.base().ref_count.load(Ordering::Relaxed)
}
fn cleanup(&self, current_ref: i64) -> bool;
#[inline]
fn is_cleanup_over(&self) -> bool {
self.base().cleanup_over.load(Ordering::Relaxed) && self.base().ref_count.load(Ordering::Relaxed) <= 0
}
}