main
py 357 lines 17.1 KB
Raw
1 from datetime import date, datetime, timedelta, timezone
2 from decimal import Decimal
3 from types import SimpleNamespace
4 from unittest.mock import AsyncMock
5 from uuid import UUID
6
7 import httpx
8 import pytest
9
10 import app.nse_daily_backfill as backfill
11 from app.market_data_population import IndiaMarketDataPopulationJobs
12 from app.models import DailyMarketBar
13 from app.nse_historical_daily import NseHistoricalDailyProvider, NseHistoricalResult, BOOTSTRAP
14 from app.persistence import SqliteResearchPersistence
15 from app.repository import ResearchRepository
16 from app.settings import Settings
17
18 NOW = datetime(2026, 9, 13, tzinfo=timezone.utc)
19 START, END = date(2026, 9, 1), date(2026, 9, 4)
20
21
22 def metadata(n=1, **changes):
23 return dict(globalInstrumentId=str(UUID(int=n)), status='ACTIVE', assetType='EQUITY', country='IN',
24 primaryExchange='NSE', currency='INR', providerMappings=[dict(provider='NSE', status='VERIFIED',
25 providerSymbol=f'S{n}', currency='INR', resolutionSource='OFFICIAL_NSE')]) | changes
26
27
28 def bar(n=1, day=START, **changes):
29 return DailyMarketBar(global_instrument_id=UUID(int=n), trading_date=day, open=Decimal('10'), high=Decimal('12'),
30 low=Decimal('9'), close=Decimal('11'), volume=0, currency='INR', provider='NSE', provider_symbol=f'S{n}',
31 source_mode='REAL', source_url='https://www.nseindia.com/history', retrieved_at=NOW).model_copy(update=changes)
32
33
34 class Provider:
35 def __init__(self):
36 self.calls = []
37 self.closed = False
38 self.failures = {}
39 self.price = Decimal('11')
40
41 async def fetch(self, key, *, start, end, **kwargs):
42 self.calls.append((key, start, end))
43 reason = self.failures.get((key.int, start), self.failures.get(key.int))
44 if isinstance(reason, Exception):
45 raise reason
46 values = [bar(key.int, day=start, close=self.price)]
47 if start != end:
48 values.append(bar(key.int, day=end, close=self.price))
49 return NseHistoricalResult(key, start, end, provider_symbol=f'S{key.int}', http_status=200,
50 status='UNAVAILABLE' if reason else 'SUCCESS', failure_reason=reason,
51 rows_parsed=0 if reason else len(values), bars=[] if reason else values)
52
53 async def aclose(self): self.closed = True
54
55
56 def setup(monkeypatch, values=None, **settings):
57 values = [metadata()] if values is None else values
58 universe = [v | {'exchange': v.get('primaryExchange')} for v in values]
59 by_id = {UUID(v['globalInstrumentId']): v for v in values}
60 orchestrator = SimpleNamespace(active_global_equities=AsyncMock(return_value=universe),
61 global_instrument_metadata=AsyncMock(side_effect=lambda key, **kw: by_id[key]))
62 store = SqliteResearchPersistence()
63 config = Settings(**settings)
64 repo = ResearchRepository(settings=config, persistence=store)
65 jobs = IndiaMarketDataPopulationJobs(repo, None, orchestrator, config, clock=lambda: NOW, sleep=AsyncMock())
66 providers = []
67 def factory(*args):
68 p = Provider(); providers.append(p); return p
69 monkeypatch.setattr(backfill, 'NseHistoricalDailyProvider', factory)
70 return jobs, store, providers
71
72
73 async def run(jobs, **kwargs):
74 return await jobs.backfill_daily_bars(identity_headers={}, start=START, end=END, **kwargs)
75
76
77 @pytest.mark.parametrize('start,end,size,count', [
78 (date(2026,1,1),date(2026,1,30),30,1), (date(2026,1,1),date(2026,1,31),30,2),
79 (date(2025,12,15),date(2026,3,5),30,3), (date(2024,2,1),date(2024,3,1),30,1),
80 (START,START,30,1), (date(2026,1,31),date(2026,2,1),1,2)])
81 def test_windows(start,end,size,count):
82 windows=backfill.plan_windows(start,end,size)
83 assert len(windows)==count and windows[0][0]==start and windows[-1][1]==end
84 assert windows==backfill.plan_windows(start,end,size)
85 assert all(0 <= (b-a).days < size for a,b in windows)
86 assert all(windows[i][1]+timedelta(days=1)==windows[i+1][0] for i in range(len(windows)-1))
87
88
89 @pytest.mark.parametrize('start,end,size',[(END,START,30),(START,END,0),(START,END,-1),(NOW,END,30)])
90 def test_invalid_windows(start,end,size):
91 with pytest.raises(ValueError): backfill.plan_windows(start,end,size)
92
93
94 def test_coverage_states_and_no_calendar_inference():
95 assert backfill.coverage_plan([],START,END,NOW,72)[0]=='NO_HISTORY'
96 assert backfill.coverage_plan([bar(day=END)],START,END,NOW,72)==('PARTIAL_HISTORY',[(START,END-timedelta(days=1))])
97 assert backfill.coverage_plan([bar(),bar(day=END)],START,END,NOW,72)==('CURRENT_HISTORY',[])
98 old=NOW-timedelta(days=5)
99 assert backfill.coverage_plan([bar(retrieved_at=old)],START,END,NOW,72)==('STALE_HISTORY',[(START+timedelta(days=1),END)])
100 # Sparse internal dates are not enough to claim a session gap.
101 assert backfill.coverage_plan([bar(),bar(day=END)],START,END,NOW,72)[0]!='GAP_DETECTED'
102 assert backfill.coverage_plan([bar(),bar(day=END)],START,END,NOW,72,force=True)[1]==[(START,END)]
103
104
105 @pytest.mark.asyncio
106 async def test_current_no_session_and_other_provider_does_not_count(monkeypatch):
107 jobs,store,providers=setup(monkeypatch)
108 store.upsert_daily_market_bars([bar(provider='OTHER'),bar(day=END,provider='OTHER')])
109 first=await run(jobs)
110 assert first['instruments'][0]['coverage_state_before']=='NO_HISTORY' and first['rows_persisted']==2
111 assert len(store.load_daily_market_bars({UUID(int=1)}))==4
112 second=await run(jobs)
113 assert second['skipped_current']==1 and second['requested_windows']==0 and len(providers)==1
114 assert providers[0].closed
115
116
117 @pytest.mark.asyncio
118 async def test_prefix_tail_and_null_rows(monkeypatch):
119 jobs,store,providers=setup(monkeypatch)
120 store.upsert_daily_market_bar(bar(day=date(2026,9,3),retrieved_at=NOW-timedelta(days=5)))
121 r=await run(jobs)
122 assert [(a,b) for _,a,b in providers[0].calls]==[(START,date(2026,9,2)),(END,END)]
123 assert r['instruments'][0]['coverage_state_before']=='PARTIAL_HISTORY'
124 assert backfill.usable_bars([None,bar(open=None)],UUID(int=1),'S1','INR',END)==[]
125
126
127 @pytest.mark.asyncio
128 async def test_batch_order_offset_and_membership(monkeypatch):
129 jobs,store,providers=setup(monkeypatch,[metadata(3),metadata(1),metadata(2)],market_data_population_batch_size=1)
130 first=await run(jobs)
131 assert first['requested_instruments']==1 and first['next_offset']==1
132 second=await run(jobs,offset=1)
133 assert second['instruments'][0]['globalInstrumentId']==str(UUID(int=2)) and second['next_offset']==2
134 third=await run(jobs,instrument_ids={UUID(int=3),UUID(int=999)})
135 assert third['requested_instruments']==1 and third['instruments'][0]['globalInstrumentId']==str(UUID(int=3))
136 assert third['next_offset'] is None
137
138
139 @pytest.mark.asyncio
140 @pytest.mark.parametrize('values',[[],[metadata(assetType='ETF')],[metadata(status='INACTIVE')],[metadata(country='US')],[metadata(primaryExchange='BSE')]])
141 async def test_empty_eligible_universe(monkeypatch,values):
142 jobs,_,providers=setup(monkeypatch,values)
143 r=await run(jobs)
144 assert r['requested_instruments']==0 and not providers and r['status']=='COMPLETED'
145
146
147 @pytest.mark.asyncio
148 @pytest.mark.parametrize('mapping',[[],[dict(provider='NSE',status='INVALID',providerSymbol='X')],
149 [dict(provider='NSE',status='VERIFIED',providerSymbol=' ')], metadata()['providerMappings']*2,
150 [metadata()['providerMappings'][0]|dict(active=False)], [metadata()['providerMappings'][0]|dict(currency='USD')]])
151 async def test_identity_failure_isolated(monkeypatch,mapping):
152 jobs,_,providers=setup(monkeypatch,[metadata(providerMappings=mapping),metadata(2)])
153 r=await run(jobs)
154 assert r['failed_instruments']==1 and r['succeeded_instruments']==1
155 assert r['instruments'][0]['failure_class']=='NO_TRUSTED_MAPPING'
156 assert [key.int for key,_,_ in providers[0].calls]==[2]
157
158
159 @pytest.mark.asyncio
160 @pytest.mark.parametrize('reason,category',[('HISTORICAL_HTTP_403','HTTP_ERROR'),('NON_CSV_RESPONSE','PARSER_FAILURE'),
161 ('EMPTY_RESPONSE','EMPTY_RESPONSE'),('HISTORICAL_TIMEOUT','PROVIDER_UNAVAILABLE'),(RuntimeError('secret-cookie'),'PROVIDER_UNAVAILABLE')])
162 async def test_provider_failure_isolation_and_cooldown(monkeypatch,reason,category):
163 jobs,store,_=setup(monkeypatch,[metadata(),metadata(2)])
164 provider=Provider(); provider.failures[1]=reason
165 monkeypatch.setattr(backfill,'NseHistoricalDailyProvider',lambda *args:provider)
166 r=await run(jobs)
167 assert r['failed_instruments']==1 and r['succeeded_instruments']==1 and r['failed_windows']==1
168 assert r['instruments'][0]['failure_class']==category and provider.closed
169 assert 'secret-cookie' not in str(r)
170 repeat=await run(jobs)
171 assert repeat['skipped_cooldown']==1 and repeat['skipped_current']==1
172
173
174 @pytest.mark.asyncio
175 async def test_throttle_stops_job_and_next_job(monkeypatch):
176 jobs,_,_=setup(monkeypatch,[metadata(),metadata(2)])
177 provider=Provider(); provider.failures[1]='HISTORICAL_HTTP_429'
178 monkeypatch.setattr(backfill,'NseHistoricalDailyProvider',lambda *args:provider)
179 r=await run(jobs)
180 assert r['failed_instruments']==1 and r['skipped_cooldown']==1 and len(provider.calls)==1
181 r=await run(jobs,force=True)
182 assert r['skipped_cooldown']==2 and len(provider.calls)==1
183
184
185 @pytest.mark.asyncio
186 async def test_per_window_persistence_later_failure_preserves_rows(monkeypatch):
187 jobs,store,_=setup(monkeypatch,[metadata(),metadata(2)],nse_historical_request_window_days=2)
188 provider=Provider(); provider.failures[(1,date(2026,9,3))]='HISTORICAL_HTTP_500'
189 monkeypatch.setattr(backfill,'NseHistoricalDailyProvider',lambda *args:provider)
190 r=await run(jobs)
191 assert r['requested_windows']==4 and r['successful_windows']==3 and r['failed_windows']==1
192 assert r['rows_received']==r['rows_accepted']==r['rows_persisted']==6
193 assert len(store.load_daily_market_bars({UUID(int=1)},provider='NSE'))==2
194 assert r['instruments'][0]['earliest_persisted_date']==START
195 assert r['first_requested_date']==START and r['last_requested_date']==END
196
197
198 @pytest.mark.asyncio
199 async def test_persistence_failure_next_instrument(monkeypatch):
200 jobs,store,_=setup(monkeypatch,[metadata(),metadata(2)])
201 original=jobs.repository.upsert_daily_market_bars_async
202 async def writer(bars):
203 if bars[0].global_instrument_id.int==1: raise RuntimeError('secret')
204 return await original(bars)
205 jobs.repository.upsert_daily_market_bars_async=writer
206 r=await run(jobs)
207 assert r['instruments'][0]['failure_class']=='PERSISTENCE_FAILURE' and r['succeeded_instruments']==1
208 assert 'secret' not in str(r)
209
210
211 @pytest.mark.asyncio
212 async def test_force_correction_idempotence(monkeypatch):
213 jobs,store,_=setup(monkeypatch)
214 provider=Provider(); monkeypatch.setattr(backfill,'NseHistoricalDailyProvider',lambda *args:provider)
215 await run(jobs)
216 provider.price=Decimal('10.5')
217 r=await run(jobs,force=True)
218 assert r['rows_persisted']==2
219 rows=store.load_daily_market_bars({UUID(int=1)},provider='NSE')
220 assert len(rows)==2 and all(b.close==Decimal('10.5') for b in rows)
221
222
223 @pytest.mark.asyncio
224 async def test_default_lookback_exchange_date_and_invalid_range(monkeypatch):
225 jobs,_,providers=setup(monkeypatch)
226 jobs._clock=lambda:datetime(2026,9,12,20,tzinfo=timezone.utc)
227 r=await jobs.backfill_daily_bars(identity_headers={})
228 assert r['target_end']==date(2026,9,13) and r['target_start']==date(2026,9,13)-timedelta(days=400)
229 assert r['requested_windows']==14 and providers[0].closed
230 bad=await jobs.backfill_daily_bars(identity_headers={},start=END,end=START)
231 assert bad['failure_reason']=='WINDOW_PLANNING_FAILURE'
232
233
234 @pytest.mark.asyncio
235 async def test_real_provider_one_session_verified_query_and_serialization(monkeypatch):
236 jobs,_,_=setup(monkeypatch,[metadata(),metadata(2)],nse_historical_request_window_days=2)
237 calls=[]
238 def handler(req):
239 calls.append(req)
240 if str(req.url)==BOOTSTRAP: return httpx.Response(200,headers={'set-cookie':'session=private; Path=/'})
241 assert req.headers['cookie']=='session=private'
242 symbol=req.url.params['symbol']; day=req.url.params['from']
243 return httpx.Response(200,text=f'Date,Symbol,Series,Open,High,Low,Close\n{day},{symbol},EQ,10,12,9,11\n')
244 async with httpx.AsyncClient(transport=httpx.MockTransport(handler)) as client:
245 sleep=AsyncMock()
246 provider=NseHistoricalDailyProvider(jobs.orchestrator,jobs.settings,client=client,sleep=sleep)
247 provider.aclose=AsyncMock()
248 monkeypatch.setattr(backfill,'NseHistoricalDailyProvider',lambda *args:provider)
249 r=await run(jobs)
250 assert len(calls)==5 and sum(str(c.url)==BOOTSTRAP for c in calls)==1
251 assert [c.url.params['symbol'] for c in calls[1:]]==['S1','S1','S2','S2']
252 assert sleep.await_count>=4 and r['successful_windows']==4
253 provider.aclose.assert_awaited_once()
254 assert 'private' not in str(r)
255
256 @pytest.mark.asyncio
257 async def test_successful_range_evidence_avoids_boundary_refetch(monkeypatch):
258 jobs,store,_=setup(monkeypatch)
259 p=Provider()
260 async def fetch(key, *, start, end, **kw):
261 p.calls.append((key,start,end))
262 return NseHistoricalResult(key,start,end,status='SUCCESS',bars=[bar(day=date(2026,9,2))],rows_parsed=1)
263 p.fetch=fetch
264 monkeypatch.setattr(backfill,'NseHistoricalDailyProvider',lambda *args:p)
265 await run(jobs)
266 r=await run(jobs)
267 assert r['skipped_current']==1 and len(p.calls)==1
268 jobs._clock=lambda: NOW+timedelta(hours=73)
269 r=await run(jobs)
270 assert r['requested_windows']>0 and len(p.calls)>1
271
272
273 @pytest.mark.asyncio
274 async def test_cooldown_expiry_and_identity_unavailable(monkeypatch):
275 jobs,_,_=setup(monkeypatch)
276 original=jobs.orchestrator.global_instrument_metadata.side_effect
277 jobs.orchestrator.global_instrument_metadata.side_effect=ValueError('Cookie: SECRET')
278 r=await run(jobs)
279 assert r['instruments'][0]['failure_class']=='IDENTITY_UNAVAILABLE'
280 assert 'SECRET' not in str(r)
281 jobs.orchestrator.global_instrument_metadata.side_effect=original
282 assert (await run(jobs))['skipped_cooldown']==1
283 jobs._clock=lambda: NOW+timedelta(hours=13)
284 assert (await run(jobs))['succeeded_instruments']==1
285
286
287 @pytest.mark.asyncio
288 async def test_session_closes_on_cancel(monkeypatch):
289 import asyncio
290 jobs,_,_=setup(monkeypatch)
291 p=Provider()
292 p.fetch=AsyncMock(side_effect=asyncio.CancelledError())
293 monkeypatch.setattr(backfill,'NseHistoricalDailyProvider',lambda *args:p)
294 with pytest.raises(asyncio.CancelledError): await run(jobs)
295 assert p.closed
296
297
298 @pytest.mark.asyncio
299 async def test_population_lock_serializes_backfill_jobs(monkeypatch):
300 import asyncio
301 jobs,_,providers=setup(monkeypatch)
302 a,b=await asyncio.gather(run(jobs),run(jobs))
303 assert a['succeeded_instruments']==1 and b['skipped_current']==1 and len(providers)==1
304
305
306 @pytest.mark.asyncio
307 async def test_window_planning_failure_and_universe_failure(monkeypatch):
308 jobs,_,providers=setup(monkeypatch)
309 r=await jobs.backfill_daily_bars(identity_headers={},offset=-1)
310 assert r['failure_reason']=='WINDOW_PLANNING_FAILURE' and not providers
311 jobs.orchestrator.active_global_equities.side_effect=RuntimeError('secret')
312 r=await run(jobs)
313 assert r['failure_reason']=='UNIVERSE_UNAVAILABLE' and 'secret' not in str(r)
314
315
316 @pytest.mark.asyncio
317 async def test_empty_window_does_not_block_later_history(monkeypatch):
318 jobs,store,_=setup(monkeypatch,nse_historical_request_window_days=2)
319 p=Provider(); p.failures[(1,START)]='NO_VALID_HISTORY'
320 monkeypatch.setattr(backfill,'NseHistoricalDailyProvider',lambda *args:p)
321 r=await run(jobs)
322 assert r['requested_windows']==2 and r['failed_windows']==1 and r['successful_windows']==1
323 assert r['failed_instruments']==1 and r['rows_persisted']==2
324 assert len(store.load_daily_market_bars({UUID(int=1)},provider='NSE'))==2
325
326
327 @pytest.mark.asyncio
328 async def test_partial_parser_rows_persist_but_not_marked_complete(monkeypatch):
329 jobs,store,_=setup(monkeypatch)
330 p=Provider()
331 async def fetch(key,*,start,end,**kw):
332 return NseHistoricalResult(key,start,end,status='SUCCESS',rows_parsed=2,rows_rejected=1,bars=[bar()])
333 p.fetch=fetch; monkeypatch.setattr(backfill,'NseHistoricalDailyProvider',lambda *args:p)
334 r=await run(jobs)
335 assert r['failed_windows']==1 and r['rows_persisted']==1 and r['rows_received']==2 and r['rows_accepted']==1
336 assert r['instruments'][0]['failure_class']=='PARSER_FAILURE'
337 assert jobs._daily_bar_completed_ranges[(UUID(int=1),'S1','INR')]==[]
338
339 @pytest.mark.asyncio
340 async def test_missing_persisted_rows_override_request_memory(monkeypatch):
341 jobs,store,providers=setup(monkeypatch)
342 await run(jobs)
343 jobs.repository.daily_market_bars_for_instruments=AsyncMock(return_value={UUID(int=1):None})
344 again=await run(jobs)
345 assert again['instruments'][0]['coverage_state_before']=='NO_HISTORY'
346 assert again['requested_windows']==1 and len(providers)==2
347
348
349 @pytest.mark.asyncio
350 async def test_rejected_rows_are_parser_failure_not_empty(monkeypatch):
351 jobs,_,_=setup(monkeypatch,nse_historical_request_window_days=2)
352 p=Provider()
353 p.fetch=AsyncMock(return_value=NseHistoricalResult(UUID(int=1),START,END,
354 failure_reason='NO_VALID_HISTORY', rows_parsed=1, rows_rejected=1))
355 monkeypatch.setattr(backfill,'NseHistoricalDailyProvider',lambda *args:p)
356 result=await run(jobs)
357 assert result['requested_windows']==1 and result['instruments'][0]['failure_class']=='PARSER_FAILURE'