From 55653762210080e9b9dbd1a5a1a133eba8057f67 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micha=C5=82=20Garapich?= Date: Mon, 18 May 2026 22:28:10 +0200 Subject: [PATCH 1/4] perf: move log receiver processing to a worker thread UDP parsing, regex game-event matching, and MongoDB log writes now run on a dedicated worker thread (node:worker_threads) with its own MongoClient and a second OTEL SDK instance for metrics. The main event loop only calls events.emit() when the worker posts a matched game event via postMessage, eliminating per-log-line synchronous work from the main thread. Co-Authored-By: Claude Sonnet 4.6 --- src/events.ts | 5 - src/games/log-message-queue.ts | 30 +++ src/games/plugins/game-log-collector.ts | 27 -- src/games/plugins/match-event-listener.ts | 242 ------------------ src/log-receiver/log-receiver.worker.ts | 290 ++++++++++++++++++++++ src/log-receiver/plugins/index.ts | 104 ++++++-- src/log-receiver/worker-message.ts | 20 ++ 7 files changed, 417 insertions(+), 301 deletions(-) create mode 100644 src/games/log-message-queue.ts delete mode 100644 src/games/plugins/game-log-collector.ts delete mode 100644 src/games/plugins/match-event-listener.ts create mode 100644 src/log-receiver/log-receiver.worker.ts create mode 100644 src/log-receiver/worker-message.ts diff --git a/src/events.ts b/src/events.ts index 234da5834..ecf0c73f4 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' @@ -72,10 +71,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..033199ce7 --- /dev/null +++ b/src/games/log-message-queue.ts @@ -0,0 +1,30 @@ +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 + } + } + + clear(key: string): void { + this.queues.delete(key) + } +} + +export const logMessageQueue = new LogMessageQueue() diff --git a/src/games/plugins/game-log-collector.ts b/src/games/plugins/game-log-collector.ts deleted file mode 100644 index 40c900bc8..000000000 --- a/src/games/plugins/game-log-collector.ts +++ /dev/null @@ -1,27 +0,0 @@ -import fp from 'fastify-plugin' -import { events } from '../../events' -import { findOne } from '../find-one' -import type { GameNumber } from '../../database/models/game.model' -import { gameLogSink } from '../game-log-sink' - -export default fp( - // eslint-disable-next-line @typescript-eslint/require-await - async () => { - events.on('gamelog:message', ({ message }) => { - gameLogSink.push(message) - }) - events.on('match:restarted', async ({ gameNumber }) => { - await pruneLogs(gameNumber) - }) - }, - { name: 'game log collector', encapsulate: true }, -) - -async function pruneLogs(number: GameNumber) { - const game = await findOne({ number }, ['logSecret']) - if (!game.logSecret) { - return - } - - await gameLogSink.clear(game.logSecret) -} diff --git a/src/games/plugins/match-event-listener.ts b/src/games/plugins/match-event-listener.ts deleted file mode 100644 index 6a7455921..000000000 --- a/src/games/plugins/match-event-listener.ts +++ /dev/null @@ -1,242 +0,0 @@ -import fp from 'fastify-plugin' -import type { GameNumber } from '../../database/models/game.model' -import { events } from '../../events' -import type { SteamId64 } from '../../shared/types/steam-id-64' -import type { Tf2Team } from '../../shared/types/tf2-team' -import SteamID from 'steamid' -import { collections } from '../../database/collections' -import { logger } from '../../logger' -import { meter } from '../../otel' -import { ValueType } from '@opentelemetry/api' - -interface GameEvent { - /* name of the game event */ - name: string - - /* the event is triggered if a log line matches this regex */ - regex: RegExp - - /* handle the event being triggered */ - handle: (number: GameNumber, matches: RegExpMatchArray) => void -} - -// converts 'Red' and 'Blue' to valid team names -const fixTeamName = (teamName: string): Tf2Team => teamName.toLowerCase().substring(0, 3) as Tf2Team - -const gameEvents: GameEvent[] = [ - { - // TODO rename to "round start" - name: 'match started', - regex: /^[\d/\s-:]+World triggered "Round_Start"$/, - handle: gameNumber => events.emit('match:started', { gameNumber }), - }, - { - name: 'game restarted', - regex: /^\d{2}\/\d{2}\/\d{4}\s-\s\d{2}:\d{2}:\d{2}:\srcon from ".+": command "exec etf2l_.+"$/, - handle: gameNumber => events.emit('match:restarted', { gameNumber }), - }, - { - name: 'round win', - // https://regex101.com/r/41LfKS/2 - regex: - /^\d{2}\/\d{2}\/\d{4}\s-\s\d{2}:\d{2}:\d{2}:\sWorld triggered "Round_Win" \(winner "(.+)"\)$/, - handle: (gameNumber, matches) => { - if (matches[1]) { - const winner = fixTeamName(matches[1]) - events.emit('match:roundWon', { gameNumber, winner }) - } - }, - }, - { - name: 'round length', - // https://regex101.com/r/mvOYMz/3 - regex: - /^\d{2}\/\d{2}\/\d{4}\s-\s\d{2}:\d{2}:\d{2}:\sWorld triggered "Round_Length" \(seconds "([\d.]+)"\)$/, - handle: (gameNumber, matches) => { - if (matches[1]) { - const seconds = parseFloat(matches[1]) - events.emit('match:roundLength', { gameNumber, lengthMs: seconds * 1000 }) - } - }, - }, - { - // TODO rename to "game over" - name: 'match ended', - regex: /^[\d/\s-:]+World triggered "Game_Over" reason ".*"$/, - handle: gameNumber => events.emit('match:ended', { gameNumber }), - }, - { - name: 'logs uploaded', - regex: /^[\d/\s-:]+\[TFTrue\].+\shttp:\/\/logs\.tf\/(\d+)\..*$/, - handle: (gameNumber, matches) => { - const logsUrl = `http://logs.tf/${matches[1]}` - events.emit('match/logs:uploaded', { gameNumber, logsUrl }) - }, - }, - { - name: 'player connected', - // https://regex101.com/r/uyPW8m/5 - 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})"$/, - handle: (gameNumber, matches) => { - if (!matches[5] || !matches[6]) { - return - } - const steamId = new SteamID(matches[5]) - if (steamId.isValid()) { - events.emit('match/player:connected', { - gameNumber, - steamId: steamId.getSteamID64() as SteamId64, - ipAddress: matches[6], - }) - } - }, - }, - { - name: 'player joined team', - // https://regex101.com/r/yzX9zG/1 - regex: - /^(\d{2}\/\d{2}\/\d{4})\s-\s(\d{2}:\d{2}:\d{2}):\s"(.+)<(\d+)><(\[.[^\]]+\])><(.+)>"\sjoined\steam\s"(.+)"/, - handle: (gameNumber, matches) => { - if (!matches[5] || !matches[7]) { - return - } - const steamId = new SteamID(matches[5]) - if (steamId.isValid()) { - events.emit('match/player:joinedTeam', { - gameNumber, - steamId: steamId.getSteamID64() as SteamId64, - team: fixTeamName(matches[7]), - }) - } - }, - }, - { - name: 'player disconnected', - // https://regex101.com/r/x4AMTG/1 - regex: - /^(\d{2}\/\d{2}\/\d{4})\s-\s(\d{2}:\d{2}:\d{2}):\s"(.+)<(\d+)><(\[.[^\]]+\])><(.[^>]+)>"\sdisconnected\s\(reason\s"(.[^"]+)"\)$/, - handle: (gameNumber, matches) => { - if (!matches[5] || !matches[6]) { - return - } - const steamId = new SteamID(matches[5]) - if (steamId.isValid()) { - events.emit('match/player:disconnected', { - gameNumber, - steamId: steamId.getSteamID64() as SteamId64, - }) - } - }, - }, - { - name: 'score reported', - // https://regex101.com/r/ZD6eLb/1 - regex: /^[\d/\s\-:]+Team "(.[^"]+)" current score "(\d)" with "(\d)" players$/, - handle: (gameNumber, matches) => { - const [, teamName, score] = matches - if (teamName && score) { - events.emit('match/score:reported', { - gameNumber, - teamName: fixTeamName(teamName), - score: Number(score), - }) - } - }, - }, - { - name: 'final score reported', - // https://regex101.com/r/RAUdTe/1 - regex: /^[\d/\s\-:]+Team "(.[^"]+)" final score "(\d)" with "(\d)" players$/, - handle: (gameNumber, matches) => { - const [, teamName, score] = matches - if (teamName && score) { - events.emit('match/score:final', { - gameNumber, - team: fixTeamName(teamName), - score: Number(score), - }) - } - }, - }, - { - name: 'demo uploaded', - // https://regex101.com/r/JLGRYa/2 - regex: /^[\d/\s-:]+\[demos\.tf\]:\sSTV\savailable\sat:\s(.+)$/, - handle: (gameNumber, matches) => { - const demoUrl = matches[1] - if (demoUrl) { - events.emit('match/demo:uploaded', { gameNumber, demoUrl }) - } - }, - }, - { - name: 'player said', - // https://regex101.com/r/zpFkkA/1 - regex: - /^(\d{2}\/\d{2}\/\d{4})\s-\s(\d{2}:\d{2}:\d{2}):\s"(.+)<(\d+)><(\[.[^\]]+\])><(.[^>]+)>"\ssay\s"(.+)"$/, - handle: (gameNumber, matches) => { - if (!matches[5] || !matches[7]) { - return - } - const steamId = new SteamID(matches[5]) - if (steamId.isValid()) { - const message = matches[7] - events.emit('match/player:said', { - gameNumber, - steamId: steamId.getSteamID64() as SteamId64, - message, - }) - } - }, - }, -] - -const eventCounter = meter.createCounter('tf2pickup.games.events.count', { - description: 'Game events that come from the gameserver', - unit: '1', - valueType: ValueType.INT, -}) - -type EventHandledResult = - | { - handled: true - gameNumber: GameNumber - } - | { - handled: false - } - -async function testForGameEvent(message: string, logSecret: string): Promise { - for (const gameEvent of gameEvents) { - const matches = message.match(gameEvent.regex) - if (matches) { - const game = await collections.games.findOne({ logSecret }, { projection: { number: 1 } }) - if (game === null) { - logger.error({ message }, `error handling game event: no such game`) - return { handled: false } - } - gameEvent.handle(game.number, matches) - return { handled: true, gameNumber: game.number } - } - } - - return { handled: false } -} - -export default fp( - // eslint-disable-next-line @typescript-eslint/require-await - async () => { - events.on('gamelog:message', async ({ message }) => { - const result = await testForGameEvent(message.payload, message.password) - eventCounter.add(1, { - 'tf2pickup.games.event.handled': result.handled, - ...(result.handled ? { 'tf2pickup.game.number': result.gameNumber } : {}), - }) - }) - }, - { - name: 'match event listener', - encapsulate: true, - }, -) diff --git a/src/log-receiver/log-receiver.worker.ts b/src/log-receiver/log-receiver.worker.ts new file mode 100644 index 000000000..9f951dba6 --- /dev/null +++ b/src/log-receiver/log-receiver.worker.ts @@ -0,0 +1,290 @@ +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, ControlMessage } 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 (msg: ControlMessage) => { + if (msg.type === 'shutdown') { + socket.close() + await mongoClient.close() + await sdk.shutdown() + process.exit(0) + } +}) diff --git a/src/log-receiver/plugins/index.ts b/src/log-receiver/plugins/index.ts index 02792f5e7..14c919090 100644 --- a/src/log-receiver/plugins/index.ts +++ b/src/log-receiver/plugins/index.ts @@ -1,42 +1,92 @@ import fp from 'fastify-plugin' -import { createSocket } from 'node:dgram' -import { parseLogMessage } from '../parse-log-message' -import { logger } from '../../logger' -import { environment } from '../../environment' +import { Worker } from 'node:worker_threads' import { events } from '../../events' -import { meter } from '../../otel' +import { logger } from '../../logger' +import type { WorkerMessage, ControlMessage } from '../worker-message' 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 worker = new Worker(new URL('../log-receiver.worker.js', import.meta.url), { + execArgv: [...process.execArgv], }) - 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 + 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 } }) - socket.on('listening', () => { - const address = socket.address() - logger.info(`log receiver listening at ${address.address}:${address.port}`) + worker.on('error', error => { + logger.error({ error }, 'log receiver worker error') }) - socket.bind(environment.LOG_RELAY_PORT, '0.0.0.0') - - app.addHook('onClose', (_, done) => { - socket.close(done) + app.addHook('onClose', async () => { + await new Promise((resolve, reject) => { + worker.postMessage({ type: 'shutdown' } satisfies ControlMessage) + worker.once('exit', resolve) + worker.once('error', reject) + }) }) }, { diff --git a/src/log-receiver/worker-message.ts b/src/log-receiver/worker-message.ts new file mode 100644 index 000000000..33727c347 --- /dev/null +++ b/src/log-receiver/worker-message.ts @@ -0,0 +1,20 @@ +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 type ControlMessage = { type: 'shutdown' } From 1f01972bab518977ff1e66c67b075b9a982bc5ab Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micha=C5=82=20Garapich?= Date: Thu, 21 May 2026 15:54:29 +0200 Subject: [PATCH 2/4] format --- src/log-receiver/log-receiver.worker.ts | 21 +++++++++++++++++---- src/log-receiver/worker-message.ts | 7 ++++++- 2 files changed, 23 insertions(+), 5 deletions(-) diff --git a/src/log-receiver/log-receiver.worker.ts b/src/log-receiver/log-receiver.worker.ts index 9f951dba6..b340f3ded 100644 --- a/src/log-receiver/log-receiver.worker.ts +++ b/src/log-receiver/log-receiver.worker.ts @@ -78,7 +78,8 @@ const gameEvents: WorkerGameEvent[] = [ { 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 "(.+)"\)$/, + 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]) } @@ -166,7 +167,12 @@ const gameEvents: WorkerGameEvent[] = [ buildMessage: (gameNumber, matches) => { const [, teamName, score] = matches if (!teamName || !score) return null - return { type: 'match/score:reported', gameNumber, teamName: fixTeamName(teamName), score: Number(score) } + return { + type: 'match/score:reported', + gameNumber, + teamName: fixTeamName(teamName), + score: Number(score), + } }, }, { @@ -176,7 +182,12 @@ const gameEvents: WorkerGameEvent[] = [ buildMessage: (gameNumber, matches) => { const [, teamName, score] = matches if (!teamName || !score) return null - return { type: 'match/score:final', gameNumber, team: fixTeamName(teamName), score: Number(score) } + return { + type: 'match/score:final', + gameNumber, + team: fixTeamName(teamName), + score: Number(score), + } }, }, { @@ -250,7 +261,9 @@ socket.on('message', (message, rinfo) => { messageCount.add(1, { source_ip: rinfo.address }) const logMessage = parseLogMessage(message) - logQueue.enqueue(logMessage.password, () => pushLogMessage(logMessage.payload, logMessage.password)) + logQueue.enqueue(logMessage.password, () => + pushLogMessage(logMessage.payload, logMessage.password), + ) matchGameEvent(logMessage.payload, logMessage.password).then( workerMessage => { diff --git a/src/log-receiver/worker-message.ts b/src/log-receiver/worker-message.ts index 33727c347..8c85229e5 100644 --- a/src/log-receiver/worker-message.ts +++ b/src/log-receiver/worker-message.ts @@ -9,7 +9,12 @@ export type WorkerMessage = | { 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: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 } From 37324e79f7ee083f4fc4db6833266abed6bcc747 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micha=C5=82=20Garapich?= Date: Thu, 21 May 2026 16:05:02 +0200 Subject: [PATCH 3/4] feat: gate worker thread log receiver behind configuration entry Add `log_receiver.use_worker_thread` boolean config (default: false) so the worker thread path can be tested per-deployment before rolling out globally. Restores the original inline UDP + gamelog:message path as the default. Also fixes pre-existing lint errors in worker-message.ts and log-receiver.worker.ts. Co-Authored-By: Claude Sonnet 4.6 --- .../models/configuration-entry.model.ts | 6 + src/events.ts | 5 + src/games/plugins/game-log-collector.ts | 27 ++ src/games/plugins/match-event-listener.ts | 242 ++++++++++++++++++ src/log-receiver/log-receiver.worker.ts | 14 +- src/log-receiver/plugins/index.ts | 199 ++++++++------ src/log-receiver/worker-message.ts | 4 +- 7 files changed, 411 insertions(+), 86 deletions(-) create mode 100644 src/games/plugins/game-log-collector.ts create mode 100644 src/games/plugins/match-event-listener.ts diff --git a/src/database/models/configuration-entry.model.ts b/src/database/models/configuration-entry.model.ts index 634955881..cfcc09f0c 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('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 ecf0c73f4..54e523edf 100644 --- a/src/events.ts +++ b/src/events.ts @@ -16,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 } diff --git a/src/games/plugins/game-log-collector.ts b/src/games/plugins/game-log-collector.ts new file mode 100644 index 000000000..40c900bc8 --- /dev/null +++ b/src/games/plugins/game-log-collector.ts @@ -0,0 +1,27 @@ +import fp from 'fastify-plugin' +import { events } from '../../events' +import { findOne } from '../find-one' +import type { GameNumber } from '../../database/models/game.model' +import { gameLogSink } from '../game-log-sink' + +export default fp( + // eslint-disable-next-line @typescript-eslint/require-await + async () => { + events.on('gamelog:message', ({ message }) => { + gameLogSink.push(message) + }) + events.on('match:restarted', async ({ gameNumber }) => { + await pruneLogs(gameNumber) + }) + }, + { name: 'game log collector', encapsulate: true }, +) + +async function pruneLogs(number: GameNumber) { + const game = await findOne({ number }, ['logSecret']) + if (!game.logSecret) { + return + } + + await gameLogSink.clear(game.logSecret) +} diff --git a/src/games/plugins/match-event-listener.ts b/src/games/plugins/match-event-listener.ts new file mode 100644 index 000000000..6a7455921 --- /dev/null +++ b/src/games/plugins/match-event-listener.ts @@ -0,0 +1,242 @@ +import fp from 'fastify-plugin' +import type { GameNumber } from '../../database/models/game.model' +import { events } from '../../events' +import type { SteamId64 } from '../../shared/types/steam-id-64' +import type { Tf2Team } from '../../shared/types/tf2-team' +import SteamID from 'steamid' +import { collections } from '../../database/collections' +import { logger } from '../../logger' +import { meter } from '../../otel' +import { ValueType } from '@opentelemetry/api' + +interface GameEvent { + /* name of the game event */ + name: string + + /* the event is triggered if a log line matches this regex */ + regex: RegExp + + /* handle the event being triggered */ + handle: (number: GameNumber, matches: RegExpMatchArray) => void +} + +// converts 'Red' and 'Blue' to valid team names +const fixTeamName = (teamName: string): Tf2Team => teamName.toLowerCase().substring(0, 3) as Tf2Team + +const gameEvents: GameEvent[] = [ + { + // TODO rename to "round start" + name: 'match started', + regex: /^[\d/\s-:]+World triggered "Round_Start"$/, + handle: gameNumber => events.emit('match:started', { gameNumber }), + }, + { + name: 'game restarted', + regex: /^\d{2}\/\d{2}\/\d{4}\s-\s\d{2}:\d{2}:\d{2}:\srcon from ".+": command "exec etf2l_.+"$/, + handle: gameNumber => events.emit('match:restarted', { gameNumber }), + }, + { + name: 'round win', + // https://regex101.com/r/41LfKS/2 + regex: + /^\d{2}\/\d{2}\/\d{4}\s-\s\d{2}:\d{2}:\d{2}:\sWorld triggered "Round_Win" \(winner "(.+)"\)$/, + handle: (gameNumber, matches) => { + if (matches[1]) { + const winner = fixTeamName(matches[1]) + events.emit('match:roundWon', { gameNumber, winner }) + } + }, + }, + { + name: 'round length', + // https://regex101.com/r/mvOYMz/3 + regex: + /^\d{2}\/\d{2}\/\d{4}\s-\s\d{2}:\d{2}:\d{2}:\sWorld triggered "Round_Length" \(seconds "([\d.]+)"\)$/, + handle: (gameNumber, matches) => { + if (matches[1]) { + const seconds = parseFloat(matches[1]) + events.emit('match:roundLength', { gameNumber, lengthMs: seconds * 1000 }) + } + }, + }, + { + // TODO rename to "game over" + name: 'match ended', + regex: /^[\d/\s-:]+World triggered "Game_Over" reason ".*"$/, + handle: gameNumber => events.emit('match:ended', { gameNumber }), + }, + { + name: 'logs uploaded', + regex: /^[\d/\s-:]+\[TFTrue\].+\shttp:\/\/logs\.tf\/(\d+)\..*$/, + handle: (gameNumber, matches) => { + const logsUrl = `http://logs.tf/${matches[1]}` + events.emit('match/logs:uploaded', { gameNumber, logsUrl }) + }, + }, + { + name: 'player connected', + // https://regex101.com/r/uyPW8m/5 + 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})"$/, + handle: (gameNumber, matches) => { + if (!matches[5] || !matches[6]) { + return + } + const steamId = new SteamID(matches[5]) + if (steamId.isValid()) { + events.emit('match/player:connected', { + gameNumber, + steamId: steamId.getSteamID64() as SteamId64, + ipAddress: matches[6], + }) + } + }, + }, + { + name: 'player joined team', + // https://regex101.com/r/yzX9zG/1 + regex: + /^(\d{2}\/\d{2}\/\d{4})\s-\s(\d{2}:\d{2}:\d{2}):\s"(.+)<(\d+)><(\[.[^\]]+\])><(.+)>"\sjoined\steam\s"(.+)"/, + handle: (gameNumber, matches) => { + if (!matches[5] || !matches[7]) { + return + } + const steamId = new SteamID(matches[5]) + if (steamId.isValid()) { + events.emit('match/player:joinedTeam', { + gameNumber, + steamId: steamId.getSteamID64() as SteamId64, + team: fixTeamName(matches[7]), + }) + } + }, + }, + { + name: 'player disconnected', + // https://regex101.com/r/x4AMTG/1 + regex: + /^(\d{2}\/\d{2}\/\d{4})\s-\s(\d{2}:\d{2}:\d{2}):\s"(.+)<(\d+)><(\[.[^\]]+\])><(.[^>]+)>"\sdisconnected\s\(reason\s"(.[^"]+)"\)$/, + handle: (gameNumber, matches) => { + if (!matches[5] || !matches[6]) { + return + } + const steamId = new SteamID(matches[5]) + if (steamId.isValid()) { + events.emit('match/player:disconnected', { + gameNumber, + steamId: steamId.getSteamID64() as SteamId64, + }) + } + }, + }, + { + name: 'score reported', + // https://regex101.com/r/ZD6eLb/1 + regex: /^[\d/\s\-:]+Team "(.[^"]+)" current score "(\d)" with "(\d)" players$/, + handle: (gameNumber, matches) => { + const [, teamName, score] = matches + if (teamName && score) { + events.emit('match/score:reported', { + gameNumber, + teamName: fixTeamName(teamName), + score: Number(score), + }) + } + }, + }, + { + name: 'final score reported', + // https://regex101.com/r/RAUdTe/1 + regex: /^[\d/\s\-:]+Team "(.[^"]+)" final score "(\d)" with "(\d)" players$/, + handle: (gameNumber, matches) => { + const [, teamName, score] = matches + if (teamName && score) { + events.emit('match/score:final', { + gameNumber, + team: fixTeamName(teamName), + score: Number(score), + }) + } + }, + }, + { + name: 'demo uploaded', + // https://regex101.com/r/JLGRYa/2 + regex: /^[\d/\s-:]+\[demos\.tf\]:\sSTV\savailable\sat:\s(.+)$/, + handle: (gameNumber, matches) => { + const demoUrl = matches[1] + if (demoUrl) { + events.emit('match/demo:uploaded', { gameNumber, demoUrl }) + } + }, + }, + { + name: 'player said', + // https://regex101.com/r/zpFkkA/1 + regex: + /^(\d{2}\/\d{2}\/\d{4})\s-\s(\d{2}:\d{2}:\d{2}):\s"(.+)<(\d+)><(\[.[^\]]+\])><(.[^>]+)>"\ssay\s"(.+)"$/, + handle: (gameNumber, matches) => { + if (!matches[5] || !matches[7]) { + return + } + const steamId = new SteamID(matches[5]) + if (steamId.isValid()) { + const message = matches[7] + events.emit('match/player:said', { + gameNumber, + steamId: steamId.getSteamID64() as SteamId64, + message, + }) + } + }, + }, +] + +const eventCounter = meter.createCounter('tf2pickup.games.events.count', { + description: 'Game events that come from the gameserver', + unit: '1', + valueType: ValueType.INT, +}) + +type EventHandledResult = + | { + handled: true + gameNumber: GameNumber + } + | { + handled: false + } + +async function testForGameEvent(message: string, logSecret: string): Promise { + for (const gameEvent of gameEvents) { + const matches = message.match(gameEvent.regex) + if (matches) { + const game = await collections.games.findOne({ logSecret }, { projection: { number: 1 } }) + if (game === null) { + logger.error({ message }, `error handling game event: no such game`) + return { handled: false } + } + gameEvent.handle(game.number, matches) + return { handled: true, gameNumber: game.number } + } + } + + return { handled: false } +} + +export default fp( + // eslint-disable-next-line @typescript-eslint/require-await + async () => { + events.on('gamelog:message', async ({ message }) => { + const result = await testForGameEvent(message.payload, message.password) + eventCounter.add(1, { + 'tf2pickup.games.event.handled': result.handled, + ...(result.handled ? { 'tf2pickup.game.number': result.gameNumber } : {}), + }) + }) + }, + { + name: 'match event listener', + encapsulate: true, + }, +) diff --git a/src/log-receiver/log-receiver.worker.ts b/src/log-receiver/log-receiver.worker.ts index b340f3ded..8a1ab684d 100644 --- a/src/log-receiver/log-receiver.worker.ts +++ b/src/log-receiver/log-receiver.worker.ts @@ -13,7 +13,7 @@ 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, ControlMessage } from './worker-message' +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' @@ -293,11 +293,9 @@ socket.on('listening', () => { socket.bind(environment.LOG_RELAY_PORT, '0.0.0.0') -parentPort.on('message', async (msg: ControlMessage) => { - if (msg.type === 'shutdown') { - socket.close() - await mongoClient.close() - await sdk.shutdown() - process.exit(0) - } +parentPort.on('message', async () => { + socket.close() + await mongoClient.close() + await sdk.shutdown() + process.exit(0) }) diff --git a/src/log-receiver/plugins/index.ts b/src/log-receiver/plugins/index.ts index 14c919090..a5befb5ec 100644 --- a/src/log-receiver/plugins/index.ts +++ b/src/log-receiver/plugins/index.ts @@ -1,93 +1,138 @@ import fp from 'fastify-plugin' +import { createSocket } from 'node:dgram' import { Worker } from 'node:worker_threads' import { events } from '../../events' import { logger } from '../../logger' +import { environment } from '../../environment' +import { meter } from '../../otel' +import { parseLogMessage } from '../parse-log-message' +import { configuration } from '../../configuration' import type { WorkerMessage, ControlMessage } from '../worker-message' export default fp( async app => { - const worker = new Worker(new URL('../log-receiver.worker.js', import.meta.url), { - execArgv: [...process.execArgv], - }) + const useWorkerThread = await configuration.get('log_receiver.use_worker_thread') - 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, + 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 new Promise((resolve, reject) => { + worker.postMessage({ type: 'shutdown' } satisfies ControlMessage) + worker.once('exit', resolve) + worker.once('error', reject) + }) + }) + } 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, }) - break - case 'match/demo:uploaded': - events.emit('match/demo:uploaded', { gameNumber: msg.gameNumber, demoUrl: msg.demoUrl }) - break - } - }) + 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}`) + }) - worker.on('error', error => { - logger.error({ error }, 'log receiver worker error') - }) + socket.bind(environment.LOG_RELAY_PORT, '0.0.0.0') - app.addHook('onClose', async () => { - await new Promise((resolve, reject) => { - worker.postMessage({ type: 'shutdown' } satisfies ControlMessage) - worker.once('exit', resolve) - worker.once('error', reject) + 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 index 8c85229e5..26b6b38eb 100644 --- a/src/log-receiver/worker-message.ts +++ b/src/log-receiver/worker-message.ts @@ -22,4 +22,6 @@ export type WorkerMessage = | { type: 'match/score:final'; gameNumber: GameNumber; team: Tf2Team; score: number } | { type: 'match/demo:uploaded'; gameNumber: GameNumber; demoUrl: string } -export type ControlMessage = { type: 'shutdown' } +export interface ControlMessage { + type: 'shutdown' +} From 1c702b75f9646f4f35112c3ae0991131d40c6d86 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micha=C5=82=20Garapich?= Date: Wed, 27 May 2026 11:22:56 +0200 Subject: [PATCH 4/4] fix: correct drain ordering and Node.js 20 compatibility in log receiver worker Co-Authored-By: Claude Sonnet 4.6 --- src/database/models/configuration-entry.model.ts | 2 +- src/games/log-message-queue.ts | 4 ++++ src/log-receiver/log-receiver.worker.ts | 12 ++++++++---- src/log-receiver/plugins/index.ts | 14 +++++++++----- 4 files changed, 22 insertions(+), 10 deletions(-) diff --git a/src/database/models/configuration-entry.model.ts b/src/database/models/configuration-entry.model.ts index cfcc09f0c..3847ea245 100644 --- a/src/database/models/configuration-entry.model.ts +++ b/src/database/models/configuration-entry.model.ts @@ -252,7 +252,7 @@ export const configurationSchema = z.discriminatedUnion('key', [ key: z.literal('log_receiver.use_worker_thread'), value: z.boolean().default(false), }) - .describe('Process game server logs in a dedicated worker thread'), + .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/games/log-message-queue.ts b/src/games/log-message-queue.ts index 033199ce7..54a6e422b 100644 --- a/src/games/log-message-queue.ts +++ b/src/games/log-message-queue.ts @@ -22,6 +22,10 @@ export class LogMessageQueue { } } + async drain(): Promise { + await Promise.all([...this.queues.keys()].map(key => this.waitForCompletion(key))) + } + clear(key: string): void { this.queues.delete(key) } diff --git a/src/log-receiver/log-receiver.worker.ts b/src/log-receiver/log-receiver.worker.ts index 8a1ab684d..e42e88447 100644 --- a/src/log-receiver/log-receiver.worker.ts +++ b/src/log-receiver/log-receiver.worker.ts @@ -294,8 +294,12 @@ socket.on('listening', () => { socket.bind(environment.LOG_RELAY_PORT, '0.0.0.0') parentPort.on('message', async () => { - socket.close() - await mongoClient.close() - await sdk.shutdown() - process.exit(0) + 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 a5befb5ec..a998f9af2 100644 --- a/src/log-receiver/plugins/index.ts +++ b/src/log-receiver/plugins/index.ts @@ -8,6 +8,8 @@ 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( async app => { @@ -96,11 +98,13 @@ export default fp( }) app.addHook('onClose', async () => { - await new Promise((resolve, reject) => { - worker.postMessage({ type: 'shutdown' } satisfies ControlMessage) - worker.once('exit', resolve) - worker.once('error', reject) - }) + 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', {