197 lines
9.7 KiB
TypeScript
197 lines
9.7 KiB
TypeScript
import type { FastifyInstance, FastifyReply, FastifyRequest } from "fastify";
|
|
import { z } from "zod";
|
|
import { isPrincipalContext, requirePermission } from "../auth/authorization.js";
|
|
import { databaseConfigurationSchema as configSchema } from "../catalog/configuration-schema.js";
|
|
import { CatalogService, type CatalogSecretName } from "../catalog/service.js";
|
|
import { WorkspaceRegistryError } from "../workspaces/git-repository.js";
|
|
import {
|
|
CatalogConflictError,
|
|
CatalogOperationInProgressError,
|
|
CatalogUnavailableError,
|
|
type CatalogRepository,
|
|
type DatabaseConfigurationInput,
|
|
} from "../catalog/types.js";
|
|
import type { CatalogOperationCoordinator } from "../catalog/operation-coordinator.js";
|
|
|
|
const idSchema = z.uuid();
|
|
const updateSchema = configSchema.extend({ version: z.number().int().positive() });
|
|
const secretNames = [
|
|
"password",
|
|
"apiKey",
|
|
"sshPrivateKey",
|
|
"sshPrivateKeyPassphrase",
|
|
"sshKnownHosts",
|
|
"tlsCa",
|
|
] as const;
|
|
const secretsSchema = z.object({
|
|
version: z.number().int().positive(),
|
|
values: z.partialRecord(z.enum(secretNames), z.string().min(1).max(65_536)).refine((values) => Object.keys(values).length > 0),
|
|
}).strict();
|
|
const versionQuery = z.object({ version: z.coerce.number().int().positive() });
|
|
const metricsQuery = z.object({ databaseId: z.uuid().optional() }).strict();
|
|
|
|
function safeError(reply: FastifyReply, error: unknown) {
|
|
if (error instanceof CatalogUnavailableError) {
|
|
return reply.code(503).send({ code: "catalog_unavailable", message: "Database catalog is unavailable." });
|
|
}
|
|
if (error instanceof CatalogConflictError) {
|
|
return reply.code(409).send({ code: "database_conflict", message: "This workspace already has a database configuration." });
|
|
}
|
|
if (error instanceof CatalogOperationInProgressError) {
|
|
return reply.code(409).send({ code: "database_operation_in_progress", message: "A database operation is already in progress." });
|
|
}
|
|
if (error instanceof z.ZodError) {
|
|
return reply.code(400).send({ code: "database_invalid", message: "Database configuration is invalid." });
|
|
}
|
|
if (error instanceof WorkspaceRegistryError) {
|
|
return reply.code(400).send({ code: "database_invalid", message: "Database configuration is invalid." });
|
|
}
|
|
return reply.code(500).send({ code: "database_operation_failed", message: "Database operation failed." });
|
|
}
|
|
|
|
function manage(request: FastifyRequest, reply: FastifyReply) {
|
|
return isPrincipalContext(requirePermission(request, reply, "database.manage"));
|
|
}
|
|
|
|
export function catalogDatabaseRoutes(
|
|
app: FastifyInstance,
|
|
deps: { repository: CatalogRepository; service: CatalogService; operations?: CatalogOperationCoordinator },
|
|
): void {
|
|
const mutate = async <T>(databaseId: string, operation: () => Promise<T>): Promise<T> => (
|
|
deps.operations ? await deps.operations.run(databaseId, operation) : await operation()
|
|
);
|
|
const activeSyncRun = async (databaseId: string) => {
|
|
const run = (await deps.repository.listSyncRuns(databaseId, 5)).find((candidate) =>
|
|
["queued", "running", "awaiting_confirmation", "applying"].includes(candidate.state));
|
|
if (!run) return undefined;
|
|
const {
|
|
observedSnapshot: _snapshot,
|
|
leaseOwner: _leaseOwner,
|
|
leaseExpiresAt: _leaseExpiresAt,
|
|
...summary
|
|
} = run;
|
|
return summary;
|
|
};
|
|
app.get("/catalog/status", async (request, reply) => {
|
|
if (!manage(request, reply)) return reply;
|
|
return { available: await deps.repository.available() };
|
|
});
|
|
|
|
app.get("/catalog/metrics", async (request, reply) => {
|
|
if (!manage(request, reply)) return reply;
|
|
try {
|
|
const { databaseId } = metricsQuery.parse(request.query);
|
|
const metrics = await deps.repository.getCatalogMetrics(databaseId);
|
|
if (!metrics) {
|
|
return reply.code(404).send({
|
|
code: "database_not_found",
|
|
message: "Database configuration was not found.",
|
|
});
|
|
}
|
|
return metrics;
|
|
} catch (error) { return safeError(reply, error); }
|
|
});
|
|
|
|
app.get("/catalog/databases", async (request, reply) => {
|
|
if (!manage(request, reply)) return reply;
|
|
try {
|
|
const rows = await deps.service.list();
|
|
return await Promise.all(rows.map(async (row) => (
|
|
row.id ? { ...row, activeSyncRun: await activeSyncRun(row.id) } : row
|
|
)));
|
|
} catch (error) { return safeError(reply, error); }
|
|
});
|
|
|
|
app.get("/catalog/databases/:id", async (request, reply) => {
|
|
if (!manage(request, reply)) return reply;
|
|
try {
|
|
const id = idSchema.parse((request.params as { id?: unknown }).id);
|
|
const database = await deps.repository.get(id);
|
|
if (!database) return reply.code(404).send({ code: "database_not_found", message: "Database configuration was not found." });
|
|
return {
|
|
...database,
|
|
configured: true,
|
|
secrets: deps.service.configuredSecrets(database.workspaceId),
|
|
activeSyncRun: await activeSyncRun(database.id),
|
|
};
|
|
} catch (error) { return safeError(reply, error); }
|
|
});
|
|
|
|
app.post("/catalog/databases", async (request, reply) => {
|
|
if (!manage(request, reply)) return reply;
|
|
try {
|
|
const input = configSchema.parse(request.body) as DatabaseConfigurationInput;
|
|
const created = await deps.repository.create(await deps.service.normalizeInput(input));
|
|
return reply.code(201).send({ ...created, configured: true, secrets: deps.service.configuredSecrets(created.workspaceId) });
|
|
} catch (error) { return safeError(reply, error); }
|
|
});
|
|
|
|
app.patch("/catalog/databases/:id", async (request, reply) => {
|
|
if (!manage(request, reply)) return reply;
|
|
try {
|
|
const id = idSchema.parse((request.params as { id?: unknown }).id);
|
|
const { version, ...input } = updateSchema.parse(request.body);
|
|
const current = await deps.repository.get(id);
|
|
if (!current) return reply.code(404).send({ code: "database_not_found", message: "Database configuration was not found." });
|
|
if (current.workspaceId !== input.workspaceId) {
|
|
return reply.code(400).send({ code: "database_invalid", message: "Database configuration is invalid." });
|
|
}
|
|
const updated = await mutate(id, async () => await deps.repository.update(
|
|
id, version, await deps.service.normalizeInput(input as DatabaseConfigurationInput),
|
|
));
|
|
if (!updated) return reply.code(409).send({ code: "database_stale", message: "Database configuration changed. Reload and try again." });
|
|
return { ...updated, configured: true, secrets: deps.service.configuredSecrets(updated.workspaceId) };
|
|
} catch (error) { return safeError(reply, error); }
|
|
});
|
|
|
|
app.put("/catalog/databases/:id/secrets", async (request, reply) => {
|
|
if (!manage(request, reply)) return reply;
|
|
if (!isPrincipalContext(requirePermission(request, reply, "workspace.secrets.manage"))) return reply;
|
|
try {
|
|
const id = idSchema.parse((request.params as { id?: unknown }).id);
|
|
const { version, values } = secretsSchema.parse(request.body);
|
|
const database = await deps.repository.get(id);
|
|
if (!database) return reply.code(404).send({ code: "database_not_found", message: "Database configuration was not found." });
|
|
if (database.version !== version) return reply.code(409).send({ code: "database_stale", message: "Database configuration changed. Reload and try again." });
|
|
const updated = await mutate(id, async () => {
|
|
const touched = await deps.repository.touch(id, version);
|
|
if (touched) deps.service.replaceSecrets(database.workspaceId, values as Partial<Record<CatalogSecretName, string>>);
|
|
return touched;
|
|
});
|
|
if (!updated) return reply.code(409).send({ code: "database_stale", message: "Database configuration changed. Reload and try again." });
|
|
return { ...updated, configured: true, secrets: deps.service.configuredSecrets(updated.workspaceId) };
|
|
} catch (error) { return safeError(reply, error); }
|
|
});
|
|
|
|
app.post("/catalog/databases/:id/test", async (request, reply) => {
|
|
if (!manage(request, reply)) return reply;
|
|
try {
|
|
const id = idSchema.parse((request.params as { id?: unknown }).id);
|
|
const { version } = z.object({ version: z.number().int().positive() }).strict().parse(request.body);
|
|
const database = await deps.repository.get(id);
|
|
if (!database) return reply.code(404).send({ code: "database_not_found", message: "Database configuration was not found." });
|
|
if (database.version !== version) return reply.code(409).send({ code: "database_stale", message: "Database configuration changed. Reload and try again." });
|
|
const tested = await deps.service.test(database);
|
|
if (!tested) return reply.code(409).send({ code: "database_stale", message: "Database configuration changed. Reload and try again." });
|
|
return { ...tested, configured: true, secrets: deps.service.configuredSecrets(tested.workspaceId) };
|
|
} catch (error) { return safeError(reply, error); }
|
|
});
|
|
|
|
app.delete("/catalog/databases/:id", async (request, reply) => {
|
|
if (!manage(request, reply)) return reply;
|
|
try {
|
|
const id = idSchema.parse((request.params as { id?: unknown }).id);
|
|
const { version } = versionQuery.parse(request.query);
|
|
const database = await deps.repository.get(id);
|
|
if (!database) return reply.code(404).send({ code: "database_not_found", message: "Database configuration was not found." });
|
|
const deleted = await mutate(id, async () => {
|
|
const removed = await deps.repository.delete(id, version);
|
|
if (removed) deps.service.forgetSecrets(database.workspaceId);
|
|
return removed;
|
|
});
|
|
if (!deleted) return reply.code(409).send({ code: "database_stale", message: "Database configuration changed. Reload and try again." });
|
|
return reply.code(204).send();
|
|
} catch (error) { return safeError(reply, error); }
|
|
});
|
|
}
|