diff --git a/apps/api/src/db.ts b/apps/api/src/db.ts index 745dc61..b6c1807 100644 --- a/apps/api/src/db.ts +++ b/apps/api/src/db.ts @@ -1,12 +1,12 @@ -import { Pool, type QueryResult, type QueryResultRow } from "pg"; -import pino from "pino"; -import client from "prom-client"; -import { env } from "./config/env.js"; +import { Pool, type QueryResult, type QueryResultRow } from 'pg'; +import pino from 'pino'; +import client from 'prom-client'; +import { env } from './config/env.js'; -const logger = pino({ name: "db" }); +const logger = pino({ name: 'db' }); const url = new URL(env.DATABASE_URL); -const sslRequired = url.searchParams.get("sslmode") === "require"; +const sslRequired = url.searchParams.get('sslmode') === 'require'; const basePool = new Pool({ connectionString: env.DATABASE_URL, @@ -19,44 +19,44 @@ const basePool = new Pool({ // ── Pool monitoring (#264) ──────────────────────────────────────────────────── export const pgPoolEventsTotal = new client.Counter({ - name: "pg_pool_events_total", - help: "Total count of PostgreSQL pool lifecycle events", - labelNames: ["event"], + name: 'pg_pool_events_total', + help: 'Total count of PostgreSQL pool lifecycle events', + labelNames: ['event'], }); // Gauges read the pool's own counters on scrape rather than being tracked by // hand — `pool.totalCount`/`idleCount`/`waitingCount` are always the source // of truth, so there's no risk of manual increment/decrement drift. new client.Gauge({ - name: "pg_pool_active", - help: "Number of PostgreSQL pool clients currently checked out (in use)", + name: 'pg_pool_active', + help: 'Number of PostgreSQL pool clients currently checked out (in use)', collect() { this.set(basePool.totalCount - basePool.idleCount); }, }); new client.Gauge({ - name: "pg_pool_idle", - help: "Number of idle PostgreSQL pool clients available for reuse", + name: 'pg_pool_idle', + help: 'Number of idle PostgreSQL pool clients available for reuse', collect() { this.set(basePool.idleCount); }, }); new client.Gauge({ - name: "pg_pool_waiting", - help: "Number of queued requests waiting for a PostgreSQL pool client", + name: 'pg_pool_waiting', + help: 'Number of queued requests waiting for a PostgreSQL pool client', collect() { this.set(basePool.waitingCount); }, }); -basePool.on("connect", () => pgPoolEventsTotal.inc({ event: "connect" })); -basePool.on("acquire", () => pgPoolEventsTotal.inc({ event: "acquire" })); -basePool.on("remove", () => pgPoolEventsTotal.inc({ event: "remove" })); -basePool.on("error", (err) => { - pgPoolEventsTotal.inc({ event: "error" }); - logger.error({ err }, "PostgreSQL pool error (idle client)"); +basePool.on('connect', () => pgPoolEventsTotal.inc({ event: 'connect' })); +basePool.on('acquire', () => pgPoolEventsTotal.inc({ event: 'acquire' })); +basePool.on('remove', () => pgPoolEventsTotal.inc({ event: 'remove' })); +basePool.on('error', (err) => { + pgPoolEventsTotal.inc({ event: 'error' }); + logger.error({ err }, 'PostgreSQL pool error (idle client)'); }); // Alert when the pool is under sustained exhaustion pressure: waitingCount > 5 @@ -78,7 +78,7 @@ setInterval(() => { waitingAlertFired = true; logger.error( { waiting, thresholdSeconds: WAITING_ALERT_DURATION_MS / 1000 }, - `PostgreSQL pool exhaustion: ${waiting} requests waiting for >${WAITING_ALERT_DURATION_MS / 1000}s`, + `PostgreSQL pool exhaustion: ${waiting} requests waiting for >${WAITING_ALERT_DURATION_MS / 1000}s` ); } } else { @@ -90,16 +90,16 @@ setInterval(() => { // ── Prometheus metrics (#373) ───────────────────────────────────────────────── export const dbQueryDurationSeconds = new client.Histogram({ - name: "db_query_duration_seconds", - help: "Duration of PostgreSQL queries in seconds", - labelNames: ["query_name"], + name: 'db_query_duration_seconds', + help: 'Duration of PostgreSQL queries in seconds', + labelNames: ['query_name'], buckets: [0.01, 0.05, 0.1, 0.25, 0.5, 1, 2.5, 5, 10], }); export const dbSlowQueriesTotal = new client.Counter({ - name: "db_slow_queries_total", - help: "Total number of slow PostgreSQL queries", - labelNames: ["threshold"], + name: 'db_slow_queries_total', + help: 'Total number of slow PostgreSQL queries', + labelNames: ['threshold'], }); // ── Query timing wrapper ────────────────────────────────────────────────────── @@ -110,21 +110,23 @@ const SQL_TRUNCATE_LEN = 500; function sanitizeSql(sql: string): string { return sql - .replace(/\$\d+/g, "?") + .replace(/\$\d+/g, '?') .replace(/'[^']*'/g, "'?'") .slice(0, SQL_TRUNCATE_LEN); } function inferQueryName(sql: string): string { const s = sql.trim().toUpperCase(); - const m = s.match(/^(SELECT|INSERT|UPDATE|DELETE|CREATE|DROP|ALTER)\s+(?:INTO\s+|FROM\s+|TABLE\s+)?(\w+)?/); - return m ? `${m[1]!.toLowerCase()}${m[2] ? `_${m[2]!.toLowerCase()}` : ""}` : "unknown"; + const m = s.match( + /^(SELECT|INSERT|UPDATE|DELETE|CREATE|DROP|ALTER)\s+(?:INTO\s+|FROM\s+|TABLE\s+)?(\w+)?/ + ); + return m ? `${m[1]!.toLowerCase()}${m[2] ? `_${m[2]!.toLowerCase()}` : ''}` : 'unknown'; } async function timedQuery( sql: string, values?: unknown[], - queryName?: string, + queryName?: string ): Promise> { const start = Date.now(); const name = queryName ?? inferQueryName(sql); @@ -136,16 +138,16 @@ async function timedQuery( endTimer(); if (durationMs >= SLOW_ERROR_MS) { - dbSlowQueriesTotal.inc({ threshold: "2000ms" }); + dbSlowQueriesTotal.inc({ threshold: '2000ms' }); logger.error( { query: sanitizeSql(sql), durationMs, rowCount: result.rowCount, caller: name }, - "critically slow query", + 'critically slow query' ); } else if (durationMs >= SLOW_WARN_MS) { - dbSlowQueriesTotal.inc({ threshold: "500ms" }); + dbSlowQueriesTotal.inc({ threshold: '500ms' }); logger.warn( { query: sanitizeSql(sql), durationMs, rowCount: result.rowCount, caller: name }, - "slow query", + 'slow query' ); } @@ -217,6 +219,17 @@ export async function migrate(): Promise { -- this index adds kind as the second column for efficient type-specific lookups. CREATE INDEX IF NOT EXISTS idx_contract_events_importer_kind ON contract_events(importer_id, kind, created_at DESC); + -- #245: BRIN index on contract_events.created_at for time-range queries. + -- A B-tree index stores pointers for every single row and becomes extremely large at scale. + -- A BRIN (Block Range Index) index summarizes block ranges (minimum/maximum timestamps per range of pages), + -- resulting in a footprint that is orders of magnitude smaller (typically ~1:800 or 99.8% smaller). + -- pages_per_range is set to 32 (down from default 128) to provide finer search granularity, + -- which is highly effective for chronologically ordered event logs under high-volume ingestion. + -- If contract_events is partitioned by month (Issue #228), parent indexes automatically propagate + -- to all child partitions. + CREATE INDEX IF NOT EXISTS idx_contract_events_created_at_brin ON contract_events USING BRIN (created_at) WITH (pages_per_range = 32); + + ALTER TABLE contract_events ADD COLUMN IF NOT EXISTS ledger_sequence INTEGER; ALTER TABLE contract_events ADD COLUMN IF NOT EXISTS event_index INTEGER; @@ -689,9 +702,9 @@ export async function migrate(): Promise { CREATE INDEX IF NOT EXISTS idx_refresh_tokens_hash ON refresh_tokens(token_hash) WHERE revoked_at IS NULL; `, undefined, - "migrate_schema", + 'migrate_schema' ); - console.log("[migrate] schema ready"); + console.log('[migrate] schema ready'); } // ── importer_metrics_mv (#251) ──────────────────────────────────────────────── @@ -718,14 +731,18 @@ export async function getImporterMetrics(): Promise { compliance_rate: string; topup_count_30d: number; refreshed_at: Date; - }>("SELECT * FROM importer_metrics_mv WHERE singleton_id = 1", undefined, "select_importer_metrics_mv"); + }>( + 'SELECT * FROM importer_metrics_mv WHERE singleton_id = 1', + undefined, + 'select_importer_metrics_mv' + ); const row = result.rows[0]; if (!row) { return { totalImporters: 0, - totalBondValue: "0", - avgBalance: "0", + totalBondValue: '0', + avgBalance: '0', complianceRate: 100, topupCount30d: 0, refreshedAt: new Date(0).toISOString(), @@ -748,17 +765,17 @@ export async function getImporterMetrics(): Promise { */ export async function refreshImporterMetrics(): Promise { await timedQuery( - "REFRESH MATERIALIZED VIEW CONCURRENTLY importer_metrics_mv", + 'REFRESH MATERIALIZED VIEW CONCURRENTLY importer_metrics_mv', undefined, - "refresh_importer_metrics_mv", + 'refresh_importer_metrics_mv' ); } export async function getLastProcessedLedger(): Promise { const result = await timedQuery<{ last_processed_ledger: number }>( - "SELECT last_processed_ledger FROM indexer_state WHERE id = $1", - ["default"], - "select_indexer_state", + 'SELECT last_processed_ledger FROM indexer_state WHERE id = $1', + ['default'], + 'select_indexer_state' ); if (!result.rowCount || result.rowCount === 0) { return null; @@ -773,8 +790,8 @@ export async function updateLastProcessedLedger(ledger: number): Promise { ON CONFLICT (id) DO UPDATE SET last_processed_ledger = EXCLUDED.last_processed_ledger, updated_at = now()`, - ["default", ledger], - "upsert_indexer_state", + ['default', ledger], + 'upsert_indexer_state' ); } @@ -782,15 +799,21 @@ export async function updateLastProcessedLedger(ledger: number): Promise { * Pings the database to check if it's alive. */ export async function ping(): Promise { - await pool.query("SELECT 1"); + await pool.query('SELECT 1'); } /** * Returns all bonds that have been registered on-chain. */ -export async function getActiveBonds(): Promise<{ bondId: string; stellarAddress: string; dbBalance: string }[]> { - const result = await pool.query<{ bond_id: string; stellar_address: string; collateral_balance: string }>( - "SELECT bond_id, stellar_address, collateral_balance FROM importers WHERE registered_on_chain_tx IS NOT NULL" +export async function getActiveBonds(): Promise< + { bondId: string; stellarAddress: string; dbBalance: string }[] +> { + const result = await pool.query<{ + bond_id: string; + stellar_address: string; + collateral_balance: string; + }>( + 'SELECT bond_id, stellar_address, collateral_balance FROM importers WHERE registered_on_chain_tx IS NOT NULL' ); return result.rows.map((row) => ({ bondId: row.bond_id, @@ -804,54 +827,60 @@ export async function recordAuthenticationAttempt( success: boolean, userId?: string, ipAddress?: string, - userAgent?: string, + userAgent?: string ): Promise { await timedQuery( `INSERT INTO authentication_attempts (email, success, user_id, ip_address, user_agent) VALUES ($1, $2, $3, $4, $5)`, [email, success, userId ?? null, ipAddress ?? null, userAgent ?? null], - "insert_auth_attempt", + 'insert_auth_attempt' ); } -export async function getFailedAuthAttempts(email: string, withinMinutes: number = 30): Promise { +export async function getFailedAuthAttempts( + email: string, + withinMinutes: number = 30 +): Promise { const result = await timedQuery<{ count: string }>( `SELECT COUNT(*) as count FROM authentication_attempts WHERE email = $1 AND success = FALSE AND attempted_at > now() - INTERVAL '${withinMinutes} minutes'`, [email], - "count_failed_auth_attempts", + 'count_failed_auth_attempts' ); - return parseInt(result.rows[0]?.count ?? "0", 10); + return parseInt(result.rows[0]?.count ?? '0', 10); } -export async function lockAccountTemporarily(userId: string, durationMinutes: number = 30): Promise { +export async function lockAccountTemporarily( + userId: string, + durationMinutes: number = 30 +): Promise { await timedQuery( `UPDATE users SET locked_until = now() + INTERVAL '${durationMinutes} minutes' WHERE id = $1`, [userId], - "lock_account", + 'lock_account' ); } export async function recordSecurityIncident( - severity: "P0" | "P1" | "P2" | "P3", + severity: 'P0' | 'P1' | 'P2' | 'P3', description: string, - affectedScope?: string, + affectedScope?: string ): Promise { const incidentId = `INC-${Date.now()}-${Math.random().toString(36).slice(2, 9)}`; const result = await timedQuery<{ id: string }>( `INSERT INTO security_incidents (incident_id, severity, description, affected_scope) VALUES ($1, $2, $3, $4) RETURNING id`, [incidentId, severity, description, affectedScope ?? null], - "insert_security_incident", + 'insert_security_incident' ); - return result.rows[0]?.id ?? ""; + return result.rows[0]?.id ?? ''; } export async function createDataErasureRequest( userId: string, - importerId?: string, + importerId?: string ): Promise { const requestId = `ERASE-${Date.now()}-${Math.random().toString(36).slice(2, 9)}`; const slaDealine = new Date(Date.now() + 30 * 24 * 60 * 60 * 1000).toISOString(); @@ -860,9 +889,9 @@ export async function createDataErasureRequest( VALUES ($1, $2, $3, $4, ARRAY['legal_name', 'ein', 'email']) RETURNING id`, [requestId, userId, importerId ?? null, slaDealine], - "insert_erasure_request", + 'insert_erasure_request' ); - return result.rows[0]?.id ?? ""; + return result.rows[0]?.id ?? ''; } // ── SOC 2 CC6 — Session management (#306) ──────────────────────────────────── @@ -872,13 +901,13 @@ const SESSION_INACTIVITY_MINUTES = 15; export async function createSession( userId: string, ipAddress?: string, - userAgent?: string, + userAgent?: string ): Promise { const result = await timedQuery<{ id: string }>( `INSERT INTO user_sessions (user_id, ip_address, user_agent) VALUES ($1, $2, $3) RETURNING id`, [userId, ipAddress ?? null, userAgent ?? null], - "insert_user_session", + 'insert_user_session' ); return result.rows[0]!.id; } @@ -890,24 +919,24 @@ export async function validateSession(sessionId: string): Promise { AND revoked_at IS NULL AND last_activity > now() - INTERVAL '${SESSION_INACTIVITY_MINUTES} minutes'`, [sessionId], - "validate_user_session", + 'validate_user_session' ); return (result.rowCount ?? 0) > 0; } export function touchSession(sessionId: string): void { timedQuery( - "UPDATE user_sessions SET last_activity = now() WHERE id = $1 AND revoked_at IS NULL", + 'UPDATE user_sessions SET last_activity = now() WHERE id = $1 AND revoked_at IS NULL', [sessionId], - "touch_user_session", + 'touch_user_session' ).catch(() => {}); } export async function revokeSession(sessionId: string): Promise { await timedQuery( - "UPDATE user_sessions SET revoked_at = now() WHERE id = $1", + 'UPDATE user_sessions SET revoked_at = now() WHERE id = $1', [sessionId], - "revoke_user_session", + 'revoke_user_session' ); } @@ -917,9 +946,9 @@ export async function getActiveSessionCount(userId: string): Promise { WHERE user_id = $1 AND revoked_at IS NULL AND last_activity > now() - INTERVAL '${SESSION_INACTIVITY_MINUTES} minutes'`, [userId], - "count_active_sessions", + 'count_active_sessions' ); - return parseInt(result.rows[0]?.count ?? "0", 10); + return parseInt(result.rows[0]?.count ?? '0', 10); } export async function revokeOldestSession(userId: string): Promise { @@ -932,7 +961,7 @@ export async function revokeOldestSession(userId: string): Promise { LIMIT 1 )`, [userId], - "revoke_oldest_session", + 'revoke_oldest_session' ); } @@ -1011,7 +1040,7 @@ export async function revokeRefreshToken(tokenHash: string): Promise { } export async function getStaleAccounts( - days: number, + days: number ): Promise> { const result = await timedQuery<{ id: string; email: string; last_login: string | null }>( `SELECT u.id, u.email, @@ -1024,7 +1053,7 @@ export async function getStaleAccounts( OR MAX(a.attempted_at) < now() - ($1::integer * INTERVAL '1 day') ORDER BY last_login ASC NULLS FIRST`, [days], - "select_stale_accounts", + 'select_stale_accounts' ); return result.rows; } diff --git a/apps/api/src/migrate.ts b/apps/api/src/migrate.ts index ce3747c..99e5d23 100644 --- a/apps/api/src/migrate.ts +++ b/apps/api/src/migrate.ts @@ -16,3 +16,5 @@ const { migrate, pool } = await import("./db.js"); await migrate(); await pool.end(); console.log("Migrations complete."); +export {}; + diff --git a/apps/api/src/routes/admin.ts b/apps/api/src/routes/admin.ts index 866cbb3..5c0978c 100644 --- a/apps/api/src/routes/admin.ts +++ b/apps/api/src/routes/admin.ts @@ -1,6 +1,6 @@ import { Router, type Request, type Response } from "express"; import { z } from "zod"; -import { pool } from "../db.js"; +import { pool, getStaleAccounts } from "../db.js"; import { authMiddleware, requireRole, privacyReacceptanceGate, tosReacceptanceGate, type AuthedRequest } from "../auth.js"; import { platformKeypair, oracleKeypair } from "../stellar.js"; import { bustHtsCache } from "../services/hts-rate-validator.js"; diff --git a/apps/api/src/routes/importers.ts b/apps/api/src/routes/importers.ts index e97d003..18dd248 100644 --- a/apps/api/src/routes/importers.ts +++ b/apps/api/src/routes/importers.ts @@ -26,50 +26,50 @@ const CreateImporterSchema = z.object({ businessState: z.string().length(2).toUpperCase().optional(), }); -importersRouter.post("/", async (req: Request, res: Response) => { +importersRouter.post('/', async (req: Request, res: Response) => { const user = (req as AuthedRequest).user; - if (user.role !== "importer") { - res.status(403).json({ error: "only importer accounts can register" }); + if (user.role !== 'importer') { + res.status(403).json({ error: 'only importer accounts can register' }); return; } const parse = CreateImporterSchema.safeParse(req.body); if (!parse.success) { - res.status(400).json({ error: "invalid input", details: parse.error.issues }); + res.status(400).json({ error: 'invalid input', details: parse.error.issues }); return; } const { legalName, ein, bondId, initialRequiredCollateral, businessState } = parse.data; const ofacClear = await screenImporterEntity(legalName, ein); if (!ofacClear) { - res.status(403).json({ error: "Importer failed OFAC sanctions screening" }); + res.status(403).json({ error: 'Importer failed OFAC sanctions screening' }); return; } - const existing = await pool.query("SELECT id FROM importers WHERE user_id = $1", [user.id]); + const existing = await pool.query('SELECT id FROM importers WHERE user_id = $1', [user.id]); if (existing.rowCount && existing.rowCount > 0) { - res.status(409).json({ error: "importer already registered for this user" }); + res.status(409).json({ error: 'importer already registered for this user' }); return; } const kp = Keypair.random(); const amlRes = await screenWalletAddress(kp.publicKey()); - if (amlRes.riskScore === "HIGH") { - res.status(403).json({ error: "Wallet address flagged as high risk by AML provider" }); + if (amlRes.riskScore === 'HIGH') { + res.status(403).json({ error: 'Wallet address flagged as high risk by AML provider' }); return; } const bondValidation = validateBondForm301({ principalLegalName: legalName, principalEin: ein, - bondTypeCode: "02", + bondTypeCode: '02', bondAmount: BigInt(initialRequiredCollateral), }); if (!bondValidation.valid) { res.status(422).json({ - error: "Bond validation failed", + error: 'Bond validation failed', details: bondValidation.errors, }); return; @@ -79,7 +79,7 @@ importersRouter.post("/", async (req: Request, res: Response) => { `INSERT INTO importers (user_id, legal_name, ein, bond_id, stellar_address, stellar_secret_encrypted, business_state) VALUES ($1, $2, $3, $4, $5, $6, $7) RETURNING id, legal_name, ein, bond_id, stellar_address, created_at`, - [user.id, legalName, ein ?? null, bondId, kp.publicKey(), kp.secret(), businessState ?? "CA"], + [user.id, legalName, ein ?? null, bondId, kp.publicKey(), kp.secret(), businessState ?? 'CA'] ); const importer = inserted.rows[0]!; @@ -87,7 +87,21 @@ importersRouter.post("/", async (req: Request, res: Response) => { `INSERT INTO bond_records (importer_id, bond_id, bond_type_code, principal_legal_name, principal_ein, surety_company_name, surety_fein, bond_amount, cbp_minimum_required, effective_date, template_version, cbp_regulation_revision_date, state_code) VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13)`, - [importer.id, bondId, "02", legalName, ein ?? null, "TBD", "TBD", initialRequiredCollateral, bondValidation.minimumRequired.toString(), new Date(), "1.0", new Date(), businessState ?? "CA"], + [ + importer.id, + bondId, + '02', + legalName, + ein ?? null, + 'TBD', + 'TBD', + initialRequiredCollateral, + bondValidation.minimumRequired.toString(), + new Date(), + '1.0', + new Date(), + businessState ?? 'CA', + ] ); // Fund the importer account via friendbot (testnet only) @@ -95,7 +109,7 @@ importersRouter.post("/", async (req: Request, res: Response) => { const friendbotRes = await fetch(`https://friendbot.stellar.org/?addr=${kp.publicKey()}`); if (!friendbotRes.ok) throw new Error(`friendbot ${friendbotRes.status}`); } catch (err) { - console.error("[importers] friendbot fund failed:", err); + console.error('[importers] friendbot fund failed:', err); } // Register importer on-chain. Platform admin signs. @@ -103,10 +117,10 @@ importersRouter.post("/", async (req: Request, res: Response) => { platformKeypair, kp.publicKey(), BigInt(bondId), - BigInt(initialRequiredCollateral), + BigInt(initialRequiredCollateral) ); - await pool.query("UPDATE importers SET registered_on_chain_tx = $1 WHERE id = $2", [ + await pool.query('UPDATE importers SET registered_on_chain_tx = $1 WHERE id = $2', [ onChain.txHash, importer.id, ]); @@ -114,7 +128,7 @@ importersRouter.post("/", async (req: Request, res: Response) => { `INSERT INTO contract_events (importer_id, kind, tx_hash, ledger_sequence, event_index) VALUES ($1, $2, $3, $4, $5) ON CONFLICT (ledger_sequence, event_index) DO NOTHING`, - [importer.id, "register", onChain.txHash, onChain.ledgerSequence, onChain.applicationOrder], + [importer.id, 'register', onChain.txHash, onChain.ledgerSequence, onChain.applicationOrder] ); await logAudit(user.id, "register", importer.id, { legalName, bondId }); @@ -134,20 +148,20 @@ importersRouter.post("/", async (req: Request, res: Response) => { }); }); -importersRouter.get("/", async (req: Request, res: Response) => { +importersRouter.get('/', async (req: Request, res: Response) => { const user = (req as AuthedRequest).user; let r; - if (user.role === "surety_admin") { + if (user.role === 'surety_admin') { r = await pool.query( `SELECT i.id, i.legal_name, i.bond_id, i.stellar_address, i.created_at, u.email FROM importers i JOIN users u ON u.id = i.user_id - ORDER BY i.created_at DESC`, + ORDER BY i.created_at DESC` ); } else { r = await pool.query( `SELECT i.id, i.legal_name, i.bond_id, i.stellar_address, i.created_at FROM importers i WHERE i.user_id = $1`, - [user.id], + [user.id] ); } res.json({ importers: r.rows }); @@ -157,33 +171,105 @@ importersRouter.get("/", async (req: Request, res: Response) => { // (a materialized view refreshed on a 5-minute schedule — see // jobs/refresh-importer-metrics.ts) instead of live GROUP BY queries. // Registered before "/:id" so Express doesn't treat "stats" as an :id param. -importersRouter.get("/stats", async (req: Request, res: Response) => { +importersRouter.get('/stats', async (req: Request, res: Response) => { const user = (req as AuthedRequest).user; - if (user.role !== "surety_admin") { - res.status(403).json({ error: "surety admin only" }); + if (user.role !== 'surety_admin') { + res.status(403).json({ error: 'surety admin only' }); return; } const metrics = await getImporterMetrics(); res.json({ metrics }); }); +const AdminEventsQuerySchema = z.object({ + from: z.string().datetime({ message: "Invalid 'from' timestamp, must be ISO 8601" }).optional(), + to: z.string().datetime({ message: "Invalid 'to' timestamp, must be ISO 8601" }).optional(), + limit: z.coerce.number().int().positive().max(500).default(500), + offset: z.coerce.number().int().nonnegative().default(0), +}); + +// #245: Admin events endpoint for time-range reporting. +// Fetches contract events within a specified range, optimized via the BRIN index. +importersRouter.get('/admin/events', async (req: Request, res: Response) => { + const user = (req as AuthedRequest).user; + if (user.role !== 'surety_admin') { + res.status(403).json({ error: 'surety admin only' }); + return; + } + + const parse = AdminEventsQuerySchema.safeParse(req.query); + if (!parse.success) { + res.status(400).json({ error: 'invalid query parameters', details: parse.error.issues }); + return; + } + + const { limit, offset, from, to } = parse.data; + + const queryParams: any[] = []; + let sql = ` + SELECT id, importer_id, kind, amount, tx_hash, created_at, ledger_sequence, event_index + FROM contract_events + `; + const conditions: string[] = []; + + if (from) { + queryParams.push(from); + conditions.push(`created_at >= $${queryParams.length}`); + } + if (to) { + queryParams.push(to); + conditions.push(`created_at <= $${queryParams.length}`); + } + + if (conditions.length > 0) { + sql += ' WHERE ' + conditions.join(' AND '); + } + + sql += ' ORDER BY created_at DESC, id DESC'; + + queryParams.push(limit); + sql += ` LIMIT $${queryParams.length}`; + + queryParams.push(offset); + sql += ` OFFSET $${queryParams.length}`; + + try { + const result = await pool.query(sql, queryParams); + const events = result.rows.map((row) => ({ + id: row.id, + importerId: row.importer_id, + kind: row.kind, + amount: row.amount, + txHash: row.tx_hash, + txUrl: row.tx_hash ? explorerTx(row.tx_hash) : null, + createdAt: row.created_at, + ledgerSequence: row.ledger_sequence, + eventIndex: row.event_index, + })); + res.json({ events }); + } catch (err: any) { + console.error('[importers] Failed to query admin events:', err); + res.status(500).json({ error: 'failed to retrieve events' }); + } +}); + async function loadImporterFor(req: Request, importerId: string) { const user = (req as AuthedRequest).user; - if (user.role === "surety_admin") { - const r = await pool.query("SELECT * FROM importers WHERE id = $1", [importerId]); + if (user.role === 'surety_admin') { + const r = await pool.query('SELECT * FROM importers WHERE id = $1', [importerId]); return r.rows[0] ?? null; } - const r = await pool.query("SELECT * FROM importers WHERE id = $1 AND user_id = $2", [ + const r = await pool.query('SELECT * FROM importers WHERE id = $1 AND user_id = $2', [ importerId, user.id, ]); return r.rows[0] ?? null; } -importersRouter.get("/:id", async (req: Request, res: Response) => { - const importer = await loadImporterFor(req, String(req.params.id ?? "")); +importersRouter.get('/:id', async (req: Request, res: Response) => { + const importer = await loadImporterFor(req, String(req.params.id ?? '')); if (!importer) { - res.status(404).json({ error: "not found" }); + res.status(404).json({ error: 'not found' }); return; } const acct = await contractClient.getAccount(importer.stellar_address); @@ -217,8 +303,8 @@ importersRouter.get("/:id", async (req: Request, res: Response) => { // events share the same created_at timestamp. function decodeEventsCursor(raw: string): { createdAt: string; id: string } | null { try { - const decoded = Buffer.from(raw, "base64").toString("utf8"); - const sep = decoded.lastIndexOf("|"); + const decoded = Buffer.from(raw, 'base64').toString('utf8'); + const sep = decoded.lastIndexOf('|'); if (sep === -1) return null; return { createdAt: decoded.slice(0, sep), id: decoded.slice(sep + 1) }; } catch { @@ -227,7 +313,7 @@ function decodeEventsCursor(raw: string): { createdAt: string; id: string } | nu } function encodeEventsCursor(createdAt: Date, id: string): string { - return Buffer.from(`${createdAt.toISOString()}|${id}`, "utf8").toString("base64"); + return Buffer.from(`${createdAt.toISOString()}|${id}`, 'utf8').toString('base64'); } const EventsQuerySchema = z.object({ @@ -235,16 +321,16 @@ const EventsQuerySchema = z.object({ limit: z.coerce.number().int().positive().max(100).default(20), }); -importersRouter.get("/:id/events", async (req: Request, res: Response) => { - const importer = await loadImporterFor(req, String(req.params.id ?? "")); +importersRouter.get('/:id/events', async (req: Request, res: Response) => { + const importer = await loadImporterFor(req, String(req.params.id ?? '')); if (!importer) { - res.status(404).json({ error: "not found" }); + res.status(404).json({ error: 'not found' }); return; } const parse = EventsQuerySchema.safeParse(req.query); if (!parse.success) { - res.status(400).json({ error: "invalid query", details: parse.error.issues }); + res.status(400).json({ error: 'invalid query', details: parse.error.issues }); return; } const { limit } = parse.data; @@ -253,7 +339,7 @@ importersRouter.get("/:id/events", async (req: Request, res: Response) => { if (parse.data.cursor) { cursor = decodeEventsCursor(parse.data.cursor); if (!cursor) { - res.status(400).json({ error: "invalid cursor" }); + res.status(400).json({ error: 'invalid cursor' }); return; } } @@ -263,13 +349,13 @@ importersRouter.get("/:id/events", async (req: Request, res: Response) => { `SELECT id, kind, amount, tx_hash, created_at FROM contract_events WHERE importer_id = $1 AND (created_at, id) < ($2::timestamptz, $3::uuid) ORDER BY created_at DESC, id DESC LIMIT $4`, - [importer.id, cursor.createdAt, cursor.id, limit], + [importer.id, cursor.createdAt, cursor.id, limit] ) : await pool.query( `SELECT id, kind, amount, tx_hash, created_at FROM contract_events WHERE importer_id = $1 ORDER BY created_at DESC, id DESC LIMIT $2`, - [importer.id, limit], + [importer.id, limit] ); const events = rows.rows.map((e) => ({ @@ -288,10 +374,10 @@ importersRouter.get("/:id/events", async (req: Request, res: Response) => { res.json({ events, nextCursor }); }); -importersRouter.get("/:id/collateral-status", async (req: Request, res: Response) => { - const importer = await loadImporterFor(req, String(req.params.id ?? "")); +importersRouter.get('/:id/collateral-status', async (req: Request, res: Response) => { + const importer = await loadImporterFor(req, String(req.params.id ?? '')); if (!importer) { - res.status(404).json({ error: "not found" }); + res.status(404).json({ error: 'not found' }); return; } const acct = await contractClient.getAccount(importer.stellar_address); @@ -305,7 +391,6 @@ importersRouter.get("/:id/collateral-status", async (req: Request, res: Response }); }); - // --- Synthetic CBP tariff CSV upload — recomputes required_collateral on-chain --- const TariffLineItemSchema = z.object({ @@ -323,12 +408,12 @@ importersRouter.post("/:id/upload-tariff-csv", async (req: Request, res: Respons const user = (req as AuthedRequest).user; const importer = await loadImporterFor(req, String(req.params.id ?? "")); if (!importer) { - res.status(404).json({ error: "not found" }); + res.status(404).json({ error: 'not found' }); return; } const parse = TariffUploadSchema.safeParse(req.body); if (!parse.success) { - res.status(400).json({ error: "invalid input", details: parse.error.issues }); + res.status(400).json({ error: 'invalid input', details: parse.error.issues }); return; } @@ -339,12 +424,12 @@ importersRouter.post("/:id/upload-tariff-csv", async (req: Request, res: Respons parse.data.lineItems.map((item) => ({ hts_code: item.htsCode, declared_rate: item.dutyRate, - })), + })) ); if (htsValidation.hasBlockingErrors) { res.status(422).json({ - error: "HTS rate validation failed: one or more line items are underreported", + error: 'HTS rate validation failed: one or more line items are underreported', flagged: htsValidation.blocking.map((r) => ({ htsCode: r.hts_code, declaredRate: r.declared_rate, @@ -372,16 +457,16 @@ importersRouter.post("/:id/upload-tariff-csv", async (req: Request, res: Respons for (const item of parse.data.lineItems) { const cbpRes = await lookupCbpDutyRate(item.htsCode); const cbpRate = cbpRes.dutyRate ?? item.dutyRate; - + const deviation = Math.abs(cbpRate - item.dutyRate); - if (cbpRate > 0 && deviation / cbpRate > 0.10) { + if (cbpRate > 0 && deviation / cbpRate > 0.1) { validationReport.push({ htsCode: item.htsCode, reportedRate: item.dutyRate, cbpRate: cbpRate, deviation, }); - if (env.CBP_VALIDATION_MODE !== "warn") { + if (env.CBP_VALIDATION_MODE !== 'warn') { hasBlockError = true; } } @@ -389,7 +474,7 @@ importersRouter.post("/:id/upload-tariff-csv", async (req: Request, res: Respons } if (hasBlockError) { - res.status(422).json({ error: "CBP validation failed", report: validationReport }); + res.status(422).json({ error: 'CBP validation failed', report: validationReport }); return; } @@ -406,17 +491,29 @@ importersRouter.post("/:id/upload-tariff-csv", async (req: Request, res: Respons importer.stellar_address, requiredStroops, env.PRICE_ORACLE_CONTRACT_ID, - false, + false ); await pool.query( - "INSERT INTO tariff_uploads (importer_id, filename, annual_duty_total, computed_required_collateral, applied_tx) VALUES ($1, $2, $3, $4, $5)", - [importer.id, parse.data.filename ?? null, annualDutyTotal, requiredStroops.toString(), onChain.txHash], + 'INSERT INTO tariff_uploads (importer_id, filename, annual_duty_total, computed_required_collateral, applied_tx) VALUES ($1, $2, $3, $4, $5)', + [ + importer.id, + parse.data.filename ?? null, + annualDutyTotal, + requiredStroops.toString(), + onChain.txHash, + ] ); await pool.query( `INSERT INTO contract_events (importer_id, kind, amount, tx_hash, ledger_sequence, event_index) VALUES ($1, 'required_changed', $2, $3, $4, $5) ON CONFLICT (ledger_sequence, event_index) DO NOTHING`, - [importer.id, requiredStroops.toString(), onChain.txHash, onChain.ledgerSequence, onChain.applicationOrder], + [ + importer.id, + requiredStroops.toString(), + onChain.txHash, + onChain.ledgerSequence, + onChain.applicationOrder, + ] ); await logAudit(user.id, "apply_tariff_upload", importer.id, { filename: parse.data.filename, annualDutyTotal, requiredStroops: requiredStroops.toString() }); @@ -431,15 +528,13 @@ importersRouter.post("/:id/upload-tariff-csv", async (req: Request, res: Respons }); } catch (err: any) { const errMsg = String(err); - if (errMsg.includes("Error(Contract, #13)") || errMsg.includes("RateLimitExceeded")) { + if (errMsg.includes('Error(Contract, #13)') || errMsg.includes('RateLimitExceeded')) { const retryAfter = Math.ceil(Date.now() / 1000) + 86400; - res.status(429) - .set("Retry-After", String(retryAfter)) - .json({ - error: "rate limit exceeded", - retryAfter, - message: "collateral requirements can only be updated once per 24 hours", - }); + res.status(429).set('Retry-After', String(retryAfter)).json({ + error: 'rate limit exceeded', + retryAfter, + message: 'collateral requirements can only be updated once per 24 hours', + }); return; } throw err; @@ -448,21 +543,21 @@ importersRouter.post("/:id/upload-tariff-csv", async (req: Request, res: Respons const DepositSchema = z.object({ amountStroops: z.string().regex(/^\d+$/), - bucket: z.enum(["collateral", "reserve"]), + bucket: z.enum(['collateral', 'reserve']), }); importersRouter.post("/:id/deposit", async (req: Request, res: Response) => { const user = (req as AuthedRequest).user; const importer = await loadImporterFor(req, String(req.params.id ?? "")); if (!importer) { - res.status(404).json({ error: "not found" }); + res.status(404).json({ error: 'not found' }); return; } // #312: block collateral deposits until KYC is approved (CIP compliance) - if (importer.kyc_status !== "approved") { + if (importer.kyc_status !== 'approved') { res.status(403).json({ - error: "KYC approval required before collateral deposits", + error: 'KYC approval required before collateral deposits', kycStatus: importer.kyc_status, }); return; @@ -470,18 +565,18 @@ importersRouter.post("/:id/deposit", async (req: Request, res: Response) => { const parse = DepositSchema.safeParse(req.body); if (!parse.success) { - res.status(400).json({ error: "invalid input" }); + res.status(400).json({ error: 'invalid input' }); return; } const amlRes = await screenWalletAddress(importer.stellar_address); - if (amlRes.riskScore === "HIGH") { - res.status(403).json({ error: "Transaction blocked pending AML review" }); + if (amlRes.riskScore === 'HIGH') { + res.status(403).json({ error: 'Transaction blocked pending AML review' }); return; } const jobId = await enqueueTxSubmit({ - method: "deposit", + method: 'deposit', importerId: importer.id, keypairSecret: importer.stellar_secret_encrypted, args: { @@ -495,14 +590,14 @@ importersRouter.post("/:id/deposit", async (req: Request, res: Response) => { res.status(202).json({ jobId, statusUrl: `/importers/${importer.id}/tx-status/${jobId}` }); }); -importersRouter.post("/:id/auto-top-up", async (req: Request, res: Response) => { - const importer = await loadImporterFor(req, String(req.params.id ?? "")); +importersRouter.post('/:id/auto-top-up', async (req: Request, res: Response) => { + const importer = await loadImporterFor(req, String(req.params.id ?? '')); if (!importer) { - res.status(404).json({ error: "not found" }); + res.status(404).json({ error: 'not found' }); return; } const jobId = await enqueueTxSubmit({ - method: "auto_top_up", + method: 'auto_top_up', importerId: importer.id, platformKey: true, args: { @@ -520,23 +615,23 @@ importersRouter.post("/:id/withdraw", async (req: Request, res: Response) => { const user = (req as AuthedRequest).user; const importer = await loadImporterFor(req, String(req.params.id ?? "")); if (!importer) { - res.status(404).json({ error: "not found" }); + res.status(404).json({ error: 'not found' }); return; } const parse = WithdrawSchema.safeParse(req.body); if (!parse.success) { - res.status(400).json({ error: "invalid input" }); + res.status(400).json({ error: 'invalid input' }); return; } const amlRes = await screenWalletAddress(importer.stellar_address); - if (amlRes.riskScore === "HIGH") { - res.status(403).json({ error: "Transaction blocked pending AML review" }); + if (amlRes.riskScore === 'HIGH') { + res.status(403).json({ error: 'Transaction blocked pending AML review' }); return; } const jobId = await enqueueTxSubmit({ - method: "withdraw", + method: 'withdraw', importerId: importer.id, keypairSecret: importer.stellar_secret_encrypted, args: { @@ -553,55 +648,63 @@ importersRouter.post("/:id/withdraw", async (req: Request, res: Response) => { const YieldSchema = z.object({ amountStroops: z.string().regex(/^\d+$/) }); -importersRouter.post("/:id/accrue-yield", requireLicenseVerified, async (req: Request, res: Response) => { - const user = (req as AuthedRequest).user; - if (user.role !== "surety_admin") { - res.status(403).json({ error: "surety admin only" }); - return; - } - const importer = await loadImporterFor(req, String(req.params.id ?? "")); - if (!importer) { - res.status(404).json({ error: "not found" }); - return; - } - const parse = YieldSchema.safeParse(req.body); - if (!parse.success) { - res.status(400).json({ error: "invalid input" }); - return; +importersRouter.post( + '/:id/accrue-yield', + requireLicenseVerified, + async (req: Request, res: Response) => { + const user = (req as AuthedRequest).user; + if (user.role !== 'surety_admin') { + res.status(403).json({ error: 'surety admin only' }); + return; + } + const importer = await loadImporterFor(req, String(req.params.id ?? '')); + if (!importer) { + res.status(404).json({ error: 'not found' }); + return; + } + const parse = YieldSchema.safeParse(req.body); + if (!parse.success) { + res.status(400).json({ error: 'invalid input' }); + return; + } + const jobId = await enqueueTxSubmit({ + method: 'accrue_yield', + importerId: importer.id, + platformKey: true, + args: { + importerAddress: importer.stellar_address, + amountStroops: parse.data.amountStroops, + }, + }); + res.status(202).json({ jobId, statusUrl: `/importers/${importer.id}/tx-status/${jobId}` }); } - const jobId = await enqueueTxSubmit({ - method: "accrue_yield", - importerId: importer.id, - platformKey: true, - args: { - importerAddress: importer.stellar_address, - amountStroops: parse.data.amountStroops, - }, - }); - res.status(202).json({ jobId, statusUrl: `/importers/${importer.id}/tx-status/${jobId}` }); -}); +); -importersRouter.post("/:id/clawback", requireLicenseVerified, async (req: Request, res: Response) => { - const user = (req as AuthedRequest).user; - if (user.role !== "surety_admin") { - res.status(403).json({ error: "surety admin only" }); - return; - } - const importer = await loadImporterFor(req, String(req.params.id ?? "")); - if (!importer) { - res.status(404).json({ error: "not found" }); - return; +importersRouter.post( + '/:id/clawback', + requireLicenseVerified, + async (req: Request, res: Response) => { + const user = (req as AuthedRequest).user; + if (user.role !== 'surety_admin') { + res.status(403).json({ error: 'surety admin only' }); + return; + } + const importer = await loadImporterFor(req, String(req.params.id ?? '')); + if (!importer) { + res.status(404).json({ error: 'not found' }); + return; + } + const jobId = await enqueueTxSubmit({ + method: 'clawback', + importerId: importer.id, + suretyKey: true, + args: { + importerAddress: importer.stellar_address, + }, + }); + res.status(202).json({ jobId, statusUrl: `/importers/${importer.id}/tx-status/${jobId}` }); } - const jobId = await enqueueTxSubmit({ - method: "clawback", - importerId: importer.id, - suretyKey: true, - args: { - importerAddress: importer.stellar_address, - }, - }); - res.status(202).json({ jobId, statusUrl: `/importers/${importer.id}/tx-status/${jobId}` }); -}); +); // ── Issue #335: Oracle data reconciliation endpoint ─────────────────────────── @@ -609,43 +712,46 @@ const VerifyOracleSchema = z.object({ as_of_date: z.string().datetime().optional(), }); -importersRouter.post("/:id/verify-oracle-data", async (req: Request, res: Response) => { +importersRouter.post('/:id/verify-oracle-data', async (req: Request, res: Response) => { const user = (req as AuthedRequest).user; - const importerId = String(req.params.id ?? ""); + const importerId = String(req.params.id ?? ''); // Accessible by: the importer themselves, surety_admin, or platform admin (surety_admin covers both) let importer: Record | null = null; - if (user.role === "surety_admin") { - const r = await pool.query("SELECT * FROM importers WHERE id = $1", [importerId]); + if (user.role === 'surety_admin') { + const r = await pool.query('SELECT * FROM importers WHERE id = $1', [importerId]); importer = r.rows[0] ?? null; } else { - const r = await pool.query("SELECT * FROM importers WHERE id = $1 AND user_id = $2", [importerId, user.id]); + const r = await pool.query('SELECT * FROM importers WHERE id = $1 AND user_id = $2', [ + importerId, + user.id, + ]); importer = r.rows[0] ?? null; } if (!importer) { - res.status(404).json({ error: "not found" }); + res.status(404).json({ error: 'not found' }); return; } const parse = VerifyOracleSchema.safeParse(req.body); if (!parse.success) { - res.status(400).json({ error: "invalid input", details: parse.error.issues }); + res.status(400).json({ error: 'invalid input', details: parse.error.issues }); return; } // Fetch latest tariff upload for this importer const uploadQ = parse.data.as_of_date ? await pool.query( - "SELECT * FROM tariff_uploads WHERE importer_id = $1 AND created_at <= $2 ORDER BY created_at DESC LIMIT 1", - [importerId, parse.data.as_of_date], + 'SELECT * FROM tariff_uploads WHERE importer_id = $1 AND created_at <= $2 ORDER BY created_at DESC LIMIT 1', + [importerId, parse.data.as_of_date] ) : await pool.query( - "SELECT * FROM tariff_uploads WHERE importer_id = $1 ORDER BY created_at DESC LIMIT 1", - [importerId], + 'SELECT * FROM tariff_uploads WHERE importer_id = $1 ORDER BY created_at DESC LIMIT 1', + [importerId] ); if (!uploadQ.rowCount || uploadQ.rowCount === 0) { - res.status(404).json({ error: "no tariff CSV data found for this importer" }); + res.status(404).json({ error: 'no tariff CSV data found for this importer' }); return; } const upload = uploadQ.rows[0]!; @@ -655,8 +761,8 @@ importersRouter.post("/:id/verify-oracle-data", async (req: Request, res: Respon const computed = BigInt(Math.round(annualDuty * 0.1 * 0.5 * 1e7)); // CSV hash — hash the stored annual_duty_total + filename as a stable fingerprint - const csvFingerprint = `${upload.filename ?? ""}:${upload.annual_duty_total}`; - const csvHash = createHash("sha256").update(csvFingerprint).digest("hex"); + const csvFingerprint = `${upload.filename ?? ''}:${upload.annual_duty_total}`; + const csvHash = createHash('sha256').update(csvFingerprint).digest('hex'); // Fetch on-chain value const onChainStr = await getRequiredCollateralOnChain(importer.stellar_address as string); @@ -664,9 +770,12 @@ importersRouter.post("/:id/verify-oracle-data", async (req: Request, res: Respon const computedNum = Number(computed); const onChainNum = Number(onChain); - const deviationPct = onChainNum === 0 - ? (computedNum === 0 ? 0 : 100) - : Math.abs(computedNum - onChainNum) / onChainNum * 100; + const deviationPct = + onChainNum === 0 + ? computedNum === 0 + ? 0 + : 100 + : (Math.abs(computedNum - onChainNum) / onChainNum) * 100; const match = deviationPct <= 1.0; @@ -676,7 +785,13 @@ importersRouter.post("/:id/verify-oracle-data", async (req: Request, res: Respon `INSERT INTO oracle_alerts (importer_id, old_value, new_value, pct_change, tx_hash) VALUES ($1, $2, $3, $4, $5) ON CONFLICT DO NOTHING`, - [importerId, onChainStr, computed.toString(), deviationPct.toFixed(2), "reconciliation_failure"], + [ + importerId, + onChainStr, + computed.toString(), + deviationPct.toFixed(2), + 'reconciliation_failure', + ] ); } @@ -690,16 +805,16 @@ importersRouter.post("/:id/verify-oracle-data", async (req: Request, res: Respon }); }); -importersRouter.get("/:id/tx-status/:jobId", async (req: Request, res: Response) => { - const importer = await loadImporterFor(req, String(req.params.id ?? "")); +importersRouter.get('/:id/tx-status/:jobId', async (req: Request, res: Response) => { + const importer = await loadImporterFor(req, String(req.params.id ?? '')); if (!importer) { - res.status(404).json({ error: "not found" }); + res.status(404).json({ error: 'not found' }); return; } - const job = await txSubmitQueue.getJob(String(req.params.jobId ?? "")); + const job = await txSubmitQueue.getJob(String(req.params.jobId ?? '')); if (!job) { - res.status(404).json({ error: "job not found" }); + res.status(404).json({ error: 'job not found' }); return; } @@ -708,9 +823,9 @@ importersRouter.get("/:id/tx-status/:jobId", async (req: Request, res: Response) const result = job.returnvalue; const failedReason = job.failedReason; - if (state === "completed") { + if (state === 'completed') { res.json({ state, result }); - } else if (state === "failed") { + } else if (state === 'failed') { res.status(400).json({ state, error: failedReason }); } else { res.json({ state, progress }); diff --git a/apps/api/src/services/oracle-event-listener.ts b/apps/api/src/services/oracle-event-listener.ts index 099f8f5..7b8ab30 100644 --- a/apps/api/src/services/oracle-event-listener.ts +++ b/apps/api/src/services/oracle-event-listener.ts @@ -22,19 +22,19 @@ * The listener handles both event shapes. */ -import pino from "pino"; -import * as Sentry from "@sentry/node"; -import { rpc, scValToNative, xdr } from "@stellar/stellar-sdk"; -import { pool } from "../db.js"; -import { env } from "../config/env.js"; -import { createRpcServer } from "../lib/soroban/rpcClient.js"; +import pino from 'pino'; +import * as Sentry from '@sentry/node'; +import { rpc, scValToNative } from '@stellar/stellar-sdk'; +import { pool } from '../db.js'; +import { env } from '../config/env.js'; +import { createRpcServer } from '../lib/soroban/rpcClient.js'; -const logger = pino({ name: "oracle-event-listener" }); +const logger = pino({ name: 'oracle-event-listener' }); // ── Constants ───────────────────────────────────────────────────────────────── /** Postgres primary key used in listener_state. */ -const STATE_KEY = "oracle_event_listener"; +const STATE_KEY = 'oracle_event_listener'; /** How often to poll for new events (ms). */ const POLL_INTERVAL_MS = 12_000; // ~2 Stellar ledgers at 5 s/ledger @@ -65,8 +65,8 @@ export interface OracleFeedRow { export async function getListenerState(): Promise { const r = await pool.query<{ last_ledger_sequence: number }>( - "SELECT last_ledger_sequence FROM listener_state WHERE id = $1", - [STATE_KEY], + 'SELECT last_ledger_sequence FROM listener_state WHERE id = $1', + [STATE_KEY] ); return r.rows[0]?.last_ledger_sequence ?? null; } @@ -78,7 +78,7 @@ export async function setListenerState(ledgerSequence: number): Promise { ON CONFLICT (id) DO UPDATE SET last_ledger_sequence = EXCLUDED.last_ledger_sequence, updated_at = now()`, - [STATE_KEY, ledgerSequence], + [STATE_KEY, ledgerSequence] ); } @@ -99,7 +99,7 @@ interface ParsedOracleEvent { * Attempt to extract oracle event data from a raw Soroban event. * Returns null when the event does not match either expected shape. */ -function parseOracleEvent(event: rpc.Api.EventRecord): ParsedOracleEvent | null { +function parseOracleEvent(event: any): ParsedOracleEvent | null { try { // Topics are XDR-encoded ScVal strings in the API response. const topics = event.topic; // ScVal[] @@ -108,14 +108,14 @@ function parseOracleEvent(event: rpc.Api.EventRecord): ParsedOracleEvent | null // Topic[0] is the event name symbol. const topicSymbol = scValToNative(topics[0]!) as unknown; - const isNormal = topicSymbol === "required"; - const isEmergency = topicSymbol === "EmergencyOracleUpdate"; + const isNormal = topicSymbol === 'required'; + const isEmergency = topicSymbol === 'EmergencyOracleUpdate'; if (!isNormal && !isEmergency) return null; // Topic[1] is the importer Address ScVal. const importerAddress = scValToNative(topics[1]!) as string; - if (!importerAddress || typeof importerAddress !== "string") return null; + if (!importerAddress || typeof importerAddress !== 'string') return null; // Data is a tuple ScVal. const dataVal = event.value; // ScVal @@ -128,7 +128,7 @@ function parseOracleEvent(event: rpc.Api.EventRecord): ParsedOracleEvent | null const newRequired = BigInt(String(dataNative[1])); // Emergency events have (old, new, ts, caller) in the data tuple. - let callerAddress = ""; + let callerAddress = ''; if (isEmergency && dataNative.length >= 4) { callerAddress = String(dataNative[3]); } @@ -143,7 +143,7 @@ function parseOracleEvent(event: rpc.Api.EventRecord): ParsedOracleEvent | null ledgerSequence: event.ledger, }; } catch (err) { - logger.warn({ err, eventId: event.id }, "Failed to parse oracle event"); + logger.warn({ err, eventId: event.id }, 'Failed to parse oracle event'); return null; } } @@ -158,8 +158,8 @@ function parseOracleEvent(event: rpc.Api.EventRecord): ParsedOracleEvent | null export async function insertOracleFeedRow(parsed: ParsedOracleEvent): Promise { // Resolve importer_id (nullable — address may not exist in the DB yet). const importerRow = await pool.query<{ id: string }>( - "SELECT id FROM importers WHERE stellar_address = $1", - [parsed.importerAddress], + 'SELECT id FROM importers WHERE stellar_address = $1', + [parsed.importerAddress] ); const importerId: string | null = importerRow.rows[0]?.id ?? null; @@ -187,7 +187,7 @@ export async function insertOracleFeedRow(parsed: ParsedOracleEvent): Promise { fromLedger = Math.max(1, currentLedger - INITIAL_LOOKBACK_LEDGERS); logger.info( { fromLedger, currentLedger }, - "[oracle-listener] No checkpoint found — replaying from initial lookback", + '[oracle-listener] No checkpoint found — replaying from initial lookback' ); } @@ -221,22 +221,19 @@ export async function pollOracleEvents(rpcServer: rpc.Server): Promise { // Cap the window to avoid overloading the RPC node. const toLedger = Math.min(currentLedger, fromLedger + MAX_LEDGER_WINDOW); - logger.debug( - { fromLedger, toLedger, currentLedger }, - "[oracle-listener] Polling events", - ); + logger.debug({ fromLedger, toLedger, currentLedger }, '[oracle-listener] Polling events'); const response = await rpcServer.getEvents({ startLedger: fromLedger + 1, filters: [ { - type: "contract", + type: 'contract', contractIds: [env.TARIFF_SHIELD_CONTRACT_ID], topics: [ // Normal oracle update: ["required", ] - ["required", "*"], + ['required', '*'], // Emergency oracle update: ["EmergencyOracleUpdate", ] - ["EmergencyOracleUpdate", "*"], + ['EmergencyOracleUpdate', '*'], ], }, ], @@ -262,7 +259,7 @@ export async function pollOracleEvents(rpcServer: rpc.Server): Promise { } catch (err) { logger.error( { err, txHash: parsed.txHash, importer: parsed.importerAddress }, - "[oracle-listener] Failed to insert feed row", + '[oracle-listener] Failed to insert feed row' ); Sentry.captureException(err); } @@ -273,7 +270,7 @@ export async function pollOracleEvents(rpcServer: rpc.Server): Promise { if (inserted > 0 || skipped > 0) { logger.info( { fromLedger: fromLedger + 1, toLedger, inserted, skipped }, - "[oracle-listener] Poll cycle complete", + '[oracle-listener] Poll cycle complete' ); } } @@ -284,11 +281,11 @@ let intervalId: NodeJS.Timeout | null = null; export async function startOracleEventListener(): Promise { if (intervalId) { - logger.warn("[oracle-listener] Already running"); + logger.warn('[oracle-listener] Already running'); return; } - logger.info("[oracle-listener] Starting oracle price feed event listener"); + logger.info('[oracle-listener] Starting oracle price feed event listener'); const rpcServer = createRpcServer(env.STELLAR_RPC_URL); @@ -296,7 +293,7 @@ export async function startOracleEventListener(): Promise { try { await pollOracleEvents(rpcServer); } catch (err) { - logger.error({ err }, "[oracle-listener] First poll failed"); + logger.error({ err }, '[oracle-listener] First poll failed'); Sentry.captureException(err); } @@ -304,7 +301,7 @@ export async function startOracleEventListener(): Promise { try { await pollOracleEvents(rpcServer); } catch (err) { - logger.error({ err }, "[oracle-listener] Poll cycle error"); + logger.error({ err }, '[oracle-listener] Poll cycle error'); Sentry.captureException(err); } }, POLL_INTERVAL_MS); @@ -314,6 +311,6 @@ export function stopOracleEventListener(): void { if (intervalId) { clearInterval(intervalId); intervalId = null; - logger.info("[oracle-listener] Stopped"); + logger.info('[oracle-listener] Stopped'); } } diff --git a/apps/api/src/stellar.ts b/apps/api/src/stellar.ts index fd36c09..4631dc5 100644 --- a/apps/api/src/stellar.ts +++ b/apps/api/src/stellar.ts @@ -1,11 +1,11 @@ -import { Keypair } from "@stellar/stellar-sdk"; -import { TariffShieldClient } from "@tariffshield/sdk"; -import client from "prom-client"; -import { trace, SpanStatusCode } from "@opentelemetry/api"; -import { env } from "./config/env.js"; -import { createRpcServer } from "./lib/soroban/rpcClient.js"; +import { Keypair } from '@stellar/stellar-sdk'; +import { TariffShieldClient } from '@tariffshield/sdk'; +import client from 'prom-client'; +import { trace, SpanStatusCode } from '@opentelemetry/api'; +import { env } from './config/env.js'; +import { createRpcServer } from './lib/soroban/rpcClient.js'; -const tracer = trace.getTracer("tariffshield-stellar"); +const tracer = trace.getTracer('tariffshield-stellar'); // #339 — general admin (registration, withdrawals, upgrades) export const platformKeypair = Keypair.fromSecret(env.PLATFORM_STELLAR_SECRET); @@ -21,15 +21,15 @@ export const emergencyOracleKeypair = env.EMERGENCY_ADMIN_SECRET : platformKeypair; export const sorobanRpcCallsTotal = new client.Counter({ - name: "soroban_rpc_calls_total", - help: "Total number of Soroban RPC calls made", - labelNames: ["method", "success"], + name: 'soroban_rpc_calls_total', + help: 'Total number of Soroban RPC calls made', + labelNames: ['method', 'success'], }); export const sorobanRpcDurationSeconds = new client.Histogram({ - name: "soroban_rpc_duration_seconds", - help: "Duration of Soroban RPC calls in seconds", - labelNames: ["method"], + name: 'soroban_rpc_duration_seconds', + help: 'Duration of Soroban RPC calls in seconds', + labelNames: ['method'], buckets: [0.005, 0.01, 0.025, 0.05, 0.1, 0.25, 0.5, 1, 2.5, 5], }); @@ -45,27 +45,27 @@ const baseClient = new TariffShieldClient({ export const contractClient = new Proxy(baseClient, { get(target, prop, receiver) { const original = Reflect.get(target, prop, receiver); - if (typeof original === "function") { + if (typeof original === 'function') { return async (...args: any[]) => { const methodName = String(prop); return tracer.startActiveSpan(`soroban.rpc.${methodName}`, async (span) => { span.setAttributes({ - "soroban.method": methodName, - "soroban.network": env.STELLAR_NETWORK_PASSPHRASE, + 'soroban.method': methodName, + 'soroban.network': env.STELLAR_NETWORK_PASSPHRASE, }); const start = process.hrtime(); try { const result = await original.apply(target, args); const diff = process.hrtime(start); const duration = diff[0] + diff[1] / 1e9; - sorobanRpcCallsTotal.inc({ method: methodName, success: "true" }); + sorobanRpcCallsTotal.inc({ method: methodName, success: 'true' }); sorobanRpcDurationSeconds.observe({ method: methodName }, duration); span.setStatus({ code: SpanStatusCode.OK }); return result; } catch (err) { const diff = process.hrtime(start); const duration = diff[0] + diff[1] / 1e9; - sorobanRpcCallsTotal.inc({ method: methodName, success: "false" }); + sorobanRpcCallsTotal.inc({ method: methodName, success: 'false' }); sorobanRpcDurationSeconds.observe({ method: methodName }, duration); span.setStatus({ code: SpanStatusCode.ERROR, message: String(err) }); throw err; @@ -84,25 +84,25 @@ export const explorerTx = (hash: string): string => export async function getCurrentLedgerSequence(): Promise { const server = createRpcServer(env.STELLAR_RPC_URL); - const methodName = "getLatestLedger"; + const methodName = 'getLatestLedger'; return tracer.startActiveSpan(`soroban.rpc.${methodName}`, async (span) => { span.setAttributes({ - "soroban.method": methodName, - "soroban.network": env.STELLAR_NETWORK_PASSPHRASE, + 'soroban.method': methodName, + 'soroban.network': env.STELLAR_NETWORK_PASSPHRASE, }); const start = process.hrtime(); try { const latest = await server.getLatestLedger(); const diff = process.hrtime(start); const duration = diff[0] + diff[1] / 1e9; - sorobanRpcCallsTotal.inc({ method: methodName, success: "true" }); + sorobanRpcCallsTotal.inc({ method: methodName, success: 'true' }); sorobanRpcDurationSeconds.observe({ method: methodName }, duration); span.setStatus({ code: SpanStatusCode.OK }); return latest.sequence; } catch (err) { const diff = process.hrtime(start); const duration = diff[0] + diff[1] / 1e9; - sorobanRpcCallsTotal.inc({ method: methodName, success: "false" }); + sorobanRpcCallsTotal.inc({ method: methodName, success: 'false' }); sorobanRpcDurationSeconds.observe({ method: methodName }, duration); span.setStatus({ code: SpanStatusCode.ERROR, message: String(err) }); throw err; @@ -131,16 +131,29 @@ export async function getBondOnChain(stellarAddress: string): Promise { return acct.collateralBalance.toString(); } +/** + * Retrieves the current required collateral for a bond directly from the Soroban contract. + * @param stellarAddress The importer's Stellar address. + * @returns The on-chain required collateral as a string. + */ +export async function getRequiredCollateralOnChain(stellarAddress: string): Promise { + const acct = await contractClient.getAccount(stellarAddress); + return acct.requiredCollateral.toString(); +} + /** * Emergency override for collateral requirements, bypassing staleness and rate limits (#332). */ -export async function emergencySetRequiredCollateral(importer: string, newRequired: bigint): Promise { +export async function emergencySetRequiredCollateral( + importer: string, + newRequired: bigint +): Promise { await contractClient.setRequiredCollateral( - emergencyOracleKeypair, + [emergencyOracleKeypair], importer, newRequired, undefined, true, // bypassRateLimit - true, // emergency + true // emergency ); } diff --git a/packages/sdk/src/index.ts b/packages/sdk/src/index.ts index 4f8e54e..03de925 100644 --- a/packages/sdk/src/index.ts +++ b/packages/sdk/src/index.ts @@ -130,8 +130,12 @@ export class TariffShieldClient { bypassRateLimit?: boolean, emergency?: boolean, ): Promise> { + const primarySigner = signers[0]; + if (!primarySigner) { + throw new Error("At least one signer is required for setRequiredCollateral"); + } const args = [ - addressToScVal(signer.publicKey()), + addressToScVal(primarySigner.publicKey()), addressToScVal(importer), nativeToScVal(newRequired, { type: "i128" }), ]; @@ -145,7 +149,7 @@ export class TariffShieldClient { args.push(nativeToScVal(bypassRateLimit ?? false, { type: "bool" })); args.push(nativeToScVal(emergency ?? false, { type: "bool" })); - return this.invokeAndSubmit(signer, "set_required_collateral", args); + return this.invokeAndSubmitMulti(signers, "set_required_collateral", args, primarySigner); } async autoTopUp(signer: Keypair, importer: string): Promise> {