diff --git a/comebackhere-backend/package.json b/comebackhere-backend/package.json index 3007d1d..8972d55 100644 --- a/comebackhere-backend/package.json +++ b/comebackhere-backend/package.json @@ -15,6 +15,8 @@ "express": "^5.2.1", "ioredis": "^6.0.0", "mongodb": "^7.6.0", + "pino": "^10.3.1", + "prom-client": "^15.1.3", "rate-limiter-flexible": "^11.2.1", "stellar-sdk": "^12.1.0", "swagger-jsdoc": "^6.2.8", @@ -32,6 +34,7 @@ "@typescript-eslint/eslint-plugin": "^8.69.0", "@typescript-eslint/parser": "^8.70.1", "eslint": "^8.57.0", + "pino-pretty": "^13.1.3", "supertest": "^7.0.0", "tsx": "^4.7.0", "typescript": "^6.0.2", diff --git a/comebackhere-backend/src/app.ts b/comebackhere-backend/src/app.ts index fad2520..4266f14 100644 --- a/comebackhere-backend/src/app.ts +++ b/comebackhere-backend/src/app.ts @@ -19,7 +19,15 @@ import { errorHandler, notFoundHandler } from "./middleware/errorHandler.js" import { createCorsMiddleware } from "./middleware/cors.js" import { parseCorsOrigins } from "./lib/env.js" import { openapiSpec } from "./openapi.js" -import { renderMetrics } from "./lib/metrics.js" +import { connectMongo } from "./db/mongo.js" +import { pingIndexerRedis } from "./indexer.js" +import { buildSorobanClient } from "./lib/soroban.js" +import { + httpRequestDuration, + metricsEnabled, + metricsRegistry, + renderMetrics, +} from "./lib/metrics.js" /** Maximum accepted JSON body size; larger requests get a 413 envelope. */ export const JSON_BODY_LIMIT = "100kb" @@ -69,6 +77,20 @@ export function createApp(options: CreateAppOptions = {}) { // line, downstream call and error envelope can reference the same // correlation ID — including body-parsing errors. app.use(correlationIdMiddleware) + app.use((req, res, next) => { + if (req.path !== "/metrics") { + const startedAt = process.hrtime.bigint() + res.on("finish", () => { + const routePath = req.route ? String(req.route.path) : "unmatched" + const route = `${req.baseUrl}${routePath}` || "/" + httpRequestDuration.observe( + { method: req.method, route, status_code: String(res.statusCode) }, + Number(process.hrtime.bigint() - startedAt) / 1_000_000_000, + ) + }) + } + next() + }) app.use((req, res, next) => req.path === "/api-docs" || req.path.startsWith("/api-docs/") ? swaggerHelmet(req, res, next) @@ -83,12 +105,52 @@ export function createApp(options: CreateAppOptions = {}) { // ── Health ────────────────────────────────────────────────────────────────── app.get("/health", (_req, res) => res.json({ status: "ok" })) + app.get("/health/ready", async (_req, res) => { + const check = async (probe: () => Promise): Promise => { + let timer: ReturnType | undefined + try { + return await Promise.race([ + probe().then(() => true, () => false), + new Promise((resolve) => { + timer = setTimeout(() => resolve(false), 1_500) + timer.unref?.() + }), + ]) + } finally { + if (timer) clearTimeout(timer) + } + } - // ── Prometheus metrics ────────────────────────────────────────────────────── - app.get("/metrics", (_req, res) => { - res.type("text/plain; version=0.0.4").send(renderMetrics()) + const rpcUrl = process.env.SOROBAN_RPC_URL + const [mongo, redis, sorobanRpc] = await Promise.all([ + check(async () => (await connectMongo()).command({ ping: 1 })), + check(async () => { + if ((await pingIndexerRedis()) !== "PONG") throw new Error("Redis ping failed") + }), + check(async () => { + if (!rpcUrl) throw new Error("Soroban RPC is not configured") + const getHealth = buildSorobanClient(rpcUrl).getHealth + if (!getHealth) throw new Error("Soroban RPC health is unsupported") + await getHealth() + }), + ]) + const dependencies = { mongo, redis, sorobanRpc } + const ready = Object.values(dependencies).every(Boolean) + res.status(ready ? 200 : 503).json({ + status: ready ? "ok" : "error", + dependencies: Object.fromEntries( + Object.entries(dependencies).map(([name, healthy]) => [name, healthy ? "ok" : "unavailable"]), + ), + }) }) + // ── Prometheus metrics ────────────────────────────────────────────────────── + if (metricsEnabled()) { + app.get("/metrics", async (_req, res) => { + res.type(metricsRegistry.contentType).send(await renderMetrics()) + }) + } + // ── OpenAPI spec (Issue #218) ─────────────────────────────────────────────── // Raw JSON spec at a stable, machine-readable URL app.get("/api-docs/swagger.json", (_req, res) => { diff --git a/comebackhere-backend/src/db/mongo.ts b/comebackhere-backend/src/db/mongo.ts index b5d031e..916a984 100644 --- a/comebackhere-backend/src/db/mongo.ts +++ b/comebackhere-backend/src/db/mongo.ts @@ -1,4 +1,5 @@ import { MongoClient, type Db, type Collection, MongoServerSelectionError } from "mongodb" +import { logger } from "../lib/logger.js" export type InvoiceStatus = "Pending" | "Paid" | "Expired" | "Cancelled" | "RefundRequested" | "Released" @@ -105,6 +106,47 @@ export const MAX_PAGE_SIZE = 100 let client: MongoClient | null = null let db: Db | null = null +let connecting: Promise | null = null +const indexesByDb = new WeakMap>() + +export function ensureIndexes(database: Db): Promise { + const existing = indexesByDb.get(database) + if (existing) return existing + + const indexes: Array<{ collection: string; keys: Record; unique?: boolean }> = [ + { collection: "settlements", keys: { id: 1 }, unique: true }, + { collection: "settlements", keys: { status: 1 } }, + { collection: "invoices", keys: { invoice_id: 1 }, unique: true }, + { collection: "invoices", keys: { status: 1 } }, + { collection: "invoices", keys: { merchant_address: 1 } }, + { collection: "invoices", keys: { status: 1, merchant_address: 1 } }, + { collection: "invoices", keys: { created_at: -1 } }, + { collection: "invoices", keys: { created_at: -1, invoice_id: -1 } }, + { collection: "indexer_cursors", keys: { _id: 1 }, unique: true }, + { collection: "invoice_events", keys: { event_id: 1 }, unique: true }, + { collection: "invoice_events", keys: { invoice_id: 1, ledger: 1 } }, + { collection: "webhook_pending_deliveries", keys: { _id: 1 }, unique: true }, + { collection: "webhook_dead_letters", keys: { _id: 1 }, unique: true }, + { collection: "webhook_dead_letters", keys: { failed_at: -1 } }, + { collection: "compliance_audit", keys: { event_id: 1 }, unique: true }, + { collection: "compliance_audit", keys: { address: 1, ledger: -1 } }, + { collection: "compliance_audit", keys: { event_type: 1, ledger: -1 } }, + { collection: "compliance_audit", keys: { ledger: -1 } }, + ] + + const creating = Promise.all( + indexes.map(async ({ collection, keys, unique }) => { + const indexName = await database.collection(collection).createIndex(keys, unique ? { unique } : {}) + logger.info({ collection, indexName }, "MongoDB index ensured") + }), + ).then(() => undefined) + const result = creating.catch((error: unknown) => { + indexesByDb.delete(database) + throw error + }) + indexesByDb.set(database, result) + return result +} // --------------------------------------------------------------------------- // #210 — Connection options: explicit pool size and timeouts so a slow or @@ -134,8 +176,17 @@ const MONGO_OPTIONS = { socketTimeoutMS: 45_000, } -export async function connectMongo(): Promise { - if (db) return db +export function connectMongo(): Promise { + if (db) return ensureIndexes(db).then(() => db!) + if (connecting) return connecting + + connecting = connectMongoOnce().finally(() => { + connecting = null + }) + return connecting +} + +async function connectMongoOnce(): Promise { const uri = process.env.MONGODB_URI ?? "mongodb://localhost:27017" const dbName = process.env.MONGODB_DB ?? "comebackhere" @@ -145,6 +196,7 @@ export async function connectMongo(): Promise { try { await client.connect() } catch (err) { + client = null // Provide a clear, actionable error message rather than letting the raw // driver error bubble up silently. const message = @@ -153,7 +205,7 @@ export async function connectMongo(): Promise { `Original error: ${err.message}` : `Failed to connect to MongoDB: ${err instanceof Error ? err.message : String(err)}` - console.error(`[mongo] ${message}`) + logger.error({ errorName: err instanceof Error ? err.name : "UnknownError" }, "MongoDB connection failed") // Re-throw so callers (routes, startup health-checks) can respond with 5xx. throw Object.assign(new Error(message), { status: 503 }) } @@ -163,43 +215,18 @@ export async function connectMongo(): Promise { // Attach a top-level error handler so an unexpected mid-run topology // failure is logged clearly rather than crashing the process silently. client.on("error", (err: Error) => { - console.error("[mongo] client error", err.message) + logger.error({ errorName: err.name }, "MongoDB client error") }) client.on("close", () => { - console.warn("[mongo] connection closed — subsequent requests will reconnect") + logger.warn("MongoDB connection closed; subsequent requests will reconnect") // Reset cached references so the next call to connectMongo() re-establishes // the connection instead of returning a stale db handle. db = null client = null }) - const settlements = db.collection("settlements") - await settlements.createIndex({ id: 1 }, { unique: true }) - await settlements.createIndex({ status: 1 }) - - const invoices = db.collection("invoices") - await invoices.createIndex({ invoice_id: 1 }, { unique: true }) - await invoices.createIndex({ status: 1 }) - await invoices.createIndex({ merchant_address: 1 }) - await invoices.createIndex({ status: 1, merchant_address: 1 }) - await invoices.createIndex({ created_at: -1 }) - await invoices.createIndex({ created_at: -1, invoice_id: -1 }) - - const cursors = db.collection("indexer_cursors") - await cursors.createIndex({ _id: 1 }, { unique: true }) - - const invoiceEvents = db.collection("invoice_events") - await invoiceEvents.createIndex({ event_id: 1 }, { unique: true }) - await invoiceEvents.createIndex({ invoice_id: 1, ledger: 1 }) - - await db.collection("webhook_dead_letters").createIndex({ failed_at: -1 }) - - const complianceAudit = db.collection("compliance_audit") - await complianceAudit.createIndex({ event_id: 1 }, { unique: true }) - await complianceAudit.createIndex({ address: 1, ledger: -1 }) - await complianceAudit.createIndex({ event_type: 1, ledger: -1 }) - await complianceAudit.createIndex({ ledger: -1 }) + await ensureIndexes(db) return db } @@ -240,4 +267,5 @@ export async function closeMongo(): Promise { export function _resetMongoSingleton(): void { client = null db = null + connecting = null } diff --git a/comebackhere-backend/src/index.ts b/comebackhere-backend/src/index.ts index 89b0ba0..2d36be6 100644 --- a/comebackhere-backend/src/index.ts +++ b/comebackhere-backend/src/index.ts @@ -11,17 +11,18 @@ import { createApp } from "./app.js" import { startTreasuryIndexer, stopTreasuryIndexer } from "./services/treasury-indexer.js" import { stopIndexer } from "./indexer.js" import { stopComplianceIndexer } from "./services/compliance-indexer.js" -import { closeMongo } from "./db/mongo.js" +import { closeMongo, connectMongo, ensureIndexes } from "./db/mongo.js" import { webhookDeliveryQueue } from "./services/webhook-delivery.js" import { createShutdownHandler, resolveWebhookDrainTimeout } from "./shutdown.js" import { validateEnv } from "./lib/env.js" +import { logger } from "./lib/logger.js" import type { Server } from "http" // Fail fast on missing variables or malformed Stellar ids, naming the variable. try { validateEnv(process.env) } catch (err) { - console.error(`[startup] ${err instanceof Error ? err.message : err}`) + logger.fatal({ errorName: err instanceof Error ? err.name : "UnknownError" }, "Environment validation failed") process.exit(1) } @@ -37,16 +38,19 @@ const WEBHOOK_DRAIN_TIMEOUT_MS = resolveWebhookDrainTimeout( const app = createApp() startTreasuryIndexer() +void connectMongo() + .then((database) => ensureIndexes(database)) + .catch((err: unknown) => { + logger.error({ errorName: err instanceof Error ? err.name : "UnknownError" }, "MongoDB startup/index initialization failed") + }) + // Retry deliveries that a previous process persisted during shutdown. webhookDeliveryQueue.resumePending().catch((err: unknown) => { - console.error( - "[webhook] could not resume persisted deliveries:", - err instanceof Error ? err.message : err, - ) + logger.error({ errorName: err instanceof Error ? err.name : "UnknownError" }, "Could not resume persisted webhook deliveries") }) const server: Server = app.listen(Number(PORT), () => { - console.log(`comebackhere-backend listening on port ${PORT}`) + logger.info({ port: Number(PORT) }, "Backend listening") }) // --------------------------------------------------------------------------- @@ -60,20 +64,11 @@ const shutdown = createShutdownHandler({ stopTreasuryIndexer() stopIndexer() stopComplianceIndexer() - console.log("[shutdown] indexers stopped") - - // 3. Close MongoDB connection. - await closeMongo() - console.log("[shutdown] MongoDB connection closed") - - clearTimeout(hardTimeout) - console.log("[shutdown] clean exit") - process.exit(0) - } catch (err) { - console.error("[shutdown] error during shutdown:", err) - process.exit(1) - } -} + }, + closeMongo, + shutdownTimeoutMs: SHUTDOWN_TIMEOUT_MS, + webhookDrainTimeoutMs: WEBHOOK_DRAIN_TIMEOUT_MS, +}) process.on("SIGTERM", () => void shutdown("SIGTERM")) process.on("SIGINT", () => void shutdown("SIGINT")) diff --git a/comebackhere-backend/src/indexer.ts b/comebackhere-backend/src/indexer.ts index ae38f02..b04b2e0 100644 --- a/comebackhere-backend/src/indexer.ts +++ b/comebackhere-backend/src/indexer.ts @@ -49,7 +49,8 @@ import { type InvoiceStatus, } from "./db/mongo.js" import { getOldestRetainedLedger, ledgerFromPagingToken, parseRetentionError } from "./lib/soroban.js" -import { counter, gauge } from "./lib/metrics.js" +import { counter, gauge, indexerLedgerLag } from "./lib/metrics.js" +import { logger } from "./lib/logger.js" // --------------------------------------------------------------------------- // Types @@ -138,15 +139,11 @@ export function createRedisClient(redisUrl?: string): Redis { attempt = times if (times > 50) { // After 50 retries (~30 min with cap) give up so operators notice. - console.error( - `[indexer] Redis retry limit reached after ${times} attempts — stopping reconnect` - ) + logger.error({ attempts: times }, "Redis retry limit reached; stopping reconnect") return null } const delay = backoffDelayMs(times - 1) - console.warn( - `[indexer] Redis reconnect attempt ${times} — waiting ${delay} ms` - ) + logger.warn({ attempt: times, delayMs: delay }, "Redis reconnect scheduled") return delay }, // Do not flood logs when commands queue during a disconnect. @@ -156,22 +153,27 @@ export function createRedisClient(redisUrl?: string): Redis { }) client.on("connect", () => { - console.log("[indexer] Redis connected") + logger.info("Indexer Redis connected") attempt = 0 }) client.on("reconnecting", (ms: number) => { - console.warn(`[indexer] Redis reconnecting in ${ms} ms (attempt ${attempt})`) + logger.warn({ delayMs: ms, attempt }, "Indexer Redis reconnecting") }) client.on("error", (err: Error) => { // Log but do not crash — the indexer continues polling Soroban. - console.error(`[indexer] Redis error: ${err.message}`) + logger.error({ errorName: err.name }, "Indexer Redis error") }) return client } +export function pingIndexerRedis(): Promise { + if (!redisClient) redisClient = createRedisClient() + return redisClient.ping() +} + // --------------------------------------------------------------------------- // Cursor read / write (with Redis fallback to in-memory) // --------------------------------------------------------------------------- @@ -186,7 +188,7 @@ export async function loadCursor(): Promise { return stored } } catch (err) { - console.warn("[indexer] could not read cursor from Redis — using in-memory cursor", err) + logger.warn({ errorName: err instanceof Error ? err.name : "UnknownError" }, "Could not read Redis cursor; using in-memory cursor") } } return memCursor @@ -200,7 +202,7 @@ export async function saveCursor(next: string): Promise { await redisClient.set(INDEXER_CURSOR_KEY, next) } catch (err) { // Non-fatal: in-memory cursor is still updated, so polling continues. - console.warn("[indexer] could not save cursor to Redis — using in-memory fallback", err) + logger.warn({ errorName: err instanceof Error ? err.name : "UnknownError" }, "Could not save Redis cursor; using in-memory fallback") } } } @@ -262,9 +264,14 @@ export function eventIdOf(event: { id?: string; pagingToken?: string }): string * crashed after the event was stored but before it was marked applied. */ export function persistTransition(transition: InvoiceStateTransition): void { - console.log( - `[indexer] ${transition.event_type} invoice_id=${transition.invoice_id}` + - ` ledger=${transition.ledger} tx=${transition.transaction_hash}` + logger.info( + { + eventType: transition.event_type, + invoiceId: transition.invoice_id, + ledger: transition.ledger, + transactionHash: transition.transaction_hash, + }, + "Invoice state transition indexed", ) } @@ -386,10 +393,9 @@ export async function recordRetentionGap( toLedger: number, ): Promise { const missing = Math.max(0, toLedger - fromLedger + 1) - console.error( - `[indexer] RETENTION GAP: ledgers ${fromLedger}-${toLedger} (${missing} ledgers) are no longer ` + - `retained by the RPC node; events in this range were NOT indexed. Resuming from ledger ` + - `${toLedger + 1}. See docs/troubleshooting.md#indexer-retention-gaps to backfill.` + logger.error( + { fromLedger, toLedger, missingLedgers: missing, resumeLedger: toLedger + 1 }, + "Indexer retention gap detected; events in the missing range were not indexed", ) const labels = { indexer: "invoice" } @@ -538,6 +544,10 @@ export async function pollOnce( if (db) await saveMongoCursor(db, nextToken, lastLedger) if (nextToken) await saveCursor(nextToken) + indexerLedgerLag.set( + { indexer: "invoice" }, + Math.max(0, (response?.latestLedger ?? lastLedger) - lastLedger), + ) return applied } @@ -606,16 +616,14 @@ export async function startIndexer(options?: { const initialCursor = saved ? `ledger=${saved.last_ledger} token=${saved.paging_token ?? "none"}` : await loadCursor() - console.log( - `[indexer] starting — contract=${contractId} cursor=${initialCursor} interval=${pollIntervalMs}ms` - ) + logger.info({ contractId, cursor: initialCursor, pollIntervalMs }, "Invoice indexer starting") const loop = async () => { if (stopped) return try { await pollOnce(rpc, contractId!, database) } catch (err) { - const handler = options?.onError ?? ((e) => console.error("[indexer] poll error", e)) + const handler = options?.onError ?? ((e) => logger.error({ errorName: e instanceof Error ? e.name : "UnknownError" }, "Invoice indexer poll failed")) handler(err) } if (!stopped) { @@ -631,7 +639,7 @@ export async function startIndexer(options?: { // Run as standalone entry point if (import.meta.url === new URL(process.argv[1], import.meta.url).href) { startIndexer().catch((err) => { - console.error("[indexer] fatal", err) + logger.fatal({ errorName: err instanceof Error ? err.name : "UnknownError" }, "Invoice indexer failed") process.exit(1) }) } diff --git a/comebackhere-backend/src/lib/cache.ts b/comebackhere-backend/src/lib/cache.ts index f827285..e10903f 100644 --- a/comebackhere-backend/src/lib/cache.ts +++ b/comebackhere-backend/src/lib/cache.ts @@ -1,4 +1,5 @@ import Redis from "ioredis" +import { logger } from "./logger.js" let _redis: Redis | null = null @@ -93,7 +94,7 @@ export function memoryCacheSet(key: string, value: unknown, ttlMs: number): void */ export function invalidateCacheKey(key: string, reason = "manual"): boolean { const existed = _memoryCache.delete(key) - console.log(`[cache] invalidated key=${key} reason=${reason} evicted=${existed}`) + logger.info({ key, reason, evicted: existed }, "Cache key invalidated") return existed } diff --git a/comebackhere-backend/src/lib/logger.ts b/comebackhere-backend/src/lib/logger.ts new file mode 100644 index 0000000..6f2067c --- /dev/null +++ b/comebackhere-backend/src/lib/logger.ts @@ -0,0 +1,31 @@ +import pino from "pino" + +const redactPaths = [ + "ADMIN_KEY", + "SIGNER_SECRET_KEY", + "WEBHOOK_SECRET", + "API_KEY", + "*.adminKey", + "*.signerSecretKey", + "*.webhookSecret", + "*.apiKey", + "req.headers.authorization", + "req.headers.x-admin-key", + "req.headers.x-api-key", + "err.message", + "error.message", +] + +export const logger = pino({ + level: process.env.LOG_LEVEL ?? "info", + base: { service: "comebackhere-backend", correlationId: null }, + redact: { paths: redactPaths, censor: "[REDACTED]" }, + ...(process.env.NODE_ENV === "development" + ? { + transport: { + target: "pino-pretty", + options: { colorize: true, translateTime: "SYS:standard", ignore: "pid,hostname" }, + }, + } + : {}), +}) \ No newline at end of file diff --git a/comebackhere-backend/src/lib/metrics.ts b/comebackhere-backend/src/lib/metrics.ts index e03e3ed..5c5da82 100644 --- a/comebackhere-backend/src/lib/metrics.ts +++ b/comebackhere-backend/src/lib/metrics.ts @@ -1,80 +1,71 @@ -/** - * Minimal in-process metrics registry rendered in the Prometheus text - * exposition format at GET /metrics. Only counters and gauges are needed - * today; swap for prom-client if histograms become necessary. - */ +import { collectDefaultMetrics, Counter, Gauge, Histogram, Registry } from "prom-client" type Labels = Record -interface Metric { - name: string - help: string - type: "counter" | "gauge" - values: Map -} +export const metricsRegistry = new Registry() +collectDefaultMetrics({ register: metricsRegistry }) -const registry = new Map() +export const httpRequestDuration = new Histogram({ + name: "http_request_duration_seconds", + help: "HTTP request duration in seconds", + labelNames: ["method", "route", "status_code"], + buckets: [0.005, 0.01, 0.025, 0.05, 0.1, 0.25, 0.5, 1, 2.5, 5, 10], + registers: [metricsRegistry], +}) -function labelKey(labels: Labels): string { - return Object.keys(labels) - .sort() - .map((k) => `${k}=${labels[k]}`) - .join(",") -} +export const indexerLedgerLag = new Gauge({ + name: "indexer_ledger_lag", + help: "Number of ledgers the indexer is behind the latest Soroban ledger", + labelNames: ["indexer"], + registers: [metricsRegistry], +}) -function register(name: string, help: string, type: Metric["type"]): Metric { - let metric = registry.get(name) - if (!metric) { - metric = { name, help, type, values: new Map() } - registry.set(name, metric) - } - return metric -} +export const webhookDeliveryOutcomes = new Counter({ + name: "webhook_delivery_total", + help: "Webhook delivery outcomes", + labelNames: ["status"], + registers: [metricsRegistry], +}) + +const retentionLabels = ["indexer"] as const -function makeMetric(name: string, help: string, type: Metric["type"]) { - const metric = register(name, help, type) +export function counter(name: string, help: string) { + const metric = new Counter({ + name, + help, + labelNames: name.startsWith("indexer_retention_") ? [...retentionLabels] : [], + registers: [metricsRegistry], + }) return { - get(labels: Labels = {}): number { - return metric.values.get(labelKey(labels))?.value ?? 0 - }, - set(value: number, labels: Labels = {}): void { - metric.values.set(labelKey(labels), { labels, value }) - }, inc(labels: Labels = {}, by = 1): void { - const key = labelKey(labels) - const current = metric.values.get(key)?.value ?? 0 - metric.values.set(key, { labels, value: current + by }) + metric.inc(labels, by) }, } } -export function counter(name: string, help: string) { - const { get, inc } = makeMetric(name, help, "counter") - return { get, inc } -} - export function gauge(name: string, help: string) { - return makeMetric(name, help, "gauge") + const metric = new Gauge({ + name, + help, + labelNames: name.startsWith("indexer_retention_") ? [...retentionLabels] : [], + registers: [metricsRegistry], + }) + return { + set(value: number, labels: Labels = {}): void { + metric.set(labels, value) + }, + } } -function escapeLabel(value: string): string { - return value.replace(/\\/g, "\\\\").replace(/"/g, '\\"').replace(/\n/g, "\\n") +export function metricsEnabled(): boolean { + return process.env.METRICS_ENABLED?.toLowerCase() !== "false" } -export function renderMetrics(): string { - const lines: string[] = [] - for (const metric of registry.values()) { - lines.push(`# HELP ${metric.name} ${metric.help}`) - lines.push(`# TYPE ${metric.name} ${metric.type}`) - for (const { labels, value } of metric.values.values()) { - const pairs = Object.entries(labels).map(([k, v]) => `${k}="${escapeLabel(v)}"`) - lines.push(`${metric.name}${pairs.length ? `{${pairs.join(",")}}` : ""} ${value}`) - } - } - return lines.join("\n") + "\n" +export function renderMetrics(): Promise { + return metricsRegistry.metrics() } -/** Exported for tests. */ +/** Exported for tests and process-local instrumentation resets. */ export function _resetMetrics(): void { - for (const metric of registry.values()) metric.values.clear() + metricsRegistry.resetMetrics() } diff --git a/comebackhere-backend/src/middleware/correlationId.ts b/comebackhere-backend/src/middleware/correlationId.ts index 717aa51..cb7742b 100644 --- a/comebackhere-backend/src/middleware/correlationId.ts +++ b/comebackhere-backend/src/middleware/correlationId.ts @@ -14,11 +14,12 @@ * The ID is stored on `res.locals.requestId` so route handlers and other * middleware can include it in log lines: * - * console.log(`[requestId=${res.locals.requestId}] processing invoice`) + * res.locals.logger.info({ invoiceId }, "Processing invoice") */ import { type Request, type Response, type NextFunction } from "express" import { randomUUID } from "node:crypto" +import { logger } from "../lib/logger.js" export function correlationIdMiddleware( req: Request, @@ -34,6 +35,7 @@ export function correlationIdMiddleware( // Make the ID available to downstream handlers and logging res.locals.requestId = requestId + res.locals.logger = logger.child({ correlationId: requestId }) // Echo the ID on the response so callers can correlate client-side res.setHeader("X-Request-Id", requestId) diff --git a/comebackhere-backend/src/middleware/errorHandler.ts b/comebackhere-backend/src/middleware/errorHandler.ts index 5136cf3..86eb309 100644 --- a/comebackhere-backend/src/middleware/errorHandler.ts +++ b/comebackhere-backend/src/middleware/errorHandler.ts @@ -6,6 +6,7 @@ import { PayloadTooLargeError, parseContractErrorCode, } from "../lib/errors.js" +import { logger } from "../lib/logger.js" /** * Standard error envelope returned by every route: @@ -93,7 +94,8 @@ export function errorHandler(err: unknown, _req: Request, res: Response, next: N const correlationId = typeof res.locals.requestId === "string" ? res.locals.requestId : null if (appError.status >= 500) { - console.error(`[requestId=${correlationId}] ${appError.code}: ${appError.message}`) + const requestLogger = res.locals.logger ?? logger.child({ correlationId }) + requestLogger.error({ code: appError.code }, "Request failed") } res.status(appError.status).json(buildErrorEnvelope(appError, correlationId)) diff --git a/comebackhere-backend/src/routes/analytics.ts b/comebackhere-backend/src/routes/analytics.ts index cf1d620..9ece040 100644 --- a/comebackhere-backend/src/routes/analytics.ts +++ b/comebackhere-backend/src/routes/analytics.ts @@ -257,7 +257,10 @@ router.get("/metrics", validateQuery(analyticsQuerySchema), async (req: Request, res.json(analyticsData) } catch (error) { - console.error("Error fetching analytics metrics:", error) + res.locals.logger.error( + { errorName: error instanceof Error ? error.name : "UnknownError" }, + "Analytics metrics request failed", + ) res.status(500).json({ error: "Failed to fetch analytics metrics" }) } diff --git a/comebackhere-backend/src/routes/compliance.ts b/comebackhere-backend/src/routes/compliance.ts index 0a011ab..3a0424f 100644 --- a/comebackhere-backend/src/routes/compliance.ts +++ b/comebackhere-backend/src/routes/compliance.ts @@ -240,8 +240,8 @@ router.post("/block", validateBody(blockBodySchema), asyncHandler(async (req: Re signerSecret: "SIGNER_SECRET_KEY", }) - // Audit log — admin identity + timestamp - console.log(`[compliance] block_address admin="${adminKey}" address="${address}" ts="${new Date().toISOString()}"`) + // The admin key is a credential and must never be included in logs. + res.locals.logger.info({ address }, "Compliance address block requested") const client = buildSorobanClient(env.rpcUrl) const result = await callComplianceOp( diff --git a/comebackhere-backend/src/routes/invoices.ts b/comebackhere-backend/src/routes/invoices.ts index 0a0fece..31de43d 100644 --- a/comebackhere-backend/src/routes/invoices.ts +++ b/comebackhere-backend/src/routes/invoices.ts @@ -342,7 +342,10 @@ router.get("/export.csv", async (req: Request, res: Response) => { } else { // Mid-stream failure: abort so the client sees a truncated download // rather than a file that silently looks complete. - console.error("[invoices] CSV export failed mid-stream:", message) + res.locals.logger.error( + { errorName: err instanceof Error ? err.name : "UnknownError" }, + "Invoice CSV export failed mid-stream", + ) res.destroy(err instanceof Error ? err : new Error(message)) } } finally { diff --git a/comebackhere-backend/src/routes/treasury.ts b/comebackhere-backend/src/routes/treasury.ts index a7091df..ff27b84 100644 --- a/comebackhere-backend/src/routes/treasury.ts +++ b/comebackhere-backend/src/routes/treasury.ts @@ -210,10 +210,14 @@ export async function executeSettlementWithBalanceCheck( env.networkPassphrase, ) - console.log( - `[execute-settlement] settlement_id=${body.settlement_id} ` + - `required=${settlement.amount.toString()} available=${balance.toString()} ` + - `token=${tokenContract}`, + res.locals.logger.info( + { + settlementId: body.settlement_id, + required: settlement.amount.toString(), + available: balance.toString(), + token: tokenContract, + }, + "Executing settlement balance check", ) if (balance < settlement.amount) { diff --git a/comebackhere-backend/src/services/compliance-indexer.ts b/comebackhere-backend/src/services/compliance-indexer.ts index b07600a..730b395 100644 --- a/comebackhere-backend/src/services/compliance-indexer.ts +++ b/comebackhere-backend/src/services/compliance-indexer.ts @@ -1,6 +1,8 @@ import { SorobanRpc, xdr } from "stellar-sdk" import { buildSorobanClient, type SorobanClient } from "../lib/soroban.js" import { connectMongo, getCursorsCollection, getComplianceAuditCollection, type ComplianceAuditRecord, type ComplianceAuditStatus } from "../db/mongo.js" +import { indexerLedgerLag } from "../lib/metrics.js" +import { logger } from "../lib/logger.js" const CURSOR_ID = "compliance_audit_events" const EVENT_LIMIT = 100 @@ -90,6 +92,10 @@ export async function processComplianceIndexerBatch( }, { upsert: true }, ) + indexerLedgerLag.set( + { indexer: "compliance" }, + Math.max(0, (response.latestLedger ?? cursor.last_ledger) - (response.events?.at(-1)?.ledger ?? response.latestLedger ?? cursor.last_ledger)), + ) return processed } @@ -102,7 +108,7 @@ export function startComplianceIndexer(): void { const client = buildSorobanClient(rpcUrl) const tick = async () => { try { await processComplianceIndexerBatch(client, contractId, await connectMongo()) } - catch (err) { console.error("[compliance-indexer] error:", err instanceof Error ? err.message : err) } + catch (err) { logger.error({ errorName: err instanceof Error ? err.name : "UnknownError" }, "Compliance indexer failed") } } void tick() timer = setInterval(() => void tick(), POLL_INTERVAL_MS) diff --git a/comebackhere-backend/src/services/treasury-indexer.ts b/comebackhere-backend/src/services/treasury-indexer.ts index d2e6e93..eb51344 100644 --- a/comebackhere-backend/src/services/treasury-indexer.ts +++ b/comebackhere-backend/src/services/treasury-indexer.ts @@ -11,6 +11,8 @@ import { } from "../db/mongo.js" import { dispatchWebhook } from "./webhooks.js" import { invalidateBalanceCache } from "../lib/cache.js" +import { indexerLedgerLag } from "../lib/metrics.js" +import { logger } from "../lib/logger.js" const CURSOR_ID = "treasury_settlement_events" const POLL_INTERVAL_MS = 5_000 @@ -229,7 +231,7 @@ export async function processIndexerBatch( // De-duplicate: skip events we have already applied (reorg / replay protection) if (processedIds.has(eventId)) { - console.log(`[treasury-indexer] skipping duplicate event id=${eventId}`) + logger.debug({ eventId }, "Skipping duplicate treasury event") continue } @@ -250,10 +252,7 @@ export async function processIndexerBatch( token, tx_hash: txHash, }).catch((err: unknown) => { - console.error( - "[treasury-indexer] webhook dispatch failed (settlement_proposed):", - err instanceof Error ? err.message : err, - ) + logger.error({ errorName: err instanceof Error ? err.name : "UnknownError", eventType: "settlement_proposed" }, "Treasury webhook dispatch failed") }) } } else if (eventType === "settlement_approved") { @@ -271,10 +270,7 @@ export async function processIndexerBatch( approval_weight: newWeight.toString(), tx_hash: txHash, }).catch((err: unknown) => { - console.error( - "[treasury-indexer] webhook dispatch failed (settlement_approved):", - err instanceof Error ? err.message : err, - ) + logger.error({ errorName: err instanceof Error ? err.name : "UnknownError", eventType: "settlement_approved" }, "Treasury webhook dispatch failed") }) } } else if (eventType === "settlement_executed") { @@ -292,10 +288,7 @@ export async function processIndexerBatch( settlement_id: settlementId, tx_hash: txHash, }).catch((err: unknown) => { - console.error( - "[treasury-indexer] webhook dispatch failed (settlement_executed):", - err instanceof Error ? err.message : err, - ) + logger.error({ errorName: err instanceof Error ? err.name : "UnknownError", eventType: "settlement_executed" }, "Treasury webhook dispatch failed") }) } } @@ -307,11 +300,13 @@ export async function processIndexerBatch( const lastLedger = response.latestLedger ?? cursor.last_ledger await saveCursor(database, lastPagingToken, lastLedger, newEventIds) + indexerLedgerLag.set( + { indexer: "treasury" }, + Math.max(0, (response.latestLedger ?? lastLedger) - lastLedger), + ) if (processed > 0) { - console.log( - `[treasury-indexer] processed ${processed} event(s); cursor ledger=${lastLedger}`, - ) + logger.info({ processed, cursorLedger: lastLedger }, "Treasury indexer batch processed") } return processed @@ -326,9 +321,7 @@ export function startTreasuryIndexer(): void { const treasuryContractId = process.env.TREASURY_CONTRACT_ID if (!rpcUrl || !treasuryContractId) { - console.warn( - "[treasury-indexer] skipped: SOROBAN_RPC_URL and TREASURY_CONTRACT_ID required", - ) + logger.warn("Treasury indexer skipped; RPC URL and contract ID are required") return } @@ -339,13 +332,13 @@ export function startTreasuryIndexer(): void { const database = await connectMongo() await processIndexerBatch(client, treasuryContractId, database) } catch (err) { - console.error("[treasury-indexer] error:", err instanceof Error ? err.message : err) + logger.error({ errorName: err instanceof Error ? err.name : "UnknownError" }, "Treasury indexer failed") } } void tick() indexerTimer = setInterval(() => void tick(), POLL_INTERVAL_MS) - console.log("[treasury-indexer] started") + logger.info("Treasury indexer started") } export function stopTreasuryIndexer(): void { diff --git a/comebackhere-backend/src/services/webhook-delivery.ts b/comebackhere-backend/src/services/webhook-delivery.ts index fb15405..fe363a6 100644 --- a/comebackhere-backend/src/services/webhook-delivery.ts +++ b/comebackhere-backend/src/services/webhook-delivery.ts @@ -16,6 +16,8 @@ import { connectMongo } from "../db/mongo.js" import { getWebhookRetryConfig, type WebhookRetryConfig } from "../lib/env.js" +import { webhookDeliveryOutcomes } from "../lib/metrics.js" +import { logger } from "../lib/logger.js" // --------------------------------------------------------------------------- // Types @@ -196,6 +198,7 @@ export async function deliverWebhook( if (statusCode >= 200 && statusCode < 300) { record.status = "delivered" + webhookDeliveryOutcomes.inc({ status: "delivered" }) return record } @@ -215,10 +218,10 @@ export async function deliverWebhook( } record.status = "failed" - console.error( - `[webhook] delivery failed after ${record.attempts} attempt(s) ` + - `key=${record.idempotency_key} requestId=${record.request_id ?? "-"} ` + - `endpoint=${endpoint} last_error=${record.last_error}`, + webhookDeliveryOutcomes.inc({ status: "failed" }) + logger.error( + { attempts: record.attempts, requestId: record.request_id, lastStatusCode: record.last_status_code }, + "webhook delivery failed", ) return record } @@ -444,9 +447,7 @@ export class WebhookDeliveryQueue { const job: WebhookDeliveryJob = { endpoint, payload, attempts, attempt_history: attemptHistory } if (!this.accepting) { - console.warn( - `[webhook] queue closed — deferring key=${payload.idempotency_key} for retry after restart`, - ) + logger.warn("Webhook queue closed; delivery deferred for retry after restart") if (this.abort.signal.aborted) { // Drain already finished; persist straight away. this.persist([job]) @@ -483,7 +484,7 @@ export class WebhookDeliveryQueue { try { await this.deadLetterStore.save(record) } catch (err: unknown) { - console.error("[webhook] failed to save dead letter:", err instanceof Error ? err.message : err) + logger.error({ errorName: err instanceof Error ? err.name : "UnknownError" }, "Failed to save webhook dead letter") } } if (record.status !== "pending") this.active.delete(entry) @@ -498,7 +499,7 @@ export class WebhookDeliveryQueue { stopAccepting(): void { if (this.accepting) { this.accepting = false - console.log("[webhook] queue stopped accepting new deliveries") + logger.info("Webhook queue stopped accepting deliveries") } } @@ -509,9 +510,7 @@ export class WebhookDeliveryQueue { async drain(timeoutMs: number): Promise { this.stopAccepting() const startedWith = this.active.size - console.log( - `[webhook] draining ${startedWith} in-flight deliver${startedWith === 1 ? "y" : "ies"} (timeout ${timeoutMs}ms)`, - ) + logger.info({ inFlight: startedWith, timeoutMs }, "Draining webhook deliveries") let timer: ReturnType | undefined const timedOut = await Promise.race([ @@ -531,9 +530,9 @@ export class WebhookDeliveryQueue { if (unfinished.length > 0) await this.persist(unfinished) - console.log( - `[webhook] drain ${timedOut ? "timed out" : "complete"}: ` + - `${completed} finished, ${unfinished.length} persisted for retry`, + logger.info( + { timedOut, completed, persisted: unfinished.length }, + "Webhook delivery drain finished", ) return { completed, persisted: unfinished.length, timedOut } } @@ -543,7 +542,7 @@ export class WebhookDeliveryQueue { const jobs = await this.store.takeAll() for (const job of jobs) this.enqueue(job.endpoint, job.payload, job.attempts, job.attempt_history) if (jobs.length > 0) { - console.log(`[webhook] resumed ${jobs.length} persisted deliver${jobs.length === 1 ? "y" : "ies"}`) + logger.info({ resumed: jobs.length }, "Persisted webhook deliveries resumed") } return jobs.length } @@ -567,9 +566,9 @@ export class WebhookDeliveryQueue { try { await this.store.save(jobs) } catch (err) { - console.error( - `[webhook] failed to persist ${jobs.length} unfinished deliver${jobs.length === 1 ? "y" : "ies"}:`, - err instanceof Error ? err.message : err, + logger.error( + { count: jobs.length, errorName: err instanceof Error ? err.name : "UnknownError" }, + "Failed to persist unfinished webhook deliveries", ) } } diff --git a/comebackhere-backend/src/shutdown.ts b/comebackhere-backend/src/shutdown.ts index b11ef48..c28324f 100644 --- a/comebackhere-backend/src/shutdown.ts +++ b/comebackhere-backend/src/shutdown.ts @@ -14,6 +14,7 @@ import type { Server } from "http" import type { WebhookDeliveryQueue } from "./services/webhook-delivery.js" +import { logger } from "./lib/logger.js" export interface ShutdownDeps { server: Pick @@ -43,9 +44,7 @@ export function resolveWebhookDrainTimeout( const ceiling = Math.max(0, shutdownTimeoutMs - marginMs) if (!Number.isFinite(configuredMs) || configuredMs < 0) return ceiling if (configuredMs > ceiling) { - console.warn( - `[shutdown] WEBHOOK_DRAIN_TIMEOUT_MS=${configuredMs} exceeds shutdown budget; using ${ceiling}ms`, - ) + logger.warn({ configuredMs, effectiveMs: ceiling }, "Webhook drain timeout exceeds shutdown budget") return ceiling } return configuredMs @@ -59,11 +58,11 @@ export function createShutdownHandler(deps: ShutdownDeps): (signal: string) => P if (shuttingDown) return shuttingDown = true - console.log(`[shutdown] received ${signal} — starting graceful shutdown`) + logger.info({ signal }, "Starting graceful shutdown") // Hard-timeout safety net: if clean shutdown takes too long, force exit. const hardTimeout = setTimeout(() => { - console.error("[shutdown] hard timeout reached — forcing exit") + logger.error("Hard shutdown timeout reached; forcing exit") exit(1) }, deps.shutdownTimeoutMs) // Allow the process to exit even if the timer is still pending. @@ -77,25 +76,25 @@ export function createShutdownHandler(deps: ShutdownDeps): (signal: string) => P await new Promise((resolve, reject) => { deps.server.close((err) => (err ? reject(err) : resolve())) }) - console.log("[shutdown] HTTP server closed") + logger.info("HTTP server closed") // 3. Stop indexer poll loops. deps.stopIndexers() - console.log("[shutdown] indexers stopped") + logger.info("Indexers stopped") // 4. Drain webhook deliveries; unfinished ones are persisted for retry. await deps.webhookQueue.drain(deps.webhookDrainTimeoutMs) // 5. Close MongoDB connection. await deps.closeMongo() - console.log("[shutdown] MongoDB connection closed") + logger.info("MongoDB connection closed") clearTimeout(hardTimeout) - console.log("[shutdown] clean exit") + logger.info("Clean shutdown complete") exit(0) } catch (err) { clearTimeout(hardTimeout) - console.error("[shutdown] error during shutdown:", err) + logger.error({ errorName: err instanceof Error ? err.name : "UnknownError" }, "Shutdown failed") exit(1) } } diff --git a/docs/troubleshooting.md b/docs/troubleshooting.md index c8af464..a3cc168 100644 --- a/docs/troubleshooting.md +++ b/docs/troubleshooting.md @@ -2,6 +2,23 @@ This guide covers the most common problems developers encounter during local setup and how to resolve them. +## Health and Metrics + +`GET /health` is a liveness check: it reports whether the backend process can +serve requests and does not contact dependencies. `GET /health/ready` is a +readiness check: it returns `200` only when MongoDB, Redis, and Soroban RPC +respond within 1.5 seconds; otherwise it returns `503` with each dependency +marked `ok` or `unavailable`. + +`GET /metrics` exposes Prometheus text metrics and is enabled by default. Set +`METRICS_ENABLED=false` to disable the endpoint. Request latency is recorded in +`http_request_duration_seconds` (route templates, method, and status code), +indexer lag in `indexer_ledger_lag` (ledger count, labelled by indexer), and +webhook outcomes in `webhook_delivery_total` (labelled by `delivered` or +`failed`). Standard Node.js process metrics are also included. Logs are +structured JSON; `LOG_LEVEL` controls verbosity, and `pino-pretty` is enabled +only when `NODE_ENV=development`. + ## Soroban RPC Connection Errors ### "Soroban RPC not reachable" or connection refused on port 8000