Changeset: e7cbea4b9dd6 for MonetDB
URL: http://dev.monetdb.org/hg/MonetDB?cmd=changeset;node=e7cbea4b9dd6
Modified Files:
sql/backends/monet5/datacell/50_datacell.sql
sql/backends/monet5/datacell/Tests/scenario00.sql
sql/backends/monet5/datacell/Tests/scenario01.sql
sql/backends/monet5/datacell/Tests/scenario02.sql
sql/backends/monet5/datacell/actuator.mx
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:
Cleanup of mode/protocol administartion and use
Fixing various smaller bugs and simplifying the datacell interface
to make the default settings easier to use.
diffs (truncated from 584 to 300 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
@@ -25,10 +25,10 @@
external name datacell.inventory;
create procedure datacell.receptor(tbl string, host string, portid integer)
- external name receptor."start";
+ external name datacell.receptor;
create procedure datacell.emitter(tbl string, host string, portid integer)
- external name emitter."start";
+ external name datacell.emitter;
create procedure datacell.mode(tbl string, mode string)
external name datacell.mode;
diff --git a/sql/backends/monet5/datacell/Tests/scenario00.sql
b/sql/backends/monet5/datacell/Tests/scenario00.sql
--- a/sql/backends/monet5/datacell/Tests/scenario00.sql
+++ b/sql/backends/monet5/datacell/Tests/scenario00.sql
@@ -1,6 +1,7 @@
-- Scenario to exercise the datacell implementation
-- using a single receptor and emitter
-- The sensor data is simple passed to the actuator.
+-- it is closest to the web description
set optimizer='datacell_pipe';
@@ -9,44 +10,24 @@
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');
+call datacell.receptor('datacell.bsktin','localhost',50500);
--- 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');
+call datacell.emitter('datacell.bsktout','localhost',50600);
--- externally, activate the sensor leaving some in the basket
+call datacell.query('datacell.pass', 'insert into datacell.bsktout select *
from datacell.bsktin;');
+
+call datacell.resume();
+call datacell.dump();
+
+-- externally, activate the sensor
--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/Tests/scenario01.sql
b/sql/backends/monet5/datacell/Tests/scenario01.sql
new file mode 100644
--- /dev/null
+++ b/sql/backends/monet5/datacell/Tests/scenario01.sql
@@ -0,0 +1,53 @@
+-- Scenario to exercise the datacell implementation
+-- using a single receptor and emitter
+-- The sensor data is simple passed to the actuator.
+-- this is the extended version of scenario00
+
+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/Tests/scenario02.sql
b/sql/backends/monet5/datacell/Tests/scenario02.sql
new file mode 100644
--- /dev/null
+++ b/sql/backends/monet5/datacell/Tests/scenario02.sql
@@ -0,0 +1,33 @@
+-- Scenario to exercise the datacell implementation
+-- using a single receptor and emitter
+-- Monitor the aggregation level
+
+set optimizer='datacell_pipe';
+
+create table datacell.bsktin(
+ id integer,
+ tag timestamp,
+ payload integer
+);
+create table datacell.bsktout( tag timestamp, cnt integer);
+
+call datacell.receptor('datacell.bsktin','localhost',50500);
+
+call datacell.emitter('datacell.bsktout','localhost',50600);
+
+call datacell.query('datacell.pass', 'insert into datacell.bsktout select
now(), count(*) from datacell.bsktin;');
+
+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.bsktin;
+drop table datacell.bsktout;
+
diff --git a/sql/backends/monet5/datacell/actuator.mx
b/sql/backends/monet5/datacell/actuator.mx
--- a/sql/backends/monet5/datacell/actuator.mx
+++ b/sql/backends/monet5/datacell/actuator.mx
@@ -313,7 +313,6 @@
consumeStream(ac);
}
- /* Handle TCP protocol */
if (server && (err = socket_server_connect(&sockfd, port))) {
mnstr_printf(ACout, "ACTUATOR:start server:%s\n", err);
return 0;
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
@@ -355,6 +355,7 @@
{
int i;
for ( i = 1; i < bsktLimit; i++)
+ if ( baskets[i].name)
BSKTdrop(ret, &baskets[i].name);
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
@@ -29,11 +29,11 @@
comment "Initialize a new basket based on a specific table definition in the
datacell schema");
pattern emitter(tbl:str, host:str, port:int)
-address DCemitterNew
+address DCemitter
comment "Define a emitter based on a basket table.";
pattern receptor(tbl:str, host:str, port:int)
-address DCreceptorNew
+address DCreceptor
comment "Define a receptor based on a basket table..";
pattern mode(tbl:str, arg:str)
@@ -53,7 +53,7 @@
comment "Pause a receptor or emitter";
pattern resume(obj:str):void
-address DCresume
+address DCresumeObject
comment "Pause a receptor or emitter";
pattern query(name:str,def:str):void
@@ -79,7 +79,7 @@
comment "(Re)start the petrinet scheduler.";
pattern resume()
-address DCresumeScheduler
+address DCresume
comment "Resume the petrinet scheduler.";
pattern postlude()
@@ -121,9 +121,11 @@
/* #define _DEBUG_DATACELL debug this module */
datacell_export str DCprelude(Client cntxt, MalBlkPtr mb, MalStkPtr stk,
InstrPtr pci);
+datacell_export str DCreceptor(Client cntxt, MalBlkPtr mb, MalStkPtr stk,
InstrPtr pci);
datacell_export str DCregister(Client cntxt, MalBlkPtr mb, MalStkPtr stk,
InstrPtr pci);
datacell_export str DCremove(Client cntxt, MalBlkPtr mb, MalStkPtr stk,
InstrPtr pci);
datacell_export str DCpause(Client cntxt, MalBlkPtr mb, MalStkPtr stk,
InstrPtr pci);
+datacell_export str DCresumeObject(Client cntxt, MalBlkPtr mb, MalStkPtr stk,
InstrPtr pci);
datacell_export str DCresume(Client cntxt, MalBlkPtr mb, MalStkPtr stk,
InstrPtr pci);
datacell_export str DCmode(Client cntxt, MalBlkPtr mb, MalStkPtr stk, InstrPtr
pci);
datacell_export str DCprotocol(Client cntxt, MalBlkPtr mb, MalStkPtr stk,
InstrPtr pci);
@@ -194,6 +196,32 @@
}
str
+DCreceptor(Client cntxt, MalBlkPtr mb, MalStkPtr stk, InstrPtr pci)
+{
+ int *ret = (int *) getArgReference(stk,pci,0);
+ str *tbl = (str *) getArgReference(stk,pci,1);
+ str *host = (str *) getArgReference(stk,pci,2);
+ int *port = (int *) getArgReference(stk,pci,3);
+ int idx = BSKTlocate(*tbl);
+ if ( idx == 0)
+ BSKTregister(cntxt,mb,stk,pci);
+ return DCreceptorNew(ret,tbl,host,port);
+}
+
+str
+DCemitter(Client cntxt, MalBlkPtr mb, MalStkPtr stk, InstrPtr pci)
+{
+ int *ret = (int *) getArgReference(stk,pci,0);
+ str *tbl = (str *) getArgReference(stk,pci,1);
+ str *host = (str *) getArgReference(stk,pci,2);
+ int *port = (int *) getArgReference(stk,pci,3);
+ int idx = BSKTlocate(*tbl);
+ if ( idx == 0)
+ BSKTregister(cntxt,mb,stk,pci);
+ return DCemitterNew(ret,tbl,host,port);
+}
+
+str
DCregister(Client cntxt, MalBlkPtr mb, MalStkPtr stk, InstrPtr pci)
{
return BSKTregister(cntxt,mb,stk,pci);
@@ -217,7 +245,7 @@
}
str
-DCresume(Client cntxt, MalBlkPtr mb, MalStkPtr stk, InstrPtr pci)
+DCresumeObject(Client cntxt, MalBlkPtr mb, MalStkPtr stk, InstrPtr pci)
{
int idx, ret= 0;
str tbl= *(str*) getArgReference(stk, pci,1);
@@ -234,6 +262,15 @@
}
str
+DCresume(Client cntxt, MalBlkPtr mb, MalStkPtr stk, InstrPtr pci)
+{
+ int ret;
+ RCresume(&ret);
+ EMresume(&ret);
+ return DCresumeScheduler(cntxt,mb,stk,pci);
+}
+
+str
_______________________________________________
Checkin-list mailing list
[email protected]
http://mail.monetdb.org/mailman/listinfo/checkin-list