use async_trait::async_trait;
use lazy_static::lazy_static;
use std::{
collections::HashMap,
sync::{Arc, RwLock, Weak},
};
use tokio::sync::broadcast::{self, Sender};
use tokio::sync::Notify;
use crate::{log_e, log_w, StatsigRuntime};
use super::{
observability_client_adapter::ObservabilityEvent, sdk_errors_observer::ErrorBoundaryEvent,
};
const TAG: &str = stringify!(OpsStats);
lazy_static! {
pub static ref OPS_STATS: OpsStats = OpsStats::new();
}
pub struct OpsStats {
instances_map: RwLock<HashMap<String, Weak<OpsStatsForInstance>>>, }
impl OpsStats {
pub fn new() -> Self {
OpsStats {
instances_map: HashMap::new().into(),
}
}
pub fn get_for_instance(&self, sdk_key: &str) -> Arc<OpsStatsForInstance> {
match self.instances_map.read() {
Ok(read_guard) => {
if let Some(instance) = read_guard.get(sdk_key) {
if let Some(instance) = instance.upgrade() {
return instance.clone();
}
}
}
Err(e) => {
log_e!(TAG, "Failed to get read guard: {}", e);
}
}
let instance = Arc::new(OpsStatsForInstance::new());
match self.instances_map.write() {
Ok(mut write_guard) => {
write_guard.insert(sdk_key.into(), Arc::downgrade(&instance));
}
Err(e) => {
log_e!(TAG, "Failed to get write guard: {}", e);
}
}
instance
}
}
#[derive(Clone)]
pub enum OpsStatsEvent {
ObservabilityEvent(ObservabilityEvent),
SDKErrorEvent(ErrorBoundaryEvent),
}
pub struct OpsStatsForInstance {
sender: Sender<OpsStatsEvent>,
shutdown_notify: Arc<Notify>,
}
impl Default for OpsStatsForInstance {
fn default() -> Self {
Self::new()
}
}
impl OpsStatsForInstance {
pub fn new() -> Self {
let (tx, _) = broadcast::channel(1000);
OpsStatsForInstance {
sender: tx,
shutdown_notify: Arc::new(Notify::new()),
}
}
pub fn log(&self, event: OpsStatsEvent) {
match self.sender.send(event) {
Ok(_) => {}
Err(e) => {
log_w!(
"OpsStats Message Queue",
"Dropping ops stats event {}",
e.to_string()
);
}
}
}
pub fn log_error(&self, error: ErrorBoundaryEvent) {
self.log(OpsStatsEvent::SDKErrorEvent(error))
}
pub fn subscribe(
&self,
runtime: Arc<StatsigRuntime>,
observer: Weak<dyn OpsStatsEventObserver>,
) {
let mut rx = self.sender.subscribe();
let shutdown_notify = self.shutdown_notify.clone();
let _ = runtime.spawn("opts_stats_listen_for", |rt_shutdown_notify| async move {
loop {
tokio::select! {
event = rx.recv() => {
let observer = match observer.upgrade() {
Some(observer) => observer,
None => break,
};
if let Ok(event) = event {
observer.handle_event(event).await;
}
}
_ = rt_shutdown_notify.notified() => {
break;
}
_ = shutdown_notify.notified() => {
break;
}
}
}
});
}
}
impl Drop for OpsStatsForInstance {
fn drop(&mut self) {
self.shutdown_notify.notify_waiters();
}
}
#[async_trait]
pub trait OpsStatsEventObserver: Send + Sync + 'static {
async fn handle_event(&self, event: OpsStatsEvent);
}