diff --git a/ROADMAP.md b/ROADMAP.md index 355fd7d..1a6a238 100644 --- a/ROADMAP.md +++ b/ROADMAP.md @@ -1,8 +1,8 @@ # World Intel MCP — Feature Parity Roadmap **Benchmark**: [koala73/worldmonitor](https://github.com/koala73/worldmonitor) -**Updated**: 2026-02-26 -**Current tools**: 89 (88 intel + 1 status) +**Updated**: 2026-03-08 +**Current tools**: 109 (108 intel + 1 status) --- @@ -16,6 +16,51 @@ --- +## 0. Current Assessment / Gap Report + +| Area | Finding | Status | Action | +|------|---------|--------|--------| +| MCP tool parity | 109 tools declared in `TOOLS`; 109 routed in `_dispatch()` | :white_check_mark: | Keep as an invariant | +| Optional vector runtime | Missing `qdrant-client` / `fastembed` previously surfaced as runtime failures | :white_check_mark: Fixed | Vector features now degrade cleanly and report availability | +| Base-environment test run | `pytest -q` fails collection without dev extras because `respx` is not installed | :yellow_circle: | Run `pip install -e ".[dev]"` before full-suite validation | +| Core verification | 77 infrastructure/vector tests pass in the base environment | :white_check_mark: | `test_cache.py`, `test_analysis.py`, `test_vector_store.py` | +| Documentation drift | Prior roadmap documented 89 tools while the codebase exposed 109 | :white_check_mark: Updated below | Keep roadmap synced with phase increments | +| Maintainability | `src/world_intel_mcp/server.py` is ~2.5k lines and remains the main refactor target | :yellow_circle: | Split tool registry and dispatch by domain | + +### Implemented Addendum Missing From Prior Roadmap + +#### Financial Intelligence Extensions (12 tools) + +| Tool | Purpose | Status | +|------|---------|--------| +| `intel_forex_rates` | Latest FX rates by base + symbol filters | :white_check_mark: | +| `intel_forex_timeseries` | Historical FX series with configurable lookback | :white_check_mark: | +| `intel_major_crosses` | Major crosses + DXY proxy snapshot | :white_check_mark: | +| `intel_yield_curve` | Treasury curve + inversion analysis | :white_check_mark: | +| `intel_bond_indices` | Bond ETF summary (AGG, TLT, HYG, LQD, TIP) | :white_check_mark: | +| `intel_earnings_calendar` | Upcoming earnings calendar | :white_check_mark: | +| `intel_earnings_surprise` | Historical earnings surprise analysis | :white_check_mark: | +| `intel_sec_filings` | Full-text SEC EDGAR filing search | :white_check_mark: | +| `intel_company_filings` | Company-specific 10-K / 10-Q / 8-K retrieval | :white_check_mark: | +| `intel_recent_8k` | Recent material-event 8-K stream | :white_check_mark: | +| `intel_company_profile` | Composite company enrichment profile | :white_check_mark: | +| `intel_macro_composite` | Weighted market regime / macro composite score | :white_check_mark: | + +#### Vector & Cross-Domain Analytics (8 tools) + +| Tool | Purpose | Status | +|------|---------|--------| +| `intel_semantic_search` | Natural-language search across accumulated intelligence | :white_check_mark: | +| `intel_similar_events` | Similarity search against historical events | :white_check_mark: | +| `intel_timeline` | Chronological timeline from vector store history | :white_check_mark: | +| `intel_vector_stats` | Qdrant collection statistics | :white_check_mark: | +| `intel_collect` | On-demand collection cycle for vector population | :white_check_mark: | +| `intel_cross_correlate` | Cross-category correlation for a topic/query | :white_check_mark: | +| `intel_domain_summary` | Per-category summary of stored intelligence | :white_check_mark: | +| `intel_trend_detection` | Recent-vs-baseline activity surge/drop detection | :white_check_mark: | + +--- + ## 1. Data Sources — Complete Inventory ### Markets & Economics (13 tools) @@ -299,18 +344,34 @@ BTC technical analysis with SMA-50/200, Mayer Multiple, golden/death cross signa ### Phase 13: Country Dossier, Full Tool Exposure & Feed Expansion (+5 = 89 tools) `intel_country_dossier`, `intel_traffic_flow`, `intel_traffic_incidents`, `intel_aviation_domestic`, `intel_webcams` -Comprehensive country intelligence dossier aggregating 6 sources in parallel (economy, markets, elections, sanctions, news, security). Exposed 4 previously hidden source functions: TomTom traffic flow (20 cities) and incidents (5 regions), OpenSky global air traffic snapshot, Windy public webcams. RSS feeds expanded from 100 to 119 across 24 categories (+6 new categories: central_asia, arctic, maritime, space, nuclear, climate). 49 CLI commands, 53 tests. +Comprehensive country intelligence dossier aggregating 6 sources in parallel (economy, markets, elections, sanctions, news, security). Exposed 4 previously hidden source functions: TomTom traffic flow (20 cities) and incidents (5 regions), OpenSky global air traffic snapshot, Windy public webcams. RSS feeds expanded from 100 to 119 across 24 categories (+6 new categories: central_asia, arctic, maritime, space, nuclear, climate). 49 CLI commands; the repo test suite has since grown well beyond this phase snapshot. + +### Phase 14: Financial Intelligence Extensions (+12 = 101 tools) +`intel_forex_rates`, `intel_forex_timeseries`, `intel_major_crosses`, `intel_yield_curve`, `intel_bond_indices`, `intel_earnings_calendar`, `intel_earnings_surprise`, `intel_sec_filings`, `intel_company_filings`, `intel_recent_8k`, `intel_company_profile`, `intel_macro_composite` + +Added business / market-intelligence depth: FX rates and timeseries, bond curves and ETF indices, earnings calendar and surprise analysis, SEC EDGAR search and company filings, composite company enrichment, and a weighted macro-composite market regime layer. + +### Phase 15: Vector Intelligence (+5 = 106 tools) +`intel_semantic_search`, `intel_similar_events`, `intel_timeline`, `intel_vector_stats`, `intel_collect` + +Qdrant-backed semantic retrieval added across all fetched intelligence, plus timeline reconstruction, store statistics, and on-demand collection. Optional dependencies remain behind `.[vector]` and now degrade cleanly when unavailable. + +### Phase 16: Cross-Domain Analytics (+3 = 109 tools) +`intel_cross_correlate`, `intel_domain_summary`, `intel_trend_detection` + +Added historical cross-category correlation, stored-data summarization, and recent-vs-baseline trend detection on top of the vector archive for early-warning and activity-shift analysis. --- ## Summary -| Category | Have | Benchmark | Coverage | -|----------|------|-----------|----------| -| Data source tools | 89 | 42 | **212%** | -| Analysis engines | 20 | 15 | **133%** | -| Static datasets | 18 | 12 | **150%** | -| RSS feeds | 119 | 150+ | **79%** | -| Strategic synthesis | Posture + brief + fleet + dossier + exposure + USNI | Dashboard-only | **Exceeds** | +| Category | Current | Notes | +|----------|---------|-------| +| Total MCP tools | 109 | 108 intelligence tools + `intel_status` | +| Tool parity | 109 / 109 | `TOOLS` and `_dispatch()` are aligned | +| Static datasets | 18 | Bases, ports, pipelines, nuclear, cables, datacenters, spaceports, minerals, exchanges, trade routes, cloud regions, financial centers | +| RSS feeds | 119 | 24 categories | +| Tests in repo | 186 | Full suite requires `.[dev]`; 77 core tests validated in the base environment | +| Primary remaining gap | Architecture | `server.py` monolith remains the main refactor target | -**Bottom line**: 89 tools across 30+ domains, over 2x WorldMonitor benchmark in tool count, 33% more analysis engines, and 50% more static datasets. 119 RSS feeds across 24 categories. All phases 1-13 complete. Live Starlette dashboard with 39 SSE streams, 14 map layers (with trade route markers), and data freshness monitoring. +**Bottom line**: 109 tools across 30+ domains, with the roadmap now aligned to the live MCP registry. The main remaining gaps are full-environment test bootstrapping (`.[dev]`) and continued modularization of the monolithic `server.py` tool registry/dispatcher. diff --git a/src/world_intel_mcp/server.py b/src/world_intel_mcp/server.py index 86a4722..937d7ea 100644 --- a/src/world_intel_mcp/server.py +++ b/src/world_intel_mcp/server.py @@ -93,11 +93,14 @@ breaker = CircuitBreaker(failure_threshold=3, cooldown_seconds=300) # Vector store — optional, degrades gracefully if Qdrant unavailable. _vector_store = None try: - from .vector_store import VectorStore + from .vector_store import VectorStore, vector_dependencies_available - _vector_store = VectorStore(enabled=True) -except Exception: - logger.info("Vector store unavailable (qdrant_client or Qdrant not installed)") + if vector_dependencies_available(): + _vector_store = VectorStore(enabled=True) + else: + logger.info("Vector store unavailable (qdrant_client / fastembed not installed)") +except Exception as exc: + logger.info("Vector store unavailable: %s", exc) fetcher = Fetcher(cache=cache, breaker=breaker, vector_store=_vector_store) @@ -2361,7 +2364,10 @@ async def _dispatch(name: str, arguments: dict[str, Any]) -> Any: # System case "intel_status": - vs_stats = {} + vs_stats = { + "enabled": False, + "error": 'Vector store dependencies not installed. Install with `pip install -e ".[vector]"`.', + } if _vector_store: vs_stats = await _vector_store.collection_stats() return { diff --git a/src/world_intel_mcp/vector_store.py b/src/world_intel_mcp/vector_store.py index 1083c4c..0f56615 100644 --- a/src/world_intel_mcp/vector_store.py +++ b/src/world_intel_mcp/vector_store.py @@ -13,7 +13,9 @@ import hashlib import json import logging import time +from dataclasses import dataclass from datetime import datetime, timezone +from importlib.util import find_spec from typing import Any logger = logging.getLogger("world-intel-mcp.vector-store") @@ -112,12 +114,62 @@ DOMAIN_CATEGORIES = { "reddit": "Social Signals", } +try: + from qdrant_client.models import ( + Distance, + FieldCondition, + Filter, + MatchValue, + PayloadSchemaType, + PointStruct, + Range, + VectorParams, + ) +except ImportError: + Distance = None + PayloadSchemaType = None + VectorParams = None + + @dataclass(slots=True) + class MatchValue: + value: Any + + @dataclass(slots=True) + class Range: + gte: float | None = None + lte: float | None = None + + @dataclass(slots=True) + class FieldCondition: + key: str + match: Any = None + range: Any = None + + @dataclass(slots=True) + class Filter: + must: list[Any] + + @dataclass(slots=True) + class PointStruct: + id: int + vector: list[float] + payload: dict[str, Any] + + +def vector_dependencies_available() -> bool: + return find_spec("fastembed") is not None and find_spec("qdrant_client") is not None + def _get_embed_model(): """Lazy-load the FastEmbed model (ONNX, ~45MB, no torch required).""" global _embed_model if _embed_model is None: - from fastembed import TextEmbedding + try: + from fastembed import TextEmbedding + except ImportError as exc: + raise RuntimeError( + 'Vector store dependencies not installed. Install with `pip install -e ".[vector]"`.' + ) from exc _embed_model = TextEmbedding(EMBEDDING_MODEL) logger.info("Loaded embedding model: %s", EMBEDDING_MODEL) @@ -128,8 +180,13 @@ def _get_qdrant(): """Lazy-load the Qdrant client and ensure collection exists.""" global _qdrant_client if _qdrant_client is None: + if not vector_dependencies_available() or any( + dep is None for dep in (Distance, VectorParams, PayloadSchemaType) + ): + raise RuntimeError( + 'Vector store dependencies not installed. Install with `pip install -e ".[vector]"`.' + ) from qdrant_client import QdrantClient - from qdrant_client.models import Distance, VectorParams _qdrant_client = QdrantClient(url=QDRANT_URL, timeout=10) @@ -144,8 +201,6 @@ def _get_qdrant(): ), ) # Create payload indexes for efficient filtering - from qdrant_client.models import PayloadSchemaType - _qdrant_client.create_payload_index( collection_name=COLLECTION_NAME, field_name="domain", @@ -317,8 +372,6 @@ class VectorStore: def _store_sync(self, domain: str, data: Any, timestamp: float) -> None: """Synchronous store operation (runs in thread pool).""" - from qdrant_client.models import PointStruct - client = _get_qdrant() text = _data_to_text(domain, data) if len(text) < 20: @@ -441,9 +494,19 @@ class VectorStore: Returns: Dict with results list, each containing score, domain, text, datetime. """ - return await asyncio.to_thread( - self._search_sync, query, limit, domain, category, hours - ) + filters = {"domain": domain, "category": category, "hours": hours} + try: + return await asyncio.to_thread( + self._search_sync, query, limit, domain, category, hours + ) + except Exception as exc: + return { + "error": str(exc), + "query": query, + "results": [], + "count": 0, + "filters": filters, + } def _search_sync( self, @@ -453,8 +516,6 @@ class VectorStore: category: str | None, hours: float | None, ) -> dict: - from qdrant_client.models import FieldCondition, Filter, MatchValue, Range - client = _get_qdrant() vector = _embed_text(query) @@ -516,7 +577,16 @@ class VectorStore: hours: float | None = None, ) -> dict: """Find historically similar events/data to a given text.""" - return await asyncio.to_thread(self._similar_sync, domain, text, limit, hours) + try: + return await asyncio.to_thread(self._similar_sync, domain, text, limit, hours) + except Exception as exc: + return { + "error": str(exc), + "reference_domain": domain, + "reference_text": text[:200], + "similar": [], + "count": 0, + } def _similar_sync( self, @@ -525,8 +595,6 @@ class VectorStore: limit: int, hours: float | None, ) -> dict: - from qdrant_client.models import FieldCondition, Filter, MatchValue, Range - client = _get_qdrant() vector = _embed_text(text) @@ -570,9 +638,19 @@ class VectorStore: limit: int = 50, ) -> dict: """Get chronological timeline of stored intelligence.""" - return await asyncio.to_thread( - self._timeline_sync, domain, category, hours, limit - ) + filters = {"domain": domain, "category": category} + try: + return await asyncio.to_thread( + self._timeline_sync, domain, category, hours, limit + ) + except Exception as exc: + return { + "error": str(exc), + "hours": hours, + "entries": [], + "count": 0, + "filters": filters, + } def _timeline_sync( self, @@ -581,8 +659,6 @@ class VectorStore: hours: float, limit: int, ) -> dict: - from qdrant_client.models import FieldCondition, Filter, MatchValue, Range - client = _get_qdrant() cutoff = time.time() - (hours * 3600) @@ -668,9 +744,19 @@ class VectorStore: Given a topic, searches all domains and groups results by category, showing how different intelligence streams relate to the same event. """ - return await asyncio.to_thread( - self._correlate_sync, query, hours, limit_per_domain - ) + try: + return await asyncio.to_thread( + self._correlate_sync, query, hours, limit_per_domain + ) + except Exception as exc: + return { + "error": str(exc), + "query": query, + "hours": hours, + "domains_found": 0, + "correlations": [], + "total_signals": 0, + } def _correlate_sync( self, @@ -678,8 +764,6 @@ class VectorStore: hours: float, limit_per_domain: int, ) -> dict: - from qdrant_client.models import FieldCondition, Filter, MatchValue, Range - client = _get_qdrant() vector = _embed_text(query) @@ -745,11 +829,18 @@ class VectorStore: async def domain_summary(self, hours: float = 24.0) -> dict: """Get per-domain summary of stored intelligence.""" - return await asyncio.to_thread(self._domain_summary_sync, hours) + try: + return await asyncio.to_thread(self._domain_summary_sync, hours) + except Exception as exc: + return { + "error": str(exc), + "hours": hours, + "total_data_points": 0, + "categories": 0, + "summary": [], + } def _domain_summary_sync(self, hours: float) -> dict: - from qdrant_client.models import FieldCondition, Filter, MatchValue, Range - client = _get_qdrant() cutoff = time.time() - (hours * 3600) @@ -843,9 +934,20 @@ class VectorStore: Compares data point density in the recent window against the baseline to identify surges or drops in intelligence activity. """ - return await asyncio.to_thread( - self._trend_sync, category, recent_hours, baseline_hours - ) + try: + return await asyncio.to_thread( + self._trend_sync, category, recent_hours, baseline_hours + ) + except Exception as exc: + return { + "error": str(exc), + "recent_window_hours": recent_hours, + "baseline_window_hours": baseline_hours, + "categories_analyzed": 0, + "surges": 0, + "drops": 0, + "trends": [], + } def _trend_sync( self, @@ -853,8 +955,6 @@ class VectorStore: recent_hours: float, baseline_hours: float, ) -> dict: - from qdrant_client.models import FieldCondition, Filter, MatchValue, Range - client = _get_qdrant() now = time.time() recent_cutoff = now - (recent_hours * 3600)