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

1import logging 

2from collections.abc import Iterable 

3from copy import copy 

4from pathlib import Path 

5from typing import ClassVar 

6 

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 

27 

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 

34 

35VCARD4_NAMESPACE = "urn:xmpp:vcard4" 

36 

37 

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 

44 

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 

52 

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 

77 

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 

86 

87 

88class PubSubComponent(NamedLockMixin, BasePlugin): 

89 xmpp: "AnyGateway" 

90 

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 

101 

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) 

109 

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 ) 

132 

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

139 

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] 

156 

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 

163 

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 

168 

169 from_ = p.get_from() 

170 features = await self.__get_features(p) 

171 

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

192 

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) 

197 

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 

206 

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 ) 

215 

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] 

219 

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 

228 

229 if contact is None: 

230 contact = await self.__get_contact(stanza) 

231 

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 

237 

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) 

243 

244 if contact is None: 

245 contact = await self.__get_contact(stanza) 

246 

247 if contact.name is not None: 

248 return get_user_nick(contact.name) 

249 else: 

250 return UserNick() 

251 

252 def __reply_with( 

253 self, iq: Iq, content: AvatarData | AvatarMetadata | None, item_id: str | None 

254 ) -> None: 

255 requested_items = iq["pubsub"]["items"] 

256 

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

265 

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) 

270 

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) 

275 

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) 

287 

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

305 

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 

321 

322 item = EventItem() 

323 if data: 

324 item.set_payload(data.xml) 

325 for k, v in kwargs.items(): 

326 item[k] = v 

327 

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 

333 

334 event = Event() 

335 event.append(items) 

336 

337 msg = self.xmpp.Message() 

338 msg.set_type("headline") 

339 msg.set_from(from_) 

340 msg.append(event) 

341 

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

351 

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) 

362 

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 ) 

376 

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 ) 

392 

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 

400 

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 = [] 

407 

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) 

417 

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) 

437 

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 

445 

446 form = msg["pubsub_event"]["configuration"]["form"] 

447 form["type"] = "result" 

448 

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) 

455 

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

468 

469 msg.send() 

470 

471 

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 

477 

478 

479log = logging.getLogger(__name__) 

480register_plugin(PubSubComponent)