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

1""" 

2This module extends slixmpp.ComponentXMPP to make writing new LegacyClients easier 

3""" 

4 

5from __future__ import annotations 

6 

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) 

30 

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 

40 

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 

76 

77T = TypeVar("T") 

78P = ParamSpec("P") 

79 

80 

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. 

90 

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. 

94 

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`. 

98 

99 Abstract methods related to the registration process must be overriden 

100 for a functional :term:`Legacy Module`: 

101 

102 - :meth:`.validate` 

103 - :meth:`.validate_two_factor_code` 

104 - :meth:`.get_qr_text` 

105 - :meth:`.confirm_qr` 

106 

107 NB: Not all of these must be overridden, it depends on the 

108 :attr:`REGISTRATION_TYPE`. 

109 

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. 

115 

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.: 

118 

119 .. code-block:: python 

120 

121 self.send_presence( 

122 pfrom="somebody@component.example.com", 

123 pto="someonwelse@anotherexample.com", 

124 ) 

125 

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. 

129 

130 """ 

131 

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 """ 

140 

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). 

160 

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 """ 

166 

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" 

172 

173 ALTERNATIVE_REGISTRATION_FLOWS: ClassVar[list[AltRegistrationFlow]] = [] 

174 

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 ] 

219 

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 """ 

235 

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. 

245 

246 Plugins should implement search by overriding :meth:`.BaseSession.search` 

247 (restricted to registered users). 

248 

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 """ 

261 

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 """ 

268 

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 """ 

275 

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 """ 

289 

290 mtype: MessageTypes = "chat" 

291 is_group = False 

292 _can_send_carbon = False 

293 store: SlidgeStore 

294 

295 session_cls: type[SessionType_co] 

296 """Concrete :class:`.BaseSession` subclass of this legacy module. 

297 

298 Derived automatically from the generic parameter, e.g., 

299 ``class Gateway(BaseGateway[Session])``. 

300 """ 

301 

302 def __init_subclass__(cls, **kwargs: object) -> None: 

303 super().__init_subclass__(**kwargs) 

304 derive_wired_class(cls, BaseGateway, "session_cls") 

305 

306 COMMANDS: ClassVar[tuple[type[CommandBase], ...]] = BUILTIN_COMMANDS 

307 """Commands (adhoc and chatbot) this gateway provides. 

308 

309 Subclasses may override this with additional custom commands. 

310 E.g.: ``COMMANDS = BaseGateway.COMMANDS + commands_from_module(command)``. 

311 """ 

312 

313 http: aiohttp.ClientSession 

314 avatar: CachedAvatar | None = None 

315 

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 

363 

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 ]() 

374 

375 self.register_plugins() 

376 self.__setup_session_cls() 

377 

378 self.get_session_from_stanza = self.session_cls.from_stanza 

379 self.get_session_from_user = self.session_cls.from_user 

380 

381 self.__register_slixmpp_events() 

382 self.__register_slixmpp_api() 

383 self.roster.set_backend(RosterBackend(self)) 

384 

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) 

388 

389 # with this we receive user avatar updates 

390 self.plugin["xep_0030"].add_feature("urn:xmpp:avatar:metadata+notify") 

391 

392 self.plugin["xep_0030"].add_feature("urn:xmpp:chat-markers:0") 

393 

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") 

409 

410 self.__adhoc_handler = AdhocProvider(self) 

411 self.__chat_commands_handler = ChatCommandProvider(self) 

412 

413 self.__dispatcher = SessionDispatcher(self) 

414 self.get_caps_ver = self.__dispatcher.get_caps_ver 

415 

416 self.__register_commands() 

417 

418 def __setup_session_cls(self) -> None: 

419 contact_cls = self.session_cls.roster_cls.contact_cls 

420 

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) 

430 

431 self.session_cls.xmpp = self 

432 

433 @property 

434 def _uploader(self) -> AttachmentUploader: 

435 # the component itself is not bound to any user session 

436 return AttachmentUploader(self, None) 

437 

438 async def kill_session(self, jid: JID) -> None: 

439 await self.session_cls.kill_by_jid(jid) 

440 

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 

446 

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) 

480 

481 for i, alt in enumerate(self.ALTERNATIVE_REGISTRATION_FLOWS): 

482 

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 

491 

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) 

496 

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 

502 

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() 

518 

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) 

529 

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 

540 

541 return wrapped 

542 

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 ) 

552 

553 @property # type:ignore[override] 

554 def jid(self) -> JID: 

555 # Override to avoid slixmpp deprecation warnings. 

556 return self.boundjid 

557 

558 @jid.setter 

559 def jid(self, jid: JID) -> None: 

560 raise RuntimeError 

561 

562 async def __on_session_start(self, event: object) -> None: 

563 log.debug("Gateway session start: %s", event) 

564 

565 await self.__setup_attachments() 

566 

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) 

571 

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 

586 

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 ) 

612 

613 log.info("Slidge has successfully started") 

614 

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"] 

653 

654 self.__expire_attachments_task = self.loop.create_task( 

655 _loop(expire_coro, 3600 * 24) 

656 ) 

657 

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 ) 

664 

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 ) 

670 

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) 

688 

689 return to_remove 

690 

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) 

701 

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() 

706 

707 def infer_real_domain(self) -> JID: 

708 return JID(re.sub(r"^.*?\.", "", self.xmpp.boundjid.bare)) 

709 

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 

715 

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 ) 

737 

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 

742 

743 aff = OwnerAffiliation() 

744 aff["jid"] = self.boundjid.bare 

745 aff["affiliation"] = "member" 

746 iq_affiliation["pubsub_owner"]["affiliations"].append(aff) 

747 

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 ) 

762 

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 

779 

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") 

796 

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 

807 

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 ) 

834 

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 

845 

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 ) 

853 

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) 

875 

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 

881 

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) 

911 

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) 

922 

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 

931 

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 ) 

941 

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 

949 

950 reg["instructions"] = self.REGISTRATION_INSTRUCTIONS 

951 

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 ) 

961 

962 reply = iq.reply() 

963 reply.set_payload(reg) 

964 return reply # type:ignore[no-any-return] 

965 

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) 

984 

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. 

991 

992 Should raise the appropriate :class:`slixmpp.exceptions.XMPPError` 

993 if the registration does not allow to continue the registration process. 

994 

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. 

999 

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`. 

1005 

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 

1011 

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 

1017 

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. 

1023 

1024 Should raise the appropriate :class:`slixmpp.exceptions.XMPPError` 

1025 if the login fails, and return successfully otherwise. 

1026 

1027 Only used when :attr:`REGISTRATION_TYPE` is 

1028 :attr:`.RegistrationType.TWO_FACTOR_CODE`. 

1029 

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 

1035 

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 

1040 

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`. 

1046 

1047 Only used in when :attr:`BaseGateway.REGISTRATION_TYPE` is 

1048 :attr:`.RegistrationType.QRCODE`. 

1049 

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 

1055 

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. 

1065 

1066 Only used in when :attr:`BaseGateway.REGISTRATION_TYPE` is 

1067 :attr:`.RegistrationType.QRCODE`. 

1068 

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) 

1081 

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) 

1088 

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. 

1098 

1099 You shouldn't need to call this directly bust instead use 

1100 :meth:`.BaseSession.input` to directly target a user. 

1101 

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 ) 

1110 

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 

1118 

1119 You shouldn't need to call directly bust instead use 

1120 :meth:`.BaseSession.send_qr` to directly target a user. 

1121 

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 

1137 

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 

1153 

1154 

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] 

1207 

1208 

1209async def _loop(coro: Callable[[], Awaitable[None]], sleep: int) -> None: 

1210 while True: 

1211 await coro() 

1212 await asyncio.sleep(sleep) 

1213 

1214 

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) 

1221 

1222 

1223LOG_STRIP_ELEMENTS = ["data", "binval"] 

1224 

1225log = logging.getLogger(__name__)