import threading import time import random import logging from datetime import datetime from sqlalchemy.orm import Session from backend.app.database import SessionLocal from backend.app import crud, models, schemas from backend.app.services.events import event_bus from backend.app.services.pricer import pricer from backend.app.services.anomaly import anomaly_detector from backend.app.services.forecaster import forecaster logging.basicConfig(level=logging.INFO) logger = logging.getLogger("MercurySimulator") class OrderSimulator: def __init__(self): self.is_running = False self.speed = 1.0 # Time multiplier self.demand_type = "normal" # normal, surge, seasonal, flash_sale self._thread = None self._lock = threading.Lock() # Expanded Enterprise State Parameters self.demand_multiplier = 1.0 self.seasonality_level = "medium" # low, medium, high self.supply_level = "medium" # low, medium, high self.flash_sale_active = False self.competitor_discount_active = False self.scenario = "Normal Demand" self.orders_per_minute = 2500 # Track engine ticks self.tick_counter = 0 def start(self, speed: float = 1.0, demand_type: str = "normal", demand_multiplier: float = 1.0, seasonality_level: str = "medium", supply_level: str = "medium", flash_sale_active: bool = False, competitor_discount_active: bool = False, scenario: str = "Normal Demand"): with self._lock: self.speed = speed self.demand_type = demand_type self.demand_multiplier = demand_multiplier self.seasonality_level = seasonality_level self.supply_level = supply_level self.flash_sale_active = flash_sale_active self.competitor_discount_active = competitor_discount_active self.scenario = scenario self.orders_per_minute = int(2500 * demand_multiplier) if not self.is_running: self.is_running = True self._thread = threading.Thread(target=self._run_loop, daemon=True) self._thread.start() logger.info(f"Simulator STARTED (Speed: {speed}x, Scenario: {scenario})") def update_params(self, speed: float = None, demand_type: str = None, demand_multiplier: float = None, seasonality_level: str = None, supply_level: str = None, flash_sale_active: bool = None, competitor_discount_active: bool = None, scenario: str = None): with self._lock: if speed is not None: self.speed = speed if demand_type is not None: self.demand_type = demand_type if demand_multiplier is not None: self.demand_multiplier = demand_multiplier if seasonality_level is not None: self.seasonality_level = seasonality_level if supply_level is not None: self.supply_level = supply_level if flash_sale_active is not None: self.flash_sale_active = flash_sale_active if competitor_discount_active is not None: self.competitor_discount_active = competitor_discount_active if scenario is not None: self.scenario = scenario self.orders_per_minute = int(2500 * self.demand_multiplier * self.speed) logger.info(f"Simulator params updated. Scenario: {self.scenario}, Orders/Min: {self.orders_per_minute}") def stop(self): with self._lock: if self.is_running: self.is_running = False logger.info("Simulator STOPPED") def _run_loop(self): while True: # Check if stopped with self._lock: if not self.is_running: break speed_val = self.speed demand_val = self.demand_type demand_mult = self.demand_multiplier # Sleep interval scaled by speed and demand multiplier # Base interval: 5 seconds. sleep_time = max(0.1, (5.0 / speed_val) / demand_mult) # Add random variance time.sleep(sleep_time * random.uniform(0.7, 1.3)) # Generate simulated order try: self._generate_simulated_order(demand_val) except Exception as e: logger.error(f"Simulator error generating order: {e}") # Periodically tick optimization engines in background self.tick_counter += 1 if self.tick_counter % 5 == 0: # every 5 ticks run pricing self._tick_pricing_engine() if self.tick_counter % 10 == 0: # every 10 ticks run anomalies self._tick_anomaly_engine() if self.tick_counter % 20 == 0: # every 20 ticks run forecasting self._tick_forecasting_engine() def _generate_simulated_order(self, demand_type: str): db: Session = SessionLocal() try: products = db.query(models.Product).all() if not products: return items_to_order = [] order_type = "normal" # Handle Supply constraints: if supply is low, reduce random replenishment amounts later # and deplete stock faster by purchasing higher quantity margins qty_multiplier = 2 if self.supply_level == "low" else 1 if demand_type == "surge": order_type = "surge" num_items = random.randint(2, 4) chosen_prods = random.sample(products, k=min(num_items, len(products))) for prod in chosen_prods: qty = random.randint(3, 8) * qty_multiplier items_to_order.append({"product_id": prod.id, "quantity": qty}) elif demand_type == "flash_sale": order_type = "flash_sale" # Focus on Wireless Headphones target_prod = db.query(models.Product).filter(models.Product.sku == "NV-HDPH-01").first() if target_prod: qty = random.randint(5, 12) * qty_multiplier items_to_order.append({"product_id": target_prod.id, "quantity": qty}) if random.random() < 0.3: other_prods = [p for p in products if p.sku != "NV-HDPH-01"] if other_prods: items_to_order.append({"product_id": random.choice(other_prods).id, "quantity": random.randint(1, 2)}) elif demand_type == "seasonal": order_type = "seasonal" # Focus on Monitors (tech gifting) gift_prod = db.query(models.Product).filter(models.Product.sku == "NV-MONI-05").first() if gift_prod: qty = random.randint(3, 7) * qty_multiplier items_to_order.append({"product_id": gift_prod.id, "quantity": qty}) num_items = random.randint(1, 2) other_prods = [p for p in products if p.sku != "NV-MONI-05"] if other_prods: chosen = random.sample(other_prods, k=min(num_items, len(other_prods))) for prod in chosen: items_to_order.append({"product_id": prod.id, "quantity": random.randint(1, 3)}) else: # Normal Demand num_items = random.randint(1, 2) chosen_prods = random.sample(products, k=min(num_items, len(products))) for prod in chosen_prods: qty = random.randint(1, 3) * qty_multiplier items_to_order.append({"product_id": prod.id, "quantity": qty}) # Formulate order request schemas items_schemas = [] for item in items_to_order: # Validate if stock is available total_stock = db.query(func.sum(models.Inventory.available_stock)).filter( models.Inventory.product_id == item["product_id"] ).scalar() or 0 if total_stock > item["quantity"]: items_schemas.append(schemas.OrderItemCreate( product_id=item["product_id"], quantity=item["quantity"] )) if not items_schemas: # All products out of stock, trigger manual replenishment/restock simulated action self._simulate_warehouse_restock(db) return order_in = schemas.OrderCreate(items=items_schemas, demand_type=order_type) db_order = crud.create_order(db, order_in) # Serialize for Event Bus serialized_items = [] total_rev = 0.0 for it in db_order.items: serialized_items.append({ "product_id": it.product_id, "sku": it.product.sku, "quantity": it.quantity, "price_per_unit": it.price_per_unit, "revenue": it.revenue }) total_rev += it.revenue event_bus.publish("orders", { "order_id": db_order.id, "order_num": db_order.order_num, "demand_type": db_order.demand_type, "items": serialized_items, "total_revenue": total_rev, "timestamp": str(db_order.timestamp) }) finally: db.close() def generate_bulk_orders(self, count: int): """ Generates a massive bulk demand surge in a rapid simulated iteration. """ logger.info(f"Bulk generating {count} simulated order units of demand...") db: Session = SessionLocal() try: products = db.query(models.Product).all() if not products: return # To generate the required volume rapidly without DB locking or freezing, # we place 100 scaled orders. scale = max(1, count // 150) for o_idx in range(100): items_schemas = [] # Pick 1-3 random products chosen = random.sample(products, k=random.randint(1, 3)) for prod in chosen: qty = random.randint(2, 5) * scale # Ensure availability total_stock = db.query(func.sum(models.Inventory.available_stock)).filter( models.Inventory.product_id == prod.id ).scalar() or 0 # Force restock if low to guarantee order generation succeeds if total_stock < qty: inventories = db.query(models.Inventory).filter(models.Inventory.product_id == prod.id).all() for inv in inventories: crud.update_inventory_stock(db, inv.id, qty * 2, "restock", "Pre-bulk Generation Seed") db.commit() total_stock = qty * 2 items_schemas.append(schemas.OrderItemCreate( product_id=prod.id, quantity=qty )) order_in = schemas.OrderCreate(items=items_schemas, demand_type=self.demand_type) crud.create_order(db, order_in) db.commit() # Force immediate recalculation of all engines pricer.run_dynamic_pricing(db) forecaster.run_forecasts(db) anomaly_detector.run_anomaly_detection(db) logger.info("Bulk generation recalculations complete.") finally: db.close() def _simulate_warehouse_restock(self, db: Session): """ Simulates delivery of stock from suppliers when warehouses run empty. """ logger.info("Out of stock detected across SKUs. Simulating supplier replenishment restocks...") inventories = db.query(models.Inventory).all() # Under low supply settings, suppliers deliver smaller quantities multiplier = 0.4 if self.supply_level == "low" else (1.6 if self.supply_level == "high" else 1.0) for inv in inventories: if inv.current_stock <= inv.safety_stock: restock_qty = int(random.randint(150, 300) * multiplier) crud.update_inventory_stock( db, inventory_id=inv.id, qty_change=restock_qty, change_type="restock", reason="Simulated Supplier Delivery" ) # Publish restock event event_bus.publish("inventory", { "product_id": inv.product_id, "product_sku": inv.product.sku, "event_type": "restocked", "restock_qty": restock_qty, "warehouse_name": inv.warehouse.name, "timestamp": str(datetime.now()) }) def _tick_pricing_engine(self): db = SessionLocal() try: pricer.run_dynamic_pricing(db) finally: db.close() def _tick_anomaly_engine(self): db = SessionLocal() try: anomaly_detector.run_anomaly_detection(db) finally: db.close() def _tick_forecasting_engine(self): db = SessionLocal() try: forecaster.run_forecasts(db) finally: db.close() from sqlalchemy import func # Singleton instance order_simulator = OrderSimulator()