Changeset: f30b78a24118 for MonetDB
URL: https://dev.monetdb.org/hg/MonetDB/rev/f30b78a24118
Modified Files:
        sql/backends/monet5/sql.c
        sql/backends/monet5/wlr.c
        sql/storage/bat/bat_storage.c
        sql/test/miscellaneous/Tests/transaction_isolation.SQL.py
Branch: Jul2021
Log Message:

Throw better error message on update conflicts and small cleanup


diffs (289 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
@@ -1726,10 +1726,11 @@ mvc_append_wrap(Client cntxt, MalBlkPtr 
        const char *cname = *getArgReference_str(stk, pci, 4);
        lng pos = *(lng*)getArgReference_lng(stk, pci, 5);
        ptr ins = getArgReference(stk, pci, 6);
-       int tpe = getArgType(mb, pci, 6), err = 0;
+       int tpe = getArgType(mb, pci, 6), log_res = LOG_OK;
        sql_schema *s;
        sql_table *t;
        sql_column *c;
+       sql_idx *i;
        BAT *b = 0;
 
        *res = 0;
@@ -1762,18 +1763,14 @@ mvc_append_wrap(Client cntxt, MalBlkPtr 
                BATmsync(b);
        sqlstore *store = m->session->tr->store;
        if (cname[0] != '%' && (c = mvc_bind_column(m, t, cname)) != NULL) {
-               if (store->storage_api.append_col(m->session->tr, c, 
(size_t)pos, ins, tpe, 1) != LOG_OK)
-                       err = 1;
-       } else if (cname[0] == '%') {
-               sql_idx *i = mvc_bind_idx(m, s, cname + 1);
-               if (i && store->storage_api.append_idx(m->session->tr, i, 
(size_t)pos, ins, tpe, 1) != LOG_OK)
-                       err = 1;
+               log_res = store->storage_api.append_col(m->session->tr, c, 
(size_t)pos, ins, tpe, 1);
+       } else if (cname[0] == '%' && (i = mvc_bind_idx(m, s, cname + 1)) != 
NULL) {
+               log_res = store->storage_api.append_idx(m->session->tr, i, 
(size_t)pos, ins, tpe, 1);
        }
-       if (err)
-               throw(SQL, "sql.append", SQLSTATE(42S02) "append failed");
-       if (b) {
+       if (b)
                BBPunfix(b->batCacheid);
-       }
+       if (log_res != LOG_OK) /* the conflict case should never happen, but 
leave it here */
+               throw(SQL, "sql.append", SQLSTATE(42000) "Append failed%s", 
log_res == LOG_CONFLICT ? " due to conflict with another transaction" : "");
        return MAL_SUCCEED;
 }
 
@@ -1790,10 +1787,11 @@ mvc_update_wrap(Client cntxt, MalBlkPtr 
        bat Tids = *getArgReference_bat(stk, pci, 5);
        bat Upd = *getArgReference_bat(stk, pci, 6);
        BAT *tids, *upd;
-       int tpe = getArgType(mb, pci, 6), err = 0;
+       int tpe = getArgType(mb, pci, 6), log_res = LOG_OK;
        sql_schema *s;
        sql_table *t;
        sql_column *c;
+       sql_idx *i;
 
        *res = 0;
        if ((msg = getSQLContext(cntxt, mb, &m, NULL)) != NULL)
@@ -1833,17 +1831,14 @@ mvc_update_wrap(Client cntxt, MalBlkPtr 
                BATmsync(tids);
        sqlstore *store = m->session->tr->store;
        if (cname[0] != '%' && (c = mvc_bind_column(m, t, cname)) != NULL) {
-               if (store->storage_api.update_col(m->session->tr, c, tids, upd, 
TYPE_bat) != LOG_OK)
-                       err = 1;
-       } else if (cname[0] == '%') {
-               sql_idx *i = mvc_bind_idx(m, s, cname + 1);
-               if (i && store->storage_api.update_idx(m->session->tr, i, tids, 
upd, TYPE_bat) != LOG_OK)
-                       err = 1;
+               log_res = store->storage_api.update_col(m->session->tr, c, 
tids, upd, TYPE_bat);
+       } else if (cname[0] == '%' && (i = mvc_bind_idx(m, s, cname + 1)) != 
NULL) {
+               log_res = store->storage_api.update_idx(m->session->tr, i, 
tids, upd, TYPE_bat);
        }
        BBPunfix(tids->batCacheid);
        BBPunfix(upd->batCacheid);
-       if (err)
-               throw(SQL, "sql.update", SQLSTATE(42S02) "update failed");
+       if (log_res != LOG_OK)
+               throw(SQL, "sql.update", SQLSTATE(42000) "Update failed%s", 
log_res == LOG_CONFLICT ? " due to conflict with another transaction" : "");
        return MAL_SUCCEED;
 }
 
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
@@ -851,11 +851,12 @@ str
 WLRappend(Client cntxt, MalBlkPtr mb, MalStkPtr stk, InstrPtr pci)
 {
        str sname, tname, cname;
-       int tpe,i;
+       int tpe,i, log_res = LOG_OK;
        mvc *m=NULL;
        sql_schema *s;
        sql_table *t;
        sql_column *c;
+       sql_idx *idx;
        BAT *ins = 0;
        str msg = MAL_SUCCEED;
 
@@ -910,12 +911,12 @@ WLRappend(Client cntxt, MalBlkPtr mb, Ma
        }
 
        if (cname[0] != '%' && (c = mvc_bind_column(m, t, cname)) != NULL) {
-               store->storage_api.append_col(m->session->tr, c, (size_t)pos, 
ins, TYPE_bat, 1);
-       } else if (cname[0] == '%') {
-               sql_idx *i = mvc_bind_idx(m, s, cname + 1);
-               if (i)
-                       store->storage_api.append_idx(m->session->tr, i, 
(size_t)pos, ins, tpe, 1);
+               log_res = store->storage_api.append_col(m->session->tr, c, 
(size_t)pos, ins, TYPE_bat, 1);
+       } else if (cname[0] == '%' && (idx = mvc_bind_idx(m, s, cname + 1)) != 
NULL) {
+               log_res = store->storage_api.append_idx(m->session->tr, idx, 
(size_t)pos, ins, tpe, 1);
        }
+       if (log_res != LOG_OK) /* the conflict case should never happen, but 
leave it here */
+               msg = createException(MAL, "WLRappend", SQLSTATE(42000) "Append 
failed%s", log_res == LOG_CONFLICT ? " due to conflict with another 
transaction" : "");
 cleanup:
        BBPunfix(((BAT *) ins)->batCacheid);
        return msg;
@@ -993,10 +994,11 @@ WLRupdate(Client cntxt, MalBlkPtr mb, Ma
        sql_schema *s;
        sql_table *t;
        sql_column *c;
+       sql_idx *i;
        BAT *upd = 0, *tids=0;
        str msg= MAL_SUCCEED;
        oid o;
-       int tpe = getArgType(mb,pci,5);
+       int tpe = getArgType(mb,pci,5), log_res = LOG_OK;
 
        if( cntxt->wlc_kind == WLC_ROLLBACK || cntxt->wlc_kind == WLC_ERROR)
                return msg;
@@ -1062,13 +1064,12 @@ WLRupdate(Client cntxt, MalBlkPtr mb, Ma
        BATmsync(tids);
        BATmsync(upd);
        if (cname[0] != '%' && (c = mvc_bind_column(m, t, cname)) != NULL) {
-               store->storage_api.update_col(m->session->tr, c, tids, upd, 
TYPE_bat);
-       } else if (cname[0] == '%') {
-               sql_idx *i = mvc_bind_idx(m, s, cname + 1);
-               if (i)
-                       store->storage_api.update_idx(m->session->tr, i, tids, 
upd, TYPE_bat);
+               log_res = store->storage_api.update_col(m->session->tr, c, 
tids, upd, TYPE_bat);
+       } else if (cname[0] == '%' && (i = mvc_bind_idx(m, s, cname + 1)) != 
NULL) {
+               log_res = store->storage_api.update_idx(m->session->tr, i, 
tids, upd, TYPE_bat);
        }
-
+       if (log_res != LOG_OK)
+               msg = createException(MAL, "WLRupdate", "Update failed%s", 
log_res == LOG_CONFLICT ? " due to conflict with another transaction" : "");
 cleanup:
        BBPunfix(((BAT *) tids)->batCacheid);
        BBPunfix(((BAT *) upd)->batCacheid);
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
@@ -997,18 +997,21 @@ destroy_delta(sql_delta *b, bool recursi
 }
 
 static sql_delta *
-bind_col_data(sql_trans *tr, sql_column *c, bool update)
+bind_col_data(sql_trans *tr, sql_column *c, bool *update_conflict)
 {
        sql_delta *obat = ATOMIC_PTR_GET(&c->data);
 
        if (isTempTable(c->t))
                obat = temp_col_timestamp_delta(tr, c);
 
-       if (obat->cs.ts == tr->tid || !update)
+       if (obat->cs.ts == tr->tid || !update_conflict) /* on append there are 
no conflicts */
                return obat;
-       if ((!tr->parent || !tr_version_of_parent(tr, obat->cs.ts)) && 
obat->cs.ts >= TRANSACTION_ID_BASE && !isTempTable(c->t))
+       if ((!tr->parent || !tr_version_of_parent(tr, obat->cs.ts)) && 
obat->cs.ts >= TRANSACTION_ID_BASE && !isTempTable(c->t)) {
                /* abort */
+               if (update_conflict)
+                       *update_conflict = true;
                return NULL;
+       }
        assert(!isTempTable(c->t));
        obat = timestamp_delta(tr, ATOMIC_PTR_GET(&c->data));
        sql_delta* bat = ZNEW(sql_delta);
@@ -1045,10 +1048,11 @@ update_col_execute(sql_trans *tr, sql_de
 static int
 update_col(sql_trans *tr, sql_column *c, void *tids, void *upd, int tpe)
 {
+       bool update_conflict = false;
        sql_delta *delta, *odelta = ATOMIC_PTR_GET(&c->data);
 
-       if ((delta = bind_col_data(tr, c, true)) == NULL)
-               return LOG_ERR;
+       if ((delta = bind_col_data(tr, c, &update_conflict)) == NULL)
+               return update_conflict ? LOG_CONFLICT : LOG_ERR;
 
        assert(delta && delta->cs.ts == tr->tid);
        if ((!inTransaction(tr, c->t) && (odelta != delta || isTempTable(c->t)) 
&& isGlobal(c->t)) || (!isNew(c->t) && isLocalTemp(c->t)))
@@ -1058,25 +1062,28 @@ update_col(sql_trans *tr, sql_column *c,
 }
 
 static sql_delta *
-bind_idx_data(sql_trans *tr, sql_idx *i, bool update)
+bind_idx_data(sql_trans *tr, sql_idx *i, bool *update_conflict)
 {
        sql_delta *obat = ATOMIC_PTR_GET(&i->data);
 
        if (isTempTable(i->t))
                obat = temp_idx_timestamp_delta(tr, i);
 
-       if (obat->cs.ts == tr->tid || !update)
+       if (obat->cs.ts == tr->tid || !update_conflict)
                return obat;
-       if ((!tr->parent || !tr_version_of_parent(tr, obat->cs.ts)) && 
obat->cs.ts >= TRANSACTION_ID_BASE && !isTempTable(i->t))
+       if ((!tr->parent || !tr_version_of_parent(tr, obat->cs.ts)) && 
obat->cs.ts >= TRANSACTION_ID_BASE && !isTempTable(i->t)) {
                /* abort */
+               if (update_conflict)
+                       *update_conflict = true;
                return NULL;
+       }
        assert(!isTempTable(i->t));
        obat = timestamp_delta(tr, ATOMIC_PTR_GET(&i->data));
        sql_delta* bat = ZNEW(sql_delta);
        if(!bat)
                return NULL;
        bat->cs.refcnt = 1;
-       if(dup_bat(tr, i->t, obat, bat, (oid_index(i->type))?TYPE_oid:TYPE_lng) 
== LOG_ERR)
+       if(dup_bat(tr, i->t, obat, bat, (oid_index(i->type))?TYPE_oid:TYPE_lng) 
!= LOG_OK)
                return NULL;
        bat->cs.ts = tr->tid;
        /* only one writer else abort */
@@ -1092,10 +1099,11 @@ bind_idx_data(sql_trans *tr, sql_idx *i,
 static int
 update_idx(sql_trans *tr, sql_idx * i, void *tids, void *upd, int tpe)
 {
+       bool update_conflict = false;
        sql_delta *delta, *odelta = ATOMIC_PTR_GET(&i->data);
 
-       if ((delta = bind_idx_data(tr, i, true)) == NULL)
-               return LOG_ERR;
+       if ((delta = bind_idx_data(tr, i, &update_conflict)) == NULL)
+               return update_conflict ? LOG_CONFLICT : LOG_ERR;
 
        assert(delta && delta->cs.ts == tr->tid);
        if ((!inTransaction(tr, i->t) && (odelta != delta || isTempTable(i->t)) 
&& isGlobal(i->t)) || (!isNew(i->t) && isLocalTemp(i->t)))
@@ -1239,7 +1247,7 @@ append_col(sql_trans *tr, sql_column *c,
        sql_delta *delta, *odelta = ATOMIC_PTR_GET(&c->data);
        int in_transaction = segments_in_transaction(tr, c->t);
 
-       if ((delta = bind_col_data(tr, c, false)) == NULL)
+       if ((delta = bind_col_data(tr, c, NULL)) == NULL)
                return LOG_ERR;
 
        assert(delta && (!isTempTable(c->t) || delta->cs.ts == tr->tid));
@@ -1256,7 +1264,7 @@ append_idx(sql_trans *tr, sql_idx * i, s
        sql_delta *delta, *odelta = ATOMIC_PTR_GET(&i->data);
        int in_transaction = segments_in_transaction(tr, i->t);
 
-       if ((delta = bind_idx_data(tr, i, false)) == NULL)
+       if ((delta = bind_idx_data(tr, i, NULL)) == NULL)
                return LOG_ERR;
 
        assert(delta && (!isTempTable(i->t) || delta->cs.ts == tr->tid));
@@ -2276,7 +2284,7 @@ clear_col(sql_trans *tr, sql_column *c)
 {
        sql_delta *delta, *odelta = ATOMIC_PTR_GET(&c->data);
 
-       if ((delta = bind_col_data(tr, c, false)) == NULL)
+       if ((delta = bind_col_data(tr, c, NULL)) == NULL)
                return BUN_NONE;
        if ((!inTransaction(tr, c->t) && (odelta != delta || isTempTable(c->t)) 
&& isGlobal(c->t)) || (!isNew(c->t) && isLocalTemp(c->t)))
                trans_add(tr, &c->base, delta, &tc_gc_col, &commit_update_col, 
isLocalTemp(c->t)?NULL:&log_update_col);
@@ -2292,7 +2300,7 @@ clear_idx(sql_trans *tr, sql_idx *i)
 
        if (!isTable(i->t) || (hash_index(i->type) && list_length(i->columns) 
<= 1) || !idx_has_column(i->type))
                return 0;
-       if ((delta = bind_idx_data(tr, i, false)) == NULL)
+       if ((delta = bind_idx_data(tr, i, NULL)) == NULL)
                return BUN_NONE;
        if ((!inTransaction(tr, i->t) && (odelta != delta || isTempTable(i->t)) 
&& isGlobal(i->t)) || (!isNew(i->t) && isLocalTemp(i->t)))
                trans_add(tr, &i->base, delta, &tc_gc_idx, &commit_update_idx, 
isLocalTemp(i->t)?NULL:&log_update_idx);
diff --git a/sql/test/miscellaneous/Tests/transaction_isolation.SQL.py 
b/sql/test/miscellaneous/Tests/transaction_isolation.SQL.py
--- a/sql/test/miscellaneous/Tests/transaction_isolation.SQL.py
+++ b/sql/test/miscellaneous/Tests/transaction_isolation.SQL.py
@@ -153,8 +153,15 @@ with SQLTestCase() as mdb1:
 
         mdb1.execute('start transaction;').assertSucceeded()
         mdb2.execute('start transaction;').assertSucceeded()
+        mdb1.execute('update integers set i = 2 where i = 
1;').assertRowCount(1)
+        mdb2.execute('update integers set i = 2 where i = 
1;').assertFailed(err_code="42000", err_message="Update failed due to conflict 
with another transaction")
+        mdb1.execute('commit;').assertSucceeded()
+        mdb2.execute('rollback;').assertSucceeded()
+
+        mdb1.execute('start transaction;').assertSucceeded()
+        mdb2.execute('start transaction;').assertSucceeded()
         mdb1.execute('delete from integers where i = 9;').assertRowCount(1)
-        mdb2.execute('truncate integers').assertFailed(err_code="42000", 
err_message="Table clear failed due to conflict with another transaction")
+        mdb2.execute('truncate integers;').assertFailed(err_code="42000", 
err_message="Table clear failed due to conflict with another transaction")
         mdb1.execute('commit;').assertSucceeded()
         mdb2.execute('rollback;').assertSucceeded()
 
_______________________________________________
checkin-list mailing list
[email protected]
https://www.monetdb.org/mailman/listinfo/checkin-list

Reply via email to