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#[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
48pub 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 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 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 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 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 pub async fn flush(&self) {
136 self.dispatcher.flush().await;
137 }
138
139 pub(crate) fn flush_blocking(&self) {
141 futures::executor::block_on(async {
142 self.flush().await;
143 });
144 }
145}