Changeset: 073e3ea7a180 for MonetDB
URL: http://dev.monetdb.org/hg/MonetDB?cmd=changeset;node=073e3ea7a180
Modified Files:
        monetdb5/mal/mal_client.c
        monetdb5/mal/mal_client.h
        monetdb5/modules/mal/wlcr.c
        monetdb5/modules/mal/wlcr.h
        monetdb5/modules/mal/wlcr.mal
        sql/backends/monet5/sql_wlcr.c
        sql/backends/monet5/sql_wlcr.h
        sql/scripts/60_wlcr.sql
        sql/test/wlcr/Tests/wlc01.py
        sql/test/wlcr/Tests/wlr20.py
        sql/test/wlcr/Tests/wlr30.py
Branch: wlcr
Log Message:

Introduce drift window management
Master and replica can be set to a limited drift.
In the master it implies publishing a new log file regularly.
Setting the drift to zero creates a log file for each transaction.
The drift management requires a separate thread.


diffs (truncated from 761 to 300 lines):

diff --git a/monetdb5/mal/mal_client.c b/monetdb5/mal/mal_client.c
--- a/monetdb5/mal/mal_client.c
+++ b/monetdb5/mal/mal_client.c
@@ -244,9 +244,7 @@ MCinitClientRecord(Client c, oid user, b
        c->error_row = c->error_fld = c->error_msg = c->error_input = NULL;
        c->wlcr_kind = 0;
        c->wlcr_mode = 0;
-       c->wlcr_threshold = 0;
        c->wlcr = NULL;
-       c->wlcr_replaylog = NULL;
 #ifndef HAVE_EMBEDDED /* no authentication in embedded mode */
        {
                str msg = AUTHgetUsername(&c->username, c);
@@ -403,11 +401,7 @@ freeClient(Client c)
                        freeMalBlk(c->wlcr);
                c->wlcr_kind = 0;
                c->wlcr_mode = 0;
-               c->wlcr_threshold = 0;
                c->wlcr = NULL;
-               if( c->wlcr_replaylog)
-                       GDKfree(c->wlcr_replaylog);
-               c->wlcr_replaylog = NULL;
        }
        if (t)
                THRdel(t);  /* you may perform suicide */
diff --git a/monetdb5/mal/mal_client.h b/monetdb5/mal/mal_client.h
--- a/monetdb5/mal/mal_client.h
+++ b/monetdb5/mal/mal_client.h
@@ -173,12 +173,13 @@ typedef struct CLIENT {
        Workset inprogress[THREADS];
        /*
         * The workload for replication/replay is saved initially as a MAL 
block.
+        * It is split into the capturing part (wlc) and the replay part (wlr).
+        * This allows a single server to act as both a master and a replica.
         */
-       int wlcr_kind;  
-       int wlcr_mode;
-       int wlcr_threshold;
-       str wlcr_replaylog;
+       int wlcr_kind;  // used by master to characterise the compound 
transaction
+       int wlcr_mode;  // used by replica to control rerunning the transaction
        MalBlkPtr wlcr;
+
        /*      
         *      Errors during copy into are collected in a user specific column 
set
         */
diff --git a/monetdb5/modules/mal/wlcr.c b/monetdb5/modules/mal/wlcr.c
--- a/monetdb5/modules/mal/wlcr.c
+++ b/monetdb5/modules/mal/wlcr.c
@@ -31,9 +31,10 @@
  * A configuration file is added to keep track on the status of the master.
  * It contains the following key=value pairs:
  *             snapshot=<path to a binary snapshot>
- *             archive=<path to a wlcr log files>
+ *             logs=<path to the wlcr log directory>
  *             start=<first batch file to be applied>
  *             last=<last batch file to be applied>
+ *             drift=<maximal delay before transactions are seen globally, in 
seconds>
  *
  * Every replica should start off with a copy of binary snapshot identified by 
'snapshot'
  * by default stored in .../dbfarm/dbname/master/bat. An alternative path can 
be given
@@ -54,6 +55,15 @@
  * The threshold setting is not saved because it is a client session specific 
action.
  * The default for a production system version is set to -1
  *
+ * A transaction log is owned by the master. He decides when the log may be 
globally
+ * used. There are several triggers for this. A new transaction log is created 
when
+ * the system has been collecting logs for some time (drift).
+ * The problem here is that we should ensure that the log file is closed even 
if there
+ * are no transactions running. After closing, the replicas can see from the
+ * master configuration file that a new batch is available.
+ * The maximum drift can be set using a SQL command. Setting it to zero leads 
to a
+ * log file per transaction.
+ *
  * A more secure way to set a database into master mode is to use the command
  *      monetdb master <dbname> [ <optional snapshot path>]
  * which locks the database, takes a save copy, initializes the state chance. 
@@ -65,12 +75,12 @@
  * master and replica, provided we start with a fresh database.
  *
  * Processing the log files starts in the background using the call.
- * CALL setreplica("mastername")
+ * CALL replicate("mastername")
  * It will iterate through the log files, applying all transactions.
  * Queries are simply ignored unless needed as replacement for catalog actions.
  *
  * The alternative is to replay only the query log
- * CALL replay("dbname",threshold)
+ * CALL replicate("dbname",threshold)
  * In this mode all pure queries are executed under the credentials of the 
query owner
  * for which the reported threshold exceeds the argument[TODO].
  * It excludes catalog and update queries.
@@ -96,14 +106,18 @@ static MT_Lock     wlcr_lock MT_LOCK_INI
 static char *wlcr_name[]= {"","query","update","catalog","ignore"};
 
 static str wlcr_snapshot= 0; // The name of the snapshot against the logs work
-static str wlcr_archive = 0; // The location in the global file store for the 
logs
-static char wlcr_time[26];      // The timestamp of the last committed 
transaction.
+static str wlcr_logs = 0;      // The location in the global file store for 
the logs
+static char wlcr_time[26];     // The timestamp of the last committed 
transaction.
 static stream *wlcr_fd = 0;
+static int wlcr_start = 0;     // time stamp of first transaction in log file
+int wlcr_threshold;
 
-str wlcr_dbname = 0;  // The current database name
-int wlcr_firstbatch = 0;       // first log file  associated with snapshot
-int wlcr_lastbatch = 0;        // last log file identifier 
-int wlcr_tid = 0;      // last transaction id
+// These properties are needed by the replica to direct the roll-forward.
+str wlcr_dbname = 0;           // The master database name
+int wlcr_firstbatch = 0;       // first log file  associated with the snapshot
+int wlcr_batches = 0;          // identifier of next batch
+int wlcr_drift = 10;   // maximal period covered by a single log file in 
seconds
+int wlcr_tid = 0;                      // transaction id of next to be 
processed
 
 /* The database snapshots are binary copies of the dbfarm/database/bat
  * New snapshots are created currently using the 'monetdb snapshot <db>' 
command
@@ -115,7 +129,7 @@ int wlcr_tid = 0;   // last transaction id
 int
 WLCused(void)
 {
-       return wlcr_archive != NULL;
+       return wlcr_logs != NULL;
 }
 
 /* The master configuration file is a simple key=value table */
@@ -123,20 +137,22 @@ str WLCgetConfig(void){
        char path[PATHLENGTH];
        FILE *fd;
 
-       snprintf(path,PATHLENGTH,"%s%cwlcr.config", wlcr_archive, DIR_SEP);
+       snprintf(path,PATHLENGTH,"%s%cwlc.config", wlcr_logs, DIR_SEP);
        fd = fopen(path,"r");
        if( fd == NULL)
                throw(MAL,"wlcr.getConfig","Could not access %s\n",path);
        while( fgets(path, PATHLENGTH, fd) ){
                path[strlen(path)-1] = 0;
-               if( strncmp("archive=", path,8) == 0)
-                       wlcr_archive = GDKstrdup(path + 8);
+               if( strncmp("logs=", path,5) == 0)
+                       wlcr_logs = GDKstrdup(path + 5);
                if( strncmp("snapshot=", path,9) == 0)
                        wlcr_snapshot = GDKstrdup(path + 9);
                if( strncmp("firstbatch=", path,11) == 0)
                        wlcr_firstbatch = atoi(path+ 11);
-               if( strncmp("lastbatch=", path, 10) == 0)
-                       wlcr_lastbatch = atoi(path+ 10);
+               if( strncmp("batches=", path, 8) == 0)
+                       wlcr_batches = atoi(path+ 8);
+               if( strncmp("drift=", path, 6) == 0)
+                       wlcr_drift = atoi(path+ 6);
        }
        fclose(fd);
        return MAL_SUCCEED;
@@ -148,21 +164,22 @@ str WLCsetConfig(void){
        FILE *fd;
 
        // default setting for the archive directory is the db itself
-       if( wlcr_archive == NULL){
-               snprintf(path,PATHLENGTH,"%s%c", wlcr_archive, DIR_SEP);
-               wlcr_archive = GDKstrdup(path);
+       if( wlcr_logs == NULL){
+               snprintf(path,PATHLENGTH,"%s%c", wlcr_logs, DIR_SEP);
+               wlcr_logs = GDKstrdup(path);
        }
 
-       snprintf(path,PATHLENGTH,"%s%cwlcr.config", wlcr_archive, DIR_SEP);
+       snprintf(path,PATHLENGTH,"%s%cwlc.config", wlcr_logs, DIR_SEP);
        fd = fopen(path,"w");
        if( fd == NULL)
                throw(MAL,"wlcr.setConfig","Could not access %s\n",path);
        if( wlcr_snapshot)
                fprintf(fd,"snapshot=%s\n", wlcr_snapshot);
-       if( wlcr_archive)
-               fprintf(fd,"archive=%s\n", wlcr_archive);
+       if( wlcr_logs)
+               fprintf(fd,"logs=%s\n", wlcr_logs);
        fprintf(fd,"firstbatch=%d\n", wlcr_firstbatch);
-       fprintf(fd,"lastbatch=%d\n", wlcr_lastbatch );
+       fprintf(fd,"batches=%d\n", wlcr_batches );
+       fprintf(fd,"drift=%d\n", wlcr_drift );
        fclose(fd);
        return MAL_SUCCEED;
 }
@@ -174,36 +191,75 @@ WLCsetlogger(void)
 {
        char path[PATHLENGTH];
 
-       if( wlcr_archive == NULL)
+       if( wlcr_logs == NULL)
                throw(MAL,"wlcr.setlogger","Path not initalized");
        MT_lock_set(&wlcr_lock);
-       snprintf(path,PATHLENGTH,"%s%c%s_%012d", wlcr_archive, DIR_SEP, 
wlcr_dbname, wlcr_lastbatch);
+       snprintf(path,PATHLENGTH,"%s%c%s_%012d", wlcr_logs, DIR_SEP, 
wlcr_dbname, wlcr_batches);
        wlcr_fd = open_wastream(path);
        if( wlcr_fd == 0){
                MT_lock_unset(&wlcr_lock);
+               GDKerror("wlcr.logger:Could not create %s\n",path);
                throw(MAL,"wlcr.logger","Could not create %s\n",path);
        }
 
-       wlcr_lastbatch++;
+       wlcr_batches++;
        wlcr_tid = 0;
+       wlcr_start = GDKms()/1000;
        WLCsetConfig();
        MT_lock_unset(&wlcr_lock);
        return MAL_SUCCEED;
 }
 
+static void
+WLCcloselogger(void)
+{
+       if( wlcr_fd == NULL)
+               return ;
+       close_stream(wlcr_fd);
+       wlcr_fd= NULL;
+       wlcr_tid = 0;
+       wlcr_start = 0;
+       WLCsetConfig();
+}
+
+/*
+ * The WLCRlogger process ensures that log files are properly closed
+ * and released when their drift time window has expired.
+ */
+
+static MT_Id wlcr_logger;
+
+static void
+WLCRlogger(void *arg)
+{
+       (void) arg;
+       while(1){
+               if( wlcr_logs && wlcr_fd ){
+                       if (wlcr_start + wlcr_drift < GDKms() / 1000){
+                               MT_lock_set(&wlcr_lock);
+                               WLCcloselogger();
+                               MT_lock_unset(&wlcr_lock);
+                       }
+                       MT_sleep_ms( (wlcr_drift? wlcr_drift:1 ) * 1000);
+               } else
+               if( wlcr_drift)
+                               MT_sleep_ms( wlcr_drift * 1000);
+               else
+                               MT_sleep_ms(  10  * 1000);
+       }
+}
 /*
  * The existence of the master directory should be checked upon server restart.
- * A new batch file should be created as a result.
- * Upon exit we should check the log file size. If empty we need not safe it. 
[TODO]
+ * Then the master record information should be set and the WLClogger started.
  */
 str 
 WLCinit(Client cntxt)
 {
        char path[PATHLENGTH];
-       str pathname;
+       str pathname, msg= MAL_SUCCEED;
        FILE *fd;
 
-       if( wlcr_archive){
+       if( wlcr_logs){
 #ifdef _WLC_DEBUG_
                mnstr_printf(cntxt->fdout,"#WLC already running\n");
 #else
@@ -212,16 +268,21 @@ WLCinit(Client cntxt)
        } else{
                // use default location for archive
                pathname = GDKfilepath(0,0,"master",0);
-               snprintf(path, PATHLENGTH,"%s%cwlcr.config", pathname, DIR_SEP);
+               snprintf(path, PATHLENGTH,"%s%cwlc.config", pathname, DIR_SEP);
 
                fd = fopen(path,"r");
-               if( fd == NULL) // no master mode
+               if( fd == NULL) // not in master mode
                        return MAL_SUCCEED;
                fclose(fd);
+               // we are in master mode
                wlcr_dbname = GDKgetenv("gdk_dbname");
-               wlcr_archive = pathname;
-               (void) WLCgetConfig();
-               return WLCsetlogger();
+               wlcr_logs = pathname;
+               msg =  WLCgetConfig();
+               if( msg)
+                       GDKerror("%s",msg);
+               if (MT_create_thread(&wlcr_logger, WLCRlogger , (void*) 0, 
MT_THR_JOINABLE) < 0) {
+                GDKerror("wlcr.logger thread could not be spawned");
+        }
        }
        return MAL_SUCCEED;
 }
@@ -229,16 +290,7 @@ WLCinit(Client cntxt)
 str 
 WLCexit(void)
 {
-       size_t sz;
-
_______________________________________________
checkin-list mailing list
[email protected]
https://www.monetdb.org/mailman/listinfo/checkin-list

Reply via email to