"""Celery tasks: sync_inventory (every 5min), sync_tracking, recalc_prices. Broker/backend from REDIS_URL (default redis://redis:6379/0 inside Docker). """ import os from datetime import datetime, timedelta, timezone from celery import Celery REDIS_URL = os.getenv("REDIS_URL", "redis://redis:6379/0") celery_app = Celery("polaris", broker=REDIS_URL, backend=REDIS_URL) celery_app.conf.timezone = "UTC" celery_app.conf.beat_schedule = { "sync-inventory-every-5m": { "task": "app.tasks.sync_inventory", "schedule": 300.0, }, "sync-tracking-every-5m": { "task": "app.tasks.sync_tracking", "schedule": 300.0, }, "recalc-prices-hourly": { "task": "app.tasks.recalc_prices", "schedule": 3600.0, }, } @celery_app.task(name="app.tasks.sync_inventory") def sync_inventory(supplier_id=None): from app.database import SessionLocal from app.engines.inventory import sync_inventory as _sync db = SessionLocal() try: return _sync(db) finally: db.close() @celery_app.task(name="app.tasks.sync_tracking") def sync_tracking(): """MVP tracking progression: - confirmed + tracking -> shipped - shipped, older than 7d -> delivered """ from app.database import SessionLocal from app.models import AuditLog, Order db = SessionLocal() try: stats = {"checked": 0, "shipped": 0, "delivered": 0} now = datetime.now(timezone.utc) for order in db.query(Order).filter( Order.status.in_(["confirmed", "shipped"]) ).all(): stats["checked"] += 1 if order.status == "confirmed" and order.tracking: order.status = "shipped" stats["shipped"] += 1 db.add( AuditLog( actor="tasks.sync_tracking", action="order_shipped", entity="order", entity_id=order.order_number, detail={"tracking": order.tracking}, ) ) elif order.status == "shipped": age = now - (order.updated_at or now) if age > timedelta(days=7): order.status = "delivered" stats["delivered"] += 1 db.add( AuditLog( actor="tasks.sync_tracking", action="order_delivered", entity="order", entity_id=order.order_number, ) ) db.commit() return stats finally: db.close() @celery_app.task(name="app.tasks.recalc_prices") def recalc_prices(): from app.database import SessionLocal from app.engines.pricing import calculate_price from app.models import PriceHistory, Product db = SessionLocal() try: updated = 0 for product in db.query(Product).filter(Product.status != "sold").all(): quote = calculate_price(db, product.cost) if quote is None: continue # no active rule -> leave alone, never publish invalid product.retail_price = quote["retail"] if product.min_price is None or quote["retail"] < product.min_price: product.min_price = quote["retail"] if product.max_price is None or quote["retail"] > product.max_price: product.max_price = quote["retail"] db.add(PriceHistory(product_id=product.id, price=quote["retail"], reason="auto")) updated += 1 db.commit() return {"recalculated": updated} finally: db.close()