"""Schedules and one sensor. EVERY SCHEDULE SHIPS STOPPED (`DefaultScheduleStatus.STOPPED`). That is the whole point on an 8GB development machine: defining a schedule should cost nothing until someone decides they want it. A schedule that starts running the moment `dagster dev` is launched would have the daemon waking up, importing sentence-transformers and hitting Postgres on a laptop that is also running the API, the frontend, Ollama and Postgres itself. Start them from the Dagster UI when you want them; the definitions are here so that turning one on is a click rather than a code change. The cadences are staggered so two heavy jobs never overlap. """ # 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. from dagster import ( DefaultScheduleStatus, DefaultSensorStatus, RunRequest, ScheduleDefinition, SensorEvaluationContext, SkipReason, sensor, ) from orchestration.jobs import ( batch_ingestion_job, catalog_ingestion_job, embedding_refresh_job, ml_training_job, nutrition_enrichment_job, ) from orchestration.partitions import active_brand_names # 02:00 - catalog first, so everything downstream sees fresh rows. daily_catalog_refresh = ScheduleDefinition( name="daily_catalog_refresh", job=catalog_ingestion_job, cron_schedule="0 2 * * *", default_status=DefaultScheduleStatus.STOPPED, description="Re-run seed intake, validation, enrichment and storage nightly.", ) # 03:00 - after ingestion has had an hour, embed whatever arrived. daily_embedding_refresh = ScheduleDefinition( name="daily_embedding_refresh", job=embedding_refresh_job, cron_schedule="0 3 * * *", default_status=DefaultScheduleStatus.STOPPED, description="Embed rows stored since the last run and verify the index.", ) # 04:00 - the slowest job, and the only one that leaves the machine. daily_nutrition_refresh = ScheduleDefinition( name="daily_nutrition_refresh", job=nutrition_enrichment_job, cron_schedule="0 4 * * *", default_status=DefaultScheduleStatus.STOPPED, description="Fetch nutrition for products that still lack it, then refit.", ) # Weekly, not daily: the training data is a deterministic simulation, so a # nightly refit would spend CPU reproducing almost exactly the same models. weekly_ml_retrain = ScheduleDefinition( name="weekly_ml_retrain", job=ml_training_job, cron_schedule="0 5 * * 0", default_status=DefaultScheduleStatus.STOPPED, description="Sunday 05:00 - rebuild the store dataset and refit the served models.", ) ALL_SCHEDULES = [ daily_catalog_refresh, daily_embedding_refresh, daily_nutrition_refresh, weekly_ml_retrain, ] @sensor( job=catalog_ingestion_job, minimum_interval_seconds=60, default_status=DefaultSensorStatus.STOPPED, description="Re-ingest a brand when its seed catalog file changes on disk.", ) def seed_catalog_sensor(context: SensorEvaluationContext): """Watch seed-file mtimes; request a run for the brands that changed. Kept as simple as it can be on purpose - a stat() of a handful of files once a minute. No file-watcher process, no message broker, no inotify. The brief asked for lightweight event handling and this is the cheapest thing that actually works; anything more would be infrastructure to maintain for a directory that changes a few times a day at most. The cursor is the newest mtime seen. On first evaluation it records the current state and requests nothing, so enabling the sensor does not immediately trigger a full rebuild of every brand. """ from app.services.brand_registry import resolve_parent_brand from app.services.brand_sync import SEED_DIR, seed_catalog_paths if not SEED_DIR.exists(): return SkipReason("Seed directory {} does not exist.".format(SEED_DIR)) partitions = {name: resolve_parent_brand(name) for name in active_brand_names()} if not partitions: return SkipReason("No active brands to watch.") previous = float(context.cursor) if context.cursor else None newest = previous or 0.0 changed = set() # Only the active directory: an archived catalog belongs to a brand the app # is not serving, so a change to one is not an event worth acting on. for path in seed_catalog_paths(SEED_DIR): if path.parent != SEED_DIR: continue try: mtime = path.stat().st_mtime except OSError: continue newest = max(newest, mtime) if previous is None or mtime <= previous: continue for name, parent in partitions.items(): if parent.split()[0].lower() in path.name.lower() or name.lower() in path.name.lower(): changed.add(name) context.update_cursor(str(newest)) if previous is None: return SkipReason( "First evaluation - recorded the current file state without " "triggering a rebuild." ) if not changed: return SkipReason("No active seed catalog changed since the last check.") context.log.info("Seed files changed for: %s", ", ".join(sorted(changed))) return [ RunRequest(run_key="{}-{}".format(brand, newest), partition_key=brand) for brand in sorted(changed) ] @sensor( job=batch_ingestion_job, minimum_interval_seconds=60, default_status=DefaultSensorStatus.STOPPED, description="Run a batch of uploaded spreadsheets once one is staged and waiting.", ) def batch_upload_sensor(context: SensorEvaluationContext): """Watch BATCH_UPLOAD_DIR for batches the API staged but did not run. There is NO schedule for batch ingestion, deliberately. A batch exists because a person uploaded files; there is nothing to do on a cron, and a schedule that woke up hourly to find nothing would be pure cost on a machine that is already sharing itself with Postgres, the API and Vite. Like seed_catalog_sensor this is as simple as it can be - a directory listing once a minute, no watcher process, no broker. The cursor is the set of batch ids already requested, so a batch is launched once even though it stays on disk afterwards. NOTE ON THE NORMAL PRODUCTION PATH: nothing here is involved. The API runs a staged batch itself, on a bounded worker thread, and this sensor exists for the development machine where Dagster is the thing driving the work. It is now safe to run alongside a live API. Both executors read the same directory, but a batch carries a `runner` naming which one owns it, and this sensor claims only `runner == "dagster"` - the batches an admin sent here from the orchestration tab. It still ships STOPPED, because a development machine should not start ingesting because a directory happened to have something in it. """ from app.core import batch_ingest root = batch_ingest.batch_root() if not root.exists(): return SkipReason("Batch upload directory {} does not exist.".format(root)) # Only batches an admin explicitly handed to Dagster. See the same filter, # and the reason for it, in assets/batch_catalog.py:_pick_batch_id. waiting = [ m for m in batch_ingest.list_manifests() if m.status in (batch_ingest.QUEUED, batch_ingest.INTERRUPTED) and m.runner == batch_ingest.RUNNER_DAGSTER ] if not waiting: return SkipReason("No staged batch is waiting for Dagster.") already = set(filter(None, (context.cursor or "").split(","))) fresh = [m for m in waiting if m.batch_id not in already] if not fresh: return SkipReason( "{} staged batch(es), all already requested.".format(len(waiting)) ) # Keep the cursor bounded - it is a string in the Dagster instance, not a # log, and only the recent tail is ever consulted. context.update_cursor(",".join(list(already | {m.batch_id for m in fresh})[-200:])) context.log.info( "Requesting runs for staged batch(es): %s", ", ".join(m.batch_id for m in fresh), ) return [ RunRequest( run_key=m.batch_id, run_config={"ops": {"batch_manifest": {"config": {"batch_id": m.batch_id}}}}, ) for m in fresh ] ALL_SENSORS = [seed_catalog_sensor, batch_upload_sensor]