Skip to main content

ag_store/
activity.rs

1//! Session-activity persistence adapters and query helpers.
2
3use async_trait::async_trait;
4use sqlx::SqlitePool;
5
6use crate::DbError;
7
8/// Session-activity persistence boundary used by app orchestration and tests.
9#[cfg_attr(test, mockall::automock)]
10#[async_trait]
11pub trait ActivityRepository: Send + Sync {
12    #[cfg(any(test, feature = "test-utils"))]
13    /// Rebuilds `session_activity` rows from current `session.created_at`.
14    async fn backfill_session_activity_from_sessions(&self) -> Result<(), DbError>;
15
16    #[cfg(any(test, feature = "test-utils"))]
17    /// Deletes all rows from `session_activity`.
18    async fn clear_session_activity(&self) -> Result<(), DbError>;
19
20    /// Persists one session-creation activity event at a specific Unix
21    /// timestamp.
22    async fn insert_session_creation_activity_at(
23        &self,
24        session_id: &str,
25        timestamp_seconds: i64,
26    ) -> Result<(), DbError>;
27
28    /// Loads persisted session-creation timestamps for clock-aware activity
29    /// aggregation.
30    async fn load_session_activity_timestamps(&self) -> Result<Vec<i64>, DbError>;
31}
32
33/// `SQLite` implementation of [`ActivityRepository`].
34#[derive(Clone)]
35pub(crate) struct SqliteActivityRepository(SqlitePool);
36
37impl SqliteActivityRepository {
38    /// Creates an activity repository backed by the provided pool.
39    pub(crate) fn new(pool: SqlitePool) -> Self {
40        Self(pool)
41    }
42}
43
44/// Row returned when loading one session activity timestamp.
45struct TimestampValueRow {
46    created_at: i64,
47}
48
49#[async_trait]
50impl ActivityRepository for SqliteActivityRepository {
51    #[cfg(any(test, feature = "test-utils"))]
52    async fn backfill_session_activity_from_sessions(&self) -> Result<(), DbError> {
53        sqlx::query!(
54            r"
55INSERT INTO session_activity (session_id, created_at)
56SELECT id, created_at
57FROM session
58"
59        )
60        .execute(&self.0)
61        .await?;
62
63        Ok(())
64    }
65
66    #[cfg(any(test, feature = "test-utils"))]
67    async fn clear_session_activity(&self) -> Result<(), DbError> {
68        sqlx::query!(
69            r"
70DELETE FROM session_activity
71"
72        )
73        .execute(&self.0)
74        .await?;
75
76        Ok(())
77    }
78
79    async fn insert_session_creation_activity_at(
80        &self,
81        session_id: &str,
82        timestamp_seconds: i64,
83    ) -> Result<(), DbError> {
84        sqlx::query!(
85            r"
86INSERT INTO session_activity (session_id, created_at)
87VALUES (?, ?)
88ON CONFLICT(session_id) DO NOTHING
89",
90            session_id,
91            timestamp_seconds
92        )
93        .execute(&self.0)
94        .await?;
95
96        Ok(())
97    }
98
99    async fn load_session_activity_timestamps(&self) -> Result<Vec<i64>, DbError> {
100        let rows = sqlx::query_as!(
101            TimestampValueRow,
102            r"
103SELECT created_at
104FROM session_activity
105ORDER BY id
106",
107        )
108        .fetch_all(&self.0)
109        .await?;
110
111        Ok(rows.into_iter().map(|row| row.created_at).collect())
112    }
113}