commit b7827798bf35fc0f34dd7fd8ef67546c34427a1c
parent 9eb879e9f54d9bd02849baf1182f81ef82581dcc
Author: ling0x <ling0x@users.noreply.github.com>
Date: Thu, 27 Aug 2026 17:44:18 +0100
use repository pattern to separate sql logics and add logging to SQLx
Diffstat:
19 files changed, 2404 insertions(+), 1444 deletions(-)
diff --git a/distributed_observability/Cargo.toml b/distributed_observability/Cargo.toml
@@ -42,6 +42,7 @@ sqlx = { version = "0.8.6", features = [
"json",
"migrate",
"bigdecimal",
+ "migrate",
] }
# Redis
diff --git a/distributed_observability/inventory/src/db/mod.rs b/distributed_observability/inventory/src/db/mod.rs
@@ -22,6 +22,12 @@ use sqlx::{
};
use std::str::FromStr;
+pub mod pricing;
+pub mod repository;
+
+pub use pricing::*;
+pub use repository::*;
+
/// Database connection pool wrapper
///
/// Wraps sqlx's PgPool with custom configuration for the inventory service.
diff --git a/distributed_observability/inventory/src/db/pricing.rs b/distributed_observability/inventory/src/db/pricing.rs
@@ -0,0 +1,187 @@
+//! Pricing repository functions with OpenTelemetry instrumentation
+//!
+//! This module contains database operations for product pricing,
+//! separated from HTTP handlers.
+
+use bigdecimal::BigDecimal;
+use sqlx::{PgPool, Postgres, QueryBuilder};
+use tracing::instrument;
+use uuid::Uuid;
+
+use crate::models::{PricingQueryParams, ProductPricing};
+
+/// List pricing with pagination and filtering
+///
+/// Queries active pricing records with optional filters for
+/// product UUID, price range, and discount presence.
+#[instrument(
+ name = "SELECT product_pricing",
+ skip(pool, params),
+ fields(
+ otel.kind = "client",
+ db.system.name = "postgresql",
+ db.namespace = "inventory",
+ db.operation.name = "SELECT",
+ db.collection.name = "product_pricing",
+ db.query.text = "SELECT * FROM product_pricing WHERE is_active = true ... LIMIT ... OFFSET ...",
+ otelmart.page = page,
+ otelmart.page_size = page_size
+ )
+)]
+pub async fn list_pricing(
+ pool: &PgPool,
+ params: &PricingQueryParams,
+ page: i32,
+ page_size: i32,
+ offset: i32,
+) -> Result<(Vec<ProductPricing>, i64), sqlx::Error> {
+ // Build count query
+ let mut count_builder: QueryBuilder<Postgres> =
+ QueryBuilder::new("SELECT COUNT(*) FROM product_pricing WHERE is_active = true");
+
+ // Apply filters (all use AND since WHERE is_active is already present)
+ apply_pricing_filters(&mut count_builder, params);
+
+ let total_count: i64 = count_builder.build_query_scalar().fetch_one(pool).await?;
+
+ // Build main query with same filters
+ let mut query_builder: QueryBuilder<Postgres> =
+ QueryBuilder::new("SELECT * FROM product_pricing WHERE is_active = true");
+
+ apply_pricing_filters(&mut query_builder, params);
+
+ // Add ordering and pagination
+ query_builder.push(" ORDER BY updated_at DESC LIMIT ");
+ query_builder.push_bind(page_size);
+ query_builder.push(" OFFSET ");
+ query_builder.push_bind(offset);
+
+ let pricing = query_builder.build_query_as().fetch_all(pool).await?;
+
+ Ok((pricing, total_count))
+}
+
+/// Get pricing for a specific product by UUID
+#[instrument(
+ name = "SELECT product_pricing",
+ skip(pool),
+ fields(
+ otel.kind = "client",
+ db.system.name = "postgresql",
+ db.namespace = "inventory",
+ db.operation.name = "SELECT",
+ db.collection.name = "product_pricing",
+ db.query.text = "SELECT * FROM product_pricing WHERE product_uuid = $1 AND is_active = true",
+ otelmart.product.uuid = %product_uuid
+ )
+)]
+pub async fn get_pricing_by_product(
+ pool: &PgPool,
+ product_uuid: Uuid,
+) -> Result<Option<ProductPricing>, sqlx::Error> {
+ sqlx::query_as::<_, ProductPricing>(
+ r#"
+ SELECT * FROM product_pricing
+ WHERE product_uuid = $1 AND is_active = true
+ "#,
+ )
+ .bind(product_uuid)
+ .fetch_optional(pool)
+ .await
+}
+
+/// Deactivate existing pricing and insert new pricing for a product
+///
+/// First deactivates any active pricing record, then inserts a new one.
+#[instrument(
+ name = "UPSERT product_pricing",
+ skip(pool),
+ fields(
+ db.system.name = "postgresql",
+ db.namespace = "inventory",
+ db.operation.name = "INSERT",
+ db.collection.name = "product_pricing",
+ db.query.text = "INSERT INTO product_pricing (...) VALUES (...) RETURNING *",
+ otelmart.product.uuid = %product_uuid
+ )
+)]
+pub async fn upsert_pricing(
+ pool: &PgPool,
+ product_uuid: Uuid,
+ final_price: BigDecimal,
+ initial_price: Option<BigDecimal>,
+ currency: String,
+ price_valid_from: Option<chrono::DateTime<chrono::Utc>>,
+ price_valid_until: Option<chrono::DateTime<chrono::Utc>>,
+) -> Result<ProductPricing, sqlx::Error> {
+ // Deactivate existing active pricing
+ if let Err(e) = sqlx::query(
+ r#"
+ UPDATE product_pricing
+ SET is_active = false
+ WHERE product_uuid = $1 AND is_active = true
+ "#,
+ )
+ .bind(product_uuid)
+ .execute(pool)
+ .await
+ {
+ tracing::error!(error = %e, "Error deactivating old pricing");
+ }
+
+ // Insert new pricing
+ sqlx::query_as::<_, ProductPricing>(
+ r#"
+ INSERT INTO product_pricing (
+ product_uuid,
+ final_price,
+ initial_price,
+ currency,
+ price_valid_from,
+ price_valid_until,
+ is_active
+ )
+ VALUES ($1, $2, $3, $4, $5, $6, true)
+ RETURNING *
+ "#,
+ )
+ .bind(product_uuid)
+ .bind(final_price)
+ .bind(initial_price)
+ .bind(currency)
+ .bind(price_valid_from)
+ .bind(price_valid_until)
+ .fetch_one(pool)
+ .await
+}
+
+/// Apply pricing filters to a query builder
+///
+/// Shared filter logic for count and main queries.
+fn apply_pricing_filters<'a>(
+ query: &mut QueryBuilder<'a, Postgres>,
+ params: &'a PricingQueryParams,
+) {
+ if let Some(product_uuid) = params.product_uuid {
+ query.push(" AND product_uuid = ");
+ query.push_bind(product_uuid);
+ }
+
+ if let Some(ref min_price) = params.min_price {
+ query.push(" AND final_price >= ");
+ query.push_bind(min_price);
+ }
+
+ if let Some(ref max_price) = params.max_price {
+ query.push(" AND final_price <= ");
+ query.push_bind(max_price);
+ }
+
+ if let Some(has_discount) = params.has_discount {
+ if has_discount {
+ query.push(" AND discount_percentage > 0");
+ } else {
+ query.push(" AND (discount_percentage IS NULL OR discount_percentage = 0)");
+ }
+ }
+}
diff --git a/distributed_observability/inventory/src/db/repository.rs b/distributed_observability/inventory/src/db/repository.rs
@@ -0,0 +1,284 @@
+//! Inventory repository functions with OpenTelemetry instrumentation
+//!
+//! This module contains database operations for inventory/stock management,
+//! separated from HTTP handlers. Each function is instrumented with
+//! OpenTelemetry semantic conventions for database spans.
+
+use sqlx::{PgPool, Postgres, QueryBuilder, Row};
+use tracing::{instrument, Span};
+use uuid::Uuid;
+
+use crate::models::{InventoryQueryParams, InventoryWithPricing};
+
+/// List inventory with pagination and filtering
+///
+/// Queries the v_product_inventory_pricing view which combines
+/// inventory and pricing data. Supports filtering by stock status,
+/// product UUID, and stock quantity range.
+#[instrument(
+ name = "SELECT inventory",
+ skip(pool, params),
+ fields(
+ otel.kind = "client",
+ db.system.name = "postgresql",
+ db.namespace = "inventory",
+ db.operation.name = "SELECT",
+ db.collection.name = "product_inventory",
+ db.query.text = "SELECT * FROM v_product_inventory_pricing WHERE ... LIMIT ... OFFSET ...",
+ otelmart.page = page,
+ otelmart.page_size = page_size,
+ otelmart.stock_status = ?params.stock_status
+ )
+)]
+pub async fn list_inventory(
+ pool: &PgPool,
+ params: &InventoryQueryParams,
+ page: i32,
+ page_size: i32,
+ offset: i32,
+) -> Result<(Vec<InventoryWithPricing>, i64), sqlx::Error> {
+ // Build count query with QueryBuilder for safe parameter binding
+ let mut count_builder: QueryBuilder<Postgres> =
+ QueryBuilder::new("SELECT COUNT(*) FROM v_product_inventory_pricing");
+
+ let mut has_filters = false;
+
+ // Apply filters to count query
+ if let Some(ref status) = params.stock_status {
+ count_builder.push(if has_filters { " AND " } else { " WHERE " });
+ has_filters = true;
+ count_builder.push("stock_status = ");
+ count_builder.push_bind(status);
+ }
+
+ if let Some(product_uuid) = params.product_uuid {
+ count_builder.push(if has_filters { " AND " } else { " WHERE " });
+ has_filters = true;
+ count_builder.push("product_uuid = ");
+ count_builder.push_bind(product_uuid);
+ }
+
+ if let Some(min_stock) = params.min_stock {
+ count_builder.push(if has_filters { " AND " } else { " WHERE " });
+ has_filters = true;
+ count_builder.push("available_quantity >= ");
+ count_builder.push_bind(min_stock);
+ }
+
+ if let Some(max_stock) = params.max_stock {
+ count_builder.push(if has_filters { " AND " } else { " WHERE " });
+ count_builder.push("available_quantity <= ");
+ count_builder.push_bind(max_stock);
+ }
+
+ // Execute count query
+ let total_count: i64 = count_builder.build_query_scalar().fetch_one(pool).await?;
+
+ // Build main query with same filters
+ let mut query_builder: QueryBuilder<Postgres> =
+ QueryBuilder::new("SELECT * FROM v_product_inventory_pricing");
+
+ let mut has_filters = false;
+
+ if let Some(ref status) = params.stock_status {
+ query_builder.push(if has_filters { " AND " } else { " WHERE " });
+ has_filters = true;
+ query_builder.push("stock_status = ");
+ query_builder.push_bind(status);
+ }
+
+ if let Some(product_uuid) = params.product_uuid {
+ query_builder.push(if has_filters { " AND " } else { " WHERE " });
+ has_filters = true;
+ query_builder.push("product_uuid = ");
+ query_builder.push_bind(product_uuid);
+ }
+
+ if let Some(min_stock) = params.min_stock {
+ query_builder.push(if has_filters { " AND " } else { " WHERE " });
+ has_filters = true;
+ query_builder.push("available_quantity >= ");
+ query_builder.push_bind(min_stock);
+ }
+
+ if let Some(max_stock) = params.max_stock {
+ query_builder.push(if has_filters { " AND " } else { " WHERE " });
+ query_builder.push("available_quantity <= ");
+ query_builder.push_bind(max_stock);
+ }
+
+ // Add ordering and pagination
+ query_builder.push(" ORDER BY stock_status DESC, available_quantity ASC LIMIT ");
+ query_builder.push_bind(page_size);
+ query_builder.push(" OFFSET ");
+ query_builder.push_bind(offset);
+
+ // Execute main query
+ let inventory = query_builder.build_query_as().fetch_all(pool).await?;
+
+ Ok((inventory, total_count))
+}
+
+/// Get inventory for a specific product by UUID
+#[instrument(
+ name = "SELECT inventory",
+ skip(pool),
+ fields(
+ otel.kind = "client",
+ db.system.name = "postgresql",
+ db.namespace = "inventory",
+ db.operation.name = "SELECT",
+ db.collection.name = "product_inventory",
+ db.query.text = "SELECT * FROM v_product_inventory_pricing WHERE product_uuid = $1",
+ otelmart.product.uuid = %product_uuid
+ )
+)]
+pub async fn get_inventory_by_product(
+ pool: &PgPool,
+ product_uuid: Uuid,
+) -> Result<Option<InventoryWithPricing>, sqlx::Error> {
+ sqlx::query_as::<_, InventoryWithPricing>(
+ r#"
+ SELECT * FROM v_product_inventory_pricing
+ WHERE product_uuid = $1
+ "#,
+ )
+ .bind(product_uuid)
+ .fetch_optional(pool)
+ .await
+}
+
+/// Update stock quantity for a product
+///
+/// Returns the new available_quantity if the product was found, None otherwise.
+#[instrument(
+ name = "UPDATE inventory",
+ skip(pool),
+ fields(
+ otel.kind = "client",
+ db.system.name = "postgresql",
+ db.namespace = "inventory",
+ db.operation.name = "UPDATE",
+ db.collection.name = "product_inventory",
+ db.query.text = "UPDATE product_inventory SET stock_quantity = $1 ... WHERE product_uuid = $4",
+ otelmart.product.uuid = %product_uuid,
+ otelmart.quantity = quantity,
+ db.response.returned_rows = tracing::field::Empty
+ )
+)]
+pub async fn update_stock(
+ pool: &PgPool,
+ product_uuid: Uuid,
+ quantity: i32,
+ reorder_level: Option<i32>,
+ reorder_quantity: Option<i32>,
+) -> Result<Option<i32>, sqlx::Error> {
+ let result = sqlx::query(
+ r#"
+ UPDATE product_inventory
+ SET stock_quantity = $1,
+ reorder_level = COALESCE($2, reorder_level),
+ reorder_quantity = COALESCE($3, reorder_quantity),
+ last_restocked_at = CURRENT_TIMESTAMP
+ WHERE product_uuid = $4
+ RETURNING available_quantity
+ "#,
+ )
+ .bind(quantity)
+ .bind(reorder_level)
+ .bind(reorder_quantity)
+ .bind(product_uuid)
+ .fetch_optional(pool)
+ .await?;
+
+ let rows = if result.is_some() { 1 } else { 0 };
+ Span::current().record("db.response.returned_rows", rows);
+
+ Ok(result.map(|row| row.get("available_quantity")))
+}
+
+/// Reserve stock for an order using the database function
+///
+/// Calls the reserve_stock PostgreSQL function which atomically checks
+/// available stock and creates a reservation. Returns true if successful.
+#[instrument(
+ name = "SELECT reserve_stock",
+ skip(pool),
+ fields(
+ db.system.name = "postgresql",
+ db.namespace = "inventory",
+ db.operation.name = "SELECT",
+ db.collection.name = "product_inventory",
+ db.query.text = "SELECT reserve_stock($1, $2)",
+ otelmart.product.uuid = %product_uuid,
+ otelmart.quantity = quantity
+ )
+)]
+pub async fn reserve_stock(
+ pool: &PgPool,
+ product_uuid: Uuid,
+ quantity: i32,
+) -> Result<bool, sqlx::Error> {
+ sqlx::query_scalar::<_, bool>(r#"SELECT reserve_stock($1, $2)"#)
+ .bind(product_uuid)
+ .bind(quantity)
+ .fetch_one(pool)
+ .await
+}
+
+/// Release reserved stock using the database function
+#[instrument(
+ name = "SELECT release_stock",
+ skip(pool),
+ fields(
+ db.system.name = "postgresql",
+ db.namespace = "inventory",
+ db.operation.name = "SELECT",
+ db.collection.name = "product_inventory",
+ db.query.text = "SELECT release_stock($1, $2)",
+ otelmart.product.uuid = %product_uuid,
+ otelmart.quantity = quantity
+ )
+)]
+pub async fn release_stock(
+ pool: &PgPool,
+ product_uuid: Uuid,
+ quantity: i32,
+) -> Result<(), sqlx::Error> {
+ sqlx::query(r#"SELECT release_stock($1, $2)"#)
+ .bind(product_uuid)
+ .bind(quantity)
+ .execute(pool)
+ .await?;
+ Ok(())
+}
+
+/// Confirm a sale and decrease stock using the database function
+#[instrument(
+ name = "SELECT confirm_stock_sale",
+ skip(pool),
+ fields(
+ db.system.name = "postgresql",
+ db.namespace = "inventory",
+ db.operation.name = "SELECT",
+ db.collection.name = "product_inventory",
+ db.query.text = "SELECT confirm_stock_sale($1, $2, $3)",
+ otelmart.product.uuid = %product_uuid,
+ otelmart.order.uuid = %order_uuid,
+ otelmart.quantity = quantity
+ )
+)]
+pub async fn confirm_sale(
+ pool: &PgPool,
+ product_uuid: Uuid,
+ quantity: i32,
+ order_uuid: Uuid,
+) -> Result<(), sqlx::Error> {
+ sqlx::query(r#"SELECT confirm_stock_sale($1, $2, $3)"#)
+ .bind(product_uuid)
+ .bind(quantity)
+ .bind(order_uuid)
+ .execute(pool)
+ .await?;
+ Ok(())
+}
diff --git a/distributed_observability/inventory/src/handlers/inventory.rs b/distributed_observability/inventory/src/handlers/inventory.rs
@@ -1,19 +1,20 @@
//! Inventory/stock management API handlers
//!
//! This module contains the HTTP request handlers for inventory-related endpoints.
+//! Database operations are delegated to the repository layer in `db::repository`.
use axum::{
extract::{Path, Query, State},
http::StatusCode,
response::{IntoResponse, Json},
};
-use sqlx::{PgPool, Postgres, QueryBuilder, Row};
-use tracing::instrument;
+use sqlx::PgPool;
use uuid::Uuid;
+use crate::db;
use crate::models::{
- ConfirmSaleRequest, InventoryQueryParams, InventoryResponse, InventoryWithPricing,
- ReleaseStockRequest, ReserveStockRequest, StockOperationResponse, UpdateStockRequest,
+ ConfirmSaleRequest, InventoryQueryParams, InventoryResponse, ReleaseStockRequest,
+ ReserveStockRequest, StockOperationResponse, UpdateStockRequest,
};
use crate::utils::{calculate_pagination, calculate_total_pages, internal_error, not_found_error};
@@ -28,15 +29,6 @@ use crate::utils::{calculate_pagination, calculate_total_pages, internal_error,
/// - `stock_status` - Filter by status ('in_stock', 'low_stock', 'out_of_stock')
/// - `product_uuid` - Filter by specific product
/// - `min_stock` / `max_stock` - Filter by available quantity range
-#[instrument(
- name = "list_inventory",
- skip(pool),
- fields(
- page = params.page,
- page_size = params.page_size,
- stock_status = ?params.stock_status
- )
-)]
pub async fn list_inventory(
State(pool): State<PgPool>,
Query(params): Query<InventoryQueryParams>,
@@ -44,92 +36,10 @@ pub async fn list_inventory(
// Apply pagination defaults and constraints
let (page, page_size, offset) = calculate_pagination(params.page, params.page_size);
- // Build count query with QueryBuilder for safe parameter binding
- let mut count_builder: QueryBuilder<Postgres> =
- QueryBuilder::new("SELECT COUNT(*) FROM v_product_inventory_pricing");
-
- let mut has_filters = false;
-
- // Apply filters to count query
- if let Some(ref status) = params.stock_status {
- count_builder.push(if has_filters { " AND " } else { " WHERE " });
- has_filters = true;
- count_builder.push("stock_status = ");
- count_builder.push_bind(status);
- }
-
- if let Some(product_uuid) = params.product_uuid {
- count_builder.push(if has_filters { " AND " } else { " WHERE " });
- has_filters = true;
- count_builder.push("product_uuid = ");
- count_builder.push_bind(product_uuid);
- }
-
- if let Some(min_stock) = params.min_stock {
- count_builder.push(if has_filters { " AND " } else { " WHERE " });
- has_filters = true;
- count_builder.push("available_quantity >= ");
- count_builder.push_bind(min_stock);
- }
-
- if let Some(max_stock) = params.max_stock {
- count_builder.push(if has_filters { " AND " } else { " WHERE " });
- count_builder.push("available_quantity <= ");
- count_builder.push_bind(max_stock);
- }
-
- // Execute count query
- let total_count: i64 = match count_builder.build_query_scalar().fetch_one(&pool).await {
- Ok(count) => count,
- Err(e) => {
- return internal_error("Failed to count inventory", e.to_string());
- }
- };
-
- // Build main query with same filters
- let mut query_builder: QueryBuilder<Postgres> =
- QueryBuilder::new("SELECT * FROM v_product_inventory_pricing");
-
- let mut has_filters = false;
-
- // Apply same filters to main query
- if let Some(ref status) = params.stock_status {
- query_builder.push(if has_filters { " AND " } else { " WHERE " });
- has_filters = true;
- query_builder.push("stock_status = ");
- query_builder.push_bind(status);
- }
-
- if let Some(product_uuid) = params.product_uuid {
- query_builder.push(if has_filters { " AND " } else { " WHERE " });
- has_filters = true;
- query_builder.push("product_uuid = ");
- query_builder.push_bind(product_uuid);
- }
-
- if let Some(min_stock) = params.min_stock {
- query_builder.push(if has_filters { " AND " } else { " WHERE " });
- has_filters = true;
- query_builder.push("available_quantity >= ");
- query_builder.push_bind(min_stock);
- }
-
- if let Some(max_stock) = params.max_stock {
- query_builder.push(if has_filters { " AND " } else { " WHERE " });
- query_builder.push("available_quantity <= ");
- query_builder.push_bind(max_stock);
- }
-
- // Add ordering and pagination
- query_builder.push(" ORDER BY stock_status DESC, available_quantity ASC LIMIT ");
- query_builder.push_bind(page_size);
- query_builder.push(" OFFSET ");
- query_builder.push_bind(offset);
-
- // Execute main query
- let inventory: Vec<InventoryWithPricing> =
- match query_builder.build_query_as().fetch_all(&pool).await {
- Ok(inventory) => inventory,
+ // Delegate to repository layer for database operations
+ let (inventory, total_count) =
+ match db::list_inventory(&pool, ¶ms, page, page_size, offset).await {
+ Ok(result) => result,
Err(e) => {
return internal_error("Failed to fetch inventory", e.to_string());
}
@@ -152,26 +62,12 @@ pub async fn list_inventory(
///
/// # Endpoint
/// `GET /inventory/{product_uuid}`
-#[instrument(
- name = "get_inventory_by_product",
- skip(pool),
- fields(product.uuid = %product_uuid)
-)]
pub async fn get_inventory_by_product(
State(pool): State<PgPool>,
Path(product_uuid): Path<Uuid>,
) -> impl IntoResponse {
- let result = sqlx::query_as::<_, InventoryWithPricing>(
- r#"
- SELECT * FROM v_product_inventory_pricing
- WHERE product_uuid = $1
- "#,
- )
- .bind(product_uuid)
- .fetch_optional(&pool)
- .await;
-
- match result {
+ // Delegate to repository layer for database lookup
+ match db::get_inventory_by_product(&pool, product_uuid).await {
Ok(Some(inventory)) => (StatusCode::OK, Json(inventory)).into_response(),
Ok(None) => not_found_error(
"Inventory not found for product",
@@ -185,51 +81,31 @@ pub async fn get_inventory_by_product(
///
/// # Endpoint
/// `PUT /inventory/{product_uuid}`
-#[instrument(
- name = "update_stock",
- skip(pool, request),
- fields(
- product.uuid = %product_uuid,
- quantity = request.quantity
- )
-)]
pub async fn update_stock(
State(pool): State<PgPool>,
Path(product_uuid): Path<Uuid>,
Json(request): Json<UpdateStockRequest>,
) -> impl IntoResponse {
- let result = sqlx::query(
- r#"
- UPDATE product_inventory
- SET stock_quantity = $1,
- reorder_level = COALESCE($2, reorder_level),
- reorder_quantity = COALESCE($3, reorder_quantity),
- last_restocked_at = CURRENT_TIMESTAMP
- WHERE product_uuid = $4
- RETURNING available_quantity
- "#,
+ // Delegate to repository layer for stock update
+ match db::update_stock(
+ &pool,
+ product_uuid,
+ request.quantity,
+ request.reorder_level,
+ request.reorder_quantity,
)
- .bind(request.quantity)
- .bind(request.reorder_level)
- .bind(request.reorder_quantity)
- .bind(product_uuid)
- .fetch_optional(&pool)
- .await;
-
- match result {
- Ok(Some(row)) => {
- let available_quantity: Option<i32> = row.get("available_quantity");
- (
- StatusCode::OK,
- Json(StockOperationResponse {
- success: true,
- message: "Stock updated successfully".to_string(),
- product_uuid,
- available_quantity,
- }),
- )
- .into_response()
- }
+ .await
+ {
+ Ok(Some(available_quantity)) => (
+ StatusCode::OK,
+ Json(StockOperationResponse {
+ success: true,
+ message: "Stock updated successfully".to_string(),
+ product_uuid,
+ available_quantity: Some(available_quantity),
+ }),
+ )
+ .into_response(),
Ok(None) => not_found_error(
"Product not found in inventory",
serde_json::json!({"product_uuid": product_uuid.to_string()}),
@@ -242,25 +118,12 @@ pub async fn update_stock(
///
/// # Endpoint
/// `POST /inventory/reserve`
-#[instrument(
- name = "reserve_stock",
- skip(pool),
- fields(
- product.uuid = %request.product_uuid,
- quantity = request.quantity
- )
-)]
pub async fn reserve_stock(
State(pool): State<PgPool>,
Json(request): Json<ReserveStockRequest>,
) -> impl IntoResponse {
- let result = sqlx::query_scalar::<_, bool>(r#"SELECT reserve_stock($1, $2)"#)
- .bind(request.product_uuid)
- .bind(request.quantity)
- .fetch_one(&pool)
- .await;
-
- match result {
+ // Delegate to repository layer for stock reservation
+ match db::reserve_stock(&pool, request.product_uuid, request.quantity).await {
Ok(success) => {
if success {
(
@@ -294,25 +157,12 @@ pub async fn reserve_stock(
///
/// # Endpoint
/// `POST /inventory/release`
-#[instrument(
- name = "release_stock",
- skip(pool),
- fields(
- product.uuid = %request.product_uuid,
- quantity = request.quantity
- )
-)]
pub async fn release_stock(
State(pool): State<PgPool>,
Json(request): Json<ReleaseStockRequest>,
) -> impl IntoResponse {
- let result = sqlx::query(r#"SELECT release_stock($1, $2)"#)
- .bind(request.product_uuid)
- .bind(request.quantity)
- .execute(&pool)
- .await;
-
- match result {
+ // Delegate to repository layer for stock release
+ match db::release_stock(&pool, request.product_uuid, request.quantity).await {
Ok(_) => (
StatusCode::OK,
Json(StockOperationResponse {
@@ -331,27 +181,19 @@ pub async fn release_stock(
///
/// # Endpoint
/// `POST /inventory/confirm-sale`
-#[instrument(
- name = "confirm_sale",
- skip(pool),
- fields(
- product.uuid = %request.product_uuid,
- order.uuid = %request.order_uuid,
- quantity = request.quantity
- )
-)]
pub async fn confirm_sale(
State(pool): State<PgPool>,
Json(request): Json<ConfirmSaleRequest>,
) -> impl IntoResponse {
- let result = sqlx::query(r#"SELECT confirm_stock_sale($1, $2, $3)"#)
- .bind(request.product_uuid)
- .bind(request.quantity)
- .bind(request.order_uuid)
- .execute(&pool)
- .await;
-
- match result {
+ // Delegate to repository layer for sale confirmation
+ match db::confirm_sale(
+ &pool,
+ request.product_uuid,
+ request.quantity,
+ request.order_uuid,
+ )
+ .await
+ {
Ok(_) => (
StatusCode::OK,
Json(StockOperationResponse {
diff --git a/distributed_observability/inventory/src/handlers/pricing.rs b/distributed_observability/inventory/src/handlers/pricing.rs
@@ -1,17 +1,18 @@
//! Pricing API handlers
//!
//! This module contains the HTTP request handlers for pricing-related endpoints.
+//! Database operations are delegated to the repository layer in `db::pricing`.
use axum::{
extract::{Path, Query, State},
http::StatusCode,
response::{IntoResponse, Json},
};
-use sqlx::{PgPool, Postgres, QueryBuilder};
-use tracing::instrument;
+use sqlx::PgPool;
use uuid::Uuid;
-use crate::models::{PricingQueryParams, PricingResponse, ProductPricing, UpdatePricingRequest};
+use crate::db;
+use crate::models::{PricingQueryParams, PricingResponse, UpdatePricingRequest};
use crate::utils::{calculate_pagination, calculate_total_pages, internal_error, not_found_error};
/// List pricing with pagination and filtering
@@ -25,7 +26,6 @@ use crate::utils::{calculate_pagination, calculate_total_pages, internal_error,
/// - `product_uuid` - Filter by specific product
/// - `min_price` / `max_price` - Filter by price range
/// - `has_discount` - Filter by discount presence
-#[instrument(name = "list_pricing", skip(pool, params))]
pub async fn list_pricing(
State(pool): State<PgPool>,
Query(params): Query<PricingQueryParams>,
@@ -33,83 +33,14 @@ pub async fn list_pricing(
// Apply pagination defaults and constraints
let (page, page_size, offset) = calculate_pagination(params.page, params.page_size);
- // Build count query with QueryBuilder for safe parameter binding
- let mut count_builder: QueryBuilder<Postgres> =
- QueryBuilder::new("SELECT COUNT(*) FROM product_pricing WHERE is_active = true");
-
- // Apply filters to count query (all use AND since WHERE is_active is already present)
- if let Some(product_uuid) = params.product_uuid {
- count_builder.push(" AND product_uuid = ");
- count_builder.push_bind(product_uuid);
- }
-
- if let Some(ref min_price) = params.min_price {
- count_builder.push(" AND final_price >= ");
- count_builder.push_bind(min_price);
- }
-
- if let Some(ref max_price) = params.max_price {
- count_builder.push(" AND final_price <= ");
- count_builder.push_bind(max_price);
- }
-
- if let Some(has_discount) = params.has_discount {
- if has_discount {
- count_builder.push(" AND discount_percentage > 0");
- } else {
- count_builder.push(" AND (discount_percentage IS NULL OR discount_percentage = 0)");
- }
- }
-
- // Execute count query
- let total_count: i64 = match count_builder.build_query_scalar().fetch_one(&pool).await {
- Ok(count) => count,
- Err(e) => {
- return internal_error("Failed to count pricing", e.to_string());
- }
- };
-
- // Build main query with same filters
- let mut query_builder: QueryBuilder<Postgres> =
- QueryBuilder::new("SELECT * FROM product_pricing WHERE is_active = true");
-
- // Apply same filters to main query (all use AND since WHERE is_active is already present)
- if let Some(product_uuid) = params.product_uuid {
- query_builder.push(" AND product_uuid = ");
- query_builder.push_bind(product_uuid);
- }
-
- if let Some(ref min_price) = params.min_price {
- query_builder.push(" AND final_price >= ");
- query_builder.push_bind(min_price);
- }
-
- if let Some(ref max_price) = params.max_price {
- query_builder.push(" AND final_price <= ");
- query_builder.push_bind(max_price);
- }
-
- if let Some(has_discount) = params.has_discount {
- if has_discount {
- query_builder.push(" AND discount_percentage > 0");
- } else {
- query_builder.push(" AND (discount_percentage IS NULL OR discount_percentage = 0)");
- }
- }
-
- // Add ordering and pagination
- query_builder.push(" ORDER BY updated_at DESC LIMIT ");
- query_builder.push_bind(page_size);
- query_builder.push(" OFFSET ");
- query_builder.push_bind(offset);
-
- // Execute main query
- let pricing: Vec<ProductPricing> = match query_builder.build_query_as().fetch_all(&pool).await {
- Ok(pricing) => pricing,
- Err(e) => {
- return internal_error("Failed to fetch pricing", e.to_string());
- }
- };
+ // Delegate to repository layer for database operations
+ let (pricing, total_count) =
+ match db::list_pricing(&pool, ¶ms, page, page_size, offset).await {
+ Ok(result) => result,
+ Err(e) => {
+ return internal_error("Failed to fetch pricing", e.to_string());
+ }
+ };
let total_pages = calculate_total_pages(total_count, page_size);
@@ -128,22 +59,12 @@ pub async fn list_pricing(
///
/// # Endpoint
/// `GET /pricing/{product_uuid}`
-#[instrument(name = "get_pricing_by_product", skip(pool), fields(product.uuid = %product_uuid))]
pub async fn get_pricing_by_product(
State(pool): State<PgPool>,
Path(product_uuid): Path<Uuid>,
) -> impl IntoResponse {
- let result = sqlx::query_as::<_, ProductPricing>(
- r#"
- SELECT * FROM product_pricing
- WHERE product_uuid = $1 AND is_active = true
- "#,
- )
- .bind(product_uuid)
- .fetch_optional(&pool)
- .await;
-
- match result {
+ // Delegate to repository layer for database lookup
+ match db::get_pricing_by_product(&pool, product_uuid).await {
Ok(Some(pricing)) => (StatusCode::OK, Json(pricing)).into_response(),
Ok(None) => not_found_error(
"Pricing not found for product",
@@ -157,53 +78,23 @@ pub async fn get_pricing_by_product(
///
/// # Endpoint
/// `PUT /pricing/{product_uuid}`
-#[instrument(name = "upsert_pricing", skip(pool, request), fields(product.uuid = %product_uuid))]
pub async fn upsert_pricing(
State(pool): State<PgPool>,
Path(product_uuid): Path<Uuid>,
Json(request): Json<UpdatePricingRequest>,
) -> impl IntoResponse {
- // First, deactivate existing active pricing
- if let Err(e) = sqlx::query(
- r#"
- UPDATE product_pricing
- SET is_active = false
- WHERE product_uuid = $1 AND is_active = true
- "#,
+ // Delegate to repository layer for pricing upsert
+ match db::upsert_pricing(
+ &pool,
+ product_uuid,
+ request.final_price,
+ request.initial_price,
+ request.currency.unwrap_or_else(|| "USD".to_string()),
+ request.price_valid_from,
+ request.price_valid_until,
)
- .bind(product_uuid)
- .execute(&pool)
.await
{
- eprintln!("Error deactivating old pricing: {}", e);
- }
-
- // Insert new pricing
- let result = sqlx::query_as::<_, ProductPricing>(
- r#"
- INSERT INTO product_pricing (
- product_uuid,
- final_price,
- initial_price,
- currency,
- price_valid_from,
- price_valid_until,
- is_active
- )
- VALUES ($1, $2, $3, $4, $5, $6, true)
- RETURNING *
- "#,
- )
- .bind(product_uuid)
- .bind(request.final_price)
- .bind(request.initial_price)
- .bind(request.currency.unwrap_or_else(|| "USD".to_string()))
- .bind(request.price_valid_from)
- .bind(request.price_valid_until)
- .fetch_one(&pool)
- .await;
-
- match result {
Ok(pricing) => (StatusCode::OK, Json(pricing)).into_response(),
Err(e) => internal_error("Failed to update pricing", e.to_string()),
}
diff --git a/distributed_observability/orders/src/db/mod.rs b/distributed_observability/orders/src/db/mod.rs
@@ -22,6 +22,12 @@ use sqlx::{
};
use std::str::FromStr;
+pub mod repository;
+pub mod transaction;
+
+pub use repository::*;
+pub use transaction::*;
+
/// Database connection pool wrapper
///
/// Wraps sqlx's PgPool with custom configuration for the orders service.
diff --git a/distributed_observability/orders/src/db/repository.rs b/distributed_observability/orders/src/db/repository.rs
@@ -0,0 +1,647 @@
+//! Order repository functions with OpenTelemetry instrumentation
+//!
+//! This module contains database operations for orders and shipments,
+//! separated from HTTP handlers.
+//! Each function is instrumented with OpenTelemetry semantic conventions for database spans.
+
+use bigdecimal::BigDecimal;
+use chrono::NaiveDate;
+use sqlx::{PgPool, Postgres, QueryBuilder, Row, Transaction};
+use tracing::{instrument, Span};
+use uuid::Uuid;
+
+use crate::models::{
+ CreateOrderItemRequest, CreatePaymentRequest, CreateShippingAddressRequest, OrderComplete,
+ Shipment,
+};
+
+/// Result of creating an order - contains the generated IDs
+#[derive(Debug)]
+pub struct CreatedOrder {
+ pub id: i32,
+ pub uuid: Uuid,
+ pub order_number: String,
+}
+
+/// Totals calculated for an order
+#[derive(Debug, Clone)]
+pub struct OrderTotals {
+ pub subtotal: BigDecimal,
+ pub tax_amount: BigDecimal,
+ pub shipping_amount: BigDecimal,
+ pub total: BigDecimal,
+}
+
+/// Generate an order number using the database function
+#[instrument(
+ name = "SELECT generate_order_number",
+ skip(tx),
+ fields(
+ db.system.name = "postgresql",
+ db.namespace = "orders",
+ db.operation.name = "SELECT",
+ db.collection.name = "orders",
+ db.query.text = "SELECT generate_order_number()"
+ )
+)]
+pub async fn generate_order_number(
+ tx: &mut Transaction<'_, Postgres>,
+) -> Result<String, sqlx::Error> {
+ sqlx::query_scalar("SELECT generate_order_number()")
+ .fetch_one(&mut **tx)
+ .await
+}
+
+/// Create a new order record
+///
+/// Inserts the order header with customer info and totals.
+/// Returns the generated order ID, UUID, and order number.
+#[instrument(
+ name = "INSERT orders",
+ skip(tx, totals),
+ fields(
+ db.system.name = "postgresql",
+ db.namespace = "orders",
+ db.operation.name = "INSERT",
+ db.collection.name = "orders",
+ db.query.text = "INSERT INTO orders (...) VALUES (...) RETURNING id, uuid",
+ db.response.returned_rows = tracing::field::Empty
+ )
+)]
+pub async fn create_order(
+ tx: &mut Transaction<'_, Postgres>,
+ order_number: &str,
+ customer_email: &str,
+ customer_phone: Option<&str>,
+ totals: &OrderTotals,
+) -> Result<CreatedOrder, sqlx::Error> {
+ let row = sqlx::query(
+ r#"
+ INSERT INTO orders (
+ order_number, customer_email, customer_phone,
+ subtotal, tax_amount, shipping_amount, total,
+ status, payment_status
+ )
+ VALUES ($1, $2, $3, $4, $5, $6, $7, 'pending', 'pending')
+ RETURNING id, uuid
+ "#,
+ )
+ .bind(order_number)
+ .bind(customer_email)
+ .bind(customer_phone)
+ .bind(&totals.subtotal)
+ .bind(&totals.tax_amount)
+ .bind(&totals.shipping_amount)
+ .bind(&totals.total)
+ .fetch_one(&mut **tx)
+ .await?;
+
+ Span::current().record("db.response.returned_rows", 1);
+
+ Ok(CreatedOrder {
+ id: row.get("id"),
+ uuid: row.get("uuid"),
+ order_number: order_number.to_string(),
+ })
+}
+
+/// Create order items (line items) for an order
+///
+/// Inserts all order items in sequence. Each item contains a snapshot
+/// of the product data at the time of order.
+#[instrument(
+ name = "INSERT order_items",
+ skip(tx, items),
+ fields(
+ db.system.name = "postgresql",
+ db.namespace = "orders",
+ db.operation.name = "INSERT",
+ db.collection.name = "order_items",
+ db.query.text = "INSERT INTO order_items (...) VALUES (...)",
+ otelmart.order.id = order_id,
+ otelmart.items.count = items.len(),
+ db.response.returned_rows = tracing::field::Empty
+ )
+)]
+pub async fn create_order_items(
+ tx: &mut Transaction<'_, Postgres>,
+ order_id: i32,
+ items: &[CreateOrderItemRequest],
+) -> Result<(), sqlx::Error> {
+ let mut rows_inserted = 0;
+
+ for item in items {
+ let total_price = &item.unit_price * BigDecimal::from(item.quantity);
+ sqlx::query(
+ r#"
+ INSERT INTO order_items (
+ order_id, product_uuid, product_name, product_sku,
+ quantity, unit_price, total_price
+ )
+ VALUES ($1, $2, $3, $4, $5, $6, $7)
+ "#,
+ )
+ .bind(order_id)
+ .bind(item.product_uuid)
+ .bind(&item.product_name)
+ .bind(&item.product_sku)
+ .bind(item.quantity)
+ .bind(&item.unit_price)
+ .bind(&total_price)
+ .execute(&mut **tx)
+ .await?;
+
+ rows_inserted += 1;
+ }
+
+ Span::current().record("db.response.returned_rows", rows_inserted);
+ Ok(())
+}
+
+/// Create shipping address for an order
+#[instrument(
+ name = "INSERT shipping_addresses",
+ skip(tx, address),
+ fields(
+ db.system.name = "postgresql",
+ db.namespace = "orders",
+ db.operation.name = "INSERT",
+ db.collection.name = "shipping_addresses",
+ db.query.text = "INSERT INTO shipping_addresses (...) VALUES (...)",
+ otelmart.order.id = order_id,
+ db.response.returned_rows = tracing::field::Empty
+ )
+)]
+pub async fn create_shipping_address(
+ tx: &mut Transaction<'_, Postgres>,
+ order_id: i32,
+ address: &CreateShippingAddressRequest,
+) -> Result<(), sqlx::Error> {
+ sqlx::query(
+ r#"
+ INSERT INTO shipping_addresses (
+ order_id, first_name, last_name, address_line1, address_line2,
+ city, state, postal_code, country, phone
+ )
+ VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10)
+ "#,
+ )
+ .bind(order_id)
+ .bind(&address.first_name)
+ .bind(&address.last_name)
+ .bind(&address.address_line1)
+ .bind(&address.address_line2)
+ .bind(&address.city)
+ .bind(&address.state)
+ .bind(&address.postal_code)
+ .bind(&address.country)
+ .bind(&address.phone)
+ .execute(&mut **tx)
+ .await?;
+
+ Span::current().record("db.response.returned_rows", 1);
+ Ok(())
+}
+
+/// Generate a payment reference using the database function
+#[instrument(
+ name = "SELECT generate_payment_reference",
+ skip(tx),
+ fields(
+ otel.kind = "client",
+ db.system.name = "postgresql",
+ db.namespace = "orders",
+ db.operation.name = "SELECT",
+ db.collection.name = "payments"
+ )
+)]
+pub async fn generate_payment_reference(
+ tx: &mut Transaction<'_, Postgres>,
+) -> Result<String, sqlx::Error> {
+ sqlx::query_scalar("SELECT generate_payment_reference()")
+ .fetch_one(&mut **tx)
+ .await
+}
+
+/// Create payment record for an order
+///
+/// NOTE: Payment is simulated - in production, this would call a payment gateway
+/// and only insert the record after successful payment processing.
+#[instrument(
+ name = "INSERT payments",
+ skip(tx, payment),
+ fields(
+ db.system.name = "postgresql",
+ db.namespace = "orders",
+ db.operation.name = "INSERT",
+ db.collection.name = "payments",
+ db.query.text = "INSERT INTO payments (...) VALUES (...)",
+ otelmart.order.id = order_id,
+ db.response.returned_rows = tracing::field::Empty
+ )
+)]
+pub async fn create_payment(
+ tx: &mut Transaction<'_, Postgres>,
+ order_id: i32,
+ payment: &CreatePaymentRequest,
+ payment_reference: &str,
+ amount: &BigDecimal,
+) -> Result<(), sqlx::Error> {
+ sqlx::query(
+ r#"
+ INSERT INTO payments (
+ order_id, payment_method, amount, status,
+ payment_reference, card_last4, card_brand, processed_at
+ )
+ VALUES ($1, $2, $3, 'paid', $4, $5, $6, CURRENT_TIMESTAMP)
+ "#,
+ )
+ .bind(order_id)
+ .bind(&payment.payment_method)
+ .bind(amount)
+ .bind(payment_reference)
+ .bind(&payment.card_last4)
+ .bind(&payment.card_brand)
+ .execute(&mut **tx)
+ .await?;
+
+ Span::current().record("db.response.returned_rows", 1);
+ Ok(())
+}
+
+/// Update order status after successful payment
+#[instrument(
+ name = "UPDATE orders",
+ skip(tx),
+ fields(
+ otel.kind = "client",
+ db.system.name = "postgresql",
+ db.namespace = "orders",
+ db.operation.name = "UPDATE",
+ db.collection.name = "orders",
+ db.query.text = "UPDATE orders SET payment_status = $1, status = $2 WHERE id = $3",
+ otelmart.order.id = order_id,
+ db.response.returned_rows = tracing::field::Empty
+ )
+)]
+pub async fn update_order_payment_status(
+ tx: &mut Transaction<'_, Postgres>,
+ order_id: i32,
+ payment_status: &str,
+ order_status: &str,
+) -> Result<(), sqlx::Error> {
+ let result = sqlx::query(
+ r#"
+ UPDATE orders
+ SET payment_status = $1, status = $2
+ WHERE id = $3
+ "#,
+ )
+ .bind(payment_status)
+ .bind(order_status)
+ .bind(order_id)
+ .execute(&mut **tx)
+ .await?;
+
+ Span::current().record("db.response.returned_rows", result.rows_affected());
+ Ok(())
+}
+
+/// Get order by UUID
+#[instrument(
+ name = "SELECT orders",
+ skip(pool),
+ fields(
+ otel.kind = "client",
+ db.system.name = "postgresql",
+ db.namespace = "orders",
+ db.operation.name = "SELECT",
+ db.collection.name = "orders",
+ db.query.text = "SELECT * FROM v_orders_complete WHERE eid = $1",
+ otelmart.order.uuid = %uuid
+ )
+)]
+pub async fn get_order_by_uuid(
+ pool: &PgPool,
+ uuid: Uuid,
+) -> Result<Option<OrderComplete>, sqlx::Error> {
+ sqlx::query_as::<_, OrderComplete>(r#"SELECT * FROM v_orders_complete WHERE eid = $1"#)
+ .bind(uuid)
+ .fetch_optional(pool)
+ .await
+}
+
+/// Get order by email and order number (for guest users)
+#[instrument(
+ name = "SELECT orders",
+ skip(pool),
+ fields(
+ otel.kind = "client",
+ db.system.name = "postgresql",
+ db.namespace = "orders",
+ db.operation.name = "SELECT",
+ db.collection.name = "orders",
+ db.query.text = "SELECT * FROM v_orders_complete WHERE customer_email = $1 AND order_number = $2"
+ )
+)]
+pub async fn get_order_by_email_and_number(
+ pool: &PgPool,
+ email: &str,
+ order_number: &str,
+) -> Result<Option<OrderComplete>, sqlx::Error> {
+ sqlx::query_as::<_, OrderComplete>(
+ r#"
+ SELECT * FROM v_orders_complete
+ WHERE customer_email = $1 AND order_number = $2
+ "#,
+ )
+ .bind(email)
+ .bind(order_number)
+ .fetch_optional(pool)
+ .await
+}
+
+/// List orders for a user with pagination
+#[instrument(
+ name = "SELECT orders",
+ skip(pool),
+ fields(
+ otel.kind = "client",
+ db.system.name = "postgresql",
+ db.namespace = "orders",
+ db.operation.name = "SELECT",
+ db.collection.name = "orders",
+ db.query.text = "SELECT * FROM v_orders_complete WHERE customer_email = $1 ...",
+ otelmart.page = page,
+ otelmart.page_size = page_size
+ )
+)]
+pub async fn list_orders_by_email(
+ pool: &PgPool,
+ email: &str,
+ status_filter: Option<&str>,
+ payment_status_filter: Option<&str>,
+ page: i32,
+ page_size: i32,
+) -> Result<(Vec<OrderComplete>, i64), sqlx::Error> {
+ let offset = (page - 1) * page_size;
+
+ // Build count query with proper parameterized filters
+ let mut count_builder: QueryBuilder<Postgres> =
+ QueryBuilder::new("SELECT COUNT(*) FROM v_orders_complete WHERE customer_email = ");
+ count_builder.push_bind(email);
+
+ if let Some(status) = status_filter {
+ count_builder.push(" AND status = ");
+ count_builder.push_bind(status);
+ }
+
+ if let Some(payment_status) = payment_status_filter {
+ count_builder.push(" AND payment_status = ");
+ count_builder.push_bind(payment_status);
+ }
+
+ // Execute count query
+ let total_count: (i64,) = count_builder.build_query_as().fetch_one(pool).await?;
+
+ // Build main query with proper parameterized filters
+ let mut query_builder: QueryBuilder<Postgres> =
+ QueryBuilder::new("SELECT * FROM v_orders_complete WHERE customer_email = ");
+ query_builder.push_bind(email);
+
+ if let Some(status) = status_filter {
+ query_builder.push(" AND status = ");
+ query_builder.push_bind(status);
+ }
+
+ if let Some(payment_status) = payment_status_filter {
+ query_builder.push(" AND payment_status = ");
+ query_builder.push_bind(payment_status);
+ }
+
+ query_builder.push(" ORDER BY ordered_at DESC LIMIT ");
+ query_builder.push_bind(page_size);
+ query_builder.push(" OFFSET ");
+ query_builder.push_bind(offset);
+
+ // Execute main query
+ let orders: Vec<OrderComplete> = query_builder.build_query_as().fetch_all(pool).await?;
+
+ Ok((orders, total_count.0))
+}
+
+/// Get the internal order ID from a public UUID
+#[instrument(
+ name = "SELECT orders",
+ skip(pool),
+ fields(
+ otel.kind = "client",
+ db.system.name = "postgresql",
+ db.namespace = "orders",
+ db.operation.name = "SELECT",
+ db.collection.name = "orders",
+ db.query.text = "SELECT id FROM orders WHERE uuid = $1",
+ otelmart.order.uuid = %order_uuid
+ )
+)]
+pub async fn get_order_id_by_uuid(
+ pool: &PgPool,
+ order_uuid: Uuid,
+) -> Result<Option<i32>, sqlx::Error> {
+ sqlx::query_scalar("SELECT id FROM orders WHERE uuid = $1")
+ .bind(order_uuid)
+ .fetch_optional(pool)
+ .await
+}
+
+/// Create a new shipment record for an order
+///
+/// Inserts a shipment with status 'shipped' and sets shipped_at to now.
+#[instrument(
+ name = "INSERT shipments",
+ skip(pool),
+ fields(
+ db.system.name = "postgresql",
+ db.namespace = "orders",
+ db.operation.name = "INSERT",
+ db.collection.name = "shipments",
+ db.query.text = "INSERT INTO shipments (...) VALUES (...) RETURNING *",
+ otelmart.order.id = order_id,
+ db.response.returned_rows = tracing::field::Empty
+ )
+)]
+pub async fn create_shipment(
+ pool: &PgPool,
+ order_id: i32,
+ carrier: &str,
+ tracking_number: Option<&str>,
+ estimated_delivery_date: Option<NaiveDate>,
+) -> Result<Shipment, sqlx::Error> {
+ let shipment = sqlx::query_as::<_, Shipment>(
+ r#"
+ INSERT INTO shipments (
+ order_id, carrier, tracking_number,
+ estimated_delivery_date, status, shipped_at
+ )
+ VALUES ($1, $2, $3, $4, 'shipped', CURRENT_TIMESTAMP)
+ RETURNING *
+ "#,
+ )
+ .bind(order_id)
+ .bind(carrier)
+ .bind(tracking_number)
+ .bind(estimated_delivery_date)
+ .fetch_one(pool)
+ .await?;
+
+ Span::current().record("db.response.returned_rows", 1);
+ Ok(shipment)
+}
+
+/// Update order status (e.g., to 'shipped' or 'delivered')
+#[instrument(
+ name = "UPDATE orders",
+ skip(pool),
+ fields(
+ otel.kind = "client",
+ db.system.name = "postgresql",
+ db.namespace = "orders",
+ db.operation.name = "UPDATE",
+ db.collection.name = "orders",
+ db.query.text = "UPDATE orders SET status = $1 WHERE id = $2",
+ otelmart.order.id = order_id,
+ otelmart.order.status = status,
+ db.response.returned_rows = tracing::field::Empty
+ )
+)]
+pub async fn update_order_status(
+ pool: &PgPool,
+ order_id: i32,
+ status: &str,
+) -> Result<u64, sqlx::Error> {
+ let result = sqlx::query("UPDATE orders SET status = $1 WHERE id = $2")
+ .bind(status)
+ .bind(order_id)
+ .execute(pool)
+ .await?;
+
+ let rows = result.rows_affected();
+ Span::current().record("db.response.returned_rows", rows);
+ Ok(rows)
+}
+
+/// Update shipment status within a transaction
+///
+/// For 'delivered' status, also sets actual_delivery_date and delivered_at.
+/// For other statuses, only updates the status field.
+/// Returns true if a shipment was found and updated.
+#[instrument(
+ name = "UPDATE shipments",
+ skip(tx),
+ fields(
+ otel.kind = "client",
+ db.system.name = "postgresql",
+ db.namespace = "orders",
+ db.operation.name = "UPDATE",
+ db.collection.name = "shipments",
+ db.query.text = "UPDATE shipments SET status = $1 ... WHERE order_id = ...",
+ otelmart.order.id = order_id,
+ otelmart.shipment.status = status,
+ db.response.returned_rows = tracing::field::Empty
+ )
+)]
+pub async fn update_shipment_status(
+ tx: &mut Transaction<'_, Postgres>,
+ order_id: i32,
+ status: &str,
+ actual_delivery_date: Option<NaiveDate>,
+) -> Result<bool, sqlx::Error> {
+ let result = if status == "delivered" {
+ sqlx::query(
+ r#"
+ UPDATE shipments
+ SET status = $1, actual_delivery_date = $2, delivered_at = CURRENT_TIMESTAMP
+ WHERE order_id = $3
+ RETURNING id
+ "#,
+ )
+ .bind(status)
+ .bind(actual_delivery_date)
+ .bind(order_id)
+ .fetch_optional(&mut **tx)
+ .await?
+ } else {
+ sqlx::query(
+ r#"
+ UPDATE shipments
+ SET status = $1
+ WHERE order_id = $2
+ RETURNING id
+ "#,
+ )
+ .bind(status)
+ .bind(order_id)
+ .fetch_optional(&mut **tx)
+ .await?
+ };
+
+ let found = result.is_some();
+ Span::current().record("db.response.returned_rows", if found { 1 } else { 0 });
+ Ok(found)
+}
+
+/// Update order status within a transaction
+#[instrument(
+ name = "UPDATE orders",
+ skip(tx),
+ fields(
+ otel.kind = "client",
+ db.system.name = "postgresql",
+ db.namespace = "orders",
+ db.operation.name = "UPDATE",
+ db.collection.name = "orders",
+ db.query.text = "UPDATE orders SET status = $1 WHERE id = $2",
+ otelmart.order.id = order_id,
+ otelmart.order.status = status,
+ db.response.returned_rows = tracing::field::Empty
+ )
+)]
+pub async fn update_order_status_in_tx(
+ tx: &mut Transaction<'_, Postgres>,
+ order_id: i32,
+ status: &str,
+) -> Result<u64, sqlx::Error> {
+ let result = sqlx::query("UPDATE orders SET status = $1 WHERE id = $2")
+ .bind(status)
+ .bind(order_id)
+ .execute(&mut **tx)
+ .await?;
+
+ let rows = result.rows_affected();
+ Span::current().record("db.response.returned_rows", rows);
+ Ok(rows)
+}
+
+/// Get order ID from UUID within a transaction
+#[instrument(
+ name = "SELECT orders",
+ skip(tx),
+ fields(
+ otel.kind = "client",
+ db.system.name = "postgresql",
+ db.namespace = "orders",
+ db.operation.name = "SELECT",
+ db.collection.name = "orders",
+ db.query.text = "SELECT id FROM orders WHERE uuid = $1",
+ otelmart.order.uuid = %order_uuid
+ )
+)]
+pub async fn get_order_id_by_uuid_in_tx(
+ tx: &mut Transaction<'_, Postgres>,
+ order_uuid: Uuid,
+) -> Result<Option<i32>, sqlx::Error> {
+ sqlx::query_scalar("SELECT id FROM orders WHERE uuid = $1")
+ .bind(order_uuid)
+ .fetch_optional(&mut **tx)
+ .await
+}
diff --git a/distributed_observability/orders/src/db/transaction.rs b/distributed_observability/orders/src/db/transaction.rs
@@ -0,0 +1,70 @@
+//! Transaction wrapper with OpenTelemetry instrumentation
+//!
+//! Provides a `with_transaction` helper that wraps database transactions
+//! in a span that tracks commit/rollback outcomes.
+
+use sqlx::{PgPool, Postgres, Transaction};
+use std::future::Future;
+use std::pin::Pin;
+use tracing::{instrument, Span};
+
+/// Execute a closure within a database transaction with OpenTelemetry instrumentation.
+///
+/// This wrapper:
+/// - Creates a span named "TRANSACTION" with the transaction name
+/// - Tracks whether the transaction committed or rolled back
+/// - Automatically rolls back on error (via Drop)
+/// - Records `otelmart.transaction.outcome` as "commit" or "rollback"
+///
+/// # Arguments
+/// * `pool` - The database connection pool
+/// * `name` - A descriptive name for the transaction (e.g., "checkout", "update_order")
+/// * `f` - The async closure to execute within the transaction
+///
+/// # Example
+/// ```no_run
+/// let result = with_transaction(&pool, "checkout", |tx| {
+/// Box::pin(async move {
+/// db::create_order(tx, &email, &totals).await?;
+/// db::create_order_items(tx, order_id, &items).await?;
+/// Ok(order)
+/// })
+/// }).await?;
+/// ```
+#[instrument(
+ name = "TRANSACTION",
+ skip(pool, f),
+ fields(
+ otel.kind = "client",
+ db.system.name = "postgresql",
+ db.operation.name = "transaction",
+ otelmart.transaction.name = %name,
+ otelmart.transaction.outcome = tracing::field::Empty
+ )
+)]
+pub async fn with_transaction<F, T, E>(pool: &PgPool, name: &'static str, f: F) -> Result<T, E>
+where
+ F: for<'c> FnOnce(
+ &'c mut Transaction<'_, Postgres>,
+ ) -> Pin<Box<dyn Future<Output = Result<T, E>> + Send + 'c>>,
+ E: From<sqlx::Error>,
+{
+ let mut tx = pool.begin().await.map_err(E::from)?;
+
+ match f(&mut tx).await {
+ Ok(result) => {
+ tx.commit().await.map_err(E::from)?;
+ Span::current().record("otelmart.transaction.outcome", "commit");
+ Ok(result)
+ }
+ Err(e) => {
+ // Rollback is automatic on drop, but we make it explicit for clarity
+ // and to ensure the span records the outcome before the error propagates
+ if let Err(rb_err) = tx.rollback().await {
+ tracing::warn!(error = %rb_err, "explicit rollback failed");
+ }
+ Span::current().record("otelmart.transaction.outcome", "rollback");
+ Err(e)
+ }
+ }
+}
diff --git a/distributed_observability/orders/src/handlers/orders.rs b/distributed_observability/orders/src/handlers/orders.rs
@@ -6,18 +6,22 @@ use axum::{
response::{IntoResponse, Json},
};
use bigdecimal::BigDecimal;
-use sqlx::{PgPool, Row};
-use tracing::instrument;
+use sqlx::PgPool;
+use tracing::{info, warn};
use uuid::Uuid;
-use crate::models::{CreateOrderRequest, OrderComplete, OrderQueryParams, OrdersResponse};
+use crate::db::{self, with_transaction, OrderTotals};
+use crate::models::{CreateOrderRequest, OrderQueryParams, OrdersResponse};
use crate::AppState;
/// Product detail response from products service
#[derive(Debug, serde::Deserialize)]
+#[allow(dead_code)] // Fields deserialized from products service response, not all used directly
struct ProductDetail {
+ eid: Uuid,
product_name: String,
final_price: BigDecimal,
+ stock: i32,
is_active: bool,
}
@@ -45,9 +49,11 @@ struct ConfirmSaleRequest {
/// Response from inventory stock operations
#[derive(Debug, serde::Deserialize)]
+#[allow(dead_code)] // Fields deserialized from inventory service response, not all used directly
struct StockOperationResponse {
success: bool,
message: String,
+ product_uuid: Uuid,
available_quantity: Option<i32>,
}
@@ -63,32 +69,31 @@ async fn release_stock(state: &AppState, product_uuid: Uuid, quantity: i32) {
quantity,
};
- let response = state
- .http_client()
- .post(&url)
- .json(&request_body)
- .send()
- .await;
+ let request = state.http_client().post(&url).json(&request_body);
+
+ let response = request.send().await;
match response {
Ok(resp) => {
if resp.status().is_success() {
- eprintln!(
- "Successfully released {} units of stock for product {}",
- quantity, product_uuid
+ info!(
+ product_uuid = %product_uuid,
+ quantity = quantity,
+ "Successfully released stock"
);
} else {
- eprintln!(
- "Failed to release stock for product {}: status {}",
- product_uuid,
- resp.status()
+ warn!(
+ product_uuid = %product_uuid,
+ status = %resp.status(),
+ "Failed to release stock"
);
}
}
Err(e) => {
- eprintln!(
- "Error calling inventory service to release stock for product {}: {}",
- product_uuid, e
+ warn!(
+ product_uuid = %product_uuid,
+ error = %e,
+ "Error calling inventory service to release stock"
);
}
}
@@ -103,9 +108,9 @@ async fn release_all_reserved_stock(state: &AppState, reserved_items: &[(Uuid, i
return;
}
- eprintln!(
- "Rolling back stock reservations for {} items",
- reserved_items.len()
+ warn!(
+ item_count = reserved_items.len(),
+ "Rolling back stock reservations"
);
for (product_uuid, quantity) in reserved_items {
@@ -131,10 +136,9 @@ async fn confirm_sale(
order_uuid,
};
- let response = state
- .http_client()
- .post(&url)
- .json(&request_body)
+ let request = state.http_client().post(&url).json(&request_body);
+
+ let response = request
.send()
.await
.map_err(|e| format!("Failed to call inventory service: {}", e))?;
@@ -158,27 +162,30 @@ async fn confirm_all_sales(state: &AppState, order_uuid: Uuid, items: &[(Uuid, i
return;
}
- println!(
- "Confirming sales with inventory service for {} items in order {}",
- items.len(),
- order_uuid
+ info!(
+ item_count = items.len(),
+ order_uuid = %order_uuid,
+ "Confirming sales with inventory service"
);
for (product_uuid, quantity) in items {
match confirm_sale(state, *product_uuid, *quantity, order_uuid).await {
Ok(_) => {
- println!(
- "Successfully confirmed sale of {} units for product {} in order {}",
- quantity, product_uuid, order_uuid
+ info!(
+ product_uuid = %product_uuid,
+ quantity = quantity,
+ order_uuid = %order_uuid,
+ "Successfully confirmed sale"
);
}
Err(e) => {
// Log error but don't fail - order is already committed
- eprintln!(
- "Failed to confirm sale for product {} in order {}: {}",
- product_uuid, order_uuid, e
+ warn!(
+ product_uuid = %product_uuid,
+ order_uuid = %order_uuid,
+ error = %e,
+ "Failed to confirm sale - manual intervention may be required"
);
- eprintln!("WARNING: Stock reservation remains in place. Manual intervention may be required.");
}
}
}
@@ -200,13 +207,9 @@ async fn reserve_stock(
quantity,
};
- let response = match state
- .http_client()
- .post(&url)
- .json(&request_body)
- .send()
- .await
- {
+ let request = state.http_client().post(&url).json(&request_body);
+
+ let response = match request.send().await {
Ok(resp) => resp,
Err(e) => {
eprintln!("Failed to call inventory service: {}", e);
@@ -289,7 +292,9 @@ async fn validate_product(
) -> Result<ProductDetail, axum::response::Response> {
let url = format!("{}/products/{}", state.products_service_url(), product_uuid);
- let response = match state.http_client().get(&url).send().await {
+ let request = state.http_client().get(&url);
+
+ let response = match request.send().await {
Ok(resp) => resp,
Err(e) => {
eprintln!("Failed to call products service: {}", e);
@@ -358,6 +363,18 @@ async fn validate_product(
}
}
+/// Order creation error type
+#[derive(Debug)]
+pub enum OrderError {
+ Database(sqlx::Error),
+}
+
+impl From<sqlx::Error> for OrderError {
+ fn from(err: sqlx::Error) -> Self {
+ OrderError::Database(err)
+ }
+}
+
/// Create a new order
///
/// # Endpoint
@@ -365,14 +382,9 @@ async fn validate_product(
///
/// Creates a new order with items, shipping address, and payment info.
/// Automatically generates order number and calculates totals.
-#[instrument(
- name = "create_order",
- skip(state, request),
- fields(
- customer_email = %request.customer_email,
- item_count = request.items.len()
- )
-)]
+///
+/// Uses the `with_transaction` wrapper to ensure all database operations
+/// are atomic and instrumented with OpenTelemetry semantic conventions.
pub async fn create_order(
State(state): State<AppState>,
Json(request): Json<CreateOrderRequest>,
@@ -421,25 +433,7 @@ pub async fn create_order(
}
}
- // Start a transaction
- let mut tx = match state.pool().begin().await {
- Ok(tx) => tx,
- Err(e) => {
- eprintln!("Failed to start transaction: {}", e);
- // Release reserved stock before returning error
- release_all_reserved_stock(&state, &reserved_items).await;
- return (
- StatusCode::INTERNAL_SERVER_ERROR,
- Json(serde_json::json!({
- "error": "Failed to start transaction",
- "details": e.to_string()
- })),
- )
- .into_response();
- }
- };
-
- // Calculate subtotal from items
+ // Calculate totals
let subtotal: BigDecimal = request
.items
.iter()
@@ -454,253 +448,103 @@ pub async fn create_order(
let total = &subtotal + &tax_amount + &shipping_amount;
- // Generate order number
- let order_number: String = match sqlx::query_scalar("SELECT generate_order_number()")
- .fetch_one(&mut *tx)
- .await
- {
- Ok(num) => num,
- Err(e) => {
- eprintln!("Failed to generate order number: {}", e);
- // Release reserved stock before returning error
- release_all_reserved_stock(&state, &reserved_items).await;
- return (
- StatusCode::INTERNAL_SERVER_ERROR,
- Json(serde_json::json!({
- "error": "Failed to generate order number",
- "details": e.to_string()
- })),
- )
- .into_response();
- }
+ let totals = OrderTotals {
+ subtotal,
+ tax_amount,
+ shipping_amount,
+ total: total.clone(),
};
- // Insert order
- let order_result = sqlx::query(
- r#"
- INSERT INTO orders (
- order_number, customer_email, customer_phone,
- subtotal, tax_amount, shipping_amount, total,
- status, payment_status
- )
- VALUES ($1, $2, $3, $4, $5, $6, $7, 'pending', 'pending')
- RETURNING id, uuid
- "#,
- )
- .bind(&order_number)
- .bind(&request.customer_email)
- .bind(&request.customer_phone)
- .bind(&subtotal)
- .bind(&tax_amount)
- .bind(&shipping_amount)
- .bind(&total)
- .fetch_one(&mut *tx)
+ // Clone data needed for the transaction closure
+ let customer_email = request.customer_email.clone();
+ let customer_phone = request.customer_phone.clone();
+ let items = request.items.clone();
+ let shipping_address = request.shipping_address.clone();
+ let payment = request.payment.clone();
+
+ // Execute all database operations within an instrumented transaction
+ let result = with_transaction(state.pool(), "checkout", |tx| {
+ // Move owned values into the closure
+ let customer_email = customer_email.clone();
+ let customer_phone = customer_phone.clone();
+ let items = items.clone();
+ let shipping_address = shipping_address.clone();
+ let payment = payment.clone();
+ let totals = totals.clone();
+ let total = total.clone();
+
+ Box::pin(async move {
+ // Generate order number
+ let order_number = db::generate_order_number(tx).await?;
+
+ // Create order record
+ let created_order = db::create_order(
+ tx,
+ &order_number,
+ &customer_email,
+ customer_phone.as_deref(),
+ &totals,
+ )
+ .await?;
+
+ // Create order items
+ db::create_order_items(tx, created_order.id, &items).await?;
+
+ // Create shipping address
+ db::create_shipping_address(tx, created_order.id, &shipping_address).await?;
+
+ // Generate payment reference and create payment
+ let payment_reference = db::generate_payment_reference(tx).await?;
+ db::create_payment(tx, created_order.id, &payment, &payment_reference, &total).await?;
+
+ // Update order status to processing after successful payment
+ db::update_order_payment_status(tx, created_order.id, "paid", "processing").await?;
+
+ Ok::<_, OrderError>(created_order)
+ })
+ })
.await;
- let (order_id, order_uuid): (i32, Uuid) = match order_result {
- Ok(row) => (row.get("id"), row.get("uuid")),
- Err(e) => {
- eprintln!("Failed to create order: {}", e);
- let _ = tx.rollback().await;
- // Release reserved stock before returning error
- release_all_reserved_stock(&state, &reserved_items).await;
- return (
- StatusCode::INTERNAL_SERVER_ERROR,
- Json(serde_json::json!({
- "error": "Failed to create order",
- "details": e.to_string()
- })),
- )
- .into_response();
- }
- };
+ match result {
+ Ok(created_order) => {
+ // Transaction committed successfully!
+ // Now confirm the sale with inventory service to convert reserved stock to sold
+ confirm_all_sales(&state, created_order.uuid, &reserved_items).await;
- // Insert order items
- for item in &request.items {
- let total_price = &item.unit_price * BigDecimal::from(item.quantity);
- if let Err(e) = sqlx::query(
- r#"
- INSERT INTO order_items (
- order_id, product_uuid, product_name, product_sku,
- quantity, unit_price, total_price
- )
- VALUES ($1, $2, $3, $4, $5, $6, $7)
- "#,
- )
- .bind(order_id)
- .bind(item.product_uuid)
- .bind(&item.product_name)
- .bind(&item.product_sku)
- .bind(item.quantity)
- .bind(&item.unit_price)
- .bind(&total_price)
- .execute(&mut *tx)
- .await
- {
- eprintln!("Failed to insert order item: {}", e);
- let _ = tx.rollback().await;
- // Release reserved stock before returning error
- release_all_reserved_stock(&state, &reserved_items).await;
- return (
- StatusCode::INTERNAL_SERVER_ERROR,
+ (
+ StatusCode::CREATED,
Json(serde_json::json!({
- "error": "Failed to insert order item",
- "details": e.to_string()
+ "success": true,
+ "order_number": created_order.order_number,
+ "order_uuid": created_order.uuid,
+ "message": "Order created successfully"
})),
)
- .into_response();
+ .into_response()
}
- }
-
- // Insert shipping address
- if let Err(e) = sqlx::query(
- r#"
- INSERT INTO shipping_addresses (
- order_id, first_name, last_name, address_line1, address_line2,
- city, state, postal_code, country, phone
- )
- VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10)
- "#,
- )
- .bind(order_id)
- .bind(&request.shipping_address.first_name)
- .bind(&request.shipping_address.last_name)
- .bind(&request.shipping_address.address_line1)
- .bind(&request.shipping_address.address_line2)
- .bind(&request.shipping_address.city)
- .bind(&request.shipping_address.state)
- .bind(&request.shipping_address.postal_code)
- .bind(&request.shipping_address.country)
- .bind(&request.shipping_address.phone)
- .execute(&mut *tx)
- .await
- {
- eprintln!("Failed to insert shipping address: {}", e);
- let _ = tx.rollback().await;
- // Release reserved stock before returning error
- release_all_reserved_stock(&state, &reserved_items).await;
- return (
- StatusCode::INTERNAL_SERVER_ERROR,
- Json(serde_json::json!({
- "error": "Failed to insert shipping address",
- "details": e.to_string()
- })),
- )
- .into_response();
- }
- // Generate payment reference
- let payment_reference: String = match sqlx::query_scalar("SELECT generate_payment_reference()")
- .fetch_one(&mut *tx)
- .await
- {
- Ok(ref_id) => ref_id,
Err(e) => {
- eprintln!("Failed to generate payment reference: {}", e);
- let _ = tx.rollback().await;
- // Release reserved stock before returning error
+ // Transaction failed (rolled back automatically)
+ // Release reserved stock
release_all_reserved_stock(&state, &reserved_items).await;
- return (
+
+ let error_msg = match e {
+ OrderError::Database(db_err) => {
+ tracing::error!(error = %db_err, "Database error during order creation");
+ format!("Database error: {}", db_err)
+ }
+ };
+
+ (
StatusCode::INTERNAL_SERVER_ERROR,
Json(serde_json::json!({
- "error": "Failed to generate payment reference",
- "details": e.to_string()
+ "error": "Failed to create order",
+ "details": error_msg
})),
)
- .into_response();
+ .into_response()
}
- };
-
- // Insert payment
- // NOTE: Payment is simulated - in production, this would call a payment gateway
- // If payment fails, we need to release reserved stock
- if let Err(e) = sqlx::query(
- r#"
- INSERT INTO payments (
- order_id, payment_method, amount, status,
- payment_reference, card_last4, card_brand, processed_at
- )
- VALUES ($1, $2, $3, 'paid', $4, $5, $6, CURRENT_TIMESTAMP)
- "#,
- )
- .bind(order_id)
- .bind(&request.payment.payment_method)
- .bind(&total)
- .bind(&payment_reference)
- .bind(&request.payment.card_last4)
- .bind(&request.payment.card_brand)
- .execute(&mut *tx)
- .await
- {
- eprintln!("Failed to insert payment: {}", e);
- let _ = tx.rollback().await;
- // Payment failed - release reserved stock
- release_all_reserved_stock(&state, &reserved_items).await;
- return (
- StatusCode::INTERNAL_SERVER_ERROR,
- Json(serde_json::json!({
- "error": "Failed to process payment",
- "details": e.to_string()
- })),
- )
- .into_response();
- }
-
- // Update order payment status
- if let Err(e) = sqlx::query(
- r#"
- UPDATE orders
- SET payment_status = 'paid', status = 'processing'
- WHERE id = $1
- "#,
- )
- .bind(order_id)
- .execute(&mut *tx)
- .await
- {
- eprintln!("Failed to update order status: {}", e);
- let _ = tx.rollback().await;
- // Release reserved stock before returning error
- release_all_reserved_stock(&state, &reserved_items).await;
- return (
- StatusCode::INTERNAL_SERVER_ERROR,
- Json(serde_json::json!({
- "error": "Failed to update order status",
- "details": e.to_string()
- })),
- )
- .into_response();
}
-
- // Commit transaction
- if let Err(e) = tx.commit().await {
- eprintln!("Failed to commit transaction: {}", e);
- // Transaction commit failed - release reserved stock
- release_all_reserved_stock(&state, &reserved_items).await;
- return (
- StatusCode::INTERNAL_SERVER_ERROR,
- Json(serde_json::json!({
- "error": "Failed to commit transaction",
- "details": e.to_string()
- })),
- )
- .into_response();
- }
-
- // Transaction committed successfully!
- // Now confirm the sale with inventory service to convert reserved stock to sold
- confirm_all_sales(&state, order_uuid, &reserved_items).await;
-
- (
- StatusCode::CREATED,
- Json(serde_json::json!({
- "success": true,
- "order_number": order_number,
- "order_uuid": order_uuid,
- "message": "Order created successfully"
- })),
- )
- .into_response()
}
/// List orders with pagination and filtering
@@ -711,14 +555,6 @@ pub async fn create_order(
/// Supports two modes:
/// 1. Authenticated user: Requires X-User-Email header, returns all orders for that user
/// 2. Guest user: Requires both email and order_number query params, returns single order
-#[instrument(
- name = "list_orders",
- skip(state, headers),
- fields(
- page = params.page,
- page_size = params.page_size
- )
-)]
pub async fn list_orders(
State(state): State<AppState>,
headers: axum::http::HeaderMap,
@@ -758,84 +594,42 @@ async fn list_user_orders(
) -> axum::response::Response {
let page = params.page.unwrap_or(1).max(1);
let page_size = params.page_size.unwrap_or(20).clamp(1, 100);
- let offset = (page - 1) * page_size;
-
- let mut query = String::from("SELECT * FROM v_orders_complete WHERE customer_email = $1");
- let mut count_query =
- String::from("SELECT COUNT(*) FROM v_orders_complete WHERE customer_email = $1");
-
- // Apply additional filters if provided
- if let Some(ref status) = params.status {
- let filter = format!(" AND status = '{}'", status.replace('\'', "''"));
- query.push_str(&filter);
- count_query.push_str(&filter);
- }
+ let result = db::list_orders_by_email(
+ pool,
+ user_email,
+ params.status.as_deref(),
+ params.payment_status.as_deref(),
+ page,
+ page_size,
+ )
+ .await;
- if let Some(ref payment_status) = params.payment_status {
- let filter = format!(
- " AND payment_status = '{}'",
- payment_status.replace('\'', "''")
- );
- query.push_str(&filter);
- count_query.push_str(&filter);
- }
+ match result {
+ Ok((orders, total_count)) => {
+ let total_pages = ((total_count as f64) / (page_size as f64)).ceil() as i32;
- query.push_str(&format!(
- " ORDER BY ordered_at DESC LIMIT {} OFFSET {}",
- page_size, offset
- ));
+ let response = OrdersResponse {
+ orders,
+ total_count,
+ page,
+ page_size,
+ total_pages,
+ };
- // Execute count query
- let total_count: (i64,) = match sqlx::query_as(&count_query)
- .bind(user_email)
- .fetch_one(pool)
- .await
- {
- Ok(count) => count,
- Err(e) => {
- eprintln!("Database error counting orders: {}", e);
- return (
- StatusCode::INTERNAL_SERVER_ERROR,
- Json(serde_json::json!({
- "error": "Failed to count orders",
- "details": e.to_string()
- })),
- )
- .into_response();
+ (StatusCode::OK, Json(response)).into_response()
}
- };
-
- // Execute main query
- let orders: Vec<OrderComplete> = match sqlx::query_as(&query)
- .bind(user_email)
- .fetch_all(pool)
- .await
- {
- Ok(orders) => orders,
Err(e) => {
- eprintln!("Database error fetching orders: {}", e);
- return (
+ tracing::error!(error = %e, "Database error fetching orders");
+ (
StatusCode::INTERNAL_SERVER_ERROR,
Json(serde_json::json!({
"error": "Failed to fetch orders",
"details": e.to_string()
})),
)
- .into_response();
+ .into_response()
}
- };
-
- let total_pages = ((total_count.0 as f64) / (page_size as f64)).ceil() as i32;
-
- let response = OrdersResponse {
- orders,
- total_count: total_count.0,
- page,
- page_size,
- total_pages,
- };
-
- (StatusCode::OK, Json(response)).into_response()
+ }
}
/// Get a single order for a guest user
@@ -844,16 +638,7 @@ async fn get_guest_order(
email: &str,
order_number: &str,
) -> axum::response::Response {
- let result = sqlx::query_as::<_, OrderComplete>(
- r#"
- SELECT * FROM v_orders_complete
- WHERE customer_email = $1 AND order_number = $2
- "#,
- )
- .bind(email)
- .bind(order_number)
- .fetch_optional(pool)
- .await;
+ let result = db::get_order_by_email_and_number(pool, email, order_number).await;
match result {
Ok(Some(order)) => {
@@ -875,7 +660,7 @@ async fn get_guest_order(
)
.into_response(),
Err(e) => {
- eprintln!("Database error fetching order: {}", e);
+ tracing::error!(error = %e, "Database error fetching order");
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(serde_json::json!({
@@ -896,14 +681,7 @@ pub async fn get_order_by_id(
State(state): State<AppState>,
Path(uuid): Path<Uuid>,
) -> impl IntoResponse {
- let result = sqlx::query_as::<_, OrderComplete>(
- r#"
- SELECT * FROM v_orders_complete WHERE eid = $1
- "#,
- )
- .bind(uuid)
- .fetch_optional(state.pool())
- .await;
+ let result = db::get_order_by_uuid(state.pool(), uuid).await;
match result {
Ok(Some(order)) => (StatusCode::OK, Json(order)).into_response(),
@@ -916,7 +694,7 @@ pub async fn get_order_by_id(
)
.into_response(),
Err(e) => {
- eprintln!("Database error fetching order: {}", e);
+ tracing::error!(error = %e, "Database error fetching order");
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(serde_json::json!({
diff --git a/distributed_observability/orders/src/handlers/shipments.rs b/distributed_observability/orders/src/handlers/shipments.rs
@@ -1,90 +1,72 @@
//! Shipment tracking API handlers
+//!
+//! This module contains the HTTP request handlers for shipment-related endpoints.
+//! Database operations are delegated to the repository layer in `db::shipments`.
use axum::{
extract::{Path, State},
http::StatusCode,
response::{IntoResponse, Json},
};
-use tracing::instrument;
use uuid::Uuid;
-use crate::models::{CreateShipmentRequest, Shipment, UpdateShipmentStatusRequest};
+use crate::db;
+use crate::models::{CreateShipmentRequest, UpdateShipmentStatusRequest};
use crate::AppState;
/// Create a shipment for an order
///
/// # Endpoint
/// `POST /orders/{order_uuid}/shipment`
-#[instrument(name = "create_shipment", skip(state, request), fields(order.uuid = %order_uuid))]
pub async fn create_shipment(
State(state): State<AppState>,
Path(order_uuid): Path<Uuid>,
Json(request): Json<CreateShipmentRequest>,
) -> impl IntoResponse {
- // Get order_id from UUID
- let order_id: Option<i32> = match sqlx::query_scalar("SELECT id FROM orders WHERE uuid = $1")
- .bind(order_uuid)
- .fetch_optional(state.pool())
- .await
- {
- Ok(id) => id,
- Err(e) => {
- eprintln!("Database error fetching order: {}", e);
+ // Look up the internal order ID from the public UUID
+ let order_id = match db::get_order_id_by_uuid(state.pool(), order_uuid).await {
+ Ok(Some(id)) => id,
+ Ok(None) => {
return (
- StatusCode::INTERNAL_SERVER_ERROR,
+ StatusCode::NOT_FOUND,
Json(serde_json::json!({
- "error": "Failed to fetch order",
- "details": e.to_string()
+ "error": "Order not found",
+ "order_uuid": order_uuid.to_string()
})),
)
.into_response();
}
- };
-
- let order_id = match order_id {
- Some(id) => id,
- None => {
+ Err(e) => {
+ tracing::error!(error = %e, "Database error fetching order");
return (
- StatusCode::NOT_FOUND,
+ StatusCode::INTERNAL_SERVER_ERROR,
Json(serde_json::json!({
- "error": "Order not found",
- "order_uuid": order_uuid.to_string()
+ "error": "Failed to fetch order",
+ "details": e.to_string()
})),
)
.into_response();
}
};
- // Create shipment
- let result = sqlx::query_as::<_, Shipment>(
- r#"
- INSERT INTO shipments (
- order_id, carrier, tracking_number,
- estimated_delivery_date, status, shipped_at
- )
- VALUES ($1, $2, $3, $4, 'shipped', CURRENT_TIMESTAMP)
- RETURNING *
- "#,
+ // Delegate shipment creation to the repository layer
+ match db::create_shipment(
+ state.pool(),
+ order_id,
+ &request.carrier,
+ request.tracking_number.as_deref(),
+ request.estimated_delivery_date,
)
- .bind(order_id)
- .bind(&request.carrier)
- .bind(&request.tracking_number)
- .bind(request.estimated_delivery_date)
- .fetch_one(state.pool())
- .await;
-
- match result {
+ .await
+ {
Ok(shipment) => {
// Update order status to shipped
- let _ = sqlx::query("UPDATE orders SET status = 'shipped' WHERE id = $1")
- .bind(order_id)
- .execute(state.pool())
- .await;
+ let _ = db::update_order_status(state.pool(), order_id, "shipped").await;
(StatusCode::CREATED, Json(shipment)).into_response()
}
Err(e) => {
- eprintln!("Database error creating shipment: {}", e);
+ tracing::error!(error = %e, "Database error creating shipment");
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(serde_json::json!({
@@ -101,17 +83,16 @@ pub async fn create_shipment(
///
/// # Endpoint
/// `PUT /orders/{order_uuid}/shipment/status`
-#[instrument(name = "update_shipment_status", skip(state, request), fields(order.uuid = %order_uuid))]
pub async fn update_shipment_status(
State(state): State<AppState>,
Path(order_uuid): Path<Uuid>,
Json(request): Json<UpdateShipmentStatusRequest>,
) -> impl IntoResponse {
- // Start transaction
+ // Start a transaction for the multi-step update
let mut tx = match state.pool().begin().await {
Ok(tx) => tx,
Err(e) => {
- eprintln!("Failed to start transaction: {}", e);
+ tracing::error!(error = %e, "Failed to start transaction");
return (
StatusCode::INTERNAL_SERVER_ERROR,
Json(serde_json::json!({
@@ -123,82 +104,51 @@ pub async fn update_shipment_status(
}
};
- // Get order_id
- let order_id: Option<i32> = match sqlx::query_scalar("SELECT id FROM orders WHERE uuid = $1")
- .bind(order_uuid)
- .fetch_optional(&mut *tx)
- .await
- {
- Ok(id) => id,
- Err(e) => {
- eprintln!("Database error fetching order: {}", e);
+ // Look up the internal order ID within the transaction
+ let order_id = match db::get_order_id_by_uuid_in_tx(&mut tx, order_uuid).await {
+ Ok(Some(id)) => id,
+ Ok(None) => {
let _ = tx.rollback().await;
return (
- StatusCode::INTERNAL_SERVER_ERROR,
+ StatusCode::NOT_FOUND,
Json(serde_json::json!({
- "error": "Failed to fetch order",
- "details": e.to_string()
+ "error": "Order not found",
+ "order_uuid": order_uuid.to_string()
})),
)
.into_response();
}
- };
-
- let order_id = match order_id {
- Some(id) => id,
- None => {
+ Err(e) => {
+ tracing::error!(error = %e, "Database error fetching order");
let _ = tx.rollback().await;
return (
- StatusCode::NOT_FOUND,
+ StatusCode::INTERNAL_SERVER_ERROR,
Json(serde_json::json!({
- "error": "Order not found",
- "order_uuid": order_uuid.to_string()
+ "error": "Failed to fetch order",
+ "details": e.to_string()
})),
)
.into_response();
}
};
- // Update shipment
- let query = if request.status == "delivered" {
- sqlx::query(
- r#"
- UPDATE shipments
- SET status = $1, actual_delivery_date = $2, delivered_at = CURRENT_TIMESTAMP
- WHERE order_id = $3
- RETURNING id
- "#,
- )
- .bind(&request.status)
- .bind(request.actual_delivery_date)
- .bind(order_id)
- } else {
- sqlx::query(
- r#"
- UPDATE shipments
- SET status = $1
- WHERE order_id = $2
- RETURNING id
- "#,
- )
- .bind(&request.status)
- .bind(order_id)
- };
-
- let result = query.fetch_optional(&mut *tx).await;
-
- match result {
- Ok(Some(_)) => {
- // Update order status if delivered
+ // Delegate shipment status update to the repository layer
+ match db::update_shipment_status(
+ &mut tx,
+ order_id,
+ &request.status,
+ request.actual_delivery_date,
+ )
+ .await
+ {
+ Ok(true) => {
+ // If delivered, also update the order status
if request.status == "delivered" {
- let _ = sqlx::query("UPDATE orders SET status = 'delivered' WHERE id = $1")
- .bind(order_id)
- .execute(&mut *tx)
- .await;
+ let _ = db::update_order_status_in_tx(&mut tx, order_id, "delivered").await;
}
if let Err(e) = tx.commit().await {
- eprintln!("Failed to commit transaction: {}", e);
+ tracing::error!(error = %e, "Failed to commit transaction");
return (
StatusCode::INTERNAL_SERVER_ERROR,
Json(serde_json::json!({
@@ -219,7 +169,7 @@ pub async fn update_shipment_status(
)
.into_response()
}
- Ok(None) => {
+ Ok(false) => {
let _ = tx.rollback().await;
(
StatusCode::NOT_FOUND,
@@ -231,7 +181,7 @@ pub async fn update_shipment_status(
.into_response()
}
Err(e) => {
- eprintln!("Database error updating shipment: {}", e);
+ tracing::error!(error = %e, "Database error updating shipment");
let _ = tx.rollback().await;
(
StatusCode::INTERNAL_SERVER_ERROR,
diff --git a/distributed_observability/orders/src/models/order.rs b/distributed_observability/orders/src/models/order.rs
@@ -8,6 +8,7 @@ use uuid::Uuid;
/// Complete order with all related data (from v_orders_complete view)
#[derive(Debug, Serialize, FromRow)]
+#[allow(dead_code)] // id field needed for FromRow but skipped in JSON
pub struct OrderComplete {
#[serde(skip)]
pub id: i32,
@@ -42,7 +43,7 @@ pub struct OrderComplete {
}
/// Request to create a new order
-#[derive(Debug, Deserialize)]
+#[derive(Debug, Clone, Deserialize)]
pub struct CreateOrderRequest {
pub customer_email: String,
pub customer_phone: Option<String>,
@@ -52,7 +53,7 @@ pub struct CreateOrderRequest {
}
/// Order item in create request
-#[derive(Debug, Deserialize)]
+#[derive(Debug, Clone, Deserialize)]
pub struct CreateOrderItemRequest {
pub product_uuid: Uuid,
pub product_name: String,
@@ -62,7 +63,7 @@ pub struct CreateOrderItemRequest {
}
/// Shipping address in create request
-#[derive(Debug, Deserialize)]
+#[derive(Debug, Clone, Deserialize)]
pub struct CreateShippingAddressRequest {
pub first_name: String,
pub last_name: String,
@@ -76,7 +77,7 @@ pub struct CreateShippingAddressRequest {
}
/// Payment info in create request
-#[derive(Debug, Deserialize)]
+#[derive(Debug, Clone, Deserialize)]
pub struct CreatePaymentRequest {
pub payment_method: String,
pub card_last4: Option<String>,
@@ -84,7 +85,7 @@ pub struct CreatePaymentRequest {
}
/// Query parameters for listing orders
-#[derive(Debug, Deserialize)]
+#[derive(Debug, Clone, Deserialize)]
pub struct OrderQueryParams {
pub page: Option<i32>,
pub page_size: Option<i32>,
diff --git a/distributed_observability/otelmart/src/db/mod.rs b/distributed_observability/otelmart/src/db/mod.rs
@@ -1,5 +1,9 @@
use anyhow::Result;
-use sqlx::{PgPool, postgres::PgPoolOptions};
+use sqlx::{postgres::PgPoolOptions, PgPool};
+
+pub mod users;
+
+pub use users::*;
#[derive(Clone)]
pub struct Database {
diff --git a/distributed_observability/otelmart/src/db/users.rs b/distributed_observability/otelmart/src/db/users.rs
@@ -0,0 +1,392 @@
+//! User profile and address repository functions with OpenTelemetry instrumentation
+//!
+//! This module contains database operations for user profiles and addresses,
+//! separated from HTTP handlers. Each function is instrumented with
+//! OpenTelemetry semantic conventions for database spans.
+
+use chrono::NaiveDate;
+use sqlx::PgPool;
+use tracing::{instrument, Span};
+use uuid::Uuid;
+
+use crate::models::{AddAddressRequest, UserAddress};
+
+/// Profile data returned from the database
+#[derive(Debug)]
+pub struct ProfileData {
+ pub eid: Uuid,
+ pub avatar_url: Option<String>,
+ pub date_of_birth: Option<NaiveDate>,
+ pub email_notifications: bool,
+ pub marketing_emails: bool,
+}
+
+/// Get user profile by user ID
+#[instrument(
+ name = "SELECT user_profiles",
+ skip(pool),
+ fields(
+ otel.kind = "client",
+ db.system.name = "postgresql",
+ db.namespace = "users",
+ db.operation.name = "SELECT",
+ db.collection.name = "user_profiles",
+ db.query.text = "SELECT uuid, avatar_url, date_of_birth, ... FROM user_profiles WHERE user_id = $1",
+ otelmart.user.id = user_id
+ )
+)]
+pub async fn get_profile(pool: &PgPool, user_id: i32) -> Result<Option<ProfileData>, sqlx::Error> {
+ let row = sqlx::query_as::<_, (Uuid, Option<String>, Option<NaiveDate>, bool, bool)>(
+ r#"
+ SELECT
+ uuid,
+ avatar_url,
+ date_of_birth,
+ email_notifications,
+ marketing_emails
+ FROM user_profiles
+ WHERE user_id = $1
+ "#,
+ )
+ .bind(user_id)
+ .fetch_optional(pool)
+ .await?;
+
+ Ok(row.map(
+ |(eid, avatar_url, date_of_birth, email_notifications, marketing_emails)| ProfileData {
+ eid,
+ avatar_url,
+ date_of_birth,
+ email_notifications,
+ marketing_emails,
+ },
+ ))
+}
+
+/// Check if a profile exists for a user
+#[instrument(
+ name = "SELECT user_profiles",
+ skip(pool),
+ fields(
+ otel.kind = "client",
+ db.system.name = "postgresql",
+ db.namespace = "users",
+ db.operation.name = "SELECT",
+ db.collection.name = "user_profiles",
+ db.query.text = "SELECT id FROM user_profiles WHERE user_id = $1",
+ otelmart.user.id = user_id
+ )
+)]
+pub async fn profile_exists(pool: &PgPool, user_id: i32) -> Result<bool, sqlx::Error> {
+ let row = sqlx::query_as::<_, (i32,)>("SELECT id FROM user_profiles WHERE user_id = $1")
+ .bind(user_id)
+ .fetch_optional(pool)
+ .await?;
+
+ Ok(row.is_some())
+}
+
+/// Update an existing user profile
+#[instrument(
+ name = "UPDATE user_profiles",
+ skip(pool),
+ fields(
+ otel.kind = "client",
+ db.system.name = "postgresql",
+ db.namespace = "users",
+ db.operation.name = "UPDATE",
+ db.collection.name = "user_profiles",
+ db.query.text = "UPDATE user_profiles SET ... WHERE user_id = $5 RETURNING ...",
+ otelmart.user.id = user_id,
+ db.response.returned_rows = tracing::field::Empty
+ )
+)]
+pub async fn update_profile(
+ pool: &PgPool,
+ user_id: i32,
+ avatar_url: Option<&str>,
+ date_of_birth: Option<NaiveDate>,
+ email_notifications: Option<bool>,
+ marketing_emails: Option<bool>,
+) -> Result<ProfileData, sqlx::Error> {
+ let (eid, avatar, dob, email_notif, marketing) =
+ sqlx::query_as::<_, (Uuid, Option<String>, Option<NaiveDate>, bool, bool)>(
+ r#"
+ UPDATE user_profiles
+ SET avatar_url = COALESCE($1, avatar_url),
+ date_of_birth = COALESCE($2, date_of_birth),
+ email_notifications = COALESCE($3, email_notifications),
+ marketing_emails = COALESCE($4, marketing_emails),
+ updated_at = NOW()
+ WHERE user_id = $5
+ RETURNING
+ uuid,
+ avatar_url,
+ date_of_birth,
+ email_notifications,
+ marketing_emails
+ "#,
+ )
+ .bind(avatar_url)
+ .bind(date_of_birth)
+ .bind(email_notifications)
+ .bind(marketing_emails)
+ .bind(user_id)
+ .fetch_one(pool)
+ .await?;
+
+ Span::current().record("db.response.returned_rows", 1);
+
+ Ok(ProfileData {
+ eid,
+ avatar_url: avatar,
+ date_of_birth: dob,
+ email_notifications: email_notif,
+ marketing_emails: marketing,
+ })
+}
+
+/// Create a new user profile
+#[instrument(
+ name = "INSERT user_profiles",
+ skip(pool),
+ fields(
+ db.system.name = "postgresql",
+ db.namespace = "users",
+ db.operation.name = "INSERT",
+ db.collection.name = "user_profiles",
+ db.query.text = "INSERT INTO user_profiles (...) VALUES (...) RETURNING ...",
+ otelmart.user.id = user_id,
+ db.response.returned_rows = tracing::field::Empty
+ )
+)]
+pub async fn create_profile(
+ pool: &PgPool,
+ user_id: i32,
+ avatar_url: Option<&str>,
+ date_of_birth: Option<NaiveDate>,
+ email_notifications: bool,
+ marketing_emails: bool,
+) -> Result<ProfileData, sqlx::Error> {
+ let (eid, avatar, dob, email_notif, marketing) =
+ sqlx::query_as::<_, (Uuid, Option<String>, Option<NaiveDate>, bool, bool)>(
+ r#"
+ INSERT INTO user_profiles (
+ user_id,
+ avatar_url,
+ date_of_birth,
+ email_notifications,
+ marketing_emails
+ ) VALUES ($1, $2, $3, $4, $5)
+ RETURNING
+ uuid,
+ avatar_url,
+ date_of_birth,
+ email_notifications,
+ marketing_emails
+ "#,
+ )
+ .bind(user_id)
+ .bind(avatar_url)
+ .bind(date_of_birth)
+ .bind(email_notifications)
+ .bind(marketing_emails)
+ .fetch_one(pool)
+ .await?;
+
+ Span::current().record("db.response.returned_rows", 1);
+
+ Ok(ProfileData {
+ eid,
+ avatar_url: avatar,
+ date_of_birth: dob,
+ email_notifications: email_notif,
+ marketing_emails: marketing,
+ })
+}
+
+/// List active addresses for a user
+#[instrument(
+ name = "SELECT user_addresses",
+ skip(pool),
+ fields(
+ otel.kind = "client",
+ db.system.name = "postgresql",
+ db.namespace = "users",
+ db.operation.name = "SELECT",
+ db.collection.name = "user_addresses",
+ db.query.text = "SELECT * FROM user_addresses WHERE user_id = $1 AND is_active = true ORDER BY ...",
+ otelmart.user.id = user_id,
+ db.response.returned_rows = tracing::field::Empty
+ )
+)]
+pub async fn list_addresses(pool: &PgPool, user_id: i32) -> Result<Vec<UserAddress>, sqlx::Error> {
+ let addresses = sqlx::query_as::<_, UserAddress>(
+ r#"
+ SELECT * FROM user_addresses
+ WHERE user_id = $1 AND is_active = true
+ ORDER BY is_default DESC, created_at DESC
+ "#,
+ )
+ .bind(user_id)
+ .fetch_all(pool)
+ .await?;
+
+ Span::current().record("db.response.returned_rows", addresses.len() as i64);
+ Ok(addresses)
+}
+
+/// Add a new address for a user
+#[instrument(
+ name = "INSERT user_addresses",
+ skip(pool, req),
+ fields(
+ db.system.name = "postgresql",
+ db.namespace = "users",
+ db.operation.name = "INSERT",
+ db.collection.name = "user_addresses",
+ db.query.text = "INSERT INTO user_addresses (...) VALUES (...) RETURNING *",
+ otelmart.user.id = user_id,
+ db.response.returned_rows = tracing::field::Empty
+ )
+)]
+pub async fn add_address(
+ pool: &PgPool,
+ user_id: i32,
+ req: &AddAddressRequest,
+) -> Result<UserAddress, sqlx::Error> {
+ let address = sqlx::query_as::<_, UserAddress>(
+ r#"
+ INSERT INTO user_addresses (
+ user_id, address_type, address_label,
+ first_name, last_name,
+ address_line1, address_line2,
+ city, state, postal_code, country,
+ phone, is_default
+ ) VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13)
+ RETURNING *
+ "#,
+ )
+ .bind(user_id)
+ .bind(&req.address_type)
+ .bind(&req.address_label)
+ .bind(&req.first_name)
+ .bind(&req.last_name)
+ .bind(&req.address_line1)
+ .bind(&req.address_line2)
+ .bind(&req.city)
+ .bind(&req.state)
+ .bind(&req.postal_code)
+ .bind(&req.country)
+ .bind(&req.phone)
+ .bind(req.is_default.unwrap_or(false))
+ .fetch_one(pool)
+ .await?;
+
+ Span::current().record("db.response.returned_rows", 1);
+ Ok(address)
+}
+
+/// Update an existing address (only if it belongs to the user and is active)
+#[instrument(
+ name = "UPDATE user_addresses",
+ skip(pool, req),
+ fields(
+ otel.kind = "client",
+ db.system.name = "postgresql",
+ db.namespace = "users",
+ db.operation.name = "UPDATE",
+ db.collection.name = "user_addresses",
+ db.query.text = "UPDATE user_addresses SET ... WHERE uuid = $13 AND user_id = $14 RETURNING *",
+ otelmart.address.uuid = %address_uuid,
+ otelmart.user.id = user_id,
+ db.response.returned_rows = tracing::field::Empty
+ )
+)]
+pub async fn update_address(
+ pool: &PgPool,
+ user_id: i32,
+ address_uuid: Uuid,
+ req: &AddAddressRequest,
+) -> Result<Option<UserAddress>, sqlx::Error> {
+ let address = sqlx::query_as::<_, UserAddress>(
+ r#"
+ UPDATE user_addresses
+ SET address_type = $1,
+ address_label = $2,
+ first_name = $3,
+ last_name = $4,
+ address_line1 = $5,
+ address_line2 = $6,
+ city = $7,
+ state = $8,
+ postal_code = $9,
+ country = $10,
+ phone = $11,
+ is_default = $12,
+ updated_at = NOW()
+ WHERE uuid = $13 AND user_id = $14 AND is_active = true
+ RETURNING *
+ "#,
+ )
+ .bind(&req.address_type)
+ .bind(&req.address_label)
+ .bind(&req.first_name)
+ .bind(&req.last_name)
+ .bind(&req.address_line1)
+ .bind(&req.address_line2)
+ .bind(&req.city)
+ .bind(&req.state)
+ .bind(&req.postal_code)
+ .bind(&req.country)
+ .bind(&req.phone)
+ .bind(req.is_default.unwrap_or(false))
+ .bind(address_uuid)
+ .bind(user_id)
+ .fetch_optional(pool)
+ .await?;
+
+ Span::current().record(
+ "db.response.returned_rows",
+ if address.is_some() { 1 } else { 0 },
+ );
+ Ok(address)
+}
+
+/// Soft-delete an address (set is_active = false)
+#[instrument(
+ name = "UPDATE user_addresses",
+ skip(pool),
+ fields(
+ otel.kind = "client",
+ db.system.name = "postgresql",
+ db.namespace = "users",
+ db.operation.name = "UPDATE",
+ db.collection.name = "user_addresses",
+ db.query.text = "UPDATE user_addresses SET is_active = false WHERE uuid = $1 AND user_id = $2",
+ otelmart.address.uuid = %address_uuid,
+ otelmart.user.id = user_id,
+ db.response.returned_rows = tracing::field::Empty
+ )
+)]
+pub async fn delete_address(
+ pool: &PgPool,
+ user_id: i32,
+ address_uuid: Uuid,
+) -> Result<u64, sqlx::Error> {
+ let result = sqlx::query(
+ r#"
+ UPDATE user_addresses
+ SET is_active = false, updated_at = NOW()
+ WHERE uuid = $1 AND user_id = $2 AND is_active = true
+ "#,
+ )
+ .bind(address_uuid)
+ .bind(user_id)
+ .execute(pool)
+ .await?;
+
+ let rows = result.rows_affected();
+ Span::current().record("db.response.returned_rows", rows);
+ Ok(rows)
+}
diff --git a/distributed_observability/otelmart/src/handlers/users.rs b/distributed_observability/otelmart/src/handlers/users.rs
@@ -1,21 +1,25 @@
+//! User profile and address API handlers
+//!
+//! This module contains the HTTP request handlers for user profile and address endpoints.
+//! Database operations are delegated to the repository layer in `db::users`.
+
use axum::{
extract::{Path, State},
http::StatusCode,
response::IntoResponse,
Json,
};
-use tracing::instrument;
use uuid::Uuid;
use crate::{
- auth,
- models::{AddAddressRequest, AddressResponse, UpdateProfileRequest, UserAddress},
+ auth, db,
+ models::{AddAddressRequest, AddressResponse, UpdateProfileRequest},
AppState,
};
use super::auth::ErrorResponse;
-#[instrument(name = "get_profile", skip(state, headers))]
+/// Get user profile
pub async fn get_profile(
State(state): State<AppState>,
headers: axum::http::HeaderMap,
@@ -25,40 +29,9 @@ pub async fn get_profile(
if let Some(token) = token {
match auth::validate_session(state.db.pool(), &token).await {
Ok(Some(user)) => {
- // Get user profile (if exists)
- let profile = sqlx::query_as::<
- _,
- (
- uuid::Uuid,
- Option<String>,
- Option<chrono::NaiveDate>,
- bool,
- bool,
- ),
- >(
- r#"
- SELECT
- uuid,
- avatar_url,
- date_of_birth,
- email_notifications,
- marketing_emails
- FROM user_profiles
- WHERE user_id = $1
- "#,
- )
- .bind(user.id)
- .fetch_optional(state.db.pool())
- .await;
-
- match profile {
- Ok(Some((
- eid,
- avatar_url,
- date_of_birth,
- email_notifications,
- marketing_emails,
- ))) => (
+ // Delegate profile lookup to the repository layer
+ match db::users::get_profile(state.db.pool(), user.id).await {
+ Ok(Some(profile)) => (
StatusCode::OK,
Json(serde_json::json!({
"user_id": user.uuid,
@@ -67,11 +40,11 @@ pub async fn get_profile(
"last_name": user.last_name,
"phone": user.phone,
"is_verified": user.is_verified,
- "profile_eid": eid,
- "avatar_url": avatar_url,
- "date_of_birth": date_of_birth,
- "email_notifications": email_notifications,
- "marketing_emails": marketing_emails,
+ "profile_eid": profile.eid,
+ "avatar_url": profile.avatar_url,
+ "date_of_birth": profile.date_of_birth,
+ "email_notifications": profile.email_notifications,
+ "marketing_emails": profile.marketing_emails,
})),
)
.into_response(),
@@ -130,7 +103,7 @@ pub async fn get_profile(
}
}
-#[instrument(name = "update_profile", skip(state, headers, req))]
+/// Update or create user profile
pub async fn update_profile(
State(state): State<AppState>,
headers: axum::http::HeaderMap,
@@ -141,147 +114,82 @@ pub async fn update_profile(
if let Some(token) = token {
match auth::validate_session(state.db.pool(), &token).await {
Ok(Some(user)) => {
- // Check if profile exists
- let profile_exists =
- sqlx::query_as::<_, (i32,)>("SELECT id FROM user_profiles WHERE user_id = $1")
- .bind(user.id)
- .fetch_optional(state.db.pool())
- .await;
-
- match profile_exists {
- Ok(Some(_)) => {
- // Update existing profile
- let result = sqlx::query_as::<
- _,
- (
- uuid::Uuid,
- Option<String>,
- Option<chrono::NaiveDate>,
- bool,
- bool,
- ),
- >(
- r#"
- UPDATE user_profiles
- SET avatar_url = COALESCE($1, avatar_url),
- date_of_birth = COALESCE($2, date_of_birth),
- email_notifications = COALESCE($3, email_notifications),
- marketing_emails = COALESCE($4, marketing_emails),
- updated_at = NOW()
- WHERE user_id = $5
- RETURNING
- uuid,
- avatar_url,
- date_of_birth,
- email_notifications,
- marketing_emails
- "#,
+ // Check if profile exists via repository
+ let exists = match db::users::profile_exists(state.db.pool(), user.id).await {
+ Ok(exists) => exists,
+ Err(e) => {
+ return (
+ StatusCode::INTERNAL_SERVER_ERROR,
+ Json(ErrorResponse {
+ error: e.to_string(),
+ }),
)
- .bind(&req.avatar_url)
- .bind(&req.date_of_birth)
- .bind(&req.email_notifications)
- .bind(&req.marketing_emails)
- .bind(user.id)
- .fetch_one(state.db.pool())
- .await;
-
- match result {
- Ok((
- eid,
- avatar_url,
- date_of_birth,
- email_notifications,
- marketing_emails,
- )) => (
- StatusCode::OK,
- Json(serde_json::json!({
- "eid": eid,
- "avatar_url": avatar_url,
- "date_of_birth": date_of_birth,
- "email_notifications": email_notifications,
- "marketing_emails": marketing_emails,
- })),
- )
- .into_response(),
- Err(e) => (
- StatusCode::INTERNAL_SERVER_ERROR,
- Json(ErrorResponse {
- error: e.to_string(),
- }),
- )
- .into_response(),
- }
+ .into_response();
}
- Ok(None) => {
- // Create new profile
- let result = sqlx::query_as::<
- _,
- (
- uuid::Uuid,
- Option<String>,
- Option<chrono::NaiveDate>,
- bool,
- bool,
- ),
- >(
- r#"
- INSERT INTO user_profiles (
- user_id,
- avatar_url,
- date_of_birth,
- email_notifications,
- marketing_emails
- ) VALUES ($1, $2, $3, $4, $5)
- RETURNING
- uuid,
- avatar_url,
- date_of_birth,
- email_notifications,
- marketing_emails
- "#,
- )
- .bind(user.id)
- .bind(&req.avatar_url)
- .bind(&req.date_of_birth)
- .bind(req.email_notifications.unwrap_or(true))
- .bind(req.marketing_emails.unwrap_or(false))
- .fetch_one(state.db.pool())
- .await;
+ };
- match result {
- Ok((
- eid,
- avatar_url,
- date_of_birth,
- email_notifications,
- marketing_emails,
- )) => (
- StatusCode::CREATED,
- Json(serde_json::json!({
- "eid": eid,
- "avatar_url": avatar_url,
- "date_of_birth": date_of_birth,
- "email_notifications": email_notifications,
- "marketing_emails": marketing_emails,
- })),
- )
- .into_response(),
- Err(e) => (
- StatusCode::INTERNAL_SERVER_ERROR,
- Json(ErrorResponse {
- error: e.to_string(),
- }),
- )
- .into_response(),
- }
+ if exists {
+ // Update existing profile via repository
+ match db::users::update_profile(
+ state.db.pool(),
+ user.id,
+ req.avatar_url.as_deref(),
+ req.date_of_birth,
+ req.email_notifications,
+ req.marketing_emails,
+ )
+ .await
+ {
+ Ok(profile) => (
+ StatusCode::OK,
+ Json(serde_json::json!({
+ "eid": profile.eid,
+ "avatar_url": profile.avatar_url,
+ "date_of_birth": profile.date_of_birth,
+ "email_notifications": profile.email_notifications,
+ "marketing_emails": profile.marketing_emails,
+ })),
+ )
+ .into_response(),
+ Err(e) => (
+ StatusCode::INTERNAL_SERVER_ERROR,
+ Json(ErrorResponse {
+ error: e.to_string(),
+ }),
+ )
+ .into_response(),
}
- Err(e) => (
- StatusCode::INTERNAL_SERVER_ERROR,
- Json(ErrorResponse {
- error: e.to_string(),
- }),
+ } else {
+ // Create new profile via repository
+ match db::users::create_profile(
+ state.db.pool(),
+ user.id,
+ req.avatar_url.as_deref(),
+ req.date_of_birth,
+ req.email_notifications.unwrap_or(true),
+ req.marketing_emails.unwrap_or(false),
)
- .into_response(),
+ .await
+ {
+ Ok(profile) => (
+ StatusCode::CREATED,
+ Json(serde_json::json!({
+ "eid": profile.eid,
+ "avatar_url": profile.avatar_url,
+ "date_of_birth": profile.date_of_birth,
+ "email_notifications": profile.email_notifications,
+ "marketing_emails": profile.marketing_emails,
+ })),
+ )
+ .into_response(),
+ Err(e) => (
+ StatusCode::INTERNAL_SERVER_ERROR,
+ Json(ErrorResponse {
+ error: e.to_string(),
+ }),
+ )
+ .into_response(),
+ }
}
}
Ok(None) => (
@@ -310,7 +218,7 @@ pub async fn update_profile(
}
}
-#[instrument(name = "get_addresses", skip(state, headers))]
+/// Get all user addresses
pub async fn get_addresses(
State(state): State<AppState>,
headers: axum::http::HeaderMap,
@@ -320,18 +228,8 @@ pub async fn get_addresses(
if let Some(token) = token {
match auth::validate_session(state.db.pool(), &token).await {
Ok(Some(user)) => {
- let addresses = sqlx::query_as::<_, UserAddress>(
- r#"
- SELECT * FROM user_addresses
- WHERE user_id = $1 AND is_active = true
- ORDER BY is_default DESC, created_at DESC
- "#,
- )
- .bind(user.id)
- .fetch_all(state.db.pool())
- .await;
-
- match addresses {
+ // Delegate address listing to the repository layer
+ match db::users::list_addresses(state.db.pool(), user.id).await {
Ok(addresses) => {
let response: Vec<AddressResponse> =
addresses.into_iter().map(AddressResponse::from).collect();
@@ -373,7 +271,7 @@ pub async fn get_addresses(
}
}
-#[instrument(name = "add_address", skip(state, headers, req))]
+/// Add a new user address
pub async fn add_address(
State(state): State<AppState>,
headers: axum::http::HeaderMap,
@@ -384,35 +282,8 @@ pub async fn add_address(
if let Some(token) = token {
match auth::validate_session(state.db.pool(), &token).await {
Ok(Some(user)) => {
- let address = sqlx::query_as::<_, UserAddress>(
- r#"
- INSERT INTO user_addresses (
- user_id, address_type, address_label,
- first_name, last_name,
- address_line1, address_line2,
- city, state, postal_code, country,
- phone, is_default
- ) VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13)
- RETURNING *
- "#,
- )
- .bind(user.id)
- .bind(&req.address_type)
- .bind(&req.address_label)
- .bind(&req.first_name)
- .bind(&req.last_name)
- .bind(&req.address_line1)
- .bind(&req.address_line2)
- .bind(&req.city)
- .bind(&req.state)
- .bind(&req.postal_code)
- .bind(&req.country)
- .bind(&req.phone)
- .bind(req.is_default.unwrap_or(false))
- .fetch_one(state.db.pool())
- .await;
-
- match address {
+ // Delegate address creation to the repository layer
+ match db::users::add_address(state.db.pool(), user.id, &req).await {
Ok(address) => {
let response = AddressResponse::from(address);
(StatusCode::CREATED, Json(response)).into_response()
@@ -452,7 +323,7 @@ pub async fn add_address(
}
}
-#[instrument(name = "update_address", skip(state, headers, req), fields(address.uuid = %id))]
+/// Update an existing user address
pub async fn update_address(
State(state): State<AppState>,
headers: axum::http::HeaderMap,
@@ -464,45 +335,8 @@ pub async fn update_address(
if let Some(token) = token {
match auth::validate_session(state.db.pool(), &token).await {
Ok(Some(user)) => {
- // Update address only if it belongs to the user
- let address = sqlx::query_as::<_, UserAddress>(
- r#"
- UPDATE user_addresses
- SET address_type = $1,
- address_label = $2,
- first_name = $3,
- last_name = $4,
- address_line1 = $5,
- address_line2 = $6,
- city = $7,
- state = $8,
- postal_code = $9,
- country = $10,
- phone = $11,
- is_default = $12,
- updated_at = NOW()
- WHERE uuid = $13 AND user_id = $14 AND is_active = true
- RETURNING *
- "#,
- )
- .bind(&req.address_type)
- .bind(&req.address_label)
- .bind(&req.first_name)
- .bind(&req.last_name)
- .bind(&req.address_line1)
- .bind(&req.address_line2)
- .bind(&req.city)
- .bind(&req.state)
- .bind(&req.postal_code)
- .bind(&req.country)
- .bind(&req.phone)
- .bind(req.is_default.unwrap_or(false))
- .bind(id)
- .bind(user.id)
- .fetch_optional(state.db.pool())
- .await;
-
- match address {
+ // Delegate address update to the repository layer
+ match db::users::update_address(state.db.pool(), user.id, id, &req).await {
Ok(Some(address)) => {
let response = AddressResponse::from(address);
(StatusCode::OK, Json(response)).into_response()
@@ -549,7 +383,7 @@ pub async fn update_address(
}
}
-#[instrument(name = "delete_address", skip(state, headers), fields(address.uuid = %id))]
+/// Soft delete a user address
pub async fn delete_address(
State(state): State<AppState>,
headers: axum::http::HeaderMap,
@@ -560,22 +394,10 @@ pub async fn delete_address(
if let Some(token) = token {
match auth::validate_session(state.db.pool(), &token).await {
Ok(Some(user)) => {
- // Soft delete: set is_active to false
- let result = sqlx::query(
- r#"
- UPDATE user_addresses
- SET is_active = false, updated_at = NOW()
- WHERE uuid = $1 AND user_id = $2 AND is_active = true
- "#,
- )
- .bind(id)
- .bind(user.id)
- .execute(state.db.pool())
- .await;
-
- match result {
- Ok(result) => {
- if result.rows_affected() > 0 {
+ // Delegate soft-delete to the repository layer
+ match db::users::delete_address(state.db.pool(), user.id, id).await {
+ Ok(rows) => {
+ if rows > 0 {
(
StatusCode::OK,
Json(serde_json::json!({
diff --git a/distributed_observability/products/src/db/mod.rs b/distributed_observability/products/src/db/mod.rs
@@ -21,6 +21,13 @@ use sqlx::{
};
use std::str::FromStr;
+pub mod ratings;
+pub mod repository;
+
+// Re-export all repository functions for easy access
+pub use ratings::*;
+pub use repository::*;
+
/// Database connection pool wrapper
///
/// Wraps sqlx's PgPool with custom configuration for the products service.
diff --git a/distributed_observability/products/src/db/ratings.rs b/distributed_observability/products/src/db/ratings.rs
@@ -0,0 +1,76 @@
+//! Rating repository functions with OpenTelemetry instrumentation
+//!
+//! This module contains database operations for product ratings,
+//! separated from HTTP handlers.
+
+use sqlx::PgPool;
+use tracing::instrument;
+use uuid::Uuid;
+
+use crate::models::{Rating, UpsertRatingRequest};
+
+/// Check if a product exists and is not deleted
+#[instrument(
+ name = "SELECT products",
+ skip(pool),
+ fields(
+ db.system.name = "postgresql",
+ db.namespace = "products",
+ db.operation.name = "SELECT",
+ db.collection.name = "products",
+ db.query.text = "SELECT EXISTS(SELECT 1 FROM products WHERE uuid = $1 AND deleted_at IS NULL)",
+ otelmart.product.uuid = %product_uuid
+ )
+)]
+pub async fn product_exists(pool: &PgPool, product_uuid: Uuid) -> Result<bool, sqlx::Error> {
+ let result: Option<(bool,)> = sqlx::query_as(
+ "SELECT EXISTS(SELECT 1 FROM products WHERE uuid = $1 AND deleted_at IS NULL)",
+ )
+ .bind(product_uuid)
+ .fetch_optional(pool)
+ .await?;
+
+ Ok(result.map(|x| x.0).unwrap_or(false))
+}
+
+/// Create or update a product rating
+///
+/// Uses PostgreSQL's ON CONFLICT clause for atomic upsert.
+/// The unique constraint on (product_id, user_id) ensures one rating per user per product.
+#[instrument(
+ name = "UPSERT ratings",
+ skip(pool, payload),
+ fields(
+ db.system.name = "postgresql",
+ db.namespace = "products",
+ db.operation.name = "INSERT",
+ db.collection.name = "ratings",
+ db.query.text = "INSERT INTO ratings (...) VALUES (...) ON CONFLICT DO UPDATE ... RETURNING *",
+ otelmart.product.uuid = %product_uuid,
+ otelmart.rating.value = payload.rating
+ )
+)]
+pub async fn upsert_rating(
+ pool: &PgPool,
+ product_uuid: Uuid,
+ payload: &UpsertRatingRequest,
+) -> Result<Rating, sqlx::Error> {
+ sqlx::query_as::<_, Rating>(
+ r#"
+ INSERT INTO ratings (product_id, user_id, rating, review)
+ VALUES ($1, $2, $3, $4)
+ ON CONFLICT (product_id, user_id)
+ DO UPDATE SET
+ rating = EXCLUDED.rating,
+ review = EXCLUDED.review,
+ updated_at = NOW()
+ RETURNING id, uuid, product_id, user_id, rating, review, created_at, updated_at
+ "#,
+ )
+ .bind(product_uuid)
+ .bind(payload.user_id)
+ .bind(payload.rating)
+ .bind(&payload.review)
+ .fetch_one(pool)
+ .await
+}
diff --git a/distributed_observability/products/src/db/repository.rs b/distributed_observability/products/src/db/repository.rs
@@ -0,0 +1,300 @@
+//! Product repository functions with OpenTelemetry instrumentation
+//!
+//! This module contains database operations for products, separated from HTTP handlers.
+//! Each function is instrumented with OpenTelemetry semantic conventions for database spans.
+
+use sqlx::{PgPool, Postgres, QueryBuilder};
+use tracing::instrument;
+use uuid::Uuid;
+
+use crate::models::{ProductDetail, ProductQueryParams, ProductWithRating};
+
+/// Fetch a single product by UUID with full details
+///
+/// Joins with categories and ratings tables to return complete product information.
+/// Returns None if the product doesn't exist or has been deleted.
+#[instrument(
+ name = "SELECT products",
+ skip(pool),
+ fields(
+ db.system.name = "postgresql",
+ db.namespace = "products",
+ db.operation.name = "SELECT",
+ db.collection.name = "products",
+ db.query.text = "SELECT p.*, c.*, AVG(r.rating), COUNT(r.id) FROM products p LEFT JOIN ... WHERE p.uuid = $1",
+ otelmart.product.uuid = %uuid
+ )
+)]
+pub async fn get_product_by_uuid(
+ pool: &PgPool,
+ uuid: Uuid,
+) -> Result<Option<ProductDetail>, sqlx::Error> {
+ sqlx::query_as::<_, ProductDetail>(
+ r#"
+ SELECT
+ p.id,
+ p.uuid,
+ p.asin,
+ p.sku,
+ p.gtin,
+ p.product_name,
+ p.brand,
+ p.description,
+ p.url,
+ p.price,
+ p.initial_price,
+ p.discount,
+ p.currency,
+ p.stock_quantity,
+ p.sizes,
+ p.colors,
+ p.image_url,
+ p.available_for_delivery,
+ p.available_for_pickup,
+ p.free_returns,
+ p.is_active,
+ p.created_at,
+ p.updated_at,
+ p.data_timestamp,
+ p.category_id,
+ c.name as category_name,
+ c.uuid as category_uuid,
+ c.slug as category_slug,
+ COALESCE(AVG(r.rating)::FLOAT8, NULL) as average_rating,
+ COUNT(r.id)::BIGINT as rating_count,
+ CAST(p.id AS TEXT) as product_id_str,
+ get_root_category_name(p.category_id) as root_category_name,
+ p.deleted_at
+ FROM products p
+ LEFT JOIN categories c ON p.category_id = c.id
+ LEFT JOIN ratings r ON p.id = r.product_id
+ WHERE p.uuid = $1 AND p.deleted_at IS NULL
+ GROUP BY p.id, p.uuid, p.asin, p.sku, p.gtin, p.product_name, p.brand, p.description,
+ p.url, p.price, p.initial_price, p.discount, p.currency, p.stock_quantity, p.sizes, p.colors, p.image_url,
+ p.available_for_delivery, p.available_for_pickup, p.free_returns,
+ p.is_active, p.created_at, p.updated_at, p.data_timestamp, p.category_id,
+ p.deleted_at, c.name, c.uuid, c.slug
+ "#,
+ )
+ .bind(uuid)
+ .fetch_optional(pool)
+ .await
+}
+
+/// List products with pagination, filtering, and rating aggregation
+///
+/// Uses a CTE to aggregate product data with ratings, then applies
+/// filters including rating-based filters on the CTE result.
+/// Returns products and a total count for pagination.
+#[instrument(
+ name = "SELECT products",
+ skip(pool, params),
+ fields(
+ otel.kind = "client",
+ db.system.name = "postgresql",
+ db.namespace = "products",
+ db.operation.name = "SELECT",
+ db.collection.name = "products",
+ db.query.text = "SELECT * FROM products WHERE ... ORDER BY ... LIMIT ... OFFSET ...",
+ otelmart.page = page,
+ otelmart.page_size = page_size,
+ otelmart.category.id = ?params.category_id,
+ otelmart.brand = ?params.brand
+ )
+)]
+pub async fn list_products(
+ pool: &PgPool,
+ params: &ProductQueryParams,
+ page: i32,
+ page_size: i32,
+ offset: i32,
+) -> Result<(Vec<ProductWithRating>, i64), sqlx::Error> {
+ // Check if we have rating filters (determines count query strategy)
+ let has_rating_filters =
+ params.rating_gt.is_some() || params.rating_lt.is_some() || params.rating_eq.is_some();
+
+ // Get total count - optimized based on whether we have rating filters
+ let total_count = if has_rating_filters {
+ count_with_ratings(pool, params).await?
+ } else {
+ count_simple(pool, params).await?
+ };
+
+ // Build main query with CTE for product data and ratings
+ let mut main_query = QueryBuilder::<Postgres>::new(
+ r#"
+ WITH product_data AS (
+ SELECT
+ p.id,
+ p.uuid,
+ p.product_name,
+ p.brand,
+ p.description,
+ p.price,
+ p.initial_price,
+ p.discount,
+ p.stock_quantity,
+ p.image_url,
+ p.category_id,
+ c.name as category_name,
+ p.created_at,
+ p.updated_at,
+ COALESCE(AVG(r.rating)::FLOAT8, NULL) as average_rating,
+ COUNT(r.id)::BIGINT as rating_count
+ FROM products p
+ LEFT JOIN categories c ON p.category_id = c.id
+ LEFT JOIN ratings r ON p.id = r.product_id
+ WHERE p.deleted_at IS NULL AND p.is_active = true AND p.stock_quantity > 0
+ GROUP BY p.id, p.uuid, p.product_name, p.brand, p.description, p.price,
+ p.initial_price, p.discount, p.stock_quantity, p.image_url, p.category_id, c.name,
+ p.created_at, p.updated_at
+ )
+ SELECT * FROM product_data WHERE 1=1
+ "#,
+ );
+
+ // Apply filters with safe parameter binding
+ apply_product_filters(&mut main_query, params, "", true);
+
+ // Add ordering and pagination
+ main_query.push(" ORDER BY updated_at DESC LIMIT ");
+ main_query.push_bind(page_size);
+ main_query.push(" OFFSET ");
+ main_query.push_bind(offset);
+
+ // Execute main query
+ let products = main_query.build_query_as().fetch_all(pool).await?;
+
+ Ok((products, total_count))
+}
+
+/// Simple count query without rating joins (faster for non-rating filters)
+async fn count_simple(pool: &PgPool, params: &ProductQueryParams) -> Result<i64, sqlx::Error> {
+ let mut count_builder = QueryBuilder::<Postgres>::new(
+ "SELECT COUNT(*) FROM products p WHERE p.deleted_at IS NULL AND p.is_active = true AND p.stock_quantity > 0"
+ );
+
+ // Apply non-rating filters
+ apply_product_filters(&mut count_builder, params, "p.", false);
+
+ count_builder.build_query_scalar().fetch_one(pool).await
+}
+
+/// Count query with rating aggregation (for rating-based filters)
+async fn count_with_ratings(
+ pool: &PgPool,
+ params: &ProductQueryParams,
+) -> Result<i64, sqlx::Error> {
+ let mut count_builder = QueryBuilder::<Postgres>::new(
+ r#"
+ SELECT COUNT(DISTINCT p.id) FROM products p
+ LEFT JOIN ratings r ON p.id = r.product_id
+ WHERE p.deleted_at IS NULL AND p.is_active = true AND p.stock_quantity > 0
+ "#,
+ );
+
+ // Apply non-rating filters to WHERE clause
+ apply_product_filters(&mut count_builder, params, "p.", false);
+
+ // Add GROUP BY for ratings aggregation
+ count_builder.push(" GROUP BY p.id");
+
+ // Apply rating filters with HAVING clause
+ let mut has_having = false;
+
+ if let Some(rating_gt) = params.rating_gt {
+ count_builder.push(if has_having { " AND " } else { " HAVING " });
+ has_having = true;
+ count_builder.push("AVG(r.rating) > ");
+ count_builder.push_bind(rating_gt);
+ }
+
+ if let Some(rating_lt) = params.rating_lt {
+ count_builder.push(if has_having { " AND " } else { " HAVING " });
+ has_having = true;
+ count_builder.push("AVG(r.rating) < ");
+ count_builder.push_bind(rating_lt);
+ }
+
+ if let Some(rating_eq) = params.rating_eq {
+ count_builder.push(if has_having { " AND " } else { " HAVING " });
+ count_builder.push("AVG(r.rating) = ");
+ count_builder.push_bind(rating_eq);
+ }
+
+ count_builder.build_query_scalar().fetch_one(pool).await
+}
+
+/// Apply product filters to a query builder
+///
+/// Centralizes filter logic to avoid duplication across count and main queries.
+///
+/// # Arguments
+/// * `query` - The QueryBuilder to append filters to
+/// * `params` - Query parameters containing filter values
+/// * `column_prefix` - Prefix for column names ("p." for table alias, "" for CTE columns)
+/// * `include_ratings` - Whether to include rating filters (only on CTE results)
+fn apply_product_filters<'a>(
+ query: &mut QueryBuilder<'a, Postgres>,
+ params: &'a ProductQueryParams,
+ column_prefix: &str,
+ include_ratings: bool,
+) {
+ // Name filter - case-insensitive partial match
+ if let Some(ref name) = params.name {
+ query.push(format!(" AND {}product_name ILIKE ", column_prefix).as_str());
+ query.push_bind(format!("%{}%", name));
+ }
+
+ // Category filter - exact match
+ if let Some(category_id) = params.category_id {
+ query.push(format!(" AND {}category_id = ", column_prefix).as_str());
+ query.push_bind(category_id);
+ }
+
+ // Brand filter - case-insensitive partial match
+ if let Some(ref brand) = params.brand {
+ query.push(format!(" AND {}brand ILIKE ", column_prefix).as_str());
+ query.push_bind(format!("%{}%", brand));
+ }
+
+ // Date range filters on updated_at
+ if let Some(start_date) = params.start_date {
+ query.push(format!(" AND {}updated_at >= ", column_prefix).as_str());
+ query.push_bind(start_date);
+ }
+
+ if let Some(end_date) = params.end_date {
+ query.push(format!(" AND {}updated_at <= ", column_prefix).as_str());
+ query.push_bind(end_date);
+ }
+
+ // Price range filters
+ if let Some(ref min_price) = params.min_price {
+ query.push(format!(" AND {}price >= ", column_prefix).as_str());
+ query.push_bind(min_price);
+ }
+
+ if let Some(ref max_price) = params.max_price {
+ query.push(format!(" AND {}price <= ", column_prefix).as_str());
+ query.push_bind(max_price);
+ }
+
+ // Rating filters (only for main query on CTE columns)
+ if include_ratings {
+ if let Some(rating_gt) = params.rating_gt {
+ query.push(" AND average_rating > ");
+ query.push_bind(rating_gt);
+ }
+
+ if let Some(rating_lt) = params.rating_lt {
+ query.push(" AND average_rating < ");
+ query.push_bind(rating_lt);
+ }
+
+ if let Some(rating_eq) = params.rating_eq {
+ query.push(" AND average_rating = ");
+ query.push_bind(rating_eq);
+ }
+ }
+}
diff --git a/distributed_observability/products/src/handlers/products.rs b/distributed_observability/products/src/handlers/products.rs
@@ -1,18 +1,18 @@
//! Product API handlers
//!
//! This module contains the HTTP request handlers for product-related endpoints.
-//! Handlers use sqlx for database queries with the products schema.
+//! Handlers delegate database operations to the repository layer (db module).
use axum::{
extract::{Path, Query, State},
http::StatusCode,
response::{IntoResponse, Json},
};
-use sqlx::{PgPool, Postgres, QueryBuilder};
-use tracing::instrument;
+use sqlx::PgPool;
use uuid::Uuid;
-use crate::models::{ProductDetail, ProductQueryParams, ProductWithRating, ProductsResponse};
+use crate::db;
+use crate::models::{ProductQueryParams, ProductsResponse};
use crate::utils::{calculate_pagination, calculate_total_pages, internal_error, not_found_error};
/// List products with pagination and filtering
@@ -33,23 +33,8 @@ use crate::utils::{calculate_pagination, calculate_total_pages, internal_error,
/// - `rating_gt` / `rating_lt` / `rating_eq` - Rating filters
/// - `min_price` / `max_price` - Price range filters
///
-/// # Example
-/// ```
-/// GET /products?page=2&page_size=50&brand=TechPro&min_price=100
-/// ```
-///
/// # Response
/// Returns a `ProductsResponse` with products array and pagination metadata.
-#[instrument(
- name = "list_products",
- skip(pool),
- fields(
- page = params.page,
- page_size = params.page_size,
- category_id = ?params.category_id,
- brand = ?params.brand
- )
-)]
pub async fn list_products(
State(pool): State<PgPool>,
Query(params): Query<ProductQueryParams>,
@@ -57,328 +42,39 @@ pub async fn list_products(
// Apply pagination defaults and constraints
let (page, page_size, offset) = calculate_pagination(params.page, params.page_size);
- // Check if we have rating filters (determines count query strategy)
- let has_rating_filters =
- params.rating_gt.is_some() || params.rating_lt.is_some() || params.rating_eq.is_some();
-
- // Build count query - optimized based on whether we have rating filters
- let total_count = if has_rating_filters {
- // Complex count: need to aggregate ratings and apply HAVING clause
- match build_count_with_ratings(&pool, ¶ms).await {
- Ok(count) => count,
- Err(response) => return response,
- }
- } else {
- // Simple count: no need for ratings JOIN or GROUP BY
- match build_simple_count(&pool, ¶ms).await {
- Ok(count) => count,
- Err(response) => return response,
- }
- };
-
- // Build main query with CTE for product data and ratings
- let mut main_query = QueryBuilder::<Postgres>::new(
- r#"
- WITH product_data AS (
- SELECT
- p.id,
- p.uuid,
- p.product_name,
- p.brand,
- p.description,
- p.price,
- p.initial_price,
- p.discount,
- p.stock_quantity,
- p.image_url,
- p.category_id,
- c.name as category_name,
- p.created_at,
- p.updated_at,
- COALESCE(AVG(r.rating)::FLOAT8, NULL) as average_rating,
- COUNT(r.id)::BIGINT as rating_count
- FROM products p
- LEFT JOIN categories c ON p.category_id = c.id
- LEFT JOIN ratings r ON p.id = r.product_id
- WHERE p.deleted_at IS NULL AND p.is_active = true AND p.stock_quantity > 0
- GROUP BY p.id, p.uuid, p.product_name, p.brand, p.description, p.price,
- p.initial_price, p.discount, p.stock_quantity, p.image_url, p.category_id, c.name,
- p.created_at, p.updated_at
- )
- SELECT * FROM product_data WHERE 1=1
- "#,
- );
-
- // Apply filters with safe parameter binding
- apply_filters(&mut main_query, ¶ms);
-
- // Add ordering and pagination
- main_query.push(" ORDER BY updated_at DESC LIMIT ");
- main_query.push_bind(page_size);
- main_query.push(" OFFSET ");
- main_query.push_bind(offset);
-
- // Execute main query to fetch products
- let products: Vec<ProductWithRating> = match main_query.build_query_as().fetch_all(&pool).await
- {
- Ok(products) => products,
- Err(e) => {
- return internal_error("Failed to fetch products", e.to_string());
- }
- };
-
- // Calculate total pages
- let total_pages = calculate_total_pages(total_count, page_size);
-
- // Build response
- let response = ProductsResponse {
- products,
- total_count,
- page,
- page_size,
- total_pages,
- };
-
- (StatusCode::OK, Json(response)).into_response()
-}
-
-/// Apply product filters to a query builder
-///
-/// This helper function centralizes filter logic to avoid duplication across
-/// count queries and main queries.
-///
-/// # Arguments
-/// * `query` - The QueryBuilder to apply filters to
-/// * `params` - Query parameters containing filter values
-/// * `column_prefix` - Prefix for column names ("p." for table alias, "" for CTE columns)
-/// * `include_ratings` - Whether to include rating filters (only for main query on CTE)
-fn apply_product_filters<'a>(
- query: &mut QueryBuilder<'a, Postgres>,
- params: &'a ProductQueryParams,
- column_prefix: &str,
- include_ratings: bool,
-) {
- // Name filter - case-insensitive partial match
- if let Some(ref name) = params.name {
- query.push(format!(" AND {}product_name ILIKE ", column_prefix).as_str());
- query.push_bind(format!("%{}%", name));
- }
-
- // Category filter - exact match
- if let Some(category_id) = params.category_id {
- query.push(format!(" AND {}category_id = ", column_prefix).as_str());
- query.push_bind(category_id);
- }
-
- // Brand filter - case-insensitive partial match
- if let Some(ref brand) = params.brand {
- query.push(format!(" AND {}brand ILIKE ", column_prefix).as_str());
- query.push_bind(format!("%{}%", brand));
- }
-
- // Date range filters on updated_at
- if let Some(start_date) = params.start_date {
- query.push(format!(" AND {}updated_at >= ", column_prefix).as_str());
- query.push_bind(start_date);
- }
-
- if let Some(end_date) = params.end_date {
- query.push(format!(" AND {}updated_at <= ", column_prefix).as_str());
- query.push_bind(end_date);
- }
-
- // Price range filters
- if let Some(ref min_price) = params.min_price {
- query.push(format!(" AND {}price >= ", column_prefix).as_str());
- query.push_bind(min_price);
- }
-
- if let Some(ref max_price) = params.max_price {
- query.push(format!(" AND {}price <= ", column_prefix).as_str());
- query.push_bind(max_price);
- }
-
- // Rating filters (only for main query on CTE columns)
- if include_ratings {
- if let Some(rating_gt) = params.rating_gt {
- query.push(" AND average_rating > ");
- query.push_bind(rating_gt);
- }
+ // Delegate to repository layer for database operations
+ match db::list_products(&pool, ¶ms, page, page_size, offset).await {
+ Ok((products, total_count)) => {
+ let total_pages = calculate_total_pages(total_count, page_size);
- if let Some(rating_lt) = params.rating_lt {
- query.push(" AND average_rating < ");
- query.push_bind(rating_lt);
- }
+ let response = ProductsResponse {
+ products,
+ total_count,
+ page,
+ page_size,
+ total_pages,
+ };
- if let Some(rating_eq) = params.rating_eq {
- query.push(" AND average_rating = ");
- query.push_bind(rating_eq);
+ (StatusCode::OK, Json(response)).into_response()
}
+ Err(e) => internal_error("Failed to fetch products", e.to_string()),
}
}
-/// Build simple count query without rating joins (faster)
-async fn build_simple_count(
- pool: &PgPool,
- params: &ProductQueryParams,
-) -> Result<i64, axum::response::Response> {
- let mut count_builder = QueryBuilder::<Postgres>::new(
- "SELECT COUNT(*) FROM products p WHERE p.deleted_at IS NULL AND p.is_active = true AND p.stock_quantity > 0"
- );
-
- // Apply non-rating filters using helper
- apply_product_filters(&mut count_builder, params, "p.", false);
-
- match count_builder.build_query_scalar().fetch_one(pool).await {
- Ok(count) => Ok(count),
- Err(e) => Err(internal_error("Failed to count products", e.to_string())),
- }
-}
-
-/// Build count query with rating aggregation (for rating filters)
-async fn build_count_with_ratings(
- pool: &PgPool,
- params: &ProductQueryParams,
-) -> Result<i64, axum::response::Response> {
- let mut count_builder = QueryBuilder::<Postgres>::new(
- r#"
- SELECT COUNT(DISTINCT p.id) FROM products p
- LEFT JOIN ratings r ON p.id = r.product_id
- WHERE p.deleted_at IS NULL AND p.is_active = true AND p.stock_quantity > 0
- "#,
- );
-
- // Apply non-rating filters to WHERE clause using helper
- apply_product_filters(&mut count_builder, params, "p.", false);
-
- // Add GROUP BY for ratings aggregation
- count_builder.push(" GROUP BY p.id");
-
- // Apply rating filters with HAVING clause
- let mut has_having = false;
-
- if let Some(rating_gt) = params.rating_gt {
- count_builder.push(if has_having { " AND " } else { " HAVING " });
- has_having = true;
- count_builder.push("AVG(r.rating) > ");
- count_builder.push_bind(rating_gt);
- }
-
- if let Some(rating_lt) = params.rating_lt {
- count_builder.push(if has_having { " AND " } else { " HAVING " });
- has_having = true;
- count_builder.push("AVG(r.rating) < ");
- count_builder.push_bind(rating_lt);
- }
-
- if let Some(rating_eq) = params.rating_eq {
- count_builder.push(if has_having { " AND " } else { " HAVING " });
- count_builder.push("AVG(r.rating) = ");
- count_builder.push_bind(rating_eq);
- }
-
- // This returns a count of distinct products, but we need to wrap it in another SELECT COUNT(*)
- // Actually, COUNT(DISTINCT p.id) already gives us the total count directly
- match count_builder.build_query_scalar().fetch_one(pool).await {
- Ok(count) => Ok(count),
- Err(e) => Err(internal_error("Failed to count products", e.to_string())),
- }
-}
-
-/// Apply filters to the main query
-fn apply_filters<'a>(query: &mut QueryBuilder<'a, Postgres>, params: &'a ProductQueryParams) {
- // Use helper with empty prefix (CTE columns) and include ratings
- apply_product_filters(query, params, "", true);
-}
-
/// Get detailed product information by UUID
///
/// # Endpoint
/// `GET /products/{uuid}`
///
-/// # Path Parameters
-/// - `uuid` - The product UUID
-///
-/// # Response
-/// Returns a `ProductDetail` with full product information including:
-/// - All product fields (name, description, price, stock, etc.)
-/// - Category information (name, uuid, slug)
-/// - Aggregated rating data (average rating, review count)
-/// - Additional metadata (timestamps, availability flags, etc.)
-///
/// # Errors
/// - `404 NOT FOUND` - Product not found or has been deleted
/// - `500 INTERNAL SERVER ERROR` - Database error
-///
-/// # Example
-/// ```
-/// GET /products/550e8400-e29b-41d4-a716-446655440000
-/// ```
-#[instrument( // Automatic span creation
- name = "get_product_by_id",
- skip(pool), // Don't log the entire AppState (too verbose)
- fields(product.uuid = %uuid) // Attach product UUID as span attribute
-)]
pub async fn get_product_by_id(
State(pool): State<PgPool>,
Path(uuid): Path<Uuid>,
) -> impl IntoResponse {
- tracing::info!("Fetching product");
-
- // Query product with all details, joined with category and aggregated ratings
- // Uses LEFT JOINs to handle products without categories or ratings
- let result = sqlx::query_as::<_, ProductDetail>(
- r#"
- SELECT
- p.id,
- p.uuid,
- p.asin,
- p.sku,
- p.gtin,
- p.product_name,
- p.brand,
- p.description,
- p.url,
- p.price,
- p.initial_price,
- p.discount,
- p.currency,
- p.stock_quantity,
- p.sizes,
- p.colors,
- p.image_url,
- p.available_for_delivery,
- p.available_for_pickup,
- p.free_returns,
- p.is_active,
- p.created_at,
- p.updated_at,
- p.data_timestamp,
- p.category_id,
- c.name as category_name,
- c.uuid as category_uuid,
- c.slug as category_slug,
- COALESCE(AVG(r.rating)::FLOAT8, NULL) as average_rating,
- COUNT(r.id)::BIGINT as rating_count,
- CAST(p.id AS TEXT) as product_id_str,
- get_root_category_name(p.category_id) as root_category_name,
- p.deleted_at
- FROM products p
- LEFT JOIN categories c ON p.category_id = c.id
- LEFT JOIN ratings r ON p.id = r.product_id
- WHERE p.uuid = $1 AND p.deleted_at IS NULL
- GROUP BY p.id, p.uuid, p.asin, p.sku, p.gtin, p.product_name, p.brand, p.description,
- p.url, p.price, p.initial_price, p.discount, p.currency, p.stock_quantity, p.sizes, p.colors, p.image_url,
- p.available_for_delivery, p.available_for_pickup, p.free_returns,
- p.is_active, p.created_at, p.updated_at, p.data_timestamp, p.category_id,
- p.deleted_at, c.name, c.uuid, c.slug
- "#,
- )
- .bind(uuid)
- .fetch_optional(&pool)
- .await;
-
- match result {
+ // Delegate to repository layer for database query
+ match db::get_product_by_uuid(&pool, uuid).await {
Ok(Some(mut product)) => {
// Set the string representation of the product ID
product.set_product_id();