users.rs (11484B)
1 //! User profile and address repository functions with OpenTelemetry instrumentation 2 //! 3 //! This module contains database operations for user profiles and addresses, 4 //! separated from HTTP handlers. Each function is instrumented with 5 //! OpenTelemetry semantic conventions for database spans. 6 7 use chrono::NaiveDate; 8 use sqlx::PgPool; 9 use tracing::{instrument, Span}; 10 use uuid::Uuid; 11 12 use crate::models::{AddAddressRequest, UserAddress}; 13 14 /// Profile data returned from the database 15 #[derive(Debug)] 16 pub struct ProfileData { 17 pub eid: Uuid, 18 pub avatar_url: Option<String>, 19 pub date_of_birth: Option<NaiveDate>, 20 pub email_notifications: bool, 21 pub marketing_emails: bool, 22 } 23 24 /// Get user profile by user ID 25 #[instrument( 26 name = "SELECT user_profiles", 27 skip(pool), 28 fields( 29 otel.kind = "client", 30 db.system.name = "postgresql", 31 db.namespace = "users", 32 db.operation.name = "SELECT", 33 db.collection.name = "user_profiles", 34 db.query.text = "SELECT uuid, avatar_url, date_of_birth, ... FROM user_profiles WHERE user_id = $1", 35 otelmart.user.id = user_id 36 ) 37 )] 38 pub async fn get_profile(pool: &PgPool, user_id: i32) -> Result<Option<ProfileData>, sqlx::Error> { 39 let row = sqlx::query_as::<_, (Uuid, Option<String>, Option<NaiveDate>, bool, bool)>( 40 r#" 41 SELECT 42 uuid, 43 avatar_url, 44 date_of_birth, 45 email_notifications, 46 marketing_emails 47 FROM user_profiles 48 WHERE user_id = $1 49 "#, 50 ) 51 .bind(user_id) 52 .fetch_optional(pool) 53 .await?; 54 55 Ok(row.map( 56 |(eid, avatar_url, date_of_birth, email_notifications, marketing_emails)| ProfileData { 57 eid, 58 avatar_url, 59 date_of_birth, 60 email_notifications, 61 marketing_emails, 62 }, 63 )) 64 } 65 66 /// Check if a profile exists for a user 67 #[instrument( 68 name = "SELECT user_profiles", 69 skip(pool), 70 fields( 71 otel.kind = "client", 72 db.system.name = "postgresql", 73 db.namespace = "users", 74 db.operation.name = "SELECT", 75 db.collection.name = "user_profiles", 76 db.query.text = "SELECT id FROM user_profiles WHERE user_id = $1", 77 otelmart.user.id = user_id 78 ) 79 )] 80 pub async fn profile_exists(pool: &PgPool, user_id: i32) -> Result<bool, sqlx::Error> { 81 let row = sqlx::query_as::<_, (i32,)>("SELECT id FROM user_profiles WHERE user_id = $1") 82 .bind(user_id) 83 .fetch_optional(pool) 84 .await?; 85 86 Ok(row.is_some()) 87 } 88 89 /// Update an existing user profile 90 #[instrument( 91 name = "UPDATE user_profiles", 92 skip(pool), 93 fields( 94 otel.kind = "client", 95 db.system.name = "postgresql", 96 db.namespace = "users", 97 db.operation.name = "UPDATE", 98 db.collection.name = "user_profiles", 99 db.query.text = "UPDATE user_profiles SET ... WHERE user_id = $5 RETURNING ...", 100 otelmart.user.id = user_id, 101 db.response.returned_rows = tracing::field::Empty 102 ) 103 )] 104 pub async fn update_profile( 105 pool: &PgPool, 106 user_id: i32, 107 avatar_url: Option<&str>, 108 date_of_birth: Option<NaiveDate>, 109 email_notifications: Option<bool>, 110 marketing_emails: Option<bool>, 111 ) -> Result<ProfileData, sqlx::Error> { 112 let (eid, avatar, dob, email_notif, marketing) = 113 sqlx::query_as::<_, (Uuid, Option<String>, Option<NaiveDate>, bool, bool)>( 114 r#" 115 UPDATE user_profiles 116 SET avatar_url = COALESCE($1, avatar_url), 117 date_of_birth = COALESCE($2, date_of_birth), 118 email_notifications = COALESCE($3, email_notifications), 119 marketing_emails = COALESCE($4, marketing_emails), 120 updated_at = NOW() 121 WHERE user_id = $5 122 RETURNING 123 uuid, 124 avatar_url, 125 date_of_birth, 126 email_notifications, 127 marketing_emails 128 "#, 129 ) 130 .bind(avatar_url) 131 .bind(date_of_birth) 132 .bind(email_notifications) 133 .bind(marketing_emails) 134 .bind(user_id) 135 .fetch_one(pool) 136 .await?; 137 138 Span::current().record("db.response.returned_rows", 1); 139 140 Ok(ProfileData { 141 eid, 142 avatar_url: avatar, 143 date_of_birth: dob, 144 email_notifications: email_notif, 145 marketing_emails: marketing, 146 }) 147 } 148 149 /// Create a new user profile 150 #[instrument( 151 name = "INSERT user_profiles", 152 skip(pool), 153 fields( 154 db.system.name = "postgresql", 155 db.namespace = "users", 156 db.operation.name = "INSERT", 157 db.collection.name = "user_profiles", 158 db.query.text = "INSERT INTO user_profiles (...) VALUES (...) RETURNING ...", 159 otelmart.user.id = user_id, 160 db.response.returned_rows = tracing::field::Empty 161 ) 162 )] 163 pub async fn create_profile( 164 pool: &PgPool, 165 user_id: i32, 166 avatar_url: Option<&str>, 167 date_of_birth: Option<NaiveDate>, 168 email_notifications: bool, 169 marketing_emails: bool, 170 ) -> Result<ProfileData, sqlx::Error> { 171 let (eid, avatar, dob, email_notif, marketing) = 172 sqlx::query_as::<_, (Uuid, Option<String>, Option<NaiveDate>, bool, bool)>( 173 r#" 174 INSERT INTO user_profiles ( 175 user_id, 176 avatar_url, 177 date_of_birth, 178 email_notifications, 179 marketing_emails 180 ) VALUES ($1, $2, $3, $4, $5) 181 RETURNING 182 uuid, 183 avatar_url, 184 date_of_birth, 185 email_notifications, 186 marketing_emails 187 "#, 188 ) 189 .bind(user_id) 190 .bind(avatar_url) 191 .bind(date_of_birth) 192 .bind(email_notifications) 193 .bind(marketing_emails) 194 .fetch_one(pool) 195 .await?; 196 197 Span::current().record("db.response.returned_rows", 1); 198 199 Ok(ProfileData { 200 eid, 201 avatar_url: avatar, 202 date_of_birth: dob, 203 email_notifications: email_notif, 204 marketing_emails: marketing, 205 }) 206 } 207 208 /// List active addresses for a user 209 #[instrument( 210 name = "SELECT user_addresses", 211 skip(pool), 212 fields( 213 otel.kind = "client", 214 db.system.name = "postgresql", 215 db.namespace = "users", 216 db.operation.name = "SELECT", 217 db.collection.name = "user_addresses", 218 db.query.text = "SELECT * FROM user_addresses WHERE user_id = $1 AND is_active = true ORDER BY ...", 219 otelmart.user.id = user_id, 220 db.response.returned_rows = tracing::field::Empty 221 ) 222 )] 223 pub async fn list_addresses(pool: &PgPool, user_id: i32) -> Result<Vec<UserAddress>, sqlx::Error> { 224 let addresses = sqlx::query_as::<_, UserAddress>( 225 r#" 226 SELECT * FROM user_addresses 227 WHERE user_id = $1 AND is_active = true 228 ORDER BY is_default DESC, created_at DESC 229 "#, 230 ) 231 .bind(user_id) 232 .fetch_all(pool) 233 .await?; 234 235 Span::current().record("db.response.returned_rows", addresses.len() as i64); 236 Ok(addresses) 237 } 238 239 /// Add a new address for a user 240 #[instrument( 241 name = "INSERT user_addresses", 242 skip(pool, req), 243 fields( 244 db.system.name = "postgresql", 245 db.namespace = "users", 246 db.operation.name = "INSERT", 247 db.collection.name = "user_addresses", 248 db.query.text = "INSERT INTO user_addresses (...) VALUES (...) RETURNING *", 249 otelmart.user.id = user_id, 250 db.response.returned_rows = tracing::field::Empty 251 ) 252 )] 253 pub async fn add_address( 254 pool: &PgPool, 255 user_id: i32, 256 req: &AddAddressRequest, 257 ) -> Result<UserAddress, sqlx::Error> { 258 let address = sqlx::query_as::<_, UserAddress>( 259 r#" 260 INSERT INTO user_addresses ( 261 user_id, address_type, address_label, 262 first_name, last_name, 263 address_line1, address_line2, 264 city, state, postal_code, country, 265 phone, is_default 266 ) VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13) 267 RETURNING * 268 "#, 269 ) 270 .bind(user_id) 271 .bind(&req.address_type) 272 .bind(&req.address_label) 273 .bind(&req.first_name) 274 .bind(&req.last_name) 275 .bind(&req.address_line1) 276 .bind(&req.address_line2) 277 .bind(&req.city) 278 .bind(&req.state) 279 .bind(&req.postal_code) 280 .bind(&req.country) 281 .bind(&req.phone) 282 .bind(req.is_default.unwrap_or(false)) 283 .fetch_one(pool) 284 .await?; 285 286 Span::current().record("db.response.returned_rows", 1); 287 Ok(address) 288 } 289 290 /// Update an existing address (only if it belongs to the user and is active) 291 #[instrument( 292 name = "UPDATE user_addresses", 293 skip(pool, req), 294 fields( 295 otel.kind = "client", 296 db.system.name = "postgresql", 297 db.namespace = "users", 298 db.operation.name = "UPDATE", 299 db.collection.name = "user_addresses", 300 db.query.text = "UPDATE user_addresses SET ... WHERE uuid = $13 AND user_id = $14 RETURNING *", 301 otelmart.address.uuid = %address_uuid, 302 otelmart.user.id = user_id, 303 db.response.returned_rows = tracing::field::Empty 304 ) 305 )] 306 pub async fn update_address( 307 pool: &PgPool, 308 user_id: i32, 309 address_uuid: Uuid, 310 req: &AddAddressRequest, 311 ) -> Result<Option<UserAddress>, sqlx::Error> { 312 let address = sqlx::query_as::<_, UserAddress>( 313 r#" 314 UPDATE user_addresses 315 SET address_type = $1, 316 address_label = $2, 317 first_name = $3, 318 last_name = $4, 319 address_line1 = $5, 320 address_line2 = $6, 321 city = $7, 322 state = $8, 323 postal_code = $9, 324 country = $10, 325 phone = $11, 326 is_default = $12, 327 updated_at = NOW() 328 WHERE uuid = $13 AND user_id = $14 AND is_active = true 329 RETURNING * 330 "#, 331 ) 332 .bind(&req.address_type) 333 .bind(&req.address_label) 334 .bind(&req.first_name) 335 .bind(&req.last_name) 336 .bind(&req.address_line1) 337 .bind(&req.address_line2) 338 .bind(&req.city) 339 .bind(&req.state) 340 .bind(&req.postal_code) 341 .bind(&req.country) 342 .bind(&req.phone) 343 .bind(req.is_default.unwrap_or(false)) 344 .bind(address_uuid) 345 .bind(user_id) 346 .fetch_optional(pool) 347 .await?; 348 349 Span::current().record( 350 "db.response.returned_rows", 351 if address.is_some() { 1 } else { 0 }, 352 ); 353 Ok(address) 354 } 355 356 /// Soft-delete an address (set is_active = false) 357 #[instrument( 358 name = "UPDATE user_addresses", 359 skip(pool), 360 fields( 361 otel.kind = "client", 362 db.system.name = "postgresql", 363 db.namespace = "users", 364 db.operation.name = "UPDATE", 365 db.collection.name = "user_addresses", 366 db.query.text = "UPDATE user_addresses SET is_active = false WHERE uuid = $1 AND user_id = $2", 367 otelmart.address.uuid = %address_uuid, 368 otelmart.user.id = user_id, 369 db.response.returned_rows = tracing::field::Empty 370 ) 371 )] 372 pub async fn delete_address( 373 pool: &PgPool, 374 user_id: i32, 375 address_uuid: Uuid, 376 ) -> Result<u64, sqlx::Error> { 377 let result = sqlx::query( 378 r#" 379 UPDATE user_addresses 380 SET is_active = false, updated_at = NOW() 381 WHERE uuid = $1 AND user_id = $2 AND is_active = true 382 "#, 383 ) 384 .bind(address_uuid) 385 .bind(user_id) 386 .execute(pool) 387 .await?; 388 389 let rows = result.rows_affected(); 390 Span::current().record("db.response.returned_rows", rows); 391 Ok(rows) 392 }