From 010bb0957d11fe8512d0e2d9ddb6fda1971a48b5 Mon Sep 17 00:00:00 2001 From: mahyargdz Date: Wed, 16 Jul 2025 12:52:18 +0330 Subject: [PATCH] chore: fix --- src/core/filters/ws-exception.filter.ts | 246 -------------- src/main.ts | 4 - src/modules/auth/guards/ws-auth.guard.ts | 240 ------------- src/modules/email/WEBSOCKET_USAGE.md | 318 ------------------ .../email/constants/email-events.constant.ts | 37 +- src/modules/email/email.gateway.ts | 227 ------------- src/modules/email/email.module.ts | 7 +- .../interfaces/email-events.interface.ts | 60 +++- .../services/email-notification.service.ts | 158 --------- src/modules/email/services/email.service.ts | 107 ------ .../email/services/websocket-auth.service.ts | 240 ------------- 11 files changed, 58 insertions(+), 1586 deletions(-) delete mode 100644 src/core/filters/ws-exception.filter.ts delete mode 100644 src/modules/auth/guards/ws-auth.guard.ts delete mode 100644 src/modules/email/WEBSOCKET_USAGE.md delete mode 100644 src/modules/email/email.gateway.ts delete mode 100644 src/modules/email/services/email-notification.service.ts delete mode 100644 src/modules/email/services/websocket-auth.service.ts diff --git a/src/core/filters/ws-exception.filter.ts b/src/core/filters/ws-exception.filter.ts deleted file mode 100644 index f911f01..0000000 --- a/src/core/filters/ws-exception.filter.ts +++ /dev/null @@ -1,246 +0,0 @@ -import { ArgumentsHost, Catch, HttpException, Logger, WsExceptionFilter as NestWsExceptionFilter } from "@nestjs/common"; -import { WsException } from "@nestjs/websockets"; -import { Socket } from "socket.io"; - -import { WebSocketMessage } from "../../common/enums/message.enum"; -import { WEBSOCKET_EVENTS } from "../../modules/email/constants/email-events.constant"; - -export interface WebSocketErrorResponse { - error: string; - message: string; - timestamp: string; - event?: string; - statusCode?: number; -} - -/** - * Global WebSocket exception filter - * Catches and handles all WebSocket-related errors - */ -@Catch() -export class WsExceptionFilter implements NestWsExceptionFilter { - private readonly logger = new Logger(WsExceptionFilter.name); - - catch(exception: Error, host: ArgumentsHost): void { - const client: Socket = host.switchToWs().getClient(); - const data = host.switchToWs().getData(); - - const errorResponse = this.createErrorResponse(exception, data); - - // Log the error for monitoring - this.logError(exception, client, data); - - // Emit error to client - client.emit(WEBSOCKET_EVENTS.CONNECTION_ERROR, errorResponse); - - // Handle specific error types - this.handleSpecificErrors(exception, client); - } - - private createErrorResponse(exception: Error, data?: any): WebSocketErrorResponse { - let error: string; - let message: string; - let statusCode: number | undefined; - - if (exception instanceof WsException) { - const wsError = exception.getError(); - - if (typeof wsError === "string") { - error = wsError; - message = wsError; - } else if (typeof wsError === "object" && wsError !== null) { - error = (wsError as any).error || "WebSocket Error"; - message = (wsError as any).message || "An error occurred"; - statusCode = (wsError as any).statusCode; - } else { - error = "WebSocket Error"; - message = "An unknown WebSocket error occurred"; - } - } else if (exception instanceof HttpException) { - statusCode = exception.getStatus(); - const response = exception.getResponse(); - - if (typeof response === "string") { - error = response; - message = response; - } else if (typeof response === "object" && response !== null) { - error = (response as any).error || "HTTP Error"; - message = (response as any).message || exception.message; - } else { - error = "HTTP Error"; - message = exception.message; - } - } else { - // Generic error handling - error = "Internal Server Error"; - message = this.getSafeErrorMessage(exception); - } - - return { - error, - message, - timestamp: new Date().toISOString(), - event: data?.event || "unknown", - statusCode, - }; - } - - private logError(exception: Error, client: Socket, data?: any): void { - const errorContext = { - clientId: client.id, - userId: client.data?.user?.id, - sessionId: client.data?.sessionId, - event: data?.event, - errorType: exception.constructor.name, - errorMessage: exception.message, - stack: exception.stack, - timestamp: new Date().toISOString(), - }; - - if (exception instanceof WsException) { - this.logger.warn("WebSocket exception occurred", errorContext); - } else if (exception instanceof HttpException) { - this.logger.warn("HTTP exception in WebSocket context", { - ...errorContext, - statusCode: exception.getStatus(), - }); - } else { - this.logger.error("Unexpected error in WebSocket context", errorContext); - } - } - - private handleSpecificErrors(exception: Error, client: Socket): void { - if (exception instanceof WsException) { - const wsError = exception.getError(); - - // Handle authentication errors - if (this.isAuthenticationError(wsError)) { - this.handleAuthenticationError(client); - return; - } - - // Handle permission errors - if (this.isPermissionError(wsError)) { - this.handlePermissionError(client); - return; - } - - // Handle rate limiting errors - if (this.isRateLimitError(wsError)) { - this.handleRateLimitError(client); - return; - } - } - - // Handle connection errors - if (this.isConnectionError(exception)) { - this.handleConnectionError(client); - return; - } - - // Handle validation errors - if (this.isValidationError(exception)) { - this.handleValidationError(client); - return; - } - } - - private isAuthenticationError(error: any): boolean { - if (typeof error === "string") { - return error.includes("authentication") || error.includes("token") || error.includes("unauthorized"); - } - if (typeof error === "object" && error !== null) { - const errorMessage = (error as any).message || (error as any).error || ""; - return errorMessage.includes("authentication") || errorMessage.includes("token") || errorMessage.includes("unauthorized"); - } - return false; - } - - private isPermissionError(error: any): boolean { - if (typeof error === "string") { - return error.includes("permission") || error.includes("access denied") || error.includes("forbidden"); - } - if (typeof error === "object" && error !== null) { - const errorMessage = (error as any).message || (error as any).error || ""; - return errorMessage.includes("permission") || errorMessage.includes("access denied") || errorMessage.includes("forbidden"); - } - return false; - } - - private isRateLimitError(error: any): boolean { - if (typeof error === "string") { - return error.includes("rate limit") || error.includes("too many requests"); - } - if (typeof error === "object" && error !== null) { - const errorMessage = (error as any).message || (error as any).error || ""; - return errorMessage.includes("rate limit") || errorMessage.includes("too many requests"); - } - return false; - } - - private isConnectionError(exception: Error): boolean { - return exception.message.includes("connection") || exception.message.includes("disconnect") || exception.message.includes("timeout"); - } - - private isValidationError(exception: Error): boolean { - return exception.message.includes("validation") || exception.message.includes("invalid") || exception.name.includes("ValidationError"); - } - - private handleAuthenticationError(client: Socket): void { - client.emit(WEBSOCKET_EVENTS.AUTHENTICATION_FAILED, { - message: WebSocketMessage.WS_AUTHENTICATION_FAILED, - timestamp: new Date().toISOString(), - reconnectAllowed: true, - }); - - // Disconnect after a short delay to allow client to handle the error - setTimeout(() => { - if (client.connected) { - client.disconnect(true); - } - }, 500); - } - - private handlePermissionError(client: Socket): void { - client.emit(WEBSOCKET_EVENTS.CONNECTION_ERROR, { - error: WebSocketMessage.PERMISSION_DENIED, - message: WebSocketMessage.PERMISSION_DENIED, - timestamp: new Date().toISOString(), - }); - } - - private handleRateLimitError(client: Socket): void { - client.emit(WEBSOCKET_EVENTS.CONNECTION_ERROR, { - error: WebSocketMessage.WS_RATE_LIMIT_EXCEEDED, - message: WebSocketMessage.WS_RATE_LIMIT_EXCEEDED, - timestamp: new Date().toISOString(), - retryAfter: 60000, // 1 minute - }); - } - - private handleConnectionError(client: Socket): void { - client.emit(WEBSOCKET_EVENTS.CONNECTION_ERROR, { - error: WebSocketMessage.WS_CONNECTION_ERROR, - message: WebSocketMessage.WS_CONNECTION_ERROR, - timestamp: new Date().toISOString(), - reconnectAllowed: true, - }); - } - - private handleValidationError(client: Socket): void { - client.emit(WEBSOCKET_EVENTS.CONNECTION_ERROR, { - error: WebSocketMessage.INVALID_DATA, - message: WebSocketMessage.INVALID_DATA, - timestamp: new Date().toISOString(), - }); - } - - private getSafeErrorMessage(exception: Error): string { - // In production, don't expose internal error details - if (process.env.NODE_ENV === "production") { - return WebSocketMessage.INTERNAL_ERROR; - } - - return exception.message || WebSocketMessage.INTERNAL_ERROR; - } -} diff --git a/src/main.ts b/src/main.ts index 17acb39..8f2aaba 100644 --- a/src/main.ts +++ b/src/main.ts @@ -3,7 +3,6 @@ import { Logger, ValidationPipe } from "@nestjs/common"; import { ConfigService } from "@nestjs/config"; import { NestFactory } from "@nestjs/core"; import { FastifyAdapter, NestFastifyApplication } from "@nestjs/platform-fastify"; -import { IoAdapter } from "@nestjs/platform-socket.io"; import { AppModule } from "./app.module"; import { getSwaggerDocument } from "./configs/swagger.config"; @@ -26,9 +25,6 @@ async function bootstrap() { const app = await NestFactory.create(AppModule, fastifyAdapter); - // Setup WebSocket adapter - app.useWebSocketAdapter(new IoAdapter(app)); - app.useGlobalPipes(new ValidationPipe({ transform: true, whitelist: true })); app.useGlobalInterceptors(new ResponseInterceptor(), new PaginationInterceptor()); diff --git a/src/modules/auth/guards/ws-auth.guard.ts b/src/modules/auth/guards/ws-auth.guard.ts deleted file mode 100644 index 84370eb..0000000 --- a/src/modules/auth/guards/ws-auth.guard.ts +++ /dev/null @@ -1,240 +0,0 @@ -import { CanActivate, ExecutionContext, Injectable, Logger } from "@nestjs/common"; -import { WsException } from "@nestjs/websockets"; -import { Socket } from "socket.io"; - -import { WebSocketMessage } from "../../../common/enums/message.enum"; -import { WebSocketAuthService } from "../../email/services/websocket-auth.service"; - -/** - * WebSocket authentication guard - * Protects WebSocket endpoints by verifying JWT tokens - */ -@Injectable() -export class WsAuthGuard implements CanActivate { - private readonly logger = new Logger(WsAuthGuard.name); - - constructor(private readonly authService: WebSocketAuthService) {} - - async canActivate(context: ExecutionContext): Promise { - try { - const client: Socket = context.switchToWs().getClient(); - const data = context.switchToWs().getData(); - - // Skip authentication for certain events - if (this.shouldSkipAuthentication(data?.event)) { - return true; - } - - // Check if client is already authenticated - if (this.isClientAuthenticated(client)) { - // Update last activity - this.authService.updateLastActivity(client); - return true; - } - - // Attempt to authenticate client - const authResult = await this.authService.authenticateClient(client); - - if (!authResult.success) { - this.logger.warn("WebSocket authentication failed in guard", { - clientId: client.id, - event: data?.event, - error: authResult.error, - }); - - throw new WsException({ - error: authResult.error || WebSocketMessage.WS_AUTHENTICATION_FAILED, - message: authResult.error || WebSocketMessage.WS_AUTHENTICATION_FAILED, - statusCode: 401, - }); - } - - this.logger.log("WebSocket authentication successful in guard", { - clientId: client.id, - userId: authResult.user?.id, - event: data?.event, - }); - - return true; - } catch (error) { - this.logger.error("WebSocket guard error", { - error: error instanceof Error ? error.message : String(error), - context: context.getType(), - }); - - if (error instanceof WsException) { - throw error; - } - - throw new WsException({ - error: WebSocketMessage.WS_AUTHENTICATION_FAILED, - message: WebSocketMessage.WS_AUTHENTICATION_FAILED, - statusCode: 401, - }); - } - } - - /** - * Checks if client is already authenticated - */ - private isClientAuthenticated(client: Socket): boolean { - return !!(client.data?.user?.id && client.data?.connectedAt); - } - - /** - * Determines if authentication should be skipped for certain events - */ - private shouldSkipAuthentication(event?: string): boolean { - const skipEvents = ["connect", "disconnect", "ping", "pong", "error", "connection_error", "authenticate"]; - - return skipEvents.includes(event || ""); - } -} - -// /** -// * Optional: Role-based WebSocket guard -// * Extends authentication to include role-based access control -// */ -// @Injectable() -// export class WsRoleGuard implements CanActivate { -// private readonly logger = new Logger(WsRoleGuard.name); -// private readonly requiredRoles: string[] = []; - -// constructor( -// private readonly authService: WebSocketAuthService, -// requiredRoles: string[] = [], -// ) { -// this.requiredRoles = requiredRoles; -// } - -// async canActivate(context: ExecutionContext): Promise { -// try { -// const client: Socket = context.switchToWs().getClient(); -// const data = context.switchToWs().getData(); - -// // First check authentication -// if (!this.isClientAuthenticated(client)) { -// throw new WsException({ -// error: WebSocketMessage.WS_AUTHENTICATION_REQUIRED, -// message: WebSocketMessage.WS_AUTHENTICATION_REQUIRED, -// statusCode: 401, -// }); -// } - -// // Then check roles if required -// if (this.requiredRoles.length > 0) { -// const user = client.data.user; -// const hasRequiredRole = this.checkUserRoles(user.permissions || [], this.requiredRoles); - -// if (!hasRequiredRole && !user.isAdmin) { -// this.logger.warn("WebSocket role check failed", { -// clientId: client.id, -// userId: user.id, -// userRoles: user.permissions, -// requiredRoles: this.requiredRoles, -// event: data?.event, -// }); - -// throw new WsException({ -// error: WebSocketMessage.PERMISSION_DENIED, -// message: WebSocketMessage.PERMISSION_DENIED, -// statusCode: 403, -// }); -// } -// } - -// return true; -// } catch (error) { -// this.logger.error("WebSocket role guard error", { -// error: error instanceof Error ? error.message : String(error), -// requiredRoles: this.requiredRoles, -// }); - -// if (error instanceof WsException) { -// throw error; -// } - -// throw new WsException({ -// error: WebSocketMessage.PERMISSION_DENIED, -// message: WebSocketMessage.PERMISSION_DENIED, -// statusCode: 403, -// }); -// } -// } - -// private isClientAuthenticated(client: Socket): boolean { -// return !!(client.data?.user?.id && client.data?.connectedAt); -// } - -// private checkUserRoles(userRoles: string[], requiredRoles: string[]): boolean { -// return requiredRoles.some((role) => userRoles.includes(role)); -// } -// } - -// /** -// * Room-based WebSocket guard -// * Checks if user has permission to access specific rooms -// */ -// @Injectable() -// export class WsRoomGuard implements CanActivate { -// private readonly logger = new Logger(WsRoomGuard.name); - -// constructor(private readonly authService: WebSocketAuthService) {} - -// async canActivate(context: ExecutionContext): Promise { -// try { -// const client: Socket = context.switchToWs().getClient(); -// const data = context.switchToWs().getData(); - -// // Check authentication first -// if (!this.isClientAuthenticated(client)) { -// throw new WsException({ -// error: WebSocketMessage.WS_AUTHENTICATION_REQUIRED, -// message: WebSocketMessage.WS_AUTHENTICATION_REQUIRED, -// statusCode: 401, -// }); -// } - -// // Check room access if roomId is provided in data -// if (data?.roomId) { -// const user = client.data.user; -// const canAccess = this.authService.canAccessRoom(user, data.roomId); - -// if (!canAccess) { -// this.logger.warn("WebSocket room access denied", { -// clientId: client.id, -// userId: user.id, -// roomId: data.roomId, -// event: data?.event, -// }); - -// throw new WsException({ -// error: WebSocketMessage.PERMISSION_DENIED, -// message: WebSocketMessage.UNAUTHORIZED_ROOM_ACCESS, -// statusCode: 403, -// }); -// } -// } - -// return true; -// } catch (error) { -// this.logger.error("WebSocket room guard error", { -// error: error instanceof Error ? error.message : String(error), -// }); - -// if (error instanceof WsException) { -// throw error; -// } - -// throw new WsException({ -// error: WebSocketMessage.PERMISSION_DENIED, -// message: WebSocketMessage.PERMISSION_DENIED, -// statusCode: 403, -// }); -// } -// } - -// private isClientAuthenticated(client: Socket): boolean { -// return !!(client.data?.user?.id && client.data?.connectedAt); -// } -// } diff --git a/src/modules/email/WEBSOCKET_USAGE.md b/src/modules/email/WEBSOCKET_USAGE.md deleted file mode 100644 index 32468ee..0000000 --- a/src/modules/email/WEBSOCKET_USAGE.md +++ /dev/null @@ -1,318 +0,0 @@ -# Email WebSocket Gateway Usage - -This document explains how to use the WebSocket gateway for real-time email notifications. - -## Overview - -The Email WebSocket Gateway provides real-time notifications for email events such as: - -- New message arrivals -- Message status updates (read/unread, flagged/unflagged) -- Message movements (archive, favorite, delete) -- Connection status updates - -## Connection - -### WebSocket URL - -``` -ws://localhost:3000/email -``` - -### Authentication - -The gateway requires JWT authentication. Pass the token in one of these ways: - -**Method 1: Auth object** - -```javascript -const socket = io("ws://localhost:3000/email", { - auth: { - token: "your-jwt-token", - }, -}); -``` - -**Method 2: Query parameter** - -```javascript -const socket = io("ws://localhost:3000/email?token=your-jwt-token"); -``` - -## Client-Side Implementation - -### Basic Setup (JavaScript) - -```javascript -import { io } from "socket.io-client"; - -const socket = io("ws://localhost:3000/email", { - auth: { - token: localStorage.getItem("authToken"), // Your JWT token - }, -}); - -// Connection events -socket.on("connect", () => { - console.log("Connected to email gateway"); -}); - -socket.on("connection_established", (data) => { - console.log("Connection established:", data); -}); - -socket.on("connection_error", (error) => { - console.error("Connection error:", error); -}); - -socket.on("disconnect", () => { - console.log("Disconnected from email gateway"); -}); -``` - -### Listening for Email Events - -```javascript -// New message notification -socket.on("new_message", (message) => { - console.log("New message received:", message); - // Update UI with new message - updateInboxUI(message); - showNotification(`New email from ${message.from.name}: ${message.subject}`); -}); - -// Message status update -socket.on("message_status_update", (update) => { - console.log("Message status updated:", update); - // Update message status in UI - updateMessageStatus(update.messageId, update.status); -}); - -// Message moved -socket.on("message_moved", (moveData) => { - console.log("Message moved:", moveData); - // Update UI to reflect message movement - moveMessageInUI(moveData.messageId, moveData.fromMailboxName, moveData.toMailboxName); -}); - -// Message deleted -socket.on("message_deleted", (deleteData) => { - console.log("Message deleted:", deleteData); - // Remove message from UI - removeMessageFromUI(deleteData.messageId); -}); -``` - -### Joining/Leaving User Rooms - -```javascript -// Join user room (usually done automatically on connection) -socket.emit("join_user_room", { - userId: "user-id", - clientId: socket.id, - timestamp: new Date().toISOString(), -}); - -// Leave user room -socket.emit("leave_user_room", { - userId: "user-id", - clientId: socket.id, - timestamp: new Date().toISOString(), -}); -``` - -### React Hook Example - -```javascript -import { useEffect, useState } from "react"; -import { io } from "socket.io-client"; - -export function useEmailWebSocket(authToken) { - const [socket, setSocket] = useState(null); - const [isConnected, setIsConnected] = useState(false); - const [newMessages, setNewMessages] = useState([]); - - useEffect(() => { - if (!authToken) return; - - const newSocket = io("ws://localhost:3000/email", { - auth: { token: authToken }, - }); - - newSocket.on("connect", () => { - setIsConnected(true); - }); - - newSocket.on("disconnect", () => { - setIsConnected(false); - }); - - newSocket.on("new_message", (message) => { - setNewMessages((prev) => [...prev, message]); - }); - - newSocket.on("message_status_update", (update) => { - // Handle message status updates - console.log("Message status update:", update); - }); - - setSocket(newSocket); - - return () => { - newSocket.disconnect(); - }; - }, [authToken]); - - return { socket, isConnected, newMessages }; -} -``` - -## Event Types - -### Server to Client Events - -#### `new_message` - -Emitted when a new message arrives. - -```typescript -interface EmailNotificationPayload { - messageId: number; - userId: string; - from: { name?: string; address: string }; - to: Array<{ name?: string; address: string }>; - subject: string; - preview?: string; - hasAttachments: boolean; - received: string; - mailboxId: string; - mailboxName: string; - isRead: boolean; - isFlagged: boolean; - size: number; -} -``` - -#### `message_status_update` - -Emitted when a message's status changes. - -```typescript -interface EmailStatusUpdatePayload { - messageId: number; - userId: string; - status: "read" | "unread" | "flagged" | "unflagged" | "archived" | "deleted"; - mailboxId: string; - timestamp: string; -} -``` - -#### `message_moved` - -Emitted when a message is moved between mailboxes. - -```typescript -interface EmailMovePayload { - messageId: number; - userId: string; - fromMailboxId: string; - toMailboxId: string; - fromMailboxName: string; - toMailboxName: string; - timestamp: string; -} -``` - -#### `message_deleted` - -Emitted when a message is deleted. - -```typescript -interface EmailDeletePayload { - messageId: number; - userId: string; - mailboxId: string; - mailboxName: string; - timestamp: string; -} -``` - -### Client to Server Events - -#### `join_user_room` - -Join a user's room to receive notifications. - -```typescript -interface ClientJoinPayload { - userId: string; - clientId: string; - timestamp: string; -} -``` - -#### `leave_user_room` - -Leave a user's room to stop receiving notifications. - -## Error Handling - -```javascript -socket.on("connection_error", (error) => { - console.error("WebSocket connection error:", error); - - switch (error.error) { - case "Authentication required": - // Redirect to login - break; - case "Invalid token": - // Refresh token and reconnect - break; - case "Authentication failed": - // Show error message - break; - default: - console.error("Unknown error:", error); - } -}); -``` - -## Best Practices - -1. **Authentication**: Always include a valid JWT token -2. **Reconnection**: Handle connection drops gracefully -3. **Error Handling**: Implement proper error handling for all events -4. **Memory Management**: Clean up event listeners when components unmount -5. **Rate Limiting**: Don't overwhelm the server with too many connections - -## Testing - -You can test the WebSocket gateway using the simulation endpoint: - -```javascript -// Call the simulation endpoint from your backend -await emailNotificationService.simulateNewMessage("user-id", { - subject: "Test Message", - from: { name: "Test User", address: "test@example.com" }, -}); -``` - -## Troubleshooting - -### Common Issues - -1. **Connection Refused**: Check if the server is running and WebSocket dependencies are installed -2. **Authentication Failed**: Verify JWT token is valid and not expired -3. **No Events Received**: Ensure you're subscribed to the correct user room -4. **CORS Issues**: Configure CORS settings in the gateway properly - -### Debug Mode - -Enable debug mode to see detailed logs: - -```javascript -const socket = io("ws://localhost:3000/email", { - auth: { token: authToken }, - debug: true, -}); -``` diff --git a/src/modules/email/constants/email-events.constant.ts b/src/modules/email/constants/email-events.constant.ts index e11d524..2a63446 100644 --- a/src/modules/email/constants/email-events.constant.ts +++ b/src/modules/email/constants/email-events.constant.ts @@ -2,46 +2,35 @@ export const WEBSOCKET_EVENTS = { // Connection Events CONNECTION_ESTABLISHED: "connection_established", CONNECTION_ERROR: "connection_error", - USER_JOINED: "user_joined", - USER_LEFT: "user_left", // Email Events - Server to Client - NEW_MESSAGE: "new_message", - MESSAGE_STATUS_UPDATE: "message_status_update", - MESSAGE_MOVED: "message_moved", - MESSAGE_DELETED: "message_deleted", + NEW_EMAIL: "new_email", + EMAIL_READ: "email_read", + EMAIL_UNREAD: "email_unread", + EMAIL_FLAGGED: "email_flagged", + EMAIL_UNFLAGGED: "email_unflagged", + EMAIL_DELETED: "email_deleted", + EMAIL_MOVED: "email_moved", + EMAIL_SENT: "email_sent", - // Room Management - Client to Server - JOIN_USER_ROOM: "join_user_room", - LEAVE_USER_ROOM: "leave_user_room", - JOIN_SESSION: "join_session", - LEAVE_SESSION: "leave_session", - - // Authentication Events - AUTHENTICATE: "authenticate", - AUTHENTICATION_SUCCESS: "authentication_success", - AUTHENTICATION_FAILED: "authentication_failed", + // Mailbox Events + MAILBOX_UPDATED: "mailbox_updated", + UNREAD_COUNT_UPDATED: "unread_count_updated", // System Events PING: "ping", PONG: "pong", - RECONNECT: "reconnect", - FORCE_DISCONNECT: "force_disconnect", } as const; export const WEBSOCKET_ROOMS = { - USER: (userId: string) => `user_${userId}`, - SESSION: (sessionId: string) => `session_${sessionId}`, - BROADCAST: "broadcast", + USER: (userId: string) => `user:${userId}`, } as const; export const WEBSOCKET_NAMESPACES = { EMAIL: "/email", - NOTIFICATIONS: "/notifications", - ADMIN: "/admin", } as const; -// Type definitions for better TypeScript support +// Type definitions export type WebSocketEvent = (typeof WEBSOCKET_EVENTS)[keyof typeof WEBSOCKET_EVENTS]; export type WebSocketRoom = string; export type WebSocketNamespace = (typeof WEBSOCKET_NAMESPACES)[keyof typeof WEBSOCKET_NAMESPACES]; diff --git a/src/modules/email/email.gateway.ts b/src/modules/email/email.gateway.ts deleted file mode 100644 index f983aa0..0000000 --- a/src/modules/email/email.gateway.ts +++ /dev/null @@ -1,227 +0,0 @@ -import { Logger, UseFilters, UseGuards, UsePipes, ValidationPipe } from "@nestjs/common"; -import { - ConnectedSocket, - MessageBody, - OnGatewayConnection, - OnGatewayDisconnect, - OnGatewayInit, - SubscribeMessage, - WebSocketGateway, - WebSocketServer, -} from "@nestjs/websockets"; -import { Server, Socket } from "socket.io"; - -import { WEBSOCKET_EVENTS } from "./constants/email-events.constant"; -import { - ClientJoinPayload, - EmailDeletePayload, - EmailMovePayload, - EmailNotificationPayload, - EmailStatusUpdatePayload, -} from "./interfaces/email-events.interface"; -import { ConnectedUserInfo } from "./interfaces/websocket.interface"; -import { WebSocketAuthService } from "./services/websocket-auth.service"; -import { WebSocketMessage } from "../../common/enums/message.enum"; -import { WsExceptionFilter } from "../../core/filters/ws-exception.filter"; -import { WsAuthGuard } from "../auth/guards/ws-auth.guard"; - -@WebSocketGateway({ - cors: { origin: "*", credentials: true }, - namespace: "/email", -}) -@UseFilters(WsExceptionFilter) -@UseGuards(WsAuthGuard) -@UsePipes(new ValidationPipe()) -export class EmailGateway implements OnGatewayInit, OnGatewayConnection, OnGatewayDisconnect { - private readonly logger = new Logger(EmailGateway.name); - - // Connection state management - private readonly userConnections: Map> = new Map(); // userId -> Set of socketIds - private readonly connectedUsers = new Map(); - private readonly userSessions = new Map>(); - - @WebSocketServer() - server: Server; - constructor(private readonly authService: WebSocketAuthService) {} - - afterInit(server: Server): void { - this.logger.log("WebSocket Gateway initialized", { - namespace: "/email", - cors: true, - serverInstance: !!server, - }); - } - - /** - * Handles new WebSocket connections with authentication - * Implements fail-fast principle for security - */ - async handleConnection(client: Socket): Promise { - const connectionContext = { - clientId: client.id, - clientIP: client.handshake.address, - userAgent: client.handshake.headers["user-agent"], - timestamp: new Date().toISOString(), - }; - - this.logger.log("New WebSocket connection attempt", connectionContext); - - try { - const authResult = await this.authService.authenticateClient(client); - - if (!authResult.success) { - this.authService.handleAuthenticationFailure(client, authResult.error!); - return; - } - - const user = authResult.user!; - this.registerUserConnection(client.id, user.id); - this.authService.emitAuthenticationSuccess(client, user); - - this.logger.log("WebSocket connection established", { - ...connectionContext, - userId: user.id, - authenticated: true, - }); - } catch (error) { - this.logger.error("WebSocket connection error", { - ...connectionContext, - error: error instanceof Error ? error.message : String(error), - }); - - this.authService.handleAuthenticationFailure(client, WebSocketMessage.WS_CONNECTION_FAILED); - } - } - - /** - * Handles client disconnections with cleanup - */ - handleDisconnect(client: Socket): void { - const userInfo = this.connectedUsers.get(client.id); - - this.logger.log("WebSocket client disconnected", { - clientId: client.id, - userId: userInfo?.userId, - sessionId: userInfo?.sessionId, - }); - - if (userInfo) { - this.cleanupUserConnection(client, userInfo); - } - - this.connectedUsers.delete(client.id); - } - - @SubscribeMessage("join_user_room") - async handleJoinUserRoom(@ConnectedSocket() client: Socket, @MessageBody() data: ClientJoinPayload) { - try { - const userId = client.data.userId; - - if (userId !== data.userId) { - client.emit("connection_error", { error: "Unauthorized", timestamp: new Date().toISOString() }); - return; - } - - await client.join(`user:${userId}`); - this.logger.log(`Client ${client.id} joined room for user ${userId}`); - - return { success: true, message: "Joined user room successfully" }; - } catch (error) { - this.logger.error(`Error joining user room: ${error instanceof Error ? error.message : "Unknown error"}`); - client.emit("connection_error", { error: "Failed to join room", timestamp: new Date().toISOString() }); - } - } - - @SubscribeMessage("leave_user_room") - async handleLeaveUserRoom(@ConnectedSocket() client: Socket, @MessageBody() data: ClientJoinPayload) { - try { - const userId = client.data.userId; - - if (userId !== data.userId) { - client.emit("connection_error", { error: "Unauthorized", timestamp: new Date().toISOString() }); - return; - } - - await client.leave(`user:${userId}`); - this.logger.log(`Client ${client.id} left room for user ${userId}`); - - return { success: true, message: "Left user room successfully" }; - } catch (error) { - this.logger.error(`Error leaving user room: ${error instanceof Error ? error.message : "Unknown error"}`); - client.emit("connection_error", { error: "Failed to leave room", timestamp: new Date().toISOString() }); - } - } - - // Methods to emit events to clients - async notifyNewMessage(payload: EmailNotificationPayload) { - this.logger.log(`Emitting new message notification to user ${payload.userId}`); - this.server.to(`user:${payload.userId}`).emit("new_message", payload); - } - - async notifyMessageStatusUpdate(payload: EmailStatusUpdatePayload) { - this.logger.log(`Emitting message status update to user ${payload.userId}`); - this.server.to(`user:${payload.userId}`).emit("message_status_update", payload); - } - - async notifyMessageMoved(payload: EmailMovePayload) { - this.logger.log(`Emitting message moved notification to user ${payload.userId}`); - this.server.to(`user:${payload.userId}`).emit("message_moved", payload); - } - - async notifyMessageDeleted(payload: EmailDeletePayload) { - this.logger.log(`Emitting message deleted notification to user ${payload.userId}`); - this.server.to(`user:${payload.userId}`).emit("message_deleted", payload); - } - - // ============================================================================ - // PRIVATE HELPER METHODS - // ============================================================================ - /** - * Cleans up user connection on disconnect - */ - private cleanupUserConnection(client: Socket, userInfo: ConnectedUserInfo): void { - // Remove from user sessions - const userSockets = this.userSessions.get(userInfo.userId); - if (userSockets) { - userSockets.delete(client.id); - if (userSockets.size === 0) { - this.userSessions.delete(userInfo.userId); - } - } - - // Leave session room if joined - if (userInfo.sessionId) { - client.leave(`session_${userInfo.sessionId}`); - - // Notify others in session - this.server.to(`session_${userInfo.sessionId}`).emit(WEBSOCKET_EVENTS.USER_LEFT, { - userId: userInfo.userId, - sessionId: userInfo.sessionId, - }); - } - } - - /** - * Registers a new user connection - */ - private registerUserConnection(clientId: string, userId: string): void { - this.connectedUsers.set(clientId, { userId }); - - if (!this.userSessions.has(userId)) { - this.userSessions.set(userId, new Set()); - } - this.userSessions.get(userId)!.add(clientId); - } - - public getUserConnections(userId: string): number { - return this.userConnections.get(userId)?.size || 0; - } - - public isUserConnected(userId: string): boolean { - return this.userConnections.has(userId) && this.userConnections.get(userId)!.size > 0; - } - - public getAllConnectedUsers(): string[] { - return Array.from(this.userConnections.keys()); - } -} diff --git a/src/modules/email/email.module.ts b/src/modules/email/email.module.ts index f899b2e..2ece2ba 100644 --- a/src/modules/email/email.module.ts +++ b/src/modules/email/email.module.ts @@ -3,12 +3,9 @@ import { Module } from "@nestjs/common"; import { JwtModule } from "@nestjs/jwt"; import { EmailController } from "./email.controller"; -import { EmailGateway } from "./email.gateway"; import { EmailHeadersService } from "./services/email-headers.service"; -import { EmailNotificationService } from "./services/email-notification.service"; import { EmailService } from "./services/email.service"; import { MailboxResolverService } from "./services/mailbox-resolver.service"; -import { WebSocketAuthService } from "./services/websocket-auth.service"; import { jwtConfig } from "../../configs/jwt.config"; import { MailServerModule } from "../mail-server/mail-server.module"; import { TemplatesModule } from "../templates/templates.module"; @@ -18,7 +15,7 @@ import { UsersModule } from "../users/users.module"; @Module({ imports: [MailServerModule, UsersModule, TemplatesModule, JwtModule.registerAsync(jwtConfig()), MikroOrmModule.forFeature([User])], controllers: [EmailController], - providers: [EmailService, MailboxResolverService, EmailGateway, EmailNotificationService, WebSocketAuthService, EmailHeadersService], - exports: [EmailService, MailboxResolverService, EmailGateway, EmailNotificationService, WebSocketAuthService, EmailHeadersService], + providers: [EmailService, MailboxResolverService, EmailHeadersService], + exports: [EmailService, MailboxResolverService, EmailHeadersService], }) export class EmailModule {} diff --git a/src/modules/email/interfaces/email-events.interface.ts b/src/modules/email/interfaces/email-events.interface.ts index a26d332..50ddf60 100644 --- a/src/modules/email/interfaces/email-events.interface.ts +++ b/src/modules/email/interfaces/email-events.interface.ts @@ -1,4 +1,4 @@ -export interface EmailNotificationPayload { +export interface EmailPayload { messageId: number; userId: string; from: { @@ -12,7 +12,7 @@ export interface EmailNotificationPayload { subject: string; preview?: string; hasAttachments: boolean; - received: string; + timestamp: string; mailboxId: string; mailboxName: string; isRead: boolean; @@ -20,10 +20,9 @@ export interface EmailNotificationPayload { size: number; } -export interface EmailStatusUpdatePayload { +export interface EmailStatusPayload { messageId: number; userId: string; - status: "read" | "unread" | "flagged" | "unflagged" | "archived" | "deleted"; mailboxId: string; timestamp: string; } @@ -42,28 +41,55 @@ export interface EmailDeletePayload { messageId: number; userId: string; mailboxId: string; - mailboxName: string; timestamp: string; } -export interface ClientJoinPayload { +export interface EmailSentPayload { + messageId: number; userId: string; - clientId: string; + to: Array<{ + name?: string; + address: string; + }>; + subject: string; + timestamp: string; + mailboxId: string; +} + +export interface MailboxUpdatePayload { + userId: string; + mailboxId: string; + mailboxName: string; + unreadCount: number; + totalCount: number; + timestamp: string; +} + +export interface UnreadCountPayload { + userId: string; + totalUnread: number; + mailboxCounts: Record; timestamp: string; } export interface EmailWebSocketEvents { - // Client to Server events - join_user_room: ClientJoinPayload; - leave_user_room: ClientJoinPayload; - // Server to Client events - new_message: EmailNotificationPayload; - message_status_update: EmailStatusUpdatePayload; - message_moved: EmailMovePayload; - message_deleted: EmailDeletePayload; - user_joined: ClientJoinPayload; - user_left: ClientJoinPayload; + new_email: EmailPayload; + email_read: EmailStatusPayload; + email_unread: EmailStatusPayload; + email_flagged: EmailStatusPayload; + email_unflagged: EmailStatusPayload; + email_deleted: EmailDeletePayload; + email_moved: EmailMovePayload; + email_sent: EmailSentPayload; + mailbox_updated: MailboxUpdatePayload; + unread_count_updated: UnreadCountPayload; + + // Connection events connection_established: { userId: string; timestamp: string }; connection_error: { error: string; timestamp: string }; + + // System events + ping: { timestamp: string }; + pong: { timestamp: string }; } diff --git a/src/modules/email/services/email-notification.service.ts b/src/modules/email/services/email-notification.service.ts deleted file mode 100644 index 43ae0b4..0000000 --- a/src/modules/email/services/email-notification.service.ts +++ /dev/null @@ -1,158 +0,0 @@ -import { Injectable, Logger } from "@nestjs/common"; -import { firstValueFrom } from "rxjs"; - -import { MailServerService } from "../../mail-server/services/mail-server.service"; -import { EmailGateway } from "../email.gateway"; -import { MailboxResolverService } from "./mailbox-resolver.service"; -import { EmailNotificationPayload } from "../interfaces/email-events.interface"; - -@Injectable() -export class EmailNotificationService { - private readonly logger = new Logger(EmailNotificationService.name); - - constructor( - private readonly emailGateway: EmailGateway, - private readonly mailServerService: MailServerService, - private readonly mailboxResolverService: MailboxResolverService, - ) {} - - /** - * Handle new incoming message and emit WebSocket notification - * This method can be called by webhook handlers, polling services, etc. - */ - async handleNewMessage(userId: string, messageId: number, mailboxId?: string) { - try { - this.logger.log(`Processing new message notification for user ${userId}, message ${messageId}`); - - // Get the mailbox ID if not provided - const targetMailboxId = mailboxId || (await this.mailboxResolverService.getInboxMailboxId(userId)); - - // Fetch the message details - const message = await firstValueFrom(this.mailServerService.messages.getMessage(userId, targetMailboxId, messageId)); - - // Determine mailbox name - const mailboxName = await this.getMailboxName(userId, targetMailboxId); - - // Create notification payload - const notificationPayload: EmailNotificationPayload = { - messageId, - userId, - from: { - name: message.from?.name || message.from?.address?.split("@")[0] || "Unknown", - address: message.from?.address || "unknown@example.com", - }, - to: - message.to?.map((recipient) => ({ - name: recipient.name || recipient.address?.split("@")[0] || "Unknown", - address: recipient.address || "unknown@example.com", - })) || [], - subject: message.subject || "No Subject", - preview: this.extractPreview( - typeof message.text === "string" ? message.text : Array.isArray(message.html) ? message.html.join("") : message.html || "", - ), - hasAttachments: Array.isArray(message.attachments) && message.attachments.length > 0, - received: message.idate || new Date().toISOString(), - mailboxId: targetMailboxId, - mailboxName, - isRead: message.seen || false, - isFlagged: message.flagged || false, - size: message.size || 0, - }; - - // Emit WebSocket notification - await this.emailGateway.notifyNewMessage(notificationPayload); - - this.logger.log(`New message notification sent for user ${userId}, message ${messageId}`); - - return notificationPayload; - } catch (error) { - this.logger.error( - `Failed to process new message notification for user ${userId}, message ${messageId}: ${ - error instanceof Error ? error.message : "Unknown error" - }`, - ); - throw error; - } - } - - /** - * Handle bulk new messages (useful for initial sync or batch processing) - */ - async handleBulkNewMessages(userId: string, messageIds: number[], mailboxId?: string) { - this.logger.log(`Processing bulk new message notifications for user ${userId}, ${messageIds.length} messages`); - - const results = []; - for (const messageId of messageIds) { - try { - const result = await this.handleNewMessage(userId, messageId, mailboxId); - results.push(result); - } catch (error) { - this.logger.error(`Failed to process message ${messageId} in bulk notification: ${error}`); - } - } - - return results; - } - - /** - * Simulate new message arrival (for testing purposes) - */ - async simulateNewMessage(userId: string, messageData: Partial) { - this.logger.log(`Simulating new message for user ${userId}`); - - const simulatedPayload: EmailNotificationPayload = { - messageId: Math.floor(Math.random() * 10000), - userId, - from: { - name: "Test Sender", - address: "test@example.com", - }, - to: [ - { - name: "Test Recipient", - address: "recipient@example.com", - }, - ], - subject: "Test Message", - preview: "This is a test message for WebSocket notification", - hasAttachments: false, - received: new Date().toISOString(), - mailboxId: "inbox", - mailboxName: "Inbox", - isRead: false, - isFlagged: false, - size: 1024, - ...messageData, - }; - - await this.emailGateway.notifyNewMessage(simulatedPayload); - - return simulatedPayload; - } - - private async getMailboxName(userId: string, mailboxId: string): Promise { - try { - const mailboxIds = await this.mailboxResolverService.getUserMailboxIds(userId); - - // Check which mailbox this ID corresponds to - if (mailboxIds.inbox === mailboxId) return "Inbox"; - if (mailboxIds.sent === mailboxId) return "Sent"; - if (mailboxIds.drafts === mailboxId) return "Drafts"; - if (mailboxIds.trash === mailboxId) return "Trash"; - if (mailboxIds.junk === mailboxId) return "Junk"; - if (mailboxIds.archive === mailboxId) return "Archive"; - if (mailboxIds.favorite === mailboxId) return "Favorite"; - - return "Unknown"; - } catch (error) { - this.logger.error(`Failed to get mailbox name for ${mailboxId}: ${error}`); - return "Unknown"; - } - } - - private extractPreview(content: string): string { - // Remove HTML tags and get first 150 characters - const textContent = content.replace(/<[^>]*>/g, "").trim(); - return textContent.length > 150 ? textContent.substring(0, 150) + "..." : textContent; - } -} diff --git a/src/modules/email/services/email.service.ts b/src/modules/email/services/email.service.ts index 48125e7..854ec34 100644 --- a/src/modules/email/services/email.service.ts +++ b/src/modules/email/services/email.service.ts @@ -12,9 +12,7 @@ import { UserRepository } from "../../users/repositories/user.repository"; import { BulkActionDto, BulkActionType } from "../DTO/bulk-actions.dto"; import { MessageListQueryDto, SearchMessagesQueryDto } from "../DTO/email-query.dto"; import { EmailType, SendEmailDto, UpdateDraftDto } from "../DTO/send-email.dto"; -import { EmailGateway } from "../email.gateway"; import { Priority } from "../enums/email-header.enum"; -import { EmailDeletePayload, EmailMovePayload, EmailStatusUpdatePayload } from "../interfaces/email-events.interface"; @Injectable() export class EmailService { @@ -23,7 +21,6 @@ export class EmailService { constructor( private readonly mailServerService: MailServerService, private readonly mailboxResolverService: MailboxResolverService, - private readonly emailGateway: EmailGateway, private readonly templateProcessorService: TemplateProcessorService, private readonly userRepository: UserRepository, private readonly emailHeadersService: EmailHeadersService, @@ -176,16 +173,6 @@ export class EmailService { const result = await firstValueFrom(this.mailServerService.messages.deleteMessage(userId, messageLocation.mailboxId, messageId)); - // Emit WebSocket notification - const payload: EmailDeletePayload = { - messageId, - userId, - mailboxId: messageLocation.mailboxId, - mailboxName: messageLocation.mailboxName, - timestamp: new Date().toISOString(), - }; - await this.emailGateway.notifyMessageDeleted(payload); - this.logger.log(`Message ${messageId} deleted successfully from ${messageLocation.mailboxName}`); return { @@ -212,16 +199,6 @@ export class EmailService { this.mailServerService.messages.updateMessage(userId, messageLocation.mailboxId, messageId, { seen: true }), ); - // Emit WebSocket notification - const payload: EmailStatusUpdatePayload = { - messageId, - userId, - status: "read", - mailboxId: messageLocation.mailboxId, - timestamp: new Date().toISOString(), - }; - await this.emailGateway.notifyMessageStatusUpdate(payload); - this.logger.log(`Message ${messageId} marked as seen successfully in ${messageLocation.mailboxName}`); return { @@ -481,18 +458,6 @@ export class EmailService { }), ); - // Emit WebSocket notification - const payload: EmailMovePayload = { - messageId, - userId, - fromMailboxId: messageLocation.mailboxId, - toMailboxId: archiveMailboxId, - fromMailboxName: messageLocation.mailboxName, - toMailboxName: MailboxEnum.ARCHIVE, - timestamp: new Date().toISOString(), - }; - await this.emailGateway.notifyMessageMoved(payload); - this.logger.log(`Message ${messageId} moved to archive successfully from ${messageLocation.mailboxName}`); return { @@ -527,18 +492,6 @@ export class EmailService { }), ); - // Emit WebSocket notification - const payload: EmailMovePayload = { - messageId, - userId, - fromMailboxId: messageLocation.mailboxId, - toMailboxId: favoriteMailboxId, - fromMailboxName: messageLocation.mailboxName, - toMailboxName: MailboxEnum.FAVORITE, - timestamp: new Date().toISOString(), - }; - await this.emailGateway.notifyMessageMoved(payload); - this.logger.log(`Message ${messageId} moved to favorite successfully from ${messageLocation.mailboxName}`); return { @@ -572,18 +525,6 @@ export class EmailService { }), ); - // Emit WebSocket notification - const payload: EmailMovePayload = { - messageId, - userId, - fromMailboxId: messageLocation.mailboxId, - toMailboxId: trashMailboxId, - fromMailboxName: messageLocation.mailboxName, - toMailboxName: MailboxEnum.TRASH, - timestamp: new Date().toISOString(), - }; - await this.emailGateway.notifyMessageMoved(payload); - this.logger.log(`Message ${messageId} moved to trash successfully from ${messageLocation.mailboxName}`); return { @@ -613,18 +554,6 @@ export class EmailService { }), ); - // Emit WebSocket notification - const payload: EmailMovePayload = { - messageId, - userId, - fromMailboxId: archiveMailboxId, - toMailboxId: inboxMailboxId, - fromMailboxName: "Archive", - toMailboxName: "Inbox", - timestamp: new Date().toISOString(), - }; - await this.emailGateway.notifyMessageMoved(payload); - this.logger.log(`Message ${messageId} moved from archive to inbox successfully`); return { @@ -654,18 +583,6 @@ export class EmailService { }), ); - // Emit WebSocket notification - const payload: EmailMovePayload = { - messageId, - userId, - fromMailboxId: favoriteMailboxId, - toMailboxId: inboxMailboxId, - fromMailboxName: "Favorite", - toMailboxName: "Inbox", - timestamp: new Date().toISOString(), - }; - await this.emailGateway.notifyMessageMoved(payload); - this.logger.log(`Message ${messageId} moved from favorite to inbox successfully`); return { @@ -695,18 +612,6 @@ export class EmailService { }), ); - // Emit WebSocket notification - const payload: EmailMovePayload = { - messageId, - userId, - fromMailboxId: trashMailboxId, - toMailboxId: inboxMailboxId, - fromMailboxName: "Trash", - toMailboxName: "Inbox", - timestamp: new Date().toISOString(), - }; - await this.emailGateway.notifyMessageMoved(payload); - this.logger.log(`Message ${messageId} moved from trash to inbox successfully`); return { @@ -741,18 +646,6 @@ export class EmailService { }), ); - // Emit WebSocket notification - const payload: EmailMovePayload = { - messageId, - userId, - fromMailboxId: messageLocation.mailboxId, - toMailboxId: junkMailboxId, - fromMailboxName: messageLocation.mailboxName, - toMailboxName: MailboxEnum.Junk, - timestamp: new Date().toISOString(), - }; - await this.emailGateway.notifyMessageMoved(payload); - this.logger.log(`Message ${messageId} moved to junk successfully from ${messageLocation.mailboxName}`); return { diff --git a/src/modules/email/services/websocket-auth.service.ts b/src/modules/email/services/websocket-auth.service.ts deleted file mode 100644 index 8682143..0000000 --- a/src/modules/email/services/websocket-auth.service.ts +++ /dev/null @@ -1,240 +0,0 @@ -import { Injectable, Logger } from "@nestjs/common"; -import { JwtService } from "@nestjs/jwt"; -import { Socket } from "socket.io"; - -import { WebSocketMessage } from "../../../common/enums/message.enum"; -import { WEBSOCKET_EVENTS, WEBSOCKET_ROOMS } from "../constants/email-events.constant"; -import { AuthenticatedSocket, WebSocketAuthResult, WebSocketConnectionContext, WebSocketUser } from "../interfaces/websocket.interface"; - -@Injectable() -export class WebSocketAuthService { - private readonly logger = new Logger(WebSocketAuthService.name); - private readonly connectionTimeouts = new Map(); - - constructor(private readonly jwtService: JwtService) {} - - /** - * Authenticates a WebSocket client using JWT token - */ - async authenticateClient(client: Socket): Promise { - try { - const token = this.extractToken(client); - - if (!token) { - return { - success: false, - error: WebSocketMessage.WS_AUTHENTICATION_REQUIRED, - }; - } - - const decoded = await this.verifyToken(token); - if (!decoded) { - return { - success: false, - error: WebSocketMessage.WS_INVALID_TOKEN, - }; - } - - const user = this.extractUserFromToken(decoded); - if (!user) { - return { - success: false, - error: WebSocketMessage.WS_INVALID_TOKEN, - }; - } - - // Set user data on socket - client.data = { - user, - connectedAt: new Date(), - lastActivity: new Date(), - sessionId: user.sessionId || this.generateSessionId(), - }; - - // Join user to their personal room - await client.join(WEBSOCKET_ROOMS.USER(user.id)); - - return { - success: true, - user, - }; - } catch (error) { - this.logger.error("WebSocket authentication error", { - clientId: client.id, - error: error instanceof Error ? error.message : String(error), - }); - - return { - success: false, - error: WebSocketMessage.WS_AUTHENTICATION_FAILED, - }; - } - } - - /** - * Handles authentication failure by emitting error and disconnecting client - */ - handleAuthenticationFailure(client: Socket, errorMessage: string): void { - const context: WebSocketConnectionContext = { - clientId: client.id, - clientIP: client.handshake.address, - userAgent: client.handshake.headers["user-agent"], - timestamp: new Date().toISOString(), - }; - - this.logger.warn("WebSocket authentication failed", { - ...context, - errorMessage, - }); - - // Emit error to client - client.emit(WEBSOCKET_EVENTS.CONNECTION_ERROR, { - error: errorMessage, - timestamp: new Date().toISOString(), - retryAllowed: this.shouldAllowRetry(client), - }); - - // Set timeout before disconnecting to allow client to handle error - const timeout = setTimeout(() => { - if (client.connected) { - client.disconnect(true); - } - this.connectionTimeouts.delete(client.id); - }, 1000); - - this.connectionTimeouts.set(client.id, timeout); - } - - /** - * Emits successful authentication event to client - */ - emitAuthenticationSuccess(client: Socket, user: WebSocketUser): void { - const authSocket = client as AuthenticatedSocket; - - client.emit(WEBSOCKET_EVENTS.CONNECTION_ESTABLISHED, { - userId: user.id, - sessionId: authSocket.data.sessionId, - timestamp: new Date().toISOString(), - permissions: user.permissions || [], - }); - - this.logger.log("WebSocket authentication successful", { - clientId: client.id, - userId: user.id, - sessionId: authSocket.data.sessionId, - }); - } - - /** - * Verifies if user has permission to access a specific room - */ - canAccessRoom(user: WebSocketUser, roomId: string): boolean { - // Basic room access control - if (roomId.startsWith("user_")) { - const targetUserId = roomId.replace("user_", ""); - return user.id === targetUserId || user.isAdmin === true; - } - - if (roomId.startsWith("session_")) { - return true; // Sessions are generally accessible - } - - if (roomId === "broadcast") { - return user.isAdmin === true; - } - - // Default deny for unknown room patterns - return false; - } - - /** - * Updates last activity timestamp for authenticated socket - */ - updateLastActivity(client: Socket): void { - if (client.data?.user) { - client.data.lastActivity = new Date(); - } - } - - /** - * Cleans up connection timeouts on disconnect - */ - cleanupConnection(client: Socket): void { - const timeout = this.connectionTimeouts.get(client.id); - if (timeout) { - clearTimeout(timeout); - this.connectionTimeouts.delete(client.id); - } - } - - // Private helper methods - - private extractToken(client: Socket): string | null { - // Try to get token from auth object first (recommended) - const authToken = client.handshake.auth?.token; - if (authToken && typeof authToken === "string") { - return authToken; - } - - // Fallback to query parameter - const queryToken = client.handshake.query?.token; - if (queryToken && typeof queryToken === "string") { - return queryToken; - } - - // Try headers as last resort - const headerToken = client.handshake.headers.authorization; - if (headerToken && typeof headerToken === "string") { - return headerToken.replace("Bearer ", ""); - } - - return null; - } - - private async verifyToken(token: string): Promise { - try { - return await this.jwtService.verifyAsync(token); - } catch (error) { - if (error instanceof Error) { - if (error.name === "TokenExpiredError") { - throw new Error(WebSocketMessage.WS_TOKEN_EXPIRED); - } - if (error.name === "JsonWebTokenError") { - throw new Error(WebSocketMessage.WS_INVALID_TOKEN); - } - if (error.name === "NotBeforeError") { - throw new Error(WebSocketMessage.WS_INVALID_TOKEN); - } - } - throw new Error(WebSocketMessage.WS_INVALID_TOKEN); - } - } - - private extractUserFromToken(decoded: any): WebSocketUser | null { - try { - return { - id: decoded.wildduckUserId || decoded.sub || decoded.id, - email: decoded.email, - permissions: decoded.permissions || [], - isAdmin: decoded.isAdmin || false, - businessId: decoded.businessId, - sessionId: decoded.sessionId, - }; - } catch (error) { - this.logger.error("Failed to extract user from token", { - error: error instanceof Error ? error.message : String(error), - }); - return null; - } - } - - private shouldAllowRetry(_client: Socket): boolean { - // Implement retry logic based on client IP or other factors - // For now, always allow retry - return true; - } - - private generateSessionId(): string { - return `ws_${Date.now()}_${Math.random().toString(36).substr(2, 9)}`; - } -}