Changeset: 568e0ace05f4 for MonetDB
URL: http://dev.monetdb.org/hg/MonetDB?cmd=changeset;node=568e0ace05f4
Modified Files:
        sql/backends/monet5/datacell/Tests/scenario05.sql
        sql/backends/monet5/datacell/basket.mx
        sql/backends/monet5/datacell/datacell.mx
        sql/backends/monet5/datacell/opt_datacell.mx
        sql/backends/monet5/datacell/petrinet.mx
Branch: default
Log Message:

add a heart beat query
The optimizer moves the window/threshold and beat to the front
provided it consists of constants. In case it is later
detected we are too early, a singled exception is sent
to the petrinet controller.


diffs (294 lines):

diff --git a/sql/backends/monet5/datacell/Tests/scenario05.sql 
b/sql/backends/monet5/datacell/Tests/scenario05.sql
new file mode 100644
--- /dev/null
+++ b/sql/backends/monet5/datacell/Tests/scenario05.sql
@@ -0,0 +1,36 @@
+-- Scenario to exercise the datacell implementation
+-- using a single receptor and emitter
+-- this is the extended version of scenario00
+-- with sliding beatdow and a 2 seconds delay
+
+create schema datacell;
+set optimizer='datacell_pipe';
+
+create table datacell.beatin(
+    id integer,
+    tag timestamp,
+    payload integer
+);
+create table datacell.beatout( tag timestamp, mi integer, ma integer, su 
bigint);
+
+call datacell.receptor('datacell.beatin','localhost',50500);
+
+call datacell.emitter('datacell.beatout','localhost',50600);
+
+call datacell.query('datacell.mavgbeat', 'insert into datacell.beatout select 
now(), min(payload), 
+       max(payload), sum(payload) from datacell.beatin where 
datacell.beat(\'datacell.beatin\',2000) and 
datacell.window(\'datacell.beatin\',10,1);');
+
+call datacell.resume();
+call datacell.dump();
+
+-- externally, activate the sensor 
+--sensor --host=localhost --port=50500 --events=100 --columns=3 --delay=1
+-- externally, activate the actuator server to listen
+-- actuator 
+
+
+-- wrapup
+call datacell.postlude();
+drop table datacell.beatin;
+drop table datacell.beatout;
+
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
@@ -44,12 +44,12 @@
 address BSKTgrab
 comment "Take a snapshot of the basket, destroying the origin when flg is set";
 
-pattern pass(tbl:str, cols:any...)
-address BSKTpass
+pattern update(tbl:str, cols:any...)
+address BSKTupdate
 comment "Dump the new tuples into the basket";
 
-pattern pass(tbl:str, cols:bat[:oid,:any]...)
-address BSKTpass
+pattern update(tbl:str, cols:bat[:oid,:any]...)
+address BSKTupdate
 comment "Dump the new tuples into the basket";
 
 command drop(tbl:str):void
@@ -60,10 +60,14 @@
 address BSKTthreshold
 comment "Set an acceptance threshold of N events before inspecting";
 
-command window(tbl:str,N:int, S:int):bit
+command window{unsafe}(tbl:str,N:int, S:int):bit
 address BSKTwindow
 comment "Use a window of N event and slide S afterwards";
 
+command beat(tbl:str,N:int):bit
+address BSKTbeat
+comment "Set an delay to N milliseconds";
+
 command reset():void
 address BSKTreset
 comment "Remove all baskets";
@@ -100,9 +104,9 @@
        MT_Lock lock;
        str name;       /* table that represents the basket */
        int threshold ; /* bound to determine scheduling eligibility */
-       int delay;      /* heartbeat to schedule next inspection */
        int winsize, winslide; /* sliding window operations */
        lng lastseen;
+       lng beat;       /* milliseconds delay */
        int colcount;
        str *cols;
        BAT **primary;  
@@ -118,8 +122,9 @@
 datacell_export int BSKTlocate(str tbl);
 datacell_export str BSKTdump(int *ret);
 datacell_export str BSKTgrab(Client cntxt, MalBlkPtr mb, MalStkPtr stk, 
InstrPtr pci);
-datacell_export str BSKTpass(Client cntxt, MalBlkPtr mb, MalStkPtr stk, 
InstrPtr pci);
+datacell_export str BSKTupdate(Client cntxt, MalBlkPtr mb, MalStkPtr stk, 
InstrPtr pci);
 datacell_export str BSKTthreshold(int *ret, str *tbl, int *sz);
+datacell_export str BSKTbeat(int *ret, str *tbl, int *sz);
 datacell_export str BSKTwindow(int *ret, str *tbl, int *sz, int *slide);
 
 datacell_export str BSKTlock(int *ret, str *tbl, int *delay);
@@ -129,7 +134,7 @@
 datacell_export str BSKTnewbasket(sql_schema *s, sql_table *t, sql_trans *tr);
 datacell_export void BSKTelements(str nme, str buf, str *schema, str *tbl);
 datacell_export InstrPtr BSKTgrabInstruction(MalBlkPtr mb, str tbl);
-datacell_export InstrPtr BSKTpassInstruction(MalBlkPtr mb, str tbl);
+datacell_export InstrPtr BSKTupdateInstruction(MalBlkPtr mb, str tbl);
 datacell_export void BSKTtolower(char *src);
 
 datacell_export BSKTbasketRec *baskets;
@@ -384,12 +389,13 @@
 
        for ( bskt = 0; bskt < bsktLimit; bskt++) 
        if ( baskets[bskt].name){
-               mnstr_printf(GDKout, "#baskets[%2d] %s columns %d threshold %d 
window=[%d,%d] events " BUNFMT "\n", bskt,
+               mnstr_printf(GDKout, "#baskets[%2d] %s columns %d threshold %d 
window=[%d,%d] beat %d events " BUNFMT "\n", bskt,
                        baskets[bskt].name,
                        baskets[bskt].colcount,
                        baskets[bskt].threshold,
                        baskets[bskt].winsize,
                        baskets[bskt].winslide,
+                       baskets[bskt].beat,
                        (baskets[bskt].primary[0]? 
BATcount(baskets[bskt].primary[0]): 0));
        }
        (void)ret;
@@ -420,6 +426,9 @@
 
                /* take care of sliding windows */
                if ( baskets[bskt].winsize ){
+                       /* we may be too early */
+                       if ( BATcount(b) < (BUN) baskets[bskt].winsize)
+                               break;
                        bn = BATcopy(b, b->htype, b->ttype,TRUE);
                        v = BATslice(bn, baskets[bskt].winslide,BATcount(bn));
                        BATclear(b);
@@ -436,11 +445,13 @@
                BBPkeepref(*ret);
        }
        mal_unset_lock(baskets[bskt].lock, "unlock basket");
+       if ( i != baskets[bskt].colcount)
+               throw(MAL,"basket.grab","too early");
        return MAL_SUCCEED;
 }
 
 str
-BSKTpass(Client cntxt, MalBlkPtr mb, MalStkPtr stk, InstrPtr pci)
+BSKTupdate(Client cntxt, MalBlkPtr mb, MalStkPtr stk, InstrPtr pci)
 {
        str tbl;
        int bskt, i, j, ret;
@@ -452,9 +463,9 @@
 
        bskt = BSKTlocate(tbl);
        if (bskt == 0 )
-               throw(MAL,"basket.pass","Basket not found");
+               throw(MAL,"basket.update","Basket not found");
        if ( baskets[bskt].colcount != pci->argc - 2)
-               throw(MAL,"basket.pass","Non-matching arguments");
+               throw(MAL,"basket.update","Non-matching arguments");
 
        /* copy the content of the temporary BATs into the basket */
        mal_set_lock(baskets[bskt].lock, "lock basket");
@@ -494,7 +505,7 @@
 }
 
 InstrPtr
-BSKTpassInstruction(MalBlkPtr mb, str tbl)
+BSKTupdateInstruction(MalBlkPtr mb, str tbl)
 {
        int i, j, bskt;
        InstrPtr p;
@@ -506,7 +517,7 @@
        p = newInstruction(mb, ASSIGNsymbol);
        getArg(p,0)= newTmpVariable(mb,TYPE_any);
        getModuleId(p) = basketRef;
-       getFunctionId(p) = putName("pass",4);
+       getFunctionId(p) = putName("update",6);
        p = pushStr(mb,p, tbl);
        for ( i=0; i < baskets[bskt].colcount; i++) {
                b = baskets[bskt].primary[i];
@@ -553,4 +564,22 @@
        *ret = TRUE;
        return MAL_SUCCEED;
 }
+
+str
+BSKTbeat(int *ret, str *tbl, int *sz)
+{
+       int bskt;
+       bskt = BSKTlocate(*tbl);
+       if (bskt == 0 )
+               throw(MAL,"basket.beat","Basket not found");
+       if ( *sz < 0)
+               throw(MAL,"basket.beat","Illegal value");
+       if ( *sz < baskets[bskt].winsize)
+               throw(MAL,"basket.threshold","Threshold smaller then window 
size");
+       baskets[bskt].beat = *sz;
+       if ( baskets[bskt].beat + baskets[bskt].lastseen < GDKusec() )
+               throw(MAL,"basket.heat","too early");
+       *ret = TRUE;
+       return MAL_SUCCEED;
+}
 @}
diff --git a/sql/backends/monet5/datacell/datacell.mx 
b/sql/backends/monet5/datacell/datacell.mx
--- a/sql/backends/monet5/datacell/datacell.mx
+++ b/sql/backends/monet5/datacell/datacell.mx
@@ -80,11 +80,11 @@
 address DCdump
 comment "Dump receptor/emitter status";
 
-command threshold(bskt:str, mi:int):bit
+command threshold{unsafe}(bskt:str, mi:int):bit
 address DCthreshold 
 comment "Prefilter on bucket content";
 
-command window(bskt:str, size:int, slide:int):bit
+command window{unsafe}(bskt:str, size:int, slide:int):bit
 address DCwindow
 comment "Slide the basket with a limited number of events";
 
@@ -460,14 +460,9 @@
 }
 
 str
-DCbeat(int *ret, str *bskt, int *t)
+DCbeat(int *ret, str *bskt, int *beat)
 {
-       int idx;
-       idx = BSKTlocate(*bskt);
-       if ( idx == 0)
-               throw(SQL, "datacell.beat", "Basket not found");
-       *ret = TRUE;
-       baskets[idx].delay = *t;
+       return BSKTbeat(ret,bskt,beat);
        return MAL_SUCCEED;
 }
 @}
diff --git a/sql/backends/monet5/datacell/opt_datacell.mx 
b/sql/backends/monet5/datacell/opt_datacell.mx
--- a/sql/backends/monet5/datacell/opt_datacell.mx
+++ b/sql/backends/monet5/datacell/opt_datacell.mx
@@ -96,6 +96,25 @@
                        /* inject transaction start */
                        newFcnCall(mb,sqlRef,putName("transaction",11));
                
+               if ( getModuleId(p) == datacellRef && getFunctionId(p) == 
putName("window",6) &&
+                       isVarConstant(mb, getArg(p,1)) && isVarConstant(mb, 
getArg(p,2)) && isVarConstant(mb,getArg(p,3)) ){
+                       /* let's move the window to the start of the block  
when it consists of constants*/
+                       pushInstruction(mb,p);
+                       for ( j = mb->stop-1; j >2; j--)
+                               mb->stmt[j]= mb->stmt[j-1];
+                       mb->stmt[j] = p;
+                       continue;
+               }
+               if ( getModuleId(p) == datacellRef &&  ( getFunctionId(p) == 
putName("threshold",9)  || getFunctionId(p) == putName("beat",4) ) &&
+                       isVarConstant(mb, getArg(p,1)) && isVarConstant(mb, 
getArg(p,2)) ){
+                       /* let's move the threshold/beat to the start of the 
block  when it consists of constants*/
+                       pushInstruction(mb,p);
+                       for ( j = mb->stop-1; j >2; j--)
+                               mb->stmt[j]= mb->stmt[j-1];
+                       mb->stmt[j] = p;
+                       continue;
+               }
+
                if ( p->token == ENDsymbol){
                        /* a good place to commit the SQL transaction */
                        for ( j = 0; j < a; j++)
@@ -188,7 +207,7 @@
                                if ( j == a ){
                                        if ( a == maxbasket) return 0;
                                        /* grab the basket tables instruction */
-                                       qa[a]= BSKTpassInstruction(mb,buf);
+                                       qa[a]= BSKTupdateInstruction(mb,buf);
                                        appends[a++]= GDKstrdup(buf);
                                        actions = 1;
                                }
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
@@ -530,7 +530,7 @@
                                                break;
                                        }
                                        /* check heart beat delays */
-                                       if ( baskets[idx].lastseen + 
baskets[idx].delay > now){
+                                       if ( baskets[idx].lastseen + 
baskets[idx].beat > now){
                                                pnet[i].enabled = 0;
                                                break;
                                        }
@@ -571,7 +571,7 @@
                                t= GDKusec();
                                msg = reenterMAL(cntxt, mb, pnet[i].pc, 
pnet[i].pc + 1, glb, 0, 0);
                                pnet[i].cycletime += GDKusec() - t;
-                               if ( msg != MAL_SUCCEED){
+                               if ( msg != MAL_SUCCEED && !strstr(msg,"too 
early") ){
                                        char buf[BUFSIZ];
                                        if ( pnet[i].error == NULL ) {
                                                snprintf(buf,BUFSIZ-1,"Query %s 
failed:%s", pnet[i].fcn, msg);
_______________________________________________
Checkin-list mailing list
[email protected]
http://mail.monetdb.org/mailman/listinfo/checkin-list

Reply via email to