Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -189,20 +189,48 @@ export async function up(knex: Knex): Promise<void> {
);
`);

// 1-hour bridge volume continuous aggregate
await knex.raw(`
CREATE MATERIALIZED VIEW IF NOT EXISTS bridge_hourly_volume_rollup
WITH (timescaledb.continuous) AS
SELECT
time_bucket('1 hour', created_at) AS bucket,
bridge_name AS bridge_id,
SUM(amount) AS total_volume,
COUNT(*) AS transaction_count,
AVG(amount) AS avg_amount,
MIN(amount) AS min_amount,
MAX(amount) AS max_amount
FROM bridge_transactions
GROUP BY bucket, bridge_id
WITH NO DATA;
`);

await knex.raw(`
SELECT add_continuous_aggregate_policy('bridge_hourly_volume_rollup',
start_offset => INTERVAL '7 days',
end_offset => INTERVAL '1 hour',
schedule_interval => INTERVAL '30 minutes',
if_not_exists => TRUE
);
`);

// Create performance indexes on materialized views
await knex.raw(`CREATE INDEX IF NOT EXISTS prices_hourly_symbol_bucket_idx ON prices_hourly (symbol, bucket DESC);`);
await knex.raw(`CREATE INDEX IF NOT EXISTS prices_daily_symbol_bucket_idx ON prices_daily (symbol, bucket DESC);`);
await knex.raw(`CREATE INDEX IF NOT EXISTS health_scores_hourly_symbol_bucket_idx ON health_scores_hourly (symbol, bucket DESC);`);
await knex.raw(`CREATE INDEX IF NOT EXISTS health_scores_daily_symbol_bucket_idx ON health_scores_daily (symbol, bucket DESC);`);
await knex.raw(`CREATE INDEX IF NOT EXISTS liquidity_hourly_symbol_bucket_idx ON liquidity_hourly (symbol, bucket DESC);`);
await knex.raw(`CREATE INDEX IF NOT EXISTS liquidity_daily_symbol_bucket_idx ON liquidity_daily (symbol, bucket DESC);`);
await knex.raw(`CREATE INDEX IF NOT EXISTS bridge_hourly_volume_rollup_idx ON bridge_hourly_volume_rollup (bridge_id, bucket DESC);`);
} catch {
// Continuous aggregates require TimescaleDB 2.x extension.
// In environments without TimescaleDB, migration proceeds gracefully without failing.
}
}

export async function down(knex: Knex): Promise<void> {
await knex.raw(`DROP MATERIALIZED VIEW IF EXISTS bridge_hourly_volume_rollup CASCADE;`);
await knex.raw(`DROP MATERIALIZED VIEW IF EXISTS liquidity_daily CASCADE;`);
await knex.raw(`DROP MATERIALIZED VIEW IF EXISTS liquidity_hourly CASCADE;`);
await knex.raw(`DROP MATERIALIZED VIEW IF EXISTS health_scores_daily CASCADE;`);
Expand Down
86 changes: 86 additions & 0 deletions backend/src/services/analytics.service.ts
Original file line number Diff line number Diff line change
Expand Up @@ -746,4 +746,90 @@ export class AnalyticsService {
{ bypassCache, tags: ["analytics", "historical"], ttl: CacheTTL.ANALYTICS }
);
}

/**
* Get hourly volume rollup for a bridge using TimescaleDB continuous aggregate
*/
async getBridgeHourlyVolumeRollup(
bridgeId: string,
startDate: Date,
endDate: Date = new Date(),
bypassCache: boolean = false
): Promise<Array<{
bucket: string;
bridgeId: string;
totalVolume: string;
transactionCount: number;
avgAmount: string;
minAmount: string;
maxAmount: string;
}>> {
const cacheKey = CacheService.generateKey(
"analytics",
`bridge_rollup:${bridgeId}:${startDate.toISOString()}:${endDate.toISOString()}`
);

return CacheService.getOrSet(
cacheKey,
async () => {
try {
const results = await knex("bridge_hourly_volume_rollup")
.select(
"bucket",
"bridge_id as bridgeId",
"total_volume as totalVolume",
"transaction_count as transactionCount",
"avg_amount as avgAmount",
"min_amount as minAmount",
"max_amount as maxAmount"
)
.where("bridge_id", bridgeId)
.where("bucket", ">=", startDate)
.where("bucket", "<=", endDate)
.orderBy("bucket", "asc");

if (results && results.length > 0) {
return results.map((row: any) => ({
bucket: row.bucket,
bridgeId: row.bridgeId,
totalVolume: row.totalVolume?.toString() || "0",
transactionCount: Number(row.transactionCount || 0),
avgAmount: row.avgAmount?.toString() || "0",
minAmount: row.minAmount?.toString() || "0",
maxAmount: row.maxAmount?.toString() || "0",
}));
}
} catch {
// Fallback to raw bridge_transactions aggregation if continuous aggregate is not present
}

const rawResults = await knex("bridge_transactions")
.select(
knex.raw("DATE_TRUNC('hour', created_at) as bucket"),
"bridge_name as bridgeId",
knex.raw("SUM(amount) as total_volume"),
knex.raw("COUNT(*) as transaction_count"),
knex.raw("AVG(amount) as avg_amount"),
knex.raw("MIN(amount) as min_amount"),
knex.raw("MAX(amount) as max_amount")
)
.where("bridge_name", bridgeId)
.where("created_at", ">=", startDate)
.where("created_at", "<=", endDate)
.groupBy("bucket", "bridge_name")
.orderBy("bucket", "asc");

return rawResults.map((row: any) => ({
bucket: row.bucket,
bridgeId: row.bridgeId,
totalVolume: row.total_volume?.toString() || "0",
transactionCount: Number(row.transaction_count || 0),
avgAmount: row.avg_amount?.toString() || "0",
minAmount: row.min_amount?.toString() || "0",
maxAmount: row.max_amount?.toString() || "0",
}));
},
{ bypassCache, tags: ["analytics", "bridge_rollup"], ttl: CacheTTL.ANALYTICS }
);
}
}
Loading
Loading