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 });
};