exercises

Unnamed repository; edit this file 'description' to name the repository.
Log | Files | Refs | README

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 }