diff --git a/src/database/models/configuration-entry.model.ts b/src/database/models/configuration-entry.model.ts index 634955881..3847ea245 100644 --- a/src/database/models/configuration-entry.model.ts +++ b/src/database/models/configuration-entry.model.ts @@ -247,6 +247,12 @@ export const configurationSchema = z.discriminatedUnion('key', [ value: z.number().default(minutesToMilliseconds(5)), }) .describe('TF2 Quick Server creation timeout'), + z + .object({ + key: z.literal('log_receiver.use_worker_thread'), + value: z.boolean().default(false), + }) + .describe('Experimental: process game server logs in a dedicated worker thread'), z.object({ key: z.literal('twitchtv.promoted_streams'), value: z diff --git a/src/events.ts b/src/events.ts index 234da5834..54e523edf 100644 --- a/src/events.ts +++ b/src/events.ts @@ -8,7 +8,6 @@ import type { GameModel, GameNumber } from './database/models/game.model' import type { PlayerBan, PlayerModel } from './database/models/player.model' import type { MapPoolEntry } from './database/models/map-pool-entry.model' import type { StaticGameServerModel } from './database/models/static-game-server.model' -import type { LogMessage } from './log-receiver/parse-log-message' import type { Tf2Team } from './shared/types/tf2-team' import type { StreamModel } from './database/models/stream.model' import type { Bot } from './shared/types/bot' @@ -17,9 +16,14 @@ import type { Configuration } from './database/models/configuration-entry.model' import type { MumbleClientStatus } from './mumble/status' import type { ChatMessageModel } from './database/models/chat-message.model' import type { GameSlotId } from './shared/types/game-slot-id' +import type { LogMessage } from './log-receiver/parse-log-message' import type { WithId } from 'mongodb' export interface Events { + 'gamelog:message': { + message: LogMessage + } + 'chat:messageDeleted': { messageId: string } @@ -72,10 +76,6 @@ export interface Events { game: GameModel } - 'gamelog:message': { - message: LogMessage - } - 'match:started': { gameNumber: GameNumber } diff --git a/src/games/log-message-queue.ts b/src/games/log-message-queue.ts new file mode 100644 index 000000000..54a6e422b --- /dev/null +++ b/src/games/log-message-queue.ts @@ -0,0 +1,34 @@ +import { noop } from 'es-toolkit' + +/** + * Queue that ensures operations for the same key are executed sequentially. + * Different keys can be processed concurrently. + */ +export class LogMessageQueue { + private readonly queues = new Map>() + + enqueue(key: string, operation: () => Promise): void { + const previous = this.queues.get(key) ?? Promise.resolve() + const current = previous.then(operation) + + // Store a caught version to prevent unhandled rejection and to not break the chain + this.queues.set(key, current.catch(noop)) + } + + async waitForCompletion(key: string): Promise { + const pending = this.queues.get(key) + if (pending) { + await pending + } + } + + async drain(): Promise { + await Promise.all([...this.queues.keys()].map(key => this.waitForCompletion(key))) + } + + clear(key: string): void { + this.queues.delete(key) + } +} + +export const logMessageQueue = new LogMessageQueue() diff --git a/src/log-receiver/log-receiver.worker.ts b/src/log-receiver/log-receiver.worker.ts new file mode 100644 index 000000000..e42e88447 --- /dev/null +++ b/src/log-receiver/log-receiver.worker.ts @@ -0,0 +1,305 @@ +import { NodeSDK } from '@opentelemetry/sdk-node' +import { PeriodicExportingMetricReader } from '@opentelemetry/sdk-metrics' +import { resourceFromAttributes } from '@opentelemetry/resources' +import { ATTR_SERVICE_NAME, ATTR_SERVICE_VERSION } from '@opentelemetry/semantic-conventions' +import { OTLPMetricExporter } from '@opentelemetry/exporter-metrics-otlp-proto' +import { metrics, ValueType } from '@opentelemetry/api' +import { createSocket } from 'node:dgram' +import { parentPort } from 'node:worker_threads' +import { MongoClient, MongoError } from 'mongodb' +import SteamID from 'steamid' +import { environment } from '../environment' +import { version } from '../version' +import { parseLogMessage } from './parse-log-message' +import { LogMessageQueue } from '../games/log-message-queue' +import { hideIpAddresses } from '../utils/hide-ip-addresses' +import type { WorkerMessage } from './worker-message' +import type { GameNumber } from '../database/models/game.model' +import type { GameLogsModel } from '../database/models/game-logs.model' +import type { SteamId64 } from '../shared/types/steam-id-64' +import type { Tf2Team } from '../shared/types/tf2-team' + +if (!parentPort) { + throw new Error('log-receiver.worker must be run as a worker thread') +} + +const sdk = new NodeSDK({ + resource: resourceFromAttributes({ + [ATTR_SERVICE_NAME]: process.env['OTEL_SERVICE_NAME'] ?? 'tf2pickup.org', + [ATTR_SERVICE_VERSION]: version, + }), + metricReader: new PeriodicExportingMetricReader({ + exporter: new OTLPMetricExporter(), + }), +}) +sdk.start() + +const meter = metrics.getMeter('tf2pickup.server', version) +const messageCount = meter.createCounter('tf2pickup.log_receiver.message.count', { + description: 'Messages coming to the log receiver', + unit: '1', +}) +const eventCount = meter.createCounter('tf2pickup.games.events.count', { + description: 'Game events that come from the gameserver', + unit: '1', + valueType: ValueType.INT, +}) + +const mongoClient = new MongoClient(environment.MONGODB_URI) +await mongoClient.connect() +const db = mongoClient.db() +const gameLogs = db.collection('gamelogs') +const games = db.collection<{ number: GameNumber; logSecret?: string }>('games') + +const logQueue = new LogMessageQueue() + +const fixTeamName = (teamName: string): Tf2Team => teamName.toLowerCase().substring(0, 3) as Tf2Team + +interface WorkerGameEvent { + name: string + keyword: string + regex: RegExp + buildMessage: (gameNumber: GameNumber, matches: RegExpMatchArray) => WorkerMessage | null +} + +const gameEvents: WorkerGameEvent[] = [ + { + name: 'match started', + keyword: 'Round_Start', + regex: /^[\d/\s-:]+World triggered "Round_Start"$/, + buildMessage: gameNumber => ({ type: 'match:started', gameNumber }), + }, + { + name: 'game restarted', + keyword: 'exec etf2l', + regex: /^\d{2}\/\d{2}\/\d{4}\s-\s\d{2}:\d{2}:\d{2}:\srcon from ".+": command "exec etf2l_.+"$/, + buildMessage: gameNumber => ({ type: 'match:restarted', gameNumber }), + }, + { + name: 'round win', + keyword: 'Round_Win', + regex: + /^\d{2}\/\d{2}\/\d{4}\s-\s\d{2}:\d{2}:\d{2}:\sWorld triggered "Round_Win" \(winner "(.+)"\)$/, + buildMessage: (gameNumber, matches) => { + if (!matches[1]) return null + return { type: 'match:roundWon', gameNumber, winner: fixTeamName(matches[1]) } + }, + }, + { + name: 'round length', + keyword: 'Round_Length', + regex: + /^\d{2}\/\d{2}\/\d{4}\s-\s\d{2}:\d{2}:\d{2}:\sWorld triggered "Round_Length" \(seconds "([\d.]+)"\)$/, + buildMessage: (gameNumber, matches) => { + if (!matches[1]) return null + return { type: 'match:roundLength', gameNumber, lengthMs: parseFloat(matches[1]) * 1000 } + }, + }, + { + name: 'match ended', + keyword: 'Game_Over', + regex: /^[\d/\s-:]+World triggered "Game_Over" reason ".*"$/, + buildMessage: gameNumber => ({ type: 'match:ended', gameNumber }), + }, + { + name: 'logs uploaded', + keyword: 'logs.tf', + regex: /^[\d/\s-:]+\[TFTrue\].+\shttp:\/\/logs\.tf\/(\d+)\..*$/, + buildMessage: (gameNumber, matches) => { + if (!matches[1]) return null + return { type: 'match/logs:uploaded', gameNumber, logsUrl: `http://logs.tf/${matches[1]}` } + }, + }, + { + name: 'player connected', + keyword: 'connected, address', + regex: + /^(\d{2}\/\d{2}\/\d{4})\s-\s(\d{2}:\d{2}:\d{2}):\s"(.+)<(\d+)><(\[.[^\]]+\])><>"\sconnected,\saddress\s"(\d{1,3}\.\d{1,3}\.\d{1,3}\.\d{1,3}):(\d{1,5})"$/, + buildMessage: (gameNumber, matches) => { + if (!matches[5] || !matches[6]) return null + const steamId = new SteamID(matches[5]) + if (!steamId.isValid()) return null + return { + type: 'match/player:connected', + gameNumber, + steamId: steamId.getSteamID64() as SteamId64, + ipAddress: matches[6], + } + }, + }, + { + name: 'player joined team', + keyword: 'joined team', + regex: + /^(\d{2}\/\d{2}\/\d{4})\s-\s(\d{2}:\d{2}:\d{2}):\s"(.+)<(\d+)><(\[.[^\]]+\])><(.+)>"\sjoined\steam\s"(.+)"/, + buildMessage: (gameNumber, matches) => { + if (!matches[5] || !matches[7]) return null + const steamId = new SteamID(matches[5]) + if (!steamId.isValid()) return null + return { + type: 'match/player:joinedTeam', + gameNumber, + steamId: steamId.getSteamID64() as SteamId64, + team: fixTeamName(matches[7]), + } + }, + }, + { + name: 'player disconnected', + keyword: 'disconnected', + regex: + /^(\d{2}\/\d{2}\/\d{4})\s-\s(\d{2}:\d{2}:\d{2}):\s"(.+)<(\d+)><(\[.[^\]]+\])><(.[^>]+)>"\sdisconnected\s\(reason\s"(.[^"]+)"\)$/, + buildMessage: (gameNumber, matches) => { + if (!matches[5]) return null + const steamId = new SteamID(matches[5]) + if (!steamId.isValid()) return null + return { + type: 'match/player:disconnected', + gameNumber, + steamId: steamId.getSteamID64() as SteamId64, + } + }, + }, + { + name: 'score reported', + keyword: 'current score', + regex: /^[\d/\s\-:]+Team "(.[^"]+)" current score "(\d)" with "(\d)" players$/, + buildMessage: (gameNumber, matches) => { + const [, teamName, score] = matches + if (!teamName || !score) return null + return { + type: 'match/score:reported', + gameNumber, + teamName: fixTeamName(teamName), + score: Number(score), + } + }, + }, + { + name: 'final score reported', + keyword: 'final score', + regex: /^[\d/\s\-:]+Team "(.[^"]+)" final score "(\d)" with "(\d)" players$/, + buildMessage: (gameNumber, matches) => { + const [, teamName, score] = matches + if (!teamName || !score) return null + return { + type: 'match/score:final', + gameNumber, + team: fixTeamName(teamName), + score: Number(score), + } + }, + }, + { + name: 'demo uploaded', + keyword: '[demos.tf]', + regex: /^[\d/\s-:]+\[demos\.tf\]:\sSTV\savailable\sat:\s(.+)$/, + buildMessage: (gameNumber, matches) => { + if (!matches[1]) return null + return { type: 'match/demo:uploaded', gameNumber, demoUrl: matches[1] } + }, + }, + { + name: 'player said', + keyword: '" say "', + regex: + /^(\d{2}\/\d{2}\/\d{4})\s-\s(\d{2}:\d{2}:\d{2}):\s"(.+)<(\d+)><(\[.[^\]]+\])><(.[^>]+)>"\ssay\s"(.+)"$/, + buildMessage: (gameNumber, matches) => { + if (!matches[5] || !matches[7]) return null + const steamId = new SteamID(matches[5]) + if (!steamId.isValid()) return null + return { + type: 'match/player:said', + gameNumber, + steamId: steamId.getSteamID64() as SteamId64, + message: matches[7], + } + }, + }, +] + +async function pushLogMessage(payload: string, logSecret: string): Promise { + const logLine = hideIpAddresses(payload) + try { + await gameLogs.findOneAndUpdate({ logSecret }, { $push: { logs: logLine } }, { upsert: true }) + } catch (error) { + if (!(error instanceof MongoError)) throw error + if (error.code === 11000) { + await gameLogs.findOneAndUpdate({ logSecret }, { $push: { logs: logLine } }, { upsert: true }) + } else { + throw error + } + } +} + +async function matchGameEvent(payload: string, logSecret: string): Promise { + for (const gameEvent of gameEvents) { + if (!payload.includes(gameEvent.keyword)) continue + const matches = payload.match(gameEvent.regex) + if (!matches) continue + + const game = await games.findOne({ logSecret }, { projection: { number: 1 } }) + if (!game) return null + + const message = gameEvent.buildMessage(game.number, matches) + if (!message) return null + + if (message.type === 'match:restarted') { + logQueue.clear(logSecret) + await gameLogs.deleteOne({ logSecret }) + } + + return message + } + return null +} + +const socket = createSocket('udp4') + +socket.on('message', (message, rinfo) => { + try { + messageCount.add(1, { source_ip: rinfo.address }) + const logMessage = parseLogMessage(message) + + logQueue.enqueue(logMessage.password, () => + pushLogMessage(logMessage.payload, logMessage.password), + ) + + matchGameEvent(logMessage.payload, logMessage.password).then( + workerMessage => { + if (workerMessage) { + parentPort!.postMessage(workerMessage) + eventCount.add(1, { + 'tf2pickup.games.event.handled': true, + 'tf2pickup.game.number': workerMessage.gameNumber, + }) + } else { + eventCount.add(1, { 'tf2pickup.games.event.handled': false }) + } + }, + () => { + eventCount.add(1, { 'tf2pickup.games.event.handled': false }) + }, + ) + } catch { + // invalid message, ignore + } +}) + +socket.on('listening', () => { + const address = socket.address() + console.log(`log receiver worker listening at ${address.address}:${address.port}`) +}) + +socket.bind(environment.LOG_RELAY_PORT, '0.0.0.0') + +parentPort.on('message', async () => { + try { + socket.close() + await logQueue.drain() + await mongoClient.close() + await sdk.shutdown() + } finally { + process.exit(0) + } +}) diff --git a/src/log-receiver/plugins/index.ts b/src/log-receiver/plugins/index.ts index 02792f5e7..a998f9af2 100644 --- a/src/log-receiver/plugins/index.ts +++ b/src/log-receiver/plugins/index.ts @@ -1,43 +1,142 @@ import fp from 'fastify-plugin' import { createSocket } from 'node:dgram' -import { parseLogMessage } from '../parse-log-message' +import { Worker } from 'node:worker_threads' +import { events } from '../../events' import { logger } from '../../logger' import { environment } from '../../environment' -import { events } from '../../events' import { meter } from '../../otel' +import { parseLogMessage } from '../parse-log-message' +import { configuration } from '../../configuration' +import type { WorkerMessage, ControlMessage } from '../worker-message' +import { withTimeout } from 'es-toolkit' +import { secondsToMilliseconds } from 'date-fns' export default fp( - // eslint-disable-next-line @typescript-eslint/require-await async app => { - const messageCount = meter.createCounter('tf2pickup.log_receiver.message.count', { - description: 'Messages coming to the log receiver', - unit: '1', - }) - const socket = createSocket('udp4') - - socket.on('message', (message, rinfo) => { - try { - messageCount.add(1, { - source_ip: rinfo.address, - }) - const logMessage = parseLogMessage(message) - events.emit('gamelog:message', { message: logMessage }) - // eslint-disable-next-line @typescript-eslint/no-unused-vars - } catch (error) { - // empty - } - }) - - socket.on('listening', () => { - const address = socket.address() - logger.info(`log receiver listening at ${address.address}:${address.port}`) - }) - - socket.bind(environment.LOG_RELAY_PORT, '0.0.0.0') - - app.addHook('onClose', (_, done) => { - socket.close(done) - }) + const useWorkerThread = await configuration.get('log_receiver.use_worker_thread') + + if (useWorkerThread) { + const worker = new Worker(new URL('../log-receiver.worker.js', import.meta.url), { + execArgv: [...process.execArgv], + }) + + worker.on('message', (msg: WorkerMessage) => { + switch (msg.type) { + case 'match:started': + events.emit('match:started', { gameNumber: msg.gameNumber }) + break + case 'match:restarted': + events.emit('match:restarted', { gameNumber: msg.gameNumber }) + break + case 'match:ended': + events.emit('match:ended', { gameNumber: msg.gameNumber }) + break + case 'match:roundWon': + events.emit('match:roundWon', { gameNumber: msg.gameNumber, winner: msg.winner }) + break + case 'match:roundLength': + events.emit('match:roundLength', { gameNumber: msg.gameNumber, lengthMs: msg.lengthMs }) + break + case 'match/logs:uploaded': + events.emit('match/logs:uploaded', { + gameNumber: msg.gameNumber, + logsUrl: msg.logsUrl, + }) + break + case 'match/player:connected': + events.emit('match/player:connected', { + gameNumber: msg.gameNumber, + steamId: msg.steamId, + ipAddress: msg.ipAddress, + }) + break + case 'match/player:joinedTeam': + events.emit('match/player:joinedTeam', { + gameNumber: msg.gameNumber, + steamId: msg.steamId, + team: msg.team, + }) + break + case 'match/player:disconnected': + events.emit('match/player:disconnected', { + gameNumber: msg.gameNumber, + steamId: msg.steamId, + }) + break + case 'match/player:said': + events.emit('match/player:said', { + gameNumber: msg.gameNumber, + steamId: msg.steamId, + message: msg.message, + }) + break + case 'match/score:reported': + events.emit('match/score:reported', { + gameNumber: msg.gameNumber, + teamName: msg.teamName, + score: msg.score, + }) + break + case 'match/score:final': + events.emit('match/score:final', { + gameNumber: msg.gameNumber, + team: msg.team, + score: msg.score, + }) + break + case 'match/demo:uploaded': + events.emit('match/demo:uploaded', { + gameNumber: msg.gameNumber, + demoUrl: msg.demoUrl, + }) + break + } + }) + + worker.on('error', error => { + logger.error({ error }, 'log receiver worker error') + }) + + app.addHook('onClose', async () => { + await withTimeout(async () => { + await new Promise((resolve, reject) => { + worker.postMessage({ type: 'shutdown' } satisfies ControlMessage) + worker.once('exit', resolve) + worker.once('error', reject) + }) + }, secondsToMilliseconds(60)) + }) + } else { + const messageCount = meter.createCounter('tf2pickup.log_receiver.message.count', { + description: 'Messages coming to the log receiver', + unit: '1', + }) + const socket = createSocket('udp4') + + socket.on('message', (message, rinfo) => { + try { + messageCount.add(1, { + source_ip: rinfo.address, + }) + const logMessage = parseLogMessage(message) + events.emit('gamelog:message', { message: logMessage }) + // eslint-disable-next-line @typescript-eslint/no-unused-vars + } catch (error) { + // empty + } + }) + + socket.on('listening', () => { + const address = socket.address() + logger.info(`log receiver listening at ${address.address}:${address.port}`) + }) + + socket.bind(environment.LOG_RELAY_PORT, '0.0.0.0') + + app.addHook('onClose', (_, done) => { + socket.close(done) + }) + } }, { name: 'log receiver', diff --git a/src/log-receiver/worker-message.ts b/src/log-receiver/worker-message.ts new file mode 100644 index 000000000..26b6b38eb --- /dev/null +++ b/src/log-receiver/worker-message.ts @@ -0,0 +1,27 @@ +import type { GameNumber } from '../database/models/game.model' +import type { SteamId64 } from '../shared/types/steam-id-64' +import type { Tf2Team } from '../shared/types/tf2-team' + +export type WorkerMessage = + | { type: 'match:started'; gameNumber: GameNumber } + | { type: 'match:restarted'; gameNumber: GameNumber } + | { type: 'match:ended'; gameNumber: GameNumber } + | { type: 'match:roundWon'; gameNumber: GameNumber; winner: Tf2Team } + | { type: 'match:roundLength'; gameNumber: GameNumber; lengthMs: number } + | { type: 'match/logs:uploaded'; gameNumber: GameNumber; logsUrl: string } + | { + type: 'match/player:connected' + gameNumber: GameNumber + steamId: SteamId64 + ipAddress: string + } + | { type: 'match/player:joinedTeam'; gameNumber: GameNumber; steamId: SteamId64; team: Tf2Team } + | { type: 'match/player:disconnected'; gameNumber: GameNumber; steamId: SteamId64 } + | { type: 'match/player:said'; gameNumber: GameNumber; steamId: SteamId64; message: string } + | { type: 'match/score:reported'; gameNumber: GameNumber; teamName: Tf2Team; score: number } + | { type: 'match/score:final'; gameNumber: GameNumber; team: Tf2Team; score: number } + | { type: 'match/demo:uploaded'; gameNumber: GameNumber; demoUrl: string } + +export interface ControlMessage { + type: 'shutdown' +}