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

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, avatar_cache 

30from ..db.models import GatewayUser, Space 

31from ..util.lock import NamedLockMixin 

32from ..util.types import AnyGateway, AnySession 

33 

34VCARD4_NAMESPACE = "urn:xmpp:vcard4" 

35 

36 

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 

43 

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 

51 

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 

76 

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 

85 

86 

87class PubSubComponent(NamedLockMixin, BasePlugin): 

88 xmpp: "AnyGateway" 

89 

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 

100 

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) 

108 

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 ) 

131 

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

138 

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] 

155 

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 

162 

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 

167 

168 from_ = p.get_from() 

169 features = await self.__get_features(p) 

170 

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

191 

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) 

196 

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 

205 

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 ) 

214 

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] 

218 

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 

227 

228 if contact is None: 

229 contact = await self.__get_contact(stanza) 

230 

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 

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

269 

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) 

273 

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) 

284 

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

302 

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 

318 

319 item = EventItem() 

320 if data: 

321 item.set_payload(data.xml) 

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

323 item[k] = v 

324 

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 

330 

331 event = Event() 

332 event.append(items) 

333 

334 msg = self.xmpp.Message() 

335 msg.set_type("headline") 

336 msg.set_from(from_) 

337 msg.append(event) 

338 

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

348 

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) 

359 

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 ) 

373 

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 ) 

389 

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 

397 

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

404 

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) 

414 

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) 

434 

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 

442 

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

444 form["type"] = "result" 

445 

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) 

452 

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

465 

466 msg.send() 

467 

468 

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 

474 

475 

476log = logging.getLogger(__name__) 

477register_plugin(PubSubComponent)