301 lines
24 KiB
TypeScript
301 lines
24 KiB
TypeScript
import type { FastifyInstance } from 'fastify';
|
||
import { z } from 'zod';
|
||
import { randomUUID, createHash } from 'node:crypto';
|
||
import { ACTION_CAPABILITY, ACTIONS, loadInstance, type ActionName } from '@kc/connectors';
|
||
import { ConnectorError, type ChildKind } from '@kc/connector-sdk';
|
||
import { decrypt, encrypt, randomToken } from '../../core/crypto.js';
|
||
import { CHILD_MANAGE, CUSTOMER_ACTIONS } from '../../core/actions.js';
|
||
import { rl } from '../../core/config.js';
|
||
import { DESTRUCTIVE_ACTIONS } from '@kc/connector-sdk';
|
||
import { one, query, run } from '../../core/db.js';
|
||
import { audit } from '../../core/audit.js';
|
||
import { enqueue } from '../../core/jobs.js';
|
||
import { clientIp, requireAuth, requirePermission, type AuthContext } from '../../core/auth.js';
|
||
import { AppError, badRequest, forbidden, notFound } from '../../core/errors.js';
|
||
import { can, canInOrg } from '../../core/policy.js';
|
||
import type { KcModule } from '../../core/module.js';
|
||
|
||
const STALE_FACTOR = 3;
|
||
const json = <T>(v: unknown, d: T): T => (v == null ? d : typeof v === 'string' ? JSON.parse(v) : (v as T));
|
||
|
||
/** Ressource + Instanzzustand; "stale" = Provider gestört oder Daten älter als 3 Abgleichintervalle. */
|
||
async function loadResource(id: string) {
|
||
return one('SELECT r.*, i.connector_key, i.name AS instance_name, i.health, i.health_message, i.sync_interval_sec, i.capabilities_json, i.enabled FROM resources r JOIN connector_instances i ON i.id = r.instance_id WHERE r.id = ?', [id]);
|
||
}
|
||
function view(r: any, staff: boolean, a?: AuthContext) {
|
||
const stale = r.health !== 'ok' || Date.now() - new Date(r.synced_at).getTime() > STALE_FACTOR * r.sync_interval_sec * 1000;
|
||
return {
|
||
id: r.id, type: r.type, name: r.name, state: r.state, orgId: r.org_id, validFrom: r.valid_from, validUntil: r.valid_until, syncedAt: r.synced_at, missing: !!r.missing_since,
|
||
canReveal: a ? canReveal(r, a) : undefined,
|
||
canResetPassword: a ? canResetPassword(r, a) : undefined,
|
||
stale, staleReason: r.health !== 'ok' ? (r.health_message ?? 'Der Dienst ist derzeit nicht erreichbar.') : stale ? 'Die Daten sind älter als erwartet.' : null,
|
||
...(staff ? { instance: r.instance_name, connector: r.connector_key, externalRef: r.external_ref } : {}),
|
||
};
|
||
}
|
||
/** Erlaubte Aktionen = Connector-Fähigkeit ∧ Rolle ∧ Ressourcenregel (Kunden nur freigegebene). */
|
||
function allowedActions(r: any, a: AuthContext): ActionName[] {
|
||
const caps = json<string[]>(r.capabilities_json, []);
|
||
const supported = ACTIONS.filter((x) => caps.includes(ACTION_CAPABILITY[x]) && !DESTRUCTIVE_ACTIONS.has(x));
|
||
if (r.health !== 'ok' || !r.enabled) return [];
|
||
if (can(a.principal, 'resources.write')) return supported;
|
||
const orgOk = r.org_id && canInOrg(a.principal, r.org_id, 'resources.manage', 'resources.write');
|
||
return orgOk ? supported.filter((x) => json<string[]>(r.customer_actions, []).includes(x)) : [];
|
||
}
|
||
const childCapsTop = (r: any): string[] => json<string[]>(r.capabilities_json, []);
|
||
const canLoginTop = (r: any, a: AuthContext): boolean => childCapsTop(r).includes('sso.login') && r.health === 'ok' && !!r.enabled && (can(a.principal, 'resources.write') || (!!r.org_id && canInOrg(a.principal, r.org_id, 'resources.manage', 'resources.write') && json<string[]>(r.customer_actions, []).includes('panel.login')));
|
||
/** Zugangsdaten/Schlüssel anzeigen: Personal mit Schreibrecht oder Inhaber/Admin der zugehörigen Organisation; Anbieter muss es unterstützen und erreichbar sein. */
|
||
function canReveal(r: any, a: AuthContext): boolean {
|
||
const caps = json<string[]>(r.capabilities_json, []);
|
||
if (!caps.includes('secret.reveal') || !r.enabled) return false;
|
||
return can(a.principal, 'resources.write') || (!!r.org_id && canInOrg(a.principal, r.org_id, 'resources.manage', 'resources.write'));
|
||
}
|
||
/** Panel-Passwort neu vergeben/einsehen: nur Personal mit Schreibrecht (kein Kunden-Selfservice, ändert ein echtes Zugangsdatum). */
|
||
function canResetPassword(r: any, a: AuthContext): boolean {
|
||
const caps = json<string[]>(r.capabilities_json, []);
|
||
return caps.includes('panel.password_reset') && !!r.enabled && r.health === 'ok' && can(a.principal, 'resources.write');
|
||
}
|
||
function access(a: AuthContext, r: any): boolean {
|
||
return can(a.principal, 'resources.read') || (!!r.org_id && canInOrg(a.principal, r.org_id, 'resources.read', 'resources.read'));
|
||
}
|
||
|
||
export const resourcesModule: KcModule = {
|
||
name: 'resources',
|
||
permissions: {
|
||
staff: { support: ['resources.read'], accounting: ['resources.read'], admin: ['resources.read', 'resources.write'], superadmin: ['resources.read', 'resources.write'] },
|
||
org: { owner: ['resources.read', 'resources.manage'], admin: ['resources.read', 'resources.manage'], member: ['resources.read'] },
|
||
},
|
||
register(app: FastifyInstance) {
|
||
app.get('/resources', async (req) => {
|
||
const a = requireAuth(req);
|
||
const q = z.object({ org: z.string().uuid().optional(), unassigned: z.enum(['1']).optional(), type: z.string().max(30).optional() }).parse(req.query);
|
||
const staff = can(a.principal, 'resources.read');
|
||
const orgs = staff ? (q.org ? [q.org] : null) : a.principal.memberships.map((m) => m.orgId);
|
||
if (orgs && orgs.length === 0) return [];
|
||
const where: string[] = []; const params: unknown[] = [];
|
||
if (orgs) { where.push(`r.org_id IN (${orgs.map(() => '?').join(',')})`); params.push(...orgs); }
|
||
if (q.unassigned && staff) where.push('r.org_id IS NULL');
|
||
if (q.type) { where.push('r.type = ?'); params.push(q.type); }
|
||
const rows = await query(`SELECT r.*, i.connector_key, i.name AS instance_name, i.health, i.health_message, i.sync_interval_sec, i.capabilities_json, i.enabled FROM resources r JOIN connector_instances i ON i.id = r.instance_id ${where.length ? 'WHERE ' + where.join(' AND ') : ''} ORDER BY r.name LIMIT 500`, params);
|
||
return rows.map((r) => view(r, staff, a));
|
||
});
|
||
|
||
app.get('/resources/:id', async (req) => {
|
||
const a = requireAuth(req);
|
||
const { id } = z.object({ id: z.string().uuid() }).parse(req.params);
|
||
const r = await loadResource(id);
|
||
if (!r || !access(a, r)) throw notFound();
|
||
const staff = can(a.principal, 'resources.read');
|
||
const jobs = await query("SELECT id, status, last_error, attempts, created_at, updated_at, JSON_VALUE(payload, '$.action') AS action FROM jobs WHERE type = 'connector.execute' AND JSON_VALUE(payload, '$.resourceId') = ? ORDER BY created_at DESC LIMIT 10", [id]);
|
||
const data = json<{ limits?: object; usage?: object; details?: object }>(r.data_json, {});
|
||
return { ...view(r, staff), limits: data.limits ?? {}, usage: data.usage ?? {}, details: data.details ?? {}, allowedActions: allowedActions(r, a), canReveal: canReveal(r, a), canResetPassword: canResetPassword(r, a), hasChildren: childCapsTop(r).includes('children.read'), canLogin: canLoginTop(r, a), customerActions: staff ? json(r.customer_actions, []) : undefined, jobs: jobs.map((j) => ({ id: j.id, action: j.action, status: j.status, error: j.last_error, attempts: j.attempts, createdAt: j.created_at, updatedAt: j.updated_at })) };
|
||
});
|
||
|
||
app.post('/resources/:id/actions', async (req, reply) => {
|
||
const a = requireAuth(req);
|
||
const { id } = z.object({ id: z.string().uuid() }).parse(req.params);
|
||
const b = z.object({ action: z.enum(['suspend', 'unsuspend', 'extend']), days: z.number().int().min(1).max(3650).optional() }).parse(req.body);
|
||
const r = await loadResource(id);
|
||
if (!r || !access(a, r)) throw notFound();
|
||
if (!allowedActions(r, a).includes(b.action)) throw forbidden('Diese Aktion ist für diese Ressource nicht verfügbar', 'ACTION_NOT_ALLOWED');
|
||
const params: Record<string, unknown> = {};
|
||
if (b.action === 'extend') {
|
||
if (!b.days) throw badRequest('Anzahl Tage fehlt');
|
||
const from = r.valid_until && new Date(r.valid_until) > new Date() ? new Date(r.valid_until) : new Date();
|
||
params.until = new Date(from.getTime() + b.days * 86400000).toISOString(); // absoluter Zielwert => bei Wiederholung wirkungsgleich
|
||
params.days = b.days;
|
||
}
|
||
// Idempotenz: Header vom Client (pro Bestätigungsdialog), sonst Hash aus Ressource/Aktion/Parameter im Minutenfenster
|
||
const hdr = req.headers['idempotency-key'];
|
||
const key = `act:${typeof hdr === 'string' && /^[\w-]{8,100}$/.test(hdr) ? hdr : createHash('sha256').update(`${id}|${b.action}|${JSON.stringify(params)}|${Math.floor(Date.now() / 60000)}`).digest('hex').slice(0, 40)}`;
|
||
const existing = await one('SELECT id FROM jobs WHERE idempotency_key = ?', [key]);
|
||
const jobId = existing?.id ?? randomUUID();
|
||
if (!existing) {
|
||
try {
|
||
await run('INSERT INTO jobs (id, type, payload, idempotency_key, correlation_id) VALUES (?,?,?,?,?)', [jobId, 'connector.execute', JSON.stringify({ resourceId: id, action: b.action, params, actorUserId: a.user.id }), key, req.correlationId]);
|
||
} catch (e) {
|
||
if ((e as { code?: string }).code !== 'ER_DUP_ENTRY') throw e; // parallele Doppelanfrage: bestehenden Auftrag zurückgeben
|
||
const dup = await one('SELECT id FROM jobs WHERE idempotency_key = ?', [key]);
|
||
return reply.code(202).send({ jobId: dup!.id, duplicate: true });
|
||
}
|
||
await audit({ actorType: 'user', actorId: a.user.id, orgId: r.org_id, action: `resource.${b.action}.request`, resourceType: 'resource', resourceId: id, connector: r.connector_key, correlationId: req.correlationId, ip: clientIp(req), after: { jobId, params } });
|
||
}
|
||
return reply.code(202).send({ jobId, duplicate: !!existing });
|
||
});
|
||
|
||
/** Schlüssel/Zugangsdaten auf Abruf: live beim Anbieter gelesen, nie gespeichert oder protokolliert (nur DASS abgerufen wurde). */
|
||
app.post('/resources/:id/reveal', { config: rl(10, '1 minute') }, async (req, reply) => {
|
||
const a = requireAuth(req);
|
||
const { id } = z.object({ id: z.string().uuid() }).parse(req.params);
|
||
const r = await loadResource(id);
|
||
if (!r || !access(a, r)) throw notFound();
|
||
if (!canReveal(r, a)) throw forbidden('Der Zugriff auf Zugangsdaten ist für diese Ressource nicht möglich', 'REVEAL_FORBIDDEN');
|
||
let items: { label: string; value: string }[];
|
||
try {
|
||
const { connector, ctx } = await loadInstance(r.instance_id, req.correlationId);
|
||
if (!connector.reveal) throw badRequest('Nicht unterstützt', 'NOT_SUPPORTED');
|
||
items = await connector.reveal(ctx, r.external_ref);
|
||
} catch (e) {
|
||
if (e instanceof ConnectorError) throw new AppError(502, 'CONNECTOR_ERROR', `Beim Anbieter konnte nichts gelesen werden: ${e.userMessage}`);
|
||
throw e;
|
||
}
|
||
await audit({ actorType: 'user', actorId: a.user.id, orgId: r.org_id, action: 'resource.reveal', resourceType: 'resource', resourceId: id, connector: r.connector_key, correlationId: req.correlationId, ip: clientIp(req), after: { labels: items.map((i) => i.label) } });
|
||
reply.header('cache-control', 'no-store');
|
||
return { items, hideAfterSec: 60 };
|
||
});
|
||
|
||
// ---- Panel-Zugangsdaten: bestehende Passwörter sind bei Anbietern grundsätzlich nicht lesbar; hier wird
|
||
// stattdessen ein NEUES Passwort erzeugt, beim Anbieter gesetzt und bei uns verschlüsselt hinterlegt. ----
|
||
app.get('/resources/:id/panel-credentials', async (req) => {
|
||
const a = requireAuth(req); const { id } = z.object({ id: z.string().uuid() }).parse(req.params);
|
||
const r = await loadResource(id); if (!r || !access(a, r)) throw notFound();
|
||
const c = await one('SELECT set_at FROM panel_credentials WHERE resource_id = ?', [id]);
|
||
return { canReset: canResetPassword(r, a), set: !!c, setAt: c?.set_at ?? null };
|
||
});
|
||
app.post('/resources/:id/panel-credentials/reset', async (req, reply) => {
|
||
const a = requireAuth(req); const { id } = z.object({ id: z.string().uuid() }).parse(req.params);
|
||
const r = await loadResource(id); if (!r || !access(a, r)) throw notFound();
|
||
if (!canResetPassword(r, a)) throw forbidden('Zugangsdaten können für diese Ressource nicht neu vergeben werden', 'RESET_FORBIDDEN');
|
||
const password = randomToken(15); // ~20 Zeichen, base64url – druckbar, keine Sonderzeichen, die Formulare/Shells stören könnten
|
||
const secretEnc = encrypt(JSON.stringify({ password }));
|
||
const hdr = req.headers['idempotency-key'];
|
||
const key = `pwreset:${typeof hdr === 'string' && /^[\w-]{8,100}$/.test(hdr) ? hdr : createHash('sha256').update(`${id}|${Math.floor(Date.now() / 10000)}`).digest('hex').slice(0, 40)}`;
|
||
const existing = await one('SELECT id FROM jobs WHERE idempotency_key = ?', [key]);
|
||
const jobId = existing?.id ?? randomUUID();
|
||
if (!existing) {
|
||
await run('INSERT INTO jobs (id, type, payload, idempotency_key, correlation_id) VALUES (?,?,?,?,?)',
|
||
[jobId, 'connector.execute', JSON.stringify({ resourceId: id, action: 'reset_password', secretEnc, actorUserId: a.user.id, destructive: true }), key, req.correlationId]);
|
||
await audit({ actorType: 'user', actorId: a.user.id, orgId: r.org_id, action: 'resource.reset_password.request', resourceType: 'resource', resourceId: id, connector: r.connector_key, correlationId: req.correlationId, ip: clientIp(req), after: { jobId } });
|
||
}
|
||
return reply.code(202).send({ jobId, duplicate: !!existing });
|
||
});
|
||
app.post('/resources/:id/panel-credentials/reveal', { config: rl(10, '1 minute') }, async (req, reply) => {
|
||
const a = requireAuth(req); const { id } = z.object({ id: z.string().uuid() }).parse(req.params);
|
||
const r = await loadResource(id); if (!r || !access(a, r)) throw notFound();
|
||
if (!canResetPassword(r, a)) throw forbidden('Zugangsdaten können für diese Ressource nicht eingesehen werden', 'REVEAL_FORBIDDEN');
|
||
const c = await one('SELECT secret_enc FROM panel_credentials WHERE resource_id = ?', [id]);
|
||
if (!c) throw notFound('Es wurde noch kein Passwort vergeben.');
|
||
const password = decrypt(c.secret_enc);
|
||
await audit({ actorType: 'user', actorId: a.user.id, orgId: r.org_id, action: 'resource.panel_credentials.reveal', resourceType: 'resource', resourceId: id, connector: r.connector_key, correlationId: req.correlationId, ip: clientIp(req) });
|
||
reply.header('cache-control', 'no-store');
|
||
return { items: [{ label: 'Panel-Passwort', value: password }], hideAfterSec: 60 };
|
||
});
|
||
|
||
// ---- Hosting: Unterobjekte (Domains, Postfächer, Datenbanken, FTP, SSL) und Panel-Login ----------
|
||
const KINDS = ['domain', 'email', 'database', 'ftp', 'certificate'] as const;
|
||
const kindParam = z.object({ id: z.string().uuid(), kind: z.enum(KINDS) });
|
||
const childCaps = (r: any): string[] => json<string[]>(r.capabilities_json, []);
|
||
/** Schreiben erlaubt: Personal mit Schreibrecht ODER Inhaber/Admin der Organisation, wenn das Produkt es für diese Art freigibt. */
|
||
const canManageKind = (r: any, a: AuthContext, kind: string): boolean => {
|
||
if (kind === 'certificate') return false;
|
||
if (r.health !== 'ok' || !r.enabled || !childCaps(r).includes('children.write')) return false;
|
||
if (can(a.principal, 'resources.write')) return true;
|
||
const need = CHILD_MANAGE[kind];
|
||
return !!need && !!r.org_id && canInOrg(a.principal, r.org_id, 'resources.manage', 'resources.write') && json<string[]>(r.customer_actions, []).includes(need);
|
||
};
|
||
const providerError = (e: unknown): never => { if (e instanceof ConnectorError) throw new AppError(e.code === 'NOT_FOUND' ? 404 : 502, 'CONNECTOR_ERROR', e.code === 'INVALID_INPUT' ? e.message : `Der Anbieter meldet: ${e.userMessage}`); throw e; };
|
||
|
||
app.get('/resources/:id/children', async (req) => {
|
||
const a = requireAuth(req);
|
||
const { id } = z.object({ id: z.string().uuid() }).parse(req.params);
|
||
const r = await loadResource(id);
|
||
if (!r || !access(a, r)) throw notFound();
|
||
if (!childCaps(r).includes('children.read')) return { kinds: [], panelLogin: false };
|
||
try {
|
||
const { connector, ctx } = await loadInstance(r.instance_id, req.correlationId);
|
||
const kinds = await connector.children!.kinds(ctx, r.external_ref);
|
||
return { kinds: kinds.map((k) => ({ kind: k.kind, canWrite: k.canWrite && canManageKind(r, a, k.kind) })), panelLogin: canLogin(r, a) };
|
||
} catch (e) { return providerError(e); }
|
||
});
|
||
app.get('/resources/:id/children/:kind', { config: rl(60, '1 minute') }, async (req) => {
|
||
const a = requireAuth(req);
|
||
const { id, kind } = kindParam.parse(req.params);
|
||
const r = await loadResource(id);
|
||
if (!r || !access(a, r) || !childCaps(r).includes('children.read')) throw notFound();
|
||
try { const { connector, ctx } = await loadInstance(r.instance_id, req.correlationId); return await connector.children!.list(ctx, r.external_ref, kind as ChildKind); }
|
||
catch (e) { return providerError(e); }
|
||
});
|
||
/** Änderung an einem Unterobjekt: läuft als persistenter Auftrag; Passwörter nur verschlüsselt im Auftrag, nach der Ausführung entfernt. */
|
||
app.post('/resources/:id/children/:kind', { config: rl(30, '1 minute') }, async (req, reply) => {
|
||
const a = requireAuth(req);
|
||
const { id, kind } = kindParam.parse(req.params);
|
||
const b = z.object({ op: z.enum(['create', 'update', 'delete']), id: z.string().max(40).optional(), data: z.record(z.string(), z.unknown()).default({}), password: z.string().max(128).optional() }).parse(req.body);
|
||
const r = await loadResource(id);
|
||
if (!r || !access(a, r)) throw notFound();
|
||
if (!canManageKind(r, a, kind)) throw forbidden('Diese Änderung ist für diese Ressource nicht möglich', 'ACTION_NOT_ALLOWED');
|
||
if ((b.op === 'update' || b.op === 'delete') && !b.id) throw badRequest('Objekt fehlt', 'ID_REQUIRED');
|
||
const needsPw = b.op === 'create' && ['email', 'database', 'ftp'].includes(kind);
|
||
if (needsPw && !b.password) throw badRequest('Bitte ein Passwort angeben.', 'PASSWORD_REQUIRED');
|
||
if (b.password && (b.password.length < 12 || /[\0\r\n]/.test(b.password))) throw badRequest('Das Passwort muss mindestens 12 Zeichen lang sein.', 'WEAK_PASSWORD');
|
||
if (JSON.stringify(b.data).length > 4000) throw badRequest('Eingabe zu groß', 'TOO_LARGE');
|
||
const hdr = req.headers['idempotency-key'];
|
||
const key = `child:${typeof hdr === 'string' && /^[\w-]{8,100}$/.test(hdr) ? hdr : createHash('sha256').update(`${id}|${kind}|${JSON.stringify(b.data)}|${b.op}|${b.id ?? ''}|${Math.floor(Date.now() / 60000)}`).digest('hex').slice(0, 40)}`;
|
||
const existing = await one('SELECT id FROM jobs WHERE idempotency_key = ?', [key]);
|
||
if (existing) return reply.code(202).send({ jobId: existing.id, duplicate: true });
|
||
const jobId = randomUUID();
|
||
const payload = { resourceId: id, kind, op: b.op, id: b.id, data: b.data, actorUserId: a.user.id, destructive: b.op === 'delete', ...(b.password ? { secretEnc: encrypt(JSON.stringify({ password: b.password })) } : {}) };
|
||
try { await run('INSERT INTO jobs (id, type, payload, idempotency_key, correlation_id) VALUES (?,?,?,?,?)', [jobId, 'connector.child', JSON.stringify(payload), key, req.correlationId]); }
|
||
catch (e) { if ((e as { code?: string }).code !== 'ER_DUP_ENTRY') throw e; const d = await one('SELECT id FROM jobs WHERE idempotency_key = ?', [key]); return reply.code(202).send({ jobId: d!.id, duplicate: true }); }
|
||
await audit({ actorType: 'user', actorId: a.user.id, orgId: r.org_id, action: `resource.child.${kind}.${b.op}.request`, resourceType: 'resource', resourceId: id, connector: r.connector_key, correlationId: req.correlationId, ip: clientIp(req), after: { jobId, id: b.id, data: b.data } });
|
||
return reply.code(202).send({ jobId, duplicate: false });
|
||
});
|
||
|
||
const canLogin = (r: any, a: AuthContext): boolean => {
|
||
if (!childCaps(r).includes('sso.login') || r.health !== 'ok' || !r.enabled) return false;
|
||
if (can(a.principal, 'resources.write')) return true;
|
||
return !!r.org_id && canInOrg(a.principal, r.org_id, 'resources.manage', 'resources.write') && json<string[]>(r.customer_actions, []).includes('panel.login');
|
||
};
|
||
/** Panel-Login: kurzlebiger Link, nie gespeichert, Abruf im Audit. */
|
||
app.post('/resources/:id/login', { config: rl(10, '1 minute') }, async (req, reply) => {
|
||
const a = requireAuth(req);
|
||
const { id } = z.object({ id: z.string().uuid() }).parse(req.params);
|
||
const r = await loadResource(id);
|
||
if (!r || !access(a, r)) throw notFound();
|
||
if (!canLogin(r, a)) throw forbidden('Der Panel-Login ist für diese Ressource nicht möglich', 'LOGIN_FORBIDDEN');
|
||
let out: { url: string; validForSec: number };
|
||
try { const { connector, ctx } = await loadInstance(r.instance_id, req.correlationId); out = await connector.loginUrl!(ctx, r.external_ref); } catch (e) { return providerError(e); }
|
||
await audit({ actorType: 'user', actorId: a.user.id, orgId: r.org_id, action: 'resource.login', resourceType: 'resource', resourceId: id, connector: r.connector_key, correlationId: req.correlationId, ip: clientIp(req) });
|
||
reply.header('cache-control', 'no-store'); return out;
|
||
});
|
||
|
||
app.get('/jobs/:id', async (req) => {
|
||
const a = requireAuth(req);
|
||
const { id } = z.object({ id: z.string().uuid() }).parse(req.params);
|
||
const j = await one("SELECT id, type, status, attempts, max_attempts, last_error, created_at, updated_at, JSON_VALUE(payload, '$.resourceId') AS resource_id FROM jobs WHERE id = ?", [id]);
|
||
if (!j || !j.resource_id) throw notFound();
|
||
const r = await loadResource(j.resource_id);
|
||
if (!r || !access(a, r)) throw notFound();
|
||
return { id: j.id, status: j.status, attempts: j.attempts, maxAttempts: j.max_attempts, error: j.last_error, createdAt: j.created_at, updatedAt: j.updated_at };
|
||
});
|
||
|
||
app.get('/admin/jobs', async (req) => {
|
||
requirePermission(req, 'jobs.read');
|
||
const q = z.object({ status: z.string().max(30).optional() }).parse(req.query);
|
||
return (await query('SELECT id, type, status, attempts, max_attempts, last_error, run_at, created_at, updated_at, correlation_id FROM jobs WHERE (? IS NULL OR status = ?) ORDER BY created_at DESC LIMIT 200', [q.status ?? null, q.status ?? null]))
|
||
.map((j) => ({ id: j.id, type: j.type, status: j.status, attempts: j.attempts, maxAttempts: j.max_attempts, error: j.last_error, runAt: j.run_at, createdAt: j.created_at, correlationId: j.correlation_id }));
|
||
});
|
||
/** Kontrollierte manuelle Wiederholung: nur für Jobs in needs_review/failed, nie für destruktive Aktionen. */
|
||
app.post('/admin/jobs/:id/retry', async (req) => {
|
||
const a = requirePermission(req, 'resources.write');
|
||
const { id } = z.object({ id: z.string().uuid() }).parse(req.params);
|
||
const j = await one("SELECT id, status, (JSON_VALUE(payload, '$.destructive') IN ('1','true') OR JSON_VALUE(payload, '$.action') = 'terminate') AS destructive FROM jobs WHERE id = ?", [id]);
|
||
if (!j) throw notFound();
|
||
if (!['needs_review', 'failed'].includes(j.status)) throw badRequest('Nur fehlgeschlagene Aufträge können wiederholt werden', 'NOT_RETRYABLE');
|
||
if (Number(j.destructive) === 1) throw forbidden('Destruktive Aufträge werden nicht automatisch wiederholt; bitte Ergebnis beim Provider prüfen', 'DESTRUCTIVE');
|
||
await run("UPDATE jobs SET status='retrying', attempts=0, run_at=UTC_TIMESTAMP(3), last_error=NULL WHERE id = ?", [id]);
|
||
await audit({ actorType: 'user', actorId: a.user.id, action: 'job.retry', resourceType: 'job', resourceId: id, correlationId: req.correlationId, ip: clientIp(req) });
|
||
return { status: 'ok' };
|
||
});
|
||
|
||
app.patch('/admin/resources/:id', async (req) => {
|
||
const a = requirePermission(req, 'resources.write');
|
||
const { id } = z.object({ id: z.string().uuid() }).parse(req.params);
|
||
const b = z.object({ orgId: z.string().uuid().nullable().optional(), customerActions: z.array(z.enum(CUSTOMER_ACTIONS)).optional() }).parse(req.body);
|
||
const r = await loadResource(id);
|
||
if (!r) throw notFound();
|
||
if (b.orgId && !(await one('SELECT 1 AS x FROM organizations WHERE id = ?', [b.orgId]))) throw badRequest('Kunde nicht gefunden');
|
||
await run('UPDATE resources SET org_id = ?, customer_actions = ? WHERE id = ?', [b.orgId === undefined ? r.org_id : b.orgId, JSON.stringify(b.customerActions ?? json(r.customer_actions, [])), id]);
|
||
await audit({ actorType: 'user', actorId: a.user.id, orgId: b.orgId ?? r.org_id, action: 'resource.update', resourceType: 'resource', resourceId: id, connector: r.connector_key, correlationId: req.correlationId, ip: clientIp(req), before: { orgId: r.org_id, customerActions: json(r.customer_actions, []) }, after: b });
|
||
return { status: 'ok' };
|
||
});
|
||
},
|
||
};
|