use std::ffi::CString;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex};
use crate::client_conductor::ClientConductor;
use crate::concurrent::atomic_buffer::AtomicBuffer;
use crate::concurrent::atomic_counter::AtomicCounter;
use crate::utils::errors::AeronError;
pub struct Counter {
atomic_counter: AtomicCounter,
client_conductor: Arc<Mutex<ClientConductor>>,
registration_id: i64,
is_closed: AtomicBool,
}
unsafe impl Send for Counter {}
unsafe impl Sync for Counter {}
impl Counter {
pub fn new(
client_conductor: Arc<Mutex<ClientConductor>>,
buffer: AtomicBuffer,
registration_id: i64,
counter_id: i32,
) -> Self {
Self {
atomic_counter: AtomicCounter::new(buffer, counter_id),
client_conductor,
registration_id,
is_closed: AtomicBool::from(false),
}
}
pub fn registration_id(&self) -> i64 {
self.registration_id
}
pub fn is_closed(&self) -> bool {
self.is_closed.load(Ordering::SeqCst)
}
pub fn close(&self) {
self.is_closed.store(true, Ordering::SeqCst);
}
pub fn state(&self) -> Result<i32, AeronError> {
let cc = self.client_conductor.lock().expect("Mutex poisoned");
let cr = cc.counters_reader()?;
cr.counter_state(self.atomic_counter.id())
}
pub fn label(&self) -> Result<CString, AeronError> {
let cc = self.client_conductor.lock().expect("Mutex poisoned");
let cr = cc.counters_reader()?;
cr.counter_label(self.atomic_counter.id())
}
pub fn id(&self) -> i32 {
self.atomic_counter.id()
}
}
impl Drop for Counter {
fn drop(&mut self) {
let _ignored = self
.client_conductor
.lock()
.expect("Mutex poisoned")
.release_counter(self.registration_id);
}
}