This commit was merged in pull request #26.
This commit is contained in:
@@ -36,6 +36,24 @@ 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.
|
||||
|
||||
@@ -220,11 +238,14 @@ class XmppChannel:
|
||||
"""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, puis attend l'un des événements
|
||||
``session_start``, ``failed_auth`` ou ``disconnected`` sous un
|
||||
timeout unique avant d'envoyer un message direct ``chat`` au
|
||||
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``. Aucun secret n'est journalisé (SEC-XMPP-02).
|
||||
``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),
|
||||
@@ -244,55 +265,132 @@ class XmppChannel:
|
||||
logger.warning("Destinataire XMPP manquant.")
|
||||
return False
|
||||
|
||||
# Create typed client
|
||||
client = ClientXMPP(
|
||||
jid_str,
|
||||
self.settings.password.get_secret_value() if self.settings.password else "",
|
||||
)
|
||||
|
||||
# Configure TLS BEFORE connect
|
||||
if self.settings.use_tls:
|
||||
# TLS direct (port 5223 typically)
|
||||
client.enable_direct_tls = True
|
||||
client.enable_starttls = False
|
||||
else:
|
||||
# STARTTLS (port 5222 typically)
|
||||
client.enable_starttls = True
|
||||
client.enable_direct_tls = False
|
||||
|
||||
# Register handlers
|
||||
# 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)
|
||||
|
||||
client.add_event_handler("session_start", on_session_start)
|
||||
client.add_event_handler("failed_auth", on_failed_auth)
|
||||
client.add_event_handler("disconnected", on_disconnected)
|
||||
|
||||
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 = client.connect(self.settings.host, self.settings.port)
|
||||
await connect_future # connect() returns a Future, not a coroutine
|
||||
connect_future = asyncio.ensure_future(
|
||||
client.connect(self.settings.host, self.settings.port)
|
||||
)
|
||||
|
||||
# Wait for one of the three events under a single timeout
|
||||
try:
|
||||
success = await asyncio.wait_for(session_future, timeout=self.settings.timeout)
|
||||
except TimeoutError:
|
||||
logger.warning("Délai d'attente de session XMPP dépassé.")
|
||||
return False
|
||||
# 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
|
||||
|
||||
if not success:
|
||||
logger.warning("Échec d'authentification ou déconnexion XMPP.")
|
||||
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
|
||||
@@ -300,19 +398,39 @@ class XmppChannel:
|
||||
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)
|
||||
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:
|
||||
try:
|
||||
disconnect_future = client.disconnect()
|
||||
await disconnect_future
|
||||
except Exception as cleanup_exc:
|
||||
logger.debug(
|
||||
"Erreur lors de la déconnexion XMPP: %s", redact_exception(cleanup_exc)
|
||||
)
|
||||
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:
|
||||
@@ -363,7 +481,7 @@ class SyncXmppChannel:
|
||||
try:
|
||||
return asyncio.run(self._channel.send_async(message))
|
||||
except Exception as exc:
|
||||
redacted = redact_exception(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
|
||||
|
||||
Reference in New Issue
Block a user