-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathsubscriber.cpp
More file actions
130 lines (103 loc) · 4.07 KB
/
Copy pathsubscriber.cpp
File metadata and controls
130 lines (103 loc) · 4.07 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
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
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
#include "helpers.h"
using namespace std;
void usage(char *file)
{
fprintf(stderr, "Usage: %s <ID_CLIENT> <IP_SERVER> <PORT_SERVER>\n", file);
exit(0);
}
int main(int argc, char *argv[])
{
if (argc != 4) {
usage(argv[0]);
}
int sockfd, n, ret;
struct sockaddr_in serv_addr;
char buffer[BUFLEN];
// Disable buffering for stdout
setvbuf(stdout, NULL, _IONBF, BUFSIZ);
sockfd = socket(AF_INET, SOCK_STREAM, 0);
DIE(sockfd < 0, "socket");
// Enable Nagle's algorithm (default) - no need to set TCP_NODELAY as per typical requirements unless specified
serv_addr.sin_family = AF_INET;
serv_addr.sin_port = htons(atoi(argv[3]));
ret = inet_aton(argv[2], &serv_addr.sin_addr);
DIE(ret == 0, "inet_aton");
ret = connect(sockfd, (struct sockaddr*) &serv_addr, sizeof(serv_addr));
DIE(ret < 0, "connect");
// Send client ID immediately after connection
ret = send(sockfd, argv[1], strlen(argv[1]), 0);
DIE(ret < 0, "send");
struct pollfd fds[2];
// standard input
fds[0].fd = STDIN_FILENO;
fds[0].events = POLLIN;
// server socket
fds[1].fd = sockfd;
fds[1].events = POLLIN;
while (1) {
ret = poll(fds, 2, -1);
DIE(ret < 0, "poll");
if (fds[0].revents & POLLIN) {
// Read from stdin
memset(buffer, 0, BUFLEN);
// fgets(buffer, BUFLEN - 1, stdin);
int n_read = read(STDIN_FILENO, buffer, BUFLEN - 1);
if (n_read <= 0) break; // Error or EOF
if (strncmp(buffer, "exit", 4) == 0) {
break;
}
// Send command to server
n = send(sockfd, buffer, strlen(buffer), 0);
DIE(n < 0, "send");
// Provide feedback
if (strncmp(buffer, "subscribe", 9) == 0) {
// Check if it's "subscribe" or "unsubscribe" (unsubscribe has 'u' at index 0, handled below or strict check?)
// unsubscribe also starts with u.. wait. "subscribe" is inside "unsubscribe" string but strncmp checks prefix.
// "unsubscribe" vs "subscribe".
// "unsubscribe" does NOT match strncmp(.., "subscribe", 9) because 'u' != 's'.
// Parse topic to print
char temp[BUFLEN];
strcpy(temp, buffer);
// remove newline
string cmd(temp);
if (!cmd.empty() && cmd.back() == '\n') cmd.pop_back();
// split
size_t sp = cmd.find(' ');
if(sp != string::npos) {
string topic = cmd.substr(sp+1); // rest is topic (can be wildcards)
printf("Subscribed to topic %s\n", topic.c_str());
fflush(stdout);
}
} else if (strncmp(buffer, "unsubscribe", 11) == 0) {
char temp[BUFLEN];
strcpy(temp, buffer);
// remove newline
string cmd(temp);
if (!cmd.empty() && cmd.back() == '\n') cmd.pop_back();
// split
size_t sp = cmd.find(' ');
if(sp != string::npos) {
string topic = cmd.substr(sp+1);
printf("Unsubscribed from topic %s\n", topic.c_str());
fflush(stdout);
}
}
}
if (fds[1].revents & POLLIN) {
// Read from server
memset(buffer, 0, BUFLEN);
n = recv(sockfd, buffer, BUFLEN, 0);
DIE(n < 0, "recv");
if (n == 0) {
// Server closed connection
break;
}
// Print message received from server
printf("%s", buffer);
fflush(stdout);
}
}
shutdown(sockfd, SHUT_RDWR);
close(sockfd);
return 0;
}