From 8fda3d7b361db160dfcb07f259293be2f49d8aa3 Mon Sep 17 00:00:00 2001 From: nicktrn <55853254+nicktrn@users.noreply.github.com> Date: Wed, 20 Mar 2024 13:17:41 +0000 Subject: [PATCH] structured logs for socket connections --- apps/coordinator/src/index.ts | 6 +- apps/webapp/app/services/logger.server.ts | 11 ++ apps/webapp/app/v3/handleSocketIo.server.ts | 10 +- apps/webapp/app/v3/sharedSocketConnection.ts | 33 ++--- .../cli-v3/src/workers/prod/entry-point.ts | 12 +- packages/core/src/v3/zodMessageHandler.ts | 9 +- packages/core/src/v3/zodNamespace.ts | 119 +++++++++++++++--- packages/core/src/v3/zodSocket.ts | 27 ++-- 8 files changed, 163 insertions(+), 64 deletions(-) diff --git a/apps/coordinator/src/index.ts b/apps/coordinator/src/index.ts index aad90d21d..385f0b2a8 100644 --- a/apps/coordinator/src/index.ts +++ b/apps/coordinator/src/index.ts @@ -344,7 +344,7 @@ class TaskCoordinator { function setSocketDataFromHeader(dataKey: keyof typeof socket.data, headerName: string) { const value = socket.handshake.headers[headerName]; if (!value) { - logger(`missing required header: ${headerName}`); + logger.error("missing required header", { headerName }); throw new Error("missing header"); } 0; @@ -361,12 +361,12 @@ class TaskCoordinator { setSocketDataFromHeader("deploymentId", "x-trigger-deployment-id"); setSocketDataFromHeader("deploymentVersion", "x-trigger-deployment-version"); } catch (error) { - logger(error); + logger.error("setSocketDataFromHeader error", { error }); socket.disconnect(true); return; } - logger("success", socket.data); + logger.debug("success", socket.data); next(); }, diff --git a/apps/webapp/app/services/logger.server.ts b/apps/webapp/app/services/logger.server.ts index c2a413e2d..e38365072 100644 --- a/apps/webapp/app/services/logger.server.ts +++ b/apps/webapp/app/services/logger.server.ts @@ -30,3 +30,14 @@ export const workerLogger = new Logger( return fields ? { ...fields } : {}; } ); + +export const socketLogger = new Logger( + "socket", + (process.env.APP_LOG_LEVEL ?? "debug") as LogLevel, + [], + sensitiveDataReplacer, + () => { + const fields = currentFieldsStore.getStore(); + return fields ? { ...fields } : {}; + } +); diff --git a/apps/webapp/app/v3/handleSocketIo.server.ts b/apps/webapp/app/v3/handleSocketIo.server.ts index 6a2189660..d4601cb68 100644 --- a/apps/webapp/app/v3/handleSocketIo.server.ts +++ b/apps/webapp/app/v3/handleSocketIo.server.ts @@ -132,14 +132,14 @@ function createSharedQueueConsumerNamespace(io: Server) { clientMessages: ClientToSharedQueueMessages, serverMessages: SharedQueueToClientMessages, onConnection: async (socket, handler, sender, logger) => { - const sharedSocketConnection = new SharedSocketConnection( - sharedQueue.namespace, + const sharedSocketConnection = new SharedSocketConnection({ + namespace: sharedQueue.namespace, socket, - logger - ); + logger, + }); sharedSocketConnection.onClose.attach((closeEvent) => { - logger("Socket closed", { closeEvent }); + logger.info("Socket closed", { closeEvent }); }); await sharedSocketConnection.initialize(); diff --git a/apps/webapp/app/v3/sharedSocketConnection.ts b/apps/webapp/app/v3/sharedSocketConnection.ts index 1e7992dfc..ac6d0fb76 100644 --- a/apps/webapp/app/v3/sharedSocketConnection.ts +++ b/apps/webapp/app/v3/sharedSocketConnection.ts @@ -1,5 +1,6 @@ import { MessageCatalogToSocketIoEvents, + StructuredLogger, ZodMessageHandler, ZodMessageSender, clientWebsocketMessages, @@ -11,6 +12,18 @@ import { logger } from "~/services/logger.server"; import { SharedQueueConsumer } from "./marqs/sharedQueueConsumer.server"; import { DisconnectReason, Namespace, Socket } from "socket.io"; +interface SharedSocketConnectionOptions { + namespace: Namespace< + MessageCatalogToSocketIoEvents, + MessageCatalogToSocketIoEvents + >; + socket: Socket< + MessageCatalogToSocketIoEvents, + MessageCatalogToSocketIoEvents + >; + logger?: StructuredLogger; +} + export class SharedSocketConnection { public id: string; public onClose: Evt = new Evt(); @@ -19,17 +32,7 @@ export class SharedSocketConnection { private _sharedConsumer: SharedQueueConsumer; private _messageHandler: ZodMessageHandler; - constructor( - namespace: Namespace< - MessageCatalogToSocketIoEvents, - MessageCatalogToSocketIoEvents - >, - private socket: Socket< - MessageCatalogToSocketIoEvents, - MessageCatalogToSocketIoEvents - >, - logger?: (...args: any[]) => void - ) { + constructor(opts: SharedSocketConnectionOptions) { this.id = randomUUID(); this._sender = new ZodMessageSender({ @@ -38,7 +41,7 @@ export class SharedSocketConnection { return new Promise((resolve, reject) => { try { const { type, ...payload } = message; - namespace.emit(type, payload as any); + opts.namespace.emit(type, payload as any); resolve(); } catch (err) { reject(err); @@ -52,8 +55,8 @@ export class SharedSocketConnection { nextTickInterval: 1000, }); - socket.on("disconnect", this.#handleClose.bind(this)); - socket.on("error", this.#handleError.bind(this)); + opts.socket.on("disconnect", this.#handleClose.bind(this)); + opts.socket.on("error", this.#handleError.bind(this)); this._messageHandler = new ZodMessageHandler({ schema: clientWebsocketMessages, @@ -78,7 +81,7 @@ export class SharedSocketConnection { }, }, }); - this._messageHandler.registerHandlers(this.socket, logger); + this._messageHandler.registerHandlers(opts.socket, opts.logger ?? logger); } async initialize() { diff --git a/packages/cli-v3/src/workers/prod/entry-point.ts b/packages/cli-v3/src/workers/prod/entry-point.ts index 510e61904..f518e6c26 100644 --- a/packages/cli-v3/src/workers/prod/entry-point.ts +++ b/packages/cli-v3/src/workers/prod/entry-point.ts @@ -252,15 +252,15 @@ class ProdWorker { }); if (success) { - logger("indexing done, shutting down.."); + logger.info("indexing done, shutting down.."); process.exit(0); } else { - logger("indexing failure, shutting down.."); + logger.info("indexing failure, shutting down.."); process.exit(1); } } catch (e) { if (e instanceof UncaughtExceptionError) { - logger("uncaught exception", e.originalError.message); + logger.error("uncaught exception", { message: e.originalError.message }); socket.emit("INDEXING_FAILED", { version: "v1", @@ -272,7 +272,7 @@ class ProdWorker { }, }); } else if (e instanceof Error) { - logger("error", e.message); + logger.error("error", { message: e.message }); socket.emit("INDEXING_FAILED", { version: "v1", @@ -284,7 +284,7 @@ class ProdWorker { }, }); } else if (typeof e === "string") { - logger("string error", e); + logger.error("string error", { message: e }); socket.emit("INDEXING_FAILED", { version: "v1", @@ -295,7 +295,7 @@ class ProdWorker { }, }); } else { - logger("unknown error", e); + logger.error("unknown error", { error: e }); socket.emit("INDEXING_FAILED", { version: "v1", diff --git a/packages/core/src/v3/zodMessageHandler.ts b/packages/core/src/v3/zodMessageHandler.ts index b4a64c1b6..3814a9159 100644 --- a/packages/core/src/v3/zodMessageHandler.ts +++ b/packages/core/src/v3/zodMessageHandler.ts @@ -1,4 +1,5 @@ import { z } from "zod"; +import { StructuredLogger } from "./zodNamespace"; export type ZodMessageValueSchema> = | z.ZodFirstPartySchemaTypes @@ -92,17 +93,17 @@ export class ZodMessageHandler }; } - public registerHandlers(emitter: EventEmitterLike, logger?: (...args: any[]) => void) { - const log = logger ?? console.log; + public registerHandlers(emitter: EventEmitterLike, logger?: StructuredLogger) { + const log = logger ?? console; if (!this.#handlers) { - log("No handlers provided"); + log.info("No handlers provided"); return; } for (const eventName of Object.keys(this.#schema)) { emitter.on(eventName, async (message: any, callback?: any): Promise => { - log(`handling ${eventName}`, message); + log.info(`handling ${eventName}`, message); let ack; diff --git a/packages/core/src/v3/zodNamespace.ts b/packages/core/src/v3/zodNamespace.ts index d6dd5976c..1cecf0577 100644 --- a/packages/core/src/v3/zodNamespace.ts +++ b/packages/core/src/v3/zodNamespace.ts @@ -25,6 +25,87 @@ export type ZodNamespaceSocket< z.infer >; +type StructuredArgs = (Record | undefined)[]; + +export interface StructuredLogger { + log: (message: string, ...args: StructuredArgs) => any; + error: (message: string, ...args: StructuredArgs) => any; + warn: (message: string, ...args: StructuredArgs) => any; + info: (message: string, ...args: StructuredArgs) => any; + debug: (message: string, ...args: StructuredArgs) => any; + child: (fields: Record) => StructuredLogger; +} + +export enum LogLevel { + "log", + "error", + "warn", + "info", + "debug", +} + +export class SimpleStructuredLogger implements StructuredLogger { + constructor( + private name: string, + private level: LogLevel = ["1", "true"].includes(process.env.DEBUG ?? "") + ? LogLevel.debug + : LogLevel.info, + private fields?: Record + ) {} + + child(fields: Record, level?: LogLevel) { + return new SimpleStructuredLogger(this.name, level, { ...this.fields, ...fields }); + } + + log(message: string, ...args: StructuredArgs) { + if (this.level < LogLevel.log) return; + + this.#structuredLog(console.log, message, "log", ...args); + } + + error(message: string, ...args: StructuredArgs) { + if (this.level < LogLevel.error) return; + + this.#structuredLog(console.error, message, "error", ...args); + } + + warn(message: string, ...args: StructuredArgs) { + if (this.level < LogLevel.warn) return; + + this.#structuredLog(console.warn, message, "warn", ...args); + } + + info(message: string, ...args: StructuredArgs) { + if (this.level < LogLevel.info) return; + + this.#structuredLog(console.info, message, "info", ...args); + } + + debug(message: string, ...args: StructuredArgs) { + if (this.level < LogLevel.debug) return; + + this.#structuredLog(console.debug, message, "debug", ...args); + } + + #structuredLog( + loggerFunction: (message: string, ...args: any[]) => void, + message: string, + level: string, + ...args: Array | undefined> + ) { + const structuredLog = { + ...args, + ...this.fields, + timestamp: new Date(), + name: this.name, + message, + level, + }; + + loggerFunction(JSON.stringify(structuredLog)); + } +} + interface ZodNamespaceOptions< TClientMessages extends ZodSocketMessageCatalogSchema, TServerMessages extends ZodSocketMessageCatalogSchema, @@ -38,32 +119,33 @@ interface ZodNamespaceOptions< socketData?: TSocketData; handlers?: ZodSocketMessageHandlers; authToken?: string; + logger?: StructuredLogger; preAuth?: ( socket: ZodNamespaceSocket, next: (err?: ExtendedError) => void, - logger: (...args: any[]) => void + logger: StructuredLogger ) => Promise; postAuth?: ( socket: ZodNamespaceSocket, next: (err?: ExtendedError) => void, - logger: (...args: any[]) => void + logger: StructuredLogger ) => Promise; onConnection?: ( socket: ZodNamespaceSocket, handler: ZodSocketMessageHandler, sender: ZodMessageSender, - logger: (...args: any[]) => void + logger: StructuredLogger ) => Promise; onDisconnect?: ( socket: ZodNamespaceSocket, reason: DisconnectReason, description: any, - logger: (...args: any[]) => void + logger: StructuredLogger ) => Promise; onError?: ( socket: ZodNamespaceSocket, err: Error, - logger: (...args: any[]) => void + logger: StructuredLogger ) => Promise; } @@ -73,6 +155,7 @@ export class ZodNamespace< TSocketData extends z.ZodObject = any, TServerSideEvents extends EventsMap = DefaultEventsMap, > { + #logger: StructuredLogger; #handler: ZodSocketMessageHandler; sender: ZodMessageSender; @@ -87,6 +170,8 @@ export class ZodNamespace< constructor( opts: ZodNamespaceOptions ) { + this.#logger = opts.logger ?? new SimpleStructuredLogger(opts.name); + this.#handler = new ZodSocketMessageHandler({ schema: opts.clientMessages, handlers: opts.handlers, @@ -114,7 +199,7 @@ export class ZodNamespace< if (opts.preAuth) { this.namespace.use(async (socket, next) => { - const logger = createLogger(`[${opts.name}][${socket.id}][preAuth]`); + const logger = this.#logger.child({ socketId: socket.id, socketStage: "preAuth" }); if (typeof opts.preAuth === "function") { await opts.preAuth(socket, next, logger); @@ -124,21 +209,21 @@ export class ZodNamespace< if (opts.authToken) { this.namespace.use((socket, next) => { - const logger = createLogger(`[${opts.name}][${socket.id}][auth]`); + const logger = this.#logger.child({ socketId: socket.id, socketStage: "auth" }); const { auth } = socket.handshake; if (!("token" in auth)) { - logger("no token"); + logger.error("no token"); return socket.disconnect(true); } if (auth.token !== opts.authToken) { - logger("invalid token"); + logger.error("invalid token"); return socket.disconnect(true); } - logger("success"); + logger.info("success"); next(); }); @@ -146,7 +231,7 @@ export class ZodNamespace< if (opts.postAuth) { this.namespace.use(async (socket, next) => { - const logger = createLogger(`[${opts.name}][${socket.id}][postAuth]`); + const logger = this.#logger.child({ socketId: socket.id, socketStage: "auth" }); if (typeof opts.postAuth === "function") { await opts.postAuth(socket, next, logger); @@ -155,13 +240,13 @@ export class ZodNamespace< } this.namespace.on("connection", async (socket) => { - const logger = createLogger(`[${opts.name}][${socket.id}]`); - logger("connection"); + const logger = this.#logger.child({ socketId: socket.id, socketStage: "connection" }); + logger.info("connected"); this.#handler.registerHandlers(socket, logger); socket.on("disconnect", async (reason, description) => { - logger("disconnect", { reason, description }); + logger.info("disconnect", { reason, description }); if (opts.onDisconnect) { await opts.onDisconnect(socket, reason, description, logger); @@ -169,7 +254,7 @@ export class ZodNamespace< }); socket.on("error", async (error) => { - logger("error", error); + logger.error("error", { error }); if (opts.onError) { await opts.onError(socket, error, logger); @@ -186,7 +271,3 @@ export class ZodNamespace< return this.namespace.fetchSockets(); } } - -function createLogger(prefix: string) { - return (...args: any[]) => console.log(prefix, ...args); -} diff --git a/packages/core/src/v3/zodSocket.ts b/packages/core/src/v3/zodSocket.ts index 68d489796..02b1f5172 100644 --- a/packages/core/src/v3/zodSocket.ts +++ b/packages/core/src/v3/zodSocket.ts @@ -1,6 +1,7 @@ import { io, Socket } from "socket.io-client"; import { z } from "zod"; import { EventEmitterLike, ZodMessageValueSchema } from "./zodMessageHandler"; +import { LogLevel, SimpleStructuredLogger, StructuredLogger } from "./zodNamespace"; export interface ZodSocketMessageCatalogSchema { [key: string]: @@ -137,17 +138,17 @@ export class ZodSocketMessageHandler void) { - const log = logger ?? console.log; + public registerHandlers(emitter: EventEmitterLike, logger?: StructuredLogger) { + const log = logger ?? console; if (!this.#handlers) { - log("No handlers provided"); + log.info("No handlers provided"); return; } for (const eventName of Object.keys(this.#handlers)) { emitter.on(eventName, async (message: any, callback?: any): Promise => { - log(`handling ${eventName}`, message); + log.info(`handling ${eventName}`, { message, hasCallback: !!callback }); let ack; @@ -267,18 +268,18 @@ interface ZodSocketConnectionOptions< socket: ZodSocket, handler: ZodSocketMessageHandler, sender: ZodSocketMessageSender, - logger: (...args: any[]) => void + logger: StructuredLogger ) => Promise; onDisconnect?: ( socket: ZodSocket, reason: Socket.DisconnectReason, description: any, - logger: (...args: any[]) => void + logger: StructuredLogger ) => Promise; onError?: ( socket: ZodSocket, err: Error, - logger: (...args: any[]) => void + logger: StructuredLogger ) => Promise; } @@ -290,7 +291,7 @@ export class ZodSocketConnection< socket: ZodSocket; #handler: ZodSocketMessageHandler; - #logger: (...args: any[]) => void; + #logger: StructuredLogger; constructor(opts: ZodSocketConnectionOptions) { this.socket = io(`ws://${opts.host}:${opts.port}/${opts.namespace}`, { @@ -301,7 +302,9 @@ export class ZodSocketConnection< extraHeaders: opts.extraHeaders, }); - this.#logger = createLogger(`[${opts.namespace}][${this.socket.id}]`); + this.#logger = new SimpleStructuredLogger(opts.namespace, LogLevel.info, { + socketId: this.socket.id, + }); this.#handler = new ZodSocketMessageHandler({ schema: opts.serverMessages, @@ -315,7 +318,7 @@ export class ZodSocketConnection< }); this.socket.on("connect_error", async (error) => { - this.#logger(`connect_error: ${error}`); + this.#logger.error(`connect_error: ${error}`); if (opts.onError) { await opts.onError(this.socket, error, this.#logger); @@ -323,7 +326,7 @@ export class ZodSocketConnection< }); this.socket.on("connect", async () => { - this.#logger("connect"); + this.#logger.info("connect"); if (opts.onConnection) { await opts.onConnection(this.socket, this.#handler, this.#sender, this.#logger); @@ -331,7 +334,7 @@ export class ZodSocketConnection< }); this.socket.on("disconnect", async (reason, description) => { - this.#logger("disconnect"); + this.#logger.info("disconnect"); if (opts.onDisconnect) { await opts.onDisconnect(this.socket, reason, description, this.#logger);