114 lines
3.7 KiB
Python
114 lines
3.7 KiB
Python
"""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()
|