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