1use super::map_db_error;
4use super::parse_helpers::parse_decimal as parse_decimal_with_context;
5use chrono::{DateTime, Datelike, Duration, NaiveDate, NaiveTime, Utc};
6use r2d2::Pool;
7use r2d2_sqlite::SqliteConnectionManager;
8use rusqlite::ToSql;
9use rust_decimal::Decimal;
10use stateset_core::{
11 AnalyticsQuery, AnalyticsRepository, CustomerMetrics, DemandForecast, FulfillmentMetrics,
12 InventoryHealth, InventoryMovement, LowStockItem, OrderStatusBreakdown, ProductId,
13 ProductPerformance, Result, ReturnMetrics, ReturnReasonCount, RevenueByPeriod, RevenueForecast,
14 SalesSummary, TimeGranularity, TimePeriod, TopCustomer, TopProduct, TopReturnedProduct, Trend,
15 validate_batch_size,
16};
17use uuid::Uuid;
18
19const ISO_WEEK_EXPR: &str = "strftime('%Y', created_at, '-3 days', 'weekday 4') || '-W' || \
25 substr('0' || ((strftime('%j', created_at, '-3 days', 'weekday 4') - 1) / 7 + 1), -2, 2)";
26
27#[derive(Debug)]
28pub struct SqliteAnalyticsRepository {
29 pool: Pool<SqliteConnectionManager>,
30}
31
32impl SqliteAnalyticsRepository {
33 #[must_use]
34 pub const fn new(pool: Pool<SqliteConnectionManager>) -> Self {
35 Self { pool }
36 }
37
38 fn conn(&self) -> Result<r2d2::PooledConnection<SqliteConnectionManager>> {
39 self.pool.get().map_err(|e| stateset_core::CommerceError::DatabaseError(e.to_string()))
40 }
41
42 const fn start_of_day(date: NaiveDate) -> DateTime<Utc> {
43 DateTime::from_naive_utc_and_offset(date.and_time(NaiveTime::MIN), Utc)
44 }
45
46 fn end_of_day(date: NaiveDate) -> DateTime<Utc> {
47 Self::start_of_day(date) + Duration::days(1) - Duration::seconds(1)
48 }
49
50 fn first_day_of_month(date: NaiveDate) -> NaiveDate {
51 date.with_day(1).unwrap_or(date)
52 }
53
54 fn first_day_of_year(date: NaiveDate) -> NaiveDate {
55 date.with_month(1)
56 .and_then(|d| d.with_day(1))
57 .unwrap_or_else(|| Self::first_day_of_month(date))
58 }
59
60 fn all_time_start() -> DateTime<Utc> {
61 if let Some(date) = NaiveDate::from_ymd_opt(2000, 1, 1) {
62 return Self::start_of_day(date);
63 }
64 if let Some(epoch) = DateTime::<Utc>::from_timestamp(0, 0) {
65 return epoch;
66 }
67 Utc::now()
68 }
69
70 fn get_date_range(&self, query: &AnalyticsQuery) -> (DateTime<Utc>, DateTime<Utc>) {
72 let now = Utc::now();
73 let period = query.period.unwrap_or(TimePeriod::Last30Days);
74
75 match period {
76 TimePeriod::Today => (Self::start_of_day(now.date_naive()), now),
77 TimePeriod::Yesterday => {
78 let yesterday = now - Duration::days(1);
79 (
80 Self::start_of_day(yesterday.date_naive()),
81 Self::end_of_day(yesterday.date_naive()),
82 )
83 }
84 TimePeriod::Last7Days => (now - Duration::days(7), now),
85 TimePeriod::Last30Days => (now - Duration::days(30), now),
86 TimePeriod::ThisMonth => {
87 let start = Self::start_of_day(Self::first_day_of_month(now.date_naive()));
88 (start, now)
89 }
90 TimePeriod::LastMonth => {
91 let this_month_start = Self::first_day_of_month(now.date_naive());
92 let last_month_end = this_month_start - Duration::days(1);
93 let last_month_start = Self::first_day_of_month(last_month_end);
94 (Self::start_of_day(last_month_start), Self::end_of_day(last_month_end))
95 }
96 TimePeriod::ThisQuarter | TimePeriod::LastQuarter => {
97 (now - Duration::days(90), now)
99 }
100 TimePeriod::ThisYear => {
101 let start = Self::start_of_day(Self::first_day_of_year(now.date_naive()));
102 (start, now)
103 }
104 TimePeriod::LastYear => (now - Duration::days(365), now),
105 TimePeriod::AllTime => (Self::all_time_start(), now),
106 TimePeriod::Custom => {
107 if let Some(ref range) = query.date_range {
108 (range.start.unwrap_or(now - Duration::days(30)), range.end.unwrap_or(now))
109 } else {
110 (now - Duration::days(30), now)
111 }
112 }
113 _ => (now - Duration::days(30), now),
114 }
115 }
116}
117
118fn parse_decimal_value(value: &str, field: &str) -> Result<Decimal> {
119 parse_decimal_with_context(value, "analytics", field)
120}
121
122impl AnalyticsRepository for SqliteAnalyticsRepository {
123 fn get_sales_summary(&self, query: AnalyticsQuery) -> Result<SalesSummary> {
124 let conn = self.conn()?;
125 let (start, end) = self.get_date_range(&query);
126 let start_str = start.to_rfc3339();
127 let end_str = end.to_rfc3339();
128
129 let mut stmt = conn
131 .prepare(
132 r"
133 SELECT
134 decimal_sum(total_amount) as revenue,
135 COUNT(*) as order_count,
136 -- avg_order is an average, so the float coercion in the
137 -- built-in SUM/divide is immaterial here; only the exact
138 -- reconciled totals use decimal_sum.
139 CAST(COALESCE(SUM(total_amount) / NULLIF(COUNT(*), 0), 0) AS TEXT) as avg_order,
140 COUNT(DISTINCT customer_id) as unique_customers
141 FROM orders
142 WHERE created_at >= ?1 AND created_at <= ?2
143 AND status NOT IN ('cancelled', 'refunded')
144 ",
145 )
146 .map_err(map_db_error)?;
147
148 let (revenue, order_count, avg_order, unique_customers): (String, i64, String, i64) = stmt
149 .query_row([&start_str, &end_str], |row| {
150 Ok((row.get(0)?, row.get(1)?, row.get(2)?, row.get(3)?))
151 })
152 .map_err(map_db_error)?;
153
154 let items_sold: i64 = conn
156 .query_row(
157 r"
158 SELECT COALESCE(SUM(oi.quantity), 0)
159 FROM order_items oi
160 JOIN orders o ON oi.order_id = o.id
161 WHERE o.created_at >= ?1 AND o.created_at <= ?2
162 AND o.status NOT IN ('cancelled', 'refunded')
163 ",
164 [&start_str, &end_str],
165 |row| row.get(0),
166 )
167 .unwrap_or(0);
168
169 let period_duration = end - start;
171 let prev_end = start;
172 let prev_start = prev_end - period_duration;
173 let prev_start_str = prev_start.to_rfc3339();
174 let prev_end_str = prev_end.to_rfc3339();
175
176 let (prev_revenue, prev_order_count): (String, i64) = conn
178 .query_row(
179 r"
180 SELECT
181 decimal_sum(total_amount) as revenue,
182 COUNT(*) as order_count
183 FROM orders
184 WHERE created_at >= ?1 AND created_at < ?2
185 AND status NOT IN ('cancelled', 'refunded')
186 ",
187 [&prev_start_str, &prev_end_str],
188 |row| Ok((row.get(0)?, row.get(1)?)),
189 )
190 .unwrap_or(("0".to_string(), 0));
191
192 let current_revenue = parse_decimal_value(&revenue, "revenue")?;
193 let previous_revenue = parse_decimal_value(&prev_revenue, "previous_revenue")?;
194
195 let revenue_change_percent = if previous_revenue != Decimal::ZERO {
197 Some(((current_revenue - previous_revenue) / previous_revenue) * Decimal::from(100))
198 } else if current_revenue != Decimal::ZERO {
199 Some(Decimal::from(100)) } else {
201 Some(Decimal::ZERO)
202 };
203
204 let order_count_change_percent = if prev_order_count > 0 {
205 let change =
206 ((order_count - prev_order_count) as f64 / prev_order_count as f64) * 100.0;
207 Decimal::from_f64_retain(change)
208 } else if order_count > 0 {
209 Some(Decimal::from(100))
210 } else {
211 Some(Decimal::ZERO)
212 };
213
214 Ok(SalesSummary {
215 total_revenue: current_revenue,
216 order_count: order_count as u64,
217 average_order_value: parse_decimal_value(&avg_order, "average_order_value")?,
218 items_sold: items_sold as u64,
219 unique_customers: unique_customers as u64,
220 revenue_change_percent,
221 order_count_change_percent,
222 period_start: Some(start),
223 period_end: Some(end),
224 })
225 }
226
227 fn get_revenue_by_period(&self, query: AnalyticsQuery) -> Result<Vec<RevenueByPeriod>> {
228 let conn = self.conn()?;
229 let (start, end) = self.get_date_range(&query);
230 let start_str = start.to_rfc3339();
231 let end_str = end.to_rfc3339();
232
233 let granularity = query.granularity.unwrap_or(TimeGranularity::Day);
234 let period_expr = match granularity {
235 TimeGranularity::Hour => "strftime('%Y-%m-%d %H:00', created_at)".to_string(),
236 TimeGranularity::Day => "strftime('%Y-%m-%d', created_at)".to_string(),
237 TimeGranularity::Week => ISO_WEEK_EXPR.to_string(),
238 TimeGranularity::Month => "strftime('%Y-%m', created_at)".to_string(),
239 TimeGranularity::Quarter => {
240 "strftime('%Y', created_at) || '-Q' || ((CAST(strftime('%m', created_at) AS INTEGER) - 1) / 3 + 1)".to_string()
243 }
244 TimeGranularity::Year => "strftime('%Y', created_at)".to_string(),
245 _ => "strftime('%Y-%m-%d', created_at)".to_string(),
246 };
247
248 let mut stmt = conn
249 .prepare(&format!(
250 r"
251 SELECT
252 {period_expr} as period,
253 decimal_sum(total_amount) as revenue,
254 COUNT(*) as order_count,
255 MIN(created_at) as period_start
256 FROM orders
257 WHERE created_at >= ?1 AND created_at <= ?2
258 AND status NOT IN ('cancelled', 'refunded')
259 GROUP BY {period_expr}
260 ORDER BY period
261 "
262 ))
263 .map_err(map_db_error)?;
264
265 let rows = stmt
266 .query_map([&start_str, &end_str], |row| {
267 let period: String = row.get(0)?;
268 let revenue: String = row.get(1)?;
269 let order_count: i64 = row.get(2)?;
270 let period_start: String = row.get(3)?;
271 Ok((period, revenue, order_count, period_start))
272 })
273 .map_err(map_db_error)?;
274
275 let mut results = Vec::new();
276 for row in rows {
277 let (period, revenue, order_count, period_start) = row.map_err(map_db_error)?;
278 let revenue = parse_decimal_value(&revenue, "revenue")?;
279 results.push(RevenueByPeriod {
280 period,
281 revenue,
282 order_count: order_count as u64,
283 period_start: DateTime::parse_from_rfc3339(&period_start)
284 .map(|dt| dt.with_timezone(&Utc))
285 .unwrap_or(start),
286 });
287 }
288
289 Ok(results)
290 }
291
292 fn get_top_products(&self, query: AnalyticsQuery) -> Result<Vec<TopProduct>> {
293 let conn = self.conn()?;
294 let (start, end) = self.get_date_range(&query);
295 let start_str = start.to_rfc3339();
296 let end_str = end.to_rfc3339();
297 let limit = i64::from(query.limit.unwrap_or(10));
298
299 let mut stmt = conn
300 .prepare(
301 r"
302 SELECT
303 MAX(oi.product_id) as product_id,
304 oi.sku,
305 MAX(oi.name) as name,
306 SUM(oi.quantity) as units_sold,
307 decimal_sum(oi.total) as revenue,
308 COUNT(DISTINCT oi.order_id) as order_count,
309 CAST(COALESCE(AVG(oi.unit_price), 0) AS TEXT) as avg_price
310 FROM order_items oi
311 JOIN orders o ON oi.order_id = o.id
312 WHERE o.created_at >= ?1 AND o.created_at <= ?2
313 AND o.status NOT IN ('cancelled', 'refunded')
314 GROUP BY oi.sku
315 ORDER BY revenue DESC
316 LIMIT ?3
317 ",
318 )
319 .map_err(map_db_error)?;
320
321 let rows = stmt
322 .query_map([&start_str as &dyn rusqlite::ToSql, &end_str, &limit], |row| {
323 let product_id: Option<String> = row.get(0)?;
324 let sku: String = row.get(1)?;
325 let name: String = row.get(2)?;
326 let units_sold: i64 = row.get(3)?;
327 let revenue: String = row.get(4)?;
328 let order_count: i64 = row.get(5)?;
329 let avg_price: String = row.get(6)?;
330 Ok((product_id, sku, name, units_sold, revenue, order_count, avg_price))
331 })
332 .map_err(map_db_error)?;
333
334 let mut results = Vec::new();
335 for row in rows {
336 let (product_id, sku, name, units_sold, revenue, order_count, avg_price) =
337 row.map_err(map_db_error)?;
338 let revenue = parse_decimal_value(&revenue, "revenue")?;
339 let average_price = parse_decimal_value(&avg_price, "average_price")?;
340 results.push(TopProduct {
341 product_id: product_id.and_then(|s| Uuid::parse_str(&s).ok().map(ProductId::from)),
342 sku,
343 name,
344 units_sold: units_sold as u64,
345 revenue,
346 order_count: order_count as u64,
347 average_price,
348 });
349 }
350
351 Ok(results)
352 }
353
354 fn get_product_performance(&self, query: AnalyticsQuery) -> Result<Vec<ProductPerformance>> {
355 let top_products = self.get_top_products(query)?;
357 Ok(top_products
358 .into_iter()
359 .map(|p| ProductPerformance {
360 product_id: p.product_id.unwrap_or_default(),
361 sku: p.sku,
362 name: p.name,
363 units_sold: p.units_sold,
364 revenue: p.revenue,
365 previous_units_sold: 0,
366 previous_revenue: Decimal::ZERO,
367 units_growth_percent: Decimal::ZERO,
368 revenue_growth_percent: Decimal::ZERO,
369 })
370 .collect())
371 }
372
373 fn get_customer_metrics(&self, query: AnalyticsQuery) -> Result<CustomerMetrics> {
374 let conn = self.conn()?;
375 let (start, end) = self.get_date_range(&query);
376 let start_str = start.to_rfc3339();
377 let end_str = end.to_rfc3339();
378
379 let total_customers: i64 =
381 conn.query_row("SELECT COUNT(*) FROM customers", [], |row| row.get(0)).unwrap_or(0);
382
383 let new_customers: i64 = conn
385 .query_row(
386 "SELECT COUNT(*) FROM customers WHERE created_at >= ?1 AND created_at <= ?2",
387 [&start_str, &end_str],
388 |row| row.get(0),
389 )
390 .unwrap_or(0);
391
392 let returning_customers: i64 = conn
394 .query_row(
395 r"
396 SELECT COUNT(*) FROM (
397 SELECT customer_id FROM orders
398 GROUP BY customer_id
399 HAVING COUNT(*) > 1
400 )
401 ",
402 [],
403 |row| row.get(0),
404 )
405 .unwrap_or(0);
406
407 let avg_ltv: String = conn
409 .query_row(
410 r"
411 SELECT CAST(COALESCE(AVG(total), 0) AS TEXT) FROM (
412 SELECT customer_id, SUM(total_amount) as total
413 FROM orders
414 WHERE status NOT IN ('cancelled', 'refunded')
415 GROUP BY customer_id
416 )
417 ",
418 [],
419 |row| row.get(0),
420 )
421 .unwrap_or_else(|_| "0".to_string());
422
423 let avg_orders: String = conn
425 .query_row(
426 r"
427 SELECT CAST(COALESCE(AVG(cnt), 0) AS TEXT) FROM (
428 SELECT customer_id, COUNT(*) as cnt
429 FROM orders
430 GROUP BY customer_id
431 )
432 ",
433 [],
434 |row| row.get(0),
435 )
436 .unwrap_or_else(|_| "0".to_string());
437
438 let average_lifetime_value = parse_decimal_value(&avg_ltv, "average_lifetime_value")?;
439 let average_orders_per_customer =
440 parse_decimal_value(&avg_orders, "average_orders_per_customer")?;
441
442 Ok(CustomerMetrics {
443 total_customers: total_customers as u64,
444 new_customers: new_customers as u64,
445 returning_customers: returning_customers as u64,
446 average_lifetime_value,
447 average_orders_per_customer,
448 retention_rate_percent: None,
449 })
450 }
451
452 fn get_top_customers(&self, query: AnalyticsQuery) -> Result<Vec<TopCustomer>> {
453 let conn = self.conn()?;
454 let (start, end) = self.get_date_range(&query);
455 let start_str = start.to_rfc3339();
456 let end_str = end.to_rfc3339();
457 let limit = i64::from(query.limit.unwrap_or(10));
458
459 let mut stmt = conn
460 .prepare(
461 r"
462 SELECT
463 c.id,
464 c.email,
465 COALESCE(c.first_name || ' ' || c.last_name, c.email) as name,
466 decimal_sum(o.total_amount) as total_spent,
467 COUNT(o.id) as order_count,
468 CAST(COALESCE(AVG(o.total_amount), 0) AS TEXT) as avg_order,
469 MIN(o.created_at) as first_order,
470 MAX(o.created_at) as last_order
471 FROM customers c
472 LEFT JOIN orders o ON c.id = o.customer_id
473 AND o.status NOT IN ('cancelled', 'refunded')
474 AND o.created_at >= ?1 AND o.created_at <= ?2
475 GROUP BY c.id
476 ORDER BY total_spent DESC
477 LIMIT ?3
478 ",
479 )
480 .map_err(map_db_error)?;
481
482 let rows = stmt
483 .query_map([&start_str as &dyn rusqlite::ToSql, &end_str, &limit], |row| {
484 let id: String = row.get(0)?;
485 let email: String = row.get(1)?;
486 let name: String = row.get(2)?;
487 let total_spent: String = row.get(3)?;
488 let order_count: i64 = row.get(4)?;
489 let avg_order: String = row.get(5)?;
490 let first_order: Option<String> = row.get(6)?;
491 let last_order: Option<String> = row.get(7)?;
492 Ok((id, email, name, total_spent, order_count, avg_order, first_order, last_order))
493 })
494 .map_err(map_db_error)?;
495
496 let mut results = Vec::new();
497 for row in rows {
498 let (id, email, name, total_spent, order_count, avg_order, first_order, last_order) =
499 row.map_err(map_db_error)?;
500 let total_spent = parse_decimal_value(&total_spent, "total_spent")?;
501 let average_order_value = parse_decimal_value(&avg_order, "average_order_value")?;
502 results.push(TopCustomer {
503 customer_id: Uuid::parse_str(&id).unwrap_or_default(),
504 email,
505 name,
506 total_spent,
507 order_count: order_count as u64,
508 average_order_value,
509 first_order_date: first_order
510 .and_then(|s| DateTime::parse_from_rfc3339(&s).ok())
511 .map(|dt| dt.with_timezone(&Utc)),
512 last_order_date: last_order
513 .and_then(|s| DateTime::parse_from_rfc3339(&s).ok())
514 .map(|dt| dt.with_timezone(&Utc)),
515 });
516 }
517
518 Ok(results)
519 }
520
521 fn get_inventory_health(&self) -> Result<InventoryHealth> {
522 let conn = self.conn()?;
523
524 let total_skus: i64 = conn
525 .query_row("SELECT COUNT(*) FROM inventory_items", [], |row| row.get(0))
526 .unwrap_or(0);
527
528 let (in_stock, low_stock, out_of_stock): (i64, i64, i64) = conn
530 .query_row(
531 r"
532 SELECT
533 SUM(CASE WHEN ib.on_hand > COALESCE(ii.reorder_point, 10) THEN 1 ELSE 0 END),
534 SUM(CASE WHEN ib.on_hand <= COALESCE(ii.reorder_point, 10) AND ib.on_hand > 0 THEN 1 ELSE 0 END),
535 SUM(CASE WHEN ib.on_hand <= 0 THEN 1 ELSE 0 END)
536 FROM inventory_items ii
537 LEFT JOIN inventory_balances ib ON ii.id = ib.item_id
538 ",
539 [],
540 |row| Ok((row.get(0).unwrap_or(0), row.get(1).unwrap_or(0), row.get(2).unwrap_or(0))),
541 )
542 .unwrap_or((0, 0, 0));
543
544 let total_value: String = conn
547 .query_row(
548 r"
549 SELECT decimal_sum_product(ib.on_hand, COALESCE(pv.cost_price, pv.price, 0))
550 FROM inventory_items ii
551 LEFT JOIN inventory_balances ib ON ii.id = ib.item_id
552 LEFT JOIN product_variants pv ON ii.sku = pv.sku
553 ",
554 [],
555 |row| row.get(0),
556 )
557 .unwrap_or_else(|_| "0".to_string());
558
559 let total_value = parse_decimal_value(&total_value, "total_value")?;
560
561 Ok(InventoryHealth {
562 total_skus: total_skus as u64,
563 in_stock_skus: in_stock as u64,
564 low_stock_skus: low_stock as u64,
565 out_of_stock_skus: out_of_stock as u64,
566 total_value,
567 turnover_ratio: None,
568 })
569 }
570
571 fn get_low_stock_items(&self, threshold: Option<Decimal>) -> Result<Vec<LowStockItem>> {
572 let conn = self.conn()?;
573 let threshold = threshold.unwrap_or(Decimal::from(10));
574
575 let mut stmt = conn
582 .prepare(
583 r"
584 SELECT
585 ii.sku,
586 ii.name,
587 decimal_sum(ib.quantity_on_hand) as on_hand,
588 decimal_sum(ib.quantity_allocated) as allocated,
589 MAX(ib.reorder_point) as reorder_point
590 FROM inventory_items ii
591 LEFT JOIN inventory_balances ib ON ii.id = ib.item_id
592 GROUP BY ii.id, ii.sku, ii.name
593 ",
594 )
595 .map_err(map_db_error)?;
596
597 let rows = stmt
598 .query_map([], |row| {
599 let sku: String = row.get(0)?;
600 let name: String = row.get(1)?;
601 let on_hand: String = row.get(2)?;
602 let allocated: String = row.get(3)?;
603 let reorder_point: Option<String> = row.get(4)?;
604 Ok((sku, name, on_hand, allocated, reorder_point))
605 })
606 .map_err(map_db_error)?;
607
608 let mut results = Vec::new();
609 for row in rows {
610 let (sku, name, on_hand, allocated, reorder_point) = row.map_err(map_db_error)?;
611 let on_hand = parse_decimal_value(&on_hand, "on_hand")?;
612 let allocated = parse_decimal_value(&allocated, "allocated")?;
613 let available = on_hand - allocated;
614 if available > threshold {
615 continue;
616 }
617 let reorder_point =
618 reorder_point.map(|s| parse_decimal_value(&s, "reorder_point")).transpose()?;
619 results.push(LowStockItem {
620 sku,
621 name,
622 on_hand,
623 allocated,
624 available,
625 reorder_point,
626 average_daily_sales: None,
627 days_of_stock: None,
628 });
629 }
630 results.sort_by(|a, b| a.available.cmp(&b.available));
631
632 Ok(results)
633 }
634
635 fn get_inventory_movement(&self, query: AnalyticsQuery) -> Result<Vec<InventoryMovement>> {
636 let conn = self.conn()?;
637 let (start, end) = self.get_date_range(&query);
638 let start_str = start.to_rfc3339();
639 let end_str = end.to_rfc3339();
640
641 let mut stmt = conn
642 .prepare(
643 r"
644 SELECT
645 ii.sku,
646 ii.name,
647 COALESCE(SUM(CASE WHEN it.transaction_type = 'sale' THEN ABS(it.quantity) ELSE 0 END), 0) as sold,
648 COALESCE(SUM(CASE WHEN it.transaction_type = 'adjustment_in' THEN it.quantity ELSE 0 END), 0) as received,
649 COALESCE(SUM(CASE WHEN it.transaction_type = 'return' THEN it.quantity ELSE 0 END), 0) as returned,
650 COALESCE(SUM(CASE WHEN it.transaction_type IN ('adjustment_in', 'adjustment_out') THEN it.quantity ELSE 0 END), 0) as adjusted,
651 COALESCE(SUM(it.quantity), 0) as net_change
652 FROM inventory_items ii
653 LEFT JOIN inventory_transactions it ON ii.id = it.item_id
654 AND it.created_at >= ?1 AND it.created_at <= ?2
655 GROUP BY ii.id
656 HAVING net_change != 0
657 ORDER BY ABS(net_change) DESC
658 LIMIT 50
659 ",
660 )
661 .map_err(map_db_error)?;
662
663 let rows = stmt
664 .query_map([&start_str, &end_str], |row| {
665 let sku: String = row.get(0)?;
666 let name: String = row.get(1)?;
667 let sold: i64 = row.get(2)?;
668 let received: i64 = row.get(3)?;
669 let returned: i64 = row.get(4)?;
670 let adjusted: i64 = row.get(5)?;
671 let net_change: i64 = row.get(6)?;
672 Ok((sku, name, sold, received, returned, adjusted, net_change))
673 })
674 .map_err(map_db_error)?;
675
676 let mut results = Vec::new();
677 for row in rows {
678 let (sku, name, sold, received, returned, adjusted, net_change) =
679 row.map_err(map_db_error)?;
680 results.push(InventoryMovement {
681 sku,
682 name,
683 units_sold: sold as u64,
684 units_received: received as u64,
685 units_returned: returned as u64,
686 units_adjusted: adjusted,
687 net_change,
688 });
689 }
690
691 Ok(results)
692 }
693
694 fn get_order_status_breakdown(&self, query: AnalyticsQuery) -> Result<OrderStatusBreakdown> {
695 let conn = self.conn()?;
696 let (start, end) = self.get_date_range(&query);
697 let start_str = start.to_rfc3339();
698 let end_str = end.to_rfc3339();
699
700 let mut stmt = conn
701 .prepare(
702 r"
703 SELECT status, COUNT(*) as cnt
704 FROM orders
705 WHERE created_at >= ?1 AND created_at <= ?2
706 GROUP BY status
707 ",
708 )
709 .map_err(map_db_error)?;
710
711 let rows = stmt
712 .query_map([&start_str, &end_str], |row| {
713 let status: String = row.get(0)?;
714 let count: i64 = row.get(1)?;
715 Ok((status, count))
716 })
717 .map_err(map_db_error)?;
718
719 let mut breakdown = OrderStatusBreakdown::default();
720 for row in rows {
721 let (status, count) = row.map_err(map_db_error)?;
722 let count = count as u64;
723 breakdown.total += count;
724 match status.as_str() {
725 "pending" => breakdown.pending = count,
726 "confirmed" => breakdown.confirmed = count,
727 "processing" => breakdown.processing = count,
728 "shipped" => breakdown.shipped = count,
729 "delivered" => breakdown.delivered = count,
730 "cancelled" => breakdown.cancelled = count,
731 "refunded" => breakdown.refunded = count,
732 _ => {}
733 }
734 }
735
736 Ok(breakdown)
737 }
738
739 fn get_fulfillment_metrics(&self, query: AnalyticsQuery) -> Result<FulfillmentMetrics> {
740 let conn = self.conn()?;
741 let (start, end) = self.get_date_range(&query);
742 let _start_str = start.to_rfc3339();
743 let _end_str = end.to_rfc3339();
744
745 let today_start = Self::start_of_day(Utc::now().date_naive()).to_rfc3339();
747 let shipped_today: i64 = conn
748 .query_row(
749 "SELECT COUNT(*) FROM orders WHERE status = 'shipped' AND updated_at >= ?1",
750 [&today_start],
751 |row| row.get(0),
752 )
753 .unwrap_or(0);
754
755 let awaiting_shipment: i64 = conn
757 .query_row(
758 "SELECT COUNT(*) FROM orders WHERE status IN ('confirmed', 'processing')",
759 [],
760 |row| row.get(0),
761 )
762 .unwrap_or(0);
763
764 Ok(FulfillmentMetrics {
765 avg_time_to_ship_hours: None,
766 avg_time_to_deliver_hours: None,
767 on_time_shipping_percent: None,
768 on_time_delivery_percent: None,
769 shipped_today: shipped_today as u64,
770 awaiting_shipment: awaiting_shipment as u64,
771 })
772 }
773
774 fn get_return_metrics(&self, query: AnalyticsQuery) -> Result<ReturnMetrics> {
775 let conn = self.conn()?;
776 let (start, end) = self.get_date_range(&query);
777 let start_str = start.to_rfc3339();
778 let end_str = end.to_rfc3339();
779
780 let total_returns: i64 = conn
782 .query_row(
783 "SELECT COUNT(*) FROM returns WHERE created_at >= ?1 AND created_at <= ?2",
784 [&start_str, &end_str],
785 |row| row.get(0),
786 )
787 .unwrap_or(0);
788
789 let total_orders: i64 = conn
791 .query_row(
792 "SELECT COUNT(*) FROM orders WHERE created_at >= ?1 AND created_at <= ?2",
793 [&start_str, &end_str],
794 |row| row.get(0),
795 )
796 .unwrap_or(1);
797
798 let return_rate = if total_orders > 0 {
799 Decimal::from(total_returns * 100) / Decimal::from(total_orders)
800 } else {
801 Decimal::ZERO
802 };
803
804 let total_refunded: String = conn
806 .query_row(
807 "SELECT decimal_sum(refund_amount) FROM returns WHERE created_at >= ?1 AND created_at <= ?2",
808 [&start_str, &end_str],
809 |row| row.get(0),
810 )
811 .unwrap_or_else(|_| "0".to_string());
812
813 let mut stmt = conn
815 .prepare(
816 r"
817 SELECT reason, COUNT(*) as cnt
818 FROM returns
819 WHERE created_at >= ?1 AND created_at <= ?2
820 GROUP BY reason
821 ORDER BY cnt DESC
822 ",
823 )
824 .map_err(map_db_error)?;
825
826 let rows = stmt
827 .query_map([&start_str, &end_str], |row| {
828 let reason: String = row.get(0)?;
829 let count: i64 = row.get(1)?;
830 Ok((reason, count))
831 })
832 .map_err(map_db_error)?;
833
834 let mut by_reason = Vec::new();
835 for row in rows {
836 let (reason, count) = row.map_err(map_db_error)?;
837 let percentage = if total_returns > 0 {
838 Decimal::from(count * 100) / Decimal::from(total_returns)
839 } else {
840 Decimal::ZERO
841 };
842 by_reason.push(ReturnReasonCount { reason, count: count as u64, percentage });
843 }
844
845 let mut stmt = conn
847 .prepare(
848 r"
849 SELECT
850 ri.sku,
851 MAX(ri.name) as name,
852 SUM(ri.quantity) as units_returned,
853 COALESCE(
854 (SELECT SUM(oi.quantity) FROM order_items oi WHERE oi.sku = ri.sku),
855 0
856 ) as units_sold
857 FROM return_items ri
858 JOIN returns r ON ri.return_id = r.id
859 WHERE r.created_at >= ?1 AND r.created_at <= ?2
860 GROUP BY ri.sku
861 ORDER BY units_returned DESC
862 LIMIT 10
863 ",
864 )
865 .map_err(map_db_error)?;
866
867 let product_rows = stmt
868 .query_map([&start_str, &end_str], |row| {
869 let sku: String = row.get(0)?;
870 let name: String = row.get(1)?;
871 let units_returned: i64 = row.get(2)?;
872 let units_sold: i64 = row.get(3)?;
873 Ok((sku, name, units_returned, units_sold))
874 })
875 .map_err(map_db_error)?;
876
877 let mut top_returned_products = Vec::new();
878 for row in product_rows {
879 let (sku, name, units_returned, units_sold) = row.map_err(map_db_error)?;
880 let return_rate = if units_sold > 0 {
881 Decimal::from(units_returned * 100) / Decimal::from(units_sold)
882 } else {
883 Decimal::ZERO
884 };
885 top_returned_products.push(TopReturnedProduct {
886 sku,
887 name,
888 units_returned: units_returned as u64,
889 units_sold: units_sold as u64,
890 return_rate_percent: return_rate,
891 });
892 }
893
894 let total_refunded = parse_decimal_value(&total_refunded, "total_refunded")?;
895
896 Ok(ReturnMetrics {
897 total_returns: total_returns as u64,
898 return_rate_percent: return_rate,
899 total_refunded,
900 by_reason,
901 top_returned_products,
902 })
903 }
904
905 fn get_demand_forecast(
906 &self,
907 skus: Option<Vec<String>>,
908 days_ahead: u32,
909 ) -> Result<Vec<DemandForecast>> {
910 let conn = self.conn()?;
911 let days_back = 30; let start = (Utc::now() - Duration::days(days_back)).to_rfc3339();
913
914 let mut params: Vec<Box<dyn ToSql>> = vec![Box::new(start)];
916 let where_clause = match &skus {
917 Some(sku_list) if !sku_list.is_empty() => {
918 let placeholders = sku_list.iter().map(|_| "?").collect::<Vec<_>>().join(", ");
919 for sku in sku_list {
920 params.push(Box::new(sku.clone()));
921 }
922 format!("WHERE ii.sku IN ({placeholders})")
923 }
924 _ => String::new(),
925 };
926
927 let query = format!(
928 r"
929 SELECT
930 ii.sku,
931 ii.name,
932 COALESCE(SUM(CASE WHEN it.transaction_type = 'sale' THEN ABS(it.quantity) ELSE 0 END), 0) / {days_back} as avg_daily,
933 COALESCE(ib.quantity_on_hand, 0) - COALESCE(ib.quantity_allocated, 0) as current_stock
934 FROM inventory_items ii
935 LEFT JOIN inventory_balances ib ON ii.id = ib.item_id
936 LEFT JOIN inventory_transactions it ON ii.id = it.item_id AND it.created_at >= ?
937 {where_clause}
938 GROUP BY ii.id
939 HAVING avg_daily > 0 OR current_stock < 50
940 ORDER BY avg_daily DESC
941 LIMIT 50
942 "
943 );
944
945 let mut stmt = conn.prepare(&query).map_err(map_db_error)?;
946
947 let params_refs: Vec<&dyn ToSql> = params.iter().map(std::convert::AsRef::as_ref).collect();
948 let rows = stmt
949 .query_map(params_refs.as_slice(), |row| {
950 let sku: String = row.get(0)?;
951 let name: String = row.get(1)?;
952 let avg_daily: f64 = row.get(2)?;
953 let current_stock: f64 = row.get(3)?;
954 Ok((sku, name, avg_daily, current_stock))
955 })
956 .map_err(map_db_error)?;
957
958 let mut results = Vec::new();
959 for row in rows {
960 let (sku, name, avg_daily, current_stock) = row.map_err(map_db_error)?;
961 let avg_daily_dec = Decimal::from_f64_retain(avg_daily).unwrap_or(Decimal::ZERO);
962 let current_stock_dec =
963 Decimal::from_f64_retain(current_stock).unwrap_or(Decimal::ZERO);
964 let forecasted = avg_daily_dec * Decimal::from(days_ahead);
965
966 let days_until_stockout =
967 if avg_daily > 0.0 { Some((current_stock / avg_daily) as i32) } else { None };
968
969 let trend = if avg_daily > 1.0 {
971 Trend::Rising
972 } else if avg_daily < 0.5 {
973 Trend::Falling
974 } else {
975 Trend::Stable
976 };
977
978 results.push(DemandForecast {
979 sku,
980 name,
981 average_daily_demand: avg_daily_dec,
982 forecasted_demand: forecasted,
983 confidence: Decimal::new(7, 1), current_stock: current_stock_dec,
985 days_until_stockout,
986 recommended_reorder_qty: if days_until_stockout.is_some_and(|d| d < 14) {
987 Some(avg_daily_dec * Decimal::from(30)) } else {
989 None
990 },
991 recommended_reorder_date: None,
992 trend,
993 });
994 }
995
996 Ok(results)
997 }
998
999 fn get_revenue_forecast(
1000 &self,
1001 periods_ahead: u32,
1002 granularity: TimeGranularity,
1003 ) -> Result<Vec<RevenueForecast>> {
1004 let conn = self.conn()?;
1005
1006 let days_back = match granularity {
1008 TimeGranularity::Day => 90,
1009 TimeGranularity::Week => 180,
1010 TimeGranularity::Month => 365,
1011 _ => 365,
1012 };
1013
1014 let start = (Utc::now() - Duration::days(days_back)).to_rfc3339();
1015 let group_expr = match granularity {
1016 TimeGranularity::Day => "strftime('%Y-%m-%d', created_at)".to_string(),
1017 TimeGranularity::Week => ISO_WEEK_EXPR.to_string(),
1018 TimeGranularity::Month => "strftime('%Y-%m', created_at)".to_string(),
1019 _ => "strftime('%Y-%m', created_at)".to_string(),
1020 };
1021
1022 let (total_str, period_count): (String, i64) = conn
1028 .query_row(
1029 &format!(
1030 r"
1031 SELECT decimal_sum(period_revenue), COUNT(*) FROM (
1032 SELECT decimal_sum(total_amount) as period_revenue
1033 FROM orders
1034 WHERE created_at >= ?1
1035 AND status NOT IN ('cancelled', 'refunded')
1036 GROUP BY {group_expr}
1037 )
1038 "
1039 ),
1040 [&start],
1041 |row| Ok((row.get(0)?, row.get(1)?)),
1042 )
1043 .unwrap_or_else(|_| ("0".to_string(), 0));
1044
1045 let avg_revenue_dec = if period_count > 0 {
1046 parse_decimal_value(&total_str, "period_revenue")? / Decimal::from(period_count)
1047 } else {
1048 Decimal::ZERO
1049 };
1050
1051 let mut results = Vec::new();
1053 let variance = Decimal::new(15, 2); let one = Decimal::ONE;
1055 for i in 1..=periods_ahead {
1056 let period_label = format!("Period +{i}");
1057 let lower = avg_revenue_dec * (one - variance);
1058 let upper = avg_revenue_dec * (one + variance);
1059
1060 results.push(RevenueForecast {
1061 period: period_label,
1062 forecasted_revenue: avg_revenue_dec,
1063 lower_bound: lower,
1064 upper_bound: upper,
1065 confidence_level: Decimal::new(8, 1), based_on_periods: (days_back / 30) as u32,
1067 });
1068 }
1069
1070 Ok(results)
1071 }
1072
1073 fn get_sales_summary_batch(&self, queries: Vec<AnalyticsQuery>) -> Result<Vec<SalesSummary>> {
1074 validate_batch_size(&queries)?;
1075 let mut results = Vec::with_capacity(queries.len());
1076 for query in queries {
1077 results.push(self.get_sales_summary(query)?);
1078 }
1079 Ok(results)
1080 }
1081}
1082
1083#[cfg(test)]
1084mod tests {
1085 use super::*;
1086 use crate::SqliteDatabase;
1087 use rust_decimal_macros::dec;
1088 use stateset_core::{AnalyticsQuery, AnalyticsRepository, TimePeriod};
1089
1090 fn fresh_repo() -> SqliteAnalyticsRepository {
1091 SqliteDatabase::in_memory().expect("in-memory").analytics()
1092 }
1093
1094 fn last_30_days() -> AnalyticsQuery {
1095 AnalyticsQuery::new().period(TimePeriod::Last30Days)
1096 }
1097
1098 #[test]
1099 fn iso_week_expr_matches_postgres_iso_labels() {
1100 let db = SqliteDatabase::in_memory().expect("in-memory");
1106 let conn = db.pool().get().expect("conn");
1107 let sql = format!("SELECT {ISO_WEEK_EXPR} FROM (SELECT ?1 AS created_at)");
1108 for (date, expected) in [
1109 ("2023-01-01", "2022-W52"), ("2023-01-02", "2023-W01"), ("2023-12-31", "2023-W52"), ("2024-12-30", "2025-W01"), ] {
1114 let got: String = conn.query_row(&sql, [date], |row| row.get(0)).expect("query");
1115 assert_eq!(got, expected, "ISO week label for {date}");
1116 }
1117 }
1118
1119 #[test]
1120 fn empty_db_sales_summary_is_zero() {
1121 let repo = fresh_repo();
1122 let summary = repo.get_sales_summary(last_30_days()).expect("ok");
1123 assert_eq!(summary.total_revenue, dec!(0));
1124 assert_eq!(summary.order_count, 0);
1125 }
1126
1127 #[test]
1128 fn empty_db_revenue_by_period_is_empty_or_zero() {
1129 let repo = fresh_repo();
1130 let rows = repo.get_revenue_by_period(last_30_days()).expect("ok");
1131 assert!(rows.iter().all(|r| r.revenue == dec!(0) && r.order_count == 0));
1133 }
1134
1135 #[test]
1136 fn empty_db_top_products_is_empty() {
1137 let repo = fresh_repo();
1138 let rows = repo.get_top_products(last_30_days()).expect("ok");
1139 assert!(rows.is_empty());
1140 }
1141
1142 #[test]
1143 fn empty_db_product_performance_is_empty() {
1144 let repo = fresh_repo();
1145 let rows = repo.get_product_performance(last_30_days()).expect("ok");
1146 assert!(rows.is_empty());
1147 }
1148
1149 #[test]
1150 fn empty_db_customer_metrics_zero() {
1151 let repo = fresh_repo();
1152 let m = repo.get_customer_metrics(last_30_days()).expect("ok");
1153 assert_eq!(m.total_customers, 0);
1154 assert_eq!(m.new_customers, 0);
1155 }
1156
1157 #[test]
1158 fn empty_db_top_customers_is_empty() {
1159 let repo = fresh_repo();
1160 let rows = repo.get_top_customers(last_30_days()).expect("ok");
1161 assert!(rows.is_empty());
1162 }
1163
1164 #[test]
1165 fn empty_db_inventory_health_zero() {
1166 let repo = fresh_repo();
1167 let h = repo.get_inventory_health().expect("ok");
1168 assert_eq!(h.total_skus, 0);
1169 assert_eq!(h.total_value, dec!(0));
1170 }
1171
1172 #[test]
1173 fn empty_db_low_stock_items_is_empty() {
1174 let repo = fresh_repo();
1175 let rows = repo.get_low_stock_items(Some(dec!(10))).expect("ok");
1176 assert!(rows.is_empty());
1177 }
1178
1179 fn seed_item_with_balances(repo: &SqliteAnalyticsRepository, sku: &str, quantities: &[&str]) {
1182 let conn = repo.conn().expect("conn");
1183 conn.execute(
1184 "INSERT INTO inventory_items (sku, name) VALUES (?1, ?2)",
1185 rusqlite::params![sku, format!("Item {sku}")],
1186 )
1187 .expect("insert item");
1188 let item_id = conn.last_insert_rowid();
1189 for (i, qty) in quantities.iter().enumerate() {
1190 let location_id = (i + 1) as i64;
1191 conn.execute(
1192 "INSERT OR IGNORE INTO inventory_locations (id, name, code) VALUES (?1, ?2, ?3)",
1193 rusqlite::params![
1194 location_id,
1195 format!("Loc {location_id}"),
1196 format!("LOC-{location_id}")
1197 ],
1198 )
1199 .expect("insert location");
1200 conn.execute(
1201 "INSERT INTO inventory_balances (item_id, location_id, quantity_on_hand)
1202 VALUES (?1, ?2, ?3)",
1203 rusqlite::params![item_id, location_id, qty],
1204 )
1205 .expect("insert balance");
1206 }
1207 }
1208
1209 #[test]
1210 fn low_stock_items_sum_quantities_exactly() {
1211 let repo = fresh_repo();
1212 seed_item_with_balances(&repo, "LOW-A", &["0.1", "0.2"]);
1215
1216 let rows = repo.get_low_stock_items(Some(dec!(0.3))).expect("ok");
1217 assert_eq!(rows.len(), 1);
1218 assert_eq!(rows[0].sku, "LOW-A");
1219 assert_eq!(rows[0].on_hand, dec!(0.3));
1220 assert_eq!(rows[0].available, dec!(0.3));
1221 }
1222
1223 #[test]
1224 fn low_stock_threshold_comparison_is_exact_not_float() {
1225 let repo = fresh_repo();
1226 seed_item_with_balances(&repo, "LOW-B", &["10.000000000000000001"]);
1229
1230 let rows = repo.get_low_stock_items(Some(dec!(10))).expect("ok");
1231 assert!(rows.iter().all(|r| r.sku != "LOW-B"), "10.000000000000000001 > 10 exactly");
1232 }
1233
1234 #[test]
1235 fn low_stock_items_sort_numerically_not_lexicographically() {
1236 let repo = fresh_repo();
1237 seed_item_with_balances(&repo, "LOW-TEN", &["10"]);
1238 seed_item_with_balances(&repo, "LOW-TWO", &["2"]);
1239
1240 let rows = repo.get_low_stock_items(Some(dec!(100))).expect("ok");
1241 let skus: Vec<&str> = rows.iter().map(|r| r.sku.as_str()).collect();
1242 assert_eq!(skus, vec!["LOW-TWO", "LOW-TEN"], "2 sorts before 10 numerically");
1243 }
1244
1245 #[test]
1246 fn empty_db_inventory_movement_is_empty() {
1247 let repo = fresh_repo();
1248 let rows = repo.get_inventory_movement(last_30_days()).expect("ok");
1249 assert!(rows.is_empty());
1250 }
1251
1252 #[test]
1253 fn empty_db_order_status_breakdown_returns_zero_counts() {
1254 let repo = fresh_repo();
1255 let b = repo.get_order_status_breakdown(last_30_days()).expect("ok");
1256 assert_eq!(b.pending, 0);
1257 assert_eq!(b.shipped, 0);
1258 assert_eq!(b.delivered, 0);
1259 }
1260
1261 #[test]
1262 fn empty_db_fulfillment_metrics_zero_workload() {
1263 let repo = fresh_repo();
1264 let m = repo.get_fulfillment_metrics(last_30_days()).expect("ok");
1265 assert_eq!(m.shipped_today, 0);
1266 assert_eq!(m.awaiting_shipment, 0);
1267 }
1268
1269 #[test]
1270 fn empty_db_return_metrics_zero() {
1271 let repo = fresh_repo();
1272 let m = repo.get_return_metrics(last_30_days()).expect("ok");
1273 assert_eq!(m.total_returns, 0);
1274 }
1275
1276 #[test]
1277 fn batch_query_returns_one_summary_per_input() {
1278 let repo = fresh_repo();
1279 let queries = vec![last_30_days(), last_30_days(), last_30_days()];
1280 let summaries = repo.get_sales_summary_batch(queries).expect("ok");
1281 assert_eq!(summaries.len(), 3);
1282 assert!(summaries.iter().all(|s| s.total_revenue == dec!(0)));
1283 }
1284}