Phase4
This commit is contained in:
54
apps/worker/src/agent-bridge.ts
Normal file
54
apps/worker/src/agent-bridge.ts
Normal file
@@ -0,0 +1,54 @@
|
||||
import { randomUUID } from 'node:crypto';
|
||||
|
||||
import { validateConfig } from '@hexahost/config';
|
||||
|
||||
import { logger } from './logger';
|
||||
import {
|
||||
publishNodeCommand,
|
||||
waitForNodeResponse,
|
||||
type NodeResponseMessage,
|
||||
} from './redis';
|
||||
|
||||
export interface AgentRequestResult extends NodeResponseMessage {
|
||||
result?: unknown;
|
||||
}
|
||||
|
||||
export async function sendAgentRequest(
|
||||
nodeId: string,
|
||||
type: string,
|
||||
payload: Record<string, unknown>,
|
||||
timeoutMs = 120_000,
|
||||
): Promise<AgentRequestResult> {
|
||||
const messageId = randomUUID();
|
||||
|
||||
await publishNodeCommand({
|
||||
nodeId,
|
||||
envelope: {
|
||||
protocolVersion: 1,
|
||||
messageId,
|
||||
type,
|
||||
timestamp: new Date().toISOString(),
|
||||
payload: {
|
||||
...payload,
|
||||
generation: payload['generation'] ?? 0,
|
||||
},
|
||||
},
|
||||
});
|
||||
|
||||
const response = await waitForNodeResponse(messageId, timeoutMs);
|
||||
|
||||
if (response === null) {
|
||||
logger.error({ nodeId, type, messageId }, 'Agent request timed out');
|
||||
return {
|
||||
success: false,
|
||||
errorCode: 'AGENT_TIMEOUT',
|
||||
errorMessage: 'Agent request timed out',
|
||||
};
|
||||
}
|
||||
|
||||
return response as AgentRequestResult;
|
||||
}
|
||||
|
||||
export function loadStorageEnv(): void {
|
||||
validateConfig();
|
||||
}
|
||||
632
apps/worker/src/handlers/server-backups.ts
Normal file
632
apps/worker/src/handlers/server-backups.ts
Normal file
@@ -0,0 +1,632 @@
|
||||
import { Job } from 'bullmq';
|
||||
import { z } from 'zod';
|
||||
|
||||
import { prisma } from '@hexahost/database';
|
||||
import {
|
||||
backupObjectKey,
|
||||
createPresignedGetUrl,
|
||||
createPresignedPutUrl,
|
||||
createS3Client,
|
||||
deleteStoredObject,
|
||||
loadStorageConfig,
|
||||
} from '@hexahost/storage';
|
||||
|
||||
import { sendAgentRequest } from '../agent-bridge';
|
||||
import { logger } from '../logger';
|
||||
import { transitionServerStatus } from '../server-state';
|
||||
import { getBackupsQueue } from '../queues';
|
||||
|
||||
const BACKUP_TIMEOUT_MS = 20 * 60 * 1000;
|
||||
|
||||
const createBackupJobSchema = z.object({
|
||||
backupId: z.string().uuid(),
|
||||
serverId: z.string().uuid(),
|
||||
previousStatus: z.enum(['RUNNING', 'STOPPED']),
|
||||
});
|
||||
|
||||
const restoreBackupJobSchema = z.object({
|
||||
restoreId: z.string().uuid(),
|
||||
serverId: z.string().uuid(),
|
||||
backupId: z.string().uuid(),
|
||||
});
|
||||
|
||||
const applyWorldUploadJobSchema = z.object({
|
||||
uploadId: z.string().uuid(),
|
||||
serverId: z.string().uuid(),
|
||||
});
|
||||
|
||||
const retentionJobSchema = z.object({
|
||||
serverId: z.string().uuid(),
|
||||
});
|
||||
|
||||
interface ArchiveResult {
|
||||
localPath: string;
|
||||
sha256: string;
|
||||
sizeBytes: number;
|
||||
worldName: string;
|
||||
}
|
||||
|
||||
function getStorage() {
|
||||
const config = loadStorageConfig();
|
||||
return {
|
||||
config,
|
||||
client: createS3Client(config),
|
||||
};
|
||||
}
|
||||
|
||||
async function loadServer(serverId: string) {
|
||||
const server = await prisma.gameServer.findUnique({
|
||||
where: { id: serverId },
|
||||
include: { backupSchedule: true },
|
||||
});
|
||||
|
||||
if (!server?.nodeId) {
|
||||
throw new Error(`Server ${serverId} is not assigned to a node`);
|
||||
}
|
||||
|
||||
return server;
|
||||
}
|
||||
|
||||
async function markBackupFailed(
|
||||
backupId: string,
|
||||
reason: string,
|
||||
): Promise<void> {
|
||||
await prisma.serverBackup.update({
|
||||
where: { id: backupId },
|
||||
data: {
|
||||
status: 'FAILED',
|
||||
failureReason: reason,
|
||||
completedAt: new Date(),
|
||||
},
|
||||
});
|
||||
}
|
||||
|
||||
async function runBackupArchive(
|
||||
serverId: string,
|
||||
nodeId: string,
|
||||
backupId: string,
|
||||
worldName: string,
|
||||
generation: number,
|
||||
): Promise<ArchiveResult> {
|
||||
const response = await sendAgentRequest(
|
||||
nodeId,
|
||||
'server.world.archive',
|
||||
{
|
||||
serverId,
|
||||
generation,
|
||||
worldName,
|
||||
archiveId: backupId,
|
||||
},
|
||||
BACKUP_TIMEOUT_MS,
|
||||
);
|
||||
|
||||
if (!response.success) {
|
||||
throw new Error(response.errorMessage ?? 'World archive failed');
|
||||
}
|
||||
|
||||
return response.result as ArchiveResult;
|
||||
}
|
||||
|
||||
async function uploadArchiveToS3(
|
||||
serverId: string,
|
||||
nodeId: string,
|
||||
backupId: string,
|
||||
archive: ArchiveResult,
|
||||
generation: number,
|
||||
): Promise<void> {
|
||||
const { config, client } = getStorage();
|
||||
const s3Key = backupObjectKey(serverId, backupId);
|
||||
const uploadUrl = await createPresignedPutUrl(
|
||||
client,
|
||||
config.S3_BUCKET,
|
||||
s3Key,
|
||||
60 * 60,
|
||||
);
|
||||
|
||||
const response = await sendAgentRequest(
|
||||
nodeId,
|
||||
'server.storage.upload',
|
||||
{
|
||||
serverId,
|
||||
generation,
|
||||
localPath: archive.localPath,
|
||||
uploadUrl,
|
||||
},
|
||||
BACKUP_TIMEOUT_MS,
|
||||
);
|
||||
|
||||
if (!response.success) {
|
||||
throw new Error(response.errorMessage ?? 'Backup upload failed');
|
||||
}
|
||||
|
||||
await prisma.serverBackup.update({
|
||||
where: { id: backupId },
|
||||
data: {
|
||||
s3Key,
|
||||
sha256: archive.sha256,
|
||||
sizeBytes: BigInt(archive.sizeBytes),
|
||||
},
|
||||
});
|
||||
}
|
||||
|
||||
async function enforceRetention(serverId: string, retentionCount: number): Promise<void> {
|
||||
const backups = await prisma.serverBackup.findMany({
|
||||
where: {
|
||||
serverId,
|
||||
status: 'AVAILABLE',
|
||||
type: { in: ['MANUAL', 'SCHEDULED'] },
|
||||
},
|
||||
orderBy: { createdAt: 'desc' },
|
||||
});
|
||||
|
||||
const extras = backups.slice(retentionCount);
|
||||
if (extras.length === 0) {
|
||||
return;
|
||||
}
|
||||
|
||||
const { config, client } = getStorage();
|
||||
|
||||
for (const backup of extras) {
|
||||
if (backup.s3Key) {
|
||||
try {
|
||||
await deleteStoredObject(client, config.S3_BUCKET, backup.s3Key);
|
||||
} catch (error) {
|
||||
logger.warn({ backupId: backup.id, error }, 'Failed to delete S3 backup object');
|
||||
}
|
||||
}
|
||||
|
||||
await prisma.serverBackup.update({
|
||||
where: { id: backup.id },
|
||||
data: { status: 'DELETED' },
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
export async function processCreateBackupJob(job: Job): Promise<{ status: string }> {
|
||||
if (job.name !== 'create-backup') {
|
||||
return { status: 'ignored' };
|
||||
}
|
||||
|
||||
const { backupId, serverId, previousStatus } = createBackupJobSchema.parse(job.data);
|
||||
const server = await loadServer(serverId);
|
||||
|
||||
await prisma.serverBackup.update({
|
||||
where: { id: backupId },
|
||||
data: { status: 'RUNNING' },
|
||||
});
|
||||
|
||||
try {
|
||||
if (previousStatus === 'RUNNING') {
|
||||
const prepare = await sendAgentRequest(
|
||||
server.nodeId!,
|
||||
'server.backup.prepare',
|
||||
{ serverId, generation: server.version },
|
||||
BACKUP_TIMEOUT_MS,
|
||||
);
|
||||
if (!prepare.success) {
|
||||
throw new Error(prepare.errorMessage ?? 'Backup prepare failed');
|
||||
}
|
||||
}
|
||||
|
||||
const archive = await runBackupArchive(
|
||||
serverId,
|
||||
server.nodeId!,
|
||||
backupId,
|
||||
server.activeWorldName,
|
||||
server.version,
|
||||
);
|
||||
|
||||
await uploadArchiveToS3(
|
||||
serverId,
|
||||
server.nodeId!,
|
||||
backupId,
|
||||
archive,
|
||||
server.version,
|
||||
);
|
||||
|
||||
if (previousStatus === 'RUNNING') {
|
||||
const release = await sendAgentRequest(
|
||||
server.nodeId!,
|
||||
'server.backup.release',
|
||||
{ serverId, generation: server.version },
|
||||
BACKUP_TIMEOUT_MS,
|
||||
);
|
||||
if (!release.success) {
|
||||
throw new Error(release.errorMessage ?? 'Backup release failed');
|
||||
}
|
||||
}
|
||||
|
||||
await prisma.serverBackup.update({
|
||||
where: { id: backupId },
|
||||
data: {
|
||||
status: 'AVAILABLE',
|
||||
completedAt: new Date(),
|
||||
},
|
||||
});
|
||||
|
||||
await transitionServerStatus(serverId, previousStatus, {
|
||||
expectedFrom: 'BACKING_UP',
|
||||
reason: 'Backup completed',
|
||||
});
|
||||
|
||||
const retentionCount = server.backupSchedule?.retentionCount ?? 5;
|
||||
await enforceRetention(serverId, retentionCount);
|
||||
|
||||
return { status: 'available' };
|
||||
} catch (error) {
|
||||
const message = error instanceof Error ? error.message : 'Backup failed';
|
||||
await markBackupFailed(backupId, message);
|
||||
|
||||
try {
|
||||
if (previousStatus === 'RUNNING') {
|
||||
await sendAgentRequest(
|
||||
server.nodeId!,
|
||||
'server.backup.release',
|
||||
{ serverId, generation: server.version },
|
||||
60_000,
|
||||
);
|
||||
}
|
||||
await transitionServerStatus(serverId, previousStatus, {
|
||||
expectedFrom: 'BACKING_UP',
|
||||
reason: 'Backup failed',
|
||||
errorMessage: message,
|
||||
});
|
||||
} catch (restoreError) {
|
||||
logger.error({ restoreError, serverId }, 'Failed to restore server status after backup error');
|
||||
await transitionServerStatus(serverId, 'ERROR', {
|
||||
expectedFrom: 'BACKING_UP',
|
||||
reason: 'Backup failed and status recovery failed',
|
||||
errorMessage: message,
|
||||
});
|
||||
}
|
||||
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
export async function processRestoreBackupJob(job: Job): Promise<{ status: string }> {
|
||||
if (job.name !== 'restore-backup') {
|
||||
return { status: 'ignored' };
|
||||
}
|
||||
|
||||
const { restoreId, serverId, backupId } = restoreBackupJobSchema.parse(job.data);
|
||||
const server = await loadServer(serverId);
|
||||
const backup = await prisma.serverBackup.findUniqueOrThrow({
|
||||
where: { id: backupId },
|
||||
});
|
||||
|
||||
if (backup.status !== 'AVAILABLE' || !backup.s3Key) {
|
||||
throw new Error('Backup is not available for restore');
|
||||
}
|
||||
|
||||
await prisma.backupRestore.update({
|
||||
where: { id: restoreId },
|
||||
data: { status: 'RUNNING' },
|
||||
});
|
||||
|
||||
const safetyBackup = await prisma.serverBackup.create({
|
||||
data: {
|
||||
serverId,
|
||||
type: 'PRE_RESTORE',
|
||||
status: 'RUNNING',
|
||||
label: 'Pre-restore safety backup',
|
||||
},
|
||||
});
|
||||
|
||||
await prisma.backupRestore.update({
|
||||
where: { id: restoreId },
|
||||
data: { safetyBackupId: safetyBackup.id },
|
||||
});
|
||||
|
||||
try {
|
||||
const safetyArchive = await runBackupArchive(
|
||||
serverId,
|
||||
server.nodeId!,
|
||||
safetyBackup.id,
|
||||
server.activeWorldName,
|
||||
server.version,
|
||||
);
|
||||
await uploadArchiveToS3(
|
||||
serverId,
|
||||
server.nodeId!,
|
||||
safetyBackup.id,
|
||||
safetyArchive,
|
||||
server.version,
|
||||
);
|
||||
await prisma.serverBackup.update({
|
||||
where: { id: safetyBackup.id },
|
||||
data: {
|
||||
status: 'AVAILABLE',
|
||||
completedAt: new Date(),
|
||||
},
|
||||
});
|
||||
} catch (error) {
|
||||
const message = error instanceof Error ? error.message : 'Safety backup failed';
|
||||
await markBackupFailed(safetyBackup.id, message);
|
||||
throw error;
|
||||
}
|
||||
|
||||
const { config, client } = getStorage();
|
||||
const downloadUrl = await createPresignedGetUrl(
|
||||
client,
|
||||
config.S3_BUCKET,
|
||||
backup.s3Key,
|
||||
60 * 60,
|
||||
);
|
||||
const localArchive = `/var/lib/hgc/servers/${serverId}/.hgc/restore/${restoreId}.tar.gz`;
|
||||
|
||||
try {
|
||||
const download = await sendAgentRequest(
|
||||
server.nodeId!,
|
||||
'server.storage.download',
|
||||
{
|
||||
serverId,
|
||||
generation: server.version,
|
||||
downloadUrl,
|
||||
localPath: localArchive,
|
||||
},
|
||||
BACKUP_TIMEOUT_MS,
|
||||
);
|
||||
|
||||
if (!download.success) {
|
||||
throw new Error(download.errorMessage ?? 'Backup download failed');
|
||||
}
|
||||
|
||||
const replace = await sendAgentRequest(
|
||||
server.nodeId!,
|
||||
'server.world.replace',
|
||||
{
|
||||
serverId,
|
||||
generation: server.version,
|
||||
worldName: server.activeWorldName,
|
||||
archivePath: localArchive,
|
||||
},
|
||||
BACKUP_TIMEOUT_MS,
|
||||
);
|
||||
|
||||
if (!replace.success) {
|
||||
throw new Error(replace.errorMessage ?? 'World replace failed');
|
||||
}
|
||||
|
||||
await prisma.backupRestore.update({
|
||||
where: { id: restoreId },
|
||||
data: {
|
||||
status: 'COMPLETED',
|
||||
completedAt: new Date(),
|
||||
},
|
||||
});
|
||||
|
||||
await transitionServerStatus(serverId, 'STOPPED', {
|
||||
expectedFrom: 'RESTORING',
|
||||
reason: 'Restore completed',
|
||||
});
|
||||
|
||||
return { status: 'completed' };
|
||||
} catch (error) {
|
||||
const message = error instanceof Error ? error.message : 'Restore failed';
|
||||
await prisma.backupRestore.update({
|
||||
where: { id: restoreId },
|
||||
data: {
|
||||
status: 'FAILED',
|
||||
failureReason: message,
|
||||
completedAt: new Date(),
|
||||
},
|
||||
});
|
||||
|
||||
await transitionServerStatus(serverId, 'STOPPED', {
|
||||
expectedFrom: 'RESTORING',
|
||||
reason: 'Restore failed',
|
||||
errorMessage: message,
|
||||
});
|
||||
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
export async function processApplyWorldUploadJob(job: Job): Promise<{ status: string }> {
|
||||
if (job.name !== 'apply-world-upload') {
|
||||
return { status: 'ignored' };
|
||||
}
|
||||
|
||||
const { uploadId, serverId } = applyWorldUploadJobSchema.parse(job.data);
|
||||
const server = await loadServer(serverId);
|
||||
const upload = await prisma.worldUpload.findUniqueOrThrow({
|
||||
where: { id: uploadId },
|
||||
});
|
||||
|
||||
await prisma.worldUpload.update({
|
||||
where: { id: uploadId },
|
||||
data: { status: 'VALIDATING' },
|
||||
});
|
||||
|
||||
const safetyBackup = await prisma.serverBackup.create({
|
||||
data: {
|
||||
serverId,
|
||||
type: 'PRE_WORLD_REPLACE',
|
||||
status: 'RUNNING',
|
||||
label: `Pre-replace backup for upload ${uploadId}`,
|
||||
},
|
||||
});
|
||||
|
||||
try {
|
||||
const safetyArchive = await runBackupArchive(
|
||||
serverId,
|
||||
server.nodeId!,
|
||||
safetyBackup.id,
|
||||
server.activeWorldName,
|
||||
server.version,
|
||||
);
|
||||
await uploadArchiveToS3(
|
||||
serverId,
|
||||
server.nodeId!,
|
||||
safetyBackup.id,
|
||||
safetyArchive,
|
||||
server.version,
|
||||
);
|
||||
await prisma.serverBackup.update({
|
||||
where: { id: safetyBackup.id },
|
||||
data: {
|
||||
status: 'AVAILABLE',
|
||||
completedAt: new Date(),
|
||||
},
|
||||
});
|
||||
} catch (error) {
|
||||
const message = error instanceof Error ? error.message : 'Safety backup failed';
|
||||
await markBackupFailed(safetyBackup.id, message);
|
||||
throw error;
|
||||
}
|
||||
|
||||
const { config, client } = getStorage();
|
||||
const downloadUrl = await createPresignedGetUrl(
|
||||
client,
|
||||
config.S3_BUCKET,
|
||||
upload.s3Key,
|
||||
60 * 60,
|
||||
);
|
||||
const localArchive = `/var/lib/hgc/servers/${serverId}/.hgc/uploads/${uploadId}.tar.gz`;
|
||||
|
||||
try {
|
||||
const download = await sendAgentRequest(
|
||||
server.nodeId!,
|
||||
'server.storage.download',
|
||||
{
|
||||
serverId,
|
||||
generation: server.version,
|
||||
downloadUrl,
|
||||
localPath: localArchive,
|
||||
},
|
||||
BACKUP_TIMEOUT_MS,
|
||||
);
|
||||
|
||||
if (!download.success) {
|
||||
throw new Error(download.errorMessage ?? 'Upload download failed');
|
||||
}
|
||||
|
||||
const replace = await sendAgentRequest(
|
||||
server.nodeId!,
|
||||
'server.world.replace',
|
||||
{
|
||||
serverId,
|
||||
generation: server.version,
|
||||
worldName: upload.worldName,
|
||||
archivePath: localArchive,
|
||||
},
|
||||
BACKUP_TIMEOUT_MS,
|
||||
);
|
||||
|
||||
if (!replace.success) {
|
||||
throw new Error(replace.errorMessage ?? 'World replace failed');
|
||||
}
|
||||
|
||||
await prisma.$transaction([
|
||||
prisma.worldUpload.update({
|
||||
where: { id: uploadId },
|
||||
data: {
|
||||
status: 'COMPLETED',
|
||||
completedAt: new Date(),
|
||||
},
|
||||
}),
|
||||
prisma.gameServer.update({
|
||||
where: { id: serverId },
|
||||
data: { activeWorldName: upload.worldName },
|
||||
}),
|
||||
]);
|
||||
|
||||
return { status: 'completed' };
|
||||
} catch (error) {
|
||||
const message = error instanceof Error ? error.message : 'World upload apply failed';
|
||||
await prisma.worldUpload.update({
|
||||
where: { id: uploadId },
|
||||
data: {
|
||||
status: 'FAILED',
|
||||
failureReason: message,
|
||||
completedAt: new Date(),
|
||||
},
|
||||
});
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
export async function processRetentionJob(job: Job): Promise<{ status: string }> {
|
||||
if (job.name !== 'enforce-retention') {
|
||||
return { status: 'ignored' };
|
||||
}
|
||||
|
||||
const { serverId } = retentionJobSchema.parse(job.data);
|
||||
const server = await prisma.gameServer.findUnique({
|
||||
where: { id: serverId },
|
||||
include: { backupSchedule: true },
|
||||
});
|
||||
|
||||
if (!server) {
|
||||
return { status: 'ignored' };
|
||||
}
|
||||
|
||||
await enforceRetention(serverId, server.backupSchedule?.retentionCount ?? 5);
|
||||
return { status: 'ok' };
|
||||
}
|
||||
|
||||
export async function processScheduledBackupsTick(): Promise<void> {
|
||||
const now = new Date();
|
||||
const schedules = await prisma.backupSchedule.findMany({
|
||||
where: {
|
||||
enabled: true,
|
||||
OR: [{ nextRunAt: null }, { nextRunAt: { lte: now } }],
|
||||
},
|
||||
});
|
||||
|
||||
for (const schedule of schedules) {
|
||||
const server = await prisma.gameServer.findUnique({
|
||||
where: { id: schedule.serverId },
|
||||
});
|
||||
|
||||
if (!server?.nodeId || !['RUNNING', 'STOPPED'].includes(server.status)) {
|
||||
continue;
|
||||
}
|
||||
|
||||
const backup = await prisma.serverBackup.create({
|
||||
data: {
|
||||
serverId: schedule.serverId,
|
||||
type: 'SCHEDULED',
|
||||
status: 'PENDING',
|
||||
label: 'Scheduled backup',
|
||||
},
|
||||
});
|
||||
|
||||
await transitionServerStatus(schedule.serverId, 'BACKING_UP', {
|
||||
expectedFrom: server.status as 'RUNNING' | 'STOPPED',
|
||||
reason: 'Scheduled backup started',
|
||||
});
|
||||
|
||||
await getBackupsQueue().add('create-backup', {
|
||||
backupId: backup.id,
|
||||
serverId: schedule.serverId,
|
||||
previousStatus: server.status,
|
||||
});
|
||||
|
||||
const nextRunAt = new Date(now.getTime() + schedule.intervalHours * 60 * 60 * 1000);
|
||||
await prisma.backupSchedule.update({
|
||||
where: { id: schedule.id },
|
||||
data: {
|
||||
lastRunAt: now,
|
||||
nextRunAt,
|
||||
},
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
export async function processBackupJob(job: Job): Promise<{ status: string }> {
|
||||
switch (job.name) {
|
||||
case 'create-backup':
|
||||
return processCreateBackupJob(job);
|
||||
case 'restore-backup':
|
||||
return processRestoreBackupJob(job);
|
||||
case 'apply-world-upload':
|
||||
return processApplyWorldUploadJob(job);
|
||||
case 'enforce-retention':
|
||||
return processRetentionJob(job);
|
||||
default:
|
||||
logger.warn({ jobName: job.name }, 'Unknown backup job');
|
||||
return { status: 'ignored' };
|
||||
}
|
||||
}
|
||||
@@ -4,6 +4,7 @@ import { validateConfig } from '@hexahost/config';
|
||||
import { prisma } from '@hexahost/database';
|
||||
|
||||
import { processNotificationJob } from './handlers/notifications';
|
||||
import { processBackupJob, processScheduledBackupsTick } from './handlers/server-backups';
|
||||
import { processLifecycleJob } from './handlers/server-lifecycle';
|
||||
import { processProvisionServerJob } from './handlers/server-provisioning';
|
||||
import { startHealthServer } from './health';
|
||||
@@ -11,6 +12,7 @@ import { logger } from './logger';
|
||||
import { closeRedisConnections } from './redis';
|
||||
import {
|
||||
QUEUE_NOTIFICATIONS,
|
||||
QUEUE_SERVER_BACKUPS,
|
||||
QUEUE_SERVER_LIFECYCLE,
|
||||
QUEUE_SERVER_PROVISIONING,
|
||||
WORKER_QUEUES,
|
||||
@@ -37,6 +39,10 @@ function createQueueWorker(
|
||||
return processProvisionServerJob(job);
|
||||
}
|
||||
|
||||
if (queueName === QUEUE_SERVER_BACKUPS) {
|
||||
return processBackupJob(job);
|
||||
}
|
||||
|
||||
if (queueName === QUEUE_SERVER_LIFECYCLE) {
|
||||
return processLifecycleJob(job);
|
||||
}
|
||||
@@ -75,9 +81,17 @@ async function bootstrap(): Promise<void> {
|
||||
|
||||
logger.info({ queues: WORKER_QUEUES }, 'Worker started');
|
||||
|
||||
const scheduleTimer = setInterval(() => {
|
||||
void processScheduledBackupsTick().catch((error: Error) => {
|
||||
logger.error({ err: error }, 'Scheduled backup tick failed');
|
||||
});
|
||||
}, 60_000);
|
||||
|
||||
const shutdown = async (signal: string): Promise<void> => {
|
||||
logger.info({ signal }, 'Shutting down worker');
|
||||
|
||||
clearInterval(scheduleTimer);
|
||||
|
||||
await Promise.all(workers.map((worker) => worker.close()));
|
||||
await closeRedisConnections();
|
||||
await prisma.$disconnect();
|
||||
|
||||
@@ -1,11 +1,33 @@
|
||||
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;
|
||||
export const QUEUE_NOTIFICATIONS = 'notifications' as const;
|
||||
|
||||
import { Queue } from 'bullmq';
|
||||
|
||||
import { validateConfig } from '@hexahost/config';
|
||||
|
||||
export const WORKER_QUEUES = [
|
||||
QUEUE_SERVER_BACKUPS,
|
||||
QUEUE_SERVER_LIFECYCLE,
|
||||
QUEUE_SERVER_PROVISIONING,
|
||||
QUEUE_NOTIFICATIONS,
|
||||
] as const;
|
||||
|
||||
export type WorkerQueueName = (typeof WORKER_QUEUES)[number];
|
||||
|
||||
let backupsQueue: Queue | null = null;
|
||||
|
||||
export function getBackupsQueue(): Queue {
|
||||
if (!backupsQueue) {
|
||||
const config = validateConfig();
|
||||
backupsQueue = new Queue(QUEUE_SERVER_BACKUPS, {
|
||||
connection: {
|
||||
url: config.REDIS_URL,
|
||||
maxRetriesPerRequest: null,
|
||||
},
|
||||
});
|
||||
}
|
||||
|
||||
return backupsQueue;
|
||||
}
|
||||
|
||||
@@ -21,6 +21,7 @@ export interface NodeCommandMessage {
|
||||
|
||||
export interface NodeResponseMessage {
|
||||
success: boolean;
|
||||
result?: unknown;
|
||||
resultCode?: string;
|
||||
errorCode?: string;
|
||||
errorMessage?: string;
|
||||
|
||||
Reference in New Issue
Block a user