added new features 3.0

This commit is contained in:
leon
2026-10-10 15:51:30 +02:00
parent 84d8697f23
commit 21e2ccefb2
41 changed files with 978 additions and 226 deletions

View File

@@ -1,4 +1,4 @@
import { Body, Controller, Get, HttpCode, Inject, Post, Req, Res, UseGuards } from '@nestjs/common';
import { Body, Controller, ForbiddenException, Get, HttpCode, Inject, NotFoundException, Param, Post, Query, Req, Res, UseGuards } from '@nestjs/common';
import type { Request, Response } from 'express';
import { APP_CONFIG, type AppConfig } from '../config/config.tokens';
import { CurrentUser } from '../common/decorators/current-user.decorator';
@@ -8,6 +8,8 @@ import type { AuthenticatedRequest } from './authenticated-request';
import { AuthService } from './auth.service';
import { CsrfGuard } from './guards/csrf.guard';
import { SessionGuard } from './guards/session.guard';
import { SessionService } from './session.service';
import { requireSameOrigin } from './request-origin';
import { loginSchema, type LoginDto } from '../users/user.types';
import type { AuthUser } from '../users/user.types';
@@ -38,6 +40,7 @@ function toAuthUserResponse(user: AuthUser): AuthUserResponse {
export class AuthController {
constructor(
private readonly authService: AuthService,
private readonly sessionService: SessionService,
@Inject(APP_CONFIG) private readonly config: AppConfig,
) {}
@@ -49,6 +52,7 @@ export class AuthController {
@Req() request: AuthenticatedRequest & Request,
@Res({ passthrough: true }) response: Response,
): Promise<{ user: AuthUserResponse }> {
requireSameOrigin(request, this.config.marketplace.publicUrl);
const result = await this.authService.login({
username: body.username,
password: body.password,
@@ -57,14 +61,14 @@ export class AuthController {
// 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, {
response.cookie(this.config.security.cookieSecure ? '__Host-mpm_session' : 'mpm_session', result.sessionToken, {
httpOnly: true,
secure: this.config.security.cookieSecure,
sameSite: 'lax',
path: '/',
maxAge: cookieMaxAgeMs,
});
response.cookie('mpm_csrf', result.session.csrfToken, {
response.cookie(this.config.security.cookieSecure ? '__Host-mpm_csrf' : 'mpm_csrf', result.session.csrfToken, {
httpOnly: false,
secure: this.config.security.cookieSecure,
sameSite: 'lax',
@@ -75,6 +79,63 @@ export class AuthController {
return { user: toAuthUserResponse(result.user) };
}
/** Navigationspunkt auf dem Plattformhost für einen eigenen Modul-Origin. */
@UseGuards(SessionGuard)
@Get('module-open/:slug')
async openModule(
@Param('slug') slug: string,
@CurrentUser() user: AuthUser,
@Req() request: AuthenticatedRequest & Request,
@Res() response: Response,
): Promise<void> {
if (!/^[a-z0-9][a-z0-9-]{2,100}$/.test(slug) || !request.session) {
throw new NotFoundException('Modul nicht gefunden');
}
const ticket = await this.sessionService.createModuleAccessTicket(request.session.id, user.id, slug);
response.setHeader('Cache-Control', 'no-store');
response.setHeader('Referrer-Policy', 'no-referrer');
response.redirect(303, `${this.config.modulePublicOrigin}/__mpm_module_handoff?ticket=${encodeURIComponent(ticket)}`);
}
/** Einmaliger Cookie-Übergang auf dem separaten Modulhost. */
@Public()
@Get('module-handoff')
async moduleHandoff(
@Query('ticket') ticket: string,
@Req() request: Request,
@Res() response: Response,
): Promise<void> {
if (request.headers.host !== new URL(this.config.modulePublicOrigin).host) {
throw new ForbiddenException('Ungültiger Modul-Host');
}
if (typeof ticket !== 'string' || !/^[A-Za-z0-9_-]{43}$/.test(ticket)) {
throw new ForbiddenException('Ungültiges Modul-Ticket');
}
const exchanged = await this.sessionService.exchangeModuleAccessTicket(
ticket,
this.config.security.sessionTtlMinutes,
);
if (!exchanged) {
throw new ForbiddenException('Modul-Ticket ist abgelaufen oder bereits verwendet');
}
// Ein Browser, der diesen Host früher als Plattformhost genutzt hat,
// darf keine alten Plattform-Cookies an Modul-JavaScript weitergeben.
response.clearCookie('mpm_session', { path: '/' });
response.clearCookie('mpm_csrf', { path: '/' });
response.clearCookie('__Host-mpm_session', { path: '/', secure: true });
response.clearCookie('__Host-mpm_csrf', { path: '/', secure: true });
response.cookie('mpm_module_session', exchanged.token, {
httpOnly: true,
secure: this.config.security.cookieSecure,
sameSite: 'lax',
path: `/${exchanged.moduleSlug}`,
maxAge: this.config.security.sessionTtlMinutes * 60 * 1000,
});
response.setHeader('Cache-Control', 'no-store');
response.setHeader('Referrer-Policy', 'no-referrer');
response.redirect(303, `${this.config.modulePublicOrigin}/${exchanged.moduleSlug}`);
}
@UseGuards(SessionGuard, CsrfGuard)
@Post('logout')
@HttpCode(200)
@@ -88,6 +149,8 @@ export class AuthController {
}
response.clearCookie('mpm_session', { path: '/' });
response.clearCookie('mpm_csrf', { path: '/' });
response.clearCookie('__Host-mpm_session', { path: '/', secure: true });
response.clearCookie('__Host-mpm_csrf', { path: '/', secure: true });
return { success: true };
}

View File

@@ -27,6 +27,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', moduleConfigurationEncryptionKey: '' },
marketplace: { publicUrl: 'http://127.0.0.1:8081', tokenEncryptionKey: '', providers: {} },
modulePublicOrigin: 'http://localhost:8081',
};
}
@@ -80,7 +81,7 @@ class MockSessionService {
},
};
async create(): Promise<{ token: string; data: SessionData }> {
async createForVerifiedPassword(): Promise<{ token: string; data: SessionData }> {
return this.createResult;
}

View File

@@ -25,6 +25,10 @@ const INVALID_CREDENTIALS_MESSAGE = 'Benutzername oder Passwort ist falsch';
*/
@Injectable()
export class AuthService {
// Ein echter Argon2-Hash für unbekannte Benutzernamen hält den teuren
// Verifikationsschritt in beiden Fehlpfaden vergleichbar.
private readonly dummyPasswordHash: Promise<string>;
constructor(
private readonly userRepository: UserRepository,
private readonly passwordHasher: PasswordHasher,
@@ -32,7 +36,9 @@ export class AuthService {
private readonly rateLimiter: RateLimiterService,
private readonly auditService: AuditService,
@Inject(APP_CONFIG) private readonly config: AppConfig,
) {}
) {
this.dummyPasswordHash = this.passwordHasher.hash('mpm-invalid-user-placeholder');
}
async login(input: {
username: string;
@@ -62,6 +68,7 @@ export class AuthService {
// Gleiches Verhalten für "unbekannter Benutzer" und "falsches Passwort"
// (keine User-Enumeration).
if (!user) {
await this.passwordHasher.verify(await this.dummyPasswordHash, input.password);
await this.auditService.record({
userId: null,
username: input.username,
@@ -96,12 +103,23 @@ export class AuthService {
throw new UnauthorizedException(INVALID_CREDENTIALS_MESSAGE);
}
await this.userRepository.updateLoginSuccess(user.id);
const { token, data } = await this.sessionService.create(
const createdSession = await this.sessionService.createForVerifiedPassword(
user.id,
user.passwordHash,
security.sessionTtlMinutes,
);
if (!createdSession) {
await this.auditService.record({
userId: user.id,
username: user.username,
action: 'LOGIN_FAILED',
details: { reason: 'PASSWORD_CHANGED_DURING_LOGIN' },
ipAddress: input.ipAddress,
});
throw new UnauthorizedException(INVALID_CREDENTIALS_MESSAGE);
}
const { token, data } = createdSession;
await this.userRepository.updateLoginSuccess(user.id);
await this.auditService.record({
userId: user.id,

View File

@@ -4,6 +4,7 @@ import { UserRepository } from '../../users/user.repository';
import type { UserRecord } from '../../users/user.types';
import { SessionService } from '../session.service';
import { SessionGuard } from './session.guard';
import type { AppConfig } from '../../config/config.tokens';
/** Erzeugt einen Benutzer-Datensatz für Tests. */
function createUserRecord(overrides: Partial<UserRecord> = {}): UserRecord {
@@ -74,6 +75,7 @@ describe('SessionGuard', () => {
sessionService as unknown as SessionService,
userRepository as unknown as UserRepository,
reflector,
{ security: { cookieSecure: false } } as AppConfig,
);
});
@@ -137,4 +139,4 @@ describe('SessionGuard', () => {
});
expect(request.session).toEqual({ id: 'session-1', csrfToken: 'csrf-token' });
});
});
});

View File

@@ -1,5 +1,6 @@
import { type CanActivate, type ExecutionContext, Injectable, UnauthorizedException } from '@nestjs/common';
import { Inject, type CanActivate, type ExecutionContext, Injectable, UnauthorizedException } from '@nestjs/common';
import { Reflector } from '@nestjs/core';
import { APP_CONFIG, type AppConfig } from '../../config/config.tokens';
import { IS_PUBLIC_KEY } from '../../common/decorators/public.decorator';
import { UserRepository } from '../../users/user.repository';
import type { AuthenticatedRequest } from '../authenticated-request';
@@ -11,7 +12,7 @@ interface RequestWithCookieHeader {
}
/** Extrahiert das Session-Cookie aus einem Request. */
export function extractSessionToken(request: RequestWithCookieHeader): string | null {
export function extractSessionToken(request: RequestWithCookieHeader, secureCookie = false): string | null {
const cookieHeader = request.headers.cookie;
if (!cookieHeader) {
return null;
@@ -19,7 +20,7 @@ export function extractSessionToken(request: RequestWithCookieHeader): string |
let sessionToken: string | null = null;
for (const part of cookieHeader.split(';')) {
const [name, ...value] = part.trim().split('=');
if (name === 'mpm_session') {
if (name === (secureCookie ? '__Host-mpm_session' : 'mpm_session')) {
// 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('='));
@@ -40,6 +41,7 @@ export class SessionGuard implements CanActivate {
private readonly sessionService: SessionService,
private readonly userRepository: UserRepository,
private readonly reflector: Reflector,
@Inject(APP_CONFIG) private readonly config: AppConfig,
) {}
async canActivate(context: ExecutionContext): Promise<boolean> {
@@ -52,7 +54,7 @@ export class SessionGuard implements CanActivate {
}
const request = context.switchToHttp().getRequest<AuthenticatedRequest>();
const token = extractSessionToken(request);
const token = extractSessionToken(request, this.config.security.cookieSecure);
if (!token) {
throw new UnauthorizedException('Nicht authentifiziert');
}

View File

@@ -0,0 +1,23 @@
import { ForbiddenException } from '@nestjs/common';
import type { Request } from 'express';
/**
* Browsers senden bei POST/PUT/PATCH/DELETE einen Origin-Header. Der Vergleich
* verhindert auch Anfragen von einer anderen Subdomain derselben Site.
*/
export function requireSameOrigin(request: Request, canonicalOrigin: string): void {
const origin = request.headers.origin;
const host = request.headers.host;
if (typeof origin !== 'string' || !host) {
throw new ForbiddenException('Ungültiger Request-Ursprung');
}
const canonical = new URL(canonicalOrigin);
// Hinter einem TLS-Reverse-Proxy sieht NestJS eventuell nur HTTP. Für den
// konfigurierten öffentlichen Host ist die veröffentlichte URL maßgeblich.
const protocol = canonical.host === host ? canonical.protocol : `${request.protocol}:`;
const expected = `${protocol}//${host}`;
if (origin !== expected) {
throw new ForbiddenException('Ungültiger Request-Ursprung');
}
}

View File

@@ -1,4 +1,4 @@
import { Injectable } from '@nestjs/common';
import { Injectable, UnauthorizedException } from '@nestjs/common';
import { createHash, randomBytes, timingSafeEqual } from 'node:crypto';
import { DatabaseService } from '../database/database.service';
@@ -17,6 +17,12 @@ interface SessionRow {
expires_at: Date;
}
interface ModuleAccessRow {
user_id: string;
module_slug: string;
platform_session_id: string;
}
/**
* Serverseitige Session-Verwaltung (Infrastructure):
* - 256-Bit-Zufalls-Token, in der DB wird nur der SHA-256-Hash gespeichert
@@ -43,6 +49,121 @@ export class SessionService {
return { token, data: this.mapRow(result.rows[0]) };
}
/**
* Erstellt die Session nur, wenn der gerade verifizierte Passwort-Hash noch
* aktuell ist. Die Zeilensperre serialisiert diesen Schritt mit Resets:
* entweder wird die Session vom Reset gelöscht oder der alte Hash abgewiesen.
*/
async createForVerifiedPassword(
userId: string,
verifiedPasswordHash: string,
ttlMinutes: number,
): Promise<{ token: string; data: SessionData } | null> {
const token = randomBytes(32).toString('base64url');
const csrfToken = randomBytes(32).toString('base64url');
return this.database.transaction(async (client) => {
const user = await client.query<{ password_hash: string; is_active: boolean }>(
'SELECT password_hash, is_active FROM users WHERE id = $1 FOR UPDATE',
[userId],
);
if (!user.rows[0]?.is_active || user.rows[0].password_hash !== verifiedPasswordHash) {
return null;
}
const result = await client.query<SessionRow>(
`INSERT INTO sessions (user_id, token_hash, csrf_token, expires_at)
VALUES ($1, $2, $3, now() + make_interval(mins => $4::int))
RETURNING id, user_id, csrf_token, expires_at`,
[userId, this.hashToken(token), csrfToken, ttlMinutes],
);
return { token, data: this.mapRow(result.rows[0]) };
});
}
/** Einmal-Ticket, das an die noch gültige Plattform-Session gebunden ist. */
async createModuleAccessTicket(
platformSessionId: string,
userId: string,
moduleSlug: string,
): Promise<string> {
const ticket = randomBytes(32).toString('base64url');
await this.database.transaction(async (client) => {
const platformSession = await client.query(
'SELECT id FROM sessions WHERE id = $1 AND user_id = $2 AND expires_at > now() FOR SHARE',
[platformSessionId, userId],
);
if (!platformSession.rows[0]) throw new UnauthorizedException('Nicht authentifiziert');
await client.query('DELETE FROM module_access_tickets WHERE expires_at <= now()');
await client.query('DELETE FROM module_sessions WHERE expires_at <= now()');
await client.query(
`INSERT INTO module_access_tickets
(ticket_hash, platform_session_id, user_id, module_slug, expires_at)
VALUES ($1, $2, $3, $4, now() + interval '60 seconds')`,
[this.hashToken(ticket), platformSessionId, userId, moduleSlug],
);
});
return ticket;
}
/** Verbraucht das Ticket atomar und legt nur für den angegebenen Modulpfad eine Session an. */
async exchangeModuleAccessTicket(
ticket: string,
ttlMinutes: number,
): Promise<{ token: string; userId: string; moduleSlug: string } | null> {
const token = randomBytes(32).toString('base64url');
return this.database.transaction(async (client) => {
const ticketHash = this.hashToken(ticket);
const ticketRow = await client.query<ModuleAccessRow>(
`SELECT user_id, module_slug, platform_session_id FROM module_access_tickets
WHERE ticket_hash = $1 AND expires_at > now()`,
[ticketHash],
);
if (!ticketRow.rows[0]) return null;
// Dieselbe Zeilensperre wie bei Passwortwechsel und Login verhindert,
// dass ein Reset nach der Prüfung eine neue Modulsession überlebt.
const activeUser = await client.query<{ id: string }>(
'SELECT id FROM users WHERE id = $1 AND is_active FOR UPDATE',
[ticketRow.rows[0].user_id],
);
if (!activeUser.rows[0]) return null;
const consumed = await client.query<ModuleAccessRow>(
`DELETE FROM module_access_tickets t
WHERE t.ticket_hash = $1 AND t.expires_at > now()
AND EXISTS (
SELECT 1 FROM sessions s
WHERE s.id = t.platform_session_id AND s.expires_at > now()
)
RETURNING t.user_id, t.module_slug, t.platform_session_id`,
[ticketHash],
);
if (!consumed.rows[0]) return null;
await client.query(
`INSERT INTO module_sessions (token_hash, platform_session_id, user_id, module_slug, expires_at)
VALUES ($1, $2, $3, $4, now() + make_interval(mins => $5::int))`,
[this.hashToken(token), consumed.rows[0].platform_session_id, consumed.rows[0].user_id, consumed.rows[0].module_slug, ttlMinutes],
);
return {
token,
userId: consumed.rows[0].user_id,
moduleSlug: consumed.rows[0].module_slug,
};
});
}
/** Modul-Cookies gelten ausschließlich für das ausgestellte Modul. */
async findValidModuleSession(token: string, moduleSlug: string): Promise<string | null> {
const result = await this.database.query<{ user_id: string }>(
`SELECT m.user_id FROM module_sessions m
JOIN sessions s ON s.id = m.platform_session_id
WHERE m.token_hash = $1 AND m.module_slug = $2
AND m.expires_at > now() AND s.expires_at > now()`,
[this.hashToken(token), moduleSlug],
);
return result.rows[0]?.user_id ?? null;
}
/** Findet eine gültige Session anhand des Klartext-Tokens. */
async findValid(token: string): Promise<SessionData | null> {
const result = await this.database.query<SessionRow>(
@@ -75,6 +196,8 @@ export class SessionService {
/** Löscht alle Sessions eines Benutzers (Deaktivierung, Passwort-Reset). */
async deleteAllForUser(userId: string, exceptSessionId?: string): Promise<void> {
await this.database.query('DELETE FROM module_sessions WHERE user_id = $1', [userId]);
await this.database.query('DELETE FROM module_access_tickets WHERE user_id = $1', [userId]);
if (exceptSessionId) {
await this.database.query('DELETE FROM sessions WHERE user_id = $1 AND id <> $2', [
userId,
@@ -87,6 +210,8 @@ export class SessionService {
/** Löscht alle abgelaufenen Sessions (Aufräumjob, später via Cron). */
async deleteExpired(): Promise<void> {
await this.database.query('DELETE FROM module_access_tickets WHERE expires_at <= now()');
await this.database.query('DELETE FROM module_sessions WHERE expires_at <= now()');
await this.database.query('DELETE FROM sessions WHERE expires_at <= now()');
}
@@ -112,4 +237,4 @@ export class SessionService {
expiresAt: row.expires_at,
};
}
}
}

View File

@@ -62,6 +62,22 @@ export interface AppConfig {
readonly adminSeed: AdminSeedConfig;
readonly runtime: RuntimeConfig;
readonly marketplace: MarketplaceConfig;
/** Eigener Browser-Host für Moduloberflächen (Host muss von MPM abweichen). */
readonly modulePublicOrigin: string;
}
function moduleOriginFor(publicUrl: string, configuredOrigin: string): string {
if (configuredOrigin) return configuredOrigin;
const platform = new URL(publicUrl);
const modules = new URL(platform.origin);
if (platform.hostname === 'localhost') {
modules.hostname = '127.0.0.1';
} else if (platform.hostname === '127.0.0.1' || platform.hostname === '[::1]') {
modules.hostname = 'localhost';
} else {
modules.hostname = `modules.${platform.hostname}`;
}
return modules.origin;
}
const booleanFromString = z
@@ -89,6 +105,7 @@ const environmentSchema = z.object({
MODULE_CONFIG_ENCRYPTION_KEY: z.string().default(''),
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'),
MODULE_PUBLIC_ORIGIN: z.string().default(''),
MARKETPLACE_TOKEN_ENCRYPTION_KEY: z.string().default(''),
GITHUB_OAUTH_CLIENT_ID: z.string().default(''),
GITHUB_OAUTH_CLIENT_SECRET: z.string().default(''),
@@ -99,6 +116,26 @@ const environmentSchema = z.object({
FORGEJO_OAUTH_CLIENT_ID: z.string().default(''),
FORGEJO_OAUTH_CLIENT_SECRET: z.string().default(''),
}).superRefine((environment, context) => {
try {
const platform = new URL(environment.MARKETPLACE_PUBLIC_URL);
const moduleOrigin = moduleOriginFor(environment.MARKETPLACE_PUBLIC_URL, environment.MODULE_PUBLIC_ORIGIN);
const modules = new URL(moduleOrigin);
if (modules.origin !== moduleOrigin || modules.hostname === platform.hostname || modules.username || modules.password || modules.search || modules.hash || modules.pathname !== '/') {
throw new Error('origin');
}
if (environment.NODE_ENV === 'production' && modules.protocol !== 'https:') {
throw new Error('https');
}
if (!/^[a-z0-9.-]+$/i.test(modules.hostname)) {
throw new Error('hostname');
}
} catch {
context.addIssue({
code: z.ZodIssueCode.custom,
path: ['MODULE_PUBLIC_ORIGIN'],
message: 'Modul-Origin benötigt einen eigenen Host (in Produktion HTTPS), ohne Pfad oder Zugangsdaten',
});
}
if (environment.NODE_ENV === 'production' && !environment.COOKIE_SECURE) {
context.addIssue({
code: z.ZodIssueCode.custom,
@@ -206,5 +243,6 @@ export function loadConfiguration(): AppConfig {
: {}),
},
},
modulePublicOrigin: moduleOriginFor(environment.MARKETPLACE_PUBLIC_URL, environment.MODULE_PUBLIC_ORIGIN),
};
}

View File

@@ -0,0 +1,32 @@
import type { Migration } from '../migration.types';
/** Einmalige Übergabe vom Plattformhost auf den getrennten Modulhost. */
export const migration013ModuleBrowserSessions: Migration = {
id: '013-module-browser-sessions',
description: 'Einmal-Tickets und getrennte Browser-Sessions für Module',
up: async (client) => {
await client.query(`
CREATE TABLE module_access_tickets (
ticket_hash TEXT PRIMARY KEY,
platform_session_id UUID NOT NULL REFERENCES sessions(id) ON DELETE CASCADE,
user_id UUID NOT NULL REFERENCES users(id) ON DELETE CASCADE,
module_slug TEXT NOT NULL,
expires_at TIMESTAMPTZ NOT NULL
)
`);
await client.query(`
CREATE TABLE module_sessions (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
token_hash TEXT NOT NULL UNIQUE,
platform_session_id UUID NOT NULL REFERENCES sessions(id) ON DELETE CASCADE,
user_id UUID NOT NULL REFERENCES users(id) ON DELETE CASCADE,
module_slug TEXT NOT NULL,
expires_at TIMESTAMPTZ NOT NULL,
created_at TIMESTAMPTZ NOT NULL DEFAULT now()
)
`);
await client.query('CREATE INDEX idx_module_sessions_user ON module_sessions(user_id)');
await client.query('CREATE INDEX idx_module_sessions_expiry ON module_sessions(expires_at)');
await client.query('CREATE INDEX idx_module_tickets_expiry ON module_access_tickets(expires_at)');
},
};

View File

@@ -10,6 +10,7 @@ import { migration009MarketplaceInstallations } from '../../modules/migrations/0
import { migration010ModuleContainers } from '../../modules/migrations/010-module-containers';
import { migration011ModuleConfiguration } from '../../modules/migrations/011-module-configuration';
import { migration012MarketplaceBranchUpdates } from '../../modules/migrations/012-marketplace-branch-updates';
import { migration013ModuleBrowserSessions } from './013-module-browser-sessions';
/** Registrierte Migrationen in aufsteigender Reihenfolge. */
export const MIGRATIONS = [
@@ -25,4 +26,5 @@ export const MIGRATIONS = [
migration010ModuleContainers,
migration011ModuleConfiguration,
migration012MarketplaceBranchUpdates,
migration013ModuleBrowserSessions,
];

View File

@@ -42,7 +42,7 @@ async function bootstrap(): Promise<void> {
.setTitle('MPM Management API')
.setDescription('Zentrale Management-API der MPM-Plattform (Auth, RBAC, Health)')
.setVersion('0.1.0')
.addCookieAuth('mpm_session')
.addCookieAuth(config.security.cookieSecure ? '__Host-mpm_session' : 'mpm_session')
.build();
const document = SwaggerModule.createDocument(app, swaggerConfig);
SwaggerModule.setup('api/docs', app, document);

View File

@@ -55,6 +55,7 @@ const MODULE_ID_PATTERN = /^[a-z][a-z0-9-]{2,63}$/;
/** URL-Slugs für das spätere Routing (/slug). */
const SLUG_PATTERN = /^[a-z0-9][a-z0-9-]{2,99}$/;
const RESERVED_SLUGS = new Set(['api', 'assets', 'login', 'admin', 'profile', '403', '404']);
/** Semantische Versionierung (major.minor.patch). */
const VERSION_PATTERN = /^\d+\.\d+\.\d+$/;
@@ -69,7 +70,9 @@ export const moduleManifestSchema = z.object({
.regex(MODULE_ID_PATTERN, 'Modul-ID muss dem Muster [a-z][a-z0-9-]{2,63} folgen'),
name: z.string().trim().min(1, 'Name ist erforderlich').max(100),
version: z.string().regex(VERSION_PATTERN, 'Version muss dem Muster major.minor.patch folgen'),
slug: z.string().regex(SLUG_PATTERN, 'Slug muss dem Muster [a-z0-9-]{3,100} folgen'),
slug: z.string()
.regex(SLUG_PATTERN, 'Slug muss dem Muster [a-z0-9-]{3,100} folgen')
.refine((slug) => !RESERVED_SLUGS.has(slug), 'Dieser Slug ist für die Plattform reserviert'),
description: z.string().max(500).default(''),
author: z.string().max(200).default(''),
runtime: z.literal('node'),

View File

@@ -1,4 +1,4 @@
import { BadRequestException, Body, Controller, Delete, Get, Param, Post, Query, Req, Res } from '@nestjs/common';
import { BadRequestException, Body, Controller, Delete, Get, Inject, 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';
@@ -6,6 +6,8 @@ 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 { extractSessionToken } from '../auth/guards/session.guard';
import { APP_CONFIG, type AppConfig } from '../config/config.tokens';
import type { AuthenticatedRequest } from '../auth/authenticated-request';
import { MarketplaceService } from './marketplace.service';
import { ModulesService } from './modules.service';
@@ -18,6 +20,7 @@ export class MarketplaceController {
private readonly marketplaceService: MarketplaceService,
private readonly sessionService: SessionService,
private readonly modulesService: ModulesService,
@Inject(APP_CONFIG) private readonly config: AppConfig,
) {}
@Get('providers')
@@ -46,11 +49,13 @@ export class MarketplaceController {
configuration: ModuleRecord['configuration']; configurationReady: boolean;
} }> {
const branch = await this.marketplaceService.defaultBranch(provider, owner, repository);
const archive = await this.marketplaceService.downloadRepositoryArchive(provider, owner, repository, branch);
const { archive, commit } = await this.marketplaceService.downloadRepositoryArchive(provider, owner, repository, branch);
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, branch);
if (await this.modulesService.findByModuleId(manifest.id)) {
throw new BadRequestException('Ein Modul mit dieser ID ist bereits installiert');
}
const module = await this.modulesService.install(archive, actor, request.ip ?? null);
await this.marketplaceService.recordInstallation(provider, owner, repository, module.id, branch, commit);
return {
module: {
id: module.id,
@@ -105,9 +110,14 @@ export class MarketplaceController {
actor,
request.ip ?? null,
report,
async () => {
report('commit', 'Neue Version wird registriert', 96);
await this.marketplaceService.commitInstalledBranch(
moduleId, update.provider, update.owner, update.repository, update.branch, update.commit,
update.previousBranch, update.previousCommit,
);
},
);
report('commit', 'Neue Version wird registriert', 96);
await this.marketplaceService.commitInstalledBranch(moduleId, update.provider, update.owner, update.repository, update.branch, update.commit);
this.marketplaceService.finishOperation(operationId, actor.id, true);
return { module };
} catch (error) {
@@ -152,7 +162,7 @@ export class MarketplaceController {
destination.searchParams.set('reason', 'callback');
} else {
try {
const token = request.cookies?.mpm_session as string | undefined;
const token = extractSessionToken(request, this.config.security.cookieSecure);
const session = token ? await this.sessionService.findValid(token) : null;
if (!session) {
destination.searchParams.set('marketplace', 'error');

View File

@@ -1,6 +1,7 @@
import {
BadGatewayException,
BadRequestException,
ConflictException,
Injectable,
InternalServerErrorException,
Logger,
@@ -52,6 +53,13 @@ export interface MarketplaceUpdatePackage {
owner: string;
repository: string;
branch: string;
previousBranch: string;
previousCommit: string | null;
}
export interface MarketplaceRepositoryArchive {
archive: Buffer;
commit: string;
}
export interface MarketplaceOperationProgress {
@@ -74,6 +82,10 @@ interface MarketplaceInstallationRow {
const MAX_MARKETPLACE_DOWNLOAD = 10 * 1024 * 1024;
function isCommitSha(value: string): boolean {
return /^(?:[0-9a-f]{40}|[0-9a-f]{64})$/i.test(value);
}
function isProvider(value: string): value is MarketplaceProvider {
return MARKETPLACE_PROVIDERS.includes(value as MarketplaceProvider);
}
@@ -305,10 +317,15 @@ export class MarketplaceService implements OnModuleInit, OnModuleDestroy {
}));
}
async recordInstallation(providerParam: string, owner: string, repository: string, moduleId: string, branch: string): Promise<void> {
async recordInstallation(providerParam: string, owner: string, repository: string, moduleId: string, branch: string, commit: string): Promise<void> {
const provider = this.requireProvider(providerParam);
const branchInfo = await this.getBranch(provider, owner, repository, branch);
const branches = await this.fetchBranches(provider, owner, repository);
if (!isCommitSha(commit)) throw new BadRequestException('Commit-ID ist ungültig');
let branches: MarketplaceBranch[] = [];
try {
branches = await this.fetchBranches(provider, owner, repository);
} catch (error) {
this.logger.warn(`Branch-Prüfung für Modul ${moduleId} nach Installation fehlgeschlagen: ${error instanceof Error ? error.message : 'unbekannter Fehler'}`);
}
await this.database.query(
`INSERT INTO marketplace_module_installations
(provider, owner, repository, module_id, installed_branch, installed_commit, available_branches, observed_branches, branches_checked_at)
@@ -317,7 +334,7 @@ export class MarketplaceService implements OnModuleInit, OnModuleDestroy {
module_id = EXCLUDED.module_id, installed_branch = EXCLUDED.installed_branch,
installed_commit = EXCLUDED.installed_commit, available_branches = '[]'::jsonb,
observed_branches = EXCLUDED.observed_branches, branches_checked_at = now()`,
[provider, owner, repository, moduleId, branch, branchInfo.commit, JSON.stringify(branches)],
[provider, owner, repository, moduleId, branch, commit, JSON.stringify(branches)],
);
}
@@ -354,7 +371,13 @@ export class MarketplaceService implements OnModuleInit, OnModuleDestroy {
};
}
async downloadRepositoryArchive(providerParam: string, owner: string, repository: string, branch?: string): Promise<Buffer> {
async downloadRepositoryArchive(
providerParam: string,
owner: string,
repository: string,
branch?: string,
pinnedCommit?: string,
): Promise<MarketplaceRepositoryArchive> {
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');
@@ -370,17 +393,19 @@ export class MarketplaceService implements OnModuleInit, OnModuleDestroy {
if (repo.private === true) throw new BadRequestException('Private Repositories werden aktuell nicht unterstuetzt');
const selectedBranch = branch ?? String(repo.default_branch ?? 'main');
if (!selectedBranch || selectedBranch.length > 200 || selectedBranch.includes('\0')) throw new BadRequestException('Branch ist ungueltig');
if (branch) await this.getBranch(provider, owner, repository, branch);
// Resolve once, then request the immutable commit instead of the movable branch ref.
const commit = pinnedCommit ?? (await this.getBranch(provider, owner, repository, selectedBranch)).commit;
if (!isCommitSha(commit)) throw new BadGatewayException('Forge hat eine ungültige Commit-ID geliefert');
const archivePath = provider === 'github'
? '/repos/' + encodeURIComponent(owner) + '/' + encodeURIComponent(repository) + '/zipball/' + encodeURIComponent(selectedBranch)
: new URL(providerConfig.baseUrl).pathname.replace(/[/]$/, '') + '/api/v1/repos/' + encodeURIComponent(owner) + '/' + encodeURIComponent(repository) + '/archive/' + encodeURIComponent(selectedBranch) + '.zip';
? '/repos/' + encodeURIComponent(owner) + '/' + encodeURIComponent(repository) + '/zipball/' + commit
: new URL(providerConfig.baseUrl).pathname.replace(/[/]$/, '') + '/api/v1/repos/' + encodeURIComponent(owner) + '/' + encodeURIComponent(repository) + '/archive/' + commit + '.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);
return { archive: await this.normalizeRepositoryArchive(archive), commit: commit.toLowerCase() };
}
async updateInstalledBranch(
@@ -408,23 +433,44 @@ export class MarketplaceService implements OnModuleInit, OnModuleDestroy {
throw new BadRequestException('Diese Branch enthält keine Änderungen gegenüber der installierten Version');
}
onProgress?.('download', 'Update-Archiv wird geladen', 28);
const archive = await this.downloadRepositoryArchive(installation.provider, installation.owner, installation.repository, branch);
const { archive } = await this.downloadRepositoryArchive(
installation.provider, installation.owner, installation.repository, branch, branchInfo.commit,
);
// Caller updates and validates the module files before this source record is advanced.
return { archive, commit: branchInfo.commit, provider: installation.provider,
owner: installation.owner, repository: installation.repository, branch };
owner: installation.owner, repository: installation.repository, branch,
previousBranch: installation.installed_branch, previousCommit: installation.installed_commit };
}
async commitInstalledBranch(moduleId: string, provider: MarketplaceProvider, owner: string, repository: string, branch: string, commit: string): Promise<void> {
await this.database.query(
async commitInstalledBranch(
moduleId: string,
provider: MarketplaceProvider,
owner: string,
repository: string,
branch: string,
commit: string,
expectedBranch: string,
expectedCommit: string | null,
): Promise<void> {
const result = await this.database.query<{ module_id: string }>(
`UPDATE marketplace_module_installations SET installed_branch = $2, installed_commit = $3,
available_branches = COALESCE((
SELECT jsonb_agg(item.value) FROM jsonb_array_elements(available_branches) AS item(value)
WHERE item.value->>'name' <> $2
), '[]'::jsonb), branches_checked_at = now()
WHERE module_id = $1 AND provider = $4 AND owner = $5 AND repository = $6`,
[moduleId, branch, commit, provider, owner, repository],
WHERE module_id = $1 AND provider = $4 AND owner = $5 AND repository = $6
AND installed_branch = $7 AND installed_commit IS NOT DISTINCT FROM $8
RETURNING module_id`,
[moduleId, branch, commit, provider, owner, repository, expectedBranch, expectedCommit],
);
await this.refreshInstallation(moduleId);
if (!result.rows.length) {
throw new ConflictException('Die installierte Branch wurde zwischenzeitlich geändert. Bitte Updates neu laden.');
}
try {
await this.refreshInstallation(moduleId);
} catch (error) {
this.logger.warn(`Branch-Prüfung für Modul ${moduleId} nach Update fehlgeschlagen: ${error instanceof Error ? error.message : 'unbekannter Fehler'}`);
}
}
private async getBranch(provider: MarketplaceProvider, owner: string, repository: string, branch: string): Promise<MarketplaceBranch> {
@@ -437,7 +483,7 @@ export class MarketplaceService implements OnModuleInit, OnModuleDestroy {
const value = await this.forgeJson<Record<string, unknown>>(provider, this.config.marketplace.providers[provider]!.baseUrl, null, basePath);
const commit = value.commit as Record<string, unknown> | undefined;
const sha = String(commit?.id ?? commit?.sha ?? '');
if (!sha) throw new NotFoundException('Branch konnte beim Forge nicht gefunden werden');
if (!isCommitSha(sha)) throw new BadGatewayException('Forge hat eine ungültige Commit-ID geliefert');
return { name: String(value.name ?? branch), commit: sha };
}

View File

@@ -1,6 +1,6 @@
import { BadRequestException, Injectable, Logger } from '@nestjs/common';
import { spawn } from 'node:child_process';
import { chmod, mkdir, readFile, rm, writeFile } from 'node:fs/promises';
import { chmod, mkdir, readFile, realpath, rm, stat, writeFile } from 'node:fs/promises';
import { tmpdir } from 'node:os';
import path from 'node:path';
import { stringify, parseDocument } from 'yaml';
@@ -15,13 +15,25 @@ const SAFE_SERVICE_KEYS = new Set([
'stop_grace_period', 'read_only', 'tty', 'stdin_open',
]);
const SAFE_BUILD_KEYS = new Set(['context', 'dockerfile', 'target', 'args']);
const MAX_MODULE_SERVICES = 8;
const MODULE_MEMORY_LIMIT = '512m';
const MODULE_CPU_LIMIT = 1;
const MODULE_PIDS_LIMIT = 256;
const COMPOSE_COMMAND_TIMEOUT_MS = 15 * 60_000;
const COMPOSE_CONTROL_TIMEOUT_MS = 2 * 60_000;
const DOCKER_COMMAND_TIMEOUT_MS = 45_000;
const COMMAND_TERMINATION_GRACE_MS = 5_000;
/** Ein Docker-CLI-Fehler mit einer für die Admin-Oberfläche bereinigten Diagnose. */
export class ModuleCommandError extends Error {
constructor(
readonly exitCode: number | null,
readonly diagnostic: string,
readonly timedOut = false,
) {
super(exitCode === null ? 'Docker-Befehl konnte nicht gestartet werden' : `Docker-Befehl endete mit Status ${exitCode}`);
super(timedOut ? 'Docker-Befehl hat das Zeitlimit überschritten' :
exitCode === null ? 'Docker-Befehl konnte nicht gestartet werden' : `Docker-Befehl endete mit Status ${exitCode}`);
this.name = 'ModuleCommandError';
}
}
@@ -83,6 +95,19 @@ export class ModuleContainerManager {
}
}
/** Removes containers from a failed start while retaining all named data volumes. */
async cleanupFailedStart(module: ModuleRecord): Promise<void> {
const { composePath, overridePath, projectName, gatewayNetwork, cleanupValues } = await this.prepare(module);
try {
await this.runCompose(module.path, projectName, composePath, overridePath, ['down', '--remove-orphans'], cleanupValues);
if (this.mpmContainer) await this.runDocker(['network', 'disconnect', '-f', gatewayNetwork, this.mpmContainer], true);
await this.runDocker(['network', 'rm', gatewayNetwork], true);
this.logger.log(`Teilweise gestarteter Stack für "${module.moduleId}" ohne Datenverlust bereinigt`);
} finally {
await rm(overridePath, { force: true });
}
}
private async prepare(module: ModuleRecord): Promise<{
composePath: string;
overridePath: string;
@@ -97,13 +122,18 @@ export class ModuleContainerManager {
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(() => {
const actualRoot = await realpath(root);
const actualCompose = await realpath(composePath).catch(() => {
throw new BadRequestException(`Compose-Datei "${module.composeFile}" wurde nicht gefunden`);
});
if (!this.isInside(actualRoot, actualCompose)) {
throw new BadRequestException('Compose-Datei liegt außerhalb des Modulpakets');
}
const source = await readFile(composePath, 'utf8');
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);
await this.validateCompose(compose, module, actualRoot, path.dirname(composePath));
const serviceMap = compose.services as Record<string, unknown>;
const moduleValues = await this.configurationService.values(module);
// Compose validates required interpolations even for stop/down. Supply
@@ -126,21 +156,28 @@ export class ModuleContainerManager {
// 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 overrideServices: Record<string, Record<string, unknown>> = {
[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}`] },
},
const overrideServices: Record<string, Record<string, unknown>> = {};
for (const serviceName of Object.keys(serviceMap)) {
overrideServices[serviceName] = {
security_opt: ['no-new-privileges:true'],
mem_limit: MODULE_MEMORY_LIMIT,
cpus: MODULE_CPU_LIMIT,
pids_limit: MODULE_PIDS_LIMIT,
};
}
overrideServices[module.appService] = {
...overrideServices[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}`] },
},
};
for (const field of module.configuration) {
@@ -164,13 +201,18 @@ export class ModuleContainerManager {
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 });
await chmod(overridePath, 0o600);
try {
await mkdir(path.dirname(overridePath), { recursive: true, mode: 0o700 });
await writeFile(overridePath, stringify(override), { mode: 0o600 });
await chmod(overridePath, 0o600);
} catch (error) {
await rm(overridePath, { force: true });
throw error;
}
return { composePath, overridePath, projectName, gatewayNetwork, moduleValues, cleanupValues };
}
private validateCompose(compose: Record<string, unknown>, module: ModuleRecord): void {
private async validateCompose(compose: Record<string, unknown>, module: ModuleRecord, root: string, composeDir: string): Promise<void> {
if (!compose || typeof compose !== 'object' || Array.isArray(compose)) {
throw new BadRequestException('Compose-Datei muss ein YAML-Objekt enthalten');
}
@@ -183,6 +225,9 @@ export class ModuleContainerManager {
throw new BadRequestException('Compose benötigt mindestens einen Service');
}
const serviceMap = services as Record<string, unknown>;
if (Object.keys(serviceMap).length > MAX_MODULE_SERVICES) {
throw new BadRequestException(`Compose darf höchstens ${MAX_MODULE_SERVICES} Services enthalten`);
}
if (!Object.hasOwn(serviceMap, module.appService!)) {
throw new BadRequestException(`Compose-Service "${module.appService}" fehlt`);
}
@@ -205,9 +250,8 @@ export class ModuleContainerManager {
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.build !== undefined) await this.validateBuild(service.build, root, composeDir);
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 &&
@@ -235,18 +279,64 @@ export class ModuleContainerManager {
return value as Record<string, unknown>;
}
private validateBuild(build: unknown, modulePath: string): void {
const context = typeof build === 'string'
? build
private async validateBuild(build: unknown, root: string, composeDir: string): Promise<void> {
const options: Record<string, unknown> | null = typeof build === 'string'
? { context: 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');
? build as Record<string, unknown>
: null;
if (!options || Object.keys(options).some((key) => !SAFE_BUILD_KEYS.has(key))) {
throw new BadRequestException('Compose-Build enthält nicht erlaubte Optionen');
}
const resolved = path.resolve(modulePath, context);
if (resolved !== modulePath && !resolved.startsWith(modulePath + path.sep)) {
throw new BadRequestException('Build-Kontext liegt außerhalb des Modulpakets');
const context = options.context ?? '.';
if (typeof context !== 'string' || !this.isStaticRelativePath(context)) {
throw new BadRequestException('Build-Kontext muss ein fester relativer Pfad sein');
}
const contextPath = path.resolve(composeDir, context);
const actualContext = await realpath(contextPath).catch(() => {
throw new BadRequestException('Build-Kontext wurde nicht gefunden');
});
if (!this.isInside(root, actualContext) || !(await stat(actualContext)).isDirectory()) {
throw new BadRequestException('Build-Kontext muss ein Verzeichnis innerhalb des Modulpakets sein');
}
const dockerfile = options.dockerfile ?? 'Dockerfile';
if (typeof dockerfile !== 'string' || !this.isStaticRelativePath(dockerfile, true)) {
throw new BadRequestException('Dockerfile muss ein fester relativer Pfad sein');
}
const actualDockerfile = await realpath(path.resolve(actualContext, dockerfile)).catch(() => {
throw new BadRequestException('Dockerfile wurde nicht gefunden');
});
if (!this.isInside(root, actualDockerfile) || !(await stat(actualDockerfile)).isFile()) {
throw new BadRequestException('Dockerfile muss eine Datei innerhalb des Modulpakets sein');
}
if (options.target !== undefined && (typeof options.target !== 'string' || !/^[A-Za-z0-9][A-Za-z0-9_.-]*$/.test(options.target))) {
throw new BadRequestException('Compose-Build-Target ist ungültig');
}
if (options.args !== undefined) this.validateBuildArgs(options.args);
}
private isStaticRelativePath(value: string, allowParent = false): boolean {
return value.length > 0 && !path.isAbsolute(value) && !value.includes('\\') && !value.includes('$') &&
!value.includes(':') && !value.includes('#') && !value.startsWith('~') &&
(allowParent || !value.split('/').includes('..'));
}
private isInside(root: string, target: string): boolean {
return target === root || target.startsWith(root + path.sep);
}
private validateBuildArgs(args: unknown): void {
if (Array.isArray(args)) {
if (!args.every((arg) => typeof arg === 'string' && /^[A-Za-z_][A-Za-z0-9_]*(=.*)?$/.test(arg))) {
throw new BadRequestException('Compose-Build-Argumente sind ungültig');
}
return;
}
if (!args || typeof args !== 'object' || Object.entries(args).some(([key, value]) =>
!/^[A-Za-z_][A-Za-z0-9_]*$/.test(key) ||
(value !== null && !['string', 'number', 'boolean'].includes(typeof value)))) {
throw new BadRequestException('Compose-Build-Argumente sind ungültig');
}
}
@@ -265,17 +355,21 @@ export class ModuleContainerManager {
}
private runCompose(cwd: string, project: string, composePath: string, overridePath: string, args: string[], config: Record<string, string>): Promise<void> {
return this.run('docker-compose', ['-p', project, '-f', composePath, '-f', overridePath, ...args], cwd, false, config);
const timeoutMs = args[0] === 'up' ? COMPOSE_COMMAND_TIMEOUT_MS : COMPOSE_CONTROL_TIMEOUT_MS;
return this.run('docker-compose', ['-p', project, '-f', composePath, '-f', overridePath, ...args], cwd, false, config, timeoutMs);
}
private runDocker(args: string[], ignoreFailure = false): Promise<void> {
return this.run('docker', args, process.cwd(), ignoreFailure);
return this.run('docker', args, process.cwd(), ignoreFailure, {}, DOCKER_COMMAND_TIMEOUT_MS);
}
private run(command: string, args: string[], cwd: string, ignoreFailure = false, extraEnv: Record<string, string> = {}): Promise<void> {
private run(command: string, args: string[], cwd: string, ignoreFailure = false, extraEnv: Record<string, string> = {}, timeoutMs = DOCKER_COMMAND_TIMEOUT_MS): Promise<void> {
return new Promise((resolve, reject) => {
const child = spawn(command, args, {
cwd,
// Give Docker Compose and its subprocesses one process group so a
// timeout can stop the whole operation before cleanup starts.
detached: process.platform !== 'win32',
env: {
...extraEnv,
PATH: process.env.PATH ?? '/usr/local/sbin:/usr/local/bin:/usr/sbin:/usr/bin:/sbin:/bin',
@@ -293,18 +387,52 @@ export class ModuleContainerManager {
child.stdout.on('data', (chunk: string) => { stdout = keepTail(stdout, chunk); });
child.stderr.setEncoding('utf8');
child.stderr.on('data', (chunk: string) => { stderr = keepTail(stderr, chunk); });
let processStartFailed = false;
let settled = false;
let timedOut = false;
let terminationTimer: NodeJS.Timeout | undefined;
const stopProcessGroup = (signal: NodeJS.Signals): void => {
if (process.platform !== 'win32' && child.pid) {
try {
process.kill(-child.pid, signal);
return;
} catch {
// The process group may already have exited.
}
}
child.kill(signal);
};
const timeout = setTimeout(() => {
if (settled) return;
timedOut = true;
terminationTimer = setTimeout(() => stopProcessGroup('SIGKILL'), COMMAND_TERMINATION_GRACE_MS);
terminationTimer.unref();
stopProcessGroup('SIGTERM');
}, timeoutMs);
child.once('error', (error) => {
if (timedOut) return;
clearTimeout(timeout);
if (terminationTimer) clearTimeout(terminationTimer);
if (settled) return;
settled = true;
if (ignoreFailure) resolve();
else {
processStartFailed = true;
const errorCode = (error as NodeJS.ErrnoException).code ?? 'unbekannt';
this.logger.error(`${command} konnte nicht gestartet werden (${errorCode})`);
reject(new ModuleCommandError(null, `Der Befehl „${command}“ konnte nicht gestartet werden. Prüfe, ob Docker auf dem System verfügbar ist.`));
}
});
child.once('close', (code) => {
if (processStartFailed) return;
clearTimeout(timeout);
if (terminationTimer) clearTimeout(terminationTimer);
if (settled) return;
settled = true;
if (timedOut) {
const duration = timeoutMs < 60_000 ? `${timeoutMs / 1_000} Sekunden` : `${timeoutMs / 60_000} Minuten`;
const diagnostic = `${command} hat das Zeitlimit von ${duration} überschritten.`;
this.logger.error(diagnostic);
reject(new ModuleCommandError(null, diagnostic, true));
return;
}
if (code === 0 || ignoreFailure) resolve();
else {
const diagnostic = this.sanitizeDiagnostic(`${stdout}\n${stderr}`, extraEnv);

View File

@@ -54,7 +54,11 @@ function createUserRecord(overrides: Partial<UserRecord> = {}): UserRecord {
function createRequest(url: string, cookie?: string): Request {
return {
url,
headers: cookie ? { cookie } : {},
method: 'GET',
headers: {
host: 'localhost:8081',
...(cookie ? { cookie } : {}),
},
} as unknown as Request;
}
@@ -79,8 +83,8 @@ function createResponse(): Response & { sentStatus: number; sentBody: unknown }
class MockSessionService {
public session: SessionData | null = null;
async findValid(): Promise<SessionData | null> {
return this.session;
async findValidModuleSession(): Promise<string | null> {
return this.session?.userId ?? null;
}
}
@@ -135,6 +139,7 @@ describe('ModuleGatewayMiddleware', () => {
sessionService as unknown as SessionService,
userRepository as unknown as UserRepository,
permissionsService as unknown as ModulePermissionsService,
{ modulePublicOrigin: 'http://localhost:8081' } as never,
);
sessionService.session = {
id: 'session-1',
@@ -163,7 +168,7 @@ describe('ModuleGatewayMiddleware', () => {
it('antwortet 401 bei ungültiger Session', async () => {
sessionService.session = null;
const request = createRequest('/api/v1/gateway/demo/', 'mpm_session=invalid');
const request = createRequest('/api/v1/gateway/demo/', 'mpm_module_session=invalid');
const response = createResponse();
await middleware.use(request, response, makeNext());
@@ -172,7 +177,7 @@ describe('ModuleGatewayMiddleware', () => {
it('antwortet 401 bei deaktiviertem Benutzer', async () => {
userRepository.user = createUserRecord({ isActive: false });
const request = createRequest('/api/v1/gateway/demo/', 'mpm_session=valid');
const request = createRequest('/api/v1/gateway/demo/', 'mpm_module_session=valid');
const response = createResponse();
await middleware.use(request, response, makeNext());
@@ -181,7 +186,7 @@ describe('ModuleGatewayMiddleware', () => {
it('antwortet 404 bei unbekanntem Modul-Slug', async () => {
moduleRepository.module = null;
const request = createRequest('/api/v1/gateway/demo/', 'mpm_session=valid');
const request = createRequest('/api/v1/gateway/demo/', 'mpm_module_session=valid');
const response = createResponse();
await middleware.use(request, response, makeNext());
@@ -190,7 +195,7 @@ describe('ModuleGatewayMiddleware', () => {
it('antwortet 503 bei gestopptem Modul', async () => {
moduleRepository.module = createModuleRecord({ status: 'STOPPED' });
const request = createRequest('/api/v1/gateway/demo/', 'mpm_session=valid');
const request = createRequest('/api/v1/gateway/demo/', 'mpm_module_session=valid');
const response = createResponse();
await middleware.use(request, response, makeNext());
@@ -199,7 +204,7 @@ describe('ModuleGatewayMiddleware', () => {
it('antwortet 503 bei deaktiviertem Modul', async () => {
moduleRepository.module = createModuleRecord({ enabled: false });
const request = createRequest('/api/v1/gateway/demo/', 'mpm_session=valid');
const request = createRequest('/api/v1/gateway/demo/', 'mpm_module_session=valid');
const response = createResponse();
await middleware.use(request, response, makeNext());
@@ -208,7 +213,7 @@ describe('ModuleGatewayMiddleware', () => {
it('antwortet 403 für USER ohne Berechtigung (fail-closed)', async () => {
permissionsService.hasAccessResult = false;
const request = createRequest('/api/v1/gateway/demo/', 'mpm_session=valid');
const request = createRequest('/api/v1/gateway/demo/', 'mpm_module_session=valid');
const response = createResponse();
await middleware.use(request, response, makeNext());
@@ -217,7 +222,7 @@ describe('ModuleGatewayMiddleware', () => {
it('leitet USER-Requests mit GRANTED-Berechtigung an den Proxy weiter', async () => {
permissionsService.hasAccessResult = true;
const request = createRequest('/api/v1/gateway/demo/health', 'mpm_session=valid');
const request = createRequest('/api/v1/gateway/demo/health', 'mpm_module_session=valid');
const response = createResponse();
const proxySpy = jest
@@ -235,7 +240,7 @@ describe('ModuleGatewayMiddleware', () => {
it('leitet ADMIN-Requests an den Modul-Proxy weiter', async () => {
userRepository.user = createUserRecord({ role: 'ADMIN' });
const request = createRequest('/api/v1/gateway/demo/health', 'mpm_session=valid');
const request = createRequest('/api/v1/gateway/demo/health', 'mpm_module_session=valid');
const response = createResponse();
// proxy.web würde einen echten Request starten – hier nur prüfen,

View File

@@ -1,8 +1,9 @@
import { Injectable, type NestMiddleware } from '@nestjs/common';
import { ForbiddenException, Inject, Injectable, type NestMiddleware } from '@nestjs/common';
import type { Request, Response, NextFunction } from 'express';
import httpProxy from 'http-proxy';
import { APP_CONFIG, type AppConfig } from '../config/config.tokens';
import { SessionService } from '../auth/session.service';
import { extractSessionToken } from '../auth/guards/session.guard';
import { requireSameOrigin } from '../auth/request-origin';
import { UserRepository } from '../users/user.repository';
import { ModuleRepository } from './module.repository';
import { ModulePermissionsService } from './module-permissions.service';
@@ -10,6 +11,23 @@ import { ModuleIdentityService } from './module-identity.service';
/** Gateway-Pfad-Präfix für interne Nginx-Weiterleitung. */
const GATEWAY_PREFIX = '/api/v1/gateway/';
const STATE_CHANGING_METHODS = new Set(['POST', 'PUT', 'PATCH', 'DELETE']);
function moduleCookie(request: Request): string | null {
let value: string | null = null;
for (const part of (request.headers.cookie ?? '').split(';')) {
const [name, ...parts] = part.trim().split('=');
if (name === 'mpm_module_session') value = decodeURIComponent(parts.join('='));
}
return value;
}
function cookiesForModule(cookieHeader: string | undefined): string | undefined {
const remaining = (cookieHeader ?? '').split(';').map((part) => part.trim()).filter((part) =>
part && !/^(mpm_module_session|mpm_session|mpm_csrf|__Host-mpm_session|__Host-mpm_csrf)=/.test(part),
);
return remaining.length ? remaining.join('; ') : undefined;
}
/**
* Modul-Gateway (Phase 4/5): Dynamisches Routing /slug → Modul-Prozess.
@@ -35,6 +53,7 @@ export class ModuleGatewayMiddleware implements NestMiddleware {
private readonly sessionService: SessionService,
private readonly userRepository: UserRepository,
private readonly permissionsService: ModulePermissionsService,
@Inject(APP_CONFIG) private readonly config: AppConfig,
private readonly identityService: ModuleIdentityService = new ModuleIdentityService(),
) {
this.proxy = httpProxy.createProxyServer({
@@ -69,6 +88,21 @@ export class ModuleGatewayMiddleware implements NestMiddleware {
proxyRequest.setHeader('Content-Length', Buffer.byteLength(body));
proxyRequest.write(body);
});
// Modul-Cookies dürfen weder für die ganze Parent-Domain noch für andere
// Module gelten. Eigene App-Cookies bleiben innerhalb des Modulpfads nutzbar.
this.proxy.on('proxyRes', (proxyResponse, request) => {
const slug = (request as Request & { mpmModuleSlug?: string }).mpmModuleSlug;
const cookies = proxyResponse.headers['set-cookie'];
if (!slug || !cookies) return;
proxyResponse.headers['set-cookie'] = cookies.map((cookie) => {
const [nameAndValue, ...attributes] = cookie.split(';');
const restrictedAttributes = attributes.filter((attribute) =>
!/^\s*(domain|path|samesite)\s*=/i.test(attribute),
);
return `${nameAndValue};${restrictedAttributes.join(';')}; Path=/${slug}; SameSite=Lax`;
});
});
}
async use(request: Request, response: Response, next: NextFunction): Promise<void> {
@@ -77,26 +111,45 @@ export class ModuleGatewayMiddleware implements NestMiddleware {
return;
}
// Die API ist auf dem Modulhost nicht sichtbar. Das Gateway darf nur
// Requests vom dedizierten Browser-Origin weiterleiten.
if (request.headers.host !== new URL(this.config.modulePublicOrigin).host) {
response.status(403).json({ statusCode: 403, message: 'Ungültiger Modul-Host' });
return;
}
if (STATE_CHANGING_METHODS.has(request.method)) {
try {
requireSameOrigin(request, this.config.modulePublicOrigin);
} catch (error) {
if (error instanceof ForbiddenException) {
response.status(403).json({ statusCode: 403, message: error.message });
return;
}
throw error;
}
}
// Slug aus dem Gateway-Pfad extrahieren: /api/v1/gateway/<slug>/<rest>
const pathAfterPrefix = request.url.slice(GATEWAY_PREFIX.length);
const gatewayUrl = new URL(request.url, 'http://gateway.internal');
const pathAfterPrefix = gatewayUrl.pathname.slice(GATEWAY_PREFIX.length);
const slashIndex = pathAfterPrefix.indexOf('/');
const slug = slashIndex === -1 ? pathAfterPrefix : pathAfterPrefix.slice(0, slashIndex);
const modulePath = slashIndex === -1 ? '/' : pathAfterPrefix.slice(slashIndex);
const modulePath = (slashIndex === -1 ? '/' : pathAfterPrefix.slice(slashIndex)) + gatewayUrl.search;
// 1. Authentifizierung: Session aus Cookie laden
const token = extractSessionToken(request);
const token = moduleCookie(request);
if (!token) {
response.status(401).json({ statusCode: 401, message: 'Nicht authentifiziert' });
return;
}
const session = await this.sessionService.findValid(token);
if (!session) {
const userId = await this.sessionService.findValidModuleSession(token, slug);
if (!userId) {
response.status(401).json({ statusCode: 401, message: 'Nicht authentifiziert' });
return;
}
const user = await this.userRepository.findById(session.userId);
const user = await this.userRepository.findById(userId);
if (!user || !user.isActive) {
response.status(401).json({ statusCode: 401, message: 'Nicht authentifiziert' });
return;
@@ -139,8 +192,11 @@ export class ModuleGatewayMiddleware implements NestMiddleware {
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 as Request & { mpmModuleSlug?: string }).mpmModuleSlug = slug;
// Nur eigene App-Cookies, nie Plattform- oder Modul-Gateway-Cookies weiterreichen.
const appCookies = cookiesForModule(request.headers.cookie);
if (appCookies) request.headers.cookie = appCookies;
else delete request.headers.cookie;
this.proxy.web(request, response, {
target: `http://${module.composeFile && module.appService ? `mpm-${module.moduleId}` : '127.0.0.1'}:${module.internalPort}`,

View File

@@ -8,6 +8,10 @@ import { parseDocument } from 'yaml';
/** Maximale Größe eines Modul-Pakets (10 MB). */
const MAX_PACKAGE_SIZE_BYTES = 10 * 1024 * 1024;
const MAX_EXTRACTED_SIZE_BYTES = 50 * 1024 * 1024;
const MAX_ARCHIVE_ENTRIES = 2000;
const MAX_SINGLE_FILE_BYTES = 20 * 1024 * 1024;
const MAX_METADATA_FILE_BYTES = 1024 * 1024;
/** Dateien, die in einem Modul-Paket erwartet werden. */
const REQUIRED_MANIFEST_FILE = 'module.json';
@@ -28,20 +32,19 @@ export class ModuleInstaller {
/** Validiert ein hochgeladenes Paket und gibt das Manifest zurück. */
async validatePackage(buffer: Buffer): Promise<ModuleManifest> {
if (buffer.length === 0) {
throw new BadRequestException('Paket ist leer');
}
if (buffer.length > MAX_PACKAGE_SIZE_BYTES) {
throw new BadRequestException('Paket ist zu groß (maximal 10 MB)');
}
this.assertCompressedSize(buffer);
const AdmZip = (await import('adm-zip')).default;
const zip = new AdmZip(buffer);
this.assertSafeArchive(zip.getEntries(), path.resolve('/module-package'));
const manifestEntry = zip.getEntry(REQUIRED_MANIFEST_FILE);
if (!manifestEntry) {
throw new BadRequestException(`Paket enthält keine ${REQUIRED_MANIFEST_FILE}`);
}
if (manifestEntry.header.size > MAX_METADATA_FILE_BYTES) {
throw new BadRequestException(`${REQUIRED_MANIFEST_FILE} ist zu groß`);
}
let manifestJson: unknown;
try {
@@ -75,10 +78,14 @@ export class ModuleInstaller {
manifest.composeFile.split('/').includes('..')) {
throw new BadRequestException('composeFile muss ein relativer Pfad innerhalb des Modulpakets sein');
}
if (!zip.getEntry(manifest.composeFile)) {
const composeEntry = zip.getEntry(manifest.composeFile);
if (!composeEntry) {
throw new BadRequestException(`Container-Konfiguration ${manifest.composeFile} fehlt im Paket`);
}
const composeDocument = parseDocument(zip.getEntry(manifest.composeFile)!.getData().toString('utf8'), { uniqueKeys: true });
if (composeEntry.header.size > MAX_METADATA_FILE_BYTES) {
throw new BadRequestException('Compose-Datei ist zu groß');
}
const composeDocument = parseDocument(composeEntry.getData().toString('utf8'), { uniqueKeys: true });
if (composeDocument.errors.length) {
throw new BadRequestException('Compose-Datei enthält keine gültige Service-Definition');
}
@@ -106,23 +113,14 @@ export class ModuleInstaller {
manifest: ModuleManifest,
modulesDir: string,
): Promise<{ directory: string; manifest: ModuleManifest }> {
this.assertCompressedSize(buffer);
const directory = path.join(modulesDir, manifest.id);
// Zip-Slip-Schutz: Alle Einträge müssen innerhalb des Zielverzeichnisses liegen.
const AdmZip = (await import('adm-zip')).default;
const zip = new AdmZip(buffer);
const resolvedDirectory = path.resolve(directory);
for (const entry of zip.getEntries()) {
const entryName = entry.entryName;
if (entryName.startsWith('/') || entryName.includes('..') || /^[A-Za-z]:/.test(entryName)) {
throw new BadRequestException(`Unsicherer Pfad im Paket: ${entryName}`);
}
const resolvedEntry = path.resolve(resolvedDirectory, entryName);
if (!resolvedEntry.startsWith(resolvedDirectory + path.sep)) {
throw new BadRequestException(`Unsicherer Pfad im Paket: ${entryName}`);
}
}
this.assertSafeArchive(zip.getEntries(), resolvedDirectory);
// Bestehende Installation entfernen (Update-Szenario).
await rm(directory, { recursive: true, force: true });
@@ -148,6 +146,7 @@ export class ModuleInstaller {
manifest: ModuleManifest,
modulesDir: string,
): Promise<{ directory: string; backupDirectory: string }> {
this.assertCompressedSize(buffer);
const directory = path.join(modulesDir, manifest.id);
const suffix = randomUUID();
const stagingDirectory = path.join(modulesDir, `.update-${manifest.id}-${suffix}`);
@@ -155,14 +154,7 @@ export class ModuleInstaller {
const AdmZip = (await import('adm-zip')).default;
const zip = new AdmZip(buffer);
const resolvedStage = path.resolve(stagingDirectory);
for (const entry of zip.getEntries()) {
const entryName = entry.entryName;
if (entryName.startsWith('/') || entryName.includes('..') || entryName.includes('\\') || /^[A-Za-z]:/.test(entryName)) {
throw new BadRequestException(`Unsicherer Pfad im Paket: ${entryName}`);
}
const resolvedEntry = path.resolve(resolvedStage, entryName);
if (!resolvedEntry.startsWith(resolvedStage + path.sep)) throw new BadRequestException(`Unsicherer Pfad im Paket: ${entryName}`);
}
this.assertSafeArchive(zip.getEntries(), resolvedStage);
try {
await mkdir(stagingDirectory, { recursive: true });
zip.extractAllTo(resolvedStage, true);
@@ -208,4 +200,48 @@ export class ModuleInstaller {
return null;
}
}
private assertCompressedSize(buffer: Buffer): void {
if (buffer.length === 0) throw new BadRequestException('Paket ist leer');
if (buffer.length > MAX_PACKAGE_SIZE_BYTES) {
throw new BadRequestException('Paket ist zu groß (maximal 10 MB)');
}
}
private assertSafeArchive(
entries: ReadonlyArray<{ entryName: string; header: { size: number; attr: number } }>,
destination: string,
): void {
if (entries.length > MAX_ARCHIVE_ENTRIES) {
throw new BadRequestException(`Paket enthält zu viele Dateien (maximal ${MAX_ARCHIVE_ENTRIES})`);
}
let extractedSize = 0;
const seenTargets = new Set<string>();
for (const entry of entries) {
const name = entry.entryName;
if (!name || name.startsWith('/') || name.includes('\\') || name.includes('\0') ||
/^[A-Za-z]:/.test(name) || name.split('/').some((part) => part === '..' || part === '.')) {
throw new BadRequestException(`Unsicherer Pfad im Paket: ${name}`);
}
const resolvedEntry = path.resolve(destination, name);
if (!resolvedEntry.startsWith(destination + path.sep)) {
throw new BadRequestException(`Unsicherer Pfad im Paket: ${name}`);
}
if (seenTargets.has(resolvedEntry)) {
throw new BadRequestException(`Doppelter Pfad im Paket: ${name}`);
}
seenTargets.add(resolvedEntry);
if (((entry.header.attr >>> 16) & 0xf000) === 0xa000) {
throw new BadRequestException(`Symbolischer Link im Paket ist nicht erlaubt: ${name}`);
}
const size = entry.header.size;
if (!Number.isSafeInteger(size) || size < 0 || size > MAX_SINGLE_FILE_BYTES) {
throw new BadRequestException('Paket enthält eine zu große oder ungültige Datei');
}
extractedSize += size;
if (extractedSize > MAX_EXTRACTED_SIZE_BYTES) {
throw new BadRequestException('Entpacktes Paket ist zu groß (maximal 50 MB)');
}
}
}
}

View File

@@ -0,0 +1,19 @@
/** Serializes filesystem and Docker changes for each installed module. */
export class ModuleOperationLock {
private readonly pending = new Map<string, Promise<void>>();
async run<T>(moduleId: string, operation: () => Promise<T>): Promise<T> {
const previous = this.pending.get(moduleId);
let release!: () => void;
const current = new Promise<void>((resolve) => { release = resolve; });
this.pending.set(moduleId, current);
if (previous) await previous;
try {
return await operation();
} finally {
if (this.pending.get(moduleId) === current) this.pending.delete(moduleId);
release();
}
}
}

View File

@@ -52,7 +52,7 @@ export class ModuleProcessManager implements OnModuleDestroy {
// Compose kann beim Build oder beim Start teilweise Container angelegt
// haben. Bereinige den Stack, bevor der ursprüngliche Fehler zurückgeht.
try {
await this.containerManager.remove(module);
await this.containerManager.cleanupFailedStart(module);
} catch {
this.logger.error(`Teilweise gestarteter Container-Stack für "${module.moduleId}" konnte nicht bereinigt werden`);
}

View File

@@ -195,6 +195,7 @@ const TEST_CONFIG: AppConfig = {
adminSeed: { username: 'admin', email: 'admin@example.com', password: 'password-123' },
runtime: { modulesDir: '/data/modules', logsDir: '/data/logs', moduleConfigurationEncryptionKey: '' },
marketplace: { publicUrl: 'http://127.0.0.1:8081', tokenEncryptionKey: '', providers: {} },
modulePublicOrigin: 'http://localhost:8081',
};
describe('ModulesService', () => {

View File

@@ -15,8 +15,9 @@ import { ModuleInstaller } from './module-installer';
import { ModuleProcessManager } from './module-process-manager';
import { ModuleCommandError } from './module-container-manager';
import { ModuleRepository } from './module.repository';
import type { ModuleRecord } from './manifest.types';
import type { ModuleManifest, ModuleRecord } from './manifest.types';
import { ModuleConfigurationService } from './module-configuration.service';
import { ModuleOperationLock } from './module-operation-lock';
/**
* Modul-Verwaltung (Application-Layer): Lifecycle-Logik für Module.
@@ -32,6 +33,7 @@ import { ModuleConfigurationService } from './module-configuration.service';
@Injectable()
export class ModulesService {
private readonly logger = new Logger(ModulesService.name);
private readonly operationLock = new ModuleOperationLock();
constructor(
private readonly moduleRepository: ModuleRepository,
@@ -64,6 +66,15 @@ export class ModulesService {
input: { values?: unknown; clearKeys?: unknown },
actor: ActingUser,
ipAddress: string | null,
) {
return this.withModuleLock(id, () => this.saveConfigurationUnlocked(id, input, actor, ipAddress));
}
private async saveConfigurationUnlocked(
id: string,
input: { values?: unknown; clearKeys?: unknown },
actor: ActingUser,
ipAddress: string | null,
) {
const module = await this.getById(id);
const result = await this.configurationService.save(module, input);
@@ -73,9 +84,11 @@ export class ModulesService {
action: AUDIT_ACTIONS.MODULE_CONFIG_UPDATED,
details: { moduleId: module.moduleId, keys: result.changedKeys },
ipAddress,
}).catch((auditError: unknown) => {
this.logger.error(`Konfiguration von ${module.moduleId} gespeichert, aber Audit konnte nicht gespeichert werden: ${auditError instanceof Error ? auditError.message : String(auditError)}`);
});
if (result.changedKeys.length && module.status === 'RUNNING') {
await this.restart(id, actor, ipAddress);
await this.restartUnlocked(id, actor, ipAddress);
}
return this.configurationService.state(await this.getById(id));
}
@@ -87,7 +100,15 @@ export class ModulesService {
ipAddress: string | null,
): Promise<ModuleRecord> {
const manifest = await this.installer.validatePackage(packageBuffer);
return this.operationLock.run(manifest.id, () => this.installUnlocked(packageBuffer, manifest, actor, ipAddress));
}
private async installUnlocked(
packageBuffer: Buffer,
manifest: ModuleManifest,
actor: ActingUser,
ipAddress: string | null,
): Promise<ModuleRecord> {
const [existingId, existingSlug, existingPort] = await Promise.all([
this.moduleRepository.findByModuleId(manifest.id),
this.moduleRepository.findBySlug(manifest.slug),
@@ -136,6 +157,8 @@ export class ModulesService {
action: AUDIT_ACTIONS.MODULE_INSTALLED,
details: { moduleId: manifest.id, version: manifest.version, slug: manifest.slug },
ipAddress,
}).catch((auditError: unknown) => {
this.logger.error(`Installation von ${manifest.id} erfolgreich, aber Audit konnte nicht gespeichert werden: ${auditError instanceof Error ? auditError.message : String(auditError)}`);
});
return module;
}
@@ -150,6 +173,18 @@ export class ModulesService {
actor: ActingUser,
ipAddress: string | null,
onProgress?: (phase: string, message: string, progress: number) => void,
afterApplied?: () => Promise<void>,
): Promise<ModuleRecord> {
return this.withModuleLock(id, () => this.updateFromMarketplaceUnlocked(id, packageBuffer, actor, ipAddress, onProgress, afterApplied));
}
private async updateFromMarketplaceUnlocked(
id: string,
packageBuffer: Buffer,
actor: ActingUser,
ipAddress: string | null,
onProgress?: (phase: string, message: string, progress: number) => void,
afterApplied?: () => Promise<void>,
): Promise<ModuleRecord> {
onProgress?.('validation', 'Update-Paket wird geprüft', 42);
const current = await this.getById(id);
@@ -166,7 +201,7 @@ export class ModulesService {
const wasRunning = current.status === 'RUNNING';
if (wasRunning) {
onProgress?.('stopping', 'Laufendes Modul wird gestoppt', 55);
await this.stop(id, actor, ipAddress);
await this.stopUnlocked(id, actor, ipAddress);
}
let replacement: { directory: string; backupDirectory: string } | undefined;
@@ -177,7 +212,11 @@ export class ModulesService {
author: manifest.author, configuration: manifest.configuration };
const configuration = await this.configurationService.state(candidate);
await this.moduleRepository.updateManifest(id, manifest, configuration.ready);
if (wasRunning) await this.start(id, actor, ipAddress, onProgress);
if (wasRunning) await this.startUnlocked(id, actor, ipAddress, onProgress);
const updated = await this.getById(id);
// Source metadata belongs to the same serialized operation. A failed
// metadata write still has a backup available for the rollback below.
await afterApplied?.();
await this.installer.finalizeReplacement(replacement.backupDirectory).catch((cleanupError: unknown) => {
this.logger.warn(`Alte Moduldateien für ${current.moduleId} konnten nicht bereinigt werden: ${cleanupError instanceof Error ? cleanupError.message : String(cleanupError)}`);
});
@@ -187,8 +226,10 @@ export class ModulesService {
action: AUDIT_ACTIONS.MODULE_UPDATED,
details: { moduleId: current.moduleId, fromVersion: current.version, toVersion: manifest.version },
ipAddress,
}).catch((auditError: unknown) => {
this.logger.error(`Update von ${current.moduleId} erfolgreich, aber Audit konnte nicht gespeichert werden: ${auditError instanceof Error ? auditError.message : String(auditError)}`);
});
return await this.getById(id);
return updated;
} catch (error) {
if (replacement) {
try {
@@ -198,13 +239,13 @@ export class ModulesService {
await this.installer.rollbackReplacement(replacement.directory, replacement.backupDirectory);
await this.moduleRepository.updateManifest(id, previousManifest, current.configurationReady);
await this.moduleRepository.updateStatus(id, wasRunning ? 'STOPPED' : current.status);
if (wasRunning) await this.start(id, actor, ipAddress);
if (wasRunning) await this.startUnlocked(id, actor, ipAddress);
} catch (rollbackError) {
this.logger.error(`Rollback des Modulupdates für ${current.moduleId} fehlgeschlagen: ${rollbackError instanceof Error ? rollbackError.message : String(rollbackError)}`);
throw new InternalServerErrorException('Update fehlgeschlagen; die vorherige Modulversion konnte nicht vollständig wiederhergestellt werden. Plattform-Logs prüfen.');
}
} else if (wasRunning) {
try { await this.start(id, actor, ipAddress); } catch { /* Preserve the original update error. */ }
try { await this.startUnlocked(id, actor, ipAddress); } catch { /* Preserve the original update error. */ }
}
throw error;
}
@@ -220,6 +261,15 @@ export class ModulesService {
actor: ActingUser,
ipAddress: string | null,
onProgress?: (phase: string, message: string, progress: number) => void,
): Promise<ModuleRecord> {
return this.withModuleLock(id, () => this.startUnlocked(id, actor, ipAddress, onProgress));
}
private async startUnlocked(
id: string,
actor: ActingUser,
ipAddress: string | null,
onProgress?: (phase: string, message: string, progress: number) => void,
): Promise<ModuleRecord> {
const module = await this.getById(id);
this.assertEnabled(module);
@@ -286,6 +336,10 @@ export class ModulesService {
/** Stoppt ein Modul (RUNNING → STOPPING → STOPPED). */
async stop(id: string, actor: ActingUser, ipAddress: string | null): Promise<ModuleRecord> {
return this.withModuleLock(id, () => this.stopUnlocked(id, actor, ipAddress));
}
private async stopUnlocked(id: string, actor: ActingUser, ipAddress: string | null): Promise<ModuleRecord> {
const module = await this.getById(id);
if (module.status === 'STOPPED' || module.status === 'STOPPING') {
return module;
@@ -311,11 +365,15 @@ export class ModulesService {
/** Startet ein Modul neu (Stop + Start). */
async restart(id: string, actor: ActingUser, ipAddress: string | null): Promise<ModuleRecord> {
return this.withModuleLock(id, () => this.restartUnlocked(id, actor, ipAddress));
}
private async restartUnlocked(id: string, actor: ActingUser, ipAddress: string | null): Promise<ModuleRecord> {
const module = await this.getById(id);
if (module.status === 'RUNNING' || module.status === 'STARTING') {
await this.stop(id, actor, ipAddress);
await this.stopUnlocked(id, actor, ipAddress);
}
return this.start(id, actor, ipAddress);
return this.startUnlocked(id, actor, ipAddress);
}
/** Aktiviert oder deaktiviert ein Modul (DISABLED-Zustand). */
@@ -324,6 +382,15 @@ export class ModulesService {
enabled: boolean,
actor: ActingUser,
ipAddress: string | null,
): Promise<ModuleRecord> {
return this.withModuleLock(id, () => this.setEnabledUnlocked(id, enabled, actor, ipAddress));
}
private async setEnabledUnlocked(
id: string,
enabled: boolean,
actor: ActingUser,
ipAddress: string | null,
): Promise<ModuleRecord> {
const module = await this.getById(id);
if (module.enabled === enabled) {
@@ -331,7 +398,7 @@ export class ModulesService {
}
if (!enabled && (module.status === 'RUNNING' || module.status === 'STARTING')) {
await this.stop(id, actor, ipAddress);
await this.stopUnlocked(id, actor, ipAddress);
}
await this.moduleRepository.updateEnabled(id, enabled);
@@ -346,16 +413,22 @@ export class ModulesService {
action: enabled ? AUDIT_ACTIONS.MODULE_ENABLED : AUDIT_ACTIONS.MODULE_DISABLED,
details: { moduleId: module.moduleId },
ipAddress,
}).catch((auditError: unknown) => {
this.logger.error(`Status von ${module.moduleId} geändert, aber Audit konnte nicht gespeichert werden: ${auditError instanceof Error ? auditError.message : String(auditError)}`);
});
return (await this.moduleRepository.findById(id)) ?? module;
}
/** Entfernt ein Modul vollständig (Prozess, Dateien, Registry). */
async remove(id: string, actor: ActingUser, ipAddress: string | null): Promise<{ cleanupWarning?: string }> {
return this.withModuleLock(id, () => this.removeUnlocked(id, actor, ipAddress));
}
private async removeUnlocked(id: string, actor: ActingUser, ipAddress: string | null): Promise<{ cleanupWarning?: string }> {
const module = await this.getById(id);
if (module.status === 'RUNNING' || module.status === 'STARTING') {
await this.stop(id, actor, ipAddress);
await this.stopUnlocked(id, actor, ipAddress);
}
try {
@@ -383,6 +456,8 @@ export class ModulesService {
action: AUDIT_ACTIONS.MODULE_REMOVED,
details: { moduleId: module.moduleId },
ipAddress,
}).catch((auditError: unknown) => {
this.logger.error(`Modul ${module.moduleId} entfernt, aber Audit konnte nicht gespeichert werden: ${auditError instanceof Error ? auditError.message : String(auditError)}`);
});
return cleanupWarning ? { cleanupWarning } : {};
}
@@ -400,6 +475,11 @@ export class ModulesService {
}
}
private async withModuleLock<T>(id: string, operation: () => Promise<T>): Promise<T> {
const module = await this.getById(id);
return this.operationLock.run(module.moduleId, operation);
}
private lifecycleFailure(
action: 'start' | 'stop' | 'remove',
module: ModuleRecord,
@@ -437,6 +517,8 @@ export class ModulesService {
action,
details: { moduleId: module.moduleId, version: module.version },
ipAddress,
}).catch((auditError: unknown) => {
this.logger.error(`Modulaktion ${action} für ${module.moduleId} erfolgreich, aber Audit konnte nicht gespeichert werden: ${auditError instanceof Error ? auditError.message : String(auditError)}`);
});
}
}

View File

@@ -1,6 +1,5 @@
import { UnauthorizedException } from '@nestjs/common';
import { AUDIT_ACTIONS, AuditService } from '../audit/audit.service';
import { SessionService } from '../auth/session.service';
import { PasswordHasher } from './password-hasher';
import { UserRepository } from './user.repository';
import type { UserRecord } from './user.types';
@@ -28,23 +27,14 @@ function createUserRecord(overrides: Partial<UserRecord> = {}): UserRecord {
/** Mock des UserRepository. */
class MockUserRepository {
public user: UserRecord | null = createUserRecord();
public updatedPasswords: Array<{ id: string; hash: string }> = [];
public updatedPasswords: Array<{ id: string; hash: string; exceptSessionId?: string; expectedHash?: string }> = [];
async findById(id: string): Promise<UserRecord | null> {
return this.user && this.user.id === id ? this.user : null;
}
async updatePassword(id: string, hash: string): Promise<void> {
this.updatedPasswords.push({ id, hash });
}
}
/** Mock des SessionService. */
class MockSessionService {
public deletedForUser: Array<{ userId: string; exceptSessionId?: string }> = [];
async deleteAllForUser(userId: string, exceptSessionId?: string): Promise<void> {
this.deletedForUser.push({ userId, exceptSessionId });
async updatePassword(id: string, hash: string, exceptSessionId?: string, expectedHash?: string): Promise<void> {
this.updatedPasswords.push({ id, hash, exceptSessionId, expectedHash });
}
}
@@ -59,20 +49,17 @@ class MockAuditService {
describe('ProfileService', () => {
let userRepository: MockUserRepository;
let sessionService: MockSessionService;
let auditService: MockAuditService;
let profileService: ProfileService;
let passwordHasher: PasswordHasher;
beforeEach(async () => {
userRepository = new MockUserRepository();
sessionService = new MockSessionService();
auditService = new MockAuditService();
passwordHasher = new PasswordHasher();
profileService = new ProfileService(
userRepository as unknown as UserRepository,
passwordHasher,
sessionService as unknown as SessionService,
auditService as unknown as AuditService,
);
@@ -101,9 +88,8 @@ describe('ProfileService', () => {
);
expect(userRepository.updatedPasswords).toHaveLength(1);
expect(sessionService.deletedForUser).toEqual([
{ userId: 'user-1', exceptSessionId: 'session-1' },
]);
expect(userRepository.updatedPasswords[0].exceptSessionId).toBe('session-1');
expect(userRepository.updatedPasswords[0].expectedHash).toBe(userRepository.user?.passwordHash);
expect(auditService.records.at(-1)?.action).toBe(AUDIT_ACTIONS.USER_PASSWORD_CHANGED);
});
@@ -118,4 +104,4 @@ describe('ProfileService', () => {
).rejects.toThrow(UnauthorizedException);
expect(userRepository.updatedPasswords).toHaveLength(0);
});
});
});

View File

@@ -1,6 +1,5 @@
import { Injectable, UnauthorizedException } from '@nestjs/common';
import { AUDIT_ACTIONS, AuditService } from '../audit/audit.service';
import { SessionService } from '../auth/session.service';
import { PasswordHasher } from './password-hasher';
import { UserRepository } from './user.repository';
import type { ChangePasswordDto, UserRecord } from './user.types';
@@ -16,7 +15,6 @@ export class ProfileService {
constructor(
private readonly userRepository: UserRepository,
private readonly passwordHasher: PasswordHasher,
private readonly sessionService: SessionService,
private readonly auditService: AuditService,
) {}
@@ -48,8 +46,7 @@ export class ProfileService {
}
const newPasswordHash = await this.passwordHasher.hash(input.newPassword);
await this.userRepository.updatePassword(userId, newPasswordHash);
await this.sessionService.deleteAllForUser(userId, currentSessionId);
await this.userRepository.updatePassword(userId, newPasswordHash, currentSessionId, user.passwordHash);
await this.auditService.record({
userId,
@@ -58,4 +55,4 @@ export class ProfileService {
ipAddress,
});
}
}
}

View File

@@ -1,4 +1,4 @@
import { BadRequestException, Injectable } from '@nestjs/common';
import { BadRequestException, Injectable, UnauthorizedException } from '@nestjs/common';
import { DatabaseService } from '../database/database.service';
import { PasswordHasher } from './password-hasher';
import type { CreateUserDto, RoleName, UpdateUserDto, UserRecord } from './user.types';
@@ -146,11 +146,31 @@ export class UserRepository {
}
/** Setzt einen neuen Passwort-Hash. */
async updatePassword(id: string, passwordHash: string): Promise<void> {
await this.database.query(
'UPDATE users SET password_hash = $2, updated_at = now() WHERE id = $1',
[id, passwordHash],
);
async updatePassword(
id: string,
passwordHash: string,
exceptSessionId?: string,
expectedPasswordHash?: string,
): Promise<void> {
await this.database.transaction(async (client) => {
// Login hält dieselbe Benutzerzeile bis zur Session-Anlage gesperrt.
const updated = await client.query(
`UPDATE users SET password_hash = $2, updated_at = now()
WHERE id = $1 AND ($3::text IS NULL OR password_hash = $3)`,
[id, passwordHash, expectedPasswordHash ?? null],
);
if (updated.rowCount !== 1) {
throw new UnauthorizedException('Passwort wurde zwischenzeitlich geändert');
}
await client.query(
exceptSessionId
? 'DELETE FROM sessions WHERE user_id = $1 AND id <> $2'
: 'DELETE FROM sessions WHERE user_id = $1',
exceptSessionId ? [id, exceptSessionId] : [id],
);
await client.query('DELETE FROM module_sessions WHERE user_id = $1', [id]);
await client.query('DELETE FROM module_access_tickets WHERE user_id = $1', [id]);
});
}
/** Setzt Fehlversuchs-Zähler und Sperre zurück (bei Aktivierung). */

View File

@@ -225,7 +225,7 @@ describe('UsersService', () => {
null,
);
expect(userRepository.updatedPasswords).toHaveLength(1);
expect(sessionService.deletedSessionsForUser).toEqual(['user-1']);
expect(sessionService.deletedSessionsForUser).toEqual([]);
expect(auditService.records.at(-1)?.action).toBe(AUDIT_ACTIONS.USER_PASSWORD_RESET);
});
});
@@ -250,4 +250,4 @@ describe('UsersService', () => {
).rejects.toThrow(BadRequestException);
});
});
});
});

View File

@@ -135,7 +135,6 @@ export class UsersService {
const user = await this.findById(id);
const passwordHash = await this.passwordHasher.hash(input.newPassword);
await this.userRepository.updatePassword(id, passwordHash);
await this.sessionService.deleteAllForUser(id);
await this.auditService.record({
userId: actor.id,
username: actor.username,
@@ -171,4 +170,4 @@ export class UsersService {
const activeAdminCount = await this.userRepository.countActiveAdmins();
return activeAdminCount <= 1;
}
}
}