Phase5
This commit is contained in:
171
apps/worker/src/handlers/server-addons.ts
Normal file
171
apps/worker/src/handlers/server-addons.ts
Normal file
@@ -0,0 +1,171 @@
|
||||
import { Job } from 'bullmq';
|
||||
import { z } from 'zod';
|
||||
|
||||
import {
|
||||
getModrinthVersion,
|
||||
isVersionCompatible,
|
||||
pickPrimaryFile,
|
||||
resolveModrinthLoader,
|
||||
} from '@hexahost/catalog';
|
||||
import { prisma } from '@hexahost/database';
|
||||
|
||||
import { sendAgentRequest } from '../agent-bridge';
|
||||
import { logger } from '../logger';
|
||||
|
||||
const installAddonJobSchema = z.object({
|
||||
addonId: z.string().uuid(),
|
||||
serverId: z.string().uuid(),
|
||||
});
|
||||
|
||||
const removeAddonJobSchema = z.object({
|
||||
addonId: z.string().uuid(),
|
||||
serverId: z.string().uuid(),
|
||||
});
|
||||
|
||||
const ADDON_TIMEOUT_MS = 10 * 60 * 1000;
|
||||
|
||||
export async function processInstallAddonJob(job: Job): Promise<{ status: string }> {
|
||||
if (job.name !== 'install-addon') {
|
||||
return { status: 'ignored' };
|
||||
}
|
||||
|
||||
const { addonId, serverId } = installAddonJobSchema.parse(job.data);
|
||||
const addon = await prisma.installedAddon.findUniqueOrThrow({
|
||||
where: { id: addonId },
|
||||
});
|
||||
const server = await prisma.gameServer.findUniqueOrThrow({
|
||||
where: { id: serverId },
|
||||
});
|
||||
|
||||
if (!server.nodeId) {
|
||||
throw new Error(`Server ${serverId} is not assigned to a node`);
|
||||
}
|
||||
|
||||
await prisma.installedAddon.update({
|
||||
where: { id: addonId },
|
||||
data: { status: 'INSTALLING' },
|
||||
});
|
||||
|
||||
try {
|
||||
const version = await getModrinthVersion(addon.versionId);
|
||||
const loader = resolveModrinthLoader(server.softwareFamily);
|
||||
|
||||
if (!loader || !isVersionCompatible(version, server.minecraftVersion, loader)) {
|
||||
throw new Error('Selected version is not compatible with this server');
|
||||
}
|
||||
|
||||
const file = pickPrimaryFile(version);
|
||||
if (!file) {
|
||||
throw new Error('No installable file found for this version');
|
||||
}
|
||||
|
||||
const response = await sendAgentRequest(
|
||||
server.nodeId,
|
||||
'server.addon.install',
|
||||
{
|
||||
serverId,
|
||||
generation: server.version,
|
||||
downloadUrl: file.url,
|
||||
relativePath: addon.filePath,
|
||||
sha512: file.hashes['sha512'] ?? addon.fileSha512 ?? '',
|
||||
},
|
||||
ADDON_TIMEOUT_MS,
|
||||
);
|
||||
|
||||
if (!response.success) {
|
||||
throw new Error(response.errorMessage ?? 'Addon install failed');
|
||||
}
|
||||
|
||||
await prisma.installedAddon.update({
|
||||
where: { id: addonId },
|
||||
data: {
|
||||
status: 'INSTALLED',
|
||||
installedAt: new Date(),
|
||||
fileSha512: file.hashes['sha512'] ?? addon.fileSha512,
|
||||
},
|
||||
});
|
||||
|
||||
return { status: 'installed' };
|
||||
} catch (error) {
|
||||
const message = error instanceof Error ? error.message : 'Addon install failed';
|
||||
await prisma.installedAddon.update({
|
||||
where: { id: addonId },
|
||||
data: {
|
||||
status: 'FAILED',
|
||||
failureReason: message,
|
||||
},
|
||||
});
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
export async function processRemoveAddonJob(job: Job): Promise<{ status: string }> {
|
||||
if (job.name !== 'remove-addon') {
|
||||
return { status: 'ignored' };
|
||||
}
|
||||
|
||||
const { addonId, serverId } = removeAddonJobSchema.parse(job.data);
|
||||
const addon = await prisma.installedAddon.findUniqueOrThrow({
|
||||
where: { id: addonId },
|
||||
});
|
||||
const server = await prisma.gameServer.findUniqueOrThrow({
|
||||
where: { id: serverId },
|
||||
});
|
||||
|
||||
if (!server.nodeId) {
|
||||
throw new Error(`Server ${serverId} is not assigned to a node`);
|
||||
}
|
||||
|
||||
await prisma.installedAddon.update({
|
||||
where: { id: addonId },
|
||||
data: { status: 'REMOVING' },
|
||||
});
|
||||
|
||||
try {
|
||||
const response = await sendAgentRequest(
|
||||
server.nodeId,
|
||||
'server.addon.remove',
|
||||
{
|
||||
serverId,
|
||||
generation: server.version,
|
||||
relativePath: addon.filePath,
|
||||
},
|
||||
ADDON_TIMEOUT_MS,
|
||||
);
|
||||
|
||||
if (!response.success) {
|
||||
throw new Error(response.errorMessage ?? 'Addon remove failed');
|
||||
}
|
||||
|
||||
await prisma.installedAddon.update({
|
||||
where: { id: addonId },
|
||||
data: {
|
||||
status: 'REMOVED',
|
||||
},
|
||||
});
|
||||
|
||||
return { status: 'removed' };
|
||||
} catch (error) {
|
||||
const message = error instanceof Error ? error.message : 'Addon remove failed';
|
||||
await prisma.installedAddon.update({
|
||||
where: { id: addonId },
|
||||
data: {
|
||||
status: 'FAILED',
|
||||
failureReason: message,
|
||||
},
|
||||
});
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
export async function processAddonJob(job: Job): Promise<{ status: string }> {
|
||||
switch (job.name) {
|
||||
case 'install-addon':
|
||||
return processInstallAddonJob(job);
|
||||
case 'remove-addon':
|
||||
return processRemoveAddonJob(job);
|
||||
default:
|
||||
logger.warn({ jobName: job.name }, 'Unknown addon job');
|
||||
return { status: 'ignored' };
|
||||
}
|
||||
}
|
||||
@@ -4,6 +4,7 @@ import { validateConfig } from '@hexahost/config';
|
||||
import { prisma } from '@hexahost/database';
|
||||
|
||||
import { processNotificationJob } from './handlers/notifications';
|
||||
import { processAddonJob } from './handlers/server-addons';
|
||||
import { processBackupJob, processScheduledBackupsTick } from './handlers/server-backups';
|
||||
import { processLifecycleJob } from './handlers/server-lifecycle';
|
||||
import { processProvisionServerJob } from './handlers/server-provisioning';
|
||||
@@ -12,6 +13,7 @@ import { logger } from './logger';
|
||||
import { closeRedisConnections } from './redis';
|
||||
import {
|
||||
QUEUE_NOTIFICATIONS,
|
||||
QUEUE_SERVER_ADDONS,
|
||||
QUEUE_SERVER_BACKUPS,
|
||||
QUEUE_SERVER_LIFECYCLE,
|
||||
QUEUE_SERVER_PROVISIONING,
|
||||
@@ -43,6 +45,10 @@ function createQueueWorker(
|
||||
return processBackupJob(job);
|
||||
}
|
||||
|
||||
if (queueName === QUEUE_SERVER_ADDONS) {
|
||||
return processAddonJob(job);
|
||||
}
|
||||
|
||||
if (queueName === QUEUE_SERVER_LIFECYCLE) {
|
||||
return processLifecycleJob(job);
|
||||
}
|
||||
|
||||
@@ -1,3 +1,4 @@
|
||||
export const QUEUE_SERVER_ADDONS = 'server-addons' as const;
|
||||
export const QUEUE_SERVER_BACKUPS = 'server-backups' as const;
|
||||
export const QUEUE_SERVER_LIFECYCLE = 'server-lifecycle' as const;
|
||||
export const QUEUE_SERVER_PROVISIONING = 'server-provisioning' as const;
|
||||
@@ -8,6 +9,7 @@ import { Queue } from 'bullmq';
|
||||
import { validateConfig } from '@hexahost/config';
|
||||
|
||||
export const WORKER_QUEUES = [
|
||||
QUEUE_SERVER_ADDONS,
|
||||
QUEUE_SERVER_BACKUPS,
|
||||
QUEUE_SERVER_LIFECYCLE,
|
||||
QUEUE_SERVER_PROVISIONING,
|
||||
@@ -16,6 +18,22 @@ export const WORKER_QUEUES = [
|
||||
|
||||
export type WorkerQueueName = (typeof WORKER_QUEUES)[number];
|
||||
|
||||
let addonsQueue: Queue | null = null;
|
||||
|
||||
export function getAddonsQueue(): Queue {
|
||||
if (!addonsQueue) {
|
||||
const config = validateConfig();
|
||||
addonsQueue = new Queue(QUEUE_SERVER_ADDONS, {
|
||||
connection: {
|
||||
url: config.REDIS_URL,
|
||||
maxRetriesPerRequest: null,
|
||||
},
|
||||
});
|
||||
}
|
||||
|
||||
return addonsQueue;
|
||||
}
|
||||
|
||||
let backupsQueue: Queue | null = null;
|
||||
|
||||
export function getBackupsQueue(): Queue {
|
||||
|
||||
Reference in New Issue
Block a user