Coverage for slidge/core/pubsub.py: 87%
281 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
1import logging
2from collections.abc import Iterable
3from copy import copy
4from pathlib import Path
5from typing import ClassVar
7from slixmpp import (
8 JID,
9 CoroutineCallback,
10 ElementBase,
11 Iq,
12 Presence,
13 StanzaPath,
14 register_stanza_plugin,
15)
16from slixmpp.exceptions import IqError, IqTimeout, XMPPError
17from slixmpp.plugins.base import BasePlugin, register_plugin
18from slixmpp.plugins.xep_0060.stanza import Event, EventItem, EventItems, Item, Items
19from slixmpp.plugins.xep_0084.stanza import Data as AvatarData
20from slixmpp.plugins.xep_0084.stanza import Info
21from slixmpp.plugins.xep_0084.stanza import MetaData as AvatarMetadata
22from slixmpp.plugins.xep_0172.stanza import UserNick
23from slixmpp.plugins.xep_0292.stanza import VCard4
24from slixmpp.plugins.xep_0402.stanza import Conference
25from slixmpp.types import JidStr, OptJidStr
26from sqlalchemy.orm.exc import DetachedInstanceError
28from ..contact import LegacyContact
29from ..db.avatar import CachedAvatar
30from ..db.models import GatewayUser, Space
31from ..util.lock import NamedLockMixin
32from ..util.types import AnyGateway, AnySession
33from .dispatcher.util import exceptions_to_xmpp_errors
35VCARD4_NAMESPACE = "urn:xmpp:vcard4"
38class PepAvatar:
39 def __init__(self) -> None:
40 self.metadata: AvatarMetadata | None = None
41 self.id: str | None = None
42 self.url: str | None = None
43 self._avatar_data_path: Path | None = None
45 @property
46 def data(self) -> AvatarData | None:
47 if self._avatar_data_path is None:
48 return None
49 data = AvatarData()
50 data.set_value(self._avatar_data_path.read_bytes())
51 return data
53 def set_avatar_from_cache(self, cached_avatar: CachedAvatar) -> None:
54 metadata = AvatarMetadata()
55 self.id = cached_avatar.hash
56 if cached_avatar.path.exists():
57 metadata.add_info(
58 id=cached_avatar.hash,
59 itype="image/png",
60 ibytes=cached_avatar.path.stat().st_size,
61 height=str(cached_avatar.height),
62 width=str(cached_avatar.width),
63 )
64 if cached_avatar.stored.http_url:
65 self.url = cached_avatar.stored.http_url
66 metadata.add_info(
67 id=cached_avatar.stored.http_id,
68 itype=cached_avatar.stored.http_type,
69 ibytes=str(cached_avatar.stored.http_bytes),
70 url=cached_avatar.stored.http_url,
71 height=str(cached_avatar.stored.http_height),
72 width=str(cached_avatar.stored.http_width),
73 )
74 self.metadata = metadata
75 if cached_avatar.hash:
76 self._avatar_data_path = cached_avatar.path
78 @property
79 def http_metadata(self) -> Info | None:
80 if not self.metadata:
81 return None
82 items = self.metadata["items"]
83 if len(items) == 2:
84 return items[-1] # type:ignore[no-any-return]
85 return None
88class PubSubComponent(NamedLockMixin, BasePlugin):
89 xmpp: "AnyGateway"
91 name = "pubsub"
92 description = "Pubsub component"
93 dependencies: ClassVar[set[str]] = {
94 "xep_0030",
95 "xep_0060",
96 "xep_0115",
97 "xep_0163",
98 }
99 default_config: ClassVar[dict[str, str | None]] = {"component_name": None}
100 component_name: str
102 def __init__(self, *a: object, **kw: object) -> None:
103 super().__init__(*a, **kw)
104 # slixmpp usually does this via self.xmpp["xep_0163"].register_pep(),
105 # in the session_bind hook of plugins. However, for components,
106 # session_bind is never called.
107 register_stanza_plugin(EventItem, UserNick)
108 register_stanza_plugin(EventItem, Conference)
110 def plugin_init(self) -> None:
111 self.xmpp.register_handler(
112 CoroutineCallback(
113 "pubsub_get_avatar_data",
114 StanzaPath(f"iq@type=get/pubsub/items@node={AvatarData.namespace}"),
115 self._get_avatar_data, # type:ignore
116 )
117 )
118 self.xmpp.register_handler(
119 CoroutineCallback(
120 "pubsub_get_avatar_metadata",
121 StanzaPath(f"iq@type=get/pubsub/items@node={AvatarMetadata.namespace}"),
122 self._get_avatar_metadata, # type:ignore
123 )
124 )
125 self.xmpp.register_handler(
126 CoroutineCallback(
127 "pubsub_get_vcard",
128 StanzaPath(f"iq@type=get/pubsub/items@node={VCARD4_NAMESPACE}"),
129 self._get_vcard, # type:ignore
130 )
131 )
133 disco = self.xmpp.plugin["xep_0030"]
134 disco.add_identity("pubsub", "pep", self.component_name)
135 disco.add_identity("account", "registered", self.component_name)
136 disco.add_feature("http://jabber.org/protocol/pubsub#event")
137 disco.add_feature("http://jabber.org/protocol/pubsub#retrieve-items")
138 disco.add_feature("http://jabber.org/protocol/pubsub#persistent-items")
140 async def __get_features(self, presence: Presence) -> list[str]:
141 from_ = presence.get_from()
142 ver_string = presence["caps"]["ver"]
143 if ver_string:
144 info = await self.xmpp.plugin["xep_0115"].get_caps(from_)
145 else:
146 info = None
147 if info is None:
148 async with self.lock(from_):
149 try:
150 iq = await self.xmpp.plugin["xep_0030"].get_info(from_)
151 except (IqError, IqTimeout):
152 log.debug("Could not get disco#info of %s, ignoring", from_)
153 return []
154 info = iq["disco_info"]
155 return info["features"] # type:ignore[no-any-return]
157 async def on_presence_available(
158 self, p: Presence, contact: LegacyContact | None
159 ) -> None:
160 if "muc_join" in p:
161 log.debug("Ignoring MUC presence here")
162 return
164 to = p.get_to()
165 # we don't want to push anything for contacts that are not in the user's roster
166 if to != self.xmpp.boundjid.bare and (contact is None or not contact.is_friend):
167 return
169 from_ = p.get_from()
170 features = await self.__get_features(p)
172 if AvatarMetadata.namespace + "+notify" in features:
173 try:
174 pep_avatar = await self._get_authorized_avatar(p, contact)
175 except XMPPError:
176 pass
177 else:
178 if pep_avatar.metadata is not None:
179 await self.__broadcast(
180 data=pep_avatar.metadata,
181 from_=p.get_to().bare,
182 to=from_,
183 id=pep_avatar.metadata["info"]["id"],
184 )
185 if UserNick.namespace + "+notify" in features:
186 try:
187 pep_nick = await self._get_authorized_nick(p, contact)
188 except XMPPError:
189 pass
190 else:
191 await self.__broadcast(data=pep_nick, from_=p.get_to(), to=from_)
193 if contact is not None and VCARD4_NAMESPACE + "+notify" in features:
194 vcard = await contact.get_vcard()
195 if vcard is not None:
196 await self.broadcast_vcard_event(p.get_to(), from_, vcard)
198 async def broadcast_vcard_event(self, from_: JID, to: JID, vcard: VCard4) -> None:
199 item = Item()
200 item.namespace = VCARD4_NAMESPACE
201 item["id"] = "current"
202 # vcard: VCard4 = await self.xmpp["xep_0292_provider"].get_vcard(from_, to)
203 # The vcard content should NOT be in this event according to the spec:
204 # https://xmpp.org/extensions/xep-0292.html#sect-idm45669698174224
205 # but movim expects it to be here, and I guess it does not hurt
207 log.debug("Broadcast vcard4 event: %s", vcard)
208 await self.__broadcast(
209 data=vcard,
210 from_=JID(from_).bare,
211 to=to,
212 id="current",
213 node=VCARD4_NAMESPACE,
214 )
216 async def __get_contact(self, stanza: Iq | Presence) -> LegacyContact:
217 session = self.xmpp.get_session_from_stanza(stanza)
218 return await session.contacts.by_jid(stanza.get_to()) # type:ignore[no-any-return]
220 async def _get_authorized_avatar(
221 self, stanza: Iq | Presence, contact: LegacyContact | None = None
222 ) -> PepAvatar:
223 if stanza.get_to() == self.xmpp.boundjid.bare:
224 item = PepAvatar()
225 if self.xmpp.avatar is not None:
226 item.set_avatar_from_cache(self.xmpp.avatar)
227 return item
229 if contact is None:
230 contact = await self.__get_contact(stanza)
232 item = PepAvatar()
233 cached_avatar = contact.get_cached_avatar()
234 if cached_avatar is not None:
235 item.set_avatar_from_cache(cached_avatar)
236 return item
238 async def _get_authorized_nick(
239 self, stanza: Iq | Presence, contact: LegacyContact | None = None
240 ) -> UserNick:
241 if stanza.get_to() == self.xmpp.boundjid.bare:
242 return get_user_nick(self.xmpp.COMPONENT_NAME)
244 if contact is None:
245 contact = await self.__get_contact(stanza)
247 if contact.name is not None:
248 return get_user_nick(contact.name)
249 else:
250 return UserNick()
252 def __reply_with(
253 self, iq: Iq, content: AvatarData | AvatarMetadata | None, item_id: str | None
254 ) -> None:
255 requested_items = iq["pubsub"]["items"]
257 if len(requested_items) == 0:
258 self._reply_with_payload(iq, content, item_id)
259 else:
260 for item in requested_items:
261 if item["id"] == item_id:
262 self._reply_with_payload(iq, content, item_id)
263 return
264 raise XMPPError("item-not-found")
266 @exceptions_to_xmpp_errors
267 async def _get_avatar_data(self, iq: Iq) -> None:
268 pep_avatar = await self._get_authorized_avatar(iq)
269 self.__reply_with(iq, pep_avatar.data, pep_avatar.id)
271 @exceptions_to_xmpp_errors
272 async def _get_avatar_metadata(self, iq: Iq) -> None:
273 pep_avatar = await self._get_authorized_avatar(iq)
274 self.__reply_with(iq, pep_avatar.metadata, pep_avatar.id)
276 @exceptions_to_xmpp_errors
277 async def _get_vcard(self, iq: Iq) -> None:
278 # this is not the proper way that clients should retrieve VCards, but
279 # gajim does it this way.
280 # https://xmpp.org/extensions/xep-0292.html#sect-idm45669698174224
281 session = self.xmpp.get_session_from_stanza(iq)
282 contact = await session.contacts.by_jid(iq.get_to())
283 vcard = await contact.get_vcard()
284 if vcard is None:
285 raise XMPPError("item-not-found")
286 self._reply_with_payload(iq, vcard, "current", VCARD4_NAMESPACE)
288 @staticmethod
289 def _reply_with_payload(
290 iq: Iq,
291 payload: AvatarMetadata | AvatarData | VCard4 | None,
292 id_: str | None,
293 namespace: str | None = None,
294 ) -> None:
295 result = iq.reply()
296 item = Item()
297 if payload:
298 item.set_payload(payload.xml)
299 item["id"] = id_
300 result["pubsub"]["items"]["node"] = (
301 namespace if namespace else payload.namespace
302 )
303 result["pubsub"]["items"].append(item)
304 result.send()
306 async def __broadcast(
307 self,
308 data: ElementBase | None,
309 from_: JidStr,
310 to: OptJidStr = None,
311 items: EventItems | None = None,
312 **kwargs: object,
313 ) -> None:
314 from_ = JID(from_)
315 if from_ != self.xmpp.boundjid.bare and to is not None:
316 to = JID(to)
317 session = self.xmpp.get_session_from_jid(to)
318 if session is None:
319 return
320 await session.ready
322 item = EventItem()
323 if data:
324 item.set_payload(data.xml)
325 for k, v in kwargs.items():
326 item[k] = v
328 if items is None:
329 items = EventItems()
330 items.append(item)
331 assert data is not None
332 items["node"] = kwargs.get("node") or data.namespace
334 event = Event()
335 event.append(items)
337 msg = self.xmpp.Message()
338 msg.set_type("headline")
339 msg.set_from(from_)
340 msg.append(event)
342 if to is None:
343 with self.xmpp.store.session() as orm:
344 for u in orm.query(GatewayUser).all():
345 new_msg = copy(msg)
346 new_msg.set_to(u.jid.bare)
347 new_msg.send()
348 else:
349 msg.set_to(to)
350 msg.send()
352 async def broadcast_avatar(
353 self, from_: JidStr, to: JidStr, cached_avatar: CachedAvatar | None
354 ) -> None:
355 if cached_avatar is None:
356 await self.__broadcast(AvatarMetadata(), from_, to)
357 else:
358 pep_avatar = PepAvatar()
359 pep_avatar.set_avatar_from_cache(cached_avatar)
360 assert pep_avatar.metadata is not None
361 await self.__broadcast(pep_avatar.metadata, from_, to, id=pep_avatar.id)
363 def broadcast_nick(
364 self,
365 user_jid: JID,
366 jid: JidStr,
367 nick: str | None = None,
368 ) -> None:
369 jid = JID(jid)
370 nickname = get_user_nick(nick)
371 log.debug("New nickname: %s", nickname)
372 self.xmpp.loop.create_task(
373 self.__broadcast(nickname, jid, user_jid.bare),
374 name=f"broadcast nick of {jid}",
375 )
377 def broadcast_space(
378 self, session: AnySession, space: Space, node: str, item_ids: Iterable[str] = ()
379 ) -> None:
380 items = EventItems()
381 items["node"] = node
382 self.set_space_items(space, items, item_ids)
383 session.create_task(
384 self.__broadcast(
385 None,
386 self.xmpp.boundjid.bare,
387 session.user_jid,
388 items=items,
389 ),
390 name=f"broadcast space {node!r}",
391 )
393 def set_space_items(
394 self,
395 space: Space,
396 collection: Items | EventItems,
397 item_ids: Iterable[str] = (),
398 ) -> None:
399 item_cls = Item if isinstance(collection, Items) else EventItem
401 try:
402 rooms = space.rooms
403 except DetachedInstanceError:
404 # if the caller does not care about rooms, it should ensure this
405 # has been loaded
406 rooms = []
408 for room in rooms:
409 if item_ids and str(room.jid) not in item_ids:
410 continue
411 item = item_cls()
412 item["id"] = room.jid
413 item.enable("conference")
414 if room.name:
415 item["conference"]["name"] = room.name
416 collection.append(item)
418 for attr, avatar in [("avatar", space.avatar), ("banner", space.banner)]:
419 ns = f"urn:xmpp:spaces:{attr}:metadata:0"
420 if avatar is None:
421 continue
422 if item_ids and ns not in item_ids:
423 continue
424 item = item_cls()
425 item["id"] = ns
426 meta = self.xmpp.plugin["xep_0084"].stanza.MetaData()
427 meta.add_info(
428 id=avatar.http_id,
429 itype=avatar.http_type,
430 ibytes=str(avatar.http_bytes),
431 height=str(avatar.http_height),
432 width=str(avatar.http_width),
433 url=avatar.http_url,
434 )
435 item.append(meta)
436 collection.append(item)
438 def broadcast_space_metadata(
439 self, session: AnySession, space: Space, node: str
440 ) -> None:
441 msg = self.xmpp.make_message(
442 mtype="headline", mto=session.user_jid, mfrom=self.xmpp.boundjid.bare
443 )
444 msg["pubsub_event"]["configuration"]["node"] = node
446 form = msg["pubsub_event"]["configuration"]["form"]
447 form["type"] = "result"
449 form.add_field(
450 var="FORM_TYPE",
451 ftype="hidden",
452 value="http://jabber.org/protocol/pubsub#node_config",
453 )
454 form.add_field(var="pubsub#title", value=space.name)
456 # MUSTs
457 form.add_field(var="pubsub#type", value="urn:xmpp:spaces:0")
458 form.add_field(var="pubsub#notify_retract", value="true")
459 form.add_field(var="pubsub#persist_items", value="true")
460 form.add_field(var="pubsub#purge_offline", value="false")
461 # SHOULDs
462 form.add_field(var="pubsub#notify_sub", value="true")
463 form.add_field(var="pubsub#notify_config", value="true")
464 form.add_field(var="pubsub#notify_delete", value="true")
465 form.add_field(var="pubsub#publish_model", value="publishers")
466 # MUST (private spaces)
467 form.add_field(var="pubsub#access_model", value="authorize")
469 msg.send()
472def get_user_nick(nick: str | None = None) -> UserNick:
473 user_nick = UserNick()
474 if nick is not None:
475 user_nick["nick"] = nick
476 return user_nick
479log = logging.getLogger(__name__)
480register_plugin(PubSubComponent)