diff --git a/__tests__/lib/settlement/settlement.test.ts b/__tests__/lib/settlement/settlement.test.ts new file mode 100644 index 00000000..4585462f --- /dev/null +++ b/__tests__/lib/settlement/settlement.test.ts @@ -0,0 +1,123 @@ +// @vitest-environment node +import { describe, expect, it, beforeEach } from "vitest" +import mongoose from "mongoose" + +import { + determineSafeActions, + isValidTransition, +} from "@/lib/settlement/state-machine" +import { getRailSettlementConfig } from "@/lib/settlement/config" +import { CanonicalSettlementState } from "@/models/SettlementRecord" + +describe("Settlement State Machine & Configuration", () => { + it("validates allowed canonical state transitions correctly", () => { + expect(isValidTransition("initiated", "provider-pending")).toBe(true) + expect(isValidTransition("provider-pending", "confirmed")).toBe(true) + expect(isValidTransition("confirmed", "reversed")).toBe(true) + expect(isValidTransition("confirmed", "disputed")).toBe(true) + expect(isValidTransition("disputed", "confirmed")).toBe(true) + expect(isValidTransition("disputed", "reversed")).toBe(true) + + // Invalid terminal state transitions + expect(isValidTransition("reversed", "confirmed")).toBe(false) + expect(isValidTransition("failed", "confirmed")).toBe(false) + expect(isValidTransition("expired", "observed")).toBe(false) + }) + + it("returns appropriate safe operator actions for each state", () => { + expect(determineSafeActions("initiated", false)).toContain("RETRY_VERIFICATION") + expect(determineSafeActions("provider-pending", true)).toContain("MARK_EXPIRED") + expect(determineSafeActions("confirmed", false)).toContain("POST_REVERSAL") + expect(determineSafeActions("disputed", false)).toContain("RESOLVE_DISPUTE_CONFIRM") + }) + + it("provides correct rail settlement configuration for production and development", () => { + const paystackConfig = getRailSettlementConfig("paystack", "production") + expect(paystackConfig.finalityThreshold).toBe(1) + expect(paystackConfig.pendingTimeoutMs).toBe(15 * 60 * 1000) + + const stellarConfig = getRailSettlementConfig("stellar", "production") + expect(stellarConfig.finalityThreshold).toBe(3) + + const devStellarConfig = getRailSettlementConfig("stellar", "development") + expect(devStellarConfig.finalityThreshold).toBe(1) + }) +}) + +describe("Settlement Logic & Scenarios", () => { + it("handles out-of-order webhook state transitions gracefully", () => { + // Valid transition from initiated to confirmed directly when webhook arrives early + expect(isValidTransition("initiated", "confirmed")).toBe(true) + // Valid transition from provider-pending to observed + expect(isValidTransition("provider-pending", "observed")).toBe(true) + }) + + it("prevents double-crediting on duplicate reference events", () => { + const seenReferences = new Set() + const ref = "PAYSTACK_REF_DUPLICATE_123" + + function processPayment(reference: string) { + if (seenReferences.has(reference)) { + return { alreadyProcessed: true, credited: false } + } + seenReferences.add(reference) + return { alreadyProcessed: false, credited: true } + } + + const firstCall = processPayment(ref) + const secondCall = processPayment(ref) + + expect(firstCall.alreadyProcessed).toBe(false) + expect(firstCall.credited).toBe(true) + + expect(secondCall.alreadyProcessed).toBe(true) + expect(secondCall.credited).toBe(false) + }) + + it("calculates reversal journal deductions without creating negative balance", () => { + const currentAvailable = 3000 + const reversalAmount = 5000 + + let deductedFromAvailable = 0 + let deductedFromHeld = 0 + let deductedFromPending = 0 + + if (currentAvailable >= reversalAmount) { + deductedFromAvailable = reversalAmount + } else { + deductedFromAvailable = currentAvailable + const remainder = reversalAmount - currentAvailable + deductedFromHeld = remainder + } + + const newAvailable = Math.max(currentAvailable - deductedFromAvailable, 0) + expect(newAvailable).toBe(0) + expect(deductedFromAvailable).toBe(3000) + expect(deductedFromHeld).toBe(2000) + }) + + it("flags stuck transactions when pending time exceeds rail threshold", () => { + const now = new Date() + const createdAt = new Date(now.getTime() - 30 * 60 * 1000) // 30 mins ago + const pendingTimeoutMs = 15 * 60 * 1000 // 15 mins + + const ageMs = now.getTime() - createdAt.getTime() + const isStuck = ageMs > pendingTimeoutMs + + expect(isStuck).toBe(true) + }) + + it("handles split settlement cross-rail correlation mapping", () => { + const settlementMapping = { + settlementId: "STL-999", + providerReference: "PAYSTACK_TRANSFER_888", + stellarHash: "0xSTELLARHASH777", + ledgerJournalId: "JOURNAL_666", + userTransactionId: "TX_555", + } + + expect(settlementMapping.providerReference).toBe("PAYSTACK_TRANSFER_888") + expect(settlementMapping.stellarHash).toBe("0xSTELLARHASH777") + expect(settlementMapping.ledgerJournalId).toBe("JOURNAL_666") + }) +}) diff --git a/app/api/admin/settlement/timeline/route.ts b/app/api/admin/settlement/timeline/route.ts new file mode 100644 index 00000000..6240831b --- /dev/null +++ b/app/api/admin/settlement/timeline/route.ts @@ -0,0 +1,120 @@ +import { NextResponse } from "next/server" +import { z } from "zod" + +import dbConnect from "@/lib/dbConnect" +import { finalizeAuthenticatedResponse, requireAuthenticatedUser } from "@/lib/api/route-guard" +import { parseJsonBody } from "@/lib/api/validation" +import SettlementRecord from "@/models/SettlementRecord" +import { + evaluateFinalityTimeouts, + transitionSettlementState, +} from "@/lib/settlement/settlement-service" + +const postSchema = z.object({ + action: z.enum(["EVALUATE_TIMEOUTS", "FORCE_CONFIRM", "POST_REVERSAL", "MARK_EXPIRED", "RETRY_VERIFICATION"]), + settlementId: z.string().optional(), + providerReference: z.string().optional(), + reason: z.string().optional(), +}) + +export async function GET(request: Request) { + try { + const authContext = await requireAuthenticatedUser(request, ["admin"]) + if ("response" in authContext) return authContext.response + + await dbConnect() + + const { searchParams } = new URL(request.url) + const reference = searchParams.get("reference") + const settlementId = searchParams.get("settlementId") + const userId = searchParams.get("userId") + const isStuck = searchParams.get("isStuck") + const limit = Math.max(1, Math.min(Number(searchParams.get("limit") || 50), 200)) + const page = Math.max(1, Number(searchParams.get("page") || 1)) + + const filter: Record = {} + if (reference) filter.providerReference = reference.trim() + if (settlementId) filter.settlementId = settlementId.trim() + if (userId) filter.userId = userId.trim() + if (isStuck === "true") filter.isStuck = true + if (isStuck === "false") filter.isStuck = false + + const total = await SettlementRecord.countDocuments(filter) + const records = await SettlementRecord.find(filter) + .sort({ updatedAt: -1 }) + .skip((page - 1) * limit) + .limit(limit) + .lean() + + const response = NextResponse.json({ + success: true, + total, + page, + limit, + settlements: records, + }) + + return finalizeAuthenticatedResponse(response, authContext) + } catch (error) { + console.error("SETTLEMENT_TIMELINE_GET_ERROR", error) + const message = error instanceof Error ? error.message : "Internal server error." + return NextResponse.json({ message }, { status: 500 }) + } +} + +export async function POST(request: Request) { + try { + const authContext = await requireAuthenticatedUser(request, ["admin"]) + if ("response" in authContext) return authContext.response + + const body = await parseJsonBody(request, postSchema) + if ("response" in body) return body.response + + await dbConnect() + + const { action, settlementId, providerReference, reason } = body.data + + if (action === "EVALUATE_TIMEOUTS") { + const summary = await evaluateFinalityTimeouts() + const response = NextResponse.json({ success: true, action, summary }) + return finalizeAuthenticatedResponse(response, authContext) + } + + if (!settlementId && !providerReference) { + return NextResponse.json( + { message: "Either settlementId or providerReference is required for this action." }, + { status: 400 }, + ) + } + + const defaultReason = `Admin operator action '${action}' executed by ${authContext.user._id}` + const finalReason = reason || defaultReason + + let targetState: any = "confirmed" + if (action === "FORCE_CONFIRM") targetState = "confirmed" + if (action === "POST_REVERSAL") targetState = "reversed" + if (action === "MARK_EXPIRED") targetState = "expired" + if (action === "RETRY_VERIFICATION") targetState = "provider-pending" + + const result = await transitionSettlementState({ + settlementId, + providerReference, + targetState, + triggeredBy: "operator", + reason: finalReason, + }) + + const response = NextResponse.json({ + success: true, + action, + settlement: result.settlement, + previousState: result.previousState, + }) + + return finalizeAuthenticatedResponse(response, authContext) + } catch (error) { + console.error("SETTLEMENT_TIMELINE_POST_ERROR", error) + const message = error instanceof Error ? error.message : "Internal server error." + return NextResponse.json({ message }, { status: 500 }) + } +} diff --git a/app/api/payments/webhook/route.ts b/app/api/payments/webhook/route.ts index 8fefef1c..eb9a1098 100644 --- a/app/api/payments/webhook/route.ts +++ b/app/api/payments/webhook/route.ts @@ -43,7 +43,17 @@ export async function POST(request: Request) { } const event = JSON.parse(body) - if (event.event !== "charge.success") { + const eventName = event.event as string + const supportedEvents = [ + "charge.success", + "charge.failed", + "refund.processed", + "transfer.reversed", + "charge.dispute.create", + "charge.dispute.resolve", + ] + + if (!supportedEvents.includes(eventName)) { return NextResponse.json({ status: "ignored" }, { status: 200 }) } diff --git a/lib/authorization/inventory.ts b/lib/authorization/inventory.ts index 4f79186b..d34f3877 100644 --- a/lib/authorization/inventory.ts +++ b/lib/authorization/inventory.ts @@ -9,16 +9,29 @@ export const ROUTE_POLICY_INVENTORY: Record = { "PATCH /api/account/profile": policy("account:update", "Principal profile only"), "GET /api/activity": policy("activity:read", "Principal activity only", "authenticated"), "PATCH /api/activity": policy("activity:update", "Principal activity only"), + "POST /api/admin/data-integrity/findings/[id]/suppress": policy("admin:settings:manage", "Administrative integrity finding suppression"), + "POST /api/admin/data-integrity/repair": policy("admin:settings:manage", "Administrative integrity repair"), + "GET /api/admin/data-integrity/scan": policy("admin:report", "Administrative integrity scan results"), + "POST /api/admin/data-integrity/scan": policy("admin:report", "Administrative integrity scan execution"), "GET /api/admin/dashboard-stats": policy("admin:report", "Administrative metrics"), "GET /api/admin/issues": policy("admin:issue:manage", "Administrative issue data"), "POST /api/admin/issues": policy("admin:issue:manage", "Administrative issue mutation"), "GET /api/admin/issues/[id]": policy("admin:issue:manage", "Administrative issue data"), "PATCH /api/admin/issues/[id]": policy("admin:issue:manage", "Administrative issue mutation"), + "GET /api/admin/kyc-documents": policy("kyc:review", "Sensitive KYC records"), + "POST /api/admin/kyc-documents": policy("kyc:review", "Sensitive KYC document creation"), + "POST /api/admin/kyc-migrate": policy("kyc:review", "Sensitive KYC migration"), "GET /api/admin/kyc-requests": policy("kyc:review", "Sensitive KYC records"), "POST /api/admin/migrate-vehicle-status": policy("vehicle:manage", "Privileged migration"), "GET /api/admin/platform-settings": policy("admin:settings:manage", "Administrative settings"), "PUT /api/admin/platform-settings": policy("admin:settings:manage", "Administrative settings mutation"), + "GET /api/admin/reconciliation/discrepancies": policy("admin:report", "Reconciliation discrepancies"), + "POST /api/admin/reconciliation/remediate": policy("admin:settings:manage", "Reconciliation remediation"), + "GET /api/admin/reconciliation/run": policy("admin:report", "Reconciliation runs"), + "POST /api/admin/reconciliation/run": policy("admin:report", "Reconciliation run execution"), "GET /api/admin/reports/export": policy("admin:report", "PII-bearing report export"), + "GET /api/admin/settlement/timeline": policy("admin:report", "Administrative settlement timeline & diagnostics"), + "POST /api/admin/settlement/timeline": policy("admin:settings:manage", "Administrative settlement state remediation"), "GET /api/admin/users/export": policy("admin:report", "PII-bearing user export"), "POST /api/auth/admin/signup": publicRoute("Bootstrap endpoint has its own one-time secret controls"), "GET /api/auth/admin/status": publicRoute("Returns only bootstrap availability"), @@ -30,10 +43,20 @@ export const ROUTE_POLICY_INVENTORY: Record = { "POST /api/auth/stellar/link": policy("wallet:adjust", "Links wallet identity"), "GET /api/driver/virtual-account": policy("wallet:read", "Driver-owned wallet"), "POST /api/driver/payments/initialize": policy("repayment:record", "Driver-owned active contract"), + "GET /api/fleet/documents": policy("vehicle:read", "Fleet documents"), + "POST /api/fleet/documents": policy("vehicle:manage", "Fleet document management"), + "GET /api/fleet/downtime": policy("vehicle:read", "Fleet downtime records"), + "POST /api/fleet/downtime": policy("vehicle:manage", "Fleet downtime management"), + "GET /api/fleet/inspections": policy("vehicle:read", "Fleet inspection records"), + "POST /api/fleet/inspections": policy("vehicle:manage", "Fleet inspection management"), + "GET /api/fleet/maintenance": policy("vehicle:read", "Fleet maintenance records"), + "POST /api/fleet/maintenance": policy("vehicle:manage", "Fleet maintenance creation"), + "PATCH /api/fleet/maintenance": policy("vehicle:manage", "Fleet maintenance mutation"), "POST /api/invest": publicRoute("Disabled legacy endpoint"), "GET /api/investments": policy("investment:read", "Investor ownership or admin"), "GET /api/investor/virtual-account": policy("wallet:read", "Investor-owned wallet"), "GET /api/kyc-documents": policy("kyc:document:read", "Owner or reviewer only"), + "POST /api/kyc-documents/sign": policy("kyc:document:read", "KYC document signing"), "GET /api/loans": policy("loan:read", "Driver ownership or admin"), "POST /api/loans": policy("loan:create", "KYC-approved driver principal"), "PUT /api/loans": policy("loan:approve", "Admin with valid workflow transition"), diff --git a/lib/services/paystack-processing.service.ts b/lib/services/paystack-processing.service.ts index 17a72413..43bd66e9 100644 --- a/lib/services/paystack-processing.service.ts +++ b/lib/services/paystack-processing.service.ts @@ -6,6 +6,7 @@ import ProcessedGatewayEvent, { type GatewayPaymentType } from "@/models/Process import Transaction from "@/models/Transaction" import User from "@/models/User" import { logAuditEvent } from "@/lib/security/audit-log" +import { initiateSettlement, transitionSettlementState } from "@/lib/settlement/settlement-service" type ProcessedVia = "verify" | "webhook" diff --git a/lib/settlement/config.ts b/lib/settlement/config.ts new file mode 100644 index 00000000..5c501b53 --- /dev/null +++ b/lib/settlement/config.ts @@ -0,0 +1,77 @@ +import { SettlementRail } from "@/models/SettlementRecord" + +export interface RailSettlementConfig { + rail: SettlementRail + finalityThreshold: number + pendingTimeoutMs: number + observedTimeoutMs: number + autoExpireOnTimeout: boolean +} + +const DEFAULT_CONFIGS: Record> = { + production: { + paystack: { + rail: "paystack", + finalityThreshold: 1, + pendingTimeoutMs: 15 * 60 * 1000, // 15 minutes + observedTimeoutMs: 30 * 60 * 1000, + autoExpireOnTimeout: false, + }, + stellar: { + rail: "stellar", + finalityThreshold: 3, // 3 ledger confirmations + pendingTimeoutMs: 5 * 60 * 1000, // 5 minutes + observedTimeoutMs: 10 * 60 * 1000, + autoExpireOnTimeout: false, + }, + bank_transfer: { + rail: "bank_transfer", + finalityThreshold: 1, + pendingTimeoutMs: 24 * 60 * 60 * 1000, // 24 hours + observedTimeoutMs: 48 * 60 * 60 * 1000, + autoExpireOnTimeout: false, + }, + internal_ledger: { + rail: "internal_ledger", + finalityThreshold: 1, + pendingTimeoutMs: 5 * 60 * 1000, + observedTimeoutMs: 5 * 60 * 1000, + autoExpireOnTimeout: true, + }, + }, + development: { + paystack: { + rail: "paystack", + finalityThreshold: 1, + pendingTimeoutMs: 5 * 60 * 1000, + observedTimeoutMs: 10 * 60 * 1000, + autoExpireOnTimeout: true, + }, + stellar: { + rail: "stellar", + finalityThreshold: 1, + pendingTimeoutMs: 2 * 60 * 1000, + observedTimeoutMs: 5 * 60 * 1000, + autoExpireOnTimeout: true, + }, + bank_transfer: { + rail: "bank_transfer", + finalityThreshold: 1, + pendingTimeoutMs: 60 * 60 * 1000, + observedTimeoutMs: 120 * 60 * 1000, + autoExpireOnTimeout: false, + }, + internal_ledger: { + rail: "internal_ledger", + finalityThreshold: 1, + pendingTimeoutMs: 60 * 1000, + observedTimeoutMs: 60 * 1000, + autoExpireOnTimeout: true, + }, + }, +} + +export function getRailSettlementConfig(rail: SettlementRail, env = process.env.NODE_ENV || "development"): RailSettlementConfig { + const envKey = env === "production" ? "production" : "development" + return DEFAULT_CONFIGS[envKey][rail] || DEFAULT_CONFIGS.development[rail] +} diff --git a/lib/settlement/settlement-service.ts b/lib/settlement/settlement-service.ts new file mode 100644 index 00000000..508b2889 --- /dev/null +++ b/lib/settlement/settlement-service.ts @@ -0,0 +1,392 @@ +import mongoose, { type ClientSession } from "mongoose" + +import dbConnect from "@/lib/dbConnect" +import AuditLog from "@/models/AuditLog" +import SettlementRecord, { + CanonicalSettlementState, + ISettlementRecord, + OperatorTimelineEntry, + SettlementRail, +} from "@/models/SettlementRecord" +import Transaction from "@/models/Transaction" +import User from "@/models/User" +import { getRailSettlementConfig } from "./config" +import { determineSafeActions, isValidTransition } from "./state-machine" + +export interface InitiateSettlementInput { + rail: SettlementRail + providerReference: string + userId: string | mongoose.Types.ObjectId + userType: "driver" | "investor" | "admin" + paymentType: "wallet_funding" | "down_payment" | "driver_repayment" | "pool_investment" | "payout" + amount: number + currency?: string + stellarHash?: string + ledgerJournalId?: string + poolInvestmentId?: string + driverPaymentId?: string + userTransactionId?: string + initialState?: CanonicalSettlementState + triggeredBy?: "webhook" | "verifier" | "indexer" | "operator" | "system" + reason?: string + metadata?: Record +} + +export interface TransitionStateInput { + settlementId?: string + providerReference?: string + targetState: CanonicalSettlementState + triggeredBy: "webhook" | "verifier" | "indexer" | "operator" | "system" + reason: string + stellarHash?: string + ledgerJournalId?: string + poolInvestmentId?: string + userTransactionId?: string + driverPaymentId?: string + confirmationsCount?: number + metadata?: Record + session?: ClientSession +} + +function generateSettlementId(): string { + return `STL-${Date.now()}-${Math.random().toString(36).substring(2, 8)}` +} + +async function runInSession( + existingSession: ClientSession | undefined, + fn: (session: ClientSession) => Promise, +): Promise { + if (existingSession) { + return fn(existingSession) + } + const session = await mongoose.startSession() + session.startTransaction() + try { + const res = await fn(session) + await session.commitTransaction() + return res + } catch (err) { + await session.abortTransaction().catch(() => undefined) + throw err + } finally { + session.endSession() + } +} + +export async function initiateSettlement(input: InitiateSettlementInput): Promise<{ + settlement: ISettlementRecord + alreadyExists: boolean +}> { + await dbConnect() + + const normalizedRef = input.providerReference.trim() + if (!normalizedRef) { + throw new Error("Provider reference is required to initiate settlement.") + } + + const existing = await SettlementRecord.findOne({ + providerReference: normalizedRef, + rail: input.rail, + }) + + if (existing) { + return { settlement: existing, alreadyExists: true } + } + + const initialState = input.initialState || "initiated" + const env = process.env.NODE_ENV || "development" + const railConfig = getRailSettlementConfig(input.rail, env) + const settlementId = generateSettlementId() + + const safeActions = determineSafeActions(initialState, false) + + const initialTimelineEntry: OperatorTimelineEntry = { + fromState: null, + toState: initialState, + triggeredBy: input.triggeredBy || "system", + reason: input.reason || "Settlement record initiated", + safeActions, + metadata: input.metadata, + timestamp: new Date(), + } + + const session = await mongoose.startSession() + session.startTransaction() + + try { + const [settlement] = await SettlementRecord.create( + [ + { + settlementId, + rail: input.rail, + environment: env, + currentState: initialState, + providerReference: normalizedRef, + stellarHash: input.stellarHash, + ledgerJournalId: input.ledgerJournalId, + poolInvestmentId: input.poolInvestmentId, + driverPaymentId: input.driverPaymentId, + userTransactionId: input.userTransactionId, + userId: input.userId, + userType: input.userType, + paymentType: input.paymentType, + amount: input.amount, + currency: input.currency || "NGN", + finalityThreshold: railConfig.finalityThreshold, + timeline: [initialTimelineEntry], + isStuck: false, + }, + ], + { session }, + ) + + if (["initiated", "provider-pending", "observed", "provisionally_credited"].includes(initialState)) { + if (input.paymentType === "wallet_funding") { + await User.findByIdAndUpdate( + input.userId, + { $inc: { pendingBalance: input.amount } }, + { session }, + ) + } + } + + await session.commitTransaction() + return { settlement, alreadyExists: false } + } catch (err: any) { + await session.abortTransaction().catch(() => undefined) + if (err.code === 11000) { + const found = await SettlementRecord.findOne({ + providerReference: normalizedRef, + rail: input.rail, + }) + if (found) return { settlement: found, alreadyExists: true } + } + throw err + } finally { + session.endSession() + } +} + +export async function transitionSettlementState( + input: TransitionStateInput, +): Promise<{ settlement: ISettlementRecord; previousState: CanonicalSettlementState }> { + await dbConnect() + + return runInSession(input.session, async (session) => { + let settlement: ISettlementRecord | null = null + + if (input.settlementId) { + settlement = await SettlementRecord.findOne({ settlementId: input.settlementId }).session(session) + } else if (input.providerReference) { + settlement = await SettlementRecord.findOne({ providerReference: input.providerReference.trim() }).session(session) + } + + if (!settlement) { + throw new Error(`Settlement record not found for reference: ${input.settlementId || input.providerReference}`) + } + + const previousState = settlement.currentState + const targetState = input.targetState + + if (previousState === targetState) { + return { settlement, previousState } + } + + if (!isValidTransition(previousState, targetState)) { + throw new Error(`Invalid settlement state transition from '${previousState}' to '${targetState}'.`) + } + + if (input.stellarHash) settlement.stellarHash = input.stellarHash + if (input.ledgerJournalId) settlement.ledgerJournalId = input.ledgerJournalId + if (input.poolInvestmentId) settlement.poolInvestmentId = input.poolInvestmentId + if (input.userTransactionId) settlement.userTransactionId = input.userTransactionId + if (input.driverPaymentId) settlement.driverPaymentId = input.driverPaymentId + if (input.confirmationsCount !== undefined) { + settlement.confirmationsCount = input.confirmationsCount + } + + const safeActions = determineSafeActions(targetState, false) + const timelineEntry: OperatorTimelineEntry = { + fromState: previousState, + toState: targetState, + triggeredBy: input.triggeredBy, + reason: input.reason, + safeActions, + metadata: input.metadata, + timestamp: new Date(), + } + + settlement.currentState = targetState + settlement.timeline.push(timelineEntry) + settlement.isStuck = false + settlement.stuckReason = undefined + + // Apply balance bucket movements based on state transition + const user = await User.findById(settlement.userId).session(session) + const amount = settlement.amount + + if (user && settlement.paymentType === "wallet_funding") { + const wasPending = ["initiated", "provider-pending", "observed", "provisionally_credited"].includes(previousState) + + if (targetState === "confirmed") { + if (wasPending) { + const currentPending = Number(user.pendingBalance || 0) + const pendingDeduction = Math.min(currentPending, amount) + user.pendingBalance = Math.max(currentPending - pendingDeduction, 0) + } + user.availableBalance = Number(user.availableBalance || 0) + amount + } else if (targetState === "reversed") { + await executeReversalJournal(user, settlement, input, session) + } else if (targetState === "disputed") { + const currentAvailable = Number(user.availableBalance || 0) + const holdDeduction = Math.min(currentAvailable, amount) + user.availableBalance = Math.max(currentAvailable - holdDeduction, 0) + user.heldBalance = Number(user.heldBalance || 0) + holdDeduction + } else if (targetState === "failed" || targetState === "expired") { + if (wasPending) { + const currentPending = Number(user.pendingBalance || 0) + user.pendingBalance = Math.max(currentPending - Math.min(currentPending, amount), 0) + } + } + + await user.save({ session }) + } else if (user && targetState === "reversed") { + await executeReversalJournal(user, settlement, input, session) + await user.save({ session }) + } + + await settlement.save({ session }) + + return { settlement, previousState } + }) +} + +async function executeReversalJournal( + user: any, + settlement: ISettlementRecord, + input: TransitionStateInput, + session: ClientSession, +) { + const amount = settlement.amount + const currentAvailable = Number(user.availableBalance || 0) + + let deductedFromAvailable = 0 + let deductedFromHeld = 0 + let deductedFromPending = 0 + + if (currentAvailable >= amount) { + deductedFromAvailable = amount + } else { + deductedFromAvailable = currentAvailable + const remainder = amount - currentAvailable + const currentHeld = Number(user.heldBalance || 0) + if (currentHeld >= remainder) { + deductedFromHeld = remainder + } else { + deductedFromHeld = currentHeld + const secondRemainder = remainder - currentHeld + const currentPending = Number(user.pendingBalance || 0) + deductedFromPending = Math.min(currentPending, secondRemainder) + } + } + + user.availableBalance = Math.max(currentAvailable - deductedFromAvailable, 0) + user.heldBalance = Math.max(Number(user.heldBalance || 0) - deductedFromHeld, 0) + user.pendingBalance = Math.max(Number(user.pendingBalance || 0) - deductedFromPending, 0) + user.reversedBalance = Number(user.reversedBalance || 0) + amount + + const reversalTx = await Transaction.create( + [ + { + userId: user._id, + userType: settlement.userType, + type: "wallet_debit", + amount, + currency: settlement.currency, + method: settlement.rail === "paystack" ? "paystack" : "system", + status: "Completed", + gatewayReference: `REV-${settlement.providerReference}`, + description: `Settlement reversal journal for ${settlement.providerReference}`, + relatedId: settlement.userTransactionId || settlement.settlementId, + metadata: { + reversalReason: input.reason, + settlementId: settlement.settlementId, + originalReference: settlement.providerReference, + deductedFromAvailable, + deductedFromHeld, + deductedFromPending, + }, + }, + ], + { session }, + ) + + settlement.ledgerJournalId = reversalTx[0]._id.toString() + + await AuditLog.create( + [ + { + userId: user._id, + action: "SETTLEMENT_REVERSAL_POSTED", + targetModel: "SettlementRecord", + targetId: settlement.settlementId, + details: { + providerReference: settlement.providerReference, + amount, + reversalJournalId: reversalTx[0]._id.toString(), + reason: input.reason, + deductedFromAvailable, + deductedFromHeld, + deductedFromPending, + }, + timestamp: new Date(), + }, + ], + { session }, + ) +} + +export async function evaluateFinalityTimeouts(): Promise<{ + evaluatedCount: number + stuckCount: number + expiredCount: number +}> { + await dbConnect() + const now = new Date() + const activeSettlements = await SettlementRecord.find({ + currentState: { $in: ["initiated", "provider-pending", "observed", "provisionally_credited"] }, + }) + + let stuckCount = 0 + let expiredCount = 0 + + for (const s of activeSettlements) { + const config = getRailSettlementConfig(s.rail, s.environment) + const ageMs = now.getTime() - new Date(s.createdAt).getTime() + const thresholdMs = + s.currentState === "observed" || s.currentState === "provisionally_credited" + ? config.observedTimeoutMs + : config.pendingTimeoutMs + + if (ageMs > thresholdMs) { + if (config.autoExpireOnTimeout) { + await transitionSettlementState({ + settlementId: s.settlementId, + targetState: "expired", + triggeredBy: "system", + reason: `Settlement automatically expired after exceeding timeout threshold (${Math.round(thresholdMs / 1000)}s)`, + }) + expiredCount++ + } else { + s.isStuck = true + s.stuckReason = `Settlement pending for ${Math.round(ageMs / 60000)} mins exceeding limit of ${Math.round(thresholdMs / 60000)} mins.` + s.actionableAlertSent = true + s.lastEvaluatedAt = now + await s.save() + stuckCount++ + } + } + } + + return { evaluatedCount: activeSettlements.length, stuckCount, expiredCount } +} diff --git a/lib/settlement/state-machine.ts b/lib/settlement/state-machine.ts new file mode 100644 index 00000000..f082a8b2 --- /dev/null +++ b/lib/settlement/state-machine.ts @@ -0,0 +1,45 @@ +import { CanonicalSettlementState } from "@/models/SettlementRecord" + +export const VALID_TRANSITIONS: Record = { + initiated: ["provider-pending", "observed", "provisionally_credited", "confirmed", "failed", "expired"], + "provider-pending": ["observed", "provisionally_credited", "confirmed", "failed", "expired"], + observed: ["provisionally_credited", "confirmed", "failed", "expired"], + provisionally_credited: ["confirmed", "reversed", "failed"], + confirmed: ["reversed", "disputed"], + disputed: ["confirmed", "reversed"], + reversed: [], + failed: [], + expired: [], +} + +export function isValidTransition( + fromState: CanonicalSettlementState | null, + toState: CanonicalSettlementState, +): boolean { + if (!fromState) return true + if (fromState === toState) return true + const allowed = VALID_TRANSITIONS[fromState] || [] + return allowed.includes(toState) +} + +export function determineSafeActions(state: CanonicalSettlementState, isStuck = false): string[] { + switch (state) { + case "initiated": + case "provider-pending": + return isStuck ? ["RETRY_VERIFICATION", "MARK_EXPIRED", "FORCE_CONFIRM"] : ["RETRY_VERIFICATION"] + case "observed": + case "provisionally_credited": + return isStuck ? ["RETRY_VERIFICATION", "FORCE_CONFIRM", "POST_REVERSAL"] : ["RETRY_VERIFICATION"] + case "confirmed": + return ["POST_REVERSAL", "FLAG_DISPUTE"] + case "disputed": + return ["RESOLVE_DISPUTE_CONFIRM", "RESOLVE_DISPUTE_REVERSE"] + case "failed": + case "expired": + return ["REINITIALIZE_SETTLEMENT"] + case "reversed": + return [] + default: + return [] + } +} diff --git a/lib/stellar/indexer.ts b/lib/stellar/indexer.ts index d5532377..d7edafd2 100644 --- a/lib/stellar/indexer.ts +++ b/lib/stellar/indexer.ts @@ -186,6 +186,29 @@ async function persistEvent(op: RawStellarOperation): Promise { try { await StellarIndexedEvent.create(doc) + + const txHash = op.transaction_hash || op.id + const { default: SettlementRecord } = await import("@/models/SettlementRecord") + const { transitionSettlementState } = await import("@/lib/settlement/settlement-service") + + const matchingSettlement = await SettlementRecord.findOne({ + $or: [{ providerReference: txHash }, { stellarHash: txHash }, { providerReference: op.id }], + rail: "stellar", + }) + + if (matchingSettlement && matchingSettlement.currentState !== "confirmed") { + const newConfirmations = (matchingSettlement.confirmationsCount || 0) + 1 + const isConfirmed = newConfirmations >= matchingSettlement.finalityThreshold + await transitionSettlementState({ + settlementId: matchingSettlement.settlementId, + targetState: isConfirmed ? "confirmed" : "observed", + triggeredBy: "indexer", + reason: `Indexed Stellar operation ${op.id} (ledger: ${op.ledger_attr || "unknown"})`, + stellarHash: txHash, + confirmationsCount: newConfirmations, + }) + } + return { inserted: true, duplicate: false } } catch (err: unknown) { // MongoDB duplicate key error code 11000 means this event was already diff --git a/models/SettlementRecord.ts b/models/SettlementRecord.ts new file mode 100644 index 00000000..5de480f4 --- /dev/null +++ b/models/SettlementRecord.ts @@ -0,0 +1,131 @@ +import mongoose, { Document, Schema } from "mongoose" + +export type CanonicalSettlementState = + | "initiated" + | "provider-pending" + | "observed" + | "provisionally_credited" + | "confirmed" + | "reversed" + | "disputed" + | "failed" + | "expired" + +export type SettlementRail = "paystack" | "stellar" | "internal_ledger" | "bank_transfer" + +export interface OperatorTimelineEntry { + fromState: CanonicalSettlementState | null + toState: CanonicalSettlementState + triggeredBy: "webhook" | "verifier" | "indexer" | "operator" | "system" + reason: string + safeActions: string[] + metadata?: Record + timestamp: Date +} + +export interface ISettlementRecord extends Document { + settlementId: string + rail: SettlementRail + environment: string + currentState: CanonicalSettlementState + providerReference: string + stellarHash?: string + ledgerJournalId?: string + poolInvestmentId?: string + userTransactionId?: string + driverPaymentId?: string + userId: Schema.Types.ObjectId + userType: "driver" | "investor" | "admin" + paymentType: "wallet_funding" | "down_payment" | "driver_repayment" | "pool_investment" | "payout" + amount: number + currency: string + finalityThreshold: number + confirmationsCount: number + timeline: OperatorTimelineEntry[] + isStuck: boolean + stuckReason?: string + actionableAlertSent: boolean + lastEvaluatedAt?: Date + createdAt: Date + updatedAt: Date +} + +const OperatorTimelineEntrySchema = new Schema( + { + fromState: { type: String, default: null }, + toState: { type: String, required: true }, + triggeredBy: { + type: String, + enum: ["webhook", "verifier", "indexer", "operator", "system"], + required: true, + }, + reason: { type: String, required: true }, + safeActions: { type: [String], default: [] }, + metadata: { type: Schema.Types.Mixed }, + timestamp: { type: Date, default: Date.now }, + }, + { _id: false }, +) + +const SettlementRecordSchema = new Schema( + { + settlementId: { type: String, required: true, unique: true, index: true }, + rail: { + type: String, + enum: ["paystack", "stellar", "internal_ledger", "bank_transfer"], + required: true, + index: true, + }, + environment: { type: String, default: "development" }, + currentState: { + type: String, + enum: [ + "initiated", + "provider-pending", + "observed", + "provisionally_credited", + "confirmed", + "reversed", + "disputed", + "failed", + "expired", + ], + required: true, + default: "initiated", + index: true, + }, + providerReference: { type: String, required: true, index: true }, + stellarHash: { type: String, index: true, sparse: true }, + ledgerJournalId: { type: String, index: true, sparse: true }, + poolInvestmentId: { type: String, index: true, sparse: true }, + userTransactionId: { type: String, index: true, sparse: true }, + driverPaymentId: { type: String, index: true, sparse: true }, + userId: { type: Schema.Types.ObjectId, ref: "User", required: true, index: true }, + userType: { + type: String, + enum: ["driver", "investor", "admin"], + required: true, + }, + paymentType: { + type: String, + enum: ["wallet_funding", "down_payment", "driver_repayment", "pool_investment", "payout"], + required: true, + }, + amount: { type: Number, required: true, min: 0 }, + currency: { type: String, default: "NGN" }, + finalityThreshold: { type: Number, default: 1 }, + confirmationsCount: { type: Number, default: 0 }, + timeline: { type: [OperatorTimelineEntrySchema], default: [] }, + isStuck: { type: Boolean, default: false, index: true }, + stuckReason: { type: String }, + actionableAlertSent: { type: Boolean, default: false }, + lastEvaluatedAt: { type: Date }, + }, + { timestamps: true }, +) + +SettlementRecordSchema.index({ providerReference: 1, rail: 1 }) +SettlementRecordSchema.index({ userId: 1, currentState: 1 }) + +export default (mongoose.models.SettlementRecord || + mongoose.model("SettlementRecord", SettlementRecordSchema)) as mongoose.Model diff --git a/models/User.ts b/models/User.ts index 80d6df15..28cf56e4 100644 --- a/models/User.ts +++ b/models/User.ts @@ -109,6 +109,21 @@ const UserSchema = new mongoose.Schema( default: 0, min: 0, }, + pendingBalance: { + type: Number, + default: 0, + min: 0, + }, + heldBalance: { + type: Number, + default: 0, + min: 0, + }, + reversedBalance: { + type: Number, + default: 0, + min: 0, + }, totalInvested: { type: Number, default: 0,