"""Helpers for single-day SPBT result analysis notebooks.""" from __future__ import annotations import json import math import sqlite3 from typing import Any import pandas as pd SELECTOR_PAIRS_COLUMNS = ("pair_name", "mr_score") TRADING_INSTRUCTIONS_COLUMNS = ("time_ns", "tstamp", "data") INITIAL_THEO_CAPITAL_USD = 10_000.0 def parse_mr_score_final(raw_score: Any) -> tuple[float | None, str]: """Parse the JSON mr_score.final value, preserving parse status.""" if raw_score is None: return None, "missing_mr_score" try: parsed = json.loads(raw_score) except (TypeError, json.JSONDecodeError): return None, "malformed_json" if not isinstance(parsed, dict): return None, "unexpected_json_type" if "final" not in parsed or parsed["final"] in (None, ""): return None, "missing_final" final_value = parsed["final"] if isinstance(final_value, bool): return None, "non_numeric_final" try: numeric_final = float(final_value) except (TypeError, ValueError): return None, "non_numeric_final" if not math.isfinite(numeric_final): return None, "non_finite_final" return numeric_final, "ok" def validate_selector_pairs_table(conn: sqlite3.Connection) -> None: """Raise an actionable error if selector_pairs lacks required columns.""" table_info = conn.execute("PRAGMA table_info(selector_pairs)").fetchall() if not table_info: raise ValueError("SQLite database is missing required table: selector_pairs") existing_columns = {row[1] for row in table_info} missing_columns = set(SELECTOR_PAIRS_COLUMNS) - existing_columns if missing_columns: missing = ", ".join(sorted(missing_columns)) raise ValueError(f"selector_pairs is missing required column(s): {missing}") def rank_selector_pairs(selector_pairs: pd.DataFrame) -> pd.DataFrame: """Rank selector pairs by dense descending mr_score.final.""" missing_columns = set(SELECTOR_PAIRS_COLUMNS) - set(selector_pairs.columns) if missing_columns: missing = ", ".join(sorted(missing_columns)) raise ValueError(f"selector_pairs dataframe is missing column(s): {missing}") ranked = selector_pairs.loc[:, list(SELECTOR_PAIRS_COLUMNS)].copy() parsed_scores = ranked["mr_score"].map(parse_mr_score_final) ranked["mr_score_final"] = [score for score, _status in parsed_scores] ranked["mr_score_parse_status"] = [status for _score, status in parsed_scores] ranked["pair_rank"] = ( ranked["mr_score_final"].rank(method="dense", ascending=False).astype("Int64") ) return ( ranked.loc[ :, [ "pair_rank", "pair_name", "mr_score_final", "mr_score_parse_status", "mr_score", ], ] .sort_values(["pair_rank", "pair_name"], na_position="last", kind="mergesort") .reset_index(drop=True) ) def load_selector_pair_rankings(conn: sqlite3.Connection) -> pd.DataFrame: """Load selector_pairs from SQLite and return dense-ranked pair rows.""" validate_selector_pairs_table(conn) selector_pairs = pd.read_sql_query( "SELECT pair_name, mr_score FROM selector_pairs", conn, ) return rank_selector_pairs(selector_pairs) def validate_trading_instructions_table(conn: sqlite3.Connection) -> None: """Raise an actionable error if trading_instructions lacks required columns.""" table_info = conn.execute("PRAGMA table_info(trading_instructions)").fetchall() if not table_info: raise ValueError("SQLite database is missing required table: trading_instructions") existing_columns = {row[1] for row in table_info} missing_columns = set(TRADING_INSTRUCTIONS_COLUMNS) - existing_columns if missing_columns: missing = ", ".join(sorted(missing_columns)) raise ValueError( f"trading_instructions is missing required column(s): {missing}" ) def load_trading_instructions(conn: sqlite3.Connection) -> pd.DataFrame: """Load the full trading_instructions table ordered by timestamp.""" validate_trading_instructions_table(conn) return pd.read_sql_query( "SELECT * FROM trading_instructions ORDER BY time_ns, rowid", conn, ) def pair_assets_and_quote(pair_name: str) -> tuple[tuple[str, ...], str]: """Parse a pair name like ADA:USD-BTC:USD into assets and quote asset.""" pair_legs = pair_name.split("-") if len(pair_legs) != 2: raise ValueError(f"Pair name must contain exactly two legs: {pair_name}") assets: list[str] = [] quote_assets: list[str] = [] for leg in pair_legs: parts = leg.split(":") if len(parts) != 2 or not all(parts): raise ValueError(f"Pair leg must use ASSET:QUOTE form: {leg}") assets.append(parts[0]) quote_assets.append(parts[1]) distinct_quote_assets = set(quote_assets) if len(distinct_quote_assets) != 1: raise ValueError(f"Pair legs must use the same quote asset: {pair_name}") return tuple(assets), quote_assets[0] def _parse_instruction_data(raw_data: Any) -> dict[str, Any] | None: if raw_data is None: return None try: parsed = json.loads(raw_data) except (TypeError, json.JSONDecodeError): return None return parsed if isinstance(parsed, dict) else None def _finite_float(value: Any, field_name: str, pair_name: str) -> float: if isinstance(value, bool): raise ValueError(f"{field_name} for {pair_name} must be numeric, got bool") try: numeric_value = float(value) except (TypeError, ValueError) as exc: raise ValueError( f"{field_name} for {pair_name} must be numeric, got {value!r}" ) from exc if not math.isfinite(numeric_value): raise ValueError(f"{field_name} for {pair_name} must be finite") return numeric_value def _reference_price(asset_data: Any, asset: str, pair_name: str) -> float: if not isinstance(asset_data, dict): raise ValueError(f"Asset data for {asset} in {pair_name} must be a JSON object") reference_price = _finite_float( asset_data.get("reference_price"), f"reference_price[{asset}]", pair_name, ) if reference_price <= 0: raise ValueError(f"reference_price[{asset}] for {pair_name} must be positive") return reference_price def _strength(asset_data: Any, asset: str, pair_name: str) -> float: if not isinstance(asset_data, dict): raise ValueError(f"Asset data for {asset} in {pair_name} must be a JSON object") return _finite_float(asset_data.get("strength"), f"strength[{asset}]", pair_name) def _sort_trading_instructions(trd_inst_df: pd.DataFrame) -> pd.DataFrame: sort_columns = [column for column in ("time_ns", "tstamp") if column in trd_inst_df] ordered = trd_inst_df.copy() ordered["_input_order"] = range(len(ordered)) return ordered.sort_values( [*sort_columns, "_input_order"], kind="mergesort", ).drop(columns="_input_order") def _matching_pair_instructions( pair_name: str, trd_inst_df: pd.DataFrame, ) -> list[dict[str, Any]]: pair_assets, quote_asset = pair_assets_and_quote(pair_name) pair_asset_set = set(pair_assets) if "data" not in trd_inst_df.columns: raise ValueError("trading instructions dataframe is missing column: data") selected_instructions: list[dict[str, Any]] = [] for raw_data in _sort_trading_instructions(trd_inst_df)["data"]: parsed = _parse_instruction_data(raw_data) if parsed is None or parsed.get("quote_asset") != quote_asset: continue assets = parsed.get("assets") if not isinstance(assets, dict) or set(assets) != pair_asset_set: continue selected_instructions.append(parsed) return selected_instructions def calculate_pair_theo_ret( pair_name: str, trd_inst_df: pd.DataFrame, ) -> dict[str, float | str]: """Calculate realized and unrealized TheoRet percentages for one pair. Repeated TARGET actions replace the previous open theoretical position. CLOSE actions liquidate the currently open position. HOLD and unknown actions are ignored. Returned PnL values are percentages of the fixed 10,000 USD theoretical capital base. """ pair_assets, _quote_asset = pair_assets_and_quote(pair_name) realized_pnl_usd = 0.0 open_position: dict[str, dict[str, float]] | None = None for instruction in _matching_pair_instructions(pair_name, trd_inst_df): action = instruction.get("action") assets_data = instruction["assets"] if action == "TARGET": open_position = {} for asset in pair_assets: asset_data = assets_data[asset] quantity = INITIAL_THEO_CAPITAL_USD * _strength( asset_data, asset, pair_name, ) entry_price = _reference_price(asset_data, asset, pair_name) open_position[asset] = { "quantity": quantity, "entry_value": quantity * entry_price, "latest_value": quantity * entry_price, } elif action == "CLOSE": if open_position is None: continue for asset in pair_assets: asset_data = assets_data[asset] close_value = ( open_position[asset]["quantity"] * _reference_price(asset_data, asset, pair_name) ) realized_pnl_usd += close_value - open_position[asset]["entry_value"] open_position = None unrealized_pnl_usd = 0.0 if open_position is not None: unrealized_pnl_usd = sum( position["latest_value"] - position["entry_value"] for position in open_position.values() ) return { "pair_name": pair_name, "realized_pnl": realized_pnl_usd / INITIAL_THEO_CAPITAL_USD * 100.0, "unrealized_pnl": unrealized_pnl_usd / INITIAL_THEO_CAPITAL_USD * 100.0, } def calculate_ranked_pairs_theo_ret( selector_pair_rankings: pd.DataFrame, trd_inst_df: pd.DataFrame, ) -> pd.DataFrame: """Calculate TheoRet percentages for every ranked selector pair.""" required_columns = {"pair_name", "pair_rank"} missing_columns = required_columns - set(selector_pair_rankings.columns) if missing_columns: missing = ", ".join(sorted(missing_columns)) raise ValueError(f"selector pair rankings missing column(s): {missing}") records = [] for row in selector_pair_rankings.itertuples(index=False): pair_result = calculate_pair_theo_ret(row.pair_name, trd_inst_df) records.append( { "pair_name": pair_result["pair_name"], "mr_ranking": row.pair_rank, "realized_pnl": pair_result["realized_pnl"], "unrealized_pnl": pair_result["unrealized_pnl"], } ) return ( pd.DataFrame.from_records( records, columns=["pair_name", "mr_ranking", "realized_pnl", "unrealized_pnl"], ) .sort_values(["mr_ranking", "pair_name"], na_position="last", kind="mergesort") .reset_index(drop=True) )