224 lines
8.6 KiB
Python
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]
|