Changeset: ccdf741f5076 for MonetDB
URL: https://dev.monetdb.org/hg/MonetDB/rev/ccdf741f5076
Modified Files:
        gdk/gdk_bat.c
        gdk/gdk_batop.c
Branch: Jul2021
Log Message:

Move hashlock around a little so we don't get a race condition.
If BATappend gets called in one thread on a bat without a hash, it locks
the hash lock, appends the data, and maybe updates the hash (no update,
since no hash), then unlocks the hash lock, locks the heap lock and
updates the bat count.  It could happen that another thread locks the
hash lock as soon as it was unlocked here, and creates a hash based on
the old count.  Then the next call to BATappend sees a hash, but the
hash size and bat count don't match.  By rearranging the locks (keeping
the hash lock locked until after the bat count was updated), we prevent
this scenario.


diffs (170 lines):

diff --git a/gdk/gdk_bat.c b/gdk/gdk_bat.c
--- a/gdk/gdk_bat.c
+++ b/gdk/gdk_bat.c
@@ -1162,6 +1162,7 @@ BUNappendmulti(BAT *b, const void *value
                b->tsorted = b->trevsorted = b->tkey = false;
        }
        MT_lock_unset(&b->theaplock);
+       MT_rwlock_wrlock(&b->thashlock);
        if (values && b->ttype) {
                int (*atomcmp) (const void *, const void *) = 
ATOMcompare(b->ttype);
                const void *atomnil = ATOMnilptr(b->ttype);
@@ -1183,6 +1184,7 @@ BUNappendmulti(BAT *b, const void *value
                                t = ((void **) values)[i];
                                gdk_return rc = tfastins_nocheckVAR(b, p, t);
                                if (rc != GDK_SUCCEED) {
+                                       MT_rwlock_wrunlock(&b->thashlock);
                                        return rc;
                                }
                                if (vbase != b->tvheap->base) {
@@ -1217,7 +1219,6 @@ BUNappendmulti(BAT *b, const void *value
                                }
                                p++;
                        }
-                       MT_rwlock_wrlock(&b->thashlock);
                        if (b->thash) {
                                p -= count;
                                for (BUN i = 0; i < count; i++) {
@@ -1226,7 +1227,6 @@ BUNappendmulti(BAT *b, const void *value
                                        p++;
                                }
                        }
-                       MT_rwlock_wrunlock(&b->thashlock);
                } else if (ATOMstorage(b->ttype) == TYPE_msk) {
                        minpos = maxpos = BUN_NONE;
                        minvalp = maxvalp = NULL;
@@ -1236,7 +1236,6 @@ BUNappendmulti(BAT *b, const void *value
                                p++;
                        }
                } else {
-                       MT_rwlock_wrlock(&b->thashlock);
                        for (BUN i = 0; i < count; i++) {
                                t = (void *) ((char *) values + (i << 
b->tshift));
                                gdk_return rc = tfastins_nocheckFIX(b, p, t);
@@ -1266,7 +1265,6 @@ BUNappendmulti(BAT *b, const void *value
                                }
                                p++;
                        }
-                       MT_rwlock_wrunlock(&b->thashlock);
                }
                MT_lock_set(&b->theaplock);
                if (minpos == BUN_NONE) {
@@ -1285,7 +1283,6 @@ BUNappendmulti(BAT *b, const void *value
                }
                MT_lock_unset(&b->theaplock);
        } else {
-               MT_rwlock_wrlock(&b->thashlock);
                for (BUN i = 0; i < count; i++) {
                        gdk_return rc = tfastins_nocheck(b, p, t);
                        if (rc != GDK_SUCCEED) {
@@ -1297,11 +1294,11 @@ BUNappendmulti(BAT *b, const void *value
                        }
                        p++;
                }
-               MT_rwlock_wrunlock(&b->thashlock);
        }
        MT_lock_set(&b->theaplock);
        BATsetcount(b, p);
        MT_lock_unset(&b->theaplock);
+       MT_rwlock_wrunlock(&b->thashlock);
 
        IMPSdestroy(b); /* no support for inserts in imprints yet */
        OIDXdestroy(b);
@@ -1586,14 +1583,11 @@ BUNinplacemulti(BAT *b, const oid *posit
                        }
                        if (b->twidth < SIZEOF_VAR_T &&
                            (b->twidth <= 2 ? _d - GDK_VAROFFSET : _d) >= 
((size_t) 1 << (8 << b->tshift))) {
-                               /* doesn't fit in current heap, upgrade
-                                * it, can't keep hashlock while doing
-                                * so */
-                               MT_rwlock_wrunlock(&b->thashlock);
+                               /* doesn't fit in current heap, upgrade it */
                                if (GDKupgradevarheap(b, _d, 0, bi.count) != 
GDK_SUCCEED) {
+                                       MT_rwlock_wrunlock(&b->thashlock);
                                        return GDK_FAIL;
                                }
-                               MT_rwlock_wrlock(&b->thashlock);
                        }
                        /* reinitialize iterator after possible heap upgrade */
                        bi = bat_iterator_nolock(b);
diff --git a/gdk/gdk_batop.c b/gdk/gdk_batop.c
--- a/gdk/gdk_batop.c
+++ b/gdk/gdk_batop.c
@@ -332,13 +332,13 @@ insert_string_bat(BAT *b, BAT *n, struct
                        r++;
                }
        }
+       MT_rwlock_wrlock(&b->thashlock);
        MT_lock_set(&b->theaplock);
        BATsetcount(b, oldcnt + ci->ncand);
        assert(b->batCapacity >= b->batCount);
        MT_lock_unset(&b->theaplock);
        bat_iterator_end(&ni);
        /* maintain hash */
-       MT_rwlock_wrlock(&b->thashlock);
        for (r = oldcnt, cnt = BATcount(b); b->thash && r < cnt; r++) {
                HASHappend_locked(b, r, b->tvheap->base + VarHeapVal(Tloc(b, 
0), r, b->twidth));
        }
@@ -410,11 +410,11 @@ append_varsized_bat(BAT *b, BAT *n, stru
                                *dst++ = src[canditer_next(ci) - hseq];
                        }
                }
+               MT_rwlock_wrlock(&b->thashlock);
                MT_lock_set(&b->theaplock);
                BATsetcount(b, BATcount(b) + ci->ncand);
                MT_lock_unset(&b->theaplock);
                /* maintain hash table */
-               MT_rwlock_wrlock(&b->thashlock);
                for (BUN i = BATcount(b) - ci->ncand;
                     b->thash && i < BATcount(b);
                     i++) {
@@ -474,10 +474,10 @@ append_varsized_bat(BAT *b, BAT *n, stru
                        r++;
                }
        }
-       MT_rwlock_wrunlock(&b->thashlock);
        MT_lock_set(&b->theaplock);
        BATsetcount(b, r);
        MT_lock_unset(&b->theaplock);
+       MT_rwlock_wrunlock(&b->thashlock);
        bat_iterator_end(&ni);
        return GDK_SUCCEED;
 }
@@ -915,10 +915,10 @@ BATappend2(BAT *b, BAT *n, BAT *s, bool 
                                r++;
                        }
                }
-               MT_rwlock_wrunlock(&b->thashlock);
                MT_lock_set(&b->theaplock);
                BATsetcount(b, b->batCount + ci.ncand);
                MT_lock_unset(&b->theaplock);
+               MT_rwlock_wrunlock(&b->thashlock);
        }
 
   doreturn:
@@ -1309,16 +1309,10 @@ BATappend_or_update(BAT *b, BAT *p, cons
                        }
                        if (b->twidth < SIZEOF_VAR_T &&
                            (b->twidth <= 2 ? d - GDK_VAROFFSET : d) >= 
((size_t) 1 << (8 << b->tshift))) {
-                               /* doesn't fit in current heap, upgrade
-                                * it, can't keep hashlock while doing
-                                * so */
-                               MT_rwlock_wrunlock(&b->thashlock);
-                               locked = false;
+                               /* doesn't fit in current heap, upgrade it */
                                if (GDKupgradevarheap(b, d, 0, MAX(updid, 
b->batCount)) != GDK_SUCCEED) {
                                        goto bailout;
                                }
-                               MT_rwlock_wrlock(&b->thashlock);
-                               locked = true;
                        }
                        /* in case ATOMreplaceVAR and/or
                         * GDKupgradevarheap replaces a heap, we need to
@@ -1351,6 +1345,7 @@ BATappend_or_update(BAT *b, BAT *p, cons
                b->tvheap->dirty = true;
                MT_lock_unset(&b->theaplock);
        } else if (ATOMstorage(b->ttype) == TYPE_msk) {
+               assert(b->thash == NULL);
                HASHdestroy(b); /* hash doesn't make sense for msk */
                for (BUN i = 0; i < ni.count; i++) {
                        oid updid;
_______________________________________________
checkin-list mailing list -- [email protected]
To unsubscribe send an email to [email protected]

Reply via email to