repository.rs (18906B)
1 //! Order repository functions with OpenTelemetry instrumentation 2 //! 3 //! This module contains database operations for orders and shipments, 4 //! separated from HTTP handlers. 5 //! Each function is instrumented with OpenTelemetry semantic conventions for database spans. 6 7 use bigdecimal::BigDecimal; 8 use chrono::NaiveDate; 9 use sqlx::{PgPool, Postgres, QueryBuilder, Row, Transaction}; 10 use tracing::{instrument, Span}; 11 use uuid::Uuid; 12 13 use crate::models::{ 14 CreateOrderItemRequest, CreatePaymentRequest, CreateShippingAddressRequest, OrderComplete, 15 Shipment, 16 }; 17 18 /// Result of creating an order - contains the generated IDs 19 #[derive(Debug)] 20 pub struct CreatedOrder { 21 pub id: i32, 22 pub uuid: Uuid, 23 pub order_number: String, 24 } 25 26 /// Totals calculated for an order 27 #[derive(Debug, Clone)] 28 pub struct OrderTotals { 29 pub subtotal: BigDecimal, 30 pub tax_amount: BigDecimal, 31 pub shipping_amount: BigDecimal, 32 pub total: BigDecimal, 33 } 34 35 /// Generate an order number using the database function 36 #[instrument( 37 name = "SELECT generate_order_number", 38 skip(tx), 39 fields( 40 db.system.name = "postgresql", 41 db.namespace = "orders", 42 db.operation.name = "SELECT", 43 db.collection.name = "orders", 44 db.query.text = "SELECT generate_order_number()" 45 ) 46 )] 47 pub async fn generate_order_number( 48 tx: &mut Transaction<'_, Postgres>, 49 ) -> Result<String, sqlx::Error> { 50 sqlx::query_scalar("SELECT generate_order_number()") 51 .fetch_one(&mut **tx) 52 .await 53 } 54 55 /// Create a new order record 56 /// 57 /// Inserts the order header with customer info and totals. 58 /// Returns the generated order ID, UUID, and order number. 59 #[instrument( 60 name = "INSERT orders", 61 skip(tx, totals), 62 fields( 63 db.system.name = "postgresql", 64 db.namespace = "orders", 65 db.operation.name = "INSERT", 66 db.collection.name = "orders", 67 db.query.text = "INSERT INTO orders (...) VALUES (...) RETURNING id, uuid", 68 db.response.returned_rows = tracing::field::Empty 69 ) 70 )] 71 pub async fn create_order( 72 tx: &mut Transaction<'_, Postgres>, 73 order_number: &str, 74 customer_email: &str, 75 customer_phone: Option<&str>, 76 totals: &OrderTotals, 77 ) -> Result<CreatedOrder, sqlx::Error> { 78 let row = sqlx::query( 79 r#" 80 INSERT INTO orders ( 81 order_number, customer_email, customer_phone, 82 subtotal, tax_amount, shipping_amount, total, 83 status, payment_status 84 ) 85 VALUES ($1, $2, $3, $4, $5, $6, $7, 'pending', 'pending') 86 RETURNING id, uuid 87 "#, 88 ) 89 .bind(order_number) 90 .bind(customer_email) 91 .bind(customer_phone) 92 .bind(&totals.subtotal) 93 .bind(&totals.tax_amount) 94 .bind(&totals.shipping_amount) 95 .bind(&totals.total) 96 .fetch_one(&mut **tx) 97 .await?; 98 99 Span::current().record("db.response.returned_rows", 1); 100 101 Ok(CreatedOrder { 102 id: row.get("id"), 103 uuid: row.get("uuid"), 104 order_number: order_number.to_string(), 105 }) 106 } 107 108 /// Create order items (line items) for an order 109 /// 110 /// Inserts all order items in sequence. Each item contains a snapshot 111 /// of the product data at the time of order. 112 #[instrument( 113 name = "INSERT order_items", 114 skip(tx, items), 115 fields( 116 db.system.name = "postgresql", 117 db.namespace = "orders", 118 db.operation.name = "INSERT", 119 db.collection.name = "order_items", 120 db.query.text = "INSERT INTO order_items (...) VALUES (...)", 121 otelmart.order.id = order_id, 122 otelmart.items.count = items.len(), 123 db.response.returned_rows = tracing::field::Empty 124 ) 125 )] 126 pub async fn create_order_items( 127 tx: &mut Transaction<'_, Postgres>, 128 order_id: i32, 129 items: &[CreateOrderItemRequest], 130 ) -> Result<(), sqlx::Error> { 131 let mut rows_inserted = 0; 132 133 for item in items { 134 let total_price = &item.unit_price * BigDecimal::from(item.quantity); 135 sqlx::query( 136 r#" 137 INSERT INTO order_items ( 138 order_id, product_uuid, product_name, product_sku, 139 quantity, unit_price, total_price 140 ) 141 VALUES ($1, $2, $3, $4, $5, $6, $7) 142 "#, 143 ) 144 .bind(order_id) 145 .bind(item.product_uuid) 146 .bind(&item.product_name) 147 .bind(&item.product_sku) 148 .bind(item.quantity) 149 .bind(&item.unit_price) 150 .bind(&total_price) 151 .execute(&mut **tx) 152 .await?; 153 154 rows_inserted += 1; 155 } 156 157 Span::current().record("db.response.returned_rows", rows_inserted); 158 Ok(()) 159 } 160 161 /// Create shipping address for an order 162 #[instrument( 163 name = "INSERT shipping_addresses", 164 skip(tx, address), 165 fields( 166 db.system.name = "postgresql", 167 db.namespace = "orders", 168 db.operation.name = "INSERT", 169 db.collection.name = "shipping_addresses", 170 db.query.text = "INSERT INTO shipping_addresses (...) VALUES (...)", 171 otelmart.order.id = order_id, 172 db.response.returned_rows = tracing::field::Empty 173 ) 174 )] 175 pub async fn create_shipping_address( 176 tx: &mut Transaction<'_, Postgres>, 177 order_id: i32, 178 address: &CreateShippingAddressRequest, 179 ) -> Result<(), sqlx::Error> { 180 sqlx::query( 181 r#" 182 INSERT INTO shipping_addresses ( 183 order_id, first_name, last_name, address_line1, address_line2, 184 city, state, postal_code, country, phone 185 ) 186 VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10) 187 "#, 188 ) 189 .bind(order_id) 190 .bind(&address.first_name) 191 .bind(&address.last_name) 192 .bind(&address.address_line1) 193 .bind(&address.address_line2) 194 .bind(&address.city) 195 .bind(&address.state) 196 .bind(&address.postal_code) 197 .bind(&address.country) 198 .bind(&address.phone) 199 .execute(&mut **tx) 200 .await?; 201 202 Span::current().record("db.response.returned_rows", 1); 203 Ok(()) 204 } 205 206 /// Generate a payment reference using the database function 207 #[instrument( 208 name = "SELECT generate_payment_reference", 209 skip(tx), 210 fields( 211 otel.kind = "client", 212 db.system.name = "postgresql", 213 db.namespace = "orders", 214 db.operation.name = "SELECT", 215 db.collection.name = "payments" 216 ) 217 )] 218 pub async fn generate_payment_reference( 219 tx: &mut Transaction<'_, Postgres>, 220 ) -> Result<String, sqlx::Error> { 221 sqlx::query_scalar("SELECT generate_payment_reference()") 222 .fetch_one(&mut **tx) 223 .await 224 } 225 226 /// Create payment record for an order 227 /// 228 /// NOTE: Payment is simulated - in production, this would call a payment gateway 229 /// and only insert the record after successful payment processing. 230 #[instrument( 231 name = "INSERT payments", 232 skip(tx, payment), 233 fields( 234 db.system.name = "postgresql", 235 db.namespace = "orders", 236 db.operation.name = "INSERT", 237 db.collection.name = "payments", 238 db.query.text = "INSERT INTO payments (...) VALUES (...)", 239 otelmart.order.id = order_id, 240 db.response.returned_rows = tracing::field::Empty 241 ) 242 )] 243 pub async fn create_payment( 244 tx: &mut Transaction<'_, Postgres>, 245 order_id: i32, 246 payment: &CreatePaymentRequest, 247 payment_reference: &str, 248 amount: &BigDecimal, 249 ) -> Result<(), sqlx::Error> { 250 sqlx::query( 251 r#" 252 INSERT INTO payments ( 253 order_id, payment_method, amount, status, 254 payment_reference, card_last4, card_brand, processed_at 255 ) 256 VALUES ($1, $2, $3, 'paid', $4, $5, $6, CURRENT_TIMESTAMP) 257 "#, 258 ) 259 .bind(order_id) 260 .bind(&payment.payment_method) 261 .bind(amount) 262 .bind(payment_reference) 263 .bind(&payment.card_last4) 264 .bind(&payment.card_brand) 265 .execute(&mut **tx) 266 .await?; 267 268 Span::current().record("db.response.returned_rows", 1); 269 Ok(()) 270 } 271 272 /// Update order status after successful payment 273 #[instrument( 274 name = "UPDATE orders", 275 skip(tx), 276 fields( 277 otel.kind = "client", 278 db.system.name = "postgresql", 279 db.namespace = "orders", 280 db.operation.name = "UPDATE", 281 db.collection.name = "orders", 282 db.query.text = "UPDATE orders SET payment_status = $1, status = $2 WHERE id = $3", 283 otelmart.order.id = order_id, 284 db.response.returned_rows = tracing::field::Empty 285 ) 286 )] 287 pub async fn update_order_payment_status( 288 tx: &mut Transaction<'_, Postgres>, 289 order_id: i32, 290 payment_status: &str, 291 order_status: &str, 292 ) -> Result<(), sqlx::Error> { 293 let result = sqlx::query( 294 r#" 295 UPDATE orders 296 SET payment_status = $1, status = $2 297 WHERE id = $3 298 "#, 299 ) 300 .bind(payment_status) 301 .bind(order_status) 302 .bind(order_id) 303 .execute(&mut **tx) 304 .await?; 305 306 Span::current().record("db.response.returned_rows", result.rows_affected()); 307 Ok(()) 308 } 309 310 /// Get order by UUID 311 #[instrument( 312 name = "SELECT orders", 313 skip(pool), 314 fields( 315 otel.kind = "client", 316 db.system.name = "postgresql", 317 db.namespace = "orders", 318 db.operation.name = "SELECT", 319 db.collection.name = "orders", 320 db.query.text = "SELECT * FROM v_orders_complete WHERE eid = $1", 321 otelmart.order.uuid = %uuid 322 ) 323 )] 324 pub async fn get_order_by_uuid( 325 pool: &PgPool, 326 uuid: Uuid, 327 ) -> Result<Option<OrderComplete>, sqlx::Error> { 328 sqlx::query_as::<_, OrderComplete>(r#"SELECT * FROM v_orders_complete WHERE eid = $1"#) 329 .bind(uuid) 330 .fetch_optional(pool) 331 .await 332 } 333 334 /// Get order by email and order number (for guest users) 335 #[instrument( 336 name = "SELECT orders", 337 skip(pool), 338 fields( 339 otel.kind = "client", 340 db.system.name = "postgresql", 341 db.namespace = "orders", 342 db.operation.name = "SELECT", 343 db.collection.name = "orders", 344 db.query.text = "SELECT * FROM v_orders_complete WHERE customer_email = $1 AND order_number = $2" 345 ) 346 )] 347 pub async fn get_order_by_email_and_number( 348 pool: &PgPool, 349 email: &str, 350 order_number: &str, 351 ) -> Result<Option<OrderComplete>, sqlx::Error> { 352 sqlx::query_as::<_, OrderComplete>( 353 r#" 354 SELECT * FROM v_orders_complete 355 WHERE customer_email = $1 AND order_number = $2 356 "#, 357 ) 358 .bind(email) 359 .bind(order_number) 360 .fetch_optional(pool) 361 .await 362 } 363 364 /// List orders for a user with pagination 365 #[instrument( 366 name = "SELECT orders", 367 skip(pool), 368 fields( 369 otel.kind = "client", 370 db.system.name = "postgresql", 371 db.namespace = "orders", 372 db.operation.name = "SELECT", 373 db.collection.name = "orders", 374 db.query.text = "SELECT * FROM v_orders_complete WHERE customer_email = $1 ...", 375 otelmart.page = page, 376 otelmart.page_size = page_size 377 ) 378 )] 379 pub async fn list_orders_by_email( 380 pool: &PgPool, 381 email: &str, 382 status_filter: Option<&str>, 383 payment_status_filter: Option<&str>, 384 page: i32, 385 page_size: i32, 386 ) -> Result<(Vec<OrderComplete>, i64), sqlx::Error> { 387 let offset = (page - 1) * page_size; 388 389 // Build count query with proper parameterized filters 390 let mut count_builder: QueryBuilder<Postgres> = 391 QueryBuilder::new("SELECT COUNT(*) FROM v_orders_complete WHERE customer_email = "); 392 count_builder.push_bind(email); 393 394 if let Some(status) = status_filter { 395 count_builder.push(" AND status = "); 396 count_builder.push_bind(status); 397 } 398 399 if let Some(payment_status) = payment_status_filter { 400 count_builder.push(" AND payment_status = "); 401 count_builder.push_bind(payment_status); 402 } 403 404 // Execute count query 405 let total_count: (i64,) = count_builder.build_query_as().fetch_one(pool).await?; 406 407 // Build main query with proper parameterized filters 408 let mut query_builder: QueryBuilder<Postgres> = 409 QueryBuilder::new("SELECT * FROM v_orders_complete WHERE customer_email = "); 410 query_builder.push_bind(email); 411 412 if let Some(status) = status_filter { 413 query_builder.push(" AND status = "); 414 query_builder.push_bind(status); 415 } 416 417 if let Some(payment_status) = payment_status_filter { 418 query_builder.push(" AND payment_status = "); 419 query_builder.push_bind(payment_status); 420 } 421 422 query_builder.push(" ORDER BY ordered_at DESC LIMIT "); 423 query_builder.push_bind(page_size); 424 query_builder.push(" OFFSET "); 425 query_builder.push_bind(offset); 426 427 // Execute main query 428 let orders: Vec<OrderComplete> = query_builder.build_query_as().fetch_all(pool).await?; 429 430 Ok((orders, total_count.0)) 431 } 432 433 /// Get the internal order ID from a public UUID 434 #[instrument( 435 name = "SELECT orders", 436 skip(pool), 437 fields( 438 otel.kind = "client", 439 db.system.name = "postgresql", 440 db.namespace = "orders", 441 db.operation.name = "SELECT", 442 db.collection.name = "orders", 443 db.query.text = "SELECT id FROM orders WHERE uuid = $1", 444 otelmart.order.uuid = %order_uuid 445 ) 446 )] 447 pub async fn get_order_id_by_uuid( 448 pool: &PgPool, 449 order_uuid: Uuid, 450 ) -> Result<Option<i32>, sqlx::Error> { 451 sqlx::query_scalar("SELECT id FROM orders WHERE uuid = $1") 452 .bind(order_uuid) 453 .fetch_optional(pool) 454 .await 455 } 456 457 /// Create a new shipment record for an order 458 /// 459 /// Inserts a shipment with status 'shipped' and sets shipped_at to now. 460 #[instrument( 461 name = "INSERT shipments", 462 skip(pool), 463 fields( 464 db.system.name = "postgresql", 465 db.namespace = "orders", 466 db.operation.name = "INSERT", 467 db.collection.name = "shipments", 468 db.query.text = "INSERT INTO shipments (...) VALUES (...) RETURNING *", 469 otelmart.order.id = order_id, 470 db.response.returned_rows = tracing::field::Empty 471 ) 472 )] 473 pub async fn create_shipment( 474 pool: &PgPool, 475 order_id: i32, 476 carrier: &str, 477 tracking_number: Option<&str>, 478 estimated_delivery_date: Option<NaiveDate>, 479 ) -> Result<Shipment, sqlx::Error> { 480 let shipment = sqlx::query_as::<_, Shipment>( 481 r#" 482 INSERT INTO shipments ( 483 order_id, carrier, tracking_number, 484 estimated_delivery_date, status, shipped_at 485 ) 486 VALUES ($1, $2, $3, $4, 'shipped', CURRENT_TIMESTAMP) 487 RETURNING * 488 "#, 489 ) 490 .bind(order_id) 491 .bind(carrier) 492 .bind(tracking_number) 493 .bind(estimated_delivery_date) 494 .fetch_one(pool) 495 .await?; 496 497 Span::current().record("db.response.returned_rows", 1); 498 Ok(shipment) 499 } 500 501 /// Update order status (e.g., to 'shipped' or 'delivered') 502 #[instrument( 503 name = "UPDATE orders", 504 skip(pool), 505 fields( 506 otel.kind = "client", 507 db.system.name = "postgresql", 508 db.namespace = "orders", 509 db.operation.name = "UPDATE", 510 db.collection.name = "orders", 511 db.query.text = "UPDATE orders SET status = $1 WHERE id = $2", 512 otelmart.order.id = order_id, 513 otelmart.order.status = status, 514 db.response.returned_rows = tracing::field::Empty 515 ) 516 )] 517 pub async fn update_order_status( 518 pool: &PgPool, 519 order_id: i32, 520 status: &str, 521 ) -> Result<u64, sqlx::Error> { 522 let result = sqlx::query("UPDATE orders SET status = $1 WHERE id = $2") 523 .bind(status) 524 .bind(order_id) 525 .execute(pool) 526 .await?; 527 528 let rows = result.rows_affected(); 529 Span::current().record("db.response.returned_rows", rows); 530 Ok(rows) 531 } 532 533 /// Update shipment status within a transaction 534 /// 535 /// For 'delivered' status, also sets actual_delivery_date and delivered_at. 536 /// For other statuses, only updates the status field. 537 /// Returns true if a shipment was found and updated. 538 #[instrument( 539 name = "UPDATE shipments", 540 skip(tx), 541 fields( 542 otel.kind = "client", 543 db.system.name = "postgresql", 544 db.namespace = "orders", 545 db.operation.name = "UPDATE", 546 db.collection.name = "shipments", 547 db.query.text = "UPDATE shipments SET status = $1 ... WHERE order_id = ...", 548 otelmart.order.id = order_id, 549 otelmart.shipment.status = status, 550 db.response.returned_rows = tracing::field::Empty 551 ) 552 )] 553 pub async fn update_shipment_status( 554 tx: &mut Transaction<'_, Postgres>, 555 order_id: i32, 556 status: &str, 557 actual_delivery_date: Option<NaiveDate>, 558 ) -> Result<bool, sqlx::Error> { 559 let result = if status == "delivered" { 560 sqlx::query( 561 r#" 562 UPDATE shipments 563 SET status = $1, actual_delivery_date = $2, delivered_at = CURRENT_TIMESTAMP 564 WHERE order_id = $3 565 RETURNING id 566 "#, 567 ) 568 .bind(status) 569 .bind(actual_delivery_date) 570 .bind(order_id) 571 .fetch_optional(&mut **tx) 572 .await? 573 } else { 574 sqlx::query( 575 r#" 576 UPDATE shipments 577 SET status = $1 578 WHERE order_id = $2 579 RETURNING id 580 "#, 581 ) 582 .bind(status) 583 .bind(order_id) 584 .fetch_optional(&mut **tx) 585 .await? 586 }; 587 588 let found = result.is_some(); 589 Span::current().record("db.response.returned_rows", if found { 1 } else { 0 }); 590 Ok(found) 591 } 592 593 /// Update order status within a transaction 594 #[instrument( 595 name = "UPDATE orders", 596 skip(tx), 597 fields( 598 otel.kind = "client", 599 db.system.name = "postgresql", 600 db.namespace = "orders", 601 db.operation.name = "UPDATE", 602 db.collection.name = "orders", 603 db.query.text = "UPDATE orders SET status = $1 WHERE id = $2", 604 otelmart.order.id = order_id, 605 otelmart.order.status = status, 606 db.response.returned_rows = tracing::field::Empty 607 ) 608 )] 609 pub async fn update_order_status_in_tx( 610 tx: &mut Transaction<'_, Postgres>, 611 order_id: i32, 612 status: &str, 613 ) -> Result<u64, sqlx::Error> { 614 let result = sqlx::query("UPDATE orders SET status = $1 WHERE id = $2") 615 .bind(status) 616 .bind(order_id) 617 .execute(&mut **tx) 618 .await?; 619 620 let rows = result.rows_affected(); 621 Span::current().record("db.response.returned_rows", rows); 622 Ok(rows) 623 } 624 625 /// Get order ID from UUID within a transaction 626 #[instrument( 627 name = "SELECT orders", 628 skip(tx), 629 fields( 630 otel.kind = "client", 631 db.system.name = "postgresql", 632 db.namespace = "orders", 633 db.operation.name = "SELECT", 634 db.collection.name = "orders", 635 db.query.text = "SELECT id FROM orders WHERE uuid = $1", 636 otelmart.order.uuid = %order_uuid 637 ) 638 )] 639 pub async fn get_order_id_by_uuid_in_tx( 640 tx: &mut Transaction<'_, Postgres>, 641 order_uuid: Uuid, 642 ) -> Result<Option<i32>, sqlx::Error> { 643 sqlx::query_scalar("SELECT id FROM orders WHERE uuid = $1") 644 .bind(order_uuid) 645 .fetch_optional(&mut **tx) 646 .await 647 }