Merge branch 'master' of https://gitapp.workolik.com/nearle_daily/catalogue_backend
This commit is contained in:
167
app/api/routers/mcp_info.py
Normal file
167
app/api/routers/mcp_info.py
Normal file
@@ -0,0 +1,167 @@
|
||||
"""
|
||||
REST inspector for the MCP server.
|
||||
|
||||
The MCP endpoint itself speaks streamable HTTP with session handling, which a
|
||||
browser cannot usefully talk to without a full MCP client implementation. Rather
|
||||
than ship one in React, these two endpoints expose what the admin page needs -
|
||||
what tools exist, and what one returns - as ordinary JSON over the REST API the
|
||||
frontend already authenticates against.
|
||||
|
||||
Admin-only. The tool inventory describes the whole catalog API surface, and the
|
||||
invoke endpoint executes a tool; neither is something to leave open just because
|
||||
the tools happen to be read-only.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
from typing import Any, Dict, List, Optional
|
||||
|
||||
from fastapi import APIRouter, Body, Depends, HTTPException, Request
|
||||
|
||||
from app.api.deps import require_admin
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
router = APIRouter(prefix="/mcp", tags=["mcp"])
|
||||
|
||||
|
||||
def _server():
|
||||
"""The FastMCP instance, or None when the optional dependency is absent."""
|
||||
try:
|
||||
from app.mcp_server import mcp
|
||||
|
||||
return mcp
|
||||
except Exception: # pragma: no cover - depends on fastmcp being installed
|
||||
return None
|
||||
|
||||
|
||||
def _public_url(request: Request, path: str) -> str:
|
||||
"""
|
||||
The URL an external MCP client should connect to.
|
||||
|
||||
Built from the forwarded headers rather than a configured constant so it is
|
||||
right on whichever host the request arrived on - the API's own domain, the
|
||||
frontend's nginx, or localhost in development. uvicorn runs with
|
||||
--proxy-headers, so request.url already reflects the external scheme.
|
||||
"""
|
||||
base = str(request.base_url).rstrip("/")
|
||||
return f"{base}{path}"
|
||||
|
||||
|
||||
@router.get("/info", dependencies=[Depends(require_admin)])
|
||||
async def mcp_info(request: Request) -> Dict[str, Any]:
|
||||
"""Describe the MCP server and list its tools, for the admin page."""
|
||||
mcp = _server()
|
||||
if mcp is None:
|
||||
return {
|
||||
"enabled": False,
|
||||
"reason": (
|
||||
"The fastmcp package is not installed, or failed to import. It "
|
||||
"requires Python 3.10 or newer."
|
||||
),
|
||||
"tools": [],
|
||||
}
|
||||
|
||||
from app.mcp_server import MCP_PATH
|
||||
|
||||
try:
|
||||
# run_middleware=False: the auth middleware reads an Authorization header
|
||||
# from an MCP request context that does not exist on this REST call. The
|
||||
# caller is already established as an admin by the route dependency.
|
||||
tools = await mcp.list_tools(run_middleware=False)
|
||||
except Exception as exc: # noqa: BLE001
|
||||
logger.exception("Could not list MCP tools")
|
||||
raise HTTPException(status_code=500, detail=f"Could not list MCP tools: {exc}")
|
||||
|
||||
described: List[Dict[str, Any]] = []
|
||||
for tool in sorted(tools, key=lambda t: t.name):
|
||||
# Server-side FunctionTool exposes its JSON Schema as `parameters`.
|
||||
# `inputSchema` is the wire-format name, present on the client-side Tool
|
||||
# an MCP client receives - checked second so this keeps working if a
|
||||
# future version converges on one name.
|
||||
schema = getattr(tool, "parameters", None) or getattr(tool, "inputSchema", None) or {}
|
||||
properties = schema.get("properties", {}) or {}
|
||||
required = set(schema.get("required", []) or [])
|
||||
described.append(
|
||||
{
|
||||
"name": tool.name,
|
||||
"description": (tool.description or "").strip(),
|
||||
"tags": sorted(getattr(tool, "tags", set()) or []),
|
||||
"parameters": [
|
||||
{
|
||||
"name": pname,
|
||||
"type": pinfo.get("type") or "any",
|
||||
"description": (pinfo.get("description") or "").strip(),
|
||||
"required": pname in required,
|
||||
"default": pinfo.get("default"),
|
||||
}
|
||||
for pname, pinfo in properties.items()
|
||||
],
|
||||
}
|
||||
)
|
||||
|
||||
return {
|
||||
"enabled": True,
|
||||
"name": mcp.name,
|
||||
"version": getattr(mcp, "version", None),
|
||||
"url": _public_url(request, MCP_PATH),
|
||||
"transport": "http",
|
||||
"auth": "Bearer token from POST /api/auth/login",
|
||||
"tool_count": len(described),
|
||||
"tools": described,
|
||||
}
|
||||
|
||||
|
||||
@router.post("/tools/{tool_name}", dependencies=[Depends(require_admin)])
|
||||
async def invoke_tool(
|
||||
tool_name: str,
|
||||
arguments: Optional[Dict[str, Any]] = Body(default=None),
|
||||
) -> Dict[str, Any]:
|
||||
"""
|
||||
Run one MCP tool and return its result - the admin page's "try it" console.
|
||||
|
||||
Exists so the server can be checked from the browser without installing an
|
||||
MCP client. Every tool is read-only, so this executes the same call an AI
|
||||
client would make, with the same code path.
|
||||
"""
|
||||
mcp = _server()
|
||||
if mcp is None:
|
||||
raise HTTPException(status_code=503, detail="The MCP server is not available.")
|
||||
|
||||
try:
|
||||
# Same reasoning as above: this request authenticated through the REST
|
||||
# guard, so the MCP auth middleware would have no headers to inspect.
|
||||
tools = {t.name: t for t in await mcp.list_tools(run_middleware=False)}
|
||||
except Exception as exc: # noqa: BLE001
|
||||
raise HTTPException(status_code=500, detail=f"Could not list MCP tools: {exc}")
|
||||
|
||||
if tool_name not in tools:
|
||||
raise HTTPException(
|
||||
status_code=404,
|
||||
detail=f"No MCP tool named '{tool_name}'. Known tools: {', '.join(sorted(tools))}.",
|
||||
)
|
||||
|
||||
try:
|
||||
result = await mcp.call_tool(tool_name, arguments or {}, run_middleware=False)
|
||||
except Exception as exc: # noqa: BLE001
|
||||
# A tool raising ToolError is a normal outcome worth showing verbatim -
|
||||
# "no nutrition data for this product" is an answer, not a server fault.
|
||||
return {"tool": tool_name, "ok": False, "error": str(exc)}
|
||||
|
||||
return {"tool": tool_name, "ok": True, "result": _jsonable(result)}
|
||||
|
||||
|
||||
def _jsonable(result: Any) -> Any:
|
||||
"""Reduce a ToolResult to something FastAPI can serialise."""
|
||||
for attr in ("structured_content", "data"):
|
||||
value = getattr(result, attr, None)
|
||||
if value is not None:
|
||||
return value
|
||||
|
||||
content = getattr(result, "content", None)
|
||||
if content is None:
|
||||
return result
|
||||
out = []
|
||||
for block in content:
|
||||
text = getattr(block, "text", None)
|
||||
out.append(text if text is not None else str(block))
|
||||
return out
|
||||
@@ -1,6 +1,7 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import io
|
||||
import json
|
||||
import uuid
|
||||
import logging
|
||||
import pandas as pd
|
||||
|
||||
@@ -7,6 +7,7 @@ from pydantic import BaseModel, Field
|
||||
from fastapi import APIRouter, Depends, HTTPException, File, UploadFile
|
||||
|
||||
from app.api.deps import require_permission
|
||||
from app.infrastructure.settings import SEED_CATALOG_DIR
|
||||
|
||||
from app.services.vector_store import (
|
||||
upsert_brand_products,
|
||||
|
||||
@@ -19,7 +19,7 @@ sys.path.append(str(Path(__file__).parent.parent))
|
||||
|
||||
from app.services.ollama_service import fetch_brand_catalog_with_gemini, fetch_brand_catalog_exhaustive, fetch_product_details
|
||||
from app.services.image_search import find_all_image_urls, find_product_quantity_openfacts
|
||||
from app.infrastructure.settings import USE_OLLAMA
|
||||
from app.infrastructure.settings import DATA_DIR, USE_OLLAMA
|
||||
from app.services.embeddings_service import embed_texts
|
||||
from app.services.vector_store import ensure_brand_schema, upsert_brand_products, get_existing_product_image_id
|
||||
from app.services.s3_service import s3_service
|
||||
@@ -948,20 +948,30 @@ class ProductCatalogEngine:
|
||||
return catalog
|
||||
|
||||
def save_catalog(self, catalog: Dict[str, Any], filename: str = None) -> str:
|
||||
"""Save catalog to JSON file inside data/ folder"""
|
||||
"""Save catalog to JSON file inside DATA_DIR.
|
||||
|
||||
Resolved against DATA_DIR rather than a bare relative "data/" path: the
|
||||
latter depends on the process's working directory, so the same call
|
||||
landed in a different place depending on whether the app was started
|
||||
from the repo root, from backend/, or by the container's uvicorn - and
|
||||
only one of those is the directory with a volume mounted on it.
|
||||
"""
|
||||
if not filename:
|
||||
brand = catalog.get('brand', 'unknown')
|
||||
storage_brand = resolve_parent_brand(brand)
|
||||
safe_brand = storage_brand.replace(' ', '_')
|
||||
ts = catalog.get('generation_timestamp', 'latest')
|
||||
filename = f"data/catalog_{safe_brand}_{ts}.json"
|
||||
|
||||
Path("data").mkdir(exist_ok=True)
|
||||
with open(filename, 'w', encoding='utf-8') as f:
|
||||
filename = f"catalog_{safe_brand}_{ts}.json"
|
||||
|
||||
path = Path(filename)
|
||||
if not path.is_absolute():
|
||||
path = DATA_DIR / path
|
||||
path.parent.mkdir(parents=True, exist_ok=True)
|
||||
with open(path, 'w', encoding='utf-8') as f:
|
||||
json.dump(catalog, f, indent=2, ensure_ascii=False)
|
||||
|
||||
logger.info(f"💾 Catalog saved to: {filename}")
|
||||
return filename
|
||||
|
||||
logger.info(f"💾 Catalog saved to: {path}")
|
||||
return str(path)
|
||||
|
||||
# Global engine instance
|
||||
catalog_engine = ProductCatalogEngine()
|
||||
|
||||
146
app/infrastructure/persistence.py
Normal file
146
app/infrastructure/persistence.py
Normal file
@@ -0,0 +1,146 @@
|
||||
"""
|
||||
First-run restore of the assets bundled into the container image.
|
||||
|
||||
The problem this solves
|
||||
----------------------
|
||||
Three directories are written to at runtime - the seed catalogs appended to by
|
||||
``POST /api/user/products/add``, the generated catalogs from the ingestion
|
||||
pipeline, and the ``*.joblib`` bundles the training endpoints produce. In a
|
||||
container all three sit inside the image, so a redeploy silently discards
|
||||
every one of them. They need to be on a volume.
|
||||
|
||||
Mounting a volume over them introduces the opposite problem. A *named* Docker
|
||||
volume is seeded from the image the first time it is used, but a *bind* mount
|
||||
starts empty and merely hides what the image had underneath. Dokploy offers
|
||||
both and neither announces which you picked, so a bind mount on ``/app/data``
|
||||
would leave the API running with no seed catalogs: the next product added would
|
||||
write a JSON file containing only that product, and every model would report as
|
||||
untrained.
|
||||
|
||||
The fix is to keep a pristine copy inside the image at a path nobody mounts
|
||||
(``BUNDLED_ASSETS_DIR``, populated by the Dockerfile) and top up the writable
|
||||
directory from it on startup.
|
||||
|
||||
Existing files are never overwritten. That is the whole contract: the bundle
|
||||
supplies what is missing, and anything the running app has already written wins
|
||||
over the copy baked into the image. Without that rule every redeploy would
|
||||
revert user-added products back to the bundled catalog.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
import shutil
|
||||
from pathlib import Path
|
||||
from typing import Tuple
|
||||
|
||||
from app.infrastructure.settings import (
|
||||
BUNDLED_ASSETS_DIR,
|
||||
DATA_DIR,
|
||||
MODEL_ARTIFACTS_DIR,
|
||||
SEED_CATALOG_DIR,
|
||||
)
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
# (subdirectory under BUNDLED_ASSETS_DIR, writable destination)
|
||||
_BUNDLES: Tuple[Tuple[str, Path], ...] = (
|
||||
("seed_catalogs", SEED_CATALOG_DIR),
|
||||
("artifacts", MODEL_ARTIFACTS_DIR),
|
||||
)
|
||||
|
||||
|
||||
def _restore_one(source: Path, destination: Path) -> int:
|
||||
"""Copy files missing from ``destination``. Returns how many were copied."""
|
||||
if not source.is_dir():
|
||||
return 0
|
||||
|
||||
destination.mkdir(parents=True, exist_ok=True)
|
||||
copied = 0
|
||||
for item in sorted(source.iterdir()):
|
||||
if not item.is_file():
|
||||
continue
|
||||
target = destination / item.name
|
||||
if target.exists():
|
||||
continue
|
||||
try:
|
||||
# copy2 rather than copy: it preserves mtime, so "is this artifact
|
||||
# older than the data it was trained on" stays answerable.
|
||||
shutil.copy2(item, target)
|
||||
copied += 1
|
||||
except OSError as exc:
|
||||
logger.warning("Could not restore %s -> %s: %s", item, target, exc)
|
||||
return copied
|
||||
|
||||
|
||||
def restore_bundled_assets() -> None:
|
||||
"""
|
||||
Top up the writable directories from the image's read-only bundle.
|
||||
|
||||
Safe to call on every boot: it is a no-op once the volume is populated, and
|
||||
a no-op outside Docker where BUNDLED_ASSETS_DIR does not exist.
|
||||
"""
|
||||
# Create these regardless. A volume mounted at DATA_DIR arrives empty, and
|
||||
# the ingestion pipeline writes into it without creating it first.
|
||||
for path in (DATA_DIR, SEED_CATALOG_DIR, MODEL_ARTIFACTS_DIR):
|
||||
try:
|
||||
path.mkdir(parents=True, exist_ok=True)
|
||||
except OSError as exc:
|
||||
logger.error(
|
||||
"Cannot create writable directory %s: %s. Uploads and trained "
|
||||
"models will fail to save - check the volume's permissions.",
|
||||
path,
|
||||
exc,
|
||||
)
|
||||
|
||||
if not BUNDLED_ASSETS_DIR.is_dir():
|
||||
logger.debug(
|
||||
"No bundled asset directory at %s - nothing to restore.", BUNDLED_ASSETS_DIR
|
||||
)
|
||||
return
|
||||
|
||||
for name, destination in _BUNDLES:
|
||||
copied = _restore_one(BUNDLED_ASSETS_DIR / name, destination)
|
||||
if copied:
|
||||
logger.info(
|
||||
"Restored %d bundled file(s) into %s (first run on this volume).",
|
||||
copied,
|
||||
destination,
|
||||
)
|
||||
|
||||
_warn_if_not_persistent()
|
||||
|
||||
|
||||
def _warn_if_not_persistent() -> None:
|
||||
"""
|
||||
Point out that the writable directories are still inside the image.
|
||||
|
||||
Only reachable when BUNDLED_ASSETS_DIR exists, i.e. in the container. If the
|
||||
destinations were never mounted, everything written to them is lost on the
|
||||
next redeploy - which looks exactly like the app quietly ignoring uploads,
|
||||
hours later and with nothing in the logs to connect it to.
|
||||
"""
|
||||
unmounted = [p for p in (SEED_CATALOG_DIR, MODEL_ARTIFACTS_DIR) if not _is_mount(p)]
|
||||
if unmounted:
|
||||
logger.warning(
|
||||
"These directories are written at runtime but do not look like "
|
||||
"mount points: %s. Anything saved there - products added through "
|
||||
"the UI, retrained models - will be discarded on the next "
|
||||
"redeploy. Mount a volume on each (see backend/README.md).",
|
||||
", ".join(str(p) for p in unmounted),
|
||||
)
|
||||
|
||||
|
||||
def _is_mount(path: Path) -> bool:
|
||||
"""
|
||||
Whether ``path`` sits on a different device than its parent.
|
||||
|
||||
A mounted volume shows up as a device-number change. Falls back to True on
|
||||
error so a probe failure produces silence rather than a false alarm telling
|
||||
somebody their correctly-mounted volume is broken.
|
||||
"""
|
||||
try:
|
||||
if path.is_mount():
|
||||
return True
|
||||
return path.stat().st_dev != path.parent.stat().st_dev
|
||||
except OSError:
|
||||
return True
|
||||
@@ -55,6 +55,52 @@ def _require(name: str, *, feature_flag: str) -> str:
|
||||
return value
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Writable data directories (persistence)
|
||||
# ---------------------------------------------------------------------------
|
||||
# Everything the running app WRITES lives under one of these three paths. They
|
||||
# are settings rather than hard-coded paths because in a container they must be
|
||||
# mounted on a volume - otherwise every product added through the UI and every
|
||||
# retrained model is discarded the next time the image is redeployed.
|
||||
#
|
||||
# DATA_DIR generated catalogs (catalog_engine.save_catalog)
|
||||
# SEED_CATALOG_DIR per-brand JSON catalogs, appended to by
|
||||
# POST /api/user/products/add and /upload-file
|
||||
# MODEL_ARTIFACTS_DIR *.joblib bundles written by the training endpoints
|
||||
#
|
||||
# See BUNDLED_ASSETS_DIR below for how the read-only copies shipped inside the
|
||||
# image get into these directories the first time a volume is mounted.
|
||||
_BACKEND_ROOT = Path(__file__).resolve().parents[2]
|
||||
|
||||
|
||||
def _dir(name: str, default: Path) -> Path:
|
||||
raw = os.getenv(name, "").strip()
|
||||
return Path(raw).expanduser() if raw else default
|
||||
|
||||
|
||||
DATA_DIR = _dir("DATA_DIR", _BACKEND_ROOT / "data")
|
||||
SEED_CATALOG_DIR = _dir("SEED_CATALOG_DIR", DATA_DIR / "seed_catalogs")
|
||||
MODEL_ARTIFACTS_DIR = _dir(
|
||||
"MODEL_ARTIFACTS_DIR", _BACKEND_ROOT / "app" / "intelligence" / "artifacts"
|
||||
)
|
||||
|
||||
# Pristine copies of the bundled seed catalogs and pre-trained models, placed
|
||||
# here by the Dockerfile at a path that is never itself mounted over.
|
||||
#
|
||||
# This exists because the two ways of mounting a volume behave differently, and
|
||||
# the difference is silent. Docker copies the image's content into a *named*
|
||||
# volume the first time it is used, but a *bind* mount starts empty and simply
|
||||
# hides whatever the image had at that path. Mounting a bind mount on /app/data
|
||||
# would therefore leave the app with no seed catalogs at all: the next product
|
||||
# added would write a fresh JSON file containing only that one product, and the
|
||||
# ML endpoints would report no trained models.
|
||||
#
|
||||
# So the image keeps a second, unmounted copy, and restore_bundled_assets()
|
||||
# (app/infrastructure/persistence.py) fills in whatever the writable directory
|
||||
# is missing at startup. Empty/absent outside Docker, where nothing is mounted
|
||||
# and the defaults above already point at the real files.
|
||||
BUNDLED_ASSETS_DIR = _dir("BUNDLED_ASSETS_DIR", Path("/app/.bundled"))
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Ollama (local LLM)
|
||||
# ---------------------------------------------------------------------------
|
||||
@@ -81,6 +127,14 @@ DB_PORT = os.getenv("DB_PORT", "5432")
|
||||
DB_NAME = os.getenv("DB_NAME", "pgvector")
|
||||
DB_USER = os.getenv("DB_USER", "postgres")
|
||||
DB_PASSWORD = _require("DB_PASSWORD", feature_flag="USE_PGVECTOR") if USE_PGVECTOR else os.getenv("DB_PASSWORD", "")
|
||||
# How long to wait for the TCP connect before giving up. Matters more than it
|
||||
# looks: a host that DROPS packets (a firewall, a typo'd DB_HOST) otherwise
|
||||
# blocks until the OS timeout - about 130 seconds on Linux - and every request
|
||||
# that touches the database inherits that wait, including /api/health. A short
|
||||
# ceiling turns "the database is unreachable" into a fast, honest error instead
|
||||
# of a hung worker and a container the platform decides is unhealthy.
|
||||
DB_CONNECT_TIMEOUT_SECONDS = int(os.getenv("DB_CONNECT_TIMEOUT_SECONDS", "5"))
|
||||
|
||||
DATABASE_URL = os.getenv(
|
||||
"DATABASE_URL",
|
||||
f"postgresql://{DB_USER}:{DB_PASSWORD}@{DB_HOST}:{DB_PORT}/{DB_NAME}",
|
||||
|
||||
@@ -13,9 +13,15 @@ from datetime import datetime, timezone
|
||||
from pathlib import Path
|
||||
from typing import Any, Dict, List, Optional
|
||||
|
||||
from app.infrastructure.settings import MODEL_ARTIFACTS_DIR
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
ARTIFACTS_DIR = Path(__file__).resolve().parent / "artifacts"
|
||||
# From settings so it can be pointed at a mounted volume: retraining writes
|
||||
# *.joblib here, and inside the image those are discarded on the next redeploy,
|
||||
# silently reverting every model to the version baked in at build time.
|
||||
# Defaults to this package's own artifacts/ directory.
|
||||
ARTIFACTS_DIR = MODEL_ARTIFACTS_DIR
|
||||
ARTIFACTS_DIR.mkdir(parents=True, exist_ok=True)
|
||||
|
||||
|
||||
|
||||
98
app/main.py
98
app/main.py
@@ -9,6 +9,7 @@ Run with:
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
import os
|
||||
import threading
|
||||
from contextlib import asynccontextmanager
|
||||
from pathlib import Path
|
||||
@@ -18,11 +19,12 @@ from fastapi.middleware.cors import CORSMiddleware
|
||||
from fastapi.staticfiles import StaticFiles
|
||||
from fastapi.responses import FileResponse
|
||||
|
||||
from app.infrastructure.persistence import restore_bundled_assets
|
||||
from app.infrastructure.settings import API_CORS_ORIGINS
|
||||
from app.api.routers import health, brands, search, 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
|
||||
from app.api.routers import auth, user_products, admin_train, mcp_info
|
||||
from app.services.store_db import ensure_store_intelligence_schema
|
||||
from app.services.nutrition_db import ensure_nutrition_schema
|
||||
|
||||
@@ -32,6 +34,18 @@ logging.basicConfig(
|
||||
)
|
||||
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"
|
||||
|
||||
|
||||
@asynccontextmanager
|
||||
async def lifespan(_app: FastAPI):
|
||||
"""
|
||||
@@ -43,6 +57,15 @@ async def lifespan(_app: FastAPI):
|
||||
healthcheck pass while that settles.
|
||||
"""
|
||||
|
||||
# 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)
|
||||
|
||||
def _async_init():
|
||||
try:
|
||||
ensure_store_intelligence_schema()
|
||||
@@ -58,7 +81,17 @@ async def lifespan(_app: FastAPI):
|
||||
logger.warning("Startup background init warning: %s", e)
|
||||
|
||||
threading.Thread(target=_async_init, daemon=True).start()
|
||||
yield
|
||||
|
||||
# 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.
|
||||
if _mcp_app is not None:
|
||||
async with _mcp_app.lifespan(_mcp_app):
|
||||
yield
|
||||
else:
|
||||
yield
|
||||
|
||||
|
||||
app = FastAPI(
|
||||
@@ -98,6 +131,24 @@ app.add_middleware(
|
||||
allow_headers=["*"],
|
||||
)
|
||||
|
||||
# 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),
|
||||
)
|
||||
|
||||
app.include_router(health.router, prefix="/api")
|
||||
app.include_router(auth.router, prefix="/api")
|
||||
app.include_router(user_products.router, prefix="/api")
|
||||
@@ -116,10 +167,44 @@ 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(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],
|
||||
)
|
||||
|
||||
# Serve built frontend static files if dist exists (single-port unified deployment)
|
||||
FRONTEND_DIST = Path(__file__).resolve().parents[2] / "frontend" / "dist"
|
||||
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}")
|
||||
@@ -130,7 +215,10 @@ if FRONTEND_DIST.exists() and (FRONTEND_DIST / "assets").exists():
|
||||
# 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.
|
||||
if full_path.startswith(("api", "docs", "redoc", "openapi.json")):
|
||||
# "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():
|
||||
|
||||
522
app/mcp_server.py
Normal file
522
app/mcp_server.py
Normal file
@@ -0,0 +1,522 @@
|
||||
"""
|
||||
MCP server exposing the catalog as tools for AI clients.
|
||||
|
||||
Mounted onto the FastAPI app at /mcp (see app/main.py), so it ships in the same
|
||||
container and answers on the same host - mcp.nearle.ai.in/mcp.
|
||||
|
||||
Scope: reads only. Every tool here maps to a service function the REST API
|
||||
already exposes through a GET. None of the write or compute endpoints - catalog
|
||||
generation, ML training, the upload endpoints, product creation - are reachable
|
||||
through MCP, because a tool list is chosen from by a model rather than by a
|
||||
person, and "retrain the models" is not a thing to leave one tool call away.
|
||||
Adding a write tool later is deliberate work, not an oversight to correct.
|
||||
|
||||
Authentication reuses the access token from POST /api/auth/login. There is no
|
||||
separate MCP credential: the same Principal, the same expiry, the same
|
||||
revocation story as the REST API, and nothing new to keep secret. Clients send
|
||||
it as an Authorization: Bearer header, which _AuthMiddleware checks once for
|
||||
every tool call and tool listing.
|
||||
|
||||
The tools are hand-written rather than generated from the OpenAPI schema. All
|
||||
63 routes would technically work, but a model picks a tool by reading its
|
||||
description, and 63 near-identical auto-generated entries is a worse starting
|
||||
point than a dozen written to be chosen between.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
from typing import Any, Dict, List, Optional
|
||||
|
||||
from fastmcp import FastMCP
|
||||
from fastmcp.exceptions import ToolError
|
||||
from fastmcp.server.dependencies import get_http_headers
|
||||
from fastmcp.server.middleware import CallNext, Middleware, MiddlewareContext
|
||||
|
||||
from app.infrastructure.security import AuthError, Principal, anonymous_principal, decode_access_token
|
||||
from app.infrastructure.settings import AUTH_ENABLED
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
MCP_PATH = "/mcp"
|
||||
|
||||
INSTRUCTIONS = """\
|
||||
Read-only access to an Indian FMCG product catalog: products and brands, \
|
||||
nutrition facts and health scores, per-store inventory and pricing, discounts, \
|
||||
and sales analytics.
|
||||
|
||||
Products are identified by a (brand, image_id) pair - image_id is the SKU-like \
|
||||
product key, not an image URL. Get one from search_products or \
|
||||
list_brand_products before calling any tool that takes an image_id.
|
||||
|
||||
Nutrition figures are per 100g unless the product states otherwise, and health \
|
||||
scores run 0-100, higher being better. Prices are in INR.
|
||||
"""
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Authentication
|
||||
# ---------------------------------------------------------------------------
|
||||
def _principal_from_headers() -> Principal:
|
||||
"""
|
||||
Resolve the caller from the Authorization header of the current request.
|
||||
|
||||
Raises ToolError rather than AuthError: ToolError is what reaches the client
|
||||
as a readable message instead of an opaque internal failure.
|
||||
"""
|
||||
if not AUTH_ENABLED:
|
||||
return anonymous_principal()
|
||||
|
||||
# include={"authorization"} is required, not optional: get_http_headers()
|
||||
# strips the Authorization header by default, so that it is not forwarded to
|
||||
# downstream services by accident. Without this the header is invisible here
|
||||
# and every request looks unauthenticated - including valid ones.
|
||||
headers = get_http_headers(include={"authorization"})
|
||||
raw = headers.get("authorization", "") # keys are lowercased by the helper
|
||||
scheme, _, token = raw.partition(" ")
|
||||
|
||||
if scheme.lower() != "bearer" or not token.strip():
|
||||
raise ToolError(
|
||||
"Not authenticated. Send an Authorization: Bearer <token> header. "
|
||||
"Get a token from POST /api/auth/login."
|
||||
)
|
||||
try:
|
||||
return decode_access_token(token.strip())
|
||||
except AuthError as exc:
|
||||
raise ToolError(str(exc)) from exc
|
||||
|
||||
|
||||
class _AuthMiddleware(Middleware):
|
||||
"""
|
||||
One authentication check covering every tool.
|
||||
|
||||
Sitting here rather than at the top of each tool means a tool added later
|
||||
cannot be left unguarded by forgetting a line - the same reason the REST
|
||||
side puts its guard in a route dependency instead of the handler body.
|
||||
|
||||
Listing is checked as well as calling. An unauthenticated client learning
|
||||
the exact shape of the catalog API is a smaller problem than an
|
||||
unauthenticated call, but it is not nothing, and there is no reason to
|
||||
publish it.
|
||||
"""
|
||||
|
||||
async def on_call_tool(self, context: MiddlewareContext, call_next: CallNext):
|
||||
principal = _principal_from_headers()
|
||||
name = getattr(context.message, "name", "?")
|
||||
logger.info("MCP tool call %r by %s (%s)", name, principal.username, principal.kind)
|
||||
return await call_next(context)
|
||||
|
||||
async def on_list_tools(self, context: MiddlewareContext, call_next: CallNext):
|
||||
_principal_from_headers()
|
||||
return await call_next(context)
|
||||
|
||||
|
||||
mcp: FastMCP = FastMCP(
|
||||
name="nearle-catalogue",
|
||||
instructions=INSTRUCTIONS,
|
||||
version="1.0.0",
|
||||
middleware=[_AuthMiddleware()],
|
||||
)
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Helpers
|
||||
# ---------------------------------------------------------------------------
|
||||
def _guard(what: str, fn, *args, **kwargs):
|
||||
"""
|
||||
Run a service call, turning failures into a ToolError the client can read.
|
||||
|
||||
Nearly every tool here reads Postgres, so the common failure is the database
|
||||
being unreachable. Left to propagate, that surfaces to the model as an
|
||||
unexplained error; named, the model can say what went wrong instead of
|
||||
retrying or inventing an answer.
|
||||
"""
|
||||
try:
|
||||
return fn(*args, **kwargs)
|
||||
except Exception as exc: # noqa: BLE001 - deliberately broad, see docstring
|
||||
logger.exception("MCP tool %s failed", what)
|
||||
raise ToolError(f"{what} failed: {exc}") from exc
|
||||
|
||||
|
||||
def _clip(n: Optional[int], default: int, maximum: int) -> int:
|
||||
"""Keep result counts sane - a model will happily ask for 10000 rows."""
|
||||
if n is None:
|
||||
return default
|
||||
return max(1, min(int(n), maximum))
|
||||
|
||||
|
||||
def _slim(product: Dict[str, Any]) -> Dict[str, Any]:
|
||||
"""
|
||||
Drop the fields a model has no use for.
|
||||
|
||||
Embeddings especially: a 384-float vector per product would swamp the
|
||||
context window and says nothing a model can reason about.
|
||||
"""
|
||||
return {
|
||||
k: v
|
||||
for k, v in product.items()
|
||||
if k not in {"embedding", "embedding_text", "distance"} and v not in (None, "", [], {})
|
||||
}
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Catalog
|
||||
# ---------------------------------------------------------------------------
|
||||
@mcp.tool(
|
||||
annotations={"readOnlyHint": True},
|
||||
tags={"catalog"},
|
||||
)
|
||||
def search_products(
|
||||
query: str,
|
||||
brand: Optional[str] = None,
|
||||
category: Optional[str] = None,
|
||||
limit: Optional[int] = 10,
|
||||
) -> List[Dict[str, Any]]:
|
||||
"""Search the catalog by meaning, not keywords.
|
||||
|
||||
The main entry point: use this to turn a description ("sugar-free biscuits",
|
||||
"something to replace butter") into concrete products. Returns each match
|
||||
with its brand and image_id, which the nutrition, store and recommendation
|
||||
tools take as input.
|
||||
|
||||
Args:
|
||||
query: What to look for, in natural language.
|
||||
brand: Restrict to one brand. Omit to search everything.
|
||||
category: Restrict to one category, e.g. "Dairy".
|
||||
limit: How many products to return (1-50).
|
||||
"""
|
||||
from app.services.rag_service import retrieve
|
||||
|
||||
top_k = _clip(limit, 10, 50)
|
||||
found = _guard("search_products", retrieve, query, brand=brand, top_k=top_k, category=category)
|
||||
return [
|
||||
_slim(
|
||||
{
|
||||
"brand": p.brand,
|
||||
"image_id": p.image_id,
|
||||
"title": p.title,
|
||||
"category": p.category,
|
||||
"description": p.description,
|
||||
"price_range": p.price_range,
|
||||
"selling_price": p.final_selling_price or p.selling_price,
|
||||
"size_variants": p.size_variants,
|
||||
"barcode": p.barcode,
|
||||
}
|
||||
)
|
||||
for p in found
|
||||
]
|
||||
|
||||
|
||||
@mcp.tool(annotations={"readOnlyHint": True}, tags={"catalog"})
|
||||
def list_brands() -> List[str]:
|
||||
"""List every brand in the catalog.
|
||||
|
||||
Useful for grounding a brand name before passing it to another tool, since
|
||||
brand names must match the catalog's spelling.
|
||||
"""
|
||||
from app.services.vector_store import list_available_brands
|
||||
|
||||
return _guard("list_brands", list_available_brands)
|
||||
|
||||
|
||||
@mcp.tool(annotations={"readOnlyHint": True}, tags={"catalog"})
|
||||
def list_categories(brand: str) -> List[str]:
|
||||
"""List the product categories a brand sells.
|
||||
|
||||
Args:
|
||||
brand: Brand name, as returned by list_brands.
|
||||
"""
|
||||
from app.services.vector_store import list_categories_for_brand
|
||||
|
||||
return _guard("list_categories", list_categories_for_brand, brand)
|
||||
|
||||
|
||||
@mcp.tool(annotations={"readOnlyHint": True}, tags={"catalog"})
|
||||
def list_brand_products(
|
||||
brand: str,
|
||||
category: Optional[str] = None,
|
||||
limit: Optional[int] = 25,
|
||||
) -> List[Dict[str, Any]]:
|
||||
"""List a brand's products, newest catalog entries first.
|
||||
|
||||
Browsing, as opposed to search_products' semantic matching. Prefer
|
||||
search_products when the user described what they want rather than naming a
|
||||
brand outright.
|
||||
|
||||
Args:
|
||||
brand: Brand name, as returned by list_brands.
|
||||
category: Restrict to one category.
|
||||
limit: How many products to return (1-100).
|
||||
"""
|
||||
from app.services.vector_store import get_products_by_brand
|
||||
|
||||
rows = _guard(
|
||||
"list_brand_products",
|
||||
get_products_by_brand,
|
||||
brand,
|
||||
limit=_clip(limit, 25, 100),
|
||||
category=category,
|
||||
)
|
||||
return [_slim(r) for r in rows]
|
||||
|
||||
|
||||
@mcp.tool(annotations={"readOnlyHint": True}, tags={"catalog"})
|
||||
def get_product(brand: str, image_id: str) -> Dict[str, Any]:
|
||||
"""Get the full catalog record for one product.
|
||||
|
||||
Args:
|
||||
brand: Brand name.
|
||||
image_id: The product key from search_products or list_brand_products.
|
||||
"""
|
||||
from app.services.vector_store import get_product_by_image_id
|
||||
|
||||
row = _guard("get_product", get_product_by_image_id, brand, image_id)
|
||||
if not row:
|
||||
raise ToolError(
|
||||
f"No product {image_id!r} for brand {brand!r}. "
|
||||
"Check the pair with search_products."
|
||||
)
|
||||
return _slim(row)
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Nutrition
|
||||
# ---------------------------------------------------------------------------
|
||||
@mcp.tool(annotations={"readOnlyHint": True}, tags={"nutrition"})
|
||||
def get_nutrition(brand: str, image_id: str) -> Dict[str, Any]:
|
||||
"""Get nutrition facts, health score and dietary flags for one product.
|
||||
|
||||
Covers macros and key micronutrients per 100g, a 0-100 health score, diet
|
||||
tags (vegetarian, high-protein, ...) and declared allergens.
|
||||
|
||||
Args:
|
||||
brand: Brand name.
|
||||
image_id: The product key from search_products.
|
||||
"""
|
||||
from app.services.nutrition_db import get_full_nutrition
|
||||
|
||||
row = _guard("get_nutrition", get_full_nutrition, brand, image_id)
|
||||
if not row:
|
||||
raise ToolError(
|
||||
f"No nutrition data for {brand}/{image_id}. Not every catalog "
|
||||
"product has been enriched yet."
|
||||
)
|
||||
return _slim(row)
|
||||
|
||||
|
||||
@mcp.tool(annotations={"readOnlyHint": True}, tags={"nutrition"})
|
||||
def find_healthier_alternatives(
|
||||
brand: str,
|
||||
image_id: str,
|
||||
limit: Optional[int] = 5,
|
||||
) -> List[Dict[str, Any]]:
|
||||
"""Find comparable products with a better health score.
|
||||
|
||||
Answers "what should I buy instead of this?" - matches are drawn from the
|
||||
same category, so the suggestion is a real substitute rather than merely
|
||||
something healthier.
|
||||
|
||||
Args:
|
||||
brand: Brand name of the product to improve on.
|
||||
image_id: Product key of the product to improve on.
|
||||
limit: How many alternatives to return (1-20).
|
||||
"""
|
||||
from app.services.nutrition_alternatives_service import find_alternatives
|
||||
|
||||
return _guard(
|
||||
"find_healthier_alternatives",
|
||||
find_alternatives,
|
||||
brand,
|
||||
image_id,
|
||||
top_k=_clip(limit, 5, 20),
|
||||
)
|
||||
|
||||
|
||||
@mcp.tool(annotations={"readOnlyHint": True}, tags={"nutrition"})
|
||||
def filter_products_by_nutrition(
|
||||
sort_by: str = "health_score",
|
||||
order: str = "desc",
|
||||
category: Optional[str] = None,
|
||||
diet_tag: Optional[str] = None,
|
||||
exclude_allergen: Optional[str] = None,
|
||||
limit: Optional[int] = 20,
|
||||
) -> List[Dict[str, Any]]:
|
||||
"""Rank and filter products by a nutritional measure.
|
||||
|
||||
The tool for questions shaped like "highest protein snacks", "lowest sugar
|
||||
dairy", or "gluten-free products without nuts".
|
||||
|
||||
Args:
|
||||
sort_by: Field to rank by - health_score, protein_g, total_sugar_g,
|
||||
dietary_fiber_g, sodium_mg, calories_kcal, total_fat_g.
|
||||
order: "desc" for highest first, "asc" for lowest first.
|
||||
category: Restrict to one category.
|
||||
diet_tag: Keep only products carrying this tag, e.g. "High Protein".
|
||||
exclude_allergen: Drop products declaring this allergen, e.g. "Dairy".
|
||||
limit: How many products to return (1-100).
|
||||
"""
|
||||
from app.services.nutrition_db import query_products
|
||||
|
||||
return _guard(
|
||||
"filter_products_by_nutrition",
|
||||
query_products,
|
||||
sort_by=sort_by,
|
||||
order=order,
|
||||
category=category,
|
||||
diet_tag=diet_tag,
|
||||
exclude_allergen=exclude_allergen,
|
||||
limit=_clip(limit, 20, 100),
|
||||
)
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Stores
|
||||
# ---------------------------------------------------------------------------
|
||||
@mcp.tool(annotations={"readOnlyHint": True}, tags={"stores"})
|
||||
def list_stores() -> List[Dict[str, Any]]:
|
||||
"""List the retail stores, with city and tier.
|
||||
|
||||
Call this first for anything store-specific - the other store tools need a
|
||||
store_id from here.
|
||||
"""
|
||||
from app.services.store_db import list_stores as _list
|
||||
|
||||
return _guard("list_stores", _list)
|
||||
|
||||
|
||||
@mcp.tool(annotations={"readOnlyHint": True}, tags={"stores"})
|
||||
def get_store_inventory(
|
||||
store_id: str,
|
||||
category: Optional[str] = None,
|
||||
in_stock_only: bool = False,
|
||||
limit: Optional[int] = 25,
|
||||
) -> List[Dict[str, Any]]:
|
||||
"""List what one store carries, with stock levels and its own pricing.
|
||||
|
||||
Prices vary per store, so this is the tool for "what does it cost at X",
|
||||
not the catalog price_range.
|
||||
|
||||
Args:
|
||||
store_id: Store identifier from list_stores.
|
||||
category: Restrict to one category.
|
||||
in_stock_only: Drop products currently out of stock.
|
||||
limit: How many products to return (1-100).
|
||||
"""
|
||||
from app.services.store_db import get_store_products
|
||||
|
||||
return _guard(
|
||||
"get_store_inventory",
|
||||
get_store_products,
|
||||
store_id,
|
||||
category=category,
|
||||
in_stock_only=in_stock_only,
|
||||
limit=_clip(limit, 25, 100),
|
||||
)
|
||||
|
||||
|
||||
@mcp.tool(annotations={"readOnlyHint": True}, tags={"stores"})
|
||||
def get_store_discounts(store_id: str, limit: Optional[int] = 25) -> List[Dict[str, Any]]:
|
||||
"""List the discounts currently allocated at one store.
|
||||
|
||||
Args:
|
||||
store_id: Store identifier from list_stores.
|
||||
limit: How many discounts to return (1-100).
|
||||
"""
|
||||
from app.services.store_db import get_latest_discounts
|
||||
|
||||
return _guard(
|
||||
"get_store_discounts", get_latest_discounts, store_id, limit=_clip(limit, 25, 100)
|
||||
)
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Analytics
|
||||
# ---------------------------------------------------------------------------
|
||||
@mcp.tool(annotations={"readOnlyHint": True}, tags={"analytics"})
|
||||
def get_trending(
|
||||
window: str = "weekly",
|
||||
scope: str = "overall",
|
||||
scope_value: Optional[str] = None,
|
||||
limit: Optional[int] = 10,
|
||||
) -> List[Dict[str, Any]]:
|
||||
"""List what is selling fastest right now.
|
||||
|
||||
Args:
|
||||
window: Period to measure over - "daily", "weekly" or "monthly".
|
||||
scope: "overall", or "store"/"category" to narrow it.
|
||||
scope_value: The store_id or category name, when scope is not "overall".
|
||||
limit: How many products to return (1-50).
|
||||
"""
|
||||
from app.services.store_db import get_trending as _trending
|
||||
|
||||
return _guard(
|
||||
"get_trending",
|
||||
_trending,
|
||||
window_label=window,
|
||||
scope=scope,
|
||||
scope_value=scope_value,
|
||||
top_k=_clip(limit, 10, 50),
|
||||
)
|
||||
|
||||
|
||||
@mcp.tool(annotations={"readOnlyHint": True}, tags={"analytics"})
|
||||
def get_top_products(
|
||||
by: str = "revenue",
|
||||
limit: Optional[int] = 10,
|
||||
ascending: bool = False,
|
||||
) -> List[Dict[str, Any]]:
|
||||
"""Rank products across every store by a sales measure.
|
||||
|
||||
Args:
|
||||
by: What to rank on - "revenue", "units" or "orders".
|
||||
limit: How many products to return (1-50).
|
||||
ascending: True ranks worst-performing first.
|
||||
"""
|
||||
from app.services.analytics_service import top_products
|
||||
|
||||
return _guard(
|
||||
"get_top_products", top_products, by=by, limit=_clip(limit, 10, 50), ascending=ascending
|
||||
)
|
||||
|
||||
|
||||
@mcp.tool(annotations={"readOnlyHint": True}, tags={"analytics"})
|
||||
def get_store_analytics(store_id: str) -> Dict[str, Any]:
|
||||
"""Get one store's sales dashboard - revenue, orders, top sellers, stock health.
|
||||
|
||||
Args:
|
||||
store_id: Store identifier from list_stores.
|
||||
"""
|
||||
from app.services.analytics_service import store_dashboard
|
||||
|
||||
return _guard("get_store_analytics", store_dashboard, store_id)
|
||||
|
||||
|
||||
@mcp.tool(annotations={"readOnlyHint": True}, tags={"catalog"})
|
||||
def get_recommendations(
|
||||
brand: str,
|
||||
image_id: str,
|
||||
limit: Optional[int] = 5,
|
||||
) -> List[Dict[str, Any]]:
|
||||
"""List products frequently bought together with this one.
|
||||
|
||||
Co-purchase based, so this answers "what else goes in the basket", which is
|
||||
a different question from find_healthier_alternatives' "what instead".
|
||||
|
||||
Args:
|
||||
brand: Brand name.
|
||||
image_id: Product key from search_products.
|
||||
limit: How many recommendations to return (1-20).
|
||||
"""
|
||||
from app.services.recommendation_service import recommend_for_product
|
||||
|
||||
return _guard(
|
||||
"get_recommendations",
|
||||
recommend_for_product,
|
||||
brand,
|
||||
image_id,
|
||||
top_k=_clip(limit, 5, 20),
|
||||
)
|
||||
|
||||
|
||||
def build_http_app(path: str = MCP_PATH):
|
||||
"""The ASGI app to mount on FastAPI. See app/main.py for the lifespan wiring."""
|
||||
return mcp.http_app(path="/", transport="http", stateless_http=True)
|
||||
@@ -42,6 +42,8 @@ actually resolves to real image bytes above a minimum size - this is what
|
||||
stops garbage/placeholder/expired URLs from reaching the S3 upload step
|
||||
and failing there silently.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
from typing import Optional, List
|
||||
import requests
|
||||
from urllib.parse import urlparse
|
||||
|
||||
@@ -1,3 +1,5 @@
|
||||
from __future__ import annotations
|
||||
|
||||
from typing import List, Dict, Any, Optional
|
||||
import json
|
||||
import re
|
||||
|
||||
@@ -9,7 +9,10 @@ import re
|
||||
import time
|
||||
import psycopg
|
||||
|
||||
from app.infrastructure.settings import DATABASE_URL, USE_PGVECTOR, DB_HOST, DB_PORT, DB_NAME, DB_USER, DB_PASSWORD
|
||||
from app.infrastructure.settings import (
|
||||
DATABASE_URL, USE_PGVECTOR, DB_HOST, DB_PORT, DB_NAME, DB_USER, DB_PASSWORD,
|
||||
DB_CONNECT_TIMEOUT_SECONDS,
|
||||
)
|
||||
from app.services.brand_registry import BRAND_ALIASES, resolve_parent_brand
|
||||
|
||||
|
||||
@@ -93,7 +96,15 @@ def _connect() -> Optional[psycopg.Connection]:
|
||||
dbname=DB_NAME,
|
||||
user=DB_USER,
|
||||
password=DB_PASSWORD,
|
||||
autocommit=True
|
||||
autocommit=True,
|
||||
# Without this, a host that DROPS packets rather than refusing them
|
||||
# - a firewall, a wrong DB_HOST - blocks here until the OS gives up,
|
||||
# which is around 130 seconds on Linux. Every caller of _connect()
|
||||
# inherits that: /api/health stops answering, the container's
|
||||
# healthcheck times out, and the platform pulls the service out of
|
||||
# its load balancer. "Database unreachable" then presents as a Bad
|
||||
# Gateway on every route, including ones that never touch the DB.
|
||||
connect_timeout=DB_CONNECT_TIMEOUT_SECONDS,
|
||||
)
|
||||
except Exception as e:
|
||||
logger.error(f"Vector DB connection failed: {e}")
|
||||
|
||||
Reference in New Issue
Block a user