Changeset: 594e5418e80b for MonetDB
URL: https://dev.monetdb.org/hg/MonetDB/rev/594e5418e80b
Modified Files:
        monetdb5/modules/mal/pipeline.c
        monetdb5/modules/mal/pipeline.h
Branch: pp_hashjoin
Log Message:

rename structs and move them to header file


diffs (truncated from 310 to 300 lines):

diff --git a/monetdb5/modules/mal/pipeline.c b/monetdb5/modules/mal/pipeline.c
--- a/monetdb5/modules/mal/pipeline.c
+++ b/monetdb5/modules/mal/pipeline.c
@@ -117,19 +117,8 @@ sleep_ns( int ns)
 #endif
 }
 
-typedef struct pp_counter_t {
-       struct pipeline_io s;
-
-       MT_Lock l;
-       int nr;
-       int current;
-       bool sync;
-       int scnt;
-       int *cur; /* nr per worker */
-} pp_counter;
-
 static void
-counter_free(pp_counter *c)
+counter_free(struct pipeline_counter *c)
 {
        GDKfree(c->cur);
        MT_lock_destroy(&c->l);
@@ -137,7 +126,7 @@ counter_free(pp_counter *c)
 }
 
 static int
-sync_counter_done(pp_counter *c, int wid, int nr_workers, int redo)
+sync_counter_done(struct pipeline_counter *c, int wid, int nr_workers, int 
redo)
 {
        (void)redo;
        int res = 0, cur;
@@ -168,7 +157,7 @@ sync_counter_done(pp_counter *c, int wid
 }
 
 static int
-counter_done(pp_counter *c, int wid, int nr_workers, int redo)
+counter_done(struct pipeline_counter *c, int wid, int nr_workers, int redo)
 {
        (void)redo;
        int res = 0, cur;
@@ -200,14 +189,14 @@ PPcounter_get(Client cntxt, MalBlkPtr mb
        BAT *b = BATdescriptor(cb);
        if (!b)
                throw(MAL, "pipeline.counter_get", SQLSTATE(HY002) 
RUNTIME_OBJECT_MISSING);
-       pp_counter *c = (pp_counter*)b->pl_io;
+       struct pipeline_counter *c = (struct pipeline_counter*)b->pl_io;
        if (!c) {
                BBPunfix(b->batCacheid);
                throw(MAL, "pipeline.counter_get", SQLSTATE(HY002) "Missing 
source sink");
        }
-       if (c->s.type != PIPELINE_IO_COUNTER) {
+       if (c->pl_io.type != PIPELINE_IO_COUNTER) {
                BBPunfix(b->batCacheid);
-               throw(MAL, "pipeline.counter_get", SQLSTATE(HY002) "Invalid 
source type %d, expected %d", c->s.type, PIPELINE_IO_COUNTER);
+               throw(MAL, "pipeline.counter_get", SQLSTATE(HY002) "Invalid 
source type %d, expected %d", c->pl_io.type, PIPELINE_IO_COUNTER);
        }
        if (c->sync && c->scnt != (int)p->p->nr_workers) {
                MT_lock_set(&p->p->l);
@@ -223,7 +212,7 @@ PPcounter_get(Client cntxt, MalBlkPtr mb
                MT_lock_unset(&p->p->l);
        }
        if (!c->cur)
-               c->s.done(c, p->wid, p->p->nr_workers, false);
+               c->pl_io.done(c, p->wid, p->p->nr_workers, false);
        *cur = c->cur[p->wid];
        if (*cur >= c->nr)
                p->seqnr = -2;
@@ -273,7 +262,7 @@ PPcounter(Client cntxt, MalBlkPtr mb, Ma
        if (pci->argc == 3)
                sync = *getArgReference_bit(stk, pci, 2);
 
-       pp_counter *c = (pp_counter*)GDKzalloc(sizeof(pp_counter));
+       struct pipeline_counter *c = (struct 
pipeline_counter*)GDKzalloc(sizeof(struct pipeline_counter));
        if (!c) {
                throw(SQL, "pipeline.counter",  SQLSTATE(HY013) 
MAL_MALLOC_FAIL);
        }
@@ -284,9 +273,9 @@ PPcounter(Client cntxt, MalBlkPtr mb, Ma
        }
 
        b->pl_io = (struct pipeline_io*)c;
-       c->s.type = PIPELINE_IO_COUNTER;
-       c->s.destroy = (pipeline_io_destroy)&counter_free;
-       c->s.done = (pipeline_io_done)&counter_done;
+       c->pl_io.type = PIPELINE_IO_COUNTER;
+       c->pl_io.destroy = (pipeline_io_destroy)&counter_free;
+       c->pl_io.done = (pipeline_io_done)&counter_done;
        c->current = 0;
        c->cur = NULL;
        c->nr = nr;
@@ -294,7 +283,7 @@ PPcounter(Client cntxt, MalBlkPtr mb, Ma
        if (sync) {
                c->sync = true;
                c->scnt = 0;
-               c->s.done = (pipeline_io_done)&sync_counter_done;
+               c->pl_io.done = (pipeline_io_done)&sync_counter_done;
        }
        *rb = b->batCacheid;
        BBPkeepref(b);
@@ -322,19 +311,8 @@ PPdone(Client cntxt, MalBlkPtr mb, MalSt
        return MAL_SUCCEED;
 }
 
-typedef struct pp_concat_t {
-       struct pipeline_io s;
-
-       MT_Lock l;
-       int current;
-       int max;
-       bool started;
-       int *cur;
-       BAT *srcs[];
-} pp_concat;
-
 static void
-concat_free( pp_concat *pcat )
+concat_free(struct pipeline_concat *pcat )
 {
        MT_lock_destroy(&pcat->l);
        for(int i = 0; i<pcat->max; i++) {
@@ -346,10 +324,10 @@ concat_free( pp_concat *pcat )
 }
 
 static int
-concat_done( pp_concat *c, int wid, int nr_workers, bool redo )
+concat_done(struct pipeline_concat *c, int wid, int nr_workers, bool redo )
 {
        int res = 1;
-       assert(c->s.type == PIPELINE_IO_CONCAT);
+       assert(c->pl_io.type == PIPELINE_IO_CONCAT);
        MT_lock_set(&c->l);
        if (!c->started) {
                c->cur = (int*)GDKzalloc(sizeof(int) * nr_workers);
@@ -376,10 +354,10 @@ concat_done( pp_concat *c, int wid, int 
 }
 
 static int
-concat_next( pp_concat *c, int wid)
+concat_next(struct pipeline_concat *c, int wid)
 {
        int res = 1;
-       assert(c->s.type == PIPELINE_IO_CONCAT);
+       assert(c->pl_io.type == PIPELINE_IO_CONCAT);
        MT_lock_set(&c->l);
        assert(c->started);
        BAT *sb = c->srcs[c->cur[wid]];
@@ -395,10 +373,10 @@ concat_next( pp_concat *c, int wid)
 }
 
 static BAT*
-concat_next_bat( pp_concat *c, int wid)
+concat_next_bat(struct pipeline_concat *c, int wid)
 {
        BAT *res = NULL;
-       assert(c->s.type == PIPELINE_IO_CONCAT);
+       assert(c->pl_io.type == PIPELINE_IO_CONCAT);
        MT_lock_set(&c->l);
        assert(c->started);
        BAT *sb = c->srcs[c->cur[wid]];
@@ -432,10 +410,10 @@ PPconcat_block(Client cntxt, MalBlkPtr m
        BAT *b = BATdescriptor(cb);
        if (!b)
                throw(MAL, "pipeline.concat_block", SQLSTATE(HY002) 
RUNTIME_OBJECT_MISSING);
-       pp_concat *pcat = (pp_concat*)b->pl_io;
-       if (pcat->s.type != PIPELINE_IO_CONCAT) {
+       struct pipeline_concat *pcat = (struct pipeline_concat*)b->pl_io;
+       if (pcat->pl_io.type != PIPELINE_IO_CONCAT) {
                BBPreclaim(b);
-               throw(MAL, "pipeline.concat_block", SQLSTATE(HY002) "Invalid 
type for a concat source %d", pcat->s.type);
+               throw(MAL, "pipeline.concat_block", SQLSTATE(HY002) "Invalid 
type for a concat source %d", pcat->pl_io.type);
        }
        MT_lock_set(&pcat->l);
        assert(pcat->cur);
@@ -461,11 +439,11 @@ PPconcat_add(Client cntxt, MalBlkPtr mb,
                BBPreclaim(i);
                throw(MAL, "pipeline.concat_add", SQLSTATE(HY002) 
RUNTIME_OBJECT_MISSING);
        }
-       pp_concat *pcat = (pp_concat*)b->pl_io;
-       if (pcat->s.type != PIPELINE_IO_CONCAT) {
+       struct pipeline_concat *pcat = (struct pipeline_concat*)b->pl_io;
+       if (pcat->pl_io.type != PIPELINE_IO_CONCAT) {
                BBPreclaim(b);
                BBPreclaim(i);
-               throw(MAL, "pipeline.concat_add", SQLSTATE(HY002) "Invalid type 
for a concat source %d", pcat->s.type);
+               throw(MAL, "pipeline.concat_add", SQLSTATE(HY002) "Invalid type 
for a concat source %d", pcat->pl_io.type);
        }
        if (pcat->current >= pcat->max) {
                BBPreclaim(b);
@@ -485,7 +463,7 @@ PPconcat(Client cntxt, MalBlkPtr mb, Mal
        (void)mb;
        bat *rb = getArgReference_bat(stk, pci, 0);
        int nr = *getArgReference_int(stk, pci, 1);
-       pp_concat *pcat = (pp_concat*)GDKzalloc(sizeof(pp_concat) + (nr+1) * 
sizeof(struct pipeline_io*) );
+       struct pipeline_concat *pcat = (struct 
pipeline_concat*)GDKzalloc(sizeof(struct pipeline_concat) + (nr+1) * 
sizeof(struct pipeline_io*) );
 
        if (!pcat)
                throw(SQL, "pipeline.concat",  SQLSTATE(HY013) MAL_MALLOC_FAIL);
@@ -496,11 +474,11 @@ PPconcat(Client cntxt, MalBlkPtr mb, Mal
                throw(SQL, "pipeline.concat",  SQLSTATE(HY013) MAL_MALLOC_FAIL);
        }
        b->pl_io = (struct pipeline_io*)pcat;
-       pcat->s.type = PIPELINE_IO_CONCAT;
-       pcat->s.destroy = (pipeline_io_destroy)&concat_free;
-       pcat->s.done = (pipeline_io_done)&concat_done;
-       pcat->s.next = (pipeline_io_next)&concat_next;
-       pcat->s.next_bat = (pipeline_io_next_bat)&concat_next_bat;
+       pcat->pl_io.type = PIPELINE_IO_CONCAT;
+       pcat->pl_io.destroy = (pipeline_io_destroy)&concat_free;
+       pcat->pl_io.done = (pipeline_io_done)&concat_done;
+       pcat->pl_io.next = (pipeline_io_next)&concat_next;
+       pcat->pl_io.next_bat = (pipeline_io_next_bat)&concat_next_bat;
        pcat->current = 0;
        pcat->max = nr;
        pcat->started = false;
@@ -510,19 +488,13 @@ PPconcat(Client cntxt, MalBlkPtr mb, Mal
        return MAL_SUCCEED;
 }
 
-typedef struct pp_resultset_t {
-       struct pipeline_io s;
-       ATOMIC_TYPE claimed;
-       MT_Lock l;
-} pp_resultset;
-
 static str
 PPresultset(Client cntxt, MalBlkPtr mb, MalStkPtr stk, InstrPtr pci)
 {
        (void)cntxt;
        (void)mb;
        bat *rb = getArgReference_bat(stk, pci, 0);
-       pp_resultset *prs = (pp_resultset*)GDKzalloc(sizeof(pp_resultset));
+       struct pipeline_resultset *prs = (struct 
pipeline_resultset*)GDKzalloc(sizeof(struct pipeline_resultset));
 
        if (!prs)
                throw(SQL, "pipeline.resultset",  SQLSTATE(HY013) 
MAL_MALLOC_FAIL);
@@ -533,7 +505,7 @@ PPresultset(Client cntxt, MalBlkPtr mb, 
                throw(SQL, "pipeline.resultset",  SQLSTATE(HY013) 
MAL_MALLOC_FAIL);
        }
        b->pl_io = (struct pipeline_io*)prs;
-       prs->s.destroy = (pipeline_io_destroy)&GDKfree;
+       prs->pl_io.destroy = (pipeline_io_destroy)&GDKfree;
        MT_lock_init(&prs->l, "resultset");
        *rb = b->batCacheid;
        BBPkeepref(b);
@@ -550,7 +522,7 @@ PPclaim(Client cntxt, MalBlkPtr mb, MalS
 
        BAT *b = BATdescriptor(rb);
        if (b) {
-               pp_resultset *rs = (pp_resultset*)b->pl_io;
+               struct pipeline_resultset *rs = (struct 
pipeline_resultset*)b->pl_io;
                *res = ATOMIC_ADD(&rs->claimed, cnt);
                BBPreclaim(b);
        }
@@ -575,7 +547,7 @@ PPidentity(Client cntxt, MalBlkPtr mb, M
        BBPreclaim(b);
        b = BATdescriptor(rb);
        if (b) {
-               pp_resultset *rs = (pp_resultset*)b->pl_io;
+               struct pipeline_resultset *rs = (struct 
pipeline_resultset*)b->pl_io;
                offset = ATOMIC_ADD(&rs->claimed, cnt);
                BBPreclaim(b);
 
@@ -610,7 +582,7 @@ PPappend(Client cntxt, MalBlkPtr mb, Mal
                BBPreclaim(r);
                throw(MAL, "bat.append", SQLSTATE(HY002) 
RUNTIME_OBJECT_MISSING);
        }
-       pp_resultset *pp_rs = (pp_resultset*)r->pl_io;
+       struct pipeline_resultset *pp_rs = (struct pipeline_resultset*)r->pl_io;
        (void)pp_rs;
 
        if (i && (i->ttype == TYPE_msk || mask_cand(i))) {
diff --git a/monetdb5/modules/mal/pipeline.h b/monetdb5/modules/mal/pipeline.h
--- a/monetdb5/modules/mal/pipeline.h
+++ b/monetdb5/modules/mal/pipeline.h
@@ -36,7 +36,33 @@
 #define PIPELINE_IO_COUNTER    8
 #define PIPELINE_IO_CONCAT     9
 #define PIPELINE_IO_PARQUET    10
-#define PIPELINE_IO_MPARQUET   10
+#define PIPELINE_IO_MPARQUET   11
+
+struct pipeline_counter {
+       struct pipeline_io pl_io;
+       MT_Lock l;
+       int nr;
+       int current;
+       bool sync;
+       int scnt;
+       int *cur; /* nr per worker */
+};
+
+struct pipeline_concat {
+       struct pipeline_io pl_io;
+       MT_Lock l;
+       int current;
+       int max;
+       bool started;
+       int *cur;
+       BAT *srcs[];
_______________________________________________
checkin-list mailing list -- [email protected]
To unsubscribe send an email to [email protected]

Reply via email to