Coverage for src/lexigram/admin/di/sub_providers/realtime.py: 0%

53 statements  

« prev     ^ index     » next       coverage.py v7.15.4, created at 2026-08-24 23:18 +0800

1"""Admin realtime sub-provider — WebSocket, SSE, collaborative editing, events.""" 

2 

3from __future__ import annotations 

4 

5from typing import TYPE_CHECKING, Any 

6 

7from lexigram.contracts.core.health import HealthCheckResult, HealthStatus 

8from lexigram.logging import get_logger 

9 

10if TYPE_CHECKING: 

11 from lexigram.admin.config import AdminConfig 

12 from lexigram.contracts.core.di import ( 

13 ContainerRegistrarProtocol, 

14 ContainerResolverProtocol, 

15 ) 

16 

17logger = get_logger(__name__) 

18 

19 

20class AdminRealtimeSubProvider: 

21 """Manages admin realtime infrastructure: WebSocket, SSE, collaborative editing. 

22 

23 Registers realtime services: WebSocket manager, event hub, SSE manager. 

24 """ 

25 

26 def __init__(self, config: AdminConfig, **kwargs: object) -> None: 

27 self._config = config 

28 self._kwargs = kwargs 

29 self._initialized = False 

30 self._hub: Any = None 

31 self._inbox_bridge_wired = False 

32 

33 @property 

34 def config(self) -> AdminConfig: 

35 """Return current admin config.""" 

36 return self._config 

37 

38 async def register(self, container: ContainerRegistrarProtocol) -> None: 

39 """Register realtime services: WebSocket manager, event hub, SSE manager, collaborative.""" 

40 from lexigram.admin.realtime.subject_hub import SubjectAdminEventHub 

41 from lexigram.admin.realtime.ws_handler_registry import WSMessageTypeRegistry 

42 from lexigram.admin.services.collaborative import CollaborativeEditingService 

43 from lexigram.admin.services.realtime import RealtimeService 

44 from lexigram.contracts.core.stores import LockStoreProtocol 

45 

46 realtime_svc = RealtimeService() 

47 

48 container.singleton(WSMessageTypeRegistry, WSMessageTypeRegistry()) 

49 container.singleton(RealtimeService, realtime_svc) 

50 container.singleton(SubjectAdminEventHub, SubjectAdminEventHub()) 

51 

52 # Lock store is provided externally through container bindings. 

53 lock_store: LockStoreProtocol | None = None 

54 

55 # Pre-build instance to avoid container trying to resolve RealtimeService 

56 # from TYPE_CHECKING-only annotation on CollaborativeEditingService.__init__. 

57 container.singleton( 

58 CollaborativeEditingService, 

59 CollaborativeEditingService( 

60 lock_store=lock_store, # type: ignore[arg-type] 

61 realtime_service=realtime_svc, 

62 ), 

63 ) 

64 

65 async def boot(self, container: ContainerResolverProtocol) -> None: 

66 """Boot realtime services: initialize WebSocket and SSE connections.""" 

67 from lexigram.admin.realtime.subject_hub import SubjectAdminEventHub 

68 from lexigram.contracts.notification.inbox import INBOX_SENT_HOOK 

69 from lexigram.hooks.ambient import register_action 

70 

71 try: 

72 self._hub = await container.resolve(SubjectAdminEventHub) 

73 except Exception: # noqa: BLE001 — hub is optional 

74 self._hub = None 

75 logger.debug("admin_realtime.no_event_hub") 

76 

77 if self._hub is not None and not self._inbox_bridge_wired: 

78 register_action(INBOX_SENT_HOOK, self._publish_inbox_sent) 

79 self._inbox_bridge_wired = True 

80 logger.info("admin_realtime.inbox_bridge_wired", hook=INBOX_SENT_HOOK) 

81 

82 self._initialized = True 

83 

84 async def _publish_inbox_sent(self, **kwargs: Any) -> None: 

85 """Push an inbox message to the SSE hub for its recipient.""" 

86 if self._hub is None: 

87 return 

88 

89 message = kwargs.get("message") 

90 metadata = getattr(message, "metadata", None) 

91 level = metadata.get("level", "info") if isinstance(metadata, dict) else "info" 

92 

93 await self._hub.publish_notification( 

94 title=kwargs.get("title", "Notification"), 

95 message=kwargs.get("body", ""), 

96 level=str(level), 

97 target_users=[kwargs["user_id"]] if kwargs.get("user_id") else None, 

98 ) 

99 

100 async def shutdown(self) -> None: 

101 """Shut down realtime services.""" 

102 self._initialized = False 

103 

104 async def health_check(self, timeout: float = 5.0) -> HealthCheckResult: 

105 """Return realtime infrastructure health status.""" 

106 return HealthCheckResult( 

107 component="admin_realtime", 

108 status=HealthStatus.HEALTHY if self._initialized else HealthStatus.UNKNOWN, 

109 message="Admin realtime operational" 

110 if self._initialized 

111 else "Not yet initialized", 

112 ) 

113 

114 

115__all__ = ["AdminRealtimeSubProvider"]