feat: orchestrer le pipeline M11
Co-authored-by: Codex/gpt-5.6-terra <codex-gpt-5-6-terra@agents.invalid>
This commit is contained in:
@@ -0,0 +1,19 @@
|
||||
"""Étapes isolées utilisées par l'orchestrateur du pipeline."""
|
||||
|
||||
from pronote_sync.pipeline.steps.caldav_sync import caldav_sync_step
|
||||
from pronote_sync.pipeline.steps.compare import compare_step
|
||||
from pronote_sync.pipeline.steps.fetch import fetch_step
|
||||
from pronote_sync.pipeline.steps.fetch_blog import fetch_blog_step
|
||||
from pronote_sync.pipeline.steps.normalize import normalize_step
|
||||
from pronote_sync.pipeline.steps.send import send_step
|
||||
from pronote_sync.pipeline.steps.synthesis import synthesis_step
|
||||
|
||||
__all__ = [
|
||||
"caldav_sync_step",
|
||||
"compare_step",
|
||||
"fetch_blog_step",
|
||||
"fetch_step",
|
||||
"normalize_step",
|
||||
"send_step",
|
||||
"synthesis_step",
|
||||
]
|
||||
|
||||
37
pronote_sync/pipeline/steps/caldav_sync.py
Normal file
37
pronote_sync/pipeline/steps/caldav_sync.py
Normal file
@@ -0,0 +1,37 @@
|
||||
"""Étape d'appel à la synchronisation CalDAV."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from typing import Protocol
|
||||
|
||||
from pronote_sync.config.settings import Settings
|
||||
from pronote_sync.models.pronote import PronoteData
|
||||
from pronote_sync.models.sync import CalDAVSyncResult
|
||||
|
||||
|
||||
class CalDAVSynchronizer(Protocol):
|
||||
"""Protocole injectable de synchronisation CalDAV."""
|
||||
|
||||
def __call__(self, data: PronoteData, settings: Settings) -> CalDAVSyncResult:
|
||||
"""Synchronise les données Pronote vers CalDAV.
|
||||
|
||||
:param data: Données Pronote normalisées.
|
||||
:param settings: Configuration effective de l'exécution.
|
||||
:return: Résultat de la synchronisation.
|
||||
:rtype: CalDAVSyncResult
|
||||
"""
|
||||
...
|
||||
|
||||
|
||||
def caldav_sync_step(
|
||||
synchronizer: CalDAVSynchronizer, data: PronoteData, settings: Settings
|
||||
) -> CalDAVSyncResult:
|
||||
"""Exécute la synchronisation CalDAV injectée.
|
||||
|
||||
:param synchronizer: Service de synchronisation injecté.
|
||||
:param data: Données Pronote normalisées.
|
||||
:param settings: Configuration effective de l'exécution.
|
||||
:return: Résultat CalDAV.
|
||||
:rtype: CalDAVSyncResult
|
||||
"""
|
||||
return synchronizer(data, settings)
|
||||
20
pronote_sync/pipeline/steps/compare.py
Normal file
20
pronote_sync/pipeline/steps/compare.py
Normal file
@@ -0,0 +1,20 @@
|
||||
"""Étape de comparaison de l'agenda réel avec l'agenda théorique."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from pronote_sync.models.diff import AgendaDiff
|
||||
from pronote_sync.models.pronote import PronoteData
|
||||
from pronote_sync.sync.diff import AgendaComparator
|
||||
|
||||
|
||||
def compare_step(comparator: AgendaComparator | None, data: PronoteData) -> AgendaDiff:
|
||||
"""Compare l'agenda ou retourne un diff vide si la comparaison est désactivée.
|
||||
|
||||
:param comparator: Comparateur configuré, ou ``None`` sans agenda théorique.
|
||||
:param data: Données Pronote normalisées.
|
||||
:return: Diff d'agenda pour la date cible.
|
||||
:rtype: AgendaDiff
|
||||
"""
|
||||
if comparator is None:
|
||||
return AgendaDiff(target_date=data.target_date)
|
||||
return comparator.compare(data.lessons, data.target_date)
|
||||
126
pronote_sync/pipeline/steps/fetch.py
Normal file
126
pronote_sync/pipeline/steps/fetch.py
Normal file
@@ -0,0 +1,126 @@
|
||||
"""Étape de récupération des données Pronote pour une exécution du pipeline."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from dataclasses import dataclass
|
||||
from datetime import date
|
||||
|
||||
from pronote_sync.errors import PipelineCriticalError, PipelineWarning
|
||||
from pronote_sync.models.agenda import Lesson, SchoolEvent
|
||||
from pronote_sync.models.homework import Homework
|
||||
from pronote_sync.models.message import Message
|
||||
from pronote_sync.sources.pronote.fallback import PronoteFetcherProtocol
|
||||
from pronote_sync.utils.redaction import redact_exception
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class FetchedPronoteData:
|
||||
"""Représente les données brutes récupérées pendant une exécution.
|
||||
|
||||
:ivar lessons: Cours récupérés depuis la source sélectionnée.
|
||||
:ivar homeworks: Devoirs destinés à la date cible.
|
||||
:ivar school_events: Événements scolaires récupérés avec l'agenda.
|
||||
:ivar messages: Messages et informations Pronote disponibles.
|
||||
:ivar target_date: Date cible du digest.
|
||||
"""
|
||||
|
||||
lessons: list[Lesson]
|
||||
homeworks: list[Homework]
|
||||
school_events: list[SchoolEvent]
|
||||
messages: list[Message]
|
||||
target_date: date
|
||||
|
||||
|
||||
def resolve_target_date(
|
||||
today: date, lessons: list[Lesson], school_events: list[SchoolEvent]
|
||||
) -> date:
|
||||
"""Détermine la date cible du digest à partir de l'agenda disponible.
|
||||
|
||||
La règle privilégie J+1 lorsqu'il contient des cours. Si la journée en
|
||||
cours contient des cours mais pas J+1, le prochain cours connu est choisi.
|
||||
Sans cours correspondant, J+1 est conservé, y compris pendant les vacances.
|
||||
|
||||
:param today: Date de référence de l'exécution.
|
||||
:param lessons: Cours récupérés pour la fenêtre de synchronisation.
|
||||
:param school_events: Événements scolaires récupérés (réservés aux évolutions
|
||||
du libellé de jour sans cours).
|
||||
:return: Date cible du digest.
|
||||
:rtype: date
|
||||
"""
|
||||
del school_events
|
||||
tomorrow = date.fromordinal(today.toordinal() + 1)
|
||||
lesson_dates = {lesson.start.date() for lesson in lessons}
|
||||
if tomorrow in lesson_dates:
|
||||
return tomorrow
|
||||
if today in lesson_dates:
|
||||
future_dates = sorted(day for day in lesson_dates if day > today)
|
||||
if future_dates:
|
||||
return future_dates[0]
|
||||
return tomorrow
|
||||
|
||||
|
||||
def _fetch_optional_messages(
|
||||
fetcher: PronoteFetcherProtocol,
|
||||
) -> tuple[list[Message], list[PipelineWarning]]:
|
||||
"""Récupère les messages et informations sans bloquer le pipeline.
|
||||
|
||||
:param fetcher: Fetcher Pronote configuré.
|
||||
:return: Messages disponibles et avertissements éventuels.
|
||||
:rtype: tuple[list[Message], list[PipelineWarning]]
|
||||
"""
|
||||
messages: list[Message] = []
|
||||
warnings: list[PipelineWarning] = []
|
||||
for step, method in (
|
||||
("fetch_messages", fetcher.fetch_messages),
|
||||
("fetch_informations", fetcher.fetch_informations),
|
||||
):
|
||||
try:
|
||||
messages.extend(method())
|
||||
except Exception as exc:
|
||||
warnings.append(
|
||||
PipelineWarning(
|
||||
f"Récupération non critique échouée : {redact_exception(exc)}",
|
||||
step=step,
|
||||
)
|
||||
)
|
||||
return messages, warnings
|
||||
|
||||
|
||||
def fetch_step(
|
||||
fetcher: PronoteFetcherProtocol, *, today: date | None = None
|
||||
) -> tuple[FetchedPronoteData, list[PipelineWarning]]:
|
||||
"""Récupère les données Pronote critiques et les compléments dégradables.
|
||||
|
||||
L'agenda et les devoirs sont critiques : leur échec empêche de produire un
|
||||
digest fiable et est donc propagé comme :class:`PipelineCriticalError`.
|
||||
Les messages et informations sont facultatifs ; leur échec produit un
|
||||
avertissement et une liste partielle reste valide.
|
||||
|
||||
:param fetcher: Fetcher Pronote configuré.
|
||||
:param today: Date de référence, injectée par les tests ; J courant par défaut.
|
||||
:return: Données récupérées et avertissements non critiques.
|
||||
:rtype: tuple[FetchedPronoteData, list[PipelineWarning]]
|
||||
:raises PipelineCriticalError: Si l'agenda ou les devoirs ne sont pas disponibles.
|
||||
"""
|
||||
try:
|
||||
lessons, school_events = fetcher.fetch_agenda()
|
||||
target_date = resolve_target_date(today or date.today(), lessons, school_events)
|
||||
homeworks = fetcher.fetch_homework(target_date)
|
||||
except PipelineCriticalError:
|
||||
raise
|
||||
except Exception as exc:
|
||||
raise PipelineCriticalError(
|
||||
f"Récupération Pronote impossible : {redact_exception(exc)}", step="fetch"
|
||||
) from None
|
||||
|
||||
messages, warnings = _fetch_optional_messages(fetcher)
|
||||
return (
|
||||
FetchedPronoteData(
|
||||
lessons=lessons,
|
||||
homeworks=homeworks,
|
||||
school_events=school_events,
|
||||
messages=messages,
|
||||
target_date=target_date,
|
||||
),
|
||||
warnings,
|
||||
)
|
||||
32
pronote_sync/pipeline/steps/fetch_blog.py
Normal file
32
pronote_sync/pipeline/steps/fetch_blog.py
Normal file
@@ -0,0 +1,32 @@
|
||||
"""Étape de récupération non bloquante des articles RSS du collège."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from pronote_sync.models.blog import BlogArticle
|
||||
from pronote_sync.sources.blog.rss import BlogRSSClient
|
||||
from pronote_sync.sources.blog.state import BlogRSSState
|
||||
from pronote_sync.utils.redaction import redact_exception
|
||||
|
||||
|
||||
def fetch_blog_step(client: BlogRSSClient | None, state: BlogRSSState | None) -> list[BlogArticle]:
|
||||
"""Récupère les articles RSS nouveaux en conservant l'état du client.
|
||||
|
||||
:param client: Client RSS configuré, ou ``None`` lorsque le blog est désactivé.
|
||||
:param state: État de déduplication et de cache HTTP associé au run.
|
||||
:return: Nouveaux articles du blog.
|
||||
:rtype: list[BlogArticle]
|
||||
:raises RuntimeError: Si la récupération RSS injectée échoue.
|
||||
"""
|
||||
if client is None or state is None:
|
||||
return []
|
||||
try:
|
||||
etag, last_modified = state.get_cache_headers()
|
||||
result = client.fetch_and_parse(
|
||||
known_guids=state.get_known_guids(), etag=etag, last_modified=last_modified
|
||||
)
|
||||
if not result.not_modified:
|
||||
state.add_guids(article.id for article in result.articles)
|
||||
state.update_cache_headers(result.etag, result.last_modified)
|
||||
return list(result.articles)
|
||||
except Exception as exc:
|
||||
raise RuntimeError(f"Récupération du blog échouée : {redact_exception(exc)}") from None
|
||||
31
pronote_sync/pipeline/steps/normalize.py
Normal file
31
pronote_sync/pipeline/steps/normalize.py
Normal file
@@ -0,0 +1,31 @@
|
||||
"""Étape de normalisation et d'ordonnancement déterministe des données Pronote."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from datetime import datetime
|
||||
|
||||
from pronote_sync.models.pronote import PronoteData
|
||||
from pronote_sync.pipeline.steps.fetch import FetchedPronoteData
|
||||
|
||||
|
||||
def normalize_step(fetched: FetchedPronoteData, *, generated_at: datetime) -> PronoteData:
|
||||
"""Construit le contrat ``PronoteData`` dans un ordre déterministe.
|
||||
|
||||
:param fetched: Données brutes produites par :func:`fetch_step`.
|
||||
:param generated_at: Horodatage de l'exécution fourni par l'orchestrateur.
|
||||
:return: Données Pronote normalisées.
|
||||
:rtype: PronoteData
|
||||
"""
|
||||
return PronoteData(
|
||||
lessons=sorted(fetched.lessons, key=lambda lesson: (lesson.start, lesson.id)),
|
||||
homeworks=sorted(
|
||||
fetched.homeworks, key=lambda homework: (homework.due_on, homework.subject, homework.id)
|
||||
),
|
||||
school_events=sorted(
|
||||
fetched.school_events,
|
||||
key=lambda event: (event.from_date, event.to_date, event.kind.value, event.label),
|
||||
),
|
||||
messages=sorted(fetched.messages, key=lambda message: (message.date, message.id)),
|
||||
target_date=fetched.target_date,
|
||||
generated_at=generated_at,
|
||||
)
|
||||
17
pronote_sync/pipeline/steps/send.py
Normal file
17
pronote_sync/pipeline/steps/send.py
Normal file
@@ -0,0 +1,17 @@
|
||||
"""Étape d'envoi du digest sur le canal de notification."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from pronote_sync.channels.protocol import Channel
|
||||
from pronote_sync.models.xmpp import XmppMessage
|
||||
|
||||
|
||||
def send_step(channel: Channel, message: XmppMessage) -> bool:
|
||||
"""Envoie le digest et retourne le statut fourni par le canal.
|
||||
|
||||
:param channel: Canal de sortie configuré.
|
||||
:param message: Digest XMPP à transmettre.
|
||||
:return: ``True`` si l'envoi a réussi, ``False`` sinon.
|
||||
:rtype: bool
|
||||
"""
|
||||
return channel.send(message)
|
||||
21
pronote_sync/pipeline/steps/synthesis.py
Normal file
21
pronote_sync/pipeline/steps/synthesis.py
Normal file
@@ -0,0 +1,21 @@
|
||||
"""Étape de génération optionnelle de synthèse IA."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from pronote_sync.models.synthesis import SynthesisInput, SynthesisResult
|
||||
from pronote_sync.synthesis.provider import SynthesisProvider
|
||||
|
||||
|
||||
def synthesis_step(
|
||||
provider: SynthesisProvider | None, input_data: SynthesisInput
|
||||
) -> SynthesisResult | None:
|
||||
"""Génère une synthèse lorsque le fournisseur IA est activé.
|
||||
|
||||
:param provider: Fournisseur IA optionnel.
|
||||
:param input_data: Données à synthétiser.
|
||||
:return: Synthèse produite, ou ``None`` si le fournisseur est désactivé.
|
||||
:rtype: SynthesisResult | None
|
||||
"""
|
||||
if provider is None:
|
||||
return None
|
||||
return provider.generate(input_data)
|
||||
Reference in New Issue
Block a user