"""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 ( 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) ] ALL_SENSORS = [seed_catalog_sensor]