| 1 | """V13 projection storage. All cycle writes publish atomically on the shared DB. |
| 2 | |
| 3 | History and snapshot methods only INSERT; projection methods alone use UPDATE. |
| 4 | The payload preserves complete evidence alongside queryable identity/version columns. |
| 5 | """ |
| 6 | import json |
| 7 | |
| 8 | SNAPSHOT_FIELDS = { |
| 9 | **dict.fromkeys('scanner_version rule_engine_version technical_version sector_version ranker_version price_as_of'.split(), 'TEXT'), |
| 10 | **dict.fromkeys(('opportunity_score opportunity_confidence score_coverage rule_engine_score rank_position ' |
| 11 | 'technical_score sector_score valuation_score quality_score growth_score balance_sheet_score ' |
| 12 | 'quarterly_score catalyst_score news_score shareholding_score governance_score current_price').split(), 'REAL'), |
| 13 | 'rank_eligible': 'INTEGER', |
| 14 | **dict.fromkeys('suppression_reasons top_positive_reasons top_negative_reasons evidence_state'.split(), 'JSON')} |
| 15 | HISTORY_FIELDS = { |
| 16 | **dict.fromkeys('rule_engine_version ranker_version new_investor_action existing_holder_action short_term_action long_term_action'.split(), 'TEXT'), |
| 17 | **dict.fromkeys(('price_at_recommendation opportunity_score confidence coverage short_entry_low short_entry_high ' |
| 18 | 'short_target_1 short_target_2 short_invalidation long_entry_low long_entry_high ' |
| 19 | 'long_fair_value long_target long_invalidation').split(), 'REAL'), |
| 20 | **dict.fromkeys('top_positive_reasons top_negative_reasons evidence_snapshot'.split(), 'JSON')} |
| 21 | STATE_FIELDS = { |
| 22 | **dict.fromkeys('short_term_state long_term_state current_short_action current_long_action lifecycle_status'.split(), 'TEXT'), |
| 23 | **dict.fromkeys('price_at_recommendation current_price short_target_distance_pct long_target_distance_pct'.split(), 'REAL')} |
| 24 | EXTRA_FIELDS = {'global_opportunity_snapshot': SNAPSHOT_FIELDS, 'stock_recommendation_history': HISTORY_FIELDS, |
| 25 | 'recommendation_current_state': STATE_FIELDS} |
| 26 | |
| 27 | SCHEMA = ''' |
| 28 | CREATE TABLE IF NOT EXISTS global_opportunity_snapshot ( |
| 29 | snapshot_id TEXT PRIMARY KEY, cycle_id TEXT NOT NULL, global_instrument_id TEXT NOT NULL, |
| 30 | market TEXT NOT NULL, generated_at TEXT NOT NULL, payload TEXT NOT NULL, |
| 31 | UNIQUE(cycle_id, global_instrument_id) |
| 32 | ); |
| 33 | CREATE INDEX IF NOT EXISTS opportunity_snapshot_instrument ON global_opportunity_snapshot(global_instrument_id, generated_at); |
| 34 | CREATE TABLE IF NOT EXISTS stock_recommendation_history ( |
| 35 | recommendation_id TEXT PRIMARY KEY, snapshot_id TEXT NOT NULL REFERENCES global_opportunity_snapshot(snapshot_id), |
| 36 | global_instrument_id TEXT NOT NULL, generated_at TEXT NOT NULL, |
| 37 | recommendation_engine_version TEXT NOT NULL, fingerprint TEXT NOT NULL, payload TEXT NOT NULL |
| 38 | ); |
| 39 | CREATE INDEX IF NOT EXISTS recommendation_history_instrument ON stock_recommendation_history(global_instrument_id, generated_at); |
| 40 | CREATE TABLE IF NOT EXISTS recommendation_current_state ( |
| 41 | global_instrument_id TEXT PRIMARY KEY, |
| 42 | latest_recommendation_id TEXT NOT NULL REFERENCES stock_recommendation_history(recommendation_id), |
| 43 | updated_at TEXT NOT NULL, payload TEXT NOT NULL |
| 44 | ); |
| 45 | CREATE TABLE IF NOT EXISTS global_opportunity_top_selection ( |
| 46 | cycle_id TEXT PRIMARY KEY, generated_at TEXT NOT NULL, market TEXT NOT NULL, payload TEXT NOT NULL |
| 47 | ); |
| 48 | CREATE TABLE IF NOT EXISTS recommendation_backtest_run ( |
| 49 | backtest_id TEXT PRIMARY KEY, generated_at TEXT NOT NULL, payload TEXT NOT NULL |
| 50 | ); |
| 51 | ''' |
| 52 | |
| 53 | for _table, _fields in EXTRA_FIELDS.items(): |
| 54 | _start = SCHEMA.index('CREATE TABLE IF NOT EXISTS ' + _table) |
| 55 | _payload = SCHEMA.index('payload TEXT NOT NULL', _start) |
| 56 | SCHEMA = SCHEMA[:_payload] + ', '.join(k + ' ' + ('TEXT' if v == 'JSON' else v) for k, v in _fields.items()) + ', ' + SCHEMA[_payload:] |
| 57 | |
| 58 | for _table in ('global_opportunity_snapshot', 'stock_recommendation_history', 'global_opportunity_top_selection', 'recommendation_backtest_run'): |
| 59 | for _operation in ('UPDATE', 'DELETE'): |
| 60 | SCHEMA += f'''CREATE TRIGGER IF NOT EXISTS immutable_{_table}_{_operation.lower()} |
| 61 | BEFORE {_operation} ON {_table} BEGIN SELECT RAISE(ABORT, 'RECOMMENDATION_HISTORY_IS_IMMUTABLE'); END;\n''' |
| 62 | |
| 63 | |
| 64 | def decode(row): |
| 65 | if row is None: |
| 66 | return None |
| 67 | value = row['payload'] |
| 68 | return json.loads(value) if isinstance(value, str) else value |
| 69 | |
| 70 | |
| 71 | class OpportunityPersistenceMixin: |
| 72 | def recommendation_history(self, instrument_id=None): |
| 73 | sql = 'SELECT payload FROM stock_recommendation_history' |
| 74 | params = () |
| 75 | if instrument_id is not None: |
| 76 | sql += ' WHERE global_instrument_id = ?' |
| 77 | params = (str(instrument_id),) |
| 78 | return [decode(r) for r in self._connection.execute(sql + ' ORDER BY generated_at, recommendation_id', params).fetchall()] |
| 79 | |
| 80 | def recommendation_states(self): |
| 81 | return [decode(r) for r in self._connection.execute( |
| 82 | 'SELECT payload FROM recommendation_current_state ORDER BY global_instrument_id').fetchall()] |
| 83 | |
| 84 | def opportunity_current(self): |
| 85 | selection = decode(self._connection.execute( |
| 86 | 'SELECT payload FROM global_opportunity_top_selection ORDER BY generated_at DESC, cycle_id DESC LIMIT 1').fetchone()) |
| 87 | if selection is None: |
| 88 | return dict(generated_at=None, best_buy_today=None, top_short_term=[], top_long_term=[], previous_recommendations=[]) |
| 89 | snapshots = {s['global_instrument_id']: s for s in self.opportunity_snapshots(selection['cycle_id'])} |
| 90 | states = {s['global_instrument_id']: s for s in self.recommendation_states()} |
| 91 | def hydrate(card): |
| 92 | key = card['global_instrument_id'] |
| 93 | state = states.get(key, {}) |
| 94 | latest_id = state.get('latest_recommendation_id') |
| 95 | history = decode(self._connection.execute( |
| 96 | 'SELECT payload FROM stock_recommendation_history WHERE recommendation_id = ?', (latest_id,)).fetchone()) if latest_id else None |
| 97 | # Current public score/price come from the snapshot/projection, not an older recommendation. |
| 98 | current = snapshots.get(key, {}) |
| 99 | scores = {k: current[k] for k in ('opportunity_score', 'opportunity_confidence', 'score_coverage', 'rule_engine_score') if k in current} |
| 100 | return {**card, **current, **(history or {}), **scores, **state} |
| 101 | for field in ('top_short_term', 'top_long_term', 'previous_recommendations'): |
| 102 | selection[field] = [hydrate(c) for c in selection[field]] |
| 103 | if selection['best_buy_today']: |
| 104 | selection['best_buy_today'] = hydrate(selection['best_buy_today']) |
| 105 | return selection |
| 106 | |
| 107 | def opportunity_snapshots(self, cycle_id): |
| 108 | return [decode(r) for r in self._connection.execute( |
| 109 | 'SELECT payload FROM global_opportunity_snapshot WHERE cycle_id = ? ORDER BY global_instrument_id', (cycle_id,)).fetchall()] |
| 110 | |
| 111 | def _insert_opportunity(self, table, columns, value): |
| 112 | # Table/column names are internal constants, never caller SQL. |
| 113 | extra = EXTRA_FIELDS.get(table, {}) |
| 114 | columns = (*columns, *extra) |
| 115 | def stored(c): |
| 116 | v = value.get(c) |
| 117 | return json.dumps(v, sort_keys=True) if extra.get(c) == 'JSON' else int(v) if isinstance(v, bool) else v |
| 118 | self._connection.execute(f"INSERT INTO {table} ({','.join(columns)},payload) VALUES ({','.join('?' for _ in range(len(columns)+1))})", |
| 119 | tuple(stored(c) for c in columns) + (json.dumps(value, sort_keys=True, allow_nan=False),)) |
| 120 | |
| 121 | def publish_opportunity_cycle(self, snapshots, recommendations, states, selection, *, build=None): |
| 122 | with self._connection: |
| 123 | for s in snapshots: |
| 124 | self._insert_opportunity('global_opportunity_snapshot', |
| 125 | ('snapshot_id', 'cycle_id', 'global_instrument_id', 'market', 'generated_at'), s) |
| 126 | if build is not None: |
| 127 | persisted = self.opportunity_snapshots(selection['cycle_id']) |
| 128 | recommendations, states, selection = build(persisted) |
| 129 | for r in recommendations: |
| 130 | self._insert_opportunity('stock_recommendation_history', |
| 131 | ('recommendation_id', 'snapshot_id', 'global_instrument_id', 'generated_at', 'recommendation_engine_version', 'fingerprint'), r) |
| 132 | for s in states: |
| 133 | cols = ('global_instrument_id', 'latest_recommendation_id', 'updated_at', *STATE_FIELDS, 'payload') |
| 134 | self._connection.execute(f"INSERT INTO recommendation_current_state ({','.join(cols)}) VALUES ({','.join('?' for _ in cols)}) " |
| 135 | + 'ON CONFLICT(global_instrument_id) DO UPDATE SET ' + ','.join(f'{c}=excluded.{c}' for c in cols[1:]), |
| 136 | tuple(s.get(c) for c in cols[:-1]) + (json.dumps(s, sort_keys=True),)) |
| 137 | self._insert_opportunity('global_opportunity_top_selection', ('cycle_id', 'generated_at', 'market'), selection) |
| 138 | return selection |
| 139 | |
| 140 | def save_backtest(self, result): |
| 141 | with self._connection: |
| 142 | self._insert_opportunity('recommendation_backtest_run', ('backtest_id', 'generated_at'), result) |
| 143 | return result |
| 144 | |
| 145 | def backtests(self): |
| 146 | return [decode(r) for r in self._connection.execute( |
| 147 | 'SELECT payload FROM recommendation_backtest_run ORDER BY generated_at DESC, backtest_id DESC').fetchall()] |