| 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' |