-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathNonCausalMulticast.cpp
More file actions
240 lines (219 loc) · 5.72 KB
/
Copy pathNonCausalMulticast.cpp
File metadata and controls
240 lines (219 loc) · 5.72 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
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
#include <stdio.h>
#include <stdlib.h>
#include <unistd.h>
#include <pthread.h>
#include <iostream>
#include <sys/types.h>
#include <sys/time.h>
#include <sys/socket.h>
#include <netinet/in.h>
#include <arpa/inet.h>
#include <time.h>
#include <string.h>
#include <unistd.h>
#define PORT 12345
#define IP "225.0.0.37"
using namespace std;
// Variable declaration
void *sender(void *);
void *receiver(void *);
int multicast_sock, ptp_sock;
int logical_clock, process, p, m, send_clock, send_port;
struct sockaddr_in p_addr, port_master, m_addr, m_r_addr;
struct ip_mreq mr;
char * msg_input;
char mesg_bcast[256];
struct message_multicast{
int processid;
char rv_msg[256];
int arr[3];
};
struct message_multicast message_array[50];
int arr1[3] = {0,0,0};
int n = 3;
int main(int argc, char *argv[])
{
u_int on = 1;
char message1[256];
char message2[256];
char message3[256];
char message4[256];
char buffer1[256];
int m1, m2, m3;
int initial_value, drift, send_drift;
unsigned int rank = atoi(argv[1]);
msg_input = argv[2];
pthread_t sender_th, receiver_th;
//Multicast socket creation
multicast_sock = socket(AF_INET, SOCK_DGRAM, 0);
if(multicast_sock < 0)
{
cout << "Cannot create socket" << endl;
return 0;
}
//Assigning destination address
memset((char*) &m_addr, 0, sizeof(m_addr));
m_addr.sin_family = AF_INET;
m_addr.sin_addr.s_addr = inet_addr(IP);
m_addr.sin_port = htons(PORT);
//set up destination address on receiver side
memset((char*) &m_r_addr, 0, sizeof(m_r_addr));
m_r_addr.sin_family = AF_INET;
m_r_addr.sin_addr.s_addr =htonl(INADDR_ANY);
m_r_addr.sin_port = htons(PORT);
//Allowing multiple sockets to use the same PORT number
if(setsockopt(multicast_sock, SOL_SOCKET, SO_REUSEADDR, &on, sizeof(on)) < 0)
{
cout << "Allowing multiple sockets to use the same PORT number failed" << endl;
return 0;
}
//Binding to destination address
m = bind(multicast_sock, (struct sockaddr *) &m_r_addr, sizeof(m_r_addr));
if(m < 0)
{
cout << "Cannot bind" << endl;
return 0;
}
//use setsockopt() to request the kernel to join the multicast group
mr.imr_multiaddr.s_addr = inet_addr(IP);
mr.imr_interface.s_addr = htonl(INADDR_ANY);
if(setsockopt(multicast_sock, IPPROTO_IP, IP_ADD_MEMBERSHIP, &mr, sizeof(mr)) < 0)
{
cout << "Error in setsockopt" << endl;
return 0;
}
//creating point to point socket
ptp_sock = socket(AF_INET, SOCK_DGRAM, 0);
//Assign destination address
p_addr.sin_family = AF_INET;
p_addr.sin_addr.s_addr = INADDR_ANY;
p_addr.sin_port = htons(INADDR_ANY);
//Request to use same address
if(setsockopt(ptp_sock, SOL_SOCKET, SO_REUSEADDR, &on, sizeof(on)) < 0)
{
cout << "Cannot reuse address" << endl;
return 0;
}
//Binding point to point socket
p = bind(ptp_sock, (struct sockaddr *) &p_addr, sizeof(p_addr));
if(p < 0)
{
cout << "Cannot bind point to point socket" << endl;
return 0;
}
struct timeval timeout;
timeout.tv_sec = 5;
timeout.tv_usec = 0;
if(setsockopt(ptp_sock,SOL_SOCKET,SO_RCVTIMEO,(char *)&timeout,sizeof(timeout))<0)
{
cout << "Timeout eror" << endl;
}
//Create threads for sender and reciever
if(pthread_create(&sender_th,NULL,sender,&rank)<0)
{
cout << "Cannot create a thread for sender" << endl;
return 1;
}
if(pthread_create(&receiver_th,NULL,receiver,&rank)<0)
{
cout << "Cannot create a thread for receiver" << endl;
return 1;
}
pthread_join(sender_th,NULL);
pthread_join(receiver_th,NULL);
return 0;
}
//thread for sending the messages to all the processes
void *sender(void *r)
{
unsigned int rank1 = *((unsigned int *) r);
int p1 = 0, p2 = 0, p3 = 0, p4 = 0;
int process;
cout << "Rank: " << rank1 << endl;
cout << "Message Input: " << msg_input << endl;
while(1)
{
for (int i = 0; i < 4; i++)
{
if(i == rank1)
{
if(i == 0)
{
p1++;
process = 1;
cout << "Process: " << process;
arr1[i] = p1;
}
if(i == 1)
{
p2++;
process = 2;
cout << "Process: " << process;
arr1[i] = p2;
}
if(i == 2)
{
p3++;
process = 3;
cout << "Process: " << process;
arr1[i] = p3;
}
if(i == 3)
{
p4++;
process = 4;
cout << "Process: " << process;
arr1[i] = p4;
}
}
}
char buffer1[256];
sprintf(buffer1,"%d,%d,%d,%d,%d,%s", process, arr1[0], arr1[1], arr1[2], arr1[3], msg_input);
cout << "Message: " << buffer1 << endl;
int send1 = sendto(multicast_sock, buffer1, sizeof(buffer1), 0, (struct sockaddr *) &m_addr, sizeof(m_addr));
sleep(10);
}
return 0;
}
//Thread for receiving the messages from all the processes
void *receiver(void *r1)
{
char buffer2[256];
int count = 0;
int processid;
char rv_msg[256];
int arr[4];
int num = 0;
int varr1[4] = {0,0,0,0}, varr2[4] = {0,0,0,0}, varr3[4] = {0,0,0,0}, varr4[4] = {0,0,0,0};
char *variable, *message, *arr2[4];
socklen_t addr=sizeof(m_r_addr);
while(recvfrom(multicast_sock, buffer2, sizeof(buffer2), 0, (struct sockaddr *) &m_r_addr, &addr) >= 0)
{
count++;
cout << endl ;
cout << "----------Non Causal ordering----------"<< endl;
cout << "Message received: " << buffer2 << endl;
cout << "Count: " << count << endl;
message = strtok(buffer2, ",");
if(message != NULL)
{
processid = atoi(message);
cout << "Processid: " << processid << endl;
message = strtok(NULL, ",");
arr[0] = atoi(message);
message = strtok(NULL, ",");
arr[1] = atoi(message);
message = strtok(NULL, ",");
arr[2] = atoi(message);
message = strtok(NULL, ",");
arr[3] = atoi(message);
message = strtok(NULL, " ,");
strcpy(rv_msg, message);
cout << endl ;
cout << "Message: " << rv_msg << endl;
cout << "******Vector value: " << arr[0] << "," << arr[1] <<"," << arr[2]<<"," << arr[3] << endl;
}
sleep(10);
bzero(buffer2, 256);
}
}