Spaces:
Running
Running
File size: 3,257 Bytes
e8a6607 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 | import { Server as SocketIOServer, Socket } from 'socket.io';
import { Server as HTTPSServer } from 'https';
import http, { Server as HttpServer } from "http";
import dotenv from 'dotenv';
import logger from '../../logger';
import { createAdapter } from "@socket.io/redis-adapter";
import { getPubSubRedisInfra } from '../../redis';
import { createClient } from "redis";
const envFile: string = process.env.NODE_ENV === 'production' ? '.env.production' : '.env.development';
dotenv.config({ path: envFile });
const jwt_access_secret = process.env.JWT_ACCESS_SECRET || '';
let io: SocketIOServer;
const allowedOrigins = [
process.env.FRONTEND_URL,
];
export const initializeSocket = async (server: HTTPSServer | HttpServer): Promise<SocketIOServer> => {
io = new SocketIOServer(server, {
cors: {
origin: allowedOrigins[0],
methods: ['GET', 'POST'],
credentials: true || false,
}
});
// const pubRedisInfra = getPubSubRedisInfra('pub-redis-infra');
// const subRedisInfra = getPubSubRedisInfra('sub-redis-infra');
// // For kafka service emit event using a redis client to listen
// io.adapter(createAdapter(pubRedisInfra.redis, subRedisInfra.redis));
const subscriber = createClient();
await subscriber.connect();
subscriber.subscribe("ai_response", (message: any) => { // This event is generating from python ai service
const data = JSON.parse(message);
// data.event = "gptChatRes"
io.to(`user:${data.userId}`).emit(data.event, {
result: data.payload.result,
msgSession: data.msg_session
}
);
});
io.on('connection', (socket: Socket) => { // in-built listener
logger.debug(`β
A user connected: ${socket.id}`);
socket.on('disconnect', async() => { // in-built listener
logger.debug(`β
A user disconnected: ${socket.id}`);
});
socket.on('admin-login', async(adminEmail: string) => { // Custom listener
logger.info(`β
Admin logined : ${adminEmail}`); //
});
socket.on('admin-logout', async(adminEmail: string) => { // Custom listener
// removeAdmin(socket.id);
logger.info(`β
Admin logged out : ${adminEmail}`); //
});
socket.on('join-user-room', (userId: string) => {
socket.join(`user:${userId}`);
logger.debug(`β
A user joined the socket user-room: ${userId}`);
});
socket.on('leave-user-room', (userId: string) => {
socket.leave(`user:${userId}`);
logger.debug(`β
A user left the socket user-room: ${userId}`);
});
//Test function
socket.on('hello', (callback) => {
logger.debug('Received hello');
callback('world');
});
});
return io;
};
export const broadcastGptChatRes = async (
gptChatRes: any
): Promise<void> => {
if (!io) {
logger.error('Socket.io is not initialized');
return;
}
logger.debug('gptChatRes');
console.log("Emitting:", `user:${gptChatRes.userId}`);
io.to(`user:${gptChatRes.userId}`).emit('gptChatRes', { gptChatRes: gptChatRes });
}; |