Files
kuber/lib/database.ts
T

409 lines
10 KiB
TypeScript

import type { KubernetesObject, V1Secret } from "@kubernetes/client-node";
import { randomUUID } from "node:crypto";
import type { ComposeSpecification, Service } from "../schema/docker.d";
import { LABELS } from "../const";
import { deleteResource, applyResource } from "./apply";
import { objectApi } from "./k8s";
export const DATABASE_NAMESPACE = "database";
export const DATABASE_CLUSTER = "postgres";
export const DATABASE_HOST = `c.${DATABASE_NAMESPACE}.svc.cluster.local`;
export const DATABASE_PORT = 5432;
export const DATABASE_PROJECT_LABEL = "kuber.dev/project";
export const DATABASE_SERVICE_LABEL = "kuber.dev/service";
type ManagedRole = {
bypassrls: boolean;
connectionLimit: number;
createdb: boolean;
createrole: boolean;
ensure: "present";
inherit: boolean;
login: boolean;
name: string;
passwordSecret: {
name: string;
};
replication: boolean;
superuser: boolean;
};
export type PostgresClaim = {
service: string;
username: string;
database: string;
secretName: string;
};
type RoleCredentials = {
username: string;
password: string;
};
function toSecretName(username: string) {
return `${DATABASE_CLUSTER}-${username}`;
}
function decodeSecretValue(value: string | undefined): string | undefined {
return value ? Buffer.from(value, "base64").toString("utf8") : undefined;
}
function encodeConnectionComponent(value: string): string {
return encodeURIComponent(value);
}
export function buildDatabaseUrl(
claim: PostgresClaim,
credentials: RoleCredentials,
): string {
return `postgresql://${encodeConnectionComponent(credentials.username)}:${encodeConnectionComponent(credentials.password)}@${DATABASE_HOST}:${DATABASE_PORT}/${encodeConnectionComponent(claim.database)}`;
}
function parsePostgresVolumeString(
entry: string,
): Omit<PostgresClaim, "service" | "secretName"> | undefined {
if (!entry.startsWith("postgresql:")) return;
const parts = entry.split(":");
if (parts.length !== 2) {
throw new Error(
`Invalid postgres volume ${entry}. Use postgresql:<name> or postgresql:<user>/<database>.`,
);
}
const target = parts[1]?.trim();
if (!target) {
throw new Error(
`Invalid postgres volume ${entry}. Use postgresql:<name> or postgresql:<user>/<database>.`,
);
}
const segments = target.split("/");
if (
segments.length > 2 ||
segments.some((segment) => segment.trim() === "")
) {
throw new Error(
`Invalid postgres volume ${entry}. Use postgresql:<name> or postgresql:<user>/<database>.`,
);
}
const username = segments[0]!;
const database = segments[1] ?? username;
return { username, database };
}
export function isPostgresVolumeEntry(
entry: NonNullable<Service["volumes"]>[number],
): boolean {
return (
typeof entry === "string" && parsePostgresVolumeString(entry) !== undefined
);
}
export function getServicePostgresClaim(
serviceName: string,
service: Service,
): PostgresClaim | undefined {
const claims =
service.volumes?.flatMap((entry) => {
if (typeof entry !== "string") return [];
const claim = parsePostgresVolumeString(entry);
return claim
? [
{
...claim,
service: serviceName,
secretName: toSecretName(claim.username),
} satisfies PostgresClaim,
]
: [];
}) ?? [];
if (claims.length > 1) {
throw new Error(
`Service ${serviceName} declares multiple postgres volumes. Only zero or one postgresql:<...> entry is allowed per service.`,
);
}
return claims[0];
}
export function getComposePostgresClaims(
compose: ComposeSpecification,
): PostgresClaim[] {
const claims = Object.entries(compose.services ?? {}).flatMap(
([serviceName, service]) => {
const claim = getServicePostgresClaim(serviceName, service);
return claim ? [claim] : [];
},
);
const ownersByDatabase = new Map<string, string>();
for (const claim of claims) {
const owner = ownersByDatabase.get(claim.database);
if (owner && owner !== claim.username) {
throw new Error(
`Database ${claim.database} is claimed by both ${owner} and ${claim.username}. A database can only have one owner.`,
);
}
ownersByDatabase.set(claim.database, claim.username);
}
return claims;
}
async function readObject<T>(
resource: KubernetesObject,
): Promise<T | undefined> {
try {
return (await objectApi.read(resource as never)) as T;
} catch (error) {
if (
error &&
typeof error === "object" &&
"code" in error &&
error.code === 404
) {
return;
}
throw error;
}
}
async function ensureRoleSecret(
claim: PostgresClaim,
): Promise<RoleCredentials> {
const existing = await readObject<V1Secret>({
apiVersion: "v1",
kind: "Secret",
metadata: {
name: claim.secretName,
namespace: DATABASE_NAMESPACE,
},
});
const username = claim.username;
const password = decodeSecretValue(existing?.data?.password) ?? randomUUID();
await applyResource({
apiVersion: "v1",
kind: "Secret",
metadata: {
name: claim.secretName,
namespace: DATABASE_NAMESPACE,
},
type: existing?.type ?? "Opaque",
stringData: {
username,
password,
},
} satisfies V1Secret);
return { username, password };
}
function toManagedRole(claim: PostgresClaim): ManagedRole {
return {
bypassrls: false,
connectionLimit: -1,
createdb: false,
createrole: false,
ensure: "present",
inherit: true,
login: true,
name: claim.username,
passwordSecret: {
name: claim.secretName,
},
replication: false,
superuser: false,
};
}
async function reconcileManagedRoles(claims: PostgresClaim[]): Promise<void> {
if (claims.length === 0) return;
const cluster = await readObject<
KubernetesObject & { spec?: { managed?: { roles?: ManagedRole[] } } }
>({
apiVersion: "postgresql.cnpg.io/v1",
kind: "Cluster",
metadata: {
name: DATABASE_CLUSTER,
namespace: DATABASE_NAMESPACE,
},
});
if (!cluster) {
throw new Error(
`CNPG cluster ${DATABASE_CLUSTER} was not found in namespace ${DATABASE_NAMESPACE}.`,
);
}
const roles = new Map(
(cluster.spec?.managed?.roles ?? []).map((role) => [role.name, role]),
);
for (const claim of claims) {
roles.set(claim.username, toManagedRole(claim));
}
try {
await applyResource({
apiVersion: "postgresql.cnpg.io/v1",
kind: "Cluster",
metadata: {
name: DATABASE_CLUSTER,
namespace: DATABASE_NAMESPACE,
},
spec: {
managed: {
roles: [...roles.values()],
},
},
});
} catch (error) {
throw new Error(
`Failed to reconcile managed roles on ${DATABASE_NAMESPACE}/${DATABASE_CLUSTER}: ${error instanceof Error ? error.message : String(error)}`,
);
}
}
async function reconcileDatabases(
project: string,
claims: PostgresClaim[],
): Promise<void> {
const uniqueDatabases = new Map<string, PostgresClaim>();
for (const claim of claims) {
uniqueDatabases.set(`${claim.database}:${claim.username}`, claim);
}
for (const claim of uniqueDatabases.values()) {
try {
await applyResource({
apiVersion: "postgresql.cnpg.io/v1",
kind: "Database",
metadata: {
name: claim.database,
namespace: DATABASE_NAMESPACE,
labels: {
...LABELS,
[DATABASE_PROJECT_LABEL]: project,
[DATABASE_SERVICE_LABEL]: claim.service,
},
},
spec: {
cluster: {
name: DATABASE_CLUSTER,
},
databaseReclaimPolicy: "retain",
ensure: "present",
name: claim.database,
owner: claim.username,
},
});
} catch (error) {
throw new Error(
`Failed to reconcile database ${claim.database} owned by ${claim.username}: ${error instanceof Error ? error.message : String(error)}`,
);
}
}
}
export async function reconcilePostgresClaims(
project: string,
compose: ComposeSpecification,
): Promise<Record<string, Record<string, string>>> {
const claims = getComposePostgresClaims(compose);
if (claims.length === 0) return {};
const credentialsBySecret = new Map<string, RoleCredentials>();
for (const claim of claims) {
if (credentialsBySecret.has(claim.secretName)) continue;
credentialsBySecret.set(claim.secretName, await ensureRoleSecret(claim));
}
await reconcileManagedRoles(claims);
await reconcileDatabases(project, claims);
return Object.fromEntries(
claims.map((claim) => {
const credentials = credentialsBySecret.get(claim.secretName);
if (!credentials) {
throw new Error(`Missing credentials for ${claim.secretName}`);
}
return [
claim.service,
{
DATABASE_URL: buildDatabaseUrl(claim, credentials),
},
];
}),
);
}
export async function getRoleCredentials(
username: string,
): Promise<RoleCredentials> {
const secret = await readObject<V1Secret>({
apiVersion: "v1",
kind: "Secret",
metadata: {
name: toSecretName(username),
namespace: DATABASE_NAMESPACE,
},
});
const resolvedUsername = decodeSecretValue(secret?.data?.username);
const password = decodeSecretValue(secret?.data?.password);
if (!resolvedUsername || !password) {
throw new Error(
`Managed role secret ${toSecretName(username)} was not found or is missing credentials.`,
);
}
return {
username: resolvedUsername,
password,
};
}
export async function listManagedDatabaseResources(
project: string,
): Promise<KubernetesObject[]> {
const result = await objectApi.list(
"postgresql.cnpg.io/v1",
"Database",
DATABASE_NAMESPACE,
undefined,
undefined,
undefined,
undefined,
`${Object.entries(LABELS)
.map(([key, value]) => `${key}=${value}`)
.join(",")},${DATABASE_PROJECT_LABEL}=${project}`,
);
return result.items.map((item) => ({
...item,
apiVersion: item.apiVersion ?? "postgresql.cnpg.io/v1",
kind: item.kind ?? "Database",
metadata: {
...item.metadata,
namespace: item.metadata?.namespace ?? DATABASE_NAMESPACE,
},
}));
}
export async function deleteManagedDatabases(project: string): Promise<void> {
const resources = await listManagedDatabaseResources(project);
for (const resource of resources) {
await deleteResource(resource);
}
}