diff --git a/api/.env.example b/api/.env.example index 518a7940..c85070bd 100644 --- a/api/.env.example +++ b/api/.env.example @@ -25,6 +25,8 @@ NGNMARKET_API_BASE_URL=https://api.ngnmarket.com/v1/ NGNMARKET_API_KEY=value RAPID_API_BASE_URL=https://tradingview-data1.p.rapidapi.com/ RAPID_API_KEY=value +NGXPULSE_API_BASE_URL=https://www.ngxpulse.ng +NGXPULSE_API_KEY=value # Scheduler SCHEDULER_ENABLED=true diff --git a/api/.env.prod.example b/api/.env.prod.example index f7408f61..0ddb1117 100644 --- a/api/.env.prod.example +++ b/api/.env.prod.example @@ -24,6 +24,8 @@ NGNMARKET_API_BASE_URL=https://api.ngnmarket.com/v1/ NGNMARKET_API_KEY=value RAPID_API_BASE_URL=https://tradingview-data1.p.rapidapi.com/ RAPID_API_KEY=value +NGXPULSE_API_BASE_URL=https://www.ngxpulse.ng +NGXPULSE_API_KEY=value # Scheduler SCHEDULER_ENABLED=true diff --git a/api/adapters/inbound/http/asset_routes.py b/api/adapters/inbound/http/asset_routes.py index d99539c4..0ec2b038 100644 --- a/api/adapters/inbound/http/asset_routes.py +++ b/api/adapters/inbound/http/asset_routes.py @@ -9,6 +9,7 @@ from sqlalchemy.ext.asyncio import AsyncSession from adapters.outbound.market_data.ngnmarket_adapter import NgnMarketAdapter +from adapters.outbound.market_data.ngxpulse_adapter import NgxPulseAdapter from adapters.outbound.market_data.tiingo_adapter import TiingoAdapter from adapters.outbound.market_data.tradingview_adapter import TradingviewAdapter from adapters.outbound.market_data.yfinance_adapter import YFinanceAdapter @@ -25,7 +26,7 @@ def _normalize_provider(provider: str) -> str: normalized = (provider or "yfinance").strip().lower() return ( normalized - if normalized in {"yfinance", "tiingo", "ngnmarket", "tradingview"} + if normalized in {"yfinance", "tiingo", "ngnmarket", "ngxpulse", "tradingview"} else "yfinance" ) @@ -48,7 +49,9 @@ async def _fetch_fx_history(from_ccy: str, start: date, end: date) -> dict[date, rate = item.get("rate") if date_str and rate is not None: try: - result[date.fromisoformat(date_str)] = round(float(rate) * 100) + result[date.fromisoformat(date_str)] = round( + float(rate) * 100 + ) except Exception: pass if result: @@ -60,14 +63,20 @@ async def _fetch_fx_history(from_ccy: str, start: date, end: date) -> dict[date, t = yf.Ticker(f"{from_ccy}USD=X") data = t.history(start=start, end=end) if not data.empty: - return {idx.date(): round(float(row["Close"]) * 100) for idx, row in data.iterrows()} + return { + idx.date(): round(float(row["Close"]) * 100) + for idx, row in data.iterrows() + } except Exception: pass try: t = yf.Ticker(f"USD{from_ccy}=X") data = t.history(start=start, end=end) if not data.empty: - return {idx.date(): -round(float(row["Close"]) * 100) for idx, row in data.iterrows()} + return { + idx.date(): -round(float(row["Close"]) * 100) + for idx, row in data.iterrows() + } except Exception: pass return {} @@ -96,6 +105,7 @@ async def get_price_history_batch( tiingo = TiingoAdapter() tradingview = TradingviewAdapter() ngnmarket = NgnMarketAdapter() + ngxpulse = NgxPulseAdapter() def _pick_adapter(ticker: str): provider = _normalize_provider(provider_map.get(ticker.upper(), "yfinance")) @@ -105,6 +115,8 @@ def _pick_adapter(ticker: str): return tiingo if provider == "ngnmarket": return ngnmarket + if provider == "ngxpulse": + return ngxpulse return yfinance async def _fetch(ticker: str): @@ -196,6 +208,9 @@ async def validate_ticker( elif selected == "ngnmarket": adapter = NgnMarketAdapter() metadata = await adapter.get_asset_metadata(ticker, currency, asset_class) + elif selected == "ngxpulse": + adapter = NgxPulseAdapter() + metadata = await adapter.get_asset_metadata(ticker, currency, asset_class) elif selected == "tradingview": adapter = TradingviewAdapter() metadata = await adapter.get_asset_metadata(ticker, currency) @@ -236,6 +251,8 @@ async def get_price_history( adapter = TiingoAdapter() elif provider == "ngnmarket": adapter = NgnMarketAdapter() + elif provider == "ngxpulse": + adapter = NgxPulseAdapter() else: adapter = YFinanceAdapter() diff --git a/api/adapters/outbound/market_data/ngxpulse_adapter.py b/api/adapters/outbound/market_data/ngxpulse_adapter.py new file mode 100644 index 00000000..fd2d137d --- /dev/null +++ b/api/adapters/outbound/market_data/ngxpulse_adapter.py @@ -0,0 +1,534 @@ +import asyncio +import logging +from datetime import date, datetime, timedelta, timezone +from typing import Optional + +import aiohttp + +from adapters.outbound.market_data.price_cache import PriceCache +from domain.value_objects.money import AssetMetadata, Currency +from infrastructure.config import settings + +logger = logging.getLogger(__name__) + + +class NgxPulseAdapter: + """Adapter for the NGXPulse market data API. + + Endpoints: + - GET /api/ngxdata/stocks — 150+ NGX equities snapshot + - GET /api/ngxdata/etfs — NGX ETF universe snapshot + - GET /api/ngxdata/indices/{code}/history — index daily value series + + Auth: X-API-Key header + Base URL: https://www.ngxpulse.ng + + Rate limits: + - 10 requests per minute + - 100 requests per day + """ + + _shared_session: Optional[aiohttp.ClientSession] = None + _session_lock: Optional[asyncio.Lock] = None + _shared_timeout = aiohttp.ClientTimeout(total=15) + + # Rate limits (per-process rolling windows) + _MAX_REQUESTS_PER_MINUTE = 8 # headroom below the 10/min cap + _MAX_REQUESTS_PER_DAY = 90 # headroom below the 100/day cap + + # In-memory cache for the bulk stocks / ETFs endpoints + _CACHE_TTL = timedelta(minutes=30) + + def __init__(self) -> None: + self._base_url = (settings.NGXPULSE_API_BASE_URL or "").rstrip("/") + self._api_key = settings.NGXPULSE_API_KEY + self._enabled = bool(self._base_url and self._api_key) + + # Rate-limit windows + self._minute_window_start: datetime = datetime.now(timezone.utc).replace( + second=0, microsecond=0 + ) + self._minute_count = 0 + + self._day_window_start: date = datetime.now(timezone.utc).date() + self._day_count = 0 + + # Cached bulk responses + self._stocks_cache: Optional[list[dict]] = None + self._stocks_cached_at: Optional[datetime] = None + + self._etfs_cache: Optional[dict] = None + self._etfs_cached_at: Optional[datetime] = None + + # Valkey-backed cache for individual ticker results (shared + # across processes, same infrastructure as the other adapters). + self._cache = PriceCache() + + # ------------------------------------------------------------------ + # Public API — asset metadata + # ------------------------------------------------------------------ + + async def get_asset_metadata( + self, ticker: str, currency: str = "NGN", asset_class: str = "" + ) -> Optional[AssetMetadata]: + """Resolve metadata for *ticker* from the stocks or ETFs snapshots. + + Individual ticker results are cached in Valkey via PriceCache + so repeated lookups don't re-scan the bulk snapshots. + """ + if not self._enabled: + return None + + symbol = self._extract_symbol(ticker) + if not symbol: + return None + + # Check the shared Valkey cache first (same infra as other adapters). + cached = await self._cache.get_metadata(ticker) + if cached is not None: + return cached + + # Try stocks first, then ETFs (scans the in-memory bulk cache). + detail = await self._lookup_stock(symbol) + resolved_class = "stock" + if not detail: + detail = await self._lookup_etf(symbol) + resolved_class = "etf" + + if not detail: + return None + + try: + resolved_currency = self._safe_currency(currency) + metadata = AssetMetadata( + ticker=ticker, + name=detail.get("name") or symbol, + asset_class=resolved_class, + currency=resolved_currency, + exchange="NGX", + sector=detail.get("sector"), + industry=None, + country="NG", + isin=detail.get("isin"), + ) + # Persist to shared cache for subsequent lookups. + await self._cache.set_metadata(metadata) + return metadata + except Exception as e: + logger.debug("NgxPulse metadata mapping failed for %s: %s", ticker, e) + return None + + # ------------------------------------------------------------------ + # Public API — price history + # ------------------------------------------------------------------ + + async def get_price_history( + self, ticker: str, start: date, end: date + ) -> list[tuple[date, int]]: + """Return daily index values for an NGX index code. + + Individual stock/ETF price history is not provided by this API; + only index history is available. Checks the stocks/ETFs cache + first to avoid wasting API calls on non-index tickers. + """ + if not self._enabled: + return [] + + symbol = self._extract_symbol(ticker) + if not symbol: + return [] + + # Skip index history call if this ticker is a known stock or ETF — + # NGXPulse only provides historical data for indices. + if await self._lookup_stock(symbol) or await self._lookup_etf(symbol): + return [] + + # Try to fetch as an index history + index_data = await self.get_index_history( + code=symbol, + from_date=start, + to_date=end, + ) + if not index_data: + return [] + + history = index_data.get("history") or [] + result = [] + for entry in history: + date_str = entry.get("date") + value = entry.get("value") + if not date_str or value is None: + continue + try: + d = date.fromisoformat(date_str) + except Exception: + continue + if start <= d <= end: + result.append((d, round(float(value) * 100))) + + if result: + logger.info( + "NgxPulse: price history OK for %s — %d data points", + ticker, + len(result), + ) + else: + logger.debug( + "NgxPulse: no price history data for %s (%s → %s)", + ticker, + start, + end, + ) + return sorted(result) + + async def get_current_price( + self, ticker: str, currency: str = "NGN" + ) -> tuple[date, int]: + """Return (date, price×100) for *ticker*. + + Checks the shared PriceCache first, then falls back to the + in-memory stocks/ETFs snapshots, then the index history endpoint. + Results are persisted to PriceCache on success. + """ + if not self._enabled: + return (date.today(), 0) + + symbol = self._extract_symbol(ticker) + if not symbol: + return (date.today(), 0) + + # Check the shared Valkey cache first. + cached = await self._cache.get_price(ticker) + if cached is not None: + # A sentinel (price=0) means not found — skip upstream. + if cached[1] == 0 and cached[0] == date.min: + return (date.today(), 0) + return cached + + # Try stocks snapshot + stock = await self._lookup_stock(symbol) + if stock: + price = stock.get("current_price") + if price is not None: + try: + price_cents = round(float(price) * 100) + await self._cache.set_price(ticker, date.today(), price_cents) + return (date.today(), price_cents) + except Exception: + pass + + # Try ETFs snapshot + etf = await self._lookup_etf(symbol) + if etf: + price = etf.get("close") or etf.get("previous_close") + if price is not None: + try: + price_cents = round(float(price) * 100) + await self._cache.set_price(ticker, date.today(), price_cents) + return (date.today(), price_cents) + except Exception: + pass + + # Fallback: try index history (real API call — check rate limit) + if not self._can_make_request(): + await self._cache.set_price(ticker, date.min, 0) + return (date.today(), 0) + + chart = await self._get(f"/api/ngxdata/indices/{symbol}/history") + if isinstance(chart, dict) and chart.get("success") and chart.get("history"): + history_list = chart["history"] + if not isinstance(history_list, list) or not history_list: + await self._cache.set_price(ticker, date.min, 0) + return (date.today(), 0) + last = history_list[-1] + value = last.get("value") + if value is not None: + try: + price_cents = round(float(value) * 100) + await self._cache.set_price(ticker, date.today(), price_cents) + return (date.today(), price_cents) + except Exception: + pass + + # Cache sentinel so repeated lookups don't hit the API again. + await self._cache.set_price(ticker, date.min, 0) + return (date.today(), 0) + + # ------------------------------------------------------------------ + # Public API — bulk data + # ------------------------------------------------------------------ + + async def get_stocks(self) -> Optional[list[dict]]: + """Return the full NGX equities list (cached for 30 min).""" + now = datetime.now(timezone.utc) + if ( + self._stocks_cache is not None + and self._stocks_cached_at is not None + and now - self._stocks_cached_at <= self._CACHE_TTL + ): + return self._stocks_cache + + if not self._can_make_request(): + return None + + payload = await self._get("/api/ngxdata/stocks") + if isinstance(payload, list): + self._stocks_cache = payload + self._stocks_cached_at = now + logger.info("NgxPulse: stocks snapshot fetched — %d entries", len(payload)) + return payload + + # The API may return a wrapped dict — try several unwrapping strategies. + if isinstance(payload, dict): + # Strategy 1: {"success": true, "data": [...], ...} + data = payload.get("data") + if isinstance(data, list): + self._stocks_cache = data + self._stocks_cached_at = now + logger.info( + "NgxPulse: stocks snapshot fetched (wrapped by 'data') — %d entries", + len(data), + ) + return data + + # Strategy 2: {"success": true, "stocks": [...], ...} etc. + for key, val in payload.items(): + if isinstance(val, list) and len(val) > 0: + self._stocks_cache = val + self._stocks_cached_at = now + logger.info( + "NgxPulse: stocks snapshot found under key '%s' — %d entries", + key, + len(val), + ) + return val + + # Strategy 3: log the actual keys for debugging + logger.warning( + "NgxPulse: unexpected stocks response dict — keys=%s", + list(payload.keys())[:10], + ) + return None + + logger.warning("NgxPulse: unexpected stocks response type: %s", type(payload)) + return None + + async def get_etfs(self) -> Optional[dict]: + """Return the full NGX ETF snapshot (cached for 30 min).""" + now = datetime.now(timezone.utc) + if ( + self._etfs_cache is not None + and self._etfs_cached_at is not None + and now - self._etfs_cached_at <= self._CACHE_TTL + ): + return self._etfs_cache + + if not self._can_make_request(): + return None + + payload = await self._get("/api/ngxdata/etfs") + if isinstance(payload, dict) and payload.get("success"): + self._etfs_cache = payload + self._etfs_cached_at = now + count = payload.get("count", 0) + logger.info("NgxPulse: ETFs snapshot fetched — %d entries", count) + return payload + + logger.warning("NgxPulse: unexpected ETFs response: %s", payload) + return None + + async def get_index_history( + self, + code: str, + from_date: Optional[date] = None, + to_date: Optional[date] = None, + ) -> Optional[dict]: + """Return daily value series for an NGX index (ASI, ngx-30, etc.).""" + if not self._can_make_request(): + return None + + params = {} + if from_date: + params["from"] = from_date.isoformat() + if to_date: + params["to"] = to_date.isoformat() + + payload = await self._get( + f"/api/ngxdata/indices/{code.lower()}/history", params=params + ) + if isinstance(payload, dict) and payload.get("success"): + logger.info( + "NgxPulse: index history OK for %s — %d entries", + code, + payload.get("count", 0), + ) + return payload + + logger.debug("NgxPulse: no index history for %s", code) + return None + + # ------------------------------------------------------------------ + # Rate limiting + # ------------------------------------------------------------------ + + def _can_make_request(self) -> bool: + now = datetime.now(timezone.utc) + + minute_start = now.replace(second=0, microsecond=0) + if minute_start != self._minute_window_start: + self._minute_window_start = minute_start + self._minute_count = 0 + + day_start = now.date() + if day_start != self._day_window_start: + self._day_window_start = day_start + self._day_count = 0 + + if self._minute_count >= self._MAX_REQUESTS_PER_MINUTE: + logger.warning( + "NgxPulse per-minute request limit reached (%d/%d); skipping request", + self._minute_count, + self._MAX_REQUESTS_PER_MINUTE, + ) + return False + + if self._day_count >= self._MAX_REQUESTS_PER_DAY: + logger.warning( + "NgxPulse daily request limit reached (%d/%d); skipping request", + self._day_count, + self._MAX_REQUESTS_PER_DAY, + ) + return False + + self._minute_count += 1 + self._day_count += 1 + return True + + # ------------------------------------------------------------------ + # Internal helpers + # ------------------------------------------------------------------ + + async def _lookup_stock(self, symbol: str) -> Optional[dict]: + """Find a stock by symbol from the cached stocks snapshot.""" + stocks = await self.get_stocks() + if not stocks: + return None + + upper = symbol.upper() + for stock in stocks: + if str(stock.get("symbol", "")).upper() == upper: + return stock + return None + + async def _lookup_etf(self, symbol: str) -> Optional[dict]: + """Find an ETF by symbol or canonical_symbol from the cached ETF snapshot.""" + etfs_payload = await self.get_etfs() + if not etfs_payload: + return None + + data = etfs_payload.get("data") or [] + upper = symbol.upper() + for etf in data: + sym = str(etf.get("symbol", "")).upper() + canonical = str(etf.get("canonical_symbol", "")).upper() + if sym == upper or canonical == upper: + return etf + return None + + async def _get( + self, path: str, params: Optional[dict] = None + ) -> Optional[dict | list]: + if not self._enabled: + return None + + url = f"{self._base_url}{path}" + headers = { + "X-API-Key": self._api_key, + "Content-Type": "application/json", + } + + logger.debug("NgxPulse: GET %s params=%s", path, params) + try: + session = await self._get_shared_session() + async with session.get(url, params=params, headers=headers) as response: + status_code = response.status + body_text = await response.text() + + if status_code >= 400: + self._log_api_error(path, status_code, body_text) + return None + + try: + data = await response.json() + logger.debug("NgxPulse: GET %s → %s OK", path, status_code) + return data + except Exception: + logger.warning("NgxPulse GET %s returned non-JSON response", path) + return None + except Exception as e: + logger.warning("NgxPulse: GET %s failed: %s", path, e) + return None + + @classmethod + async def _get_shared_session(cls) -> aiohttp.ClientSession: + if cls._session_lock is None: + cls._session_lock = asyncio.Lock() + + async with cls._session_lock: + if cls._shared_session is None or cls._shared_session.closed: + cls._shared_session = aiohttp.ClientSession(timeout=cls._shared_timeout) + + return cls._shared_session + + @classmethod + async def close_shared_session(cls) -> None: + if cls._session_lock is None: + cls._session_lock = asyncio.Lock() + + async with cls._session_lock: + if cls._shared_session and not cls._shared_session.closed: + await cls._shared_session.close() + cls._shared_session = None + + @staticmethod + def _extract_symbol(ticker: str) -> str: + raw = (ticker or "").strip() + if not raw: + return "" + + for separator in (":", "/", "|"): + if separator in raw: + left, right = raw.split(separator, 1) + if left.strip().upper() in {"NGX", "NGXPULSE"}: + return right.strip().upper() + return raw.upper() + + @staticmethod + def _safe_currency(value: str) -> Currency: + try: + return Currency(value.upper()) + except ValueError: + return Currency.NGN + + @staticmethod + def _log_api_error(path: str, status_code: int, body_text: str) -> None: + try: + import json + + payload = json.loads(body_text) + message = payload.get("message", "") + error_flag = payload.get("error") + logger.warning( + "NgxPulse error on %s: status=%s error=%s message=%s", + path, + status_code, + error_flag, + message or body_text[:200], + ) + except Exception: + logger.warning( + "NgxPulse error on %s: status=%s body=%.200s", + path, + status_code, + body_text, + ) diff --git a/api/application/analytics/analytics_interactor.py b/api/application/analytics/analytics_interactor.py index 5fbd8cae..edea7f90 100644 --- a/api/application/analytics/analytics_interactor.py +++ b/api/application/analytics/analytics_interactor.py @@ -5,6 +5,7 @@ from sqlalchemy.ext.asyncio import AsyncSession from adapters.outbound.market_data.ngnmarket_adapter import NgnMarketAdapter +from adapters.outbound.market_data.ngxpulse_adapter import NgxPulseAdapter from adapters.outbound.market_data.price_cache import PriceCache from adapters.outbound.market_data.tiingo_adapter import TiingoAdapter from adapters.outbound.market_data.tradingview_adapter import TradingviewAdapter @@ -35,6 +36,7 @@ def __init__(self, session: AsyncSession, portfolio_base_currency: Currency): self.yfinance = YFinanceAdapter() self.tiingo = TiingoAdapter() self.ngnmarket = NgnMarketAdapter() + self.ngxpulse = NgxPulseAdapter() self.tradingview = TradingviewAdapter() self.base_currency = portfolio_base_currency @@ -43,7 +45,8 @@ def _normalize_provider(provider: str | None) -> str: normalized = (provider or "yfinance").strip().lower() return ( normalized - if normalized in {"yfinance", "tiingo", "ngnmarket", "tradingview"} + if normalized + in {"yfinance", "tiingo", "ngnmarket", "ngxpulse", "tradingview"} else "yfinance" ) @@ -62,6 +65,9 @@ async def _get_price_from_provider(self, symbol: str, provider: str): return (date.today(), round(float(price) * 100)) return (date.today(), 0) + if selected == "ngxpulse": + return await self.ngxpulse.get_current_price(symbol) + if selected == "tradingview": return await self.tradingview.get_current_price(symbol) diff --git a/api/application/trades/csv_import_interactor.py b/api/application/trades/csv_import_interactor.py index 56b75647..d2f10776 100644 --- a/api/application/trades/csv_import_interactor.py +++ b/api/application/trades/csv_import_interactor.py @@ -7,6 +7,7 @@ from sqlalchemy.ext.asyncio import AsyncSession from adapters.outbound.market_data.ngnmarket_adapter import NgnMarketAdapter +from adapters.outbound.market_data.ngxpulse_adapter import NgxPulseAdapter from adapters.outbound.market_data.tiingo_adapter import TiingoAdapter from adapters.outbound.market_data.tradingview_adapter import TradingviewAdapter from adapters.outbound.market_data.yfinance_adapter import YFinanceAdapter @@ -25,6 +26,7 @@ def __init__(self, session: AsyncSession): self.yfinance = YFinanceAdapter() self.tiingo = TiingoAdapter() self.ngnmarket = NgnMarketAdapter() + self.ngxpulse = NgxPulseAdapter() self.tradingview = TradingviewAdapter() @staticmethod @@ -32,11 +34,14 @@ def _normalize_provider(provider: str | None) -> str: normalized = (provider or "yfinance").strip().lower() return ( normalized - if normalized in {"yfinance", "tiingo", "ngnmarket", "tradingview"} + if normalized + in {"yfinance", "tiingo", "ngnmarket", "ngxpulse", "tradingview"} else "yfinance" ) - async def _get_asset_metadata(self, ticker: str, currency: Currency, provider: str, asset_class: str = ""): + async def _get_asset_metadata( + self, ticker: str, currency: Currency, provider: str, asset_class: str = "" + ): selected = self._normalize_provider(provider) if selected == "tiingo": @@ -46,13 +51,17 @@ async def _get_asset_metadata(self, ticker: str, currency: Currency, provider: s metadata = await self.tradingview.get_asset_metadata(ticker, currency.value) if metadata: return metadata - metadata = await self.ngnmarket.get_asset_metadata(ticker, currency.value, asset_class) + metadata = await self.ngnmarket.get_asset_metadata( + ticker, currency.value, asset_class + ) if metadata: return metadata return await self.yfinance.get_asset_metadata(ticker, currency.value) if selected == "ngnmarket": - metadata = await self.ngnmarket.get_asset_metadata(ticker, currency.value, asset_class) + metadata = await self.ngnmarket.get_asset_metadata( + ticker, currency.value, asset_class + ) if metadata: return metadata metadata = await self.tradingview.get_asset_metadata(ticker, currency.value) @@ -70,7 +79,23 @@ async def _get_asset_metadata(self, ticker: str, currency: Currency, provider: s metadata = await self.tiingo.get_asset_metadata(ticker, currency.value) if metadata: return metadata - metadata = await self.ngnmarket.get_asset_metadata(ticker, currency.value, asset_class) + metadata = await self.ngnmarket.get_asset_metadata( + ticker, currency.value, asset_class + ) + if metadata: + return metadata + return await self.yfinance.get_asset_metadata(ticker, currency.value) + + if selected == "ngxpulse": + metadata = await self.ngxpulse.get_asset_metadata( + ticker, currency.value, asset_class + ) + if metadata: + return metadata + metadata = await self.tiingo.get_asset_metadata(ticker, currency.value) + if metadata: + return metadata + metadata = await self.tradingview.get_asset_metadata(ticker, currency.value) if metadata: return metadata return await self.yfinance.get_asset_metadata(ticker, currency.value) diff --git a/api/application/trades/trade_interactor.py b/api/application/trades/trade_interactor.py index 3d5ceca5..426ad5ac 100644 --- a/api/application/trades/trade_interactor.py +++ b/api/application/trades/trade_interactor.py @@ -6,6 +6,7 @@ from sqlalchemy.ext.asyncio import AsyncSession from adapters.outbound.market_data.ngnmarket_adapter import NgnMarketAdapter +from adapters.outbound.market_data.ngxpulse_adapter import NgxPulseAdapter from adapters.outbound.market_data.tiingo_adapter import TiingoAdapter from adapters.outbound.market_data.tradingview_adapter import TradingviewAdapter from adapters.outbound.market_data.yfinance_adapter import YFinanceAdapter @@ -30,6 +31,7 @@ def __init__(self, session: AsyncSession): self.yfinance = YFinanceAdapter() self.tiingo = TiingoAdapter() self.ngnmarket = NgnMarketAdapter() + self.ngxpulse = NgxPulseAdapter() self.tradingview = TradingviewAdapter() @staticmethod @@ -37,7 +39,8 @@ def _normalize_provider(provider: str | None) -> str: normalized = (provider or "yfinance").strip().lower() return ( normalized - if normalized in {"yfinance", "tiingo", "ngnmarket", "tradingview"} + if normalized + in {"yfinance", "tiingo", "ngnmarket", "ngxpulse", "tradingview"} else "yfinance" ) @@ -53,6 +56,9 @@ async def _get_asset_metadata(self, ticker: str, currency: Currency, provider: s if selected == "tradingview": return await self.tradingview.get_asset_metadata(ticker, currency.value) + if selected == "ngxpulse": + return await self.ngxpulse.get_asset_metadata(ticker, currency.value) + return await self.yfinance.get_asset_metadata(ticker, currency.value) async def _resolve_asset( diff --git a/api/infrastructure/config/env.py b/api/infrastructure/config/env.py index f3be0c27..b423c541 100644 --- a/api/infrastructure/config/env.py +++ b/api/infrastructure/config/env.py @@ -60,6 +60,10 @@ class EnvironSettings(BaseModel): "RAPID_API_BASE_URL", "https://tradingview-data1.p.rapidapi.com/" ) RAPID_API_KEY: str = os.getenv("RAPID_API_KEY", "") + NGXPULSE_API_BASE_URL: str = os.getenv( + "NGXPULSE_API_BASE_URL", "https://www.ngxpulse.ng" + ) + NGXPULSE_API_KEY: str = os.getenv("NGXPULSE_API_KEY", "") # Scheduler SCHEDULER_ENABLED: bool = os.getenv("SCHEDULER_ENABLED", "false").lower() == "true" diff --git a/api/pyproject.toml b/api/pyproject.toml index e67d320b..b670d6aa 100644 --- a/api/pyproject.toml +++ b/api/pyproject.toml @@ -1,6 +1,6 @@ [project] name = "folio-api" -version = "1.22.3" +version = "1.23.0" description = "Add your description here" readme = "README.md" requires-python = ">=3.13" diff --git a/api/tests/unit/adapters/outbound/market_data/test_ngxpulse_adapter.py b/api/tests/unit/adapters/outbound/market_data/test_ngxpulse_adapter.py new file mode 100644 index 00000000..b8222b6b --- /dev/null +++ b/api/tests/unit/adapters/outbound/market_data/test_ngxpulse_adapter.py @@ -0,0 +1,801 @@ +from datetime import date, datetime, timedelta, timezone + +import aiohttp +import pytest + +from adapters.outbound.market_data import ngxpulse_adapter as ngx_module +from domain.value_objects.money import Currency + +# --------------------------------------------------------------------------- +# Helpers +# --------------------------------------------------------------------------- + + +def _enable_adapter(monkeypatch) -> ngx_module.NgxPulseAdapter: + monkeypatch.setattr( + ngx_module.settings, + "NGXPULSE_API_BASE_URL", + "https://www.ngxpulse.ng", + raising=False, + ) + monkeypatch.setattr( + ngx_module.settings, + "NGXPULSE_API_KEY", + "test-ngxpulse-key", + raising=False, + ) + return ngx_module.NgxPulseAdapter() + + +def _reset_rate_limits(adapter: ngx_module.NgxPulseAdapter) -> None: + adapter._minute_count = 0 + adapter._day_count = 0 + adapter._minute_window_start = datetime.now(timezone.utc).replace( + second=0, microsecond=0 + ) + adapter._day_window_start = datetime.now(timezone.utc).date() + + +# --------------------------------------------------------------------------- +# get_asset_metadata +# --------------------------------------------------------------------------- + + +@pytest.mark.asyncio +async def test_get_asset_metadata_from_stocks(monkeypatch): + adapter = _enable_adapter(monkeypatch) + + async def fake_lookup_stock(symbol: str): + return { + "symbol": symbol, + "name": "Dangote Cement Plc", + "sector": "Industrial", + } + + async def fake_lookup_etf(symbol: str): + return None + + monkeypatch.setattr(adapter, "_lookup_stock", fake_lookup_stock) + monkeypatch.setattr(adapter, "_lookup_etf", fake_lookup_etf) + + metadata = await adapter.get_asset_metadata("DANGCEM", "NGN") + + assert metadata is not None + assert metadata.ticker == "DANGCEM" + assert metadata.name == "Dangote Cement Plc" + assert metadata.asset_class == "stock" + assert metadata.currency == Currency.NGN + assert metadata.exchange == "NGX" + assert metadata.sector == "Industrial" + assert metadata.country == "NG" + + +@pytest.mark.asyncio +async def test_get_asset_metadata_from_etfs(monkeypatch): + adapter = _enable_adapter(monkeypatch) + + async def fake_lookup_stock(symbol: str): + return None + + async def fake_lookup_etf(symbol: str): + return { + "symbol": "NEWGOLD", + "name": "NewGold ETF", + "isin": "NGNEWGOLD001", + "sector": "Exchange Traded Fund", + } + + monkeypatch.setattr(adapter, "_lookup_stock", fake_lookup_stock) + monkeypatch.setattr(adapter, "_lookup_etf", fake_lookup_etf) + + metadata = await adapter.get_asset_metadata("NEWGOLD", "NGN", asset_class="etf") + + assert metadata is not None + assert metadata.ticker == "NEWGOLD" + assert metadata.name == "NewGold ETF" + assert metadata.asset_class == "etf" + assert metadata.isin == "NGNEWGOLD001" + assert metadata.exchange == "NGX" + assert metadata.country == "NG" + + +@pytest.mark.asyncio +async def test_get_asset_metadata_not_found(monkeypatch): + adapter = _enable_adapter(monkeypatch) + + async def fake_lookup_stock(symbol: str): + return None + + async def fake_lookup_etf(symbol: str): + return None + + monkeypatch.setattr(adapter, "_lookup_stock", fake_lookup_stock) + monkeypatch.setattr(adapter, "_lookup_etf", fake_lookup_etf) + + metadata = await adapter.get_asset_metadata("UNKNOWN", "NGN") + assert metadata is None + + +@pytest.mark.asyncio +async def test_get_asset_metadata_disabled(monkeypatch): + monkeypatch.setattr(ngx_module.settings, "NGXPULSE_API_KEY", "", raising=False) + + adapter = ngx_module.NgxPulseAdapter() + metadata = await adapter.get_asset_metadata("DANGCEM", "NGN") + assert metadata is None + + +@pytest.mark.asyncio +async def test_get_asset_metadata_stock_preferred_over_etf(monkeypatch): + """When a symbol appears in both stocks and ETFs, stocks wins.""" + adapter = _enable_adapter(monkeypatch) + + async def fake_lookup_stock(symbol: str): + return {"symbol": symbol, "name": "Stock Name", "sector": "Industrial"} + + async def fake_lookup_etf(symbol: str): + return {"symbol": symbol, "name": "ETF Name", "isin": "NGETFX0001"} + + monkeypatch.setattr(adapter, "_lookup_stock", fake_lookup_stock) + monkeypatch.setattr(adapter, "_lookup_etf", fake_lookup_etf) + + metadata = await adapter.get_asset_metadata("DUAL", "NGN") + + assert metadata is not None + assert metadata.asset_class == "stock" + assert metadata.name == "Stock Name" + + +# --------------------------------------------------------------------------- +# get_price_history +# --------------------------------------------------------------------------- + + +@pytest.mark.asyncio +async def test_get_price_history_from_index(monkeypatch): + adapter = _enable_adapter(monkeypatch) + + async def fake_lookup_stock(symbol: str): + return None + + async def fake_lookup_etf(symbol: str): + return None + + async def fake_get_index_history(code, from_date=None, to_date=None): + return { + "success": True, + "code": "ASI", + "count": 2, + "history": [ + {"date": "2026-01-02", "value": 26867.79}, + {"date": "2026-01-03", "value": 27000.50}, + ], + } + + monkeypatch.setattr(adapter, "_lookup_stock", fake_lookup_stock) + monkeypatch.setattr(adapter, "_lookup_etf", fake_lookup_etf) + monkeypatch.setattr(adapter, "get_index_history", fake_get_index_history) + + history = await adapter.get_price_history("ASI", date(2026, 1, 1), date(2026, 1, 5)) + + assert len(history) == 2 + assert history[0] == (date(2026, 1, 2), 2686779) + assert history[1] == (date(2026, 1, 3), 2700050) + + +@pytest.mark.asyncio +async def test_get_price_history_skips_known_stock(monkeypatch): + adapter = _enable_adapter(monkeypatch) + index_called = False + + async def fake_lookup_stock(symbol: str): + return {"symbol": "DANGCEM", "name": "Dangote Cement"} + + async def fake_lookup_etf(symbol: str): + return None + + async def fake_get_index_history(code, from_date=None, to_date=None): + nonlocal index_called + index_called = True + return None + + monkeypatch.setattr(adapter, "_lookup_stock", fake_lookup_stock) + monkeypatch.setattr(adapter, "_lookup_etf", fake_lookup_etf) + monkeypatch.setattr(adapter, "get_index_history", fake_get_index_history) + + history = await adapter.get_price_history( + "DANGCEM", date(2026, 1, 1), date(2026, 1, 5) + ) + + assert history == [] + assert not index_called # no API call wasted + + +@pytest.mark.asyncio +async def test_get_price_history_skips_known_etf(monkeypatch): + adapter = _enable_adapter(monkeypatch) + index_called = False + + async def fake_lookup_stock(symbol: str): + return None + + async def fake_lookup_etf(symbol: str): + return {"symbol": "NEWGOLD", "name": "NewGold ETF"} + + async def fake_get_index_history(code, from_date=None, to_date=None): + nonlocal index_called + index_called = True + return None + + monkeypatch.setattr(adapter, "_lookup_stock", fake_lookup_stock) + monkeypatch.setattr(adapter, "_lookup_etf", fake_lookup_etf) + monkeypatch.setattr(adapter, "get_index_history", fake_get_index_history) + + history = await adapter.get_price_history( + "NEWGOLD", date(2026, 1, 1), date(2026, 1, 5) + ) + + assert history == [] + assert not index_called + + +@pytest.mark.asyncio +async def test_get_price_history_disabled(monkeypatch): + monkeypatch.setattr(ngx_module.settings, "NGXPULSE_API_KEY", "", raising=False) + + adapter = ngx_module.NgxPulseAdapter() + history = await adapter.get_price_history("ASI", date(2026, 1, 1), date(2026, 1, 5)) + assert history == [] + + +# --------------------------------------------------------------------------- +# get_current_price +# --------------------------------------------------------------------------- + + +@pytest.mark.asyncio +async def test_get_current_price_from_stocks(monkeypatch): + adapter = _enable_adapter(monkeypatch) + + async def fake_lookup_stock(symbol: str): + return {"symbol": symbol, "current_price": 665.00} + + async def fake_lookup_etf(symbol: str): + return None + + monkeypatch.setattr(adapter, "_lookup_stock", fake_lookup_stock) + monkeypatch.setattr(adapter, "_lookup_etf", fake_lookup_etf) + + price_date, price_cents = await adapter.get_current_price("DANGCEM") + assert price_cents == 66500 + + +@pytest.mark.asyncio +async def test_get_current_price_from_etfs(monkeypatch): + adapter = _enable_adapter(monkeypatch) + + async def fake_lookup_stock(symbol: str): + return None + + async def fake_lookup_etf(symbol: str): + return {"symbol": symbol, "close": 58000, "previous_close": 57500} + + monkeypatch.setattr(adapter, "_lookup_stock", fake_lookup_stock) + monkeypatch.setattr(adapter, "_lookup_etf", fake_lookup_etf) + + price_date, price_cents = await adapter.get_current_price("NEWGOLD") + assert price_cents == 5800000 # 58000.00 * 100 + + +@pytest.mark.asyncio +async def test_get_current_price_etf_fallback_to_previous_close(monkeypatch): + """When 'close' is missing, use 'previous_close'.""" + adapter = _enable_adapter(monkeypatch) + + async def fake_lookup_stock(symbol: str): + return None + + async def fake_lookup_etf(symbol: str): + return {"symbol": symbol, "previous_close": 57500} + + monkeypatch.setattr(adapter, "_lookup_stock", fake_lookup_stock) + monkeypatch.setattr(adapter, "_lookup_etf", fake_lookup_etf) + + price_date, price_cents = await adapter.get_current_price("NEWGOLD") + assert price_cents == 5750000 + + +@pytest.mark.asyncio +async def test_get_current_price_fallback_to_index(monkeypatch): + adapter = _enable_adapter(monkeypatch) + _reset_rate_limits(adapter) + + async def fake_lookup_stock(symbol: str): + return None + + async def fake_lookup_etf(symbol: str): + return None + + async def fake_get(path, params=None): + return { + "success": True, + "history": [{"date": "2026-06-01", "value": 72000.00}], + } + + monkeypatch.setattr(adapter, "_lookup_stock", fake_lookup_stock) + monkeypatch.setattr(adapter, "_lookup_etf", fake_lookup_etf) + monkeypatch.setattr(adapter, "_get", fake_get) + + price_date, price_cents = await adapter.get_current_price("ASI") + assert price_cents == 7200000 + + +@pytest.mark.asyncio +async def test_get_current_price_not_found(monkeypatch): + adapter = _enable_adapter(monkeypatch) + _reset_rate_limits(adapter) + + async def fake_lookup_stock(symbol: str): + return None + + async def fake_lookup_etf(symbol: str): + return None + + async def fake_get(path, params=None): + return {"success": False} + + monkeypatch.setattr(adapter, "_lookup_stock", fake_lookup_stock) + monkeypatch.setattr(adapter, "_lookup_etf", fake_lookup_etf) + monkeypatch.setattr(adapter, "_get", fake_get) + + price_date, price_cents = await adapter.get_current_price("UNKNOWN") + assert price_cents == 0 + + +@pytest.mark.asyncio +async def test_get_current_price_disabled(monkeypatch): + monkeypatch.setattr(ngx_module.settings, "NGXPULSE_API_KEY", "", raising=False) + + adapter = ngx_module.NgxPulseAdapter() + price_date, price_cents = await adapter.get_current_price("DANGCEM") + assert price_cents == 0 + + +# --------------------------------------------------------------------------- +# get_stocks +# --------------------------------------------------------------------------- + + +@pytest.mark.asyncio +async def test_get_stocks_bare_list_response(monkeypatch): + adapter = _enable_adapter(monkeypatch) + _reset_rate_limits(adapter) + + async def fake_get(path, params=None): + return [ + {"symbol": "DANGCEM", "name": "Dangote Cement"}, + {"symbol": "MTNN", "name": "MTN Nigeria"}, + ] + + monkeypatch.setattr(adapter, "_get", fake_get) + + stocks = await adapter.get_stocks() + assert stocks is not None + assert len(stocks) == 2 + assert stocks[0]["symbol"] == "DANGCEM" + + +@pytest.mark.asyncio +async def test_get_stocks_wrapped_dict_response(monkeypatch): + adapter = _enable_adapter(monkeypatch) + _reset_rate_limits(adapter) + + async def fake_get(path, params=None): + return { + "success": True, + "data": [{"symbol": "DANGCEM"}, {"symbol": "MTNN"}], + } + + monkeypatch.setattr(adapter, "_get", fake_get) + + stocks = await adapter.get_stocks() + assert stocks is not None + assert len(stocks) == 2 + + +@pytest.mark.asyncio +async def test_get_stocks_scans_dict_values_for_list(monkeypatch): + """When the dict doesn't use key 'data', scan all values for a list.""" + adapter = _enable_adapter(monkeypatch) + _reset_rate_limits(adapter) + + async def fake_get(path, params=None): + return {"error": False, "results": [{"symbol": "DANGCEM"}]} + + monkeypatch.setattr(adapter, "_get", fake_get) + + stocks = await adapter.get_stocks() + assert stocks is not None + assert len(stocks) == 1 + + +@pytest.mark.asyncio +async def test_get_stocks_cached(monkeypatch): + adapter = _enable_adapter(monkeypatch) + _reset_rate_limits(adapter) + call_count = 0 + + async def fake_get(path, params=None): + nonlocal call_count + call_count += 1 + return [{"symbol": f"STOCK{call_count}"}] + + monkeypatch.setattr(adapter, "_get", fake_get) + + stocks1 = await adapter.get_stocks() + stocks2 = await adapter.get_stocks() + + assert call_count == 1 # second call served from cache + assert stocks1 is stocks2 + assert stocks1[0]["symbol"] == "STOCK1" + + +@pytest.mark.asyncio +async def test_get_stocks_rate_limited(monkeypatch): + adapter = _enable_adapter(monkeypatch) + adapter._minute_count = adapter._MAX_REQUESTS_PER_MINUTE # exhausted + + stocks = await adapter.get_stocks() + assert stocks is None + + +# --------------------------------------------------------------------------- +# get_etfs +# --------------------------------------------------------------------------- + + +@pytest.mark.asyncio +async def test_get_etfs_success(monkeypatch): + adapter = _enable_adapter(monkeypatch) + _reset_rate_limits(adapter) + + async def fake_get(path, params=None): + return { + "success": True, + "count": 1, + "data": [{"symbol": "NEWGOLD", "name": "NewGold ETF"}], + } + + monkeypatch.setattr(adapter, "_get", fake_get) + + etfs = await adapter.get_etfs() + assert etfs is not None + assert etfs["count"] == 1 + assert etfs["data"][0]["symbol"] == "NEWGOLD" + + +@pytest.mark.asyncio +async def test_get_etfs_cached(monkeypatch): + adapter = _enable_adapter(monkeypatch) + _reset_rate_limits(adapter) + call_count = 0 + + async def fake_get(path, params=None): + nonlocal call_count + call_count += 1 + return {"success": True, "count": 1, "data": []} + + monkeypatch.setattr(adapter, "_get", fake_get) + + etfs1 = await adapter.get_etfs() + etfs2 = await adapter.get_etfs() + + assert call_count == 1 + assert etfs1 is etfs2 + + +@pytest.mark.asyncio +async def test_get_etfs_rate_limited(monkeypatch): + adapter = _enable_adapter(monkeypatch) + adapter._minute_count = adapter._MAX_REQUESTS_PER_MINUTE + + etfs = await adapter.get_etfs() + assert etfs is None + + +# --------------------------------------------------------------------------- +# get_index_history +# --------------------------------------------------------------------------- + + +@pytest.mark.asyncio +async def test_get_index_history_success(monkeypatch): + adapter = _enable_adapter(monkeypatch) + _reset_rate_limits(adapter) + + async def fake_get(path, params=None): + assert "asi" in path + return { + "success": True, + "code": "ASI", + "count": 2, + "history": [ + {"date": "2020-01-02", "value": 26867.79}, + {"date": "2020-01-03", "value": 27000.00}, + ], + } + + monkeypatch.setattr(adapter, "_get", fake_get) + + result = await adapter.get_index_history("ASI") + assert result is not None + assert result["count"] == 2 + assert result["history"][0]["value"] == 26867.79 + + +@pytest.mark.asyncio +async def test_get_index_history_with_date_params(monkeypatch): + adapter = _enable_adapter(monkeypatch) + _reset_rate_limits(adapter) + + captured_params = {} + + async def fake_get(path, params=None): + nonlocal captured_params + captured_params = params or {} + return {"success": True, "count": 1, "history": []} + + monkeypatch.setattr(adapter, "_get", fake_get) + + await adapter.get_index_history( + "ASI", from_date=date(2026, 1, 1), to_date=date(2026, 6, 1) + ) + + assert captured_params["from"] == "2026-01-01" + assert captured_params["to"] == "2026-06-01" + + +@pytest.mark.asyncio +async def test_get_index_history_rate_limited(monkeypatch): + adapter = _enable_adapter(monkeypatch) + adapter._day_count = adapter._MAX_REQUESTS_PER_DAY + + result = await adapter.get_index_history("ASI") + assert result is None + + +# --------------------------------------------------------------------------- +# Rate limiting — _can_make_request +# --------------------------------------------------------------------------- + + +def test_can_make_request_allows_under_limit(): + adapter = ngx_module.NgxPulseAdapter() + _reset_rate_limits(adapter) + + for _ in range(adapter._MAX_REQUESTS_PER_MINUTE): + assert adapter._can_make_request() + + # Now at limit + assert not adapter._can_make_request() + + +def test_can_make_request_day_limit(): + adapter = ngx_module.NgxPulseAdapter() + _reset_rate_limits(adapter) + adapter._day_count = adapter._MAX_REQUESTS_PER_DAY + + assert not adapter._can_make_request() + + +def test_can_make_request_minute_window_resets(): + adapter = ngx_module.NgxPulseAdapter() + adapter._minute_count = adapter._MAX_REQUESTS_PER_MINUTE + adapter._minute_window_start = datetime.now(timezone.utc).replace( + second=0, microsecond=0 + ) - timedelta(minutes=2) + + # Window is in the past, so it should reset + assert adapter._can_make_request() + assert adapter._minute_count == 1 + + +def test_can_make_request_day_window_resets(): + adapter = ngx_module.NgxPulseAdapter() + adapter._day_count = adapter._MAX_REQUESTS_PER_DAY + adapter._day_window_start = date.today() - timedelta(days=1) + + assert adapter._can_make_request() + assert adapter._day_count == 1 + + +# --------------------------------------------------------------------------- +# _extract_symbol +# --------------------------------------------------------------------------- + + +def test_extract_symbol_ngx_prefix(): + assert ngx_module.NgxPulseAdapter._extract_symbol("NGX:DANGCEM") == "DANGCEM" + assert ( + ngx_module.NgxPulseAdapter._extract_symbol("NGX/STANBICETF30") == "STANBICETF30" + ) + assert ngx_module.NgxPulseAdapter._extract_symbol("NGX|ASI") == "ASI" + + +def test_extract_symbol_ngxpulse_prefix(): + assert ngx_module.NgxPulseAdapter._extract_symbol("NGXPULSE:DANGCEM") == "DANGCEM" + assert ngx_module.NgxPulseAdapter._extract_symbol("NGXPULSE/NEWGOLD") == "NEWGOLD" + + +def test_extract_symbol_bare_ticker(): + assert ngx_module.NgxPulseAdapter._extract_symbol("DANGCEM") == "DANGCEM" + assert ngx_module.NgxPulseAdapter._extract_symbol("asi") == "ASI" + + +def test_extract_symbol_empty(): + assert ngx_module.NgxPulseAdapter._extract_symbol("") == "" + assert ngx_module.NgxPulseAdapter._extract_symbol(" ") == "" + + +# --------------------------------------------------------------------------- +# _safe_currency +# --------------------------------------------------------------------------- + + +def test_safe_currency_valid(): + assert ngx_module.NgxPulseAdapter._safe_currency("NGN") == Currency.NGN + assert ngx_module.NgxPulseAdapter._safe_currency("usd") == Currency.USD + + +def test_safe_currency_invalid_falls_back_to_ngn(): + assert ngx_module.NgxPulseAdapter._safe_currency("XYZ") == Currency.NGN + + +# --------------------------------------------------------------------------- +# _log_api_error +# --------------------------------------------------------------------------- + + +def test_log_api_error_json_body(caplog): + import logging + + caplog.set_level(logging.WARNING) + + body = '{"error": true, "status": 401, "message": "Invalid or missing API key"}' + ngx_module.NgxPulseAdapter._log_api_error("/test", 401, body) + + assert "NgxPulse error" in caplog.text + assert "401" in caplog.text + assert "Invalid or missing API key" in caplog.text + + +def test_log_api_error_non_json_body(caplog): + import logging + + caplog.set_level(logging.WARNING) + + body = "Server Error" + ngx_module.NgxPulseAdapter._log_api_error("/test", 500, body) + + assert "NgxPulse error" in caplog.text + assert "500" in caplog.text + assert "Server Error" in caplog.text + + +# --------------------------------------------------------------------------- +# Shared session +# --------------------------------------------------------------------------- + + +@pytest.mark.asyncio +async def test_shared_session_reused(monkeypatch): + monkeypatch.setattr( + ngx_module.settings, + "NGXPULSE_API_BASE_URL", + "https://www.ngxpulse.ng", + raising=False, + ) + monkeypatch.setattr( + ngx_module.settings, + "NGXPULSE_API_KEY", + "test-ngxpulse-key", + raising=False, + ) + + adapter_a = ngx_module.NgxPulseAdapter() + adapter_b = ngx_module.NgxPulseAdapter() + + session_a = await adapter_a._get_shared_session() + session_b = await adapter_b._get_shared_session() + + assert session_a is session_b + assert isinstance(session_a, aiohttp.ClientSession) + + await ngx_module.NgxPulseAdapter.close_shared_session() + + +@pytest.mark.asyncio +async def test_shared_session_close_and_recreate(monkeypatch): + monkeypatch.setattr( + ngx_module.settings, + "NGXPULSE_API_BASE_URL", + "https://www.ngxpulse.ng", + raising=False, + ) + monkeypatch.setattr( + ngx_module.settings, + "NGXPULSE_API_KEY", + "test-ngxpulse-key", + raising=False, + ) + + adapter = ngx_module.NgxPulseAdapter() + session_a = await adapter._get_shared_session() + + await ngx_module.NgxPulseAdapter.close_shared_session() + + session_b = await adapter._get_shared_session() + assert session_b is not session_a + + await ngx_module.NgxPulseAdapter.close_shared_session() + + +# --------------------------------------------------------------------------- +# _lookup_stock / _lookup_etf +# --------------------------------------------------------------------------- + + +@pytest.mark.asyncio +async def test_lookup_stock_finds_by_symbol_case_insensitive(monkeypatch): + adapter = _enable_adapter(monkeypatch) + _reset_rate_limits(adapter) + + async def fake_get(path, params=None): + return [{"symbol": "DANGCEM", "name": "Dangote Cement"}] + + monkeypatch.setattr(adapter, "_get", fake_get) + + result = await adapter._lookup_stock("dangcem") + assert result is not None + assert result["symbol"] == "DANGCEM" + + +@pytest.mark.asyncio +async def test_lookup_stock_not_found(monkeypatch): + adapter = _enable_adapter(monkeypatch) + _reset_rate_limits(adapter) + + async def fake_get(path, params=None): + return [{"symbol": "MTNN", "name": "MTN Nigeria"}] + + monkeypatch.setattr(adapter, "_get", fake_get) + + result = await adapter._lookup_stock("DANGCEM") + assert result is None + + +@pytest.mark.asyncio +async def test_lookup_etf_finds_by_symbol_or_canonical(monkeypatch): + adapter = _enable_adapter(monkeypatch) + _reset_rate_limits(adapter) + + async def fake_get(path, params=None): + return { + "success": True, + "data": [ + { + "symbol": "NEWGOLD", + "canonical_symbol": "NGNEWGOLD", + "name": "NewGold ETF", + } + ], + } + + monkeypatch.setattr(adapter, "_get", fake_get) + + # Match by symbol + assert await adapter._lookup_etf("NEWGOLD") is not None + # Match by canonical_symbol + assert await adapter._lookup_etf("NGNEWGOLD") is not None + # No match + assert await adapter._lookup_etf("UNKNOWN") is None diff --git a/api/uv.lock b/api/uv.lock index a35896ce..4f8ac0b5 100644 --- a/api/uv.lock +++ b/api/uv.lock @@ -618,7 +618,7 @@ wheels = [ [[package]] name = "folio-api" -version = "1.22.3" +version = "1.23.0" source = { virtual = "." } dependencies = [ { name = "aiohttp" }, diff --git a/web/package-lock.json b/web/package-lock.json index c53f55c0..cd25c433 100644 --- a/web/package-lock.json +++ b/web/package-lock.json @@ -1,12 +1,12 @@ { "name": "folio-web", - "version": "1.22.3", + "version": "1.23.0", "lockfileVersion": 3, "requires": true, "packages": { "": { "name": "folio-web", - "version": "1.22.3", + "version": "1.23.0", "dependencies": { "@sveltejs/adapter-node": "^5.5.4", "axios": "^1.15.2", diff --git a/web/package.json b/web/package.json index f204c9a5..75afc6aa 100644 --- a/web/package.json +++ b/web/package.json @@ -1,6 +1,6 @@ { "name": "folio-web", - "version": "1.22.3", + "version": "1.23.0", "type": "module", "scripts": { "dev": "svelte-kit sync && vite", diff --git a/web/src/lib/api/controllers/asset.controller.ts b/web/src/lib/api/controllers/asset.controller.ts index 2ecb428e..a117036f 100644 --- a/web/src/lib/api/controllers/asset.controller.ts +++ b/web/src/lib/api/controllers/asset.controller.ts @@ -13,7 +13,7 @@ import type { } from "../types"; export class AssetController { - constructor(private client: AxiosInstance) {} + constructor(private client: AxiosInstance) { } async searchAssets(query: string): Promise { const response = await this.client.get("/assets/search", { @@ -30,6 +30,7 @@ export class AssetController { | "yfinance" | "tiingo" | "ngnmarket" + | "ngxpulse" | "tradingview" = "yfinance", currency: string = "USD", assetClass: string = "", diff --git a/web/src/lib/api/types/requests.ts b/web/src/lib/api/types/requests.ts index fb60ee9f..0d9531ea 100644 --- a/web/src/lib/api/types/requests.ts +++ b/web/src/lib/api/types/requests.ts @@ -15,6 +15,7 @@ export type MarketDataProvider = | "yfinance" | "tiingo" | "ngnmarket" + | "ngxpulse" | "tradingview"; export interface CreateTradeRequest { diff --git a/web/src/lib/components/TradeForm.svelte b/web/src/lib/components/TradeForm.svelte index 2d40d5df..f46886fc 100644 --- a/web/src/lib/components/TradeForm.svelte +++ b/web/src/lib/components/TradeForm.svelte @@ -23,6 +23,7 @@ | "yfinance" | "tiingo" | "ngnmarket" + | "ngxpulse" | "tradingview"; } = { ticker: "", @@ -48,6 +49,7 @@ { label: "Yahoo Finance", value: "yfinance" }, { label: "Tiingo", value: "tiingo" }, { label: "NGNMarket", value: "ngnmarket" }, + { label: "NGXPulse", value: "ngxpulse" }, { label: "TradingView", value: "tradingview" }, ]; diff --git a/web/src/routes/+page.svelte b/web/src/routes/+page.svelte index aa719ed4..053f9c33 100644 --- a/web/src/routes/+page.svelte +++ b/web/src/routes/+page.svelte @@ -30,11 +30,11 @@ name: string; value: number; percent: number; - currentPrice: number; - dayChange: number; - weekAvg: number; - monthAvg: number; - monthChange: number; + currentPrice: number | null; + dayChange: number | null; + weekAvg: number | null; + monthAvg: number | null; + monthChange: number | null; } let stats: PortfolioStats = { id: "", @@ -55,8 +55,18 @@ })) ?? []; const monthOrder = [ - "Jan", "Feb", "Mar", "Apr", "May", "Jun", - "Jul", "Aug", "Sep", "Oct", "Nov", "Dec", + "Jan", + "Feb", + "Mar", + "Apr", + "May", + "Jun", + "Jul", + "Aug", + "Sep", + "Oct", + "Nov", + "Dec", ]; const getTimeFromLabel = (label: string) => { @@ -92,17 +102,23 @@ async function loadDashboardData(force = false) { try { loading = true; - let analyticsList = force ? null : await getCachedListAnalytics("1y", "USD"); + let analyticsList = force + ? null + : await getCachedListAnalytics("1y", "USD"); if (!analyticsList) { try { - analyticsList = await portfolioController.listPortfolioAnalytics({ - timeframe: "1y", - in_currency: "USD", - }); + analyticsList = + await portfolioController.listPortfolioAnalytics({ + timeframe: "1y", + in_currency: "USD", + }); await setCachedListAnalytics("1y", "USD", analyticsList); } catch (apiError) { if (force) { - analyticsList = await getCachedListAnalytics("1y", "USD"); + analyticsList = await getCachedListAnalytics( + "1y", + "USD", + ); } if (!analyticsList) throw apiError; } @@ -128,8 +144,14 @@ current_value: 0, total_gain_loss: 0, allocation: [] as { label: string; value: number }[], - performance_history: [] as { name: string; value: number }[], - contribution_history: [] as { name: string; value: number }[], + performance_history: [] as { + name: string; + value: number; + }[], + contribution_history: [] as { + name: string; + value: number; + }[], top_holdings: [] as { ticker: string; name?: string; @@ -143,7 +165,8 @@ for (const item of analyticsData.allocation) { allocationMap.set( item.label, - (allocationMap.get(item.label) || 0) + Number(item.value || 0), + (allocationMap.get(item.label) || 0) + + Number(item.value || 0), ); } const allocation = Array.from(allocationMap.entries()).map( @@ -154,23 +177,31 @@ for (const point of analyticsData.performance_history) { performanceMap.set( point.name, - (performanceMap.get(point.name) || 0) + Number(point.value || 0), + (performanceMap.get(point.name) || 0) + + Number(point.value || 0), ); } const performance_history = Array.from(performanceMap.entries()) .map(([name, value]) => ({ name, value })) - .sort((a, b) => getTimeFromLabel(a.name) - getTimeFromLabel(b.name)); + .sort( + (a, b) => + getTimeFromLabel(a.name) - getTimeFromLabel(b.name), + ); const contributionMap = new Map(); for (const point of analyticsData.contribution_history) { contributionMap.set( point.name, - (contributionMap.get(point.name) || 0) + Number(point.value || 0), + (contributionMap.get(point.name) || 0) + + Number(point.value || 0), ); } const contribution_history = Array.from(contributionMap.entries()) .map(([name, value]) => ({ name, value })) - .sort((a, b) => getTimeFromLabel(a.name) - getTimeFromLabel(b.name)); + .sort( + (a, b) => + getTimeFromLabel(a.name) - getTimeFromLabel(b.name), + ); const topHoldingsByWeightMap = new Map< string, @@ -184,7 +215,9 @@ }); } const totalValue = analyticsData.current_value || 0; - const topHoldingsByWeight = Array.from(topHoldingsByWeightMap.entries()) + const topHoldingsByWeight = Array.from( + topHoldingsByWeightMap.entries(), + ) .map(([ticker, { value, name }]) => ({ ticker, name, @@ -203,7 +236,8 @@ current_value: currentValue, cost_basis: costBasis, gain_loss: gainLoss, - return_percent: costBasis > 0 ? (gainLoss / costBasis) * 100 : 0, + return_percent: + costBasis > 0 ? (gainLoss / costBasis) * 100 : 0, allocation, performance_history, contribution_history, @@ -219,7 +253,12 @@ } async function loadPerformanceData( - holdings: { ticker: string; name: string; value: number; percent: number }[], + holdings: { + ticker: string; + name: string; + value: number; + percent: number; + }[], force = false, ) { if (!holdings.length) return; @@ -231,11 +270,19 @@ const startStr = start.toISOString().slice(0, 10); const endStr = end.toISOString().slice(0, 10); - const holdingByTicker = Object.fromEntries(holdings.map((h) => [h.ticker, h])); + const holdingByTicker = Object.fromEntries( + holdings.map((h) => [h.ticker, h]), + ); const tickers = holdings.map((h) => h.ticker); - const batchKey = buildBatchPriceHistoryKey(tickers, startStr, endStr); - let batch = force ? null : await getCachedBatchPriceHistory(batchKey); + const batchKey = buildBatchPriceHistoryKey( + tickers, + startStr, + endStr, + ); + let batch = force + ? null + : await getCachedBatchPriceHistory(batchKey); if (!batch) { try { batch = await assetController.getBatchPriceHistory({ @@ -252,56 +299,62 @@ } } - const results: PerformanceHolding[] = batch.results.map((history) => { - const h = holdingByTicker[history.ticker]; - const prices = (history.data || []) - .map((d) => Number(d.close)) - .filter((p) => p > 0 && Number.isFinite(p)); + const results: PerformanceHolding[] = batch.results.map( + (history) => { + const h = holdingByTicker[history.ticker]; + const prices = (history.data || []) + .map((d) => Number(d.close)) + .filter((p) => p > 0 && Number.isFinite(p)); + + if (prices.length < 2) { + const fallback = prices[prices.length - 1] || null; + return { + ticker: history.ticker, + name: h?.name ?? history.ticker, + value: h?.value ?? 0, + percent: h?.percent ?? 0, + currentPrice: fallback, + dayChange: null, + weekAvg: fallback, + monthAvg: fallback, + monthChange: null, + }; + } + + const currentPrice = prices[prices.length - 1]; + const prevPrice = prices[prices.length - 2]; + const dayChange = + prevPrice > 0 + ? ((currentPrice - prevPrice) / prevPrice) * 100 + : 0; + + const weekPrices = prices.slice(-7); + const weekAvg = + weekPrices.reduce((a, b) => a + b, 0) / + weekPrices.length; + + const monthAvg = + prices.reduce((a, b) => a + b, 0) / prices.length; + + const oldestPrice = prices[0]; + const monthChange = + oldestPrice > 0 + ? ((currentPrice - oldestPrice) / oldestPrice) * 100 + : 0; - if (prices.length < 2) { - const fallback = prices[prices.length - 1] || 0; return { ticker: history.ticker, name: h?.name ?? history.ticker, value: h?.value ?? 0, percent: h?.percent ?? 0, - currentPrice: fallback, - dayChange: 0, - weekAvg: fallback, - monthAvg: fallback, - monthChange: 0, + currentPrice, + dayChange, + weekAvg, + monthAvg, + monthChange, }; - } - - const currentPrice = prices[prices.length - 1]; - const prevPrice = prices[prices.length - 2]; - const dayChange = - prevPrice > 0 ? ((currentPrice - prevPrice) / prevPrice) * 100 : 0; - - const weekPrices = prices.slice(-7); - const weekAvg = - weekPrices.reduce((a, b) => a + b, 0) / weekPrices.length; - - const monthAvg = prices.reduce((a, b) => a + b, 0) / prices.length; - - const oldestPrice = prices[0]; - const monthChange = - oldestPrice > 0 - ? ((currentPrice - oldestPrice) / oldestPrice) * 100 - : 0; - - return { - ticker: history.ticker, - name: h?.name ?? history.ticker, - value: h?.value ?? 0, - percent: h?.percent ?? 0, - currentPrice, - dayChange, - weekAvg, - monthAvg, - monthChange, - }; - }); + }, + ); performanceHoldings = results; } finally { @@ -326,8 +379,14 @@
-

{today}

-

+

+ {today} +

+

Dashboard

@@ -368,57 +427,109 @@ {:else} -
+
-
-

Total Value

-
- +
+

+ Total Value +

+
+
-

Current portfolio value

+

+ Current portfolio value +

-
-

Cost Basis

-
- +
+

+ Cost Basis +

+
+
-

Total amount invested

+

+ Total amount invested +

-
-

Gain / Loss

-
+

+ Gain / Loss +

+
- + class:text-negative={!isPositive} + > +
-

Unrealized P&L

+

+ Unrealized P&L +

-
-

Total Return

-
+

+ Total Return +

+
+ class:text-negative={!isPositive} + >
-

Return on investment

+

+ Return on investment +

- +
- + -
+
- +
- + {#if performanceLoading}
{#each Array(10) as _} @@ -471,69 +594,200 @@ - - - + - - - {#each performanceHoldings as h} - + - - - - - + {/each} @@ -541,7 +795,9 @@
+ Ticker + Price + Value + Day 7D Avg 30D Avg + 30D Chg
-
- +
+
- + {h.ticker} {#if h.name && h.name !== h.ticker} - + {h.name} {/if} - - of portfolio + + of portfolio
- + + {#if h.currentPrice != null} + + {:else} + + {/if} + {#if h.value > 0} + + {:else} + + {/if} + {#if h.dayChange != null} + = + 0} + class:text-negative={h.dayChange < + 0} + > + + + {:else} + + {/if} + {#if h.weekAvg != null} + + {:else} + + {/if} - = 0} - class:text-negative={h.monthChange < 0}> - - + + {#if h.monthAvg != null} + + {:else} + + {/if} + + {#if h.monthChange != null} + = + 0} + class:text-negative={h.monthChange < + 0} + > + + + {:else} + + {/if}
{:else} -

+

No holdings data available.

{/if} @@ -559,8 +815,18 @@ class="h-10 w-10 rounded-full bg-card border border-border text-foreground flex items-center justify-center hover:bg-muted transition-all shadow-lg" title="New Trade" > - - + + - - - + + + - - + + {/if}
diff --git a/web/src/routes/trades/import/+page.svelte b/web/src/routes/trades/import/+page.svelte index b07b1399..d6c4fcb6 100644 --- a/web/src/routes/trades/import/+page.svelte +++ b/web/src/routes/trades/import/+page.svelte @@ -47,7 +47,7 @@ key: "market_data_provider", label: "Market Data Provider", required: false, - hint: "yfinance, tiingo, ngnmarket, tradingview", + hint: "yfinance, tiingo, ngnmarket, ngxpulse, tradingview", }, { key: "fees", label: "Fees", required: false, hint: "" }, { key: "notes", label: "Notes", required: false, hint: "" },