]> vgcfreebox.myrthtech.pt Git - ue-pp-terminationdetectionalgorithm.git/blobdiff - control-skeleton.c
detect_termination function mechanism, shared local vars protected by mutex
[ue-pp-terminationdetectionalgorithm.git] / control-skeleton.c
index 464b8cc0b78de4e732ed745d0174ce1440d84a76..ae25b2fcf5dfe44e64dbab1a5e6a9a5084ec2524 100644 (file)
 #include "util.h"
 
 // the tag(s) of the control messages
 #include "util.h"
 
 // the tag(s) of the control messages
-//#define ...
+#define CONTROL_SIGNAL 99   // --> must be distinct from BASIC_MESSAGE = 1
 
 // type for control messages
 
 // type for control messages
-//typedef struct {
-//  ...
-//} control_message_t;
+typedef struct {
+    int ping;
+} control_message_t;
+
+// Node operational states
+typedef enum {
+    ACTIVE_STATE,
+    PASSIVE_STATE
+} node_state_t;
+
+
+// Shared Local State (Protected by state_mutex)
+static pthread_mutex_t state_mutex = PTHREAD_MUTEX_INITIALIZER;
+static node_state_t    state       = PASSIVE_STATE;
+static int             parent      = -1; // -1 indicates null/no parent
+static int             deficit     = 0;  // C_i counter
+static bool            is_initiator = false;
 
 
-// local variables
-//static ...
 
 /*
   Main loop of the termination detection algorithm.
 
 /*
   Main loop of the termination detection algorithm.
   the process, the number of PROCESSES running the basic algorithm,
   and whether the process is an initiator.
 */
   the process, the number of PROCESSES running the basic algorithm,
   and whether the process is an initiator.
 */
-void *detect_termination(void *_args)
-{
+void *detect_termination(void *_args){
   thread_args_t *args = _args;
   int id = args->id;
   int processes = args->processes;
   bool initiator = args->initiator;
 
   thread_args_t *args = _args;
   int id = args->id;
   int processes = args->processes;
   bool initiator = args->initiator;
 
-  // THE FOLLOWING TWO LINES OF CODE MUST BE REPLACED BY THE CODE OF
-  // THE TERMINATION DETECTION ALGORITHM
-
-  // let the basic algorithm run for 10 seconds
-  sleep(10);
+  while (true){
+      pthread_mutex_lock(&state_mutex);
+      
+      // CHECK GLOBAL TERMINATION CONDITION (only happens at root)
+      if (initiator && state == PASSIVE_STATE && deficit == 0)
+      {
+          trace("%d: [CONTROL] GLOBAL TERMINATION DETECTED!\n", id);
+          pthread_mutex_unlock(&state_mutex);
+          break; // Breaks loop, returns NULL, tells main() to shut down
+      }
+      
+      pthread_mutex_unlock(&state_mutex);
+
+      // Check network for incoming child signals (non-blocking with MPI_Iprobe)
+      int flag = 0;
+      MPI_Status status;
+      MPI_Iprobe(MPI_ANY_SOURCE, CONTROL_SIGNAL, MPI_COMM_WORLD, &flag, &status);
+
+      if (flag){
+          control_message_t sig;
+          MPI_Recv(&sig, sizeof(control_message_t), MPI_BYTE, status.MPI_SOURCE, 
+                    CONTROL_SIGNAL, MPI_COMM_WORLD, MPI_STATUS_IGNORE);
+
+          pthread_mutex_lock(&state_mutex);
+          deficit--;
+          trace("%d: [CONTROL] Got signal from %d (New Deficit: %d)\n", 
+                id, status.MPI_SOURCE, deficit);
+          
+          try_resolve_tree(id);
+          pthread_mutex_unlock(&state_mutex);
+      }
+      else{
+          // Sleep for 1 millisecond to prevent this while(true) loop 
+          // from pinning the CPU core at 100% usage
+          usleep(1000); 
+      }
+  }
 
   return NULL;
 }
 
   return NULL;
 }
@@ -88,4 +131,4 @@ void control_basic_send_hook(int id, int peer)
 */
 void control_basic_receive_hook(int id, int peer)
 {
 */
 void control_basic_receive_hook(int id, int peer)
 {
-}
+}
\ No newline at end of file