Coverage for src/lexigram/admin/di/mount/contributors.py: 77%

137 statements  

« prev     ^ index     » next       coverage.py v7.15.4, created at 2026-08-21 14:56 +0800

1"""Mount phases for contributors, routing, SSE and app-state exposure.""" 

2 

3from __future__ import annotations 

4 

5from typing import TYPE_CHECKING, Any 

6 

7from lexigram.logging import get_logger 

8 

9if TYPE_CHECKING: 

10 from lexigram.admin.di.bundle_provider import AdminProvider 

11 from lexigram.admin.di.mount.context import MountContext 

12 

13_log = get_logger(__name__) 

14 

15 

16class AdminMountContributorsMixin: 

17 """Mount phases that wire contributors, the router and app state.""" 

18 

19 async def _mount_contributors( 

20 self: AdminProvider, resolver: Any, ctx: MountContext 

21 ) -> None: 

22 """Discover contributor resources and wire their data sources. 

23 

24 Contributor discovery is best-effort; failures are recorded in 

25 ``_mount_failures`` without aborting the mount. 

26 

27 Args: 

28 resolver: The DI resolver for contributor resolution. 

29 ctx: Mount pipeline state (resources/contributors populated). 

30 """ 

31 from lexigram.admin.contributors.registry import ContributorRegistry 

32 from lexigram.admin.contributors.resource_collector import ResourceCollector 

33 from lexigram.admin.dashboard.naming_policy import NamingPolicy 

34 

35 contributors: list = [] 

36 contributor_registry: Any | None = None 

37 try: 

38 contributor_registry = await resolver.resolve( 

39 ContributorRegistry, 

40 bypass_visibility=True, 

41 ) 

42 if contributor_registry is None: 

43 raise ValueError("Contributor registry unavailable") 

44 contributors = list(contributor_registry.get_all()) 

45 except Exception as exc: 

46 _log.warning("admin.contributors_discovery_failed", exc_info=True) 

47 self._mount_failures["contributor_discovery"] = str(exc) 

48 

49 ctx.contributors = contributors 

50 ctx.contributor_registry = contributor_registry 

51 

52 try: 

53 naming = NamingPolicy(mode=self._config.contributor_collision_mode) 

54 collector = ResourceCollector(naming_policy=naming) 

55 contributor_resources = collector.collect(contributors) 

56 for resource_cls in contributor_resources: 

57 name = ( 

58 getattr(resource_cls, "name", None) 

59 or resource_cls.__name__.replace("Resource", "").lower() 

60 ) 

61 try: 

62 ctx.resources[name] = await resolver.resolve( 

63 resource_cls, 

64 bypass_visibility=True, 

65 ) 

66 except Exception: # noqa: BLE001 — fall back to direct 

67 try: 

68 ctx.resources[name] = resource_cls() 

69 except Exception as inner: # noqa: BLE001 — skip, log warning 

70 _log.warning( 

71 "admin.contributor_resource_resolution_failed", 

72 resource=resource_cls.__name__, 

73 contributor=name, 

74 error=str(inner), 

75 ) 

76 self._mount_failures[f"contributor_resource:{name}"] = str( 

77 inner 

78 ) 

79 

80 # Wire data sources to resolved resources that declare _data_source_class 

81 user_permission_inventory: Any = None 

82 user_resource_cls: Any = None 

83 try: 

84 from lexigram.admin.rbac.inventory import PermissionInventoryService 

85 from lexigram.admin.resources.users import UserResource 

86 

87 user_permission_inventory = await resolver.resolve( 

88 PermissionInventoryService, 

89 bypass_visibility=True, 

90 ) 

91 user_resource_cls = UserResource 

92 except Exception: # noqa: BLE001 — best-effort wiring 

93 _log.debug("admin.user_permission_inventory_unavailable") 

94 

95 for name, resource in list(ctx.resources.items()): 

96 if ( 

97 user_permission_inventory is not None 

98 and user_resource_cls is not None 

99 and isinstance(resource, user_resource_cls) 

100 ): 

101 try: 

102 resource.permission_inventory = user_permission_inventory 

103 _log.debug( 

104 "admin.user_permission_inventory_wired", resource=name 

105 ) 

106 except Exception: # noqa: BLE001 — best-effort wiring 

107 _log.debug( 

108 "admin.user_permission_inventory_wiring_failed", 

109 resource=name, 

110 ) 

111 dsc = getattr(type(resource), "_data_source_class", None) 

112 if dsc is not None and hasattr(resource, "set_data_source"): 

113 try: 

114 ds = await resolver.resolve(dsc, bypass_visibility=True) 

115 resource.set_data_source(ds) 

116 _log.debug("admin.data_source_wired", resource=name) 

117 except Exception: 

118 _log.debug("admin.data_source_wiring_failed", resource=name) 

119 

120 # Wrap data source with search wrappers when the resource 

121 # has a searchable spec and the search engine is available. 

122 search_spec = resource.search_spec() 

123 if search_spec and search_spec.index_name: 

124 try: 

125 from lexigram.contracts.search import SearchEngineProtocol 

126 

127 search_engine = await resolver.resolve( 

128 SearchEngineProtocol, 

129 bypass_visibility=True, 

130 ) 

131 

132 from lexigram.admin.integrations.search_query import ( 

133 SearchQueryDataSourceWrapper, 

134 ) 

135 from lexigram.admin.integrations.search_sync import ( 

136 SearchSyncDataSourceWrapper, 

137 ) 

138 

139 fallback_to_like = getattr( 

140 self._config.integrations.search, 

141 "fallback_to_like", 

142 True, 

143 ) 

144 query_wrapped = SearchQueryDataSourceWrapper( 

145 ds, 

146 search_engine, 

147 search_spec.index_name, 

148 fallback_to_like=fallback_to_like, 

149 ) 

150 wrapped = SearchSyncDataSourceWrapper( 

151 query_wrapped, search_engine, search_spec 

152 ) 

153 resource.set_data_source(wrapped) 

154 _log.debug("admin.search_wired", resource=name) 

155 except Exception: 

156 _log.debug("admin.search_wiring_failed", resource=name) 

157 

158 except Exception: # noqa: BLE001 — resource collection is non-fatal 

159 _log.warning("admin.contributors_resource_collection_failed", exc_info=True) 

160 

161 async def _mount_integration( 

162 self: AdminProvider, container: Any, ctx: MountContext 

163 ) -> None: 

164 """Build the admin router and integrate contributor routes. 

165 

166 Args: 

167 container: The root DI resolver for route integration. 

168 ctx: Mount pipeline state (``router`` populated). 

169 """ 

170 from lexigram.admin.core.routing import AdminRouter 

171 from lexigram.admin.dashboard.naming_policy import NamingPolicy 

172 from lexigram.admin.dashboard.route_integrator import RouteIntegrator 

173 

174 router = AdminRouter( 

175 config=self._config, 

176 resources=ctx.resources, 

177 controllers=ctx.controllers, 

178 middleware_stack=ctx.middlewares, 

179 authorizer=self._authorizer_service, 

180 ) 

181 ctx.router = router 

182 

183 try: 

184 naming = NamingPolicy(mode=self._config.contributor_collision_mode) 

185 integrator = RouteIntegrator( 

186 router=router, 

187 naming_policy=naming, 

188 route_prefix=self._config.prefix, 

189 container=container, 

190 ) 

191 integrator.register(ctx.contributors) 

192 except Exception as exc: # noqa: BLE001 — route integration is non-fatal 

193 _log.warning("admin.contributors_route_integration_failed", exc_info=True) 

194 self._mount_failures["route_integrator"] = str(exc) 

195 

196 async def _mount_sse_widgets( 

197 self: AdminProvider, container: Any, ctx: MountContext 

198 ) -> None: 

199 """Register the SSE endpoint for live widget delivery. 

200 

201 Args: 

202 container: The root DI resolver for hub/permission services. 

203 ctx: Mount pipeline state (``router`` read). 

204 """ 

205 router = ctx.router 

206 if router is None: 

207 return 

208 try: 

209 from lexigram.admin.dashboard.widget_stream import ( 

210 build_widget_event_stream_handler, 

211 ) 

212 from lexigram.admin.rbac.service import PermissionService 

213 from lexigram.admin.realtime.subject_hub import SubjectAdminEventHub 

214 from lexigram.contracts.web.sse import ReactiveSseBridgeProtocol 

215 

216 widget_hub: SubjectAdminEventHub = await container.resolve( 

217 SubjectAdminEventHub 

218 ) 

219 permission_service: PermissionService = await container.resolve( 

220 PermissionService 

221 ) 

222 sse_bridge = await container.resolve(ReactiveSseBridgeProtocol) 

223 

224 router.add_route( 

225 "/_sse/widgets", 

226 "GET", 

227 build_widget_event_stream_handler( 

228 widget_hub, permission_service, sse_bridge=sse_bridge 

229 ), 

230 "admin_sse_widgets", 

231 ) 

232 _log.info("admin.sse_widgets_route_registered", path="/admin/_sse/widgets") 

233 except Exception as exc: # noqa: BLE001 — SSE is optional 

234 _log.warning("admin.sse_widgets_route_skipped", reason=str(exc)) 

235 

236 async def _mount_app_state( 

237 self: AdminProvider, app: Any, ctx: MountContext 

238 ) -> None: 

239 """Mount the router and expose nav/registry state on both apps. 

240 

241 The renderer looks up request.app.state.nav_builder; request.app is 

242 the *inner* admin sub-app (not the outer Starlette app), so state is 

243 set on both. 

244 

245 Args: 

246 app: The outer Starlette application to mount the panel on. 

247 ctx: Mount pipeline state (nav/registry state read). 

248 """ 

249 router = ctx.router 

250 if router is None: 

251 return 

252 admin_app = router.mount(app) 

253 ctx.admin_app = admin_app 

254 

255 # Expose nav_builder on app state so AdminRenderer can build the sidebar. 

256 if hasattr(app, "state"): 

257 app.state.nav_builder = ctx.nav_builder 

258 if admin_app is not None and hasattr(admin_app, "state"): 

259 admin_app.state.nav_builder = ctx.nav_builder 

260 

261 # Build NavigationAssembler contributions and expose on app state. 

262 assembler_nav_items: list[dict[str, object]] = [] 

263 assembler_groups: dict[str, list[Any]] | None = None 

264 registry = ctx.contributor_registry 

265 if registry is not None and ctx.contributors: 

266 from lexigram.admin.navigation.assembler import ( 

267 NavigationAssembler, 

268 contributions_to_flat_nav, 

269 ) 

270 

271 try: 

272 assembler = NavigationAssembler( 

273 contributor_registry=registry, 

274 resource_items=[], 

275 ) 

276 grouped = await assembler.build() 

277 assembler_groups = grouped 

278 assembler_nav_items = contributions_to_flat_nav(grouped) 

279 except Exception: # noqa: BLE001 — non-fatal 

280 _log.warning("admin.navigation_assembler_prebuild_failed") 

281 if hasattr(app, "state"): 

282 app.state.assembler_nav_items = assembler_nav_items 

283 app.state.assembler_groups = assembler_groups or {} 

284 if admin_app is not None and hasattr(admin_app, "state"): 

285 admin_app.state.assembler_nav_items = assembler_nav_items 

286 admin_app.state.assembler_groups = assembler_groups or {} 

287 

288 # Expose the cluster registry on app state so nav resolution and 

289 # cluster centers resolve the active cluster per request. 

290 if ctx.cluster_registry is not None: 

291 if hasattr(app, "state"): 

292 app.state.cluster_registry = ctx.cluster_registry 

293 if admin_app is not None and hasattr(admin_app, "state"): 

294 admin_app.state.cluster_registry = ctx.cluster_registry