Changeset: f5a68727bbae for MonetDB
URL: https://dev.monetdb.org/hg/MonetDB?cmd=changeset;node=f5a68727bbae
Added Files:
        sql/test/wlcr/Tests/wlc80.py
        sql/test/wlcr/Tests/wlr80.py
Modified Files:
        sql/backends/monet5/wlr.c
        sql/test/wlcr/Tests/All
        sql/test/wlcr/Tests/wlr100.py
        sql/test/wlcr/Tests/wlr30.stable.err
        sql/test/wlcr/Tests/wlr40.py
        sql/test/wlcr/Tests/wlr50.py
        sql/test/wlcr/Tests/wlr70.py
Branch: wlcrNov2019
Log Message:

Reaching a stage where test <40 behave as expected.


diffs (truncated from 539 to 300 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,7 @@
 #define WLC_ROLLBACK 50
 #define WLC_ERROR 60
 
-#define _WLR_DEBUG_
+//#define _WLR_DEBUG_
   
 MT_Lock     wlr_lock = MT_LOCK_INITIALIZER("wlr_lock");
 
@@ -53,7 +53,7 @@ MT_Lock     wlr_lock = MT_LOCK_INITIALIZ
  */
 static char wlr_master[IDLENGTH];
 static int     wlr_batches;                            // the next file to be 
processed
-static lng     wlr_tag;                                        // the next 
transaction id to be processed
+static lng     wlr_tag;                                        // 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
@@ -112,11 +112,14 @@ WLRgetConfig(void){
                        strcpy(wlr_timelimit, line + 10);
                else
                if( strncmp("error=", line, 6) == 0) {
+                       char *s;
                        len = snprintf(wlr_error, FILENAME_MAX, "%s", line + 6);
                        if (len == -1 || len >= FILENAME_MAX) {
                                fprintf(stderr, "wlr.getConfig:error config 
value is too large");
                                goto bailout;
                        }
+                       s = strchr(wlr_error, (int) '\n');
+                       if ( s) *s = 0;
                } else{
                                fprintf(stderr, "wlr.getConfig:unknown 
configuration item '%s'", line);
                                goto bailout;
@@ -227,7 +230,8 @@ WLRprocessBatch(void *arg)
        str msg, other;
        mvc *sql;
        Symbol prev = NULL;
-       lng tag;
+       lng tag = wlr_tag;
+       char tag_read[26];                      // stop re-processing 
transactions when time limit is reached
 
        c =MCforkClient(cntxt);
        if( c == 0){
@@ -264,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 
%d  wlc %d  taglimit "LLFMT"\n",
-                       wlr_state, wlr_batches, wlc_batches, wlr_limit );
+       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 );
 #endif
        path[0]=0;
-       for( i= wlr_batches; i < wlc_batches && ! GDKexiting() && wlr_state != 
WLR_STOP; 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");
@@ -301,7 +305,7 @@ WLRprocessBatch(void *arg)
 
                // now parse the file line by line to reconstruct the WLR blocks
                do{
-                       if( c->fdin->eof){
+                       if( c->fdin->eof || !currChar(cntxt)){
                                cleanup();
                                break;
                        }
@@ -322,33 +326,34 @@ WLRprocessBatch(void *arg)
                                continue;
                        }
                        q= getInstrPtr(mb, mb->stop - 1);
+                       if( getModuleId(q) != wlrRef){
+#ifdef _WLR_DEBUG_XTRA
+                fprintf(stderr,"#unexpected instruction ");
+                               fprintInstruction(stderr, mb, 0, q, 
LIST_MAL_ALL);
+#endif
+                               cleanup();
+                               break;
+                       }
                        if( getModuleId(q) == wlrRef && getFunctionId(q) == 
transactionRef){
                                tag = getVarConstant(mb, getArg(q,1)).val.lval;
+                               snprintf(tag_read, sizeof(wlr_read), "%s", 
getVarConstant(mb, getArg(q,2)).val.sval);
+#ifdef _WLR_DEBUG_
+                               fprintf(stderr,"#do transaction tag "LLFMT" 
wlr_limit "LLFMT" wlr_tag "LLFMT"\n", tag, wlr_limit, wlr_tag);
+#endif
                                // break loop if we don't see a the next 
expected transaction
-                               if ( tag != wlr_tag){
-                                               cleanup();
-                                               break;
-                               }
-#ifdef _WLR_DEBUG_
-                               fprintf(stderr,"#redo transaction tag "LLFMT" 
wlr_limit "LLFMT" wlr_tag "LLFMT"\n", tag, wlr_limit, wlr_tag);
-#endif
-                               if ( tag < wlr_tag){
+                               if ( tag <= wlr_tag){
                                        /* skip already executed transaction 
log */
+                                       continue;
                                } else
                                if(  ( tag > wlr_limit) ||
-                                         ( wlr_timelimit[0] && 
strcmp(getVarConstant(mb, getArg(q,2)).val.sval, wlr_timelimit) > 0)){
+                                         ( wlr_timelimit[0] && 
strcmp(tag_read, wlr_timelimit) > 0)){
                                        /* stop execution of the transactions 
if your reached the limit */
-                                       resetMalBlkAndFreeInstructions(mb, 1);
-                                       trimMalVariables(mb, NULL);
-                                       bstream_destroy(c->fdin);
+                                       cleanup();
 #ifdef _WLR_DEBUG_
                                        fprintf(stderr,"#Found final 
transaction "LLFMT"("LLFMT")\n", wlr_limit, wlr_tag);
 #endif
-                                       goto wrapup;
+                                       break;
                                } 
-                       }
-                       if( getModuleId(q) == wlrRef && getFunctionId(q) == 
transactionRef ){
-                               snprintf(wlr_read, sizeof(wlr_read), "%s", 
getVarConstant(mb, getArg(q,2)).val.sval);
 #ifdef _WLR_DEBUG_
                                fprintf(stderr,"#run against tlimit %s  wlr_tag 
"LLFMT"  tag" LLFMT" \n", wlr_timelimit, wlr_tag, tag);
 #endif
@@ -372,7 +377,8 @@ WLRprocessBatch(void *arg)
                                                fprintf(stderr,"#process a 
transaction\n");
                                                fprintFunction(stderr, mb, 0, 
LIST_MAL_DEBUG | LIST_MAL_MAPI );
 #endif
-                                               wlr_tag =  tag + 1;
+                                               wlr_tag =  tag; // remember 
which transaction we executed
+                                               snprintf(wlr_read, 
sizeof(wlr_read), "%s", tag_read);
                                                msg= runMAL(c,mb,0,0);
                                                if( msg == MAL_SUCCEED){
                                                        /* at this point we 
have updated the replica, but the configuration has not been changed.
@@ -404,23 +410,27 @@ WLRprocessBatch(void *arg)
                                        fprintFunction(stderr, mb, 0, 
LIST_MAL_DEBUG );
                                }
                                cleanup();
+                               if ( wlr_tag + 1 == wlc_tag || tag == wlr_limit)
+                                               break;
                        } else
                        if ( getModuleId(q) == wlrRef && getFunctionId(q) == 
rollbackRef ){
                                cleanup();
+                               if ( wlr_tag + 1 == wlc_tag || tag == wlr_limit)
+                                               break;
                        }
                } while(wlr_state != WLR_STOP &&  mb->errors == 0 && msg == 
MAL_SUCCEED);
 #ifdef _WLR_DEBUG_
-               fprintf(stderr,"#wlr.process:processed log file '%s'\n",path);
+               fprintf(stderr,"#wlr.process:processed log file  wlr_tag 
"LLFMT" wlr_limit "LLFMT" time %s\n", wlr_tag, wlr_limit, wlr_timelimit);
 #endif
                // skip to next file when all is read
-               if( c->fdin->eof)
-                       wlr_batches++;
+               wlr_batches++;
                if( msg != MAL_SUCCEED)
                        snprintf(wlr_error, FILENAME_MAX, "%s", msg);
                WLRputConfig();
                bstream_destroy(c->fdin);
+               if ( wlr_tag == wlr_limit)
+                       break;
        }
-wrapup:
        (void) fflush(stderr);
        close_stream(c->fdout);
        SQLexitClient(c);
@@ -481,13 +491,13 @@ WLRprocessScheduler(void *arg)
                                MT_sleep_ms(duration);
                } 
                for( ; duration > 0  && wlr_state != WLR_STOP; duration -= 200){
-                       if ( wlr_tag == wlc_tag || wlr_tag > wlr_limit || 
wlr_limit == -1){
+                       if ( wlr_tag + 1 == wlc_tag || wlr_tag >= wlr_limit || 
wlr_limit == -1){
                                MT_thread_setworking("sleeping");
                                MT_sleep_ms(200);
                        }
                }
                MT_thread_setworking("processing");
-               if( WLRgetMaster() == 0 &&  wlr_limit >= 0 && wlr_tag < wlc_tag)
+               if( WLRgetMaster() == 0 &&  wlr_tag + 1 < wlc_tag && wlr_tag < 
wlr_limit &&  wlr_batches <= wlc_batches)
                        WLRprocessBatch(cntxt);
                
                /* Can not use GDKexiting(), because a test may already reach 
that point before it did anything.
@@ -503,7 +513,7 @@ WLRprocessScheduler(void *arg)
        wlr_thread = 0;
        wlr_state = WLR_WAIT;
 #ifdef _WLR_DEBUG_
-       fprintf(stderr, "#Replicator thread is stopped due to %s\n", wlr_error);
+       fprintf(stderr, "#Replicator thread is stopped \n");
 #endif
 }
 
@@ -512,13 +522,9 @@ WLRprocessScheduler(void *arg)
 str
 WLRmaster(Client cntxt, MalBlkPtr mb, MalStkPtr stk, InstrPtr pci)
 {      
-       str msg = MAL_SUCCEED;
-
-       if( WLRgetConfig())
-               throw(MAL, "sql.replicate", "No replication configuration");
-
-       msg = WLRstart(cntxt, mb, stk, pci);
-       return msg;
+       if( getArgType(mb, pci, 1) == TYPE_str)
+               return WLRstart(cntxt, mb, stk, pci);
+       throw(MAL, "wlr.master", "No master configuration");
 }
 
 str
@@ -532,6 +538,7 @@ WLRreplicate(Client cntxt, MalBlkPtr mb,
        int duration = 2000;
        lng limit;
        str msg = MAL_SUCCEED;
+       (void) cntxt;
 
        if( WLRgetConfig())
                throw(MAL, "sql.replicate", "No replication configuration");
@@ -539,9 +546,6 @@ WLRreplicate(Client cntxt, MalBlkPtr mb,
        if( !wlr_thread)
                throw(MAL, "sql.replicate", "Replicator not started, call 
setmaster() first ");
        
-       if( getArgType(mb, pci, 1) == TYPE_str)
-               msg = WLRstart(cntxt, mb, stk, pci);
-       else
        if( pci->argc == 0)
                wlr_limit = INT64_MAX;
        else
@@ -565,7 +569,7 @@ WLRreplicate(Client cntxt, MalBlkPtr mb,
        if ( limit < 0 && timelimit[0] == 0)
                throw(MAL, "sql.replicate", "Stop tag limit should be positive 
or timestamp should be set");
        if ( limit < INT64_MAX && limit >= wlc_tag)
-               throw(MAL, "sql.replicate", "Stop tag limit be less than 
wlc_tag");
+               throw(MAL, "sql.replicate", "Stop tag limit "LLFMT" be less 
than wlc_tag "LLFMT, limit, wlc_tag);
        if ( limit >= 0)
                wlr_limit = limit;
 
@@ -593,27 +597,28 @@ WLRreplicate(Client cntxt, MalBlkPtr mb,
        fprintf(stderr, "#replicate: wait until wlr_limit = "LLFMT" (tag 
"LLFMT") time %s (%s)\n", wlr_limit, wlr_tag, (wlr_timelimit[0]? 
wlr_timelimit:""), clktxt);
 #endif
 
-       while ( (wlr_tag < wlr_limit + 1)  || (wlr_timelimit[0]  && 
strncmp(clktxt, wlr_timelimit, sizeof(wlr_timelimit)) > 0)  ) {
+       while ( (wlr_tag < wlr_limit )  || (wlr_timelimit[0]  && 
strncmp(clktxt, wlr_timelimit, sizeof(wlr_timelimit)) > 0)  ) {
                if ( wlr_error[0])
                        throw(MAL, "sql.replicate", "tag "LLFMT": %s", wlr_tag, 
wlr_error);
                if ( wlr_tag == wlc_tag)
                        break;
 
 #ifdef _WLR_DEBUG_
-       fprintf(stderr, "#replicate wait state %d wlr_limit "LLFMT" (wlr_tag 
"LLFMT") wlc_tag "LLFMT"\n", wlr_state, wlr_limit, wlr_tag, wlc_tag);
+       fprintf(stderr, "#replicate wait state %d wlr_limit "LLFMT" (wlr_tag 
"LLFMT") wlc_tag "LLFMT" wlr_batches %d\n",
+               wlr_state, wlr_limit, wlr_tag, wlc_tag, wlr_batches);
        fflush(stderr);
 #endif
-               // don't make the sleep too short.
-               MT_sleep_ms( 200);
                if ( !wlr_thread ){
                        if( wlr_error[0])
                                throw(SQL,"wlr.startreplicate",SQLSTATE(42000) 
"Replicator terminated prematurely %s", wlr_error);
                        throw(SQL,"wlr.startreplicate",SQLSTATE(42000) 
"Replicator terminated prematurelys");
                }
 
-               duration -= 100;
+               duration -= 200;
                if ( duration < 0)
                        throw(SQL,"wlr.startreplicate",SQLSTATE(42000) "Timeout 
to wait for replicator to catch up ");
+               // don't make the sleep too short.
+               MT_sleep_ms( 200);
        }
 #ifdef _WLR_DEBUG_
        fprintf(stderr, "#replicate finished "LLFMT" (tag "LLFMT")\n", 
wlr_limit, wlr_tag);
diff --git a/sql/test/wlcr/Tests/All b/sql/test/wlcr/Tests/All
--- a/sql/test/wlcr/Tests/All
+++ b/sql/test/wlcr/Tests/All
@@ -35,6 +35,10 @@ wlr50
 wlc70
 wlr70
 
+# stop the master
+wlc80
+wlr80
+#
 #stop the master
 wlc100
 wlr100
diff --git a/sql/test/wlcr/Tests/wlc80.py b/sql/test/wlcr/Tests/wlc80.py
new file mode 100644
--- /dev/null
+++ b/sql/test/wlcr/Tests/wlc80.py
@@ -0,0 +1,56 @@
+from __future__ import print_function
+
+try:
+    from MonetDBtesting import process
+except ImportError:
+    import process
+import os, sys
+
+dbfarm = os.getenv('GDK_DBFARM')
+tstdb = os.getenv('TSTDB')
+
+if not tstdb or not dbfarm:
+    print('No TSTDB or GDK_DBFARM in environment')
+    sys.exit(1)
+
_______________________________________________
checkin-list mailing list
[email protected]
https://www.monetdb.org/mailman/listinfo/checkin-list

Reply via email to