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/:
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:
orchestration/.env.orchestrationpinsDB_HOST=localhost, anddefinitions.pyloads it withoverride=Truebefore anything importsapp.infrastructure.settings(which snapshots the environment once, at import time).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 |
batch_ingestion_job |
a batch of uploaded spreadsheets -> the 11 stages -> Postgres |
Five jobs rather than one because they differ in cost and cadence: ingestion is offline and cheap, embedding drives the transformer, nutrition leaves the machine, ML training depends on order history rather than catalog freshness, and batch ingestion is event-driven - it exists only when somebody uploads files.
batch_ingestion_job shares every stage with the upload path in production;
both call app/core/batch_ingest.py. Pick a batch with run config:
{"ops": {"batch_manifest": {"config": {"batch_id": "<id from the Admin UI>"}}}}
Left empty it takes the oldest batch still waiting under BATCH_UPLOAD_DIR.
From the command line:
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 |
batch_upload_sensor |
every 60s | run an uploaded batch that is staged and waiting |
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) andmax_concurrent: 2on the executor. Each subprocess that touches an asset imports sentence-transformers; four workers made the machine swap.product_embeddingsonly embeds rows whereembedding 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:
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.
This includes batch_ingestion_job. The catalog ingestion endpoints
(POST /api/admin/catalog-batch/ingest and POST /api/uploads/catalog) do not
talk to Dagster and do not need it running. They stage the uploaded files, then
execute the same functions this job wraps
(app/core/batch_ingest.py) on a single bounded worker thread inside the API
process - see app/core/batch_worker.py for why one worker rather than a
thread per upload. The job here exists so the graph is inspectable and a batch
can be re-run from the Launchpad on a development machine.
Do not turn batch_upload_sensor on next to a running API: both would pick up
the same staged batch. That is why it ships STOPPED like the rest.
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.