Changeset: 166ca3be6cd7 for MonetDB
URL: http://dev.monetdb.org/hg/MonetDB?cmd=changeset;node=166ca3be6cd7
Modified Files:
        sql/backends/monet5/datacell/50_datacell.sql
        sql/backends/monet5/datacell/basket.mx
        sql/backends/monet5/datacell/emitter.mx
        sql/backends/monet5/datacell/petrinet.mx
Branch: default
Log Message:

Add last seen to petrinet scheduler.


diffs (180 lines):

diff --git a/sql/backends/monet5/datacell/50_datacell.sql 
b/sql/backends/monet5/datacell/50_datacell.sql
--- a/sql/backends/monet5/datacell/50_datacell.sql
+++ b/sql/backends/monet5/datacell/50_datacell.sql
@@ -82,5 +82,5 @@
 external name datacell.baskets;
 
 create function datacell.queries()
-returns table( nme string, status string, cycles int, events int, time bigint, 
error string, def string)
+returns table( nme string, status string, seen timestamp, cycles int, events 
int, time bigint, error string, def string)
 external name datacell.queries;
diff --git a/sql/backends/monet5/datacell/basket.mx 
b/sql/backends/monet5/datacell/basket.mx
--- a/sql/backends/monet5/datacell/basket.mx
+++ b/sql/backends/monet5/datacell/basket.mx
@@ -109,12 +109,12 @@
        str name;       /* table that represents the basket */
        int threshold ; /* bound to determine scheduling eligibility */
        int winsize, winslide; /* sliding window operations */
-       lng seen;
        lng beat;       /* milliseconds delay */
        int colcount;
        str *cols;
        BAT **primary;  
        /* statistics */
+       timestamp seen;
        int events; /* total number of events grabbed */
        int grabs; /* number of grabs */
 } *BSKTbasket, BSKTbasketRec;
@@ -415,7 +415,7 @@
        str tbl;
        int bskt, i, *ret;
        BAT *b,*bn = 0,*v;
-       int cnt;
+       int cnt = 0;
 
        (void) cntxt;
        (void) mb;
@@ -581,7 +581,8 @@
 str
 BSKTbeat(int *ret, str *tbl, int *sz)
 {
-       int bskt;
+       int bskt, tst;
+       timestamp ts,tn;
        bskt = BSKTlocate(*tbl);
        if (bskt == 0 )
                throw(MAL,"basket.beat","Basket not found");
@@ -589,7 +590,10 @@
                throw(MAL,"basket.beat","Illegal value");
        baskets[bskt].beat = *sz;
        *ret = TRUE;
-       if ( baskets[bskt].beat + baskets[bskt].seen > GDKusec() )
+       (void) MTIMEunix_epoch(&ts);
+       (void) MTIMEtimestamp_add(&tn, &baskets[bskt].seen, 
&baskets[bskt].beat);
+       tst = tn.days < ts.days || (tn.days == ts.days && tn.msecs < ts.msecs);
+       if ( tst)
                throw(MAL,"basket.heat","too early");
        return MAL_SUCCEED;
 }
diff --git a/sql/backends/monet5/datacell/emitter.mx 
b/sql/backends/monet5/datacell/emitter.mx
--- a/sql/backends/monet5/datacell/emitter.mx
+++ b/sql/backends/monet5/datacell/emitter.mx
@@ -478,7 +478,10 @@
                        BATclear(b);
                }
                BSKTunlock(&em->lck, &em->name);
-               if ((cnt = BATcount(em->table.format[1].c[0]))) {
+               if ((cnt = BATcount(em->table.format[0].c[0]))) {
+                       MTIMEcurrent_timestamp(&baskets[em->bskt].seen);
+                       baskets[em->bskt].events = cnt;
+                       baskets[em->bskt].grabs ++;
                        if (em->status != EMLISTEN)
                                break;
 
diff --git a/sql/backends/monet5/datacell/petrinet.mx 
b/sql/backends/monet5/datacell/petrinet.mx
--- a/sql/backends/monet5/datacell/petrinet.mx
+++ b/sql/backends/monet5/datacell/petrinet.mx
@@ -133,6 +133,7 @@
        int pc;
        int status;             /* query status waiting/running/ready */
        int delay;              /* maximum delay between calls */
+       timestamp seen; /* last executed */
        int cycles;     /* number of invocations of the factory */
        int events;             /* number of events consumed */
        str error;              /* last error seen */
@@ -523,15 +524,20 @@
                                }
                                pnet[i].source[j].available = cnt = 
(int)BATcount(pnet[i].source[j].b);
                                if (cnt) {
+                                       timestamp ts, tn;
                                        /* only look at large enough baskets */
                                        if ( cnt <  baskets[idx].threshold) {
                                                pnet[i].enabled = 0;
                                                break;
                                        }
                                        /* check heart beat delays */
-                                       if ( baskets[idx].seen + 
baskets[idx].beat > now){
-                                               pnet[i].enabled = 0;
-                                               break;
+                                       if ( baskets[idx].beat) {
+                                               (void) MTIMEunix_epoch(&ts);
+                                               (void) MTIMEtimestamp_add(&tn, 
&baskets[idx].seen, &baskets[idx].beat);
+                                               if(  tn.days < ts.days || 
(tn.days == ts.days && tn.msecs < ts.msecs) ){
+                                                       pnet[i].enabled = 0;
+                                                       break;
+                                               }
                                        }
 #ifdef _DEBUG_PETRINET_
                                        mnstr_printf(cntxt->fdout, 
"#PETRINET:%d tuples for %s, source %d\n", cnt, pnet[i].name, j);
@@ -570,7 +576,7 @@
 
                                t= GDKusec();
                                msg = reenterMAL(cntxt, mb, pnet[i].pc, 
pnet[i].pc + 1, glb, 0, 0);
-                               pnet[i].time += GDKusec() - t + analysis;
+                               pnet[i].time += GDKusec() - t + analysis;       
/* keep around in microseconds */
                                if ( msg != MAL_SUCCEED && !strstr(msg,"too 
early") ){
                                        char buf[BUFSIZ];
                                        if ( pnet[i].error == NULL ) {
@@ -580,9 +586,10 @@
                                        pnet[i].enabled = -1;
                                } else {
                                        pnet[i].cycles++;
+                                       (void) 
MTIMEcurrent_timestamp(&pnet[i].seen);
                                        for (j = 0; j < pnet[i].srctop; j++) {
                                                idx = pnet[i].source[j].bskt;
-                                               baskets[idx].seen = now;
+                                               (void) 
MTIMEcurrent_timestamp(&baskets[idx].seen);
                                                pnet[i].events += 
pnet[i].source[j].available;
                                                pnet[i].source[j].available = 
0;  /* force recount */
                                        }
@@ -632,7 +639,7 @@
 str
 PNtable(int *ret)
 {
-       BAT *bn, *name, *def, *status, *cycles, *events, *time, *error;
+       BAT *bn, *name, *def, *status, *seen, *cycles, *events, *time, *error;
        int i;
 
        bn = BATnew(TYPE_str, TYPE_bat, BATTINY);
@@ -648,6 +655,9 @@
        status = BATnew(TYPE_oid,TYPE_str, BATTINY);
        if ( status == 0 ) goto wrapup;
        BATseqbase(status,0);
+       seen = BATnew(TYPE_oid,TYPE_timestamp, BATTINY);
+       if ( seen == 0 ) goto wrapup;
+       BATseqbase(seen,0);
        cycles = BATnew(TYPE_oid,TYPE_int, BATTINY);
        if ( cycles == 0 ) goto wrapup;
        BATseqbase(cycles,0);
@@ -665,6 +675,7 @@
                BUNappend(name, pnet[i].name, FALSE);
                BUNappend(def, pnet[i].def, FALSE);
                BUNappend(status, statusnames[pnet[i].status], FALSE);
+               BUNappend(seen, &pnet[i].seen, FALSE);
                BUNappend(cycles, &pnet[i].cycles, FALSE);
                BUNappend(events, &pnet[i].events, FALSE);
                BUNappend(time, &pnet[i].time, FALSE);
@@ -672,6 +683,7 @@
        }
        BUNins(bn,"nme", & name->batCacheid, FALSE);
        BUNins(bn,"status", & status->batCacheid, FALSE);
+       BUNins(bn,"seen", & seen->batCacheid, FALSE);
        BUNins(bn,"cycles", & cycles->batCacheid, FALSE);
        BUNins(bn,"events", & events->batCacheid, FALSE);
        BUNins(bn,"time", & time->batCacheid, FALSE);
@@ -683,6 +695,7 @@
        BBPreleaseref(name->batCacheid);
        BBPreleaseref(def->batCacheid);
        BBPreleaseref(status->batCacheid);
+       BBPreleaseref(seen->batCacheid);
        BBPreleaseref(cycles->batCacheid);
        BBPreleaseref(events->batCacheid);
        BBPreleaseref(time->batCacheid);
@@ -693,6 +706,7 @@
        if ( name) BBPreleaseref(name->batCacheid);
        if ( def) BBPreleaseref(def->batCacheid);
        if ( status) BBPreleaseref(status->batCacheid);
+       if ( seen) BBPreleaseref(seen->batCacheid);
        if ( cycles) BBPreleaseref(cycles->batCacheid);
        if ( events) BBPreleaseref(events->batCacheid);
        if ( time) BBPreleaseref(time->batCacheid);
_______________________________________________
Checkin-list mailing list
[email protected]
http://mail.monetdb.org/mailman/listinfo/checkin-list

Reply via email to