Skip to main content

aptabase_rs/
client.rs

1use rand::RngExt;
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
18pub fn 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(id: Option<String>) -> Self {
41        Self {
42            id: id
43                .filter(|id| !id.is_empty())
44                .unwrap_or_else(new_session_id),
45            last_touch_ts: OffsetDateTime::now_utc(),
46        }
47    }
48}
49
50/// The Aptabase client used to track events.
51pub struct AptabaseClient {
52    is_enabled: bool,
53    session: SyncMutex<TrackingSession>,
54    dispatcher: Arc<EventDispatcher>,
55    app_version: String,
56    sys_info: SystemProperties,
57}
58
59impl AptabaseClient {
60    /// Creates a new Aptabase client.
61    pub(crate) fn new(config: &Config, app_version: String) -> Self {
62        let sys_info = sys::get_info();
63
64        let is_enabled = !config.app_key.is_empty();
65        let dispatcher = Arc::new(EventDispatcher::new(config, &sys_info));
66
67        Self {
68            is_enabled,
69            dispatcher,
70            session: SyncMutex::new(TrackingSession::new(config.session_id.clone())),
71            app_version,
72            sys_info,
73        }
74    }
75
76    /// Starts the event dispatcher loop.
77    pub(crate) fn start_polling(&self, interval: Duration) {
78        let dispatcher = self.dispatcher.clone();
79
80        tokio::spawn(async move {
81            loop {
82                tokio::time::sleep(interval).await;
83                dispatcher.flush().await;
84            }
85        });
86    }
87
88    /// Returns the current session ID, creating a new one if necessary.
89    pub(crate) fn eval_session_id(&self) -> String {
90        let mut session = self.session.lock().expect("could not lock events");
91
92        let now = OffsetDateTime::now_utc();
93        if (now - session.last_touch_ts) > SESSION_TIMEOUT {
94            *session = TrackingSession::new(None);
95        } else {
96            session.last_touch_ts = now;
97        }
98
99        session.id.clone()
100    }
101
102    /// Enqueues an event to be sent to the server.
103    pub fn track_event(&self, name: &str, props: Option<Value>) -> Result<(), String> {
104        if !self.is_enabled {
105            return Ok(());
106        }
107
108        if let Some(props) = &props
109            && !matches!(props, Value::Object(_))
110        {
111            return Err(
112                "props must be `None` or the `Object` variation of `serde_json::Value`".to_owned(),
113            );
114        }
115
116        let ev = json!({
117            "timestamp": OffsetDateTime::now_utc().format(&Rfc3339).unwrap(),
118            "sessionId": self.eval_session_id(),
119            "eventName": name,
120            "systemProps": {
121                "isDebug": self.sys_info.is_debug,
122                "osName": self.sys_info.os_name,
123                "osVersion": self.sys_info.os_version,
124                "locale": self.sys_info.locale,
125                "appVersion": self.app_version,
126                "sdkVersion": concat!(env!("CARGO_PKG_NAME"), "@", env!("CARGO_PKG_VERSION"))
127            },
128            "props": props
129        });
130
131        self.dispatcher.enqueue(ev);
132
133        Ok(())
134    }
135
136    /// Flushes the event queue.
137    pub async fn flush(&self) {
138        self.dispatcher.flush().await;
139    }
140
141    /// Flushes the event queue, blocking the current thread.
142    pub(crate) fn flush_blocking(&self) {
143        futures::executor::block_on(async {
144            self.flush().await;
145        });
146    }
147}
148
149#[cfg(test)]
150mod tests {
151    use super::*;
152
153    #[test]
154    fn new_session_id_is_not_empty() {
155        // Act
156        let session_id = new_session_id();
157
158        // Assert
159        assert!(!session_id.is_empty());
160    }
161}