repository.rs (9119B)
1 //! Inventory repository functions with OpenTelemetry instrumentation 2 //! 3 //! This module contains database operations for inventory/stock management, 4 //! separated from HTTP handlers. Each function is instrumented with 5 //! OpenTelemetry semantic conventions for database spans. 6 7 use sqlx::{PgPool, Postgres, QueryBuilder, Row}; 8 use tracing::{instrument, Span}; 9 use uuid::Uuid; 10 11 use crate::models::{InventoryQueryParams, InventoryWithPricing}; 12 13 /// List inventory with pagination and filtering 14 /// 15 /// Queries the v_product_inventory_pricing view which combines 16 /// inventory and pricing data. Supports filtering by stock status, 17 /// product UUID, and stock quantity range. 18 #[instrument( 19 name = "SELECT inventory", 20 skip(pool, params), 21 fields( 22 otel.kind = "client", 23 db.system.name = "postgresql", 24 db.namespace = "inventory", 25 db.operation.name = "SELECT", 26 db.collection.name = "product_inventory", 27 db.query.text = "SELECT * FROM v_product_inventory_pricing WHERE ... LIMIT ... OFFSET ...", 28 otelmart.page = page, 29 otelmart.page_size = page_size, 30 otelmart.stock_status = ?params.stock_status 31 ) 32 )] 33 pub async fn list_inventory( 34 pool: &PgPool, 35 params: &InventoryQueryParams, 36 page: i32, 37 page_size: i32, 38 offset: i32, 39 ) -> Result<(Vec<InventoryWithPricing>, i64), sqlx::Error> { 40 // Build count query with QueryBuilder for safe parameter binding 41 let mut count_builder: QueryBuilder<Postgres> = 42 QueryBuilder::new("SELECT COUNT(*) FROM v_product_inventory_pricing"); 43 44 let mut has_filters = false; 45 46 // Apply filters to count query 47 if let Some(ref status) = params.stock_status { 48 count_builder.push(if has_filters { " AND " } else { " WHERE " }); 49 has_filters = true; 50 count_builder.push("stock_status = "); 51 count_builder.push_bind(status); 52 } 53 54 if let Some(product_uuid) = params.product_uuid { 55 count_builder.push(if has_filters { " AND " } else { " WHERE " }); 56 has_filters = true; 57 count_builder.push("product_uuid = "); 58 count_builder.push_bind(product_uuid); 59 } 60 61 if let Some(min_stock) = params.min_stock { 62 count_builder.push(if has_filters { " AND " } else { " WHERE " }); 63 has_filters = true; 64 count_builder.push("available_quantity >= "); 65 count_builder.push_bind(min_stock); 66 } 67 68 if let Some(max_stock) = params.max_stock { 69 count_builder.push(if has_filters { " AND " } else { " WHERE " }); 70 count_builder.push("available_quantity <= "); 71 count_builder.push_bind(max_stock); 72 } 73 74 // Execute count query 75 let total_count: i64 = count_builder.build_query_scalar().fetch_one(pool).await?; 76 77 // Build main query with same filters 78 let mut query_builder: QueryBuilder<Postgres> = 79 QueryBuilder::new("SELECT * FROM v_product_inventory_pricing"); 80 81 let mut has_filters = false; 82 83 if let Some(ref status) = params.stock_status { 84 query_builder.push(if has_filters { " AND " } else { " WHERE " }); 85 has_filters = true; 86 query_builder.push("stock_status = "); 87 query_builder.push_bind(status); 88 } 89 90 if let Some(product_uuid) = params.product_uuid { 91 query_builder.push(if has_filters { " AND " } else { " WHERE " }); 92 has_filters = true; 93 query_builder.push("product_uuid = "); 94 query_builder.push_bind(product_uuid); 95 } 96 97 if let Some(min_stock) = params.min_stock { 98 query_builder.push(if has_filters { " AND " } else { " WHERE " }); 99 has_filters = true; 100 query_builder.push("available_quantity >= "); 101 query_builder.push_bind(min_stock); 102 } 103 104 if let Some(max_stock) = params.max_stock { 105 query_builder.push(if has_filters { " AND " } else { " WHERE " }); 106 query_builder.push("available_quantity <= "); 107 query_builder.push_bind(max_stock); 108 } 109 110 // Add ordering and pagination 111 query_builder.push(" ORDER BY stock_status DESC, available_quantity ASC LIMIT "); 112 query_builder.push_bind(page_size); 113 query_builder.push(" OFFSET "); 114 query_builder.push_bind(offset); 115 116 // Execute main query 117 let inventory = query_builder.build_query_as().fetch_all(pool).await?; 118 119 Ok((inventory, total_count)) 120 } 121 122 /// Get inventory for a specific product by UUID 123 #[instrument( 124 name = "SELECT inventory", 125 skip(pool), 126 fields( 127 otel.kind = "client", 128 db.system.name = "postgresql", 129 db.namespace = "inventory", 130 db.operation.name = "SELECT", 131 db.collection.name = "product_inventory", 132 db.query.text = "SELECT * FROM v_product_inventory_pricing WHERE product_uuid = $1", 133 otelmart.product.uuid = %product_uuid 134 ) 135 )] 136 pub async fn get_inventory_by_product( 137 pool: &PgPool, 138 product_uuid: Uuid, 139 ) -> Result<Option<InventoryWithPricing>, sqlx::Error> { 140 sqlx::query_as::<_, InventoryWithPricing>( 141 r#" 142 SELECT * FROM v_product_inventory_pricing 143 WHERE product_uuid = $1 144 "#, 145 ) 146 .bind(product_uuid) 147 .fetch_optional(pool) 148 .await 149 } 150 151 /// Update stock quantity for a product 152 /// 153 /// Returns the new available_quantity if the product was found, None otherwise. 154 #[instrument( 155 name = "UPDATE inventory", 156 skip(pool), 157 fields( 158 otel.kind = "client", 159 db.system.name = "postgresql", 160 db.namespace = "inventory", 161 db.operation.name = "UPDATE", 162 db.collection.name = "product_inventory", 163 db.query.text = "UPDATE product_inventory SET stock_quantity = $1 ... WHERE product_uuid = $4", 164 otelmart.product.uuid = %product_uuid, 165 otelmart.quantity = quantity, 166 db.response.returned_rows = tracing::field::Empty 167 ) 168 )] 169 pub async fn update_stock( 170 pool: &PgPool, 171 product_uuid: Uuid, 172 quantity: i32, 173 reorder_level: Option<i32>, 174 reorder_quantity: Option<i32>, 175 ) -> Result<Option<i32>, sqlx::Error> { 176 let result = sqlx::query( 177 r#" 178 UPDATE product_inventory 179 SET stock_quantity = $1, 180 reorder_level = COALESCE($2, reorder_level), 181 reorder_quantity = COALESCE($3, reorder_quantity), 182 last_restocked_at = CURRENT_TIMESTAMP 183 WHERE product_uuid = $4 184 RETURNING available_quantity 185 "#, 186 ) 187 .bind(quantity) 188 .bind(reorder_level) 189 .bind(reorder_quantity) 190 .bind(product_uuid) 191 .fetch_optional(pool) 192 .await?; 193 194 let rows = if result.is_some() { 1 } else { 0 }; 195 Span::current().record("db.response.returned_rows", rows); 196 197 Ok(result.map(|row| row.get("available_quantity"))) 198 } 199 200 /// Reserve stock for an order using the database function 201 /// 202 /// Calls the reserve_stock PostgreSQL function which atomically checks 203 /// available stock and creates a reservation. Returns true if successful. 204 #[instrument( 205 name = "SELECT reserve_stock", 206 skip(pool), 207 fields( 208 db.system.name = "postgresql", 209 db.namespace = "inventory", 210 db.operation.name = "SELECT", 211 db.collection.name = "product_inventory", 212 db.query.text = "SELECT reserve_stock($1, $2)", 213 otelmart.product.uuid = %product_uuid, 214 otelmart.quantity = quantity 215 ) 216 )] 217 pub async fn reserve_stock( 218 pool: &PgPool, 219 product_uuid: Uuid, 220 quantity: i32, 221 ) -> Result<bool, sqlx::Error> { 222 sqlx::query_scalar::<_, bool>(r#"SELECT reserve_stock($1, $2)"#) 223 .bind(product_uuid) 224 .bind(quantity) 225 .fetch_one(pool) 226 .await 227 } 228 229 /// Release reserved stock using the database function 230 #[instrument( 231 name = "SELECT release_stock", 232 skip(pool), 233 fields( 234 db.system.name = "postgresql", 235 db.namespace = "inventory", 236 db.operation.name = "SELECT", 237 db.collection.name = "product_inventory", 238 db.query.text = "SELECT release_stock($1, $2)", 239 otelmart.product.uuid = %product_uuid, 240 otelmart.quantity = quantity 241 ) 242 )] 243 pub async fn release_stock( 244 pool: &PgPool, 245 product_uuid: Uuid, 246 quantity: i32, 247 ) -> Result<(), sqlx::Error> { 248 sqlx::query(r#"SELECT release_stock($1, $2)"#) 249 .bind(product_uuid) 250 .bind(quantity) 251 .execute(pool) 252 .await?; 253 Ok(()) 254 } 255 256 /// Confirm a sale and decrease stock using the database function 257 #[instrument( 258 name = "SELECT confirm_stock_sale", 259 skip(pool), 260 fields( 261 db.system.name = "postgresql", 262 db.namespace = "inventory", 263 db.operation.name = "SELECT", 264 db.collection.name = "product_inventory", 265 db.query.text = "SELECT confirm_stock_sale($1, $2, $3)", 266 otelmart.product.uuid = %product_uuid, 267 otelmart.order.uuid = %order_uuid, 268 otelmart.quantity = quantity 269 ) 270 )] 271 pub async fn confirm_sale( 272 pool: &PgPool, 273 product_uuid: Uuid, 274 quantity: i32, 275 order_uuid: Uuid, 276 ) -> Result<(), sqlx::Error> { 277 sqlx::query(r#"SELECT confirm_stock_sale($1, $2, $3)"#) 278 .bind(product_uuid) 279 .bind(quantity) 280 .bind(order_uuid) 281 .execute(pool) 282 .await?; 283 Ok(()) 284 }