import type { KubernetesObject, V1Secret } from "@kubernetes/client-node"; import type { ComposeSpecification, Service } from "../schema/docker.d"; import { LABELS } from "../const"; import { applyResource } from "./apply"; import { objectApi } from "./k8s"; export const STORAGE_NAMESPACE = "garage-system"; export const STORAGE_CLUSTER = "garage"; export const STORAGE_PROJECT_LABEL = "kuber.dev/project"; export const STORAGE_SERVICE_LABEL = "kuber.dev/service"; const GARAGE_API_VERSION = "garage.rajsingh.info/v1beta1"; const SECRET_WAIT_TIMEOUT_MS = 60_000; const SECRET_WAIT_INTERVAL_MS = 1_000; export type S3Claim = { service: string; key: string; bucket: string; }; type GarageKeyStatus = { status?: { phase?: string; secretRef?: { name?: string; namespace?: string; }; }; }; function decodeSecretValue(value: string | undefined): string | undefined { return value ? Buffer.from(value, "base64").toString("utf8") : undefined; } async function readObject( resource: KubernetesObject, ): Promise { 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; } } function parseS3VolumeString( entry: string, ): Omit | undefined { if (!entry.startsWith("s3:")) return; const parts = entry.split(":"); if (parts.length !== 2) { throw new Error( `Invalid S3 volume ${entry}. Use s3: or s3:/.`, ); } const target = parts[1]?.trim(); if (!target) { throw new Error( `Invalid S3 volume ${entry}. Use s3: or s3:/.`, ); } const segments = target.split("/"); if ( segments.length > 2 || segments.some((segment) => segment.trim() === "") ) { throw new Error( `Invalid S3 volume ${entry}. Use s3: or s3:/.`, ); } const key = segments[0]!; const bucket = segments[1] ?? key; return { key, bucket }; } export function isS3VolumeEntry( entry: NonNullable[number], ): boolean { return typeof entry === "string" && parseS3VolumeString(entry) !== undefined; } export function getServiceS3Claim( serviceName: string, service: Service, ): S3Claim | undefined { const claims = service.volumes?.flatMap((entry) => { if (typeof entry !== "string") return []; const claim = parseS3VolumeString(entry); return claim ? [{ ...claim, service: serviceName } satisfies S3Claim] : []; }) ?? []; if (claims.length > 1) { throw new Error( `Service ${serviceName} declares multiple S3 volumes. Only zero or one s3:<...> entry is allowed per service.`, ); } return claims[0]; } export function getComposeS3Claims(compose: ComposeSpecification): S3Claim[] { const claims = Object.entries(compose.services ?? {}).flatMap( ([serviceName, service]) => { const claim = getServiceS3Claim(serviceName, service); return claim ? [claim] : []; }, ); const bucketsByKey = new Map(); for (const claim of claims) { const bucket = bucketsByKey.get(claim.key); if (bucket && bucket !== claim.bucket) { throw new Error( `S3 key ${claim.key} is claimed for both ${bucket} and ${claim.bucket}. A managed key can only target one bucket.`, ); } bucketsByKey.set(claim.key, claim.bucket); } return claims; } async function reconcileBuckets( project: string, claims: S3Claim[], ): Promise { const buckets = new Map(); for (const claim of claims) buckets.set(claim.bucket, claim); for (const claim of buckets.values()) { try { await applyResource({ apiVersion: GARAGE_API_VERSION, kind: "GarageBucket", metadata: { name: claim.bucket, namespace: STORAGE_NAMESPACE, labels: { ...LABELS, [STORAGE_PROJECT_LABEL]: project, [STORAGE_SERVICE_LABEL]: claim.service, }, }, spec: { clusterRef: { name: STORAGE_CLUSTER, }, globalAlias: claim.bucket, }, }); } catch (error) { throw new Error( `Failed to reconcile S3 bucket ${claim.bucket}: ${error instanceof Error ? error.message : String(error)}`, ); } } } async function reconcileKeys( project: string, claims: S3Claim[], ): Promise { const keys = new Map(); for (const claim of claims) keys.set(claim.key, claim); for (const claim of keys.values()) { try { await applyResource({ apiVersion: GARAGE_API_VERSION, kind: "GarageKey", metadata: { name: claim.key, namespace: STORAGE_NAMESPACE, labels: { ...LABELS, [STORAGE_PROJECT_LABEL]: project, [STORAGE_SERVICE_LABEL]: claim.service, }, }, spec: { bucketPermissions: [ { bucketRef: { name: claim.bucket, }, owner: true, read: true, write: true, }, ], clusterRef: { name: STORAGE_CLUSTER, }, name: claim.key, neverExpires: true, secretTemplate: { bucketNameKey: "bucket", includeBucketName: true, }, }, }); } catch (error) { throw new Error( `Failed to reconcile S3 key ${claim.key}: ${error instanceof Error ? error.message : String(error)}`, ); } } } async function readKeyEnvironment( claim: S3Claim, ): Promise> { const deadline = Date.now() + SECRET_WAIT_TIMEOUT_MS; let lastPhase: string | undefined; while (Date.now() < deadline) { const key = await readObject({ apiVersion: GARAGE_API_VERSION, kind: "GarageKey", metadata: { name: claim.key, namespace: STORAGE_NAMESPACE, }, }); lastPhase = key?.status?.phase; const secretName = key?.status?.secretRef?.name; const secretNamespace = key?.status?.secretRef?.namespace ?? STORAGE_NAMESPACE; if (secretName) { const secret = await readObject({ apiVersion: "v1", kind: "Secret", metadata: { name: secretName, namespace: secretNamespace, }, }); const accessKeyId = decodeSecretValue(secret?.data?.["access-key-id"]); const secretAccessKey = decodeSecretValue( secret?.data?.["secret-access-key"], ); const endpoint = decodeSecretValue(secret?.data?.endpoint); const region = decodeSecretValue(secret?.data?.region); const bucket = decodeSecretValue(secret?.data?.bucket); if (accessKeyId && secretAccessKey && endpoint && region) { return { AWS_ACCESS_KEY_ID: accessKeyId, AWS_SECRET_ACCESS_KEY: secretAccessKey, AWS_ENDPOINT_URL_S3: endpoint, AWS_REGION: region, S3_BUCKET: bucket ?? claim.bucket, }; } } await new Promise((resolve) => setTimeout(resolve, SECRET_WAIT_INTERVAL_MS), ); } throw new Error( `Timed out waiting for credentials for S3 key ${claim.key} in namespace ${STORAGE_NAMESPACE}${lastPhase ? ` (last phase: ${lastPhase})` : ""}.`, ); } export async function reconcileS3Claims( project: string, compose: ComposeSpecification, ): Promise>> { const claims = getComposeS3Claims(compose); if (claims.length === 0) return {}; await reconcileBuckets(project, claims); await reconcileKeys(project, claims); const environmentsByKey = new Map>(); for (const claim of claims) { if (environmentsByKey.has(claim.key)) continue; environmentsByKey.set(claim.key, await readKeyEnvironment(claim)); } return Object.fromEntries( claims.map((claim) => { const environment = environmentsByKey.get(claim.key); if (!environment) throw new Error(`Missing credentials for ${claim.key}`); return [claim.service, environment]; }), ); } export async function listManagedStorageResources( project: string, ): Promise { const selector = `${Object.entries(LABELS) .map(([key, value]) => `${key}=${value}`) .join(",")},${STORAGE_PROJECT_LABEL}=${project}`; const resources = await Promise.all( ["GarageBucket", "GarageKey"].map(async (kind) => { const result = await objectApi.list( GARAGE_API_VERSION, kind, STORAGE_NAMESPACE, undefined, undefined, undefined, undefined, selector, ); return result.items.map((item) => ({ ...item, apiVersion: item.apiVersion ?? GARAGE_API_VERSION, kind: item.kind ?? kind, metadata: { ...item.metadata, namespace: item.metadata?.namespace ?? STORAGE_NAMESPACE, }, })); }), ); return resources.flat(); }