feat: add global opportunity orchestration
prakhar82 committed
Sep 14, 2026 at 00:10 UTC
fc3df910e3d99a25e4f16fe06627880a00509b1a
2 files changed
+350
ai/research-engine/app/global_opportunity_orchestration.py
new
+176
@@ -0,0 +1,176 @@
1
+"""Explicit engineering ranking over already-read canonical and persisted evidence.
2
+
3
+No HTTP universe adapter, readiness ensure, provider refresh, or membership reads.
4
+Pass PortfolioResearchOrchestrator.register_global_profile_metadata as the
5
+profile_hydrator: this existing synchronous method consumes supplied metadata
6
+only. The repository must already have hydrated its persisted public evidence
7
+(ResearchRepository does this at initialization). Existing V1 cache writes are
8
+preserved; this operation does not persist ranking snapshots.
9
+
10
+Use a fixed timezone-aware as_of to reproduce evidence selection/fingerprints.
11
+generated_at is the completion clock; cache_hit diagnostics may change on a warm
12
+run, but ranking content does not. All limits are explicit and at most 100.
13
+"""
14
+from __future__ import annotations
15
+
16
+from copy import deepcopy
17
+from datetime import datetime, timezone
18
+from typing import Callable, Iterable, Mapping
19
+from uuid import UUID
20
+from pydantic import Field
21
+
22
+from app.global_scanner import GlobalScanner
23
+from app.global_opportunity_ranker import GlobalOpportunityRanker
24
+from app.models import ResearchBaseModel
25
+from app.research_readiness_runtime import RepositoryResearchReadinessAdapter, ResearchReadinessRuntime, jurisdiction_for_profile
26
+from app.sector_relative_strength import SectorContext
27
+from app.stock_rule_engine import StockRuleEngineService
28
+
29
+
30
+class OpportunityEntry(ResearchBaseModel):
31
+ rank: int
32
+ global_instrument_id: UUID
33
+ symbol: str | None
34
+ company_name: str | None
35
+ sector: str | None
36
+ market: str | None
37
+ pre_score: float | None
38
+ stage_b_score: float | None
39
+ technical_score: float | None
40
+ technical_state: str
41
+ sector_score: float | None
42
+ sector_state: str
43
+ rule_engine_score: float | None
44
+ rule_engine_confidence: float
45
+ opportunity_score: float | None
46
+ opportunity_confidence: float
47
+ score_coverage: float
48
+ top_positive_reasons: list[str]
49
+ top_negative_reasons: list[str]
50
+ rule_engine_version: str
51
+ technical_feature_version: str
52
+ sector_feature_version: str
53
+ opportunity_ranker_version: str
54
+
55
+
56
+class CandidateDiagnostic(ResearchBaseModel):
57
+ global_instrument_id: UUID
58
+ status: str
59
+ failure_reason: str | None = None
60
+ cache_hit: bool | None = None
61
+ rank_eligible: bool = False
62
+ suppression_reasons: list[str] = Field(default_factory=list)
63
+
64
+
65
+class OpportunityRanking(ResearchBaseModel):
66
+ generated_at: datetime
67
+ as_of: datetime
68
+ universe_count: int
69
+ phase1_eligible_count: int
70
+ stage_b_count: int
71
+ shortlist_count: int
72
+ deep_evaluated_count: int
73
+ rank_eligible_count: int
74
+ top_n: list[OpportunityEntry]
75
+ diagnostics: list[CandidateDiagnostic]
76
+
77
+
78
+class _PersistedUniverse:
79
+ """The scanner's existing universe protocol, backed solely by supplied rows."""
80
+ def __init__(self, rows):
81
+ self.rows = rows
82
+
83
+ async def active_global_equities(self, **unused):
84
+ return self.rows
85
+
86
+
87
+class GlobalOpportunityOrchestrator:
88
+ def __init__(self, repository, persistence, *, profile_hydrator: Callable[[UUID, dict], bool], clock=None):
89
+ self.repository, self.persistence = repository, persistence
90
+ self.profile_hydrator = profile_hydrator
91
+ self.clock = clock or (lambda: datetime.now(timezone.utc))
92
+ self.readiness_adapter = RepositoryResearchReadinessAdapter(repository)
93
+ self.readiness = ResearchReadinessRuntime(repository, self.readiness_adapter, executor=None)
94
+ self.rule_engine = StockRuleEngineService(repository, self.readiness_adapter)
95
+ self.ranker = GlobalOpportunityRanker()
96
+
97
+ async def run(self, canonical_instruments: Iterable[dict], *, as_of: datetime,
98
+ sector_contexts: Mapping[UUID, SectorContext] | None = None,
99
+ shortlist_limit: int = 25, top_n: int = 10) -> OpportunityRanking:
100
+ if (type(shortlist_limit) is not int or not 1 <= shortlist_limit <= 100
101
+ or type(top_n) is not int or not 0 <= top_n <= 100):
102
+ raise ValueError('INVALID_OPPORTUNITY_LIMIT')
103
+ if as_of.tzinfo is None or as_of.utcoffset() is None:
104
+ raise ValueError('AWARE_AS_OF_REQUIRED')
105
+ as_of = as_of.astimezone(timezone.utc)
106
+ rows = deepcopy(list(canonical_instruments))
107
+ # Support both existing canonical enumeration and master-detail shapes.
108
+ for row in rows:
109
+ row.setdefault('exchange', row.get('primaryExchange'))
110
+ row.setdefault('ticker', row.get('primarySymbol'))
111
+ scanner = GlobalScanner(_PersistedUniverse(rows), self.persistence)
112
+ scan = await scanner.scan(as_of=as_of, top_n=0)
113
+ phase1 = {c.global_instrument_id:c for c in scan.candidates if c.eligible_for_deep_analysis}
114
+ stage_b = scanner.enrich_candidates(scan, sector_contexts=dict(sector_contexts or {}))
115
+ def desc(value): return -value if value is not None else float('inf')
116
+ shortlist = sorted((c for c in stage_b if c.global_instrument_id in phase1), key=lambda c: (
117
+ desc(c.stage_b_score), -c.confidence, desc(phase1[c.global_instrument_id].pre_score),
118
+ -phase1[c.global_instrument_id].confidence, str(c.global_instrument_id)))[:shortlist_limit]
119
+ metadata = {str(row.get('globalInstrumentId')):row for row in rows}
120
+ diagnostics, rules, successful = [], {}, []
121
+ evaluated = 0
122
+ for candidate in shortlist:
123
+ key = candidate.global_instrument_id
124
+ step = 'PUBLIC_EVIDENCE_UNAVAILABLE'
125
+ cache_hit = None
126
+ try:
127
+ payload = dict(metadata[str(key)])
128
+ payload.setdefault('primaryExchange', payload.get('exchange'))
129
+ payload.setdefault('primarySymbol', payload.get('ticker'))
130
+ if not self.profile_hydrator(key, payload):
131
+ raise ValueError('UNRESOLVED_PROFILE')
132
+ profile = self.repository.profile(key)
133
+ if profile.instrument_id != key:
134
+ raise ValueError('PROFILE_IDENTITY_MISMATCH')
135
+ self.readiness_adapter.remember_canonical_metadata(key, payload)
136
+ step = 'READINESS_UNAVAILABLE'
137
+ readiness = await self.readiness.read(key, jurisdiction=jurisdiction_for_profile(profile), now=as_of)
138
+ step = 'RULE_ENGINE_UNAVAILABLE'
139
+ result = await self.rule_engine.analyze(profile, readiness, allow_partial=False, now=as_of)
140
+ evaluated += 1
141
+ cache_hit = result.cache_hit
142
+ step = 'RANKER_INPUT_UNAVAILABLE'
143
+ ranked = self.ranker.score(candidate, result)
144
+ diagnostics.append(CandidateDiagnostic(global_instrument_id=key,
145
+ status='RANK_ELIGIBLE' if ranked.rank_eligible else 'SUPPRESSED', cache_hit=cache_hit,
146
+ rank_eligible=ranked.rank_eligible, suppression_reasons=ranked.eligibility_reasons))
147
+ if ranked.rank_eligible:
148
+ rules[key] = result
149
+ successful.append(candidate)
150
+ except Exception:
151
+ # Never return exception messages, headers, raw evidence or recommendations.
152
+ diagnostics.append(CandidateDiagnostic(global_instrument_id=key, status='FAILED',
153
+ failure_reason=step, cache_hit=cache_hit))
154
+ ranked = self.ranker.rank(successful, rules)
155
+ by_id = {c.global_instrument_id:c for c in successful}
156
+ entries = []
157
+ for position, result in enumerate(ranked[:top_n], 1):
158
+ key = result.global_instrument_id
159
+ initial, enriched, rule = phase1[key], by_id[key], rules[key]
160
+ entries.append(OpportunityEntry(rank=position, global_instrument_id=key,
161
+ symbol=initial.symbol, company_name=initial.company_name,
162
+ sector=enriched.sector_relative_strength_snapshot.sector, market=initial.market,
163
+ pre_score=initial.pre_score, stage_b_score=enriched.stage_b_score,
164
+ technical_score=result.technical_score, technical_state=result.technical_state,
165
+ sector_score=result.sector_score, sector_state=result.sector_state,
166
+ rule_engine_score=result.rule_engine_score, rule_engine_confidence=rule.confidence_score,
167
+ opportunity_score=result.opportunity_score, opportunity_confidence=result.opportunity_confidence,
168
+ score_coverage=result.score_coverage, top_positive_reasons=result.top_positive_reasons,
169
+ top_negative_reasons=result.top_negative_reasons, rule_engine_version=rule.rule_engine_version,
170
+ technical_feature_version=enriched.technical_feature_snapshot.feature_version,
171
+ sector_feature_version=enriched.sector_relative_strength_snapshot.feature_version,
172
+ opportunity_ranker_version=result.ranker_version))
173
+ return OpportunityRanking(generated_at=self.clock(), as_of=as_of,
174
+ universe_count=scan.total_canonical_active_equities, phase1_eligible_count=len(phase1),
175
+ stage_b_count=len(stage_b), shortlist_count=len(shortlist), deep_evaluated_count=evaluated,
176
+ rank_eligible_count=len(ranked), top_n=entries, diagnostics=diagnostics)
ai/research-engine/tests/test_global_opportunity_orchestration.py
new
+174
@@ -0,0 +1,174 @@
1
+from datetime import timedelta
2
+from unittest.mock import AsyncMock
3
+from uuid import UUID
4
+
5
+import httpx
6
+import pytest
7
+
8
+from app.global_opportunity_orchestration import GlobalOpportunityOrchestrator
9
+from app.global_scanner import GlobalScanner
10
+from app.persistence import SqliteResearchPersistence
11
+from app.portfolio_orchestration import PortfolioResearchOrchestrator
12
+from app.repository import ResearchRepository
13
+from app.settings import Settings
14
+from test_global_scanner import instrument, persisted, NOW
15
+from test_global_opportunity_ranker import inputs
16
+
17
+
18
+def setup(monkeypatch, count=3, *, stub_deep=True):
19
+ store = SqliteResearchPersistence()
20
+ repo = ResearchRepository(persistence=store)
21
+ hydrator = PortfolioResearchOrchestrator(repo, Settings(), client=object())
22
+ service = GlobalOpportunityOrchestrator(repo, store,
23
+ profile_hydrator=hydrator.register_global_profile_metadata, clock=lambda:NOW)
24
+ rows = [instrument(n) for n in range(1,count+1)]
25
+ for row in rows: persisted(store, row)
26
+ pairs = {UUID(int=n):inputs(n) for n in range(1,count+1)}
27
+ monkeypatch.setattr(GlobalScanner, 'enrich_candidates', lambda self, scan, **kwargs:
28
+ [pairs[c.global_instrument_id][0] for c in reversed(scan.candidates) if c.eligible_for_deep_analysis])
29
+ if stub_deep:
30
+ service.readiness.read = AsyncMock(return_value=object())
31
+ async def analyze(profile, readiness, **kwargs):
32
+ assert kwargs == dict(allow_partial=False, now=NOW)
33
+ return pairs[profile.instrument_id][1]
34
+ service.rule_engine.analyze = AsyncMock(side_effect=analyze)
35
+ return service, rows, pairs, store
36
+
37
+
38
+@pytest.mark.asyncio
39
+async def test_shortlist_limit_order_and_v1_only_shortlist(monkeypatch):
40
+ service, rows, pairs, _ = setup(monkeypatch, 4)
41
+ pairs[UUID(int=3)][0].stage_b_score = 100
42
+ result = await service.run(reversed(rows), as_of=NOW, shortlist_limit=2)
43
+ assert [d.global_instrument_id.int for d in result.diagnostics] == [3,1]
44
+ assert result.universe_count == result.phase1_eligible_count == result.stage_b_count == 4
45
+ assert result.shortlist_count == result.deep_evaluated_count == 2
46
+ assert service.rule_engine.analyze.await_count == 2
47
+ # Final opportunity rank uses the ranker, not shortlist order.
48
+ assert [entry.global_instrument_id.int for entry in result.top_n] == [1,3]
49
+
50
+
51
+@pytest.mark.asyncio
52
+async def test_sector_unavailable_allowed_and_fewer_than_top_n(monkeypatch):
53
+ service, rows, pairs, _ = setup(monkeypatch, 1)
54
+ stage = pairs[UUID(int=1)][0]
55
+ stage.sector_score = stage.sector_relative_strength_snapshot.relative_strength_score = None
56
+ stage.sector_relative_strength_snapshot.sector_state = 'INSUFFICIENT_DATA'
57
+ result = await service.run(rows, as_of=NOW, top_n=10)
58
+ assert result.rank_eligible_count == len(result.top_n) == 1
59
+ assert result.top_n[0].sector_score is None
60
+ assert result.top_n[0].score_coverage == 97
61
+
62
+
63
+@pytest.mark.asyncio
64
+async def test_failure_isolated_safe_and_suppressed_not_top_n(monkeypatch):
65
+ service, rows, pairs, _ = setup(monkeypatch)
66
+ pairs[UUID(int=2)][1].partial = True
67
+ async def analyze(profile, readiness, **kwargs):
68
+ if profile.instrument_id.int == 1: raise RuntimeError('secret-cookie/password')
69
+ return pairs[profile.instrument_id][1]
70
+ service.rule_engine.analyze.side_effect = analyze
71
+ result = await service.run(rows, as_of=NOW)
72
+ assert [d.status for d in result.diagnostics] == ['FAILED','SUPPRESSED','RANK_ELIGIBLE']
73
+ assert result.diagnostics[0].failure_reason == 'RULE_ENGINE_UNAVAILABLE'
74
+ assert result.deep_evaluated_count == 2 and result.rank_eligible_count == 1
75
+ assert [e.global_instrument_id.int for e in result.top_n] == [3]
76
+ assert 'secret' not in result.model_dump_json()
77
+
78
+
79
+@pytest.mark.asyncio
80
+async def test_repeat_membership_independence_uuid_ties_and_no_network(monkeypatch):
81
+ service, rows, pairs, _ = setup(monkeypatch)
82
+ def forbidden(*a, **k): pytest.fail('Provider/network/membership acquisition attempted')
83
+ import socket
84
+ import yfinance
85
+ from app.nse_historical_daily import NseHistoricalDailyProvider
86
+ from app.nse_index_history import NseIndexHistoryProvider
87
+ from app.market_data_population import IndiaMarketDataPopulationJobs
88
+ monkeypatch.setattr(socket, 'create_connection', forbidden)
89
+ monkeypatch.setattr(httpx.AsyncClient, 'request', forbidden)
90
+ monkeypatch.setattr(httpx.Client, 'request', forbidden)
91
+ monkeypatch.setattr(yfinance, 'Ticker', forbidden)
92
+ monkeypatch.setattr(NseHistoricalDailyProvider, 'fetch', forbidden)
93
+ monkeypatch.setattr(NseIndexHistoryProvider, 'fetch', forbidden)
94
+ monkeypatch.setattr(IndiaMarketDataPopulationJobs, 'backfill_daily_bars', forbidden)
95
+ service.readiness.ensure = forbidden
96
+ service.repository.holdings = forbidden
97
+ service.repository.watchlist = forbidden
98
+ first = await service.run(rows, as_of=NOW, top_n=2)
99
+ service.clock = lambda: NOW+timedelta(seconds=1)
100
+ second = await service.run([r | dict(portfolioId='irrelevant',watchlistMember=True) for r in reversed(rows)], as_of=NOW, top_n=2)
101
+ assert first.model_dump(exclude={'generated_at'}) == second.model_dump(exclude={'generated_at'})
102
+ assert [r.global_instrument_id.int for r in second.top_n] == [1,2]
103
+ assert [r.rank for r in second.top_n] == [1,2]
104
+ assert second.top_n[0].company_name == 'Company 1'
105
+ assert 'decisionSignal' not in second.model_dump_json(by_alias=True)
106
+
107
+
108
+@pytest.mark.asyncio
109
+async def test_phase1_exclusion_not_deep_evaluated(monkeypatch):
110
+ service, rows, pairs, store = setup(monkeypatch)
111
+ persisted(store, rows[0], {})
112
+ result = await service.run(rows, as_of=NOW)
113
+ assert result.phase1_eligible_count == result.stage_b_count == 2
114
+ assert all(d.global_instrument_id.int != 1 for d in result.diagnostics)
115
+
116
+
117
+@pytest.mark.asyncio
118
+async def test_empty_universe(monkeypatch):
119
+ service, _, _, store = setup(monkeypatch, 0)
120
+ queries = []
121
+ store._connection.set_trace_callback(queries.append)
122
+ result = await service.run([], as_of=NOW)
123
+ assert not result.top_n and not result.diagnostics and not queries
124
+ assert result.universe_count == result.stage_b_count == result.shortlist_count == 0
125
+ assert service.rule_engine.analyze.await_count == 0
126
+
127
+
128
+@pytest.mark.asyncio
129
+async def test_real_readiness_v1_fingerprint_cache_without_refresh(monkeypatch):
130
+ service, rows, _, store = setup(monkeypatch, 1, stub_deep=False)
131
+ def forbidden(*a, **k): pytest.fail('Network or refresh attempted')
132
+ monkeypatch.setattr(httpx.AsyncClient, 'request', forbidden)
133
+ monkeypatch.setattr(httpx.Client, 'request', forbidden)
134
+ import socket
135
+ monkeypatch.setattr(socket.socket, 'connect', forbidden)
136
+ service.readiness.ensure = forbidden
137
+ first = await service.run(rows, as_of=NOW)
138
+ second = await service.run(rows, as_of=NOW)
139
+ third = await service.run([r | dict(portfolioId='irrelevant',watchlistMember=True) for r in rows], as_of=NOW)
140
+ assert first.deep_evaluated_count == second.deep_evaluated_count == 1
141
+ assert first.diagnostics[0].cache_hit is False and second.diagnostics[0].cache_hit is True
142
+ assert first.top_n == second.top_n and second == third
143
+ assert not first.top_n # Sparse evidence cannot be promoted merely to fill Top-N.
144
+ assert store._connection.execute('SELECT count(*) FROM global_stock_rule_engine_results').fetchone()[0] == 1
145
+
146
+
147
+@pytest.mark.asyncio
148
+@pytest.mark.parametrize('kwargs', [dict(shortlist_limit=0),dict(shortlist_limit=101),dict(top_n=-1),dict(top_n=True)])
149
+async def test_invalid_limits(monkeypatch, kwargs):
150
+ service, rows, _, _ = setup(monkeypatch)
151
+ with pytest.raises(ValueError, match='INVALID_OPPORTUNITY_LIMIT'):
152
+ await service.run(rows, as_of=NOW, **kwargs)
153
+
154
+
155
+@pytest.mark.asyncio
156
+async def test_shortlist_uses_screening_confidence_then_prescore(monkeypatch):
157
+ service, rows, pairs, _ = setup(monkeypatch)
158
+ pairs[UUID(int=2)][0].confidence = 100
159
+ pairs[UUID(int=1)][0].stage_b_score = None
160
+ result = await service.run(rows, as_of=NOW)
161
+ assert [d.global_instrument_id.int for d in result.diagnostics] == [2,3,1]
162
+
163
+
164
+@pytest.mark.asyncio
165
+async def test_exact_top_n_uses_opportunity_order_not_screening_order(monkeypatch):
166
+ service, rows, pairs, _ = setup(monkeypatch)
167
+ for key, score in ((1,40),(2,90),(3,60)):
168
+ for area in pairs[UUID(int=key)][1].area_scores:
169
+ area.raw_score = score
170
+ result = await service.run(rows, as_of=NOW, top_n=2)
171
+ assert [d.global_instrument_id.int for d in result.diagnostics] == [1,2,3]
172
+ assert [e.global_instrument_id.int for e in result.top_n] == [2,3]
173
+ assert result.rank_eligible_count == 3
174
+ assert result.top_n[0].opportunity_score > result.top_n[1].opportunity_score