feat: marketplace modules and isolated container management

This commit is contained in:
ayde64
2026-10-07 23:18:02 +02:00
parent 564f7def1c
commit 7b121fa908
68 changed files with 2560 additions and 699 deletions

View File

@@ -55,20 +55,21 @@ export class AuthController {
ipAddress: request.ip ?? null,
});
const cookieMaxAgeSeconds = this.config.security.sessionTtlMinutes * 60;
// Express expects cookie maxAge in milliseconds (the DB TTL is in minutes).
const cookieMaxAgeMs = this.config.security.sessionTtlMinutes * 60 * 1000;
response.cookie('mpm_session', result.sessionToken, {
httpOnly: true,
secure: this.config.security.cookieSecure,
sameSite: 'lax',
path: '/',
maxAge: cookieMaxAgeSeconds,
maxAge: cookieMaxAgeMs,
});
response.cookie('mpm_csrf', result.session.csrfToken, {
httpOnly: false,
secure: this.config.security.cookieSecure,
sameSite: 'lax',
path: '/',
maxAge: cookieMaxAgeSeconds,
maxAge: cookieMaxAgeMs,
});
return { user: toAuthUserResponse(result.user) };
@@ -95,4 +96,4 @@ export class AuthController {
async me(@CurrentUser() user: AuthUser): Promise<{ user: AuthUserResponse }> {
return { user: toAuthUserResponse(user) };
}
}
}

View File

@@ -26,6 +26,7 @@ function createConfig(overrides: Partial<AppConfig['security']> = {}): AppConfig
},
adminSeed: { username: 'admin', email: 'admin@example.com', password: 'password-123' },
runtime: { modulesDir: '/data/modules', logsDir: '/data/logs' },
marketplace: { publicUrl: 'http://127.0.0.1:8081', tokenEncryptionKey: '', providers: {} },
};
}
@@ -52,7 +53,7 @@ function createUserRecord(overrides: Partial<UserRecord> = {}): UserRecord {
class MockUserRepository {
public findByUsernameResult: UserRecord | null = null;
public updateLoginSuccessCalls: string[] = [];
public updateLoginFailureCalls: Array<{ userId: string; attempts: number; shouldLock: boolean; lockoutMinutes: number }> = [];
public updateLoginFailureCalls: string[] = [];
async findByUsername(): Promise<UserRecord | null> {
return this.findByUsernameResult;
@@ -62,8 +63,8 @@ class MockUserRepository {
this.updateLoginSuccessCalls.push(userId);
}
async updateLoginFailure(userId: string, attempts: number, shouldLock: boolean, lockoutMinutes: number): Promise<void> {
this.updateLoginFailureCalls.push({ userId, attempts, shouldLock, lockoutMinutes });
async updateLoginFailure(userId: string): Promise<void> {
this.updateLoginFailureCalls.push(userId);
}
}
@@ -158,12 +159,12 @@ describe('AuthService', () => {
).rejects.toThrow(UnauthorizedException);
expect(userRepository.updateLoginFailureCalls).toEqual([
{ userId: 'user-1', attempts: 1, shouldLock: false, lockoutMinutes: 15 },
'user-1',
]);
expect(auditService.records.at(-1)?.action).toBe(AUDIT_ACTIONS.LOGIN_FAILED);
});
it('sperrt das Konto nach Erreichen der maximalen Fehlversuche', async () => {
it('verhindert Loginversuche nicht durch Kontosperren', async () => {
const passwordHash = await passwordHasher.hash('Sicheres-Passwort-1');
userRepository.findByUsernameResult = createUserRecord({
passwordHash,
@@ -175,9 +176,9 @@ describe('AuthService', () => {
).rejects.toThrow(UnauthorizedException);
expect(userRepository.updateLoginFailureCalls).toEqual([
{ userId: 'user-1', attempts: 3, shouldLock: true, lockoutMinutes: 15 },
'user-1',
]);
expect(auditService.records.at(-1)?.action).toBe(AUDIT_ACTIONS.LOGIN_LOCKED);
expect(auditService.records.at(-1)?.action).toBe(AUDIT_ACTIONS.LOGIN_FAILED);
});
it('lehnt gesperrte Benutzer ab', async () => {
@@ -187,11 +188,12 @@ describe('AuthService', () => {
lockedUntil: new Date(Date.now() + 60_000),
});
await expect(
authService.login({ username: 'max', password: 'Sicheres-Passwort-1', ipAddress: '127.0.0.1' },
)).rejects.toThrow(UnauthorizedException);
expect(auditService.records.at(-1)?.details).toEqual({ reason: 'ACCOUNT_LOCKED' });
const result = await authService.login({
username: 'max',
password: 'Sicheres-Passwort-1',
ipAddress: '127.0.0.1',
});
expect(result.user.username).toBe('max');
});
it('lehnt deaktivierte Benutzer ab', async () => {
@@ -226,7 +228,7 @@ describe('AuthService', () => {
expect(auditService.records.at(-1)?.details).toEqual({ reason: 'RATE_LIMITED' });
});
it('setzt das Rate-Limit-Fenster nach erfolgreichem Login zurück', async () => {
it('setzt das Rate-Limit-Fenster nach erfolgreichem Login nicht zurück', async () => {
const passwordHash = await passwordHasher.hash('Sicheres-Passwort-1');
userRepository.findByUsernameResult = createUserRecord({ passwordHash });
@@ -243,13 +245,11 @@ describe('AuthService', () => {
});
expect(result.user.username).toBe('max');
// Nach Reset ist ein neuer Login sofort wieder möglich.
const secondResult = await authService.login({
await expect(authService.login({
username: 'max',
password: 'Sicheres-Passwort-1',
ipAddress: '127.0.0.1',
});
expect(secondResult.user.username).toBe('max');
})).rejects.toThrow('Zu viele Anmeldeversuche. Bitte später erneut versuchen.');
});
});
@@ -268,4 +268,4 @@ describe('AuthService', () => {
expect(auditService.records.at(-1)?.action).toBe(AUDIT_ACTIONS.LOGOUT);
});
});
});
});

View File

@@ -20,7 +20,7 @@ const INVALID_CREDENTIALS_MESSAGE = 'Benutzername oder Passwort ist falsch';
/**
* Authentifizierungs-Logik (Domain/Application):
* Login mit Rate Limiting, Account Lockout, Argon2id-Verifikation,
* Login mit IP-basiertem Rate Limiting, Argon2id-Verifikation,
* Session-Erstellung und Audit-Logging.
*/
@Injectable()
@@ -72,32 +72,14 @@ export class AuthService {
throw new UnauthorizedException(INVALID_CREDENTIALS_MESSAGE);
}
if (user.lockedUntil && user.lockedUntil > new Date()) {
const passwordValid = await this.passwordHasher.verify(user.passwordHash, input.password);
if (!passwordValid) {
await this.userRepository.updateLoginFailure(user.id);
await this.auditService.record({
userId: user.id,
username: user.username,
action: 'LOGIN_FAILED',
details: { reason: 'ACCOUNT_LOCKED' },
ipAddress: input.ipAddress,
});
throw new UnauthorizedException(INVALID_CREDENTIALS_MESSAGE);
}
const passwordValid = await this.passwordHasher.verify(user.passwordHash, input.password);
if (!passwordValid) {
const attempts = user.failedLoginAttempts + 1;
const shouldLock = attempts >= security.loginMaxAttempts;
await this.userRepository.updateLoginFailure(
user.id,
attempts,
shouldLock,
security.loginLockoutMinutes,
);
await this.auditService.record({
userId: user.id,
username: user.username,
action: shouldLock ? 'LOGIN_FAILED_LOCKED' : 'LOGIN_FAILED',
details: { reason: 'INVALID_PASSWORD', attempts },
details: { reason: 'INVALID_PASSWORD' },
ipAddress: input.ipAddress,
});
throw new UnauthorizedException(INVALID_CREDENTIALS_MESSAGE);
@@ -115,7 +97,6 @@ export class AuthService {
}
await this.userRepository.updateLoginSuccess(user.id);
this.rateLimiter.reset(rateLimitKey);
const { token, data } = await this.sessionService.create(
user.id,
@@ -151,4 +132,4 @@ export class AuthService {
ipAddress,
});
}
}
}

View File

@@ -13,7 +13,7 @@ const STATE_CHANGING_METHODS = new Set(['POST', 'PUT', 'PATCH', 'DELETE']);
*
* Requests ohne Session (z. B. Login) sind ausgenommen: Sie besitzen
* kein Session-CSRF-Token. Das Login ist stattdessen durch Rate
* Limiting, Account Lockout und SameSite=Lax-Cookies geschützt.
* Limiting und SameSite=Lax-Cookies geschützt.
* Ungültige Sessions werden bereits vom SessionGuard mit 401 abgewiesen.
*/
@Injectable()

View File

@@ -16,13 +16,16 @@ export function extractSessionToken(request: RequestWithCookieHeader): string |
if (!cookieHeader) {
return null;
}
let sessionToken: string | null = null;
for (const part of cookieHeader.split(';')) {
const [name, ...value] = part.trim().split('=');
if (name === 'mpm_session') {
return decodeURIComponent(value.join('='));
// Browsers may send same-name cookies from an older, narrower Path
// before the current Path=/ cookie. The last value is the root cookie.
sessionToken = decodeURIComponent(value.join('='));
}
}
return null;
return sessionToken;
}
/**
@@ -79,4 +82,4 @@ export class SessionGuard implements CanActivate {
};
return true;
}
}
}

View File

@@ -13,6 +13,9 @@ interface RateLimitEntry {
@Injectable()
export class RateLimiterService {
private readonly entries = new Map<string, RateLimitEntry>();
private readonly maxEntries = 10_000;
private readonly cleanupIntervalMs = 60_000;
private lastCleanupAt = 0;
/**
* Prüft, ob ein Request innerhalb des Limits liegt.
@@ -21,6 +24,17 @@ export class RateLimiterService {
isAllowed(key: string, limit: number, windowMinutes: number): boolean {
const now = Date.now();
const windowMs = windowMinutes * 60_000;
if (now - this.lastCleanupAt >= this.cleanupIntervalMs) {
for (const [entryKey, entry] of this.entries) {
const recentTimestamps = entry.timestamps.filter((timestamp) => now - timestamp < windowMs);
if (recentTimestamps.length === 0) {
this.entries.delete(entryKey);
} else if (recentTimestamps.length !== entry.timestamps.length) {
this.entries.set(entryKey, { timestamps: recentTimestamps });
}
}
this.lastCleanupAt = now;
}
const entry = this.entries.get(key) ?? { timestamps: [] };
const recent = entry.timestamps.filter((timestamp) => now - timestamp < windowMs);
@@ -30,6 +44,10 @@ export class RateLimiterService {
}
recent.push(now);
if (!this.entries.has(key) && this.entries.size >= this.maxEntries) {
const oldestKey = this.entries.keys().next().value;
if (oldestKey !== undefined) this.entries.delete(oldestKey);
}
this.entries.set(key, { timestamps: recent });
return true;
}
@@ -38,4 +56,4 @@ export class RateLimiterService {
reset(key: string): void {
this.entries.delete(key);
}
}
}

View File

@@ -17,7 +17,9 @@ export interface SecurityConfig {
readonly sessionTtlMinutes: number;
readonly cookieSecure: boolean;
readonly behindProxy: boolean;
/** @deprecated Account lockout was removed to prevent attacker-triggered account denial. */
readonly loginMaxAttempts: number;
/** @deprecated Account lockout was removed to prevent attacker-triggered account denial. */
readonly loginLockoutMinutes: number;
readonly loginRateLimitAttempts: number;
readonly loginRateLimitWindowMinutes: number;
@@ -32,6 +34,23 @@ export interface AdminSeedConfig {
export interface RuntimeConfig {
readonly modulesDir: string;
readonly logsDir: string;
readonly moduleUidBase?: number;
}
export interface MarketplaceProviderConfig {
readonly clientId: string;
readonly clientSecret: string;
readonly baseUrl: string;
}
export interface MarketplaceConfig {
readonly publicUrl: string;
readonly tokenEncryptionKey: string;
readonly providers: {
readonly github?: MarketplaceProviderConfig;
readonly gitea?: MarketplaceProviderConfig;
readonly forgejo?: MarketplaceProviderConfig;
};
}
export interface AppConfig {
@@ -41,6 +60,7 @@ export interface AppConfig {
readonly security: SecurityConfig;
readonly adminSeed: AdminSeedConfig;
readonly runtime: RuntimeConfig;
readonly marketplace: MarketplaceConfig;
}
const booleanFromString = z
@@ -63,7 +83,80 @@ const environmentSchema = z.object({
ADMIN_EMAIL: z.string().trim().email(),
ADMIN_PASSWORD: z.string().min(10, 'ADMIN_PASSWORD muss mindestens 10 Zeichen lang sein').max(200),
MODULES_DIR: z.string().min(1).default('./data/modules'),
MODULE_DATA_DIR: z.string().min(1).default('./data/module-data'),
LOGS_DIR: z.string().min(1).default('./data/logs'),
MODULE_UID_BASE: z.coerce.number().int().min(10_000).max(64_535).optional(),
MARKETPLACE_PUBLIC_URL: z.string().url().default('http://127.0.0.1:8081'),
MARKETPLACE_TOKEN_ENCRYPTION_KEY: z.string().default(''),
GITHUB_OAUTH_CLIENT_ID: z.string().default(''),
GITHUB_OAUTH_CLIENT_SECRET: z.string().default(''),
GITEA_BASE_URL: z.string().default(''),
GITEA_OAUTH_CLIENT_ID: z.string().default(''),
GITEA_OAUTH_CLIENT_SECRET: z.string().default(''),
FORGEJO_BASE_URL: z.string().default(''),
FORGEJO_OAUTH_CLIENT_ID: z.string().default(''),
FORGEJO_OAUTH_CLIENT_SECRET: z.string().default(''),
}).superRefine((environment, context) => {
if (environment.NODE_ENV === 'production' && !environment.COOKIE_SECURE) {
context.addIssue({
code: z.ZodIssueCode.custom,
path: ['COOKIE_SECURE'],
message: 'COOKIE_SECURE muss in production auf true gesetzt sein',
});
}
if (environment.NODE_ENV === 'production' && environment.MODULE_UID_BASE === undefined) {
context.addIssue({
code: z.ZodIssueCode.custom,
path: ['MODULE_UID_BASE'],
message: 'MODULE_UID_BASE ist in production erforderlich, damit Module getrennte UIDs erhalten',
});
}
const oauthFields = [
environment.GITHUB_OAUTH_CLIENT_ID,
environment.GITHUB_OAUTH_CLIENT_SECRET,
environment.GITEA_BASE_URL,
environment.GITEA_OAUTH_CLIENT_ID,
environment.GITEA_OAUTH_CLIENT_SECRET,
environment.FORGEJO_BASE_URL,
environment.FORGEJO_OAUTH_CLIENT_ID,
environment.FORGEJO_OAUTH_CLIENT_SECRET,
];
if (oauthFields.some(Boolean) && environment.MARKETPLACE_TOKEN_ENCRYPTION_KEY.length < 32) {
context.addIssue({
code: z.ZodIssueCode.custom,
path: ['MARKETPLACE_TOKEN_ENCRYPTION_KEY'],
message: 'Bei aktivierten OAuth-Anbietern ist ein Schlüssel mit mindestens 32 Zeichen erforderlich',
});
}
for (const [provider, fields] of [
['GITHUB', [environment.GITHUB_OAUTH_CLIENT_ID, environment.GITHUB_OAUTH_CLIENT_SECRET]],
['GITEA', [environment.GITEA_BASE_URL, environment.GITEA_OAUTH_CLIENT_ID, environment.GITEA_OAUTH_CLIENT_SECRET]],
['FORGEJO', [environment.FORGEJO_BASE_URL, environment.FORGEJO_OAUTH_CLIENT_ID, environment.FORGEJO_OAUTH_CLIENT_SECRET]],
] as const) {
if (fields.some(Boolean) && fields.some((field) => !field)) {
context.addIssue({
code: z.ZodIssueCode.custom,
path: [`${provider}_OAUTH_CLIENT_ID`],
message: `OAuth-Konfiguration für ${provider} ist unvollständig`,
});
}
}
for (const [field, value] of [['GITEA_BASE_URL', environment.GITEA_BASE_URL], ['FORGEJO_BASE_URL', environment.FORGEJO_BASE_URL]] as const) {
if (!value) continue;
try {
const url = new URL(value);
if (url.protocol !== 'https:' && !(url.protocol === 'http:' && ['localhost', '127.0.0.1', '[::1]'].includes(url.hostname))) {
throw new Error('protocol');
}
if (url.username || url.password || url.search || url.hash) throw new Error('url');
} catch {
context.addIssue({
code: z.ZodIssueCode.custom,
path: [field],
message: 'Forge-URL muss HTTPS verwenden (HTTP ist nur lokal zulässig) und darf keine Zugangsdaten enthalten',
});
}
}
});
/** Lädt und validiert die Konfiguration aus den Umgebungsvariablen. */
@@ -91,6 +184,24 @@ export function loadConfiguration(): AppConfig {
runtime: {
modulesDir: path.resolve(environment.MODULES_DIR),
logsDir: path.resolve(environment.LOGS_DIR),
...(environment.MODULE_UID_BASE !== undefined
? { moduleUidBase: environment.MODULE_UID_BASE }
: {}),
},
marketplace: {
publicUrl: environment.MARKETPLACE_PUBLIC_URL.replace(/\/$/, ''),
tokenEncryptionKey: environment.MARKETPLACE_TOKEN_ENCRYPTION_KEY,
providers: {
...(environment.GITHUB_OAUTH_CLIENT_ID && environment.GITHUB_OAUTH_CLIENT_SECRET
? { github: { clientId: environment.GITHUB_OAUTH_CLIENT_ID, clientSecret: environment.GITHUB_OAUTH_CLIENT_SECRET, baseUrl: 'https://github.com' } }
: {}),
...(environment.GITEA_BASE_URL && environment.GITEA_OAUTH_CLIENT_ID && environment.GITEA_OAUTH_CLIENT_SECRET
? { gitea: { clientId: environment.GITEA_OAUTH_CLIENT_ID, clientSecret: environment.GITEA_OAUTH_CLIENT_SECRET, baseUrl: environment.GITEA_BASE_URL.replace(/\/$/, '') } }
: {}),
...(environment.FORGEJO_BASE_URL && environment.FORGEJO_OAUTH_CLIENT_ID && environment.FORGEJO_OAUTH_CLIENT_SECRET
? { forgejo: { clientId: environment.FORGEJO_OAUTH_CLIENT_ID, clientSecret: environment.FORGEJO_OAUTH_CLIENT_SECRET, baseUrl: environment.FORGEJO_BASE_URL.replace(/\/$/, '') } }
: {}),
},
},
};
}
}

View File

@@ -2,6 +2,12 @@ import { migration001CoreSchema } from './001-core-schema';
import { migration002Modules } from '../../modules/migrations/002-modules';
import { migration003ModulePermissions } from '../../modules/migrations/003-module-permissions';
import { migration004SystemSettings } from '../../settings/migrations/004-system-settings';
import { migration005ModulePortUnique } from '../../modules/migrations/005-module-port-unique';
import { migration006MarketplaceConnections } from '../../modules/migrations/006-marketplace-connections';
import { migration007MarketplaceCatalog } from '../../modules/migrations/007-marketplace-catalog';
import { migration008MarketplaceSourceBranch } from '../../modules/migrations/008-marketplace-source-branch';
import { migration009MarketplaceInstallations } from '../../modules/migrations/009-marketplace-installations';
import { migration010ModuleContainers } from '../../modules/migrations/010-module-containers';
/** Registrierte Migrationen in aufsteigender Reihenfolge. */
export const MIGRATIONS = [
@@ -9,4 +15,10 @@ export const MIGRATIONS = [
migration002Modules,
migration003ModulePermissions,
migration004SystemSettings,
];
migration005ModulePortUnique,
migration006MarketplaceConnections,
migration007MarketplaceCatalog,
migration008MarketplaceSourceBranch,
migration009MarketplaceInstallations,
migration010ModuleContainers,
];

View File

@@ -4,6 +4,7 @@ import { NestExpressApplication } from '@nestjs/platform-express';
import { DocumentBuilder, SwaggerModule } from '@nestjs/swagger';
import cookieParser from 'cookie-parser';
import helmet from 'helmet';
import type { NextFunction, Request, Response } from 'express';
import { AppModule } from './app.module';
import { AllExceptionsFilter } from './common/filters/all-exceptions.filter';
import { loadConfiguration } from './config/config.tokens';
@@ -26,6 +27,10 @@ async function bootstrap(): Promise<void> {
app.use(helmet());
app.use(cookieParser());
app.use('/api', (_request: Request, response: Response, next: NextFunction) => {
response.setHeader('Cache-Control', 'no-store');
next();
});
if (config.security.behindProxy) {
app.set('trust proxy', 1);
@@ -46,4 +51,4 @@ async function bootstrap(): Promise<void> {
logger.log(`Management-Backend läuft auf Port ${config.port}`);
}
void bootstrap();
void bootstrap();

View File

@@ -52,6 +52,8 @@ export const moduleManifestSchema = z.object({
.max(MODULE_PORT_MAX, `Port muss zwischen ${MODULE_PORT_MIN} und ${MODULE_PORT_MAX} liegen`),
healthcheck: z.string().regex(/^\/[A-Za-z0-9\-./]*$/, 'Healthcheck muss ein Pfad sein'),
apiVersion: z.literal('v1'),
composeFile: z.string().min(1).max(200).optional(),
appService: z.string().regex(/^[a-zA-Z0-9][a-zA-Z0-9_.-]{0,62}$/).optional(),
});
export type ModuleManifest = z.infer<typeof moduleManifestSchema>;
@@ -71,4 +73,6 @@ export interface ModuleRecord {
readonly enabled: boolean;
readonly createdAt: Date;
readonly updatedAt: Date;
}
readonly composeFile?: string | null;
readonly appService?: string | null;
}

View File

@@ -0,0 +1,123 @@
import { BadRequestException, Controller, Delete, Get, Param, Post, Query, Req, Res } from '@nestjs/common';
import type { Request, Response } from 'express';
import { ApiTags } from '@nestjs/swagger';
import { CurrentUser } from '../common/decorators/current-user.decorator';
import { Public } from '../common/decorators/public.decorator';
import { Roles } from '../common/decorators/roles.decorator';
import type { AuthUser } from '../users/user.types';
import { SessionService } from '../auth/session.service';
import type { AuthenticatedRequest } from '../auth/authenticated-request';
import { MarketplaceService } from './marketplace.service';
import { ModulesService } from './modules.service';
@ApiTags('Marketplace')
@Controller({ path: 'api/v1/marketplace' })
export class MarketplaceController {
constructor(
private readonly marketplaceService: MarketplaceService,
private readonly sessionService: SessionService,
private readonly modulesService: ModulesService,
) {}
@Get('providers')
@Roles('ADMIN')
providers(): Promise<Awaited<ReturnType<MarketplaceService['providers']>>> {
return this.marketplaceService.providers();
}
@Get('repositories/:provider')
@Roles('ADMIN')
repositories(@Param('provider') provider: string): ReturnType<MarketplaceService['repositories']> {
return this.marketplaceService.repositories(provider);
}
@Post('repositories/:provider/:owner/:repository/install')
@Roles('ADMIN')
async installRepository(
@Param('provider') provider: string,
@Param('owner') owner: string,
@Param('repository') repository: string,
@CurrentUser() actor: AuthUser,
@Req() request: AuthenticatedRequest,
): Promise<{ module: {
id: string; moduleId: string; name: string; slug: string; version: string; description: string;
author: string; status: string; internalPort: number; healthcheckUrl: string; enabled: boolean; createdAt: string;
} }> {
const archive = await this.marketplaceService.downloadRepositoryArchive(provider, owner, repository);
const manifest = await this.modulesService.validatePackage(archive);
const module = await this.modulesService.findByModuleId(manifest.id) ??
await this.modulesService.install(archive, actor, request.ip ?? null);
await this.marketplaceService.recordInstallation(provider, owner, repository, module.id);
return {
module: {
id: module.id,
moduleId: module.moduleId,
name: module.name,
slug: module.slug,
version: module.version,
description: module.description,
author: module.author,
status: module.status,
internalPort: module.internalPort,
healthcheckUrl: module.healthcheckUrl,
enabled: module.enabled,
createdAt: module.createdAt.toISOString(),
},
};
}
@Post('connections/:provider/start')
@Roles('ADMIN')
async startConnection(
@Param('provider') provider: string,
@CurrentUser() user: AuthUser,
@Req() request: AuthenticatedRequest,
): Promise<{ authorizationUrl: string }> {
if (!request.session?.id) throw new BadRequestException('Session konnte nicht geprüft werden');
return { authorizationUrl: await this.marketplaceService.beginConnection(provider, user.id, request.session.id) };
}
@Delete('connections/:provider')
@Roles('ADMIN')
async disconnect(@Param('provider') provider: string): Promise<{ success: true }> {
await this.marketplaceService.disconnect(provider);
return { success: true };
}
@Get('oauth/:provider/callback')
@Public()
async callback(
@Param('provider') provider: string,
@Query('code') code: string | undefined,
@Query('state') state: string | undefined,
@Query('error') error: string | undefined,
@Req() request: Request,
@Res() response: Response,
): Promise<void> {
const destination = new URL('/admin/modules', this.marketplaceService.frontendUrl());
if (error) {
destination.searchParams.set('marketplace', 'denied');
} else if (!code || !state) {
destination.searchParams.set('marketplace', 'error');
destination.searchParams.set('reason', 'callback');
} else {
try {
const token = request.cookies?.mpm_session as string | undefined;
const session = token ? await this.sessionService.findValid(token) : null;
if (!session) {
destination.searchParams.set('marketplace', 'error');
destination.searchParams.set('reason', 'session');
response.redirect(302, destination.toString());
return;
}
await this.marketplaceService.completeConnection(provider, code, state, session?.id ?? null);
destination.searchParams.set('marketplace', 'connected');
destination.searchParams.set('provider', provider);
} catch {
destination.searchParams.set('marketplace', 'error');
destination.searchParams.set('reason', 'oauth');
}
}
response.redirect(302, destination.toString());
}
}

View File

@@ -0,0 +1,411 @@
import {
BadGatewayException,
BadRequestException,
Injectable,
InternalServerErrorException,
NotFoundException,
UnauthorizedException,
} from '@nestjs/common';
import { createCipheriv, createHash, randomBytes } from 'node:crypto';
import { APP_CONFIG, type AppConfig } from '../config/config.tokens';
import { Inject } from '@nestjs/common';
import { DatabaseService } from '../database/database.service';
export const MARKETPLACE_PROVIDERS = ['github', 'gitea', 'forgejo'] as const;
export type MarketplaceProvider = (typeof MARKETPLACE_PROVIDERS)[number];
interface ProviderStatus {
provider: MarketplaceProvider;
label: string;
configured: boolean;
connected: boolean;
accountLogin: string | null;
baseUrl: string;
callbackUrl: string;
}
interface OAuthStateRow {
readonly provider: MarketplaceProvider;
readonly user_id: string;
readonly session_id: string;
readonly code_verifier: string;
}
interface ProviderIdentity {
readonly id: string | number;
readonly login?: string;
readonly username?: string;
}
const MAX_MARKETPLACE_DOWNLOAD = 10 * 1024 * 1024;
function isProvider(value: string): value is MarketplaceProvider {
return MARKETPLACE_PROVIDERS.includes(value as MarketplaceProvider);
}
function safeBaseUrl(value: string): string {
const url = new URL(value);
if (url.protocol !== 'https:' && !(url.protocol === 'http:' && ['127.0.0.1', 'localhost', '::1'].includes(url.hostname))) {
throw new BadRequestException('Forge-URL muss HTTPS verwenden (HTTP ist nur für lokale Entwicklung erlaubt)');
}
if (url.username || url.password || url.search || url.hash) {
throw new BadRequestException('Forge-URL darf keine Zugangsdaten oder URL-Parameter enthalten');
}
return url.toString().replace(/\/$/, '');
}
@Injectable()
export class MarketplaceService {
constructor(
private readonly database: DatabaseService,
@Inject(APP_CONFIG) private readonly config: AppConfig,
) {}
async providers(): Promise<ProviderStatus[]> {
const connected = await this.database.query<{ provider: MarketplaceProvider; account_login: string }>(
'SELECT provider, account_login FROM marketplace_connections',
);
const connectedByProvider = new Map(connected.rows.map((row) => [row.provider, row.account_login]));
const labels: Record<MarketplaceProvider, string> = {
github: 'GitHub', gitea: 'Gitea', forgejo: 'Forgejo',
};
return MARKETPLACE_PROVIDERS.map((provider) => {
const configured = this.config.marketplace.providers[provider];
return {
provider,
label: labels[provider],
configured: Boolean(configured),
connected: connectedByProvider.has(provider),
accountLogin: connectedByProvider.get(provider) ?? null,
baseUrl: configured?.baseUrl ?? (provider === 'github' ? 'https://github.com' : ''),
callbackUrl: this.callbackUrl(provider),
};
});
}
async beginConnection(providerParam: string, userId: string, sessionId: string): Promise<string> {
const provider = this.requireProvider(providerParam);
const providerConfig = this.config.marketplace.providers[provider];
if (!providerConfig) {
throw new BadRequestException(`${provider} ist noch nicht konfiguriert`);
}
if (this.config.marketplace.tokenEncryptionKey.length < 32) {
throw new InternalServerErrorException('MARKETPLACE_TOKEN_ENCRYPTION_KEY muss mindestens 32 Zeichen lang sein');
}
const state = randomBytes(32).toString('base64url');
const verifier = randomBytes(48).toString('base64url');
const challenge = createHash('sha256').update(verifier).digest('base64url');
const callback = this.callbackUrl(provider);
await this.database.query(
`INSERT INTO marketplace_oauth_states(state_hash, provider, user_id, session_id, code_verifier, expires_at)
VALUES ($1, $2, $3, $4, $5, now() + interval '10 minutes')`,
[this.hash(state), provider, userId, sessionId, verifier],
);
await this.database.query('DELETE FROM marketplace_oauth_states WHERE expires_at < now()');
const authorizeUrl = new URL(
provider === 'github'
? '/login/oauth/authorize'
: `${new URL(providerConfig.baseUrl).pathname.replace(/\/$/, '')}/login/oauth/authorize`,
providerConfig.baseUrl,
);
authorizeUrl.searchParams.set('client_id', providerConfig.clientId);
authorizeUrl.searchParams.set('redirect_uri', callback);
authorizeUrl.searchParams.set('response_type', 'code');
authorizeUrl.searchParams.set('state', state);
authorizeUrl.searchParams.set('code_challenge', challenge);
authorizeUrl.searchParams.set('code_challenge_method', 'S256');
authorizeUrl.searchParams.set('scope', 'read:user');
return authorizeUrl.toString();
}
async completeConnection(providerParam: string, code: string, state: string, sessionId: string | null): Promise<void> {
const provider = this.requireProvider(providerParam);
if (!code || code.length > 4096 || !state || state.length > 256 || !sessionId) {
throw new BadRequestException('OAuth-Rückgabe ist ungültig');
}
const result = await this.database.query<OAuthStateRow>(
`DELETE FROM marketplace_oauth_states
WHERE state_hash = $1 AND provider = $2 AND session_id = $3 AND expires_at > now()
RETURNING provider, user_id, session_id, code_verifier`,
[this.hash(state), provider, sessionId],
);
const oauthState = result.rows[0];
if (!oauthState) {
throw new UnauthorizedException('OAuth-Status ist ungültig oder abgelaufen. Bitte erneut verbinden.');
}
const providerConfig = this.config.marketplace.providers[provider];
if (!providerConfig) throw new BadRequestException(`${provider} ist nicht konfiguriert`);
const callback = this.callbackUrl(provider);
const token = await this.exchangeCode(provider, providerConfig, code, callback, oauthState.code_verifier);
const identity = await this.fetchIdentity(provider, providerConfig.baseUrl, token);
const login = identity.login ?? identity.username;
if (!login || identity.id === undefined || identity.id === null) {
throw new BadGatewayException('Der Forge hat keine gültige Benutzeridentität zurückgegeben');
}
const encrypted = this.encryptToken(token);
await this.database.query(
`INSERT INTO marketplace_connections
(provider, account_id, account_login, token_ciphertext, token_iv, token_tag, connected_by)
VALUES ($1, $2, $3, $4, $5, $6, $7)
ON CONFLICT (provider) DO UPDATE SET
account_id = EXCLUDED.account_id,
account_login = EXCLUDED.account_login,
token_ciphertext = EXCLUDED.token_ciphertext,
token_iv = EXCLUDED.token_iv,
token_tag = EXCLUDED.token_tag,
connected_by = EXCLUDED.connected_by,
connected_at = now()`,
[provider, String(identity.id), login, encrypted.ciphertext, encrypted.iv, encrypted.tag, oauthState.user_id],
);
}
async disconnect(providerParam: string): Promise<void> {
const provider = this.requireProvider(providerParam);
await this.database.query('DELETE FROM marketplace_connections WHERE provider = $1', [provider]);
}
async repositories(providerParam: string): Promise<Array<{ owner: string; repository: string; htmlUrl: string; description: string; defaultBranch: string; installed: boolean }>> {
const provider = this.requireProvider(providerParam);
const providerConfig = this.config.marketplace.providers[provider];
if (!providerConfig) throw new BadRequestException(`${provider} ist nicht konfiguriert`);
const connection = await this.database.query<{ account_login: string }>(
'SELECT account_login FROM marketplace_connections WHERE provider = $1', [provider],
);
const login = connection.rows[0]?.account_login;
if (!login) throw new BadRequestException(`${provider} ist nicht verbunden`);
const path = provider === 'github'
? `/users/${encodeURIComponent(login)}/repos?type=owner&sort=updated&per_page=50`
: `${new URL(providerConfig.baseUrl).pathname.replace(/\/$/, '')}/api/v1/users/${encodeURIComponent(login)}/repos?limit=50&sort=updated`;
const response = await this.forgeJson<Array<Record<string, unknown>>>(provider, providerConfig.baseUrl, null, path);
const repositories = response.filter((repo) => repo.private !== true).map((repo) => {
const owner = typeof repo.owner === 'object' && repo.owner !== null
? String((repo.owner as Record<string, unknown>).login ?? (repo.owner as Record<string, unknown>).username ?? login)
: login;
const repository = String(repo.name ?? '');
return {
owner,
repository,
htmlUrl: String(repo.html_url ?? ''),
description: String(repo.description ?? ''),
defaultBranch: String(repo.default_branch ?? 'main'),
};
}).filter((repo) => repo.repository && repo.htmlUrl);
const installed = await this.database.query<{ owner: string; repository: string }>(
'SELECT owner, repository FROM marketplace_module_installations WHERE provider = $1',
[provider],
);
const installedRepositories = new Set(installed.rows.map((row) => `${row.owner.toLowerCase()}/${row.repository.toLowerCase()}`));
return repositories.map((repo) => ({
...repo,
installed: installedRepositories.has(`${repo.owner.toLowerCase()}/${repo.repository.toLowerCase()}`),
}));
}
async recordInstallation(providerParam: string, owner: string, repository: string, moduleId: string): Promise<void> {
const provider = this.requireProvider(providerParam);
await this.database.query(
`INSERT INTO marketplace_module_installations (provider, owner, repository, module_id)
VALUES ($1, $2, $3, $4)
ON CONFLICT (provider, owner, repository) DO UPDATE SET module_id = EXCLUDED.module_id`,
[provider, owner, repository, moduleId],
);
}
async downloadRepositoryArchive(providerParam: string, owner: string, repository: string): Promise<Buffer> {
const provider = this.requireProvider(providerParam);
if (![owner, repository].every((part) => /^[A-Za-z0-9_.-]{1,100}$/.test(part))) {
throw new BadRequestException('Repository-Angabe ist ungueltig');
}
const providerConfig = this.config.marketplace.providers[provider];
if (!providerConfig) throw new BadRequestException('Forge-Anbieter ist nicht konfiguriert');
const repo = await this.forgeJson<Record<string, unknown>>(
provider,
providerConfig.baseUrl,
null,
this.repositoryApiPath(provider, owner, repository),
);
if (repo.private === true) throw new BadRequestException('Private Repositories werden aktuell nicht unterstuetzt');
const defaultBranch = String(repo.default_branch ?? 'main');
if (!defaultBranch || defaultBranch.length > 200) throw new BadRequestException('Standard-Branch ist ungueltig');
const archivePath = provider === 'github'
? '/repos/' + encodeURIComponent(owner) + '/' + encodeURIComponent(repository) + '/zipball/' + encodeURIComponent(defaultBranch)
: new URL(providerConfig.baseUrl).pathname.replace(/[/]$/, '') + '/api/v1/repos/' + encodeURIComponent(owner) + '/' + encodeURIComponent(repository) + '/archive/' + encodeURIComponent(defaultBranch) + '.zip';
const archiveUrl = provider === 'github'
? new URL(archivePath, 'https://api.github.com').toString()
: new URL(archivePath, providerConfig.baseUrl).toString();
const allowedHosts = this.downloadHosts(providerConfig.baseUrl, provider);
const headers: Record<string, string> = provider === 'github' ? { 'User-Agent': 'MPM-Module-Marketplace' } : {};
const archive = await this.downloadBounded(archiveUrl, headers, allowedHosts, MAX_MARKETPLACE_DOWNLOAD);
return this.normalizeRepositoryArchive(archive);
}
private async normalizeRepositoryArchive(archive: Buffer): Promise<Buffer> {
const AdmZip = (await import('adm-zip')).default;
const input = new AdmZip(archive);
const entries = input.getEntries();
const files = entries.filter((entry) => !entry.isDirectory);
if (!files.length || files.length > 2000) throw new BadRequestException('Repository-Archiv enthaelt keine gueltigen Moduldateien');
const firstPath = files[0].entryName;
const firstSlash = firstPath.indexOf('/');
const candidateRoot = firstSlash > 0 ? firstPath.slice(0, firstSlash) : '';
const hasRootManifest = files.some((entry) => entry.entryName === 'module.json');
const root = hasRootManifest ? '' : candidateRoot;
if (!root && !hasRootManifest) throw new BadRequestException('Repository-Archiv muss module.json im Stammverzeichnis enthalten');
let totalUncompressed = 0;
for (const entry of files) {
const name = entry.entryName;
if (name.includes(String.fromCharCode(92)) || name.startsWith('/') || name.split('/').includes('..')) throw new BadRequestException('Unsicherer Pfad im Repository-Archiv');
if (root && !name.startsWith(root + '/')) throw new BadRequestException('Repository-Archiv hat mehrere Stammverzeichnisse');
const unixType = (entry.header.attr >>> 16) & 0xf000;
if (unixType === 0xa000) throw new BadRequestException('Symlinks sind in Marketplace-Modulen nicht erlaubt');
totalUncompressed += entry.header.size;
if (totalUncompressed > 50 * 1024 * 1024) throw new BadRequestException('Repository-Archiv ist entpackt zu gross');
}
const output = new AdmZip();
for (const entry of files) {
const relative = root ? entry.entryName.slice(root.length + 1) : entry.entryName;
if (relative) output.addFile(relative, entry.getData());
}
const normalized = output.toBuffer();
if (normalized.length > MAX_MARKETPLACE_DOWNLOAD) throw new BadRequestException('Modul-Paket ueberschreitet 10 MB');
return normalized;
}
private repositoryApiPath(provider: MarketplaceProvider, owner: string, repository: string): string {
return provider === 'github'
? `/repos/${encodeURIComponent(owner)}/${encodeURIComponent(repository)}`
: `${new URL(this.config.marketplace.providers[provider]!.baseUrl).pathname.replace(/\/$/, '')}/api/v1/repos/${encodeURIComponent(owner)}/${encodeURIComponent(repository)}`;
}
private async forgeJson<T>(provider: MarketplaceProvider, baseUrl: string, token: string | null, path: string): Promise<T> {
const url = provider === 'github' ? new URL(path, 'https://api.github.com') : new URL(path, baseUrl);
const response = await fetch(url, {
headers: { Accept: 'application/json', ...this.providerHeaders(provider, token) },
signal: AbortSignal.timeout(15_000),
}).catch(() => { throw new BadGatewayException('Forge-Katalog konnte nicht geladen werden'); });
if (!response.ok) throw new BadGatewayException('Forge-Katalog konnte nicht geladen werden');
return response.json() as Promise<T>;
}
private providerHeaders(provider: MarketplaceProvider, token: string | null): Record<string, string> {
return {
...(token ? { Authorization: `Bearer ${token}` } : {}),
...(provider === 'github' ? { 'X-GitHub-Api-Version': '2022-11-28', 'User-Agent': 'MPM-Module-Marketplace' } : {}),
};
}
private downloadHosts(baseUrl: string, provider: MarketplaceProvider): string[] {
const host = new URL(baseUrl).hostname.toLowerCase();
return provider === 'github' ? [host, 'api.github.com', 'codeload.github.com'] : [host];
}
private async downloadBounded(urlValue: string, headers: Record<string, string>, allowedHosts: string[], maxBytes: number): Promise<Buffer> {
let url = new URL(urlValue);
for (let redirects = 0; redirects <= 3; redirects += 1) {
if (url.protocol !== 'https:' || !allowedHosts.includes(url.hostname.toLowerCase())) throw new BadRequestException('Forge hat eine nicht vertrauenswürdige Download-Adresse geliefert');
const response = await fetch(url, { headers, redirect: 'manual', signal: AbortSignal.timeout(30_000) });
if ([301, 302, 303, 307, 308].includes(response.status)) {
const location = response.headers.get('location');
if (!location || redirects === 3) throw new BadGatewayException('Forge-Download konnte nicht aufgelöst werden');
url = new URL(location, url);
continue;
}
if (!response.ok || !response.body) throw new BadGatewayException('Forge-Download ist fehlgeschlagen');
const size = Number(response.headers.get('content-length') ?? 0);
if (size > maxBytes) throw new BadRequestException('Forge-Paket überschreitet die erlaubte Größe');
const reader = response.body.getReader();
const chunks: Buffer[] = [];
let total = 0;
while (true) {
const chunk = await reader.read();
if (chunk.done) break;
total += chunk.value.byteLength;
if (total > maxBytes) { await reader.cancel(); throw new BadRequestException('Forge-Download überschreitet die erlaubte Größe'); }
chunks.push(Buffer.from(chunk.value));
}
return Buffer.concat(chunks, total);
}
throw new BadGatewayException('Forge-Download konnte nicht aufgelöst werden');
}
callbackUrl(provider: MarketplaceProvider): string {
return `${this.config.marketplace.publicUrl}/api/v1/marketplace/oauth/${provider}/callback`;
}
frontendUrl(): string {
return this.config.marketplace.publicUrl;
}
private requireProvider(provider: string): MarketplaceProvider {
if (!isProvider(provider)) throw new NotFoundException('Unbekannter Forge-Anbieter');
return provider;
}
private hash(value: string | Buffer): string {
return createHash('sha256').update(value).digest('hex');
}
private encryptToken(token: string): { ciphertext: string; iv: string; tag: string } {
const key = createHash('sha256').update(this.config.marketplace.tokenEncryptionKey).digest();
const iv = randomBytes(12);
const cipher = createCipheriv('aes-256-gcm', key, iv);
const ciphertext = Buffer.concat([cipher.update(token, 'utf8'), cipher.final()]);
return {
ciphertext: ciphertext.toString('base64'),
iv: iv.toString('base64'),
tag: cipher.getAuthTag().toString('base64'),
};
}
private async exchangeCode(
provider: MarketplaceProvider,
config: NonNullable<AppConfig['marketplace']['providers'][MarketplaceProvider]>,
code: string,
redirectUri: string,
verifier: string,
): Promise<string> {
const tokenUrl = new URL(
provider === 'github'
? '/login/oauth/access_token'
: `${new URL(config.baseUrl).pathname.replace(/\/$/, '')}/login/oauth/access_token`,
config.baseUrl,
);
const response = await fetch(tokenUrl, {
method: 'POST',
headers: { Accept: 'application/json', 'Content-Type': 'application/json' },
body: JSON.stringify({
client_id: config.clientId,
client_secret: config.clientSecret,
code,
grant_type: 'authorization_code',
redirect_uri: redirectUri,
code_verifier: verifier,
}),
signal: AbortSignal.timeout(12_000),
}).catch(() => { throw new BadGatewayException('Token-Austausch beim Forge ist fehlgeschlagen'); });
const body = await response.json().catch(() => null) as { access_token?: unknown; error?: unknown } | null;
if (!response.ok || !body || typeof body.access_token !== 'string') {
throw new BadGatewayException('Forge hat die OAuth-Autorisierung abgelehnt');
}
return body.access_token;
}
private async fetchIdentity(provider: MarketplaceProvider, baseUrl: string, token: string): Promise<ProviderIdentity> {
const userUrl = provider === 'github'
? 'https://api.github.com/user'
: `${safeBaseUrl(baseUrl)}/api/v1/user`;
const response = await fetch(userUrl, {
headers: {
Accept: 'application/json',
Authorization: `Bearer ${token}`,
...(provider === 'github' ? { 'X-GitHub-Api-Version': '2022-11-28', 'User-Agent': 'MPM-Module-Marketplace' } : {}),
},
signal: AbortSignal.timeout(12_000),
}).catch(() => { throw new BadGatewayException('Benutzerkonto beim Forge konnte nicht gelesen werden'); });
if (!response.ok) throw new BadGatewayException('Benutzerkonto beim Forge konnte nicht gelesen werden');
return response.json() as Promise<ProviderIdentity>;
}
}

View File

@@ -0,0 +1,12 @@
import type { Migration } from '../../database/migration.types';
/** Internal ports map to isolated module UIDs and must be unique. */
export const migration005ModulePortUnique: Migration = {
id: '005-module-port-unique',
description: 'Eindeutige interne Modul-Ports sicherstellen',
up: async (client) => {
await client.query(
'CREATE UNIQUE INDEX IF NOT EXISTS idx_modules_internal_port_unique ON modules(internal_port)',
);
},
};

View File

@@ -0,0 +1,35 @@
import type { Migration } from '../../database/migration.types';
export const migration006MarketplaceConnections: Migration = {
id: '006-marketplace-connections',
description: 'OAuth-Verbindungen für Modulquellen speichern',
up: async (client) => {
await client.query(`
CREATE TABLE marketplace_oauth_states (
state_hash TEXT PRIMARY KEY,
provider TEXT NOT NULL CHECK (provider IN ('github', 'gitea', 'forgejo')),
user_id UUID NOT NULL REFERENCES users(id) ON DELETE CASCADE,
session_id UUID NOT NULL REFERENCES sessions(id) ON DELETE CASCADE,
code_verifier TEXT NOT NULL,
expires_at TIMESTAMPTZ NOT NULL,
created_at TIMESTAMPTZ NOT NULL DEFAULT now()
)
`);
await client.query(`
CREATE INDEX idx_marketplace_oauth_states_expiry
ON marketplace_oauth_states(expires_at)
`);
await client.query(`
CREATE TABLE marketplace_connections (
provider TEXT PRIMARY KEY CHECK (provider IN ('github', 'gitea', 'forgejo')),
account_id TEXT NOT NULL,
account_login TEXT NOT NULL,
token_ciphertext TEXT NOT NULL,
token_iv TEXT NOT NULL,
token_tag TEXT NOT NULL,
connected_by UUID REFERENCES users(id) ON DELETE SET NULL,
connected_at TIMESTAMPTZ NOT NULL DEFAULT now()
)
`);
},
};

View File

@@ -0,0 +1,22 @@
import type { Migration } from '../../database/migration.types';
export const migration007MarketplaceCatalog: Migration = {
id: '007-marketplace-catalog',
description: 'Ausgewählte Repository-Quellen für den Modul-Marktplatz',
up: async (client) => {
await client.query(`
CREATE TABLE marketplace_catalog_sources (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
provider TEXT NOT NULL CHECK (provider IN ('github', 'gitea', 'forgejo')),
owner TEXT NOT NULL,
repository TEXT NOT NULL,
html_url TEXT NOT NULL,
name TEXT NOT NULL,
description TEXT NOT NULL DEFAULT '',
added_by UUID REFERENCES users(id) ON DELETE SET NULL,
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
UNIQUE (provider, owner, repository)
)
`);
},
};

View File

@@ -0,0 +1,12 @@
import type { Migration } from '../../database/migration.types';
export const migration008MarketplaceSourceBranch: Migration = {
id: '008-marketplace-source-branch',
description: 'Standard-Branch für direkt installierbare Marketplace-Quellen speichern',
up: async (client) => {
await client.query(`
ALTER TABLE marketplace_catalog_sources
ADD COLUMN default_branch TEXT NOT NULL DEFAULT 'main'
`);
},
};

View File

@@ -0,0 +1,19 @@
import type { Migration } from '../../database/migration.types';
export const migration009MarketplaceInstallations: Migration = {
id: '009-marketplace-installations',
description: 'Marketplace-Repositories installierten Modulen zuordnen',
up: async (client) => {
await client.query(`
CREATE TABLE marketplace_module_installations (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
provider TEXT NOT NULL CHECK (provider IN ('github', 'gitea', 'forgejo')),
owner TEXT NOT NULL,
repository TEXT NOT NULL,
module_id UUID NOT NULL UNIQUE REFERENCES modules(id) ON DELETE CASCADE,
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
UNIQUE (provider, owner, repository)
)
`);
},
};

View File

@@ -0,0 +1,13 @@
import type { Migration } from '../../database/migration.types';
export const migration010ModuleContainers: Migration = {
id: '010-module-containers',
description: 'Container-Compose-Konfiguration installierter Module speichern',
up: async (client) => {
await client.query(`
ALTER TABLE modules
ADD COLUMN compose_file TEXT,
ADD COLUMN app_service TEXT
`);
},
};

View File

@@ -0,0 +1,244 @@
import { BadRequestException, Injectable, Logger } from '@nestjs/common';
import { spawn } from 'node:child_process';
import { mkdir, readFile, rm, writeFile } from 'node:fs/promises';
import { tmpdir } from 'node:os';
import path from 'node:path';
import { stringify, parseDocument } from 'yaml';
import { ModuleIdentityService } from './module-identity.service';
import type { ModuleRecord } from './manifest.types';
const SAFE_SERVICE_KEYS = new Set([
'image', 'build', 'command', 'entrypoint', 'environment', 'depends_on', 'volumes',
'healthcheck', 'working_dir', 'user', 'restart', 'expose', 'networks', 'hostname',
'logging', 'mem_limit', 'cpus', 'pids_limit', 'init', 'tmpfs', 'labels',
'stop_grace_period', 'read_only', 'tty', 'stdin_open',
]);
/** Orchestriert einen isolierten Docker-Compose-Stack für jedes Modul. */
@Injectable()
export class ModuleContainerManager {
private readonly logger = new Logger('ModuleContainers');
private readonly dockerHost = process.env.MODULE_DOCKER_HOST ?? 'unix:///var/run/docker.sock';
private readonly mpmContainer = process.env.MPM_CONTAINER_NAME ?? '';
constructor(private readonly identityService: ModuleIdentityService) {}
async start(module: ModuleRecord): Promise<void> {
const { composePath, overridePath, projectName, gatewayNetwork } = await this.prepare(module);
// Recreate stopped containers and project networks before each start. This
// prevents Compose v1 from trying to reconcile stale Docker Desktop network
// defaults after a stop; named data volumes are deliberately left untouched.
await this.runCompose(module.path, projectName, composePath, overridePath, ['down', '--remove-orphans']);
if (this.mpmContainer) await this.runDocker(['network', 'disconnect', '-f', gatewayNetwork, this.mpmContainer], true);
await this.runDocker(['network', 'rm', gatewayNetwork], true);
await this.runDocker(['network', 'create', gatewayNetwork]);
await this.runCompose(module.path, projectName, composePath, overridePath, ['up', '-d', '--build', '--remove-orphans']);
if (this.mpmContainer) {
await this.runDocker(['network', 'disconnect', '-f', gatewayNetwork, this.mpmContainer], true);
await this.runDocker(['network', 'connect', gatewayNetwork, this.mpmContainer]);
}
this.logger.log(`Container-Stack für "${module.moduleId}" gestartet`);
}
async stop(module: ModuleRecord): Promise<void> {
const { composePath, overridePath, projectName, gatewayNetwork } = await this.prepare(module);
await this.runCompose(module.path, projectName, composePath, overridePath, ['stop']);
if (this.mpmContainer) await this.runDocker(['network', 'disconnect', '-f', gatewayNetwork, this.mpmContainer], true);
this.logger.log(`Container-Stack für "${module.moduleId}" gestoppt`);
}
async remove(module: ModuleRecord): Promise<void> {
const { composePath, overridePath, projectName, gatewayNetwork } = await this.prepare(module);
if (this.mpmContainer) await this.runDocker(['network', 'disconnect', '-f', gatewayNetwork, this.mpmContainer], true);
// Compose down removes every app/database container and its networks. Named
// volumes remain, so uninstalling code does not silently destroy database data.
await this.runCompose(module.path, projectName, composePath, overridePath, ['down', '--remove-orphans']);
await this.runDocker(['network', 'rm', gatewayNetwork], true);
await rm(overridePath, { force: true });
this.logger.log(`Container für "${module.moduleId}" entfernt; Datenvolumes bleiben erhalten`);
}
private async prepare(module: ModuleRecord): Promise<{
composePath: string;
overridePath: string;
projectName: string;
gatewayNetwork: string;
}> {
if (!module.composeFile || !module.appService) {
throw new BadRequestException('Dieses Modul hat keine Docker-Compose-Konfiguration');
}
const root = path.resolve(module.path);
const composePath = path.resolve(root, module.composeFile);
if (!composePath.startsWith(root + path.sep)) throw new BadRequestException('Compose-Datei liegt außerhalb des Modulpakets');
const source = await readFile(composePath, 'utf8').catch(() => {
throw new BadRequestException(`Compose-Datei "${module.composeFile}" wurde nicht gefunden`);
});
const document = parseDocument(source, { uniqueKeys: true });
if (document.errors.length) throw new BadRequestException('Compose-Datei enthält ungültiges YAML');
const compose = document.toJS() as Record<string, unknown>;
this.validateCompose(compose, module);
const projectName = `mpm-${module.moduleId}`;
const gatewayNetwork = `mpm-module-${module.moduleId}-gateway`;
// Keep generated secrets outside the package/build context so Dockerfiles
// cannot accidentally copy them into an application image.
const overridePath = path.join(tmpdir(), 'mpm-compose', `${module.moduleId}.yml`);
const override = {
version: '3.8',
services: {
[module.appService]: {
container_name: `mpm-${module.moduleId}-app`,
environment: {
PORT: String(module.internalPort),
NODE_ENV: process.env.NODE_ENV ?? 'production',
MPM_MODULE_DATA_DIR: '/var/lib/mpm-module',
MPM_MODULE_IDENTITY_KEY: this.identityService.keyForModule(module.moduleId),
},
volumes: ['mpm-runtime-data:/var/lib/mpm-module'],
networks: {
default: {},
'mpm-gateway': { aliases: [`mpm-${module.moduleId}`] },
},
security_opt: ['no-new-privileges:true'],
},
},
volumes: { 'mpm-runtime-data': {} },
networks: { 'mpm-gateway': { external: true, name: gatewayNetwork } },
};
await mkdir(path.dirname(overridePath), { recursive: true, mode: 0o700 });
await writeFile(overridePath, stringify(override), { mode: 0o600 });
return { composePath, overridePath, projectName, gatewayNetwork };
}
private validateCompose(compose: Record<string, unknown>, module: ModuleRecord): void {
if (!compose || typeof compose !== 'object' || Array.isArray(compose)) {
throw new BadRequestException('Compose-Datei muss ein YAML-Objekt enthalten');
}
const topLevel = new Set(['version', 'services', 'volumes', 'networks']);
if (Object.keys(compose).some((key) => !topLevel.has(key))) {
throw new BadRequestException('Compose darf nur services, volumes und networks enthalten');
}
const services = compose.services;
if (!services || typeof services !== 'object' || Array.isArray(services)) {
throw new BadRequestException('Compose benötigt mindestens einen Service');
}
const serviceMap = services as Record<string, unknown>;
if (!Object.hasOwn(serviceMap, module.appService!)) {
throw new BadRequestException(`Compose-Service "${module.appService}" fehlt`);
}
const definedNetworks = this.record(compose.networks, 'networks');
const definedVolumes = this.record(compose.volumes, 'volumes');
if (Object.hasOwn(definedNetworks, 'mpm-gateway') || Object.hasOwn(definedVolumes, 'mpm-runtime-data')) {
throw new BadRequestException('Compose verwendet einen für MPM reservierten Netzwerk- oder Volume-Namen');
}
for (const [name, rawService] of Object.entries(serviceMap)) {
if (!/^[a-zA-Z0-9][a-zA-Z0-9_.-]{0,62}$/.test(name) || !rawService || typeof rawService !== 'object' || Array.isArray(rawService)) {
throw new BadRequestException('Compose enthält einen ungültigen Service');
}
const service = rawService as Record<string, unknown>;
if (Object.keys(service).some((key) => !SAFE_SERVICE_KEYS.has(key))) {
throw new BadRequestException(`Compose-Service "${name}" enthält nicht erlaubte Optionen`);
}
if (service.ports !== undefined || service.privileged !== undefined || service.cap_add !== undefined ||
service.devices !== undefined || service.network_mode !== undefined || service.pid !== undefined ||
service.ipc !== undefined || service.volumes_from !== undefined || service.env_file !== undefined ||
service.container_name !== undefined || service.secrets !== undefined || service.configs !== undefined) {
throw new BadRequestException(`Compose-Service "${name}" darf keine Host- oder privilegierten Ressourcen verwenden`);
}
if (service.build !== undefined) this.validateBuild(service.build, module.path);
if (service.volumes !== undefined) this.validateVolumes(service.volumes);
service.security_opt = ['no-new-privileges:true'];
}
for (const [name, volume] of Object.entries(definedVolumes)) {
if (name === 'mpm-runtime-data' || (volume !== undefined && volume !== null &&
(!volume || typeof volume !== 'object' || Array.isArray(volume) || Object.keys(volume).length > 0))) {
throw new BadRequestException('Compose darf nur projektlokale Datenvolumes definieren');
}
}
for (const [name, network] of Object.entries(definedNetworks)) {
const networkConfig = network && typeof network === 'object' && !Array.isArray(network)
? network as Record<string, unknown>
: {};
if (name === 'mpm-gateway' || (network !== undefined && network !== null &&
(!network || typeof network !== 'object' || Array.isArray(network) ||
Object.keys(networkConfig).some((key) => !['internal', 'attachable', 'labels'].includes(key))))) {
throw new BadRequestException('Compose darf keine externen Netzwerke verwenden');
}
}
}
private record(value: unknown, label: string): Record<string, unknown> {
if (value === undefined) return {};
if (!value || typeof value !== 'object' || Array.isArray(value)) {
throw new BadRequestException(`Compose-${label} muss ein Objekt sein`);
}
return value as Record<string, unknown>;
}
private validateBuild(build: unknown, modulePath: string): void {
const context = typeof build === 'string'
? build
: build && typeof build === 'object' && !Array.isArray(build)
? String((build as Record<string, unknown>).context ?? '.')
: '';
if (!context || path.isAbsolute(context) || context.split(/[\\/]/).includes('..')) {
throw new BadRequestException('Build-Kontext muss innerhalb des Modulpakets liegen');
}
const resolved = path.resolve(modulePath, context);
if (resolved !== modulePath && !resolved.startsWith(modulePath + path.sep)) {
throw new BadRequestException('Build-Kontext liegt außerhalb des Modulpakets');
}
}
private validateVolumes(volumes: unknown): void {
if (!Array.isArray(volumes)) throw new BadRequestException('Compose-Volumes müssen als Liste angegeben werden');
for (const volume of volumes) {
if (typeof volume !== 'string') throw new BadRequestException('Compose-Volume-Angabe ist ungültig');
const parts = volume.split(':');
const [source, target, mode] = parts;
if (parts.length > 3 || !target || !target.startsWith('/') ||
(source && !/^[a-zA-Z0-9][a-zA-Z0-9_.-]{0,127}$/.test(source)) ||
(mode !== undefined && !/^(ro|rw)(,ro|,rw)?$/.test(mode))) {
throw new BadRequestException('Compose darf keine Host-Verzeichnisse mounten');
}
}
}
private runCompose(cwd: string, project: string, composePath: string, overridePath: string, args: string[]): Promise<void> {
return this.run('docker-compose', ['-p', project, '-f', composePath, '-f', overridePath, ...args], cwd);
}
private runDocker(args: string[], ignoreFailure = false): Promise<void> {
return this.run('docker', args, process.cwd(), ignoreFailure);
}
private run(command: string, args: string[], cwd: string, ignoreFailure = false): Promise<void> {
return new Promise((resolve, reject) => {
const child = spawn(command, args, {
cwd,
env: {
PATH: process.env.PATH ?? '/usr/local/sbin:/usr/local/bin:/usr/sbin:/usr/bin:/sbin:/bin',
// Do not let a root-owned /root/.docker configuration affect a child
// command started by the unprivileged backend user.
HOME: '/tmp',
DOCKER_HOST: this.dockerHost,
},
stdio: ['ignore', 'ignore', 'pipe'],
});
let stderr = '';
child.stderr.setEncoding('utf8');
child.stderr.on('data', (chunk: string) => { stderr = (stderr + chunk).slice(-2000); });
child.once('error', (error) => {
if (ignoreFailure) resolve();
else reject(new Error(`${command} konnte nicht gestartet werden: ${error.message}`));
});
child.once('close', (code) => {
if (code === 0 || ignoreFailure) resolve();
else {
this.logger.error(`${command} ${args[args.length - 1]} schlug mit Status ${code} fehl: ${stderr.trim()}`);
reject(new Error(`${command} schlug mit Status ${code} fehl`));
}
});
});
}
}

View File

@@ -0,0 +1,26 @@
import { BadRequestException } from '@nestjs/common';
import { chmod, lstat, readdir } from 'node:fs/promises';
import path from 'node:path';
/** Make module files readable to runtime users, but never writable. */
export async function secureModuleDirectory(directory: string): Promise<void> {
const rootInfo = await lstat(directory);
if (!rootInfo.isDirectory() || rootInfo.isSymbolicLink()) {
throw new BadRequestException('Modulverzeichnis muss ein echtes Verzeichnis sein');
}
await chmod(directory, 0o755);
for (const entry of await readdir(directory)) {
const entryPath = path.join(directory, entry);
const info = await lstat(entryPath);
if (info.isSymbolicLink()) {
throw new BadRequestException(`Symbolische Links sind im Modul-Paket nicht erlaubt: ${entry}`);
}
if (info.isDirectory()) {
await secureModuleDirectory(entryPath);
} else if (info.isFile()) {
await chmod(entryPath, 0o644);
} else {
throw new BadRequestException(`Nicht unterstützter Dateityp im Modul-Paket: ${entry}`);
}
}
}

View File

@@ -6,6 +6,7 @@ import { extractSessionToken } from '../auth/guards/session.guard';
import { UserRepository } from '../users/user.repository';
import { ModuleRepository } from './module.repository';
import { ModulePermissionsService } from './module-permissions.service';
import { ModuleIdentityService } from './module-identity.service';
/** Gateway-Pfad-Präfix für interne Nginx-Weiterleitung. */
const GATEWAY_PREFIX = '/api/v1/gateway/';
@@ -34,6 +35,7 @@ export class ModuleGatewayMiddleware implements NestMiddleware {
private readonly sessionService: SessionService,
private readonly userRepository: UserRepository,
private readonly permissionsService: ModulePermissionsService,
private readonly identityService: ModuleIdentityService = new ModuleIdentityService(),
) {
this.proxy = httpProxy.createProxyServer({
proxyTimeout: 30_000,
@@ -52,6 +54,21 @@ export class ModuleGatewayMiddleware implements NestMiddleware {
response.end();
}
});
// Nest/Express may already have consumed JSON request bodies before this
// middleware runs. Replay the parsed body to the module instead of leaving
// the proxied request stream empty (which makes module POST handlers hang).
this.proxy.on('proxyReq', (proxyRequest, request) => {
const contentType = request.headers['content-type']?.split(';', 1)[0].trim().toLowerCase();
const parsedBody = (request as Request).body;
if (contentType !== 'application/json' || parsedBody === undefined) {
return;
}
const body = JSON.stringify(parsedBody);
proxyRequest.setHeader('Content-Length', Buffer.byteLength(body));
proxyRequest.write(body);
});
}
async use(request: Request, response: Response, next: NextFunction): Promise<void> {
@@ -114,16 +131,19 @@ export class ModuleGatewayMiddleware implements NestMiddleware {
}
// 5. Identität sicher an das Modul übergeben (Header, nicht URL)
request.url = modulePath;
const signedIdentity = this.identityService.sign(module.moduleId, request.method, request.url, user);
request.headers['x-user-id'] = user.id;
request.headers['x-user-username'] = user.username;
request.headers['x-user-display-name'] = user.displayName;
request.headers['x-user-role'] = user.role;
request.headers['x-mpm-identity-timestamp'] = signedIdentity.timestamp;
request.headers['x-mpm-identity-signature'] = signedIdentity.signature;
// Session-Cookie niemals an das Modul weiterleiten
delete request.headers.cookie;
request.url = modulePath;
this.proxy.web(request, response, {
target: `http://127.0.0.1:${module.internalPort}`,
target: `http://${module.composeFile && module.appService ? `mpm-${module.moduleId}` : '127.0.0.1'}:${module.internalPort}`,
});
}
}
}

View File

@@ -18,7 +18,8 @@ export class ModuleHealthChecker {
/** Prüft einen Modul-Prozess über seine Healthcheck-URL. */
async check(module: ModuleRecord): Promise<ModuleHealthResult> {
const url = `http://127.0.0.1:${module.internalPort}${module.healthcheckUrl}`;
const host = module.composeFile && module.appService ? `mpm-${module.moduleId}` : '127.0.0.1';
const url = `http://${host}:${module.internalPort}${module.healthcheckUrl}`;
const start = performance.now();
try {
@@ -40,4 +41,4 @@ export class ModuleHealthChecker {
return { healthy: false, latencyMs, detail: `unreachable: ${detail}` };
}
}
}
}

View File

@@ -0,0 +1,68 @@
import { createHmac, randomBytes, timingSafeEqual } from 'node:crypto';
import { Injectable } from '@nestjs/common';
import type { IncomingHttpHeaders } from 'node:http';
import type { AuthUser } from '../users/user.types';
const TIMESTAMP_HEADER = 'x-mpm-identity-timestamp';
const SIGNATURE_HEADER = 'x-mpm-identity-signature';
const MAX_AGE_MS = 30_000;
/** Issues module-specific, request-bound identities for the local module gateway. */
@Injectable()
export class ModuleIdentityService {
private readonly masterKey = randomBytes(32);
keyForModule(moduleId: string): string {
return createHmac('sha256', this.masterKey).update(moduleId).digest('base64url');
}
sign(moduleId: string, method: string, url: string, user: AuthUser): {
timestamp: string;
signature: string;
} {
const timestamp = String(Date.now());
return {
timestamp,
signature: this.signature(this.keyForModule(moduleId), [
moduleId, method, url, timestamp, user.id, user.username, user.displayName, user.role,
]),
};
}
verify(
moduleId: string,
method: string,
url: string,
headers: IncomingHttpHeaders,
key: string,
): AuthUser | null {
const userId = this.readHeader(headers, 'x-user-id');
const username = this.readHeader(headers, 'x-user-username');
const displayName = this.readHeader(headers, 'x-user-display-name') ?? username;
const role = this.readHeader(headers, 'x-user-role');
const timestamp = this.readHeader(headers, TIMESTAMP_HEADER);
const suppliedSignature = this.readHeader(headers, SIGNATURE_HEADER);
if (!userId || !username || !displayName || !timestamp || !suppliedSignature ||
(role !== 'ADMIN' && role !== 'USER')) return null;
const timestampNumber = Number(timestamp);
if (!Number.isSafeInteger(timestampNumber) || Math.abs(Date.now() - timestampNumber) > MAX_AGE_MS) {
return null;
}
const expected = Buffer.from(this.signature(key, [
moduleId, method, url, timestamp, userId, username, displayName, role,
]), 'hex');
const supplied = Buffer.from(suppliedSignature, 'hex');
if (expected.length !== supplied.length || !timingSafeEqual(expected, supplied)) return null;
return { id: userId, username, displayName, role, email: '' };
}
private signature(key: string, fields: readonly string[]): string {
return createHmac('sha256', key).update(JSON.stringify(fields)).digest('hex');
}
private readHeader(headers: IncomingHttpHeaders, name: string): string | null {
const value = headers[name];
return typeof value === 'string' ? value : null;
}
}

View File

@@ -2,6 +2,7 @@ import { BadRequestException, Injectable, Logger } from '@nestjs/common';
import { mkdir, readFile, rm, writeFile } from 'node:fs/promises';
import path from 'node:path';
import { moduleManifestSchema, type ModuleManifest } from './manifest.types';
import { secureModuleDirectory } from './module-filesystem';
/** Maximale Größe eines Modul-Pakets (10 MB). */
const MAX_PACKAGE_SIZE_BYTES = 10 * 1024 * 1024;
@@ -64,7 +65,19 @@ export class ModuleInstaller {
});
}
return result.data;
const manifest = result.data;
if (!manifest.composeFile || !manifest.appService) {
throw new BadRequestException('module.json muss composeFile und appService für den Containerbetrieb enthalten');
}
if (manifest.composeFile.startsWith('/') || manifest.composeFile.includes('\\') ||
manifest.composeFile.split('/').includes('..')) {
throw new BadRequestException('composeFile muss ein relativer Pfad innerhalb des Modulpakets sein');
}
if (!zip.getEntry(manifest.composeFile)) {
throw new BadRequestException(`Container-Konfiguration ${manifest.composeFile} fehlt im Paket`);
}
return manifest;
}
/**
@@ -99,6 +112,7 @@ export class ModuleInstaller {
await mkdir(directory, { recursive: true });
zip.extractAllTo(resolvedDirectory, true);
await secureModuleDirectory(resolvedDirectory);
// Paket-Metadaten für spätere Diagnose speichern.
await writeFile(
@@ -127,4 +141,4 @@ export class ModuleInstaller {
return null;
}
}
}
}

View File

@@ -1,10 +1,13 @@
import { Injectable, Logger, type OnModuleDestroy } from '@nestjs/common';
import { spawn, type ChildProcess } from 'node:child_process';
import { mkdir, open } from 'node:fs/promises';
import { chmod, chown, mkdir, open } from 'node:fs/promises';
import path from 'node:path';
import { APP_CONFIG, type AppConfig } from '../config/config.tokens';
import { Inject } from '@nestjs/common';
import type { ModuleRecord } from './manifest.types';
import { MODULE_PORT_MIN, type ModuleRecord } from './manifest.types';
import { secureModuleDirectory } from './module-filesystem';
import { ModuleIdentityService } from './module-identity.service';
import { ModuleContainerManager } from './module-container-manager';
/** Laufende Modul-Prozesse im Speicher (nicht persistent). */
interface RunningProcess {
@@ -27,21 +30,63 @@ interface RunningProcess {
export class ModuleProcessManager implements OnModuleDestroy {
private readonly logger = new Logger('ModuleProcesses');
private readonly running = new Map<string, RunningProcess>();
private readonly containerModules = new Map<string, ModuleRecord>();
constructor(@Inject(APP_CONFIG) private readonly config: AppConfig) {}
constructor(
@Inject(APP_CONFIG) private readonly config: AppConfig,
private readonly identityService: ModuleIdentityService,
private readonly containerManager: ModuleContainerManager,
) {}
/** Startet einen Modul-Prozess. */
async start(module: ModuleRecord): Promise<void> {
if (this.running.has(module.moduleId)) {
if (this.running.has(module.moduleId) || this.containerModules.has(module.moduleId)) {
return;
}
if (module.composeFile && module.appService) {
await secureModuleDirectory(module.path);
await this.containerManager.start(module);
this.containerModules.set(module.moduleId, module);
return;
}
const entrypoint = path.join(module.path, 'backend', 'server.js');
const logFilePath = path.join(this.config.runtime.logsDir, `module-${module.moduleId}.log`);
await mkdir(this.config.runtime.logsDir, { recursive: true });
const logFile = await open(logFilePath, 'a');
await secureModuleDirectory(module.path);
const moduleUid =
process.platform !== 'win32' && this.config.runtime.moduleUidBase !== undefined
? this.config.runtime.moduleUidBase + module.internalPort - MODULE_PORT_MIN
: undefined;
const legacyDataDir = path.join(module.path, 'data');
await mkdir(legacyDataDir, { recursive: true, mode: 0o700 });
if (moduleUid !== undefined) await chown(legacyDataDir, moduleUid, moduleUid);
await mkdir(this.config.runtime.logsDir, { recursive: true, mode: 0o700 });
await chmod(this.config.runtime.logsDir, 0o700);
const logFile = await open(logFilePath, 'a', 0o600);
await chmod(logFilePath, 0o600);
const child = spawn(process.execPath, [entrypoint], {
// Unique UID/GID per internal port prevents modules from reading or tracing
// each other's processes. Production Compose grants only SETUID/SETGID.
// The platform backend receives SETUID/SETGID/KILL as ambient capabilities from
// supervisord. Use setpriv to change the module identity, then clear all
// inheritable/ambient capabilities before executing untrusted module code.
const moduleCommand = moduleUid !== undefined ? '/usr/bin/setpriv' : process.execPath;
const moduleArgs =
moduleUid !== undefined
? [
`--reuid=${moduleUid}`,
`--regid=${moduleUid}`,
'--clear-groups',
'--inh-caps=-all',
'--ambient-caps=-all',
'--no-new-privs',
'--',
process.execPath,
entrypoint,
]
: [entrypoint];
const child = spawn(moduleCommand, moduleArgs, {
cwd: module.path,
// Log-Datei bleibt offen: Der fd wird vom Kindprozess geerbt und
// darf erst nach Prozessende geschlossen werden.
@@ -50,6 +95,8 @@ export class ModuleProcessManager implements OnModuleDestroy {
PATH: process.env.PATH ?? '',
NODE_ENV: this.config.nodeEnv,
PORT: String(module.internalPort),
MPM_MODULE_DATA_DIR: legacyDataDir,
MPM_MODULE_IDENTITY_KEY: this.identityService.keyForModule(module.moduleId),
// Modul erhält nur seinen eigenen Kontext – keine Plattform-Secrets.
},
detached: false,
@@ -68,6 +115,12 @@ export class ModuleProcessManager implements OnModuleDestroy {
/** Stoppt einen Modul-Prozess (SIGTERM, dann SIGKILL). */
async stop(moduleId: string): Promise<void> {
const containerModule = this.containerModules.get(moduleId);
if (containerModule) {
await this.containerManager.stop(containerModule);
this.containerModules.delete(moduleId);
return;
}
const process_ = this.running.get(moduleId);
if (!process_) {
return;
@@ -93,6 +146,13 @@ export class ModuleProcessManager implements OnModuleDestroy {
this.logger.log(`Modul-Prozess "${moduleId}" gestoppt`);
}
async remove(module: ModuleRecord): Promise<void> {
if (module.composeFile && module.appService) {
await this.containerManager.remove(module);
this.containerModules.delete(module.moduleId);
}
}
/** Prüft, ob ein Modul-Prozess läuft. */
isRunning(moduleId: string): boolean {
return this.running.has(moduleId);
@@ -100,11 +160,11 @@ export class ModuleProcessManager implements OnModuleDestroy {
/** Stoppt alle Modul-Prozesse (Herunterfahren). */
async stopAll(): Promise<void> {
const moduleIds = [...this.running.keys()];
const moduleIds = [...this.running.keys(), ...this.containerModules.keys()];
await Promise.all(moduleIds.map((moduleId) => this.stop(moduleId)));
}
async onModuleDestroy(): Promise<void> {
await this.stopAll();
}
}
}

View File

@@ -17,10 +17,13 @@ interface ModuleRow {
enabled: boolean;
created_at: Date;
updated_at: Date;
compose_file: string | null;
app_service: string | null;
}
const MODULE_COLUMNS = `id, module_id, name, slug, version, description, author, path,
status, internal_port, healthcheck_url, enabled, created_at, updated_at`;
status, internal_port, healthcheck_url, enabled, created_at, updated_at,
compose_file, app_service`;
/**
* Modul-Repository (Infrastructure): Datenbankzugriffe für die Modul-Registry.
@@ -80,8 +83,9 @@ export class ModuleRepository {
async create(manifest: ModuleManifest, directory: string): Promise<ModuleRecord> {
const result = await this.database.query<ModuleRow>(
`INSERT INTO modules
(module_id, name, slug, version, description, author, path, status, internal_port, healthcheck_url)
VALUES ($1, $2, $3, $4, $5, $6, $7, 'INSTALLED', $8, $9)
(module_id, name, slug, version, description, author, path, status, internal_port, healthcheck_url,
compose_file, app_service)
VALUES ($1, $2, $3, $4, $5, $6, $7, 'INSTALLED', $8, $9, $10, $11)
RETURNING ${MODULE_COLUMNS}`,
[
manifest.id,
@@ -93,6 +97,8 @@ export class ModuleRepository {
directory,
manifest.port,
manifest.healthcheck,
manifest.composeFile ?? null,
manifest.appService ?? null,
],
);
return this.mapRow(result.rows[0]);
@@ -132,6 +138,8 @@ export class ModuleRepository {
enabled: row.enabled,
createdAt: row.created_at,
updatedAt: row.updated_at,
composeFile: row.compose_file,
appService: row.app_service,
};
}
}
}

View File

@@ -21,15 +21,21 @@ import { ModuleRepository } from './module.repository';
import { ModuleStartupRecovery } from './module-startup-recovery';
import { ModulesController } from './modules.controller';
import { ModulesService } from './modules.service';
import { ModuleIdentityService } from './module-identity.service';
import { MarketplaceController } from './marketplace.controller';
import { MarketplaceService } from './marketplace.service';
import { ModuleContainerManager } from './module-container-manager';
/** Modul-System: Installation, Lifecycle, Prozessverwaltung, Gateway. */
@Module({
imports: [ConfigModule, DatabaseModule, AuditModule],
controllers: [ModulesController, ModulePermissionsController],
controllers: [ModulesController, ModulePermissionsController, MarketplaceController],
providers: [
ModuleRepository,
ModuleInstaller,
ModuleProcessManager,
ModuleIdentityService,
ModuleContainerManager,
ModuleHealthChecker,
ModulesService,
SessionService,
@@ -39,6 +45,7 @@ import { ModulesService } from './modules.service';
ModuleStartupRecovery,
ModulePermissionRepository,
ModulePermissionsService,
MarketplaceService,
],
exports: [ModuleRepository, ModulesService, ModulePermissionsService, ModuleHealthChecker],
})
@@ -49,4 +56,4 @@ export class ModulesModule implements NestModule {
.apply(ModuleGatewayMiddleware)
.forRoutes({ path: '/api/v1/gateway/(.*)', method: RequestMethod.ALL });
}
}
}

View File

@@ -191,6 +191,7 @@ const TEST_CONFIG: AppConfig = {
},
adminSeed: { username: 'admin', email: 'admin@example.com', password: 'password-123' },
runtime: { modulesDir: '/data/modules', logsDir: '/data/logs' },
marketplace: { publicUrl: 'http://127.0.0.1:8081', tokenEncryptionKey: '', providers: {} },
};
describe('ModulesService', () => {
@@ -330,4 +331,4 @@ describe('ModulesService', () => {
expect(auditService.records.at(-1)?.action).toBe(AUDIT_ACTIONS.MODULE_REMOVED);
});
});
});
});

View File

@@ -88,6 +88,14 @@ export class ModulesService {
return module;
}
async validatePackage(packageBuffer: Buffer) {
return this.installer.validatePackage(packageBuffer);
}
async findByModuleId(moduleId: string): Promise<ModuleRecord | null> {
return this.moduleRepository.findByModuleId(moduleId);
}
/** Startet ein Modul (INSTALLED/STOPPED → STARTING → RUNNING). */
async start(id: string, actor: ActingUser, ipAddress: string | null): Promise<ModuleRecord> {
const module = await this.getById(id);
@@ -210,6 +218,8 @@ export class ModulesService {
await this.stop(id, actor, ipAddress);
}
await this.processManager.remove(module);
await this.moduleRepository.delete(id);
await this.installer.remove(this.config.runtime.modulesDir, module.moduleId);
@@ -249,4 +259,4 @@ export class ModulesService {
ipAddress,
});
}
}
}

View File

@@ -1,4 +1,4 @@
import { Injectable } from '@nestjs/common';
import { BadRequestException, Injectable } from '@nestjs/common';
import { DatabaseService } from '../database/database.service';
import { PasswordHasher } from './password-hasher';
import type { CreateUserDto, RoleName, UpdateUserDto, UserRecord } from './user.types';
@@ -108,17 +108,41 @@ export class UserRepository {
setClauses.push('updated_at = now()');
params.push(id);
const result = await this.database.query<UserRow>(
`UPDATE users
SET ${setClauses.join(', ')}
WHERE id = $${parameterIndex}
RETURNING id, username, email, password_hash, display_name,
(SELECT name FROM roles WHERE id = role_id) AS role_name,
is_active, failed_login_attempts, locked_until,
last_login_at, created_at, updated_at`,
params,
);
return this.mapRow(result.rows[0]);
return this.database.transaction(async (client) => {
await client.query('SELECT pg_advisory_xact_lock(727273)');
if (changes.role === 'USER' || changes.isActive === false) {
const current = await client.query<{ role_name: RoleName; is_active: boolean }>(
`SELECT r.name AS role_name, u.is_active
FROM users u JOIN roles r ON r.id = u.role_id
WHERE u.id = $1 FOR UPDATE OF u`,
[id],
);
if (current.rows[0]?.role_name === 'ADMIN' && current.rows[0].is_active) {
const count = await client.query<{ count: number }>(
`SELECT count(*)::int AS count
FROM users u JOIN roles r ON r.id = u.role_id
WHERE r.name = 'ADMIN' AND u.is_active`,
);
if ((count.rows[0]?.count ?? 0) <= 1) {
throw new BadRequestException(
'Der letzte aktive Administrator kann nicht herabgestuft oder deaktiviert werden',
);
}
}
}
const result = await client.query<UserRow>(
`UPDATE users
SET ${setClauses.join(', ')}
WHERE id = $${parameterIndex}
RETURNING id, username, email, password_hash, display_name,
(SELECT name FROM roles WHERE id = role_id) AS role_name,
is_active, failed_login_attempts, locked_until,
last_login_at, created_at, updated_at`,
params,
);
if (!result.rows[0]) throw new BadRequestException('Benutzer nicht gefunden');
return this.mapRow(result.rows[0]);
});
}
/** Setzt einen neuen Passwort-Hash. */
@@ -138,7 +162,26 @@ export class UserRepository {
}
async delete(id: string): Promise<void> {
await this.database.query('DELETE FROM users WHERE id = $1', [id]);
await this.database.transaction(async (client) => {
await client.query('SELECT pg_advisory_xact_lock(727273)');
const current = await client.query<{ role_name: RoleName; is_active: boolean }>(
`SELECT r.name AS role_name, u.is_active
FROM users u JOIN roles r ON r.id = u.role_id
WHERE u.id = $1 FOR UPDATE OF u`,
[id],
);
if (current.rows[0]?.role_name === 'ADMIN' && current.rows[0].is_active) {
const count = await client.query<{ count: number }>(
`SELECT count(*)::int AS count
FROM users u JOIN roles r ON r.id = u.role_id
WHERE r.name = 'ADMIN' AND u.is_active`,
);
if ((count.rows[0]?.count ?? 0) <= 1) {
throw new BadRequestException('Der letzte aktive Administrator kann nicht gelöscht werden');
}
}
await client.query('DELETE FROM users WHERE id = $1', [id]);
});
}
/** Anzahl aktiver Administratoren (Schutz vor Verlust des letzten Admins). */
@@ -163,21 +206,13 @@ export class UserRepository {
);
}
async updateLoginFailure(
userId: string,
attempts: number,
shouldLock: boolean,
lockoutMinutes: number,
): Promise<void> {
async updateLoginFailure(userId: string): Promise<void> {
await this.database.query(
`UPDATE users
SET failed_login_attempts = $2,
locked_until = CASE WHEN $3::boolean
THEN now() + make_interval(mins => $4::int)
ELSE locked_until END,
SET failed_login_attempts = failed_login_attempts + 1,
updated_at = now()
WHERE id = $1`,
[userId, attempts, shouldLock, lockoutMinutes],
[userId],
);
}
@@ -197,4 +232,4 @@ export class UserRepository {
updatedAt: row.updated_at,
};
}
}
}