Phase6
This commit is contained in:
22
apps/worker/src/allocation.ts
Normal file
22
apps/worker/src/allocation.ts
Normal file
@@ -0,0 +1,22 @@
|
||||
import { prisma } from '@hexahost/database';
|
||||
|
||||
export async function releaseServerAllocation(serverId: string): Promise<void> {
|
||||
const allocation = await prisma.gameServerAllocation.findUnique({
|
||||
where: { serverId },
|
||||
});
|
||||
|
||||
if (!allocation || allocation.status === 'RELEASED') {
|
||||
return;
|
||||
}
|
||||
|
||||
await prisma.$transaction([
|
||||
prisma.gameServerAllocation.update({
|
||||
where: { serverId },
|
||||
data: { status: 'RELEASED' },
|
||||
}),
|
||||
prisma.gameNode.updateMany({
|
||||
where: { id: allocation.nodeId, activeServers: { gt: 0 } },
|
||||
data: { activeServers: { decrement: 1 } },
|
||||
}),
|
||||
]);
|
||||
}
|
||||
105
apps/worker/src/handlers/node-health.ts
Normal file
105
apps/worker/src/handlers/node-health.ts
Normal file
@@ -0,0 +1,105 @@
|
||||
import { prisma } from '@hexahost/database';
|
||||
import { HEARTBEAT_STALE_MS } from '@hexahost/scheduler';
|
||||
|
||||
import { logger } from '../logger';
|
||||
import { transitionServerStatus } from '../server-state';
|
||||
|
||||
const ACTIVE_SERVER_STATUSES = [
|
||||
'RUNNING',
|
||||
'STARTING',
|
||||
'STOPPING',
|
||||
] as const;
|
||||
|
||||
export async function processNodeHealthTick(): Promise<void> {
|
||||
const cutoff = new Date(Date.now() - HEARTBEAT_STALE_MS);
|
||||
|
||||
const staleNodes = await prisma.gameNode.findMany({
|
||||
where: {
|
||||
status: { in: ['ONLINE', 'DRAINING', 'MAINTENANCE'] },
|
||||
OR: [{ lastHeartbeatAt: null }, { lastHeartbeatAt: { lt: cutoff } }],
|
||||
},
|
||||
select: { id: true, name: true },
|
||||
});
|
||||
|
||||
for (const node of staleNodes) {
|
||||
await markNodeUnreachable(node.id);
|
||||
logger.warn({ nodeId: node.id, name: node.name }, 'Node marked unreachable');
|
||||
}
|
||||
|
||||
const drainingNodes = await prisma.gameNode.findMany({
|
||||
where: { status: 'DRAINING' },
|
||||
include: {
|
||||
allocations: {
|
||||
where: { status: 'ACTIVE' },
|
||||
select: { id: true },
|
||||
},
|
||||
},
|
||||
});
|
||||
|
||||
for (const node of drainingNodes) {
|
||||
if (node.allocations.length === 0 && node.drainRequestedAt) {
|
||||
await prisma.gameNode.update({
|
||||
where: { id: node.id },
|
||||
data: {
|
||||
status: 'MAINTENANCE',
|
||||
drainRequestedAt: null,
|
||||
},
|
||||
});
|
||||
|
||||
logger.info({ nodeId: node.id }, 'Drain complete; node set to maintenance');
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async function markNodeUnreachable(nodeId: string): Promise<void> {
|
||||
await prisma.$transaction(async (tx) => {
|
||||
const node = await tx.gameNode.findUnique({ where: { id: nodeId } });
|
||||
if (!node || node.status === 'UNREACHABLE') {
|
||||
return;
|
||||
}
|
||||
|
||||
await tx.gameNode.update({
|
||||
where: { id: nodeId },
|
||||
data: { status: 'UNREACHABLE' },
|
||||
});
|
||||
|
||||
const affectedServers = await tx.gameServer.findMany({
|
||||
where: {
|
||||
nodeId,
|
||||
status: { in: [...ACTIVE_SERVER_STATUSES] },
|
||||
},
|
||||
select: { id: true, status: true },
|
||||
});
|
||||
|
||||
for (const server of affectedServers) {
|
||||
await tx.gameServer.update({
|
||||
where: { id: server.id },
|
||||
data: { status: 'UNKNOWN', version: { increment: 1 } },
|
||||
});
|
||||
|
||||
await tx.gameServerStateTransition.create({
|
||||
data: {
|
||||
gameServerId: server.id,
|
||||
fromStatus: server.status,
|
||||
toStatus: 'UNKNOWN',
|
||||
reason: 'Node heartbeat lost',
|
||||
errorCode: 'NODE_UNREACHABLE',
|
||||
},
|
||||
});
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
export async function reconcileNodeOnline(nodeId: string): Promise<void> {
|
||||
const node = await prisma.gameNode.findUnique({ where: { id: nodeId } });
|
||||
if (!node || node.status !== 'UNREACHABLE') {
|
||||
return;
|
||||
}
|
||||
|
||||
await prisma.gameNode.update({
|
||||
where: { id: nodeId },
|
||||
data: { status: 'ONLINE' },
|
||||
});
|
||||
|
||||
logger.info({ nodeId }, 'Node reconciled back to online');
|
||||
}
|
||||
@@ -1,11 +1,14 @@
|
||||
import { randomUUID } from 'node:crypto';
|
||||
|
||||
import { canStartOnNode } from '@hexahost/scheduler';
|
||||
import { Job } from 'bullmq';
|
||||
import { z } from 'zod';
|
||||
|
||||
import { prisma } from '@hexahost/database';
|
||||
|
||||
import { defaultQueueExpiry } from './start-queue';
|
||||
import { logger } from '../logger';
|
||||
import { loadNodeSnapshot } from '../node-snapshots';
|
||||
import { publishNodeCommand } from '../redis';
|
||||
import { transitionServerStatus } from '../server-state';
|
||||
|
||||
@@ -64,10 +67,56 @@ export async function processStartServerJob(job: Job): Promise<{ status: string
|
||||
const { serverId } = lifecycleJobSchema.parse(job.data);
|
||||
const server = await loadServerForLifecycle(serverId);
|
||||
|
||||
if (server.status !== 'STOPPED') {
|
||||
if (server.status !== 'STOPPED' && server.status !== 'QUEUED') {
|
||||
throw new Error(`Cannot start server ${serverId} from status ${server.status}`);
|
||||
}
|
||||
|
||||
const requiredRamMb = server.ramMb;
|
||||
const nodeSnapshot = await loadNodeSnapshot(server.nodeId!);
|
||||
|
||||
if (!nodeSnapshot || !canStartOnNode(nodeSnapshot, requiredRamMb)) {
|
||||
const correlationId = randomUUID();
|
||||
|
||||
await prisma.$transaction(async (tx) => {
|
||||
await tx.serverStartQueueEntry.upsert({
|
||||
where: { serverId },
|
||||
create: {
|
||||
serverId,
|
||||
nodeId: server.nodeId!,
|
||||
status: 'QUEUED',
|
||||
correlationId,
|
||||
expiresAt: defaultQueueExpiry(),
|
||||
},
|
||||
update: {
|
||||
status: 'QUEUED',
|
||||
correlationId,
|
||||
expiresAt: defaultQueueExpiry(),
|
||||
requestedAt: new Date(),
|
||||
},
|
||||
});
|
||||
|
||||
if (server.status !== 'QUEUED') {
|
||||
await tx.gameServer.update({
|
||||
where: { id: serverId },
|
||||
data: { status: 'QUEUED', version: { increment: 1 } },
|
||||
});
|
||||
|
||||
await tx.gameServerStateTransition.create({
|
||||
data: {
|
||||
gameServerId: serverId,
|
||||
fromStatus: server.status,
|
||||
toStatus: 'QUEUED',
|
||||
reason: 'Insufficient node capacity; queued for start',
|
||||
correlationId,
|
||||
},
|
||||
});
|
||||
}
|
||||
});
|
||||
|
||||
logger.info({ serverId, nodeId: server.nodeId }, 'Server start queued');
|
||||
return { status: 'queued' };
|
||||
}
|
||||
|
||||
const messageId = await publishLifecycleCommand(
|
||||
serverId,
|
||||
server.nodeId!,
|
||||
@@ -76,7 +125,7 @@ export async function processStartServerJob(job: Job): Promise<{ status: string
|
||||
);
|
||||
|
||||
await transitionServerStatus(serverId, 'STARTING', {
|
||||
expectedFrom: 'STOPPED',
|
||||
expectedFrom: ['STOPPED', 'QUEUED'],
|
||||
reason: 'Start requested',
|
||||
correlationId: messageId,
|
||||
});
|
||||
|
||||
@@ -1,10 +1,13 @@
|
||||
import { randomUUID } from 'node:crypto';
|
||||
|
||||
import { prisma } from '@hexahost/database';
|
||||
import { selectNodeForAllocation } from '@hexahost/scheduler';
|
||||
import { Job } from 'bullmq';
|
||||
import { z } from 'zod';
|
||||
|
||||
import { releaseServerAllocation } from '../allocation';
|
||||
import { logger } from '../logger';
|
||||
import { loadNodeSnapshots } from '../node-snapshots';
|
||||
import { publishNodeCommand, waitForNodeResponse } from '../redis';
|
||||
import { transitionServerStatus } from '../server-state';
|
||||
|
||||
@@ -65,10 +68,9 @@ export async function processProvisionServerJob(job: Job): Promise<{ status: str
|
||||
throw new Error(`Server not found: ${serverId}`);
|
||||
}
|
||||
|
||||
const node = await prisma.gameNode.findFirst({
|
||||
where: { status: 'ONLINE' },
|
||||
orderBy: { createdAt: 'asc' },
|
||||
});
|
||||
const nodeSnapshots = await loadNodeSnapshots();
|
||||
const ramMb = server.plan?.maxRamMb ?? server.ramMb;
|
||||
const node = selectNodeForAllocation(nodeSnapshots, ramMb);
|
||||
|
||||
if (!node) {
|
||||
await transitionServerStatus(serverId, 'ERROR', {
|
||||
@@ -80,7 +82,6 @@ export async function processProvisionServerJob(job: Job): Promise<{ status: str
|
||||
}
|
||||
|
||||
const hostPort = await findAvailableHostPort(node.id);
|
||||
const ramMb = server.plan?.maxRamMb ?? server.ramMb;
|
||||
const dataPath = buildDataPath(serverId);
|
||||
const messageId = randomUUID();
|
||||
|
||||
@@ -95,6 +96,11 @@ export async function processProvisionServerJob(job: Job): Promise<{ status: str
|
||||
},
|
||||
});
|
||||
|
||||
await tx.gameNode.update({
|
||||
where: { id: node.id },
|
||||
data: { activeServers: { increment: 1 } },
|
||||
});
|
||||
|
||||
await tx.gameServer.update({
|
||||
where: { id: serverId },
|
||||
data: {
|
||||
@@ -150,6 +156,7 @@ export async function processProvisionServerJob(job: Job): Promise<{ status: str
|
||||
}
|
||||
|
||||
if (!response.success) {
|
||||
await releaseServerAllocation(serverId);
|
||||
await transitionServerStatus(serverId, 'ERROR', {
|
||||
reason: 'Provision failed',
|
||||
correlationId: messageId,
|
||||
|
||||
126
apps/worker/src/handlers/start-queue.ts
Normal file
126
apps/worker/src/handlers/start-queue.ts
Normal file
@@ -0,0 +1,126 @@
|
||||
import { randomUUID } from 'node:crypto';
|
||||
|
||||
import { prisma } from '@hexahost/database';
|
||||
import { canStartOnNode } from '@hexahost/scheduler';
|
||||
import { Queue } from 'bullmq';
|
||||
|
||||
import { logger } from '../logger';
|
||||
import { loadNodeSnapshot } from '../node-snapshots';
|
||||
import { transitionServerStatus } from '../server-state';
|
||||
import { QUEUE_SERVER_LIFECYCLE } from '../queues';
|
||||
import { getRedisPublisher } from '../redis';
|
||||
|
||||
const START_QUEUE_EXPIRY_MS = 30 * 60 * 1000;
|
||||
|
||||
let lifecycleQueue: Queue | null = null;
|
||||
|
||||
function getLifecycleQueue(): Queue {
|
||||
if (!lifecycleQueue) {
|
||||
lifecycleQueue = new Queue(QUEUE_SERVER_LIFECYCLE, {
|
||||
connection: getRedisPublisher(),
|
||||
});
|
||||
}
|
||||
|
||||
return lifecycleQueue;
|
||||
}
|
||||
|
||||
export async function processStartQueueTick(): Promise<void> {
|
||||
const entries = await prisma.serverStartQueueEntry.findMany({
|
||||
where: { status: 'QUEUED' },
|
||||
orderBy: [{ priority: 'desc' }, { requestedAt: 'asc' }],
|
||||
include: {
|
||||
server: {
|
||||
include: { plan: true },
|
||||
},
|
||||
},
|
||||
take: 25,
|
||||
});
|
||||
|
||||
if (entries.length === 0) {
|
||||
return;
|
||||
}
|
||||
|
||||
const now = Date.now();
|
||||
|
||||
for (const entry of entries) {
|
||||
if (entry.expiresAt && entry.expiresAt.getTime() < now) {
|
||||
await expireQueueEntry(entry.serverId);
|
||||
continue;
|
||||
}
|
||||
|
||||
const requiredRamMb = entry.server.plan?.maxRamMb ?? entry.server.ramMb;
|
||||
const nodeSnapshot = await loadNodeSnapshot(entry.nodeId);
|
||||
|
||||
if (!nodeSnapshot || !canStartOnNode(nodeSnapshot, requiredRamMb)) {
|
||||
continue;
|
||||
}
|
||||
|
||||
const correlationId = entry.correlationId ?? randomUUID();
|
||||
|
||||
await prisma.$transaction(async (tx) => {
|
||||
const current = await tx.serverStartQueueEntry.findUnique({
|
||||
where: { id: entry.id },
|
||||
});
|
||||
|
||||
if (!current || current.status !== 'QUEUED') {
|
||||
return;
|
||||
}
|
||||
|
||||
await tx.serverStartQueueEntry.update({
|
||||
where: { id: entry.id },
|
||||
data: { status: 'DISPATCHED', correlationId },
|
||||
});
|
||||
});
|
||||
|
||||
await getLifecycleQueue().add(
|
||||
'start-server',
|
||||
{ serverId: entry.serverId, correlationId },
|
||||
{ jobId: correlationId },
|
||||
);
|
||||
|
||||
logger.info(
|
||||
{ serverId: entry.serverId, nodeId: entry.nodeId },
|
||||
'Dispatched queued server start',
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
async function expireQueueEntry(serverId: string): Promise<void> {
|
||||
await prisma.$transaction(async (tx) => {
|
||||
const entry = await tx.serverStartQueueEntry.findUnique({
|
||||
where: { serverId },
|
||||
});
|
||||
|
||||
if (!entry || entry.status !== 'QUEUED') {
|
||||
return;
|
||||
}
|
||||
|
||||
await tx.serverStartQueueEntry.update({
|
||||
where: { serverId },
|
||||
data: { status: 'EXPIRED' },
|
||||
});
|
||||
|
||||
const server = await tx.gameServer.findUnique({ where: { id: serverId } });
|
||||
if (server?.status === 'QUEUED') {
|
||||
await tx.gameServer.update({
|
||||
where: { id: serverId },
|
||||
data: { status: 'STOPPED', version: { increment: 1 } },
|
||||
});
|
||||
|
||||
await tx.gameServerStateTransition.create({
|
||||
data: {
|
||||
gameServerId: serverId,
|
||||
fromStatus: 'QUEUED',
|
||||
toStatus: 'STOPPED',
|
||||
reason: 'Start queue entry expired',
|
||||
},
|
||||
});
|
||||
}
|
||||
});
|
||||
|
||||
logger.info({ serverId }, 'Expired queued start request');
|
||||
}
|
||||
|
||||
export function defaultQueueExpiry(): Date {
|
||||
return new Date(Date.now() + START_QUEUE_EXPIRY_MS);
|
||||
}
|
||||
@@ -3,11 +3,13 @@ import { Job, Worker, type ConnectionOptions } from 'bullmq';
|
||||
import { validateConfig } from '@hexahost/config';
|
||||
import { prisma } from '@hexahost/database';
|
||||
|
||||
import { processNodeHealthTick } from './handlers/node-health';
|
||||
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';
|
||||
import { processStartQueueTick } from './handlers/start-queue';
|
||||
import { startHealthServer } from './health';
|
||||
import { logger } from './logger';
|
||||
import { closeRedisConnections } from './redis';
|
||||
@@ -91,6 +93,12 @@ async function bootstrap(): Promise<void> {
|
||||
void processScheduledBackupsTick().catch((error: Error) => {
|
||||
logger.error({ err: error }, 'Scheduled backup tick failed');
|
||||
});
|
||||
void processStartQueueTick().catch((error: Error) => {
|
||||
logger.error({ err: error }, 'Start queue tick failed');
|
||||
});
|
||||
void processNodeHealthTick().catch((error: Error) => {
|
||||
logger.error({ err: error }, 'Node health tick failed');
|
||||
});
|
||||
}, 60_000);
|
||||
|
||||
const shutdown = async (signal: string): Promise<void> => {
|
||||
|
||||
59
apps/worker/src/node-snapshots.ts
Normal file
59
apps/worker/src/node-snapshots.ts
Normal file
@@ -0,0 +1,59 @@
|
||||
import { prisma } from '@hexahost/database';
|
||||
import type { NodeSnapshot } from '@hexahost/scheduler';
|
||||
|
||||
function readMemoryField(
|
||||
payload: Record<string, unknown> | null | undefined,
|
||||
key: 'memoryTotalBytes' | 'memoryUsedBytes',
|
||||
): number | null {
|
||||
const value = payload?.[key];
|
||||
return typeof value === 'number' ? value : null;
|
||||
}
|
||||
|
||||
export async function loadNodeSnapshots(
|
||||
nodeIds?: string[],
|
||||
): Promise<NodeSnapshot[]> {
|
||||
const nodes = await prisma.gameNode.findMany({
|
||||
where: nodeIds ? { id: { in: nodeIds } } : undefined,
|
||||
include: {
|
||||
allocations: {
|
||||
where: { status: 'ACTIVE' },
|
||||
select: { ramMbReserved: true },
|
||||
},
|
||||
heartbeats: {
|
||||
orderBy: { createdAt: 'desc' },
|
||||
take: 1,
|
||||
select: { payload: true },
|
||||
},
|
||||
},
|
||||
});
|
||||
|
||||
return nodes.map((node) => {
|
||||
const heartbeatPayload = node.heartbeats[0]?.payload as
|
||||
| Record<string, unknown>
|
||||
| undefined;
|
||||
|
||||
const reservedRamMb = node.allocations.reduce(
|
||||
(sum, allocation) => sum + allocation.ramMbReserved,
|
||||
0,
|
||||
);
|
||||
|
||||
return {
|
||||
id: node.id,
|
||||
status: node.status,
|
||||
maxServers: node.maxServers,
|
||||
maxRamMb: node.maxRamMb,
|
||||
platformRamMb: node.platformRamMb,
|
||||
activeServers: node.activeServers,
|
||||
lastHeartbeatAt: node.lastHeartbeatAt,
|
||||
memoryTotalBytes: readMemoryField(heartbeatPayload, 'memoryTotalBytes'),
|
||||
memoryUsedBytes: readMemoryField(heartbeatPayload, 'memoryUsedBytes'),
|
||||
reservedRamMb,
|
||||
activeAllocationCount: node.allocations.length,
|
||||
};
|
||||
});
|
||||
}
|
||||
|
||||
export async function loadNodeSnapshot(nodeId: string): Promise<NodeSnapshot | null> {
|
||||
const snapshots = await loadNodeSnapshots([nodeId]);
|
||||
return snapshots[0] ?? null;
|
||||
}
|
||||
Reference in New Issue
Block a user