diff --git a/src/events/transaction.events.ts b/src/events/transaction.events.ts new file mode 100644 index 00000000..125071e3 --- /dev/null +++ b/src/events/transaction.events.ts @@ -0,0 +1,24 @@ +import { DomainEventEnvelope } from './domain-event.types'; +import { DomainEventNames, DomainEventName } from './event-names'; + +export interface TransactionCreatedPayload { + transactionId: string; + amount?: string; + assetCode?: string; + recipientAddress?: string; + agentId?: string; +} + +export class TransactionCreatedEvent implements DomainEventEnvelope { + constructor( + public readonly id: string, + public readonly name: DomainEventNames[DomainEventName.TransactionCreated], + public readonly timestamp: number, + public readonly occurredAt: Date, + public readonly aggregateType: string, + public readonly aggregateId: string, + public readonly payload: TransactionCreatedPayload, + public readonly organizationId?: string, + public readonly actorId?: string, + ) {} +} diff --git a/src/modules/risk/risk.listener.ts b/src/modules/risk/risk.listener.ts new file mode 100644 index 00000000..dee899d4 --- /dev/null +++ b/src/modules/risk/risk.listener.ts @@ -0,0 +1,57 @@ +import { Injectable, Logger } from '@nestjs/common'; +import { TypedOnEvent } from '../../events/typed-event-listener.decorator'; +import { DomainEventEnvelope } from '../../events/domain-event.types'; +import { DomainEventName } from '../../events/event-names'; +import { RiskService } from './risk.service'; +import { PrismaService } from '../../database/prisma.service'; + +/** + * Listens to transaction domain events and triggers asynchronous risk analysis + * and score persistence. + */ +@Injectable() +export class RiskListener { + private readonly logger = new Logger(RiskListener.name); + + constructor( + private readonly riskService: RiskService, + private readonly prisma: PrismaService, + ) {} + + @TypedOnEvent(DomainEventName.TransactionCreated) + async handleTransactionCreated(envelope: DomainEventEnvelope<{ transactionId: string; amount?: string; assetCode?: string; recipientAddress?: string; agentId?: string }>): Promise { + const payload = envelope.payload; + if (!payload || !payload.transactionId) { + return; + } + + try { + const assessment = await this.riskService.evaluate( + envelope.organizationId, + { + amount: payload.amount ? parseFloat(payload.amount) : 100, + recipientAddress: payload.recipientAddress, + assetCode: payload.assetCode, + agentId: payload.agentId, + }, + { + transactionId: payload.transactionId, + actorId: envelope.actorId, + }, + ); + + // Persist risk assessment score/band back to the transaction record or audit/risk metadata if applicable + await this.prisma.transaction.updateMany({ + where: { id: payload.transactionId }, + data: { + riskScore: assessment.score, + riskBand: assessment.band, + }, + }); + } catch (error) { + this.logger.error( + `Failed to evaluate risk for transaction '${payload.transactionId}': ${(error as Error).message}`, + ); + } + } +} diff --git a/src/modules/risk/risk.module.ts b/src/modules/risk/risk.module.ts index bdd401a1..ca5aa0e4 100644 --- a/src/modules/risk/risk.module.ts +++ b/src/modules/risk/risk.module.ts @@ -2,8 +2,11 @@ import { Module } from '@nestjs/common'; import { RiskController } from './risk.controller'; import { RiskService } from './risk.service'; import { RiskEngine } from './risk.engine'; +import { EventsModule } from '../../events/events.module'; +import { DatabaseModule } from '../../database/database.module'; @Module({ + imports: [EventsModule, DatabaseModule], controllers: [RiskController], providers: [RiskService, RiskEngine], exports: [RiskService, RiskEngine], diff --git a/src/modules/risk/risk.service.ts b/src/modules/risk/risk.service.ts index 42dd8324..7138b01d 100644 --- a/src/modules/risk/risk.service.ts +++ b/src/modules/risk/risk.service.ts @@ -9,6 +9,8 @@ import { DomainEventName } from '../../events/event-names'; * RiskEvaluated domain event (with full factor breakdown for audit metadata), * and is called by the transactions pipeline. */ +import { OnEvent } from '@nestjs/event-emitter'; + @Injectable() export class RiskService { constructor( @@ -16,6 +18,26 @@ export class RiskService { private readonly eventBus: EventBusService, ) {} + @OnEvent(DomainEventName.TransactionCreated, { async: true }) + async handleTransactionCreated(event: { payload?: { transactionId: string; amount?: string; assetCode?: string; recipientAddress?: string; agentId?: string }; organizationId?: string; transactionId?: string; amount?: string; assetCode?: string; recipientAddress?: string }): Promise { + const payload = event.payload || event; + const orgId = event.organizationId || 'default-org'; + const txId = payload.transactionId || event.transactionId || ''; + const amtNum = payload.amount ? parseFloat(payload.amount) : 0; + await this.evaluate( + orgId, + { + amount: isNaN(amtNum) ? 0 : amtNum, + sourceAccount: payload.agentId ?? 'unknown', + destinationAccount: payload.recipientAddress ?? 'unknown', + transactionType: 'transfer', + }, + { + transactionId: txId, + }, + ); + } + /** * Full evaluation with event emission. The emitted event payload includes * the complete factor breakdown so the audit listener captures it as metadata. diff --git a/src/modules/risk/tests/risk.listener.spec.ts b/src/modules/risk/tests/risk.listener.spec.ts new file mode 100644 index 00000000..dd4c2a8a --- /dev/null +++ b/src/modules/risk/tests/risk.listener.spec.ts @@ -0,0 +1,63 @@ +import { describe, it, expect, beforeEach, vi } from 'vitest'; +import { RiskListener } from '../risk.listener'; +import { RiskService } from '../risk.service'; +import { PrismaService } from '../../../database/prisma.service'; +import { DomainEventName } from '../../../events/event-names'; + +describe('RiskListener', () => { + let listener: RiskListener; + let riskService: RiskService; + let prisma: PrismaService; + + beforeEach(() => { + riskService = { + evaluate: vi.fn().mockResolvedValue({ + score: 15, + band: 'LOW', + factors: [], + canAutoExecute: true, + }), + } as unknown as RiskService; + + prisma = { + transaction: { + updateMany: vi.fn().mockResolvedValue({ count: 1 }), + }, + } as unknown as PrismaService; + + listener = new RiskListener(riskService, prisma); + }); + + it('should evaluate risk and persist score upon transaction.created event', async () => { + const envelope = { + id: 'evt-1', + name: DomainEventName.TransactionCreated, + organizationId: 'org-1', + actorId: 'agent-1', + timestamp: Date.now(), + aggregateType: 'transaction', + aggregateId: 'tx-1', + payload: { + transactionId: 'tx-1', + amount: '500', + assetCode: 'USDC', + recipientAddress: 'GABC...', + }, + }; + + await listener.handleTransactionCreated(envelope as never); + + expect(riskService.evaluate).toHaveBeenCalledWith( + 'org-1', + expect.objectContaining({ amount: 500, assetCode: 'USDC' }), + expect.objectContaining({ transactionId: 'tx-1', actorId: 'agent-1' }), + ); + + expect(prisma.transaction.updateMany).toHaveBeenCalledWith( + expect.objectContaining({ + where: { id: 'tx-1' }, + data: { riskScore: 15, riskBand: 'LOW' }, + }), + ); + }); +}); diff --git a/src/modules/risk/tests/risk.service.event.spec.ts b/src/modules/risk/tests/risk.service.event.spec.ts new file mode 100644 index 00000000..ad4ee79f --- /dev/null +++ b/src/modules/risk/tests/risk.service.event.spec.ts @@ -0,0 +1,63 @@ +import { describe, it, expect, beforeEach, vi } from 'vitest'; +import { RiskService } from '../risk.service'; +import { RiskEngine } from '../risk.engine'; +import { EventBusService } from '../../../events/event-bus.service'; +import { DomainEventName } from '../../../events/event-names'; +import { PrismaService } from '../../../database/prisma.service'; + +describe('RiskService Event Handling', () => { + let riskService: RiskService; + let riskEngine: RiskEngine; + let eventBus: EventBusService; + let prisma: PrismaService; + + beforeEach(() => { + riskEngine = { + assess: vi.fn().mockReturnValue({ + score: 15, + band: 'LOW', + factors: [], + canAutoExecute: true, + }), + } as unknown as RiskEngine; + + eventBus = { + emit: vi.fn().mockResolvedValue(undefined), + } as unknown as EventBusService; + + prisma = { + transaction: { + updateMany: vi.fn().mockResolvedValue({ count: 1 }), + }, + } as unknown as PrismaService; + + riskService = new RiskService(riskEngine, eventBus); + }); + + it('should evaluate risk and emit event upon transaction.created event', async () => { + const envelope = { + id: 'evt-1', + name: DomainEventName.TransactionCreated, + organizationId: 'org-1', + actorId: 'agent-1', + timestamp: Date.now(), + aggregateType: 'transaction', + aggregateId: 'tx-1', + payload: { + transactionId: 'tx-1', + amount: '500', + assetCode: 'USDC', + recipientAddress: 'GABC...', + }, + }; + + await riskService.handleTransactionCreated(envelope as never); + + expect(riskEngine.assess).toHaveBeenCalled(); + expect(eventBus.emit).toHaveBeenCalledWith( + DomainEventName.RiskEvaluated, + expect.any(Object), + expect.any(Object) + ); + }); +});