Add Dagster orchestration and reduce active brands in backend

This commit is contained in:
sriram
2026-08-20 16:39:54 +05:30
parent fbb1356e47
commit 7bf8dc6922
66 changed files with 2664 additions and 21 deletions

View File

@@ -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

View File

@@ -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

262
orchestration/README.md Normal file
View File

@@ -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 <http://127.0.0.1:3030>.
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.

10
orchestration/__init__.py Normal file
View File

@@ -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.
"""

View File

@@ -0,0 +1,3 @@
from orchestration.assets import catalog, ml, nutrition
__all__ = ["catalog", "nutrition", "ml"]

View File

@@ -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)

187
orchestration/assets/ml.py Normal file
View File

@@ -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)

View File

@@ -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)

115
orchestration/config.py Normal file
View File

@@ -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()

View File

@@ -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()),
},
)

84
orchestration/jobs.py Normal file
View File

@@ -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,
]

View File

@@ -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

152
orchestration/schedules.py Normal file
View File

@@ -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]