Changeset: 5f39f7197ef8 for MonetDB
URL: http://dev.monetdb.org/hg/MonetDB?cmd=changeset;node=5f39f7197ef8
Modified Files:
        sql/backends/monet5/datacell/Tests/poc.sql
        sql/backends/monet5/datacell/Tests/scenario00.sql
        sql/backends/monet5/datacell/basket.mx
        sql/backends/monet5/datacell/datacell.mx
        sql/backends/monet5/datacell/emitter.mx
        sql/backends/monet5/datacell/opt_datacell.mx
        sql/backends/monet5/datacell/petrinet.mx
        sql/backends/monet5/datacell/receptor.mx
Branch: default
Log Message:

Properly pass data to emitter basket.
The SQL append statements are gathered and leads to a BSKTpass call,
which under locking updates the emitter basket. Likewise, the emitter
grabs the BATs from the basket and privately, lazely pass them over
the channel. Concurrently, more events may be delivered for sending.


diffs (truncated from 577 to 300 lines):

diff --git a/sql/backends/monet5/datacell/Tests/poc.sql 
b/sql/backends/monet5/datacell/Tests/poc.sql
deleted file mode 100644
--- a/sql/backends/monet5/datacell/Tests/poc.sql
+++ /dev/null
@@ -1,54 +0,0 @@
--- A straw-man's datacell proof of concept implementation
--- This example relies on the standard SQL facilities
--- to demonstrate the DataCell functionality.
--- The specific runtime issues are handled by an optimizer
-
-create schema datacell;
-set optimizer='datacell_pipe';
-
-create table datacell.sys_in(
-    id integer,
-    tag timestamp,
-    payload integer
-);
-
--- to be used by continous queries
--- could be generated from the table definition.
-create function datacell.sys_in()
-returns table (id integer, tag timestamp, payload integer)
-begin
-       return select * from datacell.sys_in;
-end;
-
-select * from datacell.sys_in();
-
-call datacell.basket('datacell','sys_in');
-call datacell.receptor('datacell','sys_in','localhost',50500,'passive');
-call datacell.start('datacell','sys_in');
-
-create table datacell.sysOut( etag timestamp, like datacell.sys_in);
-
--- the continues query is registered and started
-create procedure datacell.query_sys_in_sysOut()
-begin
-       insert into datacell.sysOut select now(), *  from datacell.sys_in() ;
-end;
-call register_query('datacell.query_sys_in_sysOut');
-call start_query('datacell.query_sys_in_sysOut');
-
-create function datacell.sysOut()
-returns table (id integer)
-begin
-       return select * from datacell.sysOut;
-end;
-
--- sent it over the wire.
-create function datacell.emitter_sysOut()
-returns boolean;
-begin
-       select * from datacell.sysOut();
-       return true;
-end;
-call datacell.basket('datacell','sysOut');
-call datacell.emitter('datacell','sysOut','localhost:50101');
-call datacell.start('datacell','sysOut');
diff --git a/sql/backends/monet5/datacell/Tests/scenario00.sql 
b/sql/backends/monet5/datacell/Tests/scenario00.sql
new file mode 100644
--- /dev/null
+++ b/sql/backends/monet5/datacell/Tests/scenario00.sql
@@ -0,0 +1,52 @@
+-- Scenario to exercise the datacell implementation
+-- using a single receptor and emitter
+-- The sensor data is simple passed to the actuator.
+
+set optimizer='datacell_pipe';
+
+create table datacell.bsktin(
+    id integer,
+    tag timestamp,
+    payload integer
+);
+create table datacell.bsktout( like datacell.bsktin);
+
+-- initialize the baskets
+-- call datacell.prelude();
+call datacell.basket('datacell.bsktin');
+call datacell.basket('datacell.bsktout');
+
+-- initialize receptor
+call datacell.receptor('datacell.bsktin','localhost',50500);
+call datacell.mode('datacell.bsktin','passive');
+call datacell.protocol('datacell.bsktin','udp');
+call datacell.resume('datacell.bsktin');
+
+-- externally, activate the sensor leaving some in the basket
+--sensor --host=localhost --port=50500 --events=100 --columns=3 --delay=1
+
+-- initialize emitter
+call datacell.emitter('datacell.bsktout','localhost',50600);
+call datacell.mode('datacell.bsktout','active');
+call datacell.protocol('datacell.bsktout','udp');
+call datacell.resume('datacell.bsktout');
+
+-- externally, activate the actuator server to listen
+-- actuator 
+
+-- compile the continous query
+call datacell.query('datacell.pass', 'insert into datacell.bsktout select * 
from datacell.bsktin;');
+call datacell.register('datacell.pass');
+
+-- start the datacell scheduler
+call datacell.resume();
+call datacell.dump();
+
+-- wrapup
+-- stop the datacell scheduler
+call datacell.postlude();
+
+-- remove everything
+drop table datacell.bsktin;
+drop table datacell.bsktout;
+
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
@@ -40,14 +40,22 @@
 address BSKTunlock
 comment "Unlock the basket";
 
-pattern grab(tbl:str, flg:int):bat[:oid,:any]...
+pattern grab(tbl:str):bat[:oid,:any]...
 address BSKTgrab
 comment "Take a snapshot of the basket, destroying the origin when flg is set";
 
+pattern pass(tbl:str, cols:bat[:oid,:any]...)
+address BSKTpass
+comment "Dump the new tuples into the basket";
+
 command drop(tbl:str):void
 address BSKTdrop
 comment "Remove the basket";
 
+command reset():void
+address BSKTreset
+comment "Remove all baskets";
+
 command dump()
 address BSKTdump
 comment "Dump the status of the basket table";
@@ -88,11 +96,13 @@
 datacell_export str schema_default;
 datacell_export str BSKTregister(Client cntxt, MalBlkPtr mb, MalStkPtr stk, 
InstrPtr pci);
 datacell_export str BSKTdrop(int *ret, str *tbl);
+datacell_export str BSKTreset(int *ret);
 datacell_export str BSKTinventory(int *ret);
 datacell_export int BSKTmemberCount(str tbl);
 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 BSKTlock(int *ret, str *tbl, int *delay);
 datacell_export str BSKTunlock(int *ret, str *tbl);
@@ -101,6 +111,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 void BSKTtolower(char *src);
 
 datacell_export BSKTbasketRec *baskets;
@@ -346,6 +357,14 @@
 }
 
 str
+BSKTreset(int *ret)
+{
+       int i;
+       for ( i = 1; i < bsktLimit; i++)
+               BSKTdrop(ret, &baskets[i].name);
+       return MAL_SUCCEED;
+}
+str
 BSKTdump(int *ret)
 {
        int bskt;
@@ -365,13 +384,12 @@
 BSKTgrab(Client cntxt, MalBlkPtr mb, MalStkPtr stk, InstrPtr pci)
 {
        str tbl;
-       int bskt, i, *ret, flg;
+       int bskt, i, *ret;
        BAT *b;
 
        (void) cntxt;
        (void) mb;
-       tbl = *(str*) getArgReference(stk,pci, pci->argc-2);
-       flg = *(int*) getArgReference(stk,pci, pci->argc-1);
+       tbl = *(str*) getArgReference(stk,pci, pci->argc-1);
 
        bskt = BSKTlocate(tbl);
        if (bskt == 0 )
@@ -385,13 +403,41 @@
                *ret = b->batCacheid;
                BBPfix(*ret);
                BBPkeepref(*ret);
-               if ( flg) {
-                       if (  b->htype == TYPE_oid) {
-                               baskets[bskt].primary[i] = BATnew(TYPE_void, 
b->ttype, 0);
-                               BATseqbase(baskets[bskt].primary[i], 0);
-                       } else
-                               baskets[bskt].primary[i] = BATnew(b->htype, 
b->ttype, 0);
-               }
+               if (  b->htype == TYPE_oid) {
+                       baskets[bskt].primary[i] = BATnew(TYPE_void, b->ttype, 
0);
+                       BATseqbase(baskets[bskt].primary[i], 0);
+               } else
+                       baskets[bskt].primary[i] = BATnew(b->htype, b->ttype, 
0);
+       }
+       mal_unset_lock(baskets[bskt].lock, "unlock basket");
+       return MAL_SUCCEED;
+}
+
+str
+BSKTpass(Client cntxt, MalBlkPtr mb, MalStkPtr stk, InstrPtr pci)
+{
+       str tbl;
+       int bskt, i, j, ret;
+       BAT *b, *bn;
+
+       (void) cntxt;
+       (void) mb;
+       tbl = *(str*) getArgReference(stk,pci, pci->retc);
+
+       bskt = BSKTlocate(tbl);
+       if (bskt == 0 )
+               throw(MAL,"basket.pass","Basket not found");
+       if ( baskets[bskt].colcount != pci->retc - 3)
+               throw(MAL,"basket.pass","Non-matchin arguments");
+
+       /* copy the content of the temporary BATs into the basket */
+       mal_set_lock(baskets[bskt].lock, "lock basket");
+       for ( j = 2,  i=0; i < baskets[bskt].colcount; i++, j++) {
+               ret= *(int*) getArgReference(stk,pci,j);
+               b = baskets[bskt].primary[i];
+               bn = BATdescriptor(ret);
+               BATappend(b, bn, FALSE);
+               BBPreleaseref(ret);
        }
        mal_unset_lock(baskets[bskt].lock, "unlock basket");
        return MAL_SUCCEED;
@@ -420,4 +466,27 @@
        p= pushStr(mb,p,tbl);
        return p;
 }
+
+InstrPtr
+BSKTpassInstruction(MalBlkPtr mb, str tbl)
+{
+       int i, j, bskt;
+       InstrPtr p;
+       BAT *b;
+
+       bskt = BSKTlocate(tbl);
+       if (bskt == 0 )
+               return 0;
+       p = newInstruction(mb, ASSIGNsymbol);
+       getArg(p,0)= newTmpVariable(mb,TYPE_any);
+       getModuleId(p) = basketRef;
+       getFunctionId(p) = putName("pass",4);
+       p = pushStr(mb,p, tbl);
+       for ( i=0; i < baskets[bskt].colcount; i++) {
+               b = baskets[bskt].primary[i];
+               j= newTmpVariable(mb,newBatType(TYPE_oid,b->ttype));
+               p= pushArgument(mb,p, j);
+       }
+       return p;
+}
 @}
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
@@ -386,9 +386,10 @@
 DCpostlude(Client cntxt, MalBlkPtr mb, MalStkPtr stk, InstrPtr pci)
 {
        int ret=0;
-       if ( DCprepared != DCINITIALIZED && DCprepared != DCPAUSED)
-               throw(MAL,"datacell.pause","Datacell already stopped");
        PNstopScheduler(&ret);
+       RCreset(&ret);
+       EMreset(&ret);
+       BSKTreset(&ret);
        (void) cntxt;
        (void) mb;
        (void) stk;
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
@@ -25,9 +25,6 @@
 After this call it will sent tuples from basket X_p1
 to the stream Y at the localhost default port.
 
-[tag each tuple with the latency?, use as derived
-baskets? [arrival,departure,latency]]
-
 Each emitter is supported by an independent thread
_______________________________________________
Checkin-list mailing list
[email protected]
http://mail.monetdb.org/mailman/listinfo/checkin-list

Reply via email to