feat: agent-11 PR3 – Scheduler + Erinnerungs-Mails + Notification-Prefs

Geplante Erinnerungen (Feature-Parität mit Urlaubsverwaltung):

- APScheduler (AsyncIOScheduler) in der FastAPI-Lifespan; tägliche Jobs ab
  settings.reminder_hour. Redis-Tageslock gegen Doppelversand bei mehreren
  Prozessen; jeder Job mit eigener Session + RLS-Bypass.
- Drei Jobs (auch einzeln aufrufbar): offene Anträge an Genehmiger,
  Resturlaub-Verfall-Vorwarnung an Mitarbeiter, fehlende AU an HR.
- Pro-User notification_prefs (JSONB, opt-out); GET/PATCH /users/me/notification-prefs
  + ProfilePage-UI; Vertreter-Mail respektiert die Prefs.
- Manueller Trigger POST /companies/me/run-reminders (Admin) – gleiche Logik,
  firmen-scoped (testbar ohne Warten).
- Bugfix: GET-/PATCH-Urlaubskonto (update_balance) nutzte nicht existente
  Felder (base_days/carried_over_days/ip_address) → korrigiert auf
  entitled_days/carried_over/ip + company_id; available_days ergänzt.

Migration 0037 (users.notification_prefs). 188/188 Tests grün. Deployed 137 + 164.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
2026-06-23 12:21:14 +02:00
co-authored by Claude Opus 4.8
parent e8bed43570
commit 2f110df619
15 changed files with 542 additions and 9 deletions
+5
View File
@@ -25,6 +25,11 @@ class Settings(BaseSettings):
# Redis
redis_url: str = "redis://localhost:6379/0"
# Scheduler (Erinnerungs-Mails, agent-11 PR3)
scheduler_enabled: bool = True
reminder_hour: int = 7 # Server-Stunde für tägliche Erinnerungen
carryover_reminder_days: int = 14 # Vorlauf für Resturlaub-Verfall-Erinnerung
# JWT
access_token_expire_minutes: int = 30
refresh_token_expire_days: int = 30
+27
View File
@@ -0,0 +1,27 @@
"""Pro-User-Benachrichtigungseinstellungen (agent-11 PR3).
Bekannte Keys + Defaults. Ein leeres `notification_prefs`-dict bedeutet „alle an"
(opt-out-Modell). Neue Keys defaulten daher automatisch auf aktiv.
"""
from __future__ import annotations
# key -> (Label, Default-aktiv)
NOTIFICATION_TYPES: dict[str, tuple[str, bool]] = {
"substitute_assigned": ("Als Vertretung eingetragen", True),
"pending_approvals": ("Offene Anträge zur Genehmigung (Zusammenfassung)", True),
"carryover_expiry": ("Resturlaub verfällt bald", True),
"certificate_overdue": ("Fehlende AU-Bescheinigung (HR)", True),
}
def pref_enabled(user, key: str) -> bool:
"""Ob der User Benachrichtigungen vom Typ `key` erhalten möchte."""
prefs = getattr(user, "notification_prefs", None) or {}
default = NOTIFICATION_TYPES.get(key, ("", True))[1]
return bool(prefs.get(key, default))
def normalize_prefs(prefs: dict | None) -> dict:
"""Nur bekannte Keys übernehmen, auf bool casten."""
prefs = prefs or {}
return {k: bool(prefs[k]) for k in NOTIFICATION_TYPES if k in prefs}
+9
View File
@@ -41,8 +41,17 @@ async def lifespan(app: FastAPI):
"In Production sollte ALLOWED_HOSTS in .env konfiguriert sein "
"um Host-Header-Injection zu verhindern."
)
# Erinnerungs-Scheduler (agent-11 PR3) nicht im Test-Kontext starten
import sys
if settings.scheduler_enabled and "pytest" not in sys.modules:
from app.services.scheduler_service import start as start_scheduler
start_scheduler()
yield
# Shutdown
from app.services.scheduler_service import shutdown as shutdown_scheduler
shutdown_scheduler()
await engine.dispose()
+4 -1
View File
@@ -4,7 +4,7 @@ from datetime import datetime, date
from typing import TYPE_CHECKING
from sqlalchemy import Boolean, Date, DateTime, Enum, ForeignKey, String, Text, func
from sqlalchemy.dialects.postgresql import UUID
from sqlalchemy.dialects.postgresql import JSONB, UUID
from sqlalchemy.orm import Mapped, mapped_column, relationship
from app.core.database import Base
@@ -54,6 +54,9 @@ class User(Base):
entry_date: Mapped[date | None] = mapped_column(Date)
exit_date: Mapped[date | None] = mapped_column(Date)
# Pro-User-Benachrichtigungseinstellungen (agent-11 PR3). Leeres dict = alle an.
notification_prefs: Mapped[dict] = mapped_column(JSONB, nullable=False, default=dict, server_default="{}")
# Kiosk auth
kiosk_pin_hash: Mapped[str | None] = mapped_column(Text)
kiosk_qr_token: Mapped[str | None] = mapped_column(Text, unique=True)
+11 -7
View File
@@ -426,28 +426,31 @@ async def update_balance(
balance = await absence_service.get_balance(user_id, year, db)
# Alte Werte für AuditLog sichern
old_base = balance.base_days
old_special = balance.special_days
old_carried = balance.carried_over_days
old_value = {
"entitled_days": balance.entitled_days,
"special_days": balance.special_days,
"carried_over": balance.carried_over,
}
for field, value in data.model_dump(exclude_unset=True).items():
setattr(balance, field, value)
# AuditLog schreiben
db.add(AuditLog(
company_id=current_user.company_id,
user_id=current_user.id,
action="update_vacation_balance",
entity_type="vacation_balance",
entity_id=balance.id,
old_value={"base_days": old_base, "special_days": old_special, "carried_over_days": old_carried},
old_value=old_value,
new_value={
"base_days": balance.base_days,
"entitled_days": balance.entitled_days,
"special_days": balance.special_days,
"carried_over_days": balance.carried_over_days,
"carried_over": balance.carried_over,
"target_user_id": str(user_id),
"year": year,
},
ip_address=get_client_ip(request),
ip=get_client_ip(request),
))
await db.commit()
@@ -456,6 +459,7 @@ async def update_balance(
expires_at, expired = _carryover_expiry(company, year) if company else (None, False)
return VacationBalanceOut.model_validate(balance).model_copy(update={
"pending_days": pending,
"available_days": absence_service.effective_available(balance, expired),
"carried_over_expires_at": expires_at,
"carried_over_expired": expired,
})
+13
View File
@@ -39,6 +39,19 @@ async def get_my_company(current_user: CurrentUser, db: AsyncSession = Depends(g
return CompanyOut.model_validate(company)
@router.post("/me/run-reminders")
async def run_reminders_now(
current_user: User = require_role(*_admin_roles),
db: AsyncSession = Depends(get_db),
):
"""Erinnerungs-Jobs (offene Anträge, Resturlaub-Verfall, AU fehlt) sofort für die
eigene Firma ausführen. Gleiche Logik wie der tägliche Scheduler."""
from app.services.scheduler_service import run_all_reminders
result = await run_all_reminders(db, company_id=current_user.company_id)
await db.commit()
return {"sent": result}
@router.patch("/me", response_model=CompanyOut)
async def update_my_company(
data: CompanyUpdate,
+28
View File
@@ -12,6 +12,8 @@ from app.schemas.auth import MessageResponse
from app.schemas.user import (
InviteRequest,
NextPersonnelNumberResponse,
NotificationPrefsUpdate,
NotificationTypeOut,
SetKioskPinRequest,
UserImportResult,
UserImportRowResult,
@@ -79,6 +81,32 @@ async def get_me(current_user: CurrentUser):
return UserOut.model_validate(current_user)
@router.get("/me/notification-prefs", response_model=list[NotificationTypeOut])
async def get_notification_prefs(current_user: CurrentUser):
from app.core.notifications import NOTIFICATION_TYPES, pref_enabled
return [
NotificationTypeOut(key=k, label=label, enabled=pref_enabled(current_user, k))
for k, (label, _default) in NOTIFICATION_TYPES.items()
]
@router.patch("/me/notification-prefs", response_model=list[NotificationTypeOut])
async def update_notification_prefs(
data: NotificationPrefsUpdate,
current_user: CurrentUser,
db: AsyncSession = Depends(get_db),
):
from app.core.notifications import NOTIFICATION_TYPES, normalize_prefs, pref_enabled
merged = dict(current_user.notification_prefs or {})
merged.update(normalize_prefs(data.prefs))
current_user.notification_prefs = merged # neues dict → JSONB-Änderung erkannt
await db.commit()
return [
NotificationTypeOut(key=k, label=label, enabled=pref_enabled(current_user, k))
for k, (label, _default) in NOTIFICATION_TYPES.items()
]
@router.get("/next-personnel-number", response_model=NextPersonnelNumberResponse)
async def next_personnel_number(
current_user: User = require_role(*_hr_roles),
+10
View File
@@ -32,6 +32,16 @@ class UserOut(BaseModel):
exit_date: date | None = None
class NotificationTypeOut(BaseModel):
key: str
label: str
enabled: bool
class NotificationPrefsUpdate(BaseModel):
prefs: dict[str, bool]
class UserUpdate(BaseModel):
first_name: str | None = Field(None, min_length=1, max_length=100)
last_name: str | None = Field(None, min_length=1, max_length=100)
+3
View File
@@ -1073,6 +1073,9 @@ class AbsenceService:
requester = await db.get(User, absence.user_id)
if substitute is None or requester is None or not substitute.email:
return
from app.core.notifications import pref_enabled
if not pref_enabled(substitute, "substitute_assigned"):
return
from app.services.email_service import email_service
try:
await email_service.send_substitute_notification(substitute, requester, absence, db)
+59
View File
@@ -163,6 +163,65 @@ class EmailService:
cfg,
)
async def send_pending_approvals_digest(
self, approver: "User", items: list[dict], db: AsyncSession
) -> None:
"""Tägliche Zusammenfassung offener Anträge an eine genehmigende Person."""
cfg = await self._load_smtp(approver.company_id, db)
rows = "".join(
f"<li>{i['user_name']} {i['type_name']} "
f"({i['start']}{' ' + i['end'] if i['end'] != i['start'] else ''}, "
f"{i['working_days']} Tag(e))</li>"
for i in items
)
body = f"""
<h1>Offene Abwesenheitsanträge</h1>
<p>Hallo {approver.first_name}, es warten <strong>{len(items)}</strong> Anträge auf deine Genehmigung:</p>
<ul style="font-size:14px;color:#555;line-height:1.7">{rows}</ul>
<a href="{settings.frontend_url}/absences?status=pending" class="btn">Anträge prüfen</a>
"""
await self._send(
approver.email, f"{len(items)} offene Abwesenheitsanträge",
_html_wrapper("Offene Anträge", body), cfg,
)
async def send_carryover_expiry_reminder(
self, user: "User", remaining: int, expires_at, db: AsyncSession
) -> None:
cfg = await self._load_smtp(user.company_id, db)
datum = expires_at.strftime("%d.%m.%Y")
body = f"""
<h1>Resturlaub verfällt bald</h1>
<p>Hallo {user.first_name}, du hast noch <strong>{remaining} Urlaubstage</strong>,
die am <strong>{datum}</strong> verfallen.</p>
<p>Plane deinen Urlaub rechtzeitig ein, damit nichts verloren geht.</p>
<a href="{settings.frontend_url}/absences" class="btn">Urlaub planen</a>
"""
await self._send(
user.email, f"Dein Resturlaub verfällt am {datum}",
_html_wrapper("Resturlaub", body), cfg,
)
async def send_certificate_overdue_digest(
self, hr_user: "User", items: list[dict], db: AsyncSession
) -> None:
cfg = await self._load_smtp(hr_user.company_id, db)
rows = "".join(
f"<li>{i['user_name']} krank seit {i['start']}, AU fällig seit {i['due']}</li>"
for i in items
)
body = f"""
<h1>Fehlende AU-Bescheinigungen</h1>
<p>Hallo {hr_user.first_name}, für <strong>{len(items)}</strong> Krankmeldung(en)
fehlt die Arbeitsunfähigkeitsbescheinigung:</p>
<ul style="font-size:14px;color:#555;line-height:1.7">{rows}</ul>
<a href="{settings.frontend_url}/absences" class="btn">Abwesenheiten ansehen</a>
"""
await self._send(
hr_user.email, f"{len(items)} fehlende AU-Bescheinigung(en)",
_html_wrapper("AU fehlt", body), cfg,
)
async def send_test(self, cfg: SmtpConfig, to: str) -> None:
"""Test-E-Mail direkt mit übergebenem Konfigurationsobjekt."""
body = f"""
+215
View File
@@ -0,0 +1,215 @@
"""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