""" FastAPI application entry point. Run with: uvicorn app.main:app --reload --port 8000 (see backend/README.md / the project documentation for full setup steps) """ from __future__ import annotations import logging import os import threading from contextlib import asynccontextmanager from pathlib import Path from fastapi import FastAPI, HTTPException, Request from fastapi.exceptions import RequestValidationError from fastapi.encoders import jsonable_encoder from fastapi.middleware.cors import CORSMiddleware from fastapi.staticfiles import StaticFiles from fastapi.responses import FileResponse, JSONResponse from app.infrastructure.persistence import restore_bundled_assets from app.infrastructure.settings import ( API_CORS_ORIGINS, BATCH_AUTO_RESUME, BRAND_SYNC_INTERVAL_SECONDS, cleaned_env_names, ) from app.infrastructure.security import auth_config_summary from app.api.routers import health, brands, search, suggest, chat, catalog, system from app.api.routers import stores, discounts, analytics as store_analytics, trending, recommendations, store_admin from app.api.routers import nutrition, nutrition_admin, upload from app.api.routers import auth, user_products, admin_train, mcp_info from app.api.routers import batch_catalog, uploads, brand_discovery from app.services.store_db import ensure_store_intelligence_schema from app.services.nutrition_db import ensure_nutrition_schema logging.basicConfig( level=logging.INFO, format="%(asctime)s - %(name)s - %(levelname)s - %(message)s", ) logger = logging.getLogger(__name__) try: from app.mcp_server import MCP_PATH, build_http_app as build_mcp_app _mcp_app = build_mcp_app() except Exception as e: # pragma: no cover - depends on an optional dependency # fastmcp needs Python >=3.10. The image is 3.11, but an older local # interpreter should lose the MCP endpoint, not the whole API. logging.getLogger(__name__).warning("MCP server unavailable: %s", e) _mcp_app = None MCP_PATH = "/mcp" # Signals the background thread below to stop. Doubles as its sleep: waiting on # an Event rather than time.sleep() means shutdown is immediate instead of # blocking until the current BRAND_SYNC_INTERVAL_SECONDS elapses. _brand_sync_stop = threading.Event() @asynccontextmanager async def lifespan(_app: FastAPI): """ Schema checks and auto-seeding, kicked off without blocking startup. The work runs on a daemon thread rather than being awaited: it touches Postgres, which may be slow or briefly unreachable on a cold boot, and the server answering /api/health in under a second is what lets the container healthcheck pass while that settles. That same thread then stays alive to re-run the brand reconcile on an interval, so a brand table that appears in the database after boot reaches the catalog without a restart. """ # Runs before the thread below, and synchronously: the seed catalogs and # model artifacts have to be in place before the first request can read # them, and it is a handful of file copies on first boot, nothing on every # boot after that. try: restore_bundled_assets() except Exception as e: logger.warning("Could not restore bundled assets: %s", e) # Cheap, and it must happen before anyone can press Resume: a batch that a # restart cut short is still marked "running" on disk, and until it is # reconciled the UI shows it as in flight with nothing behind it. This only # rewrites manifests - it deliberately starts no work. See # batch_ingest.scan_interrupted() for why auto-resume is not the default. try: from app.core.batch_ingest import scan_interrupted interrupted = scan_interrupted() if interrupted: logger.info("Marked %d interrupted catalog batch(es).", len(interrupted)) # Retention, which otherwise only runs when the ingestion worker goes # idle. That was sufficient while every upload was queued work; it is # not now that /api/uploads/catalog stages drops for review and queues # nothing. An inbox nobody acts on would never start the worker, so the # sweep that reclaims abandoned uploads would never run either. try: from app.core.batch_ingest import purge_expired purged = purge_expired() if purged: logger.info("Purged %d expired batch upload(s) at startup.", len(purged)) except Exception as exc: # noqa: BLE001 - housekeeping must not block boot logger.warning("Could not purge expired batch uploads: %s", exc) if interrupted and BATCH_AUTO_RESUME: # Opt-in only. Importing the worker here rather than at module # scope keeps the queue and its thread out of a boot that never # needs them. from app.core import batch_worker for batch_id in interrupted: try: batch_worker.submit(batch_id) except Exception as exc: # noqa: BLE001 - a full queue is not fatal logger.warning("Could not auto-resume batch %s: %s", batch_id, exc) except Exception as e: logger.warning("Could not reconcile interrupted catalog batches: %s", e) def _async_init(): try: ensure_store_intelligence_schema() ensure_nutrition_schema() from app.api.routers.system import _run_background_auto_seed, _run_background_brand_sync _run_background_auto_seed() # Runs after the auto-seed so a cold boot has already loaded the # bundled catalogs and this finds nothing to do. On a warm boot the # auto-seed no-ops and this is what picks up a brand table or seed # file that appeared since last time. _run_background_brand_sync() except Exception as e: logger.warning("Startup background init warning: %s", e) # Everything above is a one-shot boot sweep. Keep going on an interval so # a brand table created directly in the database - bypassing the app, and # therefore bypassing the cache invalidation every write path does - still # reaches the catalog. Without this it waits for the next restart. if BRAND_SYNC_INTERVAL_SECONDS <= 0: logger.info("Periodic brand sync disabled (BRAND_SYNC_INTERVAL_SECONDS=0)") return logger.info("Periodic brand sync every %ss", BRAND_SYNC_INTERVAL_SECONDS) while not _brand_sync_stop.wait(BRAND_SYNC_INTERVAL_SECONDS): try: # Imported per iteration for the same reason as above: a failure # to import must not be what kills the loop. from app.api.routers.system import _run_background_brand_sync _run_background_brand_sync() except Exception as e: # One bad sweep - an unreachable database, a malformed seed file - # must not end the loop, or recovery needs a restart again. logger.warning("Periodic brand sync warning: %s", e) threading.Thread(target=_async_init, daemon=True).start() # The mounted MCP app carries its own lifespan, which starts the session # manager its request handler depends on. Mounting a sub-app does NOT run # that lifespan - only the outermost app's is executed - so without this # the endpoint exists, accepts a connection, and then fails on the first # message with a session manager that was never started. try: if _mcp_app is not None: async with _mcp_app.lifespan(_mcp_app): yield else: yield finally: # Releases the interval wait immediately. The thread is a daemon, so this # is not what lets the process exit - it is what stops a reconcile from # starting against a database the shutdown is already tearing down. _brand_sync_stop.set() app = FastAPI( title="Brand Product Search Engine - RAG API", description=( "Local, CPU-only RAG API over an Indian FMCG product catalog stored in pgvector. " "Includes automated multi-store intelligence, ML discount engines, and nutrition intelligence." ), version="3.2.0", lifespan=lifespan, ) # When the app is served same-origin (the frontend's nginx proxies /api/* to # this service), API_CORS_ORIGINS is empty and no cross-origin request is ever # made. It is populated only when the API is also published on its own host - # see api.{$DOMAIN} in the Caddyfile. # # A wildcard origin and credentialed requests are mutually exclusive under the # CORS spec: browsers reject `Access-Control-Allow-Origin: *` on any request # carrying credentials. Sending both is a silent misconfiguration - the server # looks configured while every browser call fails - so a wildcard turns # credentials off explicitly and says so in the log. _allow_credentials = "*" not in API_CORS_ORIGINS if not _allow_credentials: logger.warning( "API_CORS_ORIGINS contains '*': disabling allow_credentials, because " "browsers reject credentialed cross-origin requests to a wildcard origin. " "List the exact origins instead if you need cookies or Authorization headers." ) app.add_middleware( CORSMiddleware, allow_origins=API_CORS_ORIGINS, allow_credentials=_allow_credentials, allow_methods=["*"], allow_headers=["*"], ) def _json_safe(value): """Replace floats JSON cannot carry (NaN, +/-inf) so an error can be sent.""" if isinstance(value, float) and (value != value or value in (float("inf"), float("-inf"))): return str(value) if isinstance(value, dict): return {k: _json_safe(v) for k, v in value.items()} if isinstance(value, (list, tuple)): return [_json_safe(v) for v in value] return value @app.exception_handler(RequestValidationError) async def _validation_error_as_422(request: Request, exc: RequestValidationError) -> JSONResponse: """FastAPI's own 422 body echoes the rejected input. Python's JSON parser accepts `NaN` on the way in, pydantic rejects it, and the echo then fails to serialise - so a client that sent one NaN in a float field got a 500 instead of the 422 that names the field. Found by the image-vector search, where the body is 1024 floats; applies to every route.""" return JSONResponse(status_code=422, content={"detail": _json_safe(jsonable_encoder(exc.errors()))}) # A wrong origin list fails only in the browser, as an opaque "blocked by CORS" # with a perfectly healthy 200 in the server log - so state the effective list # at startup, where it can actually be compared against the frontend's URL. logger.info( "CORS allowed origins: %s", ", ".join(API_CORS_ORIGINS) if API_CORS_ORIGINS else "(none - same-origin only)", ) if API_CORS_ORIGINS and all( o.startswith(("http://localhost", "http://127.0.0.1")) for o in API_CORS_ORIGINS ): logger.warning( "API_CORS_ORIGINS lists only localhost origins (%s). If this API is " "published on a domain and the frontend is served from a different one, " "every browser request will be blocked. Set API_CORS_ORIGINS to the " "frontend's exact origin, e.g. https://catalogue.nearle.ai.in", ", ".join(API_CORS_ORIGINS), ) # The same argument as the CORS block above, for the other setting whose # misconfiguration is invisible from the outside. A wrong credential fails only # as "Invalid username or password.", which is indistinguishable from a user # mistyping - so state what the process actually loaded, at startup, where it # can be compared against the config the image was built from. # # No password and no hash is printed. `fingerprint` identifies WHICH credential # is loaded (see security.password_hash_fingerprint); `source` says whether it # came from the container's environment or from the .env file, which is the # only way to notice a deployment platform's Environment tab overriding the # image. Compare against: python scripts/make_auth_secrets.py --fingerprint _auth_cfg = auth_config_summary() logger.info( "Auth config: enabled=%s allow_any_login=%s admin_username=%r " "hash=%s/%s fingerprint=%s source=%s (username source=%s) " "api_keys=%d/%s %s", _auth_cfg["enabled"], _auth_cfg["allow_any_login"], _auth_cfg["admin_username"], "pbkdf2_sha256" if _auth_cfg["password_hash_valid"] else "INVALID", _auth_cfg["password_hash_iterations"], _auth_cfg["password_hash_fingerprint"] or "(none)", _auth_cfg["password_hash_source"], _auth_cfg["admin_username_source"], _auth_cfg["api_keys_count"], _auth_cfg["api_keys_source"], # Names, not secrets. A key added to .env.production but only restarted into # a running container never appears here - which is the whole point. [k["name"] for k in _auth_cfg["api_keys"]] or "(none)", ) if cleaned_env_names(): logger.warning( "These settings arrived wrapped in quotes or padded with whitespace and were " "cleaned before use: %s. Pasting into a deployment platform's Environment tab " "is the usual source. They work now, but the next value may not - store them " "unquoted.", ", ".join(cleaned_env_names()), ) if _auth_cfg["enabled"] and _auth_cfg["password_hash_source"] == "process-env": logger.warning( "AUTH_ADMIN_PASSWORD_HASH came from the process environment, which OVERRIDES " "the .env file (load_dotenv is called without override=True). Under " "docker-compose that is just `env_file:` and is expected. Under Dokploy it " "means the service's Environment tab is supplying this credential and the one " "baked into the image by `COPY .env.production .env` is being ignored - which " "is how a corrected password keeps failing after a redeploy." ) if _auth_cfg["enabled"] and not _auth_cfg["password_hash_valid"]: logger.error( "AUTH_ADMIN_PASSWORD_HASH is not a usable PBKDF2 digest, so EVERY sign-in " "will return 401 no matter which password is typed. It came from %s. " "Regenerate it with: python scripts/make_auth_secrets.py", _auth_cfg["password_hash_source"], ) app.include_router(health.router, prefix="/api") app.include_router(auth.router, prefix="/api") app.include_router(user_products.router, prefix="/api") app.include_router(admin_train.router, prefix="/api") app.include_router(system.router, prefix="/api") app.include_router(brands.router, prefix="/api") app.include_router(search.router, prefix="/api") app.include_router(suggest.router, prefix="/api") app.include_router(chat.router, prefix="/api") app.include_router(catalog.router, prefix="/api") app.include_router(stores.router, prefix="/api") app.include_router(discounts.router, prefix="/api") app.include_router(store_analytics.router, prefix="/api") app.include_router(trending.router, prefix="/api") app.include_router(recommendations.router, prefix="/api") app.include_router(store_admin.router, prefix="/api") app.include_router(nutrition.router, prefix="/api") app.include_router(nutrition_admin.router, prefix="/api") app.include_router(upload.router, prefix="/api") app.include_router(batch_catalog.router, prefix="/api") app.include_router(uploads.router, prefix="/api") app.include_router(brand_discovery.router, prefix="/api") app.include_router(mcp_info.router, prefix="/api") # MCP lives outside /api on purpose: it is a protocol endpoint for AI clients, # not part of the REST surface, and the catch-all SPA route below only skips # paths it recognises. Mounted rather than routed because it is a whole ASGI # app with its own request handling. if _mcp_app is not None: app.mount(MCP_PATH, _mcp_app) logger.info("MCP server mounted at %s", MCP_PATH) # Serve built frontend static files if dist exists (single-port unified # deployment). Off in the normal setup: the React app is served by its own # nginx on catalogue.nearle.ai.in and this API answers on # mcp.nearle.ai.in, so no dist/ is present here and the JSON root # handler at the bottom of this file is what responds to /. # # The candidates cover both repo layouts - the sibling checkout is named # `catalogue_frontend`, and only `frontend` was checked before, so this branch # could never activate even when a build was sitting right next to it. # FRONTEND_DIST_DIR overrides both when the build lands somewhere else. _dist_override = os.getenv("FRONTEND_DIST_DIR", "").strip() _repo_root = Path(__file__).resolve().parents[2] _dist_candidates = ( [Path(_dist_override)] if _dist_override else [ _repo_root / "catalogue_frontend" / "dist", _repo_root / "frontend" / "dist", Path(__file__).resolve().parents[1] / "frontend" / "dist", ] ) FRONTEND_DIST = next( (p for p in _dist_candidates if (p / "assets").exists()), _dist_candidates[0], ) if FRONTEND_DIST.exists() and (FRONTEND_DIST / "assets").exists(): logger.info("Serving built frontend from %s", FRONTEND_DIST) app.mount("/assets", StaticFiles(directory=str(FRONTEND_DIST / "assets")), name="assets") @app.get("/{full_path:path}") def serve_frontend(full_path: str): # API and docs paths are served by the routers registered above, so # this catch-all only sees them when the path genuinely doesn't exist. # Returning None there would answer 200 with a `null` body - an unknown # endpoint would look like a successful call to any client. 404 is the # honest answer, and it matters now that the API is also reachable # directly at api.{$DOMAIN} rather than only behind the frontend. # "mcp" is in this list for the case where the MCP app failed to load: # the mount would be absent, and without this the SPA fallback would # answer an MCP client with index.html and a 200. if full_path.startswith(("api", "docs", "redoc", "openapi.json", "mcp")): raise HTTPException(status_code=404, detail="Not found") file_path = FRONTEND_DIST / full_path if file_path.exists() and file_path.is_file(): return FileResponse(file_path) return FileResponse(FRONTEND_DIST / "index.html") else: @app.get("/") def root() -> dict: return { "service": "Brand Product Search Engine - RAG API", "docs": "/docs", "health": "/api/health", "system_status": "/api/system/status", }