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