import { Router, type Request, type Response } from "express"; import { eq } from "drizzle-orm"; import { z } from "zod"; import { db } from "../db/client.js"; import { integrations, integrationCredentials, integrationTypes } from "../db/schema.js"; import { requireAuth, requireRole } from "../auth/middleware.js"; import { recordAudit } from "../services/audit.js"; import { encryptSecret } from "../crypto.js"; import { INTEGRATION_FIELDS, resolveBaseUrl, splitIntegrationConfig, validateIntegrationConfig, } from "../integrations/fieldSchemas.js"; import { createIntegrationAdapter } from "../integrations/registry.js"; import { loadIntegrationConfig } from "../integrations/loadIntegration.js"; import { createTailscaleAdapter } from "../integrations/tailscale/adapter.js"; import { createGiteaAdapter } from "../integrations/gitea/adapter.js"; import { createDockhandAdapter } from "../integrations/dockhand/adapter.js"; import { createSemaphoreAdapter } from "../integrations/semaphore/adapter.js"; import { createProxmoxAdapter } from "../integrations/proxmox/adapter.js"; import { createSynologyAdapter } from "../integrations/synology/adapter.js"; import { asyncHandler } from "../utils/asyncHandler.js"; export const integrationsRouter = Router(); integrationsRouter.use(requireAuth); const IMPLEMENTED_TYPES = Object.keys(INTEGRATION_FIELDS) as (keyof typeof INTEGRATION_FIELDS)[]; integrationsRouter.get("/fields", (_req, res) => { res.json({ fields: INTEGRATION_FIELDS }); }); integrationsRouter.get("/", asyncHandler(async (_req, res) => { const rows = await db .select({ id: integrations.id, type: integrations.type, name: integrations.name, baseUrl: integrations.baseUrl, enabled: integrations.enabled, createdAt: integrations.createdAt, }) .from(integrations); res.json({ integrations: rows }); })); const configValueSchema = z.union([z.string(), z.boolean()]); const createIntegrationSchema = z.object({ type: z.enum(integrationTypes), name: z.string().min(1).max(200), config: z.record(configValueSchema), }); integrationsRouter.post("/", requireRole("admin"), asyncHandler(async (req, res) => { const parsed = createIntegrationSchema.safeParse(req.body); if (!parsed.success) { return res.status(400).json({ error: "invalid_body", details: parsed.error.flatten() }); } const { type, name, config } = parsed.data; if (!IMPLEMENTED_TYPES.includes(type)) { return res.status(400).json({ error: "not_implemented", message: `Integration type "${type}" isn't available yet.` }); } const missing = validateIntegrationConfig(type, config); if (missing.length > 0) { return res.status(400).json({ error: "missing_fields", fields: missing }); } const { secretFields, nonSecretFields } = splitIntegrationConfig(type, config); let credentialId: number | null = null; if (Object.keys(secretFields).length > 0) { const [cred] = await db .insert(integrationCredentials) .values({ name: `${type}:${name}`, encryptedSecret: encryptSecret(JSON.stringify(secretFields)) }) .returning(); credentialId = cred.id; } const [created] = await db .insert(integrations) .values({ type, name, baseUrl: resolveBaseUrl(type, nonSecretFields), credentialId, config: JSON.stringify(nonSecretFields), enabled: true, }) .returning(); await recordAudit({ actor: req.currentUser!, category: "integration", action: "create", targetType: "integration", targetId: created.id, detail: { type, name }, }); res.status(201).json({ integration: { id: created.id, type: created.type, name: created.name, baseUrl: created.baseUrl, enabled: created.enabled, createdAt: created.createdAt, }, }); })); const updateIntegrationSchema = z.object({ name: z.string().min(1).max(200).optional(), enabled: z.boolean().optional(), config: z.record(configValueSchema).optional(), }); integrationsRouter.patch("/:id", requireRole("admin"), asyncHandler(async (req, res) => { const id = Number(req.params.id); const parsed = updateIntegrationSchema.safeParse(req.body); if (!parsed.success) { return res.status(400).json({ error: "invalid_body", details: parsed.error.flatten() }); } const [existing] = await db.select().from(integrations).where(eq(integrations.id, id)).limit(1); if (!existing) { return res.status(404).json({ error: "not_found" }); } let credentialId = existing.credentialId; let configJson = existing.config; let baseUrl = existing.baseUrl; if (parsed.data.config) { const loaded = await loadIntegrationConfig(id); const { secretFields, nonSecretFields } = splitIntegrationConfig(existing.type, parsed.data.config); const mergedNonSecret = { ...(loaded?.config ?? {}), ...nonSecretFields }; for (const field of INTEGRATION_FIELDS[existing.type] ?? []) { if (field.secret) delete (mergedNonSecret as Record)[field.key]; } configJson = JSON.stringify(mergedNonSecret); baseUrl = resolveBaseUrl(existing.type, mergedNonSecret); if (Object.keys(secretFields).length > 0) { const existingSecrets: Record = {}; const fields = (INTEGRATION_FIELDS[existing.type] ?? []).filter((f) => f.secret); for (const f of fields) { if (loaded?.config[f.key] !== undefined) existingSecrets[f.key] = loaded.config[f.key]!; } const mergedSecrets = { ...existingSecrets, ...secretFields }; const encrypted = encryptSecret(JSON.stringify(mergedSecrets)); if (existing.credentialId) { await db.update(integrationCredentials).set({ encryptedSecret: encrypted }).where(eq(integrationCredentials.id, existing.credentialId)); } else { const [cred] = await db .insert(integrationCredentials) .values({ name: `${existing.type}:${existing.name}`, encryptedSecret: encrypted }) .returning(); credentialId = cred.id; } } } const [updated] = await db .update(integrations) .set({ name: parsed.data.name ?? existing.name, enabled: parsed.data.enabled ?? existing.enabled, credentialId, config: configJson, baseUrl, }) .where(eq(integrations.id, id)) .returning(); await recordAudit({ actor: req.currentUser!, category: "integration", action: "update", targetType: "integration", targetId: id, detail: { name: updated.name }, }); res.json({ integration: { id: updated.id, type: updated.type, name: updated.name, baseUrl: updated.baseUrl, enabled: updated.enabled, createdAt: updated.createdAt, }, }); })); integrationsRouter.get("/:id/config", requireRole("admin"), asyncHandler(async (req, res) => { const id = Number(req.params.id); const loaded = await loadIntegrationConfig(id); if (!loaded) { return res.status(404).json({ error: "not_found" }); } // Never return secret fields (API tokens/passwords) to the browser — the edit // form pre-fills only non-secret fields and leaves secret inputs blank. const nonSecretConfig: Record = { ...loaded.config }; for (const field of INTEGRATION_FIELDS[loaded.integration.type] ?? []) { if (field.secret) delete nonSecretConfig[field.key]; } res.json({ integration: { id: loaded.integration.id, type: loaded.integration.type, name: loaded.integration.name, enabled: loaded.integration.enabled, }, config: nonSecretConfig, }); })); const testExistingIntegrationSchema = z.object({ config: z.record(configValueSchema).optional() }); integrationsRouter.post("/:id/test", requireRole("admin"), asyncHandler(async (req, res) => { const id = Number(req.params.id); const parsed = testExistingIntegrationSchema.safeParse(req.body); if (!parsed.success) { return res.status(400).json({ error: "invalid_body", details: parsed.error.flatten() }); } const loaded = await loadIntegrationConfig(id); if (!loaded) { return res.status(404).json({ error: "not_found" }); } // Merge any freshly-typed fields (e.g. a replacement token) over the // already-stored, decrypted config — lets "Test connection" work during an // edit without ever sending the current secret value back to the browser. const mergedConfig = { ...loaded.config, ...(parsed.data.config ?? {}) }; const adapter = createIntegrationAdapter(loaded.integration.type, mergedConfig); const result = await adapter.ping(); res.json(result); })); integrationsRouter.delete("/:id", requireRole("admin"), asyncHandler(async (req, res) => { const id = Number(req.params.id); const [existing] = await db.select().from(integrations).where(eq(integrations.id, id)).limit(1); if (!existing) { return res.status(404).json({ error: "not_found" }); } await db.delete(integrations).where(eq(integrations.id, id)); if (existing.credentialId) { await db.delete(integrationCredentials).where(eq(integrationCredentials.id, existing.credentialId)); } await recordAudit({ actor: req.currentUser!, category: "integration", action: "delete", targetType: "integration", targetId: id, detail: { name: existing.name }, }); res.status(204).end(); })); const testIntegrationSchema = z.object({ type: z.enum(integrationTypes), config: z.record(configValueSchema), }); integrationsRouter.post("/test", requireRole("admin"), asyncHandler(async (req, res) => { const parsed = testIntegrationSchema.safeParse(req.body); if (!parsed.success) { return res.status(400).json({ error: "invalid_body", details: parsed.error.flatten() }); } if (!IMPLEMENTED_TYPES.includes(parsed.data.type)) { return res.status(400).json({ error: "not_implemented" }); } const adapter = createIntegrationAdapter(parsed.data.type, parsed.data.config); const result = await adapter.ping(); res.json(result); })); // ─── Tailscale ─────────────────────────────────────────────────────────────── /** * Loads the integration row and checks its type BEFORE constructing an * adapter — createIntegrationAdapter() throws for any not-yet-implemented * type, and that must never happen for an existing non-tailscale row (e.g. a * Proxmox row someone points at these Tailscale-only routes). */ async function requireTailscaleAdapter(req: Request, res: Response) { const id = Number(req.params.id); const loaded = await loadIntegrationConfig(id); if (!loaded) { res.status(404).json({ error: "not_found" }); return null; } if (loaded.integration.type !== "tailscale") { res.status(400).json({ error: "wrong_type" }); return null; } if (!loaded.integration.enabled) { res.status(400).json({ error: "integration_disabled" }); return null; } return { integration: loaded.integration, adapter: createTailscaleAdapter(loaded.config as any) }; } integrationsRouter.get("/:id/tailscale/devices", asyncHandler(async (req, res) => { const found = await requireTailscaleAdapter(req, res); if (!found) return; try { const devices = await found.adapter.listDevices(); res.json({ devices, summary: { total: devices.length, online: devices.filter((d) => d.online).length, unauthorized: devices.filter((d) => !d.authorized).length, }, }); } catch (err) { res.status(502).json({ error: err instanceof Error ? err.message : String(err) }); } })); const authorizeSchema = z.object({ authorized: z.boolean() }); integrationsRouter.post( "/:id/tailscale/devices/:deviceId/authorize", requireRole("operator"), asyncHandler(async (req, res) => { const found = await requireTailscaleAdapter(req, res); if (!found) return; const parsed = authorizeSchema.safeParse(req.body); if (!parsed.success) { return res.status(400).json({ error: "invalid_body", details: parsed.error.flatten() }); } try { await found.adapter.setAuthorized(req.params.deviceId, parsed.data.authorized); await recordAudit({ actor: req.currentUser!, category: "integration", action: parsed.data.authorized ? "authorize_device" : "deauthorize_device", targetType: "tailscale_device", targetId: req.params.deviceId, detail: { integrationId: found.integration.id }, }); res.status(204).end(); } catch (err) { res.status(502).json({ error: err instanceof Error ? err.message : String(err) }); } }), ); integrationsRouter.delete( "/:id/tailscale/devices/:deviceId", requireRole("operator"), asyncHandler(async (req, res) => { const found = await requireTailscaleAdapter(req, res); if (!found) return; try { await found.adapter.deleteDevice(req.params.deviceId); await recordAudit({ actor: req.currentUser!, category: "integration", action: "remove_device", targetType: "tailscale_device", targetId: req.params.deviceId, detail: { integrationId: found.integration.id }, }); res.status(204).end(); } catch (err) { res.status(502).json({ error: err instanceof Error ? err.message : String(err) }); } }), ); // ─── Gitea ─────────────────────────────────────────────────────────────────── async function requireGiteaAdapter(req: Request, res: Response) { const id = Number(req.params.id); const loaded = await loadIntegrationConfig(id); if (!loaded) { res.status(404).json({ error: "not_found" }); return null; } if (loaded.integration.type !== "gitea") { res.status(400).json({ error: "wrong_type" }); return null; } if (!loaded.integration.enabled) { res.status(400).json({ error: "integration_disabled" }); return null; } return { integration: loaded.integration, adapter: createGiteaAdapter(loaded.config as any) }; } integrationsRouter.get("/:id/gitea/repos", asyncHandler(async (req, res) => { const found = await requireGiteaAdapter(req, res); if (!found) return; try { const repos = await found.adapter.listReposWithStatus(); res.json({ repos }); } catch (err) { res.status(502).json({ error: err instanceof Error ? err.message : String(err) }); } })); integrationsRouter.post( "/:id/gitea/repos/:owner/:repo/runs/:runId/rerun-failed", requireRole("operator"), asyncHandler(async (req, res) => { const found = await requireGiteaAdapter(req, res); if (!found) return; const runId = Number(req.params.runId); if (!Number.isInteger(runId)) { return res.status(400).json({ error: "invalid_run_id" }); } try { await found.adapter.rerunFailedJobs(req.params.owner, req.params.repo, runId); await recordAudit({ actor: req.currentUser!, category: "integration", action: "rerun_failed_jobs", targetType: "gitea_workflow_run", targetId: `${req.params.owner}/${req.params.repo}#${runId}`, detail: { integrationId: found.integration.id }, }); res.status(204).end(); } catch (err) { res.status(502).json({ error: err instanceof Error ? err.message : String(err) }); } }), ); // ─── Dockhand ──────────────────────────────────────────────────────────────── async function requireDockhandAdapter(req: Request, res: Response) { const id = Number(req.params.id); const loaded = await loadIntegrationConfig(id); if (!loaded) { res.status(404).json({ error: "not_found" }); return null; } if (loaded.integration.type !== "dockhand") { res.status(400).json({ error: "wrong_type" }); return null; } if (!loaded.integration.enabled) { res.status(400).json({ error: "integration_disabled" }); return null; } return { integration: loaded.integration, adapter: createDockhandAdapter(loaded.config as any) }; } integrationsRouter.get("/:id/dockhand/containers", asyncHandler(async (req, res) => { const found = await requireDockhandAdapter(req, res); if (!found) return; try { const containers = await found.adapter.listContainers(); res.json({ containers, summary: { total: containers.length, running: containers.filter((c) => c.state === "running").length, }, }); } catch (err) { res.status(502).json({ error: err instanceof Error ? err.message : String(err) }); } })); const dockhandActions = ["start", "stop", "restart"] as const; for (const action of dockhandActions) { integrationsRouter.post( `/:id/dockhand/environments/:envId/containers/:containerId/${action}`, requireRole("operator"), asyncHandler(async (req, res) => { const found = await requireDockhandAdapter(req, res); if (!found) return; const envId = Number(req.params.envId); if (!Number.isInteger(envId)) { return res.status(400).json({ error: "invalid_environment_id" }); } try { await found.adapter[`${action}Container`](envId, req.params.containerId); await recordAudit({ actor: req.currentUser!, category: "integration", action: `${action}_container`, targetType: "dockhand_container", targetId: req.params.containerId, detail: { integrationId: found.integration.id, environmentId: envId }, }); res.status(204).end(); } catch (err) { res.status(502).json({ error: err instanceof Error ? err.message : String(err) }); } }), ); } // ─── Semaphore ─────────────────────────────────────────────────────────────── async function requireSemaphoreAdapter(req: Request, res: Response) { const id = Number(req.params.id); const loaded = await loadIntegrationConfig(id); if (!loaded) { res.status(404).json({ error: "not_found" }); return null; } if (loaded.integration.type !== "semaphore") { res.status(400).json({ error: "wrong_type" }); return null; } if (!loaded.integration.enabled) { res.status(400).json({ error: "integration_disabled" }); return null; } return { integration: loaded.integration, adapter: createSemaphoreAdapter(loaded.config as any) }; } integrationsRouter.get("/:id/semaphore/templates", asyncHandler(async (req, res) => { const found = await requireSemaphoreAdapter(req, res); if (!found) return; try { const templates = await found.adapter.listTemplatesWithStatus(); res.json({ templates, summary: { total: templates.length, failing: templates.filter((t) => t.lastTask?.status === "error").length, }, }); } catch (err) { res.status(502).json({ error: err instanceof Error ? err.message : String(err) }); } })); integrationsRouter.post( "/:id/semaphore/projects/:projectId/templates/:templateId/run", requireRole("operator"), asyncHandler(async (req, res) => { const found = await requireSemaphoreAdapter(req, res); if (!found) return; const projectId = Number(req.params.projectId); const templateId = Number(req.params.templateId); if (!Number.isInteger(projectId) || !Number.isInteger(templateId)) { return res.status(400).json({ error: "invalid_id" }); } try { const task = await found.adapter.runTemplate(projectId, templateId); await recordAudit({ actor: req.currentUser!, category: "integration", action: "run_template", targetType: "semaphore_template", targetId: templateId, detail: { integrationId: found.integration.id, projectId, taskId: task.id }, }); res.status(201).json({ task }); } catch (err) { res.status(502).json({ error: err instanceof Error ? err.message : String(err) }); } }), ); // ─── Proxmox ───────────────────────────────────────────────────────────────── async function requireProxmoxAdapter(req: Request, res: Response) { const id = Number(req.params.id); const loaded = await loadIntegrationConfig(id); if (!loaded) { res.status(404).json({ error: "not_found" }); return null; } if (loaded.integration.type !== "proxmox") { res.status(400).json({ error: "wrong_type" }); return null; } if (!loaded.integration.enabled) { res.status(400).json({ error: "integration_disabled" }); return null; } return { integration: loaded.integration, adapter: createProxmoxAdapter(loaded.config as any) }; } integrationsRouter.get("/:id/proxmox/guests", asyncHandler(async (req, res) => { const found = await requireProxmoxAdapter(req, res); if (!found) return; try { const guests = await found.adapter.listGuests(); res.json({ guests, summary: { total: guests.length, running: guests.filter((g) => g.status === "running").length, }, }); } catch (err) { res.status(502).json({ error: err instanceof Error ? err.message : String(err) }); } })); const proxmoxActions = ["start", "stop", "restart", "shutdown"] as const; for (const action of proxmoxActions) { integrationsRouter.post( `/:id/proxmox/nodes/:node/:type/:vmid/${action}`, requireRole("operator"), asyncHandler(async (req, res) => { const found = await requireProxmoxAdapter(req, res); if (!found) return; const vmid = Number(req.params.vmid); if (!Number.isInteger(vmid)) { return res.status(400).json({ error: "invalid_vmid" }); } const type = req.params.type; if (type !== "qemu" && type !== "lxc") { return res.status(400).json({ error: "invalid_type" }); } try { await found.adapter[`${action}Guest`](req.params.node, type, vmid); await recordAudit({ actor: req.currentUser!, category: "integration", action: `${action}_guest`, targetType: "proxmox_guest", targetId: vmid, detail: { integrationId: found.integration.id, node: req.params.node, guestType: type }, }); res.status(204).end(); } catch (err) { res.status(502).json({ error: err instanceof Error ? err.message : String(err) }); } }), ); } // ─── Synology ──────────────────────────────────────────────────────────────── // Read-only per the delivery plan — no start/stop/etc actions. async function requireSynologyAdapter(req: Request, res: Response) { const id = Number(req.params.id); const loaded = await loadIntegrationConfig(id); if (!loaded) { res.status(404).json({ error: "not_found" }); return null; } if (loaded.integration.type !== "synology") { res.status(400).json({ error: "wrong_type" }); return null; } if (!loaded.integration.enabled) { res.status(400).json({ error: "integration_disabled" }); return null; } return { integration: loaded.integration, adapter: createSynologyAdapter(loaded.config as any) }; } integrationsRouter.get("/:id/synology/storage", asyncHandler(async (req, res) => { const found = await requireSynologyAdapter(req, res); if (!found) return; try { const info = await found.adapter.getStorageInfo(); res.json({ ...info, summary: { volumeCount: info.volumes.length, volumesNotNormal: info.volumes.filter((v) => v.status !== "normal").length, diskCount: info.disks.length, disksNotNormal: info.disks.filter((d) => d.status !== "normal").length, }, }); } catch (err) { res.status(502).json({ error: err instanceof Error ? err.message : String(err) }); } }));