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

Merge with group-commit.


diffs (truncated from 672 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
@@ -543,9 +543,9 @@ gdk_return log_bat_transient(logger *lg,
 gdk_return log_constant(logger *lg, int type, ptr val, log_id id, lng offset, 
lng cnt);
 gdk_return log_delta(logger *lg, BAT *uid, BAT *uval, log_id id);
 gdk_return log_sequence(logger *lg, int seq, lng id);
-gdk_return log_tdone(logger *lg, ulng commit_ts);
 gdk_return log_tend(logger *lg);
-gdk_return log_tstart(logger *lg, bool flush);
+gdk_return log_tflush(logger *lg, ulng log_file_id, ulng commit_ts);
+gdk_return log_tstart(logger *lg, bool flushnow, ulng *log_file_id);
 gdk_return logger_activate(logger *lg);
 lng logger_changes(logger *lg);
 logger *logger_create(int debug, const char *fn, const char *logdir, int 
version, preversionfix_fptr prefuncp, postversionfix_fptr postfuncp, void 
*funcdata);
diff --git a/gdk/gdk_logger.c b/gdk/gdk_logger.c
--- a/gdk/gdk_logger.c
+++ b/gdk/gdk_logger.c
@@ -2048,7 +2048,12 @@ logger_load(int debug, const char *fn, c
        logbat_destroy(lg->seqs_id);
        logbat_destroy(lg->seqs_val);
        logbat_destroy(lg->dseqs);
+       ATOMIC_DESTROY(&lg->refcount);
        MT_lock_destroy(&lg->lock);
+       MT_lock_destroy(&lg->rotation_lock);
+       MT_sema_destroy(&lg->flush_queue_semaphore);
+       MT_lock_destroy(&lg->flush_lock);
+       MT_lock_destroy(&lg->flush_queue_lock);
        GDKfree(lg->fn);
        GDKfree(lg->dir);
        GDKfree(lg->local_dir);
@@ -2108,11 +2113,21 @@ logger_new(int debug, const char *fn, co
                GDKfree(lg);
                return NULL;
        }
-       MT_lock_init(&lg->lock, fn);
        if (lg->debug & 1) {
                fprintf(stderr, "#logger_new dir set to %s\n", lg->dir);
        }
 
+       ATOMIC_INIT(&lg->refcount, 0);
+       MT_lock_init(&lg->lock, fn);
+       MT_lock_init(&lg->rotation_lock, "rotation_lock");
+       MT_sema_init(&lg->flush_queue_semaphore, FLUSH_QUEUE_SIZE, 
"flush_queue_semaphore");
+       MT_lock_init(&lg->flush_lock, "flush_lock");
+       MT_lock_init(&lg->flush_queue_lock, "flush_queue_lock");
+
+       // flush variables
+       lg->flush_queue_begin = 0;
+       lg->flush_queue_length = 0;
+
        if (logger_load(debug, fn, logdir, lg, filename) == GDK_SUCCEED) {
                return lg;
        }
@@ -2156,7 +2171,12 @@ logger_destroy(logger *lg)
                logbat_destroy(lg->catalog_lid);
                logger_unlock(lg);
        }
+       ATOMIC_DESTROY(&lg->refcount);
        MT_lock_destroy(&lg->lock);
+       MT_lock_destroy(&lg->rotation_lock);
+       MT_sema_destroy(&lg->flush_queue_semaphore);
+       MT_lock_destroy(&lg->flush_lock);
+       MT_lock_destroy(&lg->flush_queue_lock);
        GDKfree(lg->fn);
        GDKfree(lg->dir);
        GDKfree(lg->buf);
@@ -2215,6 +2235,7 @@ logger_cleanup_range(logger *lg)
 gdk_return
 logger_activate(logger *lg)
 {
+       MT_lock_set(&lg->rotation_lock);
        logger_lock(lg);
        if (lg->end > 0 && lg->saved_id+1 == lg->id) {
                lg->id++;
@@ -2222,10 +2243,12 @@ logger_activate(logger *lg)
                /* start new file */
                if (logger_open_output(lg) != GDK_SUCCEED) {
                        logger_unlock(lg);
+                       MT_lock_unset(&lg->rotation_lock);
                        return GDK_FAIL;
                }
        }
        logger_unlock(lg);
+       MT_lock_unset(&lg->rotation_lock);
        return GDK_SUCCEED;
 }
 
@@ -2361,6 +2384,7 @@ log_constant(logger *lg, int type, ptr v
            (!is_row && !mnstr_writeLng(lg->output_log, nr)) ||
            (!is_row && mnstr_write(lg->output_log, &tpe, 1, 1) != 1) ||
            (!is_row && !mnstr_writeLng(lg->output_log, offset))) {
+               (void) ATOMIC_DEC(&lg->refcount);
                ok = GDK_FAIL;
                goto bailout;
        }
@@ -2507,6 +2531,7 @@ internal_log_bat(logger *lg, BAT *b, log
 
   bailout:
        if (ok != GDK_SUCCEED) {
+               (void) ATOMIC_DEC(&lg->refcount);
                const char *err = mnstr_peek_error(lg->output_log);
                TRC_CRITICAL(GDK, "write failed%s%s\n", err ? ": " : "", err ? 
err : "");
        }
@@ -2527,6 +2552,8 @@ log_bat_persists(logger *lg, BAT *b, log
 
        if (logger_add_bat(lg, b, id, -1) != GDK_SUCCEED) {
                logger_unlock(lg);
+               if (!LOG_DISABLED(lg))
+                       (void) ATOMIC_DEC(&lg->refcount);
                return GDK_FAIL;
        }
 
@@ -2536,6 +2563,7 @@ log_bat_persists(logger *lg, BAT *b, log
                if (log_write_format(lg, &l) != GDK_SUCCEED ||
                    mnstr_write(lg->output_log, &ta, 1, 1) != 1) {
                        logger_unlock(lg);
+                       (void) ATOMIC_DEC(&lg->refcount);
                        return GDK_FAIL;
                }
        }
@@ -2544,6 +2572,8 @@ log_bat_persists(logger *lg, BAT *b, log
                fprintf(stderr, "#persists id (%d) bat (%d)\n", id, 
b->batCacheid);
        gdk_return r = internal_log_bat(lg, b, id, 0, BATcount(b), 0);
        logger_unlock(lg);
+       if (r != GDK_SUCCEED)
+               (void) ATOMIC_DEC(&lg->refcount);
        return r;
 }
 
@@ -2561,6 +2591,7 @@ log_bat_transient(logger *lg, log_id id)
                if (log_write_format(lg, &l) != GDK_SUCCEED) {
                        TRC_CRITICAL(GDK, "write failed\n");
                        logger_unlock(lg);
+                       (void) ATOMIC_DEC(&lg->refcount);
                        return GDK_FAIL;
                }
        }
@@ -2571,6 +2602,8 @@ log_bat_transient(logger *lg, log_id id)
        lg->end += BATcount(BBPquickdesc(bid));
        gdk_return r =  logger_del_bat(lg, bid);
        logger_unlock(lg);
+       if (r != GDK_SUCCEED)
+               (void) ATOMIC_DEC(&lg->refcount);
        return r;
 }
 
@@ -2596,6 +2629,8 @@ log_delta(logger *lg, BAT *uid, BAT *uva
        if (BATtdense(uid)) {
                ok = internal_log_bat(lg, uval, id, uid->tseqbase, 
BATcount(uval), 1);
                logger_unlock(lg);
+               if (!LOG_DISABLED(lg) && ok != GDK_SUCCEED)
+                       (void) ATOMIC_DEC(&lg->refcount);
                return ok;
        }
 
@@ -2654,6 +2689,7 @@ log_delta(logger *lg, BAT *uid, BAT *uva
        if (ok != GDK_SUCCEED) {
                const char *err = mnstr_peek_error(lg->output_log);
                TRC_CRITICAL(GDK, "write failed%s%s\n", err ? ": " : "", err ? 
err : "");
+               (void) ATOMIC_DEC(&lg->refcount);
        }
        logger_unlock(lg);
        return ok;
@@ -2678,7 +2714,11 @@ log_bat_clear(logger *lg, int id)
 
        if (lg->debug & 1)
                fprintf(stderr, "#Logged clear %d\n", id);
-       return log_write_format(lg, &l);
+
+       gdk_return r =  log_write_format(lg, &l);
+       if(r != GDK_SUCCEED)
+               (void) ATOMIC_DEC(&lg->refcount);
+       return r;
 }
 
 #define DBLKSZ         8192
@@ -2690,60 +2730,161 @@ log_bat_clear(logger *lg, int id)
 static gdk_return
 new_logfile(logger *lg)
 {
-       lng log_large = (GDKdebug & FORCEMITOMASK)?LOG_MINI:LOG_LARGE;
        assert(!LOG_DISABLED(lg));
-       lng p;
-       p = (lng) getfilepos(getFile(lg->output_log));
-       if (p == -1)
+
+
+       const lng log_large = (GDKdebug & FORCEMITOMASK)?LOG_MINI:LOG_LARGE;
+
+       gdk_return result = GDK_SUCCEED;
+       MT_lock_set(&lg->rotation_lock);
+       const lng p = (lng) getfilepos(getFile(lg->output_log));
+       if (p == -1) {
+               MT_lock_unset(&lg->rotation_lock);
                return GDK_FAIL;
-       if (p > log_large || (lg->end*1024) > log_large) {
-               lg->id++;
-               logger_close_output(lg);
-               return logger_open_output(lg);
        }
-       return GDK_SUCCEED;
+       if (( p > log_large || (lg->end*1024) > log_large )) {
+               if (ATOMIC_GET(&lg->refcount) == 1) {
+                       lg->id++;
+                       logger_close_output(lg);
+                       result = logger_open_output(lg);
+                       lg->request_rotation = false;
+               }
+               else {
+                       // Delegate wal rotation to next writer or last flusher.
+                       lg->request_rotation = true;
+               }
+       }
+       MT_lock_unset(&lg->rotation_lock);
+       return result;
 }
 
 gdk_return
 log_tend(logger *lg)
 {
-       logformat l;
-
        if (lg->debug & 1)
                fprintf(stderr, "#log_tend %d\n", lg->tid);
 
-       l.flag = LOG_END;
-       l.id = lg->tid;
-       if (lg->flushnow) {
-               lg->flushnow = 0;
-               return logger_commit(lg);
-       }
-
        lg->end++;
        if (LOG_DISABLED(lg)) {
                return GDK_SUCCEED;
        }
 
-       if (log_write_format(lg, &l) != GDK_SUCCEED ||
-           mnstr_flush(lg->output_log, MNSTR_FLUSH_DATA) ||
-           (!(GDKdebug & NOSYNCMASK) && mnstr_fsync(lg->output_log)) ||
-           new_logfile(lg) != GDK_SUCCEED) {
-               TRC_CRITICAL(GDK, "write failed\n");
-               return GDK_FAIL;
+       gdk_return result;
+       logformat l;
+       l.flag = LOG_END;
+       l.id = lg->tid;
+
+       if ((result = log_write_format(lg, &l)) != GDK_SUCCEED)
+               (void) ATOMIC_DEC(&lg->refcount);
+       return result;
+}
+static int
+request_number_flush_queue(logger *lg) {
+       // Semaphore protects ring buffer structure in queue against overflowing
+       static int _number = 0;
+       int result;
+       MT_sema_down(&lg->flush_queue_semaphore);
+       MT_lock_set(&lg->flush_queue_lock);
+       result = ++_number;
+       const int end = (lg->flush_queue_begin + lg->flush_queue_length) % 
FLUSH_QUEUE_SIZE;
+       lg->flush_queue[end] = _number;
+       lg->flush_queue_length++;
+       MT_lock_unset(&lg->flush_queue_lock);
+
+       return result;
+}
+
+static void
+left_truncate_flush_queue(logger *lg, int limit) {
+       MT_lock_set(&lg->flush_queue_lock);
+       lg->flush_queue_begin = (lg->flush_queue_begin + limit) % 
FLUSH_QUEUE_SIZE;
+       lg->flush_queue_length -= limit;
+       MT_lock_unset(&lg->flush_queue_lock);
+
+       for (int i = 0; i < limit; i++)
+               MT_sema_up(&lg->flush_queue_semaphore);
+}
+
+static int
+number_in_flush_queue(logger *lg, int number) {
+       MT_lock_set(&lg->flush_queue_lock);
+       const int fql = lg->flush_queue_length;
+       MT_lock_unset(&lg->flush_queue_lock);
+       for (int i = 0; i < fql; i++) {
+               const int idx = (lg->flush_queue_begin + i) % FLUSH_QUEUE_SIZE;
+               if (lg->flush_queue[idx] == number) {
+                       return 1;
+               }
+       }
+       return 0;
+}
+
+static int
+flush_queue_length(logger *lg) {
+       MT_lock_set(&lg->flush_queue_lock);
+       const int fql = lg->flush_queue_length;
+       MT_lock_unset(&lg->flush_queue_lock);
+       return fql;
_______________________________________________
checkin-list mailing list -- [email protected]
To unsubscribe send an email to [email protected]

Reply via email to