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

Reply via email to