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
2 changes: 2 additions & 0 deletions apps/api/src/app/app.module.ts
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,7 @@ import { ExceptionsModule } from '../modules/exceptions/exceptions.module.js';
import { AiModule } from '../modules/ai/ai.module.js';
import { MappingsModule } from '../modules/mappings/mappings.module.js';
import { GitopsModule } from '../modules/gitops/gitops.module.js';
import { PluginsModule } from '../modules/plugins/plugins.module.js';

@Module({
imports: [
Expand Down Expand Up @@ -115,6 +116,7 @@ import { GitopsModule } from '../modules/gitops/gitops.module.js';
AiModule,
MappingsModule,
GitopsModule,
PluginsModule,
],
controllers: [AppController],
providers: [AppService, ShutdownService],
Expand Down
217 changes: 217 additions & 0 deletions apps/api/src/modules/plugins/plugins.controller.spec.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,217 @@
/* eslint-disable @typescript-eslint/unbound-method, @typescript-eslint/no-unsafe-assignment */
import { Test, TestingModule } from '@nestjs/testing';
import { PluginsController, WebhookPayloadDto } from './plugins.controller.js';
import { PluginManagerService } from '@soopa/piece-registry';
import { ConfigService } from '@nestjs/config';
import { UnauthorizedException } from '@nestjs/common';
import { QUEUE_SERVICE, QueueName } from '@soopa/queue';
import type { IQueueService } from '@soopa/queue';
import * as crypto from 'crypto';
import { describe, it, expect, beforeEach, afterEach, vi } from 'vitest';
import type { Mocked } from 'vitest';

describe('PluginsController', () => {
let controller: PluginsController;
let pluginManagerService: Mocked<PluginManagerService>;
let configService: Mocked<ConfigService>;
let queueService: Mocked<IQueueService>;
let installPieceMock: ReturnType<typeof vi.fn>;

const TEST_SECRET = 'test-secret';

beforeEach(async () => {
installPieceMock = vi.fn();

// Mock the dependencies
pluginManagerService = {
installPiece: installPieceMock,
} as unknown as Mocked<PluginManagerService>;

configService = {
get: vi.fn().mockImplementation((key: string) => {
if (key === 'NPM_WEBHOOK_SECRET') return TEST_SECRET;
return null;
}),
} as unknown as Mocked<ConfigService>;

queueService = {
send: vi.fn().mockResolvedValue(undefined),
consume: vi.fn(),
stopConsuming: vi.fn().mockResolvedValue(undefined),
} as unknown as Mocked<IQueueService>;

const module: TestingModule = await Test.createTestingModule({
controllers: [PluginsController],
providers: [
{ provide: PluginManagerService, useValue: pluginManagerService },
{ provide: ConfigService, useValue: configService },
{ provide: QUEUE_SERVICE, useValue: queueService },
],
}).compile();

controller = module.get<PluginsController>(PluginsController);
});

afterEach(() => {
vi.clearAllMocks();
});
Comment thread
coderabbitai[bot] marked this conversation as resolved.

// Helper to generate valid HMAC signature from raw body buffer
const generateValidSignature = (rawBody: Buffer): string => {
const hmac = crypto.createHmac('sha256', TEST_SECRET);
return 'sha256=' + hmac.update(rawBody).digest('hex');
};

it('should be defined', () => {
expect(controller).toBeDefined();
});

describe('handleNpmWebhook', () => {
it('should throw UnauthorizedException when NPM_WEBHOOK_SECRET is not configured', async () => {
// Override configService to return undefined for NPM_WEBHOOK_SECRET
configService.get = vi.fn().mockReturnValue(undefined);

const payload = {
name: '@soopa/piece-slack',
version: '1.0.0',
} as WebhookPayloadDto;

const rawBody = Buffer.from(JSON.stringify(payload));

await expect(
controller.handleNpmWebhook('sha256=somesignature', payload, rawBody),
).rejects.toThrow(UnauthorizedException);
expect(queueService.send).not.toHaveBeenCalled();
});

it('should throw UnauthorizedException if signature is missing but secret is configured', async () => {
const payload = {
name: '@soopa/piece-slack',
version: '1.0.0',
} as WebhookPayloadDto;

const rawBody = Buffer.from(JSON.stringify(payload));

await expect(
controller.handleNpmWebhook('', payload, rawBody),
).rejects.toThrow(UnauthorizedException);
expect(queueService.send).not.toHaveBeenCalled();
});

it('should throw UnauthorizedException if signature is invalid', async () => {
const payload = {
name: '@soopa/piece-slack',
version: '1.0.0',
} as WebhookPayloadDto;
const invalidSignature = 'sha256=invalidhash12345';
const rawBody = Buffer.from(JSON.stringify(payload));

await expect(
controller.handleNpmWebhook(invalidSignature, payload, rawBody),
).rejects.toThrow(UnauthorizedException);
expect(queueService.send).not.toHaveBeenCalled();
});

it('should ignore payload without a package name gracefully', async () => {
const payload = { version: '1.0.0' } as WebhookPayloadDto;
const rawBody = Buffer.from(JSON.stringify(payload));
const signature = generateValidSignature(rawBody);

const result = await controller.handleNpmWebhook(
signature,
payload,
rawBody,
);

expect(result).toEqual({
status: 'ignored',
reason: 'No package name found in payload',
});
expect(queueService.send).not.toHaveBeenCalled();
});

it('should ignore payload that does not belong to @soopa scope', async () => {
const payload = {
name: '@other/piece-slack',
version: '1.0.0',
} as WebhookPayloadDto;
const rawBody = Buffer.from(JSON.stringify(payload));
const signature = generateValidSignature(rawBody);

const result = await controller.handleNpmWebhook(
signature,
payload,
rawBody,
);

expect(result).toEqual({
status: 'ignored',
reason: 'Only @soopa packages are hot-loaded',
});
expect(queueService.send).not.toHaveBeenCalled();
});

it('should successfully queue a valid piece installation and return accepted', async () => {
const payload = {
name: '@soopa/piece-slack',
version: '1.2.3',
} as unknown as WebhookPayloadDto;
const rawBody = Buffer.from(JSON.stringify(payload));
const signature = generateValidSignature(rawBody);

const result = await controller.handleNpmWebhook(
signature,
payload,
rawBody,
);

expect(result).toEqual({
status: 'accepted',
message: 'Installation queued for @soopa/piece-slack@1.2.3',
});

// Verify the install event was sent to the queue
expect(queueService.send).toHaveBeenCalledWith(
QueueName.PluginInstallQueue,
expect.objectContaining({
packageName: '@soopa/piece-slack',
version: '1.2.3',
requestMetadata: expect.objectContaining({
source: 'npm-webhook',
}),
}),
);

// The controller should NOT directly call installPiece anymore
expect(installPieceMock).not.toHaveBeenCalled();
});

it('should return accepted status when queue send succeeds', async () => {
const payload = {
name: '@soopa/piece-slack',
version: '1.2.3',
} as unknown as WebhookPayloadDto;
const rawBody = Buffer.from(JSON.stringify(payload));
const signature = generateValidSignature(rawBody);

const result = await controller.handleNpmWebhook(
signature,
payload,
rawBody,
);

expect(result).toEqual({
status: 'accepted',
message: 'Installation queued for @soopa/piece-slack@1.2.3',
});

expect(queueService.send).toHaveBeenCalledWith(
QueueName.PluginInstallQueue,
expect.objectContaining({
packageName: '@soopa/piece-slack',
version: '1.2.3',
}),
);
});
});
});
142 changes: 142 additions & 0 deletions apps/api/src/modules/plugins/plugins.controller.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,142 @@
import {
Controller,
Post,
Body,
Headers,
Logger,
UnauthorizedException,
Inject,
RawBody,
} from '@nestjs/common';
import { ConfigService } from '@nestjs/config';
import { PluginManagerService } from '@soopa/piece-registry';
import { QUEUE_SERVICE, QueueName } from '@soopa/queue';
import type { IQueueService, PluginInstallEvent } from '@soopa/queue';
import * as crypto from 'crypto';
import { createZodDto } from 'nestjs-zod';
import { z } from 'zod';

const WebhookPayloadSchema = z
.object({
name: z.string().optional(),
version: z.string().optional(),
package: z
.object({
name: z.string().optional(),
version: z.string().optional(),
})
.optional(),
})
.passthrough();

export type WebhookPayload = z.infer<typeof WebhookPayloadSchema>;
export class WebhookPayloadDto extends createZodDto(WebhookPayloadSchema) {}

@Controller('api/internal/system/plugins')
export class PluginsController {
private readonly logger = new Logger(PluginsController.name);

constructor(
private readonly pluginManager: PluginManagerService,
private readonly configService: ConfigService,
@Inject(QUEUE_SERVICE) private readonly queueService: IQueueService,
) {}

/**
* Webhook endpoint triggered by the NPM Registry (e.g., Verdaccio)
* when a new version of a piece is published.
*/
@Post('webhook')
async handleNpmWebhook(
@Headers('x-npm-signature') signature: string,
@Body() payloadRaw: WebhookPayloadDto,
@RawBody() rawBody: Buffer,
) {
const payload = payloadRaw as unknown as WebhookPayload;
this.logger.log('Received NPM publish webhook event');

const webhookSecret = this.configService.get<string>('NPM_WEBHOOK_SECRET');

// Signature validation is mandatory - fail if secret is not configured
if (!webhookSecret) {
this.logger.error(
'NPM_WEBHOOK_SECRET is not configured. Rejecting webhook request.',
);
throw new UnauthorizedException('Webhook authentication not configured');
}

if (!signature) {
throw new UnauthorizedException('Missing webhook signature');
}

// Validate HMAC signature to ensure request originated from our private registry
const hmac = crypto.createHmac('sha256', webhookSecret);
const digest = 'sha256=' + hmac.update(rawBody).digest('hex');

// Use constant-time comparison to prevent timing attacks
try {
const signatureBuffer = Buffer.from(signature, 'utf8');
const digestBuffer = Buffer.from(digest, 'utf8');

// Fail early if lengths don't match
if (signatureBuffer.length !== digestBuffer.length) {
throw new Error('Signature length mismatch');
}

if (!crypto.timingSafeEqual(signatureBuffer, digestBuffer)) {
this.logger.warn(
'Invalid NPM webhook signature detected. Dropping payload.',
);
throw new UnauthorizedException('Invalid webhook signature');
}
} catch (_error) {
this.logger.warn(
'Invalid NPM webhook signature detected. Dropping payload.',
);
throw new UnauthorizedException('Invalid webhook signature');
}
Comment thread
coderabbitai[bot] marked this conversation as resolved.

// Example Verdaccio payload parsing:
// Extract the package name and version from the webhook payload.
// Ensure we only process @soopa scopes or allowed packages.
const packageName = payload?.name || payload?.package?.name;
const version = payload?.version || payload?.package?.version || 'latest';

if (!packageName) {
return { status: 'ignored', reason: 'No package name found in payload' };
}

if (!packageName.startsWith('@soopa/')) {
return {
status: 'ignored',
reason: 'Only @soopa packages are hot-loaded',
};
}

this.logger.log(
`Persisting durable install job for ${packageName}@${version}...`,
);

// Persist a durable install job to the queue before responding
// This ensures the install request is not lost if the process dies
const installEvent: PluginInstallEvent = {
packageName,
version,
requestMetadata: {
webhookReceivedAt: new Date().toISOString(),
source: 'npm-webhook',
},
};

await this.queueService.send(QueueName.PluginInstallQueue, installEvent);

this.logger.log(
`Installation job persisted to queue for ${packageName}@${version}`,
);

return {
status: 'accepted',
message: `Installation queued for ${packageName}@${version}`,
};
}
}
9 changes: 9 additions & 0 deletions apps/api/src/modules/plugins/plugins.module.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,9 @@
import { Module } from '@nestjs/common';
import { PluginsController } from './plugins.controller.js';
import { PiecesModule } from '@soopa/piece-registry';

@Module({
imports: [PiecesModule.forRoot({ anchorUrl: import.meta.url })], // Brings in PluginManagerService
controllers: [PluginsController],
})
export class PluginsModule {}
1 change: 0 additions & 1 deletion packages/pieces/platform/quickbooks/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,6 @@
"vitest": "^4.0.18"
},
"dependencies": {
"@soopa/credentials": "workspace:*",
"@soopa/piece-framework": "workspace:*",
"dayjs": "^1.11.19"
}
Expand Down
Loading
Loading