1use super::backends::{SessionBackend, SessionError};
28use async_trait::async_trait;
29use chrono::{DateTime, Utc};
30use serde::{Deserialize, Serialize};
31use std::sync::Arc;
32
33mod logger;
35pub use logger::LoggerAnalytics;
36
37#[cfg(feature = "analytics-prometheus")]
38mod prometheus;
39#[cfg(feature = "analytics-prometheus")]
40pub use self::prometheus::PrometheusAnalytics;
41
42#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
44pub enum DeletionReason {
45 Explicit,
47 Expired,
49 Invalidated,
51 Replaced,
53}
54
55#[derive(Debug, Clone)]
73pub enum SessionEvent {
74 Created {
76 session_key: String,
78 size_bytes: usize,
80 ttl_secs: Option<u64>,
82 timestamp: DateTime<Utc>,
84 },
85 Accessed {
87 session_key: String,
89 latency_ms: u64,
91 hit: bool,
93 timestamp: DateTime<Utc>,
95 },
96 Deleted {
98 session_key: String,
100 reason: DeletionReason,
102 timestamp: DateTime<Utc>,
104 },
105 Expired {
107 session_key: String,
109 age_secs: u64,
111 timestamp: DateTime<Utc>,
113 },
114}
115
116#[async_trait]
137pub trait SessionAnalytics: Send + Sync {
138 async fn record_event(&self, event: SessionEvent);
140}
141
142#[derive(Clone)]
156pub struct CompositeAnalytics {
157 backends: Vec<Arc<dyn SessionAnalytics>>,
158}
159
160impl CompositeAnalytics {
161 pub fn new() -> Self {
163 Self {
164 backends: Vec::new(),
165 }
166 }
167
168 pub fn add<A: SessionAnalytics + 'static>(&mut self, analytics: A) {
170 self.backends.push(Arc::new(analytics));
171 }
172}
173
174impl Default for CompositeAnalytics {
175 fn default() -> Self {
176 Self::new()
177 }
178}
179
180#[async_trait]
181impl SessionAnalytics for CompositeAnalytics {
182 async fn record_event(&self, event: SessionEvent) {
183 for backend in &self.backends {
184 backend.record_event(event.clone()).await;
185 }
186 }
187}
188
189#[derive(Clone)]
210pub struct InstrumentedSessionBackend<B, A> {
211 backend: B,
212 analytics: A,
213}
214
215impl<B, A> InstrumentedSessionBackend<B, A>
216where
217 B: SessionBackend + Clone,
218 A: SessionAnalytics + Clone,
219{
220 pub fn new(backend: B, analytics: A) -> Self {
233 Self { backend, analytics }
234 }
235
236 pub fn backend(&self) -> &B {
238 &self.backend
239 }
240
241 pub fn analytics(&self) -> &A {
243 &self.analytics
244 }
245}
246
247#[async_trait]
248impl<B, A> SessionBackend for InstrumentedSessionBackend<B, A>
249where
250 B: SessionBackend + Clone,
251 A: SessionAnalytics + Clone,
252{
253 async fn load<T>(&self, session_key: &str) -> Result<Option<T>, SessionError>
254 where
255 T: for<'de> Deserialize<'de> + Serialize + Send + Sync,
256 {
257 let start = std::time::Instant::now();
258 let result = self.backend.load(session_key).await;
259 let latency_ms = start.elapsed().as_millis() as u64;
260
261 let hit = result.as_ref().map(|opt| opt.is_some()).unwrap_or(false);
262
263 self.analytics
264 .record_event(SessionEvent::Accessed {
265 session_key: session_key.to_string(),
266 latency_ms,
267 hit,
268 timestamp: Utc::now(),
269 })
270 .await;
271
272 result
273 }
274
275 async fn save<T>(
276 &self,
277 session_key: &str,
278 data: &T,
279 ttl: Option<u64>,
280 ) -> Result<(), SessionError>
281 where
282 T: Serialize + Send + Sync,
283 {
284 let serialized = serde_json::to_vec(data)
286 .map_err(|e| SessionError::SerializationError(e.to_string()))?;
287 let size_bytes = serialized.len();
288
289 let result = self.backend.save(session_key, data, ttl).await;
290
291 if result.is_ok() {
292 self.analytics
293 .record_event(SessionEvent::Created {
294 session_key: session_key.to_string(),
295 size_bytes,
296 ttl_secs: ttl,
297 timestamp: Utc::now(),
298 })
299 .await;
300 }
301
302 result
303 }
304
305 async fn delete(&self, session_key: &str) -> Result<(), SessionError> {
306 let result = self.backend.delete(session_key).await;
307
308 if result.is_ok() {
309 self.analytics
310 .record_event(SessionEvent::Deleted {
311 session_key: session_key.to_string(),
312 reason: DeletionReason::Explicit,
313 timestamp: Utc::now(),
314 })
315 .await;
316 }
317
318 result
319 }
320
321 async fn exists(&self, session_key: &str) -> Result<bool, SessionError> {
322 self.backend.exists(session_key).await
323 }
324}
325
326#[cfg(test)]
327mod tests {
328 use super::*;
329 use crate::sessions::InMemorySessionBackend;
330
331 #[tokio::test]
332 async fn test_instrumented_backend_save() {
333 let backend = InMemorySessionBackend::new();
334 let analytics = LoggerAnalytics::new();
335 let instrumented = InstrumentedSessionBackend::new(backend, analytics);
336
337 let data = serde_json::json!({"key": "value"});
338
339 instrumented
340 .save("test_key", &data, Some(3600))
341 .await
342 .unwrap();
343
344 let loaded: Option<serde_json::Value> = instrumented.load("test_key").await.unwrap();
345 assert_eq!(loaded.unwrap(), data);
346 }
347
348 #[tokio::test]
349 async fn test_instrumented_backend_delete() {
350 let backend = InMemorySessionBackend::new();
351 let analytics = LoggerAnalytics::new();
352 let instrumented = InstrumentedSessionBackend::new(backend, analytics);
353
354 let data = serde_json::json!({"key": "value"});
355
356 instrumented.save("test_key", &data, None).await.unwrap();
357 assert!(instrumented.exists("test_key").await.unwrap());
358
359 instrumented.delete("test_key").await.unwrap();
360 assert!(!instrumented.exists("test_key").await.unwrap());
361 }
362
363 #[tokio::test]
364 async fn test_composite_analytics() {
365 let mut composite = CompositeAnalytics::new();
366 composite.add(LoggerAnalytics::new());
367
368 let backend = InMemorySessionBackend::new();
369 let instrumented = InstrumentedSessionBackend::new(backend, composite);
370
371 let data = serde_json::json!({"key": "value"});
372 instrumented.save("test_key", &data, None).await.unwrap();
373 }
374}