Changeset: 23362f63873f for MonetDB
URL: https://dev.monetdb.org/hg/MonetDB/rev/23362f63873f
Modified Files:
        sql/backends/monet5/sql.c
        sql/storage/bat/bat_storage.c
        sql/storage/bat/bat_storage.h
        sql/storage/store.c
Branch: Jul2021
Log Message:

beter split of mergering updates and garbage collection the delta structures. 
Solves crash in concurrent update cases.


diffs (259 lines):

diff --git a/sql/backends/monet5/sql.c b/sql/backends/monet5/sql.c
--- a/sql/backends/monet5/sql.c
+++ b/sql/backends/monet5/sql.c
@@ -2150,6 +2150,13 @@ DELTAsub(bat *result, const bat *col, co
                                BBPunfix(u->batCacheid);
                                throw(MAL, "sql.delta", SQLSTATE(HY013) 
MAL_MALLOC_FAIL);
                        }
+                       BAT *nres;
+                       if ((nres = COLcopy(res, res->ttype, true, TRANSIENT)) 
== NULL) {
+                               BBPunfix(res->batCacheid);
+                               throw(MAL, "sql.delta", GDK_EXCEPTION);
+                       }
+                       BBPunfix(res->batCacheid);
+                       res = nres;
                        ret = BATappend(res, u, cminu, true);
                        BBPunfix(u->batCacheid);
                        if (cminu)
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
@@ -841,13 +841,13 @@ older_delta( sql_delta *d, sql_trans *tr
 {
        sql_delta *o = d->next;
 
-       while (o) {
-               if (o->cs.ucnt && VALID_4_READ(o->cs.ts, tr))
+       while (o && !o->cs.merged) {
+               if (o->cs.ucnt && VALID_4_READ(o->cs.ts, tr))
                        break;
                else
                        o = o->next;
        }
-       if (o && o->cs.ucnt && VALID_4_READ(o->cs.ts, tr))
+       if (o && !o->cs.merged && o->cs.ucnt && VALID_4_READ(o->cs.ts, tr))
                return o;
        return NULL;
 }
@@ -1699,6 +1699,7 @@ append_col_execute(sql_trans *tr, sql_de
 {
        int ok = LOG_OK;
 
+       delta->cs.merged = 0;
        if (is_bat) {
                BAT *bat = incoming_data;
 
@@ -2327,7 +2328,7 @@ commit_create_col_( sql_trans *tr, sql_c
                delta->cs.ts = commit_ts;
 
                assert(delta->next == NULL);
-               if (!delta->cs.alter)
+               if (!delta->cs.alter && !delta->cs.merged)
                        ok = merge_delta(delta);
                delta->cs.alter = 0;
                if (!tr->parent)
@@ -2433,7 +2434,8 @@ commit_create_idx_( sql_trans *tr, sql_i
                delta->cs.ts = commit_ts;
 
                assert(delta->next == NULL);
-               ok = merge_delta(delta);
+               if (!delta->cs.alter && !delta->cs.merged)
+                       ok = merge_delta(delta);
                if (!tr->parent)
                        i->base.new = 0;
        }
@@ -3139,6 +3141,7 @@ merge_cs( column_storage *cs)
                bat_destroy(uv);
        }
        cs->cleared = 0;
+       cs->merged = 1;
        bat_destroy(cur);
        return ok;
 }
@@ -3146,6 +3149,8 @@ merge_cs( column_storage *cs)
 static int
 merge_delta( sql_delta *obat)
 {
+       if (obat && obat->next && !obat->cs.merged)
+               merge_delta(obat->next);
        return merge_cs(&obat->cs);
 }
 
@@ -3217,9 +3222,10 @@ commit_update_col_( sql_trans *tr, sql_c
        (void)oldest;
        if (isTempTable(c->t)) {
                if (commit_ts) { /* commit */
-                       if (c->t->commit_action == CA_COMMIT || 
c->t->commit_action == CA_PRESERVE)
-                               ok = merge_delta(delta);
-                       else /* CA_DELETE as CA_DROP's are gone already (or for 
globals are equal to a CA_DELETE) */
+                       if (c->t->commit_action == CA_COMMIT || 
c->t->commit_action == CA_PRESERVE) {
+                               if (!delta->cs.merged)
+                                       ok = merge_delta(delta);
+                       } else /* CA_DELETE as CA_DROP's are gone already (or 
for globals are equal to a CA_DELETE) */
                                clear_cs(tr, &delta->cs, true, 
isTempTable(c->t));
                } else { /* rollback */
                        if (c->t->commit_action == CA_COMMIT/* || 
c->t->commit_action == CA_PRESERVE*/)
@@ -3290,20 +3296,12 @@ commit_update_col( sql_trans *tr, sql_ch
                        o->next = d->next;
                change->cleanup = &tc_gc_rollbacked;
        } else if (ok == LOG_OK && !tr->parent) {
-               sql_delta *d = delta;
-               /* clean up and merge deltas */
-               while (delta && delta->cs.ts > oldest) {
+               /* merge deltas */
+               while (delta && delta->cs.ts > oldest)
                        delta = delta->next;
-               }
-               if (delta && delta != d) {
-                       if (delta->next) {
-                               ok = destroy_delta(delta->next, true);
-                               delta->next = NULL;
-                       }
-               }
-               if (ok == LOG_OK && delta == d && oldest == commit_ts) {
-                       lock_column(tr->store, c->base.id);
-                       ok = merge_delta(delta);
+               if (ok == LOG_OK && delta && !delta->cs.merged && delta->cs.ts 
<= oldest) {
+                       lock_column(tr->store, c->base.id); /* lock for 
concurrent updates (appends) */
+                       merge_delta(delta);
                        unlock_column(tr->store, c->base.id);
                }
        } else if (ok == LOG_OK && tr->parent) /* move delta into older and 
cleanup current save points */
@@ -3333,9 +3331,10 @@ commit_update_idx_( sql_trans *tr, sql_i
        (void)oldest;
        if (isTempTable(i->t)) {
                if (commit_ts) { /* commit */
-                       if (i->t->commit_action == CA_COMMIT || 
i->t->commit_action == CA_PRESERVE)
-                               ok = merge_delta(delta);
-                       else /* CA_DELETE as CA_DROP's are gone already */
+                       if (i->t->commit_action == CA_COMMIT || 
i->t->commit_action == CA_PRESERVE) {
+                               if (!delta->cs.merged)
+                                       ok = merge_delta(delta);
+                       } else /* CA_DELETE as CA_DROP's are gone already */
                                clear_cs(tr, &delta->cs, true, 
isTempTable(i->t));
                } else { /* rollback */
                        if (i->t->commit_action == CA_COMMIT/* || 
i->t->commit_action == CA_PRESERVE*/)
@@ -3375,20 +3374,12 @@ commit_update_idx( sql_trans *tr, sql_ch
                        o->next = d->next;
                change->cleanup = &tc_gc_rollbacked;
        } else if (ok == LOG_OK && !tr->parent) {
-               sql_delta *d = delta;
-               /* clean up and merge deltas */
-               while (delta && delta->cs.ts > oldest) {
+               /* merge deltas */
+               while (delta && delta->cs.ts > oldest)
                        delta = delta->next;
-               }
-               if (delta && delta != d) {
-                       if (delta->next) {
-                               ok = destroy_delta(delta->next, true);
-                               delta->next = NULL;
-                       }
-               }
-               if (ok == LOG_OK && delta == d && oldest == commit_ts) {
-                       lock_column(tr->store, i->base.id);
-                       ok = merge_delta(delta);
+               if (ok == LOG_OK && delta && !delta->cs.merged && delta->cs.ts 
<= oldest) {
+                       lock_column(tr->store, i->base.id); /* lock for 
concurrent updates (appends) */
+                       merge_delta(delta);
                        unlock_column(tr->store, i->base.id);
                }
        } else if (ok == LOG_OK && tr->parent) /* cleanup older save points */
@@ -3500,6 +3491,18 @@ commit_update_del( sql_trans *tr, sql_ch
        return ok;
 }
 
+static int
+gc_delta( sql_store Store, sql_change *change, ulng oldest)
+{
+       sqlstore *store = Store;
+       sql_delta *n = change->data;
+       (void)store;
+       (void)oldest;
+
+       destroy_delta(n, true);
+       return 1;
+}
+
 /* only rollback (content version) case for now */
 static int
 gc_col( sqlstore *store, sql_change *change, ulng oldest, bool cleanup)
@@ -3516,11 +3519,24 @@ gc_col( sqlstore *store, sql_change *cha
                return 0;
        sql_delta *d = (sql_delta*)change->data;
        if (d->next) {
+               assert(!cleanup);
                if (d->cs.ts > oldest)
                        return LOG_OK; /* cannot cleanup yet */
 
-               destroy_delta(d->next, true);
+               sql_delta *n = d->next;
+               if (n->cs.ucnt && !n->cs.merged) {
+                       lock_column(store, c->base.id); /* lock for concurrent 
updates (appends) */
+                       merge_delta(n);
+                       unlock_column(store, c->base.id);
+               } else if (d && d->cs.ucnt && !d->cs.merged) {
+                       lock_column(store, c->base.id); /* lock for concurrent 
updates (appends) */
+                       merge_delta(d);
+                       unlock_column(store, c->base.id);
+               }
                d->next = NULL;
+               change->cleanup = &gc_delta;
+               change->data = n;
+               return 0;
        }
        if (cleanup)
                column_destroy(store, c);
@@ -3555,11 +3571,24 @@ gc_idx( sqlstore *store, sql_change *cha
                return 0;
        sql_delta *d = (sql_delta*)change->data;
        if (d->next) {
+               assert(!cleanup);
                if (d->cs.ts > oldest)
                        return LOG_OK; /* cannot cleanup yet */
 
-               destroy_delta(d->next, true);
+               sql_delta *n = d->next;
+               if (n->cs.ucnt && !n->cs.merged) {
+                       lock_column(store, i->base.id); /* lock for concurrent 
updates (appends) */
+                       merge_delta(n);
+                       unlock_column(store, i->base.id);
+               } else if (d && d->cs.ucnt && !d->cs.merged) {
+                       lock_column(store, i->base.id); /* lock for concurrent 
updates (appends) */
+                       merge_delta(d);
+                       unlock_column(store, i->base.id);
+               }
                d->next = NULL;
+               change->cleanup = &gc_delta;
+               change->data = n;
+               return 0;
        }
        if (cleanup)
                idx_destroy(store, i);
diff --git a/sql/storage/bat/bat_storage.h b/sql/storage/bat/bat_storage.h
--- a/sql/storage/bat/bat_storage.h
+++ b/sql/storage/bat/bat_storage.h
@@ -19,6 +19,7 @@ typedef struct column_storage {
        int uvbid;              /* bat with values of updates */
        bool cleared;
        bool alter;             /* set when the delta is created for an alter 
statement */
+       bool merged;    /* only merge changes once */
        size_t ucnt;    /* number of updates */
        ulng ts;                /* version timestamp */
 } column_storage;
diff --git a/sql/storage/store.c b/sql/storage/store.c
--- a/sql/storage/store.c
+++ b/sql/storage/store.c
@@ -5707,11 +5707,10 @@ sql_trans_drop_column(sql_trans *tr, sql
                        n = nn;
                        col = next;
                } else if (col) { /* if the column to be dropped was found, 
decrease the column number for others after it */
-                       oid rid;
                        next->colnr--;
 
                        if (!isDeclaredTable(t)) {
-                               rid = store->table_api.column_find_row(tr, cid, 
&next->base.id, NULL);
+                               oid rid = store->table_api.column_find_row(tr, 
cid, &next->base.id, NULL);
                                assert(!is_oid_nil(rid));
                                if ((res = 
store->table_api.column_update_value(tr, cnr, rid, &next->colnr)))
                                        return res;
_______________________________________________
checkin-list mailing list
[email protected]
https://www.monetdb.org/mailman/listinfo/checkin-list

Reply via email to