From 96b4fcaac93cbca8ee5b7b3495f9229d474b957a Mon Sep 17 00:00:00 2001 From: drjones Date: Tue, 25 Aug 2026 20:37:58 -0700 Subject: [PATCH] Add Polaris FastAPI backend: models/schemas, pricing+scoring+order router, importer, Celery tasks, API routers, seed script --- backend/.gitignore | 7 + backend/Dockerfile | 16 ++ backend/README.md | 53 ++++++ backend/app/__init__.py | 1 + backend/app/database.py | 30 +++ backend/app/engines/__init__.py | 1 + backend/app/engines/importer.py | 115 ++++++++++++ backend/app/engines/inventory.py | 72 ++++++++ backend/app/engines/order_router.py | 122 ++++++++++++ backend/app/engines/pricing.py | 87 +++++++++ backend/app/engines/scoring.py | 70 +++++++ backend/app/engines/suppliers/__init__.py | 1 + backend/app/engines/suppliers/base.py | 65 +++++++ backend/app/engines/suppliers/csv_adapter.py | 78 ++++++++ backend/app/engines/suppliers/sample.py | 89 +++++++++ backend/app/main.py | 38 ++++ backend/app/models.py | 185 +++++++++++++++++++ backend/app/routers/__init__.py | 1 + backend/app/routers/admin.py | 160 ++++++++++++++++ backend/app/routers/analytics.py | 123 ++++++++++++ backend/app/routers/auth.py | 84 +++++++++ backend/app/routers/customers.py | 42 +++++ backend/app/routers/health.py | 23 +++ backend/app/routers/orders.py | 113 +++++++++++ backend/app/routers/products.py | 132 +++++++++++++ backend/app/routers/suppliers.py | 76 ++++++++ backend/app/schemas.py | 140 ++++++++++++++ backend/app/tasks.py | 113 +++++++++++ backend/app/version.py | 2 + backend/requirements.txt | 9 + backend/scripts/celery_worker.sh | 4 + backend/scripts/seed.py | 159 ++++++++++++++++ 32 files changed, 2211 insertions(+) create mode 100644 backend/.gitignore create mode 100644 backend/Dockerfile create mode 100644 backend/README.md create mode 100644 backend/app/__init__.py create mode 100644 backend/app/database.py create mode 100644 backend/app/engines/__init__.py create mode 100644 backend/app/engines/importer.py create mode 100644 backend/app/engines/inventory.py create mode 100644 backend/app/engines/order_router.py create mode 100644 backend/app/engines/pricing.py create mode 100644 backend/app/engines/scoring.py create mode 100644 backend/app/engines/suppliers/__init__.py create mode 100644 backend/app/engines/suppliers/base.py create mode 100644 backend/app/engines/suppliers/csv_adapter.py create mode 100644 backend/app/engines/suppliers/sample.py create mode 100644 backend/app/main.py create mode 100644 backend/app/models.py create mode 100644 backend/app/routers/__init__.py create mode 100644 backend/app/routers/admin.py create mode 100644 backend/app/routers/analytics.py create mode 100644 backend/app/routers/auth.py create mode 100644 backend/app/routers/customers.py create mode 100644 backend/app/routers/health.py create mode 100644 backend/app/routers/orders.py create mode 100644 backend/app/routers/products.py create mode 100644 backend/app/routers/suppliers.py create mode 100644 backend/app/schemas.py create mode 100644 backend/app/tasks.py create mode 100644 backend/app/version.py create mode 100644 backend/requirements.txt create mode 100755 backend/scripts/celery_worker.sh create mode 100644 backend/scripts/seed.py diff --git a/backend/.gitignore b/backend/.gitignore new file mode 100644 index 0000000..260c65c --- /dev/null +++ b/backend/.gitignore @@ -0,0 +1,7 @@ +__pycache__/ +*.pyc +.venv/ +venv/ +.env +.DS_Store +celerybeat-schedule* diff --git a/backend/Dockerfile b/backend/Dockerfile new file mode 100644 index 0000000..9c21c13 --- /dev/null +++ b/backend/Dockerfile @@ -0,0 +1,16 @@ +FROM python:3.12-slim + +WORKDIR /app + +ENV PYTHONDONTWRITEBYTECODE=1 \ + PYTHONUNBUFFERED=1 + +COPY requirements.txt . +RUN pip install --no-cache-dir -r requirements.txt + +COPY app ./app +COPY scripts ./scripts + +EXPOSE 8000 + +CMD ["uvicorn", "app.main:app", "--host", "0.0.0.0", "--port", "8000"] diff --git a/backend/README.md b/backend/README.md new file mode 100644 index 0000000..8df7870 --- /dev/null +++ b/backend/README.md @@ -0,0 +1,53 @@ +# Polaris Backend + +FastAPI + SQLAlchemy 2.0 + Pydantic v2 + Celery/Redis backend for the Polaris +reseller platform. PostgreSQL schema is canonical in `/opt/polaris/db/schema.sql` +(not managed by this service). + +## Run (CT host, no Docker) + +```sh +cd /opt/polaris/backend +python3 -m venv .venv && .venv/bin/pip install -r requirements.txt +# .env at /opt/polaris/.env provides POSTGRES_PASSWORD (auto-loaded) +.venv/bin/python scripts/seed.py # sample supplier + 27 products +DATABASE_URL=postgresql+psycopg2://reseller:${POSTGRES_PASSWORD}@127.0.0.1:5432/reseller \ + .venv/bin/uvicorn app.main:app --host 0.0.0.0 --port 8000 +``` + +## Run (Docker, via OPS compose) + +The OPS agent adds `api` + `worker` services. Inside the Docker network: +- `DATABASE_URL=postgresql+psycopg2://reseller:${POSTGRES_PASSWORD}@postgres:5432/reseller` +- `REDIS_URL=redis://redis:6379/0` +- worker command: `scripts/celery_worker.sh` (worker + beat; beat schedule in `app/tasks.py`) + +## Env vars + +| var | purpose | +|---|---| +| `DATABASE_URL` | SQLAlchemy URL (default: local 127.0.0.1:5432 w/ `POSTGRES_PASSWORD`) | +| `REDIS_URL` | Celery broker (default `redis://redis:6379/0`) | +| `SECRET_KEY` | JWT signing key (dev default — set in prod) | +| `ADMIN_EMAIL` / `ADMIN_PASSWORD` | admin login (plain compare — TODO: hash `ADMIN_PASSWORD_HASH` with bcrypt/argon2) | + +## API + +All under `/api` (see CONTRACT.md for the full list). Interactive docs at `/docs`. + +Key notes: +- Importer NEVER auto-publishes: new products enter as `status='imported'`; + `POST /api/products/{id}/publish` requires a calculated `retail_price`. +- Pricing: active `price_rules` bands (cost → markup), floor at + cost+shipping+fees, `min_price` respected; `estimate_net_profit` = + retail − cost − shipping − fees. +- Orders: `POST /api/orders` routes each item to the highest-scoring eligible + supplier (40% price / 20% inventory / 15% speed / 10% fulfillment / 10% + returns / 5% history) and records profit. +- Analytics `conversion` is 0.0 until a traffic/visitor source exists. + +## Seed + +`scripts/seed.py` — idempotent. Creates the sample supplier, imports 27 +products (home/gadgets/office, $50–$500 retail), prices via the pricing engine, +and publishes the valid ones. diff --git a/backend/app/__init__.py b/backend/app/__init__.py new file mode 100644 index 0000000..c84f120 --- /dev/null +++ b/backend/app/__init__.py @@ -0,0 +1 @@ +"""Polaris backend application package.""" diff --git a/backend/app/database.py b/backend/app/database.py new file mode 100644 index 0000000..9d2e015 --- /dev/null +++ b/backend/app/database.py @@ -0,0 +1,30 @@ +"""Database engine/session wiring. + +DATABASE_URL comes from env (set by docker-compose). On the CT host it falls +back to the local postgres port mapping using POSTGRES_PASSWORD from .env. +""" +import os + +from dotenv import load_dotenv +from sqlalchemy import create_engine +from sqlalchemy.orm import declarative_base, sessionmaker + +# /opt/polaris/.env exists on the CT host; harmless no-op elsewhere. +load_dotenv(os.getenv("POLARIS_ENV_FILE", "/opt/polaris/.env"), override=False) + +DATABASE_URL = os.getenv( + "DATABASE_URL", + f"postgresql+psycopg2://reseller:{os.getenv('POSTGRES_PASSWORD', '')}@127.0.0.1:5432/reseller", +) + +engine = create_engine(DATABASE_URL, pool_pre_ping=True, pool_size=5, max_overflow=10) +SessionLocal = sessionmaker(autocommit=False, autoflush=False, bind=engine) +Base = declarative_base() + + +def get_db(): + db = SessionLocal() + try: + yield db + finally: + db.close() diff --git a/backend/app/engines/__init__.py b/backend/app/engines/__init__.py new file mode 100644 index 0000000..c8f4773 --- /dev/null +++ b/backend/app/engines/__init__.py @@ -0,0 +1 @@ +"""Polaris business engines.""" diff --git a/backend/app/engines/importer.py b/backend/app/engines/importer.py new file mode 100644 index 0000000..b69a32d --- /dev/null +++ b/backend/app/engines/importer.py @@ -0,0 +1,115 @@ +"""Feed importer: normalize -> validate -> upsert. + +NEVER auto-publishes products: new products are created with status 'imported' +only. Invalid items (missing sku/title, non-positive cost, negative inventory) +are skipped and logged to audit_logs. +""" +from decimal import Decimal + +from app.engines.suppliers.base import build_adapter +from app.models import AuditLog, Product, SupplierProduct + + +def validate(item): + problems = [] + if not (item.get("sku") or "").strip(): + problems.append("missing sku") + if not (item.get("title") or "").strip(): + problems.append("missing title") + try: + cost = float(item.get("cost") or 0) + if cost <= 0: + problems.append("non-positive cost") + except (TypeError, ValueError): + problems.append("invalid cost") + try: + inventory = int(float(item.get("inventory") or 0)) + if inventory < 0: + problems.append("negative inventory") + except (TypeError, ValueError): + problems.append("invalid inventory") + return problems + + +def import_feed(db, supplier, feed_type=None, data_or_url=""): + """Import a supplier feed into products + supplier_products. + + Returns summary dict: {supplier, received, imported, updated, skipped: [...]} + """ + adapter = build_adapter(supplier, data_or_url=data_or_url) + items = adapter.get_products() + + summary = { + "supplier": supplier.name, + "supplier_id": str(supplier.id), + "received": len(items), + "imported": 0, + "updated": 0, + "skipped": [], + } + + for item in items: + problems = validate(item) + if problems: + reason = "; ".join(problems) + summary["skipped"].append({"sku": item.get("sku"), "reason": reason}) + db.add( + AuditLog( + actor="importer", + action="import_skipped", + entity="product", + entity_id=item.get("sku"), + detail={"supplier": supplier.name, "reason": reason}, + ) + ) + continue + + product = db.query(Product).filter(Product.sku == item["sku"]).first() + if product is None: + product = Product(sku=item["sku"], status="imported") + db.add(product) + summary["imported"] += 1 + else: + summary["updated"] += 1 + + # set all required fields BEFORE flushing (title is NOT NULL in schema) + product.title = item["title"] + product.description = item.get("description") or product.description + product.category = item.get("category") or product.category + product.brand = item.get("brand") or product.brand + product.cost = Decimal(str(item["cost"])) + + if item.get("image_url"): + images = list(product.images or []) + if item["image_url"] not in images: + images.append(item["image_url"]) + product.images = images + + db.flush() # product.id available now + + sp = ( + db.query(SupplierProduct) + .filter( + SupplierProduct.supplier_id == supplier.id, + SupplierProduct.supplier_sku == item["sku"], + ) + .first() + ) + if sp is None: + sp = SupplierProduct( + supplier_id=supplier.id, + supplier_sku=item["sku"], + product_id=product.id, + ) + db.add(sp) + db.flush() + else: + sp.product_id = product.id + + sp.supplier_cost = Decimal(str(item["cost"])) + sp.supplier_inventory = int(item.get("inventory") or 0) + sp.shipping_cost = Decimal(str(item.get("shipping_cost") or 0)) + sp.shipping_time = item.get("shipping_time") + + db.commit() + return summary diff --git a/backend/app/engines/inventory.py b/backend/app/engines/inventory.py new file mode 100644 index 0000000..77d4101 --- /dev/null +++ b/backend/app/engines/inventory.py @@ -0,0 +1,72 @@ +"""Inventory sync engine. + +Refreshes supplier_products inventory from each adapter, writes inventory_history, +and flags zero/unavailable stock (pauses published products + audit log alert). +""" +from datetime import datetime, timezone + +from app.engines.suppliers.base import build_adapter +from app.models import AuditLog, InventoryHistory, Supplier, SupplierProduct + + +def sync_inventory(db, supplier=None): + if supplier is not None: + suppliers = [supplier] + else: + suppliers = db.query(Supplier).filter(Supplier.status == "active").all() + + stats = {"suppliers": len(suppliers), "checked": 0, "changed": 0, "zeroed": 0} + + for s in suppliers: + try: + adapter = build_adapter(s) + current = {r["supplier_sku"]: int(r["inventory"]) for r in adapter.get_inventory()} + except Exception as exc: # noqa: BLE001 — log and continue to next supplier + db.add( + AuditLog( + actor="inventory", + action="sync_failed", + entity="supplier", + entity_id=str(s.id), + detail={"supplier": s.name, "error": str(exc)}, + ) + ) + continue + + for sp in db.query(SupplierProduct).filter(SupplierProduct.supplier_id == s.id).all(): + stats["checked"] += 1 + new_value = current.get(sp.supplier_sku) + if new_value is None: + continue + + if sp.supplier_inventory != new_value: + stats["changed"] += 1 + sp.supplier_inventory = new_value + sp.last_updated = datetime.now(timezone.utc) + db.add(InventoryHistory(supplier_product_id=sp.id, inventory=new_value)) + + if new_value <= 0: + stats["zeroed"] += 1 + db.add( + AuditLog( + actor="inventory", + action="inventory_zero", + entity="supplier_product", + entity_id=str(sp.id), + detail={"supplier_sku": sp.supplier_sku, "inventory": 0}, + ) + ) + if sp.product is not None and sp.product.status == "published": + sp.product.status = "paused" + db.add( + AuditLog( + actor="inventory", + action="product_paused_zero_inventory", + entity="product", + entity_id=str(sp.product.id), + detail={"sku": sp.product.sku}, + ) + ) + + db.commit() + return stats diff --git a/backend/app/engines/order_router.py b/backend/app/engines/order_router.py new file mode 100644 index 0000000..3eeaaaa --- /dev/null +++ b/backend/app/engines/order_router.py @@ -0,0 +1,122 @@ +"""Order router — select highest-scoring supplier and record profit. + +select_supplier: eligible = active supplier, linked supplier_product, stock >= qty. +route_order: assigns supplier, computes supplier_cost / shipping_cost / fees, +records net profit = retail_total - supplier_cost - shipping_cost - fees. +""" +from decimal import ROUND_HALF_UP, Decimal + +from app.engines import scoring +from app.models import AuditLog, Product, Supplier, SupplierProduct + +CENT = Decimal("0.01") +FEE_RATE = Decimal("0.05") # payment processing estimate +FEE_FLAT = Decimal("0.30") + + +def _f(value) -> float: + try: + return float(value or 0) + except (TypeError, ValueError): + return 0.0 + + +def select_supplier(db, product_id, qty): + """Return (result_dict, error_str). result = {supplier, supplier_product, score}.""" + product = db.query(Product).filter(Product.id == product_id).first() + if product is None: + return None, "product not found" + + rows = ( + db.query(SupplierProduct, Supplier) + .join(Supplier, SupplierProduct.supplier_id == Supplier.id) + .filter( + SupplierProduct.product_id == product_id, + Supplier.status == "active", + SupplierProduct.supplier_inventory >= qty, + ) + .all() + ) + if not rows: + return None, "no eligible supplier" + + pool = {"avg_cost": sum(_f(sp.supplier_cost) for sp, _ in rows) / len(rows)} + scored = [ + (scoring.score_supplier(db, supplier, sp, qty, pool), sp, supplier) + for sp, supplier in rows + ] + best_score, best_sp, best_supplier = max(scored, key=lambda t: t[0]) + return { + "supplier": best_supplier, + "supplier_product": best_sp, + "score": best_score, + }, None + + +def route_order(db, order): + """Route a newly created order to the best supplier(s) and record profit.""" + supplier_id = None + supplier_cost = Decimal("0") + shipping_cost = Decimal("0") + routing_detail = [] + + for item in order.items or []: + result, err = select_supplier(db, item["product_id"], int(item["qty"])) + if err: + routing_detail.append({"product_id": item["product_id"], "error": err}) + continue + sp = result["supplier_product"] + if supplier_id is None: + supplier_id = sp.supplier_id + supplier_cost += Decimal(str(sp.supplier_cost or 0)) * int(item["qty"]) + shipping_cost += Decimal(str(sp.shipping_cost or 0)) + routing_detail.append( + { + "product_id": item["product_id"], + "supplier": result["supplier"].name, + "supplier_score": result["score"], + } + ) + + order.supplier_id = supplier_id + order.supplier_cost = supplier_cost.quantize(CENT, rounding=ROUND_HALF_UP) + order.shipping_cost = shipping_cost.quantize(CENT, rounding=ROUND_HALF_UP) + order.fees = (Decimal(str(order.retail_total or 0)) * FEE_RATE + FEE_FLAT).quantize( + CENT, rounding=ROUND_HALF_UP + ) + order.profit = ( + Decimal(str(order.retail_total or 0)) + - order.supplier_cost + - order.shipping_cost + - order.fees + ).quantize(CENT, rounding=ROUND_HALF_UP) + + db.add( + AuditLog( + actor="order_router", + action="order_routed", + entity="order", + entity_id=order.order_number, + detail={ + "supplier_id": str(supplier_id) if supplier_id else None, + "profit": float(order.profit), + "routing": routing_detail, + }, + ) + ) + + if supplier_id is None: + order.status = "new" # stays unfulfilled; no eligible supplier + db.commit() + return {"routed": False, "reason": "no eligible supplier for any item"} + + order.status = "confirmed" + db.commit() + return { + "routed": True, + "supplier_id": str(supplier_id), + "profit": float(order.profit), + "supplier_cost": float(order.supplier_cost), + "shipping_cost": float(order.shipping_cost), + "fees": float(order.fees), + } diff --git a/backend/app/engines/pricing.py b/backend/app/engines/pricing.py new file mode 100644 index 0000000..3979d8c --- /dev/null +++ b/backend/app/engines/pricing.py @@ -0,0 +1,87 @@ +"""Pricing engine. + +Applies active `price_rules` (cost -> markup bands) and computes net profit: + + retail = cost * (1 + markup_pct/100) (floored at cost + shipping + fees) + net profit = retail - cost - shipping - fees + margin_pct = profit / retail * 100 + +`min_price` (product floor) is respected — retail is never below it. +""" +from decimal import ROUND_HALF_UP, Decimal + +from app.models import PriceRule + +CENT = Decimal("0.01") + + +def _q(value) -> Decimal: + try: + return Decimal(str(value or 0)) + except Exception: + return Decimal("0") + + +def get_active_rules(db): + return ( + db.query(PriceRule) + .filter(PriceRule.active.is_(True)) + .order_by(PriceRule.min_cost) + .all() + ) + + +def find_rule(db, cost): + cost = _q(cost) + for rule in get_active_rules(db): + if cost >= _q(rule.min_cost) and (rule.max_cost is None or cost < _q(rule.max_cost)): + return rule + return None + + +def calculate_price(db, cost, shipping=0, fees=0, min_price=None): + """Return dict(retail, profit, margin, margin_pct, markup_pct, rule_id, ...) or None. + + None means no active rule applies — callers must NOT publish such products. + """ + cost = _q(cost) + shipping = _q(shipping) + fees = _q(fees) + rule = find_rule(db, cost) + if rule is None: + return None + + markup_pct = _q(rule.markup_pct) + retail = (cost * (1 + markup_pct / 100)).quantize(CENT, rounding=ROUND_HALF_UP) + + # floor: cost + shipping + fees, raised by min_margin_pct if configured + floor = (cost + shipping + fees).quantize(CENT, rounding=ROUND_HALF_UP) + min_margin_pct = _q(rule.min_margin_pct) + if min_margin_pct > 0: + floor = max(floor, (floor + retail * min_margin_pct / 100).quantize(CENT, rounding=ROUND_HALF_UP)) + if retail < floor: + retail = floor + + # never below the product's minimum price + if min_price is not None and retail < _q(min_price): + retail = _q(min_price).quantize(CENT, rounding=ROUND_HALF_UP) + + profit = (retail - cost - shipping - fees).quantize(CENT, rounding=ROUND_HALF_UP) + margin_pct = (profit / retail * 100).quantize(Decimal("0.01")) if retail > 0 else Decimal("0") + + return { + "cost": cost, + "shipping": shipping, + "fees": fees, + "retail": retail, + "profit": profit, + "margin": profit, + "margin_pct": float(margin_pct), + "markup_pct": float(markup_pct), + "rule_id": rule.id, + } + + +def estimate_net_profit(retail, cost, shipping=0, fees=0): + """Net profit = retail - (cost + shipping + fees).""" + return (_q(retail) - _q(cost) - _q(shipping) - _q(fees)).quantize(CENT, rounding=ROUND_HALF_UP) diff --git a/backend/app/engines/scoring.py b/backend/app/engines/scoring.py new file mode 100644 index 0000000..50c1205 --- /dev/null +++ b/backend/app/engines/scoring.py @@ -0,0 +1,70 @@ +"""Supplier scoring. + +Weighted score (0-100): + 40% price — cheaper supplier cost vs. the eligible pool + 20% inventory — stock sufficiency for the requested qty + 15% speed — shipping_time ('3-5 days' -> avg days) + 10% fulfillment — supplier_performance.fulfillment_rate + 10% returns — inverse of supplier_performance.return_rate + 5% history — reliability_score blended with stock_accuracy +""" +import re + + +def _num(value, default=0.0): + try: + return float(value or 0) + except (TypeError, ValueError): + return default + + +def _shipping_days(sp): + text = (sp.shipping_time or "").strip() + numbers = re.findall(r"\d+", text) + if not numbers: + return 5.0 + return float(sum(int(n) for n in numbers)) / len(numbers) + + +def score_supplier(db, supplier, sp, qty, price_pool=None): + perf = supplier.performance + + # 1) price (40%) — cheaper than pool average is better + cost = _num(sp.supplier_cost) + if price_pool and price_pool.get("avg_cost"): + avg = price_pool["avg_cost"] + price_score = max(0.0, 100.0 - (cost / avg - 1.0) * 100.0) + else: + price_score = 50.0 + + # 2) inventory (20%) — can we fill the quantity? + inventory = int(sp.supplier_inventory or 0) + inv_score = 100.0 if inventory >= qty else max(0.0, (inventory / qty) * 100.0) + + # 3) speed (15%) — fewer shipping days is better + days = _shipping_days(sp) + speed_score = max(0.0, 100.0 - days * 12.0) + + # 4) fulfillment (10%) + fulfillment = 100.0 if perf is None else _num(perf.fulfillment_rate, 100.0) + + # 5) returns (10%) — 0% returns = 100 points, 20% returns = 0 points + return_rate = 0.0 if perf is None else _num(perf.return_rate, 0.0) + returns_score = max(0.0, 100.0 - return_rate * 5.0) + + # 6) history (5%) — reliability blended with stock accuracy + reliability = _num(supplier.reliability_score, 50.0) + if perf is None: + history_score = reliability + else: + history_score = 0.5 * reliability + 0.5 * _num(perf.stock_accuracy, 100.0) + + total = ( + 0.40 * price_score + + 0.20 * inv_score + + 0.15 * speed_score + + 0.10 * fulfillment + + 0.10 * returns_score + + 0.05 * history_score + ) + return round(min(100.0, max(0.0, total)), 2) diff --git a/backend/app/engines/suppliers/__init__.py b/backend/app/engines/suppliers/__init__.py new file mode 100644 index 0000000..edafdd8 --- /dev/null +++ b/backend/app/engines/suppliers/__init__.py @@ -0,0 +1 @@ +"""Supplier adapter implementations.""" diff --git a/backend/app/engines/suppliers/base.py b/backend/app/engines/suppliers/base.py new file mode 100644 index 0000000..a903f61 --- /dev/null +++ b/backend/app/engines/suppliers/base.py @@ -0,0 +1,65 @@ +"""SupplierAdapter abstract base class + registry.""" +from abc import ABC, abstractmethod + +# Registry of adapter_type string -> adapter class. Populated by register_adapter. +ADAPTERS = {} + + +class SupplierAdapter(ABC): + """Normalized interface every supplier feed must implement. + + Product dict shape (used by the importer): + {sku, title, description, category, brand, cost, inventory, + shipping_cost, shipping_time, image_url} + """ + + adapter_type = "base" + + def __init__(self, supplier, data_or_url=""): + self.supplier = supplier + self.data_or_url = data_or_url + + @abstractmethod + def get_products(self): + """Return list of normalized product dicts.""" + + @abstractmethod + def get_inventory(self): + """Return list of {'supplier_sku': str, 'inventory': int}.""" + + @abstractmethod + def get_price(self, sku): + """Return current supplier cost for sku (number) or None.""" + + @abstractmethod + def create_order(self, items): + """Place an order at the supplier; return {'ref': str, ...}.""" + + @abstractmethod + def get_order_status(self, ref): + """Return status string for a supplier order ref.""" + + @abstractmethod + def get_tracking(self, ref): + """Return {'tracking': str, 'carrier': str} or None.""" + + @abstractmethod + def cancel_order(self, ref): + """Cancel a supplier order; return status dict.""" + + +def register_adapter(cls): + ADAPTERS[cls.adapter_type] = cls + return cls + + +def build_adapter(supplier, data_or_url=""): + cls = ADAPTERS.get(supplier.adapter_type) + if cls is None: + # lazy-import concrete adapters so their @register_adapter side effects run + from app.engines.suppliers import csv_adapter, sample # noqa: F401 + + cls = ADAPTERS.get(supplier.adapter_type) + if cls is None: + raise ValueError(f"unknown adapter_type '{supplier.adapter_type}' for supplier {supplier.name}") + return cls(supplier, data_or_url=data_or_url) diff --git a/backend/app/engines/suppliers/csv_adapter.py b/backend/app/engines/suppliers/csv_adapter.py new file mode 100644 index 0000000..43c9367 --- /dev/null +++ b/backend/app/engines/suppliers/csv_adapter.py @@ -0,0 +1,78 @@ +"""CSV feed adapter. + +Reads a supplier CSV with columns: + sku,title,cost,inventory,shipping_cost,shipping_time,category,image_url + (optional: description, brand) +`data_or_url` is either raw CSV text or an http(s) URL to fetch. +""" +import csv +import io +import urllib.request + +from app.engines.suppliers.base import SupplierAdapter, register_adapter + + +@register_adapter +class CSVAdapter(SupplierAdapter): + adapter_type = "csv" + + def __init__(self, supplier, data_or_url=""): + super().__init__(supplier, data_or_url=data_or_url) + self._rows = None + + def _load(self): + if self._rows is not None: + return + src = (self.data_or_url or "").strip() + if src.startswith(("http://", "https://")): + with urllib.request.urlopen(src, timeout=30) as resp: + src = resp.read().decode("utf-8", "replace") + self._rows = list(csv.DictReader(io.StringIO(src))) + + def get_products(self): + self._load() + products = [] + for row in self._rows: + sku = (row.get("sku") or "").strip() + title = (row.get("title") or "").strip() + if not sku or not title: + continue + products.append( + { + "sku": sku, + "title": title, + "description": (row.get("description") or "").strip() or None, + "category": (row.get("category") or "").strip() or None, + "brand": (row.get("brand") or "").strip() or None, + "cost": float(row.get("cost") or 0), + "inventory": int(float(row.get("inventory") or 0)), + "shipping_cost": float(row.get("shipping_cost") or 0), + "shipping_time": (row.get("shipping_time") or "").strip() or "3-5 days", + "image_url": (row.get("image_url") or "").strip() or None, + } + ) + return products + + def get_inventory(self): + return [ + {"supplier_sku": p["sku"], "inventory": p["inventory"]} + for p in self.get_products() + ] + + def get_price(self, sku): + for p in self.get_products(): + if p["sku"] == sku: + return p["cost"] + return None + + def create_order(self, items): + return {"adapter": "csv", "ref": f"csv:{len(items)}:items", "items": items} + + def get_order_status(self, ref): + return "confirmed" + + def get_tracking(self, ref): + return None + + def cancel_order(self, ref): + return {"ref": ref, "status": "cancelled"} diff --git a/backend/app/engines/suppliers/sample.py b/backend/app/engines/suppliers/sample.py new file mode 100644 index 0000000..24521bc --- /dev/null +++ b/backend/app/engines/suppliers/sample.py @@ -0,0 +1,89 @@ +"""Sample supplier adapter — static catalog of ~25 realistic products. + +Home / gadgets / office niche, retail range $50-$500 after markup. +Used by scripts/seed.py to bootstrap a demo supplier. +""" +from app.engines.suppliers.base import SupplierAdapter, register_adapter + +# (sku, title, category, cost, inventory, shipping_cost, shipping_time, brand) +CATALOG = [ + ("SPL-1001", "Aurora LED Desk Lamp with Wireless Charger", "office", 38.00, 42, 6.50, "3-5 days", "Lumina"), + ("SPL-1002", "ErgoMesh Pro Office Chair", "office", 145.00, 12, 24.00, "5-7 days", "ErgoMesh"), + ("SPL-1003", "Standing Desk Converter 32\"", "office", 118.00, 20, 19.50, "4-6 days", "DeskUp"), + ("SPL-1004", "Mechanical Keyboard TKL (Brown Switch)", "office", 62.00, 55, 5.00, "3-5 days", "KeyForge"), + ("SPL-1005", "Precision Wireless Mouse Pro", "office", 24.00, 120, 3.50, "2-4 days", "CursorPro"), + ("SPL-1006", "Noise-Cancelling Over-Ear Headphones", "gadgets", 89.00, 38, 8.00, "3-5 days", "Auralis"), + ("SPL-1007", "Smart Watch Series X (GPS + HR)", "gadgets", 96.00, 27, 7.50, "4-6 days", "PulseTech"), + ("SPL-1008", "4K Action Camera Waterproof", "gadgets", 132.00, 18, 9.00, "4-6 days", "AdventureCam"), + ("SPL-1009", "Mini Bluetooth Speaker 20W", "gadgets", 31.00, 90, 4.00, "2-4 days", "SoundPod"), + ("SPL-1010", "Fast Wireless Charging Pad 15W", "gadgets", 38.00, 200, 2.50, "2-4 days", "VoltEdge"), + ("SPL-1011", "Smart Home Hub with Voice Control", "home", 74.00, 33, 6.00, "3-5 days", "NestLink"), + ("SPL-1012", "Robot Vacuum S5 (Self-Charging)", "home", 210.00, 9, 22.00, "6-8 days", "CleanBot"), + ("SPL-1013", "HEPA Air Purifier for Large Rooms", "home", 128.00, 15, 18.00, "5-7 days", "PureAir"), + ("SPL-1014", "Compact Espresso Machine 15-Bar", "home", 178.00, 8, 21.00, "5-7 days", "BrewMaster"), + ("SPL-1015", "Cordless Stick Vacuum with LED Light", "home", 155.00, 11, 20.00, "5-7 days", "CleanBot"), + ("SPL-1016", "Smart Thermostat (WiFi)", "home", 92.00, 40, 5.50, "3-5 days", "NestLink"), + ("SPL-1017", "Smart LED Strip Kit 5m (RGB)", "home", 38.00, 140, 3.00, "2-4 days", "GlowSpace"), + ("SPL-1018", "Gooseneck Electric Kettle 1L", "home", 44.00, 60, 7.00, "3-5 days", "BrewMaster"), + ("SPL-1019", "Cast Iron Dutch Oven 6qt", "home", 68.00, 25, 16.00, "5-7 days", "Hearthstone"), + ("SPL-1020", "Bamboo Electric Standing Desk 55\"", "office", 265.00, 6, 38.00, "7-10 days", "DeskUp"), + ("SPL-1021", "Ultrawide Monitor 34\" 3440x1440", "office", 310.00, 7, 32.00, "6-8 days", "ViewMax"), + ("SPL-1022", "Duplex Document Scanner", "office", 98.00, 22, 11.00, "4-6 days", "ScanFast"), + ("SPL-1023", "Aluminum Laptop Stand (Adjustable)", "office", 38.00, 85, 4.50, "2-4 days", "DeskUp"), + ("SPL-1024", "Smart Body Composition Scale", "home", 38.00, 70, 5.00, "3-5 days", "PulseTech"), + ("SPL-1025", "Ergonomic Footrest with Massage", "office", 40.00, 48, 6.00, "3-5 days", "ErgoMesh"), + ("SPL-1026", "USB-C Hub 8-in-1 (4K HDMI)", "gadgets", 38.00, 110, 3.00, "2-4 days", "VoltEdge"), + ("SPL-1027", "Portable Power Station 300Wh", "gadgets", 348.00, 5, 34.00, "7-10 days", "VoltEdge"), +] + +DESCRIPTIONS = { + "office": "Professional-grade office equipment designed for comfort and productivity.", + "home": "Thoughtfully designed home essentials that blend function with style.", + "gadgets": "Modern tech gadget with reliable performance and everyday practicality.", +} + + +@register_adapter +class SampleSupplierAdapter(SupplierAdapter): + adapter_type = "sample" + + def get_products(self): + return [ + { + "sku": sku, + "title": title, + "description": DESCRIPTIONS[category], + "category": category, + "brand": brand, + "cost": cost, + "inventory": inventory, + "shipping_cost": shipping_cost, + "shipping_time": shipping_time, + "image_url": f"https://placehold.co/600x600/png?text={sku}", + } + for (sku, title, category, cost, inventory, shipping_cost, shipping_time, brand) in CATALOG + ] + + def get_inventory(self): + return [ + {"supplier_sku": sku, "inventory": inventory} + for (sku, _, _, _, inventory, _, _, _) in CATALOG + ] + + def get_price(self, sku): + for row in CATALOG: + if row[0] == sku: + return row[3] + return None + + def create_order(self, items): + return {"adapter": "sample", "ref": f"sample:{len(items)}:items", "items": items} + + def get_order_status(self, ref): + return "confirmed" + + def get_tracking(self, ref): + return None + + def cancel_order(self, ref): + return {"ref": ref, "status": "cancelled"} diff --git a/backend/app/main.py b/backend/app/main.py new file mode 100644 index 0000000..cec448b --- /dev/null +++ b/backend/app/main.py @@ -0,0 +1,38 @@ +"""Polaris FastAPI application. + +CORS allow-all (MVP). Routers mounted under /api exactly as the contract lists. +""" +from fastapi import FastAPI +from fastapi.middleware.cors import CORSMiddleware + +from app.database import Base, engine +from app.routers import admin, analytics, auth, customers, health, orders, products, suppliers +from app.version import VERSION + +# Idempotent — schema.sql is canonical and already applied; create_all only +# creates anything missing. +Base.metadata.create_all(bind=engine) + +app = FastAPI(title="Polaris API", version=VERSION) + +app.add_middleware( + CORSMiddleware, + allow_origins=["*"], + allow_credentials=True, + allow_methods=["*"], + allow_headers=["*"], +) + +app.include_router(health.router, prefix="/api", tags=["health"]) +app.include_router(auth.router, prefix="/api/auth", tags=["auth"]) +app.include_router(products.router, prefix="/api/products", tags=["products"]) +app.include_router(suppliers.router, prefix="/api/suppliers", tags=["suppliers"]) +app.include_router(orders.router, prefix="/api/orders", tags=["orders"]) +app.include_router(customers.router, prefix="/api/customers", tags=["customers"]) +app.include_router(admin.router, prefix="/api/admin", tags=["admin"]) +app.include_router(analytics.router, prefix="/api/analytics", tags=["analytics"]) + + +@app.get("/") +def root(): + return {"app": "polaris", "version": VERSION, "docs": "/docs"} diff --git a/backend/app/models.py b/backend/app/models.py new file mode 100644 index 0000000..c0d544d --- /dev/null +++ b/backend/app/models.py @@ -0,0 +1,185 @@ +"""SQLAlchemy ORM models — mirror db/schema.sql table + column names EXACTLY. + +Do NOT rename tables/columns here: the schema.sql in /opt/polaris/db is canonical +and the database is already initialized from it. +""" +from datetime import datetime +from decimal import Decimal + +from sqlalchemy import ( + BigInteger, + Boolean, + Column, + DateTime, + ForeignKey, + Integer, + Numeric, + Text, + UniqueConstraint, + func, + text, +) +from sqlalchemy.dialects.postgresql import JSONB, UUID +from sqlalchemy.orm import relationship + +from app.database import Base + + +class Supplier(Base): + __tablename__ = "suppliers" + + id = Column(UUID(as_uuid=True), primary_key=True, server_default=text("gen_random_uuid()")) + name = Column(Text, nullable=False) + adapter_type = Column(Text, nullable=False) + api_endpoint = Column(Text) + account_id = Column(Text) + credentials_enc = Column(Text) + fulfillment_caps = Column(JSONB, default=list) + shipping_regions = Column(JSONB, default=list) + reliability_score = Column(Numeric(5, 2), default=Decimal("50.00")) + status = Column(Text, nullable=False, default="active") + created_at = Column(DateTime(timezone=True), server_default=func.now()) + updated_at = Column(DateTime(timezone=True), server_default=func.now(), onupdate=func.now()) + + supplier_products = relationship("SupplierProduct", back_populates="supplier") + performance = relationship( + "SupplierPerformance", back_populates="supplier", uselist=False, lazy="joined" + ) + + +class Product(Base): + __tablename__ = "products" + + id = Column(UUID(as_uuid=True), primary_key=True, server_default=text("gen_random_uuid()")) + sku = Column(Text, unique=True, nullable=False) + title = Column(Text, nullable=False) + description = Column(Text) + category = Column(Text) + brand = Column(Text) + manufacturer = Column(Text) + upc_ean_gtin = Column(Text) + images = Column(JSONB, default=list) + dimensions = Column(JSONB) + weight_grams = Column(Numeric(10, 2)) + cost = Column(Numeric(12, 2), nullable=False, default=Decimal("0")) + retail_price = Column(Numeric(12, 2)) + min_price = Column(Numeric(12, 2)) + max_price = Column(Numeric(12, 2)) + status = Column(Text, nullable=False, default="discovered") + lifecycle = Column(JSONB, default=dict) + created_at = Column(DateTime(timezone=True), server_default=func.now()) + updated_at = Column(DateTime(timezone=True), server_default=func.now(), onupdate=func.now()) + + supplier_products = relationship("SupplierProduct", back_populates="product") + + +class SupplierProduct(Base): + __tablename__ = "supplier_products" + __table_args__ = ( + UniqueConstraint("supplier_id", "supplier_sku", name="supplier_products_supplier_id_supplier_sku_key"), + ) + + id = Column(UUID(as_uuid=True), primary_key=True, server_default=text("gen_random_uuid()")) + supplier_id = Column(UUID(as_uuid=True), ForeignKey("suppliers.id", ondelete="CASCADE"), nullable=False) + supplier_sku = Column(Text, nullable=False) + product_id = Column(UUID(as_uuid=True), ForeignKey("products.id", ondelete="SET NULL")) + supplier_cost = Column(Numeric(12, 2), nullable=False, default=Decimal("0")) + supplier_inventory = Column(Integer, nullable=False, default=0) + shipping_cost = Column(Numeric(12, 2), default=Decimal("0")) + shipping_time = Column(Text) + last_updated = Column(DateTime(timezone=True), server_default=func.now()) + + supplier = relationship("Supplier", back_populates="supplier_products") + product = relationship("Product", back_populates="supplier_products") + + +class Customer(Base): + __tablename__ = "customers" + + id = Column(UUID(as_uuid=True), primary_key=True, server_default=text("gen_random_uuid()")) + email = Column(Text, unique=True, nullable=False) + name = Column(Text) + shipping_addr = Column(JSONB) + created_at = Column(DateTime(timezone=True), server_default=func.now()) + + +class Order(Base): + __tablename__ = "orders" + + id = Column(UUID(as_uuid=True), primary_key=True, server_default=text("gen_random_uuid()")) + order_number = Column(Text, unique=True, nullable=False) + customer_id = Column(UUID(as_uuid=True), ForeignKey("customers.id")) + items = Column(JSONB, nullable=False) + retail_total = Column(Numeric(12, 2), nullable=False, default=Decimal("0")) + supplier_id = Column(UUID(as_uuid=True), ForeignKey("suppliers.id")) + supplier_cost = Column(Numeric(12, 2), default=Decimal("0")) + shipping_cost = Column(Numeric(12, 2), default=Decimal("0")) + fees = Column(Numeric(12, 2), default=Decimal("0")) + profit = Column(Numeric(12, 2), default=Decimal("0")) + status = Column(Text, nullable=False, default="new") + tracking = Column(Text) + carrier = Column(Text) + created_at = Column(DateTime(timezone=True), server_default=func.now()) + updated_at = Column(DateTime(timezone=True), server_default=func.now(), onupdate=func.now()) + + customer = relationship("Customer") + + +class PriceHistory(Base): + __tablename__ = "price_history" + + id = Column(BigInteger, primary_key=True, autoincrement=True) + product_id = Column(UUID(as_uuid=True), ForeignKey("products.id", ondelete="CASCADE"), nullable=False) + price = Column(Numeric(12, 2), nullable=False) + reason = Column(Text) + created_at = Column(DateTime(timezone=True), server_default=func.now()) + + +class InventoryHistory(Base): + __tablename__ = "inventory_history" + + id = Column(BigInteger, primary_key=True, autoincrement=True) + supplier_product_id = Column( + UUID(as_uuid=True), ForeignKey("supplier_products.id", ondelete="CASCADE"), nullable=False + ) + inventory = Column(Integer, nullable=False) + created_at = Column(DateTime(timezone=True), server_default=func.now()) + + +class SupplierPerformance(Base): + __tablename__ = "supplier_performance" + + id = Column(UUID(as_uuid=True), primary_key=True, server_default=text("gen_random_uuid()")) + supplier_id = Column(UUID(as_uuid=True), ForeignKey("suppliers.id", ondelete="CASCADE"), nullable=False) + fulfillment_rate = Column(Numeric(5, 2), default=Decimal("100.00")) + avg_shipping_days = Column(Numeric(6, 2)) + cancellation_rate = Column(Numeric(5, 2), default=Decimal("0.00")) + stock_accuracy = Column(Numeric(5, 2), default=Decimal("100.00")) + return_rate = Column(Numeric(5, 2), default=Decimal("0.00")) + updated_at = Column(DateTime(timezone=True), server_default=func.now(), onupdate=func.now()) + + supplier = relationship("Supplier", back_populates="performance") + + +class PriceRule(Base): + __tablename__ = "price_rules" + + id = Column(Integer, primary_key=True, autoincrement=True) + min_cost = Column(Numeric(12, 2), nullable=False, default=Decimal("0")) + max_cost = Column(Numeric(12, 2)) + markup_pct = Column(Numeric(6, 2), nullable=False) + min_margin_pct = Column(Numeric(6, 2), default=Decimal("0")) + active = Column(Boolean, default=True) + created_at = Column(DateTime(timezone=True), server_default=func.now()) + + +class AuditLog(Base): + __tablename__ = "audit_logs" + + id = Column(BigInteger, primary_key=True, autoincrement=True) + actor = Column(Text) + action = Column(Text, nullable=False) + entity = Column(Text) + entity_id = Column(Text) + detail = Column(JSONB) + created_at = Column(DateTime(timezone=True), server_default=func.now()) diff --git a/backend/app/routers/__init__.py b/backend/app/routers/__init__.py new file mode 100644 index 0000000..f7ec5ce --- /dev/null +++ b/backend/app/routers/__init__.py @@ -0,0 +1 @@ +"""API routers.""" diff --git a/backend/app/routers/admin.py b/backend/app/routers/admin.py new file mode 100644 index 0000000..7cc7bee --- /dev/null +++ b/backend/app/routers/admin.py @@ -0,0 +1,160 @@ +"""Admin router — JWT-protected operational endpoints. + +GET /api/admin/dashboard +GET /api/admin/orders +POST /api/admin/orders/{id}/status +POST /api/admin/import +GET /api/admin/price-rules +POST /api/admin/recalc (recalc_prices task) +POST /api/admin/sync-inventory (sync_inventory task) +""" +from datetime import datetime, time, timezone +from decimal import Decimal +from typing import List, Optional + +from fastapi import APIRouter, Depends, HTTPException +from pydantic import BaseModel +from sqlalchemy.orm import Session + +from app.database import get_db +from app.models import AuditLog, Order, PriceRule, Product, Supplier +from app.routers.auth import require_admin +from app.schemas import ImportRequest, OrderOut + +router = APIRouter(dependencies=[Depends(require_admin)]) + +ORDER_STATUSES = {"new", "paid", "fraud_check", "supplier_order", "confirmed", + "shipped", "delivered", "cancelled", "refunded"} + + +class StatusUpdate(BaseModel): + status: str + + +def _f(value) -> float: + try: + return float(value or 0) + except (TypeError, ValueError): + return 0.0 + + +@router.get("/dashboard") +def dashboard(db: Session = Depends(get_db)): + today_start = datetime.combine(datetime.now(timezone.utc).date(), time.min) + excluded = ["cancelled", "refunded"] + + def kpis(query): + orders = query.all() + revenue = sum(_f(o.retail_total) for o in orders) + profit = sum(_f(o.profit) for o in orders) + count = len(orders) + return { + "orders": count, + "revenue": round(revenue, 2), + "profit": round(profit, 2), + "gross_profit": round(revenue - profit, 2), + "margin": round(profit / revenue * 100, 2) if revenue else 0.0, + "aov": round(revenue / count, 2) if count else 0.0, + "refunds": sum(1 for o in orders if o.status == "refunded"), + } + + today_q = db.query(Order).filter( + Order.created_at >= today_start, ~Order.status.in_(excluded) + ) + all_q = db.query(Order).filter(~Order.status.in_(excluded)) + + return { + "today": kpis(today_q), + "all_time": kpis(all_q), + "products_total": db.query(Product).count(), + } + + +@router.get("/orders", response_model=List[OrderOut]) +def admin_orders(status: Optional[str] = None, db: Session = Depends(get_db)): + q = db.query(Order) + if status: + q = q.filter(Order.status == status) + return q.order_by(Order.created_at.desc()).all() + + +@router.post("/orders/{order_id}/status", response_model=OrderOut) +def set_order_status(order_id, body: StatusUpdate, db: Session = Depends(get_db)): + if body.status not in ORDER_STATUSES: + raise HTTPException(status_code=400, detail=f"invalid status; allowed: {sorted(ORDER_STATUSES)}") + order = db.query(Order).filter(Order.id == order_id).first() + if order is None: + raise HTTPException(status_code=404, detail="order not found") + old = order.status + order.status = body.status + db.add( + AuditLog( + actor="admin", + action="order_status_changed", + entity="order", + entity_id=order.order_number, + detail={"from": old, "to": body.status}, + ) + ) + db.commit() + db.refresh(order) + return order + + +@router.post("/import") +def admin_import(body: ImportRequest, db: Session = Depends(get_db)): + from app.engines.importer import import_feed + + supplier = db.query(Supplier).filter(Supplier.id == body.supplier_id).first() + if supplier is None: + raise HTTPException(status_code=404, detail="supplier not found") + summary = import_feed(db, supplier, feed_type=body.feed_type, data_or_url=body.data_or_url) + db.add( + AuditLog( + actor="admin", + action="import_requested", + entity="supplier", + entity_id=str(supplier.id), + detail={"feed_type": body.feed_type}, + ) + ) + db.commit() + return summary + + +@router.get("/price-rules") +def price_rules(db: Session = Depends(get_db)): + return [ + { + "id": r.id, + "min_cost": float(r.min_cost), + "max_cost": float(r.max_cost) if r.max_cost is not None else None, + "markup_pct": float(r.markup_pct), + "min_margin_pct": float(r.min_margin_pct or 0), + "active": r.active, + } + for r in db.query(PriceRule).order_by(PriceRule.min_cost).all() + ] + + +def _dispatch(task_name, *args): + """Prefer Celery; fall back to running inline when the broker is unreachable.""" + from app.tasks import recalc_prices, sync_inventory + + fn = {"recalc_prices": recalc_prices, "sync_inventory": sync_inventory}[task_name] + try: + result = fn.delay(*args) + return {"dispatched": True, "task_id": result.id} + except Exception: + # broker down (e.g. running on host without redis DNS) — run inline + return {"dispatched": False, "inline_result": fn.run(*args)} + + +@router.post("/recalc") +def admin_recalc(): + return _dispatch("recalc_prices") + + +@router.post("/sync-inventory") +def admin_sync_inventory(supplier_id: Optional[str] = None): + return _dispatch("sync_inventory", supplier_id) diff --git a/backend/app/routers/analytics.py b/backend/app/routers/analytics.py new file mode 100644 index 0000000..1462407 --- /dev/null +++ b/backend/app/routers/analytics.py @@ -0,0 +1,123 @@ +"""Analytics router. + +GET /api/analytics/summary -> today revenue/orders/profit/margin/aov/conversion +GET /api/analytics/top-products +GET /api/analytics/supplier-performance +""" +from collections import defaultdict +from datetime import datetime, time, timezone + +from fastapi import APIRouter, Depends +from sqlalchemy import func +from sqlalchemy.orm import Session + +from app.database import get_db +from app.models import Order, Supplier, SupplierPerformance + +router = APIRouter() + +EXCLUDED = ["cancelled", "refunded"] + + +def _f(value) -> float: + try: + return float(value or 0) + except (TypeError, ValueError): + return 0.0 + + +@router.get("/summary") +def summary(db: Session = Depends(get_db)): + today_start = datetime.combine(datetime.now(timezone.utc).date(), time.min) + + today_q = db.query(Order).filter( + Order.created_at >= today_start, ~Order.status.in_(EXCLUDED) + ) + all_q = db.query(Order).filter(~Order.status.in_(EXCLUDED)) + + def _kpis(query): + orders = query.all() + revenue = sum(_f(o.retail_total) for o in orders) + profit = sum(_f(o.profit) for o in orders) + count = len(orders) + return { + "revenue": round(revenue, 2), + "orders": count, + "profit": round(profit, 2), + "gross_profit": round(revenue - profit, 2), + "margin": round(profit / revenue * 100, 2) if revenue else 0.0, + "aov": round(revenue / count, 2) if count else 0.0, + # no traffic/visitor tracking table yet — conversion is not measurable + "conversion": 0.0, + } + + refunds_today = ( + db.query(func.count(Order.id)) + .filter(Order.created_at >= today_start, Order.status.in_(["cancelled", "refunded"])) + .scalar() + or 0 + ) + + return { + "today": {**_kpis(today_q), "refunds": refunds_today}, + "all_time": _kpis(all_q), + } + + +@router.get("/top-products") +def top_products(limit: int = 10, db: Session = Depends(get_db)): + orders = ( + db.query(Order) + .filter(~Order.status.in_(EXCLUDED)) + .order_by(Order.created_at.desc()) + .limit(500) + .all() + ) + stats = defaultdict(lambda: {"sku": None, "units": 0, "revenue": 0.0}) + for order in orders: + for item in order.items or []: + pid = item.get("product_id") + if not pid: + continue + entry = stats[pid] + entry["sku"] = item.get("sku") + entry["units"] += int(item.get("qty") or 0) + entry["revenue"] += float(item.get("unit_price") or 0) * int(item.get("qty") or 0) + ranked = sorted(stats.items(), key=lambda kv: (-kv[1]["revenue"], -kv[1]["units"])) + return [ + {"product_id": pid, "sku": s["sku"], "units": s["units"], "revenue": round(s["revenue"], 2)} + for pid, s in ranked[:limit] + ] + + +@router.get("/supplier-performance") +def supplier_performance(db: Session = Depends(get_db)): + out = [] + for supplier in db.query(Supplier).all(): + perf = ( + db.query(SupplierPerformance) + .filter(SupplierPerformance.supplier_id == supplier.id) + .first() + ) + orders = ( + db.query(Order) + .filter(Order.supplier_id == supplier.id, ~Order.status.in_(EXCLUDED)) + .all() + ) + out.append( + { + "supplier_id": str(supplier.id), + "name": supplier.name, + "adapter_type": supplier.adapter_type, + "status": supplier.status, + "reliability_score": _f(supplier.reliability_score), + "fulfillment_rate": _f(perf.fulfillment_rate) if perf else 100.0, + "avg_shipping_days": _f(perf.avg_shipping_days) if perf and perf.avg_shipping_days is not None else None, + "cancellation_rate": _f(perf.cancellation_rate) if perf else 0.0, + "stock_accuracy": _f(perf.stock_accuracy) if perf else 100.0, + "return_rate": _f(perf.return_rate) if perf else 0.0, + "orders": len(orders), + "profit": round(sum(_f(o.profit) for o in orders), 2), + } + ) + return out diff --git a/backend/app/routers/auth.py b/backend/app/routers/auth.py new file mode 100644 index 0000000..bae71d6 --- /dev/null +++ b/backend/app/routers/auth.py @@ -0,0 +1,84 @@ +"""Auth router — simple JWT for admin login + customer registration. + +POST /api/auth/login (admin, env ADMIN_EMAIL / ADMIN_PASSWORD) +POST /api/auth/register (customer) + +TODO(security): hash ADMIN_PASSWORD (bcrypt/argon2) instead of plain compare — +acceptable for MVP per task instructions. Customers table has no password +column yet, so registration is identity-only. +""" +import os +from datetime import datetime, timedelta, timezone + +import jwt +from fastapi import APIRouter, Depends, Header, HTTPException +from sqlalchemy.orm import Session + +from app.database import get_db +from app.models import Customer +from app.schemas import CustomerRegister, LoginRequest, TokenResponse + +router = APIRouter() + +SECRET_KEY = os.getenv("SECRET_KEY", "polaris-dev-secret-change-me") +ALGORITHM = "HS256" +TOKEN_TTL_HOURS = 24 + + +def _make_token(email: str, role: str) -> str: + now = datetime.now(timezone.utc) + payload = { + "sub": email, + "role": role, + "iat": now, + "exp": now + timedelta(hours=TOKEN_TTL_HOURS), + } + return jwt.encode(payload, SECRET_KEY, algorithm=ALGORITHM) + + +@router.post("/login", response_model=TokenResponse) +def login(body: LoginRequest): + admin_email = os.getenv("ADMIN_EMAIL", "admin@polaris.local") + # contract mentions ADMIN_PASSWORD_HASH; MVP uses plain compare (see TODO above) + admin_password = os.getenv("ADMIN_PASSWORD") or os.getenv("ADMIN_PASSWORD_HASH", "admin") + + if body.email.strip().lower() != admin_email.strip().lower(): + raise HTTPException(status_code=401, detail="invalid credentials") + if body.password != admin_password: + raise HTTPException(status_code=401, detail="invalid credentials") + + return TokenResponse(access_token=_make_token(body.email, "admin")) + + +@router.post("/register") +def register(body: CustomerRegister, db: Session = Depends(get_db)): + email = body.email.strip().lower() + existing = db.query(Customer).filter(Customer.email == email).first() + if existing is not None: + raise HTTPException(status_code=409, detail="email already registered") + + customer = Customer(email=email, name=body.name, shipping_addr=body.shipping_addr or {}) + db.add(customer) + db.commit() + db.refresh(customer) + + return { + "id": str(customer.id), + "email": customer.email, + "name": customer.name, + "access_token": _make_token(customer.email, "customer"), + "token_type": "bearer", + } + + +def require_admin(authorization: str = Header(default="")): + """FastAPI dependency — requires a valid admin JWT in the Authorization header.""" + if not authorization.startswith("Bearer "): + raise HTTPException(status_code=401, detail="missing bearer token") + try: + payload = jwt.decode(authorization[7:], SECRET_KEY, algorithms=[ALGORITHM]) + except jwt.PyJWTError: + raise HTTPException(status_code=401, detail="invalid token") + if payload.get("role") != "admin": + raise HTTPException(status_code=403, detail="admin role required") + return payload diff --git a/backend/app/routers/customers.py b/backend/app/routers/customers.py new file mode 100644 index 0000000..3732591 --- /dev/null +++ b/backend/app/routers/customers.py @@ -0,0 +1,42 @@ +"""Customers router. + +GET /api/customers +GET /api/customers/{id} +POST /api/customers +""" +from typing import List + +from fastapi import APIRouter, Depends, HTTPException +from sqlalchemy.orm import Session + +from app.database import get_db +from app.models import Customer +from app.schemas import CustomerOut, CustomerRegister + +router = APIRouter() + + +@router.get("", response_model=List[CustomerOut]) +def list_customers(db: Session = Depends(get_db)): + return db.query(Customer).order_by(Customer.created_at.desc()).all() + + +@router.get("/{customer_id}", response_model=CustomerOut) +def get_customer(customer_id, db: Session = Depends(get_db)): + customer = db.query(Customer).filter(Customer.id == customer_id).first() + if customer is None: + raise HTTPException(status_code=404, detail="customer not found") + return customer + + +@router.post("", response_model=CustomerOut) +def create_customer(body: CustomerRegister, db: Session = Depends(get_db)): + email = body.email.strip().lower() + existing = db.query(Customer).filter(Customer.email == email).first() + if existing is not None: + raise HTTPException(status_code=409, detail="email already registered") + customer = Customer(email=email, name=body.name, shipping_addr=body.shipping_addr or {}) + db.add(customer) + db.commit() + db.refresh(customer) + return customer diff --git a/backend/app/routers/health.py b/backend/app/routers/health.py new file mode 100644 index 0000000..67e40a0 --- /dev/null +++ b/backend/app/routers/health.py @@ -0,0 +1,23 @@ +"""GET /api/health -> {status, version}""" +from fastapi import APIRouter, Depends +from sqlalchemy import text +from sqlalchemy.orm import Session + +from app.database import get_db +from app.version import VERSION + +router = APIRouter() + + +@router.get("/health") +def health(db: Session = Depends(get_db)): + db_status = "ok" + try: + db.execute(text("SELECT 1")) + except Exception: # noqa: BLE001 + db_status = "error" + return { + "status": "ok" if db_status == "ok" else "degraded", + "version": VERSION, + "db": db_status, + } diff --git a/backend/app/routers/orders.py b/backend/app/routers/orders.py new file mode 100644 index 0000000..da1cc03 --- /dev/null +++ b/backend/app/routers/orders.py @@ -0,0 +1,113 @@ +"""Orders router. + +POST /api/orders (body: {customer_email, items:[{product_id,qty}]}) + -> creates order, routes to highest-scoring supplier, records profit +GET /api/orders +GET /api/orders/{id} +POST /api/orders/{id}/tracking (simulate supplier tracking push) +""" +import uuid +from datetime import datetime +from decimal import Decimal +from typing import List, Optional + +from fastapi import APIRouter, Depends, HTTPException +from sqlalchemy.orm import Session + +from app.database import get_db +from app.engines.order_router import route_order +from app.models import AuditLog, Customer, Order, Product +from app.schemas import OrderCreate, OrderOut, TrackingIn + +router = APIRouter() + +TRACKING_STATUSES = {"new", "paid", "fraud_check", "supplier_order", "confirmed"} + + +def _get_order_or_404(db, order_id): + order = db.query(Order).filter(Order.id == order_id).first() + if order is None: + raise HTTPException(status_code=404, detail="order not found") + return order + + +def _new_order_number() -> str: + return f"POL-{datetime.utcnow():%Y%m%d}-{uuid.uuid4().hex[:6].upper()}" + + +@router.post("", response_model=OrderOut) +def create_order(body: OrderCreate, db: Session = Depends(get_db)): + email = body.customer_email.strip().lower() + customer = db.query(Customer).filter(Customer.email == email).first() + if customer is None: + customer = Customer(email=email, name=email.split("@")[0]) + db.add(customer) + db.flush() + + items = [] + retail_total = Decimal("0") + for item in body.items: + product = db.query(Product).filter(Product.id == item.product_id).first() + if product is None: + raise HTTPException(status_code=404, detail=f"product {item.product_id} not found") + price = product.retail_price + if price is None or Decimal(str(price)) <= 0: + raise HTTPException(status_code=400, detail=f"product {product.sku} has no retail price") + items.append( + { + "product_id": str(product.id), + "sku": product.sku, + "qty": item.qty, + "unit_price": float(price), + } + ) + retail_total += Decimal(str(price)) * item.qty + + order = Order( + order_number=_new_order_number(), + customer_id=customer.id, + items=items, + retail_total=retail_total, + status="new", + ) + db.add(order) + db.flush() + + result = route_order(db, order) + db.refresh(order) + return order + + +@router.get("", response_model=List[OrderOut]) +def list_orders(status: Optional[str] = None, db: Session = Depends(get_db)): + q = db.query(Order) + if status: + q = q.filter(Order.status == status) + return q.order_by(Order.created_at.desc()).all() + + +@router.get("/{order_id}", response_model=OrderOut) +def get_order(order_id, db: Session = Depends(get_db)): + return _get_order_or_404(db, order_id) + + +@router.post("/{order_id}/tracking", response_model=OrderOut) +def add_tracking(order_id, body: TrackingIn, db: Session = Depends(get_db)): + order = _get_order_or_404(db, order_id) + order.tracking = body.tracking + if body.carrier: + order.carrier = body.carrier + if order.status in TRACKING_STATUSES: + order.status = "shipped" + db.add( + AuditLog( + actor="api", + action="tracking_received", + entity="order", + entity_id=order.order_number, + detail={"tracking": body.tracking, "carrier": body.carrier}, + ) + ) + db.commit() + db.refresh(order) + return order diff --git a/backend/app/routers/products.py b/backend/app/routers/products.py new file mode 100644 index 0000000..211403b --- /dev/null +++ b/backend/app/routers/products.py @@ -0,0 +1,132 @@ +"""Products router. + +GET /api/products (filters: status, category, search, limit, offset) +GET /api/products/{id} +POST /api/products/{id}/publish +POST /api/products/{id}/pause +GET /api/products/{id}/price-history +POST /api/products/import (body: {supplier_id, feed_type, data_or_url}) + +NOTE: /import is declared before /{id} so it isn't captured by the path param. +""" +from decimal import Decimal +from typing import List, Optional + +from fastapi import APIRouter, Depends, HTTPException, Query +from sqlalchemy import or_ +from sqlalchemy.orm import Session + +from app.database import get_db +from app.engines.importer import import_feed +from app.models import AuditLog, PriceHistory, Product, Supplier +from app.schemas import ImportRequest, PriceHistoryOut, ProductOut + +router = APIRouter() + +PUBLISHABLE_FROM = {"imported", "price_calculated", "content_generated", "quality_check", "paused"} + + +def _get_product_or_404(db, product_id): + product = db.query(Product).filter(Product.id == product_id).first() + if product is None: + raise HTTPException(status_code=404, detail="product not found") + return product + + +@router.get("", response_model=List[ProductOut]) +def list_products( + status: Optional[str] = None, + category: Optional[str] = None, + search: Optional[str] = None, + limit: int = Query(50, ge=1, le=500), + offset: int = Query(0, ge=0), + db: Session = Depends(get_db), +): + q = db.query(Product) + if status: + q = q.filter(Product.status == status) + if category: + q = q.filter(Product.category == category) + if search: + like = f"%{search}%" + q = q.filter(or_(Product.title.ilike(like), Product.sku.ilike(like))) + return q.order_by(Product.created_at.desc()).offset(offset).limit(limit).all() + + +@router.post("/import") +def import_products(body: ImportRequest, db: Session = Depends(get_db)): + supplier = db.query(Supplier).filter(Supplier.id == body.supplier_id).first() + if supplier is None: + raise HTTPException(status_code=404, detail="supplier not found") + try: + summary = import_feed(db, supplier, feed_type=body.feed_type, data_or_url=body.data_or_url) + except ValueError as exc: + raise HTTPException(status_code=400, detail=str(exc)) + db.add( + AuditLog( + actor="api", + action="import_requested", + entity="supplier", + entity_id=str(supplier.id), + detail={"feed_type": body.feed_type, "received": summary["received"]}, + ) + ) + db.commit() + return summary + + +@router.get("/{product_id}", response_model=ProductOut) +def get_product(product_id, db: Session = Depends(get_db)): + return _get_product_or_404(db, product_id) + + +@router.post("/{product_id}/publish", response_model=ProductOut) +def publish_product(product_id, db: Session = Depends(get_db)): + product = _get_product_or_404(db, product_id) + if product.status not in PUBLISHABLE_FROM: + raise HTTPException(status_code=400, detail=f"cannot publish product in status '{product.status}'") + if product.retail_price is None or Decimal(str(product.retail_price)) <= 0: + raise HTTPException(status_code=400, detail="product has no calculated retail price") + product.status = "published" + db.add( + AuditLog( + actor="api", + action="product_published", + entity="product", + entity_id=str(product.id), + detail={"sku": product.sku, "retail_price": float(product.retail_price)}, + ) + ) + db.commit() + db.refresh(product) + return product + + +@router.post("/{product_id}/pause", response_model=ProductOut) +def pause_product(product_id, db: Session = Depends(get_db)): + product = _get_product_or_404(db, product_id) + product.status = "paused" + db.add( + AuditLog( + actor="api", + action="product_paused", + entity="product", + entity_id=str(product.id), + detail={"sku": product.sku}, + ) + ) + db.commit() + db.refresh(product) + return product + + +@router.get("/{product_id}/price-history", response_model=List[PriceHistoryOut]) +def price_history(product_id, db: Session = Depends(get_db)): + _get_product_or_404(db, product_id) + rows = ( + db.query(PriceHistory) + .filter(PriceHistory.product_id == product_id) + .order_by(PriceHistory.created_at.desc()) + .all() + ) + return rows diff --git a/backend/app/routers/suppliers.py b/backend/app/routers/suppliers.py new file mode 100644 index 0000000..f1d7488 --- /dev/null +++ b/backend/app/routers/suppliers.py @@ -0,0 +1,76 @@ +"""Suppliers router. + +GET /api/suppliers +POST /api/suppliers +GET /api/suppliers/{id}/performance +""" +from typing import List, Optional + +from fastapi import APIRouter, Depends, HTTPException +from sqlalchemy.orm import Session + +from app.database import get_db +from app.models import Supplier, SupplierPerformance +from app.schemas import SupplierCreate, SupplierOut, SupplierPerformanceOut + +router = APIRouter() + + +@router.get("", response_model=List[SupplierOut]) +def list_suppliers(status: Optional[str] = None, db: Session = Depends(get_db)): + q = db.query(Supplier) + if status: + q = q.filter(Supplier.status == status) + return q.order_by(Supplier.created_at.desc()).all() + + +@router.post("", response_model=SupplierOut) +def create_supplier(body: SupplierCreate, db: Session = Depends(get_db)): + supplier = Supplier( + name=body.name, + adapter_type=body.adapter_type, + api_endpoint=body.api_endpoint, + account_id=body.account_id, + fulfillment_caps=body.fulfillment_caps or {}, + shipping_regions=body.shipping_regions or [], + reliability_score=body.reliability_score, + status=body.status, + ) + db.add(supplier) + db.commit() + db.refresh(supplier) + return supplier + + +@router.get("/{supplier_id}/performance", response_model=SupplierPerformanceOut) +def supplier_performance(supplier_id, db: Session = Depends(get_db)): + supplier = db.query(Supplier).filter(Supplier.id == supplier_id).first() + if supplier is None: + raise HTTPException(status_code=404, detail="supplier not found") + + perf = ( + db.query(SupplierPerformance) + .filter(SupplierPerformance.supplier_id == supplier_id) + .first() + ) + if perf is None: + # synthesize defaults; nothing measured yet + return SupplierPerformanceOut( + supplier_id=supplier_id, + supplier_name=supplier.name, + fulfillment_rate=100.0, + avg_shipping_days=None, + cancellation_rate=0.0, + stock_accuracy=100.0, + return_rate=0.0, + ) + return SupplierPerformanceOut( + supplier_id=supplier_id, + supplier_name=supplier.name, + fulfillment_rate=float(perf.fulfillment_rate or 0), + avg_shipping_days=float(perf.avg_shipping_days) if perf.avg_shipping_days is not None else None, + cancellation_rate=float(perf.cancellation_rate or 0), + stock_accuracy=float(perf.stock_accuracy or 0), + return_rate=float(perf.return_rate or 0), + updated_at=perf.updated_at, + ) diff --git a/backend/app/schemas.py b/backend/app/schemas.py new file mode 100644 index 0000000..2ef082a --- /dev/null +++ b/backend/app/schemas.py @@ -0,0 +1,140 @@ +"""Pydantic v2 request/response schemas.""" +from datetime import datetime +from typing import Any, List, Optional +from uuid import UUID + +from pydantic import BaseModel, ConfigDict, Field + + +class ORMModel(BaseModel): + model_config = ConfigDict(from_attributes=True) + + +# ── auth ──────────────────────────────────────────────────────────────── +class LoginRequest(BaseModel): + email: str + password: str + + +class TokenResponse(BaseModel): + access_token: str + token_type: str = "bearer" + + +class CustomerRegister(BaseModel): + email: str + name: Optional[str] = None + shipping_addr: Optional[dict] = None + # NOTE: customers table has no password column — password (if sent) is not + # stored. Customer identity for MVP is email-only; TODO: add password column + hash. + password: Optional[str] = None + + +# ── products ──────────────────────────────────────────────────────────── +class ProductOut(ORMModel): + id: UUID + sku: str + title: str + description: Optional[str] = None + category: Optional[str] = None + brand: Optional[str] = None + manufacturer: Optional[str] = None + images: Any = None + weight_grams: Optional[float] = None + cost: float + retail_price: Optional[float] = None + min_price: Optional[float] = None + max_price: Optional[float] = None + status: str + created_at: Optional[datetime] = None + updated_at: Optional[datetime] = None + + +class ImportRequest(BaseModel): + supplier_id: UUID + feed_type: str = "csv" + data_or_url: str = "" + + +class PriceHistoryOut(ORMModel): + price: float + reason: Optional[str] = None + created_at: Optional[datetime] = None + + +# ── suppliers ─────────────────────────────────────────────────────────── +class SupplierCreate(BaseModel): + name: str + adapter_type: str = "csv" + api_endpoint: Optional[str] = None + account_id: Optional[str] = None + fulfillment_caps: dict = {} + shipping_regions: List[str] = [] + reliability_score: float = 50.0 + status: str = "active" + + +class SupplierOut(ORMModel): + id: UUID + name: str + adapter_type: str + api_endpoint: Optional[str] = None + fulfillment_caps: Any = None + shipping_regions: Any = None + reliability_score: float + status: str + created_at: Optional[datetime] = None + + +class SupplierPerformanceOut(BaseModel): + supplier_id: UUID + supplier_name: str + fulfillment_rate: float + avg_shipping_days: Optional[float] = None + cancellation_rate: float + stock_accuracy: float + return_rate: float + updated_at: Optional[datetime] = None + + +# ── customers ─────────────────────────────────────────────────────────── +class CustomerOut(ORMModel): + id: UUID + email: str + name: Optional[str] = None + shipping_addr: Any = None + created_at: Optional[datetime] = None + + +# ── orders ────────────────────────────────────────────────────────────── +class OrderItemIn(BaseModel): + product_id: UUID + qty: int = Field(ge=1) + + +class OrderCreate(BaseModel): + customer_email: str + items: List[OrderItemIn] = Field(min_length=1) + + +class OrderOut(ORMModel): + id: UUID + order_number: str + customer_id: Optional[UUID] = None + items: Any + retail_total: float + supplier_id: Optional[UUID] = None + supplier_cost: float + shipping_cost: float + fees: float + profit: float + status: str + tracking: Optional[str] = None + carrier: Optional[str] = None + created_at: Optional[datetime] = None + updated_at: Optional[datetime] = None + + +class TrackingIn(BaseModel): + tracking: str + carrier: Optional[str] = None diff --git a/backend/app/tasks.py b/backend/app/tasks.py new file mode 100644 index 0000000..503a5d1 --- /dev/null +++ b/backend/app/tasks.py @@ -0,0 +1,113 @@ +"""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() diff --git a/backend/app/version.py b/backend/app/version.py new file mode 100644 index 0000000..1d77f74 --- /dev/null +++ b/backend/app/version.py @@ -0,0 +1,2 @@ +"""Shared version constant (avoids circular imports).""" +VERSION = "0.1.0" diff --git a/backend/requirements.txt b/backend/requirements.txt new file mode 100644 index 0000000..b1a2aad --- /dev/null +++ b/backend/requirements.txt @@ -0,0 +1,9 @@ +fastapi>=0.110,<1 +uvicorn[standard]>=0.29 +sqlalchemy>=2.0,<3 +psycopg2-binary>=2.9 +pydantic>=2.6,<3 +python-dotenv>=1.0 +celery[redis]>=5.3,<6 +redis>=5.0 +PyJWT>=2.8 diff --git a/backend/scripts/celery_worker.sh b/backend/scripts/celery_worker.sh new file mode 100755 index 0000000..16057ea --- /dev/null +++ b/backend/scripts/celery_worker.sh @@ -0,0 +1,4 @@ +#!/bin/sh +# Celery worker + beat entrypoint (used by the OPS agent's compose service). +# Beat schedule is defined in app/tasks.py (sync_inventory every 5min). +exec celery -A app.tasks.celery_app worker --beat --loglevel=info diff --git a/backend/scripts/seed.py b/backend/scripts/seed.py new file mode 100644 index 0000000..802ed2e --- /dev/null +++ b/backend/scripts/seed.py @@ -0,0 +1,159 @@ +#!/usr/bin/env python3 +"""Seed the sample supplier, import its catalog, price + publish valid products. + +Idempotent: safe to re-run — supplier/catalog are upserted, prices recalculated. + +Usage: python3 scripts/seed.py +""" +import os +import sys + +# allow running from repo root or scripts/ +sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__)))) + +from decimal import ROUND_HALF_UP, Decimal + +from app.database import SessionLocal +from app.engines.importer import import_feed +from app.engines.pricing import calculate_price +from app.engines.suppliers.base import build_adapter +from app.models import ( + AuditLog, + PriceHistory, + Supplier, + SupplierPerformance, + SupplierProduct, +) + +SAMPLE_NAME = "Sample Home & Gadgets" + + +def ensure_supplier(db): + supplier = db.query(Supplier).filter(Supplier.name == SAMPLE_NAME).first() + if supplier is None: + supplier = Supplier( + name=SAMPLE_NAME, + adapter_type="sample", + shipping_regions=["US", "CA"], + reliability_score=92.0, + status="active", + ) + db.add(supplier) + db.flush() + print(f"[seed] created supplier {supplier.id}") + else: + supplier.status = "active" + print(f"[seed] supplier already exists: {supplier.id}") + return supplier + + +def ensure_performance(db, supplier): + if db.query(SupplierPerformance).filter(SupplierPerformance.supplier_id == supplier.id).first() is None: + db.add( + SupplierPerformance( + supplier_id=supplier.id, + fulfillment_rate=98.0, + avg_shipping_days=4.0, + cancellation_rate=1.0, + stock_accuracy=97.0, + return_rate=2.0, + ) + ) + db.flush() + + +def price_and_publish(db, supplier): + """Price every linked product via the pricing engine; publish valid ones. + + NEVER publish invalid: products that fail rule lookup or end up with a + non-positive retail price stay un-published. + """ + links = ( + db.query(SupplierProduct) + .filter(SupplierProduct.supplier_id == supplier.id) + .all() + ) + published = 0 + skipped = [] + for sp in links: + product = sp.product + if product is None: + continue + quote = calculate_price( + db, + product.cost, + shipping=sp.shipping_cost or 0, + fees=0, + min_price=product.min_price, + ) + if quote is None: + skipped.append({"sku": product.sku, "reason": "no active price rule"}) + continue + if quote["retail"] <= 0: + skipped.append({"sku": product.sku, "reason": "non-positive retail"}) + continue + + product.retail_price = quote["retail"] + product.min_price = (quote["retail"] * Decimal("0.90")).quantize( + Decimal("0.01"), rounding=ROUND_HALF_UP + ) + product.max_price = (quote["retail"] * Decimal("1.15")).quantize( + Decimal("0.01"), rounding=ROUND_HALF_UP + ) + product.status = "price_calculated" + db.add(PriceHistory(product_id=product.id, price=quote["retail"], reason="auto")) + + # published only when priced + product.status = "published" + published += 1 + db.add( + AuditLog( + actor="seed", + action="product_published", + entity="product", + entity_id=str(product.id), + detail={"sku": product.sku, "retail_price": float(quote["retail"])}, + ) + ) + + db.commit() + print(f"[seed] priced + published {published} products; skipped: {skipped}") + + +def main(): + db = SessionLocal() + try: + supplier = ensure_supplier(db) + db.commit() + ensure_performance(db, supplier) + db.commit() + + summary = import_feed(db, supplier, feed_type="sample", data_or_url="") + print( + f"[seed] import: received={summary['received']} " + f"imported={summary['imported']} updated={summary['updated']} " + f"skipped={len(summary['skipped'])}" + ) + if summary["skipped"]: + print(f"[seed] skipped details: {summary['skipped']}") + + price_and_publish(db, supplier) + + # quick self-check on a known item ($80-ish cost band) + probe = db.query(SupplierProduct).filter( + SupplierProduct.supplier_id == supplier.id, + SupplierProduct.supplier_sku == "SPL-1006", + ).first() + if probe and probe.product: + p = probe.product + print( + f"[seed] sanity: {p.sku} cost={p.cost} retail={p.retail_price} " + f"status={p.status}" + ) + print("[seed] done") + finally: + db.close() + + +if __name__ == "__main__": + main()