Changeset: b1076e8f33c9 for MonetDB
URL: https://dev.monetdb.org/hg/MonetDB?cmd=changeset;node=b1076e8f33c9
Added Files:
        sql/backends/monet5/Tests/cquery20.stable.err
        sql/backends/monet5/Tests/cquery20.stable.out
Modified Files:
        sql/backends/monet5/Tests/cquery20.sql
        sql/backends/monet5/cquery.mal
        sql/backends/monet5/sql_cquery.c
        sql/backends/monet5/sql_cquery.h
Branch: trails
Log Message:

Wait for concurrent actions
before you continue, because it may lead to objects freed
while being in use.


diffs (truncated from 368 to 300 lines):

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
@@ -16,8 +16,6 @@ call cqfoo(123);
 select * from cqtbl;
 
 start continuous cqfoo(321);
-select * from cqtbl;
-
 stop continuous cqfoo(321);
 --stop continuous cqfoo;
 select * from cqtbl;
diff --git a/sql/backends/monet5/Tests/cquery20.stable.err 
b/sql/backends/monet5/Tests/cquery20.stable.err
new file mode 100644
--- /dev/null
+++ b/sql/backends/monet5/Tests/cquery20.stable.err
@@ -0,0 +1,34 @@
+stderr of test 'cquery20` in directory 'sql/backends/monet5` itself:
+
+
+# 23:19:50 >  
+# 23:19:50 >  "mserver5" "--debug=10" "--set" "gdk_nr_threads=0" "--set" 
"mapi_open=true" "--set" "mapi_port=34472" "--set" 
"mapi_usock=/var/tmp/mtest-6809/.s.monetdb.34472" "--set" "monet_prompt=" 
"--forcemito" 
"--dbpath=/export/scratch1/mk/trails//Linux/var/MonetDB/mTests_sql_backends_monet5"
+# 23:19:50 >  
+
+# builtin opt  gdk_dbpath = 
/export/scratch1/mk/trails//Linux/var/monetdb5/dbfarm/demo
+# builtin opt  gdk_debug = 0
+# builtin opt  gdk_vmtrim = no
+# builtin opt  monet_prompt = >
+# builtin opt  monet_daemon = no
+# builtin opt  mapi_port = 50000
+# builtin opt  mapi_open = false
+# builtin opt  mapi_autosense = false
+# builtin opt  sql_optimizer = default_pipe
+# builtin opt  sql_debug = 0
+# cmdline opt  gdk_nr_threads = 0
+# cmdline opt  mapi_open = true
+# cmdline opt  mapi_port = 34472
+# cmdline opt  mapi_usock = /var/tmp/mtest-6809/.s.monetdb.34472
+# cmdline opt  monet_prompt = 
+# cmdline opt  gdk_dbpath = 
/export/scratch1/mk/trails//Linux/var/MonetDB/mTests_sql_backends_monet5
+# cmdline opt  gdk_debug = 536870922
+
+# 23:19:50 >  
+# 23:19:50 >  "mclient" "-lsql" "-ftest" "-Eutf-8" "-i" "-e" 
"--host=/var/tmp/mtest-6809" "--port=34472"
+# 23:19:50 >  
+
+
+# 23:19:50 >  
+# 23:19:50 >  "Done."
+# 23:19:50 >  
+
diff --git a/sql/backends/monet5/Tests/cquery20.stable.out 
b/sql/backends/monet5/Tests/cquery20.stable.out
new file mode 100644
--- /dev/null
+++ b/sql/backends/monet5/Tests/cquery20.stable.out
@@ -0,0 +1,58 @@
+stdout of test 'cquery20` in directory 'sql/backends/monet5` itself:
+
+
+# 23:19:50 >  
+# 23:19:50 >  "mserver5" "--debug=10" "--set" "gdk_nr_threads=0" "--set" 
"mapi_open=true" "--set" "mapi_port=34472" "--set" 
"mapi_usock=/var/tmp/mtest-6809/.s.monetdb.34472" "--set" "monet_prompt=" 
"--forcemito" 
"--dbpath=/export/scratch1/mk/trails//Linux/var/MonetDB/mTests_sql_backends_monet5"
+# 23:19:50 >  
+
+# MonetDB 5 server v11.28.0
+# This is an unreleased version
+# Serving database 'mTests_sql_backends_monet5', using 8 threads
+# Compiled for x86_64-unknown-linux-gnu/64bit with 128bit integers
+# Found 15.588 GiB available main-memory.
+# Copyright (c) 1993-July 2008 CWI.
+# Copyright (c) August 2008-2017 MonetDB B.V., all rights reserved
+# Visit https://www.monetdb.org/ for further information
+# Listening for connection requests on mapi:monetdb://vienna.da.cwi.nl:34472/
+# Listening for UNIX domain connection requests on 
mapi:monetdb:///var/tmp/mtest-6809/.s.monetdb.34472
+# MonetDB/GIS module loaded
+# MonetDB/SQL module loaded
+# MonetDB/Timetrails module loaded
+
+Ready.
+
+# 23:19:50 >  
+# 23:19:50 >  "mclient" "-lsql" "-ftest" "-Eutf-8" "-i" "-e" 
"--host=/var/tmp/mtest-6809" "--port=34472"
+# 23:19:50 >  
+
+#create stream table cqtbl(i integer);
+#create continuous procedure cqfoo(v integer)
+#begin
+#      insert into cqtbl values(v);
+#end;
+#select * from functions where name = 'cqfoo';
+% sys.functions,       sys.functions,  sys.functions,  sys.functions,  
sys.functions,  sys.functions,  sys.functions,  sys.functions,  sys.functions,  
sys.functions # table_name
+% id,  name,   func,   mod,    language,       type,   side_effect,    varres, 
vararg, schema_id # name
+% int, varchar,        varchar,        varchar,        int,    int,    
boolean,        boolean,        boolean,        int # type
+% 4,   5,      85,     4,      1,      1,      5,      5,      5,      4 # 
length
+[ 8434,        "cqfoo",        "create continuous procedure cqfoo(v 
integer)\nbegin\n insert into cqtbl values(v);\nend;",     "user", 2,      2,   
   true,   false,  false,  2000    ]
+#select * from cqtbl;
+% sys.cqtbl # table_name
+% i # name
+% int # type
+% 3 # length
+[ 123  ]
+#select * from cqtbl;
+% sys.cqtbl # table_name
+% i # name
+% int # type
+% 3 # length
+[ 123  ]
+[ 321  ]
+#drop procedure cqfoo;
+#drop table cqtbl;
+
+# 23:19:50 >  
+# 23:19:50 >  "Done."
+# 23:19:50 >  
+
diff --git a/sql/backends/monet5/cquery.mal b/sql/backends/monet5/cquery.mal
--- a/sql/backends/monet5/cquery.mal
+++ b/sql/backends/monet5/cquery.mal
@@ -36,21 +36,21 @@ pattern resume(mod:str, fcn:str)
 address CQresume
 comment "Activate a specific continuous query";
 pattern resume()
-address CQresume
+address CQresumeAll
 comment "Activate all continuous queries";
 
 pattern pause(mod:str, fcn:str)
 address CQpause
 comment "Deactivate a continuous query";
 pattern pause()
-address CQpause
+address CQpauseAll
 comment "Deactivate all continuous queries";
 
 pattern deregister(mod:str, fcn:str)
 address CQderegister
 comment "Remove a continuous query";
 pattern deregister()
-address CQderegister
+address CQderegisterAll
 comment "Remove all continuous queries";
 
 pattern wait(cnt:int)
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
@@ -671,9 +671,22 @@ finish:
 }
 
 str
+CQresumeAll(Client cntxt, MalBlkPtr mb, MalStkPtr stk, InstrPtr pci)
+{
+       str msg;
+       (void) cntxt;
+       (void) mb;
+       (void) stk;
+       (void) pci;
+       MT_lock_set(&ttrLock);
+       msg = CQresumeInternalRanges(0, pnettop);
+       MT_lock_unset(&ttrLock);
+       return msg;
+}
+
+str
 CQresume(Client cntxt, MalBlkPtr mb, MalStkPtr stk, InstrPtr pci)
 {
-       str msg = MAL_SUCCEED;
        int i, k =-1;
        InstrPtr q;
        (void) cntxt;
@@ -683,19 +696,16 @@ CQresume(Client cntxt, MalBlkPtr mb, Mal
        for( i=1; i < mb->stop; i++){
                q = getInstrPtr(mb,i);
 
+               if( q->token == ENDsymbol )
+                       break;
                if( getModuleId(q) == userRef){
-                       if( k != -1)
-                               throw(SQL,"cquery.resume","Ambiguous resume 
statement");
                        k = i;
+                       break;
                }
        }
-       if( k >= 0)
+       if( k >= 0 )
                return CQresumeInternal(getModuleId(getInstrPtr(mb,k)), 
getFunctionId(getInstrPtr(mb,k)));
-       //resume all
-       MT_lock_set(&ttrLock);
-       msg = CQresumeInternalRanges(0, pnettop);
-       MT_lock_unset(&ttrLock);
-       return msg;
+       throw(SQL,"cquery.resume","Continuous query not found ");
 }
 
 static str
@@ -724,7 +734,7 @@ CQpauseInternal(str modnme, str fcnnme)
        // actually wait if the query was running
        while( pnet[idx].status == CQRUNNING ){
                MT_lock_unset(&ttrLock);
-               MT_sleep_ms(10);  
+               MT_sleep_ms(5);  
                MT_lock_set(&ttrLock);
                if( pnet[idx].status == CQWAIT)
                        break;
@@ -736,9 +746,22 @@ finish:
 }
 
 str
+CQpauseAll(Client cntxt, MalBlkPtr mb, MalStkPtr stk, InstrPtr pci)
+{
+       str msg;
+       (void) cntxt;
+       (void) mb;
+       (void) stk;
+       (void) pci;
+       //pause all
+       MT_lock_set(&ttrLock);
+       msg = CQpauseInternalRanges(0, pnettop);
+       MT_lock_unset(&ttrLock);
+       return msg;
+}
+str
 CQpause(Client cntxt, MalBlkPtr mb, MalStkPtr stk, InstrPtr pci)
 {
-       str msg = MAL_SUCCEED;
        int i,k = -1;
        InstrPtr q;
        (void) cntxt;
@@ -747,20 +770,16 @@ CQpause(Client cntxt, MalBlkPtr mb, MalS
 
        for( i=1; i < mb->stop; i++){
                q = getInstrPtr(mb,i);
+               if( q->token == ENDsymbol )
+                       break;
                if( getModuleId(q) == userRef){
-                       if( k != -1)
-                               throw(SQL,"cquery.pause","Ambiguous pause 
statement");
                        k = i;
+                       break;
                }
        }
        if( k >= 0)
                return CQpauseInternal(getModuleId(getInstrPtr(mb,k)), 
getFunctionId(getInstrPtr(mb,k)));
-
-       //pause all
-       MT_lock_set(&ttrLock);
-       msg = CQpauseInternalRanges(0, pnettop);
-       MT_lock_unset(&ttrLock);
-       return msg;
+       throw(SQL,"cquery.pause","Continuous query not found ");
 }
 
 str
@@ -892,28 +911,40 @@ CQderegisterInternal(str modnme, str fcn
                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 == CQRUNNING ){
-               MT_lock_unset(&ttrLock);
-               MT_sleep_ms(10);  
-               MT_lock_set(&ttrLock);
-               if( pnet[idx].status == CQWAIT)
-                       break;
+       while( pnet[idx].status != CQDEREGISTER ){
+               MT_sleep_ms(5);
        }
+       MT_lock_set(&ttrLock);
        msg = CQderegisterInternalRanges(idx, idx+1);
 
 finish:
        MT_lock_unset(&ttrLock);
        return msg;
 }
+str
+CQderegisterAll(Client cntxt, MalBlkPtr mb, MalStkPtr stk, InstrPtr pci)
+{
+       str msg;
+       (void) cntxt;
+       (void) mb;
+       (void) stk;
+       (void) pci;
+       MT_lock_set(&ttrLock);
+       msg = CQderegisterInternalRanges(0, pnettop);
+       MT_lock_unset(&ttrLock);
+       return msg;
+}
 
 str
 CQderegister(Client cntxt, MalBlkPtr mb, MalStkPtr stk, InstrPtr pci)
 {
-       str msg = MAL_SUCCEED;
        int i, k= -1;
        InstrPtr q;
        (void) cntxt;
_______________________________________________
checkin-list mailing list
[email protected]
https://www.monetdb.org/mailman/listinfo/checkin-list

Reply via email to