Files
2026-08-29 14:49:45 +05:30

11 KiB

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:

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

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.

batch_upload_sensor is safe to run next to a live API. Each batch carries a runner field naming its owner, and the sensor (and _pick_batch_id) claim only runner == "dagster" — the batches an admin sent here from Admin → Dagster Orchestration. Anything staged for the API's own worker is left alone, so the two no longer race for the same files. It still ships STOPPED like the rest: a development machine should not begin ingesting just because a directory has something in it.

Passing an explicit {"batch_id": "..."} in the Launchpad bypasses the runner filter — that is a person naming a batch, not a poll.

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.