From 625788a8304938537ec834bd79a28329d15746d1 Mon Sep 17 00:00:00 2001 From: Marc Shade Date: Mon, 6 Apr 2026 19:07:42 -0400 Subject: [PATCH] Progressive SSE loading + Starlette 1.0 fix MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Dashboard changes: - SSE stream now emits per-source events as each API completes instead of waiting for all 47 sources to finish before rendering anything - Map and UI appear within seconds; each data layer populates progressively - Loading indicator shows feed progress (LOADING FEEDS 12/47) - Spinner hidden after first 3 sources arrive Backend (app.py): - Replaced monolithic _fetch_overview gather with asyncio.Queue-based progressive emission — each source fires as a task, results stream to the client as they complete - Fixed Starlette 1.0 incompatibility: replaced removed on_startup kwarg with asynccontextmanager lifespan pattern - Added proper vector store shutdown in lifespan teardown Frontend (index.html): - SSE handler now accumulates partial events into _accumulatedData - Each incoming source merges and triggers updateAll() progressively - Backwards compatible with legacy full-payload mode --- src/world_intel_mcp/dashboard/app.py | 123 +++++++++++++++++++++-- src/world_intel_mcp/dashboard/index.html | 38 ++++++- 2 files changed, 151 insertions(+), 10 deletions(-) diff --git a/src/world_intel_mcp/dashboard/app.py b/src/world_intel_mcp/dashboard/app.py index 4361398..fc68f66 100644 --- a/src/world_intel_mcp/dashboard/app.py +++ b/src/world_intel_mcp/dashboard/app.py @@ -7,6 +7,7 @@ All data pulled from the same source modules used by the MCP server. import asyncio import json import logging +from contextlib import asynccontextmanager from datetime import datetime, timezone from pathlib import Path @@ -310,14 +311,117 @@ async def api_overview(request): async def api_stream(request): - """SSE endpoint — pushes full overview every 30 seconds.""" + """SSE endpoint — progressive: emits each source as it completes.""" async def event_generator(): while True: try: - data = await _fetch_overview() - payload = json.dumps(data, default=str) - yield f"data: {payload}\n\n" + fetcher = _ensure_fetcher() + + # Sources to fetch, grouped so fast ones arrive first + source_defs = { + "market_quotes": lambda: markets.fetch_market_quotes(fetcher), + "crypto_quotes": lambda: markets.fetch_crypto_quotes(fetcher), + "macro_signals": lambda: markets.fetch_macro_signals(fetcher), + "sector_heatmap": lambda: markets.fetch_sector_heatmap(fetcher), + "earthquakes": lambda: seismology.fetch_earthquakes(fetcher), + "military_flights": lambda: military.fetch_military_flights(fetcher), + "cyber_threats": lambda: cyber.fetch_cyber_threats(fetcher), + "news_feed": lambda: news.fetch_news_feed(fetcher), + "trending_keywords": lambda: news.fetch_trending_keywords(fetcher), + "nav_warnings": lambda: maritime.fetch_nav_warnings(fetcher), + "internet_outages": lambda: infrastructure.fetch_internet_outages(fetcher), + "cable_health": lambda: infrastructure.fetch_cable_health(fetcher), + "wildfires": lambda: wildfire.fetch_wildfires(fetcher), + "prediction_markets": lambda: prediction.fetch_prediction_markets(fetcher), + "airport_delays": lambda: aviation.fetch_airport_delays(fetcher), + "climate_anomalies": lambda: climate.fetch_climate_anomalies(fetcher), + "energy_prices": lambda: economic.fetch_energy_prices(fetcher), + "gas_prices": lambda: economic.fetch_gas_prices(fetcher), + "residential_natgas": lambda: economic.fetch_residential_natgas_prices(fetcher), + "electricity_rates": lambda: economic.fetch_electricity_rates(fetcher), + "stablecoin_status": lambda: markets.fetch_stablecoin_status(fetcher), + "etf_flows": lambda: markets.fetch_etf_flows(fetcher), + "acled_events": lambda: conflict.fetch_acled_events(fetcher), + "ucdp_events": lambda: conflict.fetch_ucdp_events(fetcher), + "displacement": lambda: displacement.fetch_displacement_summary(fetcher), + "risk_scores": lambda: intelligence.fetch_risk_scores(fetcher), + "signal_convergence": lambda: intelligence.fetch_signal_convergence(fetcher), + "space_weather": lambda: space_weather.fetch_space_weather(fetcher), + "ai_watch": lambda: ai_watch.fetch_ai_watch(fetcher), + "disease_outbreaks": lambda: health.fetch_disease_outbreaks(fetcher), + "election_calendar": lambda: elections.fetch_election_calendar(fetcher), + "shipping_index": lambda: shipping.fetch_shipping_index(fetcher), + "social_signals": lambda: social.fetch_social_signals(fetcher), + "nuclear_monitor": lambda: nuclear.fetch_nuclear_monitor(fetcher), + "alert_digest": lambda: fetch_alert_digest(fetcher), + "weekly_trends": lambda: fetch_weekly_trends(fetcher), + "service_status": lambda: service_status.fetch_service_status(fetcher), + "strategic_posture": lambda: fetch_strategic_posture(fetcher), + "fleet_report": lambda: fetch_fleet_report(fetcher), + "usni_fleet": lambda: fetch_usni_fleet(fetcher), + "population_exposure": lambda: fetch_population_exposure(fetcher), + "domestic_flights": lambda: aviation.fetch_domestic_flights(fetcher), + "traffic_flow": lambda: traffic.fetch_traffic_flow(fetcher), + "traffic_incidents": lambda: traffic.fetch_traffic_incidents(fetcher), + "webcams": lambda: webcams.fetch_webcams(fetcher), + "btc_technicals": lambda: markets.fetch_btc_technicals(fetcher), + "central_bank_rates": lambda: fetch_central_bank_rates(fetcher), + } + + _SLOW = { + "news_feed", "trending_keywords", "alert_digest", + "weekly_trends", "strategic_posture", + } + results_q: asyncio.Queue = asyncio.Queue() + + async def _run_source(name: str, factory): + timeout = 90.0 if name in _SLOW else 45.0 + try: + data = await asyncio.wait_for(factory(), timeout=timeout) + except asyncio.TimeoutError: + data = {"error": f"timeout after {timeout}s", "_timeout": True} + except Exception as exc: + data = {"error": f"{type(exc).__name__}: {str(exc)[:120]}"} + await results_q.put((name, data)) + + # Fire all sources concurrently + tasks = [ + asyncio.create_task(_run_source(name, factory)) + for name, factory in source_defs.items() + ] + total = len(tasks) + completed = 0 + + # Yield each result as it arrives + while completed < total: + try: + name, data = await asyncio.wait_for(results_q.get(), timeout=120) + except asyncio.TimeoutError: + break + completed += 1 + payload = json.dumps( + {"_progressive": True, "_done": completed, "_total": total, name: data}, + default=str, + ) + yield f"data: {payload}\n\n" + + # Final event with metadata + meta = { + "_progressive": False, + "_complete": True, + "source_health": _breaker.status() if _breaker else {}, + "cache_stats": _cache.stats() if _cache else {}, + "cache_freshness": _cache.freshness() if _cache else {}, + "timestamp": datetime.now(timezone.utc).isoformat(), + } + yield f"data: {json.dumps(meta, default=str)}\n\n" + + # Clean up + for t in tasks: + if not t.done(): + t.cancel() + except asyncio.CancelledError: return except Exception as exc: @@ -515,12 +619,17 @@ async def api_vector_search_options(request): # --------------------------------------------------------------------------- -async def on_startup(): - """Start vector store worker if available.""" +@asynccontextmanager +async def lifespan(app): + """Start/stop vector store worker alongside the app.""" _ensure_fetcher() if _vector_store: await _vector_store.start() logger.info("Dashboard vector store worker started") + yield + if _vector_store: + await _vector_store.stop() + logger.info("Dashboard vector store worker stopped") app = Starlette( @@ -535,7 +644,7 @@ app = Starlette( Route("/api/vector-search", api_vector_search, methods=["POST"]), Route("/api/vector-search", api_vector_search_options, methods=["OPTIONS"]), ], - on_startup=[on_startup], + lifespan=lifespan, ) diff --git a/src/world_intel_mcp/dashboard/index.html b/src/world_intel_mcp/dashboard/index.html index 34a363c..00da0ac 100644 --- a/src/world_intel_mcp/dashboard/index.html +++ b/src/world_intel_mcp/dashboard/index.html @@ -3020,7 +3020,9 @@ function updateAll(data) { // ════════════ SSE CONNECTION ════════════ var reconnectDelay = 1000; +var _accumulatedData = {}; function connectSSE() { + _accumulatedData = {}; var source = new EventSource('/api/stream'); source.onopen = function() { $('#connDot').className = 'conn-dot'; @@ -3030,9 +3032,39 @@ function connectSSE() { source.onmessage = function(event) { try { var data = JSON.parse(event.data); - if (data.error && !data.market_quotes) return; - updateAll(data); - document.getElementById('loading').classList.add('gone'); + if (data.error && !data._progressive && !data._complete) return; + + if (data._progressive) { + // Merge this source into accumulated state + var keys = Object.keys(data); + for (var i = 0; i < keys.length; i++) { + if (keys[i].charAt(0) !== '_') _accumulatedData[keys[i]] = data[keys[i]]; + } + // Update progress indicator + var loading = document.getElementById('loading'); + if (loading && !loading.classList.contains('gone')) { + var label = loading.querySelector('.loading-label'); + if (label) label.textContent = 'LOADING FEEDS ' + data._done + '/' + data._total; + } + // Hide spinner after first 3 sources arrive — show map immediately + if (data._done >= 3) { + document.getElementById('loading').classList.add('gone'); + } + // Progressively update the dashboard with accumulated data + updateAll(_accumulatedData); + } else if (data._complete) { + // Final metadata event + var mk = Object.keys(data); + for (var j = 0; j < mk.length; j++) { + if (mk[j].charAt(0) !== '_') _accumulatedData[mk[j]] = data[mk[j]]; + } + updateAll(_accumulatedData); + document.getElementById('loading').classList.add('gone'); + } else { + // Legacy full-payload mode + updateAll(data); + document.getElementById('loading').classList.add('gone'); + } } catch (e) { console.error('Parse error:', e); } }; source.onerror = function() {