Files
catalogue_backend/orchestration/definitions.py
2026-08-28 07:49:56 +05:30

86 lines
3.7 KiB
Python

"""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 batch_catalog, 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, batch_catalog])
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()),
},
)