· 8 years ago · Nov 17, 2017, 04:32 AM
1/**********************************
2 * FILE NAME: Gossip_membership.cpp
3 *
4 * DESCRIPTION: Membership protocol run by this Node.
5 * Definition of MP1Node class functions.
6 **********************************/
7
8#include "MP1Node.h"
9#include "Params.h"
10#include <string>
11
12
13
14/*
15 * Note: You can change/add any functions in MP1Node.{h,cpp}
16 */
17
18//random number generator - uniform distribution
19unsigned seed = chrono::system_clock::now().time_since_epoch().count();
20default_random_engine generator (seed);
21
22//global variables
23Params const *parameter;
24unsigned int recvmlSize = 0, from_mlSize = 0, erasecount = 0, foundcount = 0;
25unsigned int erasearraycount = 0, myIDcount = 0, mlCount = 0;
26long prevtimestamp, localtimestamp;
27unsigned int fromnodeID = 0, recvnodeID = 0, myID = 0;
28Member* fromNode = new Member;
29Member* recvNode = new Member;
30short sendnodePort, recvnodePort;
31long fromnodeHB, recvnodeHB;
32Address fromaddr, recvaddr, joinaddr, joinMLaddr, myaddr;
33int z = 0, f = 0, e = 0, tick1 = 0, tick2 = 0;
34unsigned int fromID = 0;
35short fromPort = 0;
36vector<MemberListEntry>::iterator foundpos, erasepos, itr;
37vector<memberState>::iterator it;
38MemberListEntry *MLentry = new MemberListEntry;
39vector<MemberListEntry> fromnodeML, recvnodeML;
40Monitor *State = new Monitor;
41
42
43/*
44 * Overloaded Constructor of the MP1Node class
45 * You can add new members to the class if you think it
46 * is necessary for your logic to work
47 */
48MP1Node::MP1Node(Member *member, Params *params, EmulNet *emul, Log *log, Address *address) {
49
50 for(int i = 0; i < 6; i++) {
51 NULLADDR[i] = 0;
52 }
53
54 this->memberNode = member;
55 this->emulNet = emul;
56 this->log = log;
57 this->par = params;
58 this->memberNode->addr = *address;
59
60}
61
62
63/**
64 * Destructor of the MP1Node class
65 */
66MP1Node::~MP1Node() {}
67
68
69/*
70 * FUNCTION NAME: recvLoop
71 *
72 * DESCRIPTION: This function receives message from the network and pushes into the queue
73 * This function is called by a node to receive messages currently waiting for it
74 */
75int MP1Node::recvLoop() {
76 if ( memberNode->bFailed ) {
77 return false;
78 }
79 else {
80 return emulNet->ENrecv(&(memberNode->addr), enqueueWrapper, NULL, 1, &(memberNode->mp1q));
81 }
82}
83
84/*
85 * FUNCTION NAME: enqueueWrapper
86 *
87 * DESCRIPTION: Enqueue the message from Emulnet into the queue
88 */
89
90int MP1Node::enqueueWrapper(void *env, char *buff, int size) {
91 Queue *new_q = new Queue;
92 new_q->enqueue((queue<q_elt> *)env, (void *)buff, size);
93 delete new_q;
94 return true;
95}
96
97/*
98 * FUNCTION NAME: nodeStart
99 *
100 * DESCRIPTION: This function bootstraps the node
101 * All initializations routines for a member.
102 * Called by the application layer.
103 */
104void MP1Node::nodeStart(char *servaddrstr, short servport) {
105
106
107 joinaddr = getJoinAddress();
108
109
110 // Self booting routines
111 if( initThisNode(joinaddr) == -1 ) {
112 #ifdef DEBUGLOG
113 log->LOG(&memberNode->addr, "Initializing failed. Exiting Application.");
114 #endif
115 exit(1);
116 }
117
118if( !introduceSelfToGroup(joinaddr) ) {
119 finishUpThisNode();
120 #ifdef DEBUGLOG
121 log->LOG(&memberNode->addr, "Unable to join self to group. Exiting Application.");
122 #endif
123 exit(1);
124 }
125
126 return;
127}
128
129/*
130 * FUNCTION NAME: initThisNode
131 *
132 * DESCRIPTION: Find out who I am and start up
133 */
134int MP1Node::initThisNode(Address joinaddr) {
135 /*
136 * This function is partially implemented and may require changes
137 */
138 if (memberNode->bFailed) {
139 //if node has already failed , it cannot be turned on again: CRASH-STOP failure mode
140 return FAILURE;
141 }
142
143
144 if (!memberNode->inited ) {
145
146 memberNode->bFailed = false;
147 memberNode->inited = true;
148 memberNode->inGroup = false;
149 memberNode->nnb = 0;
150 memberNode->heartbeat = 0;
151 memberNode->pingCounter = TFAIL;
152 memberNode->timeOutCounter = -1;
153
154
155 //initialize my ML
156 initMemberListTable(memberNode);
157
158 #ifdef DEBUGLOG
159 if (memberNode->addr.addr[0] == 1) {
160 log->LOG(&memberNode->addr, /*"Introducer(NodeID:%d)*/ "Node initialized!", memberNode->addr.addr[0]);
161 }
162 else log->LOG(&memberNode->addr, "Node initialized!");
163 #endif
164
165 }
166
167 return SUCCESS;
168}
169
170/**
171 * FUNCTION NAME: introduceSelfToGroup
172 *
173* DESCRIPTION: Co-ordinator self-joins the distributed system and then receives JOINREQ from other nodes
174 */
175
176int MP1Node::introduceSelfToGroup(Address joinaddr) {
177 msg_struct* sendmsg = (msg_struct *) malloc(State->msgsize * sizeof(char)* sizeof(MemberListEntry) * MAX_NODES);
178 msg_struct* recmsg = (msg_struct *) malloc(State->msgsize * sizeof(char)* sizeof(MemberListEntry) * MAX_NODES);
179
180 if (!memberNode->inGroup && !memberNode->bFailed){
181
182
183 if ( 0 == memcmp(&memberNode->addr.addr, &joinaddr.addr, sizeof(memberNode->addr.addr))) {
184 // I am the INTRODUCER (first process to join the group). Boot up the group...
185 #ifdef DEBUGLOG
186 log->LOG(&memberNode->addr, /*"Introducer(NodeID:%d)*/ "Starting up group...", memberNode->addr.addr[0]);
187 #endif
188
189
190 //set memberNode in-group status to TRUE
191 memberNode->inGroup = true;
192
193 // Create membership entry for INTRODUCER and emplace in its own memberlist table
194 localtimestamp = par->getcurrtime();
195 memberNode->myPos = memberNode->memberList.begin();
196 memberNode->memberList.emplace(memberNode->myPos,makeMLentry(memberNode,localtimestamp));
197
198 //State variables
199 State->INTRODUCERML = memberNode->memberList;
200 ++State->nodejoinedcount;
201 memberNode->nnb = memberNode->memberList.size() - 1;
202
203
204 #ifdef DEBUGLOG
205 log->logNodeAdd(&joinaddr, &memberNode->addr);
206 /*char selfjoinmsg[par->MAX_MSG_SIZE];
207 sprintf(selfjoinmsg, "Introducer node requested to self-join. JOINREQ Count = %u. My list size is now:%lu. Node join count:%u,", State->joinreqcount,memberNode->memberList.size(), State->nodejoinedcount);
208 log->LOG(&memberNode->addr,selfjoinmsg);*/
209 #endif
210
211
212 }
213
214 else {
215
216 // I am NOT the INTRODUCER. So, I need to send a JOINNREQ msg to the INTRODUCER.
217 // create JOINREQ message: format of data is {struct Address myaddr}
218
219
220 sendmsg->from_msgType = JOINREQ;
221 sendmsg->from = memberNode->addr;
222 sendmsg->to = joinaddr;
223 sendmsg->mNode = memberNode;
224 sendmsg->size = sizeof(*sendmsg);
225
226 //send JOINREQ message via emulNet
227 emulNet->ENsend(&sendmsg->from, &sendmsg->to, (char *)sendmsg, sendmsg->size);
228 ++State->joinreqcount;//increment JOINREQ counter
229
230 //Log JOINREQ message sent
231 #ifdef DEBUGLOG
232 char joinreqmsg[par->MAX_MSG_SIZE];
233 sprintf(joinreqmsg, "Trying to join...");/*sent JOINREQ to Introducer. JOINREQ count = %u", State->joinreqcount);*/
234 log->LOG(&memberNode->addr, joinreqmsg);
235 #endif
236
237
238 }
239
240
241 }
242 free((void*)sendmsg);
243 free((void*)recmsg);
244 return 1;
245
246}
247
248/**
249 * FUNCTION NAME: finishUpThisNode
250 *
251 * DESCRIPTION: Wind up this node and clean up state
252 */
253void MP1Node::finishUpThisNode(){
254 /*
255 * Your code goes here
256 */
257 this->memberNode->inited = false;
258 this->memberNode->bFailed = true;
259 this->memberNode->inGroup = false;
260 this->memberNode->nnb = 0;
261 this->memberNode->heartbeat = 0;
262 this->memberNode->pingCounter = 0;
263 this->memberNode->timeOutCounter = 0;
264 initMemberListTable(memberNode);
265
266}
267
268/**
269 * FUNCTION NAME: nodeLoop
270 *
271 * DESCRIPTION: Executed periodically at each member
272 * Check your messages in queue and perform membership protocol duties
273 */
274void MP1Node::nodeLoop() {
275
276 if (memberNode->bFailed) {
277 return;
278 }
279
280
281 // Check my messages
282 checkMessages();
283
284 // Wait until you're in the group...
285 if( !memberNode->inGroup ) {
286 return;
287 }
288
289 // ...then jump in and share your responsibilites!
290 nodeLoopOps();
291
292
293 return;
294
295}
296
297/**
298 * FUNCTION NAME: checkMessages
299 *
300 * DESCRIPTION: Check messages in the queue and call the respective message handler
301 */
302void MP1Node::checkMessages() {
303
304 //number of messages in queue
305 State->num_msgsinqueue = memberNode->mp1q.size();
306
307 // Pop waiting messages from memberNode's mp1q
308 void *ptr;
309 int size;
310
311 while (!memberNode->mp1q.empty()) {
312
313 ptr = memberNode->mp1q.front().elt;
314 size = memberNode->mp1q.front().size;
315 memberNode->mp1q.pop();
316
317
318 /*#ifdef DEBUGLOG
319 char str[100] ;
320 sprintf(str,"memberNode ID: %d's queue size is %u.",memberNode->addr.addr[0],memberNode->mp1q.size());
321 log->LOG(&memberNode->addr, str);
322 #endif*/
323
324 recvCallBack((void *)memberNode, (char *)ptr, size);
325 }
326
327return;
328
329}
330
331/**
332 * FUNCTION NAME: recvCallBack
333 *
334 * DESCRIPTION: Message handler for different message types
335 */
336bool MP1Node::recvCallBack(void *env, char *data, int size ) {
337 /*
338 * Your code goes here
339 */
340
341 msg_struct* sendmsg = (msg_struct *) malloc(sizeof(*sendmsg) * sizeof(char)* sizeof(MemberListEntry) * MAX_NODES);
342 msg_struct* recmsg = (msg_struct *) malloc(sizeof(*recmsg) * sizeof(char)* sizeof(MemberListEntry) * MAX_NODES);
343 memberNode = (Member *)env;
344 recmsg = (msg_struct *)data;
345 localtimestamp = par->getcurrtime();
346 joinaddr = getJoinAddress();
347 MemberListEntry *List = new MemberListEntry;
348 MemberListEntry* fromMemberEntry = new MemberListEntry;
349
350 /* -------------------------------------
351 THREE message handling cases
352 ------------------------------------- */
353
354 /*-----------------------------------CASE 1--------------------------------------------------*/
355
356 // CASE1:If node receives JOINREQ, Introducer joins the requesting node.
357 // Adds the node to its ML. Then sends JOINREP msg to all nodes in its ML.
358
359 if (recmsg->from_msgType == JOINREQ && memberNode->inGroup && !memberNode->bFailed ) {
360
361 myaddr = memberNode->addr;
362 myID = myaddr.addr[0];
363 fromNode = recmsg->mNode;
364 fromaddr = recmsg->from;
365 fromnodeID = fromaddr.addr[0];
366
367 //Introducer creates MLentry for the JOINREQ sending node
368 *MLentry = makeMLentry(fromNode, localtimestamp);
369
370 if (sizeof(*MLentry) != 0 ) {
371
372 //introducer adds to its ML, node entry from which it received JOINREQ
373 memberNode->memberList.emplace(memberNode->memberList.begin(),*MLentry);
374
375 //State variables
376
377 memberNode->nnb = memberNode->memberList.size()-1;
378 State->INTRODUCERML = memberNode->memberList;
379 State->FROMADDR[memberNode->nnb-1] = fromaddr; //storing all JOINREQ message sending addresses
380
381 //logging addition to INTRODUCER's list
382 #ifdef DEBUGLOG
383 /*char str[par->MAX_MSG_SIZE] ;
384 sprintf(str, "Introducer ADDED %s to its ML. List size is now %lu.",formatAddress(&fromaddr).c_str(),memberNode->memberList.size());
385 log->LOG(&memberNode->addr, str);*/
386 log->logNodeAdd(&myaddr, &fromaddr);
387 #endif
388
389
390 //every time a node joins, the introducer sends JOINREP messages to all its neighboring nodes. Includes its current membership list in the message.
391
392 for (int n = 0; n < memberNode->nnb; n++) {
393
394 //construct JOINREP msg
395 sendmsg->from_msgType = JOINREP;
396 sendmsg->from = myaddr;
397 sendmsg->to = State->FROMADDR[n];
398 sendmsg->mNode = memberNode;
399 sendmsg->size = sizeof(*sendmsg);
400
401
402 //Introducer sends JOINREP msg to the JOINREQ requesting node
403 emulNet->ENsend(&sendmsg->from, &sendmsg->to , (char *) sendmsg, sendmsg->size);
404
405 //log the JOINREP message send
406
407 /*#ifdef DEBUGLOG
408 char joinrepsendmsg[par->MAX_MSG_SIZE] ;
409 sprintf(joinrepsendmsg, "Introducer sent JOINREP message and current list[size:%lu] to %s",sendmsg->mNode->memberList.size(),formatAddress(&sendmsg->to).c_str());
410 log->LOG(&memberNode->addr,joinrepsendmsg);
411 #endif*/
412
413 }
414
415
416
417
418 }
419
420
421 }
422
423 /*-----------------------------------CASE 2--------------------------------------------------*/
424
425 // CASE2: Member node receives JOINREP message from the Introducer.
426 //Node checks if its in introducer's ML. if it is, it is flagged joined
427 //Adopts introducer's ML.
428 if (recmsg->from_msgType == JOINREP && !memberNode->bFailed) {
429
430 Address myAddr;
431 memberState mlState;
432 myID = memberNode->addr.addr[0];
433 myAddr = memberNode->addr;
434 fromNode = recmsg->mNode;
435 fromaddr = recmsg->from;
436 fromnodeID = fromaddr.addr[0];
437 fromPort = fromaddr.addr[4];
438 recvnodeML = memberNode->memberList;
439
440 if(fromnodeID == joinaddr.addr[0] && myID != 0 && myID != fromID) {
441
442 fromnodeML.clear();
443 fromnodeML = fromNode->memberList;
444 from_mlSize = fromnodeML.size();
445
446 //Check if the introducer ML has the member node's ID in it.
447 if (!fromnodeML.empty() && from_mlSize != 0) foundcount = findMLEntryPos(fromnodeML, fromnodeID);
448
449 //if memberNode is IN introducer's list, flag memberNode as having joined group.
450 if(foundcount > 0) {
451 if (!memberNode->inGroup){
452 memberNode->inGroup = true;
453 ++State->nodejoinedcount;//increment node join counter
454 //log the join status
455 /*#ifdef DEBUGLOG
456 char nodejoinmsg[par->MAX_MSG_SIZE] ;
457 sprintf(nodejoinmsg, "Node present in JOINREP introducer's ML.Flagged joined! Node join count:%u.",State->nodejoinedcount);
458 log->LOG(&myAddr, nodejoinmsg);
459 #endif*/
460
461 }
462 //JOINED (inGroup) node ADOPTS introducer's current ML.
463
464 if (memberNode->inGroup) {
465
466 if (recvnodeML.empty() && !fromnodeML.empty()) {
467 memberNode->memberList = fromnodeML; //adopt introducer's list
468 recvnodeML = memberNode->memberList;
469 recvmlSize = recvnodeML.size();
470 memberNode->nnb = fromNode->nnb;
471
472 for (int i=0; i < recvmlSize; i++) {
473 recvnodeID = recvnodeML[i].id;
474 recvnodePort = recvnodeML[i].port;
475 /*if (recvnodeID != myID) {*/
476 joinMLaddr = getAnyAddress(recvnodeID,recvnodePort);
477 //log the newly joined node
478 log->logNodeAdd(&myAddr, &joinMLaddr);
479 /*}*/
480
481 }
482
483 }
484
485 else if (!recvnodeML.empty() && !fromnodeML.empty()){
486
487 for (itr = fromnodeML.begin(); itr < fromnodeML.end(); itr++) {
488 int itrID = itr->id;
489 short itrPort = itr->port;
490
491 foundcount = findMLEntryPos(recvnodeML,itrID);
492
493 if (foundcount == 0) {
494
495 joinMLaddr = getAnyAddress(itrID,itrPort);
496 //log the newly joined node
497 /*if (itrID != myID)*/ log->logNodeAdd(&myAddr, &joinMLaddr);
498
499 }
500
501 }
502
503 memberNode->memberList = fromnodeML;//adopt introducer's list
504 memberNode->nnb = fromNode->nnb;
505 }
506
507
508 if (memberNode->memberList.size() == par->EN_GPSZ) State->full_listadoptcount++;
509
510
511 }
512
513
514 }
515
516 //when nodejoincount equals par->EN_GPSZ, set allNodesJoined flag to TRUE
517 if (State->nodejoinedcount == par->EN_GPSZ && !State->allNodesJoined) {
518
519 State->allNodesJoined = true;
520
521 //log flag status
522 /*#ifdef DEBUGLOG
523 char allnodesjoinmsg[par->MAX_MSG_SIZE];
524 sprintf(allnodesjoinmsg, "ALL %u nodes have JOINED @time:%lu!",State->nodejoinedcount, localtimestamp);
525 log->LOG(&joinaddr, allnodesjoinmsg);
526 #endif*/
527
528
529 }
530
531
532 //when all nodes have adopted the full ML, sett adopt ML flag to true and log ML list adoption
533
534 if (State->full_listadoptcount == (par->EN_GPSZ - 1) && !State->allNodesAdopt){
535
536 State->allNodesAdopt = true;
537
538 /*#ifdef DEBUGLOG
539 char nodeadoptmsg[par->MAX_MSG_SIZE];
540 sprintf(nodeadoptmsg, "ALL %u joined nodes have ADOPTED introducer's full ML(size:%lu) @time:%lu!",State->full_listadoptcount, fromnodeML.size(), localtimestamp);
541 log->LOG(&joinaddr, nodeadoptmsg);
542 #endif*/
543
544 }
545
546
547
548 }
549
550
551 }
552
553
554
555 /*-----------------------------------CASE 3---------------------------------------------------*/
556
557
558 //Case 3:If ALIVE and InGroup node receives UPDATEMEMTABLE message, it updates its membership list table per GOSSIP PROTOCOL rules.
559
560
561 if (recmsg->from_msgType == UPDATEMEMTABLE && memberNode->inGroup && !memberNode->bFailed) {
562
563
564 //Step 1: comparing heartbeat status of each node in the incoming member list with the member list in state.
565 fromnodeML.clear();
566 recvnodeML.clear();
567 fromNode = recmsg->mNode;
568 fromnodeML = fromNode->memberList;
569 recvnodeML = memberNode->memberList;
570 myID = memberNode->addr.addr[0];
571 foundcount = 0;
572
573 if(!fromnodeML.empty() && !recvnodeML.empty()){
574
575 unsigned int autorecvnodeID = 0;
576 from_mlSize = fromnodeML.size();
577 recvmlSize = recvnodeML.size();
578
579 for (int i = 0; i < recvmlSize; i++){
580
581 autorecvnodeID = memberNode->memberList[i].id;
582 recvnodeHB = memberNode->memberList[i].heartbeat;
583 prevtimestamp = memberNode->memberList[i].timestamp;
584
585
586 //if in-state ID is found in from-member list: for that ID, check if from-member's HB > in-state member's HB
587
588 if(autorecvnodeID != myID){
589
590 foundcount = findMLEntryPos(fromnodeML,autorecvnodeID);
591
592 if (foundcount > 0) {
593
594 foundpos = find_if (fromnodeML.begin(), fromnodeML.end(), [autorecvnodeID] (MemberListEntry List) mutable -> bool {return List.id == autorecvnodeID;});
595 fromnodeID = foundpos->id;
596 fromnodeHB = foundpos->heartbeat;
597
598 if(fromnodeHB > State->run_time_end) exit(fromnodeHB);//error check function
599
600 //if from-member's HB is greater, UPDATE in-state member's HB and local time
601 if (fromnodeHB > recvnodeHB){
602
603 memberNode->memberList[i].setheartbeat(fromnodeHB);//update with the higher HB
604 memberNode->memberList[i].settimestamp(localtimestamp);//update time stamp
605
606 //log heart beat updated status
607 /*#ifdef DEBUGLOG
608 char HBupdatemsg[par->MAX_MSG_SIZE];
609 sprintf(HBupdatemsg, "Updated my mem list. Stats: Compared NodeID:%u. Previous HB:%ld; Updated HB:%ld. Previous time:%ld.Updated time: %ld ",
610 autorecvnodeID,recvnodeHB,memberNode->memberList[i].heartbeat, prevtimestamp, memberNode->memberList[i].timestamp);
611 log->LOG(&memberNode->addr, HBupdatemsg);
612 #endif*/
613
614 }
615
616
617 /*else {
618
619 //log heart beat NOT updated status
620 #ifdef DEBUGLOG
621 char HBnotupdatemsg[par->MAX_MSG_SIZE];
622 sprintf(HBnotupdatemsg, "NOT updated my mem list. Stats: Compared NodeID:%u. Previous HB:%ld; Updated HB:%ld. Previous time:%ld.Updated time: %ld ",
623 autorecvnodeID,recvnodeHB,memberNode->memberList[i].heartbeat, prevtimestamp, memberNode->memberList[i].timestamp);
624 log->LOG(&memberNode->addr, HBnotupdatemsg);
625 #endif
626
627
628 }*/
629
630
631
632 }
633
634
635
636
637 }
638
639
640 }
641
642 }
643
644
645
646 }
647
648
649 delete List;
650 delete fromMemberEntry;
651 free((void*)sendmsg);
652 free((void*)recmsg);
653 return true;
654}
655
656
657
658/*
659 * FUNCTION NAME: nodeLoopOps
660 *
661 * DESCRIPTION: Check if any node hasn't responded within a timeout period and then delete
662 * the nodes
663 * Propagate your membership list
664 */
665
666void MP1Node::nodeLoopOps() {
667
668 /*
669 * Your code goes here
670 */
671
672
673 unsigned int entryID = 0, failcount = 0;
674 short entryPort = 0;
675 long entryHB = 0;
676 long entryTimeLapse = 0;
677 localtimestamp = par->getcurrtime();
678 unsigned int myID = 0;
679 Address entryAddr;
680
681 myID = memberNode->addr.addr[0];
682
683 //STEP 1:Alive and JOINED (inGroup) nodes start HBing.
684 //update HB and local time in membernode's ML. Then propagate ML.
685
686 if(!memberNode->bFailed && memberNode->inGroup) {
687
688 ++memberNode->heartbeat;
689 State->HBlistcount++;//increment state variable
690
691 //update HB and TS in memberNode's member list
692 if (!memberNode->memberList.empty()) updateHBandTS(memberNode, localtimestamp);
693
694 if(memberNode->addr.addr[0] == joinaddr.addr[0]) State->INTRODUCERML = memberNode->memberList;//keeping track of introducer's ML
695
696 //logging heart beating start status
697
698 /*#ifdef DEBUGLOG
699 char HBstartmsg[par->MAX_MSG_SIZE];
700 sprintf(HBstartmsg, "HeartBeating @time:%ld. HB Count:%ld.", localtimestamp, memberNode->heartbeat);
701 log->LOG(&memberNode->addr, HBstartmsg);
702 #endif*/
703
704
705
706
707
708
709
710 //STEP 2: Check node status at TFAIL and TREMOVE after ALL nodes have adopted full ML
711 recvmlSize = memberNode->memberList.size();
712 memberState feState;
713
714 for (int j=0; j < recvmlSize; j++){
715
716 entryID = memberNode->memberList[j].id;
717 entryPort = memberNode->memberList[j].port;
718 entryHB = memberNode->memberList[j].heartbeat;
719 prevtimestamp = memberNode->memberList[j].timestamp;
720
721 if (entryID != 0 && entryID != myID){
722
723 //compute elapsed time
724 entryTimeLapse = localtimestamp - prevtimestamp;
725
726 //if elapsed time >TFAIL but < TREMOVE, node in the memberlist is marked failed but NOT erased from membership list.
727 if (entryTimeLapse > TFAIL && entryTimeLapse < TREMOVE) {
728
729 entryAddr = getAnyAddress(entryID,entryPort);
730
731
732 //log failed node's stats
733 /*#ifdef DEBUGLOG
734 char failnodemsg[par->MAX_MSG_SIZE] ;
735 sprintf(failnodemsg, "marked FAILED %s. Elapsed time(%ld) > TFAIL(%d). \nStats: latest HB: %ld; Prev time: %ld; Current time: %ld.", formatAddress(&entryAddr).c_str(), entryTimeLapse, TFAIL ,entryHB, prevtimestamp, localtimestamp);
736 log->LOG(&memberNode->addr, failnodemsg);
737 #endif*/
738
739 //if not previously recorded, record the marked failed node address and first fail time
740 if (!State->FAILEDMEMBER.empty()) failcount = findMLEntryPos(State->FAILEDMEMBER, entryID);
741
742 if ((failcount == 0) or (State->FAILEDMEMBER.empty())) {
743
744 feState.fe_addr = entryAddr;
745 feState.fe_time = localtimestamp;
746 State->FAILEDMEMBER.emplace(State->FAILEDMEMBER.begin(),feState);
747
748 }
749
750 }
751
752
753 //else if time elapses beyond TREMOVE, erase node from my membership list
754
755 else if (entryTimeLapse > TREMOVE ){
756
757
758 if ((entryID != 0) && (entryID != myID)){
759
760
761 //find the position of entryID in my list and erase that node
762
763 if(!memberNode->memberList.empty()) erasecount = findMLEntryPos(memberNode->memberList,entryID);
764
765 //if entryID is found in the member list, erase it
766 if (erasecount > 0) {
767
768 erasepos = find_if (memberNode->memberList.begin(), memberNode->memberList.end(), [entryID] (MemberListEntry List) mutable -> bool {return List.id == entryID;});
769 entryAddr = getAnyAddress(erasepos->id,erasepos->port);
770
771 //erase node
772 memberNode->memberList.erase(erasepos);
773 memberNode->nnb --;
774
775 //log erased node
776 #ifdef DEBUGLOG
777 log->logNodeRemove(&memberNode->addr,&entryAddr);
778 /*char erasednodestatsmsg[par->MAX_MSG_SIZE];
779 sprintf(erasednodestatsmsg, "Elapsed time(%ld) > TREMOVE(%d). List size after removal: %lu.\nStats: latest HB: %ld; Prev time: %ld; Current time: %ld.", entryTimeLapse,TREMOVE,memberNode->memberList.size(), entryHB, prevtimestamp,localtimestamp);
780 log->LOG(&memberNode->addr, erasednodestatsmsg);*/
781 #endif
782
783
784 //if not previously recorded, record the erased node address in the ERASEDMEMBER vector.
785
786
787 if (!State->ERASEDMEMBER.empty()) erasearraycount = findMLEntryPos(State->ERASEDMEMBER,entryID);
788
789 if ((erasearraycount == 0) or (State->ERASEDMEMBER.empty())) {
790
791 feState.fe_addr = entryAddr;
792 feState.fe_time = localtimestamp;
793 State->ERASEDMEMBER.emplace(State->ERASEDMEMBER.begin(), feState);
794
795 }
796
797 }
798
799 }
800
801
802 }
803
804 }
805
806
807 }
808
809
810
811 //after all steps are complete, infection propagate my updated ML [GOSSIP PROTOCOL].
812 //Alive and inGroup nodes call randomPropagate function which fans out msg to random nodes.
813 if(State->nodejoinedcount > 1) randomPropagate(memberNode,UPDATEMEMTABLE);
814
815
816
817 }
818
819 //At the end of time loop [TOTAL_RUNNING_TIME],log all failed and erased nodes
820 if ((par->getcurrtime() == State->run_time_end-1) && (tick2 < 1)){
821
822 for (it = State->FAILEDMEMBER.begin(); it!= State->FAILEDMEMBER.end(); ++it) {log->LOG(&it->fe_addr, "Summary:First time Marked FAILED @time:%ld.", it->fe_time);}
823 for (it = State->ERASEDMEMBER.begin(); it!= State->ERASEDMEMBER.end(); ++it) {log->LOG(&it->fe_addr, "Summary:First time REMOVED @time:%ld.",it->fe_time);}
824
825 ++tick2;
826
827 }
828
829 return;
830}
831
832/**
833 * FUNCTION NAME: isNullAddress
834 *
835 * DESCRIPTION: Function checks if the address is NULL
836 */
837
838int MP1Node::isNullAddress(Address *addr) {
839 return (memcmp(addr->addr, NULLADDR, 6) == 0 ? 1 : 0);
840}
841
842/**
843 * FUNCTION NAME: getJoinAddress
844 *
845 * DESCRIPTION: Returns the Address of the coordinator
846 */
847
848Address MP1Node::getJoinAddress() {
849 Address joinaddr;
850 memset(&joinaddr, 0, sizeof(Address));
851 joinaddr.addr[0] = 1;
852 joinaddr.addr[4] = 0;
853 return joinaddr;
854}
855
856/**
857 * FUNCTION NAME: initMemberListTable
858 *
859 * DESCRIPTION: Initialize the membership list
860 */
861void MP1Node::initMemberListTable(Member *memberNode) {
862 memberNode->memberList.clear();
863}
864
865/**
866 * FUNCTION NAME: printAddress
867 *
868 * DESCRIPTION: Print the Address
869 */
870void MP1Node::printAddress(Address *addr)
871{
872 printf("%d.%d.%d.%d:%d \n", addr->addr[0],addr->addr[1],addr->addr[2],addr->addr[3], *(short*)&addr->addr[4]);
873
874}
875
876/**
877 * FUNCTION NAME: formatAddress
878 *
879 * DESCRIPTION: format Address from NodeID
880 */
881string MP1Node::formatAddress(Address *addr)
882{
883 char str[par->MAX_MSG_SIZE];
884 sprintf(str,"%d.%d.%d.%d:%d", addr->addr[0],addr->addr[1],addr->addr[2],addr->addr[3], *(short*)&addr->addr[4]);
885 return str;
886
887}
888
889
890
891/**
892 * FUNCTION NAME: randomPropagate
893 *
894 * DESCRIPTION: Propagates message to random nodes
895 */
896
897void MP1Node::randomPropagate(Member* memNode, enum MsgTypes m_type)
898{
899 int fan=0;
900 short recvnodePort = 0;
901 msg_struct* anymsg = (msg_struct *) malloc(State->msgsize * sizeof(char)* sizeof(MemberListEntry) * MAX_NODES);
902 myID = memNode->addr.addr[0];
903
904 //random integer generator - initializing...
905 unsigned int x = 0;
906 if(State->nodejoinedcount < 2) x = 2;
907 else x = State->nodejoinedcount;
908 random_device rd;
909
910
911 //constructing anymsg to propagate
912 anymsg->from_msgType = m_type;
913 anymsg->from = memNode->addr;
914 anymsg->mNode= memNode;
915
916
917 while (fan < FANOUT){
918 unsigned int y = rd();
919 recvnodeID = (y % x)+1 ;//generate random integer as ID
920
921 recvaddr = getAnyAddress(recvnodeID, recvnodePort);
922 anymsg->to = recvaddr;
923 anymsg->size = sizeof(*anymsg);
924
925 //send msg via emulNet to fanth random node to ask it to update its membership list
926
927 if (!(recvnodeID == myID) or (recvnodeID == 0)){
928
929 emulNet->ENsend(&anymsg->from, &anymsg->to, (char *)anymsg, anymsg->size);
930
931 fan++;
932
933 /*const char* str;
934 switch (m_type) {
935 case 0:
936 str = "JOINREQ" ;
937 break;
938 case 1:
939 str = "JOINREP";
940 break;
941 case 2:
942 str = "UPDATEMEMTABLE";
943 }
944
945 #ifdef DEBUGLOG
946
947 char propagatemsg[par->MAX_MSG_SIZE] ;
948 sprintf(propagatemsg, "Propagating message type: %s to randomID:%u", str, recvnodeID);
949 log->LOG(&memNode->addr, propagatemsg);
950
951 #endif*/
952
953
954 }
955
956
957 }
958
959 free((void*)anymsg);
960 return;
961
962}
963
964/**
965 * FUNCTION NAME: getAnyAddress
966 *
967 * DESCRIPTION: makes an address object from ID and Port
968 */
969
970Address MP1Node::getAnyAddress(unsigned int ID, short Port)
971
972{
973 Address makeAddr;
974 *(int *)(&(makeAddr.addr)) = ID;
975 *(int *)(&(makeAddr.addr[1])) = 0 ;
976 *(int *)(&(makeAddr.addr[2])) = 0 ;
977 *(int *)(&(makeAddr.addr[3])) = 0 ;
978 *(short *)(&(makeAddr.addr[4])) = Port;
979 *(int *)(&(makeAddr.addr[5])) = 0 ;
980
981 return makeAddr;
982
983
984}
985
986/**
987 * FUNCTION NAME: makeMLentry
988 *
989 * DESCRIPTION: makes an MLentry
990 */
991
992MemberListEntry MP1Node::makeMLentry (Member* mNode, long tStamp)
993{
994
995
996 MLentry->setid(mNode->addr.addr[0]);
997 MLentry->setport(mNode->addr.addr[4]);
998 MLentry->setheartbeat(mNode->heartbeat);
999 MLentry->settimestamp(tStamp);
1000 return *MLentry;
1001
1002}
1003
1004/**
1005 * FUNCTION NAME: findMLEntryPos
1006 *
1007 * DESCRIPTION: finding if an ID exists in member list
1008 *
1009 */
1010
1011unsigned int MP1Node::findMLEntryPos(vector<MemberListEntry> targetML, unsigned int searchID)
1012{
1013 int fcount = 0;
1014 fcount = count_if (targetML.begin(), targetML.end(), [searchID] (MemberListEntry List) mutable -> bool {return List.id == searchID;});
1015 return fcount;
1016
1017}
1018//overloaded function
1019unsigned int MP1Node::findMLEntryPos(vector<memberState> targetMS, unsigned int searchID)
1020{
1021 int fcount = 0;
1022 fcount = count_if (targetMS.begin(), targetMS.end(), [searchID] (memberState seq_State) mutable -> bool {return seq_State.fe_addr.addr[0] == searchID;});
1023 return fcount;
1024}
1025
1026
1027
1028void MP1Node::updateHBandTS(Member* m_node, long t_stamp)
1029{
1030 unsigned int myautoID = 0;
1031 vector<MemberListEntry>::iterator myIDpos;
1032
1033
1034 //update HB and timestamp in own memberlist
1035 if (!m_node->memberList.empty()){
1036
1037 myautoID = m_node->addr.addr[0];
1038 myIDcount = count_if (m_node->memberList.begin(), m_node->memberList.end(), [myautoID](MemberListEntry List) mutable -> bool {return List.id == myautoID;});
1039
1040 if (myIDcount > 0) {
1041 myIDpos = find_if (m_node->memberList.begin(), m_node->memberList.end(), [myautoID](MemberListEntry List) mutable -> bool {return List.id == myautoID;});
1042 myIDpos->setheartbeat(m_node->heartbeat); //update heartbeatin memberNode's ML
1043 myIDpos->settimestamp(t_stamp); //Update timestamp in memberNode's ML
1044
1045 //log Heartbeat increment status
1046
1047 /*#ifdef DEBUGLOG
1048 char HBmsg[par->MAX_MSG_SIZE];
1049 sprintf(HBmsg, "Increments heartbeat to %ld @time:%ld",myIDpos->heartbeat,myIDpos->timestamp);
1050 log->LOG(&memberNode->addr, HBmsg);
1051 #endif*/
1052
1053
1054 }
1055 }
1056
1057 return;
1058
1059}