153 lines
5.6 KiB
Python
153 lines
5.6 KiB
Python
"""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]
|