main
py 1,072 lines 49.1 KB
Raw
1 from uuid import UUID, uuid4
2 import logging
3 import time
4
5 import httpx
6 from fastapi import Body, FastAPI, Header, HTTPException, Query, Request, status
7 from pydantic import BaseModel, Field
8
9 from app.models import CategoryEvidence, PortfolioResearchCompany, PortfolioResearchSummary, ReliabilityLevel, ResearchEventType, ResearchSummary
10 from app.portfolio_orchestration import (
11 GlobalInstrumentNotFoundError,
12 PortfolioResearchOrchestrator,
13 PortfolioServiceUnavailableError,
14 WatchlistNotFoundError,
15 WatchlistRegionMismatchError,
16 )
17 from app.structured_market import StructuredProviderError
18 from app.repository import ResearchRepository
19 from app.scoring import canonical_read_model_score
20 from app.sector_leaderboard import build_sector_leaderboard
21 from app.sector_performance import belongs_to_region, performance_window, rank_performers
22 from app.market_universe import IndiaMarketUniverseProvider, MarketUniverseUnavailable
23 from app.market_data_population import IndiaMarketDataPopulationJobs
24 from app.market_data_ensure import IndiaMarketDataEnsureService
25 from app.market_universe_sectors import canonical_sector_counts
26 from app.watchlists import AddWatchlistInstrumentRequest, EnsureDefaultWatchlistRequest, watchlist_research_projection
27 from app.scheduler import default_schedule_rules
28 from app.settings import Settings, configure_application_logging, reset_request_id, set_request_id
29 from app.sources import default_source_providers
30 from app.research_readiness_runtime import (
31 ExistingResearchCapabilityExecutor,
32 RepositoryResearchReadinessAdapter,
33 ResearchReadinessRuntime,
34 jurisdiction_for_profile,
35 readiness_response,
36 )
37 from app.stock_rule_engine import (
38 StockRuleEngineService,
39 analysis_eligibility_response,
40 )
41 from app.yahoo_mcp_acquisition import (
42 HttpExternalResearchToolGateway,
43 McpFirstResearchCapabilityExecutor,
44 )
45
46 settings = Settings()
47 configure_application_logging(settings)
48 repository = ResearchRepository(settings=settings)
49 portfolio_orchestrator = PortfolioResearchOrchestrator(repository, settings)
50 india_market_universe_provider = IndiaMarketUniverseProvider(portfolio_orchestrator)
51 market_data_population_jobs = IndiaMarketDataPopulationJobs(
52 repository, india_market_universe_provider, portfolio_orchestrator, settings
53 )
54 market_data_ensure_service = IndiaMarketDataEnsureService(
55 repository, india_market_universe_provider, portfolio_orchestrator,
56 market_data_population_jobs, settings,
57 )
58 research_readiness_adapter = RepositoryResearchReadinessAdapter(repository)
59 existing_research_capability_executor = ExistingResearchCapabilityExecutor(
60 repository, portfolio_orchestrator, market_data_population_jobs
61 )
62 research_readiness_runtime = ResearchReadinessRuntime(
63 repository,
64 research_readiness_adapter,
65 McpFirstResearchCapabilityExecutor(
66 existing_research_capability_executor,
67 repository,
68 HttpExternalResearchToolGateway(
69 settings.research_mcp_gateway_base_url,
70 settings.research_mcp_gateway_timeout_seconds,
71 settings.research_mcp_service_identity,
72 ),
73 enabled=settings.research_mcp_first_enabled,
74 ),
75 ensure_timeout_seconds=settings.research_readiness_ensure_timeout_seconds,
76 )
77 stock_rule_engine_service = StockRuleEngineService(repository, research_readiness_adapter)
78 app = FastAPI(title="Research Engine", version="0.3.0")
79 logger = logging.getLogger(__name__)
80
81
82 @app.get('/api/v1/research/opportunities/current')
83 async def opportunity_radar():
84 with repository._persistence_worker_lock:
85 return repository.persistence.opportunity_current()
86
87
88 @app.get('/api/v1/research/opportunities/history/{instrument_id}')
89 async def opportunity_history(instrument_id: UUID):
90 with repository._persistence_worker_lock:
91 return repository.persistence.recommendation_history(instrument_id)
92
93
94 class OpportunityCycleRequest(BaseModel):
95 top_n: int = Field(default=4, ge=2, le=4)
96 shortlist_limit: int = Field(default=25, ge=1, le=100)
97 candidate_ids: list[UUID] | None = Field(default=None, max_length=100)
98
99
100 @app.post('/api/v1/research/opportunities/cycles')
101 async def opportunity_cycle(body: OpportunityCycleRequest, request: Request):
102 from app.global_opportunity_cycle import run_global_opportunity_cycle
103 if not hasattr(repository.persistence, 'publish_opportunity_cycle'):
104 raise HTTPException(503, 'RECOMMENDATION_PERSISTENCE_REQUIRED')
105 # One in-process cycle at a time; no awaits between history read and atomic publish.
106 import asyncio
107 if not hasattr(app.state, 'opportunity_cycle_lock'):
108 app.state.opportunity_cycle_lock = asyncio.Lock()
109 async with app.state.opportunity_cycle_lock:
110 return await run_global_opportunity_cycle(repository, portfolio_orchestrator,
111 top_n=body.top_n, shortlist_limit=body.shortlist_limit, candidate_ids=body.candidate_ids,
112 identity_headers={k: v for k, v in request.headers.items()
113 if k.lower() in {'authorization', 'x-user-id', 'x-correlation-id'}})
114
115
116 class BacktestRequest(BaseModel):
117 start: str
118 end: str
119 market: str = 'NSE'
120 horizon: str = 'SHORT_TERM'
121 benchmark_id: UUID | None = None
122
123
124 @app.get('/api/v1/research/backtesting/runs')
125 async def backtest_runs():
126 with repository._persistence_worker_lock:
127 return repository.persistence.backtests()
128
129
130 @app.post('/api/v1/research/backtesting/runs')
131 async def create_backtest(body: BacktestRequest):
132 from app.recommendation_backtesting import run_backtest
133 if not hasattr(repository.persistence, 'save_backtest'):
134 raise HTTPException(503, 'RECOMMENDATION_PERSISTENCE_REQUIRED')
135 if body.market != 'NSE':
136 raise HTTPException(422, 'UNSUPPORTED_MARKET')
137 try:
138 with repository._persistence_worker_lock:
139 return run_backtest(repository.persistence, start=body.start, end=body.end,
140 horizon=body.horizon, benchmark_id=body.benchmark_id)
141 except ValueError as exc:
142 raise HTTPException(422, str(exc)) from exc
143
144
145 @app.middleware("http")
146 async def correlation_id_middleware(request: Request, call_next):
147 candidate = request.headers.get("X-Request-ID") or request.headers.get("X-Correlation-Id")
148 request_id = candidate.strip() if candidate else ""
149 if not request_id or len(request_id) > 120 or not all(
150 character.isalnum() or character in "-_.:/" for character in request_id
151 ):
152 request_id = str(uuid4())
153 token = set_request_id(request_id)
154 try:
155 response = await call_next(request)
156 finally:
157 reset_request_id(token)
158 response.headers["X-Request-ID"] = request_id
159 response.headers["X-Correlation-Id"] = request_id
160 return response
161
162
163 class ResearchReadinessEnsureRequest(BaseModel):
164 requirements: list[str] | None = Field(default=None)
165
166
167 class StockRuleEngineAnalysisRequest(BaseModel):
168 allow_partial: bool = Field(default=False, alias="allowPartial")
169
170
171 @app.get("/health")
172 def health() -> dict[str, str]:
173 return {"status": "ok", "service": settings.service_name}
174
175
176 @app.get("/providers/llm")
177 def llm_provider() -> dict[str, str]:
178 return {"provider": settings.llm_provider, "mode": "optional-not-called"}
179
180
181 @app.get("/api/v1/research/sources")
182 def sources():
183 return default_source_providers()
184
185
186 @app.get("/api/v1/research/schedule")
187 def schedule():
188 return default_schedule_rules()
189
190
191 @app.get("/api/v1/research/companies")
192 def companies():
193 return repository.list_profiles()
194
195
196 @app.get("/api/v1/research/readiness/{global_instrument_id}")
197 async def research_readiness(
198 global_instrument_id: UUID,
199 x_correlation_id: str | None = Header(default=None),
200 x_aip_user_id: str | None = Header(default=None),
201 x_aip_user_issuer: str | None = Header(default=None),
202 x_aip_user_subject: str | None = Header(default=None),
203 x_aip_user_email: str | None = Header(default=None),
204 x_aip_user_display_name: str | None = Header(default=None),
205 x_aip_user_roles: str | None = Header(default=None),
206 ):
207 """Read public-company readiness without invoking any research provider."""
208 started = time.perf_counter()
209 _require_research_user(x_aip_user_id, x_aip_user_issuer, x_aip_user_subject)
210 identity_headers = _identity_headers(
211 x_aip_user_id,
212 x_aip_user_issuer,
213 x_aip_user_subject,
214 x_aip_user_email,
215 x_aip_user_display_name,
216 x_aip_user_roles,
217 )
218 profile = await _readiness_profile(
219 global_instrument_id, x_correlation_id, identity_headers
220 )
221 result = await research_readiness_runtime.read(
222 global_instrument_id, jurisdiction=jurisdiction_for_profile(profile)
223 )
224 response = readiness_response(result, research_readiness_runtime.requirement_registry)
225 response["analysisEligibility"] = analysis_eligibility_response(result)
226 logger.info(
227 "research_flow operation=READINESS_GET globalInstrumentId=%s outcome=SUCCESS durationMs=%s",
228 global_instrument_id,
229 round((time.perf_counter() - started) * 1000),
230 )
231 return response
232
233
234 @app.post("/api/v1/research/readiness/{global_instrument_id}/ensure")
235 async def ensure_research_readiness(
236 global_instrument_id: UUID,
237 payload: ResearchReadinessEnsureRequest | None = Body(default=None),
238 x_correlation_id: str | None = Header(default=None),
239 x_aip_user_id: str | None = Header(default=None),
240 x_aip_user_issuer: str | None = Header(default=None),
241 x_aip_user_subject: str | None = Header(default=None),
242 x_aip_user_email: str | None = Header(default=None),
243 x_aip_user_display_name: str | None = Header(default=None),
244 x_aip_user_roles: str | None = Header(default=None),
245 ):
246 """Execute only stale/missing planner targets, then re-read readiness."""
247 started = time.perf_counter()
248 _require_research_user(x_aip_user_id, x_aip_user_issuer, x_aip_user_subject)
249 identity_headers = _identity_headers(
250 x_aip_user_id,
251 x_aip_user_issuer,
252 x_aip_user_subject,
253 x_aip_user_email,
254 x_aip_user_display_name,
255 x_aip_user_roles,
256 )
257 profile = await _readiness_profile(
258 global_instrument_id, x_correlation_id, identity_headers
259 )
260 try:
261 result = await research_readiness_runtime.ensure(
262 global_instrument_id,
263 jurisdiction=jurisdiction_for_profile(profile),
264 requirement_ids=payload.requirements if payload else None,
265 correlation_id=x_correlation_id,
266 identity_headers=identity_headers,
267 )
268 except ValueError as exc:
269 raise HTTPException(status_code=400, detail=str(exc)) from exc
270 response = readiness_response(
271 result.readiness,
272 research_readiness_runtime.requirement_registry,
273 ensure=result,
274 )
275 response["analysisEligibility"] = analysis_eligibility_response(result.readiness)
276 logger.info(
277 "research_flow operation=TARGETED_ENSURE globalInstrumentId=%s outcome=SUCCESS durationMs=%s planned=%s executed=%s failures=%s singleFlightReused=%s",
278 global_instrument_id,
279 round((time.perf_counter() - started) * 1000),
280 len(result.planned_requirement_ids),
281 len(result.executed_capabilities),
282 len(result.failures),
283 result.reused_single_flight,
284 )
285 return response
286
287
288 @app.post("/api/v1/research/analysis/{global_instrument_id}")
289 async def analyze_company_research(
290 global_instrument_id: UUID,
291 payload: StockRuleEngineAnalysisRequest | None = Body(default=None),
292 x_correlation_id: str | None = Header(default=None),
293 x_aip_user_id: str | None = Header(default=None),
294 x_aip_user_issuer: str | None = Header(default=None),
295 x_aip_user_subject: str | None = Header(default=None),
296 x_aip_user_email: str | None = Header(default=None),
297 x_aip_user_display_name: str | None = Header(default=None),
298 x_aip_user_roles: str | None = Header(default=None),
299 ):
300 """Calculate a versioned public-company score using durable data only."""
301 started = time.perf_counter()
302 _require_research_user(x_aip_user_id, x_aip_user_issuer, x_aip_user_subject)
303 identity_headers = _identity_headers(
304 x_aip_user_id,
305 x_aip_user_issuer,
306 x_aip_user_subject,
307 x_aip_user_email,
308 x_aip_user_display_name,
309 x_aip_user_roles,
310 )
311 profile = await _readiness_profile(
312 global_instrument_id, x_correlation_id, identity_headers
313 )
314 readiness = await research_readiness_runtime.read(
315 global_instrument_id, jurisdiction=jurisdiction_for_profile(profile)
316 )
317 result = await stock_rule_engine_service.analyze(
318 profile,
319 readiness,
320 allow_partial=payload.allow_partial if payload else False,
321 )
322 logger.info(
323 "research_flow operation=RULE_ENGINE_ANALYSIS globalInstrumentId=%s outcome=SUCCESS durationMs=%s cacheHit=%s decisionSignal=%s",
324 global_instrument_id,
325 round((time.perf_counter() - started) * 1000),
326 result.cache_hit,
327 result.decision_signal,
328 )
329 return result.model_dump(mode="json", by_alias=True)
330
331
332 @app.get("/api/v1/research/instruments/search")
333 async def search_research_instruments(
334 q: str = Query(..., min_length=3, max_length=128),
335 region: str = Query(default="INDIA"),
336 limit: int = Query(default=20, ge=1, le=20),
337 x_correlation_id: str | None = Header(default=None),
338 x_aip_user_id: str | None = Header(default=None),
339 x_aip_user_issuer: str | None = Header(default=None),
340 x_aip_user_subject: str | None = Header(default=None),
341 x_aip_user_email: str | None = Header(default=None),
342 x_aip_user_display_name: str | None = Header(default=None),
343 x_aip_user_roles: str | None = Header(default=None),
344 ) -> list[dict]:
345 """Search the global canonical instrument universe from official masters.
346
347 India uses the persisted canonical NSE-backed master; the USA and
348 EUROPE universes use the verified portfolio-service active-equity master. This
349 never proxies a browser-side NSE/Yahoo text search and never restricts results
350 to portfolio holdings or a watchlist.
351 """
352 started = time.perf_counter()
353 _require_research_user(x_aip_user_id, x_aip_user_issuer, x_aip_user_subject)
354 q = q.strip()
355 if len(q) < 3:
356 raise HTTPException(status_code=422, detail="SEARCH_QUERY_TOO_SHORT")
357 normalized_region = str(region).strip().upper()
358 if normalized_region not in {"USA", "EUROPE", "INDIA"}:
359 raise HTTPException(status_code=400, detail="UNSUPPORTED_REGION")
360 identity_headers = _identity_headers(
361 x_aip_user_id, x_aip_user_issuer, x_aip_user_subject,
362 x_aip_user_email, x_aip_user_display_name, x_aip_user_roles,
363 )
364 try:
365 results = await portfolio_orchestrator.search_instruments(
366 q,
367 normalized_region,
368 limit=limit,
369 correlation_id=x_correlation_id,
370 identity_headers=identity_headers,
371 )
372 except PortfolioServiceUnavailableError as exc:
373 raise HTTPException(
374 status_code=502,
375 detail="Portfolio service unavailable for instrument search",
376 ) from exc
377 logger.info(
378 "research_flow operation=INSTRUMENT_SEARCH region=%s query=%s matches=%s durationMs=%s",
379 normalized_region,
380 q[:64],
381 len(results),
382 round((time.perf_counter() - started) * 1000),
383 )
384 return results
385
386
387 @app.get("/api/v1/research/sector-leaderboard")
388 async def sector_leaderboard(
389 x_correlation_id: str | None = Header(default=None),
390 x_aip_user_id: str | None = Header(default=None),
391 x_aip_user_issuer: str | None = Header(default=None),
392 x_aip_user_subject: str | None = Header(default=None),
393 x_aip_user_email: str | None = Header(default=None),
394 x_aip_user_display_name: str | None = Header(default=None),
395 x_aip_user_roles: str | None = Header(default=None),
396 ):
397 candidates = []
398 identity_headers = _identity_headers(
399 x_aip_user_id, x_aip_user_issuer, x_aip_user_subject,
400 x_aip_user_email, x_aip_user_display_name, x_aip_user_roles,
401 )
402 for item in await portfolio_orchestrator.active_global_equities(
403 correlation_id=x_correlation_id, identity_headers=identity_headers
404 ):
405 try:
406 instrument_id = UUID(str(item["globalInstrumentId"]))
407 except (KeyError, ValueError):
408 continue
409 score = repository.persisted_canonical_read_model_score(instrument_id)
410 if score is None or score.overall_score is None:
411 continue
412 records = repository.structured_market_snapshots_for({instrument_id}).get(instrument_id, [])
413 sector = next((str(record.snapshot.facts["sector"].value) for record in records if record.snapshot.facts.get("sector") and record.snapshot.facts["sector"].value), None)
414 if not sector:
415 continue
416 coverage = sum(1 for value in score.category_evidence.values() if value.score is not None)
417 candidates.append({"globalInstrumentId": str(instrument_id), "companyName": item.get("canonicalName"), "ticker": item.get("ticker"), "exchange": item.get("exchange"), "country": item.get("country"), "currency": item.get("currency"), "sector": sector, "industry": next((str(record.snapshot.facts["industry"].value) for record in records if record.snapshot.facts.get("industry")), None), "score": score.overall_score, "evidenceCoverage": coverage})
418 return build_sector_leaderboard(candidates)
419
420
421 @app.get("/api/v1/research/sector-performance")
422 async def sector_performance(
423 region: str = Query(...), sector: str = Query(...), period: str = Query(...), limit: int = Query(default=5),
424 x_correlation_id: str | None = Header(default=None),
425 x_aip_user_id: str | None = Header(default=None), x_aip_user_issuer: str | None = Header(default=None),
426 x_aip_user_subject: str | None = Header(default=None), x_aip_user_email: str | None = Header(default=None),
427 x_aip_user_display_name: str | None = Header(default=None), x_aip_user_roles: str | None = Header(default=None),
428 ):
429 from app.sector_leaderboard import normalize_sector
430 normalized_sector, sector_label = normalize_sector(sector)
431 normalized_period = period.upper()
432 if normalized_period not in {"DAY", "WEEK", "MONTH", "YEAR"}:
433 raise HTTPException(status_code=400, detail="UNSUPPORTED_PERFORMANCE_PERIOD")
434 normalized_region = region.upper()
435 if normalized_region not in {"USA", "EUROPE", "INDIA"}:
436 raise HTTPException(status_code=400, detail="UNSUPPORTED_REGION")
437 identity_headers = _identity_headers(x_aip_user_id, x_aip_user_issuer, x_aip_user_subject, x_aip_user_email, x_aip_user_display_name, x_aip_user_roles)
438 if normalized_region == "INDIA":
439 try:
440 universe = [value.as_payload() for value in await india_market_universe_provider.listings(
441 correlation_id=x_correlation_id, identity_headers=identity_headers
442 )]
443 except MarketUniverseUnavailable as exc:
444 raise HTTPException(status_code=503, detail="INDIA_MARKET_UNIVERSE_UNAVAILABLE") from exc
445 else:
446 universe = await portfolio_orchestrator.active_global_equities(
447 correlation_id=x_correlation_id, identity_headers=identity_headers
448 )
449 eligible = [item for item in universe if belongs_to_region(item, normalized_region)]
450 ids = {UUID(str(item["globalInstrumentId"])) for item in eligible if item.get("globalInstrumentId")}
451 snapshots = {} if normalized_region == "INDIA" else repository.structured_market_snapshots_for(ids)
452 observations = repository.market_price_observations_for(ids)
453 candidates = []
454 for item in eligible:
455 try:
456 instrument_id = UUID(str(item["globalInstrumentId"]))
457 except (KeyError, ValueError):
458 continue
459 records = snapshots.get(instrument_id, [])
460 raw_sector = (item.get("canonicalSector") if normalized_region == "INDIA" else
461 next((str(record.snapshot.facts["sector"].value) for record in records if record.snapshot.facts.get("sector") and record.snapshot.facts["sector"].value), None))
462 if not raw_sector or normalize_sector(raw_sector)[0] != normalized_sector:
463 continue
464 window = performance_window(observations.get(instrument_id, []), normalized_period)
465 if window is None:
466 continue
467 latest, reference, performance_pct = window
468 candidates.append({"globalInstrumentId": str(instrument_id), "companyName": item.get("canonicalName"), "ticker": item.get("ticker"), "exchange": item.get("exchange"), "currency": latest.currency or item.get("currency"), "latestPrice": latest.price, "referencePrice": reference.price, "performancePct": performance_pct, "_asOf": latest.observed_at})
469 best, worst = rank_performers(candidates, limit)
470 selected = [*best, *worst]
471 return {"region": normalized_region, "sector": sector_label, "period": normalized_period,
472 "asOf": max((value["_asOf"] for value in selected), default=None),
473 "bestPerformers": [{key: value for key, value in row.items() if key != "_asOf"} for row in best],
474 "worstPerformers": [{key: value for key, value in row.items() if key != "_asOf"} for row in worst]}
475
476
477 @app.get("/api/v1/research/market-universe/sectors")
478 async def market_universe_sectors(
479 region: str = Query(...),
480 x_correlation_id: str | None = Header(default=None),
481 x_aip_user_id: str | None = Header(default=None),
482 x_aip_user_issuer: str | None = Header(default=None),
483 x_aip_user_subject: str | None = Header(default=None),
484 x_aip_user_email: str | None = Header(default=None),
485 x_aip_user_display_name: str | None = Header(default=None),
486 x_aip_user_roles: str | None = Header(default=None),
487 ):
488 """List sectors represented by the selected region's durable universe."""
489 normalized_region = str(region).strip().upper()
490 if normalized_region not in {"USA", "EUROPE", "INDIA"}:
491 raise HTTPException(status_code=400, detail="UNSUPPORTED_REGION")
492 identity_headers = _identity_headers(
493 x_aip_user_id, x_aip_user_issuer, x_aip_user_subject,
494 x_aip_user_email, x_aip_user_display_name, x_aip_user_roles,
495 )
496 try:
497 if normalized_region == "INDIA":
498 universe = [
499 value.as_payload()
500 for value in await india_market_universe_provider.listings(
501 correlation_id=x_correlation_id,
502 identity_headers=identity_headers,
503 )
504 ]
505 else:
506 universe = await portfolio_orchestrator.active_global_equities(
507 correlation_id=x_correlation_id,
508 identity_headers=identity_headers,
509 )
510 except (MarketUniverseUnavailable, PortfolioServiceUnavailableError) as exc:
511 raise HTTPException(status_code=503, detail="MARKET_UNIVERSE_UNAVAILABLE") from exc
512
513 eligible = [item for item in universe if belongs_to_region(item, normalized_region)]
514 if normalized_region == "INDIA":
515 sector_by_instrument = {
516 str(item.get("globalInstrumentId") or "").strip(): item.get("canonicalSector")
517 for item in eligible
518 }
519 else:
520 ids_by_text: dict[str, UUID] = {}
521 for item in eligible:
522 try:
523 instrument_id = UUID(str(item["globalInstrumentId"]))
524 except (KeyError, ValueError, TypeError):
525 continue
526 ids_by_text[str(instrument_id)] = instrument_id
527 snapshots = repository.structured_market_snapshots_for(set(ids_by_text.values()))
528 sector_by_instrument = {
529 identity: next(
530 (
531 str(record.snapshot.facts["sector"].value)
532 for record in snapshots.get(instrument_id, [])
533 if record.snapshot.facts.get("sector")
534 and record.snapshot.facts["sector"].value
535 ),
536 None,
537 )
538 for identity, instrument_id in ids_by_text.items()
539 }
540
541 return {
542 "region": normalized_region,
543 "sectors": canonical_sector_counts(eligible, sector_by_instrument),
544 }
545
546
547 @app.get("/api/v1/research/watchlists")
548 async def list_research_watchlists(
549 x_correlation_id: str | None = Header(default=None),
550 x_aip_user_id: str | None = Header(default=None),
551 x_aip_user_issuer: str | None = Header(default=None),
552 x_aip_user_subject: str | None = Header(default=None),
553 x_aip_user_email: str | None = Header(default=None),
554 x_aip_user_display_name: str | None = Header(default=None),
555 x_aip_user_roles: str | None = Header(default=None),
556 ):
557 _require_research_user(x_aip_user_id, x_aip_user_issuer, x_aip_user_subject)
558 try:
559 return await portfolio_orchestrator.list_watchlists(
560 correlation_id=x_correlation_id,
561 identity_headers=_identity_headers(
562 x_aip_user_id, x_aip_user_issuer, x_aip_user_subject,
563 x_aip_user_email, x_aip_user_display_name, x_aip_user_roles,
564 ),
565 )
566 except PortfolioServiceUnavailableError as exc:
567 raise HTTPException(status_code=502, detail="WATCHLIST_SERVICE_UNAVAILABLE") from exc
568
569
570 @app.post("/api/v1/research/watchlists/default/ensure")
571 async def ensure_default_research_watchlist(
572 payload: EnsureDefaultWatchlistRequest,
573 x_correlation_id: str | None = Header(default=None),
574 x_aip_user_id: str | None = Header(default=None),
575 x_aip_user_issuer: str | None = Header(default=None),
576 x_aip_user_subject: str | None = Header(default=None),
577 x_aip_user_email: str | None = Header(default=None),
578 x_aip_user_display_name: str | None = Header(default=None),
579 x_aip_user_roles: str | None = Header(default=None),
580 ):
581 _require_research_user(x_aip_user_id, x_aip_user_issuer, x_aip_user_subject)
582 try:
583 return await portfolio_orchestrator.ensure_default_watchlist(
584 payload.region,
585 correlation_id=x_correlation_id,
586 identity_headers=_identity_headers(
587 x_aip_user_id, x_aip_user_issuer, x_aip_user_subject,
588 x_aip_user_email, x_aip_user_display_name, x_aip_user_roles,
589 ),
590 )
591 except PortfolioServiceUnavailableError as exc:
592 raise HTTPException(status_code=502, detail="WATCHLIST_SERVICE_UNAVAILABLE") from exc
593
594
595 @app.post("/api/v1/research/watchlists/{watchlist_id}/instruments")
596 async def add_research_watchlist_instrument(
597 watchlist_id: UUID,
598 payload: AddWatchlistInstrumentRequest,
599 x_correlation_id: str | None = Header(default=None),
600 x_aip_user_id: str | None = Header(default=None),
601 x_aip_user_issuer: str | None = Header(default=None),
602 x_aip_user_subject: str | None = Header(default=None),
603 x_aip_user_email: str | None = Header(default=None),
604 x_aip_user_display_name: str | None = Header(default=None),
605 x_aip_user_roles: str | None = Header(default=None),
606 ):
607 _require_research_user(x_aip_user_id, x_aip_user_issuer, x_aip_user_subject)
608 try:
609 return await portfolio_orchestrator.add_watchlist_instrument(
610 watchlist_id,
611 payload.portfolio_payload(),
612 correlation_id=x_correlation_id,
613 identity_headers=_identity_headers(
614 x_aip_user_id, x_aip_user_issuer, x_aip_user_subject,
615 x_aip_user_email, x_aip_user_display_name, x_aip_user_roles,
616 ),
617 )
618 except WatchlistRegionMismatchError as exc:
619 raise HTTPException(status_code=409, detail="WATCHLIST_REGION_MISMATCH") from exc
620 except WatchlistNotFoundError as exc:
621 raise HTTPException(status_code=404, detail="WATCHLIST_NOT_FOUND") from exc
622 except PortfolioServiceUnavailableError as exc:
623 raise HTTPException(status_code=502, detail="WATCHLIST_SERVICE_UNAVAILABLE") from exc
624
625
626 @app.delete("/api/v1/research/watchlists/{watchlist_id}/instruments/{instrument_id}", status_code=204)
627 async def remove_research_watchlist_instrument(
628 watchlist_id: UUID,
629 instrument_id: UUID,
630 x_correlation_id: str | None = Header(default=None),
631 x_aip_user_id: str | None = Header(default=None),
632 x_aip_user_issuer: str | None = Header(default=None),
633 x_aip_user_subject: str | None = Header(default=None),
634 x_aip_user_email: str | None = Header(default=None),
635 x_aip_user_display_name: str | None = Header(default=None),
636 x_aip_user_roles: str | None = Header(default=None),
637 ):
638 _require_research_user(x_aip_user_id, x_aip_user_issuer, x_aip_user_subject)
639 try:
640 await portfolio_orchestrator.remove_watchlist_instrument(
641 watchlist_id, instrument_id, correlation_id=x_correlation_id,
642 identity_headers=_identity_headers(
643 x_aip_user_id, x_aip_user_issuer, x_aip_user_subject,
644 x_aip_user_email, x_aip_user_display_name, x_aip_user_roles,
645 ),
646 )
647 except WatchlistNotFoundError as exc:
648 raise HTTPException(status_code=404, detail="WATCHLIST_NOT_FOUND") from exc
649 except PortfolioServiceUnavailableError as exc:
650 raise HTTPException(status_code=502, detail="WATCHLIST_SERVICE_UNAVAILABLE") from exc
651
652
653 @app.get("/api/v1/research/watchlists/{watchlist_id}/research")
654 async def research_watchlist_presentation(
655 watchlist_id: UUID,
656 x_correlation_id: str | None = Header(default=None),
657 x_aip_user_id: str | None = Header(default=None),
658 x_aip_user_issuer: str | None = Header(default=None),
659 x_aip_user_subject: str | None = Header(default=None),
660 x_aip_user_email: str | None = Header(default=None),
661 x_aip_user_display_name: str | None = Header(default=None),
662 x_aip_user_roles: str | None = Header(default=None),
663 ):
664 _require_research_user(x_aip_user_id, x_aip_user_issuer, x_aip_user_subject)
665 identity = _identity_headers(
666 x_aip_user_id, x_aip_user_issuer, x_aip_user_subject,
667 x_aip_user_email, x_aip_user_display_name, x_aip_user_roles,
668 )
669 try:
670 watchlist = await portfolio_orchestrator.watchlist(
671 watchlist_id, correlation_id=x_correlation_id, identity_headers=identity,
672 )
673 return await watchlist_research_projection(
674 portfolio_orchestrator, watchlist,
675 correlation_id=x_correlation_id, identity_headers=identity,
676 )
677 except WatchlistNotFoundError as exc:
678 raise HTTPException(status_code=404, detail="WATCHLIST_NOT_FOUND") from exc
679 except PortfolioServiceUnavailableError as exc:
680 raise HTTPException(status_code=502, detail="WATCHLIST_SERVICE_UNAVAILABLE") from exc
681
682
683 @app.post("/api/v1/research/market-data/population", status_code=status.HTTP_202_ACCEPTED)
684 async def start_market_data_population(
685 region: str = Query(...),
686 x_correlation_id: str | None = Header(default=None),
687 x_aip_user_id: str | None = Header(default=None),
688 x_aip_user_issuer: str | None = Header(default=None),
689 x_aip_user_subject: str | None = Header(default=None),
690 x_aip_user_email: str | None = Header(default=None),
691 x_aip_user_display_name: str | None = Header(default=None),
692 x_aip_user_roles: str | None = Header(default=None),
693 ):
694 _require_market_data_admin(x_aip_user_id, x_aip_user_issuer, x_aip_user_subject, x_aip_user_roles)
695 if str(region).strip().upper() != "INDIA":
696 raise HTTPException(status_code=400, detail="MARKET_DATA_POPULATION_REGION_NOT_SUPPORTED")
697 return await market_data_population_jobs.submit(
698 identity_headers=_identity_headers(
699 x_aip_user_id, x_aip_user_issuer, x_aip_user_subject,
700 x_aip_user_email, x_aip_user_display_name, x_aip_user_roles,
701 ),
702 correlation_id=x_correlation_id,
703 )
704
705
706 @app.post("/api/v1/research/market-data/ensure")
707 async def ensure_market_data(
708 region: str = Query(...),
709 x_correlation_id: str | None = Header(default=None),
710 x_aip_user_id: str | None = Header(default=None),
711 x_aip_user_issuer: str | None = Header(default=None),
712 x_aip_user_subject: str | None = Header(default=None),
713 x_aip_user_email: str | None = Header(default=None),
714 x_aip_user_display_name: str | None = Header(default=None),
715 x_aip_user_roles: str | None = Header(default=None),
716 ):
717 _require_market_data_user(x_aip_user_id, x_aip_user_issuer, x_aip_user_subject)
718 if str(region).strip().upper() != "INDIA":
719 raise HTTPException(status_code=400, detail="MARKET_DATA_ENSURE_REGION_NOT_SUPPORTED")
720 return await market_data_ensure_service.ensure(
721 identity_headers=_identity_headers(
722 x_aip_user_id, x_aip_user_issuer, x_aip_user_subject,
723 x_aip_user_email, x_aip_user_display_name, x_aip_user_roles,
724 ),
725 correlation_id=x_correlation_id,
726 )
727
728
729 @app.get("/api/v1/research/market-data/population/{job_id}")
730 async def market_data_population_status(
731 job_id: UUID,
732 x_aip_user_id: str | None = Header(default=None),
733 x_aip_user_issuer: str | None = Header(default=None),
734 x_aip_user_subject: str | None = Header(default=None),
735 x_aip_user_roles: str | None = Header(default=None),
736 ):
737 _require_market_data_admin(x_aip_user_id, x_aip_user_issuer, x_aip_user_subject, x_aip_user_roles)
738 job = market_data_population_jobs.get(job_id)
739 if job is None:
740 raise HTTPException(status_code=404, detail="MARKET_DATA_POPULATION_JOB_NOT_FOUND")
741 return job
742
743
744 @app.get("/api/v1/research/companies/{instrument_id}")
745 async def company(instrument_id: UUID, x_correlation_id: str | None = Header(default=None), x_aip_user_id: str | None = Header(default=None), x_aip_user_issuer: str | None = Header(default=None), x_aip_user_subject: str | None = Header(default=None), x_aip_user_email: str | None = Header(default=None), x_aip_user_display_name: str | None = Header(default=None), x_aip_user_roles: str | None = Header(default=None)):
746 try:
747 return repository.profile(instrument_id)
748 except StopIteration as exc:
749 if await _resolve_restored_instrument_profile(instrument_id, x_correlation_id, _identity_headers(x_aip_user_id, x_aip_user_issuer, x_aip_user_subject, x_aip_user_email, x_aip_user_display_name, x_aip_user_roles)):
750 return repository.profile(instrument_id)
751 raise HTTPException(status_code=404, detail="COMPANY_NOT_RESOLVED") from exc
752
753
754 @app.get("/api/v1/research/companies/{instrument_id}/events")
755 async def events(
756 instrument_id: UUID,
757 event_type: ResearchEventType | None = Query(default=None, alias="eventType"),
758 impact: str | None = None,
759 reliability: ReliabilityLevel | None = None,
760 x_correlation_id: str | None = Header(default=None),
761 x_aip_user_id: str | None = Header(default=None),
762 x_aip_user_issuer: str | None = Header(default=None),
763 x_aip_user_subject: str | None = Header(default=None),
764 ):
765 await _require_profile_or_restored_instrument(instrument_id, x_correlation_id, _identity_headers(x_aip_user_id, x_aip_user_issuer, x_aip_user_subject, None, None, None))
766 return repository.events_for(instrument_id, event_type=event_type, impact=impact, reliability=reliability)
767
768
769 @app.get("/api/v1/research/companies/{instrument_id}/documents")
770 async def documents(instrument_id: UUID, x_correlation_id: str | None = Header(default=None), x_aip_user_id: str | None = Header(default=None), x_aip_user_issuer: str | None = Header(default=None), x_aip_user_subject: str | None = Header(default=None)):
771 await _require_profile_or_restored_instrument(instrument_id, x_correlation_id, _identity_headers(x_aip_user_id, x_aip_user_issuer, x_aip_user_subject, None, None, None))
772 return repository.documents_for(instrument_id)
773
774
775 @app.get("/api/v1/research/companies/{instrument_id}/summary", response_model=ResearchSummary)
776 async def summary(instrument_id: UUID, x_correlation_id: str | None = Header(default=None), x_aip_user_id: str | None = Header(default=None), x_aip_user_issuer: str | None = Header(default=None), x_aip_user_subject: str | None = Header(default=None), x_aip_user_email: str | None = Header(default=None), x_aip_user_display_name: str | None = Header(default=None), x_aip_user_roles: str | None = Header(default=None)) -> ResearchSummary:
777 identity_headers = _identity_headers(x_aip_user_id, x_aip_user_issuer, x_aip_user_subject, x_aip_user_email, x_aip_user_display_name, x_aip_user_roles)
778 if await _is_restored_etf_instrument(instrument_id, x_correlation_id, identity_headers):
779 raise HTTPException(status_code=404, detail="ETF research summary is available through portfolio research")
780 await _require_profile_or_restored_instrument(instrument_id, x_correlation_id, identity_headers)
781 return _normalized_summary_response(repository.summary(instrument_id, allow_demo=True))
782
783
784 @app.get(
785 "/api/v1/research/companies/{instrument_id}/presentation",
786 response_model=PortfolioResearchCompany,
787 )
788 async def global_company_research_presentation(
789 instrument_id: UUID,
790 region: str = Query(...),
791 x_correlation_id: str | None = Header(default=None),
792 x_aip_user_id: str | None = Header(default=None),
793 x_aip_user_issuer: str | None = Header(default=None),
794 x_aip_user_subject: str | None = Header(default=None),
795 x_aip_user_email: str | None = Header(default=None),
796 x_aip_user_display_name: str | None = Header(default=None),
797 x_aip_user_roles: str | None = Header(default=None),
798 ) -> PortfolioResearchCompany:
799 """Read a transient, non-held company row from durable global research."""
800 _require_research_user(x_aip_user_id, x_aip_user_issuer, x_aip_user_subject)
801 normalized_region = str(region).strip().upper()
802 if normalized_region not in {"USA", "EUROPE", "INDIA"}:
803 raise HTTPException(status_code=400, detail="UNSUPPORTED_REGION")
804 identity_headers = _identity_headers(
805 x_aip_user_id, x_aip_user_issuer, x_aip_user_subject,
806 x_aip_user_email, x_aip_user_display_name, x_aip_user_roles,
807 )
808 try:
809 metadata = await portfolio_orchestrator.global_instrument_metadata(
810 instrument_id,
811 correlation_id=x_correlation_id,
812 identity_headers=identity_headers,
813 )
814 canonical_identity = {
815 "country": metadata.get("country"),
816 "exchange": metadata.get("primaryExchange") or metadata.get("exchange"),
817 "mic": metadata.get("primaryMic") or metadata.get("mic"),
818 }
819 if not belongs_to_region(canonical_identity, normalized_region):
820 raise HTTPException(status_code=409, detail="INSTRUMENT_REGION_MISMATCH")
821 company = await portfolio_orchestrator.read_global_company_state(
822 instrument_id,
823 metadata=metadata,
824 correlation_id=x_correlation_id,
825 identity_headers=identity_headers,
826 )
827 company = await portfolio_orchestrator.enrich_global_company_durables(company)
828 return company.model_copy(update={
829 "evidence_coverage": _canonical_casefolded_text(company.evidence_coverage),
830 })
831 except PortfolioServiceUnavailableError as exc:
832 raise HTTPException(
833 status_code=502,
834 detail="Portfolio service unavailable for global research presentation",
835 ) from exc
836
837
838 def _normalized_summary_response(summary: ResearchSummary) -> ResearchSummary:
839 """Normalize legacy/display category aliases at the HTTP contract edge."""
840 score = canonical_read_model_score(summary.catalyst_score)
841 buckets = _canonical_casefolded_values(score.buckets)
842 evidence = _canonical_casefolded_evidence(score.category_evidence)
843 if buckets == summary.catalyst_score.buckets and evidence == summary.catalyst_score.category_evidence:
844 return summary
845 return summary.model_copy(update={
846 "catalyst_score": score.model_copy(update={"buckets": buckets, "category_evidence": evidence}),
847 })
848
849
850 def _canonical_casefolded_values(values: dict[str, int | None]) -> dict[str, int | None]:
851 groups: dict[str, list[str]] = {}
852 for key in values:
853 groups.setdefault(key.casefold(), []).append(key)
854 normalized: dict[str, int | None] = {}
855 for keys in groups.values():
856 canonical = _canonical_category_key(keys)
857 normalized[canonical] = values[canonical]
858 if normalized[canonical] is None:
859 normalized[canonical] = next((values[key] for key in keys if values[key] is not None), None)
860 return normalized
861
862
863 def _canonical_casefolded_evidence(values: dict[str, CategoryEvidence]) -> dict[str, CategoryEvidence]:
864 groups: dict[str, list[str]] = {}
865 for key in values:
866 groups.setdefault(key.casefold(), []).append(key)
867 normalized: dict[str, CategoryEvidence] = {}
868 for keys in groups.values():
869 canonical = _canonical_category_key(keys)
870 selected = values[canonical]
871 if selected.score is None:
872 selected = next((values[key] for key in keys if values[key].score is not None), selected)
873 normalized[canonical] = selected.model_copy(update={"category": canonical})
874 return normalized
875
876
877 def _canonical_casefolded_text(values: dict[str, str]) -> dict[str, str]:
878 """Apply the same legacy/canonical key rule to portfolio evidence status."""
879 groups: dict[str, list[str]] = {}
880 for key in values:
881 groups.setdefault(key.casefold(), []).append(key)
882 normalized: dict[str, str] = {}
883 for keys in groups.values():
884 canonical = _canonical_category_key(keys)
885 selected = values[canonical]
886 if not selected:
887 selected = next((values[key] for key in keys if values[key]), selected)
888 normalized[canonical] = selected
889 return normalized
890
891
892 def _normalized_portfolio_summary_response(summary: PortfolioResearchSummary) -> PortfolioResearchSummary:
893 """Normalize only API-facing legacy research-category aliases per company."""
894 companies = [
895 company.model_copy(update={
896 "evidence_coverage": _canonical_casefolded_text(company.evidence_coverage),
897 })
898 for company in summary.companies
899 ]
900 if all(company.evidence_coverage == normalized.evidence_coverage
901 for company, normalized in zip(summary.companies, companies, strict=True)):
902 return summary
903 return summary.model_copy(update={"companies": companies})
904
905
906 def _canonical_category_key(keys: list[str]) -> str:
907 # Canonical category identifiers are upper-case in the research contract;
908 # retain the original key only when a duplicate set has no such identifier.
909 return next((key for key in keys if key == key.upper()), keys[0])
910
911
912 @app.get("/api/v1/research/portfolios/{portfolio_id}/summary", response_model=PortfolioResearchSummary)
913 async def portfolio_research_summary(
914 portfolio_id: UUID,
915 x_correlation_id: str | None = Header(default=None),
916 x_aip_user_id: str | None = Header(default=None),
917 x_aip_user_issuer: str | None = Header(default=None),
918 x_aip_user_subject: str | None = Header(default=None),
919 x_aip_user_email: str | None = Header(default=None),
920 x_aip_user_display_name: str | None = Header(default=None),
921 x_aip_user_roles: str | None = Header(default=None),
922 ):
923 identity_headers = {
924 "X-AIP-User-Id": x_aip_user_id,
925 "X-AIP-User-Issuer": x_aip_user_issuer,
926 "X-AIP-User-Subject": x_aip_user_subject,
927 "X-AIP-User-Email": x_aip_user_email,
928 "X-AIP-User-Display-Name": x_aip_user_display_name,
929 "X-AIP-User-Roles": x_aip_user_roles,
930 }
931 try:
932 result = await portfolio_orchestrator.read_portfolio_summary(
933 portfolio_id,
934 correlation_id=x_correlation_id,
935 identity_headers=identity_headers,
936 )
937 normalization_started = time.perf_counter()
938 normalized = _normalized_portfolio_summary_response(result)
939 logger.info(
940 "portfolio_summary_stage stage=NORMALIZE portfolioId=%s durationMs=%s",
941 portfolio_id,
942 round((time.perf_counter() - normalization_started) * 1000),
943 )
944 return normalized
945 except httpx.HTTPStatusError as exc:
946 status_code = exc.response.status_code
947 if status_code in {401, 403, 404}:
948 raise HTTPException(status_code=status_code, detail="Portfolio research access denied") from exc
949 raise HTTPException(status_code=502, detail="Portfolio service unavailable for research summary") from exc
950
951
952 @app.post("/api/v1/research/structured-market/snapshot")
953 async def structured_market_snapshot(instrument: dict = Body(...)):
954 """Internal provider-neutral structured snapshot used by research and manual price refresh."""
955 try:
956 return await portfolio_orchestrator.structured_quote(instrument)
957 except StructuredProviderError as exc:
958 raise HTTPException(status_code=422, detail=str(exc)) from exc
959
960
961 @app.post("/internal/v1/research/instruments/resolve-provider")
962 async def resolve_provider_identity(instrument: dict = Body(...), x_aip_service_identity: str | None = Header(default=None)):
963 """Identity-only reconciliation. Never collects or persists research data."""
964 if x_aip_service_identity != "portfolio-service":
965 raise HTTPException(status_code=403, detail="INTERNAL_SERVICE_REQUIRED")
966 if (instrument.get("structuredNseCandidateSource") != "VERIFIED_NSE"
967 or not instrument.get("isin") or instrument.get("assetType") != "EQUITY"):
968 raise HTTPException(status_code=422, detail="VERIFIED_NSE_IDENTITY_REQUIRED")
969 try:
970 resolution = await portfolio_orchestrator.structured_provider.resolve_instrument(instrument)
971 return resolution.model_dump(mode="json")
972 except StructuredProviderError as exc:
973 raise HTTPException(status_code=503 if "UNAVAILABLE" in str(exc) else 422, detail=str(exc)) from exc
974
975
976 def _require_profile(instrument_id: UUID) -> None:
977 try:
978 repository.profile(instrument_id)
979 except StopIteration as exc:
980 raise HTTPException(status_code=404, detail="Research profile not found") from exc
981
982
983 async def _require_profile_or_restored_instrument(instrument_id: UUID, correlation_id: str | None = None, identity_headers: dict[str, str | None] | None = None) -> None:
984 try:
985 repository.profile(instrument_id)
986 except StopIteration as exc:
987 if await _resolve_restored_instrument_profile(instrument_id, correlation_id, identity_headers):
988 return
989 raise HTTPException(status_code=404, detail="COMPANY_NOT_RESOLVED") from exc
990
991
992 async def _readiness_profile(
993 global_instrument_id: UUID,
994 correlation_id: str | None,
995 identity_headers: dict[str, str | None],
996 ):
997 """Resolve canonical public identity without reconciliation or holding writes."""
998 try:
999 metadata = await portfolio_orchestrator.global_instrument_metadata(
1000 global_instrument_id,
1001 correlation_id=correlation_id,
1002 identity_headers=identity_headers,
1003 )
1004 if not portfolio_orchestrator.register_global_profile_metadata(
1005 global_instrument_id, metadata
1006 ):
1007 raise HTTPException(status_code=404, detail="COMPANY_NOT_RESOLVED")
1008 profile = repository.profile(global_instrument_id)
1009 except (GlobalInstrumentNotFoundError, StopIteration) as exc:
1010 raise HTTPException(status_code=404, detail="COMPANY_NOT_RESOLVED") from exc
1011 except PortfolioServiceUnavailableError as exc:
1012 raise HTTPException(
1013 status_code=502,
1014 detail="Portfolio service unavailable for global instrument lookup",
1015 ) from exc
1016 research_readiness_adapter.remember_canonical_metadata(
1017 global_instrument_id, metadata
1018 )
1019 return profile
1020
1021
1022 async def _resolve_restored_instrument_profile(instrument_id: UUID, correlation_id: str | None = None, identity_headers: dict[str, str | None] | None = None) -> bool:
1023 try:
1024 return await portfolio_orchestrator.restore_global_profile(
1025 instrument_id, correlation_id=correlation_id, identity_headers=identity_headers
1026 )
1027 except GlobalInstrumentNotFoundError:
1028 return False
1029 except PortfolioServiceUnavailableError as exc:
1030 raise HTTPException(status_code=502, detail="Portfolio service unavailable for global instrument lookup") from exc
1031
1032
1033 async def _is_restored_etf_instrument(instrument_id: UUID, correlation_id: str | None = None, identity_headers: dict[str, str | None] | None = None) -> bool:
1034 try:
1035 repository.profile(instrument_id)
1036 return False
1037 except StopIteration:
1038 if not await _resolve_restored_instrument_profile(instrument_id, correlation_id, identity_headers):
1039 return False
1040 try:
1041 repository.etf_profile(instrument_id)
1042 return True
1043 except StopIteration:
1044 return False
1045
1046
1047 def _identity_headers(user_id, issuer, subject, email, display_name, roles) -> dict[str, str | None]:
1048 return {"X-AIP-User-Id": user_id, "X-AIP-User-Issuer": issuer, "X-AIP-User-Subject": subject,
1049 "X-AIP-User-Email": email, "X-AIP-User-Display-Name": display_name, "X-AIP-User-Roles": roles}
1050
1051
1052 def _require_market_data_admin(
1053 user_id: str | None,
1054 issuer: str | None,
1055 subject: str | None,
1056 roles: str | None,
1057 ) -> None:
1058 if not (user_id and issuer and subject):
1059 raise HTTPException(status_code=401, detail="Market data population access denied")
1060 role_set = {value.strip().upper() for value in str(roles or "").split(",") if value.strip()}
1061 if "ADMIN" not in role_set:
1062 raise HTTPException(status_code=403, detail="Market data population admin role required")
1063
1064
1065 def _require_market_data_user(user_id: str | None, issuer: str | None, subject: str | None) -> None:
1066 if not (user_id and issuer and subject):
1067 raise HTTPException(status_code=401, detail="Market data ensure access denied")
1068
1069
1070 def _require_research_user(user_id: str | None, issuer: str | None, subject: str | None) -> None:
1071 if not (user_id and issuer and subject):
1072 raise HTTPException(status_code=401, detail="Research ensure access denied")