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#[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
50pub 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 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 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 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 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 pub async fn flush(&self) {
138 self.dispatcher.flush().await;
139 }
140
141 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 let session_id = new_session_id();
157
158 assert!(!session_id.is_empty());
160 }
161}