Cause racine (#25) : use_tls=True (défaut) activait le direct TLS sur le port 5222 (conventionnellement STARTTLS). Le client envoyait un ClientHello TLS sur un port attendant un stream XMPP en clair, le serveur ne voyait jamais l'identité configurée, et la session expirait après 30 s. Corrections : - Remplacer use_tls (bool) par tls_mode: Literal[direct|starttls|disabled] (défaut starttls, compatible port 5222). use_tls conservé comme alias déprécié avec DeprecationWarning. - Ajouter connect_timeout (15 s) et cleanup_timeout (10 s) distincts du timeout de session (30 s). - Gérer l'événement connection_failed de Slixmpp pour échouer rapidement au lieu d'attendre le timeout de session. - Borner await connect_future et await disconnect_future par leurs timeouts respectifs (anti-blocage). - Attendre connect_future et session_future conjointement (asyncio.wait, FIRST_COMPLETED) pour détecter connection_failed avant l'expiration du connect_timeout. - Annuler les tâches pending sur tous les chemins de retour, y compris CancelledError et Exception. - Redact tous les redact_exception avec extra_secrets=_secret_values(). - Construire ClientXMPP dans le try (contrat « never raises »). - Enrichir FakeClientXMPP avec modes connect/disconnect configurables. - 25 nouveaux tests (timeout connexion, connection_failed, cleanup bloqué, CancelledError, TLS mismatch, fuite secrets). Couverture 96 %. - Mettre à jour .env.example, GUIDE_DEV_PYTHON.md, README.LLM.md. Co-authored-by: OpenCode <opencode@antoineve.me>
488 lines
21 KiB
Python
488 lines
21 KiB
Python
"""Canal de sortie XMPP du pipeline ``pronote-sync``.
|
|
|
|
Ce module implémente le canal d'envoi de notifications XMPP : la classe
|
|
:class:`XmppChannel` envoie un message direct via ``slixmpp``
|
|
(:meth:`XmppChannel.send_async`), tandis que :class:`SyncXmppChannel`
|
|
fournit le point d'entrée synchrone unique utilisé par le pipeline. Le corps
|
|
du message est formaté en texte brut par ``_format_message`` (en-tête de date
|
|
cible puis sections emoji 📌📅📚💬📢) et chaque texte est assaini par
|
|
:func:`pronote_sync.utils.text.sanitize_plaintext` (SEC-XMPP-06).
|
|
|
|
Contrat d'erreur (D6) : le canal ne lève jamais :pyexc:`PipelineWarning` ;
|
|
en cas d'échec, il journalise la version expurgée de l'erreur et retourne
|
|
``False``. Le :pyexc:`PipelineWarning` est créé par l'étape pipeline, pas par
|
|
le canal.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import asyncio
|
|
import logging
|
|
|
|
from pydantic import SecretStr
|
|
from slixmpp import JID, ClientXMPP
|
|
|
|
from pronote_sync.config.settings import XmppSettings
|
|
from pronote_sync.models.blog import ExternalInfo
|
|
from pronote_sync.models.diff import AgendaChange, AgendaChangeType
|
|
from pronote_sync.models.homework import Homework
|
|
from pronote_sync.models.message import Message
|
|
from pronote_sync.models.xmpp import XmppMessage
|
|
from pronote_sync.utils.redaction import redact_exception, redact_secrets
|
|
from pronote_sync.utils.text import sanitize_plaintext
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
__all__ = ["XmppChannel", "SyncXmppChannel", "XmppMessage"]
|
|
|
|
|
|
async def _cancel_pending(
|
|
*futures: asyncio.Future[bool],
|
|
) -> None:
|
|
"""Annule les futures/tâches encore en attente et supprime le bruit.
|
|
|
|
À appeler avant chaque retour anticipé de :meth:`XmppChannel.send_async`
|
|
afin qu'aucune tentative de connexion ne survive au retour de la méthode.
|
|
|
|
:param futures: Futures ou tâches à annuler (les déjà terminées sont
|
|
ignorées pour la cancellation mais attendues pour purger l'attente).
|
|
:rtype: None
|
|
"""
|
|
for future in futures:
|
|
if not future.done():
|
|
future.cancel()
|
|
await asyncio.gather(*futures, return_exceptions=True)
|
|
|
|
|
|
def _secret_values(settings: XmppSettings) -> tuple[SecretStr | str, ...]:
|
|
"""Rassemble les secrets du canal XMPP pour le masquage des logs.
|
|
|
|
:param settings: Paramètres du canal XMPP.
|
|
:return: Valeurs sensibles (mot de passe, JID du bot, destinataire).
|
|
:rtype: tuple[SecretStr | str, ...]
|
|
"""
|
|
secrets: list[SecretStr | str] = []
|
|
if settings.jid is not None:
|
|
secrets.append(settings.jid)
|
|
if settings.password is not None:
|
|
secrets.append(settings.password)
|
|
if settings.to is not None:
|
|
secrets.append(settings.to)
|
|
return tuple(secrets)
|
|
|
|
|
|
def _format_synthesis(synthesis: str | None) -> str:
|
|
"""Formate la section synthèse du message XMPP.
|
|
|
|
:param synthesis: Texte de synthèse, ou ``None`` si absente.
|
|
:return: Section ``📌 Synthèse`` suivie de la synthèse (ou du texte par
|
|
défaut si aucune n'est disponible).
|
|
:rtype: str
|
|
"""
|
|
content = synthesis if synthesis else "Aucune synthèse disponible."
|
|
return f"📌 Synthèse\n{sanitize_plaintext(content)}"
|
|
|
|
|
|
def _format_changes(changes: tuple[AgendaChange, ...]) -> str:
|
|
"""Formate la section des changements d'agenda du message XMPP.
|
|
|
|
Distingue les ajouts, suppressions et modifications (U4). Pour un ajout,
|
|
les horaires du cours (``HH:MM-HH:MM``) sont inclus si le cours est
|
|
disponible.
|
|
|
|
:param changes: Liste des changements d'agenda.
|
|
:return: Section ``📅 Changements d'agenda`` avec une ligne par
|
|
changement (type, matière et détails).
|
|
:rtype: str
|
|
"""
|
|
if not changes:
|
|
body = "Aucun changement."
|
|
else:
|
|
lines: list[str] = []
|
|
for change in changes:
|
|
subject = "—"
|
|
if change.lesson is not None:
|
|
subject = change.lesson.subject
|
|
elif change.theoretical_lesson is not None:
|
|
subject = change.theoretical_lesson.subject
|
|
if change.type == AgendaChangeType.ADDED and change.lesson is not None:
|
|
times = (
|
|
f"{change.lesson.start.strftime('%H:%M')}-{change.lesson.end.strftime('%H:%M')}"
|
|
)
|
|
lines.append(f"• [Ajouté] {subject}: {change.details} ({times})")
|
|
elif change.type == AgendaChangeType.REMOVED:
|
|
lines.append(f"• [Supprimé] {subject}: {change.details}")
|
|
else:
|
|
lines.append(f"• [Modifié] {subject}: {change.details}")
|
|
body = "\n".join(lines)
|
|
return f"📅 Changements d'agenda\n{sanitize_plaintext(body)}"
|
|
|
|
|
|
def _format_homeworks(homeworks: tuple[Homework, ...]) -> str:
|
|
"""Formate la section des devoirs du message XMPP.
|
|
|
|
:param homeworks: Liste des devoirs.
|
|
:return: Section ``📚 Devoirs`` avec une ligne par devoir (matière,
|
|
texte et date d'échéance).
|
|
:rtype: str
|
|
"""
|
|
if not homeworks:
|
|
body = "Aucun devoir."
|
|
else:
|
|
lines = [
|
|
f"• {homework.subject}: {homework.text} "
|
|
f"(à rendre le {homework.due_on.strftime('%d/%m')})"
|
|
for homework in homeworks
|
|
]
|
|
body = "\n".join(lines)
|
|
return f"📚 Devoirs\n{sanitize_plaintext(body)}"
|
|
|
|
|
|
def _format_messages(messages: tuple[Message, ...]) -> str:
|
|
"""Formate la section des messages Pronote du message XMPP.
|
|
|
|
:param messages: Liste des messages/informations.
|
|
:return: Section ``💬 Messages`` avec une ligne par message (titre,
|
|
auteur et contenu) ; sans titre, seul l'auteur est affiché.
|
|
:rtype: str
|
|
"""
|
|
if not messages:
|
|
body = "Aucun message."
|
|
else:
|
|
lines: list[str] = []
|
|
for message in messages:
|
|
if message.title:
|
|
lines.append(f"• {message.title} ({message.author}): {message.content}")
|
|
else:
|
|
lines.append(f"• {message.author}: {message.content}")
|
|
body = "\n".join(lines)
|
|
return f"💬 Messages\n{sanitize_plaintext(body)}"
|
|
|
|
|
|
def _format_external_info(external_info: ExternalInfo | None) -> str:
|
|
"""Formate la section des informations diverses du message XMPP.
|
|
|
|
Regroupe uniquement les articles du blog et les autres informations
|
|
(``other_info``) : les messages Pronote (``pronote_messages``) sont
|
|
exclus car ils sont déjà transmis par la section des messages.
|
|
|
|
:param external_info: Informations externes agrégées, ou ``None``.
|
|
:return: Section ``📢 Informations diverses`` avec une ligne par élément.
|
|
:rtype: str
|
|
"""
|
|
if external_info is None:
|
|
body = "Aucune information."
|
|
else:
|
|
lines: list[str] = []
|
|
for article in external_info.blog_articles:
|
|
lines.append(f"• {article.title}: {article.content_text}")
|
|
for info in external_info.other_info:
|
|
lines.append(f"• {info}")
|
|
body = "\n".join(lines) if lines else "Aucune information."
|
|
return f"📢 Informations diverses\n{sanitize_plaintext(body)}"
|
|
|
|
|
|
class XmppChannel:
|
|
"""Canal d'envoi de messages XMPP via un compte bot dédié.
|
|
|
|
Envoie un message direct (``type="chat"``) au destinataire configuré en
|
|
utilisant :class:`slixmpp.ClientXMPP`. La connexion est établie à chaque
|
|
appel de :meth:`send_async` ; le constructeur n'effectue aucun accès
|
|
réseau.
|
|
|
|
Contrat d'erreur (D6) : :meth:`send_async` ne lève jamais
|
|
:pyexc:`PipelineWarning` ; en cas d'échec, elle journalise la version
|
|
expurgée de l'erreur et retourne ``False``. En mode ``dry_run``, aucun
|
|
client n'est créé.
|
|
|
|
:ivar settings: Paramètres XMPP (JID, mot de passe, destinataire, TLS).
|
|
:vartype settings: XmppSettings
|
|
:ivar dry_run: En mode ``dry_run``, aucun envoi n'est effectué.
|
|
:vartype dry_run: bool
|
|
"""
|
|
|
|
def __init__(self, settings: XmppSettings, dry_run: bool = False) -> None:
|
|
"""Initialise le canal XMPP sans connexion réseau.
|
|
|
|
:param settings: Paramètres de configuration du canal XMPP.
|
|
:param dry_run: Si ``True``, :meth:`send_async` journalise le message
|
|
formaté et retourne ``True`` sans se connecter.
|
|
"""
|
|
self.settings = settings
|
|
self.dry_run = dry_run
|
|
|
|
def _format_message(self, message: XmppMessage) -> str:
|
|
"""Formate un message XMPP en texte brut avec des sections emoji.
|
|
|
|
Produit le corps du message : un en-tête avec la date cible du
|
|
digest, puis les sections synthèse, changements d'agenda, devoirs,
|
|
messages et informations diverses. Chaque texte est assaini par
|
|
:func:`pronote_sync.utils.text.sanitize_plaintext` avant insertion
|
|
(SEC-XMPP-06).
|
|
|
|
:param message: Message final à formater.
|
|
:return: Corps du message en texte brut, prêt pour l'envoi.
|
|
:rtype: str
|
|
"""
|
|
sections = [
|
|
f"Digest du {message.target_date.strftime('%d/%m/%Y')}",
|
|
_format_synthesis(message.synthesis),
|
|
_format_changes(message.changes),
|
|
_format_homeworks(message.homeworks),
|
|
_format_messages(message.messages),
|
|
_format_external_info(message.external_info),
|
|
]
|
|
return "\n\n".join(sections)
|
|
|
|
async def send_async(self, message: XmppMessage) -> bool:
|
|
"""Exécute le flux asynchrone d'envoi XMPP (U2).
|
|
|
|
Connecte le client ``slixmpp`` avec un hôte et un port explicites,
|
|
configure TLS avant la connexion selon ``tls_mode`` (``direct``,
|
|
``starttls`` ou ``disabled``), attend la connexion sous
|
|
``connect_timeout`` puis l'un des événements ``session_start``,
|
|
``failed_auth``, ``connection_failed`` ou ``disconnected`` sous
|
|
``timeout`` avant d'envoyer un message direct ``chat`` au
|
|
destinataire configuré. La déconnexion est garantie par un bloc
|
|
``try/finally`` borné par ``cleanup_timeout``. Aucun secret n'est
|
|
journalisé (SEC-XMPP-02).
|
|
|
|
:param message: Message final à envoyer.
|
|
:return: ``True`` si l'envoi a réussi (ou a été simulé en dry-run),
|
|
``False`` sinon (destinataire manquant, timeout, échec
|
|
d'authentification, déconnexion ou erreur réseau).
|
|
:rtype: bool
|
|
"""
|
|
if self.dry_run:
|
|
formatted = self._format_message(message)
|
|
logger.info("XMPP dry-run: message would be sent")
|
|
return True
|
|
|
|
# Build JID with resource
|
|
jid_str = f"{self.settings.jid}/{self.settings.resource}"
|
|
recipient = JID(self.settings.to) if self.settings.to else None
|
|
if recipient is None:
|
|
logger.warning("Destinataire XMPP manquant.")
|
|
return False
|
|
|
|
# Client typed lazily: the construction is done inside the try block so that
|
|
# any error is caught and converted to ``False`` (channel contract: never raise)
|
|
client: ClientXMPP | None = None
|
|
session_future: asyncio.Future[bool] = asyncio.get_event_loop().create_future()
|
|
failure_kind = "disconnected"
|
|
# Declared before the ``try`` so the exception handlers (CancelledError and
|
|
# Exception) can cancel any task still pending from ``asyncio.wait()``
|
|
connect_future: asyncio.Future[bool] | None = None
|
|
session_task: asyncio.Future[bool] | None = None
|
|
|
|
def on_session_start(event: object) -> None:
|
|
if not session_future.done():
|
|
session_future.set_result(True)
|
|
|
|
def on_failed_auth(event: object) -> None:
|
|
nonlocal failure_kind
|
|
if not session_future.done():
|
|
failure_kind = "failed_auth"
|
|
session_future.set_result(False)
|
|
|
|
def on_connection_failed(event: object) -> None:
|
|
nonlocal failure_kind
|
|
if not session_future.done():
|
|
failure_kind = "connection_failed"
|
|
session_future.set_result(False)
|
|
|
|
def on_disconnected(event: object) -> None:
|
|
if not session_future.done():
|
|
session_future.set_result(False)
|
|
|
|
try:
|
|
# Create typed client
|
|
client = ClientXMPP(
|
|
jid_str,
|
|
self.settings.password.get_secret_value() if self.settings.password else "",
|
|
)
|
|
|
|
# Configure TLS BEFORE connect (canonical tls_mode)
|
|
match self.settings.tls_mode:
|
|
case "direct":
|
|
client.enable_direct_tls = True
|
|
client.enable_starttls = False
|
|
case "starttls":
|
|
client.enable_direct_tls = False
|
|
client.enable_starttls = True
|
|
case "disabled":
|
|
client.enable_direct_tls = False
|
|
client.enable_starttls = False
|
|
|
|
# Register handlers
|
|
client.add_event_handler("session_start", on_session_start)
|
|
client.add_event_handler("failed_auth", on_failed_auth)
|
|
client.add_event_handler("connection_failed", on_connection_failed)
|
|
client.add_event_handler("disconnected", on_disconnected)
|
|
|
|
# Connect with explicit host and port
|
|
connect_future = asyncio.ensure_future(
|
|
client.connect(self.settings.host, self.settings.port)
|
|
)
|
|
|
|
# Wrap the session future in a task so that cancelling pending tasks
|
|
# during the concurrent wait never cancels ``session_future`` itself
|
|
async def _await_session() -> bool:
|
|
return await session_future
|
|
|
|
session_task = asyncio.ensure_future(_await_session())
|
|
|
|
# Wait for the connection and the session event concurrently, bounded by
|
|
# connect_timeout as the global time limit: a ``connection_failed`` event
|
|
# can thus trigger an early return before the connect timeout expires
|
|
done, _pending = await asyncio.wait(
|
|
{connect_future, session_task},
|
|
timeout=self.settings.connect_timeout,
|
|
return_when=asyncio.FIRST_COMPLETED,
|
|
)
|
|
|
|
if session_task in done:
|
|
if not session_future.result():
|
|
# ``connection_failed``/``failed_auth``/``disconnected`` fired
|
|
# before the connection was resolved: immediate failure (fail fast)
|
|
if failure_kind == "connection_failed":
|
|
logger.warning("Échec de connexion réseau XMPP.")
|
|
else:
|
|
logger.warning("Échec d'authentification ou déconnexion XMPP.")
|
|
await _cancel_pending(connect_future, session_task)
|
|
return False
|
|
# ``session_start`` fired: the connection succeeded even if the
|
|
# connect future is still pending; proceed to send the message
|
|
await _cancel_pending(connect_future, session_task)
|
|
elif connect_future in done:
|
|
# The connection resolved: surface a connect error (redacted) if any
|
|
if not connect_future.cancelled():
|
|
connect_exc = connect_future.exception()
|
|
if connect_exc is not None and isinstance(connect_exc, Exception):
|
|
logger.warning(
|
|
"Échec de connexion XMPP : %s",
|
|
redact_exception(
|
|
connect_exc,
|
|
extra_secrets=_secret_values(self.settings),
|
|
),
|
|
)
|
|
await _cancel_pending(connect_future, session_task)
|
|
return False
|
|
|
|
# Connection established: wait for a session event under ``timeout``
|
|
try:
|
|
success = await asyncio.wait_for(
|
|
asyncio.shield(session_future), timeout=self.settings.timeout
|
|
)
|
|
except TimeoutError:
|
|
logger.warning("Délai d'attente de session XMPP dépassé.")
|
|
await _cancel_pending(connect_future, session_task)
|
|
return False
|
|
if not success:
|
|
if failure_kind == "connection_failed":
|
|
logger.warning("Échec de connexion réseau XMPP.")
|
|
else:
|
|
logger.warning("Échec d'authentification ou déconnexion XMPP.")
|
|
await _cancel_pending(connect_future, session_task)
|
|
return False
|
|
else:
|
|
# connect_timeout expired: cancel everything and fail fast
|
|
logger.warning(
|
|
"Délai de connexion XMPP dépassé (%ss).", self.settings.connect_timeout
|
|
)
|
|
await _cancel_pending(connect_future, session_task)
|
|
return False
|
|
|
|
# Send the message
|
|
formatted = self._format_message(message)
|
|
client.send_message(mto=JID(self.settings.to), mbody=formatted, mtype="chat")
|
|
return True
|
|
|
|
except asyncio.CancelledError:
|
|
# Contrat du canal : toujours retourner un booléen, même en cas
|
|
# d'annulation de la tâche appelante (cleanup exécuté par le finally).
|
|
logger.debug("Envoi XMPP annulé avant la fin de l'opération.")
|
|
pending = [f for f in (connect_future, session_task) if f is not None]
|
|
if pending:
|
|
await _cancel_pending(*pending)
|
|
return False
|
|
except Exception as exc:
|
|
redacted = redact_exception(exc, extra_secrets=_secret_values(self.settings))
|
|
extra = _secret_values(self.settings)
|
|
logger.warning("Erreur XMPP: %s", redact_secrets(redacted, extra_secrets=extra))
|
|
pending = [f for f in (connect_future, session_task) if f is not None]
|
|
if pending:
|
|
await _cancel_pending(*pending)
|
|
return False
|
|
finally:
|
|
if client is not None:
|
|
try:
|
|
disconnect_future = client.disconnect()
|
|
await asyncio.wait_for(disconnect_future, timeout=self.settings.cleanup_timeout)
|
|
except asyncio.CancelledError:
|
|
logger.debug("Déconnexion XMPP annulée.")
|
|
except TimeoutError:
|
|
logger.debug(
|
|
"Délai de déconnexion XMPP dépassé (%ss), abandon.",
|
|
self.settings.cleanup_timeout,
|
|
)
|
|
except Exception as cleanup_exc:
|
|
logger.debug(
|
|
"Erreur lors de la déconnexion XMPP : %s",
|
|
redact_exception(cleanup_exc, extra_secrets=_secret_values(self.settings)),
|
|
)
|
|
|
|
|
|
class SyncXmppChannel:
|
|
"""Point d'entrée synchrone unique du canal XMPP pour le pipeline (U3).
|
|
|
|
Enveloppe une instance de :class:`XmppChannel` pour offrir une interface
|
|
synchrone conforme au :class:`~pronote_sync.channels.protocol.Channel`.
|
|
:meth:`send` délègue à :func:`asyncio.run` et ne lève jamais : toute
|
|
erreur est journalisée de façon expurgée et convertie en retour
|
|
``False`` (D6). En mode ``dry_run``, aucun client ``slixmpp`` n'est créé.
|
|
|
|
:ivar settings: Paramètres XMPP.
|
|
:vartype settings: XmppSettings
|
|
:ivar dry_run: Mode simulation (aucun envoi réseau).
|
|
:vartype dry_run: bool
|
|
"""
|
|
|
|
def __init__(self, settings: XmppSettings, dry_run: bool = False) -> None:
|
|
"""Initialise le point d'entrée synchrone et son canal interne.
|
|
|
|
:param settings: Paramètres de configuration du canal XMPP.
|
|
:param dry_run: Si ``True``, l'envoi est simulé.
|
|
"""
|
|
self.settings = settings
|
|
self.dry_run = dry_run
|
|
self._channel = XmppChannel(settings, dry_run)
|
|
|
|
def send(self, message: XmppMessage) -> bool:
|
|
"""Envoie un message XMPP de façon synchrone et sans lever.
|
|
|
|
En mode ``dry_run``, le message formaté (expurgé de ses secrets) est
|
|
journalisé et la méthode retourne ``True`` sans créer de client XMPP.
|
|
Sinon, le flux asynchrone :meth:`XmppChannel.send_async` est exécuté
|
|
via :func:`asyncio.run` ; toute exception est journalisée sous forme
|
|
expurgée et convertie en retour ``False``. La méthode ne lève jamais
|
|
(D6).
|
|
|
|
:param message: Message final à envoyer.
|
|
:return: ``True`` si l'envoi a réussi (ou a été simulé en dry-run),
|
|
``False`` sinon.
|
|
:rtype: bool
|
|
"""
|
|
if self.dry_run:
|
|
formatted = self._channel._format_message(message)
|
|
redacted = redact_secrets(formatted, extra_secrets=_secret_values(self.settings))
|
|
logger.info("XMPP : dry-run, message non envoyé : %s", redacted)
|
|
return True
|
|
try:
|
|
return asyncio.run(self._channel.send_async(message))
|
|
except Exception as exc:
|
|
redacted = redact_exception(exc, extra_secrets=_secret_values(self.settings))
|
|
redacted = redact_secrets(redacted, extra_secrets=_secret_values(self.settings))
|
|
logger.warning("XMPP : erreur lors de l'envoi synchrone : %s", redacted)
|
|
return False
|