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