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