Files
2026-08-28 07:49:56 +05:30

135 lines
5.0 KiB
Python

"""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
class BatchConfig(Config):
"""Per-run config for the uploaded-spreadsheet batch job.
`batch_id` selects a batch staged under BATCH_UPLOAD_DIR by the API. Left
empty, the job picks the oldest batch still queued, which is what the
sensor wants and what makes a manual "run it now" from the Launchpad
convenient.
The two network flags default OFF for the same reason they do on
BrandConfig, with more force: a batch is up to twenty files, so leaving
image search on would fire tens of thousands of outbound requests from one
run.
"""
batch_id: Optional[str] = None
use_llm: bool = False
fetch_images: bool = False
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()