use chrono::{DateTime, Utc};
use serde::Serialize;
use sqlx::PgPool;
use uuid::Uuid;
use crate::error::Result;
#[derive(Debug, Clone, Serialize, sqlx::FromRow)]
pub struct DashboardMetrics {
pub total_users: i64,
pub new_users_24h: i64,
pub new_users_7d: i64,
pub total_tokens: i64,
pub new_tokens_24h: i64,
pub total_trades: i64,
pub trades_24h: i64,
pub total_volume_sats: i64,
pub volume_24h_sats: i64,
pub total_fees_sats: i64,
pub fees_24h_sats: i64,
pub pending_commitments: i64,
pub pending_kyc: i64,
}
#[derive(Debug, Clone, Serialize, sqlx::FromRow)]
pub struct TokenMetricsSummary {
pub token_id: uuid::Uuid,
pub symbol: String,
pub name: String,
pub total_supply: rust_decimal::Decimal,
pub holder_count: i64,
pub trade_count: i64,
pub trades_24h: i64,
pub total_volume_btc: rust_decimal::Decimal,
pub volume_24h_btc: rust_decimal::Decimal,
pub current_price_btc: rust_decimal::Decimal,
pub price_change_24h_pct: rust_decimal::Decimal,
}
#[derive(Debug, Clone, Serialize, sqlx::FromRow)]
pub struct UserActivitySummary {
pub user_id: uuid::Uuid,
pub username: String,
pub trade_count: i64,
pub total_volume_btc: rust_decimal::Decimal,
pub tokens_held: i64,
pub tokens_issued: i64,
pub reputation_score: i32,
pub last_activity: Option<chrono::DateTime<chrono::Utc>>,
}
#[derive(Debug, Clone, Serialize, sqlx::FromRow)]
pub struct DailyStats {
pub date: chrono::NaiveDate,
pub new_users: i64,
pub new_tokens: i64,
pub trade_count: i64,
pub volume_sats: i64,
pub fees_sats: i64,
pub active_users: i64,
}
#[derive(Debug, Clone, Serialize, sqlx::FromRow)]
pub struct PriceHistory {
pub time: DateTime<Utc>,
pub token_id: Uuid,
pub price_satoshis: i64,
pub supply: i64,
pub market_cap_satoshis: i64,
}
#[derive(Debug, Clone, Serialize, sqlx::FromRow)]
pub struct VolumeHistory {
pub time: DateTime<Utc>,
pub token_id: Uuid,
pub buy_volume_satoshis: i64,
pub sell_volume_satoshis: i64,
pub trade_count: i32,
pub unique_traders: i32,
}
#[derive(Debug, Clone, Serialize, sqlx::FromRow)]
pub struct PlatformVolumeHistory {
pub time: DateTime<Utc>,
pub total_volume_satoshis: i64,
pub trade_count: i32,
pub active_tokens: i32,
pub active_traders: i32,
pub fees_collected_satoshis: i64,
}
#[derive(Debug, Clone, Serialize, sqlx::FromRow)]
pub struct OhlcData {
pub bucket: DateTime<Utc>,
pub token_id: Uuid,
pub open_price: i64,
pub high_price: i64,
pub low_price: i64,
pub close_price: i64,
pub final_supply: i64,
pub final_market_cap: i64,
}
pub struct AnalyticsService {
pool: PgPool,
}
impl AnalyticsService {
pub fn new(pool: PgPool) -> Self {
Self { pool }
}
pub async fn create_materialized_views(&self) -> Result<()> {
sqlx::query(DASHBOARD_METRICS_VIEW_SQL)
.execute(&self.pool)
.await?;
sqlx::query(TOKEN_METRICS_VIEW_SQL)
.execute(&self.pool)
.await?;
sqlx::query(DAILY_STATS_VIEW_SQL)
.execute(&self.pool)
.await?;
sqlx::query(USER_ACTIVITY_VIEW_SQL)
.execute(&self.pool)
.await?;
tracing::info!("Created all materialized views for analytics");
Ok(())
}
pub async fn refresh_all_views(&self) -> Result<()> {
sqlx::query("REFRESH MATERIALIZED VIEW CONCURRENTLY IF EXISTS mv_dashboard_metrics")
.execute(&self.pool)
.await
.ok();
sqlx::query("REFRESH MATERIALIZED VIEW CONCURRENTLY IF EXISTS mv_token_metrics")
.execute(&self.pool)
.await
.ok();
sqlx::query("REFRESH MATERIALIZED VIEW CONCURRENTLY IF EXISTS mv_daily_stats")
.execute(&self.pool)
.await
.ok();
sqlx::query("REFRESH MATERIALIZED VIEW CONCURRENTLY IF EXISTS mv_user_activity")
.execute(&self.pool)
.await
.ok();
tracing::debug!("Refreshed all materialized views");
Ok(())
}
pub async fn get_dashboard_metrics(&self) -> Result<DashboardMetrics> {
let result =
sqlx::query_as::<_, DashboardMetrics>("SELECT * FROM mv_dashboard_metrics LIMIT 1")
.fetch_optional(&self.pool)
.await;
if let Ok(Some(metrics)) = result {
return Ok(metrics);
}
self.compute_dashboard_metrics().await
}
pub async fn compute_dashboard_metrics(&self) -> Result<DashboardMetrics> {
let metrics = sqlx::query_as::<_, DashboardMetrics>(
r#"
SELECT
(SELECT COUNT(*) FROM users) as total_users,
(SELECT COUNT(*) FROM users WHERE created_at > NOW() - INTERVAL '24 hours') as new_users_24h,
(SELECT COUNT(*) FROM users WHERE created_at > NOW() - INTERVAL '7 days') as new_users_7d,
(SELECT COUNT(*) FROM tokens WHERE status = 'active') as total_tokens,
(SELECT COUNT(*) FROM tokens WHERE created_at > NOW() - INTERVAL '24 hours') as new_tokens_24h,
(SELECT COUNT(*) FROM trades) as total_trades,
(SELECT COUNT(*) FROM trades WHERE created_at > NOW() - INTERVAL '24 hours') as trades_24h,
COALESCE((SELECT SUM((total_btc * 100000000)::bigint) FROM trades), 0) as total_volume_sats,
COALESCE((SELECT SUM((total_btc * 100000000)::bigint) FROM trades WHERE created_at > NOW() - INTERVAL '24 hours'), 0) as volume_24h_sats,
COALESCE((SELECT SUM((platform_fee * 100000000)::bigint) FROM trades), 0) as total_fees_sats,
COALESCE((SELECT SUM((platform_fee * 100000000)::bigint) FROM trades WHERE created_at > NOW() - INTERVAL '24 hours'), 0) as fees_24h_sats,
(SELECT COUNT(*) FROM output_commitments WHERE status = 'pending') as pending_commitments,
(SELECT COUNT(*) FROM kyc_applications WHERE status = 'pending') as pending_kyc
"#,
)
.fetch_one(&self.pool)
.await?;
Ok(metrics)
}
pub async fn get_top_tokens(&self, limit: i64) -> Result<Vec<TokenMetricsSummary>> {
let result = sqlx::query_as::<_, TokenMetricsSummary>(
r#"
SELECT * FROM mv_token_metrics
ORDER BY volume_24h_btc DESC
LIMIT $1
"#,
)
.bind(limit)
.fetch_all(&self.pool)
.await;
if let Ok(tokens) = result {
if !tokens.is_empty() {
return Ok(tokens);
}
}
self.compute_top_tokens(limit).await
}
async fn compute_top_tokens(&self, limit: i64) -> Result<Vec<TokenMetricsSummary>> {
let tokens = sqlx::query_as::<_, TokenMetricsSummary>(
r#"
SELECT
t.token_id,
t.symbol,
t.name,
t.total_supply,
COALESCE(h.holder_count, 0) as holder_count,
COALESCE(tr.trade_count, 0) as trade_count,
COALESCE(tr.trades_24h, 0) as trades_24h,
COALESCE(tr.total_volume_btc, 0) as total_volume_btc,
COALESCE(tr.volume_24h_btc, 0) as volume_24h_btc,
COALESCE(tr.last_price, t.base_price) as current_price_btc,
COALESCE(
CASE WHEN tr.price_24h_ago > 0
THEN ((tr.last_price - tr.price_24h_ago) / tr.price_24h_ago * 100)
ELSE 0
END,
0
) as price_change_24h_pct
FROM tokens t
LEFT JOIN (
SELECT token_id, COUNT(DISTINCT user_id) as holder_count
FROM balances
WHERE amount > 0
GROUP BY token_id
) h ON h.token_id = t.token_id
LEFT JOIN (
SELECT
token_id,
COUNT(*) as trade_count,
COUNT(*) FILTER (WHERE created_at > NOW() - INTERVAL '24 hours') as trades_24h,
SUM(total_btc) as total_volume_btc,
SUM(total_btc) FILTER (WHERE created_at > NOW() - INTERVAL '24 hours') as volume_24h_btc,
(SELECT price_btc FROM trades tr2 WHERE tr2.token_id = trades.token_id ORDER BY created_at DESC LIMIT 1) as last_price,
(SELECT price_btc FROM trades tr2 WHERE tr2.token_id = trades.token_id AND tr2.created_at < NOW() - INTERVAL '24 hours' ORDER BY created_at DESC LIMIT 1) as price_24h_ago
FROM trades
GROUP BY token_id
) tr ON tr.token_id = t.token_id
WHERE t.status = 'active'
ORDER BY COALESCE(tr.volume_24h_btc, 0) DESC
LIMIT $1
"#,
)
.bind(limit)
.fetch_all(&self.pool)
.await?;
Ok(tokens)
}
pub async fn get_daily_stats(
&self,
start_date: chrono::NaiveDate,
end_date: chrono::NaiveDate,
) -> Result<Vec<DailyStats>> {
let result = sqlx::query_as::<_, DailyStats>(
r#"
SELECT * FROM mv_daily_stats
WHERE date >= $1 AND date <= $2
ORDER BY date DESC
"#,
)
.bind(start_date)
.bind(end_date)
.fetch_all(&self.pool)
.await;
if let Ok(stats) = result {
if !stats.is_empty() {
return Ok(stats);
}
}
self.compute_daily_stats(start_date, end_date).await
}
async fn compute_daily_stats(
&self,
start_date: chrono::NaiveDate,
end_date: chrono::NaiveDate,
) -> Result<Vec<DailyStats>> {
let stats = sqlx::query_as::<_, DailyStats>(
r#"
WITH dates AS (
SELECT generate_series($1::date, $2::date, '1 day'::interval)::date as date
)
SELECT
d.date,
COALESCE(u.new_users, 0) as new_users,
COALESCE(t.new_tokens, 0) as new_tokens,
COALESCE(tr.trade_count, 0) as trade_count,
COALESCE(tr.volume_sats, 0) as volume_sats,
COALESCE(tr.fees_sats, 0) as fees_sats,
COALESCE(tr.active_users, 0) as active_users
FROM dates d
LEFT JOIN (
SELECT DATE(created_at) as date, COUNT(*) as new_users
FROM users
WHERE DATE(created_at) >= $1 AND DATE(created_at) <= $2
GROUP BY DATE(created_at)
) u ON u.date = d.date
LEFT JOIN (
SELECT DATE(created_at) as date, COUNT(*) as new_tokens
FROM tokens
WHERE DATE(created_at) >= $1 AND DATE(created_at) <= $2
GROUP BY DATE(created_at)
) t ON t.date = d.date
LEFT JOIN (
SELECT
DATE(created_at) as date,
COUNT(*) as trade_count,
SUM((total_btc * 100000000)::bigint) as volume_sats,
SUM((platform_fee * 100000000)::bigint) as fees_sats,
COUNT(DISTINCT buyer_id) + COUNT(DISTINCT seller_id) as active_users
FROM trades
WHERE DATE(created_at) >= $1 AND DATE(created_at) <= $2
GROUP BY DATE(created_at)
) tr ON tr.date = d.date
ORDER BY d.date DESC
"#,
)
.bind(start_date)
.bind(end_date)
.fetch_all(&self.pool)
.await?;
Ok(stats)
}
pub async fn get_top_users(&self, limit: i64) -> Result<Vec<UserActivitySummary>> {
let users = sqlx::query_as::<_, UserActivitySummary>(
r#"
SELECT
u.user_id,
u.username,
COALESCE(t.trade_count, 0) as trade_count,
COALESCE(t.total_volume_btc, 0) as total_volume_btc,
COALESCE(b.tokens_held, 0) as tokens_held,
COALESCE(tk.tokens_issued, 0) as tokens_issued,
u.reputation_score,
GREATEST(t.last_trade, u.created_at) as last_activity
FROM users u
LEFT JOIN (
SELECT
user_id,
COUNT(*) as trade_count,
SUM(total_btc) as total_volume_btc,
MAX(created_at) as last_trade
FROM (
SELECT buyer_id as user_id, total_btc, created_at FROM trades
UNION ALL
SELECT seller_id as user_id, total_btc, created_at FROM trades
) all_trades
GROUP BY user_id
) t ON t.user_id = u.user_id
LEFT JOIN (
SELECT user_id, COUNT(DISTINCT token_id) as tokens_held
FROM balances
WHERE amount > 0
GROUP BY user_id
) b ON b.user_id = u.user_id
LEFT JOIN (
SELECT issuer_id as user_id, COUNT(*) as tokens_issued
FROM tokens
GROUP BY issuer_id
) tk ON tk.user_id = u.user_id
ORDER BY COALESCE(t.total_volume_btc, 0) DESC
LIMIT $1
"#,
)
.bind(limit)
.fetch_all(&self.pool)
.await?;
Ok(users)
}
pub async fn drop_materialized_views(&self) -> Result<()> {
sqlx::query("DROP MATERIALIZED VIEW IF EXISTS mv_dashboard_metrics CASCADE")
.execute(&self.pool)
.await?;
sqlx::query("DROP MATERIALIZED VIEW IF EXISTS mv_token_metrics CASCADE")
.execute(&self.pool)
.await?;
sqlx::query("DROP MATERIALIZED VIEW IF EXISTS mv_daily_stats CASCADE")
.execute(&self.pool)
.await?;
sqlx::query("DROP MATERIALIZED VIEW IF EXISTS mv_user_activity CASCADE")
.execute(&self.pool)
.await?;
tracing::info!("Dropped all materialized views");
Ok(())
}
pub async fn record_price_history(
&self,
token_id: Uuid,
price_satoshis: i64,
supply: i64,
) -> Result<()> {
let market_cap_satoshis = price_satoshis.saturating_mul(supply);
sqlx::query(
r#"
INSERT INTO price_history (time, token_id, price_satoshis, supply, market_cap_satoshis)
VALUES (NOW(), $1, $2, $3, $4)
ON CONFLICT (time, token_id) DO UPDATE
SET price_satoshis = EXCLUDED.price_satoshis,
supply = EXCLUDED.supply,
market_cap_satoshis = EXCLUDED.market_cap_satoshis
"#,
)
.bind(token_id)
.bind(price_satoshis)
.bind(supply)
.bind(market_cap_satoshis)
.execute(&self.pool)
.await?;
Ok(())
}
pub async fn record_volume_history(
&self,
token_id: Uuid,
buy_volume_satoshis: i64,
sell_volume_satoshis: i64,
trade_count: i32,
unique_traders: i32,
) -> Result<()> {
sqlx::query(
r#"
INSERT INTO volume_history
(time, token_id, buy_volume_satoshis, sell_volume_satoshis, trade_count, unique_traders)
VALUES (date_trunc('hour', NOW()), $1, $2, $3, $4, $5)
ON CONFLICT (time, token_id) DO UPDATE
SET buy_volume_satoshis = volume_history.buy_volume_satoshis + EXCLUDED.buy_volume_satoshis,
sell_volume_satoshis = volume_history.sell_volume_satoshis + EXCLUDED.sell_volume_satoshis,
trade_count = volume_history.trade_count + EXCLUDED.trade_count,
unique_traders = GREATEST(volume_history.unique_traders, EXCLUDED.unique_traders)
"#,
)
.bind(token_id)
.bind(buy_volume_satoshis)
.bind(sell_volume_satoshis)
.bind(trade_count)
.bind(unique_traders)
.execute(&self.pool)
.await?;
Ok(())
}
pub async fn record_platform_volume(
&self,
total_volume_satoshis: i64,
trade_count: i32,
active_tokens: i32,
active_traders: i32,
fees_collected_satoshis: i64,
) -> Result<()> {
sqlx::query(
r#"
INSERT INTO platform_volume_history
(time, total_volume_satoshis, trade_count, active_tokens, active_traders, fees_collected_satoshis)
VALUES (date_trunc('hour', NOW()), $1, $2, $3, $4, $5)
ON CONFLICT (time) DO UPDATE
SET total_volume_satoshis = platform_volume_history.total_volume_satoshis + EXCLUDED.total_volume_satoshis,
trade_count = platform_volume_history.trade_count + EXCLUDED.trade_count,
active_tokens = GREATEST(platform_volume_history.active_tokens, EXCLUDED.active_tokens),
active_traders = GREATEST(platform_volume_history.active_traders, EXCLUDED.active_traders),
fees_collected_satoshis = platform_volume_history.fees_collected_satoshis + EXCLUDED.fees_collected_satoshis
"#,
)
.bind(total_volume_satoshis)
.bind(trade_count)
.bind(active_tokens)
.bind(active_traders)
.bind(fees_collected_satoshis)
.execute(&self.pool)
.await?;
Ok(())
}
pub async fn get_price_history(
&self,
token_id: Uuid,
start_time: DateTime<Utc>,
end_time: DateTime<Utc>,
) -> Result<Vec<PriceHistory>> {
let history = sqlx::query_as::<_, PriceHistory>(
r#"
SELECT time, token_id, price_satoshis, supply, market_cap_satoshis
FROM price_history
WHERE token_id = $1 AND time >= $2 AND time <= $3
ORDER BY time ASC
"#,
)
.bind(token_id)
.bind(start_time)
.bind(end_time)
.fetch_all(&self.pool)
.await?;
Ok(history)
}
pub async fn get_volume_history(
&self,
token_id: Uuid,
start_time: DateTime<Utc>,
end_time: DateTime<Utc>,
) -> Result<Vec<VolumeHistory>> {
let history = sqlx::query_as::<_, VolumeHistory>(
r#"
SELECT time, token_id, buy_volume_satoshis, sell_volume_satoshis, trade_count, unique_traders
FROM volume_history
WHERE token_id = $1 AND time >= $2 AND time <= $3
ORDER BY time ASC
"#,
)
.bind(token_id)
.bind(start_time)
.bind(end_time)
.fetch_all(&self.pool)
.await?;
Ok(history)
}
pub async fn get_ohlc_data(
&self,
token_id: Uuid,
start_time: DateTime<Utc>,
end_time: DateTime<Utc>,
bucket_interval: &str, ) -> Result<Vec<OhlcData>> {
if bucket_interval == "1 hour" {
if let Ok(data) = self
.get_ohlc_from_aggregate(token_id, start_time, end_time)
.await
{
if !data.is_empty() {
return Ok(data);
}
}
}
self.compute_ohlc_data(token_id, start_time, end_time, bucket_interval)
.await
}
async fn get_ohlc_from_aggregate(
&self,
token_id: Uuid,
start_time: DateTime<Utc>,
end_time: DateTime<Utc>,
) -> Result<Vec<OhlcData>> {
let data = sqlx::query_as::<_, OhlcData>(
r#"
SELECT bucket, token_id, open_price, high_price, low_price, close_price,
final_supply, final_market_cap
FROM price_history_hourly
WHERE token_id = $1 AND bucket >= $2 AND bucket <= $3
ORDER BY bucket ASC
"#,
)
.bind(token_id)
.bind(start_time)
.bind(end_time)
.fetch_all(&self.pool)
.await?;
Ok(data)
}
async fn compute_ohlc_data(
&self,
token_id: Uuid,
start_time: DateTime<Utc>,
end_time: DateTime<Utc>,
bucket_interval: &str,
) -> Result<Vec<OhlcData>> {
let data = sqlx::query_as::<_, OhlcData>(
r#"
SELECT
time_bucket($1::interval, time) AS bucket,
token_id,
first(price_satoshis, time) AS open_price,
max(price_satoshis) AS high_price,
min(price_satoshis) AS low_price,
last(price_satoshis, time) AS close_price,
last(supply, time) AS final_supply,
last(market_cap_satoshis, time) AS final_market_cap
FROM price_history
WHERE token_id = $2 AND time >= $3 AND time <= $4
GROUP BY bucket, token_id
ORDER BY bucket ASC
"#,
)
.bind(bucket_interval)
.bind(token_id)
.bind(start_time)
.bind(end_time)
.fetch_all(&self.pool)
.await?;
Ok(data)
}
pub async fn get_aggregated_volume(
&self,
token_id: Option<Uuid>,
start_time: DateTime<Utc>,
end_time: DateTime<Utc>,
bucket_interval: &str, ) -> Result<Vec<VolumeHistory>> {
let volume = if let Some(tid) = token_id {
sqlx::query_as::<_, VolumeHistory>(
r#"
SELECT
time_bucket($1::interval, time) AS time,
token_id,
sum(buy_volume_satoshis) AS buy_volume_satoshis,
sum(sell_volume_satoshis) AS sell_volume_satoshis,
sum(trade_count)::int AS trade_count,
max(unique_traders)::int AS unique_traders
FROM volume_history
WHERE token_id = $2 AND time >= $3 AND time <= $4
GROUP BY time_bucket($1::interval, time), token_id
ORDER BY time ASC
"#,
)
.bind(bucket_interval)
.bind(tid)
.bind(start_time)
.bind(end_time)
.fetch_all(&self.pool)
.await?
} else {
sqlx::query_as::<_, VolumeHistory>(
r#"
SELECT
time_bucket($1::interval, time) AS time,
'00000000-0000-0000-0000-000000000000'::uuid AS token_id,
sum(buy_volume_satoshis) AS buy_volume_satoshis,
sum(sell_volume_satoshis) AS sell_volume_satoshis,
sum(trade_count)::int AS trade_count,
sum(unique_traders)::int AS unique_traders
FROM volume_history
WHERE time >= $2 AND time <= $3
GROUP BY time_bucket($1::interval, time)
ORDER BY time ASC
"#,
)
.bind(bucket_interval)
.bind(start_time)
.bind(end_time)
.fetch_all(&self.pool)
.await?
};
Ok(volume)
}
pub async fn get_platform_volume_history(
&self,
start_time: DateTime<Utc>,
end_time: DateTime<Utc>,
) -> Result<Vec<PlatformVolumeHistory>> {
let history = sqlx::query_as::<_, PlatformVolumeHistory>(
r#"
SELECT time, total_volume_satoshis, trade_count, active_tokens,
active_traders, fees_collected_satoshis
FROM platform_volume_history
WHERE time >= $1 AND time <= $2
ORDER BY time ASC
"#,
)
.bind(start_time)
.bind(end_time)
.fetch_all(&self.pool)
.await?;
Ok(history)
}
pub async fn is_timescaledb_available(&self) -> bool {
sqlx::query_scalar::<_, bool>(
"SELECT EXISTS(SELECT 1 FROM pg_extension WHERE extname = 'timescaledb')",
)
.fetch_one(&self.pool)
.await
.unwrap_or(false)
}
pub async fn get_hypertable_info(&self, table_name: &str) -> Result<Option<String>> {
let info = sqlx::query_scalar::<_, String>(
r#"
SELECT format('Hypertable: %s, Chunks: %s, Compression: %s',
hypertable_name,
num_chunks,
compression_enabled)
FROM timescaledb_information.hypertables
WHERE hypertable_name = $1
"#,
)
.bind(table_name)
.fetch_optional(&self.pool)
.await?;
Ok(info)
}
}
const DASHBOARD_METRICS_VIEW_SQL: &str = r#"
CREATE MATERIALIZED VIEW IF NOT EXISTS mv_dashboard_metrics AS
SELECT
(SELECT COUNT(*) FROM users) as total_users,
(SELECT COUNT(*) FROM users WHERE created_at > NOW() - INTERVAL '24 hours') as new_users_24h,
(SELECT COUNT(*) FROM users WHERE created_at > NOW() - INTERVAL '7 days') as new_users_7d,
(SELECT COUNT(*) FROM tokens WHERE status = 'active') as total_tokens,
(SELECT COUNT(*) FROM tokens WHERE created_at > NOW() - INTERVAL '24 hours') as new_tokens_24h,
(SELECT COUNT(*) FROM trades) as total_trades,
(SELECT COUNT(*) FROM trades WHERE created_at > NOW() - INTERVAL '24 hours') as trades_24h,
COALESCE((SELECT SUM((total_btc * 100000000)::bigint) FROM trades), 0) as total_volume_sats,
COALESCE((SELECT SUM((total_btc * 100000000)::bigint) FROM trades WHERE created_at > NOW() - INTERVAL '24 hours'), 0) as volume_24h_sats,
COALESCE((SELECT SUM((platform_fee * 100000000)::bigint) FROM trades), 0) as total_fees_sats,
COALESCE((SELECT SUM((platform_fee * 100000000)::bigint) FROM trades WHERE created_at > NOW() - INTERVAL '24 hours'), 0) as fees_24h_sats,
(SELECT COUNT(*) FROM output_commitments WHERE status = 'pending') as pending_commitments,
COALESCE((SELECT COUNT(*) FROM kyc_applications WHERE status = 'pending'), 0) as pending_kyc;
CREATE UNIQUE INDEX IF NOT EXISTS mv_dashboard_metrics_idx ON mv_dashboard_metrics ((1));
"#;
const TOKEN_METRICS_VIEW_SQL: &str = r#"
CREATE MATERIALIZED VIEW IF NOT EXISTS mv_token_metrics AS
SELECT
t.token_id,
t.symbol,
t.name,
t.total_supply,
COALESCE(h.holder_count, 0) as holder_count,
COALESCE(tr.trade_count, 0) as trade_count,
COALESCE(tr.trades_24h, 0) as trades_24h,
COALESCE(tr.total_volume_btc, 0) as total_volume_btc,
COALESCE(tr.volume_24h_btc, 0) as volume_24h_btc,
COALESCE(tr.last_price, t.base_price) as current_price_btc,
COALESCE(
CASE WHEN tr.price_24h_ago > 0
THEN ((tr.last_price - tr.price_24h_ago) / tr.price_24h_ago * 100)
ELSE 0
END,
0
) as price_change_24h_pct
FROM tokens t
LEFT JOIN (
SELECT token_id, COUNT(DISTINCT user_id) as holder_count
FROM balances
WHERE amount > 0
GROUP BY token_id
) h ON h.token_id = t.token_id
LEFT JOIN (
SELECT
token_id,
COUNT(*) as trade_count,
COUNT(*) FILTER (WHERE created_at > NOW() - INTERVAL '24 hours') as trades_24h,
SUM(total_btc) as total_volume_btc,
SUM(total_btc) FILTER (WHERE created_at > NOW() - INTERVAL '24 hours') as volume_24h_btc,
(SELECT price_btc FROM trades tr2 WHERE tr2.token_id = trades.token_id ORDER BY created_at DESC LIMIT 1) as last_price,
(SELECT price_btc FROM trades tr2 WHERE tr2.token_id = trades.token_id AND tr2.created_at < NOW() - INTERVAL '24 hours' ORDER BY created_at DESC LIMIT 1) as price_24h_ago
FROM trades
GROUP BY token_id
) tr ON tr.token_id = t.token_id
WHERE t.status = 'active';
CREATE UNIQUE INDEX IF NOT EXISTS mv_token_metrics_token_id_idx ON mv_token_metrics (token_id);
CREATE INDEX IF NOT EXISTS mv_token_metrics_volume_idx ON mv_token_metrics (volume_24h_btc DESC);
"#;
const DAILY_STATS_VIEW_SQL: &str = r#"
CREATE MATERIALIZED VIEW IF NOT EXISTS mv_daily_stats AS
WITH dates AS (
SELECT generate_series(
(SELECT COALESCE(MIN(DATE(created_at)), CURRENT_DATE - INTERVAL '30 days') FROM users),
CURRENT_DATE,
'1 day'::interval
)::date as date
)
SELECT
d.date,
COALESCE(u.new_users, 0)::bigint as new_users,
COALESCE(t.new_tokens, 0)::bigint as new_tokens,
COALESCE(tr.trade_count, 0)::bigint as trade_count,
COALESCE(tr.volume_sats, 0)::bigint as volume_sats,
COALESCE(tr.fees_sats, 0)::bigint as fees_sats,
COALESCE(tr.active_users, 0)::bigint as active_users
FROM dates d
LEFT JOIN (
SELECT DATE(created_at) as date, COUNT(*) as new_users
FROM users
GROUP BY DATE(created_at)
) u ON u.date = d.date
LEFT JOIN (
SELECT DATE(created_at) as date, COUNT(*) as new_tokens
FROM tokens
GROUP BY DATE(created_at)
) t ON t.date = d.date
LEFT JOIN (
SELECT
DATE(created_at) as date,
COUNT(*) as trade_count,
SUM((total_btc * 100000000)::bigint) as volume_sats,
SUM((platform_fee * 100000000)::bigint) as fees_sats,
COUNT(DISTINCT buyer_id) + COUNT(DISTINCT seller_id) as active_users
FROM trades
GROUP BY DATE(created_at)
) tr ON tr.date = d.date;
CREATE UNIQUE INDEX IF NOT EXISTS mv_daily_stats_date_idx ON mv_daily_stats (date);
"#;
const USER_ACTIVITY_VIEW_SQL: &str = r#"
CREATE MATERIALIZED VIEW IF NOT EXISTS mv_user_activity AS
SELECT
u.user_id,
u.username,
COALESCE(t.trade_count, 0)::bigint as trade_count,
COALESCE(t.total_volume_btc, 0) as total_volume_btc,
COALESCE(b.tokens_held, 0)::bigint as tokens_held,
COALESCE(tk.tokens_issued, 0)::bigint as tokens_issued,
u.reputation_score,
GREATEST(t.last_trade, u.created_at) as last_activity
FROM users u
LEFT JOIN (
SELECT
user_id,
COUNT(*) as trade_count,
SUM(total_btc) as total_volume_btc,
MAX(created_at) as last_trade
FROM (
SELECT buyer_id as user_id, total_btc, created_at FROM trades
UNION ALL
SELECT seller_id as user_id, total_btc, created_at FROM trades
) all_trades
GROUP BY user_id
) t ON t.user_id = u.user_id
LEFT JOIN (
SELECT user_id, COUNT(DISTINCT token_id) as tokens_held
FROM balances
WHERE amount > 0
GROUP BY user_id
) b ON b.user_id = u.user_id
LEFT JOIN (
SELECT issuer_id as user_id, COUNT(*) as tokens_issued
FROM tokens
GROUP BY issuer_id
) tk ON tk.user_id = u.user_id;
CREATE UNIQUE INDEX IF NOT EXISTS mv_user_activity_user_id_idx ON mv_user_activity (user_id);
CREATE INDEX IF NOT EXISTS mv_user_activity_volume_idx ON mv_user_activity (total_volume_btc DESC);
"#;
#[derive(Debug, Clone)]
pub struct RefreshConfig {
pub refresh_interval_secs: u64,
pub concurrent: bool,
}
impl Default for RefreshConfig {
fn default() -> Self {
Self {
refresh_interval_secs: 300, concurrent: true,
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_refresh_config_defaults() {
let config = RefreshConfig::default();
assert_eq!(config.refresh_interval_secs, 300);
assert!(config.concurrent);
}
#[test]
fn test_price_history_creation() {
let token_id = Uuid::new_v4();
let price_history = PriceHistory {
time: Utc::now(),
token_id,
price_satoshis: 100_000,
supply: 1_000_000,
market_cap_satoshis: 100_000_000_000,
};
assert_eq!(price_history.token_id, token_id);
assert_eq!(price_history.price_satoshis, 100_000);
assert_eq!(price_history.supply, 1_000_000);
assert_eq!(price_history.market_cap_satoshis, 100_000_000_000);
}
#[test]
fn test_volume_history_creation() {
let token_id = Uuid::new_v4();
let volume_history = VolumeHistory {
time: Utc::now(),
token_id,
buy_volume_satoshis: 50_000_000,
sell_volume_satoshis: 30_000_000,
trade_count: 42,
unique_traders: 15,
};
assert_eq!(volume_history.token_id, token_id);
assert_eq!(volume_history.buy_volume_satoshis, 50_000_000);
assert_eq!(volume_history.sell_volume_satoshis, 30_000_000);
assert_eq!(volume_history.trade_count, 42);
assert_eq!(volume_history.unique_traders, 15);
}
#[test]
fn test_platform_volume_history_creation() {
let platform_volume = PlatformVolumeHistory {
time: Utc::now(),
total_volume_satoshis: 1_000_000_000,
trade_count: 1000,
active_tokens: 50,
active_traders: 200,
fees_collected_satoshis: 5_000_000,
};
assert_eq!(platform_volume.total_volume_satoshis, 1_000_000_000);
assert_eq!(platform_volume.trade_count, 1000);
assert_eq!(platform_volume.active_tokens, 50);
assert_eq!(platform_volume.active_traders, 200);
assert_eq!(platform_volume.fees_collected_satoshis, 5_000_000);
}
#[test]
fn test_ohlc_data_creation() {
let token_id = Uuid::new_v4();
let ohlc = OhlcData {
bucket: Utc::now(),
token_id,
open_price: 95_000,
high_price: 105_000,
low_price: 90_000,
close_price: 100_000,
final_supply: 1_000_000,
final_market_cap: 100_000_000_000,
};
assert_eq!(ohlc.token_id, token_id);
assert_eq!(ohlc.open_price, 95_000);
assert_eq!(ohlc.high_price, 105_000);
assert_eq!(ohlc.low_price, 90_000);
assert_eq!(ohlc.close_price, 100_000);
assert!(ohlc.high_price >= ohlc.open_price);
assert!(ohlc.high_price >= ohlc.close_price);
assert!(ohlc.low_price <= ohlc.open_price);
assert!(ohlc.low_price <= ohlc.close_price);
}
#[test]
fn test_timescaledb_structures_are_serializable() {
let token_id = Uuid::new_v4();
let time = Utc::now();
let price = PriceHistory {
time,
token_id,
price_satoshis: 100_000,
supply: 1_000_000,
market_cap_satoshis: 100_000_000_000,
};
let volume = VolumeHistory {
time,
token_id,
buy_volume_satoshis: 50_000_000,
sell_volume_satoshis: 30_000_000,
trade_count: 42,
unique_traders: 15,
};
let platform = PlatformVolumeHistory {
time,
total_volume_satoshis: 1_000_000_000,
trade_count: 1000,
active_tokens: 50,
active_traders: 200,
fees_collected_satoshis: 5_000_000,
};
let ohlc = OhlcData {
bucket: time,
token_id,
open_price: 95_000,
high_price: 105_000,
low_price: 90_000,
close_price: 100_000,
final_supply: 1_000_000,
final_market_cap: 100_000_000_000,
};
assert!(serde_json::to_string(&price).is_ok());
assert!(serde_json::to_string(&volume).is_ok());
assert!(serde_json::to_string(&platform).is_ok());
assert!(serde_json::to_string(&ohlc).is_ok());
}
#[test]
fn test_price_history_market_cap_calculation() {
let token_id = Uuid::new_v4();
let price_satoshis = 100_000i64;
let supply = 1_000_000i64;
let expected_market_cap = price_satoshis.saturating_mul(supply);
let price_history = PriceHistory {
time: Utc::now(),
token_id,
price_satoshis,
supply,
market_cap_satoshis: expected_market_cap,
};
assert_eq!(price_history.market_cap_satoshis, 100_000_000_000);
assert_eq!(
price_history.market_cap_satoshis,
price_history.price_satoshis * price_history.supply
);
}
}