Public Access
feat: orchestrator runner — run_step / run_week per-family loop
Co-Authored-By: Claude Sonnet 4.6 (1M context) <noreply@anthropic.com>
This commit is contained in:
@@ -0,0 +1,83 @@
|
||||
"""
|
||||
Orchestrator runner.
|
||||
|
||||
run_step(step_name, week_start_date) — execute one step for all families
|
||||
run_week(week_start_date) — execute all steps in sequence
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
from datetime import date, timedelta
|
||||
from typing import Optional
|
||||
|
||||
from app.database import SessionLocal
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
STEPS = ("scrape", "generate", "email", "deadline", "finalize")
|
||||
|
||||
|
||||
def _current_week_start() -> date:
|
||||
"""Return the most recent Friday (today if today is Friday)."""
|
||||
today = date.today()
|
||||
days_since_friday = (today.weekday() - 4) % 7
|
||||
return today - timedelta(days=days_since_friday)
|
||||
|
||||
|
||||
def _get_or_create_run(db, family_id, week_start_date: date):
|
||||
from app.models import WeeklyRun
|
||||
|
||||
run = (
|
||||
db.query(WeeklyRun)
|
||||
.filter(
|
||||
WeeklyRun.family_id == family_id,
|
||||
WeeklyRun.week_start_date == week_start_date,
|
||||
)
|
||||
.first()
|
||||
)
|
||||
if run is None:
|
||||
run = WeeklyRun(
|
||||
family_id=family_id,
|
||||
week_start_date=week_start_date,
|
||||
status="pending",
|
||||
)
|
||||
db.add(run)
|
||||
db.flush()
|
||||
db.refresh(run)
|
||||
return run
|
||||
|
||||
|
||||
def run_step(step_name: str, week_start_date: Optional[date] = None) -> None:
|
||||
from app.models import FamilyProfile
|
||||
from app.services.orchestrator import steps as s
|
||||
|
||||
step_fns = {
|
||||
"scrape": s.step_scrape,
|
||||
"generate": s.step_generate,
|
||||
"email": s.step_email,
|
||||
"deadline": s.step_deadline,
|
||||
"finalize": s.step_finalize,
|
||||
}
|
||||
if step_name not in step_fns:
|
||||
raise ValueError(f"Unknown step: {step_name!r}. Valid: {list(step_fns)}")
|
||||
|
||||
week = week_start_date or _current_week_start()
|
||||
db = SessionLocal()
|
||||
families = db.query(FamilyProfile).all()
|
||||
if not families:
|
||||
logger.warning("run_step(%s): no family profiles, skipping", step_name)
|
||||
return
|
||||
for family in families:
|
||||
run = _get_or_create_run(db, family.id, week)
|
||||
try:
|
||||
step_fns[step_name](run, db)
|
||||
except Exception:
|
||||
logger.exception(
|
||||
"run_step(%s) failed for family %s", step_name, family.id
|
||||
)
|
||||
|
||||
|
||||
def run_week(week_start_date: Optional[date] = None) -> None:
|
||||
week = week_start_date or _current_week_start()
|
||||
for step in STEPS:
|
||||
run_step(step, week)
|
||||
Reference in New Issue
Block a user