| 1 | """Official index JSON acquisition; reuses NSE transport, not equity CSV parsing.""" |
| 2 | from datetime import date, datetime, timezone |
| 3 | from decimal import Decimal |
| 4 | import json |
| 5 | |
| 6 | import httpx |
| 7 | from app.models import DailyMarketBar |
| 8 | from app.nse_historical_daily import NseHistoricalDailyProvider, NseHistoricalResult, BOOTSTRAP, number, trading_date |
| 9 | from app.sector_benchmarks import BENCHMARKS, benchmark_id, benchmark_identity |
| 10 | |
| 11 | ENDPOINT = 'https://www.nseindia.com/api/historicalOR/indicesHistory' |
| 12 | HISTORY_VERSION = 'NSE_INDEX_HISTORY_V1' |
| 13 | # Exact response names observed in the checked-in official response fixtures. |
| 14 | # These are contract aliases, never canonical identities or fuzzy matches. |
| 15 | RESPONSE_NAMES = {'NIFTY 500': 'NIFTY 500', 'NIFTY IT': 'NIFTY IT', |
| 16 | 'NIFTY FINANCIAL SERVICES': 'NIFTY FIN SERVICE', |
| 17 | 'NIFTY HEALTHCARE INDEX': 'NIFTY HEALTHCARE'} |
| 18 | REQUIRED = {'EOD_INDEX_NAME', 'EOD_TIMESTAMP', 'EOD_OPEN_INDEX_VAL', 'EOD_HIGH_INDEX_VAL', |
| 19 | 'EOD_LOW_INDEX_VAL', 'EOD_CLOSE_INDEX_VAL'} |
| 20 | |
| 21 | |
| 22 | def parse_index_history(content, result, currency): |
| 23 | try: |
| 24 | payload = json.loads(content.decode('utf-8-sig'), parse_float=Decimal) |
| 25 | except (ValueError, UnicodeError): |
| 26 | raise ValueError('NON_JSON_INDEX_RESPONSE') from None |
| 27 | data = payload.get('data') if isinstance(payload, dict) else None |
| 28 | if not isinstance(data, list): |
| 29 | raise ValueError('INVALID_INDEX_RESPONSE_CONTRACT') |
| 30 | if not data: |
| 31 | raise ValueError('EMPTY_RESPONSE') |
| 32 | accepted = {} |
| 33 | for row in data: |
| 34 | result.rows_parsed += 1 |
| 35 | if not isinstance(row, dict) or not REQUIRED.issubset(row): |
| 36 | raise ValueError('MISSING_INDEX_FIELDS') |
| 37 | result.headers = sorted(row) |
| 38 | if ' '.join(str(row['EOD_INDEX_NAME']).upper().split()) != RESPONSE_NAMES.get(result.provider_symbol): |
| 39 | raise ValueError('INDEX_SYMBOL_MISMATCH') |
| 40 | try: |
| 41 | day = trading_date(row['EOD_TIMESTAMP']) |
| 42 | if not result.request_from <= day <= result.request_to or day in accepted: |
| 43 | raise ValueError('INVALID_OR_DUPLICATE_INDEX_DATE') |
| 44 | bar = DailyMarketBar(global_instrument_id=result.global_instrument_id, trading_date=day, |
| 45 | **{name:number(str(row[field]), required=True) for name,field in ( |
| 46 | ('open','EOD_OPEN_INDEX_VAL'),('high','EOD_HIGH_INDEX_VAL'), |
| 47 | ('low','EOD_LOW_INDEX_VAL'),('close','EOD_CLOSE_INDEX_VAL'))}, |
| 48 | currency=currency, provider='NSE', provider_symbol=result.provider_symbol, |
| 49 | source_mode='REAL', source_url=result.source_url, retrieved_at=result.retrieved_at, |
| 50 | volume=None, turnover=None, previous_close=None) |
| 51 | except (ValueError, TypeError, KeyError): |
| 52 | raise ValueError('INVALID_INDEX_ROW') from None |
| 53 | accepted[day] = bar |
| 54 | result.bars = [accepted[day] for day in sorted(accepted)] |
| 55 | |
| 56 | |
| 57 | class NseIndexHistoryProvider: |
| 58 | def __init__(self, orchestrator, settings, *, client=None, sleep=None): |
| 59 | kwargs = {'client':client} |
| 60 | if sleep is not None: |
| 61 | kwargs['sleep'] = sleep |
| 62 | self.transport = NseHistoricalDailyProvider(orchestrator, settings, **kwargs) |
| 63 | self.orchestrator, self.settings = orchestrator, settings |
| 64 | |
| 65 | async def aclose(self): |
| 66 | await self.transport.aclose() |
| 67 | |
| 68 | async def fetch(self, global_instrument_id, *, start, end, identity_headers=None, correlation_id=None): |
| 69 | result = NseHistoricalResult(global_instrument_id, start, end, source_url=ENDPOINT) |
| 70 | stage = 'IDENTITY' |
| 71 | try: |
| 72 | keys = [key for key in BENCHMARKS if benchmark_id(key) == global_instrument_id] |
| 73 | if len(keys) != 1: |
| 74 | raise ValueError('BENCHMARK_IDENTITY_UNAVAILABLE') |
| 75 | if type(start) is not date or type(end) is not date or start > end or (end-start).days+1 > self.settings.nse_historical_request_window_days: |
| 76 | raise ValueError('INVALID_REQUEST_WINDOW') |
| 77 | key = keys[0] |
| 78 | metadata = await self.orchestrator.global_instrument_metadata(global_instrument_id, |
| 79 | identity_headers=identity_headers, correlation_id=correlation_id) |
| 80 | ref = benchmark_identity(metadata, key) |
| 81 | result.provider_symbol = BENCHMARKS[key] |
| 82 | async with self.transport._lock: |
| 83 | if not self.transport._bootstrapped: |
| 84 | stage = 'BOOTSTRAP' |
| 85 | await self.transport._get(BOOTSTRAP, result) |
| 86 | self.transport._bootstrapped = True |
| 87 | stage = 'HISTORICAL' |
| 88 | response = await self.transport._get(ENDPOINT, result, params={'indexType':result.provider_symbol, |
| 89 | 'from':start.strftime('%d-%m-%Y'),'to':end.strftime('%d-%m-%Y')}) |
| 90 | result.retrieved_at = datetime.now(timezone.utc) |
| 91 | result.source_url = str(response.url) |
| 92 | parse_index_history(response.content, result, ref.currency) |
| 93 | result.status = 'SUCCESS' |
| 94 | except httpx.HTTPStatusError as exc: |
| 95 | result.failure_reason = f'{stage}_HTTP_{exc.response.status_code}' |
| 96 | if exc.response.status_code == 403: |
| 97 | self.transport._bootstrapped = False |
| 98 | except httpx.HTTPError: |
| 99 | result.failure_reason = f'{stage}_PROVIDER_UNAVAILABLE' |
| 100 | except ValueError as exc: |
| 101 | safe = {'BENCHMARK_IDENTITY_UNAVAILABLE','INVALID_REQUEST_WINDOW','NON_JSON_INDEX_RESPONSE', |
| 102 | 'INVALID_INDEX_RESPONSE_CONTRACT','EMPTY_RESPONSE','MISSING_INDEX_FIELDS','INDEX_SYMBOL_MISMATCH','INVALID_INDEX_ROW'} |
| 103 | result.failure_reason = str(exc) if str(exc) in safe else f'{stage}_UNAVAILABLE' |
| 104 | except Exception: |
| 105 | result.failure_reason = f'{stage}_UNAVAILABLE' |
| 106 | if result.failure_reason: |
| 107 | result.rows_rejected = result.rows_parsed |
| 108 | result.bars.clear() |
| 109 | return result |