From 34bdfbe208f518883c74e4079b93809e3f7409d5 Mon Sep 17 00:00:00 2001 From: Dominika Czupich Date: Mon, 10 Aug 2026 21:41:59 +0200 Subject: [PATCH 1/3] feat: add batch read/insert database operations --- create-table.js | 68 ----------- server.ts | 247 ++++++++++++++++++++++++++++++++------ src/create-table.ts | 6 +- src/fetchers/defillama.ts | 17 +++ src/fetchers/jupiter.ts | 4 +- src/fetchers/kamino.ts | 2 +- src/fetchers/morpho.ts | 2 +- src/fetchers/save.ts | 2 +- src/types.ts | 11 ++ 9 files changed, 245 insertions(+), 114 deletions(-) delete mode 100644 create-table.js diff --git a/create-table.js b/create-table.js deleted file mode 100644 index 64c027c..0000000 --- a/create-table.js +++ /dev/null @@ -1,68 +0,0 @@ -'use strict'; - -// Run once to provision the BigQuery tables: -// node create-table.js - -const { BigQuery } = require('@google-cloud/bigquery'); -const { GoogleAuth, Impersonated } = require('google-auth-library'); - -const PROJECT = 'calmal'; -const DATASET = 'lending_poc'; -const SA = process.env.IMPERSONATE_SA || 'lending-poc@calmal.iam.gserviceaccount.com'; - -const POOL_SCHEMA = [ - { name: 'id', type: 'INT64', mode: 'REQUIRED' }, - { name: 'reservePubkey', type: 'STRING', mode: 'REQUIRED' }, - { name: 'symbol', type: 'STRING', mode: 'REQUIRED' }, - { name: 'mintAddress', type: 'STRING', mode: 'REQUIRED' }, - { name: 'lending', type: 'STRING', mode: 'REQUIRED' }, - { name: 'chain', type: 'STRING', mode: 'REQUIRED' }, - { name: 'market', type: 'STRING', mode: 'REQUIRED' }, -]; - -const SNAPSHOTS_SCHEMA = [ - { name: 'poolId', type: 'INT64', mode: 'REQUIRED' }, - { name: 'tvl', type: 'NUMERIC', mode: 'REQUIRED' }, - { name: 'utilization', type: 'NUMERIC', mode: 'REQUIRED' }, - { name: 'supplyAPY', type: 'NUMERIC', mode: 'REQUIRED' }, - { name: 'borrowRate', type: 'NUMERIC', mode: 'REQUIRED' }, - { name: 'borrowAPY', type: 'NUMERIC', mode: 'REQUIRED' }, - { name: 'totalBorrowUsd',type: 'NUMERIC', mode: 'REQUIRED' }, - { name: 'liquidityUsd', type: 'NUMERIC', mode: 'REQUIRED' }, - { name: 'fetchedAt', type: 'TIMESTAMP', mode: 'REQUIRED' }, -]; - -(async () => { - const base = new GoogleAuth(); - const sourceClient = await base.getClient(); - - const impersonated = new Impersonated({ - sourceClient, - targetPrincipal: SA, - lifetime: 300, - targetScopes: ['https://www.googleapis.com/auth/bigquery'], - }); - - const bq = new BigQuery({ projectId: PROJECT, authClient: impersonated }); - console.log(`[bq] impersonating ${SA}`); - - const ds = bq.dataset(DATASET); - - for (const name of ['pool', 'snapshots']) { - const table = ds.table(name); - const [exists] = await table.exists(); - if (exists) { - await table.delete(); - console.log(`Deleted ${DATASET}.${name}`); - } - } - - const [pool] = await ds.createTable('pool', { schema: POOL_SCHEMA }); - const [snapshots] = await ds.createTable('snapshots', { schema: SNAPSHOTS_SCHEMA }); - - console.log(`Created ${PROJECT}.${DATASET}.${pool.id}`); - console.log(`Created ${PROJECT}.${DATASET}.${snapshots.id}`); -})().catch(err => { - console.error(err.message); - process.exit(1); -}); diff --git a/server.ts b/server.ts index c72f5c5..a8c1727 100644 --- a/server.ts +++ b/server.ts @@ -9,7 +9,7 @@ import type { StandarizedMetric, PoolRow, SnapshotRow } from './src/types.js'; const PORT: number = Number(process.env.PORT) || 3000; -const POLL_MS = 5000; +const POLL_MS = 300_000; const BQ_PROJECT: string = process.env.BIGQUERY_PROJECT_ID || 'calmal'; const BQ_DATASET: string = process.env.BIGQUERY_DATASET || 'lending_poc'; @@ -19,8 +19,24 @@ const htmlContent: Buffer = fs.readFileSync(path.join(__dirname, 'public', 'inde let latestData: { fetchedAt: string } | null = null; const clients: Set = new Set(); -const seenGenericPools: Set = new Set(); let isPolling = false; +const knownPools = new Set(); +let globalBQ: BigQuery | null = null; + +const BUCKET_SECONDS = 300; + +function timeBucket(date: Date): string { + const epoch = Math.floor(date.getTime() / 1000); + const bucket = epoch - (epoch % BUCKET_SECONDS); + return new Date(bucket * 1000).toISOString(); +} + +function stableSnapshotId(poolId: number, bucketTimestamp: string): string { + return crypto.createHash('sha256') + .update(`${poolId}::${bucketTimestamp}`) + .digest('hex') + .slice(0, 16); +} function formatSSE(data: unknown): string { return `event: update\ndata: ${JSON.stringify(data)}\n\n`; @@ -38,6 +54,7 @@ function stableGenericPoolId(metric: StandarizedMetric): number { return parseInt(crypto.createHash('sha256').update(key).digest('hex').slice(0, 8), 16); } +/* async function commitGenericPool(ds: Dataset, metric: StandarizedMetric): Promise { const key = `${metric.mintAddress}:${metric.lending}:${metric.market}`; if (seenGenericPools.has(key)) return; @@ -60,42 +77,140 @@ async function commitGenericPool(ds: Dataset, metric: StandarizedMetric): Promis throw err; } } - +*/ function safeNum(v: unknown): number { const n = Number(v); return isFinite(n) ? n : 0; } -async function commitGenericSnapshot(ds: Dataset, metric: StandarizedMetric): Promise { - const tvl = safeNum(metric.tvl); - const utilF = safeNum(metric.utilization) / 100; - const borrow = parseFloat((tvl * utilF).toFixed(9)); - const liquid = parseFloat((tvl - borrow).toFixed(9)); - const row: SnapshotRow = { - poolId: stableGenericPoolId(metric), - tvl: String(tvl), - utilization: String(utilF.toFixed(9)), - supplyAPY: String(safeNum(metric.supplyAPY)), - borrowRate: String(safeNum(metric.borrowRate)), - borrowAPY: String(safeNum(metric.borrowAPY)), - totalBorrowUsd: String(borrow), - liquidityUsd: String(liquid), - fetchedAt: new Date().toISOString(), - }; - await ds.table('snapshots').insert([row]); +async function batchCommitPools(bq: BigQuery, metrics: StandarizedMetric[]): Promise { + if (metrics.length === 0) return; + + const newMetrics = metrics.filter(m => !knownPools.has(stableGenericPoolId(m))); + if (newMetrics.length === 0) return; + + for (let i = 0; i < newMetrics.length; i += 100) { + const chunk = newMetrics.slice(i, i + 100); + const params: Record = {}; + const selects = chunk.map((m, idx) => { + const id = stableGenericPoolId(m); + params[`id_${idx}`] = id; + params[`rp_${idx}`] = m.mintAddress; + params[`sym_${idx}`] = m.symbol; + params[`ma_${idx}`] = m.mintAddress; + params[`len_${idx}`] = m.lending; + params[`ch_${idx}`] = m.chain; + params[`mkt_${idx}`] = m.market; + return `SELECT @id_${idx} AS id, @rp_${idx} AS reservePubkey, @sym_${idx} AS symbol, @ma_${idx} AS mintAddress, @len_${idx} AS lending, @ch_${idx} AS chain, @mkt_${idx} AS market`; + }).join(' UNION ALL '); + + const query = ` + MERGE \`${BQ_PROJECT}.${BQ_DATASET}.pool\` AS target + USING (${selects}) AS source + ON target.id = source.id + WHEN NOT MATCHED THEN + INSERT (id, reservePubkey, symbol, mintAddress, lending, chain, market) + VALUES (source.id, source.reservePubkey, source.symbol, source.mintAddress, source.lending, source.chain, source.market) + `; + + await bq.query({ query, params }); + + for (const m of chunk) { + knownPools.add(stableGenericPoolId(m)); + } + } +} + +async function batchCommitSnapshots(bq: BigQuery, metrics: StandarizedMetric[], bucketTimestamp: string): Promise { + if (metrics.length === 0) return; + + for (let i = 0; i < metrics.length; i += 100) { + const chunk = metrics.slice(i, i + 100); + const params: Record = {}; + const selects = chunk.map((m, idx) => { + const poolId = stableGenericPoolId(m); + const snapshotId = stableSnapshotId(poolId, bucketTimestamp); + const tvl = safeNum(m.tvl); + const utilF = safeNum(m.utilization) / 100; + const borrow = parseFloat((tvl * utilF).toFixed(9)); + const liquid = parseFloat((tvl - borrow).toFixed(9)); + + params[`sid_${idx}`] = snapshotId; + params[`pid_${idx}`] = poolId; + params[`tvl_${idx}`] = tvl; + params[`uti_${idx}`] = parseFloat(utilF.toFixed(9)); + params[`sapy_${idx}`] = safeNum(m.supplyAPY); + params[`br_${idx}`] = safeNum(m.borrowRate); + params[`bapy_${idx}`] = safeNum(m.borrowAPY); + params[`tbu_${idx}`] = borrow; + params[`liq_${idx}`] = liquid; + params[`fa_${idx}`] = bucketTimestamp; + + return `SELECT @sid_${idx} AS snapshotId, @pid_${idx} AS poolId, @tvl_${idx} AS tvl, @uti_${idx} AS utilization, @sapy_${idx} AS supplyAPY, @br_${idx} AS borrowRate, @bapy_${idx} AS borrowAPY, @tbu_${idx} AS totalBorrowUsd, @liq_${idx} AS liquidityUsd, CAST(@fa_${idx} AS TIMESTAMP) AS fetchedAt`; + }).join(' UNION ALL '); + + const query = ` + MERGE \`${BQ_PROJECT}.${BQ_DATASET}.snapshots\` AS target + USING (${selects}) AS source + ON target.snapshotId = source.snapshotId + WHEN MATCHED THEN + UPDATE SET + tvl = source.tvl, utilization = source.utilization, supplyAPY = source.supplyAPY, + borrowRate = source.borrowRate, borrowAPY = source.borrowAPY, + totalBorrowUsd = source.totalBorrowUsd, liquidityUsd = source.liquidityUsd, + fetchedAt = source.fetchedAt + WHEN NOT MATCHED THEN + INSERT (snapshotId, poolId, tvl, utilization, supplyAPY, borrowRate, + borrowAPY, totalBorrowUsd, liquidityUsd, fetchedAt) + VALUES (source.snapshotId, source.poolId, source.tvl, source.utilization, source.supplyAPY, source.borrowRate, + source.borrowAPY, source.totalBorrowUsd, source.liquidityUsd, source.fetchedAt) + `; + + await bq.query({ query, params }); + } +} + +function isMetricSafe(m: StandarizedMetric): boolean { + if (m.tvl < 0 || m.tvl > 1_000_000_000_000) return false; + if (m.utilization < 0 || m.utilization > 100) return false; + if (m.supplyAPY < 0 || m.supplyAPY > 1_000_000) return false; + if (m.borrowRate < 0 || m.borrowRate > 1_000_000) return false; + return true; } -async function pollProtocols(ds: Dataset): Promise { +async function pollProtocols(bq: BigQuery | null, ds: Dataset | null, enableDebugDump: boolean = true): Promise { if (isPolling) return; isPolling = true; try { const metrics = await fetchAllProtocolMetrics(); console.log(`[protocols] fetched ${metrics.length} metrics`); - latestData = { fetchedAt: new Date().toISOString() }; + + const now = new Date(); + const bucket = timeBucket(now); + + latestData = { fetchedAt: now.toISOString() }; broadcast(latestData); - for (const metric of metrics) { - commitGenericPool(ds, metric).catch((err: Error) => console.error('[bq:generic:pool]', err.message)); - commitGenericSnapshot(ds, metric).catch((err: Error) => console.error('[bq:generic:snapshot]', err.message)); + + const safeMetrics = metrics.filter(isMetricSafe); + + if (enableDebugDump) { + const dumpPath = path.join(process.cwd(), 'debug_data.json'); + fs.writeFileSync(dumpPath, JSON.stringify({ + bucketTimestamp: bucket, + totalMetrics: safeMetrics.length, + data: safeMetrics + }, null, 2)); + console.log(`[debug] Saved ${safeMetrics.length} metrics to ${dumpPath}`); + } + + if (safeMetrics.length > 0 && bq) { + await batchCommitPools(bq, safeMetrics); + console.log(`[bq] batch merged pools`); + + await batchCommitSnapshots(bq, safeMetrics, bucket); + console.log(`[bq] batch merged ${safeMetrics.length} snapshots for bucket ${bucket}`); + } else if (safeMetrics.length > 0 && !bq) { + console.log(`[warning] Database disabled. Skipped BigQuery insert for ${safeMetrics.length} snapshots.`); } } catch (err) { console.error('[protocols] poll error:', (err as Error).message); @@ -104,6 +219,45 @@ async function pollProtocols(ds: Dataset): Promise { } } +import type { TokenDataResult } from './src/types.js'; + +async function queryBigQueryLatest(bq: BigQuery, protocol: string, symbol: string): Promise { + const query = ` + SELECT + UNIX_SECONDS(s.fetchedAt) AS date, + SUM(s.tvl) AS tvlUsd, + AVG(s.supplyAPY) AS supplyAPY, + AVG(s.borrowAPY) AS borrowAPY + FROM \`${BQ_PROJECT}.${BQ_DATASET}.pool\` p + JOIN \`${BQ_PROJECT}.${BQ_DATASET}.snapshots\` s ON p.id = s.poolId + WHERE LOWER(p.lending) = LOWER(@protocol) + AND LOWER(p.symbol) = LOWER(@symbol) + GROUP BY s.fetchedAt + ORDER BY s.fetchedAt DESC + LIMIT 1 + `; + + const [rows] = await bq.query({ + query, + params: { protocol, symbol } + }); + + if (!rows || rows.length === 0) return null; + + const history = rows.map(r => ({ + date: Number(r.date), + tvlUsd: Number(r.tvlUsd), + supplyAPY: Number(r.supplyAPY), + borrowAPY: Number(r.borrowAPY) + })); + + return { + source: "BigQuery", + poolId: null, + history + } as any; +} + // Heartbeat keeps SSE connections alive through proxies setInterval(() => { for (const res of clients) { @@ -157,7 +311,11 @@ function requestHandler(req: http.IncomingMessage, res: http.ServerResponse): vo return; } - fetchPlotData(protocol, symbol, collateral) + const fetchPromise = globalBQ + ? queryBigQueryLatest(globalBQ, protocol, symbol).then(data => data || fetchPlotData(protocol, symbol, collateral)) + : fetchPlotData(protocol, symbol, collateral); + + fetchPromise .then((data) => { if (!data || (data.history.length === 0 && !data.poolId)) { res.writeHead(200, { 'Content-Type': 'application/json; charset=utf-8', 'Access-Control-Allow-Origin': '*' }); @@ -180,26 +338,35 @@ function requestHandler(req: http.IncomingMessage, res: http.ServerResponse): vo async function main(): Promise { - const base = new GoogleAuth(); - const sourceClient = await base.getClient(); - - const impersonated = new Impersonated({ - sourceClient, - targetPrincipal: BQ_SA, - lifetime: 3600, - targetScopes: ['https://www.googleapis.com/auth/bigquery'], - }); + let bigquery: BigQuery | null = null; + let ds: Dataset | null = null; + + try { + const base = new GoogleAuth(); + const sourceClient = await base.getClient(); + + const impersonated = new Impersonated({ + sourceClient, + targetPrincipal: BQ_SA, + lifetime: 3600, + targetScopes: ['https://www.googleapis.com/auth/bigquery'], + }); - const bigquery = new BigQuery({ projectId: BQ_PROJECT, authClient: impersonated }); - console.log(`[bq] impersonating ${BQ_SA}`); + bigquery = new BigQuery({ projectId: BQ_PROJECT, authClient: impersonated }); + console.log(`[bq] impersonating ${BQ_SA}`); + ds = bigquery.dataset(BQ_DATASET); + globalBQ = bigquery; + } catch (err) { + console.warn(`[warning] Could not load Google credentials. Running in local dry-run mode. Error: ${(err as Error).message}`); + } - const ds: Dataset = bigquery.dataset(BQ_DATASET); const server = http.createServer(requestHandler); server.listen(PORT, () => { console.log(`[server] listening on port ${PORT}`); - pollProtocols(ds); - setInterval(() => pollProtocols(ds), POLL_MS); + // Change to 'true' to enable local JSON debug dumping without BQ credentials + pollProtocols(bigquery, ds, false); + setInterval(() => pollProtocols(bigquery, ds, false), POLL_MS); if (process.send) process.send('ready'); }); diff --git a/src/create-table.ts b/src/create-table.ts index 1a7a58d..3cad9e2 100644 --- a/src/create-table.ts +++ b/src/create-table.ts @@ -17,6 +17,7 @@ const POOL_SCHEMA: BQSchemaField[] = [ ]; const SNAPSHOTS_SCHEMA: BQSchemaField[] = [ + { name: 'snapshotId', type: 'STRING', mode: 'REQUIRED' }, { name: 'poolId', type: 'INT64', mode: 'REQUIRED' }, { name: 'tvl', type: 'NUMERIC', mode: 'REQUIRED' }, { name: 'utilization', type: 'NUMERIC', mode: 'REQUIRED' }, @@ -54,7 +55,10 @@ const SNAPSHOTS_SCHEMA: BQSchemaField[] = [ } const [pool] = await ds.createTable('pool', { schema: POOL_SCHEMA }); - const [snapshots] = await ds.createTable('snapshots', { schema: SNAPSHOTS_SCHEMA }); + const [snapshots] = await ds.createTable('snapshots', { + schema: SNAPSHOTS_SCHEMA, + clustering: { fields: ['poolId'] }, + }); console.log(`Created ${PROJECT}.${DATASET}.${pool.id}`); console.log(`Created ${PROJECT}.${DATASET}.${snapshots.id}`); diff --git a/src/fetchers/defillama.ts b/src/fetchers/defillama.ts index faa1613..08eb152 100644 --- a/src/fetchers/defillama.ts +++ b/src/fetchers/defillama.ts @@ -74,6 +74,17 @@ interface DefiLlamaPool { apyReward?: number; apyBaseBorrow?: number; utilization?: number; + predictions?: { + predictedClass?: string; + predictedProbability?: number; + binnedConfidence?: number; + }; + mu?: number; + sigma?: number; + apyMean30d?: number; + ilRisk?: string; + exposure?: string; + stablecoin?: boolean; } interface DefiLlamaChartEntry { @@ -137,6 +148,12 @@ export async function fetchDefiLlamaPlot( supplyAPY: parseFloat(((tokenPool.apyBase ?? 0) + (tokenPool.apyReward ?? 0)).toFixed(2)), borrowRate: parseFloat((tokenPool.apyBaseBorrow ?? 0).toFixed(2)), utilization: parseFloat((tokenPool.utilization ?? 0).toFixed(2)), + predictions: tokenPool.predictions, + mu: tokenPool.mu, + sigma: tokenPool.sigma, + apyMean30d: tokenPool.apyMean30d, + ilRisk: tokenPool.ilRisk, + exposure: tokenPool.exposure, }, }; } catch { diff --git a/src/fetchers/jupiter.ts b/src/fetchers/jupiter.ts index c128130..3c2c0bc 100644 --- a/src/fetchers/jupiter.ts +++ b/src/fetchers/jupiter.ts @@ -86,7 +86,7 @@ async function fetchEarnTokens(): Promise { const price = parseFloat(t.asset.price) || 0; const totalAssets = Number(t.totalAssets) / Math.pow(10, t.decimals); const tvlUsd = totalAssets * price; - if (tvlUsd <= 100_000) return null; + if (tvlUsd <= 10_000) return null; const mint = new PublicKey(t.assetAddress); const reserve = await fetchTokenReserve(mint); @@ -140,7 +140,7 @@ export async function fetchJupiterMetrics(): Promise { .filter((v) => { const supplyPrice = parseFloat(v.supplyToken.price) || 0; const supplyUsd = (Number(v.totalSupply) / Math.pow(10, v.supplyToken.decimals)) * supplyPrice; - return supplyUsd > 100000; + return supplyUsd > 10000; }) .map((v): StandarizedMetric => { const borrowPrice = parseFloat(v.borrowToken.price) || 0; diff --git a/src/fetchers/kamino.ts b/src/fetchers/kamino.ts index f483d60..f667bec 100644 --- a/src/fetchers/kamino.ts +++ b/src/fetchers/kamino.ts @@ -64,7 +64,7 @@ export async function fetchKaminoMetrics(): Promise { const marketName = config.name || 'isolated'; return reserves - .filter((r) => parseFloat(r.totalSupplyUsd) > 100_000) + .filter((r) => parseFloat(r.totalSupplyUsd) > 10_000) .map((r): StandarizedMetric => { const totalSupplyUsd = parseFloat(r.totalSupplyUsd); const totalBorrowUsd = parseFloat(r.totalBorrowUsd); diff --git a/src/fetchers/morpho.ts b/src/fetchers/morpho.ts index d1992ac..6efde5b 100644 --- a/src/fetchers/morpho.ts +++ b/src/fetchers/morpho.ts @@ -76,7 +76,7 @@ export async function fetchMorphoMetrics(): Promise { const markets = await fetchMorphoMarkets(); return markets - .filter((m) => m.state.supplyAssetsUsd > 100000) + .filter((m) => m.state.supplyAssetsUsd > 10000) .map((m): StandarizedMetric => { const lltvRaw = m.lltv ? Number(m.lltv) / 1e18 : null; return { diff --git a/src/fetchers/save.ts b/src/fetchers/save.ts index 6f8c8a9..11d6569 100644 --- a/src/fetchers/save.ts +++ b/src/fetchers/save.ts @@ -10,7 +10,7 @@ import type { import { symbolMatches, downsampleToDaily } from './defillama.js'; const SAVE_API = 'https://api.solend.fi'; -const MIN_TVL_USD = 100_000; +const MIN_TVL_USD = 10_000; interface SaveHistoryPoint { supplyAPY: number; diff --git a/src/types.ts b/src/types.ts index 7d6719c..ea0c484 100644 --- a/src/types.ts +++ b/src/types.ts @@ -135,6 +135,7 @@ export interface PoolRow { } export interface SnapshotRow { + snapshotId: string; poolId: number; tvl: string; utilization: string; @@ -164,6 +165,16 @@ export interface TokenSnapshot { borrowRate: number; utilization: number; protocolTotalActiveLoans?: number | null; + predictions?: { + predictedClass?: string; + predictedProbability?: number; + binnedConfidence?: number; + }; + mu?: number; + sigma?: number; + apyMean30d?: number; + ilRisk?: string; + exposure?: string; } export interface TokenDataResult { From db48c6b57cb107d3d2a5eb56a048190ed898e741 Mon Sep 17 00:00:00 2001 From: Dominika Czupich Date: Tue, 11 Aug 2026 22:59:03 +0200 Subject: [PATCH 2/3] implemented database + cache data architecture --- public/index.html | 10 ++-- server.ts | 107 ++++++++++++++++++++++++----------------- src/create-table.ts | 2 + src/fetchers/morpho.ts | 1 + src/types.ts | 8 ++- 5 files changed, 80 insertions(+), 48 deletions(-) diff --git a/public/index.html b/public/index.html index a952a31..66d9708 100644 --- a/public/index.html +++ b/public/index.html @@ -238,8 +238,12 @@