"""Geplante Erinnerungs-Jobs (agent-11 PR3). AsyncIOScheduler läuft im FastAPI-Event-Loop. Jeder Job öffnet eine eigene DB-Session mit RLS-Bypass (interner Job ohne Tenant-Kontext) und holt vor dem Versand einen Redis-Tageslock, damit bei mehreren Prozessen nicht doppelt gemailt wird. Die run_*-Funktionen sind ohne Scheduler aufrufbar (Tests + manueller Trigger über die API) und können auf eine Firma eingegrenzt werden. """ from __future__ import annotations import logging from datetime import date from sqlalchemy import select from sqlalchemy.ext.asyncio import AsyncSession from app.core.config import settings from app.core.notifications import pref_enabled from app.models.absence import Absence, AbsenceStatus from app.models.absence_type import AbsenceCategory, AbsenceType from app.models.company import Company from app.models.user import User, UserRole from app.models.vacation_balance import VacationBalance from app.services.absence_service import absence_service from app.services.email_service import email_service log = logging.getLogger(__name__) _APPROVER_ROLES = (UserRole.MANAGER, UserRole.HR, UserRole.COMPANY_ADMIN) _HR_ROLES = (UserRole.HR, UserRole.COMPANY_ADMIN) _scheduler = None def _fmt(d: date) -> str: return d.strftime("%d.%m.%Y") async def _active_companies(db: AsyncSession, company_id) -> list[Company]: q = select(Company).where(Company.is_active.is_(True)) if company_id: q = q.where(Company.id == company_id) return list((await db.scalars(q)).all()) async def _recipients(db: AsyncSession, company_id, roles) -> list[User]: return list((await db.scalars( select(User).where( User.company_id == company_id, User.is_active.is_(True), User.role.in_(roles), ) )).all()) # ── Job-Logik ──────────────────────────────────────────────────────────────── async def run_pending_approvals(db: AsyncSession, company_id=None) -> int: """Zusammenfassung offener Anträge an alle Genehmiger je Firma. Liefert # Mails.""" sent = 0 for company in await _active_companies(db, company_id): rows = (await db.execute( select(Absence, User, AbsenceType) .join(User, Absence.user_id == User.id) .join(AbsenceType, Absence.type_id == AbsenceType.id) .where(User.company_id == company.id, Absence.status == AbsenceStatus.PENDING) .order_by(Absence.start_date) )).all() if not rows: continue items = [{ "user_name": u.full_name, "type_name": t.name, "start": _fmt(a.start_date), "end": _fmt(a.end_date), "working_days": float(a.working_days), } for a, u, t in rows] for approver in await _recipients(db, company.id, _APPROVER_ROLES): if not approver.email or not pref_enabled(approver, "pending_approvals"): continue await email_service.send_pending_approvals_digest(approver, items, db) sent += 1 return sent async def run_carryover_expiry(db: AsyncSession, company_id=None) -> int: """Erinnert Mitarbeiter an bald verfallenden Resturlaub (Vorlauf konfigurierbar).""" today = date.today() year = today.year sent = 0 for company in await _active_companies(db, company_id): expires_at, expired = absence_service._carryover_expired(company, year) if expires_at is None or expired: continue if (expires_at - today).days > settings.carryover_reminder_days: continue balances = (await db.execute( select(VacationBalance, User) .join(User, VacationBalance.user_id == User.id) .where( User.company_id == company.id, User.is_active.is_(True), VacationBalance.year == year, VacationBalance.carried_over > 0, ) )).all() for bal, user in balances: expiring = max(0, bal.carried_over - bal.used_days) # Übertrag zuerst verbraucht if expiring <= 0 or not user.email or not pref_enabled(user, "carryover_expiry"): continue await email_service.send_carryover_expiry_reminder(user, expiring, expires_at, db) sent += 1 return sent async def run_certificate_overdue(db: AsyncSession, company_id=None) -> int: """Meldet HR fehlende AU-Bescheinigungen (knüpft an agent-05).""" today = date.today() sent = 0 for company in await _active_companies(db, company_id): rows = (await db.execute( select(Absence, User) .join(User, Absence.user_id == User.id) .join(AbsenceType, Absence.type_id == AbsenceType.id) .where( User.company_id == company.id, AbsenceType.category == AbsenceCategory.SICK, Absence.certificate_required_by.isnot(None), Absence.certificate_required_by < today, Absence.certificate_received_at.is_(None), Absence.status == AbsenceStatus.APPROVED, ) .order_by(Absence.start_date) )).all() if not rows: continue items = [{ "user_name": u.full_name, "start": _fmt(a.start_date), "due": _fmt(a.certificate_required_by), } for a, u in rows] for hr in await _recipients(db, company.id, _HR_ROLES): if not hr.email or not pref_enabled(hr, "certificate_overdue"): continue await email_service.send_certificate_overdue_digest(hr, items, db) sent += 1 return sent async def run_all_reminders(db: AsyncSession, company_id=None) -> dict: return { "pending_approvals": await run_pending_approvals(db, company_id), "carryover_expiry": await run_carryover_expiry(db, company_id), "certificate_overdue": await run_certificate_overdue(db, company_id), } # ── Scheduler-Verdrahtung ──────────────────────────────────────────────────── def _acquire_daily_lock(name: str) -> bool: """True, wenn dieser Prozess den heutigen Job ausführen darf (Redis SET NX). Ohne Redis: immer True (Single-Prozess-Annahme).""" from app.core.redis import get_redis_client redis = get_redis_client() if redis is None: return True key = f"reminder_lock:{name}:{date.today().isoformat()}" try: return bool(redis.set(key, "1", nx=True, ex=23 * 3600)) except Exception as exc: log.warning("Reminder-Lock fehlgeschlagen (%s) – führe trotzdem aus: %s", name, exc) return True async def _run_job(name: str, fn) -> None: if not _acquire_daily_lock(name): log.info("Reminder-Job %s bereits von anderem Prozess übernommen.", name) return from sqlalchemy import text from app.core.database import AsyncSessionLocal try: async with AsyncSessionLocal() as db: await db.execute(text("SET LOCAL app.bypass_rls = 'on'")) count = await fn(db) await db.commit() log.info("Reminder-Job %s: %s Mail(s) versendet.", name, count) except Exception as exc: log.exception("Reminder-Job %s fehlgeschlagen: %s", name, exc) def start() -> None: global _scheduler if _scheduler is not None: return try: from apscheduler.schedulers.asyncio import AsyncIOScheduler from apscheduler.triggers.cron import CronTrigger except Exception as exc: log.warning("APScheduler nicht verfügbar – Erinnerungen deaktiviert: %s", exc) return _scheduler = AsyncIOScheduler() hour = settings.reminder_hour _scheduler.add_job(_run_job, CronTrigger(hour=hour, minute=0), args=["pending_approvals", run_pending_approvals], id="pending_approvals") _scheduler.add_job(_run_job, CronTrigger(hour=hour, minute=10), args=["carryover_expiry", run_carryover_expiry], id="carryover_expiry") _scheduler.add_job(_run_job, CronTrigger(hour=hour, minute=20), args=["certificate_overdue", run_certificate_overdue], id="certificate_overdue") _scheduler.start() log.info("Reminder-Scheduler gestartet (täglich ab %02d:00 Uhr).", hour) def shutdown() -> None: global _scheduler if _scheduler is not None: _scheduler.shutdown(wait=False) _scheduler = None