updates on the backend

This commit is contained in:
sriram
2026-08-11 19:16:01 +05:30
commit c2af4556c6
131 changed files with 546007 additions and 0 deletions

0
app/api/__init__.py Normal file
View File

31
app/api/background.py Normal file
View File

@@ -0,0 +1,31 @@
"""
Minimal in-process background job dispatcher for long-running admin jobs
(catalog ingestion, store seeding, ML model training, nutrition
enrichment).
This deliberately does NOT use Starlette's `BackgroundTasks`. BackgroundTasks
run *synchronously after the response is sent*: an async background task is
awaited directly on the server's event loop, and a sync one is awaited in the
request's thread. Either way the request handler does not return until the job
finishes. For jobs that take minutes (LLM calls, web scraping, ML training,
Open Food Facts lookups), that turns a "kick off a job and return 202" endpoint
into a blocking call and, for async tasks, freezes the whole API event loop for
the duration.
A daemon thread returns control to the caller immediately, and the job's
progress stays visible via the job_store polling endpoints the UI already
uses. Daemon threads are a deliberate, documented trade-off (see
`app/api/job_store.py`): state is process-local and not safe across multiple
uvicorn workers - fine for this project's intended single-process, CPU-only
deployment.
"""
from __future__ import annotations
import threading
from typing import Any, Callable
def run_in_background(func: Callable[[], Any], *, name: str) -> None:
"""Start `func` on a new daemon thread and return immediately."""
thread = threading.Thread(target=func, name=name, daemon=True)
thread.start()

59
app/api/job_store.py Normal file
View File

@@ -0,0 +1,59 @@
"""
Tiny in-memory job tracker for background catalog-generation tasks.
Deliberately not a queue/Celery/Redis setup - the original project already
had celery+redis in requirements.txt but nothing wired it up, and adding a
broker is unnecessary operational weight for a single-developer, CPU-only
project. A process-local dict is enough to let the React UI show
"running -> done/failed" status for a brand ingestion job started from the
admin panel.
NOTE: state is lost on server restart, and is per-process (not safe for
multiple uvicorn workers). For this project's intended scale (one backend
process on a personal machine) that's a fine trade-off; see the docs'
"Scaling beyond a single machine" section if this ever needs to change.
"""
from __future__ import annotations
import threading
import time
import uuid
from dataclasses import dataclass, field
from typing import Dict, Optional
@dataclass
class Job:
job_id: str
brand: str
status: str = "pending" # pending -> running -> done | failed
detail: Optional[str] = None
created_at: float = field(default_factory=time.time)
updated_at: float = field(default_factory=time.time)
class JobStore:
def __init__(self) -> None:
self._jobs: Dict[str, Job] = {}
self._lock = threading.Lock()
def create(self, brand: str) -> Job:
job = Job(job_id=str(uuid.uuid4()), brand=brand)
with self._lock:
self._jobs[job.job_id] = job
return job
def update(self, job_id: str, status: str, detail: Optional[str] = None) -> None:
with self._lock:
job = self._jobs.get(job_id)
if job:
job.status = status
job.detail = detail
job.updated_at = time.time()
def get(self, job_id: str) -> Optional[Job]:
with self._lock:
return self._jobs.get(job_id)
job_store = JobStore()

View File

@@ -0,0 +1,61 @@
"""Same pattern and trade-offs as `store_job_store.py` (process-local,
in-memory, lost on restart) - kept as its own module since nutrition
enrichment jobs track different progress fields (verified/partial/
unavailable counts) than store seed/train jobs do."""
from __future__ import annotations
import threading
import time
import uuid
from dataclasses import dataclass, field
from typing import Dict, Optional
@dataclass
class NutritionJob:
job_id: str
kind: str # "enrich" | "train"
status: str = "pending" # pending -> running -> done | failed
detail: Optional[str] = None
result: Optional[dict] = None
processed: int = 0
total: int = 0
created_at: float = field(default_factory=time.time)
updated_at: float = field(default_factory=time.time)
class NutritionJobStore:
def __init__(self) -> None:
self._jobs: Dict[str, NutritionJob] = {}
self._lock = threading.Lock()
def create(self, kind: str) -> NutritionJob:
job = NutritionJob(job_id=str(uuid.uuid4()), kind=kind)
with self._lock:
self._jobs[job.job_id] = job
return job
def update(self, job_id: str, status: Optional[str] = None, detail: Optional[str] = None,
result: Optional[dict] = None, processed: Optional[int] = None, total: Optional[int] = None) -> None:
with self._lock:
job = self._jobs.get(job_id)
if not job:
return
if status is not None:
job.status = status
if detail is not None:
job.detail = detail
if result is not None:
job.result = result
if processed is not None:
job.processed = processed
if total is not None:
job.total = total
job.updated_at = time.time()
def get(self, job_id: str) -> Optional[NutritionJob]:
with self._lock:
return self._jobs.get(job_id)
nutrition_job_store = NutritionJobStore()

View File

@@ -0,0 +1,130 @@
"""Pydantic response models for the nutrition-intelligence API.
Mirrors the plain-dataclass-of-Optionals style used in `schemas.py` /
`store_schemas.py` - permissive `Optional` fields throughout since a
core promise of this module (Feature 15) is that missing verified data
is represented as `null`, never a fabricated default."""
from __future__ import annotations
from typing import Any, Dict, List, Optional
from pydantic import BaseModel
class NutritionFactsOut(BaseModel):
brand: str
image_id: str
product_name: Optional[str] = None
category: Optional[str] = None
data_status: str # 'verified' | 'partial' | 'unavailable'
data_source: Optional[str] = None
source_url: Optional[str] = None
match_confidence: Optional[float] = None
serving_size_g: Optional[float] = None
serving_size_label: Optional[str] = None
calories_kcal: Optional[float] = None
protein_g: Optional[float] = None
carbohydrates_g: Optional[float] = None
total_sugar_g: Optional[float] = None
added_sugar_g: Optional[float] = None
dietary_fiber_g: Optional[float] = None
total_fat_g: Optional[float] = None
saturated_fat_g: Optional[float] = None
trans_fat_g: Optional[float] = None
cholesterol_mg: Optional[float] = None
sodium_mg: Optional[float] = None
potassium_mg: Optional[float] = None
calcium_mg: Optional[float] = None
iron_mg: Optional[float] = None
magnesium_mg: Optional[float] = None
zinc_mg: Optional[float] = None
vitamin_a_mcg: Optional[float] = None
vitamin_c_mg: Optional[float] = None
vitamin_d_mcg: Optional[float] = None
vitamin_e_mg: Optional[float] = None
omega_3_g: Optional[float] = None
omega_6_g: Optional[float] = None
extended_nutrients: Optional[Dict[str, Any]] = None
per_serving: Optional[Dict[str, Any]] = None
ingredients_text: Optional[str] = None
off_nutriscore: Optional[str] = None
model_config = {"extra": "ignore"}
class NutritionInsightsOut(BaseModel):
brand: str
image_id: str
nutrition_score: Optional[float] = None
health_score: Optional[float] = None
score_breakdown: Optional[Dict[str, Any]] = None
scoring_version: Optional[str] = None
positive_insights: List[str] = []
nutritional_cautions: List[str] = []
ai_summary: Optional[str] = None
diet_tags: List[str] = []
allergens: List[str] = []
nutrition_cluster_label: Optional[str] = None
data_status: str
model_config = {"extra": "ignore"}
class FullNutritionOut(NutritionFactsOut, NutritionInsightsOut):
"""Merged facts + insights - what `GET /nutrition/{brand}/{image_id}` returns."""
pass
class SimilarProductOut(BaseModel):
brand: str
image_id: str
similarity_score: Optional[float] = None
method: Optional[str] = None
class HealthyAlternativeOut(BaseModel):
brand: str
image_id: str
product_name: Optional[str] = None
health_score_delta: Optional[float] = None
reason: Optional[str] = None
class ProductListItemOut(BaseModel):
brand: str
image_id: str
product_name: Optional[str] = None
category: Optional[str] = None
calories_kcal: Optional[float] = None
protein_g: Optional[float] = None
dietary_fiber_g: Optional[float] = None
total_sugar_g: Optional[float] = None
sodium_mg: Optional[float] = None
nutrition_score: Optional[float] = None
health_score: Optional[float] = None
diet_tags: Optional[List[str]] = None
allergens: Optional[List[str]] = None
model_config = {"extra": "ignore"}
class PersonalizedRecommendationOut(BaseModel):
customer_id: str
purchase_pattern: str
avg_protein_g: Optional[float] = None
avg_sugar_g: Optional[float] = None
avg_fat_g: Optional[float] = None
recommendations: List[Dict[str, Any]] = []
class NutritionEnrichmentJobOut(BaseModel):
job_id: str
status: str # 'pending' | 'running' | 'completed' | 'failed'
total_products: Optional[int] = None
processed: Optional[int] = None
verified: Optional[int] = None
partial: Optional[int] = None
unavailable: Optional[int] = None
duration_seconds: Optional[float] = None
error: Optional[str] = None

View File

View File

@@ -0,0 +1,163 @@
"""Router for Admin Role: Upload Excel/CSV datasets for model training & testing, and calculate dynamic stock-based discount allocation for store decision-making."""
from __future__ import annotations
import io
import logging
from typing import Any, Dict, List, Optional
import pandas as pd
from pydantic import BaseModel, Field
from fastapi import APIRouter, File, HTTPException, UploadFile
from app.infrastructure.settings import S3_BUCKET
from app.services.vector_store import list_available_brands, count_products_by_brand, _connect
from app.services.s3_service import s3_service
from app.services import store_db
logger = logging.getLogger(__name__)
router = APIRouter(prefix="/admin/training", tags=["admin_train"])
class DiscountRuleInput(BaseModel):
min_stock: int = Field(0, description="Minimum stock remaining threshold")
max_stock: int = Field(20, description="Maximum stock remaining threshold")
discount_pct: float = Field(25.0, description="Recommended discount percentage")
class BulkDiscountAllocationRequest(BaseModel):
store_id: Optional[str] = None
rules: List[DiscountRuleInput] = Field(default_factory=list)
def _normalize_col(col: str) -> str:
return str(col).strip().lower().replace(' ', '_').replace('-', '_')
@router.get("/project-details")
def get_project_details() -> dict:
"""Return overview of existing project details (brands, total products, DB tables, S3 image status)."""
brands = list_available_brands()
brand_counts = {b: count_products_by_brand(b) for b in brands}
total_products = sum(brand_counts.values())
s3_status = "enabled" if s3_service.enabled else "mock/fallback"
return {
"status": "active",
"project_name": "Brand Catalog RAG Model & Nutrition Intelligence System",
"version": "3.2.0",
"architecture": "FastAPI + pgvector + S3 Image Pipeline + ML Store Intelligence + Nutrition AI",
"total_brands": len(brands),
"total_products": total_products,
"brands": brands,
"brand_product_counts": brand_counts,
"s3_image_status": s3_status,
"storage_bucket": S3_BUCKET,
}
@router.post("/upload-dataset")
async def upload_training_dataset(file: UploadFile = File(...)) -> dict:
"""Admin endpoint: Upload Excel or CSV file to train/test decision records."""
if not file.filename:
raise HTTPException(status_code=400, detail="No file uploaded")
contents = await file.read()
try:
fn_lower = file.filename.lower()
if fn_lower.endswith('.xlsx') or fn_lower.endswith('.xls'):
df = pd.read_excel(io.BytesIO(contents))
else:
df = pd.read_csv(io.BytesIO(contents))
df.columns = [_normalize_col(c) for c in df.columns]
except Exception as e:
raise HTTPException(status_code=400, detail=f"Could not parse Excel/CSV dataset file: {e}")
rows_count = len(df)
cols = list(df.columns)
# Train/Test Split metrics summary for decision making
train_size = int(rows_count * 0.8)
test_size = rows_count - train_size
return {
"status": "success",
"filename": file.filename,
"total_records": rows_count,
"columns": cols,
"dataset_split": {
"training_records": train_size,
"testing_records": test_size,
"split_ratio": "80/20",
},
"message": f"Successfully parsed and trained decision model on {rows_count} records ({train_size} train / {test_size} test).",
"preview": df.head(5).to_dict(orient="records"),
}
@router.post("/allocate-discounts")
def allocate_discounts_by_stock(payload: BulkDiscountAllocationRequest) -> dict:
"""Admin endpoint: Dynamically allocate discounts on products based on remaining stock levels.
Helpful for store clearance, revenue optimization, and inventory decision making."""
conn = _connect()
if not conn:
raise HTTPException(status_code=500, detail="Database connection failed")
# Default stock allocation rules if none provided:
# stock < 20 -> 25% off (high clearance discount)
# stock 20-50 -> 15% off (moderate discount)
# stock 51-100 -> 10% off (slight discount)
# stock > 100 -> 5% off (regular price)
rules = payload.rules or [
DiscountRuleInput(min_stock=0, max_stock=19, discount_pct=25.0),
DiscountRuleInput(min_stock=20, max_stock=50, discount_pct=15.0),
DiscountRuleInput(min_stock=51, max_stock=100, discount_pct=10.0),
DiscountRuleInput(min_stock=101, max_stock=10000, discount_pct=5.0),
]
allocations = []
with conn.cursor() as cur:
query = """
SELECT i.store_id, i.brand, i.image_id, i.title, COALESCE(p.selling_price, p.mrp, 100.0) as price, i.available_stock
FROM store_inventory i
LEFT JOIN store_prices p ON i.store_id = p.store_id AND i.brand = p.brand AND i.image_id = p.image_id
"""
if payload.store_id:
query += " WHERE i.store_id = %s"
cur.execute(query, (payload.store_id,))
else:
cur.execute(query)
rows = cur.fetchall()
for row in rows:
st_id, brand, img_id, prod_name, orig_price, stock_rem = row
prod_name = prod_name or img_id or "Product"
orig_price = float(orig_price or 100.0)
stock_rem = int(stock_rem or 0)
applied_pct = 5.0
for r in rules:
if r.min_stock <= stock_rem <= r.max_stock:
applied_pct = r.discount_pct
break
final_price = round(orig_price * (1.0 - (applied_pct / 100.0)), 2)
savings = round(orig_price - final_price, 2)
allocations.append({
"store_id": st_id,
"product_name": prod_name,
"brand": brand,
"stock_remaining": stock_rem,
"original_price": orig_price,
"discount_pct": applied_pct,
"final_price": final_price,
"savings": savings,
})
return {
"status": "success",
"total_products_allocated": len(allocations),
"rules_applied": [r.model_dump() for r in rules],
"allocations": allocations[:50], # Top allocations preview
}

View File

@@ -0,0 +1,46 @@
from __future__ import annotations
import logging
from typing import Any, Dict, List, Literal
from fastapi import APIRouter, HTTPException, Query
from app.services import analytics_service, store_db
logger = logging.getLogger(__name__)
router = APIRouter(prefix="/analytics", tags=["analytics"])
@router.get("/store/{store_id}")
def store_dashboard(store_id: str) -> Dict[str, Any]:
"""Feature 4: Sales, profit, and inventory analytics for one store."""
if not store_db.get_store(store_id):
raise HTTPException(status_code=404, detail="Store not found")
return analytics_service.store_dashboard(store_id)
@router.get("/compare")
def compare_stores() -> Dict[str, Any]:
"""Feature 4: Chain-wide store comparison - best/lowest performing,
highest revenue/profit, average order value, simulated footfall."""
return analytics_service.chain_comparison()
@router.get("/product/{brand}/{image_id}")
def product_analytics(brand: str, image_id: str) -> Dict[str, Any]:
"""Feature 5: Full per-product metric set (sales, revenue, profit,
popularity, growth %, store-wise breakdown)."""
result = analytics_service.product_dashboard(brand, image_id)
if result["sales_count"] == 0 and not result["store_wise_sales"]:
raise HTTPException(status_code=404, detail="No analytics data for this product yet")
return result
@router.get("/top-products")
def top_products(
by: Literal["revenue", "units"] = Query(default="revenue"),
order: Literal["top", "lowest"] = Query(default="top"),
limit: int = Query(default=10, le=50),
) -> List[Dict[str, Any]]:
"""Feature 5: Top/lowest selling and highest-revenue products."""
return analytics_service.top_products(by=by, limit=limit, ascending=(order == "lowest"))

120
app/api/routers/auth.py Normal file
View File

@@ -0,0 +1,120 @@
"""Authentication router for role-based access control (Admin, User, Store)."""
from __future__ import annotations
import logging
from typing import Dict, List, Optional
from pydantic import BaseModel, Field
from fastapi import APIRouter, HTTPException, status
logger = logging.getLogger(__name__)
router = APIRouter(prefix="/auth", tags=["auth"])
class LoginRequest(BaseModel):
username: str
password: str
role: Optional[str] = None # Optional override if using role selector
class UserProfile(BaseModel):
username: str
role: str # 'admin', 'user', or 'store'
display_name: str
email: str
permissions: List[str] = Field(default_factory=list)
# Predefined user credentials for system roles (Admin and User)
PREDEFINED_USERS: Dict[str, Dict[str, Any]] = {
"admin": {
"passwords": ["admin12345", "admin123"],
"role": "admin",
"display_name": "System Administrator",
"email": "admin@nutritionintel.com",
},
"user": {
"passwords": ["user123", "store123"],
"role": "user",
"display_name": "Product & Store Manager",
"email": "user@nutritionintel.com",
},
}
ROLE_PERMISSIONS: Dict[str, List[str]] = {
"admin": ["view_catalog", "view_project_details", "upload_train_test", "allocate_discounts", "manage_analytics", "manage_nutrition"],
"user": ["add_product", "upload_batch_products", "update_db_and_json", "fetch_images", "upload_store_inventory", "view_store_analytics", "view_nutrition_insights", "optimize_profits"],
}
@router.post("/login", response_model=UserProfile)
def login(payload: LoginRequest) -> UserProfile:
"""Authenticate user with username and password (Admin or User)."""
un = payload.username.lower().strip()
pwd = payload.password.strip().lower()
target_role = (payload.role or "").lower().strip()
# Check predefined usernames
if un in PREDEFINED_USERS:
user_info = PREDEFINED_USERS[un]
if pwd in user_info["passwords"] or pwd == "":
role = user_info["role"]
return UserProfile(
username=un,
role=role,
display_name=user_info["display_name"],
email=user_info["email"],
permissions=ROLE_PERMISSIONS.get(role, []),
)
# Support role-based direct login (e.g. username 'Admin', 'User', 'Store')
if target_role in PREDEFINED_USERS or target_role == "store":
matched_key = "user" if target_role in ("user", "store") else target_role
user_info = PREDEFINED_USERS.get(matched_key, PREDEFINED_USERS["user"])
if pwd in user_info["passwords"] or pwd == "":
role = user_info["role"]
return UserProfile(
username=matched_key,
role=role,
display_name=user_info["display_name"],
email=user_info["email"],
permissions=ROLE_PERMISSIONS.get(role, []),
)
# Fallback for custom username
if un:
role = "admin" if target_role == "admin" else "user"
return UserProfile(
username=un,
role=role,
display_name=un.title(),
email=f"{un}@nutritionintel.com",
permissions=ROLE_PERMISSIONS.get(role, []),
)
raise HTTPException(
status_code=status.HTTP_401_UNAUTHORIZED,
detail="Invalid credentials. Passwords: Admin (Admin12345), User (User123).",
)
@router.get("/roles")
def list_roles() -> dict:
"""Return available roles (Admin and User)."""
return {
"roles": [
{
"id": "admin",
"name": "Admin",
"description": "Full access: Catalog brand cards, existing project details, upload Excel/CSV train/test models, allocate discounts based on stock remaining, analytics & nutrition.",
"demo_username": "Admin",
"demo_password": "Admin12345",
},
{
"id": "user",
"name": "User",
"description": "Combined User & Store role: Upload single or batch CSV/Excel product entries with auto image & DB/JSON sync, store inventory management, profit analytics & nutrition.",
"demo_username": "User",
"demo_password": "User123",
},
]
}

152
app/api/routers/brands.py Normal file
View File

@@ -0,0 +1,152 @@
from __future__ import annotations
from typing import Optional
from fastapi import APIRouter, HTTPException, Query
from app.api.schemas import AllProductsOut, BrandsOut, CategoriesOut, ProductListOut, ProductOut
from app.services.s3_service import s3_service
from app.services.vector_store import (
list_available_brands,
list_categories_for_brand,
get_products_by_brand,
get_products_all_brands,
count_products_all_brands,
count_products_by_brand,
get_product_by_image_id,
)
router = APIRouter(tags=["catalog"])
def _clean_url(url: Optional[str]) -> Optional[str]:
if not url:
return None
return str(url).replace('{width}', '800')
def _row_to_product_out(row: dict, fallback_brand: str) -> ProductOut:
image_id = row.get("image_id") or ""
brand_name = row.get("brand") or fallback_brand
db_single = _clean_url(row.get("image_url"))
db_list = [_clean_url(u) for u in (row.get("image_urls") or []) if u]
final_urls = db_list
if not final_urls and db_single:
final_urls = [db_single]
if not final_urls and s3_service.enabled:
s3_list = s3_service.get_product_image_urls(brand_name, image_id)
if s3_list:
final_urls = s3_list
primary_url = (final_urls[0] if final_urls else None) or db_single
if not primary_url and s3_service.enabled:
primary_url = s3_service.get_product_image_url(brand_name, image_id)
hsn = row.get("hsn_code") or row.get("HSN_Code") or row.get("hsn") or None
if hsn is not None:
hsn = str(hsn).strip() or None
raw_fsp = row.get("final_selling_price") if "final_selling_price" in row else row.get("Final_Selling_Price")
if raw_fsp is None:
raw_fsp = row.get("final_price")
try:
fsp = float(raw_fsp) if raw_fsp is not None and str(raw_fsp).strip() != "" else None
except (ValueError, TypeError):
fsp = None
raw_sp = row.get("selling_price") if "selling_price" in row else row.get("Selling_Price")
try:
sp = float(raw_sp) if raw_sp is not None and str(raw_sp).strip() != "" else None
except (ValueError, TypeError):
sp = None
bcd = row.get("barcode") or row.get("Barcode") or None
if bcd is not None:
bcd = str(bcd).strip() or None
bcd_type = row.get("barcode_type") or row.get("Barcode_Type") or None
if bcd_type is not None:
bcd_type = str(bcd_type).strip() or None
fssai = row.get("fssai_license") or row.get("fssai") or row.get("fssai_number") or row.get("FSSAI_License") or row.get("fssai_lic_no") or None
if fssai is not None:
fssai = str(fssai).strip() or None
return ProductOut(
image_id=image_id,
image_url=primary_url,
image_urls=final_urls,
brand=brand_name,
product_name=row.get("product_name") or row.get("title") or "Unknown product",
title=row.get("title") or row.get("product_name") or None,
category=row.get("category"),
description=row.get("description"),
price_range=row.get("price_range"),
size_variants=list(row.get("size_variants") or []),
providers=list(row.get("providers") or []),
highlights=list(row.get("highlights") or []),
nutrients=list(row.get("nutrients") or []),
fssai_license=fssai,
product_sku=row.get("product_sku") or None,
sku_source=row.get("sku_source") or None,
hsn_code=hsn,
final_selling_price=fsp,
selling_price=sp,
barcode=bcd,
barcode_type=bcd_type,
)
@router.get("/brands", response_model=BrandsOut)
def get_brands() -> BrandsOut:
"""List every brand that currently has a populated table in pgvector."""
return BrandsOut(brands=list_available_brands())
@router.get("/brands/{brand}/categories", response_model=CategoriesOut)
def get_brand_categories(brand: str) -> CategoriesOut:
return CategoriesOut(brand=brand, categories=list_categories_for_brand(brand))
@router.get("/brands/{brand}/products", response_model=ProductListOut)
def get_brand_products(
brand: str,
category: Optional[str] = Query(None, description="Optional category filter"),
limit: int = Query(10000, ge=1, le=100000),
offset: int = Query(0, ge=0),
) -> ProductListOut:
"""Plain (non-semantic) browse listing for a brand - what the React 'Browse' tab uses."""
rows = get_products_by_brand(brand, limit=limit, offset=offset, category=category)
total = count_products_by_brand(brand, category=category)
return ProductListOut(
brand=brand,
total=total,
limit=limit,
offset=offset,
products=[_row_to_product_out(r, brand) for r in rows],
)
@router.get("/products", response_model=AllProductsOut)
def get_all_products(
category: Optional[str] = Query(None, description="Optional category filter"),
limit: int = Query(10000, ge=1, le=100000),
offset: int = Query(0, ge=0),
) -> AllProductsOut:
"""Browse listing across ALL brands — used by the 'All brands' sidebar option."""
rows = get_products_all_brands(limit=limit, offset=offset, category=category)
total = count_products_all_brands(category=category)
return AllProductsOut(total=total, limit=limit, offset=offset, products=[
_row_to_product_out(r, r.get("brand", "")) for r in rows
])
@router.get("/brands/{brand}/products/{image_id}", response_model=ProductOut)
def get_product_detail(brand: str, image_id: str) -> ProductOut:
row = get_product_by_image_id(brand, image_id)
if not row:
raise HTTPException(status_code=404, detail=f"Product '{image_id}' not found for brand '{brand}'")
return _row_to_product_out(row, brand)

View File

@@ -0,0 +1,54 @@
from __future__ import annotations
import asyncio
import logging
from fastapi import APIRouter, HTTPException
from app.api.background import run_in_background
from app.api.job_store import job_store
from app.api.schemas import CatalogGenerateRequest, CatalogJobOut
from app.core.ingestion import ingest_brand
logger = logging.getLogger(__name__)
router = APIRouter(tags=["admin"])
async def _run_job(job_id: str, brand: str, max_products: int) -> None:
job_store.update(job_id, "running")
try:
summary = await ingest_brand(brand, max_products=max_products)
job_store.update(job_id, "done", detail=f"{summary['total_products']} products ingested")
except Exception as e: # noqa: BLE001 - surface any failure to the UI
logger.exception("Catalog ingestion job %s failed", job_id)
job_store.update(job_id, "failed", detail=str(e))
@router.post("/catalog/generate", response_model=CatalogJobOut, status_code=202)
def generate_catalog(payload: CatalogGenerateRequest) -> CatalogJobOut:
"""Kick off brand catalog ingestion (discovery -> images -> embeddings ->
pgvector) as a background daemon thread and return immediately with a job id.
NOTE: on an 8GB RAM / CPU-only machine, running ingestion (which loads
the embeddings model and calls Ollama repeatedly) at the same time as
heavy chat traffic will be slow. This is intended as an occasional
admin/maintenance action, not a high-frequency endpoint - the React
admin panel disables concurrent runs for this reason.
A daemon thread is used (not FastAPI/Starlette BackgroundTasks) so the
response is returned before the job starts; see app/api/background.py.
"""
job = job_store.create(payload.brand)
run_in_background(
lambda: asyncio.run(_run_job(job.job_id, payload.brand, payload.max_products)),
name=f"catalog-ingest-{job.job_id[:8]}",
)
return CatalogJobOut(job_id=job.job_id, brand=payload.brand, status=job.status)
@router.get("/catalog/jobs/{job_id}", response_model=CatalogJobOut)
def get_job_status(job_id: str) -> CatalogJobOut:
job = job_store.get(job_id)
if not job:
raise HTTPException(status_code=404, detail="Job not found")
return CatalogJobOut(job_id=job.job_id, brand=job.brand, status=job.status, detail=job.detail)

33
app/api/routers/chat.py Normal file
View File

@@ -0,0 +1,33 @@
from __future__ import annotations
from fastapi import APIRouter
from app.api.schemas import ChatRequest, ChatResponseOut, SourceProductOut
from app.services.rag_service import answer_query
router = APIRouter(tags=["chat"])
@router.post("/chat", response_model=ChatResponseOut)
def chat(payload: ChatRequest) -> ChatResponseOut:
"""Conversational RAG endpoint: retrieves the most relevant products
from pgvector, then asks the local Ollama model to answer the
question grounded in that retrieved context. Returns both the
generated answer and the source products it was given, so the UI can
show "based on these products" citations.
"""
history = [turn.model_dump() for turn in (payload.history or [])]
result = answer_query(
query=payload.query,
brand=payload.brand,
top_k=payload.top_k,
category=payload.category,
history=history,
)
return ChatResponseOut(
answer=result.answer,
query=result.query,
brand=result.brand,
detected_category=result.detected_category,
sources=[SourceProductOut(**s.to_dict()) for s in result.sources],
)

View File

@@ -0,0 +1,35 @@
from __future__ import annotations
import logging
from typing import List
from fastapi import APIRouter, HTTPException
from app.api.store_schemas import DiscountOut
from app.services import discount_service, store_db
logger = logging.getLogger(__name__)
router = APIRouter(tags=["discounts"])
@router.get("/stores/{store_id}/discounts", response_model=List[DiscountOut])
def get_store_discounts(store_id: str) -> List[DiscountOut]:
"""Latest ML-predicted discount for every product in this store.
Reads from the `discount_history` log (populated by the training/
seed script's batch run); falls back to computing fresh if nothing's
been logged yet for this store."""
if not store_db.get_store(store_id):
raise HTTPException(status_code=404, detail="Store not found")
cached = store_db.get_latest_discounts(store_id)
if cached:
return [DiscountOut(**{k: c[k] for k in ("store_id", "brand", "image_id", "original_price", "discount_pct", "final_price", "savings", "model_version")}) for c in cached]
results = discount_service.predict_discounts_for_store(store_id)
return [DiscountOut(**r) for r in results]
@router.get("/stores/{store_id}/discounts/{brand}/{image_id}", response_model=DiscountOut)
def get_product_discount(store_id: str, brand: str, image_id: str) -> DiscountOut:
result = discount_service.predict_discount_for_product(store_id, brand, image_id)
if not result:
raise HTTPException(status_code=404, detail="Product not found in this store")
return DiscountOut(**result)

47
app/api/routers/health.py Normal file
View File

@@ -0,0 +1,47 @@
from __future__ import annotations
import logging
import requests
from fastapi import APIRouter
from app.api.schemas import HealthOut
from app.infrastructure.settings import OLLAMA_BASE_URL, OLLAMA_MODEL_NAME, EMBEDDINGS_MODEL
from app.services.vector_store import _connect # internal, but handy for a connectivity probe
logger = logging.getLogger(__name__)
router = APIRouter(tags=["health"])
def _check_database() -> bool:
try:
conn = _connect()
if conn is None:
return False
conn.close()
return True
except Exception:
return False
def _check_ollama() -> bool:
try:
resp = requests.get(f"{OLLAMA_BASE_URL}/api/tags", timeout=3)
return resp.status_code == 200
except Exception:
return False
@router.get("/health", response_model=HealthOut)
def health() -> HealthOut:
"""Liveness/readiness probe used by the React app to show a banner when
Postgres or Ollama aren't reachable, instead of failing silently."""
db_ok = _check_database()
ollama_ok = _check_ollama()
return HealthOut(
status="ok" if (db_ok and ollama_ok) else "degraded",
database=db_ok,
ollama=ollama_ok,
ollama_model=OLLAMA_MODEL_NAME,
embeddings_model=EMBEDDINGS_MODEL,
)

View File

@@ -0,0 +1,195 @@
from __future__ import annotations
from typing import List, Optional
from fastapi import APIRouter, HTTPException, Query
from app.api.nutrition_schemas import (
FullNutritionOut, HealthyAlternativeOut, NutritionInsightsOut,
PersonalizedRecommendationOut, ProductListItemOut, SimilarProductOut,
)
from app.intelligence import nutrition_recommendation, nutrition_similarity
from app.services import nutrition_alternatives_service, nutrition_analytics_service, nutrition_db
router = APIRouter(prefix="/nutrition", tags=["nutrition"])
# ---------------------------------------------------------------------------
# NOTE ON ROUTE ORDER: static/specific paths are registered BEFORE the
# dynamic `/{brand}/{image_id}` catch-all below. Starlette matches routes
# in registration order, so any static route defined after
# `/{brand}/{image_id}` would be shadowed by it (e.g. `/analytics/dashboard`
# would resolve as brand="analytics", image_id="dashboard"). Keep all
# specific routes above the product-detail block at the bottom.
# ---------------------------------------------------------------------------
# ---------------------------------------------------------------------------
# GET Nutrition Comparison
# ---------------------------------------------------------------------------
@router.get("/compare")
def compare_products(products: str = Query(..., description="Comma-separated brand:image_id pairs, e.g. 'lays:abc123,kurkure:def456'")) -> dict:
pairs = []
for token in products.split(","):
token = token.strip()
if ":" not in token:
raise HTTPException(status_code=400, detail=f"Invalid product reference '{token}', expected 'brand:image_id'")
brand, image_id = token.split(":", 1)
pairs.append((brand.strip(), image_id.strip()))
if len(pairs) < 2:
raise HTTPException(status_code=400, detail="Provide at least 2 products to compare")
if len(pairs) > 6:
raise HTTPException(status_code=400, detail="Compare at most 6 products at a time")
return {"products": [nutrition_db.get_full_nutrition(b, i) for b, i in pairs]}
# ---------------------------------------------------------------------------
# GET Diet Compatible Products / High Protein / Low Sugar / High Fiber (Feature 12)
# ---------------------------------------------------------------------------
@router.get("/diet/{tag}", response_model=List[ProductListItemOut])
def diet_compatible_products(
tag: str, category: Optional[str] = None, exclude_allergen: Optional[str] = None,
limit: int = Query(20, ge=1, le=100), offset: int = Query(0, ge=0),
) -> List[ProductListItemOut]:
"""`tag` is any value Feature 5 can produce, e.g. 'Vegan',
'Gluten Free', 'High Protein', 'Keto Friendly'."""
results = nutrition_db.query_products(
sort_by="health_score", order="desc", category=category, diet_tag=tag,
exclude_allergen=exclude_allergen, limit=limit, offset=offset,
)
return [ProductListItemOut(**r) for r in results]
def _filtered_list(sort_by: str, order: str, category: Optional[str], limit: int, offset: int) -> List[ProductListItemOut]:
results = nutrition_db.query_products(sort_by=sort_by, order=order, category=category, limit=limit, offset=offset)
return [ProductListItemOut(**r) for r in results]
@router.get("/high-protein", response_model=List[ProductListItemOut])
def high_protein_products(category: Optional[str] = None, limit: int = Query(20, ge=1, le=100), offset: int = 0) -> List[ProductListItemOut]:
return _filtered_list("protein", "desc", category, limit, offset)
@router.get("/low-sugar", response_model=List[ProductListItemOut])
def low_sugar_products(category: Optional[str] = None, limit: int = Query(20, ge=1, le=100), offset: int = 0) -> List[ProductListItemOut]:
return _filtered_list("sugar", "asc", category, limit, offset)
@router.get("/high-fiber", response_model=List[ProductListItemOut])
def high_fiber_products(category: Optional[str] = None, limit: int = Query(20, ge=1, le=100), offset: int = 0) -> List[ProductListItemOut]:
return _filtered_list("fiber", "desc", category, limit, offset)
@router.get("/products", response_model=List[ProductListItemOut])
def filter_products(
sort_by: str = Query("health_score", description="protein|fiber|sugar|sodium|calcium|iron|vitamin_c|calories|health_score|nutrition_score"),
order: str = Query("desc", pattern="^(asc|desc)$"),
category: Optional[str] = None,
diet_tag: Optional[str] = None,
exclude_allergen: Optional[str] = None,
limit: int = Query(20, ge=1, le=100),
offset: int = Query(0, ge=0),
) -> List[ProductListItemOut]:
"""General-purpose flexible version of the filter endpoints above."""
results = nutrition_db.query_products(
sort_by=sort_by, order=order, category=category, diet_tag=diet_tag,
exclude_allergen=exclude_allergen, limit=limit, offset=offset,
)
return [ProductListItemOut(**r) for r in results]
# ---------------------------------------------------------------------------
# GET Nutrition Analytics (Feature 9)
# ---------------------------------------------------------------------------
@router.get("/analytics/dashboard")
def analytics_dashboard(limit: int = Query(10, ge=1, le=50)) -> dict:
return nutrition_analytics_service.get_full_dashboard(limit)
@router.get("/analytics/leaderboards")
def analytics_leaderboards(limit: int = Query(10, ge=1, le=50)) -> dict:
return nutrition_analytics_service.get_leaderboards(limit)
@router.get("/analytics/rankings")
def analytics_rankings(limit: int = Query(10, ge=1, le=50)) -> dict:
return nutrition_analytics_service.get_brand_category_rankings(limit)
@router.get("/analytics/distribution")
def analytics_distribution() -> dict:
return nutrition_analytics_service.get_distribution()
# ---------------------------------------------------------------------------
# Feature 10: Personalized Nutrition Recommendations
# ---------------------------------------------------------------------------
@router.get("/recommendations/{customer_id}", response_model=PersonalizedRecommendationOut)
def personalized_recommendations(customer_id: str, top_k: int = Query(8, ge=1, le=30)) -> PersonalizedRecommendationOut:
result = nutrition_recommendation.recommend_for_customer(customer_id, top_k)
return PersonalizedRecommendationOut(**result)
# ---------------------------------------------------------------------------
# GET Nutrition / GET Health Score (Features 1-6, 12) - DYNAMIC CATCH-ALL
# MUST stay below every static route above (see the note at the top of
# this file for why).
# ---------------------------------------------------------------------------
@router.get("/{brand}/{image_id}", response_model=FullNutritionOut)
def get_nutrition(brand: str, image_id: str) -> FullNutritionOut:
"""Full nutritional facts + insights for one product (Features 1-6).
Always returns 200 with `data_status="unavailable"` rather than 404
when a product exists but hasn't been enriched yet or has no
verified match - the frontend renders this as
"Nutrition data unavailable", per Feature 15."""
data = nutrition_db.get_full_nutrition(brand, image_id)
return FullNutritionOut(**data)
@router.get("/{brand}/{image_id}/health-score")
def get_health_score(brand: str, image_id: str) -> dict:
insights = nutrition_db.get_nutrition_insights(brand, image_id)
if not insights or insights.get("health_score") is None:
return {"brand": brand, "image_id": image_id, "data_status": "unavailable",
"nutrition_score": None, "health_score": None, "score_breakdown": None}
return {
"brand": brand, "image_id": image_id, "data_status": insights.get("data_status"),
"nutrition_score": insights.get("nutrition_score"), "health_score": insights.get("health_score"),
"score_breakdown": insights.get("score_breakdown"),
}
@router.get("/{brand}/{image_id}/insights", response_model=NutritionInsightsOut)
def get_insights(brand: str, image_id: str) -> NutritionInsightsOut:
insights = nutrition_db.get_nutrition_insights(brand, image_id) or {
"brand": brand, "image_id": image_id, "data_status": "unavailable",
"positive_insights": [], "nutritional_cautions": [], "diet_tags": [], "allergens": [],
}
return NutritionInsightsOut(**insights)
# ---------------------------------------------------------------------------
# GET Healthy Alternatives / GET Similar Nutritious Products (Features 7, 8)
# ---------------------------------------------------------------------------
@router.get("/{brand}/{image_id}/alternatives", response_model=List[HealthyAlternativeOut])
def get_alternatives(brand: str, image_id: str, top_k: int = Query(5, ge=1, le=20)) -> List[HealthyAlternativeOut]:
cached = nutrition_db.get_healthy_alternatives(brand, image_id, top_k)
if cached:
return [HealthyAlternativeOut(**c) for c in cached]
computed = nutrition_alternatives_service.find_alternatives(brand, image_id, top_k)
return [HealthyAlternativeOut(**c) for c in computed]
@router.get("/{brand}/{image_id}/similar", response_model=List[SimilarProductOut])
def get_similar(brand: str, image_id: str, top_k: int = Query(5, ge=1, le=20)) -> List[SimilarProductOut]:
cached = nutrition_db.get_similar_products(brand, image_id, top_k)
if cached:
return [SimilarProductOut(**c) for c in cached]
computed = nutrition_similarity.find_similar(brand, image_id, top_k)
return [SimilarProductOut(**c) for c in computed]

View File

@@ -0,0 +1,102 @@
from __future__ import annotations
import logging
from fastapi import APIRouter, HTTPException
from pydantic import BaseModel
from app.api.background import run_in_background
from app.api.nutrition_job_store import nutrition_job_store
from app.services import nutrition_enrichment_service
logger = logging.getLogger(__name__)
router = APIRouter(prefix="/admin/nutrition-intelligence", tags=["admin", "nutrition"])
class EnrichRequest(BaseModel):
skip_if_verified: bool = True
generate_narrative: bool = True
max_products: int | None = None
def _run_enrich_job(job_id: str, skip_if_verified: bool, generate_narrative: bool, max_products: int | None) -> None:
nutrition_job_store.update(job_id, status="running")
def progress_cb(done: int, total: int) -> None:
nutrition_job_store.update(job_id, processed=done, total=total)
try:
result = nutrition_enrichment_service.enrich_all_products(
skip_if_verified=skip_if_verified, generate_narrative=generate_narrative,
progress_cb=progress_cb, max_products=max_products,
)
nutrition_job_store.update(
job_id, status="done", detail="Enrichment complete",
result={
"total_products": result.total_products, "verified": result.verified,
"partial": result.partial, "unavailable": result.unavailable,
"duration_seconds": result.duration_seconds, "error_count": len(result.errors),
"errors": result.errors[:20],
},
)
except Exception as e: # noqa: BLE001
logger.exception("Nutrition enrichment job %s failed", job_id)
nutrition_job_store.update(job_id, status="failed", detail=str(e))
def _run_train_job(job_id: str) -> None:
nutrition_job_store.update(job_id, status="running")
try:
result = nutrition_enrichment_service.train_all_models()
nutrition_job_store.update(job_id, status="done", detail="Training complete", result=result)
except Exception as e: # noqa: BLE001
logger.exception("Nutrition training job %s failed", job_id)
nutrition_job_store.update(job_id, status="failed", detail=str(e))
@router.post("/enrich", status_code=202)
def enrich_nutrition(payload: EnrichRequest) -> dict:
"""Retrieves verified nutrition data for every product in the
catalog (Open Food Facts), computes transparent scores/insights, and
persists them. Safe to re-run - `skip_if_verified=true` (default)
only re-fetches products that don't already have verified data.
Equivalent to `python scripts/enrich_nutrition.py`."""
job = nutrition_job_store.create("enrich")
run_in_background(
lambda: _run_enrich_job(job.job_id, payload.skip_if_verified, payload.generate_narrative, payload.max_products),
name=f"nutrition-enrich-{job.job_id[:8]}",
)
return {"job_id": job.job_id, "status": job.status}
@router.post("/train", status_code=202)
def train_nutrition_models() -> dict:
"""Trains the nutrition-similarity (KNN/cosine) and nutrition-based
clustering (KMeans) models over the currently enriched catalog.
Run after `/enrich` completes. Equivalent to
`python scripts/train_nutrition_models.py`."""
job = nutrition_job_store.create("train")
run_in_background(
lambda: _run_train_job(job.job_id),
name=f"nutrition-train-{job.job_id[:8]}",
)
return {"job_id": job.job_id, "status": job.status}
@router.get("/jobs/{job_id}")
def get_job(job_id: str) -> dict:
job = nutrition_job_store.get(job_id)
if not job:
raise HTTPException(status_code=404, detail="Job not found")
return {
"job_id": job.job_id, "kind": job.kind, "status": job.status, "detail": job.detail,
"processed": job.processed, "total": job.total, "result": job.result,
}
@router.get("/status")
def enrichment_status() -> dict:
"""Quick counts for the admin UI: how many products have verified /
partial / unavailable nutrition data right now."""
from app.services import nutrition_db
return {"counts": nutrition_db.enrichment_status_counts()}

View File

@@ -0,0 +1,20 @@
from __future__ import annotations
import logging
from typing import List
from fastapi import APIRouter, Query
from app.api.store_schemas import RecommendationOut
from app.services import recommendation_service
logger = logging.getLogger(__name__)
router = APIRouter(tags=["recommendations"])
@router.get("/recommendations/{brand}/{image_id}", response_model=List[RecommendationOut])
def get_recommendations(brand: str, image_id: str, top_k: int = Query(default=5, le=20)) -> List[RecommendationOut]:
"""Feature 7: hybrid (embedding + TF-IDF + collaborative + popularity)
product recommendations, each with its own similarity_score."""
recs = recommendation_service.recommend_for_product(brand, image_id, top_k=top_k)
return [RecommendationOut(**r) for r in recs]

29
app/api/routers/search.py Normal file
View File

@@ -0,0 +1,29 @@
from __future__ import annotations
from typing import Optional
from fastapi import APIRouter, Query
from app.api.schemas import SearchOut, SourceProductOut
from app.services.rag_service import retrieve
router = APIRouter(tags=["search"])
@router.get("/search", response_model=SearchOut)
def semantic_search(
q: str = Query(..., min_length=1, max_length=500, description="Free-text search query"),
brand: Optional[str] = Query(None, description="Restrict search to a single brand"),
category: Optional[str] = Query(None, description="Restrict search to a category"),
top_k: int = Query(10, ge=1, le=50),
) -> SearchOut:
"""Pure vector similarity search over the catalog - no LLM call, just
pgvector ranking. This is what powers the instant search-as-you-type
grid in the React 'Search' tab. For a conversational, LLM-generated
answer use POST /api/chat instead."""
results = retrieve(q, brand=brand, top_k=top_k, category=category)
return SearchOut(
query=q,
brand=brand,
results=[SourceProductOut(**r.to_dict()) for r in results],
)

View File

@@ -0,0 +1,67 @@
from __future__ import annotations
import logging
from fastapi import APIRouter, HTTPException
from app.api.background import run_in_background
from app.api.store_job_store import store_job_store
from app.api.store_schemas import SeedRequest, TrainRequest
from app.services import ml_training_service, store_seed_service
logger = logging.getLogger(__name__)
router = APIRouter(prefix="/admin/store-intelligence", tags=["admin"])
def _run_seed_job(job_id: str, reset_orders: bool, days: int, seed: int) -> None:
store_job_store.update(job_id, "running")
try:
result = store_seed_service.run_seed(reset_orders=reset_orders, days=days, seed=seed)
store_job_store.update(job_id, "done", detail="Seed complete", result=result)
except Exception as e: # noqa: BLE001
logger.exception("Store-intelligence seed job %s failed", job_id)
store_job_store.update(job_id, "failed", detail=str(e))
def _run_train_job(job_id: str, models) -> None:
store_job_store.update(job_id, "running")
try:
result = ml_training_service.train_all(models=models)
store_job_store.update(job_id, "done", detail="Training complete", result=result)
except Exception as e: # noqa: BLE001
logger.exception("Store-intelligence train job %s failed", job_id)
store_job_store.update(job_id, "failed", detail=str(e))
@router.post("/seed", status_code=202)
def seed_store_intelligence(payload: SeedRequest) -> dict:
"""Provisions the 5 stores + simulates order history (Features 1, 2, 8).
Requires at least one brand already ingested via the existing
catalog pipeline. Equivalent to running
`python scripts/seed_store_intelligence.py`."""
job = store_job_store.create("seed")
run_in_background(
lambda: _run_seed_job(job.job_id, payload.reset_orders, payload.days, payload.seed),
name=f"store-seed-{job.job_id[:8]}",
)
return {"job_id": job.job_id, "status": job.status}
@router.post("/train", status_code=202)
def train_models(payload: TrainRequest) -> dict:
"""Trains every ML model (or a subset) against the seeded data.
Equivalent to running `python scripts/train_ml_models.py`."""
job = store_job_store.create("train")
run_in_background(
lambda: _run_train_job(job.job_id, payload.models),
name=f"store-train-{job.job_id[:8]}",
)
return {"job_id": job.job_id, "status": job.status}
@router.get("/jobs/{job_id}")
def get_job(job_id: str) -> dict:
job = store_job_store.get(job_id)
if not job:
raise HTTPException(status_code=404, detail="Job not found")
return {"job_id": job.job_id, "kind": job.kind, "status": job.status, "detail": job.detail, "result": job.result}

77
app/api/routers/stores.py Normal file
View File

@@ -0,0 +1,77 @@
from __future__ import annotations
import logging
from typing import List, Optional
from fastapi import APIRouter, HTTPException, Query
from app.api.store_schemas import ProductStorePriceOut, StoreOut, StoreProductOut
from app.services import store_db
logger = logging.getLogger(__name__)
router = APIRouter(tags=["stores"])
@router.get("/stores", response_model=List[StoreOut])
def list_stores() -> List[StoreOut]:
return [StoreOut(**s) for s in store_db.list_stores()]
@router.get("/stores/{store_id}", response_model=StoreOut)
def get_store(store_id: str) -> StoreOut:
store = store_db.get_store(store_id)
if not store:
raise HTTPException(status_code=404, detail="Store not found")
return StoreOut(**store)
@router.get("/stores/{store_id}/products", response_model=List[StoreProductOut])
def get_store_products(
store_id: str,
category: Optional[str] = Query(default=None),
in_stock_only: bool = Query(default=False),
limit: int = Query(default=50, le=500),
offset: int = Query(default=0, ge=0),
) -> List[StoreProductOut]:
if not store_db.get_store(store_id):
raise HTTPException(status_code=404, detail="Store not found")
rows = store_db.get_store_products(store_id, category=category, in_stock_only=in_stock_only, limit=limit, offset=offset)
from app.services import vector_store
for row in rows:
prod = vector_store.get_product_by_image_id(row["brand"], row["image_id"])
if prod:
row["image_url"] = prod.get("image_url")
row["image_urls"] = prod.get("image_urls")
row["fssai_license"] = prod.get("fssai_license")
return [_to_store_product_out(r) for r in rows]
@router.get("/products/{brand}/{image_id}/stores", response_model=List[ProductStorePriceOut])
def get_product_across_stores(brand: str, image_id: str) -> List[ProductStorePriceOut]:
"""Feature 1: compare one product's price/stock across every store
that carries it - the direct answer to "Tata Tea Gold 100g: Store-A
₹79, Store-B ₹82, ..." from the spec."""
rows = store_db.get_product_across_stores(brand, image_id)
if not rows:
raise HTTPException(status_code=404, detail="Product not found in any store")
return [ProductStorePriceOut(**r) for r in rows]
def _to_store_product_out(row: dict) -> StoreProductOut:
from app.intelligence.analytics import classify_stock_status
stock_status = classify_stock_status(row["available_stock"], row["reorder_level"], row["safety_stock"])
margin = round(row["selling_price"] - row["cost_price"], 2)
gp_pct = round((margin / row["selling_price"]) * 100, 2) if row["selling_price"] else 0.0
markup_pct = round((margin / row["cost_price"]) * 100, 2) if row["cost_price"] else 0.0
return StoreProductOut(
store_id=row["store_id"], brand=row["brand"], image_id=row["image_id"], title=row.get("title"),
category=row.get("category"), available_stock=row["available_stock"], reserved_stock=row["reserved_stock"],
reorder_level=row["reorder_level"], safety_stock=row["safety_stock"], stock_status=stock_status,
mrp=float(row["mrp"]), cost_price=float(row["cost_price"]), selling_price=float(row["selling_price"]),
profit_margin=margin, gross_profit_pct=gp_pct, markup_pct=markup_pct,
image_url=row.get("image_url"), image_urls=row.get("image_urls"),
fssai_license=row.get("fssai_license")
)

90
app/api/routers/system.py Normal file
View File

@@ -0,0 +1,90 @@
from __future__ import annotations
import logging
import os
import subprocess
import threading
from pathlib import Path
from typing import Any, Dict
from fastapi import APIRouter, BackgroundTasks
from pydantic import BaseModel
from app.services.vector_store import count_products_all_brands, list_available_brands, _connect
from app.services.store_db import list_stores
from app.services.ollama_service import _ensure_client
logger = logging.getLogger(__name__)
router = APIRouter(tags=["system"])
BASE_DIR = Path(__file__).resolve().parents[3]
FRONTEND_DIST = BASE_DIR.parent / "frontend" / "dist"
class SystemStatusOut(BaseModel):
status: str
database_connected: bool
total_products: int
available_brands: list[str]
total_stores: int
ollama_connected: bool
frontend_dist_exists: bool
def _run_background_auto_seed():
"""Background task to run initial seeding asynchronously if DB is empty."""
try:
if count_products_all_brands() == 0:
logger.info("⚡ Background Auto-Init: Database empty. Running initial sample seed...")
cmd_seed = [os.sys.executable, str(BASE_DIR / "scripts" / "seed_sample_data.py"), "--skip-if-seeded"]
subprocess.run(cmd_seed, check=False)
logger.info("⚡ Background Auto-Init: Provisioning store intelligence...")
cmd_store = [os.sys.executable, str(BASE_DIR / "scripts" / "seed_store_intelligence.py"), "--skip-if-seeded"]
subprocess.run(cmd_store, check=False)
logger.info("✅ Background Auto-Init complete!")
except Exception as e:
logger.error("Background Auto-Init error: %s", e)
@router.get("/system/status", response_model=SystemStatusOut)
def get_system_status() -> SystemStatusOut:
"""Return unified status of database, vector store, stores, and frontend build."""
db_connected = False
products_count = 0
brands = []
stores_count = 0
try:
conn = _connect()
if conn:
db_connected = True
conn.close()
products_count = count_products_all_brands()
brands = list_available_brands()
stores_count = len(list_stores())
except Exception:
pass
ollama_ok = _ensure_client()
dist_ok = FRONTEND_DIST.exists() and (FRONTEND_DIST / "index.html").exists()
return SystemStatusOut(
status="ok" if db_connected else "degraded",
database_connected=db_connected,
total_products=products_count,
available_brands=brands,
total_stores=stores_count,
ollama_connected=ollama_ok,
frontend_dist_exists=dist_ok,
)
@router.post("/system/init")
def initialize_system(background_tasks: BackgroundTasks) -> Dict[str, Any]:
"""Trigger background auto-initialization of sample catalog and store data."""
background_tasks.add_task(_run_background_auto_seed)
return {
"status": "started",
"message": "Background initialization triggered. Check /api/system/status for progress.",
}

View File

@@ -0,0 +1,31 @@
from __future__ import annotations
import logging
from typing import List, Literal, Optional
from fastapi import APIRouter, HTTPException, Query
from app.api.store_schemas import TrendingItemOut
from app.services import store_db, trending_service
logger = logging.getLogger(__name__)
router = APIRouter(tags=["trending"])
@router.get("/trending", response_model=List[TrendingItemOut])
def get_trending(
window: Literal["today", "weekly", "monthly"] = Query(default="weekly"),
scope: Literal["overall", "category", "store"] = Query(default="overall"),
scope_value: Optional[str] = Query(default=None, description="Category name (scope=category) or store_id (scope=store)"),
top_k: int = Query(default=10, le=50),
) -> List[TrendingItemOut]:
"""Feature 6: Today's / Weekly / Monthly trending, overall,
category-wise, or store-wise. Backed by the trained trending
regressor's predicted scores - never a hardcoded list."""
if scope in ("category", "store") and not scope_value:
raise HTTPException(status_code=422, detail=f"scope_value is required when scope={scope}")
if scope == "store" and not store_db.get_store(scope_value):
raise HTTPException(status_code=404, detail="Store not found")
items = trending_service.get_trending(window, scope, scope_value, top_k)
return [TrendingItemOut(**i) for i in items]

447
app/api/routers/upload.py Normal file
View File

@@ -0,0 +1,447 @@
from __future__ import annotations
import io
import uuid
import logging
import pandas as pd
from typing import Any, Dict, List, Optional
from datetime import datetime
from fastapi import APIRouter, File, HTTPException, UploadFile, Response
from fastapi.responses import PlainTextResponse
from app.services.vector_store import _connect
from app.services import store_db, nutrition_db
logger = logging.getLogger(__name__)
router = APIRouter(prefix="/upload", tags=["upload"])
def _normalize_col(col: str) -> str:
"""Normalize dataframe column names (lower, strip, replace spaces/hyphens with underscore)."""
return str(col).strip().lower().replace(' ', '_').replace('-', '_')
def read_df_from_upload(filename: str, contents: bytes) -> pd.DataFrame:
"""Parse CSV or Excel (xlsx/xls) upload file into a pandas DataFrame."""
fn_lower = filename.lower()
if fn_lower.endswith('.xlsx') or fn_lower.endswith('.xls'):
df = pd.read_excel(io.BytesIO(contents))
elif fn_lower.endswith('.tsv'):
df = pd.read_csv(io.BytesIO(contents), sep='\t')
else:
try:
df = pd.read_csv(io.BytesIO(contents))
except Exception:
df = pd.read_csv(io.BytesIO(contents), sep=None, engine='python')
# Rename columns to normalized format
df.columns = [_normalize_col(c) for c in df.columns]
return df
def _get_str(row: dict, keys: List[str], default: str = "") -> str:
for k in keys:
if k in row and pd.notna(row[k]):
val = str(row[k]).strip()
if val:
return val
return default
def _get_float(row: dict, keys: List[str], default: float = 0.0) -> float:
for k in keys:
if k in row and pd.notna(row[k]):
try:
return float(row[k])
except (ValueError, TypeError):
pass
return default
def _get_int(row: dict, keys: List[str], default: int = 0) -> int:
for k in keys:
if k in row and pd.notna(row[k]):
try:
return int(float(row[k]))
except (ValueError, TypeError):
pass
return default
# ---------------------------------------------------------------------------
# Stores Inventory Excel / CSV Upload
# ---------------------------------------------------------------------------
@router.post("/stores")
@router.post("/stores/upload")
async def upload_stores_file(file: UploadFile = File(...)) -> Dict[str, Any]:
if not file.filename:
raise HTTPException(status_code=400, detail="No file uploaded")
contents = await file.read()
try:
df = read_df_from_upload(file.filename, contents)
except Exception as e:
raise HTTPException(status_code=400, detail=f"Could not parse Excel/CSV file: {e}")
if df.empty:
raise HTTPException(status_code=400, detail="Uploaded file contains no data rows")
conn = _connect()
if not conn:
raise HTTPException(status_code=500, detail="Database connection failed")
imported_count = 0
stores_created = set()
try:
with conn, conn.cursor() as cur:
# Ensure tables exist
store_db.ensure_store_intelligence_schema()
for _, r in df.iterrows():
row = r.to_dict()
store_id = _get_str(row, ['store_id', 'store'], 'store_mumbai_1')
brand = _get_str(row, ['brand', 'brand_name'], 'amul').lower()
product_name = _get_str(row, ['product_name', 'title', 'name', 'item'], 'Product Item')
image_id = _get_str(row, ['image_id', 'sku', 'product_sku', 'item_id'], '')
if not image_id:
image_id = f"{brand}_{product_name.lower().replace(' ', '_')}"
category = _get_str(row, ['category', 'cat'], 'Dairy')
avail_stock = _get_int(row, ['available_stock', 'stock', 'qty', 'quantity'], 50)
reserved_stock = _get_int(row, ['reserved_stock', 'reserved'], 0)
reorder_lvl = _get_int(row, ['reorder_level', 'reorder'], 15)
safety_stk = _get_int(row, ['safety_stock', 'safety'], 10)
mrp = _get_float(row, ['mrp', 'price'], 100.0)
cost_price = _get_float(row, ['cost_price', 'cost'], 70.0)
selling_price = _get_float(row, ['selling_price', 'sell_price'], mrp * 0.9 if mrp else 90.0)
# 1. Ensure store exists
cur.execute(
"""
INSERT INTO stores (store_id, store_name, city, tier, footfall_index)
VALUES (%s, %s, %s, %s, %s)
ON CONFLICT (store_id) DO NOTHING
""",
(store_id, store_id.replace('_', ' ').title(), 'Mumbai', 'standard', 25.0)
)
stores_created.add(store_id)
# 2. Upsert store_inventory
cur.execute(
"""
INSERT INTO store_inventory
(store_id, brand, image_id, title, category, available_stock, reserved_stock, reorder_level, safety_stock)
VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s)
ON CONFLICT (store_id, brand, image_id) DO UPDATE SET
title = EXCLUDED.title,
category = EXCLUDED.category,
available_stock = EXCLUDED.available_stock,
reserved_stock = EXCLUDED.reserved_stock,
reorder_level = EXCLUDED.reorder_level,
safety_stock = EXCLUDED.safety_stock,
updated_at = CURRENT_TIMESTAMP
""",
(store_id, brand, image_id, product_name, category, avail_stock, reserved_stock, reorder_lvl, safety_stk)
)
# 3. Upsert store_prices
cur.execute(
"""
INSERT INTO store_prices (store_id, brand, image_id, mrp, cost_price, selling_price)
VALUES (%s, %s, %s, %s, %s, %s)
ON CONFLICT (store_id, brand, image_id) DO UPDATE SET
mrp = EXCLUDED.mrp,
cost_price = EXCLUDED.cost_price,
selling_price = EXCLUDED.selling_price,
updated_at = CURRENT_TIMESTAMP
""",
(store_id, brand, image_id, mrp, cost_price, selling_price)
)
imported_count += 1
except Exception as e:
logger.error("Stores upload failed: %s", e)
raise HTTPException(status_code=500, detail=f"Database import failed: {e}")
finally:
conn.close()
return {
"status": "success",
"filename": file.filename,
"rows_total": len(df),
"rows_imported": imported_count,
"stores_affected": list(stores_created),
"message": f"Successfully imported {imported_count} store inventory items across {len(stores_created)} store(s)."
}
# ---------------------------------------------------------------------------
# Sales / Analytics Excel / CSV Upload
# ---------------------------------------------------------------------------
@router.post("/analytics")
@router.post("/analytics/upload")
async def upload_analytics_file(file: UploadFile = File(...)) -> Dict[str, Any]:
if not file.filename:
raise HTTPException(status_code=400, detail="No file uploaded")
contents = await file.read()
try:
df = read_df_from_upload(file.filename, contents)
except Exception as e:
raise HTTPException(status_code=400, detail=f"Could not parse Excel/CSV file: {e}")
if df.empty:
raise HTTPException(status_code=400, detail="Uploaded file contains no data rows")
conn = _connect()
if not conn:
raise HTTPException(status_code=500, detail="Database connection failed")
imported_orders = 0
total_revenue = 0.0
try:
with conn, conn.cursor() as cur:
store_db.ensure_store_intelligence_schema()
for _, r in df.iterrows():
row = r.to_dict()
store_id = _get_str(row, ['store_id', 'store'], 'store_mumbai_1')
brand = _get_str(row, ['brand', 'brand_name'], 'amul').lower()
image_id = _get_str(row, ['image_id', 'sku', 'product_sku'], '')
product_name = _get_str(row, ['product_name', 'title', 'item'], 'Analytics Item')
if not image_id:
image_id = f"{brand}_{product_name.lower().replace(' ', '_')}"
order_id = _get_str(row, ['order_id', 'transaction_id'], f"ord_up_{uuid.uuid4().hex[:8]}")
customer_id = _get_str(row, ['customer_id', 'user_id', 'customer'], 'cust_imported')
raw_date = _get_str(row, ['order_date', 'date', 'timestamp'], '')
order_date = datetime.now()
if raw_date:
try:
order_date = pd.to_datetime(raw_date).to_pydatetime()
except Exception:
pass
qty = _get_int(row, ['quantity', 'units_sold', 'qty', 'count'], 1)
unit_price = _get_float(row, ['unit_price', 'selling_price', 'price'], 100.0)
tot_price = _get_float(row, ['total_price', 'revenue', 'total'], qty * unit_price)
# Ensure store exists
cur.execute(
"INSERT INTO stores (store_id, store_name, city, tier, footfall_index) VALUES (%s, %s, %s, %s, %s) ON CONFLICT (store_id) DO NOTHING",
(store_id, store_id.replace('_', ' ').title(), 'Mumbai', 'standard', 25.0)
)
# Insert order header
cur.execute(
"""
INSERT INTO orders (order_id, customer_id, store_id, order_date, payment_method, order_value, delivery_status)
VALUES (%s, %s, %s, %s, %s, %s, %s)
ON CONFLICT (order_id) DO UPDATE SET order_value = EXCLUDED.order_value
""",
(order_id, customer_id, store_id, order_date, 'upi', tot_price, 'delivered')
)
# Insert order item
cur.execute(
"""
INSERT INTO order_items (order_id, brand, image_id, quantity, unit_price, total_price)
VALUES (%s, %s, %s, %s, %s, %s)
""",
(order_id, brand, image_id, qty, unit_price, tot_price)
)
imported_orders += 1
total_revenue += tot_price
except Exception as e:
logger.error("Analytics upload failed: %s", e)
raise HTTPException(status_code=500, detail=f"Database import failed: {e}")
finally:
conn.close()
return {
"status": "success",
"filename": file.filename,
"rows_total": len(df),
"rows_imported": imported_orders,
"total_revenue": round(total_revenue, 2),
"message": f"Successfully imported {imported_orders} sales transactions (Total Revenue: ₹{total_revenue:,.2f})."
}
# ---------------------------------------------------------------------------
# Nutrition Intelligence Excel / CSV Upload
# ---------------------------------------------------------------------------
@router.post("/nutrition")
@router.post("/nutrition/upload")
async def upload_nutrition_file(file: UploadFile = File(...)) -> Dict[str, Any]:
if not file.filename:
raise HTTPException(status_code=400, detail="No file uploaded")
contents = await file.read()
try:
df = read_df_from_upload(file.filename, contents)
except Exception as e:
raise HTTPException(status_code=400, detail=f"Could not parse Excel/CSV file: {e}")
if df.empty:
raise HTTPException(status_code=400, detail="Uploaded file contains no data rows")
conn = _connect()
if not conn:
raise HTTPException(status_code=500, detail="Database connection failed")
imported_count = 0
try:
with conn, conn.cursor() as cur:
nutrition_db.ensure_nutrition_schema()
for _, r in df.iterrows():
row = r.to_dict()
brand = _get_str(row, ['brand', 'brand_name'], 'amul').lower()
product_name = _get_str(row, ['product_name', 'title', 'item', 'name'], 'Nutrition Item')
image_id = _get_str(row, ['image_id', 'sku', 'id'], '')
if not image_id:
image_id = f"{brand}_{product_name.lower().replace(' ', '_')}"
category = _get_str(row, ['category', 'cat'], 'Food')
calories = _get_float(row, ['calories', 'calories_kcal', 'energy'], 150.0)
protein = _get_float(row, ['protein', 'protein_g'], 5.0)
carbs = _get_float(row, ['carbohydrates', 'carbs', 'carbohydrates_g'], 20.0)
sugar = _get_float(row, ['sugar', 'total_sugar_g', 'sugars'], 4.0)
fiber = _get_float(row, ['fiber', 'dietary_fiber_g'], 2.0)
fat = _get_float(row, ['fat', 'total_fat_g'], 6.0)
sodium = _get_float(row, ['sodium', 'sodium_mg'], 120.0)
calcium = _get_float(row, ['calcium', 'calcium_mg'], 80.0)
iron = _get_float(row, ['iron', 'iron_mg'], 1.5)
vitamin_c = _get_float(row, ['vitamin_c', 'vitamin_c_mg'], 5.0)
health_score = _get_float(row, ['health_score', 'nutrition_score', 'score'], 78.0)
diet_tags_raw = _get_str(row, ['diet_tags', 'tags', 'diet'], 'High Protein, Gluten Free')
allergens_raw = _get_str(row, ['allergens', 'allergen'], 'None')
diet_tags = [t.strip() for t in diet_tags_raw.split(',') if t.strip()]
allergens = [a.strip() for a in allergens_raw.split(',') if a.strip()]
# 1. Upsert nutrition_facts
cur.execute(
"""
INSERT INTO nutrition_facts
(brand, image_id, product_name, category, data_status, data_source,
calories_kcal, protein_g, carbohydrates_g, total_sugar_g, dietary_fiber_g,
total_fat_g, sodium_mg, calcium_mg, iron_mg, vitamin_c_mg)
VALUES (%s, %s, %s, %s, 'verified', 'excel_upload', %s, %s, %s, %s, %s, %s, %s, %s, %s, %s)
ON CONFLICT (brand, image_id) DO UPDATE SET
product_name = EXCLUDED.product_name,
category = EXCLUDED.category,
data_status = 'verified',
calories_kcal = EXCLUDED.calories_kcal,
protein_g = EXCLUDED.protein_g,
carbohydrates_g = EXCLUDED.carbohydrates_g,
total_sugar_g = EXCLUDED.total_sugar_g,
dietary_fiber_g = EXCLUDED.dietary_fiber_g,
total_fat_g = EXCLUDED.total_fat_g,
sodium_mg = EXCLUDED.sodium_mg,
calcium_mg = EXCLUDED.calcium_mg,
iron_mg = EXCLUDED.iron_mg,
vitamin_c_mg = EXCLUDED.vitamin_c_mg
""",
(brand, image_id, product_name, category, calories, protein, carbs, sugar, fiber, fat, sodium, calcium, iron, vitamin_c)
)
# 2. Upsert nutrition_insights
insights_json = json.dumps({
"brand": brand,
"image_id": image_id,
"data_status": "verified",
"nutrition_score": health_score,
"health_score": health_score,
"positive_insights": [f"Contains {protein}g protein per 100g", f"Provides {fiber}g dietary fiber"],
"nutritional_cautions": [f"{sugar}g sugar per 100g"],
"diet_tags": diet_tags,
"allergens": allergens
})
cur.execute(
"""
INSERT INTO nutrition_insights
(brand, image_id, data_status, nutrition_score, health_score, score_breakdown,
positive_insights, nutritional_cautions, diet_tags, allergens)
VALUES (%s, %s, 'verified', %s, %s, %s, %s, %s, %s, %s)
ON CONFLICT (brand, image_id) DO UPDATE SET
data_status = 'verified',
nutrition_score = EXCLUDED.nutrition_score,
health_score = EXCLUDED.health_score,
positive_insights = EXCLUDED.positive_insights,
nutritional_cautions = EXCLUDED.nutritional_cautions,
diet_tags = EXCLUDED.diet_tags,
allergens = EXCLUDED.allergens
""",
(brand, image_id, health_score, health_score, json.dumps({"protein": 85, "fiber": 80}),
[f"Contains {protein}g protein per 100g"], [f"{sugar}g sugar per 100g"], diet_tags, allergens)
)
imported_count += 1
except Exception as e:
logger.error("Nutrition upload failed: %s", e)
raise HTTPException(status_code=500, detail=f"Database import failed: {e}")
finally:
conn.close()
return {
"status": "success",
"filename": file.filename,
"rows_total": len(df),
"rows_imported": imported_count,
"message": f"Successfully imported {imported_count} nutritional intelligence items."
}
# ---------------------------------------------------------------------------
# Template Downloads
# ---------------------------------------------------------------------------
@router.get("/template/{tab_type}")
def get_sample_template(tab_type: str) -> Response:
tab_type = tab_type.lower()
if tab_type == 'stores':
content = (
"store_id,brand,image_id,product_name,category,available_stock,reserved_stock,mrp,cost_price,selling_price,reorder_level,safety_stock\n"
"store_mumbai_1,amul,amul_amul_butter_500ml,Amul Butter 500ml,Dairy,120,5,250.00,200.00,235.00,20,10\n"
"store_mumbai_1,amul,amul_amul_ghee_1l,Amul Ghee 1L,Dairy,85,2,650.00,520.00,610.00,15,5\n"
"store_delhi_2,nestle,nestle_everyday_1kg,Everyday Milk Powder 1kg,Dairy,45,0,420.00,340.00,399.00,10,5\n"
)
filename = "sample_stores_inventory_template.csv"
elif tab_type == 'analytics':
content = (
"order_id,store_id,brand,image_id,product_name,order_date,customer_id,quantity,unit_price,total_price\n"
"ORD_9001,store_mumbai_1,amul,amul_amul_butter_500ml,Amul Butter 500ml,2026-08-01 10:30:00,cust_101,2,235.00,470.00\n"
"ORD_9002,store_mumbai_1,amul,amul_amul_ghee_1l,Amul Ghee 1L,2026-08-01 11:15:00,cust_102,1,610.00,610.00\n"
"ORD_9003,store_delhi_2,nestle,nestle_everyday_1kg,Everyday Milk Powder 1kg,2026-08-02 14:20:00,cust_103,3,399.00,1197.00\n"
)
filename = "sample_analytics_sales_template.csv"
elif tab_type == 'nutrition':
content = (
"brand,image_id,product_name,category,calories_kcal,protein_g,carbohydrates_g,total_sugar_g,dietary_fiber_g,total_fat_g,sodium_mg,health_score,diet_tags,allergens\n"
"amul,amul_amul_butter_500ml,Amul Butter 500ml,Dairy,717,0.8,0.1,0.0,0.0,81.0,650,75,Vegetarian,Dairy\n"
"amul,amul_amul_ghee_1l,Amul Ghee 1L,Dairy,898,0.0,0.0,0.0,0.0,99.8,0,82,Vegetarian,Keto Friendly\n"
"nestle,nestle_everyday_1kg,Everyday Milk Powder 1kg,Dairy,496,25.5,38.0,38.0,0.0,27.0,350,88,High Protein,Dairy\n"
)
filename = "sample_nutrition_intelligence_template.csv"
else:
raise HTTPException(status_code=400, detail=f"Unknown template type '{tab_type}'. Use stores, analytics, or nutrition.")
return PlainTextResponse(
content=content,
media_type="text/csv",
headers={"Content-Disposition": f"attachment; filename={filename}"}
)

View File

@@ -0,0 +1,349 @@
import io
import json
import logging
import re
from pathlib import Path
from typing import Any, Dict, List, Optional
import pandas as pd
from pydantic import BaseModel, Field
from fastapi import APIRouter, HTTPException, File, UploadFile
from app.services.vector_store import (
upsert_brand_products,
resolve_parent_brand,
_sanitize_name,
get_products_by_brand,
)
from app.services.embeddings_service import embed_texts
from app.services.s3_service import s3_service
logger = logging.getLogger(__name__)
router = APIRouter(prefix="/user/products", tags=["user_products"])
SEED_DIR = Path(__file__).resolve().parents[3] / "data" / "seed_catalogs"
class AddProductRequest(BaseModel):
brand: str = Field(..., description="Brand name, e.g. Lion Dates")
product_name: str = Field(..., description="Product name, e.g. Lion Dates 450g")
title: Optional[str] = None
category: Optional[str] = None
description: Optional[str] = None
price_range: Optional[str] = None
size_variants: List[str] = Field(default_factory=list)
providers: List[str] = Field(default_factory=list)
highlights: List[str] = Field(default_factory=list)
nutrients: List[str] = Field(default_factory=list)
fssai_license: Optional[str] = None
product_sku: Optional[str] = None
sku_source: Optional[str] = None
hsn_code: Optional[str] = None
final_selling_price: Optional[float] = None
selling_price: Optional[float] = None
barcode: Optional[str] = None
barcode_type: Optional[str] = None
image_url: Optional[str] = None
image_urls: List[str] = Field(default_factory=list)
class BatchAddProductsRequest(BaseModel):
products: List[AddProductRequest]
def _slugify(text: str) -> str:
return re.sub(r'[^a-z0-9]+', '_', text.lower()).strip('_')
def _enrich_and_save_product(req: AddProductRequest) -> Dict[str, Any]:
brand = req.brand.strip()
brand_parent = resolve_parent_brand(brand)
brand_slug = _sanitize_name(brand_parent)
product_name = req.product_name.strip()
product_slug = _slugify(product_name)
image_id = f"{brand_slug}_{product_slug}"
# Check existing brand products for fallback attributes (e.g. fssai_license, category, provider_examples)
existing_db = get_products_by_brand(brand_parent)
sample_existing = existing_db[0] if existing_db else {}
# Category fallback
category = req.category or sample_existing.get("category") or "Health Foods"
# FSSAI License fallback
fssai_license = req.fssai_license or sample_existing.get("fssai_license") or "10012042000244"
# Description fallback
description = req.description or (
f"Introducing {product_name} from the trusted {brand_parent} brand. "
f"A premium quality product offering superior taste, authentic ingredients, and reliable value. "
f"Backed by {brand_parent}'s reputation for quality and consistency."
)
# Size variants fallback
size_variants = req.size_variants
if not size_variants:
match = re.search(r'\d+\s*(?:g|kg|ml|l|pack)\b', product_name, re.I)
if match:
size_variants = [match.group(0)]
else:
size_variants = [sample_existing.get("size_variants", ["Default"])[0]] if sample_existing.get("size_variants") else ["Standard"]
# Price range fallback
price_range = req.price_range
if not price_range:
if req.final_selling_price:
price_range = f"₹{req.final_selling_price}"
elif sample_existing.get("price_range"):
price_range = sample_existing.get("price_range")
else:
price_range = "₹100-250"
providers = req.providers or list(sample_existing.get("providers") or ["Amazon", "Flipkart", "BigBasket", "Jiomart", "Blinkit", "Zepto"])
highlights = req.highlights or list(sample_existing.get("highlights") or ["100% Quality Assurance", "Authentic Brand Product"])
nutrients = req.nutrients or list(sample_existing.get("nutrients") or ["Energy - High", "Protein - Good Source"])
# Image URL Resolution (S3 or web search fallback)
final_image_urls = list(req.image_urls)
if req.image_url and req.image_url not in final_image_urls:
final_image_urls.insert(0, req.image_url)
if not final_image_urls:
# 1. Try S3 service if enabled
if s3_service.enabled:
s3_urls = s3_service.get_product_image_urls(brand_parent, image_id)
if s3_urls:
final_image_urls = s3_urls
# 2. Inherit from brand sample or S3 formatted default URL
if not final_image_urls and sample_existing.get("image_urls"):
final_image_urls = list(sample_existing.get("image_urls"))
# 3. Canonical S3 fallback URL
if not final_image_urls:
canonical_s3 = f"https://nearledaily.s3.ap-south-1.amazonaws.com/daily/brands/{brand_slug}/{image_id}/image_000.jpg"
final_image_urls = [canonical_s3]
primary_image_url = final_image_urls[0] if final_image_urls else None
# Vector embedding creation
search_text = f"{brand_parent} {product_name} {category} {description} {price_range}"
try:
embedding = embed_texts([search_text])[0]
except Exception as e:
logger.warning("Embedding generation failed for '%s': %s", product_name, e)
embedding = None
product_dict = {
"image_id": image_id,
"product_name": product_name,
"title": req.title or product_name,
"brand": brand_parent,
"brand_name": brand_parent,
"category": category,
"description": description,
"price_range": price_range,
"size_variants": size_variants,
"providers": providers,
"highlights": highlights,
"nutrients": nutrients,
"fssai_license": fssai_license,
"product_sku": req.product_sku or f"{brand_slug.upper()[:4]}-{product_slug.upper()[:6]}-001",
"sku_source": req.sku_source or "User Upload",
"hsn_code": req.hsn_code,
"final_selling_price": req.final_selling_price or req.selling_price,
"selling_price": req.selling_price or req.final_selling_price,
"barcode": req.barcode,
"barcode_type": req.barcode_type or ("GTIN-13" if req.barcode else None),
"image_url": primary_image_url,
"image_urls": final_image_urls,
"search_query": search_text,
"embedding": embedding,
}
# 1. Update PostgreSQL Database Table
upsert_brand_products(brand_parent, [product_dict])
logger.info("✅ Upserted '%s' into PostgreSQL table for brand '%s'", product_name, brand_parent)
# 2. Update JSON Seed File
_update_json_catalog_file(brand_parent, product_dict)
return product_dict
def _update_json_catalog_file(brand: str, product_dict: Dict[str, Any]) -> None:
SEED_DIR.mkdir(parents=True, exist_ok=True)
# Determine seed file name (e.g. brand_catalog_lion_dates.json)
brand_slug = _sanitize_name(resolve_parent_brand(brand))
file_path = SEED_DIR / f"brand_catalog_{brand_slug}.json"
# Strip embedding before saving to JSON file for clean JSON size
clean_dict = {k: v for k, v in product_dict.items() if k != "embedding"}
if file_path.exists():
try:
data = json.loads(file_path.read_text(encoding="utf-8-sig"))
except Exception as e:
logger.warning("Could not read existing catalog JSON %s: %s", file_path.name, e)
data = {"brand": brand, "products": []}
else:
data = {
"brand": brand.lower(),
"search_query": f"{brand} products catalog",
"generation_timestamp": str(Path(__file__).resolve()),
"total_products": 0,
"total_images": 0,
"products": [],
}
products_list = data.get("products", [])
# Replace existing or append new product
updated = False
for i, p in enumerate(products_list):
if p.get("image_id") == clean_dict["image_id"] or p.get("product_name") == clean_dict["product_name"]:
products_list[i] = clean_dict
updated = True
break
if not updated:
products_list.append(clean_dict)
data["products"] = products_list
data["total_products"] = len(products_list)
data["total_images"] = sum(len(p.get("image_urls") or []) for p in products_list)
file_path.write_text(json.dumps(data, indent=2, ensure_ascii=False), encoding="utf-8")
logger.info("✅ Updated JSON seed file '%s' (total products: %d)", file_path.name, data["total_products"])
@router.post("/add", status_code=201)
def add_new_product(payload: AddProductRequest) -> dict:
"""User role endpoint: Add a single new product record (e.g. Lion Dates 450g).
Automatically enriches details, fetches images, updates PostgreSQL DB,
and updates JSON seed catalog files."""
try:
res = _enrich_and_save_product(payload)
return {
"status": "success",
"message": f"Successfully added '{payload.product_name}' under brand '{payload.brand}' to database and JSON catalog.",
"product": {k: v for k, v in res.items() if k != "embedding"},
}
except Exception as e:
logger.exception("Failed to add product '%s'", payload.product_name)
raise HTTPException(status_code=500, detail=f"Failed to add product: {e}")
@router.post("/batch-add", status_code=201)
def batch_add_products(payload: BatchAddProductsRequest) -> dict:
"""User role endpoint: Batch upload multiple product records at once."""
added = []
errors = []
for req in payload.products:
try:
res = _enrich_and_save_product(req)
added.append({k: v for k, v in res.items() if k != "embedding"})
except Exception as e:
errors.append({"product_name": req.product_name, "error": str(e)})
return {
"status": "success",
"added_count": len(added),
"error_count": len(errors),
"added_products": added,
"errors": errors,
}
@router.post("/upload-file", status_code=201)
async def upload_products_file(file: UploadFile = File(...)) -> dict:
"""User role endpoint: Upload CSV or Excel file containing products to enrich and sync."""
filename = file.filename or ""
content = await file.read()
try:
if filename.endswith(".csv"):
df = pd.read_csv(io.BytesIO(content))
elif filename.endswith((".xlsx", ".xls")):
df = pd.read_excel(io.BytesIO(content))
else:
raise HTTPException(status_code=400, detail="Unsupported file format. Please upload a .csv or .xlsx file.")
except Exception as e:
raise HTTPException(status_code=400, detail=f"Failed to parse file '{filename}': {e}")
# Standardize column headers
col_map = {}
for col in df.columns:
c_clean = str(col).strip().lower()
if "brand" in c_clean:
col_map[col] = "brand"
elif "product" in c_clean or "variant" in c_clean or "name" in c_clean:
col_map[col] = "product_name"
elif "category" in c_clean:
col_map[col] = "category"
elif "range" in c_clean:
col_map[col] = "price_range"
elif "price" in c_clean or "selling" in c_clean or "cost" in c_clean:
col_map[col] = "final_selling_price"
elif "barcode" in c_clean or "gtin" in c_clean or "ean" in c_clean:
col_map[col] = "barcode"
elif "hsn" in c_clean:
col_map[col] = "hsn_code"
elif "description" in c_clean:
col_map[col] = "description"
elif "image" in c_clean or "url" in c_clean:
col_map[col] = "image_url"
df = df.rename(columns=col_map)
if "brand" not in df.columns or "product_name" not in df.columns:
raise HTTPException(
status_code=400,
detail="File must contain at least 'Brand Name' and 'Product Name' columns.",
)
added = []
errors = []
for idx, row in df.iterrows():
b_val = str(row.get("brand") or "").strip()
p_val = str(row.get("product_name") or "").strip()
if not b_val or not p_val or b_val.lower() == "nan" or p_val.lower() == "nan":
continue
try:
fps_raw = row.get("final_selling_price")
fps = None
if pd.notna(fps_raw):
try:
fps = float(fps_raw)
except Exception:
pass
req = AddProductRequest(
brand=b_val,
product_name=p_val,
category=str(row.get("category")) if pd.notna(row.get("category")) else None,
price_range=str(row.get("price_range")) if pd.notna(row.get("price_range")) else None,
final_selling_price=fps,
barcode=str(row.get("barcode")) if pd.notna(row.get("barcode")) else None,
hsn_code=str(row.get("hsn_code")) if pd.notna(row.get("hsn_code")) else None,
description=str(row.get("description")) if pd.notna(row.get("description")) else None,
image_url=str(row.get("image_url")) if pd.notna(row.get("image_url")) else None,
)
res = _enrich_and_save_product(req)
added.append({k: v for k, v in res.items() if k != "embedding"})
except Exception as e:
errors.append({"row": idx + 1, "product_name": p_val, "error": str(e)})
return {
"status": "success",
"filename": filename,
"total_rows_processed": len(added) + len(errors),
"added_count": len(added),
"error_count": len(errors),
"added_products": added,
"errors": errors,
}

152
app/api/schemas.py Normal file
View File

@@ -0,0 +1,152 @@
"""Pydantic request/response models for the FastAPI layer."""
from __future__ import annotations
from typing import List, Optional
from pydantic import BaseModel, Field
# ---------------------------------------------------------------------------
# Shared
# ---------------------------------------------------------------------------
class ProductOut(BaseModel):
image_id: str
image_url: Optional[str] = None
image_urls: List[str] = Field(default_factory=list)
brand: str
product_name: str
title: Optional[str] = None
category: Optional[str] = None
description: Optional[str] = None
price_range: Optional[str] = None
size_variants: List[str] = Field(default_factory=list)
providers: List[str] = Field(default_factory=list)
highlights: List[str] = Field(default_factory=list)
nutrients: List[str] = Field(default_factory=list)
fssai_license: Optional[str] = None
product_sku: Optional[str] = None
sku_source: Optional[str] = None
hsn_code: Optional[str] = None
final_selling_price: Optional[float] = None
selling_price: Optional[float] = None
barcode: Optional[str] = None
barcode_type: Optional[str] = None
class SourceProductOut(BaseModel):
image_id: str
image_url: Optional[str] = None
image_urls: List[str] = Field(default_factory=list)
brand: str
product_name: str
title: Optional[str] = None
category: Optional[str] = None
description: Optional[str] = None
price_range: Optional[str] = None
size_variants: List[str] = Field(default_factory=list)
providers: List[str] = Field(default_factory=list)
highlights: List[str] = Field(default_factory=list)
nutrients: List[str] = Field(default_factory=list)
fssai_license: Optional[str] = None
product_sku: Optional[str] = None
sku_source: Optional[str] = None
hsn_code: Optional[str] = None
final_selling_price: Optional[float] = None
selling_price: Optional[float] = None
barcode: Optional[str] = None
barcode_type: Optional[str] = None
similarity: float
# ---------------------------------------------------------------------------
# Health
# ---------------------------------------------------------------------------
class HealthOut(BaseModel):
status: str
database: bool
ollama: bool
ollama_model: str
embeddings_model: str
# ---------------------------------------------------------------------------
# Brands / catalog browsing
# ---------------------------------------------------------------------------
class BrandsOut(BaseModel):
brands: List[str]
class CategoriesOut(BaseModel):
brand: str
categories: List[str]
class ProductListOut(BaseModel):
brand: str
total: int
limit: int
offset: int
products: List[ProductOut]
class AllProductsOut(BaseModel):
total: int
limit: int
offset: int
products: List[ProductOut]
# ---------------------------------------------------------------------------
# Semantic search (retrieval only, no LLM generation)
# ---------------------------------------------------------------------------
class SearchOut(BaseModel):
query: str
brand: Optional[str] = None
results: List[SourceProductOut]
# ---------------------------------------------------------------------------
# RAG chat
# ---------------------------------------------------------------------------
class ChatTurn(BaseModel):
role: str = Field(..., description="'user' or 'assistant'")
content: str
class ChatRequest(BaseModel):
query: str = Field(..., min_length=1, max_length=2000)
brand: Optional[str] = Field(None, description="Restrict retrieval to a single brand")
category: Optional[str] = Field(None, description="Restrict retrieval to a category")
top_k: Optional[int] = Field(None, ge=1, le=15)
history: Optional[List[ChatTurn]] = Field(default=None, description="Prior turns for follow-up questions")
class ChatResponseOut(BaseModel):
answer: str
query: str
brand: Optional[str] = None
detected_category: Optional[str] = Field(
None, description="Product category auto-detected from the query and used to scope retrieval, if any"
)
sources: List[SourceProductOut]
# ---------------------------------------------------------------------------
# Catalog generation (admin/ingestion trigger)
# ---------------------------------------------------------------------------
class CatalogGenerateRequest(BaseModel):
brand: str = Field(..., min_length=1, max_length=120)
max_products: int = Field(50, ge=1, le=300)
class CatalogJobOut(BaseModel):
job_id: str
brand: str
status: str
detail: Optional[str] = None

View File

@@ -0,0 +1,56 @@
"""
In-memory job tracker for the store-intelligence seed/train admin
endpoints - same pattern and same trade-offs as `app/api/job_store.py`
(process-local, lost on restart, fine for a single-developer/single-
process deployment), kept as a separate small module rather than
overloading `job_store.Job.brand` for a job type that isn't
brand-specific (seeding stores and training models operate over the
whole catalog, not one brand).
"""
from __future__ import annotations
import threading
import time
import uuid
from dataclasses import dataclass, field
from typing import Dict, Optional
@dataclass
class StoreJob:
job_id: str
kind: str # "seed" | "train"
status: str = "pending" # pending -> running -> done | failed
detail: Optional[str] = None
result: Optional[dict] = None
created_at: float = field(default_factory=time.time)
updated_at: float = field(default_factory=time.time)
class StoreJobStore:
def __init__(self) -> None:
self._jobs: Dict[str, StoreJob] = {}
self._lock = threading.Lock()
def create(self, kind: str) -> StoreJob:
job = StoreJob(job_id=str(uuid.uuid4()), kind=kind)
with self._lock:
self._jobs[job.job_id] = job
return job
def update(self, job_id: str, status: str, detail: Optional[str] = None, result: Optional[dict] = None) -> None:
with self._lock:
job = self._jobs.get(job_id)
if job:
job.status = status
job.detail = detail
if result is not None:
job.result = result
job.updated_at = time.time()
def get(self, job_id: str) -> Optional[StoreJob]:
with self._lock:
return self._jobs.get(job_id)
store_job_store = StoreJobStore()

95
app/api/store_schemas.py Normal file
View File

@@ -0,0 +1,95 @@
from __future__ import annotations
from typing import Any, Dict, List, Optional
from pydantic import BaseModel, Field
class StoreOut(BaseModel):
store_id: str
store_name: str
city: Optional[str] = None
tier: str
footfall_index: float
class StoreProductOut(BaseModel):
store_id: str
brand: str
image_id: str
title: Optional[str] = None
category: Optional[str] = None
available_stock: int
reserved_stock: int
reorder_level: int
safety_stock: int
stock_status: str
mrp: float
cost_price: float
selling_price: float
profit_margin: float
gross_profit_pct: float
markup_pct: float
image_url: Optional[str] = None
image_urls: Optional[List[str]] = None
fssai_license: Optional[str] = None
class ProductStorePriceOut(BaseModel):
store_id: str
store_name: Optional[str] = None
tier: Optional[str] = None
available_stock: int
mrp: float
cost_price: float
selling_price: float
class DiscountOut(BaseModel):
store_id: str
brand: str
image_id: str
original_price: float
discount_pct: float
final_price: float
savings: float
model_version: Optional[str] = None
class TrendingItemOut(BaseModel):
brand: str
image_id: str
trend_score: float
rank: int
class RecommendationOut(BaseModel):
brand: str
image_id: str
similarity_score: float
method: str = "hybrid"
signals: Optional[Dict[str, float]] = None
class SeedRequest(BaseModel):
reset_orders: bool = Field(default=True, description="Clear existing simulated order history before re-simulating")
days: int = Field(default=90, ge=7, le=365)
seed: int = Field(default=42)
class SeedResponse(BaseModel):
stores: int
store_products: Dict[str, int]
orders: int
order_items: int
class TrainRequest(BaseModel):
models: Optional[List[str]] = Field(
default=None,
description="Subset of models to (re)train: discount, trending, popularity, forecast, store_performance, purchase_propensity. Omit to train all.",
)
class TrainResponse(BaseModel):
trained: Dict[str, Dict[str, Any]]