Changeset: 370321b0627c for MonetDB
URL: https://dev.monetdb.org/hg/MonetDB/rev/370321b0627c
Branch: default
Log Message:

Merge with pax-log.


diffs (truncated from 575 to 300 lines):

diff --git a/clients/Tests/exports.stable.out b/clients/Tests/exports.stable.out
--- a/clients/Tests/exports.stable.out
+++ b/clients/Tests/exports.stable.out
@@ -541,8 +541,10 @@ ssize_t lngToStr(str *dst, size_t *len, 
 const lng lng_nil;
 struct tm *localtime_r(const time_t *restrict, struct tm *restrict);
 gdk_return log_activate(logger *lg);
-gdk_return log_bat(logger *lg, BAT *b, log_id id, lng offset, lng cnt);
+gdk_return log_bat(logger *lg, BAT *b, log_id id, lng offset, lng cnt, lng 
total_cnt);
 gdk_return log_bat_clear(logger *lg, log_id id);
+gdk_return log_bat_group_end(logger *lg, log_id id);
+gdk_return log_bat_group_start(logger *lg, log_id id);
 gdk_return log_bat_persists(logger *lg, BAT *b, log_id id);
 gdk_return log_bat_transient(logger *lg, log_id id);
 lng log_changes(logger *lg);
diff --git a/gdk/gdk_logger.c b/gdk/gdk_logger.c
--- a/gdk/gdk_logger.c
+++ b/gdk/gdk_logger.c
@@ -34,7 +34,8 @@ static gdk_return log_del_bat(logger *lg
 #define LOG_DESTROY    6
 #define LOG_SEQ                7
 #define LOG_CLEAR      8
-#define LOG_ROW                9 /* per row relative small log entry */
+#define LOG_BAT_GROUP  9
+#define LOG_COMMIT     10
 
 #ifdef NATIVE_WIN32
 #define getfilepos _ftelli64
@@ -60,7 +61,7 @@ static const char *log_commands[] = {
        "LOG_DESTROY",
        "LOG_SEQ",
        "LOG_CLEAR",
-       "LOG_ROW",
+       "LOG_BAT_GROUP",
 };
 
 typedef struct logaction {
@@ -375,8 +376,15 @@ string_reader(logger *lg, BAT *b, lng nr
        return res;
 }
 
+
+struct offset {
+       lng os /*offset within source BAT in logfile */;
+       lng nr /*number of values to be copied*/;
+       lng od /*offset within destination BAT in database*/;
+};
+
 static log_return
-log_read_updates(logger *lg, trans *tr, logformat *l, log_id id)
+log_read_updates(logger *lg, trans *tr, logformat *l, log_id id, BAT** cands)
 {
        log_return res = LOG_OK;
        lng nr, pnr;
@@ -406,6 +414,53 @@ log_read_updates(logger *lg, trans *tr, 
                                return LOG_ERR;
                        }
                }
+
+               if (l->flag == LOG_UPDATE_CONST) {
+                       if (mnstr_readLng(lg->input_log, &offset) != 1)
+                               return LOG_ERR;
+                       if (cands) {
+                               // This const range actually represents a 
segment of candidates corresponding to updated bat entries
+
+                               if (BATcount(*cands) == 0 || lg->flushing) {
+                                       // when flushing, we only need the 
offset and count of the last segment of inserts.
+                                       assert((*cands)->ttype == TYPE_void);
+                                       BATtseqbase(*cands, (oid) offset);
+                                       BATsetcount(*cands, (BUN) nr);
+                               }
+                               else if (!lg->flushing) {
+                                       assert(BATcount(*cands) > 0);
+                                       BAT* dense = BATdense(0, (oid) offset, 
(BUN) nr);
+                                       BAT* newcands = NULL;
+                                       if (!dense ) {
+                                               res = LOG_ERR;
+                                       }
+                                       else if ((*cands)->ttype == TYPE_void) {
+                                               if ( (newcands = 
BATmergecand(*cands, dense)) ) {
+                                                       BBPreclaim(*cands);
+                                                       *cands = newcands;
+                                               }
+                                               else
+                                                       res = LOG_ERR;
+                                       }
+                                       else {
+                                               assert((*cands)->ttype == 
TYPE_oid);
+                                               assert(BATcount(*cands) > 0);
+                                               if (BATappend(*cands, dense, 
NULL, true) != GDK_SUCCEED)
+                                                       res = LOG_ERR;
+                                       }
+                                       BBPreclaim(dense);
+                               }
+
+                               // We have to read the value to update the read 
cursor
+                               size_t tlen = lg->bufsize;
+                               void *t = rt(lg->buf, &tlen, lg->input_log, 1);
+                               if (t == NULL) {
+                                       res = LOG_ERR;
+                               }
+                               return res;
+                       }
+               }
+
                if (!lg->flushing) {
                        r = COLnew(0, tpe, (BUN) nr, PERSISTENT);
                        if (r == NULL) {
@@ -416,11 +471,6 @@ log_read_updates(logger *lg, trans *tr, 
                }
 
                if (l->flag == LOG_UPDATE_CONST) {
-                       if (mnstr_readLng(lg->input_log, &offset) != 1) {
-                               if (r)
-                                       BBPreclaim(r);
-                               return LOG_ERR;
-                       }
                        size_t tlen = lg->bufsize;
                        void *t = rt(lg->buf, &tlen, lg->input_log, 1);
                        if (t == NULL) {
@@ -563,25 +613,43 @@ log_read_updates(logger *lg, trans *tr, 
                        GDKfree(hv);
                }
 
-               if (res == LOG_OK) {
-                       if (tr_grow(tr) == GDK_SUCCEED) {
-                               tr->changes[tr->nr].type =
-                                       
l->flag==LOG_UPDATE_CONST?LOG_UPDATE_BULK:l->flag;
-                               tr->changes[tr->nr].nr = pnr;
-                               tr->changes[tr->nr].tt = tpe;
-                               tr->changes[tr->nr].cid = id;
-                               tr->changes[tr->nr].offset = offset;
-                               tr->changes[tr->nr].b = r;
-                               tr->changes[tr->nr].uid = uid;
-                               tr->nr++;
-                       } else {
-                               res = LOG_ERR;
+               if (res == LOG_OK && tr_grow(tr) == GDK_SUCCEED) {
+                       tr->changes[tr->nr].type = l->flag;
+                       if (l->flag==LOG_UPDATE_BULK && offset == -1) {
+                               assert(cands); // bat r is part of a group of 
bats logged together.
+                               struct canditer ci;
+                               canditer_init(&ci, NULL, *cands);
+                               const oid first = canditer_peek(&ci);
+                               const oid last = canditer_last(&ci);
+                               offset = (lng) first;
+                               pnr = (lng) (last - first) + 1;
+                               if (!lg->flushing ) {
+                                       assert(uid == NULL);
+                                       uid = *cands;
+                                       BBPfix((*cands)->batCacheid);
+                                       tr->changes[tr->nr].type = LOG_UPDATE;
+                               }
                        }
+                       if (l->flag==LOG_UPDATE_CONST) {
+                               assert(!cands); // TODO: This might change in 
the future.
+                               tr->changes[tr->nr].type = LOG_UPDATE_BULK;
+                       }
+                       tr->changes[tr->nr].nr = pnr;
+                       tr->changes[tr->nr].tt = tpe;
+                       tr->changes[tr->nr].cid = id;
+                       tr->changes[tr->nr].offset = offset;
+                       tr->changes[tr->nr].b = r;
+                       tr->changes[tr->nr].uid = uid;
+                       tr->nr++;
+               } else {
+                       res = LOG_ERR;
                }
                if (res == LOG_ERR) {
                        if (r)
                                BBPreclaim(r);
-                       if (uid)
+                       if (cands && uid)
+                               BBPunfix((*cands)->batCacheid);
+                       else if (uid)
                                BBPreclaim(uid);
                }
        } else {
@@ -635,9 +703,8 @@ la_bat_updates(logger *lg, logaction *la
                if (b == NULL)
                        return GDK_FAIL;
        }
+       BUN cnt = 0;
        if (la->type == LOG_UPDATE_BULK) {
-               BUN cnt = 0;
-
                if (!lg->flushing) {
                        cnt = BATcount(b);
                        int is_msk = (b->ttype == TYPE_msk);
@@ -685,27 +752,16 @@ la_bat_updates(logger *lg, logaction *la
                                bat_iterator_end(&vi);
                        }
                }
-               cnt = (BUN)(la->offset + la->nr);
-               if (la_bat_update_count(lg, la->cid, cnt, tid) != GDK_SUCCEED) {
-                       if (b)
-                               logbat_destroy(b);
+       } else if (la->type == LOG_UPDATE) {
+               if (!lg->flushing && BATupdate(b, la->uid, la->b, true) != 
GDK_SUCCEED) {
                        return GDK_FAIL;
                }
-       } else if (!lg->flushing && la->type == LOG_UPDATE) {
-               BATiter vi = bat_iterator(la->b);
-               BUN p, q;
-
-               BATloop(la->b, p, q) {
-                       oid h = BUNtoid(la->uid, p);
-                       const void *t = BUNtail(vi, p);
-
-                       if (BUNreplace(b, h, t, true) != GDK_SUCCEED) {
-                               logbat_destroy(b);
-                               bat_iterator_end(&vi);
-                               return GDK_FAIL;
-                       }
-               }
-               bat_iterator_end(&vi);
+       }
+       cnt = (BUN)(la->offset + la->nr);
+       if (la_bat_update_count(lg, la->cid, cnt, tid) != GDK_SUCCEED) {
+               if (b)
+                       logbat_destroy(b);
+               return GDK_FAIL;
        }
        if (b)
                logbat_destroy(b);
@@ -1118,6 +1174,8 @@ log_read_transaction(logger *lg)
        if (!lg->flushing)
                GDKdebug &= ~CHECKMASK;
 
+       BAT* cands = NULL; // used in case of LOG_BAT_GROUP
+
        while (err == LOG_OK && (ok=log_read_format(lg, &l))) {
                if (l.flag == 0 && l.id == 0) {
                        err = LOG_EOF;
@@ -1168,8 +1226,9 @@ log_read_transaction(logger *lg)
                case LOG_UPDATE:
                        if (tr == NULL)
                                err = LOG_EOF;
-                       else
-                               err = log_read_updates(lg, tr, &l, l.id);
+                       else {
+                               err = log_read_updates(lg, tr, &l, l.id, 
cands?&cands:NULL);
+                       }
                        break;
                case LOG_CREATE:
                        if (tr == NULL)
@@ -1189,6 +1248,23 @@ log_read_transaction(logger *lg)
                        else
                                err = log_read_clear(lg, tr, l.id);
                        break;
+               case LOG_BAT_GROUP:
+                       if (tr == NULL)
+                               err = LOG_EOF;
+                       else {
+                               if (l.id > 0) {
+                                       // START OF LOG_BAT_GROUP
+                                       cands = COLnew(0, TYPE_void, 0, 
TRANSIENT);
+                                       if (!cands)
+                                               err = LOG_ERR;
+                               }
+                               else {
+                                       // END OF LOG_BAT_GROUP
+                                       BBPunfix(cands->batCacheid);
+                                       cands = NULL;
+                               }
+                       }
+                       break;
                default:
                        err = LOG_ERR;
                }
@@ -1202,6 +1278,9 @@ log_read_transaction(logger *lg)
                tr = tr_abort(lg, tr);
        if (!lg->flushing)
                GDKdebug = dbg;
+
+       if (cands)
+               GDKfree(cands);
        if (!ok)
                return LOG_EOF;
        return err;
@@ -2452,7 +2531,7 @@ string_writer(logger *lg, BAT *b, lng of
 }
 
 static gdk_return
-internal_log_bat(logger *lg, BAT *b, log_id id, lng offset, lng cnt, int 
sliced)
+internal_log_bat(logger *lg, BAT *b, log_id id, lng offset, lng cnt, int 
sliced, lng total_cnt)
 {
        bte tpe = find_type(lg, b->ttype);
        gdk_return ok = GDK_SUCCEED;
@@ -2473,13 +2552,19 @@ internal_log_bat(logger *lg, BAT *b, log
 
        gdk_return (*wt) (const void *, stream *, size_t) = 
BATatoms[b->ttype].atomWrite;
 
-       if (log_write_format(lg, &l) != GDK_SUCCEED ||
-           !mnstr_writeLng(lg->output_log, nr) ||
-           mnstr_write(lg->output_log, &tpe, 1, 1) != 1 ||
-           !mnstr_writeLng(lg->output_log, offset)) {
-               ok = GDK_FAIL;
-               goto bailout;
-       }
_______________________________________________
checkin-list mailing list -- [email protected]
To unsubscribe send an email to [email protected]

Reply via email to