86 lines
3.7 KiB
Python
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()),
|
|
},
|
|
)
|