main
py 147 lines 9.04 KB
Raw
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()]