-
Notifications
You must be signed in to change notification settings - Fork 0
/
Copy pathserver.js
94 lines (75 loc) · 2.3 KB
/
server.js
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
const express = require("express");
const { WebSocketServer } = require("ws");
const {v4: uuidv4} = require('uuid')
const {Kafka, Partitioners} = require('kafkajs')
require('dotenv').config()
const app = express()
const PORT = process.env.PORT || 3000
console.log('PORT: ', PORT)
//KAFKA
const kafka = new Kafka({brokers: ['localhost:9092']})
const producer = kafka.producer({
createPartitioner: Partitioners.LegacyPartitioner
})
const consumer = kafka.consumer({ groupId: "chat-group" });
(async () => {
try {
await producer.connect();
await consumer.connect();
await consumer.subscribe({ topic: "chat-messages", fromBeginning: true });
// Start consuming messages
await consumer.run({
eachMessage: async ({ topic, partition, message }) => {
const data = JSON.parse(message.value.toString());
console.log(`Received from Kafka: ${JSON.stringify(data)}`);
broadcast(data, null); // Pass null as the senderSocket
},
});
} catch (err) {
console.error("Error with Kafka:", err);
}
})();
// END KAFKA
app.use(express.static('public'))
const server = app.listen(PORT, () => {
console.log(`Server is running on http://localhost:${PORT}`);
});
const wss = new WebSocketServer({server})
const clients = new Map()
function generateClientName() {
return `User-${uuidv4().split('-')[0]}`
}
function broadcast(message, senderSocket) {
clients.forEach((_, client) => {
if(client !== senderSocket && client.readyState === client.OPEN){
client.send(message)
}
})
}
wss.on('connection', (ws) => {
const clientName = generateClientName()
clients.set(ws, clientName)
console.log(`${clientName} connected`)
ws.send(JSON.stringify({type: 'welcome', name: clientName}))
ws.on('message', async (message) => {
const senderName = clients.get(ws)
const broadcastMessage = JSON.stringify({
type: 'message',
name: senderName,
content: message.toString()
})
try{
await producer.send({
topic: 'chat-messages',
messages: [{ value: broadcastMessage}]
})
}catch(err){
console.error("Error sending message to Kafka:", err);
}
// broadcast(broadcastMessage, ws)
}) //Broadcast to all other clients
ws.on('close', () => {
clients.delete(ws)
console.log('Client disconnected')
})
})