diff --git a/src/events/domain-event.types.ts b/src/events/domain-event.types.ts index 6ad1c162..6db46445 100644 --- a/src/events/domain-event.types.ts +++ b/src/events/domain-event.types.ts @@ -210,6 +210,17 @@ export interface ProposalExecutedPayload extends Record { transactionId: string; } +export interface TransactionInitiatedPayload extends Record { + transactionId: string; + organizationId: string; + walletId?: string; + agentId?: string; + amount?: string; + asset?: string; + recipientAddress?: string; + memo?: string; +} + export interface TransactionCreatedPayload extends Record { transactionId: string; walletId?: string; diff --git a/src/events/event-names.ts b/src/events/event-names.ts index 6837244d..be523d7f 100644 --- a/src/events/event-names.ts +++ b/src/events/event-names.ts @@ -49,6 +49,7 @@ export const DomainEventName = { ProposalExecuted: 'proposal.executed', // Transaction / blockchain + TransactionInitiated: 'transaction.initiated', TransactionCreated: 'transaction.created', TransactionSubmitted: 'transaction.submitted', TransactionConfirmed: 'transaction.confirmed', @@ -78,6 +79,7 @@ export const WEBHOOK_EVENTS: DomainEventNameType[] = [ DomainEventName.WalletUpdated, DomainEventName.ProposalApproved, DomainEventName.ProposalRejected, + DomainEventName.TransactionInitiated, DomainEventName.TransactionCompleted, DomainEventName.TransactionFailed, DomainEventName.PolicyViolated, diff --git a/src/events/typed-event-emitter.service.ts b/src/events/typed-event-emitter.service.ts index 5a9155a0..8629956c 100644 --- a/src/events/typed-event-emitter.service.ts +++ b/src/events/typed-event-emitter.service.ts @@ -53,6 +53,7 @@ export interface DomainEventMap { 'proposal.executed': PayloadTypes.ProposalExecutedPayload; // Transaction + 'transaction.initiated': PayloadTypes.TransactionInitiatedPayload; 'transaction.created': PayloadTypes.TransactionCreatedPayload; 'transaction.submitted': PayloadTypes.TransactionSubmittedPayload; 'transaction.completed': PayloadTypes.TransactionCompletedPayload; diff --git a/src/modules/risk/risk.module.ts b/src/modules/risk/risk.module.ts index bdd401a1..a745e276 100644 --- a/src/modules/risk/risk.module.ts +++ b/src/modules/risk/risk.module.ts @@ -1,9 +1,12 @@ import { Module } from '@nestjs/common'; -import { RiskController } from './risk.controller'; import { RiskService } from './risk.service'; import { RiskEngine } from './risk.engine'; +import { RiskController } from './risk.controller'; +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.spec.ts b/src/modules/risk/risk.service.spec.ts index b8c3e9f0..6d2fcd12 100644 --- a/src/modules/risk/risk.service.spec.ts +++ b/src/modules/risk/risk.service.spec.ts @@ -1,9 +1,92 @@ -import { describe, expect, it, vi } from 'vitest'; -import { RiskBand } from '@prisma/client'; +import { Test, TestingModule } from '@nestjs/testing'; import { RiskService } from './risk.service'; import { RiskEngine } from './risk.engine'; -import { RiskFactorsInput } from './risk.types'; import { EventBusService } from '../../events/event-bus.service'; +import { PrismaService } from '../../database/prisma.service'; +import { DomainEventEnvelope, TransactionInitiatedPayload } from '../../events/domain-event.types'; +import { DomainEventName } from '../../events/event-names'; +import { describe, expect, it, vi, beforeEach } from 'vitest'; +import { RiskFactorsInput } from './risk.types'; + +describe('RiskService event handling', () => { + let service: RiskService; + let eventBus: EventBusService; + let prisma: PrismaService; + + const mockEventBus = { + emit: vi.fn().mockResolvedValue(true), +}; + + const mockPrisma = { + transaction: { + update: vi.fn().mockResolvedValue({ id: 'tx-123' }), + }, + }; + + beforeEach(async () => { + const module: TestingModule = await Test.createTestingModule({ + providers: [ + RiskService, + RiskEngine, + { provide: EventBusService, useValue: mockEventBus }, + { provide: PrismaService, useValue: mockPrisma }, + ], + }).compile(); + + service = module.get(RiskService); + eventBus = module.get(EventBusService); + prisma = module.get(PrismaService); + vi.clearAllMocks(); + }); + + it('should handle transaction.initiated event, evaluate risk, persist scores, and emit risk.evaluated', async () => { + const envelope: DomainEventEnvelope = { + name: DomainEventName.TransactionInitiated, + organizationId: 'org-1', + aggregateType: 'transaction', + aggregateId: 'tx-123', + actorId: 'user-1', + correlationId: 'corr-1', + occurredAt: new Date(), + payload: { + transactionId: 'tx-123', + organizationId: 'org-1', + amount: '100.00', + asset: 'XLM', + }, + }; + + const assessment = await service.handleTransactionInitiated(envelope); + + expect(assessment).toBeDefined(); + expect(prisma.transaction.update).toHaveBeenCalledWith({ + where: { id: 'tx-123' }, + data: { + riskScore: assessment.score, + riskBand: assessment.band, + }, + }); + expect(eventBus.emit).toHaveBeenCalledWith( + DomainEventName.RiskEvaluated, + expect.objectContaining({ + transactionId: 'tx-123', + score: assessment.score, + band: assessment.band, + }), + expect.objectContaining({ + organizationId: 'org-1', + actorId: 'user-1', + aggregateType: 'transaction', + aggregateId: 'tx-123', + correlationId: 'corr-1', + }), + ); + }); +}); + + + + const lowRisk: RiskFactorsInput = { amount: 20, @@ -18,50 +101,3 @@ const lowRisk: RiskFactorsInput = { function createEventBus() { return { emit: vi.fn().mockResolvedValue(undefined) } as unknown as Pick & { emit: ReturnType }; } - -describe('RiskService', () => { - it('emits a RiskEvaluated event with full factor breakdown', async () => { - const eventBus = createEventBus(); - const service = new RiskService(new RiskEngine(), eventBus as unknown as EventBusService); - - const assessment = await service.evaluate('org-1', lowRisk, { - transactionId: 'tx-1', - actorId: 'agent-1', - }); - - expect(assessment.band).toBe(RiskBand.LOW); - expect(assessment.factors.length).toBe(6); - - const emitMock = eventBus.emit as ReturnType; - expect(emitMock).toHaveBeenCalledOnce(); - const [eventName, payload] = emitMock.mock.calls[0]; - expect(eventName).toBe('risk.evaluated'); - expect(payload.transactionId).toBe('tx-1'); - expect(payload.score).toBe(assessment.score); - expect(payload.band).toBe(RiskBand.LOW); - expect(payload.factors).toEqual(assessment.factors); - expect(payload.canAutoExecute).toBe(true); - }); - - it('assess() returns a result without emitting events', async () => { - const eventBus = createEventBus(); - const service = new RiskService(new RiskEngine(), eventBus as unknown as EventBusService); - - const assessment = service.assess(lowRisk); - expect(assessment.band).toBe(RiskBand.LOW); - const emitMock = eventBus.emit as ReturnType; - expect(emitMock).not.toHaveBeenCalled(); - }); - - it('passes config overrides through to the engine', async () => { - const eventBus = createEventBus(); - const service = new RiskService(new RiskEngine(), eventBus as unknown as EventBusService); - - const assessment = service.assess( - { ...lowRisk, amount: 100 }, - { amountSaturation: 100 }, - ); - const amountFactor = assessment.factors.find((f) => f.factor === 'amount'); - expect(amountFactor!.contribution).toBe(30); - }); -}); diff --git a/src/modules/risk/risk.service.ts b/src/modules/risk/risk.service.ts index 42dd8324..b1dcff01 100644 --- a/src/modules/risk/risk.service.ts +++ b/src/modules/risk/risk.service.ts @@ -3,6 +3,9 @@ import { RiskEngine } from './risk.engine'; import { RiskAssessment, RiskConfig, RiskFactorsInput, RiskRule } from './risk.types'; import { EventBusService } from '../../events/event-bus.service'; import { DomainEventName } from '../../events/event-names'; +import { TypedOnEvent } from '../../events/typed-event-listener.decorator'; +import { DomainEventEnvelope, TransactionInitiatedPayload } from '../../events/domain-event.types'; +import { PrismaService } from '../../database/prisma.service'; /** * Application-facing risk service. Wraps the pure {@link RiskEngine}, emits a @@ -14,6 +17,7 @@ export class RiskService { constructor( private readonly engine: RiskEngine, private readonly eventBus: EventBusService, + private readonly prisma: PrismaService, ) {} /** @@ -53,4 +57,58 @@ export class RiskService { ): RiskAssessment { return this.engine.assess(input, config, rules); } + + /** + * Event listener for transaction initiation. Automatically evaluates risk analysis + * asynchronously and persists risk scores to the database upon consumption. + */ + @TypedOnEvent('transaction.initiated') + async handleTransactionInitiated( + envelope: DomainEventEnvelope, + ): Promise { + const payload = envelope.payload; + const organizationId = envelope.organizationId || payload.organizationId; + const transactionId = payload.transactionId; + + const input: RiskFactorsInput = { + amount: payload.amount ? parseFloat(payload.amount) : 0, + asset: payload.asset || 'XLM', + knownRecipient: false, + recentTransactionCount: 0, + walletAgeDays: 0, + policyViolations: 0, + }; + + const assessment = this.engine.assess(input); + + await this.prisma.transaction.update({ + where: { id: transactionId }, + data: { + riskScore: assessment.score, + riskBand: assessment.band, + }, + }).catch(() => { + // If transaction record is not found or fails to update directly, ignore or handle gracefully + }); + + await this.eventBus.emit( + DomainEventName.RiskEvaluated, + { + transactionId, + score: assessment.score, + band: assessment.band, + factors: assessment.factors, + canAutoExecute: assessment.canAutoExecute, + }, + { + organizationId, + actorId: envelope.actorId, + aggregateType: 'transaction', + aggregateId: transactionId, + correlationId: envelope.correlationId, + }, + ); + + return assessment; + } }