Changeset: 81ca3c9bc6ac for MonetDB
URL: https://dev.monetdb.org/hg/MonetDB?cmd=changeset;node=81ca3c9bc6ac
Modified Files:
sql/backends/monet5/Tests/cqcreate.sql
sql/backends/monet5/Tests/cqstream00.sql
sql/backends/monet5/Tests/cqstream01.sql
sql/backends/monet5/Tests/cquery00.sql
sql/backends/monet5/Tests/cquery20.sql
sql/backends/monet5/sql_cquery.c
sql/backends/monet5/sql_cquery.h
sql/backends/monet5/sql_execute.c
sql/server/rel_psm.c
sql/server/sql_mvc.h
sql/server/sql_parser.y
sql/server/sql_scan.c
Branch: trails
Log Message:
Added BEAT and CYCLES to the continuous query definition
diffs (truncated from 415 to 300 lines):
diff --git a/sql/backends/monet5/Tests/cqcreate.sql
b/sql/backends/monet5/Tests/cqcreate.sql
--- a/sql/backends/monet5/Tests/cqcreate.sql
+++ b/sql/backends/monet5/Tests/cqcreate.sql
@@ -2,7 +2,7 @@ CREATE stream TABLE testing (a int);
CREATE TABLE results (b int);
-CREATE CONTINUOUS PROCEDURE stressing() BEGIN INSERT INTO results SELECT a
FROM testing; END;
+CREATE PROCEDURE stressing() BEGIN INSERT INTO results SELECT a FROM testing;
END;
-START CONTINUOUS PROCEDURE;
-STOP CONTINUOUS PROCEDURE;
+START CONTINUOUS stressing();
+STOP CONTINUOUS stressing();
diff --git a/sql/backends/monet5/Tests/cqstream00.sql
b/sql/backends/monet5/Tests/cqstream00.sql
--- a/sql/backends/monet5/Tests/cqstream00.sql
+++ b/sql/backends/monet5/Tests/cqstream00.sql
@@ -12,7 +12,7 @@ insert into stmp2 values('2005-09-23 12:
create table result1(like stmp2);
create table result2(like stmp2);
--- CREATE CONTINUOUS QUERY cq_splitter
+-- CREATE PROCEDURE cq_splitter
create procedure cq_splitter()
begin
insert into result1 select * from stmp2 where val <12;
diff --git a/sql/backends/monet5/Tests/cqstream01.sql
b/sql/backends/monet5/Tests/cqstream01.sql
--- a/sql/backends/monet5/Tests/cqstream01.sql
+++ b/sql/backends/monet5/Tests/cqstream01.sql
@@ -7,8 +7,8 @@ insert into stmp2 values('2005-09-23 12:
insert into stmp2 values('2005-09-23 12:34:28.000',1,13.0);
insert into stmp2 values('2005-09-23 12:34:28.000',1,15.0);
--- CREATE CONTINUOUS QUERY cq_window
-create continuous procedure cq_window()
+-- CREATE procedure cq_window
+create procedure cq_window()
begin
-- The window ensures a maximal number of tuples to consider
-- Could be considered a property of the stream table
diff --git a/sql/backends/monet5/Tests/cquery00.sql
b/sql/backends/monet5/Tests/cquery00.sql
--- a/sql/backends/monet5/Tests/cquery00.sql
+++ b/sql/backends/monet5/Tests/cquery00.sql
@@ -4,7 +4,7 @@ insert into testing values(123);
create table results (a int);
-create continuous procedure myproc()
+create procedure myproc()
begin
insert into results select a from sys.testing;
END;
diff --git a/sql/backends/monet5/Tests/cquery20.sql
b/sql/backends/monet5/Tests/cquery20.sql
--- a/sql/backends/monet5/Tests/cquery20.sql
+++ b/sql/backends/monet5/Tests/cquery20.sql
@@ -3,7 +3,7 @@
create stream table cqtbl(i integer);
-- the hello example
-create continuous procedure cqfoo(v integer)
+create procedure cqfoo(v integer)
begin
insert into cqtbl values(v);
end;
diff --git a/sql/backends/monet5/sql_cquery.c b/sql/backends/monet5/sql_cquery.c
--- a/sql/backends/monet5/sql_cquery.c
+++ b/sql/backends/monet5/sql_cquery.c
@@ -467,7 +467,7 @@ CQregisterInternal(Client cntxt, str mod
fprintFunction(stderr, nmb, 0, LIST_MAL_ALL);
#endif
// and hand it over to the scheduler
- msg = CQregisterMAL(cntxt, nmb,0,0);
+ msg = CQregister(cntxt, nmb,0,0);
if( msg != MAL_SUCCEED){
freeSymbol(s);
}
@@ -556,16 +556,26 @@ CQprocedure(Client cntxt, MalBlkPtr mb,
*/
str
-CQregisterMAL(Client cntxt, MalBlkPtr mb, MalStkPtr stk, InstrPtr pci )
+CQregister(Client cntxt, MalBlkPtr mb, MalStkPtr stk, InstrPtr pci )
{
- int i;
+ mvc* sqlcontext = ((backend *) cntxt->sqlcontext)->mvc;
str msg = MAL_SUCCEED;
InstrPtr sig = getInstrPtr(mb,0),q;
Symbol s;
+ int i, cycles = sqlcontext ? sqlcontext->cycles : int_nil, heartbeats =
sqlcontext ? sqlcontext->heartbeats : 1;
(void) pci;
(void) stk;
+ if(cycles <= 0 && cycles != int_nil){
+ msg = createException(SQL,"cquery.register","The cycles value
must be positive");
+ goto finish;
+ }
+ if(heartbeats <= 0){
+ msg = createException(SQL,"cquery.register","The heartbeats
value must be positive");
+ goto finish;
+ }
+
/* extract the actual procedure call and check for duplicate*/
for(i = 1; i< mb->stop; i++){
sig= getInstrPtr(mb,i);
@@ -607,8 +617,8 @@ CQregisterMAL(Client cntxt, MalBlkPtr mb
pnet[pnettop].stmt = instruction2str(mb,stk,sig,LIST_MAL_CALL);
pnet[pnettop].mb = mb;
pnet[pnettop].stk = prepareMALstack(mb, mb->vsize);
- pnet[pnettop].cycles = int_nil;
- pnet[pnettop].beats = lng_nil;
+ pnet[pnettop].cycles = cycles;
+ pnet[pnettop].beats = heartbeats * 1000;
pnet[pnettop].run = lng_nil;
pnet[pnettop].seen = *timestamp_nil;
pnet[pnettop].status = CQWAIT;
@@ -622,16 +632,6 @@ finish:
return msg;
}
-str
-CQregister(Client cntxt, MalBlkPtr mb, MalStkPtr stk, InstrPtr pci)
-{
- str modnme = *getArgReference_str(stk, pci, 1);
- str fcnnme = *getArgReference_str(stk, pci, 2);
- (void) mb;
-
- return CQregisterInternal(cntxt, modnme, fcnnme);
-}
-
static str
CQresumeInternalRanges(int first, int last)
{
@@ -897,25 +897,21 @@ CQderegisterInternalRanges(int first, in
return MAL_SUCCEED;
}
-//The force flag is set when deleting the CQ from the SQL catalog. In that
case the CQ may be registered in the Petrinet or not
-
-str
-CQderegisterInternal(str modnme, str fcnnme, int force)
+static str
+CQderegisterInternal(str modnme, str fcnnme)
{
int idx;
str msg = MAL_SUCCEED;
MT_lock_set(&ttrLock);
idx = CQlocate(modnme, fcnnme);
- if(!force && idx == pnettop) {
+ if(idx == pnettop) {
msg = createException(SQL, "cquery.deregister", "Continuous
procedure %s.%s not accessible\n", modnme, fcnnme);
goto finish;
}
if (idx <pnettop)
pnet[idx].status = CQSTOP;
MT_lock_unset(&ttrLock);
- if(idx == pnettop)
- goto finish;
// actually wait if the query was running
while( pnet[idx].status != CQDEREGISTER ){
@@ -928,6 +924,7 @@ finish:
MT_lock_unset(&ttrLock);
return msg;
}
+
str
CQderegisterAll(Client cntxt, MalBlkPtr mb, MalStkPtr stk, InstrPtr pci)
{
@@ -961,7 +958,7 @@ CQderegister(Client cntxt, MalBlkPtr mb,
}
}
if( k>= 0 )
- return CQderegisterInternal(getModuleId(getInstrPtr(mb,k)),
getFunctionId(getInstrPtr(mb,k)), 0);
+ return CQderegisterInternal(getModuleId(getInstrPtr(mb,k)),
getFunctionId(getInstrPtr(mb,k)));
throw(SQL,"cquery.stop","Continuous query not found ");
}
diff --git a/sql/backends/monet5/sql_cquery.h b/sql/backends/monet5/sql_cquery.h
--- a/sql/backends/monet5/sql_cquery.h
+++ b/sql/backends/monet5/sql_cquery.h
@@ -72,14 +72,12 @@ sql5_export MT_Lock ttrLock;
sql5_export str CQregister(Client cntxt, MalBlkPtr mb, MalStkPtr stk, InstrPtr
pci);
sql5_export str CQprocedure(Client cntxt, MalBlkPtr mb, MalStkPtr stk,
InstrPtr pci);
-sql5_export str CQregisterMAL(Client cntxt, MalBlkPtr mb, MalStkPtr stk,
InstrPtr pci );
sql5_export str CQresume(Client cntxt, MalBlkPtr mb, MalStkPtr stk, InstrPtr
pci);
sql5_export str CQresumeAll(Client cntxt, MalBlkPtr mb, MalStkPtr stk,
InstrPtr pci);
sql5_export str CQpause(Client cntxt, MalBlkPtr mb, MalStkPtr stk, InstrPtr
pci);
sql5_export str CQpauseAll(Client cntxt, MalBlkPtr mb, MalStkPtr stk, InstrPtr
pci);
sql5_export str CQderegister(Client cntxt, MalBlkPtr mb, MalStkPtr stk,
InstrPtr pci);
sql5_export str CQderegisterAll(Client cntxt, MalBlkPtr mb, MalStkPtr stk,
InstrPtr pci);
-sql5_export str CQderegisterInternal(str modnme, str fcnnme, int force);
sql5_export str CQwait(Client cntxt, MalBlkPtr mb, MalStkPtr stk, InstrPtr
pci);
sql5_export str CQcycles(Client cntxt, MalBlkPtr mb, MalStkPtr stk, InstrPtr
pci);
sql5_export str CQheartbeat(Client cntxt, MalBlkPtr mb, MalStkPtr stk,
InstrPtr pci);
diff --git a/sql/backends/monet5/sql_execute.c
b/sql/backends/monet5/sql_execute.c
--- a/sql/backends/monet5/sql_execute.c
+++ b/sql/backends/monet5/sql_execute.c
@@ -354,31 +354,26 @@ SQLrun(Client c, backend *be, mvc *m){
case mod_start_continuous:
//mnstr_printf(c->fdout, "#Start continuous
query\n");
// hand over the wrapper command to the
scheduler
- CQregisterMAL(c,mb, 0,0);
- m->continuous = 0;
- return MAL_SUCCEED;
+ msg = CQregister(c,mb, 0,0);
break;
case mod_stop_continuous:
//mnstr_printf(c->fdout, "#Stop continuous
query\n");
- CQderegister(c,mb, 0,0);
- m->continuous = 0;
- msg = MAL_SUCCEED;
+ msg = CQderegister(c,mb, 0,0);
break;
case mod_pause_continuous:
//mnstr_printf(c->fdout, "#Pause continuous
query\n");
- CQpause(c,mb, 0,0);
- m->continuous = 0;
- msg = MAL_SUCCEED;
+ msg = CQpause(c,mb, 0,0);
break;
case mod_resume_continuous:
//mnstr_printf(c->fdout, "#Resume continuous
query\n");
- CQresume(c,mb, 0,0);
- m->continuous = 0;
- msg = MAL_SUCCEED;
+ msg = CQresume(c,mb, 0,0);
break;
default:
msg = runMAL(c, mb, 0, 0);
}
+ m->continuous = 0;
+ m->heartbeats = 0;
+ m->cycles = 0;
}
}
diff --git a/sql/server/rel_psm.c b/sql/server/rel_psm.c
--- a/sql/server/rel_psm.c
+++ b/sql/server/rel_psm.c
@@ -618,10 +618,17 @@ sequential_block (mvc *sql, sql_subtype
case SQL_CALL:
res = rel_psm_call(sql, s->data.sym);
break;
- case SQL_START_CALL:
+ case SQL_START_CALL: {
+ dlist *l = s->data.lval;
sql->continuous = mod_start_continuous;
+ sql->heartbeats = l->h->next->data.i_val;
+ sql->cycles = l->h->next->next->data.i_val;
+ res = rel_psm_call(sql, l->h->data.sym);
+ } break;
+ case SQL_RESUME_CALL: {
+ sql->continuous = mod_resume_continuous;
res = rel_psm_call(sql, s->data.sym);
- break;
+ } break;
case SQL_STOP_CALL:
sql->continuous = mod_stop_continuous;
res = rel_psm_call(sql, s->data.sym);
@@ -630,10 +637,6 @@ sequential_block (mvc *sql, sql_subtype
sql->continuous = mod_pause_continuous;
res = rel_psm_call(sql, s->data.sym);
break;
- case SQL_RESUME_CALL:
- sql->continuous = mod_resume_continuous;
- res = rel_psm_call(sql, s->data.sym);
- break;
case SQL_RETURN:
/*If it is not a function it cannot have a return
statement*/
if (!is_func)
@@ -1431,11 +1434,19 @@ rel_psm(mvc *sql, symbol *s)
ret = rel_psm_stmt(sql->sa, rel_psm_call(sql, s->data.sym));
sql->type = Q_UPDATE;
break;
- case SQL_START_CALL:
+ case SQL_START_CALL: {
+ dlist *l = s->data.lval;
sql->continuous = mod_start_continuous;
+ sql->heartbeats = l->h->next->data.i_val;
+ sql->cycles = l->h->next->next->data.i_val;
+ ret = rel_psm_stmt(sql->sa, rel_psm_call(sql, l->h->data.sym));
+ sql->type = Q_UPDATE;
+ } break;
+ case SQL_RESUME_CALL: {
+ sql->continuous = mod_resume_continuous;
ret = rel_psm_stmt(sql->sa, rel_psm_call(sql, s->data.sym));
sql->type = Q_UPDATE;
- break;
+ } break;
case SQL_STOP_CALL:
sql->continuous = mod_stop_continuous;
ret = rel_psm_stmt(sql->sa, rel_psm_call(sql, s->data.sym));
@@ -1446,11 +1457,6 @@ rel_psm(mvc *sql, symbol *s)
ret = rel_psm_stmt(sql->sa, rel_psm_call(sql, s->data.sym));
sql->type = Q_UPDATE;
break;
- case SQL_RESUME_CALL:
_______________________________________________
checkin-list mailing list
[email protected]
https://www.monetdb.org/mailman/listinfo/checkin-list