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

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions comebackhere-backend/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand All @@ -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",
Expand Down
70 changes: 66 additions & 4 deletions comebackhere-backend/src/app.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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)
Expand All @@ -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<unknown>): Promise<boolean> => {
let timer: ReturnType<typeof setTimeout> | undefined
try {
return await Promise.race([
probe().then(() => true, () => false),
new Promise<boolean>((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) => {
Expand Down
90 changes: 59 additions & 31 deletions comebackhere-backend/src/db/mongo.ts
Original file line number Diff line number Diff line change
@@ -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"

Expand Down Expand Up @@ -105,6 +106,47 @@ export const MAX_PAGE_SIZE = 100

let client: MongoClient | null = null
let db: Db | null = null
let connecting: Promise<Db> | null = null
const indexesByDb = new WeakMap<Db, Promise<void>>()

export function ensureIndexes(database: Db): Promise<void> {
const existing = indexesByDb.get(database)
if (existing) return existing

const indexes: Array<{ collection: string; keys: Record<string, 1 | -1>; 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
Expand Down Expand Up @@ -134,8 +176,17 @@ const MONGO_OPTIONS = {
socketTimeoutMS: 45_000,
}

export async function connectMongo(): Promise<Db> {
if (db) return db
export function connectMongo(): Promise<Db> {
if (db) return ensureIndexes(db).then(() => db!)
if (connecting) return connecting

connecting = connectMongoOnce().finally(() => {
connecting = null
})
return connecting
}

async function connectMongoOnce(): Promise<Db> {

const uri = process.env.MONGODB_URI ?? "mongodb://localhost:27017"
const dbName = process.env.MONGODB_DB ?? "comebackhere"
Expand All @@ -145,6 +196,7 @@ export async function connectMongo(): Promise<Db> {
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 =
Expand All @@ -153,7 +205,7 @@ export async function connectMongo(): Promise<Db> {
`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 })
}
Expand All @@ -163,43 +215,18 @@ export async function connectMongo(): Promise<Db> {
// 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<SettlementRecord>("settlements")
await settlements.createIndex({ id: 1 }, { unique: true })
await settlements.createIndex({ status: 1 })

const invoices = db.collection<InvoiceRecord>("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<IndexerCursor>("indexer_cursors")
await cursors.createIndex({ _id: 1 }, { unique: true })

const invoiceEvents = db.collection<InvoiceEventRecord>("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<ComplianceAuditRecord>("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
}
Expand Down Expand Up @@ -240,4 +267,5 @@ export async function closeMongo(): Promise<void> {
export function _resetMongoSingleton(): void {
client = null
db = null
connecting = null
}
37 changes: 16 additions & 21 deletions comebackhere-backend/src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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)
}

Expand All @@ -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")
})

// ---------------------------------------------------------------------------
Expand All @@ -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"))
Loading