riftenlabs-indexer/src/db/cauldron/tokenlist/db_utils.rs
jakobsn f05d4c6c34 Expose first_pool_ts in cached token list API and add it as a sort key
Adds the already-backfilled cached_token_metrics.first_pool_ts column to
the TokenListItemCached response (nullable, unix seconds) in all three
query paths (list_cached, search_cached, list_cached_by_ids), and accepts
by=first_pool_ts on the cached list endpoints so a frontend can fetch
"newest tokens" in one call.

NULLs (and the 0 "not backfilled" sentinel, via NULLIF) always sort last
regardless of order, so tokens without a timestamp never pollute newest
results. Includes a matching expression index and sort/serialization
tests for both the list and search paths.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-07-30 13:57:21 +02:00

370 lines
14 KiB
Rust

// 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?;
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<bool> {
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<i64>);
/// 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::<TokenID>(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<Option<FirstPoolInfo>> {
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::<TokenID>(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<u8> = row.get(0);
let txid_blob: Vec<u8> = row.get(1);
let creation_utxo = blob_to_display_hex::<OutPointHash>(&creation_utxo_blob)?;
let txid = blob_to_display_hex::<Txid>(&txid_blob)?;
let ts: i64 = row.get(2);
let height: Option<i64> = row.get(3);
Ok(Some((creation_utxo, txid, ts, height)))
} else {
Ok(None)
}
}