Skip to content
Merged
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
1 change: 1 addition & 0 deletions jest.config.js
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ module.exports = {
'^.+\\.(t|j)s$': [
'ts-jest',
{
isolatedModules: true,
tsconfig: {
types: ['node', 'jest'],
skipLibCheck: true,
Expand Down
11 changes: 9 additions & 2 deletions src/payments/payments.module.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
import { Module } from '@nestjs/common';
import { HttpModule } from '@nestjs/axios';
import { TypeOrmModule } from '@nestjs/typeorm';
import { BullModule } from '@nestjs/bull';
import { CurrencyModule } from '../currency/currency.module';
import { AuditLogModule } from '../audit-log/audit-log.module';
import { IdempotencyModule } from '../common/modules/idempotency.module';
Expand All @@ -12,10 +13,12 @@ import { Invoice } from './entities/invoice.entity';
import { Refund } from './entities/refund.entity';
import { PricingService } from './services/pricing.service';
import { PricingController } from './controllers/pricing.controller';
import { PaymentReconciliationJob } from './reconciliation/reconciliation.service';
import { PaymentReconciliationController } from './reconciliation/reconciliation.controller';
import { SubscriptionsService } from './subscriptions/subscriptions.service';
import { SubscriptionsController } from './subscriptions/subscriptions.controller';
import { SubscriptionJobProcessor } from './subscriptions/subscription-job.processor';
import { QUEUE_NAMES } from '../common/constants/queue.constants';
import { PaymentReconciliationJob } from './reconciliation/reconciliation.service';
import { PaymentReconciliationController } from './reconciliation/reconciliation.controller';
import { PaymentProviderService } from './providers/payment-provider.service';
import { StripeProvider } from './providers/stripe.provider';

Expand Down Expand Up @@ -46,12 +49,16 @@ import { StripeProvider } from './providers/stripe.provider';
OutboxModule,
HttpModule,
QueueModule,
BullModule.registerQueue({
name: QUEUE_NAMES.SUBSCRIPTIONS,
}),
],
providers: [
PricingService,
PaymentReconciliationJob,
StripeProvider,
SubscriptionsService,
SubscriptionJobProcessor,
PaymentProviderService,
{
provide: 'IPaymentProvider',
Expand Down
135 changes: 135 additions & 0 deletions src/payments/subscriptions/subscription-job.processor.spec.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,135 @@
import { Test, TestingModule } from '@nestjs/testing';
import { getRepositoryToken } from '@nestjs/typeorm';
import { Job } from 'bull';
import { SubscriptionJobProcessor } from './subscription-job.processor';
import { SubscriptionsService } from './subscriptions.service';
import { Subscription, SubscriptionStatus } from '../entities/subscription.entity';

describe('SubscriptionJobProcessor', () => {
let processor: SubscriptionJobProcessor;

const mockSubscriptionsService = {
resumeSubscription: jest.fn(),
};

const mockSubscriptionRepository = {
findOne: jest.fn(),
};

beforeEach(async () => {
jest.clearAllMocks();

const module: TestingModule = await Test.createTestingModule({
providers: [
SubscriptionJobProcessor,
{
provide: SubscriptionsService,
useValue: mockSubscriptionsService,
},
{
provide: getRepositoryToken(Subscription),
useValue: mockSubscriptionRepository,
},
],
}).compile();

processor = module.get<SubscriptionJobProcessor>(SubscriptionJobProcessor);
});

describe('handleSubscription', () => {
it('should process default subscription job', async () => {
const job = { data: { test: true } } as Job<unknown>;
const result = await processor.handleSubscription(job);
expect(result).toEqual({ success: true });
});
});

describe('handleResumeSubscription', () => {
it('should return error if subscriptionId is missing', async () => {
const job = { data: {} } as Job<any>;
const result = await processor.handleResumeSubscription(job);
expect(result).toEqual({ success: false, reason: 'Missing subscriptionId' });
expect(mockSubscriptionRepository.findOne).not.toHaveBeenCalled();
expect(mockSubscriptionsService.resumeSubscription).not.toHaveBeenCalled();
});

it('should handle missing subscription safely as no-op', async () => {
mockSubscriptionRepository.findOne.mockResolvedValue(null);
const job = { data: { subscriptionId: 'sub-missing' } } as Job<any>;

const result = await processor.handleResumeSubscription(job);

expect(result).toEqual({ success: false, reason: 'Subscription not found' });
expect(mockSubscriptionsService.resumeSubscription).not.toHaveBeenCalled();
});

it('should safely handle cancelled subscription as no-op', async () => {
mockSubscriptionRepository.findOne.mockResolvedValue({
id: 'sub-cancelled',
status: SubscriptionStatus.CANCELLED,
properties: { isPaused: true },
});
const job = { data: { subscriptionId: 'sub-cancelled' } } as Job<any>;

const result = await processor.handleResumeSubscription(job);

expect(result).toEqual({ success: false, reason: 'Subscription cancelled' });
expect(mockSubscriptionsService.resumeSubscription).not.toHaveBeenCalled();
});

it('should be idempotent and no-op if subscription is already active / not paused', async () => {
mockSubscriptionRepository.findOne.mockResolvedValue({
id: 'sub-active',
status: SubscriptionStatus.ACTIVE,
properties: { isPaused: false },
});
const job = { data: { subscriptionId: 'sub-active' } } as Job<any>;

const result = await processor.handleResumeSubscription(job);

expect(result).toEqual({ success: true, reason: 'Subscription not paused' });
expect(mockSubscriptionsService.resumeSubscription).not.toHaveBeenCalled();
});

it('should successfully resume a paused subscription', async () => {
mockSubscriptionRepository.findOne.mockResolvedValue({
id: 'sub-paused',
status: SubscriptionStatus.ACTIVE,
properties: { isPaused: true, resumeAt: '2026-09-02T16:00:00.000Z' },
});
mockSubscriptionsService.resumeSubscription.mockResolvedValue({
id: 'sub-paused',
status: SubscriptionStatus.ACTIVE,
properties: { isPaused: false },
});

const job = {
data: { subscriptionId: 'sub-paused', userId: 'user-1', reason: 'Automatic resume' },
} as Job<any>;

const result = await processor.handleResumeSubscription(job);

expect(result).toEqual({ success: true });
expect(mockSubscriptionsService.resumeSubscription).toHaveBeenCalledWith('sub-paused', {
reason: 'Automatic resume',
});
});

it('should throw error when resumeSubscription fails so Bull can retry', async () => {
mockSubscriptionRepository.findOne.mockResolvedValue({
id: 'sub-paused',
status: SubscriptionStatus.ACTIVE,
properties: { isPaused: true },
});
mockSubscriptionsService.resumeSubscription.mockRejectedValue(
new Error('Database connection failed'),
);

const job = { data: { subscriptionId: 'sub-paused' } } as Job<any>;

await expect(processor.handleResumeSubscription(job)).rejects.toThrow(
'Database connection failed',
);
});
});
});
102 changes: 56 additions & 46 deletions src/payments/subscriptions/subscription-job.processor.ts
Original file line number Diff line number Diff line change
@@ -1,86 +1,96 @@
import { Processor, Process, OnQueueActive, OnQueueCompleted, OnQueueFailed } from '@nestjs/bull';

import { Inject, Logger } from '@nestjs/common';
import { Inject, Injectable, Logger, Optional } from '@nestjs/common';
import { Job } from 'bull';
import { QUEUE_NAMES, JOB_NAMES } from '../../common/constants/queue.constants';
import { InjectRepository } from '@nestjs/typeorm';
import { Repository } from 'typeorm';
import { QUEUE_NAMES, JOB_NAMES } from '../../common/constants/queue.constants';
import { SubscriptionsService } from './subscriptions.service';
import { Subscription, SubscriptionStatus } from '../entities/subscription.entity';
import { IPaymentProvider } from '../providers/payment-provider.interface';

export interface ResumeSubscriptionJobData {
subscriptionId: string;
userId?: string;
reason?: string;
}

@Injectable()
@Processor(QUEUE_NAMES.SUBSCRIPTIONS)
export class SubscriptionJobProcessor {
private readonly logger = new Logger(SubscriptionJobProcessor.name);

constructor(
private readonly subscriptionsService: SubscriptionsService,
@InjectRepository(Subscription)
private subscriptionRepository: Repository<Subscription>,
private readonly subscriptionRepository: Repository<Subscription>,
@Optional()
@Inject('IPaymentProvider')
private paymentProvider: IPaymentProvider,
private readonly paymentProvider?: IPaymentProvider,
) {}

@Process(JOB_NAMES.PROCESS_SUBSCRIPTION)
async handleSubscription(job: Job<unknown>): Promise<unknown> {
// Process subscription job
this.logger.log('Processing subscription job:', job.data);
return { success: true };
}

@Process(JOB_NAMES.RESUME_SUBSCRIPTION)
async handleResumeSubscription(
job: Job<{ subscriptionId: string }>,
): Promise<{ success: boolean; message: string }> {
job: Job<ResumeSubscriptionJobData>,
): Promise<{ success: boolean; reason?: string; message?: string }> {
const { subscriptionId } = job.data;
this.logger.log(`Processing automatic resume for subscription: ${subscriptionId}`);

try {
this.logger.log(`Processing resume subscription job for ${subscriptionId}`);
if (!subscriptionId) {
this.logger.warn('Missing subscriptionId in resume job payload');
return { success: false, reason: 'Missing subscriptionId' };
}

const subscription = await this.subscriptionRepository.findOne({
where: { id: subscriptionId },
});
const subscription = await this.subscriptionRepository.findOne({
where: { id: subscriptionId },
});

if (!subscription) {
this.logger.error(`Subscription ${subscriptionId} not found`);
return { success: false, message: 'Subscription not found' };
}
if (!subscription) {
this.logger.warn(`Subscription ${subscriptionId} not found for auto-resume`);
return { success: false, reason: 'Subscription not found' };
}

if (subscription.status !== SubscriptionStatus.PAUSED) {
this.logger.warn(
`Subscription ${subscriptionId} is not paused (status: ${subscription.status})`,
);
return { success: false, message: 'Subscription is not paused' };
}
// Safe handling for cancelled subscriptions
if (subscription.status === SubscriptionStatus.CANCELLED) {
this.logger.log(`Subscription ${subscriptionId} is cancelled. Skipping auto-resume.`);
return { success: false, reason: 'Subscription cancelled' };
}

if (!subscription.providerSubscriptionId) {
this.logger.error(`Subscription ${subscriptionId} has no provider subscription ID`);
return { success: false, message: 'No provider subscription ID' };
}
// Idempotency check: Guard against double-resume or already resumed subscriptions
if (!subscription.properties?.isPaused && subscription.status !== SubscriptionStatus.PAUSED) {
this.logger.log(
`Subscription ${subscriptionId} is not paused (already resumed). Skipping auto-resume.`,
);
return { success: true, reason: 'Subscription not paused' };
}

// Resume at provider (Stripe) first
// Resume at provider (Stripe) if provider subscription ID is present
if (subscription.providerSubscriptionId && this.paymentProvider?.resumeSubscription) {
try {
await this.paymentProvider.resumeSubscription(subscription.providerSubscriptionId);
} catch (error) {
this.logger.error(`Failed to resume subscription ${subscriptionId} at provider`, error);
return { success: false, message: 'Provider resume failed' };
this.logger.error(
`Failed to resume subscription ${subscriptionId} at payment provider: ${(error as Error).message}`,
);
throw error;
}
}

// Resume the subscription locally only after provider succeeds
subscription.status = SubscriptionStatus.ACTIVE;
subscription.cancelAtPeriodEnd = false;
subscription.properties = {
...subscription.properties,
isPaused: false,
resumedAt: new Date(),
resumeReason: 'Scheduled automatic resume',
};

await this.subscriptionRepository.save(subscription);

this.logger.log(`Successfully resumed subscription ${subscriptionId} via scheduled job`);

return { success: true, message: 'Subscription resumed successfully' };
try {
await this.subscriptionsService.resumeSubscription(subscriptionId, {
reason: job.data.reason || 'Automatic resume from scheduled pause',
});
this.logger.log(`Successfully auto-resumed subscription ${subscriptionId}`);
return { success: true };
} catch (error) {
this.logger.error(`Failed to resume subscription ${subscriptionId}`, error);
this.logger.error(
`Failed to auto-resume subscription ${subscriptionId}: ${(error as Error).message}`,
);
throw error;
}
}
Expand Down
Loading
Loading