Initialize StreamManager with detected transport type. ENHANCED: Initializes different transport types in parallel for faster startup.
(self, namespace: str)
| 216 | logger.debug(f"Progress: {message}") |
| 217 | |
| 218 | async def _initialize_stream_manager(self, namespace: str) -> bool: |
| 219 | """Initialize StreamManager with detected transport type. |
| 220 | |
| 221 | ENHANCED: Initializes different transport types in parallel for faster startup. |
| 222 | """ |
| 223 | self.stream_manager = StreamManager() |
| 224 | |
| 225 | http_servers = self._config_loader.http_servers |
| 226 | sse_servers = self._config_loader.sse_servers |
| 227 | stdio_servers = self._config_loader.stdio_servers |
| 228 | |
| 229 | if not (http_servers or sse_servers or stdio_servers): |
| 230 | logger.info("No servers detected") |
| 231 | return True |
| 232 | |
| 233 | try: |
| 234 | # Build initialization tasks for parallel execution |
| 235 | init_tasks: list[asyncio.Task[None]] = [] |
| 236 | task_names: list[str] = [] |
| 237 | |
| 238 | # Create OAuth callback once (shared by HTTP and SSE) |
| 239 | oauth_callback = self._config_loader.create_oauth_refresh_callback( |
| 240 | http_servers, sse_servers |
| 241 | ) |
| 242 | |
| 243 | if http_servers: |
| 244 | self._report_progress( |
| 245 | f"Connecting to {len(http_servers)} HTTP server(s)..." |
| 246 | ) |
| 247 | logger.info(f"Preparing {len(http_servers)} HTTP servers for init") |
| 248 | http_dicts = [ |
| 249 | {"name": s.name, "url": s.url, "headers": s.headers or {}} |
| 250 | for s in http_servers |
| 251 | ] |
| 252 | task = asyncio.create_task( |
| 253 | self.stream_manager.initialize_with_http_streamable( |
| 254 | servers=http_dicts, |
| 255 | server_names=self.server_names, |
| 256 | initialization_timeout=self.initialization_timeout, |
| 257 | oauth_refresh_callback=oauth_callback, |
| 258 | ), |
| 259 | name="init_http", |
| 260 | ) |
| 261 | init_tasks.append(task) |
| 262 | task_names.append("HTTP") |
| 263 | |
| 264 | if sse_servers: |
| 265 | self._report_progress( |
| 266 | f"Connecting to {len(sse_servers)} SSE server(s)..." |
| 267 | ) |
| 268 | logger.info(f"Preparing {len(sse_servers)} SSE servers for init") |
| 269 | sse_dicts = [ |
| 270 | {"name": s.name, "url": s.url, "headers": s.headers or {}} |
| 271 | for s in sse_servers |
| 272 | ] |
| 273 | task = asyncio.create_task( |
| 274 | self.stream_manager.initialize_with_sse( |
| 275 | servers=sse_dicts, |