diff --git a/__tests__/api/stellar/activity.test.ts b/__tests__/api/stellar/activity.test.ts index 4fd89504..2c3980a0 100644 --- a/__tests__/api/stellar/activity.test.ts +++ b/__tests__/api/stellar/activity.test.ts @@ -59,7 +59,7 @@ describe("GET /api/stellar/activity", () => { response: NextResponse.json({ message: "Unauthorized" }, { status: 401 }), }) - const response = await GET(buildRequest()) + const response = (await GET(buildRequest()))! expect(response.status).toBe(401) }) @@ -75,7 +75,7 @@ describe("GET /api/stellar/activity", () => { issuerPublicKey: "GD123", }) - const response = await GET(buildRequest()) + const response = (await GET(buildRequest()))! const payload = await response.json() expect(response.status).toBe(200) @@ -93,7 +93,7 @@ describe("GET /api/stellar/activity", () => { network: "testnet", }) - const response = await GET(buildRequest()) + const response = (await GET(buildRequest()))! const payload = await response.json() expect(response.status).toBe(200) @@ -136,7 +136,7 @@ describe("GET /api/stellar/activity", () => { }), }) - const response = await GET(buildRequest()) + const response = (await GET(buildRequest()))! const payload = await response.json() expect(response.status).toBe(200) diff --git a/__tests__/api/stellar/sync.test.ts b/__tests__/api/stellar/sync.test.ts index 94f2fd99..70bbbc11 100644 --- a/__tests__/api/stellar/sync.test.ts +++ b/__tests__/api/stellar/sync.test.ts @@ -45,7 +45,7 @@ describe("POST /api/admin/stellar/sync", () => { response: NextResponse.json({ message: "Unauthorized" }, { status: 401 }), }) - const response = await POST(buildRequest()) + const response = (await POST(buildRequest()))! expect(response.status).toBe(401) expect(sync).not.toHaveBeenCalled() }) @@ -60,7 +60,7 @@ describe("POST /api/admin/stellar/sync", () => { lastCursor: "cursor-123", }) - const response = await POST(buildRequest()) + const response = (await POST(buildRequest()))! const payload = await response.json() expect(response.status).toBe(200) diff --git a/__tests__/auth-hardening.test.ts b/__tests__/auth-hardening.test.ts index 85b713b8..f8601099 100644 --- a/__tests__/auth-hardening.test.ts +++ b/__tests__/auth-hardening.test.ts @@ -149,7 +149,7 @@ describe("requireRecentAuth", () => { // ── Session revocation ──────────────────────────────────────────────────────── -vi.mock("../models/RevokedSession", () => ({ default: { create: vi.fn().mockResolvedValue({}), findOne: vi.fn() } }), { virtual: true }) +vi.mock("../models/RevokedSession", () => ({ default: { create: vi.fn().mockResolvedValue({}), findOne: vi.fn() } })) // Mongoose mock so the schema registration doesn't fail in unit test context. vi.mock("mongoose", async () => { @@ -276,7 +276,7 @@ describe("POST /api/auth/stellar/link — recent-auth enforcement", () => { getAuthenticatedUser.mockResolvedValue({ user: makeUser(), shouldRefreshSession: false }) mockExtractPrivyToken.mockReturnValue(null) - const response = await POST(buildRequest({ stellarPublicKey: VALID_STELLAR_KEY })) + const response = (await POST(buildRequest({ stellarPublicKey: VALID_STELLAR_KEY })))! expect(response.status).toBe(401) const body = await response.json() expect(body.code).toBe("RECENT_AUTH_REQUIRED") @@ -287,7 +287,7 @@ describe("POST /api/auth/stellar/link — recent-auth enforcement", () => { mockExtractPrivyToken.mockReturnValue("bad-token") mockVerifyPrivyToken.mockRejectedValue(new Error("JWTExpired")) - const response = await POST(buildRequest({ stellarPublicKey: VALID_STELLAR_KEY })) + const response = (await POST(buildRequest({ stellarPublicKey: VALID_STELLAR_KEY })))! expect(response.status).toBe(401) }) @@ -300,7 +300,7 @@ describe("POST /api/auth/stellar/link — recent-auth enforcement", () => { iat: Math.floor(Date.now() / 1_000) - 30, }) - const response = await POST(buildRequest({ stellarPublicKey: VALID_STELLAR_KEY })) + const response = (await POST(buildRequest({ stellarPublicKey: VALID_STELLAR_KEY })))! expect(response.status).toBe(401) expect(mockLogAuthEvent).toHaveBeenCalledWith(expect.objectContaining({ type: "privy_subject_mismatch" })) }) @@ -314,7 +314,7 @@ describe("POST /api/auth/stellar/link — recent-auth enforcement", () => { iat: Math.floor(Date.now() / 1_000) - (RECENT_AUTH_CRITICAL_MAX_AGE_SECONDS + 30), }) - const response = await POST(buildRequest({ stellarPublicKey: VALID_STELLAR_KEY })) + const response = (await POST(buildRequest({ stellarPublicKey: VALID_STELLAR_KEY })))! expect(response.status).toBe(401) expect(mockLogAuthEvent).toHaveBeenCalledWith( expect.objectContaining({ type: "high_risk_action_denied_recent_auth" }), diff --git a/app/api/admin/reports/export/route.ts b/app/api/admin/reports/export/route.ts index 467cf936..6571f1ab 100644 --- a/app/api/admin/reports/export/route.ts +++ b/app/api/admin/reports/export/route.ts @@ -1,87 +1,61 @@ -import mongoose from "mongoose" import { NextResponse } from "next/server" -import dbConnect from "@/lib/dbConnect" import { withSessionRefresh } from "@/lib/auth/current-user" import { authorizeRequest } from "@/lib/authorization/route" +import dbConnect from "@/lib/dbConnect" +import { createCsvStream } from "@/lib/exports/csv-stream" import DriverPayment from "@/models/DriverPayment" import HirePurchaseContract from "@/models/HirePurchaseContract" import Investment from "@/models/Investment" +import InvestmentPool from "@/models/InvestmentPool" import PoolInvestment from "@/models/PoolInvestment" import Transaction from "@/models/Transaction" import User from "@/models/User" import Vehicle from "@/models/Vehicle" -import InvestmentPool from "@/models/InvestmentPool" type ExportType = "deposits" | "investments" | "repayments" | "kyc" | "fleet" | "users" type RangeType = "7d" | "30d" | "90d" | "all" | "custom" -function csvEscape(value: unknown): string { - const raw = value == null ? "" : String(value) - if (raw.includes(",") || raw.includes("\"") || raw.includes("\n")) { - return `"${raw.replace(/"/g, "\"\"")}"` - } - return raw -} - -function toCsv(headers: string[], rows: Array>) { - const lines = [headers.map(csvEscape).join(",")] - for (const row of rows) { - lines.push(row.map(csvEscape).join(",")) - } - return lines.join("\n") -} +const CURSOR_BATCH_SIZE = 250 function parseRange(raw: string | null): RangeType { - if (raw === "7d" || raw === "30d" || raw === "90d" || raw === "all" || raw === "custom") return raw - return "30d" + return raw === "7d" || raw === "30d" || raw === "90d" || raw === "all" || raw === "custom" ? raw : "30d" } function buildWindow(range: RangeType, fromRaw: string | null, toRaw: string | null) { if (range === "all") return { startDate: null as Date | null, endDate: null as Date | null } - if (range === "custom") { - const fromDate = fromRaw ? new Date(fromRaw) : null - const toDate = toRaw ? new Date(toRaw) : null - if (fromDate && toDate && !Number.isNaN(fromDate.getTime()) && !Number.isNaN(toDate.getTime())) { - const endDate = new Date(toDate) + const startDate = fromRaw ? new Date(fromRaw) : null + const endDate = toRaw ? new Date(toRaw) : null + if (startDate && endDate && !Number.isNaN(startDate.getTime()) && !Number.isNaN(endDate.getTime())) { endDate.setHours(23, 59, 59, 999) - return { startDate: fromDate, endDate } + return { startDate, endDate } } } - - const days = range === "7d" ? 7 : range === "90d" ? 90 : 30 const startDate = new Date() - startDate.setDate(startDate.getDate() - days) + startDate.setDate(startDate.getDate() - (range === "7d" ? 7 : range === "90d" ? 90 : 30)) return { startDate, endDate: null as Date | null } } function dateMatch(field: string, startDate: Date | null, endDate: Date | null) { - if (!startDate && !endDate) return {} - const clause: Record = {} - if (startDate) clause.$gte = startDate - if (endDate) clause.$lte = endDate - return { [field]: clause } + const range: Record = {} + if (startDate) range.$gte = startDate + if (endDate) range.$lte = endDate + return Object.keys(range).length ? { [field]: range } : {} } -function getUserName(user: any) { +function userName(user: any) { return user?.fullName || user?.name || user?.email || "Unknown User" } -function normalizeObjectId(value: unknown) { - if (value instanceof mongoose.Types.ObjectId) return value - if (typeof value === "string" && mongoose.Types.ObjectId.isValid(value)) return new mongoose.Types.ObjectId(value) - return null -} - -function collectObjectIds(values: unknown[]) { - const map = new Map() - for (const value of values) { - const normalized = normalizeObjectId(value) - if (!normalized) continue - map.set(normalized.toString(), normalized) - } - return Array.from(map.values()) +function csvResponse(headers: string[], rows: AsyncIterable, filename: string) { + return new NextResponse(createCsvStream(headers, rows), { + headers: { + "Content-Type": "text/csv; charset=utf-8", + "Content-Disposition": `attachment; filename="${filename}"`, + "Cache-Control": "no-store", + }, + }) } export async function GET(request: Request) { @@ -89,305 +63,54 @@ export async function GET(request: Request) { const auth = await authorizeRequest(request, "admin:report", { type: "report" }) if ("response" in auth) return auth.response const { user, shouldRefreshSession } = auth - await dbConnect() const { searchParams } = new URL(request.url) const type = (searchParams.get("type") || "deposits") as ExportType - if (!["deposits", "investments", "repayments", "kyc", "fleet", "users"].includes(type)) { + if (!(["deposits", "investments", "repayments", "kyc", "fleet", "users"] as string[]).includes(type)) { return NextResponse.json({ message: "Invalid export type" }, { status: 400 }) } - const range = parseRange(searchParams.get("range")) const { startDate, endDate } = buildWindow(range, searchParams.get("from"), searchParams.get("to")) + let response: NextResponse - // ── deposits ──────────────────────────────────────────────────────────── if (type === "deposits") { - const deposits = await Transaction.find({ - type: { $in: ["deposit", "wallet_funding"] }, - status: { $in: ["Completed", "completed", "SUCCESS", "success", "Successful", "successful"] }, - ...dateMatch("timestamp", startDate, endDate), - }) - .select("userId amount method status gatewayReference timestamp") - .sort({ timestamp: -1 }) - .lean() - - const userIds = collectObjectIds(deposits.map((item: any) => item.userId)) - const users = userIds.length - ? await User.find({ _id: { $in: userIds } }).select("name fullName email").lean() - : [] - const userById = new Map(users.map((entry: any) => [entry._id.toString(), entry])) - - const csv = toCsv( - ["Date", "User", "Email", "Amount (NGN)", "Method", "Status", "Reference"], - deposits.map((item: any) => { - const userEntry = userById.get(item.userId?.toString?.() || "") - return [ - item.timestamp ? new Date(item.timestamp).toISOString() : "", - getUserName(userEntry), - userEntry?.email || "", - Number(item.amount || 0), - item.method || "unknown", - item.status || "unknown", - item.gatewayReference || "", - ] - }), - ) - - const response = new NextResponse(csv, { - status: 200, - headers: { - "Content-Type": "text/csv; charset=utf-8", - "Content-Disposition": `attachment; filename="deposits-${range}.csv"`, - }, - }) - return shouldRefreshSession ? withSessionRefresh(response, user) : response - } - - // ── investments ────────────────────────────────────────────────────────── - if (type === "investments") { - const [poolInvestments, legacyInvestments] = await Promise.all([ - PoolInvestment.find({ - status: "CONFIRMED", - ...dateMatch("createdAt", startDate, endDate), - }) - .select("userId poolId amountNgn ownershipBps txRef status createdAt") - .sort({ createdAt: -1 }) - .lean(), - Investment.find({ - status: { $in: ["Active", "Completed"] }, - ...dateMatch("date", startDate, endDate), - }) - .select("investorId vehicleId amount status date") - .sort({ date: -1 }) - .lean(), - ]) - - const userIds = collectObjectIds([ - ...poolInvestments.map((item: any) => item.userId), - ...legacyInvestments.map((item: any) => item.investorId), - ]) - const poolIds = collectObjectIds(poolInvestments.map((item: any) => item.poolId)) - const vehicleIds = collectObjectIds(legacyInvestments.map((item: any) => item.vehicleId)) - - const [users, pools, vehicles] = await Promise.all([ - userIds.length ? User.find({ _id: { $in: userIds } }).select("name fullName email").lean() : Promise.resolve([]), - poolIds.length ? InvestmentPool.find({ _id: { $in: poolIds } }).select("assetType status").lean() : Promise.resolve([]), - vehicleIds.length ? Vehicle.find({ _id: { $in: vehicleIds } }).select("name type").lean() : Promise.resolve([]), - ]) - - const userById = new Map(users.map((entry: any) => [entry._id.toString(), entry])) - const poolById = new Map(pools.map((entry: any) => [entry._id.toString(), entry])) - const vehicleById = new Map(vehicles.map((entry: any) => [entry._id.toString(), entry])) - - const poolRows = poolInvestments.map((item: any) => { - const userEntry = userById.get(item.userId?.toString?.() || "") - const pool = poolById.get(item.poolId?.toString?.() || "") - return [ - item.createdAt ? new Date(item.createdAt).toISOString() : "", - getUserName(userEntry), - userEntry?.email || "", - "pool", - pool ? `${pool.assetType} (${pool.status})` : "Pool", - Number(item.amountNgn || 0), - Number(item.ownershipBps || 0) / 100, - item.status || "unknown", - item.txRef || "", - ] - }) - - const legacyRows = legacyInvestments.map((item: any) => { - const userEntry = userById.get(item.investorId?.toString?.() || "") - const vehicle = vehicleById.get(item.vehicleId?.toString?.() || "") - return [ - item.date ? new Date(item.date).toISOString() : "", - getUserName(userEntry), - userEntry?.email || "", - "legacy", - vehicle?.name || vehicle?.type || "Vehicle", - Number(item.amount || 0), - "", - item.status || "unknown", - "", - ] - }) - - const csv = toCsv( - ["Date", "User", "Email", "Source", "Asset", "Amount (NGN)", "Ownership (%)", "Status", "Reference"], - [...poolRows, ...legacyRows], - ) - - const response = new NextResponse(csv, { - status: 200, - headers: { - "Content-Type": "text/csv; charset=utf-8", - "Content-Disposition": `attachment; filename="investments-${range}.csv"`, - }, - }) - return shouldRefreshSession ? withSessionRefresh(response, user) : response - } - - // ── repayments ─────────────────────────────────────────────────────────── - if (type === "repayments") { - const repayments = await DriverPayment.find({ - status: "CONFIRMED", - ...dateMatch("createdAt", startDate, endDate), - }) - .select("driverUserId contractId amountNgn appliedAmountNgn method paystackRef status createdAt") - .sort({ createdAt: -1 }) - .lean() - - const userIds = collectObjectIds(repayments.map((item: any) => item.driverUserId)) - const contractIds = collectObjectIds(repayments.map((item: any) => item.contractId)) - - const [users, contracts] = await Promise.all([ - userIds.length ? User.find({ _id: { $in: userIds } }).select("name fullName email").lean() : Promise.resolve([]), - contractIds.length - ? HirePurchaseContract.find({ _id: { $in: contractIds } }).select("vehicleDisplayName").lean() - : Promise.resolve([]), - ]) - - const userById = new Map(users.map((entry: any) => [entry._id.toString(), entry])) - const contractById = new Map(contracts.map((entry: any) => [entry._id.toString(), entry])) - - const csv = toCsv( - ["Date", "Driver", "Email", "Vehicle/Contract", "Amount (NGN)", "Applied Amount (NGN)", "Method", "Reference", "Status"], - repayments.map((item: any) => { - const userEntry = userById.get(item.driverUserId?.toString?.() || "") - const contract = contractById.get(item.contractId?.toString?.() || "") - return [ - item.createdAt ? new Date(item.createdAt).toISOString() : "", - getUserName(userEntry), - userEntry?.email || "", - contract?.vehicleDisplayName || "Contract", - Number(item.amountNgn || 0), - Number(item.appliedAmountNgn || 0), - item.method || "PAYSTACK", - item.paystackRef || "", - item.status || "unknown", - ] - }), - ) - - const response = new NextResponse(csv, { - status: 200, - headers: { - "Content-Type": "text/csv; charset=utf-8", - "Content-Disposition": `attachment; filename="repayments-${range}.csv"`, - }, - }) - return shouldRefreshSession ? withSessionRefresh(response, user) : response - } - - // ── kyc ────────────────────────────────────────────────────────────────── - if (type === "kyc") { - const kycStatusFilter = searchParams.get("status") || "" - const kycQuery: Record = { - kycStatus: { $ne: "none" }, - ...dateMatch("createdAt", startDate, endDate), + async function* rows(): AsyncGenerator { + const cursor = Transaction.find({ type: { $in: ["deposit", "wallet_funding"] }, status: { $in: ["Completed", "completed", "SUCCESS", "success", "Successful", "successful"] }, ...dateMatch("timestamp", startDate, endDate) }) + .select("userId amount method status gatewayReference timestamp").populate({ path: "userId", select: "name fullName email" }).sort({ timestamp: -1, _id: -1 }).lean().cursor({ batchSize: CURSOR_BATCH_SIZE }) + for await (const item of cursor as AsyncIterable) yield [item.timestamp ? new Date(item.timestamp).toISOString() : "", userName(item.userId), item.userId?.email || "", Number(item.amount || 0), item.method || "unknown", item.status || "unknown", item.gatewayReference || ""] } - if (["pending", "approved", "rejected"].includes(kycStatusFilter)) { - kycQuery.kycStatus = kycStatusFilter + response = csvResponse(["Date", "User", "Email", "Amount (NGN)", "Method", "Status", "Reference"], rows(), `deposits-${range}.csv`) + } else if (type === "investments") { + async function* rows(): AsyncGenerator { + const poolCursor = PoolInvestment.find({ status: "CONFIRMED", ...dateMatch("createdAt", startDate, endDate) }).select("userId poolId amountNgn ownershipBps txRef status createdAt").populate([{ path: "userId", select: "name fullName email" }, { path: "poolId", select: "assetType status" }]).sort({ createdAt: -1, _id: -1 }).lean().cursor({ batchSize: CURSOR_BATCH_SIZE }) + for await (const item of poolCursor as AsyncIterable) yield [item.createdAt ? new Date(item.createdAt).toISOString() : "", userName(item.userId), item.userId?.email || "", "pool", item.poolId ? `${item.poolId.assetType} (${item.poolId.status})` : "Pool", Number(item.amountNgn || 0), Number(item.ownershipBps || 0) / 100, item.status || "unknown", item.txRef || ""] + const legacyCursor = Investment.find({ status: { $in: ["Active", "Completed"] }, ...dateMatch("date", startDate, endDate) }).select("investorId vehicleId amount status date").populate([{ path: "investorId", select: "name fullName email" }, { path: "vehicleId", select: "name type" }]).sort({ date: -1, _id: -1 }).lean().cursor({ batchSize: CURSOR_BATCH_SIZE }) + for await (const item of legacyCursor as AsyncIterable) yield [item.date ? new Date(item.date).toISOString() : "", userName(item.investorId), item.investorId?.email || "", "legacy", item.vehicleId?.name || item.vehicleId?.type || "Vehicle", Number(item.amount || 0), "", item.status || "unknown", ""] } - - const kycUsers = await User.find(kycQuery) - .select("name fullName email role kycStatus kycVerified createdAt") - .sort({ createdAt: -1 }) - .lean() - - const csv = toCsv( - ["Date Joined", "Name", "Email", "Role", "KYC Status", "KYC Verified"], - kycUsers.map((u: any) => [ - u.createdAt ? new Date(u.createdAt).toISOString() : "", - getUserName(u), - u.email || "", - u.role || "", - u.kycStatus || "none", - u.kycVerified ? "Yes" : "No", - ]), - ) - - const response = new NextResponse(csv, { - status: 200, - headers: { - "Content-Type": "text/csv; charset=utf-8", - "Content-Disposition": `attachment; filename="kyc-report-${range}.csv"`, - }, - }) - return shouldRefreshSession ? withSessionRefresh(response, user) : response - } - - // ── fleet ───────────────────────────────────────────────────────────────── - if (type === "fleet") { - const vehicleStatusFilter = searchParams.get("vstatus") || "" - const vehicleQuery: Record = {} - if (["Available", "Financed", "Reserved", "Maintenance", "Retired"].includes(vehicleStatusFilter)) { - vehicleQuery.status = vehicleStatusFilter + response = csvResponse(["Date", "User", "Email", "Source", "Asset", "Amount (NGN)", "Ownership (%)", "Status", "Reference"], rows(), `investments-${range}.csv`) + } else if (type === "repayments") { + async function* rows(): AsyncGenerator { + const cursor = DriverPayment.find({ status: "CONFIRMED", ...dateMatch("createdAt", startDate, endDate) }).select("driverUserId contractId amountNgn appliedAmountNgn method paystackRef status createdAt").populate([{ path: "driverUserId", select: "name fullName email" }, { path: "contractId", select: "vehicleDisplayName" }]).sort({ createdAt: -1, _id: -1 }).lean().cursor({ batchSize: CURSOR_BATCH_SIZE }) + for await (const item of cursor as AsyncIterable) yield [item.createdAt ? new Date(item.createdAt).toISOString() : "", userName(item.driverUserId), item.driverUserId?.email || "", item.contractId?.vehicleDisplayName || "Contract", Number(item.amountNgn || 0), Number(item.appliedAmountNgn || 0), item.method || "PAYSTACK", item.paystackRef || "", item.status || "unknown"] } - - const vehicles = await Vehicle.find(vehicleQuery) - .select("name identifier type year price status fundingStatus totalFundedAmount addedDate") - .sort({ addedDate: -1 }) - .lean() - - const csv = toCsv( - ["Date Added", "Name", "Identifier", "Type", "Year", "Price (NGN)", "Status", "Funding Status", "Total Funded (NGN)"], - vehicles.map((v: any) => [ - v.addedDate ? new Date(v.addedDate).toISOString() : "", - v.name || "", - v.identifier || "", - v.type || "", - v.year || "", - Number(v.price || 0), - v.status || "", - v.fundingStatus || "", - Number(v.totalFundedAmount || 0), - ]), - ) - - const response = new NextResponse(csv, { - status: 200, - headers: { - "Content-Type": "text/csv; charset=utf-8", - "Content-Disposition": `attachment; filename="fleet-report-${range}.csv"`, - }, - }) - return shouldRefreshSession ? withSessionRefresh(response, user) : response + response = csvResponse(["Date", "Driver", "Email", "Vehicle/Contract", "Amount (NGN)", "Applied Amount (NGN)", "Method", "Reference", "Status"], rows(), `repayments-${range}.csv`) + } else if (type === "kyc") { + const query: Record = { kycStatus: { $ne: "none" }, ...dateMatch("createdAt", startDate, endDate) } + const status = searchParams.get("status") || "" + if (["pending", "approved", "rejected"].includes(status)) query.kycStatus = status + async function* rows(): AsyncGenerator { const cursor = User.find(query).select("name fullName email role kycStatus kycVerified createdAt").sort({ createdAt: -1, _id: -1 }).lean().cursor({ batchSize: CURSOR_BATCH_SIZE }); for await (const item of cursor as AsyncIterable) yield [item.createdAt ? new Date(item.createdAt).toISOString() : "", userName(item), item.email || "", item.role || "", item.kycStatus || "none", item.kycVerified ? "Yes" : "No"] } + response = csvResponse(["Date Joined", "Name", "Email", "Role", "KYC Status", "KYC Verified"], rows(), `kyc-report-${range}.csv`) + } else if (type === "fleet") { + const query: Record = {}; const status = searchParams.get("vstatus") || ""; if (["Available", "Financed", "Reserved", "Maintenance", "Retired"].includes(status)) query.status = status + async function* rows(): AsyncGenerator { const cursor = Vehicle.find(query).select("name identifier type year price status fundingStatus totalFundedAmount addedDate").sort({ addedDate: -1, _id: -1 }).lean().cursor({ batchSize: CURSOR_BATCH_SIZE }); for await (const item of cursor as AsyncIterable) yield [item.addedDate ? new Date(item.addedDate).toISOString() : "", item.name || "", item.identifier || "", item.type || "", item.year || "", Number(item.price || 0), item.status || "", item.fundingStatus || "", Number(item.totalFundedAmount || 0)] } + response = csvResponse(["Date Added", "Name", "Identifier", "Type", "Year", "Price (NGN)", "Status", "Funding Status", "Total Funded (NGN)"], rows(), `fleet-report-${range}.csv`) + } else { + const query: Record = { ...dateMatch("createdAt", startDate, endDate) }; const role = searchParams.get("role") || ""; if (["driver", "investor", "admin"].includes(role)) query.role = role + async function* rows(): AsyncGenerator { const cursor = User.find(query).select("name fullName email role kycVerified createdAt").sort({ createdAt: -1, _id: -1 }).lean().cursor({ batchSize: CURSOR_BATCH_SIZE }); for await (const item of cursor as AsyncIterable) yield [item.createdAt ? new Date(item.createdAt).toISOString() : "", userName(item), item.email || "", item.role || "", item.kycVerified ? "Yes" : "No"] } + response = csvResponse(["Date Joined", "Name", "Email", "Role", "KYC Verified"], rows(), `users-report-${range}.csv`) } - - // ── users ───────────────────────────────────────────────────────────────── - const roleFilter = searchParams.get("role") || "" - const userQuery: Record = { - ...dateMatch("createdAt", startDate, endDate), - } - if (["driver", "investor", "admin"].includes(roleFilter)) { - userQuery.role = roleFilter - } - - const platformUsers = await User.find(userQuery) - .select("name fullName email role kycVerified createdAt") - .sort({ createdAt: -1 }) - .lean() - - const csv = toCsv( - ["Date Joined", "Name", "Email", "Role", "KYC Verified"], - platformUsers.map((u: any) => [ - u.createdAt ? new Date(u.createdAt).toISOString() : "", - getUserName(u), - u.email || "", - u.role || "", - u.kycVerified ? "Yes" : "No", - ]), - ) - - const response = new NextResponse(csv, { - status: 200, - headers: { - "Content-Type": "text/csv; charset=utf-8", - "Content-Disposition": `attachment; filename="users-report-${range}.csv"`, - }, - }) - return shouldRefreshSession ? withSessionRefresh(response, user) : response + return shouldRefreshSession ? await withSessionRefresh(response, user) : response } catch (error) { console.error("ADMIN_REPORT_EXPORT_ERROR", error) return NextResponse.json({ message: "Failed to export report." }, { status: 500 }) diff --git a/app/api/recovery/[id]/execute/route.ts b/app/api/recovery/[id]/execute/route.ts index a0d7b205..958b6bf1 100644 --- a/app/api/recovery/[id]/execute/route.ts +++ b/app/api/recovery/[id]/execute/route.ts @@ -15,7 +15,7 @@ import WalletRecovery from "@/models/WalletRecovery" import WalletMigrationRecord from "@/models/WalletMigrationRecord" import User from "@/models/User" import { getAuthenticatedUser } from "@/lib/auth/current-user" -import { assertTransition, isTerminal, allFactorsVerified as _allFactorsVerified } from "@/lib/recovery/recovery-state-machine" +import { assertTransition, isTerminal } from "@/lib/recovery/recovery-state-machine" import { allFactorsVerified } from "@/lib/recovery/recovery-factors" import { sendRecoveryNotifications } from "@/lib/recovery/recovery-notifications" diff --git a/app/api/transactions/ledger/export/route.ts b/app/api/transactions/ledger/export/route.ts index 745e2ff3..da697cb1 100644 --- a/app/api/transactions/ledger/export/route.ts +++ b/app/api/transactions/ledger/export/route.ts @@ -4,6 +4,7 @@ import { z } from "zod" import { normalizeUserRole, requireAuthenticatedUser } from "@/lib/api/route-guard" import { parseSearchParams } from "@/lib/api/validation" import dbConnect from "@/lib/dbConnect" +import { createCsvStream } from "@/lib/exports/csv-stream" import { buildLedgerFilter, normalizeLedgerEntry, @@ -13,7 +14,7 @@ import { import Transaction from "@/models/Transaction" import "@/models/User" -const EXPORT_LIMIT = 5000 +const CURSOR_BATCH_SIZE = 250 const querySchema = z.object({ search: z.string().trim().max(120).default(""), @@ -27,14 +28,6 @@ const querySchema = z.object({ userId: z.string().trim().max(40).default(""), }) -function csvEscape(value: unknown): string { - const raw = value == null ? "" : String(value) - if (raw.includes(",") || raw.includes("\"") || raw.includes("\n")) { - return `"${raw.replace(/"/g, "\"\"")}"` - } - return raw -} - export async function GET(request: Request) { try { const authContext = await requireAuthenticatedUser(request, ["admin", "driver", "investor"]) @@ -50,7 +43,7 @@ export async function GET(request: Request) { await dbConnect() - const params = { ...parsed.data, page: 1, pageSize: EXPORT_LIMIT } as LedgerQueryParams + const params = { ...parsed.data, page: 1, pageSize: CURSOR_BATCH_SIZE } as LedgerQueryParams const actor: LedgerActor = { id: authContext.user._id.toString(), role } const isAdmin = role === "admin" @@ -76,11 +69,12 @@ export async function GET(request: Request) { queryFilter.gatewayReference = { $nin: Array.from(duplicateReferences) } } - const pageQuery = Transaction.find(queryFilter).sort({ timestamp: -1 }).limit(EXPORT_LIMIT) + // `_id` breaks timestamp ties, so the cursor has a deterministic order even + // while new transactions are being written after the export has started. + const pageQuery = Transaction.find(queryFilter).sort({ timestamp: -1, _id: -1 }) if (isAdmin) { pageQuery.populate({ path: "userId", select: "name fullName email role" }) } - const transactions = (await pageQuery.lean()) as Array> const headers = [ "Date", @@ -96,29 +90,29 @@ export async function GET(request: Request) { ...(isAdmin ? ["User", "User Email", "User Type"] : []), ] - const lines = [headers.map(csvEscape).join(",")] - for (const tx of transactions) { - const entry = normalizeLedgerEntry(tx, duplicateReferences) - const row = [ - entry.timestamp, - entry.type, - entry.direction, - entry.amount, - entry.currency, - entry.status, - entry.reconciliation, - entry.method ?? "", - entry.reference ?? "", - entry.description, - ...(isAdmin ? [entry.userName ?? "", entry.userEmail ?? "", entry.userType] : []), - ] - lines.push(row.map(csvEscape).join(",")) + async function* rows(): AsyncGenerator { + const cursor = pageQuery.lean().cursor({ batchSize: CURSOR_BATCH_SIZE }) + for await (const tx of cursor as AsyncIterable>) { + const entry = normalizeLedgerEntry(tx, duplicateReferences) + yield [ + entry.timestamp, + entry.type, + entry.direction, + entry.amount, + entry.currency, + entry.status, + entry.reconciliation, + entry.method ?? "", + entry.reference ?? "", + entry.description, + ...(isAdmin ? [entry.userName ?? "", entry.userEmail ?? "", entry.userType] : []), + ] + } } - const csv = lines.join("\n") const filename = `transaction-ledger-${new Date().toISOString().slice(0, 10)}.csv` - return new NextResponse(csv, { + return new NextResponse(createCsvStream(headers, rows()), { status: 200, headers: { "Content-Type": "text/csv; charset=utf-8", diff --git a/components/ui/alert.tsx b/components/ui/alert.tsx index 57996ebc..6eb02d7a 100644 --- a/components/ui/alert.tsx +++ b/components/ui/alert.tsx @@ -9,8 +9,10 @@ const alertVariants = cva( variants: { variant: { default: "bg-background text-foreground", - destructive: - "border-destructive/50 text-destructive dark:border-destructive [&>svg]:text-destructive", + destructive: + "border-destructive/50 text-destructive dark:border-destructive [&>svg]:text-destructive", + warning: + "border-amber-500/50 text-amber-700 dark:text-amber-400 [&>svg]:text-amber-600 dark:[&>svg]:text-amber-400", }, }, defaultVariants: { diff --git a/lib/exports/csv-stream.ts b/lib/exports/csv-stream.ts new file mode 100644 index 00000000..57936ca8 --- /dev/null +++ b/lib/exports/csv-stream.ts @@ -0,0 +1,43 @@ +const encoder = new TextEncoder() + +export function csvEscape(value: unknown): string { + const raw = value == null ? "" : String(value) + return raw.includes(",") || raw.includes("\"") || raw.includes("\n") + ? `"${raw.replace(/"/g, "\"\"")}"` + : raw +} + +/** + * Converts an async row source to a response body. The stream only requests the + * next row when the consumer is ready, so database cursors remain bounded by + * their configured batch size instead of the size of the export. + */ +export function createCsvStream(headers: string[], rows: AsyncIterable): ReadableStream { + const iterator = rows[Symbol.asyncIterator]() + let wroteHeader = false + + return new ReadableStream({ + async pull(controller) { + try { + if (!wroteHeader) { + wroteHeader = true + controller.enqueue(encoder.encode(`\uFEFF${headers.map(csvEscape).join(",")}\n`)) + return + } + + const next = await iterator.next() + if (next.done) { + controller.close() + return + } + controller.enqueue(encoder.encode(`${next.value.map(csvEscape).join(",")}\n`)) + } catch (error) { + await iterator.return?.() + controller.error(error) + } + }, + async cancel() { + await iterator.return?.() + }, + }) +} diff --git a/lib/security/audit-export.ts b/lib/security/audit-export.ts index 8ba49326..702b5087 100644 --- a/lib/security/audit-export.ts +++ b/lib/security/audit-export.ts @@ -3,6 +3,7 @@ import TamperEvidentAuditLog from "@/models/TamperEvidentAuditLog" import AuditCheckpoint from "@/models/AuditCheckpoint" import { verifyAuditChain } from "./audit-verification" import { buildCanonicalAuditEventData, canonicalizeEventData, computeEventHash, redactPII } from "./audit-hash" +import { createCsvStream } from "@/lib/exports/csv-stream" export interface ExportOptions { partition: string @@ -53,69 +54,73 @@ export interface ExportManifest { piiRedacted: boolean } -/** - * Export audit events with integrity manifest - */ -export async function exportAuditEvents(options: ExportOptions): Promise { - await dbConnect() - - // Build query - const query: any = { - partition: options.partition, - } - - if (!options.includeLegacy) { - query.isLegacy = false - } +const AUDIT_EXPORT_BATCH_SIZE = 250 +function buildAuditExportQuery(options: ExportOptions) { + const query: any = { partition: options.partition } + if (!options.includeLegacy) query.isLegacy = false if (options.startDate || options.endDate) { query.timestamp = {} if (options.startDate) query.timestamp.$gte = options.startDate if (options.endDate) query.timestamp.$lte = options.endDate } - if (options.startSequence !== undefined || options.endSequence !== undefined) { query.sequence = {} if (options.startSequence !== undefined) query.sequence.$gte = options.startSequence if (options.endSequence !== undefined) query.sequence.$lte = options.endSequence } + if (options.actions?.length) query.action = { $in: options.actions } + if (options.actorId) query.actorId = options.actorId + if (options.targetType) query.targetType = options.targetType + return query +} - if (options.actions && options.actions.length > 0) { - query.action = { $in: options.actions } +/** Reads audit events in bounded database batches for response/file streaming. */ +export async function* iterateAuditEvents(options: ExportOptions): AsyncGenerator { + await dbConnect() + const cursor = TamperEvidentAuditLog.find(buildAuditExportQuery(options)) + .sort({ sequence: 1, _id: 1 }) + .lean() + .cursor({ batchSize: AUDIT_EXPORT_BATCH_SIZE }) + let previousHash: string | undefined + for await (const event of cursor as AsyncIterable) { + if (!options.redactPII) { + yield event + continue + } + const redactedEvent = { + ...event, + actorIdentifier: event.actorIdentifier ? "[REDACTED]" : event.actorIdentifier, + metadata: redactPII(event.metadata), + sourceEventHash: event.eventHash, + previousHash: previousHash ?? event.previousHash, + } + const canonicalData = canonicalizeEventData(buildCanonicalAuditEventData(redactedEvent)) + const eventHash = computeEventHash(redactedEvent.previousHash + canonicalData) + previousHash = eventHash + yield { ...redactedEvent, canonicalData, eventHash } } +} - if (options.actorId) { - query.actorId = options.actorId +export function createAuditCsvStream(events: AsyncIterable): ReadableStream { + const headers = ["Sequence", "Event ID", "Timestamp", "Actor ID", "Actor Role", "Action", "Target Type", "Target ID", "Status", "Request ID", "IP Address", "Previous Hash", "Event Hash", "Is Legacy"] + async function* rows(): AsyncGenerator { + for await (const event of events) yield [event.sequence, event.eventId, event.timestamp?.toISOString?.() ?? "", event.actorId || "", event.actorRole || "", event.action, event.targetType, event.targetId || "", event.status, event.requestId || "", event.ipAddress || "", event.previousHash, event.eventHash, event.isLegacy ? "true" : "false"] } + return createCsvStream(headers, rows()) +} - if (options.targetType) { - query.targetType = options.targetType - } +/** + * Export audit events with integrity manifest + */ +export async function exportAuditEvents(options: ExportOptions): Promise { + await dbConnect() - // Fetch events - let events: any[] = await TamperEvidentAuditLog.find(query).sort({ sequence: 1 }).lean() - - if (options.redactPII) { - let previousHash = events[0]?.previousHash - events = events.map((event) => { - const redactedEvent = { - ...event, - actorIdentifier: event.actorIdentifier ? "[REDACTED]" : event.actorIdentifier, - metadata: redactPII(event.metadata), - sourceEventHash: event.eventHash, - previousHash, - } - const canonicalData = canonicalizeEventData(buildCanonicalAuditEventData(redactedEvent)) - const eventHash = computeEventHash(redactedEvent.previousHash + canonicalData) - previousHash = eventHash - - return { - ...redactedEvent, - canonicalData, - eventHash, - } - }) - } + // Build query + // Kept for the CLI's manifest-producing API. HTTP callers should use + // iterateAuditEvents/createAuditCsvStream so output is never materialized. + const events: any[] = [] + for await (const event of iterateAuditEvents(options)) events.push(event) // Fetch checkpoints if requested let checkpoints: any[] = [] diff --git a/lib/settlement/settlement-service.ts b/lib/settlement/settlement-service.ts index 508b2889..542040c8 100644 --- a/lib/settlement/settlement-service.ts +++ b/lib/settlement/settlement-service.ts @@ -361,7 +361,10 @@ export async function evaluateFinalityTimeouts(): Promise<{ let expiredCount = 0 for (const s of activeSettlements) { - const config = getRailSettlementConfig(s.rail, s.environment) + const config = getRailSettlementConfig( + s.rail, + s.environment as "development" | "production" | "test" | undefined, + ) const ageMs = now.getTime() - new Date(s.createdAt).getTime() const thresholdMs = s.currentState === "observed" || s.currentState === "provisionally_credited"