use std::ffi::c_void;
use std::mem;
use std::sync::Arc;
use std::sync::atomic::{AtomicU64, Ordering};
use crossbeam_queue::ArrayQueue;
use skippy_ffi::{
SkippyRuntimeEventReporterV1 as RawRuntimeEventReporter,
SkippyRuntimeEventV1 as RawRuntimeEvent,
};
use super::{NativeEventRecord, OperationId, RUNTIME_EVENT_V1_ABI_VERSION};
pub const MODEL_OPEN_RECORD_CAPACITY: usize = 256;
pub struct ModelOpenEventQueue {
operation_id: OperationId,
records: ArrayQueue<NativeEventRecord>,
dropped: AtomicU64,
rejected: AtomicU64,
}
impl ModelOpenEventQueue {
#[must_use]
pub fn new(operation_id: OperationId) -> Arc<Self> {
Self::with_capacity(operation_id, MODEL_OPEN_RECORD_CAPACITY)
}
#[must_use]
pub fn with_capacity(operation_id: OperationId, capacity: usize) -> Arc<Self> {
Arc::new(Self {
operation_id,
records: ArrayQueue::new(capacity.max(1)),
dropped: AtomicU64::new(0),
rejected: AtomicU64::new(0),
})
}
#[must_use]
pub fn operation_id(&self) -> OperationId {
self.operation_id
}
pub fn drain(&self, out: &mut Vec<NativeEventRecord>, max: usize) -> usize {
let mut taken = 0;
while taken < max {
let Some(record) = self.records.pop() else {
break;
};
out.push(record);
taken += 1;
}
taken
}
#[must_use]
pub fn dropped(&self) -> u64 {
self.dropped.load(Ordering::Relaxed)
}
#[must_use]
pub fn rejected(&self) -> u64 {
self.rejected.load(Ordering::Relaxed)
}
#[must_use]
pub fn len(&self) -> usize {
self.records.len()
}
#[must_use]
pub fn is_empty(&self) -> bool {
self.records.is_empty()
}
#[doc(hidden)]
pub unsafe fn deliver_for_test(&self, event: *const RawRuntimeEvent) {
unsafe { model_open_event_trampoline(event, self as *const Self as *mut c_void) };
}
fn ingest(&self, event: *const RawRuntimeEvent) {
match unsafe { NativeEventRecord::from_raw_ptr(event) } {
Ok(record) => {
if self.records.push(record).is_err() {
self.dropped.fetch_add(1, Ordering::Relaxed);
}
}
Err(_) => {
self.rejected.fetch_add(1, Ordering::Relaxed);
}
}
}
}
pub(super) unsafe extern "C" fn model_open_event_trampoline(
event: *const RawRuntimeEvent,
user_data: *mut c_void,
) {
let _ = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
if user_data.is_null() {
return;
}
let queue = unsafe { &*(user_data as *const ModelOpenEventQueue) };
queue.ingest(event);
}));
}
pub(super) struct ModelOpenEventReporterRegistration {
_queue: Arc<ModelOpenEventQueue>,
reporter: RawRuntimeEventReporter,
}
impl ModelOpenEventReporterRegistration {
pub(super) fn new(queue: &Arc<ModelOpenEventQueue>) -> Self {
let queue = Arc::clone(queue);
let reporter = RawRuntimeEventReporter {
abi_version: RUNTIME_EVENT_V1_ABI_VERSION,
struct_size: mem::size_of::<RawRuntimeEventReporter>() as u32,
callback: Some(model_open_event_trampoline),
user_data: Arc::as_ptr(&queue) as *mut c_void,
};
Self {
_queue: queue,
reporter,
}
}
pub(super) fn reporter_ptr(&self) -> *const RawRuntimeEventReporter {
&self.reporter
}
}
#[cfg(test)]
mod tests;