// Copyright (C) 2024-2026 Whiterun LLC // // This software is licensed under the GNU Affero General Public License (AGPL), version 3.0 or later. // A copy of the license can be found in the LICENSE file or at https://www.gnu.org/licenses/agpl-3.0.html use anyhow::Result; use bitcoincash::TokenID; use sqlx::{Row, SqlitePool}; use crate::db::blob::display_hex_to_blob; pub enum CachedSort { ScoreDesc, ScoreAsc, NameAsc, NameDesc, SymbolAsc, SymbolDesc, TvlDesc, TvlAsc, VolumeDesc, VolumeAsc, Change24hDesc, Change24hAsc, Change7dDesc, Change7dAsc, PriceUsdDesc, PriceUsdAsc, Change24hUsdDesc, Change24hUsdAsc, Change7dUsdDesc, Change7dUsdAsc, Apy30dDesc, Apy30dAsc, FirstPoolTsDesc, FirstPoolTsAsc, } pub fn order_clause(sort: CachedSort) -> &'static str { match sort { CachedSort::ScoreDesc => { "ORDER BY (NULLIF(score_rank,0) IS NULL) ASC, \ NULLIF(score_rank,0) ASC, \ score DESC, trade_volume DESC, token_id ASC" } CachedSort::ScoreAsc => { "ORDER BY (NULLIF(score_rank,0) IS NULL) ASC, \ NULLIF(score_rank,0) DESC, \ score ASC, trade_volume ASC, token_id ASC" } CachedSort::TvlDesc => "ORDER BY tvl_sats DESC, token_id ASC", CachedSort::TvlAsc => "ORDER BY tvl_sats ASC, token_id ASC", CachedSort::VolumeDesc => "ORDER BY trade_volume DESC, token_id ASC", CachedSort::VolumeAsc => "ORDER BY trade_volume ASC, token_id ASC", CachedSort::NameAsc => { "ORDER BY display_name IS NULL, display_name COLLATE NOCASE ASC, token_id ASC" } CachedSort::NameDesc => { "ORDER BY display_name IS NULL, display_name COLLATE NOCASE DESC, token_id ASC" } CachedSort::SymbolAsc => { "ORDER BY display_symbol IS NULL, display_symbol COLLATE NOCASE ASC, token_id ASC" } CachedSort::SymbolDesc => { "ORDER BY display_symbol IS NULL, display_symbol COLLATE NOCASE DESC, token_id ASC" } CachedSort::Change24hDesc => { "ORDER BY change_24h_bp IS NULL, change_24h_bp DESC, token_id ASC" } CachedSort::Change24hAsc => { "ORDER BY change_24h_bp IS NULL, change_24h_bp ASC, token_id ASC" } CachedSort::Change7dDesc => { "ORDER BY change_7d_bp IS NULL, change_7d_bp DESC, token_id ASC" } CachedSort::Change7dAsc => { "ORDER BY change_7d_bp IS NULL, change_7d_bp ASC, token_id ASC" } CachedSort::PriceUsdDesc => "ORDER BY price_now_usd DESC, token_id ASC", CachedSort::PriceUsdAsc => "ORDER BY price_now_usd ASC, token_id ASC", CachedSort::Change24hUsdDesc => { "ORDER BY change_24h_usd_bp IS NULL, change_24h_usd_bp DESC, token_id ASC" } CachedSort::Change24hUsdAsc => { "ORDER BY change_24h_usd_bp IS NULL, change_24h_usd_bp ASC, token_id ASC" } CachedSort::Change7dUsdDesc => { "ORDER BY change_7d_usd_bp IS NULL, change_7d_usd_bp DESC, token_id ASC" } CachedSort::Change7dUsdAsc => { "ORDER BY change_7d_usd_bp IS NULL, change_7d_usd_bp ASC, token_id ASC" } CachedSort::Apy30dDesc => "ORDER BY apy_30d_bp IS NULL, apy_30d_bp DESC, token_id ASC", CachedSort::Apy30dAsc => "ORDER BY apy_30d_bp IS NULL, apy_30d_bp ASC, token_id ASC", // NULLIF(...,0): 0 is the "not backfilled yet" sentinel — sort it with the NULLs, // last, so un-backfilled tokens never pollute "newest" results. CachedSort::FirstPoolTsDesc => { "ORDER BY NULLIF(first_pool_ts,0) IS NULL, NULLIF(first_pool_ts,0) DESC, token_id ASC" } CachedSort::FirstPoolTsAsc => { "ORDER BY NULLIF(first_pool_ts,0) IS NULL, NULLIF(first_pool_ts,0) ASC, token_id ASC" } } } /// Create the *current* index set only (no drops). pub async fn create_cached_token_metrics_indexes(pool: &SqlitePool) -> Result<()> { sqlx::query( r#" CREATE INDEX IF NOT EXISTS idx_ctm_score ON cached_token_metrics(score, trade_volume, token_id); "#, ) .execute(pool) .await?; sqlx::query( "CREATE INDEX IF NOT EXISTS idx_ctm_tvl ON cached_token_metrics(tvl_sats, token_id);", ) .execute(pool) .await?; sqlx::query("CREATE INDEX IF NOT EXISTS idx_ctm_volume ON cached_token_metrics(trade_volume, token_id);") .execute(pool).await?; sqlx::query("CREATE INDEX IF NOT EXISTS idx_ctm_price_usd ON cached_token_metrics(price_now_usd, token_id);") .execute(pool).await?; sqlx::query( "CREATE INDEX IF NOT EXISTS idx_ctm_apy ON cached_token_metrics(apy_30d_bp, token_id);", ) .execute(pool) .await?; sqlx::query("CREATE INDEX IF NOT EXISTS idx_ctm_name_ord ON cached_token_metrics((display_name IS NULL), display_name COLLATE NOCASE, token_id);") .execute(pool).await?; sqlx::query("CREATE INDEX IF NOT EXISTS idx_ctm_symbol_ord ON cached_token_metrics((display_symbol IS NULL), display_symbol COLLATE NOCASE, token_id);") .execute(pool).await?; sqlx::query("CREATE INDEX IF NOT EXISTS idx_ctm_ch24_ord ON cached_token_metrics((change_24h_bp IS NULL), change_24h_bp, token_id);") .execute(pool).await?; sqlx::query("CREATE INDEX IF NOT EXISTS idx_ctm_ch7d_ord ON cached_token_metrics((change_7d_bp IS NULL), change_7d_bp, token_id);") .execute(pool).await?; sqlx::query("CREATE INDEX IF NOT EXISTS idx_ctm_ch24_usd_ord ON cached_token_metrics((change_24h_usd_bp IS NULL), change_24h_usd_bp, token_id);") .execute(pool).await?; sqlx::query("CREATE INDEX IF NOT EXISTS idx_ctm_ch7d_usd_ord ON cached_token_metrics((change_7d_usd_bp IS NULL), change_7d_usd_bp, token_id);") .execute(pool).await?; sqlx::query("CREATE INDEX IF NOT EXISTS idx_ctm_score_rank ON cached_token_metrics(score_rank, token_id);") .execute(pool).await?; sqlx::query("CREATE INDEX IF NOT EXISTS idx_ctm_first_pool_ts_ord ON cached_token_metrics((NULLIF(first_pool_ts,0) IS NULL), NULLIF(first_pool_ts,0), token_id);") .execute(pool).await?; sqlx::query("ANALYZE;").execute(pool).await?; sqlx::query("PRAGMA optimize;").execute(pool).await?; Ok(()) } /// Indexes that speed up the aggregation path (tx → phe → pool). pub async fn create_aggregation_path_indexes(pool: &SqlitePool) -> Result<()> { sqlx::query("CREATE INDEX IF NOT EXISTS idx_tx_effective_ts ON tx(effective_timestamp, txid);") .execute(pool) .await?; sqlx::query("CREATE INDEX IF NOT EXISTS idx_phe_txid ON pool_history_entry(txid);") .execute(pool) .await?; sqlx::query("CREATE INDEX IF NOT EXISTS idx_phe_pool_txid ON pool_history_entry(pool, txid);") .execute(pool) .await?; sqlx::query("CREATE INDEX IF NOT EXISTS idx_phe_pool_timestamp ON pool_history_entry(pool, effective_timestamp);") .execute(pool).await?; sqlx::query( "CREATE INDEX IF NOT EXISTS idx_pool_creation_utxo ON pool(creation_utxo, token_id);", ) .execute(pool) .await?; sqlx::query("CREATE INDEX IF NOT EXISTS idx_pool_token_id ON pool(token_id, creation_utxo);") .execute(pool) .await?; sqlx::query("CREATE INDEX IF NOT EXISTS idx_volume_query_optimization ON pool_history_entry(txid, pool, sats_delta, token_delta);") .execute(pool).await?; sqlx::query("CREATE INDEX IF NOT EXISTS idx_tx_effective_timestamp_txid ON tx(effective_timestamp, txid);") .execute(pool).await?; sqlx::query("ANALYZE;").execute(pool).await?; sqlx::query("PRAGMA optimize;").execute(pool).await?; Ok(()) } pub async fn create_cached_token_metrics_table(pool: &SqlitePool) -> Result<()> { sqlx::query( r#" CREATE TABLE IF NOT EXISTS cached_token_metrics ( token_id BLOB PRIMARY KEY, trade_volume INTEGER NOT NULL DEFAULT 0, tvl_sats INTEGER NOT NULL DEFAULT 0, tvl_tokens INTEGER NOT NULL DEFAULT 0, score INTEGER NOT NULL DEFAULT 0, score_rank INTEGER NOT NULL DEFAULT 0, display_name TEXT, display_symbol TEXT, price_now REAL NOT NULL DEFAULT 0, price_24h REAL, price_7d REAL, change_24h_bp INTEGER, change_7d_bp INTEGER, price_now_usd REAL NOT NULL DEFAULT 0, price_24h_usd REAL, price_7d_usd REAL, change_24h_usd_bp INTEGER, change_7d_usd_bp INTEGER, apy_30d_bp INTEGER, last_trade_ts INTEGER, updated_at INTEGER NOT NULL ); "#, ) .execute(pool) .await?; add_column_if_missing(pool, "cached_token_metrics", "first_pool_ts", "INTEGER").await?; add_column_if_missing(pool, "cached_token_metrics", "bcmr_json", "TEXT").await?; add_column_if_missing(pool, "cached_token_metrics", "bcmr_well_known_json", "TEXT").await?; // Timestamp each change_* was actually measured from. For a token younger than // the window this is its first pool, not `now - window`, so consumers can label // a since-launch move as such instead of passing it off as a full 24h/7d one. add_column_if_missing( pool, "cached_token_metrics", "change_24h_anchor_ts", "INTEGER", ) .await?; add_column_if_missing( pool, "cached_token_metrics", "change_7d_anchor_ts", "INTEGER", ) .await?; // 7d aggregate-pool-price curve as a JSON array of ratios to price_now; see // db::cauldron::priceseries. Lets the token list ship a sparkline instead of // every client asking for candlesticks per row. add_column_if_missing(pool, "cached_token_metrics", "price_series_7d", "TEXT").await?; Ok(()) } async fn add_column_if_missing( pool: &SqlitePool, table: &str, column: &str, decl: &str, ) -> Result<()> { let rows = sqlx::query(&format!("PRAGMA table_info({table});")) .fetch_all(pool) .await?; let mut exists = false; for row in rows { let name: String = row.get(1); if name.eq_ignore_ascii_case(column) { exists = true; break; } } if !exists { sqlx::query(&format!("ALTER TABLE {table} ADD COLUMN {column} {decl};")) .execute(pool) .await?; } Ok(()) } async fn table_has_column(pool: &SqlitePool, table: &str, column: &str) -> anyhow::Result { let rows = sqlx::query(&format!("PRAGMA table_info({table});")) .fetch_all(pool) .await?; for row in rows { let name: String = row.get(1); if name.eq_ignore_ascii_case(column) { return Ok(true); } } Ok(false) } /// First pool info: (creation_utxo, txid, first_ts, first_height) pub type FirstPoolInfo = (String, String, i64, Option); /// Only set `first_pool_ts` when it is currently NULL or 0. pub async fn cache_first_pool_ts_if_empty( pool: &SqlitePool, token_id: &str, ts: i64, ) -> Result<()> { let token_blob = display_hex_to_blob::(token_id)?; sqlx::query( r#" INSERT INTO cached_token_metrics (token_id, updated_at) VALUES (?1, strftime('%s','now')) ON CONFLICT(token_id) DO NOTHING; "#, ) .bind(&token_blob) .execute(pool) .await?; sqlx::query( r#" UPDATE cached_token_metrics SET first_pool_ts = ?2 WHERE token_id = ?1 AND (first_pool_ts IS NULL OR first_pool_ts = 0); "#, ) .bind(&token_blob) .bind(ts) .execute(pool) .await?; Ok(()) } /// Returns (creation_utxo, txid, first_ts, first_height) for the first-ever pool of `token_id`. pub async fn db_first_pool_creation_row( pool: &SqlitePool, token_id: &str, ) -> anyhow::Result> { let has_tx_height = table_has_column(pool, "tx", "block_height").await?; let has_blockhash = table_has_column(pool, "tx", "blockhash").await?; let has_headers_key = table_has_column(pool, "headers", "key").await?; let has_headers_ht = table_has_column(pool, "headers", "height").await?; let sql = if has_tx_height { r#" SELECT p.creation_utxo, uf.txid, t.effective_timestamp, t.block_height AS block_height FROM pool AS p JOIN utxo_funding AS uf ON uf.new_utxo_hash = p.creation_utxo JOIN tx AS t ON t.txid = uf.txid WHERE p.token_id = ? ORDER BY (t.block_height IS NULL) ASC, t.block_height ASC, t.effective_timestamp ASC, uf.txid ASC, p.creation_utxo ASC LIMIT 1; "# } else if has_blockhash && has_headers_key && has_headers_ht { r#" SELECT p.creation_utxo, uf.txid, t.effective_timestamp, h.height AS block_height FROM pool AS p JOIN utxo_funding AS uf ON uf.new_utxo_hash = p.creation_utxo JOIN tx AS t ON t.txid = uf.txid LEFT JOIN headers AS h ON h.key = t.blockhash WHERE p.token_id = ? ORDER BY (t.blockhash IS NULL) ASC, (h.height IS NULL) ASC, h.height ASC, t.effective_timestamp ASC, uf.txid ASC, p.creation_utxo ASC LIMIT 1; "# } else { r#" SELECT p.creation_utxo, uf.txid, t.effective_timestamp, NULL AS block_height FROM pool AS p JOIN utxo_funding AS uf ON uf.new_utxo_hash = p.creation_utxo JOIN tx AS t ON t.txid = uf.txid WHERE p.token_id = ? ORDER BY t.effective_timestamp ASC, uf.txid ASC, p.creation_utxo ASC LIMIT 1; "# }; let token_blob = display_hex_to_blob::(token_id)?; let row = sqlx::query(sql) .bind(&token_blob) .fetch_optional(pool) .await?; if let Some(row) = row { use crate::db::blob::blob_to_display_hex; use bitcoincash::Txid; use riftenlabs_defi::chainutil::OutPointHash; let creation_utxo_blob: Vec = row.get(0); let txid_blob: Vec = row.get(1); let creation_utxo = blob_to_display_hex::(&creation_utxo_blob)?; let txid = blob_to_display_hex::(&txid_blob)?; let ts: i64 = row.get(2); let height: Option = row.get(3); Ok(Some((creation_utxo, txid, ts, height))) } else { Ok(None) } }