riftenlabs-indexer/src/rpc/tokens.rs
Dagur Valberg Johannsson f421327a5d
Switch to tokio and sqlx
This makes the application fully async; freeing web server threads to
handle new connections when waiting on SQL queries.

Additionally sqlx will allow easier move to a different database if
needed in the future.

Includes some SQL optimizations as well (slow queries more easliy identifiable
with sqlx).
2026-02-18 10:16:51 +01:00

246 lines
7.2 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 bitcoin_hashes::hex::{FromHex, ToHex};
use bitcoincash::TokenID;
use rocket::{get, State};
use serde_json::json;
use serde_json::Value;
use crate::db::cauldron::tokenlist::db_utils::{
cache_first_pool_ts_if_empty, db_first_pool_creation_row, CachedSort,
};
use crate::db::cauldron::tokenlist::list_cached::{
db_list_tokens_cached, db_list_tokens_cached_by_ids, TokenListItemCached,
};
use crate::db::search::db_search_tokens_cached;
use crate::db::DB;
use super::err::{bad_request, db_error, not_found, ApiErrorCode, CachedApiResult};
use super::response::{cached_ok, CACHE_AGGREGATE, CACHE_IMMUTABLE};
#[get("/tokens/search_by_volume?<search_query>")]
pub async fn search_by_volume(search_query: &str, db: &State<DB>) -> CachedApiResult<Vec<Value>> {
use rayon::prelude::*;
let list: Vec<(String, Option<String>, Option<String>, u64)> =
crate::db::search::search_tokens_by_volume(
&db.cauldron_r,
&db.bcmr_r,
&db.crc20_r,
search_query,
)
.await
.map_err(db_error)?;
let result: Vec<Value> = list
.into_par_iter()
.map(|(token_id, name, ticker, trade_volume)| {
json!({
"token_id": token_id,
"name": name,
"ticker": ticker,
"trade_volume": trade_volume
})
})
.collect();
Ok(cached_ok(result, CACHE_AGGREGATE))
}
#[get("/tokens/search_cached?<q>&<limit>&<offset>&<by>&<order>")]
pub async fn search_cached(
db: &State<DB>,
q: Option<String>,
limit: Option<usize>,
offset: Option<usize>,
by: Option<String>,
order: Option<String>,
) -> CachedApiResult<Value> {
let sort = parse_cached_sort(by, order);
let limit = limit.unwrap_or(250);
let offset = offset.unwrap_or(0);
let query = q.unwrap_or_default();
let items: Vec<TokenListItemCached> =
db_search_tokens_cached(&db.cauldron_r, &query, sort, limit, offset)
.await
.map_err(db_error)?;
Ok(cached_ok(json!(items), CACHE_AGGREGATE))
}
fn parse_cached_sort(by: Option<String>, order: Option<String>) -> CachedSort {
let desc = matches!(order.as_deref(), Some("desc") | Some("DESC"));
match by.as_deref() {
Some("name") => {
if desc {
CachedSort::NameDesc
} else {
CachedSort::NameAsc
}
}
Some("symbol") => {
if desc {
CachedSort::SymbolDesc
} else {
CachedSort::SymbolAsc
}
}
Some("tvl") => {
if desc {
CachedSort::TvlDesc
} else {
CachedSort::TvlAsc
}
}
Some("volume") => {
if desc {
CachedSort::VolumeDesc
} else {
CachedSort::VolumeAsc
}
}
Some("change_24h_bp") => {
if desc {
CachedSort::Change24hDesc
} else {
CachedSort::Change24hAsc
}
}
Some("change_7d_bp") => {
if desc {
CachedSort::Change7dDesc
} else {
CachedSort::Change7dAsc
}
}
// USD sorts
Some("price_usd") => {
if desc {
CachedSort::PriceUsdDesc
} else {
CachedSort::PriceUsdAsc
}
}
Some("change_24h_usd_bp") => {
if desc {
CachedSort::Change24hUsdDesc
} else {
CachedSort::Change24hUsdAsc
}
}
Some("change_7d_usd_bp") => {
if desc {
CachedSort::Change7dUsdDesc
} else {
CachedSort::Change7dUsdAsc
}
}
// APY (30d) sorts — accept a few aliases
Some("apy") | Some("apy_30d") | Some("apy_30d_bp") => {
if desc {
CachedSort::Apy30dDesc
} else {
CachedSort::Apy30dAsc
}
}
// default
_ => {
if desc {
CachedSort::ScoreDesc
} else {
CachedSort::ScoreAsc
}
}
}
}
#[get("/tokens/list_cached?<limit>&<offset>&<by>&<order>")]
pub async fn list_cached(
db: &State<DB>,
limit: Option<usize>,
offset: Option<usize>,
by: Option<String>,
order: Option<String>,
) -> CachedApiResult<Value> {
let limit = limit.unwrap_or(250);
let offset = offset.unwrap_or(0);
let sort = parse_cached_sort(by, order);
let items: Vec<TokenListItemCached> =
db_list_tokens_cached(&db.cauldron_r, limit, offset, sort)
.await
.map_err(db_error)?;
Ok(cached_ok(json!(items), CACHE_AGGREGATE))
}
#[get("/tokens/list_cached_by_ids?<ids>&<by>&<order>")]
pub async fn list_cached_by_ids(
db: &State<DB>,
ids: &str,
by: Option<String>,
order: Option<String>,
) -> CachedApiResult<Value> {
// split, trim, and normalize (lowercase is typical for hex IDs in DB)
let token_ids: Vec<String> = ids
.split(',')
.map(|s| s.trim().to_lowercase())
.filter(|s| !s.is_empty())
.collect();
if token_ids.is_empty() {
return Ok(cached_ok(json!([]), CACHE_AGGREGATE));
}
let sort = parse_cached_sort(by, order);
let items = db_list_tokens_cached_by_ids(&db.cauldron_r, &token_ids, sort)
.await
.map_err(db_error)?;
Ok(cached_ok(json!(items), CACHE_AGGREGATE))
}
#[get("/token/<token>/first_pool")]
pub async fn first_pool_creation(token: &str, dbp: &State<DB>) -> CachedApiResult<Value> {
let token = TokenID::from_hex(token).map_err(|e| {
bad_request(
ApiErrorCode::InvalidTokenId,
&format!("Invalid token id: {e}"),
)
})?;
let token_hex = token.to_hex();
match db_first_pool_creation_row(&dbp.cauldron_r, &token_hex).await {
Ok(Some((creation_utxo, txid, timestamp, block_height))) => {
// Opportunistically cache the timestamp (non-blocking)
let cauldron_w = dbp.cauldron_w.clone();
let token_hex_clone = token_hex.clone();
tokio::spawn(async move {
if let Err(e) =
cache_first_pool_ts_if_empty(&cauldron_w, &token_hex_clone, timestamp).await
{
log::warn!("Failed to cache first_pool_ts for {token_hex_clone}: {e}");
}
});
Ok(cached_ok(
json!({
"token": token_hex,
"creation_utxo": creation_utxo,
"txid": txid,
"timestamp": timestamp,
"block_height": block_height
}),
CACHE_IMMUTABLE,
))
}
Ok(None) => Err(not_found(ApiErrorCode::PoolNotFound, "No pools for token")),
Err(e) => Err(db_error(e)),
}
}