Skip to main content

aptabase_rs/
client.rs

1use rand::Rng;
2use serde_json::{Value, json};
3use std::time::{SystemTime, UNIX_EPOCH};
4use std::{
5    sync::{Arc, Mutex as SyncMutex},
6    time::Duration,
7};
8use time::{OffsetDateTime, format_description::well_known::Rfc3339};
9
10use crate::{
11    config::Config,
12    dispatcher::EventDispatcher,
13    sys::{self, SystemProperties},
14};
15
16static SESSION_TIMEOUT: Duration = Duration::from_secs(4 * 60 * 60);
17
18fn new_session_id() -> String {
19    let epoch_in_seconds = SystemTime::now()
20        .duration_since(UNIX_EPOCH)
21        .expect("time went backwards")
22        .as_secs();
23
24    let mut rng = rand::rng();
25    let random: u64 = rng.random_range(0..=99999999);
26
27    let id = epoch_in_seconds * 100_000_000 + random;
28
29    id.to_string()
30}
31
32/// A tracking session.
33#[derive(Debug, Clone)]
34pub struct TrackingSession {
35    pub id: String,
36    pub last_touch_ts: OffsetDateTime,
37}
38
39impl TrackingSession {
40    fn new() -> Self {
41        Self {
42            id: new_session_id(),
43            last_touch_ts: OffsetDateTime::now_utc(),
44        }
45    }
46}
47
48/// The Aptabase client used to track events.
49pub struct AptabaseClient {
50    is_enabled: bool,
51    session: SyncMutex<TrackingSession>,
52    dispatcher: Arc<EventDispatcher>,
53    app_version: String,
54    sys_info: SystemProperties,
55}
56
57impl AptabaseClient {
58    /// Creates a new Aptabase client.
59    pub(crate) fn new(config: &Config, app_version: String) -> Self {
60        let sys_info = sys::get_info();
61
62        let is_enabled = !config.app_key.is_empty();
63        let dispatcher = Arc::new(EventDispatcher::new(config, &sys_info));
64
65        Self {
66            is_enabled,
67            dispatcher,
68            session: SyncMutex::new(TrackingSession::new()),
69            app_version,
70            sys_info,
71        }
72    }
73
74    /// Starts the event dispatcher loop.
75    pub(crate) fn start_polling(&self, interval: Duration) {
76        let dispatcher = self.dispatcher.clone();
77
78        tokio::spawn(async move {
79            loop {
80                tokio::time::sleep(interval).await;
81                dispatcher.flush().await;
82            }
83        });
84    }
85
86    /// Returns the current session ID, creating a new one if necessary.
87    pub(crate) fn eval_session_id(&self) -> String {
88        let mut session = self.session.lock().expect("could not lock events");
89
90        let now = OffsetDateTime::now_utc();
91        if (now - session.last_touch_ts) > SESSION_TIMEOUT {
92            *session = TrackingSession::new();
93        } else {
94            session.last_touch_ts = now;
95        }
96
97        session.id.clone()
98    }
99
100    /// Enqueues an event to be sent to the server.
101    pub fn track_event(&self, name: &str, props: Option<Value>) -> Result<(), String> {
102        if !self.is_enabled {
103            return Ok(());
104        }
105
106        if let Some(props) = &props
107            && !matches!(props, Value::Object(_))
108        {
109            return Err(
110                "props must be `None` or the `Object` variation of `serde_json::Value`".to_owned(),
111            );
112        }
113
114        let ev = json!({
115            "timestamp": OffsetDateTime::now_utc().format(&Rfc3339).unwrap(),
116            "sessionId": self.eval_session_id(),
117            "eventName": name,
118            "systemProps": {
119                "isDebug": self.sys_info.is_debug,
120                "osName": self.sys_info.os_name,
121                "osVersion": self.sys_info.os_version,
122                "locale": self.sys_info.locale,
123                "appVersion": self.app_version,
124                "sdkVersion": concat!(env!("CARGO_PKG_NAME"), "@", env!("CARGO_PKG_VERSION"))
125            },
126            "props": props
127        });
128
129        self.dispatcher.enqueue(ev);
130
131        Ok(())
132    }
133
134    /// Flushes the event queue.
135    pub async fn flush(&self) {
136        self.dispatcher.flush().await;
137    }
138
139    /// Flushes the event queue, blocking the current thread.
140    pub(crate) fn flush_blocking(&self) {
141        futures::executor::block_on(async {
142            self.flush().await;
143        });
144    }
145}