1use super::{map_db_error, parse_datetime_row, parse_enum_row, parse_json_row, parse_uuid_row};
4use chrono::Utc;
5use r2d2::Pool;
6use r2d2_sqlite::SqliteConnectionManager;
7use stateset_core::{
8 ActivityLogEntry, ActivityLogFilter, ActivityLogId, ActivityLogRepository, ActorKind,
9 CommerceError, RecordActivity, Result,
10};
11use uuid::Uuid;
12
13#[derive(Debug)]
14pub struct SqliteActivityLogRepository {
15 pool: Pool<SqliteConnectionManager>,
16}
17
18impl SqliteActivityLogRepository {
19 #[must_use]
20 pub const fn new(pool: Pool<SqliteConnectionManager>) -> Self {
21 Self { pool }
22 }
23
24 fn conn(&self) -> Result<r2d2::PooledConnection<SqliteConnectionManager>> {
25 self.pool.get().map_err(|e| CommerceError::DatabaseError(e.to_string()))
26 }
27
28 fn row_to_entry(row: &rusqlite::Row<'_>) -> rusqlite::Result<ActivityLogEntry> {
29 let metadata_json: String = row.get("metadata")?;
30 Ok(ActivityLogEntry {
31 id: parse_uuid_row(&row.get::<_, String>("id")?, "activity_log", "id")?.into(),
32 subject_type: row.get("subject_type")?,
33 subject_id: parse_uuid_row(
34 &row.get::<_, String>("subject_id")?,
35 "activity_log",
36 "subject_id",
37 )?,
38 action: row.get("action")?,
39 summary: row.get("summary")?,
40 actor_kind: parse_enum_row::<ActorKind>(
41 &row.get::<_, String>("actor_kind")?,
42 "activity_log",
43 "actor_kind",
44 )?,
45 actor: row.get("actor")?,
46 metadata: parse_json_row(&metadata_json, "activity_log", "metadata")?,
47 created_at: parse_datetime_row(
48 &row.get::<_, String>("created_at")?,
49 "activity_log",
50 "created_at",
51 )?,
52 })
53 }
54}
55
56impl ActivityLogRepository for SqliteActivityLogRepository {
57 fn record(&self, input: RecordActivity) -> Result<ActivityLogEntry> {
58 let id = ActivityLogId::new();
59 let id_str = id.to_string();
60 let now_str = Utc::now().to_rfc3339();
61 let metadata_json = serde_json::to_string(&input.metadata)
62 .map_err(|e| CommerceError::DatabaseError(e.to_string()))?;
63 let conn = self.conn()?;
64 conn.execute(
65 "INSERT INTO activity_logs (id, subject_type, subject_id, action, summary, actor_kind, actor, metadata, created_at)
66 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)",
67 rusqlite::params![
68 &id_str,
69 &input.subject_type,
70 input.subject_id.to_string(),
71 &input.action,
72 &input.summary,
73 input.actor_kind.to_string(),
74 &input.actor,
75 &metadata_json,
76 &now_str,
77 ],
78 )
79 .map_err(map_db_error)?;
80 conn.query_row("SELECT * FROM activity_logs WHERE id = ?", [&id_str], Self::row_to_entry)
81 .map_err(map_db_error)
82 }
83
84 fn get(&self, id: ActivityLogId) -> Result<Option<ActivityLogEntry>> {
85 let conn = self.conn()?;
86 match conn.query_row(
87 "SELECT * FROM activity_logs WHERE id = ?",
88 [id.to_string()],
89 Self::row_to_entry,
90 ) {
91 Ok(e) => Ok(Some(e)),
92 Err(rusqlite::Error::QueryReturnedNoRows) => Ok(None),
93 Err(e) => Err(map_db_error(e)),
94 }
95 }
96
97 fn list(&self, filter: ActivityLogFilter) -> Result<Vec<ActivityLogEntry>> {
98 let conn = self.conn()?;
99 let mut sql = "SELECT * FROM activity_logs WHERE 1=1".to_string();
100 let mut params: Vec<Box<dyn rusqlite::types::ToSql>> = vec![];
101 if let Some(ref subject_type) = filter.subject_type {
102 sql.push_str(" AND subject_type = ?");
103 params.push(Box::new(subject_type.clone()));
104 }
105 if let Some(subject_id) = filter.subject_id {
106 sql.push_str(" AND subject_id = ?");
107 params.push(Box::new(subject_id.to_string()));
108 }
109 if let Some(ref action) = filter.action {
110 sql.push_str(" AND action = ?");
111 params.push(Box::new(action.clone()));
112 }
113 if let Some(actor_kind) = filter.actor_kind {
114 sql.push_str(" AND actor_kind = ?");
115 params.push(Box::new(actor_kind.to_string()));
116 }
117 sql.push_str(" ORDER BY created_at DESC");
118 crate::sqlite::append_limit_offset(&mut sql, filter.limit, filter.offset);
119 let param_refs: Vec<&dyn rusqlite::types::ToSql> =
120 params.iter().map(|p| p.as_ref()).collect();
121 let mut stmt = conn.prepare(&sql).map_err(map_db_error)?;
122 let rows = stmt
123 .query_map(param_refs.as_slice(), Self::row_to_entry)
124 .map_err(map_db_error)?
125 .collect::<std::result::Result<Vec<_>, _>>()
126 .map_err(map_db_error)?;
127 Ok(rows)
128 }
129
130 fn history_for_subject(
131 &self,
132 subject_type: &str,
133 subject_id: Uuid,
134 ) -> Result<Vec<ActivityLogEntry>> {
135 let conn = self.conn()?;
136 let mut stmt = conn
137 .prepare("SELECT * FROM activity_logs WHERE subject_type = ? AND subject_id = ? ORDER BY created_at DESC")
138 .map_err(map_db_error)?;
139 let rows = stmt
140 .query_map(rusqlite::params![subject_type, subject_id.to_string()], Self::row_to_entry)
141 .map_err(map_db_error)?
142 .collect::<std::result::Result<Vec<_>, _>>()
143 .map_err(map_db_error)?;
144 Ok(rows)
145 }
146}
147
148#[cfg(test)]
149mod tests {
150 use super::*;
151 use crate::DatabaseConfig;
152 use crate::sqlite::SqliteDatabase;
153
154 fn test_repo() -> SqliteActivityLogRepository {
155 let db = SqliteDatabase::new(&DatabaseConfig::in_memory()).expect("in-memory db");
156 SqliteActivityLogRepository::new(db.pool().clone())
157 }
158
159 fn record(repo: &SqliteActivityLogRepository, subject: Uuid, action: &str) -> ActivityLogEntry {
160 repo.record(RecordActivity {
161 subject_type: "sales_order".into(),
162 subject_id: subject,
163 action: action.into(),
164 summary: format!("did {action}"),
165 actor_kind: ActorKind::User,
166 actor: Some("alice".into()),
167 metadata: serde_json::json!({"k": "v"}),
168 })
169 .expect("record")
170 }
171
172 #[test]
173 fn record_and_get() {
174 let repo = test_repo();
175 let subject = Uuid::new_v4();
176 let e = record(&repo, subject, "created");
177 assert_eq!(e.actor_label(), "alice");
178 let fetched = repo.get(e.id).expect("get").expect("found");
179 assert_eq!(fetched.action, "created");
180 assert_eq!(fetched.metadata["k"], "v");
181 }
182
183 #[test]
184 fn history_scoped_to_subject() {
185 let repo = test_repo();
186 let a = Uuid::new_v4();
187 let b = Uuid::new_v4();
188 record(&repo, a, "created");
189 record(&repo, a, "status_changed");
190 record(&repo, b, "created");
191 let hist = repo.history_for_subject("sales_order", a).expect("history");
192 assert_eq!(hist.len(), 2);
193 assert_eq!(hist[0].action, "status_changed");
195 }
196
197 #[test]
198 fn list_filters_by_action() {
199 let repo = test_repo();
200 let subject = Uuid::new_v4();
201 record(&repo, subject, "created");
202 record(&repo, subject, "status_changed");
203 let changed = repo
204 .list(ActivityLogFilter { action: Some("status_changed".into()), ..Default::default() })
205 .expect("list");
206 assert_eq!(changed.len(), 1);
207 }
208
209 #[test]
210 fn list_filters_by_subject() {
211 let repo = test_repo();
212 let a = Uuid::new_v4();
213 record(&repo, a, "created");
214 record(&repo, Uuid::new_v4(), "created");
215 let scoped = repo
216 .list(ActivityLogFilter { subject_id: Some(a), ..Default::default() })
217 .expect("list");
218 assert_eq!(scoped.len(), 1);
219 }
220}