Changeset: 1b9c81040b65 for MonetDB
URL: https://dev.monetdb.org/hg/MonetDB/rev/1b9c81040b65
Modified Files:
        gdk/gdk_logger.c
        sql/storage/bat/bat_storage.c
Branch: pax-log
Log Message:

first working draft of logger_read_transaction with multi-offsets.


diffs (187 lines):

diff --git a/gdk/gdk_logger.c b/gdk/gdk_logger.c
--- a/gdk/gdk_logger.c
+++ b/gdk/gdk_logger.c
@@ -375,8 +375,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, struct 
offset* offsets, lng* poffset_cnt)
 {
        log_return res = LOG_OK;
        lng nr, pnr;
@@ -433,6 +440,20 @@ log_read_updates(logger *lg, trans *tr, 
                                                res = LOG_ERR;
                                }
                        }
+                       if (offsets && !*(bool*) t /*not deleted*/) {
+                               // This bat actually represents a segment of 
appended rows and we want to collect this oid range
+
+                               lng offset_cnt = *poffset_cnt;
+
+                               const lng previous_os = offset_cnt ? 
offsets[offset_cnt-1].os : 0;
+                               const lng previous_nr = offset_cnt ? 
offsets[offset_cnt-1].nr : 0;
+
+                               const lng os = previous_os + previous_nr;
+                               const lng od = offset;
+                               const struct offset new_value = {os, nr, od};
+                               offsets[offset_cnt] = new_value;
+                               (*poffset_cnt)++;
+                       }
                } else if (l->flag == LOG_UPDATE_BULK) {
                        if (mnstr_readLng(lg->input_log, &offset) != 1) {
                                if (r)
@@ -564,7 +585,32 @@ log_read_updates(logger *lg, trans *tr, 
                }
 
                if (res == LOG_OK) {
-                       if (tr_grow(tr) == GDK_SUCCEED) {
+                       if (l->flag==LOG_UPDATE_BULK && offsets && offset == 
-1) {
+                               for (lng i = 0; i < *poffset_cnt; i++) {
+                                       BAT* rs = BATslice(r, offsets->os, 
offsets->os+offsets->nr);
+                                       if (rs == NULL) {
+                                               res = LOG_ERR;
+                                               break;
+                                       }
+                                       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 = 
offsets[i].od;
+                                               tr->changes[tr->nr].b = rs;
+                                               tr->changes[tr->nr].uid = uid;
+                                               tr->nr++;
+                                       } else {
+                                               BBPreclaim(rs);
+                                               res = LOG_ERR;
+                                       }
+                               }
+                               if (res == LOG_OK)
+                                       BBPreclaim(r);
+                       }
+                       else 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;
@@ -1118,6 +1164,13 @@ logger_read_transaction(logger *lg)
        if (!lg->flushing)
                GDKdebug &= ~CHECKMASK;
 
+       // START variables used in case of LOG_TABLE
+       log_id tid = 0; /*sql table id in case of LOG_TABLE*/
+       struct offset* offsets = NULL;
+       lng nr_offsets = 0;
+       lng cap_offsets = 0;
+       // END variables used in case of LOG_TABLE
+
        while (err == LOG_OK && (ok=log_read_format(lg, &l))) {
                if (l.flag == 0 && l.id == 0) {
                        err = LOG_EOF;
@@ -1168,8 +1221,13 @@ logger_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, 
offsets, &nr_offsets);
+                               if (l.flag == LOG_UPDATE_CONST && offsets && 
nr_offsets == cap_offsets) {
+                                       cap_offsets <<= 1;
+                                       offsets = GDKrealloc(offsets, 
cap_offsets);
+                               }
+                       }
                        break;
                case LOG_CREATE:
                        if (tr == NULL)
@@ -1189,6 +1247,30 @@ logger_read_transaction(logger *lg)
                        else
                                err = log_read_clear(lg, tr, l.id);
                        break;
+               case LOG_TABLE:
+                       if (tr == NULL)
+                               err = LOG_EOF;
+                       else {
+                               if (l.id > 0) {
+                                       // START OF LOG_TABLE
+                                       assert(tid == 0 && offsets == NULL && 
nr_offsets == 0);
+                                       tid = l.id;
+                                       cap_offsets = 10;
+                                       offsets = GDKzalloc(sizeof(*offsets) * 
cap_offsets);
+                                       if (!offsets)
+                                               err = LOG_ERR;
+                               }
+                               else {
+                                       // END OF LOG_TABLE
+                                       assert(tid == -l.id);
+                                       tid = 0;
+                                       GDKfree(offsets);
+                                       offsets = NULL;
+                                       cap_offsets = 0;
+                                       nr_offsets = 0;
+                               }
+                       }
+                       break;
                default:
                        err = LOG_ERR;
                }
diff --git a/sql/storage/bat/bat_storage.c b/sql/storage/bat/bat_storage.c
--- a/sql/storage/bat/bat_storage.c
+++ b/sql/storage/bat/bat_storage.c
@@ -3865,8 +3865,6 @@ log_table_append(sql_trans *tr, sql_tabl
                return LOG_OK;
        size_t end = segs_end(segs, tr, t);
 
-       log_table_start(store->logger, t->base.id);
-
        for (node *n = ol_first_node(t->columns); n && ok; n = n->next) {
                sql_column *c = n->data;
                column_storage *cs = ATOMIC_PTR_GET(&c->data);
@@ -3920,12 +3918,21 @@ log_table_append(sql_trans *tr, sql_tabl
                }
        }
 
-       if (ok == GDK_SUCCEED)
-               ok = log_table_end(store->logger, t->base.id);
-
        return ok == GDK_SUCCEED ? LOG_OK : LOG_ERR;
 }
 
+static inline int
+tr_log_table_start(sql_trans *tr, sql_table *t) {
+       sqlstore *store = tr->store;
+       return log_table_start(store->logger, t->base.id) == GDK_SUCCEED? 
LOG_OK: LOG_ERR;
+}
+
+static inline int
+tr_log_table_end(sql_trans *tr, sql_table *t) {
+       sqlstore *store = tr->store;
+       return log_table_end(store->logger, t->base.id) == GDK_SUCCEED? LOG_OK: 
LOG_ERR;
+}
+
 static int
 log_storage(sql_trans *tr, sql_table *t, storage *s, sqlid id)
 {
@@ -3936,9 +3943,13 @@ log_storage(sql_trans *tr, sql_table *t,
        if (ok == LOG_OK)
                ok = segments2cs(tr, s->segs, &s->cs);
        if (ok == LOG_OK)
+               ok = tr_log_table_start(tr, t);
+       if (ok == LOG_OK)
                ok = log_segments(tr, s->segs, id, &nr_appends);
        if (ok == LOG_OK && !cleared)
                ok = log_table_append(tr, t, s->segs, nr_appends);
+       if (ok == LOG_OK)
+               ok = tr_log_table_end(tr, t);
        return ok;
 }
 
_______________________________________________
checkin-list mailing list -- [email protected]
To unsubscribe send an email to [email protected]

Reply via email to