Phase7
This commit is contained in:
273
apps/worker/src/handlers/idle-billing.ts
Normal file
273
apps/worker/src/handlers/idle-billing.ts
Normal file
@@ -0,0 +1,273 @@
|
||||
import { randomUUID } from 'node:crypto';
|
||||
|
||||
import { prisma } from '@hexahost/database';
|
||||
import {
|
||||
calculateUsageCredits,
|
||||
parseMinecraftListOutput,
|
||||
resolveIdlePolicy,
|
||||
} from '@hexahost/metering';
|
||||
import { Queue } from 'bullmq';
|
||||
|
||||
import { sendAgentRequest } from '../agent-bridge';
|
||||
import { logger } from '../logger';
|
||||
import { getRedisPublisher } from '../redis';
|
||||
import { QUEUE_SERVER_LIFECYCLE } from '../queues';
|
||||
|
||||
let lifecycleQueue: Queue | null = null;
|
||||
|
||||
function getLifecycleQueue(): Queue {
|
||||
if (!lifecycleQueue) {
|
||||
lifecycleQueue = new Queue(QUEUE_SERVER_LIFECYCLE, {
|
||||
connection: getRedisPublisher(),
|
||||
});
|
||||
}
|
||||
|
||||
return lifecycleQueue;
|
||||
}
|
||||
|
||||
async function fetchPlayerCount(server: {
|
||||
id: string;
|
||||
nodeId: string | null;
|
||||
version: number;
|
||||
}): Promise<number | null> {
|
||||
if (!server.nodeId) {
|
||||
return null;
|
||||
}
|
||||
|
||||
const response = await sendAgentRequest(
|
||||
server.nodeId,
|
||||
'server.command',
|
||||
{
|
||||
serverId: server.id,
|
||||
generation: server.version,
|
||||
command: 'list',
|
||||
},
|
||||
10_000,
|
||||
);
|
||||
|
||||
if (!response.success) {
|
||||
return null;
|
||||
}
|
||||
|
||||
const output =
|
||||
typeof response.result === 'object' &&
|
||||
response.result !== null &&
|
||||
'output' in response.result &&
|
||||
typeof (response.result as { output: unknown }).output === 'string'
|
||||
? (response.result as { output: string }).output
|
||||
: typeof response.result === 'string'
|
||||
? response.result
|
||||
: '';
|
||||
|
||||
return parseMinecraftListOutput(output);
|
||||
}
|
||||
|
||||
export async function processIdleShutdownTick(): Promise<void> {
|
||||
const servers = await prisma.gameServer.findMany({
|
||||
where: {
|
||||
status: 'RUNNING',
|
||||
idleShutdownEnabled: true,
|
||||
},
|
||||
include: { plan: true },
|
||||
});
|
||||
|
||||
for (const server of servers) {
|
||||
try {
|
||||
await processServerIdle(server);
|
||||
} catch (error) {
|
||||
logger.error({ err: error, serverId: server.id }, 'Idle shutdown check failed');
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async function processServerIdle(
|
||||
server: Awaited<ReturnType<typeof prisma.gameServer.findMany>>[number] & {
|
||||
plan: {
|
||||
idleTimeoutMinutes: number;
|
||||
idleGraceAfterStartMinutes: number;
|
||||
idleCountdownSeconds: number;
|
||||
} | null;
|
||||
},
|
||||
): Promise<void> {
|
||||
const activeBackup = await prisma.serverBackup.findFirst({
|
||||
where: {
|
||||
serverId: server.id,
|
||||
status: { in: ['PENDING', 'RUNNING'] },
|
||||
},
|
||||
});
|
||||
|
||||
if (activeBackup) {
|
||||
return;
|
||||
}
|
||||
|
||||
const activeRestore = await prisma.backupRestore.findFirst({
|
||||
where: {
|
||||
serverId: server.id,
|
||||
status: { in: ['PENDING', 'RUNNING'] },
|
||||
},
|
||||
});
|
||||
|
||||
if (activeRestore) {
|
||||
return;
|
||||
}
|
||||
|
||||
const policy = resolveIdlePolicy(server.plan);
|
||||
const now = Date.now();
|
||||
const playerCount = await fetchPlayerCount(server);
|
||||
|
||||
if (playerCount !== null && playerCount > 0) {
|
||||
await prisma.gameServer.update({
|
||||
where: { id: server.id },
|
||||
data: {
|
||||
lastPlayerSeenAt: new Date(),
|
||||
idleStopAt: null,
|
||||
},
|
||||
});
|
||||
return;
|
||||
}
|
||||
|
||||
const runningSince = server.runningSince?.getTime() ?? now;
|
||||
if (now - runningSince < policy.idleGraceAfterStartMinutes * 60_000) {
|
||||
return;
|
||||
}
|
||||
|
||||
const lastActive = server.lastPlayerSeenAt?.getTime() ?? runningSince;
|
||||
const idleMs = now - lastActive;
|
||||
|
||||
if (idleMs < policy.idleTimeoutMinutes * 60_000) {
|
||||
if (server.idleStopAt) {
|
||||
await prisma.gameServer.update({
|
||||
where: { id: server.id },
|
||||
data: { idleStopAt: null },
|
||||
});
|
||||
}
|
||||
return;
|
||||
}
|
||||
|
||||
if (!server.idleStopAt) {
|
||||
await prisma.gameServer.update({
|
||||
where: { id: server.id },
|
||||
data: {
|
||||
idleStopAt: new Date(now + policy.idleCountdownSeconds * 1000),
|
||||
},
|
||||
});
|
||||
logger.info({ serverId: server.id }, 'Idle stop countdown started');
|
||||
return;
|
||||
}
|
||||
|
||||
if (server.idleStopAt.getTime() <= now) {
|
||||
const correlationId = randomUUID();
|
||||
await getLifecycleQueue().add(
|
||||
'stop-server',
|
||||
{ serverId: server.id, correlationId, reason: 'idle_shutdown' },
|
||||
{ jobId: correlationId },
|
||||
);
|
||||
|
||||
await prisma.gameServer.update({
|
||||
where: { id: server.id },
|
||||
data: { idleStopAt: null },
|
||||
});
|
||||
|
||||
logger.info({ serverId: server.id }, 'Idle shutdown stop enqueued');
|
||||
}
|
||||
}
|
||||
|
||||
export async function processUsageMeteringTick(): Promise<void> {
|
||||
const servers = await prisma.gameServer.findMany({
|
||||
where: { status: 'RUNNING' },
|
||||
include: { plan: true },
|
||||
});
|
||||
|
||||
const now = new Date();
|
||||
|
||||
for (const server of servers) {
|
||||
try {
|
||||
await meterServerUsage(server, now);
|
||||
} catch (error) {
|
||||
logger.error({ err: error, serverId: server.id }, 'Usage metering failed');
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async function meterServerUsage(
|
||||
server: {
|
||||
id: string;
|
||||
userId: string;
|
||||
ramMb: number;
|
||||
lastMeteredAt: Date | null;
|
||||
runningSince: Date | null;
|
||||
},
|
||||
now: Date,
|
||||
): Promise<void> {
|
||||
const periodStart = server.lastMeteredAt ?? server.runningSince ?? now;
|
||||
const durationSeconds = Math.floor((now.getTime() - periodStart.getTime()) / 1000);
|
||||
|
||||
if (durationSeconds < 30) {
|
||||
return;
|
||||
}
|
||||
|
||||
const credits = calculateUsageCredits(server.ramMb, durationSeconds);
|
||||
const idempotencyKey = `usage:${server.id}:${periodStart.toISOString()}`;
|
||||
|
||||
const wallet = await prisma.creditWallet.findUnique({
|
||||
where: { userId: server.userId },
|
||||
});
|
||||
|
||||
if (!wallet) {
|
||||
await prisma.gameServer.update({
|
||||
where: { id: server.id },
|
||||
data: { lastMeteredAt: now },
|
||||
});
|
||||
return;
|
||||
}
|
||||
|
||||
await prisma.$transaction(async (tx) => {
|
||||
const existing = await tx.usageRecord.findUnique({
|
||||
where: { idempotencyKey },
|
||||
});
|
||||
|
||||
if (existing) {
|
||||
return;
|
||||
}
|
||||
|
||||
await tx.usageRecord.create({
|
||||
data: {
|
||||
serverId: server.id,
|
||||
userId: server.userId,
|
||||
ramMb: server.ramMb,
|
||||
durationSeconds,
|
||||
creditsCharged: credits,
|
||||
periodStart,
|
||||
periodEnd: now,
|
||||
idempotencyKey,
|
||||
},
|
||||
});
|
||||
|
||||
const nextBalance = Math.max(0, wallet.balance - credits);
|
||||
|
||||
await tx.creditWallet.update({
|
||||
where: { id: wallet.id },
|
||||
data: { balance: nextBalance },
|
||||
});
|
||||
|
||||
await tx.creditTransaction.create({
|
||||
data: {
|
||||
walletId: wallet.id,
|
||||
amount: -credits,
|
||||
type: 'USAGE_DEBIT',
|
||||
referenceType: 'server',
|
||||
referenceId: server.id,
|
||||
idempotencyKey: `debit:${idempotencyKey}`,
|
||||
metadata: {
|
||||
durationSeconds,
|
||||
ramMb: server.ramMb,
|
||||
},
|
||||
},
|
||||
});
|
||||
|
||||
await tx.gameServer.update({
|
||||
where: { id: server.id },
|
||||
data: { lastMeteredAt: now },
|
||||
});
|
||||
});
|
||||
}
|
||||
@@ -3,6 +3,7 @@ import { Job, Worker, type ConnectionOptions } from 'bullmq';
|
||||
import { validateConfig } from '@hexahost/config';
|
||||
import { prisma } from '@hexahost/database';
|
||||
|
||||
import { processIdleShutdownTick, processUsageMeteringTick } from './handlers/idle-billing';
|
||||
import { processNodeHealthTick } from './handlers/node-health';
|
||||
import { processNotificationJob } from './handlers/notifications';
|
||||
import { processAddonJob } from './handlers/server-addons';
|
||||
@@ -99,6 +100,12 @@ async function bootstrap(): Promise<void> {
|
||||
void processNodeHealthTick().catch((error: Error) => {
|
||||
logger.error({ err: error }, 'Node health tick failed');
|
||||
});
|
||||
void processIdleShutdownTick().catch((error: Error) => {
|
||||
logger.error({ err: error }, 'Idle shutdown tick failed');
|
||||
});
|
||||
void processUsageMeteringTick().catch((error: Error) => {
|
||||
logger.error({ err: error }, 'Usage metering tick failed');
|
||||
});
|
||||
}, 60_000);
|
||||
|
||||
const shutdown = async (signal: string): Promise<void> => {
|
||||
|
||||
Reference in New Issue
Block a user