Files
catalogue_backend/orchestration/schedules.py
2026-08-29 14:49:45 +05:30

224 lines
8.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 (
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]