feat: add database reconciler

This commit is contained in:
2026-08-16 09:35:34 +07:00 Unverified
parent 3f70382fd3
commit a244f49155
6 changed files with 659 additions and 105 deletions
+83
View File
@@ -0,0 +1,83 @@
import { defineCommand } from "citty";
import { ctx } from "../lib/context";
import {
DATABASE_HOST,
DATABASE_PORT,
buildDatabaseUrl,
getComposePostgresClaims,
getRoleCredentials,
} from "../lib/database";
import { toTable } from "../lib/format";
const list = defineCommand({
meta: {
name: "ls",
description: "List managed postgres claims",
},
async run() {
const claims = getComposePostgresClaims(await ctx().compose());
if (claims.length === 0) {
console.log("No managed postgres volumes declared");
return;
}
console.log(
toTable(
claims.map((claim) => ({
service: claim.service,
username: claim.username,
database: claim.database,
secret: claim.secretName,
})),
),
);
},
});
const creds = defineCommand({
meta: {
name: "creds",
description: "Print postgres credentials for a service",
},
async run({ args }) {
const service = args._[0];
if (!service) throw new Error("Service name is required");
const claim = getComposePostgresClaims(await ctx().compose()).find(
(entry) => entry.service === service,
);
if (!claim) {
throw new Error(
`Service ${service} does not declare a managed postgres volume.`,
);
}
const credentials = await getRoleCredentials(claim.username);
const url = buildDatabaseUrl(claim, credentials);
console.log(
toTable([
{
service: claim.service,
username: credentials.username,
database: claim.database,
host: DATABASE_HOST,
port: DATABASE_PORT,
secret: claim.secretName,
},
]),
);
console.log(`DATABASE_URL=${url}`);
},
});
export const db = defineCommand({
meta: {
name: "db",
description: "Inspect managed postgres databases",
},
subCommands: {
creds,
ls: list,
},
});
+5
View File
@@ -7,6 +7,7 @@ import {
listManagedResources,
sortResources,
} from "../lib/apply";
import { listManagedDatabaseResources } from "../lib/database";
export const down = defineCommand({
meta: {
@@ -30,6 +31,10 @@ export const down = defineCommand({
resource.kind !== "PersistentVolumeClaim"),
);
if (args.full) {
resources.push(...(await listManagedDatabaseResources(project)));
}
if (args.full) {
resources.push({
apiVersion: "v1",
+127 -98
View File
@@ -11,10 +11,15 @@ import {
} from "../lib/apply";
import { buildServices } from "../lib/build";
import { composeToKubernetes, type KubernetesResource } from "../lib/convert";
import {
getComposePostgresClaims,
reconcilePostgresClaims,
} from "../lib/database";
import { restartDeployment, waitForDeploymentRollout } from "../lib/shared";
type UpContext = {
compose?: ComposeSpecification;
serviceEnv?: Record<string, Record<string, string>>;
resources?: KubernetesResource[];
staleResources?: KubernetesResource[];
};
@@ -30,117 +35,141 @@ export async function runUp(build: boolean) {
const { project, compose, cwd } = ctx();
await assertManagedNamespace(project);
await new Listr<UpContext>([
{
title: "Read compose",
task: async (taskCtx, task) => {
taskCtx.compose = await compose();
task.output = `${Object.keys(taskCtx.compose.services ?? {}).length} services`;
await new Listr<UpContext>(
[
{
title: "Read compose",
task: async (taskCtx, task) => {
taskCtx.compose = await compose();
task.output = `${Object.keys(taskCtx.compose.services ?? {}).length} services`;
},
},
},
{
title: "Build images",
rendererOptions: {
outputBar: 10,
persistentOutput: true,
{
title: "Build images",
rendererOptions: {
outputBar: 10,
persistentOutput: true,
},
enabled: () => build,
skip: (taskCtx) =>
Object.values(taskCtx.compose?.services ?? {}).some(
(service) => service.build,
)
? false
: "No buildable services",
task: async (taskCtx, task) => {
const built = await buildServices(project, taskCtx.compose!, cwd, {
progress: (message) => {
task.output = message;
},
stream: task.stdout(),
});
task.output = `Built ${built} image${built === 1 ? "" : "s"}`;
},
},
enabled: () => build,
skip: (taskCtx) =>
Object.values(taskCtx.compose?.services ?? {}).some((service) => service.build)
? false
: "No buildable services",
task: async (taskCtx, task) => {
const built = await buildServices(project, taskCtx.compose!, cwd, {
progress: (message) => {
task.output = message;
},
stream: task.stdout(),
});
task.output = `Built ${built} image${built === 1 ? "" : "s"}`;
{
title: "Reconcile databases",
skip: (taskCtx) =>
getComposePostgresClaims(taskCtx.compose!).length > 0
? false
: "No managed postgres volumes",
task: async (taskCtx, task) => {
taskCtx.serviceEnv = await reconcilePostgresClaims(
project,
taskCtx.compose!,
);
task.output = `${Object.keys(taskCtx.serviceEnv).length} service${Object.keys(taskCtx.serviceEnv).length === 1 ? "" : "s"}`;
},
},
},
{
title: "Render manifests",
task: async (taskCtx, task) => {
taskCtx.resources = await composeToKubernetes(project, taskCtx.compose!, cwd);
task.output = `${taskCtx.resources.length} resources`;
{
title: "Render manifests",
task: async (taskCtx, task) => {
taskCtx.resources = await composeToKubernetes(
project,
taskCtx.compose!,
cwd,
taskCtx.serviceEnv,
);
task.output = `${taskCtx.resources.length} resources`;
},
},
},
{
title: "Plan reconciliation",
task: async (taskCtx, task) => {
taskCtx.staleResources = (await getStaleResources(
project,
taskCtx.resources!,
)) as KubernetesResource[];
task.output = `${taskCtx.staleResources.length} stale resource${taskCtx.staleResources.length === 1 ? "" : "s"}`;
{
title: "Plan reconciliation",
task: async (taskCtx, task) => {
taskCtx.staleResources = (await getStaleResources(
project,
taskCtx.resources!,
)) as KubernetesResource[];
task.output = `${taskCtx.staleResources.length} stale resource${taskCtx.staleResources.length === 1 ? "" : "s"}`;
},
},
},
{
title: "Apply resources",
task: (taskCtx, task) => {
const resources = sortResources(taskCtx.resources!);
{
title: "Apply resources",
task: (taskCtx, task) => {
const resources = sortResources(taskCtx.resources!);
task.output = `${resources.length} resources queued`;
return task.newListr(
resources.map((resource) => ({
title: `${resource.kind} ${resource.metadata?.name}`,
task: () => applyResource(resource),
})),
{ concurrent: false, exitOnError: true },
);
task.output = `${resources.length} resources queued`;
return task.newListr(
resources.map((resource) => ({
title: `${resource.kind} ${resource.metadata?.name}`,
task: () => applyResource(resource),
})),
{ concurrent: false, exitOnError: true },
);
},
},
},
{
title: "Restart deployments",
task: (taskCtx, task) => {
const deployments = getDeploymentNames(taskCtx.resources!);
{
title: "Restart deployments",
task: (taskCtx, task) => {
const deployments = getDeploymentNames(taskCtx.resources!);
task.output = `${deployments.length} deployments queued`;
return task.newListr(
deployments.map((name) => ({
title: `Deployment ${name}`,
task: () => restartDeployment(name),
})),
{ concurrent: false, exitOnError: true },
);
task.output = `${deployments.length} deployments queued`;
return task.newListr(
deployments.map((name) => ({
title: `Deployment ${name}`,
task: () => restartDeployment(name),
})),
{ concurrent: false, exitOnError: true },
);
},
},
},
{
title: "Wait for rollout",
task: (taskCtx, task) => {
const deployments = getDeploymentNames(taskCtx.resources!);
{
title: "Wait for rollout",
task: (taskCtx, task) => {
const deployments = getDeploymentNames(taskCtx.resources!);
task.output = `${deployments.length} deployments queued`;
return task.newListr(
deployments.map((name) => ({
title: `Deployment ${name}`,
task: () => waitForDeploymentRollout(name),
})),
{ concurrent: false, exitOnError: true },
);
task.output = `${deployments.length} deployments queued`;
return task.newListr(
deployments.map((name) => ({
title: `Deployment ${name}`,
task: () => waitForDeploymentRollout(name),
})),
{ concurrent: false, exitOnError: true },
);
},
},
},
{
title: "Delete stale resources",
skip: (taskCtx) =>
taskCtx.staleResources && taskCtx.staleResources.length > 0
? false
: "No stale resources",
task: (taskCtx, task) => {
const resources = sortResources(taskCtx.staleResources!).reverse();
{
title: "Delete stale resources",
skip: (taskCtx) =>
taskCtx.staleResources && taskCtx.staleResources.length > 0
? false
: "No stale resources",
task: (taskCtx, task) => {
const resources = sortResources(taskCtx.staleResources!).reverse();
task.output = `${resources.length} resources queued`;
return task.newListr(
resources.map((resource) => ({
title: `${resource.kind} ${resource.metadata?.name}`,
task: () => deleteResource(resource),
})),
{ concurrent: false, exitOnError: true },
);
task.output = `${resources.length} resources queued`;
return task.newListr(
resources.map((resource) => ({
title: `${resource.kind} ${resource.metadata?.name}`,
task: () => deleteResource(resource),
})),
{ concurrent: false, exitOnError: true },
);
},
},
},
], { rendererOptions: { collapseErrors: false } }).run();
],
{ rendererOptions: { collapseErrors: false } },
).run();
}
export const up = defineCommand({
+19 -3
View File
@@ -1,4 +1,5 @@
import { defineCommand, runMain } from "citty";
import { db } from "./command/db";
import { down } from "./command/down";
import { exec } from "./command/exec";
import { logs } from "./command/logs";
@@ -10,16 +11,31 @@ import { up } from "./command/up";
import { provideContext } from "./lib/context";
import z, { ZodError } from "zod";
function formatUnknownError(error: unknown): string {
if (error instanceof ZodError) return z.prettifyError(error);
if (error instanceof Error) {
const stack = error.stack?.split("\n").slice(0, 3).join("\n");
return stack && stack.trim().length > 0 ? stack : error.message;
}
try {
return JSON.stringify(error, null, 2);
} catch {
return String(error);
}
}
const main = defineCommand({
meta: {
name: "kuber",
version: "1.0.0",
description: "Docker Compose -> K8s translation layer",
},
subCommands: { down, exec, logs, ps, restart, start, stop, up },
subCommands: { db, down, exec, logs, ps, restart, start, stop, up },
});
provideContext(() => runMain(main)).catch((e) => {
if (e instanceof ZodError) console.error(z.prettifyError(e));
else console.error(e);
console.error(formatUnknownError(e));
process.exitCode = 1;
});
+29 -4
View File
@@ -19,6 +19,7 @@ import { basename, resolve } from "node:path";
import type { ComposeSpecification } from "../schema/docker.d";
import type { Service } from "../schema/docker.d";
import { IMAGE_REGISTRY, LABELS } from "../const";
import { getServicePostgresClaim, isPostgresVolumeEntry } from "./database";
import { toEnvVars } from "./format";
import { deepMerge } from "./shared";
@@ -133,6 +134,8 @@ function parseStringMount(
index: number,
cwd: string,
): NormalizedMount | undefined {
if (isPostgresVolumeEntry(entry)) return;
const parts = entry.split(":");
if (parts.length === 1) {
@@ -479,9 +482,14 @@ export function serviceToDeployment(
name: string,
service: Service,
cwd = process.cwd(),
extraEnv: Record<string, string> = {},
): V1Deployment {
const mounts = toMounts(service, cwd);
const ports = toPorts(service);
const hasEnvSecret =
Boolean(service.env_file) ||
Object.keys(extraEnv).length > 0 ||
Boolean(getServicePostgresClaim(name, service));
return deepMerge(
{
@@ -525,7 +533,7 @@ export function serviceToDeployment(
: service.command
? service.command.split(" ")
: undefined,
envFrom: service.env_file
envFrom: hasEnvSecret
? [{ secretRef: { name: `${name}-env` } }]
: undefined,
ports: toContainerPorts(ports),
@@ -661,8 +669,12 @@ export async function envFromToSecrets(
name: string,
service: Service,
cwd = process.cwd(),
extraEnv: Record<string, string> = {},
): Promise<V1Secret[]> {
const stringData = await readEnvFiles(service.env_file, cwd);
const stringData = {
...(await readEnvFiles(service.env_file, cwd)),
...extraEnv,
};
if (Object.keys(stringData).length === 0) return [];
return [
@@ -709,6 +721,7 @@ export async function composeToKubernetes(
project: string,
compose: ComposeSpecification,
cwd = process.cwd(),
serviceEnv: Record<string, Record<string, string>> = {},
): Promise<KubernetesResource[]> {
const resources = new Map<string, KubernetesResource>();
@@ -724,14 +737,26 @@ export async function composeToKubernetes(
resources.set(getResourceKey(configMap), configMap);
}
for (const secret of await envFromToSecrets(project, name, service, cwd)) {
for (const secret of await envFromToSecrets(
project,
name,
service,
cwd,
serviceEnv[name] ?? {},
)) {
resources.set(getResourceKey(secret), secret);
}
const svc = serviceToSvc(project, name, service);
if (svc) resources.set(getResourceKey(svc), svc);
const deployment = serviceToDeployment(project, name, service, cwd);
const deployment = serviceToDeployment(
project,
name,
service,
cwd,
serviceEnv[name] ?? {},
);
resources.set(getResourceKey(deployment), deployment);
const ingress = serviceToIngress(project, name, service);
+396
View File
@@ -0,0 +1,396 @@
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 = `${DATABASE_CLUSTER}-rw.${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));
}
await applyResource({
apiVersion: "postgresql.cnpg.io/v1",
kind: "Cluster",
metadata: {
name: DATABASE_CLUSTER,
namespace: DATABASE_NAMESPACE,
},
spec: {
managed: {
roles: [...roles.values()],
},
},
});
}
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()) {
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,
},
});
}
}
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);
}
}