kundencenter/packages/connectors/src/index.ts

254 lines
22 KiB
TypeScript
Raw Normal View History

import { randomUUID } from 'node:crypto';
import { ACTION_CAPABILITY, ConnectorError, DESTRUCTIVE_ACTIONS, type ActionName, type Capability, type ChildKind, type ChildOp, type Connector, type ConnectorContext, type NormalizedResource } from '@kc/connector-sdk';
import { licensingConnector } from '@kc/connector-licensing';
import { mockConnector } from '@kc/connector-mock';
import { keyhelpConnector } from '@kc/connector-keyhelp';
import { audit } from '@kc/platform/audit';
import { decrypt, encrypt } from '@kc/platform/crypto';
import { one, query, run } from '@kc/platform/db';
import { JobFailure, enqueue } from '@kc/platform/jobs';
import { addMonths } from '@kc/platform/contractterms';
import { CONTRACT_MACHINE, ORDER_MACHINE, transition } from '@kc/platform/statemachine';
/** Registry: neue Connectoren (KeyHelp, Plesk, …) werden hier eingetragen, sonst nirgends. */
const registry = new Map<string, Connector>([licensingConnector, keyhelpConnector, mockConnector].map((c) => [c.key, c]));
export const listConnectors = (): Connector[] => [...registry.values()];
export const getConnector = (key: string): Connector => {
const c = registry.get(key);
if (!c) throw new ConnectorError('BAD_CONFIG', `Unbekannter Connector ${key}`);
return c;
};
export const encryptSecrets = (s: Record<string, string>): string | null => (Object.keys(s).length ? encrypt(JSON.stringify(s)) : null);
export const ACTIONS: readonly ActionName[] = ['suspend', 'unsuspend', 'extend', 'terminate'];
interface Loaded { inst: any; connector: Connector; ctx: ConnectorContext }
export async function loadInstance(id: string, correlationId: string): Promise<Loaded> {
const inst = await one('SELECT * FROM connector_instances WHERE id = ?', [id]);
if (!inst) throw new JobFailure('Connector-Instanz nicht gefunden', false, 'failed');
const config = typeof inst.config_json === 'string' ? JSON.parse(inst.config_json) : inst.config_json;
const secrets = inst.secrets_enc ? JSON.parse(decrypt(inst.secrets_enc)) : {};
return { inst, connector: getConnector(inst.connector_key), ctx: { config, secrets, correlationId } };
}
const iso = (d?: string | null) => (d ? new Date(d) : null);
/** Gleicht eine Instanz mit dem Provider ab. Bei Ausfall bleiben die zuletzt bekannten Daten erhalten (mit altem Zeitstempel). */
export async function syncInstance(instanceId: string, correlationId: string): Promise<{ ok: boolean; count?: number; error?: string }> {
const { inst, connector, ctx } = await loadInstance(instanceId, correlationId);
const wasDown = inst.health === 'down';
await run('UPDATE connector_instances SET last_sync_at = UTC_TIMESTAMP(3) WHERE id = ?', [instanceId]);
try {
const health = await connector.healthCheck(ctx);
if (!health.ok) throw new ConnectorError(/unvollständig|API-Schlüssel|Fingerabdruck/.test(health.message ?? '') ? 'BAD_CONFIG' : 'UNREACHABLE', health.message);
const [caps, list] = [await connector.capabilities(ctx), await connector.listResources(ctx)];
const seen = new Set<string>();
for (const r of list) { seen.add(r.externalRef); await upsert(instanceId, r); }
await run('UPDATE resources SET missing_since = COALESCE(missing_since, UTC_TIMESTAMP(3)) WHERE instance_id = ?' + (seen.size ? ` AND external_ref NOT IN (${[...seen].map(() => '?').join(',')})` : ''), [instanceId, ...seen]);
await run("UPDATE connector_instances SET health='ok', health_message=NULL, capabilities_json=?, last_ok_at=UTC_TIMESTAMP(3), last_error=NULL WHERE id = ?", [JSON.stringify(caps), instanceId]);
if (wasDown) await audit({ actorType: 'system', action: 'connector.recovered', resourceType: 'connector', resourceId: instanceId, connector: inst.connector_key, correlationId });
return { ok: true, count: list.length };
} catch (e) {
const msg = e instanceof ConnectorError ? e.userMessage : 'Unerwarteter Fehler beim Abgleich';
let tech = String((e as Error).message);
for (let n = 0; n < 3; n++) { const m = /^[^()]{5,80}\. \((.*)\)$/s.exec(tech); if (!m) break; tech = m[1]!; }
tech = tech.slice(0, 250);
const hint = e instanceof ConnectorError ? e.hint : '';
await run("UPDATE connector_instances SET health='down', health_message=?, last_error=?, last_error_at=UTC_TIMESTAMP(3) WHERE id = ?", [msg, JSON.stringify({ tech, hint }).slice(0, 490), instanceId]);
if (!wasDown) await audit({ actorType: 'system', action: 'connector.down', resourceType: 'connector', resourceId: instanceId, connector: inst.connector_key, result: 'failure', errorClass: e instanceof ConnectorError ? e.code : 'unexpected', correlationId });
return { ok: false, error: msg };
}
}
async function upsert(instanceId: string, r: NormalizedResource): Promise<string> {
await run(
`INSERT INTO resources (id, instance_id, external_ref, type, name, state, valid_from, valid_until, data_json, synced_at)
VALUES (?,?,?,?,?,?,?,?,?, UTC_TIMESTAMP(3))
ON DUPLICATE KEY UPDATE type=VALUES(type), name=VALUES(name), state=VALUES(state), valid_from=VALUES(valid_from), valid_until=VALUES(valid_until), data_json=VALUES(data_json), synced_at=UTC_TIMESTAMP(3), missing_since=NULL`,
[randomUUID(), instanceId, r.externalRef, r.type, r.name.slice(0, 300), r.state, iso(r.validFrom), iso(r.validUntil), JSON.stringify({ limits: r.limits ?? {}, usage: r.usage ?? {}, details: r.details ?? {} })],
);
return (await one('SELECT id FROM resources WHERE instance_id = ? AND external_ref = ?', [instanceId, r.externalRef]))!.id as string;
}
export interface ExecutePayload { resourceId: string; action: ActionName; params?: Record<string, unknown>; actorUserId: string | null; destructive?: boolean }
/** Führt eine Aktion aus. Aufgerufen vom Worker im Rahmen eines persistenten Auftrags; jobId dient als Idempotenzschlüssel. */
export async function executeAction(p: ExecutePayload, jobId: string, correlationId: string): Promise<string> {
const r = await one('SELECT * FROM resources WHERE id = ?', [p.resourceId]);
if (!r) throw new JobFailure('Ressource nicht gefunden', false, 'failed');
const { inst, connector, ctx } = await loadInstance(r.instance_id, correlationId);
const caps = await connector.capabilities(ctx);
if (!caps.includes(ACTION_CAPABILITY[p.action] as Capability) || !connector.execute) throw new JobFailure('Aktion wird vom Connector nicht unterstützt', false, 'failed');
const before = { state: r.state, validUntil: r.valid_until };
try {
const out = await connector.execute(ctx, { action: p.action, externalRef: r.external_ref, params: p.params, idempotencyKey: jobId });
if (out.resource) await upsert(r.instance_id, out.resource);
const after = await one('SELECT state, valid_until FROM resources WHERE id = ?', [p.resourceId]);
await audit({ actorType: p.actorUserId ? 'user' : 'system', actorId: p.actorUserId, orgId: r.org_id, action: `resource.${p.action}`, resourceType: 'resource', resourceId: p.resourceId, connector: inst.connector_key, correlationId, before, after: { state: after?.state, validUntil: after?.valid_until, params: p.params } });
return `${p.action}: ${after?.state ?? 'ok'}`;
} catch (e) {
if (e instanceof JobFailure) throw e;
if (e instanceof ConnectorError) {
await audit({ actorType: p.actorUserId ? 'user' : 'system', actorId: p.actorUserId, orgId: r.org_id, action: `resource.${p.action}`, resourceType: 'resource', resourceId: p.resourceId, connector: inst.connector_key, result: 'failure', errorClass: e.code, correlationId });
// Destruktive Aktionen und nicht wiederholbare Fehler nie automatisch wiederholen.
if (DESTRUCTIVE_ACTIONS.has(p.action) || !e.retryable) throw new JobFailure(e.userMessage, false, DESTRUCTIVE_ACTIONS.has(p.action) ? 'needs_review' : 'failed');
throw new JobFailure(e.userMessage, true, 'needs_review', e.retryAfterMs ? Math.ceil(e.retryAfterMs / 1000) : undefined);
}
throw e;
}
}
/** Plant Abgleiche für fällige Instanzen (idempotent pro Zeitfenster). Wird regelmäßig vom Worker aufgerufen. */
export async function scheduleDueSyncs(enqueue: (type: string, payload: unknown, o: { idempotencyKey: string }) => Promise<void>): Promise<number> {
const due = await query("SELECT id, sync_interval_sec FROM connector_instances WHERE enabled = 1 AND (last_sync_at IS NULL OR last_sync_at < DATE_SUB(UTC_TIMESTAMP(3), INTERVAL sync_interval_sec SECOND))");
for (const d of due) await enqueue('connector.sync', { instanceId: d.id }, { idempotencyKey: `sync:${d.id}:${Math.floor(Date.now() / (Number(d.sync_interval_sec) * 1000))}` });
return due.length;
}
export { ACTION_CAPABILITY, DESTRUCTIVE_ACTIONS };
export type { ActionName, Capability };
export interface JobMeta { id: string; correlationId: string; attempt: number; maxAttempts: number }
/**
* Provisioniert eine freigegebene Bestellung (Auftrag `order.provision`). Je Position: Ressource beim Provider anlegen
* (falls das Produkt an eine Verbindung gebunden ist), Vertrag aktivieren. Bereits aktive Verträge werden übersprungen,
* sodass ein Neustart des Auftrags keine Doppelanlage erzeugt. Anlage wird nur bei sicher nicht gesendeter Anfrage wiederholt.
*/
export async function provisionOrder(orderId: string, meta: JobMeta): Promise<string> {
const order = await one('SELECT * FROM orders WHERE id = ?', [orderId]);
if (!order) throw new JobFailure('Bestellung nicht gefunden', false, 'failed');
if (order.status === 'completed') return 'bereits abgeschlossen';
if (order.status !== 'provisioning') throw new JobFailure(`Bestellung ist im Status ${order.status}`, false, 'failed');
const items = await query(
`SELECT oi.id AS item_id, oi.snapshot_json, c.id AS contract_id, c.number AS contract_number, c.status AS cstatus, c.renewal AS c_renewal, c.renewal_term_months AS c_renewal_months, pv.name AS product_name, pv.term_months, pv.provisioning_json,
p.connector_instance_id, p.customer_actions, o.name AS org_name, o.customer_number, ord.number AS order_number
FROM order_items oi JOIN contracts c ON c.order_item_id = oi.id JOIN product_versions pv ON pv.id = oi.product_version_id
JOIN products p ON p.id = pv.product_id JOIN orders ord ON ord.id = oi.order_id JOIN organizations o ON o.id = ord.org_id
WHERE oi.order_id = ? ORDER BY oi.id`, [orderId]);
// Kunde/Vorgang für den Anbieter (Inhaber der Organisation als Ansprechpartner)
const owner = await one("SELECT u.name, u.email FROM memberships m JOIN users u ON u.id = m.user_id WHERE m.org_id = ? AND m.role = 'owner' ORDER BY m.created_at LIMIT 1", [order.org_id]);
const fail = async (note: string, ambiguous: boolean, err: unknown) => {
await run("UPDATE orders SET status = ?, failure_note = ?, failure_ambiguous = ? WHERE id = ?", [transition(ORDER_MACHINE, 'provisioning', 'fail'), note.slice(0, 500), ambiguous ? 1 : 0, orderId]);
await audit({ actorType: 'system', orgId: order.org_id, action: 'order.provision', resourceType: 'order', resourceId: orderId, result: 'failure', errorClass: err instanceof ConnectorError ? err.code : 'unexpected', correlationId: meta.correlationId, after: { note, ambiguous } });
};
for (const it of items) {
if (it.cstatus === 'active') continue;
let resourceId: string | null = null; let ref = ''; let providerValidUntil: string | null = null;
try {
if (it.connector_instance_id) {
const { inst, connector, ctx } = await loadInstance(it.connector_instance_id, meta.correlationId);
if (!inst.enabled) throw new ConnectorError('BAD_CONFIG', 'Verbindung ist deaktiviert');
const caps = await connector.capabilities(ctx);
if (!caps.includes('lifecycle.create') || !connector.provision) throw new ConnectorError('UNSUPPORTED', 'Anlegen wird nicht unterstützt');
const params = typeof it.provisioning_json === 'string' ? JSON.parse(it.provisioning_json) : it.provisioning_json;
const out = await connector.provision(ctx, { params, label: `${it.org_name} · ${it.product_name}`, idempotencyKey: `${meta.id}:${it.item_id}`,
context: { source: 'kundencenter', customerNumber: it.customer_number, customerName: it.org_name, contactName: owner?.name ?? null, customerEmail: owner?.email ?? null, orderNumber: it.order_number, contractNumber: it.contract_number } });
ref = out.resource.externalRef; providerValidUntil = out.resource.validUntil ?? null;
try {
resourceId = await upsert(it.connector_instance_id, out.resource);
await run('UPDATE resources SET org_id = ?, customer_actions = ? WHERE id = ?', [order.org_id, JSON.stringify(typeof it.customer_actions === 'string' ? JSON.parse(it.customer_actions) : it.customer_actions), resourceId]);
} catch (dbErr) {
// Provider hat angelegt, lokale Zuordnung schlug fehl: nie automatisch wiederholen (sonst Doppelanlage)
await fail(`Beim Anbieter wurde ${ref} angelegt, die lokale Zuordnung schlug fehl. Bitte prüfen, nicht blind wiederholen.`, true, dbErr);
throw new JobFailure('Zuordnung nach Anlage fehlgeschlagen', false, 'needs_review');
}
}
// Erste Laufzeitperiode: Mindestlaufzeit, sonst (bei automatischer Verlängerung) die Verlängerungsperiode, damit auch rollierende Monatsverträge einen Verlängerungstakt haben
const now = new Date(); const months = Number(it.term_months) > 0 ? Number(it.term_months) : it.c_renewal === 'auto' ? Number(it.c_renewal_months) : 0;
const c = transition(CONTRACT_MACHINE, 'pending', 'activate');
await run('UPDATE contracts SET status = ?, started_at = ?, term_end = ?, resource_id = ? WHERE id = ? AND status = \'pending\'', [c, now, months > 0 ? addMonths(now, months) : null, resourceId, it.contract_id]);
// Laufzeit des Providerobjekts an das Vertragsende angleichen (z. B. Lizenz 365 Tage vs. 12 Kalendermonate)
if (resourceId && months > 0 && providerValidUntil) {
const target = addMonths(now, months);
if (Math.abs(new Date(providerValidUntil).getTime() - target.getTime()) > 60_000) await enqueue('connector.execute', { resourceId, action: 'extend', params: { until: target.toISOString() }, actorUserId: null }, { idempotencyKey: `align:${it.contract_id}`, correlationId: meta.correlationId });
}
await audit({ actorType: 'system', orgId: order.org_id, action: 'contract.activate', resourceType: 'contract', resourceId: it.contract_id, connector: it.connector_instance_id ? 'provisioned' : undefined, correlationId: meta.correlationId, after: { resourceId, externalRef: ref || undefined } });
} catch (e) {
if (e instanceof JobFailure) throw e;
if (e instanceof ConnectorError) {
const last = meta.attempt >= meta.maxAttempts;
if (e.retryable && e.notSent && !last) throw new JobFailure(e.userMessage, true, 'needs_review');
await fail(`${e.userMessage}${e.ambiguous ? ' Es ist unklar, ob beim Anbieter bereits etwas angelegt wurde. Bitte dort prüfen, bevor die Bestellung erneut gestartet wird.' : ''}`, e.ambiguous, e);
throw new JobFailure(e.userMessage, false, 'needs_review');
}
await fail('Unerwarteter Fehler bei der Bereitstellung', true, e);
throw new JobFailure('Unerwarteter Fehler', false, 'needs_review');
}
}
await run('UPDATE orders SET status = ?, failure_note = NULL, failure_ambiguous = 0 WHERE id = ?', [transition(ORDER_MACHINE, 'provisioning', 'complete'), orderId]);
await audit({ actorType: 'system', orgId: order.org_id, action: 'order.complete', resourceType: 'order', resourceId: orderId, correlationId: meta.correlationId });
return `${items.length} Position(en) bereitgestellt`;
}
/** Beendet Verträge zum Kündigungs-/Laufzeitende und verlängert automatisch. Idempotent; regelmäßig vom Worker aufgerufen. */
export async function processContractLifecycle(now: Date, enqueueJob: (type: string, payload: unknown, o: { idempotencyKey: string }) => Promise<void>): Promise<{ ended: number; renewed: number }> {
let ended = 0; let renewed = 0;
// 1) Gekündigte Verträge, deren Kündigung wirksam wird
const due = await query("SELECT id, org_id, resource_id, status FROM contracts WHERE status IN ('active','suspended') AND cancel_effective_at IS NOT NULL AND cancel_effective_at <= ?", [now]);
for (const c of due) {
const to = transition(CONTRACT_MACHINE, c.status, 'cancel');
const r = await run("UPDATE contracts SET status = ?, cancelled_at = ? WHERE id = ? AND status IN ('active','suspended')", [to, now, c.id]);
if (!r.affectedRows) continue;
ended++;
await audit({ actorType: 'system', orgId: c.org_id, action: 'contract.cancel.effective', resourceType: 'contract', resourceId: c.id });
// Ressource sperren (nicht löschen): Löschung ist destruktiv und bleibt eine bewusste manuelle Entscheidung
if (c.resource_id) await enqueueJob('connector.execute', { resourceId: c.resource_id, action: 'suspend', actorUserId: null }, { idempotencyKey: `contract-end:${c.id}` });
}
// 2) Laufzeitende ohne Kündigung: automatisch verlängern oder auslaufen lassen
const ended2 = await query("SELECT c.id, c.org_id, c.term_end, c.renewal, c.renewal_term_months, c.status, c.resource_id, r.valid_until AS res_until FROM contracts c LEFT JOIN resources r ON r.id = c.resource_id WHERE c.status IN ('active','suspended') AND c.term_end IS NOT NULL AND c.term_end <= ? AND c.cancel_effective_at IS NULL", [now]);
for (const c of ended2) {
if (c.renewal === 'auto' && Number(c.renewal_term_months) > 0) {
let end = new Date(c.term_end as Date); let guard = 0;
while (end <= now && guard++ < 1200) end = addMonths(end, Number(c.renewal_term_months));
await run('UPDATE contracts SET term_end = ? WHERE id = ? AND term_end = ?', [end, c.id, c.term_end]);
renewed++;
// Befristetes Providerobjekt (z. B. Monats-/Jahreslizenz) mitverlängern, sonst läuft es trotz laufendem Vertrag ab.
// HINWEIS: Bis zur Rechnungs-/Zahlungsanbindung erfolgt das ohne Zahlungsprüfung (siehe Plane).
if (c.resource_id && c.res_until) await enqueueJob('connector.execute', { resourceId: c.resource_id, action: 'extend', params: { until: end.toISOString() }, actorUserId: null }, { idempotencyKey: `renew:${c.id}:${end.toISOString()}` });
await audit({ actorType: 'system', orgId: c.org_id, action: 'contract.renew', resourceType: 'contract', resourceId: c.id, after: { termEnd: end.toISOString() } });
} else {
const to = transition(CONTRACT_MACHINE, c.status, 'expire');
const r = await run("UPDATE contracts SET status = ? WHERE id = ? AND status IN ('active','suspended')", [to, c.id]);
if (!r.affectedRows) continue;
ended++;
await audit({ actorType: 'system', orgId: c.org_id, action: 'contract.expire', resourceType: 'contract', resourceId: c.id });
if (c.resource_id) await enqueueJob('connector.execute', { resourceId: c.resource_id, action: 'suspend', actorUserId: null }, { idempotencyKey: `contract-end:${c.id}` });
}
}
return { ended, renewed };
}
export interface ChildPayload { resourceId: string; kind: ChildKind; op: ChildOp; id?: string; data?: Record<string, unknown>; secretEnc?: string; actorUserId: string | null; destructive?: boolean }
/**
* Führt eine Änderung an einem Unterobjekt aus (Auftrag `connector.child`). Geheimnisse (z. B. Passwörter) liegen nur verschlüsselt im Auftrag und
* werden nach dem Ende (Erfolg oder endgültiger Fehler) aus der Datenbank entfernt. Wiederholt wird nur, wenn die Anfrage sicher nicht ankam.
*/
export async function executeChild(p: ChildPayload, jobId: string, correlationId: string, isFinal: () => boolean): Promise<string> {
const wipe = () => run("UPDATE jobs SET payload = JSON_REMOVE(payload, '$.secretEnc') WHERE id = ?", [jobId]).catch(() => undefined);
const r = await one('SELECT * FROM resources WHERE id = ?', [p.resourceId]);
if (!r) { await wipe(); throw new JobFailure('Ressource nicht gefunden', false, 'failed'); }
const { inst, connector, ctx } = await loadInstance(r.instance_id, correlationId);
const caps = await connector.capabilities(ctx);
if (!connector.children || !caps.includes('children.write')) { await wipe(); throw new JobFailure('Der Anbieter unterstützt diese Änderung nicht', false, 'failed'); }
const secrets = p.secretEnc ? (JSON.parse(decrypt(p.secretEnc)) as Record<string, string>) : undefined;
const name = typeof p.data?.domain === 'string' ? p.data.domain : typeof p.data?.local === 'string' ? `${p.data.local}@${String(p.data.domain ?? '')}` : typeof p.data?.name === 'string' ? p.data.name : typeof p.data?.username === 'string' ? p.data.username : (p.id ?? '');
const base = { actorType: (p.actorUserId ? 'user' : 'system') as 'user' | 'system', actorId: p.actorUserId, orgId: r.org_id, resourceType: 'resource', resourceId: p.resourceId, connector: inst.connector_key, correlationId };
try {
const out = await connector.children.act(ctx, { parentRef: r.external_ref, kind: p.kind, op: p.op, id: p.id, data: p.data ?? {}, secrets, idempotencyKey: jobId });
await wipe();
await audit({ ...base, action: `resource.child.${p.kind}.${p.op}`, after: { name: out.child?.name ?? name, id: out.child?.id ?? p.id } });
return `${p.kind}: ${p.op === 'create' ? 'angelegt' : p.op === 'update' ? 'geändert' : 'gelöscht'}${out.child?.name ? ` (${out.child.name})` : ''}`;
} catch (e) {
if (e instanceof ConnectorError) {
await audit({ ...base, action: `resource.child.${p.kind}.${p.op}`, result: 'failure', errorClass: e.code, after: { name } }).catch(() => undefined);
const retry = e.retryable && e.notSent && !isFinal();
if (!retry) await wipe();
// Bei abgelehnten Eingaben ist die Meldung des Anbieters für Nutzer hilfreich (z. B. "Name schon vergeben").
const msg = e.code === 'INVALID_INPUT' || e.code === 'CONFLICT' || e.code === 'NOT_FOUND' ? e.message : e.userMessage;
throw new JobFailure(msg, retry, e.ambiguous ? 'needs_review' : 'failed');
}
await wipe(); throw e;
}
}