Spaces:
Runtime error
Runtime error
| 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() | |