Coverage for slidge/core/gateway.py: 57%
532 statements
« prev ^ index » next coverage.py v7.15.2, created at 2026-09-29 05:05 +0000
« prev ^ index » next coverage.py v7.15.2, created at 2026-09-29 05:05 +0000
1"""
2This module extends slixmpp.ComponentXMPP to make writing new LegacyClients easier
3"""
5from __future__ import annotations
7import abc
8import asyncio
9import logging
10import re
11import sys
12import tempfile
13from collections.abc import (
14 Awaitable,
15 Callable,
16 Iterable,
17 Sequence,
18)
19from copy import copy
20from datetime import UTC, datetime, timedelta
21from pathlib import Path
22from typing import (
23 Any,
24 ClassVar,
25 Concatenate,
26 Generic,
27 ParamSpec,
28 TypeVar,
29)
31import aiohttp
32from slixmpp import JID, ComponentXMPP, Iq
33from slixmpp.exceptions import IqError, IqTimeout, XMPPError
34from slixmpp.plugins.xep_0004.stanza.form import Form
35from slixmpp.plugins.xep_0060.stanza import OwnerAffiliation
36from slixmpp.plugins.xep_0356.privilege import PrivilegedIqError
37from slixmpp.types import MessageTypes
38from slixmpp.xmlstream.xmlstream import NotConnectedError
39from sqlalchemy.orm import Session as OrmSession
41import slidge.command.categories
42from slidge.command import BUILTIN_COMMANDS
43from slidge.command.adhoc import AdhocProvider
44from slidge.command.admin import Exec
45from slidge.command.base import (
46 Command,
47 CommandBase,
48 ContactCommand,
49 FormField,
50 MUCCommand,
51)
52from slidge.command.chat_command import ChatCommandProvider
53from slidge.command.register import AltRegistrationFlow, Register, RegistrationType
54from slidge.contact import LegacyContact
55from slidge.core import config
56from slidge.core.attachment_upload import AttachmentUploader
57from slidge.core.dispatcher.session_dispatcher import SessionDispatcher
58from slidge.core.mixins.avatar import convert_avatar
59from slidge.core.mixins.message import ContentMessageMixin, InviteMixin
60from slidge.core.pubsub import PubSubComponent
61from slidge.db import GatewayUser, SlidgeStore
62from slidge.db.avatar import CachedAvatar, avatar_cache
63from slidge.db.meta import JSONSerializable
64from slidge.db.models import Attachment
65from slidge.group import LegacyMUC
66from slidge.slixfix.delivery_receipt import DeliveryReceipt
67from slidge.slixfix.roster import RosterBackend
68from slidge.util.types import (
69 AnySession,
70 Avatar,
71 MessageOrPresenceTypeVar,
72 RegistrationValidationCoroutine,
73 SessionType_co,
74)
75from slidge.util.util import derive_wired_class
77T = TypeVar("T")
78P = ParamSpec("P")
81class BaseGateway(
82 ComponentXMPP,
83 InviteMixin,
84 ContentMessageMixin,
85 abc.ABC,
86 Generic[SessionType_co], # noqa: UP046 (mypy cannot infer variance with PEP 695 syntax)
87):
88 """
89 The gateway component, handling registrations and un-registrations.
91 On slidge launch, a singleton is instantiated, and it will be made available
92 to public classes such :class:`.LegacyContact` or :class:`.BaseSession` as the
93 ``.xmpp`` attribute.
95 Must be subclassed by a legacy module to set up various aspects of the XMPP
96 component behaviour, such as its display name or welcome message, via
97 class attributes :attr:`.COMPONENT_NAME` :attr:`.WELCOME_MESSAGE`.
99 Abstract methods related to the registration process must be overriden
100 for a functional :term:`Legacy Module`:
102 - :meth:`.validate`
103 - :meth:`.validate_two_factor_code`
104 - :meth:`.get_qr_text`
105 - :meth:`.confirm_qr`
107 NB: Not all of these must be overridden, it depends on the
108 :attr:`REGISTRATION_TYPE`.
110 The other methods, such as :meth:`.send_text` or :meth:`.react` are the same
111 as those of :class:`.LegacyContact` and :class:`.LegacyParticipant`, because
112 the component itself is also a "messaging actor", ie, an :term:`XMPP Entity`.
113 For these methods, you need to specify the JID of the recipient with the
114 `mto` parameter.
116 Since it inherits from :class:`slixmpp.componentxmpp.ComponentXMPP`,you also
117 have a hand on low-level XMPP interactions via slixmpp methods, e.g.:
119 .. code-block:: python
121 self.send_presence(
122 pfrom="somebody@component.example.com",
123 pto="someonwelse@anotherexample.com",
124 )
126 However, you should not need to do so often since the classes of the plugin
127 API provides higher level abstractions around most commonly needed use-cases, such
128 as sending messages, or displaying a custom status.
130 """
132 COMPONENT_NAME: str = NotImplemented
133 """Name of the component, as seen in service discovery by XMPP clients"""
134 COMPONENT_TYPE: str = ""
135 """Type of the gateway, should follow https://xmpp.org/registrar/disco-categories.html"""
136 COMPONENT_AVATAR: Avatar | Path | str | None = None
137 """
138 Path, bytes or URL used by the component as an avatar.
139 """
141 REGISTRATION_FIELDS: ClassVar[Sequence[FormField]] = [
142 FormField(var="username", label="User name", required=True),
143 FormField(var="password", label="Password", required=True, private=True),
144 ]
145 """
146 Iterable of fields presented to the gateway user when registering using :xep:`0077`
147 `extended <https://xmpp.org/extensions/xep-0077.html#extensibility>`_ by :xep:`0004`.
148 """
149 REGISTRATION_INSTRUCTIONS: str = "Enter your credentials"
150 """
151 The text presented to a user who wants to register (or modify) their
152 :term:`legacy <Legacy>` account configuration.
153 """
154 REGISTRATION_TYPE: RegistrationType = RegistrationType.SINGLE_STEP_FORM
155 """
156 This attribute determines how users register to the gateway, ie, how they
157 login to the :term:`legacy network <Legacy Network>`.
158 The credentials are then stored persistently, so this process should happen
159 once per user (unless they unregister).
161 The registration process always start with a basic data form (:xep:`0004`)
162 presented to the user.
163 But the legacy login flow might require something more sophisticated, see
164 :class:`.RegistrationType` for more details.
165 """
167 REGISTRATION_2FA_TITLE = "Enter your 2FA code"
168 REGISTRATION_2FA_INSTRUCTIONS = (
169 "You should have received something via email or SMS, or something"
170 )
171 REGISTRATION_QR_INSTRUCTIONS = "Flash this code or follow this link"
173 ALTERNATIVE_REGISTRATION_FLOWS: ClassVar[list[AltRegistrationFlow]] = []
175 PREFERENCES: ClassVar[list[FormField]] = [
176 FormField(
177 var="sync_presence",
178 label="Propagate your XMPP presence to the legacy network.",
179 value="true",
180 required=True,
181 type="boolean",
182 ),
183 FormField(
184 var="sync_avatar",
185 label="Propagate your XMPP avatar to the legacy network.",
186 value="true",
187 required=True,
188 type="boolean",
189 ),
190 FormField(
191 var="always_invite_when_adding_bookmarks",
192 label="Always send invitations to join MUCs.",
193 value="false`",
194 required=True,
195 type="boolean",
196 ),
197 FormField(
198 var="last_seen_fallback",
199 label="Use contact presence status message to show when they were last seen.",
200 value="true",
201 required=True,
202 type="boolean",
203 ),
204 FormField(
205 var="roster_push",
206 label="Add contacts to your roster.",
207 value="true",
208 required=True,
209 type="boolean",
210 ),
211 FormField(
212 var="reaction_fallback",
213 label="Receive fallback messages for reactions (for legacy XMPP clients)",
214 value="false",
215 required=True,
216 type="boolean",
217 ),
218 ]
220 ROSTER_GROUP: str = "slidge"
221 """
222 Name of the group assigned to a :class:`.LegacyContact` automagically
223 added to the :term:`User`'s roster with :meth:`.LegacyContact.add_to_roster`.
224 """
225 WELCOME_MESSAGE = (
226 "Thank you for registering. Type 'help' to list the available commands, "
227 "or just start messaging away!"
228 )
229 """
230 A welcome message displayed to users on registration.
231 This is useful notably for clients that don't consider component JIDs as a
232 valid recipient in their UI, yet still open a functional chat window on
233 incoming messages from components.
234 """
236 SEARCH_FIELDS: ClassVar[Sequence[FormField]] = [
237 FormField(var="first", label="First name", required=True),
238 FormField(var="last", label="Last name", required=True),
239 FormField(var="phone", label="Phone number", required=False),
240 ]
241 """
242 Fields used for searching items via the component, through :xep:`0055` (jabber search).
243 A common use case is to allow users to search for legacy contacts by something else than
244 their usernames, eg their phone number.
246 Plugins should implement search by overriding :meth:`.BaseSession.search`
247 (restricted to registered users).
249 If there is only one field, it can also be used via the ``jabber:iq:gateway`` protocol
250 described in :xep:`0100`. Limitation: this only works if the search request returns
251 one result item, and if this item has a 'jid' var.
252 """
253 SEARCH_TITLE: str = "Search for legacy contacts"
254 """
255 Title of the search form.
256 """
257 SEARCH_INSTRUCTIONS: str = ""
258 """
259 Instructions of the search form.
260 """
262 MARK_ALL_MESSAGES = False
263 """
264 Set this to True for :term:`legacy networks <Legacy Network>` that expects
265 read marks for *all* messages and not just the latest one that was read
266 (as most XMPP clients will only send a read mark for the latest msg).
267 """
269 PROPER_RECEIPTS = False
270 """
271 Set this to True if the legacy service provides a real equivalent of message delivery receipts
272 (:xep:`0184`), meaning that there is an event thrown when the actual device of a contact receives
273 a message. Make sure to call Contact.received() adequately if this is set to True.
274 """
276 GROUPS = False
277 """
278 This must be set to True if this gateway supports groups.
279 """
280 SPACES = False
281 """
282 This must be set to True if this gateway supports spaces, cf :xep:`0503`.
283 """
284 THREADS: bool = True
285 """
286 This must be set to True if this gateway supports threads, to enable
287 storing an XMPP/legacy thread ID mapping when necessary.
288 """
290 mtype: MessageTypes = "chat"
291 is_group = False
292 _can_send_carbon = False
293 store: SlidgeStore
295 session_cls: type[SessionType_co]
296 """Concrete :class:`.BaseSession` subclass of this legacy module.
298 Derived automatically from the generic parameter, e.g.,
299 ``class Gateway(BaseGateway[Session])``.
300 """
302 def __init_subclass__(cls, **kwargs: object) -> None:
303 super().__init_subclass__(**kwargs)
304 derive_wired_class(cls, BaseGateway, "session_cls")
306 COMMANDS: ClassVar[tuple[type[CommandBase], ...]] = BUILTIN_COMMANDS
307 """Commands (adhoc and chatbot) this gateway provides.
309 Subclasses may override this with additional custom commands.
310 E.g.: ``COMMANDS = BaseGateway.COMMANDS + commands_from_module(command)``.
311 """
313 http: aiohttp.ClientSession
314 avatar: CachedAvatar | None = None
316 def __init__(self) -> None:
317 if getattr(self, "session_cls", None) is None:
318 raise RuntimeError(
319 f"{type(self).__name__}.session_cls is not set. Your gateway"
320 " class must be parameterized with its BaseSession subclass,"
321 " e.g. 'class Gateway(BaseGateway[Session])'."
322 )
323 if config.COMPONENT_NAME:
324 self.COMPONENT_NAME = config.COMPONENT_NAME
325 if config.WELCOME_MESSAGE:
326 self.WELCOME_MESSAGE = config.WELCOME_MESSAGE
327 self.log = log
328 self.datetime_started = datetime.now(tz=UTC)
329 # FIXME: ugly hack to work with the BaseSender mixin :/
330 self.xmpp = self
331 self.default_ns = "jabber:component:accept"
332 super().__init__(
333 config.JID,
334 config.SECRET,
335 config.SERVER,
336 config.PORT,
337 plugin_whitelist=SLIXMPP_PLUGINS,
338 plugin_config={
339 "xep_0077": {
340 "form_fields": None,
341 "form_instructions": self.REGISTRATION_INSTRUCTIONS,
342 "enable_subscription": self.REGISTRATION_TYPE
343 == RegistrationType.SINGLE_STEP_FORM,
344 },
345 "xep_0100": {
346 "component_name": self.COMPONENT_NAME,
347 "type": self.COMPONENT_TYPE,
348 },
349 "xep_0184": {
350 "auto_ack": False,
351 "auto_request": False,
352 },
353 "xep_0363": {
354 "upload_service": config.UPLOAD_SERVICE,
355 },
356 },
357 fix_error_ns=True,
358 )
359 self.loop.set_exception_handler(self.__exception_handler)
360 self.loop.create_task(self.__set_http(), name="set http client")
361 self.has_crashed: bool = False
362 self.use_origin_id = False
364 if config.USER_JID_VALIDATOR is None:
365 config.USER_JID_VALIDATOR = f".*@{self.infer_real_domain()}"
366 log.info(
367 "No USER_JID_VALIDATOR was set, using '%s'.",
368 config.USER_JID_VALIDATOR,
369 )
370 self.jid_validator: re.Pattern[str] = re.compile(config.USER_JID_VALIDATOR)
371 self.qr_pending_registrations = dict[
372 str, asyncio.Future[JSONSerializable | None]
373 ]()
375 self.register_plugins()
376 self.__setup_session_cls()
378 self.get_session_from_stanza = self.session_cls.from_stanza
379 self.get_session_from_user = self.session_cls.from_user
381 self.__register_slixmpp_events()
382 self.__register_slixmpp_api()
383 self.roster.set_backend(RosterBackend(self))
385 self.register_plugin("pubsub", {"component_name": self.COMPONENT_NAME})
386 self.pubsub: PubSubComponent = self.plugin["pubsub"] # type:ignore[typeddict-item]
387 self.delivery_receipt = DeliveryReceipt(self)
389 # with this we receive user avatar updates
390 self.plugin["xep_0030"].add_feature("urn:xmpp:avatar:metadata+notify")
392 self.plugin["xep_0030"].add_feature("urn:xmpp:chat-markers:0")
394 if self.GROUPS:
395 self.plugin["xep_0030"].add_feature("http://jabber.org/protocol/muc")
396 self.plugin["xep_0030"].add_feature(self.plugin["xep_0463"].stanza.NS)
397 self.plugin["xep_0030"].add_feature("urn:xmpp:mam:2")
398 self.plugin["xep_0030"].add_feature("urn:xmpp:mam:2#extended")
399 self.plugin["xep_0030"].add_feature(self.plugin["xep_0421"].namespace)
400 self.plugin["xep_0030"].add_feature(self.plugin["xep_0317"].stanza.NS)
401 self.plugin["xep_0030"].add_identity(
402 category="conference",
403 name=self.COMPONENT_NAME,
404 itype="text",
405 jid=self.boundjid,
406 )
407 if self.SPACES:
408 self.plugin["xep_0030"].add_feature("urn:xmpp:spaces:0")
410 self.__adhoc_handler = AdhocProvider(self)
411 self.__chat_commands_handler = ChatCommandProvider(self)
413 self.__dispatcher = SessionDispatcher(self)
414 self.get_caps_ver = self.__dispatcher.get_caps_ver
416 self.__register_commands()
418 def __setup_session_cls(self) -> None:
419 contact_cls = self.session_cls.roster_cls.contact_cls
421 if contact_cls.REACTIONS_SINGLE_EMOJI:
422 form = Form()
423 form["type"] = "result"
424 form.add_field(
425 "FORM_TYPE", "hidden", value="urn:xmpp:reactions:0:restrictions"
426 )
427 form.add_field("max_reactions_per_user", value="1", type="text-single")
428 form.add_field("scope", value="domain")
429 self.plugin["xep_0128"].add_extended_info(data=form)
431 self.session_cls.xmpp = self
433 @property
434 def _uploader(self) -> AttachmentUploader:
435 # the component itself is not bound to any user session
436 return AttachmentUploader(self, None)
438 async def kill_session(self, jid: JID) -> None:
439 await self.session_cls.kill_by_jid(jid)
441 async def __set_http(self) -> None:
442 self.http = aiohttp.ClientSession()
443 if getattr(self, "_test_mode", False):
444 return
445 avatar_cache.http = self.http
447 def __register_commands(self) -> None:
448 for cls in self.COMMANDS:
449 if any(x is NotImplemented for x in [cls.CHAT_COMMAND, cls.NODE, cls.NAME]):
450 raise TypeError(
451 f"{cls} is listed in {type(self).__name__}.COMMANDS but"
452 " does not set NAME, NODE and CHAT_COMMAND."
453 )
454 if issubclass(cls, ContactCommand):
455 LegacyContact.commands[cls.NODE] = cls
456 LegacyContact.commands_chat[cls.CHAT_COMMAND] = cls
457 continue
458 if issubclass(cls, MUCCommand):
459 LegacyMUC.commands[cls.NODE] = cls
460 LegacyMUC.commands_chat[cls.CHAT_COMMAND] = cls
461 continue
462 if not issubclass(cls, Command):
463 raise TypeError(
464 f"{cls} is listed in {type(self).__name__}.COMMANDS but is"
465 " not a Command, ContactCommand or MUCCommand subclass."
466 )
467 if cls is Exec:
468 if config.DEV_MODE:
469 log.warning(r"/!\ DEV MODE ENABLED /!\\")
470 else:
471 continue
472 if cls.CATEGORY == slidge.command.categories.GROUPS and not self.GROUPS:
473 continue
474 if cls.CATEGORY == slidge.command.categories.SPACES and not self.SPACES:
475 continue
476 c = cls(self)
477 log.debug("Registering %s", cls)
478 self.__adhoc_handler.register(c)
479 self.__chat_commands_handler.register(c)
481 for i, alt in enumerate(self.ALTERNATIVE_REGISTRATION_FLOWS):
483 class AltRegister(Register):
484 NAME = alt.title
485 NODE = f"https://slidge.im/register/alt{i}"
486 CHAT_COMMAND = f"register-{i}"
487 _instructions = alt.instructions
488 _type = alt.type
489 _fields = alt.fields
490 _validate = staticmethod(alt.validate) if alt.validate else None
492 inst = AltRegister(self)
493 log.debug("Registering alternative registration flow: %s", alt.title)
494 self.__adhoc_handler.register(inst)
495 self.__chat_commands_handler.register(inst)
497 def __exception_handler(
498 self, loop: asyncio.AbstractEventLoop, context: dict[Any, Any]
499 ) -> None:
500 """
501 Called when a task created by loop.create_task() raises an Exception
503 :param loop:
504 :param context:
505 :return:
506 """
507 log.debug("Context in the exception handler: %s", context)
508 exc = context.get("exception")
509 if exc is None:
510 log.debug("No exception in this context: %s", context)
511 elif isinstance(exc, SystemExit):
512 log.debug("SystemExit called in an asyncio task")
513 else:
514 log.error("Crash in an asyncio task: %s", context)
515 log.exception("Crash in task", exc_info=exc)
516 self.has_crashed = True
517 loop.stop()
519 def __register_slixmpp_events(self) -> None:
520 self.del_event_handler("presence_subscribe", self._handle_subscribe)
521 self.del_event_handler("presence_unsubscribe", self._handle_unsubscribe)
522 self.del_event_handler("presence_subscribed", self._handle_subscribed)
523 self.del_event_handler("presence_unsubscribed", self._handle_unsubscribed)
524 self.del_event_handler(
525 "roster_subscription_request", self._handle_new_subscription
526 )
527 self.del_event_handler("presence_probe", self._handle_probe)
528 self.add_event_handler("session_start", self.__on_session_start)
530 def __register_slixmpp_api(self) -> None:
531 def with_session(
532 func: Callable[Concatenate[OrmSession, P], T], commit: bool = True
533 ) -> Callable[P, T]:
534 def wrapped(*a: P.args, **kw: P.kwargs) -> T:
535 with self.store.session() as orm:
536 res = func(orm, *a, **kw)
537 if commit:
538 orm.commit()
539 return res
541 return wrapped
543 self.plugin["xep_0231"].api.register(
544 with_session(self.store.bob.get_bob, False), "get_bob"
545 )
546 self.plugin["xep_0231"].api.register(
547 with_session(self.store.bob.set_bob), "set_bob"
548 )
549 self.plugin["xep_0231"].api.register(
550 with_session(self.store.bob.del_bob), "del_bob"
551 )
553 @property # type:ignore[override]
554 def jid(self) -> JID:
555 # Override to avoid slixmpp deprecation warnings.
556 return self.boundjid
558 @jid.setter
559 def jid(self, jid: JID) -> None:
560 raise RuntimeError
562 async def __on_session_start(self, event: object) -> None:
563 log.debug("Gateway session start: %s", event)
565 await self.__setup_attachments()
567 # prevents XMPP clients from considering the gateway as an HTTP upload
568 disco = self.plugin["xep_0030"]
569 await disco.del_feature(feature="urn:xmpp:http:upload:0", jid=self.boundjid)
570 await self.plugin["xep_0115"].update_caps(jid=self.boundjid)
572 if self.COMPONENT_AVATAR is not None:
573 log.debug("Setting gateway avatar…")
574 avatar = convert_avatar(self.COMPONENT_AVATAR, "!!---slidge---special---")
575 assert avatar is not None
576 try:
577 cached_avatar = await avatar_cache.get(avatar)
578 except Exception as e:
579 log.exception("Could not set the component avatar.", exc_info=e)
580 cached_avatar = None
581 else:
582 assert cached_avatar is not None
583 self.avatar = cached_avatar
584 else:
585 cached_avatar = None
587 with self.store.session() as orm:
588 users = orm.query(GatewayUser).all()
589 for user in users:
590 # TODO: before this, we should check if the user has removed us from their roster
591 # while we were offline and trigger unregister from there. Presence probe does not seem
592 # to work in this case, there must be another way. privileged entity could be used
593 # as last resort.
594 try:
595 await self["xep_0100"].add_component_to_roster(user.jid)
596 await self.__add_component_to_mds_whitelist(user.jid)
597 except (IqError, IqTimeout) as e:
598 # TODO: remove the user when this happens? or at least
599 # this can happen when the user has unsubscribed from the XMPP server
600 log.warning(
601 "Error with user %s, not logging them automatically",
602 user,
603 exc_info=e,
604 )
605 continue
606 session = self.session_cls.from_user(user)
607 session.create_task(self.login_wrap(session), name=f"login wrap of {user}")
608 if cached_avatar is not None:
609 await self.pubsub.broadcast_avatar(
610 self.boundjid.bare, session.user_jid, cached_avatar
611 )
613 log.info("Slidge has successfully started")
615 async def __setup_attachments(self) -> None:
616 expire_coro = self.__expire_attachments_upload
617 if config.NO_UPLOAD_PATH:
618 if config.NO_UPLOAD_URL_PREFIX is None:
619 raise RuntimeError(
620 "If you set NO_UPLOAD_PATH you must set NO_UPLOAD_URL_PREFIX too."
621 )
622 expire_coro = self.__expire_attachments_no_upload
623 elif not config.UPLOAD_SERVICE:
624 try:
625 info_iq = await self.xmpp.plugin["xep_0363"].find_upload_service(
626 self.infer_real_domain()
627 )
628 except XMPPError:
629 info_iq = None
630 log.exception(
631 "The upload service could not be automatically determine. "
632 "Attachments to XMPP will not work. "
633 "Either specify 'upload-service' or 'no-upload-path' to fix that."
634 )
635 if info_iq is None:
636 if self.REGISTRATION_TYPE == RegistrationType.QRCODE:
637 log.warning(
638 "No method was configured for attachment and slidge "
639 "could not automatically determine the JID of a usable upload service. "
640 "Users likely won't be able to register since this network uses a "
641 "QR-code based registration flow."
642 )
643 if not config.USE_ATTACHMENT_ORIGINAL_URLS:
644 log.warning(
645 "Setting USE_ATTACHMENT_ORIGINAL_URLS to True since no method was configured "
646 "for attachments and no upload service was found. NB: this does not work for all "
647 "networks, especially for the E2EE attachments."
648 )
649 config.USE_ATTACHMENT_ORIGINAL_URLS = True
650 else:
651 log.info("Auto-discovered upload service: %s", info_iq["from"])
652 config.UPLOAD_SERVICE = info_iq["from"]
654 self.__expire_attachments_task = self.loop.create_task(
655 _loop(expire_coro, 3600 * 24)
656 )
658 async def __expire_attachments_no_upload(self) -> None:
659 with self.store.session() as orm:
660 attachments = self.store.attachments.get_all(orm)
661 self.__expire_attachments_in_store(
662 await asyncio.to_thread(self.__rm_local_attachments, attachments)
663 )
665 def __rm_local_attachments(self, attachments: list[Attachment]) -> list[int]:
666 to_remove = []
667 cutoff = datetime.now(tz=UTC) - timedelta(
668 days=config.NO_UPLOAD_MAX_DAYS or config.MAM_MAX_DAYS
669 )
671 for attachment in attachments:
672 path = attachment.local_path
673 try:
674 created = datetime.fromtimestamp(path.stat().st_ctime, tz=UTC)
675 except OSError:
676 self.log.debug(
677 "mtime of %s could not be determined, clearing up the row",
678 attachment,
679 )
680 to_remove.append(attachment.id)
681 else:
682 if created > cutoff:
683 continue
684 to_remove.append(attachment.id)
685 self.log.debug("%s is too old", attachment)
686 _safe_rm_parent(path)
687 to_remove.append(attachment.id)
689 return to_remove
691 async def __expire_attachments_upload(self) -> None:
692 with self.store.session() as orm:
693 attachments = self.store.attachments.get_all(orm)
694 to_remove = []
695 for attachment in attachments:
696 async with self.http.head(attachment.url) as resp:
697 if not resp.ok:
698 self.log.debug("%s is not reachable anymore", attachment)
699 to_remove.append(attachment.id)
700 self.__expire_attachments_in_store(to_remove)
702 def __expire_attachments_in_store(self, to_remove: list[int]) -> None:
703 with self.store.session() as orm:
704 self.store.attachments.remove(orm, to_remove)
705 orm.commit()
707 def infer_real_domain(self) -> JID:
708 return JID(re.sub(r"^.*?\.", "", self.xmpp.boundjid.bare))
710 async def __add_component_to_mds_whitelist(self, user_jid: JID) -> None:
711 # Uses privileged entity to add ourselves to the whitelist of the PEP
712 # MDS node so we receive MDS events
713 iq_creation = Iq(sto=user_jid.bare, sfrom=user_jid, stype="set")
714 iq_creation["pubsub"]["create"]["node"] = self.plugin["xep_0490"].stanza.NS
716 try:
717 await self.plugin["xep_0356"].send_privileged_iq(iq_creation)
718 except PermissionError:
719 log.warning(
720 "IQ privileges not granted for pubsub namespace, we cannot "
721 "create the MDS node of %s",
722 user_jid,
723 )
724 except PrivilegedIqError as exc:
725 nested = exc.nested_error()
726 # conflict this means the node already exists, we can ignore that
727 if nested is not None and nested.condition != "conflict":
728 log.exception(
729 "Could not create the MDS node of %s", user_jid, exc_info=exc
730 )
731 except Exception as e:
732 log.exception(
733 "Error while trying to create to the MDS node of %s",
734 user_jid,
735 exc_info=e,
736 )
738 iq_affiliation = Iq(sto=user_jid.bare, sfrom=user_jid, stype="set")
739 iq_affiliation["pubsub_owner"]["affiliations"]["node"] = self.plugin[
740 "xep_0490"
741 ].stanza.NS
743 aff = OwnerAffiliation()
744 aff["jid"] = self.boundjid.bare
745 aff["affiliation"] = "member"
746 iq_affiliation["pubsub_owner"]["affiliations"].append(aff)
748 try:
749 await self.plugin["xep_0356"].send_privileged_iq(iq_affiliation)
750 except PermissionError:
751 log.warning(
752 "IQ privileges not granted for pubsub#owner namespace, we cannot "
753 "listen to the MDS events of %s",
754 user_jid,
755 )
756 except Exception as e:
757 log.exception(
758 "Error while trying to subscribe to the MDS node of %s",
759 user_jid,
760 exc_info=e,
761 )
763 async def login_wrap(self, session: AnySession) -> str:
764 session.send_gateway_status("Logging in…", show="dnd")
765 session.is_logging_in = True
766 try:
767 status = await session.login()
768 except Exception as e:
769 log.warning("Login problem for %s", session.user_jid, exc_info=e)
770 session.send_gateway_status(f"Could not login: {e}", show="dnd")
771 msg = (
772 "You are not connected to this gateway! "
773 f"Maybe this message will tell you why: {e}"
774 )
775 session.send_gateway_message(msg)
776 session.logged = False
777 session.send_gateway_status("Login failed", show="dnd")
778 return msg
780 log.info("Login success for %s", session.user_jid)
781 session.logged = True
782 session.send_gateway_status("Syncing contacts…", show="dnd")
783 with self.store.session() as orm:
784 await session.contacts._fill(orm)
785 if not (r := session.contacts.ready).done():
786 r.set_result(True)
787 if self.GROUPS:
788 session.send_gateway_status("Syncing groups…", show="dnd")
789 await session.bookmarks.fill()
790 if not (r := session.bookmarks.ready).done():
791 r.set_result(True)
792 self.send_presence(pto=session.user.jid.bare, ptype="probe")
793 if status is None:
794 status = "Logged in"
795 session.send_gateway_status(status, show="chat")
797 if session.user.preferences.get("sync_avatar", False):
798 session.create_task(
799 self.fetch_user_avatar(session), name=f"fetch avatar of {session.user}"
800 )
801 else:
802 with self.store.session(expire_on_commit=False) as orm:
803 session.user.avatar_hash = None
804 orm.add(session.user)
805 orm.commit()
806 return status
808 async def fetch_user_avatar(self, session: AnySession) -> None:
809 try:
810 iq = await self.xmpp.plugin["xep_0060"].get_items(
811 session.user_jid.bare,
812 self.xmpp.plugin["xep_0084"].stanza.MetaData.namespace,
813 ifrom=self.boundjid.bare,
814 )
815 except IqTimeout:
816 self.log.warning("Iq timeout trying to fetch user avatar")
817 return
818 except IqError as e:
819 self.log.debug("Iq error when trying to fetch user avatar: %s", e)
820 if e.condition == "item-not-found":
821 try:
822 await session.on_avatar(None, None, None, None, None)
823 except NotImplementedError:
824 pass
825 else:
826 with self.store.session(expire_on_commit=False) as orm:
827 session.user.avatar_hash = None
828 orm.add(session.user)
829 orm.commit()
830 return
831 await self.__dispatcher.on_avatar_metadata_info(
832 session, iq["pubsub"]["items"]["item"]["avatar_metadata"]["info"]
833 )
835 def _send(
836 self,
837 stanza: MessageOrPresenceTypeVar,
838 **send_kwargs: Any, # noqa:ANN401
839 ) -> MessageOrPresenceTypeVar:
840 stanza.set_from(self.boundjid.bare)
841 if mto := send_kwargs.get("mto"):
842 stanza.set_to(mto)
843 stanza.send()
844 return stanza
846 def raise_if_not_allowed_jid(self, jid: JID) -> None:
847 if not self.jid_validator.match(jid.bare):
848 raise XMPPError(
849 condition="not-allowed",
850 text="Your account is not allowed to use this gateway. "
851 "The admin controls that with the USER_JID_VALIDATOR option.",
852 )
854 def send_raw(self, data: str | bytes) -> None:
855 # overridden from XMLStream to strip base64-encoded data from the logs
856 # to make them more readable.
857 if log.isEnabledFor(level=logging.DEBUG):
858 stripped = copy(data) if isinstance(data, str) else data.decode("utf-8")
859 # there is probably a way to do that in a single RE,
860 # but since it's only for debugging, the perf penalty
861 # does not matter much
862 for el in LOG_STRIP_ELEMENTS:
863 stripped = re.sub(
864 f"(<{el}.*?>)(.*)(</{el}>)",
865 "\1[STRIPPED]\3",
866 stripped,
867 flags=re.DOTALL | re.IGNORECASE,
868 )
869 log.debug("SEND: %s", stripped)
870 if not self.transport:
871 raise NotConnectedError()
872 if isinstance(data, str):
873 data = data.encode("utf-8")
874 self.transport.write(data)
876 def get_session_from_jid(self, j: JID) -> AnySession | None:
877 try:
878 return self.session_cls.from_jid(j)
879 except XMPPError:
880 return None
882 def exception(self, exception: Exception) -> None:
883 # """
884 # Called when a task created by slixmpp's internal (eg, on slix events) raises an Exception.
885 #
886 # Stop the event loop and exit on unhandled exception.
887 #
888 # The default :class:`slixmpp.basexmpp.BaseXMPP` behaviour is just to
889 # log the exception, but we want to avoid undefined behaviour.
890 #
891 # :param exception: An unhandled :class:`Exception` object.
892 # """
893 if isinstance(exception, IqError):
894 iq = exception.iq
895 log.error("%s: %s", iq["error"]["condition"], iq["error"]["text"])
896 log.warning("You should catch IqError exceptions")
897 elif isinstance(exception, IqTimeout):
898 iq = exception.iq
899 log.error("Request timed out: %s", iq)
900 log.warning("You should catch IqTimeout exceptions")
901 elif isinstance(exception, SyntaxError):
902 # Hide stream parsing errors that occur when the
903 # stream is disconnected (they've been handled, we
904 # don't need to make a mess in the logs).
905 pass
906 else:
907 if exception:
908 log.exception(exception)
909 self.loop.stop()
910 sys.exit(1)
912 async def make_registration_form(
913 self, _jid: JID, _node: str, _ifrom: JID, iq: Iq
914 ) -> Iq:
915 self.raise_if_not_allowed_jid(iq.get_from())
916 reg = iq["register"]
917 with self.store.session() as orm:
918 user = (
919 orm.query(GatewayUser).filter_by(jid=iq.get_from().bare).one_or_none()
920 )
921 log.debug("User found: %s", user)
923 form = reg["form"]
924 form.add_field(
925 "FORM_TYPE",
926 ftype="hidden",
927 value="jabber:iq:register",
928 )
929 form["title"] = f"Registration to '{self.COMPONENT_NAME}'"
930 form["instructions"] = self.REGISTRATION_INSTRUCTIONS
932 if user is not None:
933 reg["registered"] = False
934 form.add_field(
935 "remove",
936 label="Remove my registration",
937 required=True,
938 ftype="boolean",
939 value=False,
940 )
942 for field in self.REGISTRATION_FIELDS:
943 if field.var in reg.interfaces:
944 val = None if user is None else user.get(field.var)
945 if val is None:
946 reg.add_field(field.var)
947 else:
948 reg[field.var] = val
950 reg["instructions"] = self.REGISTRATION_INSTRUCTIONS
952 for field in self.REGISTRATION_FIELDS:
953 form.add_field(
954 field.var,
955 label=field.label,
956 required=field.required,
957 ftype=field.type,
958 options=field.options,
959 value=field.value if user is None else user.get(field.var, field.value),
960 )
962 reply = iq.reply()
963 reply.set_payload(reg)
964 return reply # type:ignore[no-any-return]
966 async def user_prevalidate(
967 self,
968 ifrom: JID,
969 form_dict: JSONSerializable,
970 fields: Iterable[FormField] | None = None,
971 validate: RegistrationValidationCoroutine | None = None,
972 ) -> JSONSerializable | None:
973 # Pre validate a registration form using the content of self.REGISTRATION_FIELDS
974 # before passing it to the plugin custom validation logic
975 if fields is None:
976 fields = self.REGISTRATION_FIELDS
977 for field in fields:
978 if field.required and not form_dict.get(field.var):
979 raise ValueError(f"Missing field: '{field.label}'")
980 if validate:
981 return await validate(ifrom, form_dict)
982 else:
983 return await self.validate(ifrom, form_dict)
985 @abc.abstractmethod
986 async def validate(
987 self, user_jid: JID, registration_form: JSONSerializable
988 ) -> JSONSerializable | None:
989 """
990 Validate a user's initial registration form.
992 Should raise the appropriate :class:`slixmpp.exceptions.XMPPError`
993 if the registration does not allow to continue the registration process.
995 If :py:attr:`REGISTRATION_TYPE` is a
996 :attr:`.RegistrationType.SINGLE_STEP_FORM`,
997 this method should raise something if it wasn't possible to successfully
998 log in to the legacy service with the registration form content.
1000 It is also used for other types of :py:attr:`REGISTRATION_TYPE` too, since
1001 the first step is always a form. If :attr:`.REGISTRATION_FIELDS` is an
1002 empty list (ie, it declares no :class:`.FormField`), the "form" is
1003 effectively a confirmation dialog displaying
1004 :attr:`.REGISTRATION_INSTRUCTIONS`.
1006 :param user_jid: JID of the user that has just registered
1007 :param registration_form: A dict where keys are the :attr:`.FormField.var` attributes
1008 of the :attr:`.BaseGateway.REGISTRATION_FIELDS` iterable.
1009 This dict can be modified and will be accessible as the ``legacy_module_data``
1010 of the
1012 :return : A dict that will be stored as the persistent "legacy_module_data"
1013 for this user. If you don't return anything here, the whole registration_form
1014 content will be stored.
1015 """
1016 raise NotImplementedError
1018 async def validate_two_factor_code(
1019 self, user: GatewayUser, code: str
1020 ) -> JSONSerializable | None:
1021 """
1022 Called when the user enters their 2FA code.
1024 Should raise the appropriate :class:`slixmpp.exceptions.XMPPError`
1025 if the login fails, and return successfully otherwise.
1027 Only used when :attr:`REGISTRATION_TYPE` is
1028 :attr:`.RegistrationType.TWO_FACTOR_CODE`.
1030 :param user: The :class:`.GatewayUser` whose registration is pending
1031 Use their :attr:`.GatewayUser.bare_jid` and/or
1032 :attr:`.registration_form` attributes to get what you need.
1033 :param code: The code they entered, either via "chatbot" message or
1034 adhoc command
1036 :return : A dict which keys and values will be added to the persistent "legacy_module_data"
1037 for this user.
1038 """
1039 raise NotImplementedError
1041 async def get_qr_text(self, user: GatewayUser) -> str:
1042 """
1043 This is where slidge gets the QR code content for the QR-based
1044 registration process. It will turn it into a QR code image and send it
1045 to the not-yet-fully-registered :class:`.GatewayUser`.
1047 Only used in when :attr:`BaseGateway.REGISTRATION_TYPE` is
1048 :attr:`.RegistrationType.QRCODE`.
1050 :param user: The :class:`.GatewayUser` whose registration is pending
1051 Use their :attr:`.GatewayUser.bare_jid` and/or
1052 :attr:`.registration_form` attributes to get what you need.
1053 """
1054 raise NotImplementedError
1056 async def confirm_qr(
1057 self,
1058 user_bare_jid: str,
1059 exception: Exception | None = None,
1060 legacy_data: JSONSerializable | None = None,
1061 ) -> None:
1062 """
1063 This method is meant to be called to finalize QR code-based registration
1064 flows, once the legacy service confirms the QR flashing.
1066 Only used in when :attr:`BaseGateway.REGISTRATION_TYPE` is
1067 :attr:`.RegistrationType.QRCODE`.
1069 :param user_bare_jid: The bare JID of the almost-registered
1070 :class:`GatewayUser` instance
1071 :param exception: Optionally, an XMPPError to be raised to **not** confirm
1072 QR code flashing.
1073 :param legacy_data: dict which keys and values will be added to the persistent
1074 "legacy_module_data" for this user.
1075 """
1076 fut = self.qr_pending_registrations[user_bare_jid]
1077 if exception is None:
1078 fut.set_result(legacy_data)
1079 else:
1080 fut.set_exception(exception)
1082 async def unregister_user(
1083 self, user: GatewayUser, msg: str = "You unregistered from this gateway."
1084 ) -> None:
1085 self.send_presence(pshow="dnd", pstatus=msg, pto=user.jid)
1086 await self.xmpp.plugin["xep_0077"].api["user_remove"](None, None, user.jid) # type:ignore[call-arg]
1087 await self.xmpp.session_cls.kill_by_jid(user.jid)
1089 async def input(
1090 self,
1091 jid: JID,
1092 text: str | None = None,
1093 mtype: MessageTypes = "chat",
1094 **input_kwargs: Any, # noqa:ANN401
1095 ) -> str:
1096 """
1097 Request arbitrary user input using a simple chat message, and await the result.
1099 You shouldn't need to call this directly bust instead use
1100 :meth:`.BaseSession.input` to directly target a user.
1102 :param jid: The JID we want input from
1103 :param text: A prompt to display for the user
1104 :param mtype: Message type
1105 :return: The user's reply
1106 """
1107 return await self.__chat_commands_handler.input(
1108 jid, text, mtype=mtype, **input_kwargs
1109 )
1111 async def send_qr(
1112 self,
1113 text: str,
1114 **msg_kwargs: Any, # noqa:ANN401
1115 ) -> str | None:
1116 """
1117 Sends a QR Code to a JID
1119 You shouldn't need to call directly bust instead use
1120 :meth:`.BaseSession.send_qr` to directly target a user.
1122 :param text: The text that will be converted to a QR Code
1123 :param msg_kwargs: Optional additional arguments to pass to
1124 :meth:`.BaseGateway.send_file`, such as the recipient of the QR,
1125 code
1126 """
1127 try:
1128 import qrcode
1129 except ImportError:
1130 log.error("Slidge needs the [qr] extra to be able to generate QR codes")
1131 raise
1132 qr = qrcode.make(text)
1133 with tempfile.NamedTemporaryFile(suffix=".png") as f:
1134 qr.save(f.name)
1135 url, _msgs = await self.send_file(Path(f.name), **msg_kwargs)
1136 return url
1138 def shutdown(self) -> list[asyncio.Task[None]]:
1139 # """
1140 # Called by the slidge entrypoint on normal exit.
1141 #
1142 # Sends offline presences from all contacts of all user sessions and from
1143 # the gateway component itself.
1144 # No need to call this manually, :func:`slidge.__main__.main` should take care of it.
1145 # """
1146 log.debug("Shutting down")
1147 tasks = []
1148 with self.store.session() as orm:
1149 for user in orm.query(GatewayUser).all():
1150 tasks.append(self.session_cls.from_jid(user.jid).shutdown())
1151 self.send_presence(ptype="unavailable", pto=user.jid)
1152 return tasks
1155SLIXMPP_PLUGINS = [
1156 "xep_0030", # Service discovery
1157 "xep_0045", # Multi-User Chat
1158 "xep_0050", # Adhoc commands
1159 "xep_0054", # VCard-temp (for MUC avatars)
1160 "xep_0055", # Jabber search
1161 "xep_0059", # Result Set Management
1162 "xep_0066", # Out of Band Data
1163 "xep_0071", # XHTML-IM (for stickers and custom emojis maybe later)
1164 "xep_0077", # In-band registration
1165 "xep_0084", # User Avatar
1166 "xep_0085", # Chat state notifications
1167 "xep_0100", # Gateway interaction
1168 "xep_0106", # JID Escaping
1169 "xep_0115", # Entity capabilities
1170 "xep_0122", # Data Forms Validation
1171 "xep_0128", # Service Discovery Extensions
1172 "xep_0153", # vCard-Based Avatars (for MUC avatars)
1173 "xep_0172", # User nickname
1174 "xep_0184", # Message Delivery Receipts
1175 "xep_0199", # XMPP Ping
1176 "xep_0221", # Data Forms Media Element
1177 "xep_0231", # Bits of Binary (for stickers and custom emojis maybe later)
1178 "xep_0249", # Direct MUC Invitations
1179 "xep_0264", # Jingle Content Thumbnails
1180 "xep_0280", # Carbons
1181 "xep_0292_provider", # VCard4
1182 "xep_0308", # Last message correction
1183 "xep_0313", # Message Archive Management
1184 "xep_0317", # Hats
1185 "xep_0319", # Last User Interaction in Presence
1186 "xep_0333", # Chat markers
1187 "xep_0334", # Message Processing Hints
1188 "xep_0356", # Privileged Entity
1189 "xep_0363", # HTTP file upload
1190 "xep_0385", # Stateless in-line media sharing
1191 "xep_0402", # PEP Native Bookmarks
1192 "xep_0421", # Anonymous unique occupant identifiers for MUCs
1193 "xep_0424", # Message retraction
1194 "xep_0425", # Message moderation
1195 "xep_0444", # Message reactions
1196 "xep_0447", # Stateless File Sharing
1197 "xep_0449", # Stickers
1198 "xep_0461", # Message replies
1199 "xep_0462", # Pubsub Type Filtering
1200 "xep_0463", # MUC Affiliation Versioning
1201 "xep_0469", # Bookmark Pinning
1202 "xep_0490", # Message Displayed Synchronization
1203 "xep_0492", # Chat Notification Settings
1204 # "xep_0503", # Server-side spaces
1205 "xep_0511", # Link Metadata
1206]
1209async def _loop(coro: Callable[[], Awaitable[None]], sleep: int) -> None:
1210 while True:
1211 await coro()
1212 await asyncio.sleep(sleep)
1215def _safe_rm_parent(path: Path) -> None:
1216 try:
1217 path.unlink()
1218 path.parent.rmdir()
1219 except OSError as exc:
1220 log.warning("%s couldn't be cleanly removed: !r", exc)
1223LOG_STRIP_ELEMENTS = ["data", "binval"]
1225log = logging.getLogger(__name__)