feat: add plugin registry and server operations

This commit is contained in:
2026-08-13 17:37:18 +07:00
parent 3f09e330be
commit 71a21b7b79
39 changed files with 2201 additions and 320 deletions
+1
View File
@@ -22,6 +22,7 @@
"typescript": "^7.0.2"
},
"dependencies": {
"@aws-sdk/client-s3": "^3.1109.0",
"@elysiajs/node": "^1.4.5",
"@kubernetes/client-node": "^1.4.0",
"@minikura/api": "workspace:*",
+3 -1
View File
@@ -3,6 +3,7 @@ import { PrismaServerRepository } from "../infrastructure/repositories/prisma/se
import { PrismaUserRepository } from "../infrastructure/repositories/prisma/user.repository.impl";
import { K8sService } from "../services/k8s";
import { OperatorResourceSync } from "../services/operator-resource-sync";
import { PluginRegistryService } from "../services/plugin-registry";
import { WebSocketService } from "../services/websocket";
import { ReverseProxyService } from "./services/reverse-proxy.service";
import { ServerService } from "./services/server.service";
@@ -14,9 +15,10 @@ const reverseProxyRepo = new PrismaReverseProxyRepository();
const webSocketService = new WebSocketService();
const k8sService = new K8sService();
const operatorResourceSync = new OperatorResourceSync();
const pluginRegistryService = new PluginRegistryService();
export const userService = new UserService(userRepo);
export const serverService = new ServerService(serverRepo, k8sService);
export const reverseProxyService = new ReverseProxyService(reverseProxyRepo, k8sService);
export const wsService = webSocketService;
export { k8sService, operatorResourceSync };
export { k8sService, operatorResourceSync, pluginRegistryService };
@@ -13,6 +13,7 @@ export interface IK8sService {
getPods(): Promise<any[]>;
getPodsByLabel(labelSelector: string): Promise<any[]>;
getPodInfo(podName: string): Promise<any>;
restartPod(podName: string): Promise<void>;
getPodLogs(
podName: string,
options?: {
+3 -1
View File
@@ -11,6 +11,7 @@ import { errorHandler } from "./middleware/error-handler";
import { bootstrapRoutes } from "./routes/bootstrap";
import { k8sRoutes } from "./routes/k8s";
import { pluginRoutes } from "./routes/plugin";
import { registryDownloadRoutes, registryRoutes } from "./routes/registry";
import { reverseProxyRoutes } from "./routes/reverse-proxy";
import { serverRoutes } from "./routes/servers";
import { terminalRoutes } from "./routes/terminal";
@@ -37,15 +38,16 @@ const app = new Elysia({ adapter: node() })
})
.all("/auth/*", ({ request }) => auth.handler(request))
.use(bootstrapRoutes)
.group("/api", (app) => app.use(registryDownloadRoutes).use(terminalRoutes))
.use(authPlugin)
.group("/api", (app) =>
app
.use(userRoutes)
.use(serverRoutes)
.use(reverseProxyRoutes)
.use(registryRoutes)
.use(pluginRoutes)
.use(k8sRoutes)
.use(terminalRoutes)
);
export type App = typeof app;
@@ -83,6 +83,7 @@ export class PrismaServerRepository implements ServerRepository {
level_seed: input.level_seed ?? null,
level_type: input.level_type ?? null,
api_key: token,
running: input.running ?? true,
env_variables: input.env_variables
? {
create: input.env_variables.map((ev) => ({
@@ -135,6 +136,7 @@ export class PrismaServerRepository implements ServerRepository {
motd: input.motd,
level_seed: input.level_seed,
level_type: input.level_type,
running: input.running,
},
include: { env_variables: true },
});
+140
View File
@@ -0,0 +1,140 @@
import { Elysia, t } from "elysia";
import { pluginRegistryService, operatorResourceSync } from "../application/di-container";
import { assertAdmin, requireAuth } from "../middleware/auth-guards";
const providerSchema = t.Union([t.Literal("MODRINTH"), t.Literal("HANGAR")]);
const platformSchema = t.Union([t.Literal("PAPER"), t.Literal("FOLIA"), t.Literal("VELOCITY")]);
const versionSchema = t.Object({
provider: providerSchema,
projectId: t.String({ minLength: 1 }),
projectName: t.Optional(t.String({ minLength: 1 })),
versionId: t.String({ minLength: 1 }),
version: t.String({ minLength: 1 }),
platform: platformSchema,
minecraftVersions: t.Array(t.String()),
filename: t.String({ minLength: 1 }),
size: t.Integer({ minimum: 1 }),
sha256: t.Nullable(t.String({ minLength: 64, maxLength: 64 })),
downloadUrl: t.String({ format: "uri" }),
license: t.Nullable(t.String()),
description: t.Optional(t.String()),
author: t.Optional(t.String()),
iconUrl: t.Optional(t.Nullable(t.String())),
projectUrl: t.Optional(t.String({ format: "uri" })),
categories: t.Optional(t.Array(t.String())),
updatedAt: t.Optional(t.Nullable(t.String())),
});
export const registryDownloadRoutes = new Elysia({ prefix: "/registry/artifacts" }).get(
"/:token/download",
async ({ params }) => {
const download = await pluginRegistryService.download(params.token);
if (download.response) return download.response;
if (download.redirect) return Response.redirect(download.redirect, 302);
throw new Error("Plugin artifact download is unavailable");
}
);
export const registryRoutes = new Elysia({ prefix: "/registry" })
.use(requireAuth)
.get(
"/search",
async ({ query, user }) => {
assertAdmin(user);
return pluginRegistryService.search(
query.query ?? "",
query.platform,
query.minecraftVersion
);
},
{
query: t.Object({
query: t.Optional(t.String({ default: "" })),
platform: platformSchema,
minecraftVersion: t.Optional(t.String()),
}),
}
)
.get(
"/versions",
async ({ query, user }) => {
assertAdmin(user);
return pluginRegistryService.versions(
query.provider,
query.projectId,
query.platform,
query.minecraftVersion
);
},
{
query: t.Object({
provider: providerSchema,
projectId: t.String({ minLength: 1 }),
platform: platformSchema,
minecraftVersion: t.Optional(t.String()),
}),
}
)
.get("/artifacts", async ({ user }) => {
assertAdmin(user);
return pluginRegistryService.listArtifacts();
})
.post(
"/artifacts",
async ({ body, user }) => {
assertAdmin(user);
return pluginRegistryService.register(body);
},
{ body: versionSchema }
)
.post(
"/artifacts/upload",
async ({ body, user }) => {
assertAdmin(user);
return pluginRegistryService.upload(body.file);
},
{ body: t.Object({ file: t.File({ type: "application/java-archive" }) }) }
)
.delete("/library/:artifactId", async ({ params, user }) => {
assertAdmin(user);
const serverIds = await pluginRegistryService.deleteArtifact(params.artifactId);
await Promise.all(serverIds.map((serverId) => operatorResourceSync.syncServerById(serverId)));
return { success: true };
})
.get("/servers/:serverId/plugins", async ({ params, user }) => {
assertAdmin(user);
return pluginRegistryService.list(params.serverId);
})
.put(
"/servers/:serverId/plugins",
async ({ params, body, user }) => {
assertAdmin(user);
await pluginRegistryService.reconcile(params.serverId, body.artifactIds);
await operatorResourceSync.syncServerById(params.serverId);
return { success: true };
},
{ body: t.Object({ artifactIds: t.Array(t.String()) }) }
)
.post(
"/servers/:serverId/plugins",
async ({ params, body, user }) => {
assertAdmin(user);
const registered = await pluginRegistryService.register(body);
const artifact = await pluginRegistryService.deploy(params.serverId, registered.id);
await operatorResourceSync.syncServerById(params.serverId);
return artifact;
},
{ body: versionSchema }
)
.post("/servers/:serverId/deployments/:artifactId", async ({ params, user }) => {
assertAdmin(user);
const artifact = await pluginRegistryService.deploy(params.serverId, params.artifactId);
await operatorResourceSync.syncServerById(params.serverId);
return artifact;
})
.delete("/servers/:serverId/installations/:installationId", async ({ params, user }) => {
assertAdmin(user);
await pluginRegistryService.remove(params.serverId, params.installationId);
await operatorResourceSync.syncServerById(params.serverId);
return { success: true };
});
+27 -1
View File
@@ -1,5 +1,6 @@
import { labelKeys } from "@minikura/api";
import { Elysia } from "elysia";
import { serverService, wsService } from "../application/di-container";
import { k8sService, serverService, wsService } from "../application/di-container";
import { bearerToken, findApiKeyOwner } from "../middleware/api-key";
import { assertAdmin, requireAuth } from "../middleware/auth-guards";
import {
@@ -8,6 +9,7 @@ import {
updateServerSchema,
} from "../schemas/server.schema";
import type { WebSocketClient } from "../services/websocket";
import { operatorResourceName } from "../services/operator-resource-sync";
export const serverRoutes = new Elysia({ prefix: "/servers" })
.ws("/ws", {
@@ -79,4 +81,28 @@ export const serverRoutes = new Elysia({ prefix: "/servers" })
assertAdmin(user);
await serverService.deleteEnvVariable(params.id, params.key);
return { success: true };
})
.post("/:id/actions/start", async ({ params, user }) => {
assertAdmin(user);
return await serverService.updateServer(params.id, { running: true });
})
.post("/:id/actions/stop", async ({ params, user }) => {
assertAdmin(user);
return await serverService.updateServer(params.id, { running: false });
})
.post("/:id/actions/restart", async ({ params, user }) => {
assertAdmin(user);
const server = await serverService.getServerById(params.id);
if (!server.running) {
await serverService.updateServer(params.id, { running: true });
return { success: true };
}
const pods = await k8sService.getPodsByLabel(
`${labelKeys.serverId}=${operatorResourceName(params.id)}`
);
await Promise.all(pods.map((pod) => k8sService.restartPod(pod.name)));
return { success: true };
});
+275 -262
View File
@@ -1,12 +1,15 @@
import { isUserSuspended } from "@minikura/db";
import { getErrorMessage } from "@minikura/shared/errors";
import { Elysia } from "elysia";
import type { Elysia } from "elysia";
import WebSocket, { type RawData } from "ws";
import { k8sService } from "../application/di-container";
import { logger } from "../infrastructure/logger";
import { requireAdmin } from "../middleware/auth-guards";
import { auth } from "../middleware/auth";
type TerminalWsData = {
query?: Record<string, string>;
k8sWs?: WebSocket;
request: Request;
};
type TerminalWs = {
@@ -19,307 +22,317 @@ type TerminalMessage =
| { type: "input"; data: string }
| { type: "resize"; cols: number; rows: number };
type BunTlsOptions = {
type TlsOptions = {
rejectUnauthorized: boolean;
cert?: string;
key?: string;
ca?: string;
};
export const terminalRoutes = new Elysia({ prefix: "/terminal" }).use(requireAdmin).ws("/exec", {
open: async (ws: TerminalWs) => {
const podName = ws.data.query?.podName;
const container = ws.data.query?.container;
const shell = ws.data.query?.shell || "/bin/sh";
const mode = ws.data.query?.mode || "shell";
export const terminalRoutes = (app: Elysia) =>
app.ws("/terminal/exec", {
open: async (ws: TerminalWs) => {
const session = await auth.api.getSession({ headers: ws.data.request.headers });
const user = session?.user;
const suspended =
user &&
((user as unknown as { banned?: boolean }).banned === true ||
isUserSuspended(
user as unknown as { isSuspended: boolean; suspendedUntil: Date | null }
));
logger.debug(
`Opening terminal for pod: ${podName}, container: ${container}, shell: ${shell}, mode: ${mode}`
);
if (user?.role !== "admin" || suspended) {
ws.send(JSON.stringify({ type: "error", data: "Admin access required" }));
ws.close();
return;
}
if (!podName) {
ws.send(
JSON.stringify({
type: "error",
data: "Pod name is required",
})
const podName = ws.data.query?.podName;
const container = ws.data.query?.container;
const shell = ws.data.query?.shell || "/bin/sh";
const mode = ws.data.query?.mode || "shell";
logger.debug(
`Opening terminal for pod: ${podName}, container: ${container}, shell: ${shell}, mode: ${mode}`
);
ws.close();
return;
}
try {
if (!k8sService.isInitialized()) {
if (!podName) {
ws.send(
JSON.stringify({
type: "error",
data: "Kubernetes client not initialized",
data: "Pod name is required",
})
);
ws.close();
return;
}
const kc = k8sService.getKubeConfig();
const namespace = k8sService.getNamespace();
const cluster = kc.getCurrentCluster();
const user = kc.getCurrentUser();
if (!cluster) {
throw new Error("No current cluster configured");
}
const server = cluster.server;
const isAttach = mode === "attach";
const apiPath = isAttach
? `/api/v1/namespaces/${namespace}/pods/${podName}/attach`
: `/api/v1/namespaces/${namespace}/pods/${podName}/exec`;
const params = new URLSearchParams({
stdout: "true",
stderr: "true",
stdin: "true",
tty: "true",
});
if (!isAttach) {
params.append("command", shell);
}
if (container) {
params.append("container", container);
}
const wsUrl = `${server}${apiPath}?${params.toString()}`
.replace("https://", "wss://")
.replace("http://", "ws://");
logger.debug(`Connecting to Kubernetes: ${wsUrl}`);
const headers: Record<string, string> = {
Connection: "Upgrade",
Upgrade: "websocket",
"Sec-WebSocket-Version": "13",
"Sec-WebSocket-Key": Buffer.from(Math.random().toString())
.toString("base64")
.substring(0, 24),
"Sec-WebSocket-Protocol": "v4.channel.k8s.io",
};
if (user?.token) {
headers.Authorization = `Bearer ${user.token}`;
} else if (user?.username && user?.password) {
const auth = Buffer.from(`${user.username}:${user.password}`).toString("base64");
headers.Authorization = `Basic ${auth}`;
}
const tlsOptions: BunTlsOptions = {
rejectUnauthorized: cluster.skipTLSVerify !== true,
};
if (user?.certData) {
tlsOptions.cert = Buffer.from(user.certData, "base64").toString();
}
if (user?.keyData) {
tlsOptions.key = Buffer.from(user.keyData, "base64").toString();
}
if (cluster.caData) {
tlsOptions.ca = Buffer.from(cluster.caData, "base64").toString();
}
const wsOptions = { headers, tls: tlsOptions };
const k8sWs = new WebSocket(wsUrl, wsOptions as unknown as string | string[]);
ws.data.k8sWs = k8sWs;
k8sWs.onopen = async () => {
logger.debug(`Connected to Kubernetes ${isAttach ? "attach" : "exec"}`);
if (isAttach) {
try {
const coreApi = k8sService.getCoreApi();
const logs = await coreApi.readNamespacedPodLog({
name: podName,
namespace: namespace,
container: container,
});
if (logs) {
const lines = logs.split("\n");
for (const line of lines) {
ws.send(
JSON.stringify({
type: "output",
data: `${line}\r\n`,
})
);
}
}
ws.send(
JSON.stringify({
type: "ready",
data: "Attached to container (showing logs since start)",
})
);
} catch (logError) {
logger.error({ err: logError }, "Failed to fetch historical logs");
ws.send(
JSON.stringify({
type: "ready",
data: "Attached to container",
})
);
}
} else {
try {
if (!k8sService.isInitialized()) {
ws.send(
JSON.stringify({
type: "ready",
data: "Shell ready",
type: "error",
data: "Kubernetes client not initialized",
})
);
}
};
k8sWs.onmessage = (event: MessageEvent) => {
try {
const data = event.data;
let buffer: Uint8Array;
if (data instanceof Uint8Array) {
buffer = data;
} else if (data instanceof ArrayBuffer) {
buffer = new Uint8Array(data);
} else if (Buffer.isBuffer(data)) {
buffer = new Uint8Array(data);
} else if (data instanceof Blob) {
data.arrayBuffer().then((ab) => {
const uint8 = new Uint8Array(ab);
processBuffer(uint8);
});
return;
} else if (typeof data === "string") {
ws.send(JSON.stringify({ type: "output", data }));
return;
} else {
logger.debug(
{ dataType: typeof data, constructor: data?.constructor?.name },
"Unknown data type"
);
buffer = new Uint8Array(data);
}
processBuffer(buffer);
} catch (err) {
logger.error({ err }, "Error processing Kubernetes message");
}
};
function processBuffer(buffer: Uint8Array): void {
if (buffer.length === 0) {
ws.close();
return;
}
const channel = buffer[0];
const message = new TextDecoder().decode(buffer.slice(1));
const kc = k8sService.getKubeConfig();
const namespace = k8sService.getNamespace();
const cluster = kc.getCurrentCluster();
const user = kc.getCurrentUser();
if (channel === 1 || channel === 2) {
ws.send(JSON.stringify({ type: "output", data: message }));
} else if (channel === 3) {
logger.error({ message }, "Kubernetes error channel");
ws.send(JSON.stringify({ type: "error", data: message }));
if (!cluster) {
throw new Error("No current cluster configured");
}
}
k8sWs.onerror = (error: Event) => {
logger.error({ err: error }, "Kubernetes WebSocket error");
const message = getErrorMessage(error);
const server = cluster.server;
const isConsole = mode === "console";
const apiPath = `/api/v1/namespaces/${namespace}/pods/${podName}/exec`;
const params = new URLSearchParams({
stdout: "true",
stderr: "true",
stdin: "true",
tty: "true",
});
const command = isConsole
? [
"/bin/sh",
"-c",
'tail -n 0 -F /data/logs/latest.log & tail_pid=$!; trap "kill $tail_pid" EXIT; cat > /tmp/minikura-console',
]
: [shell];
for (const part of command) {
params.append("command", part);
}
if (container) {
params.append("container", container);
}
const wsUrl = `${server}${apiPath}?${params.toString()}`
.replace("https://", "wss://")
.replace("http://", "ws://");
logger.debug(`Connecting to Kubernetes: ${wsUrl}`);
const headers: Record<string, string> = {};
if (user?.token) {
headers.Authorization = `Bearer ${user.token}`;
} else if (user?.username && user?.password) {
const auth = Buffer.from(`${user.username}:${user.password}`).toString("base64");
headers.Authorization = `Basic ${auth}`;
}
const tlsOptions: TlsOptions = {
rejectUnauthorized: cluster.skipTLSVerify !== true,
};
if (user?.certData) {
tlsOptions.cert = Buffer.from(user.certData, "base64").toString();
}
if (user?.keyData) {
tlsOptions.key = Buffer.from(user.keyData, "base64").toString();
}
if (cluster.caData) {
tlsOptions.ca = Buffer.from(cluster.caData, "base64").toString();
}
const k8sWs = new WebSocket(wsUrl, "v4.channel.k8s.io", {
headers,
...tlsOptions,
});
ws.data.k8sWs = k8sWs;
k8sWs.on("open", async () => {
logger.debug(`Connected to Kubernetes ${isConsole ? "console" : "shell"}`);
if (isConsole) {
try {
const coreApi = k8sService.getCoreApi();
const logs = await coreApi.readNamespacedPodLog({
name: podName,
namespace: namespace,
container: container,
});
if (logs) {
const lines = logs.split("\n");
for (const line of lines) {
ws.send(
JSON.stringify({
type: "output",
data: `${line}\r\n`,
})
);
}
}
ws.send(
JSON.stringify({
type: "ready",
data: "Minecraft console ready",
})
);
} catch (logError) {
logger.error({ err: logError }, "Failed to fetch historical logs");
ws.send(
JSON.stringify({
type: "ready",
data: "Minecraft console ready",
})
);
}
} else {
ws.send(
JSON.stringify({
type: "ready",
data: "Shell ready",
})
);
}
});
k8sWs.on("message", (data: RawData) => {
try {
let buffer: Uint8Array;
if (data instanceof Uint8Array) {
buffer = data;
} else if (data instanceof ArrayBuffer) {
buffer = new Uint8Array(data);
} else if (Buffer.isBuffer(data)) {
buffer = new Uint8Array(data);
} else if (Array.isArray(data)) {
buffer = new Uint8Array(Buffer.concat(data));
} else {
logger.debug({ dataType: typeof data }, "Unknown data type");
return;
}
processBuffer(buffer);
} catch (err) {
logger.error({ err }, "Error processing Kubernetes message");
}
});
function processBuffer(buffer: Uint8Array): void {
if (buffer.length === 0) {
return;
}
const channel = buffer[0];
const message = new TextDecoder().decode(buffer.slice(1));
if (channel === 1 || channel === 2) {
ws.send(JSON.stringify({ type: "output", data: message }));
} else if (channel === 3) {
try {
const status = JSON.parse(message) as { status?: string; message?: string };
if (status.status !== "Success") {
ws.send(JSON.stringify({ type: "error", data: status.message || message }));
}
} catch {
ws.send(JSON.stringify({ type: "error", data: message }));
}
}
}
k8sWs.on("error", (error: Error) => {
logger.error({ err: error }, "Kubernetes WebSocket error");
const message = getErrorMessage(error);
ws.send(
JSON.stringify({
type: "error",
data: `Connection error: ${message}`,
})
);
});
k8sWs.on("close", (code: number, reason: Buffer) => {
const closeReason = reason.toString();
logger.debug(`Kubernetes WebSocket closed: ${code} ${closeReason}`);
ws.send(
JSON.stringify({
type: "close",
data: closeReason || `Connection closed (${code})`,
})
);
ws.close();
});
} catch (error: unknown) {
logger.error({ err: error }, "Error setting up terminal");
if (error instanceof Error) {
logger.error({ stack: error.stack }, "Error stack");
}
ws.send(
JSON.stringify({
type: "error",
data: `Connection error: ${message}`,
})
);
};
k8sWs.onclose = (event: CloseEvent) => {
logger.debug(`Kubernetes WebSocket closed: ${event.code} ${event.reason}`);
ws.send(
JSON.stringify({
type: "close",
data: event.reason || "Connection closed",
data: `Failed to connect: ${getErrorMessage(error)}`,
})
);
ws.close();
};
} catch (error: unknown) {
logger.error({ err: error }, "Error setting up terminal");
if (error instanceof Error) {
logger.error({ stack: error.stack }, "Error stack");
}
ws.send(
JSON.stringify({
type: "error",
data: `Failed to connect: ${getErrorMessage(error)}`,
})
);
ws.close();
}
},
},
message: async (ws: TerminalWs, message: unknown) => {
try {
const data = parseTerminalMessage(message);
if (!data) {
return;
message: async (ws: TerminalWs, message: unknown) => {
try {
const data = parseTerminalMessage(message);
if (!data) {
return;
}
const k8sWs = ws.data.k8sWs;
if (!k8sWs || k8sWs.readyState !== WebSocket.OPEN) {
logger.error({ readyState: k8sWs?.readyState }, "Kubernetes WebSocket not ready");
return;
}
if (data.type === "input") {
logger.debug({ input: data.data }, "Sending input to k8s");
const encoder = new TextEncoder();
const textData = encoder.encode(data.data);
const buffer = new Uint8Array(1 + textData.length);
buffer[0] = 0;
buffer.set(textData, 1);
k8sWs.send(buffer.buffer);
} else if (data.type === "resize") {
const resizeMsg = JSON.stringify({
Width: data.cols,
Height: data.rows,
});
const encoder = new TextEncoder();
const textData = encoder.encode(resizeMsg);
const buffer = new Uint8Array(1 + textData.length);
buffer[0] = 4;
buffer.set(textData, 1);
k8sWs.send(buffer.buffer);
}
} catch (error: unknown) {
logger.error({ err: error }, "Error handling terminal message");
ws.send(
JSON.stringify({
type: "error",
data: `Error: ${getErrorMessage(error)}`,
})
);
}
},
close: (ws: TerminalWs) => {
logger.debug("Client WebSocket closed");
const k8sWs = ws.data.k8sWs;
if (!k8sWs || k8sWs.readyState !== WebSocket.OPEN) {
logger.error({ readyState: k8sWs?.readyState }, "Kubernetes WebSocket not ready");
return;
if (k8sWs && k8sWs.readyState === WebSocket.OPEN) {
k8sWs.close();
}
if (data.type === "input") {
logger.debug({ input: data.data }, "Sending input to k8s");
const encoder = new TextEncoder();
const textData = encoder.encode(data.data);
const buffer = new Uint8Array(1 + textData.length);
buffer[0] = 0;
buffer.set(textData, 1);
k8sWs.send(buffer.buffer);
} else if (data.type === "resize") {
const resizeMsg = JSON.stringify({
Width: data.cols,
Height: data.rows,
});
const encoder = new TextEncoder();
const textData = encoder.encode(resizeMsg);
const buffer = new Uint8Array(1 + textData.length);
buffer[0] = 4;
buffer.set(textData, 1);
k8sWs.send(buffer.buffer);
}
} catch (error: unknown) {
logger.error({ err: error }, "Error handling terminal message");
ws.send(
JSON.stringify({
type: "error",
data: `Error: ${getErrorMessage(error)}`,
})
);
}
},
close: (ws: TerminalWs) => {
logger.debug("Client WebSocket closed");
const k8sWs = ws.data.k8sWs;
if (k8sWs && k8sWs.readyState === WebSocket.OPEN) {
k8sWs.close();
}
},
});
},
});
function parseTerminalMessage(message: unknown): TerminalMessage | null {
if (typeof message === "string") {
@@ -50,6 +50,7 @@ export const createServerSchema = z.object({
motd: z.string().optional(),
level_seed: z.string().optional(),
level_type: z.string().optional(),
running: z.boolean().optional(),
});
export const updateServerSchema = createServerSchema.omit({ id: true, type: true }).partial();
+9
View File
@@ -162,6 +162,15 @@ export class K8sService implements IK8sService {
return this.resources.getPodInfo(podName);
}
async restartPod(podName: string): Promise<void> {
this.ensureInitialized();
await this.coreApi.deleteNamespacedPod({
name: podName,
namespace: this.namespace,
propagationPolicy: "Foreground",
});
}
async getServiceInfo(serviceName: string) {
this.ensureInitialized();
return this.resources.getServiceInfo(serviceName);
@@ -3,6 +3,7 @@ import { API_GROUP } from "@minikura/api";
import { prisma, type ReverseProxyWithEnvVars, type ServerWithEnvVars } from "@minikura/db";
import { buildKubeConfig } from "@minikura/shared/kube-auth";
import { logger } from "../infrastructure/logger";
import { PluginRegistryService } from "./plugin-registry";
const API_VERSION = "v1alpha1";
const FIELD_MANAGER = "minikura-backend";
@@ -50,6 +51,7 @@ export class OperatorResourceSync {
private coreApi?: k8s.CoreV1Api;
private customObjectsApi?: k8s.CustomObjectsApi;
private syncing = false;
private readonly pluginRegistry = new PluginRegistryService();
constructor() {
try {
@@ -72,7 +74,9 @@ export class OperatorResourceSync {
this.syncing = true;
try {
const [servers, proxies] = await Promise.all([
prisma.server.findMany({ include: { env_variables: true } }),
prisma.server.findMany({
include: { env_variables: true, plugins: { include: { artifact: true } } },
}),
prisma.reverseProxyServer.findMany({ include: { env_variables: true } }),
]);
const syncResults = await Promise.allSettled([
@@ -106,10 +110,10 @@ export class OperatorResourceSync {
}
async syncServerById(id: string): Promise<void> {
this.requireClients();
if (!this.coreApi || !this.customObjectsApi) return;
const server = await prisma.server.findUnique({
where: { id },
include: { env_variables: true },
include: { env_variables: true, plugins: { include: { artifact: true } } },
});
if (server) await this.syncServer(server);
}
@@ -133,7 +137,11 @@ export class OperatorResourceSync {
await this.deleteResource("reverseproxyservers", operatorResourceName(id));
}
private async syncServer(server: ServerWithEnvVars): Promise<void> {
private async syncServer(
server: ServerWithEnvVars & {
plugins?: Array<{ enabled: boolean; artifact: { download_token: string } }>;
}
): Promise<void> {
const name = operatorResourceName(server.id);
const secretName = `mc-${name}-api-key`;
await this.upsertSecret(secretName, server.api_key, labels(server.id));
@@ -147,6 +155,7 @@ export class OperatorResourceSync {
},
spec: {
type: server.type,
running: server.running,
description: server.description ?? undefined,
listenPort: server.listen_port,
serviceType: serviceType(server.service_type),
@@ -163,7 +172,7 @@ export class OperatorResourceSync {
opts: server.jvm_opts ?? undefined,
useAikarFlags: server.use_aikar_flags,
useMeowIceFlags: server.use_meowice_flags,
heapPercent: 80,
heapPercent: 60,
},
properties: {
difficulty: server.difficulty,
@@ -175,13 +184,29 @@ export class OperatorResourceSync {
levelSeed: server.level_seed ?? undefined,
levelType: server.level_type ?? undefined,
},
env: server.env_variables.map((entry) => ({ name: entry.key, value: entry.value })),
env: this.serverEnvironment(server),
apiKeySecretRef: secretName,
},
});
await this.deleteSecret(`${name}-api-key`);
}
private serverEnvironment(
server: ServerWithEnvVars & {
plugins?: Array<{ enabled: boolean; artifact: { download_token: string } }>;
}
): Array<{ name: string; value: string }> {
const environment = new Map(server.env_variables.map((entry) => [entry.key, entry.value]));
const registryUrls = (server.plugins ?? [])
.filter((plugin) => plugin.enabled)
.map((plugin) => this.pluginRegistry.artifactUrl(plugin.artifact.download_token));
if (registryUrls.length > 0) {
const existing = environment.get("PLUGINS");
environment.set("PLUGINS", [existing, ...registryUrls].filter(Boolean).join(","));
}
return [...environment].map(([name, value]) => ({ name, value }));
}
private async syncReverseProxy(proxy: ReverseProxyWithEnvVars): Promise<void> {
const name = operatorResourceName(proxy.id);
const secretName = `rp-${name}-api-key`;
@@ -0,0 +1,541 @@
import { createHash } from "node:crypto";
import {
CreateBucketCommand,
DeleteObjectCommand,
GetObjectCommand,
HeadBucketCommand,
PutObjectCommand,
S3Client,
} from "@aws-sdk/client-s3";
import { type PluginArtifact, type PluginProvider, prisma } from "@minikura/db";
import { NotFoundError, ValidationError } from "../domain/errors/base.error";
const USER_AGENT = process.env.PLUGIN_REGISTRY_USER_AGENT || "Minikura/1.0 (plugin registry)";
const MAX_UPLOAD_BYTES = 100 * 1024 * 1024;
export type RegistryProject = {
provider: "MODRINTH" | "HANGAR";
projectId: string;
name: string;
description: string;
iconUrl: string | null;
downloads: number;
license: string | null;
author: string;
categories: string[];
minecraftVersions: string[];
updatedAt: string | null;
projectUrl: string;
};
export type RegistryVersion = {
provider: "MODRINTH" | "HANGAR";
projectId: string;
projectName?: string;
versionId: string;
version: string;
platform: string;
minecraftVersions: string[];
filename: string;
size: number;
sha256: string | null;
downloadUrl: string;
license: string | null;
description?: string;
author?: string;
iconUrl?: string | null;
projectUrl?: string;
categories?: string[];
updatedAt?: string | null;
};
type ModrinthSearchResponse = {
hits: Array<{
project_id: string;
title: string;
description: string;
icon_url?: string;
downloads: number;
license?: string;
author: string;
categories: string[];
versions: string[];
date_modified: string;
slug?: string;
}>;
};
type ModrinthVersion = {
id: string;
version_number: string;
version_type: string;
game_versions: string[];
loaders: string[];
files: Array<{
filename: string;
size: number;
url: string;
primary: boolean;
hashes: { sha512?: string; sha1?: string };
}>;
};
type HangarProject = {
namespace: { owner: string; slug: string };
name: string;
description: string;
avatarUrl?: string;
stats: { downloads: number };
settings?: { license?: { type?: string; name?: string } };
category?: string;
lastUpdated?: string;
memberNames?: string[] | null;
supportedPlatforms?: Record<string, string[]>;
};
type HangarVersion = {
id: number;
name: string;
channel: { name: string };
downloads: Record<
string,
{
fileInfo: { name: string; sizeBytes: number; sha256Hash: string };
downloadUrl: string | null;
externalUrl: string | null;
}
>;
platformDependencies: Record<string, string[]>;
};
function s3Client(): S3Client {
return new S3Client({
region: process.env.S3_REGION || "us-east-1",
endpoint: process.env.S3_ENDPOINT,
forcePathStyle: process.env.S3_FORCE_PATH_STYLE === "true",
credentials:
process.env.S3_ACCESS_KEY_ID && process.env.S3_SECRET_ACCESS_KEY
? {
accessKeyId: process.env.S3_ACCESS_KEY_ID,
secretAccessKey: process.env.S3_SECRET_ACCESS_KEY,
}
: undefined,
});
}
async function registryFetch<T>(url: string): Promise<T> {
const response = await fetch(url, { headers: { "User-Agent": USER_AGENT } });
if (!response.ok) {
throw new Error(`Plugin provider request failed (${response.status})`);
}
return (await response.json()) as T;
}
function modrinthLoader(platform: string): string {
if (platform === "VELOCITY") return "velocity";
if (platform === "FOLIA") return "folia";
return "paper";
}
export class PluginRegistryService {
async search(query: string, platform: string, minecraftVersion?: string) {
const loader = modrinthLoader(platform);
const facets = [["project_type:plugin"], [`categories:${loader}`]];
if (minecraftVersion && minecraftVersion !== "LATEST") {
facets.push([`versions:${minecraftVersion}`]);
}
const modrinthUrl = new URL("https://api.modrinth.com/v2/search");
modrinthUrl.searchParams.set("query", query);
modrinthUrl.searchParams.set("limit", "20");
modrinthUrl.searchParams.set("index", "downloads");
modrinthUrl.searchParams.set("facets", JSON.stringify(facets));
const hangarUrl = new URL("https://hangar.papermc.io/api/v1/projects");
if (query.trim()) hangarUrl.searchParams.set("query", query);
hangarUrl.searchParams.set("limit", "20");
hangarUrl.searchParams.set("sort", "-downloads");
hangarUrl.searchParams.set("platform", platform === "FOLIA" ? "PAPER" : platform);
const [modrinth, hangar] = await Promise.allSettled([
registryFetch<ModrinthSearchResponse>(modrinthUrl.toString()),
registryFetch<{ result: HangarProject[] }>(hangarUrl.toString()),
]);
const projects: RegistryProject[] = [];
if (modrinth.status === "fulfilled") {
projects.push(
...modrinth.value.hits.map((project) => ({
provider: "MODRINTH" as const,
projectId: project.project_id,
name: project.title,
description: project.description,
iconUrl: project.icon_url ?? null,
downloads: project.downloads,
license: project.license ?? null,
author: project.author,
categories: project.categories,
minecraftVersions: project.versions,
updatedAt: project.date_modified,
projectUrl: `https://modrinth.com/plugin/${project.slug ?? project.project_id}`,
}))
);
}
if (hangar.status === "fulfilled") {
projects.push(
...hangar.value.result.map((project) => ({
provider: "HANGAR" as const,
projectId: `${project.namespace.owner}/${project.namespace.slug}`,
name: project.name,
description: project.description,
iconUrl: project.avatarUrl ?? null,
downloads: project.stats.downloads,
license: project.settings?.license?.type ?? project.settings?.license?.name ?? null,
author: project.memberNames?.[0] ?? project.namespace.owner,
categories: project.category ? [project.category] : [],
minecraftVersions: Object.values(project.supportedPlatforms ?? {}).flat(),
updatedAt: project.lastUpdated ?? null,
projectUrl: `https://hangar.papermc.io/${project.namespace.owner}/${project.namespace.slug}`,
}))
);
}
return projects.sort((a, b) => b.downloads - a.downloads);
}
async versions(
provider: "MODRINTH" | "HANGAR",
projectId: string,
platform: string,
minecraftVersion?: string
): Promise<RegistryVersion[]> {
if (provider === "MODRINTH") {
const url = new URL(
`https://api.modrinth.com/v2/project/${encodeURIComponent(projectId)}/version`
);
url.searchParams.set("loaders", JSON.stringify([modrinthLoader(platform)]));
if (minecraftVersion && minecraftVersion !== "LATEST") {
url.searchParams.set("game_versions", JSON.stringify([minecraftVersion]));
}
url.searchParams.set("include_changelog", "false");
const versions = await registryFetch<ModrinthVersion[]>(url.toString());
return versions
.filter((version) => version.version_type === "release")
.flatMap((version): RegistryVersion[] => {
const file = version.files.find((candidate) => candidate.primary) ?? version.files[0];
return file
? [
{
provider,
projectId,
versionId: version.id,
version: version.version_number,
platform,
minecraftVersions: version.game_versions,
filename: file.filename,
size: file.size,
sha256: null,
downloadUrl: file.url,
license: null,
},
]
: [];
});
}
const [owner, slug] = projectId.split("/", 2);
if (!owner || !slug) throw new ValidationError("Invalid Hangar project ID");
const url = new URL(
`https://hangar.papermc.io/api/v1/projects/${encodeURIComponent(owner)}/${encodeURIComponent(slug)}/versions`
);
url.searchParams.set("limit", "100");
const response = await registryFetch<{ result: HangarVersion[] }>(url.toString());
const resolvedPlatform = platform === "FOLIA" ? "PAPER" : platform;
return response.result
.filter((version) => version.channel.name.toLowerCase() === "release")
.filter(
(version) =>
!minecraftVersion ||
minecraftVersion === "LATEST" ||
version.platformDependencies[resolvedPlatform]?.includes(minecraftVersion)
)
.flatMap((version) => {
const download = version.downloads[resolvedPlatform];
const downloadUrl = download?.downloadUrl ?? download?.externalUrl;
if (!download || !downloadUrl) return [];
return [
{
provider,
projectId,
versionId: String(version.id),
version: version.name,
platform,
minecraftVersions: version.platformDependencies[resolvedPlatform] ?? [],
filename: download.fileInfo.name,
size: download.fileInfo.sizeBytes,
sha256: download.fileInfo.sha256Hash,
downloadUrl,
license: null,
},
];
});
}
async register(input: RegistryVersion): Promise<PluginArtifact> {
this.validateProviderUrl(input.provider, input.downloadUrl);
const sha256 = input.sha256 ?? (await this.hashRemoteFile(input.downloadUrl, input.size));
const artifact = await prisma.pluginArtifact.upsert({
where: {
provider_provider_version_id_platform_filename: {
provider: input.provider,
provider_version_id: input.versionId,
platform: input.platform,
filename: input.filename,
},
},
update: {
source_url: input.downloadUrl,
sha256,
size: input.size,
name: input.projectName ?? input.projectId,
license: input.license,
description: input.description,
author: input.author,
icon_url: input.iconUrl,
project_url: input.projectUrl,
categories: input.categories ?? [],
provider_updated_at: input.updatedAt ? new Date(input.updatedAt) : null,
},
create: {
provider: input.provider,
provider_project_id: input.projectId,
provider_version_id: input.versionId,
name: input.projectName ?? input.projectId,
version: input.version,
platform: input.platform,
minecraft_versions: input.minecraftVersions,
filename: input.filename,
size: input.size,
sha256,
source_url: input.downloadUrl,
license: input.license,
},
});
return artifact;
}
async deploy(serverId: string, artifactId: string): Promise<PluginArtifact> {
const [server, artifact] = await Promise.all([
prisma.server.findUnique({ where: { id: serverId } }),
prisma.pluginArtifact.findUnique({ where: { id: artifactId } }),
]);
if (!server) throw new NotFoundError("Server", serverId);
if (!artifact) throw new NotFoundError("Plugin artifact", artifactId);
await prisma.$transaction([
prisma.serverPlugin.deleteMany({
where: {
server_id: serverId,
artifact: {
provider: artifact.provider,
provider_project_id: artifact.provider_project_id,
id: { not: artifact.id },
},
},
}),
prisma.serverPlugin.upsert({
where: { server_id_artifact_id: { server_id: serverId, artifact_id: artifact.id } },
update: { enabled: true },
create: { server_id: serverId, artifact_id: artifact.id },
}),
]);
return artifact;
}
async upload(file: File): Promise<PluginArtifact> {
if (!file.name.toLowerCase().endsWith(".jar")) {
throw new ValidationError("Plugin upload must be a JAR file");
}
if (file.size === 0 || file.size > MAX_UPLOAD_BYTES) {
throw new ValidationError("Plugin upload must be between 1 byte and 100 MiB");
}
const bytes = new Uint8Array(await file.arrayBuffer());
if (bytes[0] !== 0x50 || bytes[1] !== 0x4b) {
throw new ValidationError("Plugin upload is not a valid JAR archive");
}
const sha256 = createHash("sha256").update(bytes).digest("hex");
const bucket = this.bucket();
const client = s3Client();
await this.ensureBucket(client, bucket);
const objectKey = `artifacts/sha256/${sha256.slice(0, 2)}/${sha256}.jar`;
await client.send(
new PutObjectCommand({
Bucket: bucket,
Key: objectKey,
Body: bytes,
ContentType: "application/java-archive",
Metadata: { sha256, filename: file.name },
})
);
const artifact = await prisma.pluginArtifact.upsert({
where: {
provider_provider_version_id_platform_filename: {
provider: "UPLOAD",
provider_version_id: sha256,
platform: "SERVER",
filename: file.name,
},
},
update: { object_key: objectKey },
create: {
provider: "UPLOAD",
provider_project_id: sha256,
provider_version_id: sha256,
name: file.name.replace(/\.jar$/i, ""),
version: "uploaded",
platform: "SERVER",
minecraft_versions: [],
filename: file.name,
size: file.size,
sha256,
storage_mode: "S3",
object_key: objectKey,
},
});
return artifact;
}
async listArtifacts() {
return prisma.pluginArtifact.findMany({
include: { server_plugins: { select: { server_id: true } } },
orderBy: { created_at: "desc" },
});
}
async deleteArtifact(artifactId: string): Promise<string[]> {
const artifact = await prisma.pluginArtifact.findUnique({
where: { id: artifactId },
include: { server_plugins: { select: { server_id: true } } },
});
if (!artifact) throw new NotFoundError("Plugin artifact", artifactId);
if (artifact.storage_mode === "S3" && artifact.object_key) {
await s3Client().send(
new DeleteObjectCommand({ Bucket: this.bucket(), Key: artifact.object_key })
);
}
await prisma.$transaction([
prisma.serverPlugin.deleteMany({ where: { artifact_id: artifactId } }),
prisma.pluginArtifact.delete({ where: { id: artifactId } }),
]);
return [...new Set(artifact.server_plugins.map((plugin) => plugin.server_id))];
}
async list(serverId: string) {
return prisma.serverPlugin.findMany({
where: { server_id: serverId },
include: { artifact: true },
orderBy: { created_at: "asc" },
});
}
async reconcile(serverId: string, artifactIds: string[]): Promise<void> {
const server = await prisma.server.findUnique({ where: { id: serverId } });
if (!server) throw new NotFoundError("Server", serverId);
const uniqueIds = [...new Set(artifactIds)];
const artifacts = await prisma.pluginArtifact.findMany({ where: { id: { in: uniqueIds } } });
if (artifacts.length !== uniqueIds.length) {
throw new ValidationError("One or more plugin artifacts do not exist");
}
await prisma.$transaction([
prisma.serverPlugin.deleteMany({
where: { server_id: serverId, artifact_id: { notIn: uniqueIds } },
}),
...uniqueIds.map((artifactId) =>
prisma.serverPlugin.upsert({
where: { server_id_artifact_id: { server_id: serverId, artifact_id: artifactId } },
update: { enabled: true },
create: { server_id: serverId, artifact_id: artifactId },
})
),
]);
}
async remove(serverId: string, installationId: string): Promise<void> {
const result = await prisma.serverPlugin.deleteMany({
where: { id: installationId, server_id: serverId },
});
if (result.count === 0) throw new NotFoundError("Server plugin", installationId);
}
async download(token: string): Promise<{ redirect?: string; response?: Response }> {
const artifact = await prisma.pluginArtifact.findUnique({ where: { download_token: token } });
if (!artifact) throw new NotFoundError("Plugin artifact");
if (artifact.storage_mode === "REMOTE" && artifact.source_url) {
return { redirect: artifact.source_url };
}
if (!artifact.object_key) throw new NotFoundError("Plugin artifact object");
const object = await s3Client().send(
new GetObjectCommand({
Bucket: this.bucket(),
Key: artifact.object_key,
})
);
if (!object.Body) throw new NotFoundError("Plugin artifact object");
return {
response: new Response(object.Body.transformToWebStream(), {
headers: {
"Content-Type": object.ContentType || "application/java-archive",
"Content-Length": String(object.ContentLength ?? artifact.size),
"Content-Disposition": `attachment; filename="${artifact.filename.replaceAll('"', "")}"`,
ETag: `"${artifact.sha256}"`,
"Cache-Control": "private, max-age=300, immutable",
},
}),
};
}
artifactUrl(token: string): string {
const baseUrl =
process.env.MINIKURA_PLUGIN_DOWNLOAD_BASE_URL ||
process.env.MINIKURA_OPERATOR_BACKEND_URL ||
"http://minikura-backend:3000/api";
return `${baseUrl.replace(/\/$/, "")}/registry/artifacts/${token}/download`;
}
private async hashRemoteFile(url: string, expectedSize: number): Promise<string> {
if (expectedSize > MAX_UPLOAD_BYTES) throw new ValidationError("Plugin artifact is too large");
const response = await fetch(url, {
headers: { "User-Agent": USER_AGENT },
redirect: "follow",
});
if (!response.ok) throw new ValidationError("Unable to download plugin artifact");
const bytes = new Uint8Array(await response.arrayBuffer());
if (bytes.byteLength !== expectedSize)
throw new ValidationError("Plugin artifact size changed");
return createHash("sha256").update(bytes).digest("hex");
}
private validateProviderUrl(provider: PluginProvider, value: string): void {
const url = new URL(value);
const allowed =
provider === "MODRINTH"
? url.protocol === "https:" && url.hostname === "cdn.modrinth.com"
: provider === "HANGAR"
? url.protocol === "https:" && url.hostname === "hangarcdn.papermc.io"
: false;
if (!allowed)
throw new ValidationError("Plugin artifact URL is not from the selected provider");
}
private bucket(): string {
return process.env.S3_BUCKET || "minikura-plugins";
}
private async ensureBucket(client: S3Client, bucket: string): Promise<void> {
try {
await client.send(new HeadBucketCommand({ Bucket: bucket }));
} catch {
await client.send(new CreateBucketCommand({ Bucket: bucket }));
}
}
}