- Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathClientHandler.java
More file actions
Latest commit
72 lines (56 loc) · 2.15 KB
/
Copy pathClientHandler.java
File metadata and controls
72 lines (56 loc) · 2.15 KB
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
importcom.google.gson.Gson;
importcom.google.gson.GsonBuilder;
importjava.io.BufferedInputStream;
importjava.io.ByteArrayOutputStream;
importjava.io.IOException;
importjava.net.Socket;
/**
* Created by sameerraghuram on 4/25/17.
*/
publicclassClientHandlerimplementsRunnable {
SocketclientSocket;
Callbackscallbacks;
publicClientHandler(SocketclientSocket){
this.clientSocket = clientSocket;
this.callbacks = Callbacks.getInstance();
}
@Override
publicvoidrun() {
// Buffer
byte[] buf = newbyte[1024];
try {
BufferedInputStreambis = newBufferedInputStream(clientSocket.getInputStream());
ByteArrayOutputStreambyos = newByteArrayOutputStream();
// Read bytes from the socket input stream in 1024 increments
// Untill we reach EOF
while(bis.read(buf, 0, buf.length) != -1){
// Write bytes to a op stream
byos.write(buf, 0, buf.length);
}
// Capture the byte array as a string
Stringmessage = byos.toString("utf-8").trim();
// Convert String to JSON
Gsongson = newGsonBuilder().setPrettyPrinting().setLenient().create();
System.out.println(message);
MessagemessageJSON = gson.fromJson(message, Message.class);
if(messageJSON.HEADER.REQUEST.equals(Consts.REQ_ADD_MESSGAE)){
PublishMessagepublishMessage = gson.fromJson(message, PublishMessage.class);
StringqueueName = publishMessage.QUEUE_NAME;
StringexchangeName = publishMessage.EXCHANGE_NAME;
Stringkey = exchangeName.equals("default") ? queueName : exchangeName + "." + queueName;
callbacks.exec(key, publishMessage.PUBLISH_MESSAGE);
}
} catch (IOExceptione) {
e.printStackTrace();
} catch (Exceptione) {
e.printStackTrace();
}
finally {
try {
clientSocket.close();
} catch (IOExceptione) {
e.printStackTrace();
}
}
}
}