structured logs for socket connections
This commit is contained in:
@@ -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();
|
||||
},
|
||||
|
||||
@@ -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 } : {};
|
||||
}
|
||||
);
|
||||
|
||||
@@ -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();
|
||||
|
||||
@@ -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<typeof clientWebsocketMessages>,
|
||||
MessageCatalogToSocketIoEvents<typeof serverWebsocketMessages>
|
||||
>;
|
||||
socket: Socket<
|
||||
MessageCatalogToSocketIoEvents<typeof clientWebsocketMessages>,
|
||||
MessageCatalogToSocketIoEvents<typeof serverWebsocketMessages>
|
||||
>;
|
||||
logger?: StructuredLogger;
|
||||
}
|
||||
|
||||
export class SharedSocketConnection {
|
||||
public id: string;
|
||||
public onClose: Evt<DisconnectReason> = new Evt();
|
||||
@@ -19,17 +32,7 @@ export class SharedSocketConnection {
|
||||
private _sharedConsumer: SharedQueueConsumer;
|
||||
private _messageHandler: ZodMessageHandler<typeof clientWebsocketMessages>;
|
||||
|
||||
constructor(
|
||||
namespace: Namespace<
|
||||
MessageCatalogToSocketIoEvents<typeof clientWebsocketMessages>,
|
||||
MessageCatalogToSocketIoEvents<typeof serverWebsocketMessages>
|
||||
>,
|
||||
private socket: Socket<
|
||||
MessageCatalogToSocketIoEvents<typeof clientWebsocketMessages>,
|
||||
MessageCatalogToSocketIoEvents<typeof serverWebsocketMessages>
|
||||
>,
|
||||
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() {
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -1,4 +1,5 @@
|
||||
import { z } from "zod";
|
||||
import { StructuredLogger } from "./zodNamespace";
|
||||
|
||||
export type ZodMessageValueSchema<TDiscriminatedUnion extends z.ZodDiscriminatedUnion<any, any>> =
|
||||
| z.ZodFirstPartySchemaTypes
|
||||
@@ -92,17 +93,17 @@ export class ZodMessageHandler<TMessageCatalog extends ZodMessageCatalogSchema>
|
||||
};
|
||||
}
|
||||
|
||||
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<void> => {
|
||||
log(`handling ${eventName}`, message);
|
||||
log.info(`handling ${eventName}`, message);
|
||||
|
||||
let ack;
|
||||
|
||||
|
||||
@@ -25,6 +25,87 @@ export type ZodNamespaceSocket<
|
||||
z.infer<TSocketData>
|
||||
>;
|
||||
|
||||
type StructuredArgs = (Record<string, unknown> | 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<string, unknown>) => 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<string, unknown>
|
||||
) {}
|
||||
|
||||
child(fields: Record<string, unknown>, 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<Record<string, unknown> | 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<TClientMessages>;
|
||||
authToken?: string;
|
||||
logger?: StructuredLogger;
|
||||
preAuth?: (
|
||||
socket: ZodNamespaceSocket<TClientMessages, TServerMessages, TServerSideEvents, TSocketData>,
|
||||
next: (err?: ExtendedError) => void,
|
||||
logger: (...args: any[]) => void
|
||||
logger: StructuredLogger
|
||||
) => Promise<void>;
|
||||
postAuth?: (
|
||||
socket: ZodNamespaceSocket<TClientMessages, TServerMessages, TServerSideEvents, TSocketData>,
|
||||
next: (err?: ExtendedError) => void,
|
||||
logger: (...args: any[]) => void
|
||||
logger: StructuredLogger
|
||||
) => Promise<void>;
|
||||
onConnection?: (
|
||||
socket: ZodNamespaceSocket<TClientMessages, TServerMessages, TServerSideEvents, TSocketData>,
|
||||
handler: ZodSocketMessageHandler<TClientMessages>,
|
||||
sender: ZodMessageSender<TServerMessages>,
|
||||
logger: (...args: any[]) => void
|
||||
logger: StructuredLogger
|
||||
) => Promise<void>;
|
||||
onDisconnect?: (
|
||||
socket: ZodNamespaceSocket<TClientMessages, TServerMessages, TServerSideEvents, TSocketData>,
|
||||
reason: DisconnectReason,
|
||||
description: any,
|
||||
logger: (...args: any[]) => void
|
||||
logger: StructuredLogger
|
||||
) => Promise<void>;
|
||||
onError?: (
|
||||
socket: ZodNamespaceSocket<TClientMessages, TServerMessages, TServerSideEvents, TSocketData>,
|
||||
err: Error,
|
||||
logger: (...args: any[]) => void
|
||||
logger: StructuredLogger
|
||||
) => Promise<void>;
|
||||
}
|
||||
|
||||
@@ -73,6 +155,7 @@ export class ZodNamespace<
|
||||
TSocketData extends z.ZodObject<any, any, any> = any,
|
||||
TServerSideEvents extends EventsMap = DefaultEventsMap,
|
||||
> {
|
||||
#logger: StructuredLogger;
|
||||
#handler: ZodSocketMessageHandler<TClientMessages>;
|
||||
sender: ZodMessageSender<TServerMessages>;
|
||||
|
||||
@@ -87,6 +170,8 @@ export class ZodNamespace<
|
||||
constructor(
|
||||
opts: ZodNamespaceOptions<TClientMessages, TServerMessages, TServerSideEvents, TSocketData>
|
||||
) {
|
||||
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);
|
||||
}
|
||||
|
||||
@@ -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<TRPCCatalog extends ZodSocketMessageCatalog
|
||||
};
|
||||
}
|
||||
|
||||
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.#handlers)) {
|
||||
emitter.on(eventName, async (message: any, callback?: any): Promise<void> => {
|
||||
log(`handling ${eventName}`, message);
|
||||
log.info(`handling ${eventName}`, { message, hasCallback: !!callback });
|
||||
|
||||
let ack;
|
||||
|
||||
@@ -267,18 +268,18 @@ interface ZodSocketConnectionOptions<
|
||||
socket: ZodSocket<TServerMessages, TClientMessages>,
|
||||
handler: ZodSocketMessageHandler<TServerMessages>,
|
||||
sender: ZodSocketMessageSender<TClientMessages>,
|
||||
logger: (...args: any[]) => void
|
||||
logger: StructuredLogger
|
||||
) => Promise<void>;
|
||||
onDisconnect?: (
|
||||
socket: ZodSocket<TServerMessages, TClientMessages>,
|
||||
reason: Socket.DisconnectReason,
|
||||
description: any,
|
||||
logger: (...args: any[]) => void
|
||||
logger: StructuredLogger
|
||||
) => Promise<void>;
|
||||
onError?: (
|
||||
socket: ZodSocket<TServerMessages, TClientMessages>,
|
||||
err: Error,
|
||||
logger: (...args: any[]) => void
|
||||
logger: StructuredLogger
|
||||
) => Promise<void>;
|
||||
}
|
||||
|
||||
@@ -290,7 +291,7 @@ export class ZodSocketConnection<
|
||||
socket: ZodSocket<TServerMessages, TClientMessages>;
|
||||
|
||||
#handler: ZodSocketMessageHandler<TServerMessages>;
|
||||
#logger: (...args: any[]) => void;
|
||||
#logger: StructuredLogger;
|
||||
|
||||
constructor(opts: ZodSocketConnectionOptions<TClientMessages, TServerMessages>) {
|
||||
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);
|
||||
|
||||
Reference in New Issue
Block a user