diff --git a/apps/stage-tamagotchi/src/main/services/airi/channel-server/index.ts b/apps/stage-tamagotchi/src/main/services/airi/channel-server/index.ts index b9cc7b2bb..9a433509c 100644 --- a/apps/stage-tamagotchi/src/main/services/airi/channel-server/index.ts +++ b/apps/stage-tamagotchi/src/main/services/airi/channel-server/index.ts @@ -1,6 +1,6 @@ import { env } from 'node:process' -import { useLogg } from '@guiiai/logg' +import { LogLevelString, useLogg } from '@guiiai/logg' import { onAppBeforeQuit } from '../../../libs/bootkit/lifecycle' @@ -13,12 +13,18 @@ export async function setupServerChannel() { const serverRuntime = await import('@proj-airi/server-runtime') const { serve } = await import('h3') const { plugin: ws } = await import('crossws/server') + const app = serverRuntime.setupApp({ + logger: { + app: { level: LogLevelString.Debug }, + websocket: { level: LogLevelString.Debug }, + }, + }) try { - const serverInstance = serve(serverRuntime.app, { + const serverInstance = serve(app, { // TODO: fix types // @ts-expect-error - the .crossws property wasn't extended in types - plugins: [ws({ resolve: async req => (await serverRuntime.app.fetch(req)).crossws })], + plugins: [ws({ resolve: async req => (await app.fetch(req)).crossws })], port: env.PORT ? Number(env.PORT) : 6121, hostname: env.SERVER_RUNTIME_HOSTNAME || 'localhost', reusePort: true, diff --git a/packages/server-runtime/src/config/config.ts b/packages/server-runtime/src/config/config.ts new file mode 100644 index 000000000..53976ec94 --- /dev/null +++ b/packages/server-runtime/src/config/config.ts @@ -0,0 +1,11 @@ +import { fromEnv } from './env' + +export function optionOrEnv(option: T | undefined, envKey: string, envDefault: T, options?: { validator: (value: string) => value is T }): T +export function optionOrEnv(option: T | undefined, envKey: string, envDefault?: undefined, options?: { validator: (value: string) => value is T }): undefined | T +export function optionOrEnv(option: T | undefined, envKey: string, envDefault?: T, options?: { validator: (value: string) => value is T }): T | undefined { + if (option !== undefined) { + return option + } + + return fromEnv(envKey, envDefault, options) +} diff --git a/packages/server-runtime/src/config/env.ts b/packages/server-runtime/src/config/env.ts new file mode 100644 index 000000000..32e90f37b --- /dev/null +++ b/packages/server-runtime/src/config/env.ts @@ -0,0 +1,19 @@ +import { env } from 'node:process' + +export function fromEnv(envKey: string, envDefault?: T, options?: { validator?: (value: string) => value is T }): T | undefined { + const value = env[envKey] ?? envDefault + if (value === undefined) { + return undefined + } + + if (options?.validator) { + if (options.validator(value)) { + return value + } + else { + return undefined + } + } + + return value as T +} diff --git a/packages/server-runtime/src/config/index.ts b/packages/server-runtime/src/config/index.ts new file mode 100644 index 000000000..26158a345 --- /dev/null +++ b/packages/server-runtime/src/config/index.ts @@ -0,0 +1,2 @@ +export * from './config' +export * from './env' diff --git a/packages/server-runtime/src/index.ts b/packages/server-runtime/src/index.ts index e044503ce..d362acff5 100644 --- a/packages/server-runtime/src/index.ts +++ b/packages/server-runtime/src/index.ts @@ -2,26 +2,12 @@ import type { WebSocketEvent } from '@proj-airi/server-shared/types' import type { AuthenticatedPeer, Peer } from './types' -import { env } from 'node:process' - -import { availableLogLevelStrings, Format, LogLevel, logLevelStringToLogLevelMap, setGlobalFormat, setGlobalLogLevel, useLogg } from '@guiiai/logg' +import { availableLogLevelStrings, Format, LogLevelString, logLevelStringToLogLevelMap, useLogg } from '@guiiai/logg' import { defineWebSocketHandler, H3 } from 'h3' +import { optionOrEnv } from './config' import { WebSocketReadyState } from './types' -setGlobalFormat(Format.Pretty) -setGlobalLogLevel(LogLevel.Log) - -if (env.LOG_LEVEL) { - const level = env.LOG_LEVEL as typeof availableLogLevelStrings[number] - if (availableLogLevelStrings.includes(level)) { - setGlobalLogLevel(logLevelStringToLogLevelMap[level]) - } -} - -// cache token once -const AUTH_TOKEN = env.AUTHENTICATION_TOKEN || '' - // pre-stringified responses const RESPONSES = { authenticated: JSON.stringify({ type: 'module:authenticated', data: { authenticated: true } }), @@ -33,9 +19,24 @@ function send(peer: Peer, event: WebSocketEvent> | strin peer.send(typeof event === 'string' ? event : JSON.stringify(event)) } -function setupApp(): H3 { - const appLogger = useLogg('App').useGlobalConfig() - const websocketLogger = useLogg('WebSocket').useGlobalConfig() +export function setupApp(options?: { + auth?: { + token: string + } + logger?: { + app?: { level?: LogLevelString, format?: Format } + websocket?: { level?: LogLevelString, format?: Format } + } +}): H3 { + const authToken = optionOrEnv(options?.auth?.token, 'AUTHENTICATION_TOKEN', '') + + const appLogLevel = optionOrEnv(options?.logger?.app?.level, 'LOG_LEVEL', LogLevelString.Log, { validator: (value): value is LogLevelString => availableLogLevelStrings.includes(value as LogLevelString) }) + const appLogFormat = optionOrEnv(options?.logger?.app?.format, 'LOG_FORMAT', Format.Pretty, { validator: (value): value is Format => Object.values(Format).includes(value as Format) }) + const websocketLogLevel = options?.logger?.websocket?.level || appLogLevel || LogLevelString.Log + const websocketLogFormat = options?.logger?.websocket?.format || appLogFormat || Format.Pretty + + const appLogger = useLogg('@proj-airi/server-runtime').withLogLevel(logLevelStringToLogLevelMap[appLogLevel]).withFormat(appLogFormat) + const logger = useLogg('@proj-airi/server-runtime:websocket').withLogLevel(logLevelStringToLogLevelMap[websocketLogLevel]).withFormat(websocketLogFormat) const app = new H3({ onError: error => appLogger.withError(error).error('an error occurred'), @@ -48,20 +49,24 @@ function setupApp(): H3 { if (!peersByModule.has(name)) { peersByModule.set(name, new Map()) } + const group = peersByModule.get(name)! if (group.has(index)) { // log instead of silent overwrite - websocketLogger.withFields({ name, index }).debug('peer replaced for module') + logger.withFields({ name, index }).debug('peer replaced for module') } + group.set(index, p) } function unregisterModulePeer(p: AuthenticatedPeer) { if (!p.name) return + const group = peersByModule.get(p.name) if (group) { group.delete(p.index) + if (group.size === 0) { peersByModule.delete(p.name) } @@ -70,7 +75,7 @@ function setupApp(): H3 { app.get('/ws', defineWebSocketHandler({ open: (peer) => { - if (AUTH_TOKEN) { + if (authToken) { peers.set(peer.id, { peer, authenticated: false, name: '' }) } else { @@ -78,56 +83,77 @@ function setupApp(): H3 { peers.set(peer.id, { peer, authenticated: true, name: '' }) } - websocketLogger.withFields({ peer: peer.id, activePeers: peers.size }).log('connected') + logger.withFields({ peer: peer.id, activePeers: peers.size }).log('connected') }, message: (peer, message) => { + const authenticatedPeer = peers.get(peer.id) let event: WebSocketEvent + try { event = message.json() as WebSocketEvent } catch (err) { const errorMessage = err instanceof Error ? err.message : String(err) send(peer, { type: 'error', data: { message: `invalid JSON, error: ${errorMessage}` } }) + return } + logger.withFields({ + peer: peer.id, + peerAuthenticated: authenticatedPeer?.authenticated, + peerModule: authenticatedPeer?.name, + peerModuleIndex: authenticatedPeer?.index, + }).debug('received event') + switch (event.type) { case 'module:authenticate': { - if (AUTH_TOKEN && event.data.token !== AUTH_TOKEN) { - websocketLogger.withFields({ peer: peer.id }).debug('authentication failed') + if (authToken && event.data.token !== authToken) { + logger.withFields({ peer: peer.id, peerRemote: peer.remoteAddress, peerRequest: peer.request.url }).log('authentication failed') send(peer, { type: 'error', data: { message: 'invalid token' } }) + return } peer.send(RESPONSES.authenticated) const p = peers.get(peer.id) if (p) { - Object.assign(p, { authenticated: true }) + p.authenticated = true } + return } + case 'module:announce': { const p = peers.get(peer.id) - if (p) { - unregisterModulePeer(p) - const { name, index } = event.data as { name: string, index?: number } - if (!name || typeof name !== 'string') { - send(peer, { type: 'error', data: { message: 'the field \'name\' must be a non-empty string for event \'module:announce\'' } }) - return - } - if (typeof index !== 'undefined') { - if (!Number.isInteger(index) || index < 0) { - send(peer, { type: 'error', data: { message: 'the field \'index\' must be a non-negative integer for event \'module:announce\'' } }) - return - } - } - if (AUTH_TOKEN && !p.authenticated) { - send(peer, { type: 'error', data: { message: 'must authenticate before announcing' } }) - return - } - Object.assign(p, { name, index }) - registerModulePeer(p, p.name, p.index) + if (!p) { + return } + + unregisterModulePeer(p) + + // verify + const { name, index } = event.data as { name: string, index?: number } + if (!name || typeof name !== 'string') { + send(peer, { type: 'error', data: { message: 'the field \'name\' must be a non-empty string for event \'module:announce\'' } }) + return + } + if (typeof index !== 'undefined') { + if (!Number.isInteger(index) || index < 0) { + send(peer, { type: 'error', data: { message: 'the field \'index\' must be a non-negative integer for event \'module:announce\'' } }) + return + } + } + if (authToken && !p.authenticated) { + send(peer, { type: 'error', data: { message: 'must authenticate before announcing' } }) + return + } + + p.name = name + p.index = index + + registerModulePeer(p, name, index) + return } @@ -136,6 +162,7 @@ function setupApp(): H3 { if (moduleName === '') { send(peer, { type: 'error', data: { message: 'the field \'moduleName\' can\'t be empty for event \'ui:configure\'' } }) + return } if (typeof moduleIndex !== 'undefined') { @@ -144,6 +171,7 @@ function setupApp(): H3 { type: 'error', data: { message: 'the field \'moduleIndex\' must be a non-negative integer for event \'ui:configure\'' }, }) + return } } @@ -155,6 +183,7 @@ function setupApp(): H3 { else { send(peer, { type: 'error', data: { message: 'module not found, it hasn\'t announced itself or the name is incorrect' } }) } + return } } @@ -162,39 +191,42 @@ function setupApp(): H3 { // default case const p = peers.get(peer.id) if (!p?.authenticated) { - websocketLogger.withFields({ peer: peer.id }).debug('not authenticated') + logger.withFields({ peer: peer.id, peerRemote: peer.remoteAddress, peerRequest: peer.request.url }).debug('not authenticated') peer.send(RESPONSES.notAuthenticated) + return } const payload = JSON.stringify(event) - websocketLogger.withFields({ peer: peer.id, event: payload }).debug('broadcasting event to peers') + logger.withFields({ peer: peer.id, event }).debug('broadcasting event to peers') for (const [id, other] of peers.entries()) { if (id === peer.id) { - websocketLogger.withFields({ peer: peer.id, event: payload }).debug('not sending event to self') + logger.withFields({ peer: peer.id, event }).debug('not sending event to self') continue } + if (other.peer.readyState === WebSocketReadyState.OPEN) { - websocketLogger.withFields({ fromPeer: peer.id, toPeer: other.peer.id, event: payload }).debug('sending event to peer') + logger.withFields({ fromPeer: peer.id, toPeer: other.peer.id, event }).debug('sending event to peer') other.peer.send(payload) } else { - websocketLogger.withFields({ peer: other.peer.id }).debug('removing closed peer') + logger.withFields({ peer: other.peer.id }).debug('removing closed peer') peers.delete(id) + unregisterModulePeer(other) } } }, error: (peer, error) => { - websocketLogger.withFields({ peer: peer.id }).withError(error).error('an error occurred') + logger.withFields({ peer: peer.id }).withError(error).error('an error occurred') }, close: (peer, details) => { const p = peers.get(peer.id) if (p) unregisterModulePeer(p) - websocketLogger.withFields({ peer: peer.id, details, activePeers: peers.size }).log('closed') + logger.withFields({ peer: peer.id, peerRemote: peer.remoteAddress, details, activePeers: peers.size }).log('closed') peers.delete(peer.id) }, })) diff --git a/packages/server-shared/src/types/websocket/events.ts b/packages/server-shared/src/types/websocket/events.ts index 16e15b0f4..6e384fd95 100644 --- a/packages/server-shared/src/types/websocket/events.ts +++ b/packages/server-shared/src/types/websocket/events.ts @@ -33,7 +33,7 @@ export interface ContextMessage