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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 11 additions & 0 deletions src/events/domain-event.types.ts
Original file line number Diff line number Diff line change
Expand Up @@ -210,6 +210,17 @@ export interface ProposalExecutedPayload extends Record<string, unknown> {
transactionId: string;
}

export interface TransactionInitiatedPayload extends Record<string, unknown> {
transactionId: string;
organizationId: string;
walletId?: string;
agentId?: string;
amount?: string;
asset?: string;
recipientAddress?: string;
memo?: string;
}

export interface TransactionCreatedPayload extends Record<string, unknown> {
transactionId: string;
walletId?: string;
Expand Down
2 changes: 2 additions & 0 deletions src/events/event-names.ts
Original file line number Diff line number Diff line change
Expand Up @@ -49,6 +49,7 @@ export const DomainEventName = {
ProposalExecuted: 'proposal.executed',

// Transaction / blockchain
TransactionInitiated: 'transaction.initiated',
TransactionCreated: 'transaction.created',
TransactionSubmitted: 'transaction.submitted',
TransactionConfirmed: 'transaction.confirmed',
Expand Down Expand Up @@ -78,6 +79,7 @@ export const WEBHOOK_EVENTS: DomainEventNameType[] = [
DomainEventName.WalletUpdated,
DomainEventName.ProposalApproved,
DomainEventName.ProposalRejected,
DomainEventName.TransactionInitiated,
DomainEventName.TransactionCompleted,
DomainEventName.TransactionFailed,
DomainEventName.PolicyViolated,
Expand Down
1 change: 1 addition & 0 deletions src/events/typed-event-emitter.service.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
5 changes: 4 additions & 1 deletion src/modules/risk/risk.module.ts
Original file line number Diff line number Diff line change
@@ -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],
Expand Down
136 changes: 86 additions & 50 deletions src/modules/risk/risk.service.spec.ts
Original file line number Diff line number Diff line change
@@ -1,11 +1,94 @@
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>(RiskService);
eventBus = module.get<EventBusService>(EventBusService);
prisma = module.get<PrismaService>(PrismaService);
vi.clearAllMocks();
});

it('should handle transaction.initiated event, evaluate risk, persist scores, and emit risk.evaluated', async () => {
const envelope: DomainEventEnvelope<TransactionInitiatedPayload> = {
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 = {

Check failure on line 91 in src/modules/risk/risk.service.spec.ts

View workflow job for this annotation

GitHub Actions / build-and-test

'lowRisk' is declared but its value is never read.
amount: 20,
asset: 'USDC',
knownRecipient: true,
Expand All @@ -15,53 +98,6 @@
hourUtc: 12,
};

function createEventBus() {

Check failure on line 101 in src/modules/risk/risk.service.spec.ts

View workflow job for this annotation

GitHub Actions / build-and-test

'createEventBus' is declared but its value is never read.
return { emit: vi.fn().mockResolvedValue(undefined) } as unknown as Pick<EventBusService, 'emit'> & { emit: ReturnType<typeof vi.fn> };
}

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<typeof vi.fn>;
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<typeof vi.fn>;
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);
});
});
58 changes: 58 additions & 0 deletions src/modules/risk/risk.service.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -14,6 +17,7 @@ export class RiskService {
constructor(
private readonly engine: RiskEngine,
private readonly eventBus: EventBusService,
private readonly prisma: PrismaService,
) {}

/**
Expand Down Expand Up @@ -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<TransactionInitiatedPayload>,
): Promise<RiskAssessment> {
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;
}
}
Loading