main
py 109 lines 5.81 KB
Raw
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