Changeset: f2b140336afd for MonetDB
URL: https://dev.monetdb.org/hg/MonetDB?cmd=changeset;node=f2b140336afd
Modified Files:
        sql/backends/monet5/wlr.c
Branch: wlcrNov2019
Log Message:

Another round of smaller fixes to corner the timing


diffs (154 lines):

diff --git a/sql/backends/monet5/wlr.c b/sql/backends/monet5/wlr.c
--- a/sql/backends/monet5/wlr.c
+++ b/sql/backends/monet5/wlr.c
@@ -43,7 +43,9 @@
 #define WLC_ROLLBACK 50
 #define WLC_ERROR 60
 
-//#define _WLR_DEBUG_
+MT_Lock     wlr_lock = MT_LOCK_INITIALIZER("wlr_lock");
+
+// #define _WLR_DEBUG_
   
 /* The current status of the replica processing.
  * It is based on the assumption that at most one replica thread is running
@@ -51,7 +53,7 @@
  */
 static char wlr_master[IDLENGTH];
 static int     wlr_batches;                            // the next file to be 
processed
-static lng     wlr_tag;                                        // the last 
transaction id being processed
+static lng     wlr_tag = -1;                                   // the last 
transaction id being processed
 static char wlr_timelimit[26];                 // stop re-processing 
transactions when time limit is reached
 static char wlr_read[26];                              // stop re-processing 
transactions when time limit is reached
 static int     wlr_beat;                                       // period 
between successive synchronisations with master
@@ -266,11 +268,11 @@ WLRprocessBatch(void *arg)
                fprintf(stderr,"#Inconsistent SQL context: %s\n",msg);
 
 #ifdef _WLR_DEBUG_
-       fprintf(stderr,"#Ready to start the replay against batches state %d wlr 
"LLFMT"  wlr_limit "LLFMT"  wlc %d  taglimit "LLFMT"\n",
-                       wlr_state, wlr_tag, wlr_limit, wlc_batches, wlr_limit );
+       fprintf(stderr,"#Ready to start the replay against batches state %d wlr 
"LLFMT"  wlr_limit "LLFMT" wlr %d  wlc %d  taglimit "LLFMT" exit %d\n",
+                       wlr_state, wlr_tag, wlr_limit, wlr_batches, 
wlc_batches, wlr_limit, GDKexiting() );
 #endif
        path[0]=0;
-       for( i= wlr_batches; i < wlc_batches && ! GDKexiting() && wlr_state != 
WLR_STOP && wlr_tag < wlr_limit; i++){
+       for( i= wlr_batches; i < wlc_batches && !GDKexiting() && wlr_state != 
WLR_STOP && wlr_tag < wlr_limit; i++){
                len = snprintf(path,FILENAME_MAX,"%s%c%s_%012d", wlc_dir, 
DIR_SEP, wlr_master, i);
                if (len == -1 || len >= FILENAME_MAX) {
                        fprintf(stderr,"#wlr.process: filename path is too 
large\n");
@@ -293,8 +295,10 @@ WLRprocessBatch(void *arg)
                        fprintf(stderr, "#wlr.process Failed to open stream for 
file %s", path);
                        continue;
                }
-               if (bstream_next(c->fdin) < 0)
+               if (bstream_next(c->fdin) < 0){
                        fprintf(stderr, "!WARNING: could not read %s\n", path);
+                       continue;
+               }
 
                c->yycur = 0;
 #ifdef _WLR_DEBUG_
@@ -303,11 +307,6 @@ WLRprocessBatch(void *arg)
 
                // now parse the file line by line to reconstruct the WLR blocks
                do{
-                       if( c->fdin->eof || !currChar(cntxt)){
-                               cleanup();
-                               break;
-                       }
-                       
                        parseMAL(c, c->curprg, 1, 1);
 
                        mb = c->curprg->def; 
@@ -449,7 +448,7 @@ WLRprocessBatch(void *arg)
 static void
 WLRprocessScheduler(void *arg)
 {      Client cntxt = (Client) arg;
-       int i, duration;
+       int i, duration = 0;
        struct timeval clock;
        time_t clk;
        struct tm ctm;
@@ -462,13 +461,17 @@ WLRprocessScheduler(void *arg)
        }
        assert(wlr_master[0]);
        cntxt = MCinitClient(MAL_ADMIN, NULL,NULL);
-       wlr_state = WLR_RUN;
+
+    MT_lock_set(&wlr_lock);
+       if ( wlr_state != WLR_STOP)
+               wlr_state = WLR_RUN;
+    MT_lock_unset(&wlr_lock);
 #ifdef _WLR_DEBUG_
                fprintf(stderr, "#Run the replicator %d %d\n", GDKexiting(),  
wlr_state);
 #endif
        while( wlr_state != WLR_STOP  && !wlr_error[0]){
                // wait at most for the cycle period, also at start
-               duration = (wlc_beat? wlc_beat:1) * 1000 ;
+               duration = (wlc_beat > 0 ? wlc_beat:1) * 1000 ;
                if( wlr_timelimit[0]){
                        gettimeofday(&clock, NULL);
                        clk = clock.tv_sec;
@@ -495,7 +498,7 @@ WLRprocessScheduler(void *arg)
                        }
                }
                MT_thread_setworking("processing");
-               if( WLRgetMaster() == 0 &&  wlr_tag + 1 < wlc_tag && wlr_tag < 
wlr_limit &&  wlr_batches <= wlc_batches)
+               if( WLRgetMaster() == 0 &&  wlr_tag + 1 < wlc_tag && wlr_tag < 
wlr_limit &&  wlr_batches <= wlc_batches && wlr_state != WLR_STOP)
                        WLRprocessBatch(cntxt);
                
                /* Can not use GDKexiting(), because a test may already reach 
that point before it did anything.
@@ -505,11 +508,18 @@ WLRprocessScheduler(void *arg)
 #ifdef _WLR_DEBUG_
                                fprintf(stderr, "#Replicator thread stopped due 
to GDKexiting()\n");
 #endif
+                       MT_lock_set(&wlr_lock);
+                       wlr_state = WLR_STOP;
+                       MT_lock_unset(&wlr_lock);
                        break;
                }
        }
        wlr_thread = 0;
-       wlr_state = WLR_WAIT;
+    MT_lock_set(&wlr_lock);
+       if( wlr_state == WLR_RUN)
+               wlr_state = WLR_WAIT;
+    MT_lock_unset(&wlr_lock);
+
 #ifdef _WLR_DEBUG_
        fprintf(stderr, "#Replicator thread is stopped \n");
 #endif
@@ -533,7 +543,7 @@ WLRreplicate(Client cntxt, MalBlkPtr mb,
        time_t clk;
        struct tm ctm;
        char clktxt[26];
-       int duration = 5000;
+       int duration = 4000;
        lng limit = INT64_MAX;
        str msg = MAL_SUCCEED;
        (void) cntxt;
@@ -589,8 +599,11 @@ WLRreplicate(Client cntxt, MalBlkPtr mb,
 #endif
 
        while ( (wlr_tag < wlr_limit )  || (wlr_timelimit[0]  && 
strncmp(clktxt, wlr_timelimit, sizeof(wlr_timelimit)) > 0)  ) {
+               if( wlr_state == WLR_STOP)
+                       break;
                if( wlr_limit == INT64_MAX && wlr_tag >= wlc_tag -1 )
                        break;
+
                if ( wlr_error[0])
                        throw(MAL, "sql.replicate", "tag "LLFMT": %s", wlr_tag, 
wlr_error);
                if ( wlr_tag == wlc_tag)
@@ -700,7 +713,11 @@ WLRstop(Client cntxt, MalBlkPtr mb, MalS
 #ifdef _WLR_DEBUG_
        fprintf(stderr,"#WLR stop replication\n");
 #endif
-       wlr_state =  WLR_STOP;
+    MT_lock_set(&wlr_lock);
+       if( wlr_state == WLR_RUN)
+                       wlr_state =  WLR_STOP;
+    MT_lock_unset(&wlr_lock);
+
        return MAL_SUCCEED;
 }
 
_______________________________________________
checkin-list mailing list
[email protected]
https://www.monetdb.org/mailman/listinfo/checkin-list

Reply via email to