Skip to main content

systemprompt_users/repository/user/
operations.rs

1//! User row creation and update operations.
2//!
3//! Copyright (c) systemprompt.io — Business Source License 1.1.
4//! See <https://systemprompt.io> for licensing details.
5
6use chrono::{Duration, Utc};
7use systemprompt_identifiers::UserId;
8
9use crate::error::{Result, UserError};
10use crate::models::{User, UserRole, UserStatus, normalise_email};
11use crate::repository::UserRepository;
12
13#[derive(Debug)]
14pub struct UpdateUserParams<'a> {
15    pub email: &'a str,
16    pub full_name: Option<&'a str>,
17    pub display_name: Option<&'a str>,
18    pub status: UserStatus,
19}
20
21impl UserRepository {
22    pub async fn create(
23        &self,
24        name: &str,
25        email: &str,
26        full_name: Option<&str>,
27        display_name: Option<&str>,
28    ) -> Result<User> {
29        let now = Utc::now();
30        let id = UserId::new(uuid::Uuid::new_v4().to_string());
31        let display_name_val = display_name.or(full_name);
32        let status = UserStatus::Active.as_str();
33        let role = UserRole::User.as_str();
34        let email = normalise_email(email);
35
36        let row = sqlx::query_as!(
37            User,
38            r#"
39            INSERT INTO users (
40                id, name, email, full_name, display_name,
41                status, email_verified, roles, is_bot,
42                created_at, updated_at
43            )
44            VALUES ($1, $2, $3, $4, $5, $6, false, ARRAY[$7]::TEXT[], false, $8, $8)
45            RETURNING id, name, email, full_name, display_name, status, email_verified,
46                      roles, avatar_url, is_bot, is_scanner, created_at, updated_at
47            "#,
48            id.as_str(),
49            name,
50            email,
51            full_name,
52            display_name_val,
53            status,
54            role,
55            now
56        )
57        .fetch_one(&*self.write_pool)
58        .await?;
59
60        Ok(row)
61    }
62
63    pub async fn create_if_absent(
64        &self,
65        name: &str,
66        email: &str,
67        full_name: Option<&str>,
68        display_name: Option<&str>,
69    ) -> Result<Option<User>> {
70        let now = Utc::now();
71        let id = UserId::new(uuid::Uuid::new_v4().to_string());
72        let display_name_val = display_name.or(full_name);
73        let status = UserStatus::Active.as_str();
74        let role = UserRole::User.as_str();
75        let email = normalise_email(email);
76
77        let row = sqlx::query_as!(
78            User,
79            r#"
80            INSERT INTO users (
81                id, name, email, full_name, display_name,
82                status, email_verified, roles, is_bot,
83                created_at, updated_at
84            )
85            VALUES ($1, $2, $3, $4, $5, $6, false, ARRAY[$7]::TEXT[], false, $8, $8)
86            ON CONFLICT DO NOTHING
87            RETURNING id, name, email, full_name, display_name, status, email_verified,
88                      roles, avatar_url, is_bot, is_scanner, created_at, updated_at
89            "#,
90            id.as_str(),
91            name,
92            email,
93            full_name,
94            display_name_val,
95            status,
96            role,
97            now
98        )
99        .fetch_optional(&*self.write_pool)
100        .await?;
101
102        Ok(row)
103    }
104
105    pub async fn create_anonymous(&self, fingerprint: &str) -> Result<User> {
106        let email = normalise_email(&format!("{}@anonymous.local", fingerprint));
107
108        if let Some(existing) = sqlx::query_as!(
109            User,
110            r#"
111            SELECT id, name, email, full_name, display_name, status, email_verified,
112                   roles, avatar_url, is_bot, is_scanner, created_at, updated_at
113            FROM users
114            WHERE email = $1
115            "#,
116            email
117        )
118        .fetch_optional(&*self.pool)
119        .await?
120        {
121            return Ok(existing);
122        }
123
124        let user_id = uuid::Uuid::new_v4();
125        let id = UserId::new(user_id.to_string());
126        let name = format!("anonymous_{}", &user_id.to_string()[..8]);
127        let now = Utc::now();
128        let status = UserStatus::Active.as_str();
129        let role = UserRole::Anonymous.as_str();
130
131        let row = sqlx::query_as!(
132            User,
133            r#"
134            INSERT INTO users (
135                id, name, email, status, email_verified, roles,
136                is_bot, created_at, updated_at
137            )
138            VALUES ($1, $2, $3, $4, false, ARRAY[$5]::TEXT[], false, $6, $6)
139            ON CONFLICT (email) DO UPDATE SET updated_at = $6
140            RETURNING id, name, email, full_name, display_name, status, email_verified,
141                      roles, avatar_url, is_bot, is_scanner, created_at, updated_at
142            "#,
143            id.as_str(),
144            name,
145            email,
146            status,
147            role,
148            now
149        )
150        .fetch_one(&*self.write_pool)
151        .await?;
152
153        Ok(row)
154    }
155
156    pub async fn update_email(&self, id: &UserId, email: &str) -> Result<User> {
157        let row = sqlx::query_as!(
158            User,
159            r#"
160            UPDATE users
161            SET email = $1, email_verified = false, updated_at = $2
162            WHERE id = $3
163            RETURNING id, name, email, full_name, display_name, status, email_verified,
164                      roles, avatar_url, is_bot, is_scanner, created_at, updated_at
165            "#,
166            email,
167            Utc::now(),
168            id.as_str()
169        )
170        .fetch_optional(&*self.write_pool)
171        .await?
172        .ok_or_else(|| UserError::NotFound(id.clone()))?;
173
174        Ok(row)
175    }
176
177    pub async fn update_full_name(&self, id: &UserId, full_name: &str) -> Result<User> {
178        let row = sqlx::query_as!(
179            User,
180            r#"
181            UPDATE users
182            SET full_name = $1, updated_at = $2
183            WHERE id = $3
184            RETURNING id, name, email, full_name, display_name, status, email_verified,
185                      roles, avatar_url, is_bot, is_scanner, created_at, updated_at
186            "#,
187            full_name,
188            Utc::now(),
189            id.as_str()
190        )
191        .fetch_optional(&*self.write_pool)
192        .await?
193        .ok_or_else(|| UserError::NotFound(id.clone()))?;
194
195        Ok(row)
196    }
197
198    pub async fn update_status(&self, id: &UserId, status: UserStatus) -> Result<User> {
199        let row = sqlx::query_as!(
200            User,
201            r#"
202            UPDATE users
203            SET status = $1, updated_at = $2
204            WHERE id = $3
205            RETURNING id, name, email, full_name, display_name, status, email_verified,
206                      roles, avatar_url, is_bot, is_scanner, created_at, updated_at
207            "#,
208            status.as_str(),
209            Utc::now(),
210            id.as_str()
211        )
212        .fetch_optional(&*self.write_pool)
213        .await?
214        .ok_or_else(|| UserError::NotFound(id.clone()))?;
215
216        Ok(row)
217    }
218
219    pub async fn update_email_verified(&self, id: &UserId, verified: bool) -> Result<User> {
220        let row = sqlx::query_as!(
221            User,
222            r#"
223            UPDATE users
224            SET email_verified = $1, updated_at = $2
225            WHERE id = $3
226            RETURNING id, name, email, full_name, display_name, status, email_verified,
227                      roles, avatar_url, is_bot, is_scanner, created_at, updated_at
228            "#,
229            verified,
230            Utc::now(),
231            id.as_str()
232        )
233        .fetch_optional(&*self.write_pool)
234        .await?
235        .ok_or_else(|| UserError::NotFound(id.clone()))?;
236
237        Ok(row)
238    }
239
240    pub async fn update_display_name(&self, id: &UserId, display_name: &str) -> Result<User> {
241        let row = sqlx::query_as!(
242            User,
243            r#"
244            UPDATE users
245            SET display_name = $1, updated_at = $2
246            WHERE id = $3
247            RETURNING id, name, email, full_name, display_name, status, email_verified,
248                      roles, avatar_url, is_bot, is_scanner, created_at, updated_at
249            "#,
250            display_name,
251            Utc::now(),
252            id.as_str()
253        )
254        .fetch_optional(&*self.write_pool)
255        .await?
256        .ok_or_else(|| UserError::NotFound(id.clone()))?;
257
258        Ok(row)
259    }
260
261    pub async fn update_all_fields(
262        &self,
263        id: &UserId,
264        params: UpdateUserParams<'_>,
265    ) -> Result<User> {
266        let row = sqlx::query_as!(
267            User,
268            r#"
269            UPDATE users
270            SET email = $1, full_name = $2, display_name = $3, status = $4, updated_at = $5
271            WHERE id = $6
272            RETURNING id, name, email, full_name, display_name, status, email_verified,
273                      roles, avatar_url, is_bot, is_scanner, created_at, updated_at
274            "#,
275            params.email,
276            params.full_name,
277            params.display_name,
278            params.status.as_str(),
279            Utc::now(),
280            id.as_str()
281        )
282        .fetch_optional(&*self.write_pool)
283        .await?
284        .ok_or_else(|| UserError::NotFound(id.clone()))?;
285
286        Ok(row)
287    }
288
289    pub async fn assign_roles(&self, id: &UserId, roles: &[String]) -> Result<User> {
290        let row = sqlx::query_as!(
291            User,
292            r#"
293            UPDATE users
294            SET roles = $1, updated_at = $2
295            WHERE id = $3
296            RETURNING id, name, email, full_name, display_name, status, email_verified,
297                      roles, avatar_url, is_bot, is_scanner, created_at, updated_at
298            "#,
299            roles,
300            Utc::now(),
301            id.as_str()
302        )
303        .fetch_optional(&*self.write_pool)
304        .await?
305        .ok_or_else(|| UserError::NotFound(id.clone()))?;
306
307        Ok(row)
308    }
309
310    pub async fn delete(&self, id: &UserId) -> Result<()> {
311        let result = sqlx::query!(r#"DELETE FROM users WHERE id = $1"#, id.as_str())
312            .execute(&*self.write_pool)
313            .await?;
314
315        if result.rows_affected() == 0 {
316            return Err(UserError::NotFound(id.clone()));
317        }
318
319        Ok(())
320    }
321
322    pub async fn cleanup_old_anonymous(&self, days: i32) -> Result<u64> {
323        let cutoff = Utc::now() - Duration::days(i64::from(days));
324        let anonymous_role = UserRole::Anonymous.as_str();
325        let result = sqlx::query!(
326            r#"
327            DELETE FROM users u
328            WHERE $1 = ANY(u.roles)
329              AND u.created_at < $2
330              AND NOT EXISTS (
331                  SELECT 1
332                  FROM user_sessions s
333                  WHERE s.user_id = u.id
334                    AND s.ended_at IS NULL
335              )
336            "#,
337            anonymous_role,
338            cutoff
339        )
340        .execute(&*self.write_pool)
341        .await?;
342
343        Ok(result.rows_affected())
344    }
345
346    pub async fn count_old_anonymous(&self, days: i32) -> Result<i64> {
347        let cutoff = Utc::now() - Duration::days(i64::from(days));
348        let anonymous_role = UserRole::Anonymous.as_str();
349        let count = sqlx::query_scalar!(
350            r#"
351            SELECT COUNT(*) as "count!"
352            FROM users u
353            WHERE $1 = ANY(u.roles)
354              AND u.created_at < $2
355              AND NOT EXISTS (
356                  SELECT 1
357                  FROM user_sessions s
358                  WHERE s.user_id = u.id
359                    AND s.ended_at IS NULL
360              )
361            "#,
362            anonymous_role,
363            cutoff
364        )
365        .fetch_one(&*self.write_pool)
366        .await?;
367
368        Ok(count)
369    }
370}