diff --git a/.env.example b/.env.example index 849e92e..c06bcca 100644 --- a/.env.example +++ b/.env.example @@ -157,6 +157,16 @@ DB_PASSWORD=changeme # off entirely, leaving only the startup run and the manual endpoint. BRAND_SYNC_INTERVAL_SECONDS=300 +# --- Active brands (development working set) -------------------------------- +# Comma-separated. BLANK OR UNSET = every brand is active (the production +# default). Setting it narrows the catalog, RAG, search, MCP, analytics, +# nutrition and every Dagster asset to these brands in one place - nothing is +# deleted, the other brand_* tables just stop being discovered. +# Names resolve through the brand aliases, so "Tata" activates +# brand_hindustan_unilever exactly as ingesting Tata products would. +#ACTIVE_BRANDS=Amul,Cadbury,Hindustan Unilever + + USE_S3=true S3_ACCESS_KEY=your-do-spaces-key S3_SECRET_KEY=your-do-spaces-secret diff --git a/.gitignore b/.gitignore index 05969e2..28f1eb5 100644 --- a/.gitignore +++ b/.gitignore @@ -36,3 +36,17 @@ venv/ # OS .DS_Store Thumbs.db + +# Dagster (orchestration layer, development only) +# Run/event history, compute logs and the pickled asset outputs from the +# filesystem IO manager. Regenerated on demand; nothing here is source. +orchestration/.dagster_home/storage/ +orchestration/.dagster_home/history/ +orchestration/.dagster_home/logs/ +orchestration/.dagster_home/schedules/ +orchestration/.dagster_home/*.db +orchestration/.dagster_home/*.db-* +orchestration/.dagster_home/.telemetry/ +# dagster.yaml IS committed - it is the instance configuration. +# .env.orchestration IS committed - it holds no secrets, only the local +# database pin that keeps orchestrated writes off production. diff --git a/README.md b/README.md index f676610..34124b5 100644 --- a/README.md +++ b/README.md @@ -52,6 +52,58 @@ uvicorn app.main:app --reload --port 8000 Then open http://localhost:8000/docs for interactive API docs, or run the frontend (`../frontend/README.md`) to use the React UI. +## Active brands (development working set) + +One setting decides which brands the application and every pipeline work on: + +``` +# backend/.env +ACTIVE_BRANDS=Amul,Cadbury,Hindustan Unilever +``` + +**Blank or unset means every brand is active** - that is the production default +and the way to turn the feature off. + +Nothing is deleted when it is set. The other `brand_*` tables and their +embeddings stay in Postgres untouched; they simply stop being discovered. +Everything downstream inherits it, because the whole app funnels through two +functions in `app/services/vector_store.py` +(`list_available_brands` and `_list_brand_table_suffixes`): + +``` +ACTIVE_BRANDS + | + +-- /api/brands, /api/products, /api/search, /api/suggest + +-- RAG chat and the query_intent brand index + +-- MCP tools + +-- nutrition enrichment, store intelligence, analytics + +-- the boot auto-seed and the 300s brand reconcile + +-- every Dagster asset and partition +``` + +Going from 3 brands to 5, 10 or all of them is an edit to this one line - no +code changes. Catalogs for inactive brands live in +`data/seed_catalogs/archive/`, which the loader still reads, so re-activating a +brand does not require moving files back. + +Names resolve through the brand aliases, so `ACTIVE_BRANDS=Tata` activates +`brand_hindustan_unilever` - the same table an ingest of Tata products targets. + +## Orchestration (Dagster) + +Ingestion, enrichment, embedding and ML training can be run as a Dagster asset +graph, with lineage, retries and run history. It is a **development tool** - +FastAPI still serves every request and nothing in a request path touches it. + +``` +pip install -r requirements-orchestration.txt +DAGSTER_HOME="$(pwd)/orchestration/.dagster_home" dagster dev -m orchestration.definitions -p 3030 +``` + +See `orchestration/README.md`. Note that it pins itself to the **local** +database and refuses to write to a remote one, because `backend/.env` points at +production. + ## Authentication Reads are public; the 18 write/compute endpoints require a credential, enforced diff --git a/app/api/store_schemas.py b/app/api/store_schemas.py index d29fa64..a469d59 100644 --- a/app/api/store_schemas.py +++ b/app/api/store_schemas.py @@ -87,7 +87,13 @@ class SeedResponse(BaseModel): 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.", + description=( + "Subset of models to (re)train. Any of: discount, trending, popularity, " + "forecast, store_performance, purchase_propensity. Omit to train the " + "models that are actually served (discount, trending, popularity) - the " + "other three fit fine but nothing reads their output back, so ask for " + "them by name." + ), ) diff --git a/app/infrastructure/settings.py b/app/infrastructure/settings.py index 1e07cd6..fd64f49 100644 --- a/app/infrastructure/settings.py +++ b/app/infrastructure/settings.py @@ -152,6 +152,23 @@ DATABASE_URL = os.getenv( # 0 disables the loop, leaving the startup run and POST /api/system/brand-sync. BRAND_SYNC_INTERVAL_SECONDS = int(os.getenv("BRAND_SYNC_INTERVAL_SECONDS", "300")) +# --------------------------------------------------------------------------- +# Active brands (development working set) +# --------------------------------------------------------------------------- +# Comma-separated brand names that the application and every pipeline operate +# on. BLANK OR UNSET MEANS EVERY BRAND IS ACTIVE - that is the backwards +# compatible default and the way to switch this feature off again. +# +# Nothing is deleted when this is set: the other brand_* tables and their +# embeddings stay in the database untouched, they simply stop being discovered. +# Going from 3 brands to 5, 10 or all of them is an edit to this one line. +# +# Names are resolved through resolve_parent_brand + _sanitize_name, the same +# two steps that pick a product's storage table, so "Tata" here activates +# brand_hindustan_unilever exactly as ingesting Tata products would. +# See app/services/active_brands.py. +ACTIVE_BRANDS = os.getenv("ACTIVE_BRANDS", "") + # --------------------------------------------------------------------------- # S3 / DigitalOcean Spaces (product image storage) - optional # --------------------------------------------------------------------------- diff --git a/app/intelligence/artifacts/archive/README.md b/app/intelligence/artifacts/archive/README.md new file mode 100644 index 0000000..c3a91b3 --- /dev/null +++ b/app/intelligence/artifacts/archive/README.md @@ -0,0 +1,24 @@ +# Archived model artifacts + +These three models still have **all of their training code, CLI flags and API +options intact**. Only the pre-trained `.joblib` bundles were moved out of the +loaded directory, because nothing in the application ever reads them back. + +| Artifact | Why archived | +|---|---| +| `demand_forecast_model.joblib` | `train_forecast` fits it and writes the `demand_forecast` table, but `store_db.get_latest_demand_forecast()` has **zero callers** - no endpoint, service or MCP tool consumes the forecast. | +| `store_performance_model.joblib` | Referenced only from `ml_training_service.train_store_performance`. No inference consumer. | +| `purchase_propensity_model.joblib` | Referenced only from `ml_training_service.train_purchase_propensity`. No inference consumer. | + +They are excluded from `train_all()`'s **default** set +(`ml_training_service.PRODUCTION_MODELS`), not from the codebase. To rebuild one: + +``` +python scripts/train_ml_models.py --models store_performance +``` + +or `POST /api/admin/store-intelligence/train {"models": ["store_performance"]}`. + +Training writes to `MODEL_ARTIFACTS_DIR` (the parent directory), so a retrain +promotes the model back to production automatically - wire up an endpoint that +reads it first, or it will simply sit there unread again. diff --git a/app/intelligence/artifacts/demand_forecast_model.joblib b/app/intelligence/artifacts/archive/demand_forecast_model.joblib similarity index 100% rename from app/intelligence/artifacts/demand_forecast_model.joblib rename to app/intelligence/artifacts/archive/demand_forecast_model.joblib diff --git a/app/intelligence/artifacts/purchase_propensity_model.joblib b/app/intelligence/artifacts/archive/purchase_propensity_model.joblib similarity index 100% rename from app/intelligence/artifacts/purchase_propensity_model.joblib rename to app/intelligence/artifacts/archive/purchase_propensity_model.joblib diff --git a/app/intelligence/artifacts/store_performance_model.joblib b/app/intelligence/artifacts/archive/store_performance_model.joblib similarity index 100% rename from app/intelligence/artifacts/store_performance_model.joblib rename to app/intelligence/artifacts/archive/store_performance_model.joblib diff --git a/app/intelligence/artifacts/discount_model.joblib b/app/intelligence/artifacts/discount_model.joblib index a224a05..f119073 100644 Binary files a/app/intelligence/artifacts/discount_model.joblib and b/app/intelligence/artifacts/discount_model.joblib differ diff --git a/app/intelligence/artifacts/popularity_model.joblib b/app/intelligence/artifacts/popularity_model.joblib index 41fa8c9..acb7b85 100644 Binary files a/app/intelligence/artifacts/popularity_model.joblib and b/app/intelligence/artifacts/popularity_model.joblib differ diff --git a/app/intelligence/artifacts/trending_model.joblib b/app/intelligence/artifacts/trending_model.joblib index ab401fa..e172df6 100644 Binary files a/app/intelligence/artifacts/trending_model.joblib and b/app/intelligence/artifacts/trending_model.joblib differ diff --git a/app/services/active_brands.py b/app/services/active_brands.py new file mode 100644 index 0000000..9e0de26 --- /dev/null +++ b/app/services/active_brands.py @@ -0,0 +1,130 @@ +"""The single source of truth for which brands the application works on. + +WHY THIS EXISTS +--------------- +The dev dataset grew to 31 seed catalogs and 16 `brand_*` tables. Every +"all brands" query fans out across every table, boot parses ~19MB of seed +JSON, and 14 of those seed files have no table yet - so the next auto-seed +would silently create 14 more. That is far more than an 8GB dev machine +needs to exercise the features. + +Rather than delete data, ONE setting decides which brands participate: + + ACTIVE_BRANDS=Amul,Cadbury,Hindustan Unilever + +and every read path, pipeline, RAG query, MCP tool and Dagster asset +inherits it, because they all funnel through `vector_store`'s two brand +discovery functions (`list_available_brands` and `_list_brand_table_suffixes`) +plus `brand_sync.load_seed_catalogs`. Going back to 5, 10 or 25 brands is a +one-line `.env` change, not a code edit. + +THE EMPTY-MEANS-ALL CONTRACT +---------------------------- +An unset or blank `ACTIVE_BRANDS` disables filtering entirely, so this module +is a no-op on any deployment that does not opt in. That is deliberate: +`.env.production` leaves it unset, so production keeps serving every brand +while local development runs on three. It also means the feature can be +switched off wholesale if it ever gets in the way. + +NAMES ARE RESOLVED, NOT MATCHED +------------------------------- +Configured names go through `resolve_parent_brand` -> `_sanitize_name`, the +same two steps that decide which table a product is stored in. So +`ACTIVE_BRANDS=Tata` activates `brand_hindustan_unilever` (the "hul tata tea" +alias claims it), exactly like an ingest of that brand would. Comparing raw +strings here would have let the config and the storage layer disagree. +""" +from __future__ import annotations + +from typing import FrozenSet, Iterable, List, Optional + +from app.infrastructure.settings import ACTIVE_BRANDS as _RAW_ACTIVE_BRANDS +from app.services.brand_registry import resolve_parent_brand + +# Parsed once. `_CACHE` holds `(suffixes, display_names)`; `None` in slot 0 +# means "no filtering configured", which is different from "an empty set of +# active brands" - the latter would hide the entire catalog. +_CACHE: Optional[tuple] = None + + +def _sanitize(name: str) -> str: + """Brand name -> table suffix. + + Imported lazily from vector_store because vector_store imports THIS module + at the top level; a top-level import here would be a cycle. query_intent + already breaks the identical cycle the same way. + """ + from app.services.vector_store import _sanitize_name + + return _sanitize_name(name) + + +def _parse(raw: str) -> tuple: + names = [part.strip() for part in (raw or "").split(",")] + names = [n for n in names if n] + if not names: + return (None, []) + + suffixes = [] + display = [] + for name in names: + suffix = _sanitize(resolve_parent_brand(name)) + if not suffix or suffix in suffixes: + continue + suffixes.append(suffix) + display.append(name) + return (frozenset(suffixes), display) + + +def _load() -> tuple: + global _CACHE + if _CACHE is None: + _CACHE = _parse(_RAW_ACTIVE_BRANDS) + return _CACHE + + +def invalidate() -> None: + """Drop the parsed config. For tests that monkeypatch the raw setting.""" + global _CACHE + _CACHE = None + + +def filtering_enabled() -> bool: + """True when ACTIVE_BRANDS is set to a non-empty list.""" + return _load()[0] is not None + + +def active_brand_suffixes() -> Optional[FrozenSet[str]]: + """Table suffixes of the active brands, or None when filtering is off.""" + return _load()[0] + + +def is_active_suffix(suffix: str) -> bool: + """Whether a `brand_` table participates. True for all when off.""" + active = _load()[0] + return True if active is None else (suffix or "").lower() in active + + +def filter_suffixes(suffixes: Iterable[str]) -> List[str]: + """Keep only the active suffixes, preserving the caller's order.""" + active = _load()[0] + if active is None: + return list(suffixes) + return [s for s in suffixes if (s or "").lower() in active] + + +def is_active_brand(brand: str) -> bool: + """Whether a brand NAME (in any alias form) resolves to an active table.""" + active = _load()[0] + if active is None: + return True + return _sanitize(resolve_parent_brand(brand)) in active + + +def active_display_names() -> List[str]: + """The configured names, verbatim, for logs and Dagster partition keys. + + Empty when filtering is off - callers that need the real brand list in + that case should ask `vector_store.list_available_brands()` instead. + """ + return list(_load()[1]) diff --git a/app/services/brand_sync.py b/app/services/brand_sync.py index 3c3e84a..3ca9d32 100644 --- a/app/services/brand_sync.py +++ b/app/services/brand_sync.py @@ -29,6 +29,11 @@ from pathlib import Path from typing import Any, Dict, List, Optional from app.services.brand_registry import resolve_parent_brand +from app.services.active_brands import ( + active_brand_suffixes, + is_active_brand, + is_active_suffix, +) from app.services.vector_store import ( _connect, _list_brand_table_suffixes, @@ -45,6 +50,28 @@ logger = logging.getLogger(__name__) # app/services/brand_sync.py -> parents[2] is backend/ SEED_DIR = Path(__file__).resolve().parents[2] / "data" / "seed_catalogs" +# Catalogs for brands outside ACTIVE_BRANDS live in this subdirectory. It is a +# plain subfolder rather than a separate tree so the two stay side by side, and +# it is NOT matched by SEED_DIR.glob("*.json") - which is what keeps the boot +# auto-seed from parsing ~19MB and creating a table for every archived brand. +ARCHIVE_DIR_NAME = "archive" + + +def seed_catalog_paths(seed_dir: Path = SEED_DIR) -> List[Path]: + """Every seed catalog, active directory first, then the archive. + + The archive is searched too, deliberately: re-activating a brand must be a + one-line ACTIVE_BRANDS change, so where a file physically sits is + presentation, not policy. Filtering by brand happens in load_seed_catalogs. + """ + if not seed_dir.exists(): + return [] + paths = sorted(seed_dir.glob("*.json")) + archive = seed_dir / ARCHIVE_DIR_NAME + if archive.is_dir(): + paths.extend(sorted(archive.glob("*.json"))) + return paths + # The column list written by upsert_brand_products, minus `embedding` (handled # separately because it is huge and optional). Keeping these in sync is what # makes an exported file round-trip back through the seeder without losing @@ -95,7 +122,7 @@ def index_seed_files(seed_dir: Path = SEED_DIR) -> Dict[str, List[Path]]: index: Dict[str, List[Path]] = defaultdict(list) if not seed_dir.exists(): return {} - for path in sorted(seed_dir.glob("*.json")): + for path in seed_catalog_paths(seed_dir): data = _read_catalog(path) if data is None: continue @@ -267,25 +294,69 @@ def load_seed_catalogs(seed_dir: Path = SEED_DIR, logger.error("Seed directory not found: %s", seed_dir) return {} - files = sorted(seed_dir.glob("*.json")) + files = seed_catalog_paths(seed_dir) if only: wanted = [w.lower() for w in only] files = [f for f in files if any(w in f.name.lower() for w in wanted)] + # `only` is an explicit request for named files, so it wins over the + # ACTIVE_BRANDS filter - that is how an archived brand can still be + # re-ingested deliberately without first editing the config. + respect_active = not only + brand_products: Dict[str, List[Dict[str, Any]]] = defaultdict(list) + skipped_inactive = 0 for path in files: data = _read_catalog(path) if data is None: logger.warning("Skipping %s - no brand/products found", path.name) continue resolved = resolve_parent_brand(data["brand"]) + if respect_active and not is_active_brand(resolved): + skipped_inactive += 1 + continue brand_products[resolved].extend(data["products"]) logger.info("Read %d products from %s -> resolved brand '%s'", len(data["products"]), path.name, resolved) + if skipped_inactive: + logger.info("Skipped %d seed catalog(s) outside ACTIVE_BRANDS", skipped_inactive) + return dict(brand_products) +def load_brand_products(brand: str) -> List[Dict[str, Any]]: + """Every seed product belonging to `brand`, resolved properly. + + Use this instead of `load_seed_catalogs(only=[brand])` when you have a + BRAND NAME. `only=` is a case-insensitive substring match on the FILE NAME, + which is the right thing for `seed_sample_data.py --only amul` but the + wrong thing for a brand: + + * "Hindustan Unilever" contains a space; the file is + `brand_catalog_hindustan_unilever.json`. The substring never matches + and you silently get zero products - no error, just an empty result. + * Even with the slug, a filename match misses the other files that feed + the same table: `brand_catalog_tata.json` holds 121 Hindustan Unilever + products because the "hul tata tea" alias claims it. + + Going through `index_seed_files()` fixes both: it is keyed by the resolved + table suffix and its value is the full list of files feeding that table. + Archived catalogs are included, so an explicitly requested brand is found + whether or not it is currently active. + """ + index = index_seed_files() + products: List[Dict[str, Any]] = [] + for path in index.get(brand_slug(brand), []): + data = _read_catalog(path) + if data is None: + continue + products.extend(data["products"]) + logger.info("Read %d products from %s for brand '%s'", + len(data["products"]), path.name, brand) + return products + + def seed_brands(brand_products: Dict[str, List[Dict[str, Any]]], cleanup: bool = True) -> int: """Upsert grouped products into their brand tables. Returns the row count.""" total = 0 @@ -339,7 +410,7 @@ def _detect_collisions(seed_dir: Path = SEED_DIR) -> List[Dict[str, Any]]: by_slug: Dict[str, set] = defaultdict(set) if not seed_dir.exists(): return [] - for path in sorted(seed_dir.glob("*.json")): + for path in seed_catalog_paths(seed_dir): data = _read_catalog(path) if data is None: continue @@ -354,6 +425,9 @@ def _detect_collisions(seed_dir: Path = SEED_DIR) -> List[Dict[str, Any]]: def reconcile_brand_catalogs(dry_run: bool = False) -> Dict[str, Any]: """Repair the table <-> seed-file correspondence in both directions. + Scoped to ACTIVE_BRANDS when that is set: an archived brand is neither + exported nor seeded, in either direction. + Idempotent and non-destructive: * A populated table with no seed file gets one exported. @@ -373,6 +447,22 @@ def reconcile_brand_catalogs(dry_run: bool = False) -> Dict[str, Any]: file_index = index_seed_files() collisions = _detect_collisions() + # BOTH SIDES MUST BE NARROWED TO THE ACTIVE BRANDS, OR NEITHER. + # + # _db_brand_counts() goes through _list_brand_table_suffixes(), so under + # ACTIVE_BRANDS it only sees the active tables. index_seed_files() reads the + # archive directory too, deliberately, so that re-activating a brand needs + # only a config change. + # + # Left mismatched, every archived brand looks like "a seed file whose table + # is empty" and lands in `to_seed` - so the boot reconcile would re-seed all + # 27 archived catalogs on the next restart, recreate their tables, and + # silently undo the archiving. Verified: a dry run reported exactly that. + if active_brand_suffixes() is not None: + file_index = { + slug: paths for slug, paths in file_index.items() if is_active_suffix(slug) + } + to_export = sorted( suffix for suffix, count in db_counts.items() if count > 0 and suffix not in file_index diff --git a/app/services/ml_training_service.py b/app/services/ml_training_service.py index 79953dd..dbbb25c 100644 --- a/app/services/ml_training_service.py +++ b/app/services/ml_training_service.py @@ -24,8 +24,25 @@ from app.services.discount_service import predict_discounts_for_store logger = logging.getLogger(__name__) +# Every model this module knows how to train. Still the full set accepted by +# scripts/train_ml_models.py and POST /api/admin/store-intelligence/train. ALL_MODELS = ["discount", "trending", "popularity", "forecast", "store_performance", "purchase_propensity"] +# What train_all() trains when the caller does not name a subset. +# +# `forecast`, `store_performance` and `purchase_propensity` are omitted +# deliberately: all three fit cleanly, but nothing reads their output back. +# store_db.get_latest_demand_forecast() has zero callers, and neither of the +# other two has an inference consumer anywhere in the app or the MCP tools. +# Training them by default spent roughly half of every retrain, and ~5MB of +# artifacts, producing bundles no request would ever load. +# +# Their training code, CLI flags and API options are untouched - ask for one by +# name (`--models store_performance`) and it trains and is written to +# MODEL_ARTIFACTS_DIR exactly as before. Move a model back into this list once +# something actually serves it. +PRODUCTION_MODELS = ["discount", "trending", "popularity"] + def _load_common_frames(): order_items = store_db.get_order_items_df() @@ -144,7 +161,11 @@ def train_purchase_propensity(orders: pd.DataFrame) -> Dict: def train_all(models: Optional[List[str]] = None) -> Dict[str, Dict]: - models = models or ALL_MODELS + """Train `models`, defaulting to the models that are actually served. + + Pass an explicit list (including any of ALL_MODELS) to train more. + """ + models = models or PRODUCTION_MODELS order_items, orders, store_products = _load_common_frames() results: Dict[str, Dict] = {} diff --git a/app/services/query_intent.py b/app/services/query_intent.py index c41337e..49ed89d 100644 --- a/app/services/query_intent.py +++ b/app/services/query_intent.py @@ -205,8 +205,16 @@ def _brand_index() -> tuple[list[str], Dict[str, str]]: if _MENTION_CACHE["map"] is not None and now - _MENTION_CACHE["at"] < BRAND_MENTION_TTL_SECONDS: return _MENTION_CACHE["known"], _MENTION_CACHE["map"] - known = list(KNOWN_BRANDS) - mapping = dict(BRAND_SEARCH_MAP) + # KNOWN_BRANDS/BRAND_SEARCH_MAP are a hand-tuned static table of 16 brands. + # Under ACTIVE_BRANDS they must be narrowed too, not just added to: without + # this a query for an archived brand ("Nestle") still parses as a brand + # mention, gets routed to brand_catalog mode, and returns zero rows - the + # same silent wrong-result failure mode as fuzzy category detection. Dropped + # here, it is treated as ordinary search text instead. + from app.services.active_brands import is_active_brand + + known = [b for b in KNOWN_BRANDS if is_active_brand(b)] + mapping = {k: v for k, v in BRAND_SEARCH_MAP.items() if is_active_brand(v)} try: from app.services.vector_store import list_available_brands diff --git a/app/services/vector_store.py b/app/services/vector_store.py index 37098e3..42bab06 100644 --- a/app/services/vector_store.py +++ b/app/services/vector_store.py @@ -14,6 +14,7 @@ from app.infrastructure.settings import ( DB_CONNECT_TIMEOUT_SECONDS, ) from app.services.brand_registry import BRAND_ALIASES, resolve_parent_brand +from app.services.active_brands import filter_suffixes, is_active_suffix from app.services.s3_service import s3_service @@ -562,6 +563,12 @@ def list_available_brands() -> List[str]: # Extract brand name from table name (e.g., brand_britannia -> britannia) table_name = table[0] suffix = table_name[len("brand_"):].lower() if table_name.startswith("brand_") else table_name.lower() + # ACTIVE_BRANDS narrows the whole app here: this list feeds + # /api/brands, the MCP list_brands tool, rag_service, + # query_intent's brand index, nutrition enrichment, store_db + # and get_products_all_brands. The table itself is untouched. + if not is_active_suffix(suffix): + continue # Prefer the original display name from the brand map if available brands.append(display_name_for_suffix(suffix, brand_map)) except Exception as e: @@ -1220,12 +1227,23 @@ def _table_exists(cur, table_name: str) -> bool: return bool(row and row[0]) -def _list_brand_table_suffixes(cur) -> List[str]: - """Return brand-table suffixes (e.g. 'parle' from 'brand_parle') for every brand table.""" +def _list_brand_table_suffixes(cur, *, include_inactive: bool = False) -> List[str]: + """Return brand-table suffixes (e.g. 'parle' from 'brand_parle') for every brand table. + + Narrowed to ACTIVE_BRANDS by default. This is the shared chokepoint for + get_brand_overview, semantic_search, text_search, lexical_search, + list_all_categories and brand_sync._db_brand_counts, which is why the + active-brand setting reaches all of them without touching any of them. + + `include_inactive=True` is the escape hatch for callers that address an + archived brand deliberately - an explicit re-ingest or a backfill - so + narrowing the catalog never means losing the ability to maintain it. + """ cur.execute( """ SELECT table_name FROM information_schema.tables WHERE table_schema = 'public' AND table_name LIKE 'brand_%' """ ) - return [row[0][len("brand_"):] for row in cur.fetchall()] + suffixes = [row[0][len("brand_"):] for row in cur.fetchall()] + return suffixes if include_inactive else filter_suffixes(suffixes) diff --git a/data/brand_catalog.db b/data/brand_catalog.db deleted file mode 100644 index e69de29..0000000 diff --git a/data/seed_catalogs/archive/README.md b/data/seed_catalogs/archive/README.md new file mode 100644 index 0000000..68a68e7 --- /dev/null +++ b/data/seed_catalogs/archive/README.md @@ -0,0 +1,24 @@ +# Archived seed catalogs + +These brands are **not deleted** - they are inactive. + +The application works on the brands listed in `ACTIVE_BRANDS` (see +`backend/.env` and `app/services/active_brands.py`). Catalogs for every other +brand live here so the boot auto-seed does not parse ~19MB of JSON and does not +create a `brand_*` table for a brand nobody is developing against. + +`brand_sync.seed_catalog_paths()` reads this directory too, so **re-activating a +brand needs only a config change** - the file does not have to be moved back: + +``` +ACTIVE_BRANDS=Amul,Cadbury,Hindustan Unilever,Nestle +``` + +...then restart the API (or `POST /api/system/init`) and Nestle is seeded and +served again. + +Their `brand_*` tables and embeddings were never dropped; they are still in +Postgres, just filtered out of discovery. + +Moving a file back into the parent directory is optional tidiness, not a +requirement. diff --git a/data/seed_catalogs/brand_catalog_aachi.json b/data/seed_catalogs/archive/brand_catalog_aachi.json similarity index 100% rename from data/seed_catalogs/brand_catalog_aachi.json rename to data/seed_catalogs/archive/brand_catalog_aachi.json diff --git a/data/seed_catalogs/brand_catalog_anil.json b/data/seed_catalogs/archive/brand_catalog_anil.json similarity index 100% rename from data/seed_catalogs/brand_catalog_anil.json rename to data/seed_catalogs/archive/brand_catalog_anil.json diff --git a/data/seed_catalogs/brand_catalog_bikaji.json b/data/seed_catalogs/archive/brand_catalog_bikaji.json similarity index 100% rename from data/seed_catalogs/brand_catalog_bikaji.json rename to data/seed_catalogs/archive/brand_catalog_bikaji.json diff --git a/data/seed_catalogs/brand_catalog_britannia.json b/data/seed_catalogs/archive/brand_catalog_britannia.json similarity index 100% rename from data/seed_catalogs/brand_catalog_britannia.json rename to data/seed_catalogs/archive/brand_catalog_britannia.json diff --git a/data/seed_catalogs/brand_catalog_cavinkare.json b/data/seed_catalogs/archive/brand_catalog_cavinkare.json similarity index 100% rename from data/seed_catalogs/brand_catalog_cavinkare.json rename to data/seed_catalogs/archive/brand_catalog_cavinkare.json diff --git a/data/seed_catalogs/brand_catalog_coca_cola.json b/data/seed_catalogs/archive/brand_catalog_coca_cola.json similarity index 100% rename from data/seed_catalogs/brand_catalog_coca_cola.json rename to data/seed_catalogs/archive/brand_catalog_coca_cola.json diff --git a/data/seed_catalogs/brand_catalog_colgate_palmolive.json b/data/seed_catalogs/archive/brand_catalog_colgate_palmolive.json similarity index 100% rename from data/seed_catalogs/brand_catalog_colgate_palmolive.json rename to data/seed_catalogs/archive/brand_catalog_colgate_palmolive.json diff --git a/data/seed_catalogs/brand_catalog_dabur.json b/data/seed_catalogs/archive/brand_catalog_dabur.json similarity index 100% rename from data/seed_catalogs/brand_catalog_dabur.json rename to data/seed_catalogs/archive/brand_catalog_dabur.json diff --git a/data/seed_catalogs/brand_catalog_everest.json b/data/seed_catalogs/archive/brand_catalog_everest.json similarity index 100% rename from data/seed_catalogs/brand_catalog_everest.json rename to data/seed_catalogs/archive/brand_catalog_everest.json diff --git a/data/seed_catalogs/brand_catalog_fortune.json b/data/seed_catalogs/archive/brand_catalog_fortune.json similarity index 100% rename from data/seed_catalogs/brand_catalog_fortune.json rename to data/seed_catalogs/archive/brand_catalog_fortune.json diff --git a/data/seed_catalogs/brand_catalog_godrej.json b/data/seed_catalogs/archive/brand_catalog_godrej.json similarity index 100% rename from data/seed_catalogs/brand_catalog_godrej.json rename to data/seed_catalogs/archive/brand_catalog_godrej.json diff --git a/data/seed_catalogs/brand_catalog_grb.json b/data/seed_catalogs/archive/brand_catalog_grb.json similarity index 100% rename from data/seed_catalogs/brand_catalog_grb.json rename to data/seed_catalogs/archive/brand_catalog_grb.json diff --git a/data/seed_catalogs/brand_catalog_haldirams.json b/data/seed_catalogs/archive/brand_catalog_haldirams.json similarity index 100% rename from data/seed_catalogs/brand_catalog_haldirams.json rename to data/seed_catalogs/archive/brand_catalog_haldirams.json diff --git a/data/seed_catalogs/brand_catalog_idhayam.json b/data/seed_catalogs/archive/brand_catalog_idhayam.json similarity index 100% rename from data/seed_catalogs/brand_catalog_idhayam.json rename to data/seed_catalogs/archive/brand_catalog_idhayam.json diff --git a/data/seed_catalogs/brand_catalog_itc.json b/data/seed_catalogs/archive/brand_catalog_itc.json similarity index 100% rename from data/seed_catalogs/brand_catalog_itc.json rename to data/seed_catalogs/archive/brand_catalog_itc.json diff --git a/data/seed_catalogs/brand_catalog_kaleesuwari.json b/data/seed_catalogs/archive/brand_catalog_kaleesuwari.json similarity index 100% rename from data/seed_catalogs/brand_catalog_kaleesuwari.json rename to data/seed_catalogs/archive/brand_catalog_kaleesuwari.json diff --git a/data/seed_catalogs/brand_catalog_lion_dates.json b/data/seed_catalogs/archive/brand_catalog_lion_dates.json similarity index 100% rename from data/seed_catalogs/brand_catalog_lion_dates.json rename to data/seed_catalogs/archive/brand_catalog_lion_dates.json diff --git a/app/data/seed_catalogs/brand_catalog_lion_dates.json b/data/seed_catalogs/archive/brand_catalog_lion_dates.stray-copy.json similarity index 100% rename from app/data/seed_catalogs/brand_catalog_lion_dates.json rename to data/seed_catalogs/archive/brand_catalog_lion_dates.stray-copy.json diff --git a/data/seed_catalogs/brand_catalog_manna.json b/data/seed_catalogs/archive/brand_catalog_manna.json similarity index 100% rename from data/seed_catalogs/brand_catalog_manna.json rename to data/seed_catalogs/archive/brand_catalog_manna.json diff --git a/data/seed_catalogs/brand_catalog_marico.json b/data/seed_catalogs/archive/brand_catalog_marico.json similarity index 100% rename from data/seed_catalogs/brand_catalog_marico.json rename to data/seed_catalogs/archive/brand_catalog_marico.json diff --git a/data/seed_catalogs/brand_catalog_mdh.json b/data/seed_catalogs/archive/brand_catalog_mdh.json similarity index 100% rename from data/seed_catalogs/brand_catalog_mdh.json rename to data/seed_catalogs/archive/brand_catalog_mdh.json diff --git a/data/seed_catalogs/brand_catalog_milky_mist.json b/data/seed_catalogs/archive/brand_catalog_milky_mist.json similarity index 100% rename from data/seed_catalogs/brand_catalog_milky_mist.json rename to data/seed_catalogs/archive/brand_catalog_milky_mist.json diff --git a/data/seed_catalogs/brand_catalog_mtr.json b/data/seed_catalogs/archive/brand_catalog_mtr.json similarity index 100% rename from data/seed_catalogs/brand_catalog_mtr.json rename to data/seed_catalogs/archive/brand_catalog_mtr.json diff --git a/data/seed_catalogs/brand_catalog_naga.json b/data/seed_catalogs/archive/brand_catalog_naga.json similarity index 100% rename from data/seed_catalogs/brand_catalog_naga.json rename to data/seed_catalogs/archive/brand_catalog_naga.json diff --git a/data/seed_catalogs/brand_catalog_nestle.json b/data/seed_catalogs/archive/brand_catalog_nestle.json similarity index 100% rename from data/seed_catalogs/brand_catalog_nestle.json rename to data/seed_catalogs/archive/brand_catalog_nestle.json diff --git a/data/seed_catalogs/brand_catalog_p_and_g.json b/data/seed_catalogs/archive/brand_catalog_p_and_g.json similarity index 100% rename from data/seed_catalogs/brand_catalog_p_and_g.json rename to data/seed_catalogs/archive/brand_catalog_p_and_g.json diff --git a/data/seed_catalogs/brand_catalog_parle.json b/data/seed_catalogs/archive/brand_catalog_parle.json similarity index 100% rename from data/seed_catalogs/brand_catalog_parle.json rename to data/seed_catalogs/archive/brand_catalog_parle.json diff --git a/data/seed_catalogs/brand_catalog_pepsico.json b/data/seed_catalogs/archive/brand_catalog_pepsico.json similarity index 100% rename from data/seed_catalogs/brand_catalog_pepsico.json rename to data/seed_catalogs/archive/brand_catalog_pepsico.json diff --git a/orchestration/.dagster_home/dagster.yaml b/orchestration/.dagster_home/dagster.yaml new file mode 100644 index 0000000..ff123ea --- /dev/null +++ b/orchestration/.dagster_home/dagster.yaml @@ -0,0 +1,30 @@ +# Dagster instance config for local development. +# +# Everything here is the SQLite-backed default, written out explicitly so the +# storage location is obvious and stays inside the project rather than landing +# in a home directory. +# +# dagster-postgres is deliberately not used: it requires psycopg2-binary while +# this project is on psycopg3, and putting orchestration run/event data in the +# catalog database would mix operational metadata with product data for no gain +# at this scale. + +telemetry: + enabled: false + +run_coordinator: + module: dagster.core.run_coordinator + class: QueuedRunCoordinator + config: + # One run at a time. The machine is also hosting Postgres, Ollama, the API + # and the Vite dev server; two concurrent runs each importing + # sentence-transformers is what makes an 8GB box swap. + max_concurrent_runs: 1 + +retention: + # Keep run history small. Without this, sqlite grows without bound and the + # UI gets slower every week. + schedule: + purge_after_days: 14 + sensor: + purge_after_days: 7 diff --git a/orchestration/.env.orchestration b/orchestration/.env.orchestration new file mode 100644 index 0000000..10706d4 --- /dev/null +++ b/orchestration/.env.orchestration @@ -0,0 +1,29 @@ +# Loaded by orchestration/definitions.py BEFORE app.infrastructure.settings is +# imported, so these win over backend/.env (python-dotenv does not override a +# variable that is already in the environment, and settings.py snapshots the +# environment once at import time). +# +# WHY THIS FILE EXISTS: backend/.env points DB_HOST at the production database. +# Dagster's ingestion assets write. Without this pin, `dagster dev` would run +# pipeline experiments against the live catalog. +# +# Everything not named here still comes from backend/.env (auth, S3, Ollama, +# embeddings, the feature flags) - this file overrides the write target only. + +DB_HOST=localhost +DB_PORT=5432 +DB_NAME=pgvector +DB_USER=postgres +DB_PASSWORD=changeme + +# The development working set. Kept in step with backend/.env deliberately: +# the orchestrator should build exactly the brands the app serves. +ACTIVE_BRANDS=Amul,Cadbury,Hindustan Unilever + +# Network-touching enrichment stays off for orchestrated runs, matching the +# application defaults. Turn a stage on per-run in the Dagster UI instead of +# flipping it here, so a scheduled run never quietly starts making thousands +# of outbound requests. +ENABLE_BARCODE_LOOKUP=false +ENABLE_SKU_WEB_LOOKUP=false +ENABLE_MANUFACTURER_SITE_LOOKUP=false diff --git a/orchestration/README.md b/orchestration/README.md new file mode 100644 index 0000000..f7fbd3e --- /dev/null +++ b/orchestration/README.md @@ -0,0 +1,262 @@ +# Dagster orchestration layer + +Dagster orchestrates this project's **existing** ingestion, enrichment, +embedding and ML functions. It does not reimplement any of them, and it is +**not** on any user request path. + +``` + USER DAGSTER (dev tool, port 3030) + | | + REACT FRONTEND +-----------+-----------+ + | | | | + FASTAPI :8000 INGESTION ENRICHMENT ML + | | | | + +------+------+------+ | nutrition training + | | | | barcode evaluation +CATALOG RAG/SEARCH ANALYTICS | HSN/GST artifacts + | | | | | | + +------+------+------+ +-----+-----+ | + | | | + POSTGRESQL <------------------------ + writes --------+ + | + pgvector +``` + +FastAPI still answers every request. Dagster prepares the data those requests +read. + +--- + +## Running it + +From `backend/`: + +```bash +pip install -r requirements-orchestration.txt + +# Windows (Git Bash) +DAGSTER_HOME="$(pwd)/orchestration/.dagster_home" \ + dagster dev -m orchestration.definitions -p 3030 +``` + +Then open . + +Run it **from `backend/`** so `app.*` is importable. `definitions.py` also puts +the backend root on `sys.path`, so an IDE run configuration rooted at the repo +works too. + +### The database it writes to + +`backend/.env` points `DB_HOST` at the **production** database. Dagster's +ingestion assets write, so two independent safeguards keep them off it: + +1. `orchestration/.env.orchestration` pins `DB_HOST=localhost`, and + `definitions.py` loads it with `override=True` **before** anything imports + `app.infrastructure.settings` (which snapshots the environment once, at + import time). +2. `config.require_local_database()` re-checks the resolved host at asset + runtime and fails the run if it is not local. + +The second exists because the first is a file that can go missing or be +overridden by a shell variable. If you ever see + +``` +Refusing to run 'catalog_database': ... the resolved database is +31.97.228.132:6054, which is not local. +``` + +then step 1 did not happen. Fix the env file rather than reaching for the +override. + +To target a remote database deliberately: +`ORCHESTRATION_ALLOW_REMOTE_WRITES=true`. + +--- + +## Assets + +Twelve assets in three groups. Everything in `catalog` and `nutrition_data` is +**partitioned by brand** - one partition per name in `ACTIVE_BRANDS`. + +### catalog + +``` +active_brand + | +raw_products brand_sync.load_seed_catalogs(only=[brand]) + | +validated_products title_validator + category_units + product_validator + | +enriched_products sku_service + enrichment.pipeline (barcode, HSN/GST) + | + price_estimator +catalog_database vector_store.upsert_brand_products(cleanup=False) + | +product_embeddings embeddings_service.embed_texts, rows missing a vector + | +vector_index read-back check: every row is searchable +``` + +### nutrition + +``` +catalog_database -> nutrition_data -> nutrition_models +``` + +`nutrition_data` is per brand; `nutrition_models` is not, because "find a +healthier alternative" has to cross brands. + +### ml + +``` +training_dataset -> trained_models -> model_evaluation +``` + +Only the three **served** models are trained by default (discount, trending, +popularity). `forecast`, `store_performance` and `purchase_propensity` fit +fine but nothing reads their output, so they are opt-in by name - see +`app/intelligence/artifacts/archive/README.md`. + +### Which function each asset calls + +| Asset | Existing code it delegates to | +|---|---| +| `active_brand` | `brand_registry.resolve_parent_brand` | +| `raw_products` | `brand_sync.load_seed_catalogs` | +| `validated_products` | `title_validator.validate_and_fix_title`, `category_units.fix_or_reject_size`, `product_validator.validate_catalog` | +| `enriched_products` | `sku_service.resolve_product_sku`, `enrichment.pipeline.run_default_pipeline`, `price_estimator.estimate_price_range_for_size` | +| `catalog_database` | `vector_store.ensure_brand_schema`, `vector_store.upsert_brand_products` | +| `product_embeddings` | `embeddings_service.embed_texts` | +| `nutrition_data` | `nutrition_enrichment_service.enrich_one_product` | +| `nutrition_models` | `nutrition_enrichment_service.train_all_models` | +| `training_dataset` | `store_seed_service.run_seed` | +| `trained_models` | `ml_training_service.train_all` | + +If you find yourself writing business logic in this package, it belongs in +`app/services/` instead - so the API and the orchestrator keep sharing it. + +--- + +## Jobs + +| Job | What it does | +|---|---| +| `catalog_ingestion_job` | brand -> seed intake -> validation -> enrichment -> Postgres | +| `embedding_refresh_job` | embed rows lacking a vector, verify the index | +| `nutrition_enrichment_job` | fetch nutrition, score it, refit the two models | +| `ml_training_job` | rebuild the store dataset, fit the served models, report | + +Four jobs rather than one because they differ in cost and cadence: ingestion is +offline and cheap, embedding drives the transformer, nutrition leaves the +machine, and ML training depends on order history rather than catalog freshness. + +From the command line: + +```bash +dagster asset materialize -m orchestration.definitions \ + --select "active_brand,raw_products,validated_products,enriched_products,catalog_database" \ + --partition Amul +``` + +--- + +## Schedules and the sensor + +**All of them ship STOPPED.** Turn them on in the UI when you want them. + +| | Cadence | | +|---|---|---| +| `daily_catalog_refresh` | 02:00 | catalog first | +| `daily_embedding_refresh` | 03:00 | embed what arrived | +| `daily_nutrition_refresh` | 04:00 | slowest, and the only one making outbound calls | +| `weekly_ml_retrain` | Sun 05:00 | training data is a deterministic simulation, so nightly would repeat itself | +| `seed_catalog_sensor` | every 60s | re-ingest a brand when its seed file changes | + +Shipping them stopped is deliberate. This machine also runs Postgres, Ollama, +the API and Vite; a schedule that started itself the moment `dagster dev` +launched would turn a development tool into a background workload. A test +(`tests/test_orchestration_defs.py`) asserts they stay stopped. + +The sensor is a `stat()` over a handful of files once a minute - no watcher +process, no broker. On its first evaluation it records the current state and +requests nothing, so enabling it does not trigger a full rebuild. + +--- + +## Resource notes (8GB target) + +* `max_concurrent_runs: 1` (dagster.yaml) and `max_concurrent: 2` on the + executor. Each subprocess that touches an asset imports + sentence-transformers; four workers made the machine swap. +* `product_embeddings` only embeds rows where `embedding IS NULL`, in batches + of 32. A re-run costs almost nothing, which is what makes a daily schedule + reasonable. +* Partitioning keeps one brand's rows in memory instead of the whole catalog. +* Network enrichment (`ENABLE_BARCODE_LOOKUP`, `ENABLE_SKU_WEB_LOOKUP`, + `ENABLE_MANUFACTURER_SITE_LOOKUP`) is **off**. Turn a stage on per-run in the + UI rather than in the env file, so a schedule never silently starts making + thousands of outbound requests. +* Run history is purged after 14 days (schedules) / 7 days (sensors) so the + SQLite instance does not grow without bound. + +--- + +## Errors and retries + +`raw_products`, `enriched_products` and `product_embeddings` retry twice with +exponential backoff and jitter. Nothing retries indefinitely, and database +writes do not retry at all - a failed write here is almost never transient, so +retrying only delays a red run someone has to look at. + +Each asset records a row-count funnel as metadata, so the UI shows where rows +were lost: + +``` +raw_products 122 products +validated_products 122 in -> 118 out (3 title fixes, 4 size rejects) +enriched_products 118 rows (118 gained an HSN code) +catalog_database 118 offered -> 118 written +vector_index 122 rows, 122 embedded, 0 missing, searchable +``` + +`vector_index` is a separate asset from `catalog_database` on purpose: "the +upsert reported success" and "the table can answer a similarity query" are +different claims, and only the second is what a user experiences. A green +write beside a red index is exactly the failure that used to present as "the +catalog looks empty". + +--- + +## Docker + +Opt-in only: + +```bash +docker compose --profile orchestration up dagster +``` + +It is **not** in `docker-compose.prod.yml`. The production overlay already +allocates the whole 8GB VPS (ollama 3G, backend 2560M, postgres 1G), so there +is no headroom for a webserver plus daemon. Orchestration is a development +concern here. + +Port 3030, not Dagster's default 3000 - `serve.py` binds `PORTS=3000,8000` and +the frontend nginx also listens on 3000. + +--- + +## Gotchas + +**No `from __future__ import annotations` in this package.** Dagster resolves +the decorated signatures at definition time to validate `context` and infer +asset input types. Under PEP 563/649 the annotations arrive as strings and it +fails with *"Cannot annotate `context` parameter with type +AssetExecutionContext"*. Local Python is 3.14, which defers annotations by +default, so this is not hypothetical. + +**Import order in `definitions.py` is load-bearing.** `app.infrastructure.settings` +reads the whole environment once at import. The dotenv call has to come before +the first `app.*` import or the database pin does nothing. + +**`cleanup=False` in `catalog_database` is load-bearing.** `cleanup=True` +deletes every row in the table that is not in the batch being written, so a +partial or cancelled run would wipe the rest of the brand's catalog. diff --git a/orchestration/__init__.py b/orchestration/__init__.py new file mode 100644 index 0000000..19933cb --- /dev/null +++ b/orchestration/__init__.py @@ -0,0 +1,10 @@ +"""Dagster orchestration layer for the RAG catalog application. + +This package ORCHESTRATES the application's existing functions. It does not +reimplement any of them: every asset here is a thin wrapper that calls into +`app.core.*` / `app.services.*`, records what happened as Dagster metadata, +and lets Dagster own scheduling, retries, run history and lineage. + +FastAPI remains the application layer. Nothing in a user request path goes +through Dagster. +""" diff --git a/orchestration/assets/__init__.py b/orchestration/assets/__init__.py new file mode 100644 index 0000000..1e1f453 --- /dev/null +++ b/orchestration/assets/__init__.py @@ -0,0 +1,3 @@ +from orchestration.assets import catalog, ml, nutrition + +__all__ = ["catalog", "nutrition", "ml"] diff --git a/orchestration/assets/catalog.py b/orchestration/assets/catalog.py new file mode 100644 index 0000000..0eb9fe8 --- /dev/null +++ b/orchestration/assets/catalog.py @@ -0,0 +1,515 @@ +"""Catalog ingestion assets: brand selection -> validated -> enriched -> stored. + +Every asset delegates to a function that already exists and is already tested. +What Dagster adds is the lineage between them, per-brand partitioning, retries +on the flaky steps, and a run history that survives a restart - none of which +the daemon-thread + in-memory job store could give. + +The functions being wrapped are `app.core.store_catalog_pipeline`'s eleven +stages and the services they call, which the store-spreadsheet upload path +already uses. Wrapping rather than reimplementing is what stops the two paths +drifting: a fix to a stage fixes both. +""" +# NOTE: deliberately no `from __future__ import annotations` here. +# Dagster resolves the decorated function signatures at definition time to +# validate the `context` parameter and to infer asset input types. Under +# PEP 563/649 the annotations arrive as strings and that validation fails +# with "Cannot annotate `context` parameter with type AssetExecutionContext". +# Local Python is 3.14, which defers annotations by default, so this is not +# hypothetical. + +import time +from typing import Any, Dict, List + +from dagster import ( + AssetExecutionContext, + Backoff, + Failure, + Jitter, + MetadataValue, + RetryPolicy, + asset, +) + +from orchestration.config import BrandConfig, database_target, require_local_database +from orchestration.partitions import brand_partitions + +# Only the assets that reach the network get retried. A failed DB write is +# rarely transient here (bad schema, bad row), so retrying it just delays a red +# run that a human has to look at anyway. +NETWORK_RETRY = RetryPolicy( + max_retries=2, delay=5, backoff=Backoff.EXPONENTIAL, jitter=Jitter.PLUS_MINUS +) + +GROUP = "catalog" + + +@asset( + group_name=GROUP, + partitions_def=brand_partitions(), + description="The brand this partition builds, resolved through the alias table.", +) +def active_brand(context: AssetExecutionContext) -> str: + """Root of the lineage: one partition per brand in ACTIVE_BRANDS. + + Resolving here rather than inside each downstream asset means the UI shows + which storage table a partition actually targets. "Tata" landing in + brand_hindustan_unilever is a real case in this data set, and it should be + visible up front instead of discovered halfway through a run. + """ + from app.services.brand_registry import resolve_parent_brand + from app.services.vector_store import _sanitize_name + + brand = context.partition_key + parent = resolve_parent_brand(brand) + table = "brand_" + _sanitize_name(parent) + + context.add_output_metadata( + { + "brand": brand, + "resolved_parent": parent, + "storage_table": table, + "database": database_target(), + } + ) + return brand + + +@asset( + group_name=GROUP, + partitions_def=brand_partitions(), + retry_policy=NETWORK_RETRY, + description="Products for this brand, read from its seed catalog(s).", +) +def raw_products( + context: AssetExecutionContext, config: BrandConfig, active_brand: str +) -> List[Dict[str, Any]]: + """Source rows via `brand_sync.load_seed_catalogs(only=...)`. + + Seed files are the default source rather than live discovery because + `catalog_engine.generate_catalog` drives Ollama plus a few hundred image + searches - minutes of CPU and a lot of outbound traffic per brand. Seed + files make the whole DAG runnable offline in seconds, which is what makes + it usable as a development tool. Live discovery stays reachable through the + existing POST /api/catalog/generate and cli/ingest_brand.py. + + `load_brand_products` resolves the brand through the seed-file index rather + than matching the brand name against file names. That matters: the obvious + `load_seed_catalogs(only=[brand])` matches substrings of the FILE NAME, so + "Hindustan Unilever" (a space) never matches + `brand_catalog_hindustan_unilever.json` (an underscore) and the partition + silently ingests nothing. It also picks up every file feeding the brand's + table, including `brand_catalog_tata.json`'s 121 HUL products, and reaches + archived catalogs so an explicitly requested brand is always found. + """ + from app.services.brand_sync import load_brand_products + + products = load_brand_products(active_brand) + available = len(products) + + if config.max_products_per_brand: + products = products[: config.max_products_per_brand] + + if not products: + # Loud, not a warning. A partition that ingests zero rows and reports + # success is the worst outcome here: everything downstream materializes + # green over an empty list, and the brand looks processed when nothing + # happened. This is the exact failure the name-matching bug above + # produced before it was found. + raise Failure( + description=( + f"No seed products found for '{active_brand}'. Expected a " + f"catalog feeding table " + f"'brand_{_slug(active_brand)}' in data/seed_catalogs/ or its " + f"archive/ subdirectory." + ), + metadata={"brand": active_brand, "products": 0}, + ) + + context.add_output_metadata( + { + "brand": active_brand, + "products": len(products), + "products_available": available, + "preview": MetadataValue.md(_preview(products)), + } + ) + return products + + +@asset( + group_name=GROUP, + partitions_def=brand_partitions(), + description="Rows that passed the title, pack-size and confidence gates.", +) +def validated_products( + context: AssetExecutionContext, + active_brand: str, + raw_products: List[Dict[str, Any]], +) -> List[Dict[str, Any]]: + """Stages 3, 4 and 10 of the existing pipeline, in that order. + + `validate_catalog` is the same deterministic gate the upload path runs, so + a row rejected here would have been rejected there too. Rejections are + recorded as metadata rather than raised: a brand where 3 of 120 rows fail + is a data problem worth seeing, not a reason to fail the run and block the + other 117. + """ + from app.services.category_units import fix_or_reject_size + from app.services.product_validator import validate_catalog + from app.services.title_validator import validate_and_fix_title + + started = time.time() + prepared: List[Dict[str, Any]] = [] + title_fixes = 0 + size_rejects = 0 + + for row in raw_products: + row = dict(row) + category = row.get("category") or "" + title = row.get("title") or row.get("product_name") or "" + + fixed_title, changed, _removed = validate_and_fix_title( + title, category, brand=active_brand + ) + if changed: + title_fixes += 1 + row["title"] = fixed_title + + size = row.get("size") or "" + if size: + fixed_size, rejected_size, reason = fix_or_reject_size( + size, category, row.get("product_name") or title + ) + if rejected_size: + size_rejects += 1 + context.log.debug("Dropped %s (%s): %s", title, size, reason) + continue + row["size"] = fixed_size + + prepared.append(row) + + kept, rejected, summary = validate_catalog( + prepared, active_brand, images_checked=False + ) + + context.add_output_metadata( + { + "brand": active_brand, + "rows_in": len(raw_products), + "rows_out": len(kept), + "rejected": len(rejected), + "title_fixes": title_fixes, + "size_rejects": size_rejects, + "validator_summary": MetadataValue.json(_jsonable(summary)), + "duration_s": round(time.time() - started, 2), + } + ) + return kept + + +@asset( + group_name=GROUP, + partitions_def=brand_partitions(), + retry_policy=NETWORK_RETRY, + description="SKU, barcode, HSN/GST and pricing filled in - blanks only.", +) +def enriched_products( + context: AssetExecutionContext, + active_brand: str, + validated_products: List[Dict[str, Any]], +) -> List[Dict[str, Any]]: + """Stages 5 and 7-9, through the enrichment framework that already exists. + + FILL-ONLY-BLANKS is the contract, inherited from the upload pipeline: a + stage writes a field only when it is empty, so re-running over the same + input is a no-op and a brand's own data is never overwritten by a guess. + + The barcode and manufacturer-site lookups stay off unless a run turns them + on - see orchestration/.env.orchestration for why. + """ + import asyncio + + from app.services.enrichment.pipeline import run_default_pipeline + from app.services.price_estimator import estimate_price_range_for_size + from app.services.sku_service import resolve_product_sku + + started = time.time() + rows = [dict(r) for r in validated_products] + + skus_added = 0 + prices_added = 0 + for row in rows: + title = row.get("title") or row.get("product_name") or "" + size = row.get("size") or "" + category = row.get("category") or "" + + if not row.get("product_sku"): + row.update(resolve_product_sku(active_brand, title, size)) + skus_added += 1 + + if not row.get("price_range") and size: + try: + low, high = estimate_price_range_for_size( + size, title, active_brand, category + ) + row["price_range"] = "₹{}-{}".format(low, high) + prices_added += 1 + except Exception as exc: # noqa: BLE001 - enrichment is best effort + context.log.debug("Price estimate failed for %s: %s", title, exc) + + # Barcode + HSN/GST. Each individual stage is already never-raising; this + # only guards against the pipeline itself being unavailable, and degrades + # to "the rows we enriched so far" rather than failing the partition. + try: + rows = asyncio.run(run_default_pipeline(rows, active_brand)) + except Exception as exc: # noqa: BLE001 + context.log.warning( + "Enrichment pipeline degraded for %s (%s) - continuing with the " + "rows enriched so far.", + active_brand, + exc, + ) + + context.add_output_metadata( + { + "brand": active_brand, + "rows": len(rows), + "skus_assigned": skus_added, + "price_bands_estimated": prices_added, + "rows_with_barcode": sum(1 for r in rows if r.get("barcode")), + "rows_with_hsn": sum(1 for r in rows if r.get("hsn_code")), + "duration_s": round(time.time() - started, 2), + } + ) + return rows + + +@asset( + group_name=GROUP, + partitions_def=brand_partitions(), + description="Rows upserted into the brand's Postgres table.", +) +def catalog_database( + context: AssetExecutionContext, + active_brand: str, + enriched_products: List[Dict[str, Any]], +) -> int: + """Stage 11's persistence half: `upsert_brand_products(..., cleanup=False)`. + + TWO THINGS HERE ARE LOAD-BEARING. + + `cleanup=False`: cleanup=True deletes every row in the table that is not in + the batch being written, so a partial or filtered batch would wipe the rest + of the brand's catalog. A test already asserts this for the upload path; + the same reasoning applies with more force here, where a run can be + retried or cancelled halfway through. + + The guard: this asset writes, so it refuses a non-local database. backend/.env + points at production. See orchestration/config.py. + """ + from app.services.vector_store import ensure_brand_schema, upsert_brand_products + + target = require_local_database("catalog_database") + + if not enriched_products: + context.log.warning("Nothing to store for %s.", active_brand) + context.add_output_metadata({"brand": active_brand, "rows_written": 0}) + return 0 + + table = ensure_brand_schema(active_brand) + written = upsert_brand_products(active_brand, enriched_products, cleanup=False) + + context.add_output_metadata( + { + "brand": active_brand, + "table": table or "?", + "rows_offered": len(enriched_products), + "rows_written": written, + "database": target, + "cleanup": "False (never deletes rows outside this batch)", + } + ) + return written + + +@asset( + group_name=GROUP, + partitions_def=brand_partitions(), + retry_policy=NETWORK_RETRY, + description="384-dim vectors generated for rows that do not have one.", +) +def product_embeddings( + context: AssetExecutionContext, + config: BrandConfig, + active_brand: str, + catalog_database: int, +) -> int: + """Embed only the rows missing a vector, in batches. + + Re-embedding an entire brand on every run would be the most expensive thing + this DAG does and would change nothing. Selecting on `embedding IS NULL` + makes the asset idempotent and makes a re-run nearly free, which is what + allows the refresh schedule to be daily. + """ + from app.services.embeddings_service import embed_texts + from app.services.vector_store import _connect, _table_name + + require_local_database("product_embeddings") + + table = _table_name(active_brand) + batch = max(1, config.embed_batch_size) + started = time.time() + + conn = _connect() + if conn is None: + raise Failure(description="Database unreachable at " + database_target()) + + embedded = 0 + try: + with conn, conn.cursor() as cur: + cur.execute( + "SELECT image_id, title, product_name, description, category " + "FROM {} WHERE embedding IS NULL".format(table) + ) + pending = cur.fetchall() + + for start in range(0, len(pending), batch): + chunk = pending[start : start + batch] + texts = [ + " ".join( + str(part) + for part in (row[1], row[2], row[4], row[3]) + if part + )[:2000] + for row in chunk + ] + vectors = embed_texts(texts) + for row, vector in zip(chunk, vectors): + cur.execute( + "UPDATE {} SET embedding = %s, " + "updated_at = CURRENT_TIMESTAMP WHERE image_id = %s".format( + table + ), + (str(list(vector)), row[0]), + ) + embedded += len(chunk) + context.log.info( + "Embedded %d/%d rows for %s", embedded, len(pending), active_brand + ) + finally: + conn.close() + + context.add_output_metadata( + { + "brand": active_brand, + "table": table, + "rows_embedded": embedded, + "batch_size": batch, + "model": _embedding_model_name(), + "duration_s": round(time.time() - started, 2), + } + ) + return embedded + + +@asset( + group_name=GROUP, + partitions_def=brand_partitions(), + description="Verification: the brand is fully searchable in pgvector.", +) +def vector_index( + context: AssetExecutionContext, active_brand: str, product_embeddings: int +) -> Dict[str, Any]: + """Read-back check, deliberately a separate node from the write. + + "The upsert reported success" and "the table can answer a similarity query" + are different claims. Making the second one its own asset means the UI can + show a green `catalog_database` beside a red `vector_index` when rows land + but never become searchable - precisely the failure that used to present + to users as "the catalog looks empty". + """ + from app.services.vector_store import _connect, _table_name + + table = _table_name(active_brand) + conn = _connect() + if conn is None: + raise Failure(description="Database unreachable at " + database_target()) + + try: + with conn, conn.cursor() as cur: + cur.execute( + "SELECT count(*), count(embedding), count(DISTINCT category) " + "FROM {}".format(table) + ) + rows, embedded, categories = cur.fetchone() + finally: + conn.close() + + missing = (rows or 0) - (embedded or 0) + if rows and missing: + context.log.warning( + "%s has %d row(s) without an embedding - they will not appear in " + "semantic search.", + table, + missing, + ) + + context.add_output_metadata( + { + "table": table, + "rows": rows or 0, + "embedded": embedded or 0, + "missing_embeddings": missing, + "categories": categories or 0, + "searchable": missing == 0, + } + ) + return { + "table": table, + "rows": rows or 0, + "embedded": embedded or 0, + "missing_embeddings": missing, + "categories": categories or 0, + } + + +# --- small helpers ---------------------------------------------------------- + + +def _slug(brand: str) -> str: + from app.services.brand_registry import resolve_parent_brand + from app.services.vector_store import _sanitize_name + + return _sanitize_name(resolve_parent_brand(brand)) + + +def _embedding_model_name() -> str: + from app.infrastructure.settings import EMBEDDINGS_MODEL + + return EMBEDDINGS_MODEL + + +def _preview(products: List[Dict[str, Any]], limit: int = 5) -> str: + if not products: + return "_no rows_" + lines = ["| title | category | size |", "|---|---|---|"] + for row in products[:limit]: + lines.append( + "| {} | {} | {} |".format( + str(row.get("title") or row.get("product_name") or "")[:60], + row.get("category") or "", + row.get("size") or "", + ) + ) + return "\n".join(lines) + + +def _jsonable(value: Any) -> Any: + if isinstance(value, dict): + return {str(k): _jsonable(v) for k, v in value.items()} + if isinstance(value, (list, tuple)): + return [_jsonable(v) for v in value] + if isinstance(value, (str, int, float, bool)) or value is None: + return value + return str(value) diff --git a/orchestration/assets/ml.py b/orchestration/assets/ml.py new file mode 100644 index 0000000..8695a66 --- /dev/null +++ b/orchestration/assets/ml.py @@ -0,0 +1,187 @@ +"""Store-intelligence ML: dataset -> training -> evaluation. + +Wraps `store_seed_service` and `ml_training_service`. These are the only ML +workflows worth orchestrating: the nutrition models live in the nutrition +group next to the data they fit, and the experimental models (forecast, +store_performance, purchase_propensity) are excluded on purpose - nothing +serves their output, so scheduling them would burn CPU to write artifacts no +request reads. They remain trainable by name. +""" +# NOTE: deliberately no `from __future__ import annotations` here. +# Dagster resolves the decorated function signatures at definition time to +# validate the `context` parameter and to infer asset input types. Under +# PEP 563/649 the annotations arrive as strings and that validation fails +# with "Cannot annotate `context` parameter with type AssetExecutionContext". +# Local Python is 3.14, which defers annotations by default, so this is not +# hypothetical. + +import time +from typing import Any, Dict, List, Optional + +from dagster import AssetExecutionContext, Config, Failure, MetadataValue, asset + +from orchestration.config import require_local_database + +GROUP = "ml" + + +class TrainingConfig(Config): + """Which models to fit. + + `None` means `ml_training_service.PRODUCTION_MODELS` - the three that are + actually served. Name any of ALL_MODELS to train more. + """ + + models: Optional[List[str]] = None + + +class SeedConfig(Config): + """Synthetic store/order history parameters. + + `reset_orders` defaults to False so a scheduled run cannot silently discard + an order history someone is mid-analysis on. Turn it on per-run to rebuild + from scratch. + """ + + days: int = 90 + seed: int = 42 + reset_orders: bool = False + + +@asset( + group_name=GROUP, + description="Stores, inventory, prices and simulated order history.", +) +def training_dataset(context: AssetExecutionContext, config: SeedConfig) -> Dict[str, Any]: + """`store_seed_service.run_seed()` - the ML models' training data. + + The store intelligence models learn from order history, and this project + has no real orders, so the history is deterministically simulated + (`intelligence/order_simulation.py`, fixed seed). That is documented as + cold-start bootstrapping; the important part is that INFERENCE goes through + the trained model, not the generating formula. + + Made an explicit asset so the lineage answers "what were these models + actually fitted on" instead of leaving it implicit. + """ + from app.services.store_seed_service import run_seed + + require_local_database("training_dataset") + + started = time.time() + result = run_seed( + reset_orders=config.reset_orders, days=config.days, seed=config.seed + ) + + context.add_output_metadata( + { + "days_simulated": config.days, + "random_seed": config.seed, + "reset_orders": config.reset_orders, + "result": MetadataValue.json(_jsonable(result)), + "duration_s": round(time.time() - started, 2), + } + ) + return _jsonable(result) + + +@asset( + group_name=GROUP, + description="Fitted discount / trending / popularity models (*.joblib).", +) +def trained_models( + context: AssetExecutionContext, + config: TrainingConfig, + training_dataset: Dict[str, Any], +) -> Dict[str, Any]: + """`ml_training_service.train_all()`. + + Depends on `training_dataset` because training against an empty orders + table produces a model that fits nothing and reports success - a silent + failure the dependency turns into an ordering guarantee. + """ + from app.services.ml_training_service import PRODUCTION_MODELS, train_all + + require_local_database("trained_models") + + models = config.models or list(PRODUCTION_MODELS) + started = time.time() + results = train_all(models=models) or {} + + trained = [name for name, r in results.items() if "skipped" not in (r or {})] + skipped = { + name: (r or {}).get("skipped") for name, r in results.items() if "skipped" in (r or {}) + } + + if not trained: + # Every model skipping means there was no usable training data. That is + # a failed run, not a successful no-op - the artifacts on disk are now + # stale and nothing says so. + raise Failure( + description=( + "No model trained. Every requested model reported 'skipped', " + "which means the training frames were empty - materialize " + "training_dataset first." + ), + metadata={"requested": MetadataValue.json(models), + "skipped": MetadataValue.json(_jsonable(skipped))}, + ) + + context.add_output_metadata( + { + "requested": MetadataValue.json(models), + "trained": MetadataValue.json(trained), + "skipped": MetadataValue.json(_jsonable(skipped)), + "results": MetadataValue.json(_jsonable(results)), + "duration_s": round(time.time() - started, 2), + "note": "Restart the API process to pick up new artifacts.", + } + ) + return _jsonable(results) + + +@asset( + group_name=GROUP, + description="Reports each production artifact's metrics and sample count.", +) +def model_evaluation( + context: AssetExecutionContext, trained_models: Dict[str, Any] +) -> Dict[str, Any]: + """Read the bundles back and report what is actually on disk. + + Separate from training for the same reason `vector_index` is separate from + `catalog_database`: "training returned a metrics dict" and "a loadable + artifact exists at MODEL_ARTIFACTS_DIR" are different claims, and only the + second is what the API will serve. + """ + from app.intelligence.model_utils import artifact_path, load_bundle + + report: Dict[str, Any] = {} + for name in ("discount_model", "trending_model", "popularity_model"): + path = artifact_path(name) + if not path.exists(): + report[name] = {"present": False} + context.log.warning("Artifact missing: %s", path) + continue + bundle = load_bundle(name) + report[name] = { + "present": True, + "size_bytes": path.stat().st_size, + "n_samples": getattr(bundle, "n_samples", None), + "trained_at": getattr(bundle, "trained_at", None), + "features": len(getattr(bundle, "feature_columns", []) or []), + "metrics": _jsonable(getattr(bundle, "extra", {}) or {}), + } + + context.add_output_metadata({"report": MetadataValue.json(_jsonable(report))}) + return _jsonable(report) + + +def _jsonable(value: Any) -> Any: + if isinstance(value, dict): + return {str(k): _jsonable(v) for k, v in value.items()} + if isinstance(value, (list, tuple)): + return [_jsonable(v) for v in value] + if isinstance(value, (str, int, float, bool)) or value is None: + return value + return str(value) diff --git a/orchestration/assets/nutrition.py b/orchestration/assets/nutrition.py new file mode 100644 index 0000000..16b788f --- /dev/null +++ b/orchestration/assets/nutrition.py @@ -0,0 +1,162 @@ +"""Nutrition enrichment and the two nutrition models. + +Wraps `nutrition_enrichment_service`, which already owns the Open Food Facts +lookup, the scoring, the diet-tag classification and the persistence. Dagster +contributes the dependency on `catalog_database` (you cannot enrich products +that are not stored yet), a bounded retry around a third-party API, and a +visible record of how many products were enriched versus skipped. +""" +# NOTE: deliberately no `from __future__ import annotations` here. +# Dagster resolves the decorated function signatures at definition time to +# validate the `context` parameter and to infer asset input types. Under +# PEP 563/649 the annotations arrive as strings and that validation fails +# with "Cannot annotate `context` parameter with type AssetExecutionContext". +# Local Python is 3.14, which defers annotations by default, so this is not +# hypothetical. + +import time +from typing import Any, Dict + +from dagster import ( + AssetExecutionContext, + Backoff, + Config, + Jitter, + MetadataValue, + RetryPolicy, + asset, +) + +from orchestration.config import require_local_database +from orchestration.partitions import brand_partitions + +GROUP = "nutrition" + +# Open Food Facts is a free public API and does rate-limit. Two retries with +# exponential backoff and jitter is enough for a transient 429/timeout without +# turning a bad afternoon into a retry storm. +OFF_RETRY = RetryPolicy( + max_retries=2, delay=10, backoff=Backoff.EXPONENTIAL, jitter=Jitter.PLUS_MINUS +) + + +class NutritionConfig(Config): + """Per-run knobs. + + `max_products` exists because a first enrichment of a large brand is one + outbound request per product. Capping it makes the asset safe to try + interactively before committing to a full run. + """ + + skip_if_verified: bool = True + generate_narrative: bool = False + max_products: int = 50 + + +@asset( + group_name=GROUP, + partitions_def=brand_partitions(), + retry_policy=OFF_RETRY, + deps=["catalog_database"], + description="Nutrition facts, scores and diet tags for this brand's products.", +) +def nutrition_data( + context: AssetExecutionContext, config: NutritionConfig +) -> Dict[str, Any]: + """Per-product enrichment via `enrich_one_product`. + + The single-product entry point is used rather than `enrich_all_products` + because the latter walks every brand in the database. Here the brand comes + from the partition, which is what keeps a partitioned run doing one brand's + worth of work. + + `generate_narrative` defaults to False: the narrative is written by the + local Ollama model, and on an 8GB CPU-only machine that is by far the + slowest part of enrichment. The nutrition panel renders without it. + """ + from app.services.nutrition_enrichment_service import enrich_one_product + from app.services.vector_store import get_products_by_brand + + require_local_database("nutrition_data") + + brand = context.partition_key + started = time.time() + + products = get_products_by_brand(brand, limit=config.max_products) or [] + counts: Dict[str, int] = {} + + for product in products: + image_id = product.get("image_id") + if not image_id: + continue + try: + status = enrich_one_product( + brand, + image_id, + product.get("product_name") or product.get("title") or "", + product.get("category") or "", + skip_if_verified=config.skip_if_verified, + generate_narrative=config.generate_narrative, + ) + except Exception as exc: # noqa: BLE001 - one bad product must not + # end the partition; the count of errors is the signal. + context.log.debug("Enrichment failed for %s: %s", image_id, exc) + status = "error" + counts[status] = counts.get(status, 0) + 1 + + result = { + "brand": brand, + "products_considered": len(products), + "by_status": counts, + "duration_s": round(time.time() - started, 2), + } + context.add_output_metadata( + { + "brand": brand, + "products_considered": len(products), + "by_status": MetadataValue.json(counts), + "narrative_generated": config.generate_narrative, + "duration_s": result["duration_s"], + } + ) + return result + + +@asset( + group_name=GROUP, + deps=["nutrition_data"], + description="KNN similarity index and KMeans clusters over nutrition_facts.", +) +def nutrition_models(context: AssetExecutionContext) -> Dict[str, Any]: + """`nutrition_enrichment_service.train_all_models()`, unpartitioned. + + Deliberately NOT partitioned by brand: both models fit across the whole + nutrition_facts table, and "find me a healthier alternative" is only useful + if it can cross brands. Partitioning would have produced per-brand indexes + that answer a narrower question than the endpoint asks. + """ + from app.services.nutrition_enrichment_service import train_all_models + + require_local_database("nutrition_models") + + started = time.time() + result = train_all_models() or {} + + context.add_output_metadata( + { + "result": MetadataValue.json(_jsonable(result)), + "duration_s": round(time.time() - started, 2), + "artifacts": "nutrition_similarity.joblib, nutrition_clustering.joblib", + } + ) + return _jsonable(result) + + +def _jsonable(value: Any) -> Any: + if isinstance(value, dict): + return {str(k): _jsonable(v) for k, v in value.items()} + if isinstance(value, (list, tuple)): + return [_jsonable(v) for v in value] + if isinstance(value, (str, int, float, bool)) or value is None: + return value + return str(value) diff --git a/orchestration/config.py b/orchestration/config.py new file mode 100644 index 0000000..5d0350c --- /dev/null +++ b/orchestration/config.py @@ -0,0 +1,115 @@ +"""Run configuration and the write guard. + +THE WRITE GUARD - read this before removing it +---------------------------------------------- +`backend/.env` points DB_HOST at the PRODUCTION database. That is fine for the +API, which only reads on the request path, but Dagster's ingestion assets call +`upsert_brand_products`, which writes. + +So a plain `dagster dev` from `backend/` would, with no warning, run pipeline +experiments against the live catalog. Two independent safeguards stop that: + +1. `orchestration/.env.orchestration` pins DB_HOST/DB_PORT at the local + Postgres container, and `definitions.py` loads it BEFORE anything imports + `app.infrastructure.settings` (which snapshots os.environ at import time). +2. `require_local_database()` re-checks the host that settings actually + resolved, at asset runtime, and fails the run if it is not local. + +The second exists because the first is a file that can be missing, renamed, or +overridden by a shell variable. Belt and braces is the right amount of caution +for a guard whose failure mode is silently mutating production. + +Set ORCHESTRATION_ALLOW_REMOTE_WRITES=true to deliberately target a remote +database. It is not wired to any default. +""" +# NOTE: deliberately no `from __future__ import annotations` here. +# Dagster resolves the decorated function signatures at definition time to +# validate the `context` parameter and to infer asset input types. Under +# PEP 563/649 the annotations arrive as strings and that validation fails +# with "Cannot annotate `context` parameter with type AssetExecutionContext". +# Local Python is 3.14, which defers annotations by default, so this is not +# hypothetical. + +import os +from typing import List, Optional + +from dagster import Config, Failure + +_LOCAL_HOSTS = {"localhost", "127.0.0.1", "::1", "postgres", "host.docker.internal"} + + +def database_target() -> str: + """`host:port` the app's settings actually resolved to. For logs/metadata.""" + from app.infrastructure import settings + + return f"{settings.DB_HOST}:{settings.DB_PORT}" + + +def remote_writes_allowed() -> bool: + return os.getenv("ORCHESTRATION_ALLOW_REMOTE_WRITES", "").strip().lower() in { + "1", + "true", + "yes", + } + + +def require_local_database(what: str) -> str: + """Raise unless the resolved database is local. Returns `host:port`. + + Called by every asset that writes. Read-only assets do not call it - they + are safe against any target and are genuinely useful pointed at production. + """ + from app.infrastructure import settings + + target = database_target() + if settings.DB_HOST in _LOCAL_HOSTS or remote_writes_allowed(): + return target + + raise Failure( + description=( + f"Refusing to run '{what}': it writes to Postgres, and the resolved " + f"database is {target}, which is not local.\n\n" + "This is the production database that backend/.env points at. Either " + "start Dagster with orchestration/.env.orchestration loaded (the " + "normal path - `dagster dev` from backend/ does this via " + "definitions.py), or set ORCHESTRATION_ALLOW_REMOTE_WRITES=true if " + "you genuinely mean to write there." + ), + metadata={"resolved_database": target, "guard": "require_local_database"}, + ) + + +class BrandConfig(Config): + """Per-run overrides. Every field has a safe default. + + The network-touching stages default to OFF, matching the application's own + settings: a large run with them on fires thousands of outbound requests, + which is why the store-catalog feature disabled them in the first place. + """ + + brands: Optional[List[str]] = None + fetch_images: bool = False + use_llm: bool = False + embed_batch_size: int = 32 + max_products_per_brand: Optional[int] = None + + +def resolve_brands(cfg_brands: Optional[List[str]]) -> List[str]: + """Run config wins; otherwise ACTIVE_BRANDS; otherwise whatever the DB has. + + The fallback matters: with ACTIVE_BRANDS unset (production's default) the + orchestrator must still have something to work on rather than silently + doing nothing. + """ + if cfg_brands: + return list(cfg_brands) + + from app.services.active_brands import active_display_names + + names = active_display_names() + if names: + return names + + from app.services.vector_store import list_available_brands + + return list_available_brands() diff --git a/orchestration/definitions.py b/orchestration/definitions.py new file mode 100644 index 0000000..ab7df34 --- /dev/null +++ b/orchestration/definitions.py @@ -0,0 +1,85 @@ +"""Dagster entry point: dagster dev -m orchestration.definitions -p 3030 + +ORDER MATTERS IN THIS FILE. + +`app.infrastructure.settings` reads the entire environment once, at import +time, and every `app.*` module imports it transitively. So the orchestration +env file has to be loaded before the first `app.*` import happens, or the +process silently keeps backend/.env's production database as its write target. +That is why the dotenv call is at the top, ahead of the asset imports, and why +those imports are not hoisted. + +`orchestration/config.require_local_database()` re-checks the resolved host at +asset runtime, so if this ordering is ever broken the failure is a red run with +an explanation rather than a write to production. +""" +# NOTE: deliberately no `from __future__ import annotations` here. +# Dagster resolves the decorated function signatures at definition time to +# validate the `context` parameter and to infer asset input types. Under +# PEP 563/649 the annotations arrive as strings and that validation fails +# with "Cannot annotate `context` parameter with type AssetExecutionContext". +# Local Python is 3.14, which defers annotations by default, so this is not +# hypothetical. + +import os +import sys +from pathlib import Path + +_BACKEND_ROOT = Path(__file__).resolve().parents[1] + +# `dagster dev -m orchestration.definitions` is documented to run from backend/, +# but adding the root here means it also works from the repository root or an +# IDE run configuration. +if str(_BACKEND_ROOT) not in sys.path: + sys.path.insert(0, str(_BACKEND_ROOT)) + + +def _load_orchestration_env() -> str: + """Load orchestration/.env.orchestration with override=True. + + override=True is deliberate and is the opposite of what settings.py does. + settings.py must not clobber real container environment variables, but this + file exists precisely to beat backend/.env - without the override, DB_HOST + would already be set from a developer's shell or a previous dotenv load and + the pin would do nothing. + """ + env_path = Path(os.getenv("ORCHESTRATION_ENV_FILE", _BACKEND_ROOT / "orchestration" / ".env.orchestration")) + if not env_path.exists(): + return "not found: {} (falling back to backend/.env - writes will be refused unless it is local)".format(env_path) + try: + from dotenv import load_dotenv + except ImportError: # pragma: no cover - python-dotenv is a hard dependency + return "python-dotenv unavailable; environment left as-is" + + load_dotenv(env_path, override=True) + return str(env_path) + + +_ENV_SOURCE = _load_orchestration_env() + +# Every app.* import must come after _load_orchestration_env(). +from dagster import Definitions, load_assets_from_modules, multiprocess_executor # noqa: E402 + +from orchestration import config as orchestration_config # noqa: E402 +from orchestration.assets import catalog, ml, nutrition # noqa: E402 +from orchestration.jobs import ALL_JOBS # noqa: E402 +from orchestration.schedules import ALL_SCHEDULES, ALL_SENSORS # noqa: E402 + +_all_assets = load_assets_from_modules([catalog, nutrition, ml]) + +defs = Definitions( + assets=_all_assets, + jobs=ALL_JOBS, + schedules=ALL_SCHEDULES, + sensors=ALL_SENSORS, + # Two concurrent processes, not the CPU count. Dagster runs alongside + # Postgres, Ollama, the API and Vite on the same 8GB machine, and each + # subprocess that touches an asset imports sentence-transformers. Four + # workers was enough to make the box swap. + executor=multiprocess_executor.configured({"max_concurrent": 2}), + metadata={ + "env_file": _ENV_SOURCE, + "database": orchestration_config.database_target(), + "remote_writes_allowed": str(orchestration_config.remote_writes_allowed()), + }, +) diff --git a/orchestration/jobs.py b/orchestration/jobs.py new file mode 100644 index 0000000..d9ea729 --- /dev/null +++ b/orchestration/jobs.py @@ -0,0 +1,84 @@ +"""Jobs - four logical workflows, deliberately not one large one. + +Splitting them this way follows how the work actually differs in cost and +cadence: + +* ingestion is cheap and offline (seed files), so it can run often; +* embedding drives the sentence-transformer, so it is separated to be + retried or run alone without redoing ingestion; +* nutrition enrichment makes third-party calls and is the slowest thing here; +* ML training needs order history, which is unrelated to catalog freshness. + +Selecting a subset of assets in one giant job would express the same graph, +but you could not schedule the parts on different cadences, and a failure in +one concern would show as a failure of everything. +""" +# NOTE: deliberately no `from __future__ import annotations` here. +# Dagster resolves the decorated function signatures at definition time to +# validate the `context` parameter and to infer asset input types. Under +# PEP 563/649 the annotations arrive as strings and that validation fails +# with "Cannot annotate `context` parameter with type AssetExecutionContext". +# Local Python is 3.14, which defers annotations by default, so this is not +# hypothetical. + +from dagster import AssetSelection, define_asset_job + +# NOTE: no partitions_def on define_asset_job. Dagster infers the partitioning +# from the selected assets, and passing it explicitly is deprecated (removed in +# 2.0). The catalog and embedding jobs are still per-brand partitioned because +# every asset they select is. + +# Brand -> raw -> validated -> enriched -> Postgres. +catalog_ingestion_job = define_asset_job( + name="catalog_ingestion_job", + description=( + "Brand selection, seed intake, validation, enrichment and the " + "Postgres upsert, for one brand partition." + ), + selection=AssetSelection.assets( + "active_brand", + "raw_products", + "validated_products", + "enriched_products", + "catalog_database", + ), +) + +# Stored rows -> embeddings -> a verified vector index. +embedding_refresh_job = define_asset_job( + name="embedding_refresh_job", + description=( + "Generate embeddings for stored rows that lack one, then verify the " + "brand is fully searchable in pgvector." + ), + selection=AssetSelection.assets("product_embeddings", "vector_index"), +) + +# Stored rows -> Open Food Facts -> nutrition_facts/insights -> the 2 models. +nutrition_enrichment_job = define_asset_job( + name="nutrition_enrichment_job", + description=( + "Fetch verified nutrition for a brand's products, score them, then " + "refit the similarity and clustering models across all brands." + ), + selection=AssetSelection.assets("nutrition_data", "nutrition_models"), +) + +# Synthetic order history -> fitted models -> artifact report. +ml_training_job = define_asset_job( + name="ml_training_job", + description=( + "Rebuild the store-intelligence training data, fit the served models " + "(discount, trending, popularity) and report the artifacts on disk." + ), + selection=AssetSelection.assets( + "training_dataset", "trained_models", "model_evaluation" + ), +) + +ALL_JOBS = [ + catalog_ingestion_job, + embedding_refresh_job, + nutrition_enrichment_job, + ml_training_job, +] diff --git a/orchestration/partitions.py b/orchestration/partitions.py new file mode 100644 index 0000000..7516a42 --- /dev/null +++ b/orchestration/partitions.py @@ -0,0 +1,70 @@ +"""Brand partitioning. + +One Dagster partition per active brand. This is the natural expression of the +"Brand Selection -> Ingestion -> ..." flow, and it buys three concrete things +that a single un-partitioned asset would not: + +* per-brand lineage in the UI, so "which brands are stale" is readable at a + glance rather than buried in one run's logs; +* per-brand retry - a failed Nestle partition does not force Amul to be + rebuilt; +* bounded memory. Each run holds one brand's rows, not the whole catalog, + which is what keeps this comfortable on an 8GB machine. + +The partition set is derived from ACTIVE_BRANDS, so widening the working set +from 3 brands to 25 adds partitions with no code change. +""" +# NOTE: deliberately no `from __future__ import annotations` here. +# Dagster resolves the decorated function signatures at definition time to +# validate the `context` parameter and to infer asset input types. Under +# PEP 563/649 the annotations arrive as strings and that validation fails +# with "Cannot annotate `context` parameter with type AssetExecutionContext". +# Local Python is 3.14, which defers annotations by default, so this is not +# hypothetical. + +from typing import List + +from dagster import StaticPartitionsDefinition + + +def active_brand_names() -> List[str]: + """Brand names to partition over. + + Falls back to the brand tables that exist when ACTIVE_BRANDS is unset, and + finally to a single literal so that the definitions still LOAD when the + database is unreachable. A Dagster code location that cannot be loaded is + far harder to debug than one that loads and shows an empty run - and the + webserver imports this at startup, before anyone can fix a connection. + """ + from app.services.active_brands import active_display_names + + names = active_display_names() + if names: + return names + + try: + from app.services.vector_store import list_available_brands + + live = list_available_brands() + if live: + return live + except Exception: # noqa: BLE001 - see docstring + pass + + return ["Amul"] + + +_PARTITIONS = None + + +def brand_partitions() -> StaticPartitionsDefinition: + """Cached so every asset shares one identical partition set. + + Dagster compares partition definitions by value, but building the list once + also means a single database round trip at code-load time instead of one + per asset. + """ + global _PARTITIONS + if _PARTITIONS is None: + _PARTITIONS = StaticPartitionsDefinition(active_brand_names()) + return _PARTITIONS diff --git a/orchestration/schedules.py b/orchestration/schedules.py new file mode 100644 index 0000000..ee3780f --- /dev/null +++ b/orchestration/schedules.py @@ -0,0 +1,152 @@ +"""Schedules and one sensor. + +EVERY SCHEDULE SHIPS STOPPED (`DefaultScheduleStatus.STOPPED`). + +That is the whole point on an 8GB development machine: defining a schedule +should cost nothing until someone decides they want it. A schedule that starts +running the moment `dagster dev` is launched would have the daemon waking up, +importing sentence-transformers and hitting Postgres on a laptop that is also +running the API, the frontend, Ollama and Postgres itself. Start them from the +Dagster UI when you want them; the definitions are here so that turning one on +is a click rather than a code change. + +The cadences are staggered so two heavy jobs never overlap. +""" +# NOTE: deliberately no `from __future__ import annotations` here. +# Dagster resolves the decorated function signatures at definition time to +# validate the `context` parameter and to infer asset input types. Under +# PEP 563/649 the annotations arrive as strings and that validation fails +# with "Cannot annotate `context` parameter with type AssetExecutionContext". +# Local Python is 3.14, which defers annotations by default, so this is not +# hypothetical. + +from dagster import ( + DefaultScheduleStatus, + DefaultSensorStatus, + RunRequest, + ScheduleDefinition, + SensorEvaluationContext, + SkipReason, + sensor, +) + +from orchestration.jobs import ( + catalog_ingestion_job, + embedding_refresh_job, + ml_training_job, + nutrition_enrichment_job, +) +from orchestration.partitions import active_brand_names + +# 02:00 - catalog first, so everything downstream sees fresh rows. +daily_catalog_refresh = ScheduleDefinition( + name="daily_catalog_refresh", + job=catalog_ingestion_job, + cron_schedule="0 2 * * *", + default_status=DefaultScheduleStatus.STOPPED, + description="Re-run seed intake, validation, enrichment and storage nightly.", +) + +# 03:00 - after ingestion has had an hour, embed whatever arrived. +daily_embedding_refresh = ScheduleDefinition( + name="daily_embedding_refresh", + job=embedding_refresh_job, + cron_schedule="0 3 * * *", + default_status=DefaultScheduleStatus.STOPPED, + description="Embed rows stored since the last run and verify the index.", +) + +# 04:00 - the slowest job, and the only one that leaves the machine. +daily_nutrition_refresh = ScheduleDefinition( + name="daily_nutrition_refresh", + job=nutrition_enrichment_job, + cron_schedule="0 4 * * *", + default_status=DefaultScheduleStatus.STOPPED, + description="Fetch nutrition for products that still lack it, then refit.", +) + +# Weekly, not daily: the training data is a deterministic simulation, so a +# nightly refit would spend CPU reproducing almost exactly the same models. +weekly_ml_retrain = ScheduleDefinition( + name="weekly_ml_retrain", + job=ml_training_job, + cron_schedule="0 5 * * 0", + default_status=DefaultScheduleStatus.STOPPED, + description="Sunday 05:00 - rebuild the store dataset and refit the served models.", +) + +ALL_SCHEDULES = [ + daily_catalog_refresh, + daily_embedding_refresh, + daily_nutrition_refresh, + weekly_ml_retrain, +] + + +@sensor( + job=catalog_ingestion_job, + minimum_interval_seconds=60, + default_status=DefaultSensorStatus.STOPPED, + description="Re-ingest a brand when its seed catalog file changes on disk.", +) +def seed_catalog_sensor(context: SensorEvaluationContext): + """Watch seed-file mtimes; request a run for the brands that changed. + + Kept as simple as it can be on purpose - a stat() of a handful of files + once a minute. No file-watcher process, no message broker, no inotify. The + brief asked for lightweight event handling and this is the cheapest thing + that actually works; anything more would be infrastructure to maintain for + a directory that changes a few times a day at most. + + The cursor is the newest mtime seen. On first evaluation it records the + current state and requests nothing, so enabling the sensor does not + immediately trigger a full rebuild of every brand. + """ + from app.services.brand_registry import resolve_parent_brand + from app.services.brand_sync import SEED_DIR, seed_catalog_paths + + if not SEED_DIR.exists(): + return SkipReason("Seed directory {} does not exist.".format(SEED_DIR)) + + partitions = {name: resolve_parent_brand(name) for name in active_brand_names()} + if not partitions: + return SkipReason("No active brands to watch.") + + previous = float(context.cursor) if context.cursor else None + newest = previous or 0.0 + changed = set() + + # Only the active directory: an archived catalog belongs to a brand the app + # is not serving, so a change to one is not an event worth acting on. + for path in seed_catalog_paths(SEED_DIR): + if path.parent != SEED_DIR: + continue + try: + mtime = path.stat().st_mtime + except OSError: + continue + newest = max(newest, mtime) + if previous is None or mtime <= previous: + continue + for name, parent in partitions.items(): + if parent.split()[0].lower() in path.name.lower() or name.lower() in path.name.lower(): + changed.add(name) + + context.update_cursor(str(newest)) + + if previous is None: + return SkipReason( + "First evaluation - recorded the current file state without " + "triggering a rebuild." + ) + if not changed: + return SkipReason("No active seed catalog changed since the last check.") + + context.log.info("Seed files changed for: %s", ", ".join(sorted(changed))) + return [ + RunRequest(run_key="{}-{}".format(brand, newest), partition_key=brand) + for brand in sorted(changed) + ] + + +ALL_SENSORS = [seed_catalog_sensor] diff --git a/requirements-orchestration.txt b/requirements-orchestration.txt new file mode 100644 index 0000000..06d32c2 --- /dev/null +++ b/requirements-orchestration.txt @@ -0,0 +1,20 @@ +# Dagster orchestration layer - DEVELOPMENT ONLY. +# +# Kept in a separate file from requirements.txt on purpose: the API image +# (backend/Dockerfile) must not grow an orchestrator it never runs, and +# docker-compose.prod.yml already allocates the whole 8GB VPS (ollama 3G, +# backend 2560M, postgres 1G) with no headroom for a daemon + webserver. +# +# pip install -r requirements-orchestration.txt +# dagster dev -m orchestration.definitions -p 3030 (run from backend/) +# +# Compatibility notes: +# * dagster 1.13 declares requires_python <3.15,>=3.10, so the local 3.14 +# interpreter and the container's 3.11 are both supported. +# * dagster pins protobuf<7 while the ambient environment has 7.x. Installing +# into backend/venv shadows it for this project only. +# * dagster-postgres is deliberately NOT used - it needs psycopg2-binary and +# this project is on psycopg3. Run/event storage stays on the default +# SQLite instance under orchestration/.dagster_home. +dagster>=1.13,<1.14 +dagster-webserver>=1.13,<1.14 diff --git a/scripts/train_ml_models.py b/scripts/train_ml_models.py index 8f4acfd..3066179 100644 --- a/scripts/train_ml_models.py +++ b/scripts/train_ml_models.py @@ -8,11 +8,16 @@ Run from `backend/` AFTER scripts/seed_store_intelligence.py: python scripts/train_ml_models.py python scripts/train_ml_models.py --models discount trending # train a subset -Models trained: discount, trending, popularity, forecast (demand + -inventory), store_performance, purchase_propensity. Every trained -model is saved to app/intelligence/artifacts/*.joblib and loaded lazily -by the API on first use - restart the API process after retraining to -pick up new artifacts (see docs/CHANGES.md, "Retraining"). +By default this trains the models the API actually serves: discount, +trending, popularity. The other three - forecast (demand + inventory), +store_performance, purchase_propensity - train correctly but have no +inference consumer, so they are opt-in by name: + + python scripts/train_ml_models.py --models store_performance + +Every trained model is saved to app/intelligence/artifacts/*.joblib and +loaded lazily by the API on first use - restart the API process after +retraining to pick up new artifacts (see docs/CHANGES.md, "Retraining"). """ from __future__ import annotations @@ -32,7 +37,7 @@ def main() -> None: parser.add_argument( "--models", nargs="+", default=None, choices=["discount", "trending", "popularity", "forecast", "store_performance", "purchase_propensity"], - help="Subset of models to train (default: all)", + help="Subset of models to train (default: the served models - discount, trending, popularity)", ) args = parser.parse_args() diff --git a/tests/conftest.py b/tests/conftest.py index e98e60f..806d8ad 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -56,6 +56,17 @@ os.environ.setdefault("DB_PASSWORD", "test-password-not-real") os.environ.setdefault("USE_S3", "false") os.environ.setdefault("USE_GOOGLE_CSE", "false") +# Unconditional, NOT setdefault. The suite pins brand-name behaviour all over +# the place ("any Nestle chocolates?", the suggest ranking fixtures), and a +# developer's backend/.env now carries a real ACTIVE_BRANDS value. Letting that +# leak in would make those tests pass or fail depending on whose machine ran +# them - the same trap AUTH_ALLOW_ANY_LOGIN sprang before it was pinned here. +# +# Blank means "no brand filtering", so every existing test sees the historical +# behaviour. The filtering itself is covered by tests/test_active_brands.py, +# which sets the value explicitly and clears the parsed cache. +os.environ["ACTIVE_BRANDS"] = "" + # Auth is set unconditionally (not setdefault): the suite asserts on the real # guards, so it must never inherit a developer's AUTH_ENABLED=false. os.environ["AUTH_ENABLED"] = "true" diff --git a/tests/test_active_brands.py b/tests/test_active_brands.py new file mode 100644 index 0000000..fd22895 --- /dev/null +++ b/tests/test_active_brands.py @@ -0,0 +1,247 @@ +"""ACTIVE_BRANDS: the one setting that narrows the whole application. + +conftest.py pins ACTIVE_BRANDS="" for the rest of the suite, so this module +sets it explicitly and restores it. Everything here is pure config parsing - +no database, no network. +""" +from __future__ import annotations + +import pytest + +from app.services import active_brands as ab +from app.services.brand_sync import ARCHIVE_DIR_NAME, SEED_DIR, seed_catalog_paths + + +@pytest.fixture +def configured(monkeypatch): + """Set ACTIVE_BRANDS and drop the parsed cache on the way in and out.""" + + def _apply(raw: str): + monkeypatch.setattr(ab, "_RAW_ACTIVE_BRANDS", raw) + ab.invalidate() + return ab + + yield _apply + ab.invalidate() + + +THREE = "Amul,Cadbury,Hindustan Unilever" + + +# --- the empty-means-all contract ------------------------------------------- +# Production leaves ACTIVE_BRANDS unset. If blank ever started meaning "no +# brands are active" instead of "every brand is active", the entire catalog +# would go empty in production and the API would keep answering 200s while +# doing it - the single most expensive failure mode this project has had. + +@pytest.mark.parametrize("blank", ["", " ", ",", " , , "]) +def test_blank_config_disables_filtering_entirely(configured, blank): + cfg = configured(blank) + assert cfg.filtering_enabled() is False + assert cfg.active_brand_suffixes() is None + assert cfg.is_active_suffix("nestle") is True + assert cfg.is_active_brand("literally anything") is True + assert cfg.filter_suffixes(["amul", "nestle", "p_g"]) == ["amul", "nestle", "p_g"] + + +def test_configured_brands_narrow_the_set(configured): + cfg = configured(THREE) + assert cfg.filtering_enabled() is True + assert cfg.active_brand_suffixes() == frozenset({"amul", "cadbury", "hindustan_unilever"}) + assert cfg.active_display_names() == ["Amul", "Cadbury", "Hindustan Unilever"] + + +def test_filter_suffixes_drops_inactive_and_keeps_order(configured): + cfg = configured(THREE) + given = ["nestle", "amul", "p_g", "hindustan_unilever", "cadbury", "pepsico"] + assert cfg.filter_suffixes(given) == ["amul", "hindustan_unilever", "cadbury"] + + +# --- names are resolved, not string-matched --------------------------------- + +def test_names_resolve_through_the_brand_aliases(configured): + """Config must agree with the storage layer about which table a brand is. + + "Tata" has no table of its own - the "hul tata tea" alias routes it into + brand_hindustan_unilever, where brand_catalog_tata.json's 121 products + already live. Comparing raw strings here would have let ACTIVE_BRANDS and + resolve_parent_brand disagree, hiding those products from an active brand. + """ + cfg = configured(THREE) + assert cfg.is_active_brand("Tata") is True + assert cfg.is_active_brand("tata tea") is True + assert cfg.is_active_brand("Dove") is True # -> hindustan unilever + assert cfg.is_active_brand("cadbury dairy milk") is True + assert cfg.is_active_brand("amul butter") is True + assert cfg.is_active_brand("Nestle") is False + assert cfg.is_active_brand("P&G") is False + + +def test_configuring_a_sub_brand_activates_its_parent_table(configured): + cfg = configured("Tata") + assert cfg.active_brand_suffixes() == frozenset({"hindustan_unilever"}) + + +def test_whitespace_and_duplicates_are_tolerated(configured): + cfg = configured(" Amul , amul , Cadbury ") + assert cfg.active_brand_suffixes() == frozenset({"amul", "cadbury"}) + + +# --- the query_intent brand index must narrow too --------------------------- + +def test_brand_index_drops_inactive_static_brands(configured, monkeypatch): + """An archived brand must stop parsing as a brand mention. + + KNOWN_BRANDS/BRAND_SEARCH_MAP are a hand-tuned static table, and the live + brand list only ever *added* to it. Left alone, "any Nestle chocolates?" + would still resolve to a brand, get routed to brand-catalog mode, and come + back empty rather than falling through to ordinary search. + """ + from app.services import query_intent as qi + + configured(THREE) + monkeypatch.setattr(qi, "list_available_brands", lambda: [], raising=False) + qi.invalidate_brand_mention_cache() + monkeypatch.setattr( + "app.services.vector_store.list_available_brands", lambda: [], raising=False + ) + + known, mapping = qi._brand_index() + assert "Amul" in known + assert "Cadbury" in known + assert "Hindustan Unilever" in known + assert "Nestle" not in known + assert "Pepsico" not in known + # Synonyms follow their target brand. + assert mapping.get("hul") == "Hindustan Unilever" + assert "pepsi" not in mapping + assert "coke" not in mapping + + qi.invalidate_brand_mention_cache() + + +def test_extract_brand_mention_ignores_an_inactive_brand(configured, monkeypatch): + from app.services import query_intent as qi + + configured(THREE) + monkeypatch.setattr( + "app.services.vector_store.list_available_brands", lambda: [], raising=False + ) + qi.invalidate_brand_mention_cache() + + assert qi.extract_brand_mention("show me Amul butter") == "Amul" + assert qi.extract_brand_mention("any Nestle chocolates?") is None + + qi.invalidate_brand_mention_cache() + + +# --- the seed archive ------------------------------------------------------- + +def test_archived_catalogs_are_still_discoverable_on_disk(): + """Re-activating a brand must be a config change, not a file move. + + seed_catalog_paths() reads the archive subdirectory as well, so flipping a + name into ACTIVE_BRANDS is enough to seed and serve it again. + """ + names = {p.name for p in seed_catalog_paths(SEED_DIR)} + assert "brand_catalog_amul.json" in names + assert "brand_catalog_nestle.json" in names, "archived catalogs must stay reachable" + + active_only = {p.name for p in SEED_DIR.glob("*.json")} + assert "brand_catalog_nestle.json" not in active_only, ( + "an archived catalog must NOT be picked up by the plain glob the boot " + "auto-seed used to run over every file" + ) + assert (SEED_DIR / ARCHIVE_DIR_NAME).is_dir() + + +def test_load_seed_catalogs_only_returns_active_brands(configured): + from app.services import brand_sync + + configured(THREE) + grouped = brand_sync.load_seed_catalogs() + assert set(grouped) == {"amul", "cadbury", "hindustan unilever"} + # tata.json resolves into hindustan unilever, so its products must be there. + assert len(grouped["hindustan unilever"]) > 104 + + +def test_explicit_only_overrides_the_active_filter(configured): + """`only=` is a deliberate request, so it must still reach an archived brand.""" + from app.services import brand_sync + + configured(THREE) + grouped = brand_sync.load_seed_catalogs(only=["nestle"]) + # resolve_parent_brand returns the lowercase alias parent, not the file's + # "Nestle" display casing. + assert set(grouped) == {"nestle"} + assert len(grouped["nestle"]) == 123 + + +def test_reconcile_does_not_resurrect_archived_brands(configured, monkeypatch): + """The boot reconcile must not re-seed the brands we just archived. + + REGRESSION. `_db_brand_counts()` is narrowed by ACTIVE_BRANDS (it goes + through _list_brand_table_suffixes), while `index_seed_files()` reads the + archive directory on purpose. Left mismatched, every archived catalog looks + like "a seed file whose table is empty" and lands in `to_seed` - so the + reconcile that runs on every API start, and again every 300 seconds, would + recreate all 27 archived brand tables and undo the archiving silently. + + A dry run confirmed exactly that before the fix. + """ + from app.services import brand_sync + + configured(THREE) + + # Only the active tables are visible, which is what the real filtered + # _db_brand_counts() returns. + monkeypatch.setattr( + brand_sync, + "_db_brand_counts", + lambda: {"amul": 122, "cadbury": 105, "hindustan_unilever": 225}, + ) + + summary = brand_sync.reconcile_brand_catalogs(dry_run=True) + + seeded = summary.get("would_seed") or summary.get("seeded") or [] + exported = summary.get("would_export") or summary.get("exported") or [] + assert seeded == [], "archived brands must never be re-seeded by reconcile" + assert exported == [] + assert summary["files"] == 3, "reconcile must only consider active seed files" + + +@pytest.mark.parametrize( + "brand,expected_min", + [("Amul", 122), ("Cadbury", 105), ("Hindustan Unilever", 225), ("Tata", 225)], +) +def test_load_brand_products_resolves_multi_word_brands(brand, expected_min): + """REGRESSION: a brand name is not a filename. + + `load_seed_catalogs(only=["Hindustan Unilever"])` matches substrings of the + FILE NAME, and "hindustan unilever" (space) is not a substring of + "brand_catalog_hindustan_unilever.json" (underscore). It returned zero + products, silently - the Dagster partition for HUL reported RUN_SUCCESS + having ingested nothing. + + `load_brand_products` goes through the seed-file index instead, which is + keyed by resolved table suffix, so it also picks up brand_catalog_tata.json + (121 HUL products via the "hul tata tea" alias). + """ + from app.services.brand_sync import load_brand_products + + assert len(load_brand_products(brand)) >= expected_min + + +def test_load_brand_products_reaches_archived_catalogs(): + """An explicitly named brand must be found even while it is archived.""" + from app.services.brand_sync import load_brand_products + + assert len(load_brand_products("Nestle")) == 123 + + +def test_load_seed_catalogs_only_still_matches_filenames(): + """The `only=` filename semantics are unchanged - the CLI documents them.""" + from app.services.brand_sync import load_seed_catalogs + + grouped = load_seed_catalogs(only=["amul"]) + assert sum(len(v) for v in grouped.values()) == 122 diff --git a/tests/test_brand_registry.py b/tests/test_brand_registry.py index a69551d..0843fa7 100644 --- a/tests/test_brand_registry.py +++ b/tests/test_brand_registry.py @@ -17,15 +17,22 @@ from pathlib import Path import pytest from app.services.brand_registry import BRAND_ALIASES, resolve_parent_brand +from app.services.brand_sync import seed_catalog_paths from app.services.vector_store import _sanitize_name SEED_DIR = Path(__file__).resolve().parents[1] / "data" / "seed_catalogs" def _seed_brand_fields() -> list[str]: - """The `brand` value of every seed catalog (skipping non-catalog files).""" + """The `brand` value of every seed catalog (skipping non-catalog files). + + Deliberately covers the `archive/` subdirectory as well as the active + catalogs. Archiving a brand takes it out of the running app, but it must + not take it out of this guarantee - an archived file is re-activated by a + config change alone, and it has to land in the same table it always did. + """ brands = [] - for path in sorted(SEED_DIR.glob("*.json")): + for path in seed_catalog_paths(SEED_DIR): try: data = json.loads(path.read_text(encoding="utf-8-sig")) except Exception: @@ -36,29 +43,45 @@ def _seed_brand_fields() -> list[str]: return brands -# Every seed catalog and the table it must continue to feed. `tata` mapping to -# hindustan_unilever is not a typo: alias "hul tata tea" claims it, which is -# where all 121 of that file's products already live. +# Every seed catalog and the table it must continue to feed, active and +# archived alike. `tata` mapping to hindustan_unilever is not a typo: alias +# "hul tata tea" claims it, which is where all 121 of that file's products +# already live - and it is why brand_catalog_tata.json stays in the active +# directory while Hindustan Unilever is an active brand. EXPECTED_SEED_TABLES = { + "aachi": "brand_aachi", "amul": "brand_amul", + "anil": "brand_anil", + "bikaji": "brand_bikaji", + "britannia": "brand_britannia", "cadbury": "brand_cadbury", "cavinkare": "brand_cavinkare", "coca-cola": "brand_coca_cola", "colgate-palmolive": "brand_colgate_palmolive", "dabur": "brand_dabur", + "everest": "brand_everest", + "fortune": "brand_fortune", "godrej": "brand_godrej", "grb": "brand_grb", + "haldirams": "brand_haldirams", "hindustan unilever": "brand_hindustan_unilever", # Ingested straight into the database with no BRAND_ALIASES entry, so it # exercises the unaliased path: resolve_parent_brand returns it unchanged # and it gets its own table. Pinned here to catch the day some new alias # whole-word-matches "idhayam" and silently re-parents 24 products. "idhayam": "brand_idhayam", + "itc": "brand_itc", + "kaleesuwari": "brand_kaleesuwari", "lion dates": "brand_lion_dates", "Manna": "brand_manna", + "marico": "brand_marico", + "mdh": "brand_mdh", "milky mist": "brand_milky_mist", + "mtr": "brand_mtr", + "naga": "brand_naga", "Nestle": "brand_nestle", "p&g": "brand_p_g", + "parle": "brand_parle", "pepsico": "brand_pepsico", "tata": "brand_hindustan_unilever", } diff --git a/tests/test_orchestration_defs.py b/tests/test_orchestration_defs.py new file mode 100644 index 0000000..e93d49a --- /dev/null +++ b/tests/test_orchestration_defs.py @@ -0,0 +1,219 @@ +"""The Dagster definitions load, and the graph is the one we meant to build. + +Skipped entirely when dagster is not installed: it lives in +requirements-orchestration.txt, not requirements.txt, so the API image and CI +runs that only install the app dependencies must still get a green suite. + +Nothing here executes a run. These are structural assertions - that the code +location imports, that the lineage edges exist, that the schedules are off, +and that the write guard refuses a remote database. A broken code location is +the failure mode worth catching early, because the Dagster webserver reports +it as an opaque load error long after the change that caused it. +""" +import os + +import pytest + +dagster = pytest.importorskip("dagster", reason="orchestration extra not installed") + + +@pytest.fixture(scope="module") +def defs(): + from orchestration.definitions import defs as _defs + + return _defs + + +@pytest.fixture(scope="module") +def asset_graph(defs): + return defs.get_repository_def().asset_graph + + +def _keys(asset_graph): + return {key.to_user_string() for key in asset_graph.get_all_asset_keys()} + + +def test_code_location_loads(defs): + assert defs is not None + + +def test_every_expected_asset_exists(asset_graph): + assert _keys(asset_graph) == { + # catalog + "active_brand", + "raw_products", + "validated_products", + "enriched_products", + "catalog_database", + "product_embeddings", + "vector_index", + # nutrition + "nutrition_data", + "nutrition_models", + # ml + "training_dataset", + "trained_models", + "model_evaluation", + } + + +@pytest.mark.parametrize( + "asset_key,expected_parents", + [ + ("raw_products", {"active_brand"}), + ("validated_products", {"active_brand", "raw_products"}), + ("enriched_products", {"active_brand", "validated_products"}), + ("catalog_database", {"active_brand", "enriched_products"}), + ("product_embeddings", {"active_brand", "catalog_database"}), + ("vector_index", {"active_brand", "product_embeddings"}), + ("nutrition_data", {"catalog_database"}), + ("nutrition_models", {"nutrition_data"}), + ("trained_models", {"training_dataset"}), + ("model_evaluation", {"trained_models"}), + ], +) +def test_lineage_edges(asset_graph, asset_key, expected_parents): + """The DAG shape IS the deliverable - pin it. + + Losing an edge does not fail a run, it just silently lets an asset + materialize against stale upstream data (embedding rows that were never + written, models fitted on an empty orders table). + """ + from dagster import AssetKey + + node = asset_graph.get(AssetKey(asset_key)) + assert {k.to_user_string() for k in node.parent_keys} == expected_parents + + +def test_catalog_assets_are_partitioned_by_brand(asset_graph): + """Per-brand partitioning is what bounds memory and lets one brand retry.""" + from dagster import AssetKey + + for key in ( + "active_brand", + "raw_products", + "validated_products", + "enriched_products", + "catalog_database", + "product_embeddings", + "vector_index", + ): + assert asset_graph.get(AssetKey(key)).is_partitioned, key + + +def test_cross_brand_assets_are_not_partitioned(asset_graph): + """nutrition_models and the ML models fit across every brand at once. + + Partitioning them would produce per-brand indexes that answer a narrower + question than "find a healthier alternative" actually asks. + """ + from dagster import AssetKey + + for key in ("nutrition_models", "training_dataset", "trained_models", "model_evaluation"): + assert not asset_graph.get(AssetKey(key)).is_partitioned, key + + +def test_all_four_jobs_resolve(defs): + assert {job.name for job in defs.jobs} == { + "catalog_ingestion_job", + "embedding_refresh_job", + "nutrition_enrichment_job", + "ml_training_job", + } + + +def test_every_schedule_ships_stopped(defs): + """A schedule that auto-starts would run heavy jobs on an 8GB dev laptop. + + This is the assertion that keeps `dagster dev` from quietly becoming a + background workload. + """ + from dagster import DefaultScheduleStatus + + assert defs.schedules, "expected schedules to be defined" + for schedule in defs.schedules: + assert schedule.default_status == DefaultScheduleStatus.STOPPED, schedule.name + + +def test_sensor_ships_stopped_and_is_not_hot(defs): + from dagster import DefaultSensorStatus + + assert defs.sensors, "expected the seed-catalog sensor" + for sensor in defs.sensors: + assert sensor.default_status == DefaultSensorStatus.STOPPED, sensor.name + assert sensor.minimum_interval_seconds >= 60, sensor.name + + +def test_network_assets_retry_and_are_bounded(asset_graph): + """Retries must exist on the flaky steps and must never be unbounded.""" + from dagster import AssetKey + + for key in ("raw_products", "enriched_products", "product_embeddings"): + policy = asset_graph.get(AssetKey(key)).assets_def.op.retry_policy + assert policy is not None, key + assert 0 < policy.max_retries <= 3, key + + +# --- the write guard -------------------------------------------------------- + + +def test_guard_allows_a_local_database(monkeypatch): + from app.infrastructure import settings + from orchestration import config + + monkeypatch.setattr(settings, "DB_HOST", "localhost") + monkeypatch.setattr(settings, "DB_PORT", "5432") + assert config.require_local_database("test") == "localhost:5432" + + +def test_guard_refuses_a_remote_database(monkeypatch): + """backend/.env points at production. This is the last line of defence. + + If the env-file ordering in definitions.py is ever broken, this turns a + silent write to the live catalog into a red run with an explanation. + """ + from dagster import Failure + + from app.infrastructure import settings + from orchestration import config + + monkeypatch.setattr(settings, "DB_HOST", "31.97.228.132") + monkeypatch.setattr(settings, "DB_PORT", "6054") + monkeypatch.delenv("ORCHESTRATION_ALLOW_REMOTE_WRITES", raising=False) + + with pytest.raises(Failure) as excinfo: + config.require_local_database("catalog_database") + assert "31.97.228.132" in str(excinfo.value) + + +def test_guard_can_be_overridden_deliberately(monkeypatch): + from app.infrastructure import settings + from orchestration import config + + monkeypatch.setattr(settings, "DB_HOST", "31.97.228.132") + monkeypatch.setattr(settings, "DB_PORT", "6054") + monkeypatch.setenv("ORCHESTRATION_ALLOW_REMOTE_WRITES", "true") + assert config.require_local_database("catalog_database") == "31.97.228.132:6054" + + +# --- brand config passthrough ---------------------------------------------- + + +def test_partitions_follow_active_brands(monkeypatch): + """3 brands -> 3 partitions. Widening the working set needs no code edit.""" + from app.services import active_brands + + monkeypatch.setattr(active_brands, "_RAW_ACTIVE_BRANDS", "Amul,Cadbury,Nestle") + active_brands.invalidate() + try: + from orchestration.partitions import active_brand_names + + assert active_brand_names() == ["Amul", "Cadbury", "Nestle"] + finally: + active_brands.invalidate() + + +def test_resolve_brands_prefers_explicit_run_config(monkeypatch): + from orchestration.config import resolve_brands + + assert resolve_brands(["Britannia"]) == ["Britannia"]