feat: enhance database and storage management with Longhorn support

This commit is contained in:
2026-09-02 14:07:14 +07:00 Unverified
parent 6c8c8e0c1b
commit 26c2d0f015
12 changed files with 632 additions and 133 deletions
+8 -7
View File
@@ -9,13 +9,14 @@ const ResourceOrder = {
Namespace: 0,
GarageBucket: 1,
GarageKey: 2,
PersistentVolumeClaim: 3,
Secret: 4,
ConfigMap: 5,
Service: 6,
Deployment: 7,
Ingress: 8,
IngressRoute: 9,
StorageClass: 3,
PersistentVolumeClaim: 4,
Secret: 5,
ConfigMap: 6,
Service: 7,
Deployment: 8,
Ingress: 9,
IngressRoute: 10,
} as const;
const ManagedResources = [
+92 -7
View File
@@ -55,6 +55,20 @@ type NormalizedPort = {
type ServiceVolume = NonNullable<Service["volumes"]>[number];
type ComposeVolumes = ComposeSpecification["volumes"];
type LonghornStoragePolicy = {
diskTag?: string[];
replicaCount?: number;
dataLocality?: "none";
};
type LonghornStorageClass = KubernetesObject & {
provisioner: "driver.longhorn.io";
allowVolumeExpansion: boolean;
reclaimPolicy: "Delete";
volumeBindingMode: "Immediate";
parameters: Record<string, string>;
};
const DefaultNamedVolumeStorage = {
requestedStorage: "1Gi",
diskTag: ["fast"],
@@ -62,6 +76,32 @@ const DefaultNamedVolumeStorage = {
dataLocality: "none" as const,
};
function getLonghornStoragePolicy(
mount: NormalizedMount,
): LonghornStoragePolicy | undefined {
if (
!mount.diskTag &&
mount.replicaCount === undefined &&
!mount.dataLocality
) {
return;
}
return {
diskTag: mount.diskTag ? [...mount.diskTag].sort() : undefined,
replicaCount: mount.replicaCount,
dataLocality: mount.dataLocality,
};
}
function toLonghornStorageClassName(policy: LonghornStoragePolicy): string {
const digest = createHash("sha256")
.update(JSON.stringify(policy))
.digest("hex")
.slice(0, 12);
return `kuber-longhorn-${digest}`;
}
function toPortNumber(value: number | string | undefined): number | undefined {
if (typeof value === "number") return value;
if (!value || value.includes("-")) return;
@@ -149,7 +189,9 @@ function normalizeDiskTags(value: unknown): string[] | undefined {
function normalizeReplicaCount(value: unknown): number | undefined {
if (value === undefined || value === null || value === "") return;
const count = Number(value);
return Number.isInteger(count) ? count : undefined;
return Number.isInteger(count) && count >= 1 && count <= 20
? count
: undefined;
}
function parseStorageSpec(value: string): {
@@ -980,6 +1022,7 @@ export function volumesToPvc(
if (!claimName || claims.has(claimName)) continue;
const policy = getLonghornStoragePolicy(mount);
claims.set(claimName, {
apiVersion: "v1",
kind: "PersistentVolumeClaim",
@@ -990,23 +1033,57 @@ export function volumesToPvc(
},
spec: {
accessModes: ["ReadWriteMany"],
storageClassName: policy
? toLonghornStorageClassName(policy)
: undefined,
resources: {
requests: {
storage: mount.requestedStorage ?? "1Gi",
},
},
...(mount.diskTag ? { diskTag: mount.diskTag } : {}),
...(mount.replicaCount !== undefined
? { replicaCount: mount.replicaCount }
: {}),
...(mount.dataLocality ? { dataLocality: mount.dataLocality } : {}),
},
} as V1PersistentVolumeClaim);
});
}
return [...claims.values()];
}
export function volumesToStorageClasses(
service: Service,
cwd = process.cwd(),
volumes: ComposeVolumes = {},
): LonghornStorageClass[] {
const classes = new Map<string, LonghornStorageClass>();
for (const mount of toMounts(service, cwd, volumes)) {
const policy = getLonghornStoragePolicy(mount);
if (!policy) continue;
const name = toLonghornStorageClassName(policy);
classes.set(name, {
apiVersion: "storage.k8s.io/v1",
kind: "StorageClass",
metadata: {
name,
labels: LABELS,
},
provisioner: "driver.longhorn.io",
allowVolumeExpansion: true,
reclaimPolicy: "Delete",
volumeBindingMode: "Immediate",
parameters: {
...(policy.diskTag ? { diskSelector: policy.diskTag.join(",") } : {}),
...(policy.replicaCount !== undefined
? { numberOfReplicas: String(policy.replicaCount) }
: {}),
...(policy.dataLocality ? { dataLocality: "disabled" } : {}),
},
});
}
return [...classes.values()];
}
export function volumesToConfigMaps(
project: string,
service: Service,
@@ -1110,6 +1187,14 @@ export async function composeToKubernetes(
resources.set(getResourceKey(namespace), namespace);
for (const [name, service] of Object.entries(compose.services ?? {})) {
for (const storageClass of volumesToStorageClasses(
service,
cwd,
compose.volumes,
)) {
resources.set(getResourceKey(storageClass), storageClass);
}
for (const pvc of volumesToPvc(project, service, cwd, compose.volumes)) {
resources.set(getResourceKey(pvc), pvc);
}
+23 -7
View File
@@ -9,6 +9,7 @@ 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 REDIS_URL = "redis://redis.database.svc.cluster.local";
export const DATABASE_PROJECT_LABEL = "kuber.dev/project";
export const DATABASE_SERVICE_LABEL = "kuber.dev/service";
@@ -35,7 +36,7 @@ export type PostgresClaim = {
secretName: string;
};
type RoleCredentials = {
export type RoleCredentials = {
username: string;
password: string;
};
@@ -59,6 +60,16 @@ export function buildDatabaseUrl(
return `postgresql://${encodeConnectionComponent(credentials.username)}:${encodeConnectionComponent(credentials.password)}@${DATABASE_HOST}:${DATABASE_PORT}/${encodeConnectionComponent(claim.database)}`;
}
export function buildPostgresEnvironment(
claim: PostgresClaim,
credentials: RoleCredentials,
): Record<string, string> {
return {
DATABASE_URL: buildDatabaseUrl(claim, credentials),
REDIS_URL,
};
}
function parsePostgresVolumeString(
entry: string,
): Omit<PostgresClaim, "service" | "secretName"> | undefined {
@@ -314,6 +325,16 @@ async function reconcileDatabases(
}
}
export async function reconcilePostgresClaim(
project: string,
claim: PostgresClaim,
): Promise<RoleCredentials> {
const credentials = await ensureRoleSecret(claim);
await reconcileManagedRoles([claim]);
await reconcileDatabases(project, [claim]);
return credentials;
}
export async function reconcilePostgresClaims(
project: string,
compose: ComposeSpecification,
@@ -337,12 +358,7 @@ export async function reconcilePostgresClaims(
throw new Error(`Missing credentials for ${claim.secretName}`);
}
return [
claim.service,
{
DATABASE_URL: buildDatabaseUrl(claim, credentials),
},
];
return [claim.service, buildPostgresEnvironment(claim, credentials)];
}),
);
}
+230
View File
@@ -0,0 +1,230 @@
import type {
HttpLibrary,
RequestContext,
} from "@kubernetes/client-node/dist/gen/http/http.js";
import { ResponseContext } from "@kubernetes/client-node/dist/gen/http/http.js";
import { from } from "@kubernetes/client-node/dist/gen/rxjsStub.js";
import http from "node:http";
import https from "node:https";
type TransportOptions = {
maxConcurrent: number;
minIntervalMs: number;
maxRetries: number;
baseRetryMs: number;
maxRetryMs: number;
random: () => number;
};
const DefaultTransportOptions: TransportOptions = {
maxConcurrent: 4,
minIntervalMs: 100,
maxRetries: 3,
baseRetryMs: 250,
maxRetryMs: 5000,
random: Math.random,
};
function abortError(): Error {
const error = new Error("Request aborted");
error.name = "AbortError";
return error;
}
export function delay(ms: number, signal?: AbortSignal): Promise<void> {
if (signal?.aborted) return Promise.reject(abortError());
return new Promise((resolve, reject) => {
const timer = setTimeout(() => {
signal?.removeEventListener("abort", onAbort);
resolve();
}, ms);
const onAbort = () => {
clearTimeout(timer);
reject(abortError());
};
signal?.addEventListener("abort", onAbort, { once: true });
});
}
export function parseRetryAfter(
value: string | undefined,
now = Date.now(),
): number | undefined {
if (!value) return;
const seconds = Number(value);
if (Number.isFinite(seconds) && seconds >= 0) return seconds * 1000;
const date = Date.parse(value);
if (Number.isNaN(date)) return;
return Math.max(0, date - now);
}
class RequestLimiter {
private active = 0;
private nextRequestAt = 0;
private readonly queue: Array<{
signal?: AbortSignal;
onAbort?: () => void;
resolve: (release: () => void) => void;
reject: (error: Error) => void;
}> = [];
constructor(
private readonly maxConcurrent: number,
private readonly minIntervalMs: number,
) {}
acquire(signal?: AbortSignal): Promise<() => void> {
if (signal?.aborted) return Promise.reject(abortError());
return new Promise((resolve, reject) => {
const entry: (typeof this.queue)[number] = { signal, resolve, reject };
entry.onAbort = () => {
const index = this.queue.indexOf(entry);
if (index === -1) return;
this.queue.splice(index, 1);
reject(abortError());
};
signal?.addEventListener("abort", entry.onAbort, { once: true });
this.queue.push(entry);
this.drain();
});
}
private drain(): void {
while (this.active < this.maxConcurrent && this.queue.length > 0) {
const entry = this.queue.shift()!;
entry.signal?.removeEventListener("abort", entry.onAbort!);
if (entry.signal?.aborted) {
entry.reject(abortError());
continue;
}
this.active += 1;
const now = Date.now();
const waitMs = Math.max(0, this.nextRequestAt - now);
this.nextRequestAt =
Math.max(now, this.nextRequestAt) + this.minIntervalMs;
void delay(waitMs, entry.signal)
.then(() => {
let released = false;
entry.resolve(() => {
if (released) return;
released = true;
this.active -= 1;
this.drain();
});
})
.catch((error) => {
this.active -= 1;
entry.reject(
error instanceof Error ? error : new Error(String(error)),
);
this.drain();
});
}
}
}
function sendOnce(request: RequestContext): Promise<ResponseContext> {
return new Promise((resolve, reject) => {
const url = new URL(request.getUrl());
const transport = url.protocol === "http:" ? http : https;
const signal = request.getSignal();
const req = transport.request(
url,
{
method: request.getHttpMethod().toString(),
headers: request.getHeaders(),
agent: request.getAgent() as never,
},
(response) => {
const chunks: Buffer[] = [];
response.on("data", (chunk) => {
chunks.push(Buffer.isBuffer(chunk) ? chunk : Buffer.from(chunk));
});
response.on("end", () => {
signal?.removeEventListener("abort", onAbort);
const buffer = Buffer.concat(chunks);
const headers: Record<string, string> = {};
for (const [key, value] of Object.entries(response.headers)) {
if (Array.isArray(value)) headers[key] = value.join(", ");
else if (value !== undefined) headers[key] = String(value);
}
resolve(
new ResponseContext(response.statusCode ?? 0, headers, {
text: async () => buffer.toString("utf8"),
binary: async () => buffer,
}),
);
});
},
);
const onAbort = () => req.destroy(abortError());
if (signal?.aborted) onAbort();
else signal?.addEventListener("abort", onAbort, { once: true });
req.on("error", (error) => {
signal?.removeEventListener("abort", onAbort);
reject(error);
});
const body = request.getBody();
if (body !== undefined && body !== null) {
req.write(body as string | Uint8Array);
}
req.end();
});
}
export function createKubernetesHttpLibrary(
overrides: Partial<TransportOptions> = {},
): HttpLibrary {
const options = { ...DefaultTransportOptions, ...overrides };
const limiter = new RequestLimiter(
options.maxConcurrent,
options.minIntervalMs,
);
return {
send(request) {
const result = (async () => {
for (let attempt = 0; ; attempt += 1) {
const release = await limiter.acquire(request.getSignal());
let response: ResponseContext;
try {
response = await sendOnce(request);
} finally {
release();
}
if (
response.httpStatusCode !== 429 ||
attempt >= options.maxRetries
) {
return response;
}
const retryAfter = parseRetryAfter(response.headers["retry-after"]);
const exponential = Math.min(
options.maxRetryMs,
options.baseRetryMs * 2 ** attempt,
);
const backoff =
retryAfter ?? exponential * (0.8 + options.random() * 0.4);
await delay(backoff, request.getSignal());
}
})();
return from(result);
},
};
}
+4 -64
View File
@@ -6,73 +6,12 @@ import {
KubernetesObjectApi,
Log,
} from "@kubernetes/client-node";
import { ResponseContext } from "@kubernetes/client-node/dist/gen/http/http.js";
import { from } from "@kubernetes/client-node/dist/gen/rxjsStub.js";
import http from "node:http";
import https from "node:https";
import { createKubernetesHttpLibrary } from "./k8s-http";
const kc = new KubeConfig();
kc.loadFromDefault();
const bunHttpLibrary = {
send(request: {
getUrl(): string;
getHttpMethod(): { toString(): string };
getBody(): unknown;
getHeaders(): Record<string, string>;
getSignal(): AbortSignal | undefined;
getAgent(): http.Agent | https.Agent | undefined;
}) {
const result = new Promise<ResponseContext>((resolve, reject) => {
const url = new URL(request.getUrl());
const transport = url.protocol === "http:" ? http : https;
const req = transport.request(
url,
{
method: request.getHttpMethod().toString(),
headers: request.getHeaders(),
agent: request.getAgent(),
},
(response) => {
const chunks: Buffer[] = [];
response.on("data", (chunk) => {
chunks.push(Buffer.isBuffer(chunk) ? chunk : Buffer.from(chunk));
});
response.on("end", () => {
const buffer = Buffer.concat(chunks);
const headers: Record<string, string> = {};
for (const [key, value] of Object.entries(response.headers)) {
if (Array.isArray(value)) headers[key] = value.join(", ");
else if (value !== undefined) headers[key] = String(value);
}
resolve(
new ResponseContext(response.statusCode ?? 0, headers, {
text: async () => buffer.toString("utf8"),
binary: async () => buffer,
}),
);
});
},
);
request.getSignal()?.addEventListener("abort", () => {
req.destroy(new Error("Request aborted"));
});
req.on("error", reject);
const body = request.getBody();
if (body !== undefined && body !== null) req.write(body as string | Uint8Array);
req.end();
});
return from(result);
},
};
const bunHttpLibrary = createKubernetesHttpLibrary();
const makeApiClient = kc.makeApiClient.bind(kc);
kc.makeApiClient = ((apiClientType) => {
@@ -80,7 +19,8 @@ kc.makeApiClient = ((apiClientType) => {
api?: { configuration?: { httpApi?: typeof bunHttpLibrary } };
configuration?: { httpApi?: typeof bunHttpLibrary };
};
if (client.api?.configuration) client.api.configuration.httpApi = bunHttpLibrary;
if (client.api?.configuration)
client.api.configuration.httpApi = bunHttpLibrary;
if (client.configuration) client.configuration.httpApi = bunHttpLibrary;
return client;
}) as typeof kc.makeApiClient;