main
py 500 lines 20.3 KB
Raw
1 from __future__ import annotations
2
3 import asyncio
4 from dataclasses import replace
5 from datetime import datetime, timedelta, timezone
6 from decimal import Decimal
7 from types import SimpleNamespace
8 from uuid import uuid4
9
10 import pytest
11 import httpx
12 from fastapi import HTTPException
13
14 import app.main as main
15 from app.market_data_ensure import IndiaMarketDataEnsureService
16 from app.market_universe import MarketUniverseInstrument, MarketUniverseUnavailable
17 from app.portfolio_orchestration import PortfolioServiceUnavailableError
18 from app.settings import Settings
19
20
21 NOW = datetime(2026, 9, 7, 12, tzinfo=timezone.utc)
22
23
24 def _instrument(*, retrieved_at=NOW - timedelta(hours=1), sector="Technology"):
25 return MarketUniverseInstrument(
26 global_instrument_id=uuid4(), ticker="NONHELD", company_name="Non-held Limited",
27 isin="INE000A01001", exchange="NSE", mic="XNSE", country="IN", currency="INR",
28 asset_type="EQUITY", status="ACTIVE", region="INDIA", canonical_sector=sector,
29 official_industry="Information Technology", source="NSE_INDICES_NIFTY500",
30 retrieved_at=retrieved_at,
31 )
32
33
34 def _prices(latest=NOW - timedelta(hours=1)):
35 return [
36 SimpleNamespace(observed_at=latest - timedelta(days=370), price=Decimal("100")),
37 SimpleNamespace(observed_at=latest, price=Decimal("110")),
38 ]
39
40
41 class _Universe:
42 def __init__(self, values=None, error=False):
43 self.values = values or []
44 self.error = error
45 self.calls = []
46
47 async def listings(self, **kwargs):
48 self.calls.append(kwargs)
49 if self.error:
50 raise MarketUniverseUnavailable("fixture")
51 return self.values
52
53
54 class _SequencedUniverse:
55 def __init__(self, responses):
56 self.responses = responses
57 self.calls = []
58
59 async def listings(self, **kwargs):
60 self.calls.append(kwargs)
61 response = self.responses[min(len(self.calls) - 1, len(self.responses) - 1)]
62 if isinstance(response, Exception):
63 raise response
64 return response
65
66
67 class _Repository:
68 def __init__(self, grouped=None, *, error=False):
69 self.grouped = grouped or {}
70 self.error = error
71 self.calls = 0
72
73 async def market_price_coverage_for_instruments(self, ids):
74 self.calls += 1
75 if self.error:
76 raise RuntimeError("fixture coverage unavailable")
77 result = {}
78 for instrument_id in ids:
79 values = sorted(self.grouped.get(instrument_id, []), key=lambda value: value.observed_at)
80 usable = [value for value in values if value.price is not None and value.price > 0]
81 if usable:
82 result[instrument_id] = (usable[0].observed_at, usable[-1].observed_at, len(usable))
83 return result
84
85
86 class _Jobs:
87 def __init__(self, *, active=None, latest=None, submit_failures=0):
88 self.active_job = active
89 self.latest_job = latest
90 self.submit_failures = submit_failures
91 self.submit_calls = []
92
93 def active(self):
94 return self.active_job
95
96 def latest(self):
97 return self.latest_job
98
99 async def submit(self, **kwargs):
100 self.submit_calls.append(kwargs)
101 if self.submit_failures:
102 self.submit_failures -= 1
103 raise RuntimeError("fixture population submit failed")
104 self.active_job = {"jobId": "population-job", "status": "QUEUED", "region": "INDIA"}
105 return self.active_job
106
107
108 class _Orchestrator:
109 def __init__(self, *, block=False, fail=False, failure_cause=None):
110 self.calls = []
111 self.started = asyncio.Event()
112 self.release = asyncio.Event()
113 self.block = block
114 self.fail = fail
115 self.failure_cause = failure_cause
116
117 async def refresh_india_nifty500_reference(self, **kwargs):
118 self.calls.append(kwargs)
119 self.started.set()
120 if self.block:
121 await self.release.wait()
122 if self.fail:
123 failure = PortfolioServiceUnavailableError("fixture")
124 if self.failure_cause is not None:
125 raise failure from self.failure_cause
126 raise failure
127 return {"activeUniverse": 500}
128
129
130 def _settings():
131 return Settings(
132 market_data_nifty_freshness_hours=12,
133 market_data_historical_freshness_hours=72,
134 market_data_ensure_retry_cooldown_minutes=15,
135 market_data_population_retry_cooldown_hours=12,
136 market_data_internal_user_id="00000000-0000-0000-0000-000000000777",
137 market_data_internal_issuer="test-internal",
138 market_data_internal_subject="research-engine-test",
139 )
140
141
142 def _service(instrument, *, prices=None, jobs=None, orchestrator=None, universe=None):
143 universe = universe or _Universe([instrument] if instrument else [])
144 repository = _Repository({instrument.global_instrument_id: prices or []} if instrument else {})
145 jobs = jobs or _Jobs()
146 orchestrator = orchestrator or _Orchestrator()
147 service = IndiaMarketDataEnsureService(
148 repository, universe, orchestrator, jobs, _settings(), clock=lambda: NOW
149 )
150 return service, universe, repository, jobs, orchestrator
151
152
153 @pytest.mark.asyncio
154 async def test_fresh_universe_and_history_cause_no_background_work():
155 instrument = _instrument()
156 service, _, _, jobs, orchestrator = _service(instrument, prices=_prices())
157 result = await service.ensure(identity_headers={"X-AIP-User-Id": "browser-user"})
158 assert result["universe"]["status"] == "FRESH"
159 assert result["historicalPrices"]["status"] == "FRESH"
160 assert result["historicalPrices"]["coveredInstruments"] == 1
161 assert orchestrator.calls == []
162 assert jobs.submit_calls == []
163
164
165 @pytest.mark.asyncio
166 async def test_stale_universe_starts_exactly_one_global_refresh_for_simultaneous_ensures_and_returns_promptly():
167 instrument = _instrument(retrieved_at=NOW - timedelta(hours=13))
168 orchestrator = _Orchestrator(block=True)
169 service, _, _, _, _ = _service(instrument, prices=_prices(), orchestrator=orchestrator)
170
171 first, second = await asyncio.wait_for(asyncio.gather(
172 service.ensure(identity_headers={"X-AIP-User-Id": "browser-a", "X-AIP-User-Roles": "USER"}),
173 service.ensure(identity_headers={"X-AIP-User-Id": "browser-b", "X-AIP-User-Roles": "ADMIN"}),
174 ), timeout=0.2)
175 assert first["universe"]["status"] == second["universe"]["status"] == "REFRESH_STARTED"
176 await orchestrator.started.wait()
177 assert len(orchestrator.calls) == 1
178 internal_headers = orchestrator.calls[0]["identity_headers"]
179 assert internal_headers == {
180 "X-AIP-User-Id": "00000000-0000-0000-0000-000000000777",
181 "X-AIP-User-Issuer": "test-internal",
182 "X-AIP-User-Subject": "research-engine-test",
183 "X-AIP-User-Roles": "ADMIN",
184 }
185 assert internal_headers["X-AIP-User-Id"] not in {"browser-a", "browser-b"}
186 orchestrator.release.set()
187 await service.wait_for_reference_refresh()
188
189
190 @pytest.mark.asyncio
191 async def test_empty_universe_starts_reference_refresh_and_failed_refresh_observes_retry_cooldown():
192 orchestrator = _Orchestrator(fail=True)
193 service, _, _, _, _ = _service(None, orchestrator=orchestrator)
194 first = await service.ensure(identity_headers={"X-AIP-User-Id": "browser"})
195 assert first["universe"]["status"] == "REFRESH_STARTED"
196 await service.wait_for_reference_refresh()
197 second = await service.ensure(identity_headers={"X-AIP-User-Id": "browser"})
198 assert second["universe"]["status"] == "UNAVAILABLE"
199 assert len(orchestrator.calls) == 1
200
201
202 @pytest.mark.asyncio
203 async def test_successful_stale_reference_refresh_chains_nonblocking_population_only_when_history_needs_it():
204 instrument = _instrument(retrieved_at=NOW - timedelta(hours=13))
205 jobs = _Jobs()
206 service, _, _, jobs, orchestrator = _service(instrument, prices=[], jobs=jobs)
207 result = await service.ensure(identity_headers={"X-AIP-User-Id": "browser"})
208 assert result["universe"]["status"] == "REFRESH_STARTED"
209 assert jobs.submit_calls == []
210 await service.wait_for_reference_refresh()
211 assert len(orchestrator.calls) == 1
212 assert len(jobs.submit_calls) == 1
213
214
215 @pytest.mark.asyncio
216 async def test_post_refresh_rereads_fresh_universe_and_submits_stale_history_once(caplog):
217 caplog.set_level("INFO", logger="app.market_data_ensure")
218 stale = _instrument(retrieved_at=NOW - timedelta(hours=13))
219 fresh = replace(stale, retrieved_at=NOW)
220 universe = _SequencedUniverse([[stale], [fresh]])
221 repository = _Repository()
222 jobs = _Jobs()
223 orchestrator = _Orchestrator()
224 service = IndiaMarketDataEnsureService(
225 repository, universe, orchestrator, jobs, _settings(), clock=lambda: NOW
226 )
227
228 result = await service.ensure(identity_headers={"X-AIP-User-Id": "browser"})
229 assert result["universe"]["status"] == "REFRESH_STARTED"
230 await service.wait_for_reference_refresh()
231
232 assert len(universe.calls) == 2
233 assert universe.calls[1]["identity_headers"]["X-AIP-User-Roles"] == "ADMIN"
234 assert len(jobs.submit_calls) == 1
235 assert "event=REFERENCE_REFRESH_COMPLETED" in caplog.text
236 assert "event=POST_REFRESH_UNIVERSE_READ_COMPLETED instrumentCount=1" in caplog.text
237 assert "event=HISTORICAL_FRESHNESS_RESULT status=STALE" in caplog.text
238 assert "event=POPULATION_SUBMIT_STARTED" in caplog.text
239 assert "event=POPULATION_SUBMIT_COMPLETED jobId=population-job jobStatus=QUEUED" in caplog.text
240
241
242 @pytest.mark.asyncio
243 async def test_post_refresh_universe_read_failure_is_observable_and_later_ensure_retries(caplog):
244 caplog.set_level("INFO", logger="app.market_data_ensure")
245 stale = _instrument(retrieved_at=NOW - timedelta(hours=13))
246 fresh = replace(stale, retrieved_at=NOW)
247 universe = _SequencedUniverse([
248 [stale],
249 MarketUniverseUnavailable("fixture post-refresh read failure"),
250 [fresh],
251 ])
252 jobs = _Jobs()
253 orchestrator = _Orchestrator()
254 service = IndiaMarketDataEnsureService(
255 _Repository(), universe, orchestrator, jobs, _settings(), clock=lambda: NOW
256 )
257
258 await service.ensure(identity_headers={"X-AIP-User-Id": "browser"})
259 await service.wait_for_reference_refresh()
260
261 assert jobs.submit_calls == []
262 assert "event=POST_REFRESH_UNIVERSE_READ_FAILED" in caplog.text
263 assert "errorIdentifier=POST_REFRESH_UNIVERSE_UNAVAILABLE" in caplog.text
264
265 retried = await service.ensure(identity_headers={"X-AIP-User-Id": "browser"})
266 assert retried["historicalPrices"]["status"] == "POPULATION_STARTED"
267 assert len(jobs.submit_calls) == 1
268 assert len(orchestrator.calls) == 1
269
270
271 @pytest.mark.asyncio
272 async def test_exact_year_span_is_compatible_with_historical_freshness_requirement():
273 instrument = _instrument()
274 prices = [
275 SimpleNamespace(observed_at=NOW - timedelta(days=365), price=Decimal("100")),
276 SimpleNamespace(observed_at=NOW, price=Decimal("110")),
277 ]
278 service, _, _, jobs, orchestrator = _service(instrument, prices=prices)
279
280 result = await service.ensure(identity_headers={"X-AIP-User-Id": "browser-user"})
281
282 assert result["historicalPrices"]["status"] == "FRESH"
283 assert result["historicalPrices"]["coveredInstruments"] == 1
284 assert jobs.submit_calls == []
285 assert orchestrator.calls == []
286
287
288 @pytest.mark.asyncio
289 async def test_reference_refresh_timeout_logs_underlying_safe_transport_class(caplog):
290 caplog.set_level("WARNING", logger="app.market_data_ensure")
291 orchestrator = _Orchestrator(
292 fail=True,
293 failure_cause=httpx.ReadTimeout("fixture refresh timeout"),
294 )
295 service, _, _, _, _ = _service(None, orchestrator=orchestrator)
296
297 await service.ensure(identity_headers={"X-AIP-User-Id": "browser"})
298 await service.wait_for_reference_refresh()
299
300 assert "event=REFERENCE_REFRESH_FAILED exceptionClass=ReadTimeout" in caplog.text
301 assert "fixture refresh timeout" not in caplog.text
302
303
304 @pytest.mark.asyncio
305 async def test_post_refresh_unavailable_historical_coverage_is_explicit_and_does_not_submit(caplog):
306 caplog.set_level("INFO", logger="app.market_data_ensure")
307 stale = _instrument(retrieved_at=NOW - timedelta(hours=13))
308 fresh = replace(stale, retrieved_at=NOW)
309 jobs = _Jobs()
310 service = IndiaMarketDataEnsureService(
311 _Repository(error=True),
312 _SequencedUniverse([[stale], [fresh]]),
313 _Orchestrator(),
314 jobs,
315 _settings(),
316 clock=lambda: NOW,
317 )
318
319 await service.ensure(identity_headers={"X-AIP-User-Id": "browser"})
320 await service.wait_for_reference_refresh()
321
322 assert jobs.submit_calls == []
323 assert "event=HISTORICAL_COVERAGE_FAILED" in caplog.text
324 assert "event=HISTORICAL_FRESHNESS_RESULT status=UNAVAILABLE" in caplog.text
325 assert "errorIdentifier=HISTORICAL_COVERAGE_UNAVAILABLE" in caplog.text
326
327
328 @pytest.mark.asyncio
329 async def test_post_refresh_active_population_is_reused_without_duplicate_submit(caplog):
330 caplog.set_level("INFO", logger="app.market_data_ensure")
331 stale = _instrument(retrieved_at=NOW - timedelta(hours=13))
332 active = {"jobId": "running-job", "status": "RUNNING", "region": "INDIA"}
333 jobs = _Jobs(active=active)
334 service = IndiaMarketDataEnsureService(
335 _Repository(), _Universe([stale]), _Orchestrator(), jobs, _settings(), clock=lambda: NOW
336 )
337
338 await service.ensure(identity_headers={"X-AIP-User-Id": "browser"})
339 await service.wait_for_reference_refresh()
340
341 assert jobs.submit_calls == []
342 assert "event=POPULATION_ACTIVE_REUSED jobId=running-job jobStatus=RUNNING" in caplog.text
343
344
345 @pytest.mark.asyncio
346 async def test_post_refresh_recent_terminal_job_preserves_population_cooldown(caplog):
347 caplog.set_level("INFO", logger="app.market_data_ensure")
348 stale = _instrument(retrieved_at=NOW - timedelta(hours=13))
349 recent = {
350 "jobId": "recent-job",
351 "status": "COMPLETED",
352 "completedAt": NOW - timedelta(hours=1),
353 }
354 jobs = _Jobs(latest=recent)
355 service = IndiaMarketDataEnsureService(
356 _Repository(), _Universe([stale]), _Orchestrator(), jobs, _settings(), clock=lambda: NOW
357 )
358
359 await service.ensure(identity_headers={"X-AIP-User-Id": "browser"})
360 await service.wait_for_reference_refresh()
361
362 assert jobs.submit_calls == []
363 assert "event=POPULATION_RETRY_DEFERRED jobId=recent-job jobStatus=COMPLETED" in caplog.text
364
365
366 @pytest.mark.asyncio
367 async def test_post_refresh_submit_failure_is_observable_without_poisoning_refresh_retry_state(caplog):
368 caplog.set_level("INFO", logger="app.market_data_ensure")
369 stale = _instrument(retrieved_at=NOW - timedelta(hours=13))
370 fresh = replace(stale, retrieved_at=NOW)
371 universe = _SequencedUniverse([[stale], [fresh], [fresh]])
372 jobs = _Jobs(submit_failures=1)
373 orchestrator = _Orchestrator()
374 service = IndiaMarketDataEnsureService(
375 _Repository(), universe, orchestrator, jobs, _settings(), clock=lambda: NOW
376 )
377
378 await service.ensure(identity_headers={"X-AIP-User-Id": "browser"})
379 await service.wait_for_reference_refresh()
380
381 assert len(jobs.submit_calls) == 1
382 assert service._last_refresh_failure_at is None
383 assert "event=POPULATION_HANDOFF_FAILED" in caplog.text
384 assert "errorIdentifier=POPULATION_HANDOFF_FAILED" in caplog.text
385
386 retried = await service.ensure(identity_headers={"X-AIP-User-Id": "browser"})
387 assert retried["historicalPrices"]["status"] == "POPULATION_STARTED"
388 assert len(jobs.submit_calls) == 2
389 assert len(orchestrator.calls) == 1
390
391
392 @pytest.mark.asyncio
393 @pytest.mark.parametrize("prices", [[], [SimpleNamespace(observed_at=NOW, price=Decimal("100"))], _prices(NOW - timedelta(hours=73))])
394 async def test_missing_incomplete_or_stale_history_starts_population(prices, caplog):
395 caplog.set_level("INFO", logger="app.market_data_ensure")
396 instrument = _instrument()
397 service, _, _, jobs, orchestrator = _service(instrument, prices=prices)
398 result = await service.ensure(identity_headers={"X-AIP-User-Id": "browser"}, correlation_id="ensure-test")
399 assert result["universe"]["status"] == "FRESH"
400 assert result["historicalPrices"]["status"] == "POPULATION_STARTED"
401 assert result["historicalPrices"]["jobId"] == "population-job"
402 assert len(jobs.submit_calls) == 1 and orchestrator.calls == []
403 assert jobs.submit_calls[0]["identity_headers"]["X-AIP-User-Roles"] == "ADMIN"
404 assert "event=POPULATION_SUBMIT_STARTED" in caplog.text
405 assert "event=POPULATION_SUBMIT_COMPLETED jobId=population-job jobStatus=QUEUED" in caplog.text
406
407
408 @pytest.mark.asyncio
409 async def test_partial_universe_history_is_stale_until_all_eligible_candidates_are_covered():
410 first, second = _instrument(), _instrument()
411 universe = _Universe([first, second])
412 repository = _Repository({first.global_instrument_id: _prices()})
413 jobs = _Jobs()
414 service = IndiaMarketDataEnsureService(
415 repository, universe, _Orchestrator(), jobs, _settings(), clock=lambda: NOW
416 )
417 result = await service.ensure(identity_headers={"X-AIP-User-Id": "browser"})
418 assert result["historicalPrices"]["status"] == "POPULATION_STARTED"
419 assert result["historicalPrices"]["eligibleInstruments"] == 2
420 assert result["historicalPrices"]["coveredInstruments"] == 1
421
422
423 @pytest.mark.asyncio
424 async def test_running_population_is_reused_and_recent_completed_job_prevents_login_storm(caplog):
425 caplog.set_level("INFO", logger="app.market_data_ensure")
426 instrument = _instrument()
427 running = {"jobId": "running-job", "status": "RUNNING", "region": "INDIA"}
428 jobs = _Jobs(active=running)
429 service, _, _, jobs, _ = _service(instrument, prices=[], jobs=jobs)
430 result = await service.ensure(identity_headers={"X-AIP-User-Id": "browser"})
431 assert result["historicalPrices"] ["status"] == "RUNNING"
432 assert result["historicalPrices"]["jobId"] == "running-job"
433 assert jobs.submit_calls == []
434 assert "event=POPULATION_ACTIVE_REUSED jobId=running-job jobStatus=RUNNING" in caplog.text
435
436 jobs.active_job = None
437 jobs.latest_job = {"jobId": "recent", "status": "COMPLETED", "completedAt": NOW - timedelta(hours=1)}
438 caplog.clear()
439 result = await service.ensure(identity_headers={"X-AIP-User-Id": "browser"})
440 assert result["historicalPrices"]["status"] == "STALE"
441 assert jobs.submit_calls == []
442 assert "event=POPULATION_RETRY_DEFERRED jobId=recent jobStatus=COMPLETED" in caplog.text
443
444
445 @pytest.mark.asyncio
446 async def test_unavailable_universe_starts_reference_refresh_instead_of_terminating_ensure():
447 jobs = _Jobs()
448 orchestrator = _Orchestrator()
449 service = IndiaMarketDataEnsureService(
450 _Repository(), _Universe(error=True), orchestrator, jobs, _settings(), clock=lambda: NOW
451 )
452
453 result = await service.ensure(identity_headers={"X-AIP-User-Id": "browser"})
454
455 assert result["universe"]["status"] == "REFRESH_STARTED"
456 assert result["historicalPrices"]["status"] == "UNAVAILABLE"
457
458 await service.wait_for_reference_refresh()
459
460 assert len(orchestrator.calls) == 1
461 internal_headers = orchestrator.calls[0]["identity_headers"]
462 assert internal_headers == {
463 "X-AIP-User-Id": "00000000-0000-0000-0000-000000000777",
464 "X-AIP-User-Issuer": "test-internal",
465 "X-AIP-User-Subject": "research-engine-test",
466 "X-AIP-User-Roles": "ADMIN",
467 }
468
469 # This fixture deliberately keeps the universe unavailable even after the
470 # reference refresh, so population must not start against an unreadable
471 # universe.
472 assert jobs.submit_calls == []
473
474
475 @pytest.mark.asyncio
476 async def test_ensure_route_accepts_normal_authenticated_user_without_admin(monkeypatch):
477 observed = {}
478
479 class Ensure:
480 async def ensure(self, **kwargs):
481 observed.update(kwargs)
482 return {"region": "INDIA", "universe": {"status": "FRESH"}, "historicalPrices": {"status": "FRESH"}}
483
484 monkeypatch.setattr(main, "market_data_ensure_service", Ensure())
485 result = await main.ensure_market_data(
486 region="INDIA", x_correlation_id="correlation", x_aip_user_id="user-id",
487 x_aip_user_issuer="issuer", x_aip_user_subject="subject", x_aip_user_email="user@example.test",
488 x_aip_user_display_name="User", x_aip_user_roles="USER",
489 )
490 assert result["universe"]["status"] == "FRESH"
491 assert observed["identity_headers"]["X-AIP-User-Roles"] == "USER"
492 assert observed["correlation_id"] == "correlation"
493
494 with pytest.raises(HTTPException) as unauthorized:
495 await main.ensure_market_data(
496 region="INDIA", x_correlation_id=None, x_aip_user_id="user-id",
497 x_aip_user_issuer=None, x_aip_user_subject="subject", x_aip_user_email=None,
498 x_aip_user_display_name=None, x_aip_user_roles="USER",
499 )
500 assert unauthorized.value.status_code == 401