chore: fix
This commit is contained in:
@@ -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<Socket>();
|
|
||||||
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;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -3,7 +3,6 @@ import { Logger, ValidationPipe } from "@nestjs/common";
|
|||||||
import { ConfigService } from "@nestjs/config";
|
import { ConfigService } from "@nestjs/config";
|
||||||
import { NestFactory } from "@nestjs/core";
|
import { NestFactory } from "@nestjs/core";
|
||||||
import { FastifyAdapter, NestFastifyApplication } from "@nestjs/platform-fastify";
|
import { FastifyAdapter, NestFastifyApplication } from "@nestjs/platform-fastify";
|
||||||
import { IoAdapter } from "@nestjs/platform-socket.io";
|
|
||||||
|
|
||||||
import { AppModule } from "./app.module";
|
import { AppModule } from "./app.module";
|
||||||
import { getSwaggerDocument } from "./configs/swagger.config";
|
import { getSwaggerDocument } from "./configs/swagger.config";
|
||||||
@@ -26,9 +25,6 @@ async function bootstrap() {
|
|||||||
|
|
||||||
const app = await NestFactory.create<NestFastifyApplication>(AppModule, fastifyAdapter);
|
const app = await NestFactory.create<NestFastifyApplication>(AppModule, fastifyAdapter);
|
||||||
|
|
||||||
// Setup WebSocket adapter
|
|
||||||
app.useWebSocketAdapter(new IoAdapter(app));
|
|
||||||
|
|
||||||
app.useGlobalPipes(new ValidationPipe({ transform: true, whitelist: true }));
|
app.useGlobalPipes(new ValidationPipe({ transform: true, whitelist: true }));
|
||||||
|
|
||||||
app.useGlobalInterceptors(new ResponseInterceptor(), new PaginationInterceptor());
|
app.useGlobalInterceptors(new ResponseInterceptor(), new PaginationInterceptor());
|
||||||
|
|||||||
@@ -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<boolean> {
|
|
||||||
try {
|
|
||||||
const client: Socket = context.switchToWs().getClient<Socket>();
|
|
||||||
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<boolean> {
|
|
||||||
// try {
|
|
||||||
// const client: Socket = context.switchToWs().getClient<Socket>();
|
|
||||||
// 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<boolean> {
|
|
||||||
// try {
|
|
||||||
// const client: Socket = context.switchToWs().getClient<Socket>();
|
|
||||||
// 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);
|
|
||||||
// }
|
|
||||||
// }
|
|
||||||
@@ -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,
|
|
||||||
});
|
|
||||||
```
|
|
||||||
@@ -2,46 +2,35 @@ export const WEBSOCKET_EVENTS = {
|
|||||||
// Connection Events
|
// Connection Events
|
||||||
CONNECTION_ESTABLISHED: "connection_established",
|
CONNECTION_ESTABLISHED: "connection_established",
|
||||||
CONNECTION_ERROR: "connection_error",
|
CONNECTION_ERROR: "connection_error",
|
||||||
USER_JOINED: "user_joined",
|
|
||||||
USER_LEFT: "user_left",
|
|
||||||
|
|
||||||
// Email Events - Server to Client
|
// Email Events - Server to Client
|
||||||
NEW_MESSAGE: "new_message",
|
NEW_EMAIL: "new_email",
|
||||||
MESSAGE_STATUS_UPDATE: "message_status_update",
|
EMAIL_READ: "email_read",
|
||||||
MESSAGE_MOVED: "message_moved",
|
EMAIL_UNREAD: "email_unread",
|
||||||
MESSAGE_DELETED: "message_deleted",
|
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
|
// Mailbox Events
|
||||||
JOIN_USER_ROOM: "join_user_room",
|
MAILBOX_UPDATED: "mailbox_updated",
|
||||||
LEAVE_USER_ROOM: "leave_user_room",
|
UNREAD_COUNT_UPDATED: "unread_count_updated",
|
||||||
JOIN_SESSION: "join_session",
|
|
||||||
LEAVE_SESSION: "leave_session",
|
|
||||||
|
|
||||||
// Authentication Events
|
|
||||||
AUTHENTICATE: "authenticate",
|
|
||||||
AUTHENTICATION_SUCCESS: "authentication_success",
|
|
||||||
AUTHENTICATION_FAILED: "authentication_failed",
|
|
||||||
|
|
||||||
// System Events
|
// System Events
|
||||||
PING: "ping",
|
PING: "ping",
|
||||||
PONG: "pong",
|
PONG: "pong",
|
||||||
RECONNECT: "reconnect",
|
|
||||||
FORCE_DISCONNECT: "force_disconnect",
|
|
||||||
} as const;
|
} as const;
|
||||||
|
|
||||||
export const WEBSOCKET_ROOMS = {
|
export const WEBSOCKET_ROOMS = {
|
||||||
USER: (userId: string) => `user_${userId}`,
|
USER: (userId: string) => `user:${userId}`,
|
||||||
SESSION: (sessionId: string) => `session_${sessionId}`,
|
|
||||||
BROADCAST: "broadcast",
|
|
||||||
} as const;
|
} as const;
|
||||||
|
|
||||||
export const WEBSOCKET_NAMESPACES = {
|
export const WEBSOCKET_NAMESPACES = {
|
||||||
EMAIL: "/email",
|
EMAIL: "/email",
|
||||||
NOTIFICATIONS: "/notifications",
|
|
||||||
ADMIN: "/admin",
|
|
||||||
} as const;
|
} as const;
|
||||||
|
|
||||||
// Type definitions for better TypeScript support
|
// Type definitions
|
||||||
export type WebSocketEvent = (typeof WEBSOCKET_EVENTS)[keyof typeof WEBSOCKET_EVENTS];
|
export type WebSocketEvent = (typeof WEBSOCKET_EVENTS)[keyof typeof WEBSOCKET_EVENTS];
|
||||||
export type WebSocketRoom = string;
|
export type WebSocketRoom = string;
|
||||||
export type WebSocketNamespace = (typeof WEBSOCKET_NAMESPACES)[keyof typeof WEBSOCKET_NAMESPACES];
|
export type WebSocketNamespace = (typeof WEBSOCKET_NAMESPACES)[keyof typeof WEBSOCKET_NAMESPACES];
|
||||||
|
|||||||
@@ -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<string, Set<string>> = new Map(); // userId -> Set of socketIds
|
|
||||||
private readonly connectedUsers = new Map<string, ConnectedUserInfo>();
|
|
||||||
private readonly userSessions = new Map<string, Set<string>>();
|
|
||||||
|
|
||||||
@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<void> {
|
|
||||||
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());
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -3,12 +3,9 @@ import { Module } from "@nestjs/common";
|
|||||||
import { JwtModule } from "@nestjs/jwt";
|
import { JwtModule } from "@nestjs/jwt";
|
||||||
|
|
||||||
import { EmailController } from "./email.controller";
|
import { EmailController } from "./email.controller";
|
||||||
import { EmailGateway } from "./email.gateway";
|
|
||||||
import { EmailHeadersService } from "./services/email-headers.service";
|
import { EmailHeadersService } from "./services/email-headers.service";
|
||||||
import { EmailNotificationService } from "./services/email-notification.service";
|
|
||||||
import { EmailService } from "./services/email.service";
|
import { EmailService } from "./services/email.service";
|
||||||
import { MailboxResolverService } from "./services/mailbox-resolver.service";
|
import { MailboxResolverService } from "./services/mailbox-resolver.service";
|
||||||
import { WebSocketAuthService } from "./services/websocket-auth.service";
|
|
||||||
import { jwtConfig } from "../../configs/jwt.config";
|
import { jwtConfig } from "../../configs/jwt.config";
|
||||||
import { MailServerModule } from "../mail-server/mail-server.module";
|
import { MailServerModule } from "../mail-server/mail-server.module";
|
||||||
import { TemplatesModule } from "../templates/templates.module";
|
import { TemplatesModule } from "../templates/templates.module";
|
||||||
@@ -18,7 +15,7 @@ import { UsersModule } from "../users/users.module";
|
|||||||
@Module({
|
@Module({
|
||||||
imports: [MailServerModule, UsersModule, TemplatesModule, JwtModule.registerAsync(jwtConfig()), MikroOrmModule.forFeature([User])],
|
imports: [MailServerModule, UsersModule, TemplatesModule, JwtModule.registerAsync(jwtConfig()), MikroOrmModule.forFeature([User])],
|
||||||
controllers: [EmailController],
|
controllers: [EmailController],
|
||||||
providers: [EmailService, MailboxResolverService, EmailGateway, EmailNotificationService, WebSocketAuthService, EmailHeadersService],
|
providers: [EmailService, MailboxResolverService, EmailHeadersService],
|
||||||
exports: [EmailService, MailboxResolverService, EmailGateway, EmailNotificationService, WebSocketAuthService, EmailHeadersService],
|
exports: [EmailService, MailboxResolverService, EmailHeadersService],
|
||||||
})
|
})
|
||||||
export class EmailModule {}
|
export class EmailModule {}
|
||||||
|
|||||||
@@ -1,4 +1,4 @@
|
|||||||
export interface EmailNotificationPayload {
|
export interface EmailPayload {
|
||||||
messageId: number;
|
messageId: number;
|
||||||
userId: string;
|
userId: string;
|
||||||
from: {
|
from: {
|
||||||
@@ -12,7 +12,7 @@ export interface EmailNotificationPayload {
|
|||||||
subject: string;
|
subject: string;
|
||||||
preview?: string;
|
preview?: string;
|
||||||
hasAttachments: boolean;
|
hasAttachments: boolean;
|
||||||
received: string;
|
timestamp: string;
|
||||||
mailboxId: string;
|
mailboxId: string;
|
||||||
mailboxName: string;
|
mailboxName: string;
|
||||||
isRead: boolean;
|
isRead: boolean;
|
||||||
@@ -20,10 +20,9 @@ export interface EmailNotificationPayload {
|
|||||||
size: number;
|
size: number;
|
||||||
}
|
}
|
||||||
|
|
||||||
export interface EmailStatusUpdatePayload {
|
export interface EmailStatusPayload {
|
||||||
messageId: number;
|
messageId: number;
|
||||||
userId: string;
|
userId: string;
|
||||||
status: "read" | "unread" | "flagged" | "unflagged" | "archived" | "deleted";
|
|
||||||
mailboxId: string;
|
mailboxId: string;
|
||||||
timestamp: string;
|
timestamp: string;
|
||||||
}
|
}
|
||||||
@@ -42,28 +41,55 @@ export interface EmailDeletePayload {
|
|||||||
messageId: number;
|
messageId: number;
|
||||||
userId: string;
|
userId: string;
|
||||||
mailboxId: string;
|
mailboxId: string;
|
||||||
mailboxName: string;
|
|
||||||
timestamp: string;
|
timestamp: string;
|
||||||
}
|
}
|
||||||
|
|
||||||
export interface ClientJoinPayload {
|
export interface EmailSentPayload {
|
||||||
|
messageId: number;
|
||||||
userId: string;
|
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<string, number>;
|
||||||
timestamp: string;
|
timestamp: string;
|
||||||
}
|
}
|
||||||
|
|
||||||
export interface EmailWebSocketEvents {
|
export interface EmailWebSocketEvents {
|
||||||
// Client to Server events
|
|
||||||
join_user_room: ClientJoinPayload;
|
|
||||||
leave_user_room: ClientJoinPayload;
|
|
||||||
|
|
||||||
// Server to Client events
|
// Server to Client events
|
||||||
new_message: EmailNotificationPayload;
|
new_email: EmailPayload;
|
||||||
message_status_update: EmailStatusUpdatePayload;
|
email_read: EmailStatusPayload;
|
||||||
message_moved: EmailMovePayload;
|
email_unread: EmailStatusPayload;
|
||||||
message_deleted: EmailDeletePayload;
|
email_flagged: EmailStatusPayload;
|
||||||
user_joined: ClientJoinPayload;
|
email_unflagged: EmailStatusPayload;
|
||||||
user_left: ClientJoinPayload;
|
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_established: { userId: string; timestamp: string };
|
||||||
connection_error: { error: string; timestamp: string };
|
connection_error: { error: string; timestamp: string };
|
||||||
|
|
||||||
|
// System events
|
||||||
|
ping: { timestamp: string };
|
||||||
|
pong: { timestamp: string };
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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<EmailNotificationPayload>) {
|
|
||||||
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<string> {
|
|
||||||
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;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -12,9 +12,7 @@ import { UserRepository } from "../../users/repositories/user.repository";
|
|||||||
import { BulkActionDto, BulkActionType } from "../DTO/bulk-actions.dto";
|
import { BulkActionDto, BulkActionType } from "../DTO/bulk-actions.dto";
|
||||||
import { MessageListQueryDto, SearchMessagesQueryDto } from "../DTO/email-query.dto";
|
import { MessageListQueryDto, SearchMessagesQueryDto } from "../DTO/email-query.dto";
|
||||||
import { EmailType, SendEmailDto, UpdateDraftDto } from "../DTO/send-email.dto";
|
import { EmailType, SendEmailDto, UpdateDraftDto } from "../DTO/send-email.dto";
|
||||||
import { EmailGateway } from "../email.gateway";
|
|
||||||
import { Priority } from "../enums/email-header.enum";
|
import { Priority } from "../enums/email-header.enum";
|
||||||
import { EmailDeletePayload, EmailMovePayload, EmailStatusUpdatePayload } from "../interfaces/email-events.interface";
|
|
||||||
|
|
||||||
@Injectable()
|
@Injectable()
|
||||||
export class EmailService {
|
export class EmailService {
|
||||||
@@ -23,7 +21,6 @@ export class EmailService {
|
|||||||
constructor(
|
constructor(
|
||||||
private readonly mailServerService: MailServerService,
|
private readonly mailServerService: MailServerService,
|
||||||
private readonly mailboxResolverService: MailboxResolverService,
|
private readonly mailboxResolverService: MailboxResolverService,
|
||||||
private readonly emailGateway: EmailGateway,
|
|
||||||
private readonly templateProcessorService: TemplateProcessorService,
|
private readonly templateProcessorService: TemplateProcessorService,
|
||||||
private readonly userRepository: UserRepository,
|
private readonly userRepository: UserRepository,
|
||||||
private readonly emailHeadersService: EmailHeadersService,
|
private readonly emailHeadersService: EmailHeadersService,
|
||||||
@@ -176,16 +173,6 @@ export class EmailService {
|
|||||||
|
|
||||||
const result = await firstValueFrom(this.mailServerService.messages.deleteMessage(userId, messageLocation.mailboxId, messageId));
|
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}`);
|
this.logger.log(`Message ${messageId} deleted successfully from ${messageLocation.mailboxName}`);
|
||||||
|
|
||||||
return {
|
return {
|
||||||
@@ -212,16 +199,6 @@ export class EmailService {
|
|||||||
this.mailServerService.messages.updateMessage(userId, messageLocation.mailboxId, messageId, { seen: true }),
|
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}`);
|
this.logger.log(`Message ${messageId} marked as seen successfully in ${messageLocation.mailboxName}`);
|
||||||
|
|
||||||
return {
|
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}`);
|
this.logger.log(`Message ${messageId} moved to archive successfully from ${messageLocation.mailboxName}`);
|
||||||
|
|
||||||
return {
|
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}`);
|
this.logger.log(`Message ${messageId} moved to favorite successfully from ${messageLocation.mailboxName}`);
|
||||||
|
|
||||||
return {
|
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}`);
|
this.logger.log(`Message ${messageId} moved to trash successfully from ${messageLocation.mailboxName}`);
|
||||||
|
|
||||||
return {
|
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`);
|
this.logger.log(`Message ${messageId} moved from archive to inbox successfully`);
|
||||||
|
|
||||||
return {
|
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`);
|
this.logger.log(`Message ${messageId} moved from favorite to inbox successfully`);
|
||||||
|
|
||||||
return {
|
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`);
|
this.logger.log(`Message ${messageId} moved from trash to inbox successfully`);
|
||||||
|
|
||||||
return {
|
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}`);
|
this.logger.log(`Message ${messageId} moved to junk successfully from ${messageLocation.mailboxName}`);
|
||||||
|
|
||||||
return {
|
return {
|
||||||
|
|||||||
@@ -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<string, NodeJS.Timeout>();
|
|
||||||
|
|
||||||
constructor(private readonly jwtService: JwtService) {}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Authenticates a WebSocket client using JWT token
|
|
||||||
*/
|
|
||||||
async authenticateClient(client: Socket): Promise<WebSocketAuthResult> {
|
|
||||||
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<any> {
|
|
||||||
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)}`;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
Reference in New Issue
Block a user