Changeset: 8dcc85a4a027 for MonetDB
URL: https://dev.monetdb.org/hg/MonetDB?cmd=changeset;node=8dcc85a4a027
Modified Files:
        sql/backends/monet5/Tests/cquery20.py
        sql/backends/monet5/sql_cquery.c
        sql/include/sql_catalog.h
        sql/server/rel_schema.c
        sql/server/sql_parser.y
        sql/storage/sql_storage.h
        sql/storage/store.c
Branch: trails
Log Message:

We want local temp stream tables instead of global ones.


diffs (289 lines):

diff --git a/sql/backends/monet5/Tests/cquery20.py 
b/sql/backends/monet5/Tests/cquery20.py
--- a/sql/backends/monet5/Tests/cquery20.py
+++ b/sql/backends/monet5/Tests/cquery20.py
@@ -17,32 +17,29 @@ def client(input):
     sys.stderr.write(err)
 
 script1 = '''\
-create temporary stream table sta (a int);
+create temporary stream table sta (a int);\
+select count(*) from streams;\
+create stream table stb (a int);\
+select count(*) from streams;
 '''
 
 script2 = '''\
-create stream table stb (a int);
+select count(*) from streams;\
 '''
 
 script3 = '''\
+select count(*) from streams;\
+drop table stb;\
 select count(*) from streams;
 '''
 
-script4 = '''\
-drop table stb;
-'''
-
 def main():
     s = process.server(args = [], stdin = process.PIPE, stdout = process.PIPE, 
stderr = process.PIPE)
     client(script1)
-    client(script3)
     client(script2)
-    client(script3)
     server_stop(s)
     s = process.server(args = [], stdin = process.PIPE, stdout = process.PIPE, 
stderr = process.PIPE)
     client(script3)
-    client(script4)
-    client(script3)
     server_stop(s)
 
 if __name__ == '__main__':
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
@@ -69,24 +69,45 @@ static BAT *CQ_id_error = 0;
 static CQnode *pnet = 0;
 static int pnetLimit = 0, pnettop = 0;
 
-#define CQ_SQL_QUERY_SIZE  1024
+#define CQ_SCHEDULER_CLIENTID     0
+#define CQ_SQL_QUERY_SIZE      1024
 
 #define SET_HEARTBEATS(X) (X != HEARTBEAT_NIL) ? X : HEARTBEAT_NIL /* minimal 
1 ms */
 
-#define ALL_ROOT_CHECK(cntxt, malcal, name)                                    
                                        \
-       do {                                                                    
                                           \
-               smvc = ((backend *) cntxt->sqlcontext)->mvc;                    
                                               \
-               if(!smvc)                                                       
                                               \
-                       throw(SQL,malcal,SQLSTATE(42000) "##name##ALL 
CONTINUOUS: SQL clients only\n");                            \
-                else if (smvc->user_id != USER_MONETDB && smvc->role_id != 
ROLE_SYSADMIN)                                     \
-                       throw(SQL,malcal,SQLSTATE(42000) "##name##ALL 
CONTINUOUS: insufficient privileges for the current user\n");\
-       } while(0);
+#define ALL_ROOT_CHECK(cntxt, malcal, name)                                    
                                    \
+do {                                                                           
                                    \
+       smvc = ((backend *) cntxt->sqlcontext)->mvc;                            
                                       \
+       if(!smvc)                                                               
                                       \
+               throw(SQL,malcal,SQLSTATE(42000) "##name##ALL CONTINUOUS: SQL 
clients only\n");                            \
+        else if (smvc->user_id != USER_MONETDB && smvc->role_id != 
ROLE_SYSADMIN)                                     \
+               throw(SQL,malcal,SQLSTATE(42000) "##name##ALL CONTINUOUS: 
insufficient privileges for the current user\n");\
+} while(0);
+
+#define CLEAN_BASKETS(idx)                                                     
               \
+do {                                                                           
               \
+       for(int m=0; m< MAXSTREAMS && pnet[idx].baskets[m]; m++) {              
                  \
+               int found = 0;                                                  
                      \
+               str sch = baskets[pnet[idx].baskets[m]].table->s->base.name;    
                      \
+               str tbl = baskets[pnet[idx].baskets[m]].table->base.name;       
                      \
+               for(int n=0; n < pnettop && !found; n++) {                      
                      \
+                       if(n != idx) {                                          
                          \
+                               for(int o=0; o< MAXSTREAMS && 
pnet[n].baskets[o] && !found; o++) {            \
+                                       if (strcmp(sch, 
baskets[pnet[n].baskets[o]].table->s->base.name) == 0 &&  \
+                                               strcmp(tbl, 
baskets[pnet[n].baskets[o]].table->base.name) == 0)       \
+                                               found = 1;                      
                                      \
+                               }                                               
                              \
+                       }                                                       
                          \
+               }                                                               
                      \
+               if(!found) {                                                    
                      \
+                       BSKTclean(pnet[idx].baskets[m]);                        
                          \
+               }                                                               
                      \
+       }                                                                       
                  \
+} while(0);
 
 static void
 CQfree(int idx)
 {
-       int i, j, k, found;
-       str sch, tbl;
+       int i;
        InstrPtr p;
 
        if( pnet[idx].mb) {
@@ -99,24 +120,8 @@ CQfree(int idx)
        if(pnet[idx].error)
                GDKfree(pnet[idx].error);
        GDKfree(pnet[idx].alias);
-       //try delete the baskets
-       for( j=0; j< MAXSTREAMS && pnet[idx].baskets[j]; j++) {
-               found = 0;
-               sch = baskets[pnet[idx].baskets[j]].table->s->base.name;
-               tbl = baskets[pnet[idx].baskets[j]].table->base.name;
-               for( i=0; i < pnettop && !found; i++) {
-                       if(i != idx) {
-                               for( k=0; k< MAXSTREAMS && pnet[i].baskets[k] 
&& !found; k++) {
-                                       if (strcmp(sch, 
baskets[pnet[i].baskets[k]].table->s->base.name) == 0 &&
-                                               strcmp(tbl, 
baskets[pnet[i].baskets[k]].table->base.name) == 0)
-                                               found = 1;
-                               }
-                       }
-               }
-               if(!found) {
-                       BSKTclean(pnet[idx].baskets[j]);
-               }
-       }
+       //clean the baskets if so
+       CLEAN_BASKETS(idx)
        // compact the pnet table
        for(i=idx; i<pnettop-1; i++)
                pnet[i] = pnet[i+1];
@@ -438,7 +443,7 @@ CQregister(Client cntxt, str sname, str 
 
        //Find the UDF
        f = sql_find_func(m->sa, s, fname, argc > 0 ? argc : -1, (which & 
mod_continuous_function) ? F_FUNC : F_PROC, NULL);
-       if(!f && which & mod_continuous_function) { //If an UDF returns a 
table, it gets compiled into a F_UNION
+       if(!f && (which & mod_continuous_function)) { //If an UDF returns a 
table, it gets compiled into a F_UNION
                f = sql_find_func(m->sa, s, fname, argc > 0 ? argc : -1, 
F_UNION, NULL);
        }
        if(!f) {
@@ -597,6 +602,7 @@ CQregister(Client cntxt, str sname, str 
                FREE_CQ_MB(unlock)
        }
        if((msg = CQanalysis(cntxt, sym->def, pnettop)) != MAL_SUCCEED) {
+               CLEAN_BASKETS(pnettop)
                FREE_CQ_MB(unlock)
        }
        if(heartbeats != HEARTBEAT_NIL) {
@@ -604,12 +610,14 @@ CQregister(Client cntxt, str sname, str 
                        if(baskets[pnet[pnettop].baskets[i]].window != 
DEFAULT_TABLE_WINDOW) {
                                msg = createException(SQL, "cquery.register",
                                                                          
SQLSTATE(42000) "Heartbeat ignored, a window constraint exists\n");
+                               CLEAN_BASKETS(pnettop)
                                FREE_CQ_MB(unlock)
                        }
                }
        }
 
        if((pnet[pnettop].stk = prepareMALstack(mb, mb->vsize)) == NULL) { 
//prepare MAL stack
+               CLEAN_BASKETS(pnettop)
                CQ_MALLOC_FAIL(unlock)
        }
 
@@ -1326,7 +1334,7 @@ CQstartScheduler(void)
                throw(MAL, "cquery.startScheduler",SQLSTATE(HY001) "Could not 
initialize CQscheduler\n");
        }
 
-       cntxt = MCinitClient(0,bin,fout);
+       cntxt = MCinitClient(CQ_SCHEDULER_CLIENTID,bin,fout);
        if( cntxt == NULL) {
                bstream_destroy(cntxt->fdin);
                mnstr_destroy(cntxt->fdout);
diff --git a/sql/include/sql_catalog.h b/sql/include/sql_catalog.h
--- a/sql/include/sql_catalog.h
+++ b/sql/include/sql_catalog.h
@@ -129,7 +129,7 @@ typedef enum temp_t {
        SQL_DECLARED_TABLE = 3, /* variable inside a stored procedure */
        SQL_MERGE_TABLE = 4,
        SQL_PERSISTED_STREAM = 5,
-       SQL_GLOBAL_TEMP_STREAM = 6,
+       SQL_LOCAL_TEMP_STREAM = 6,
        SQL_REMOTE = 7,
        SQL_REPLICA_TABLE = 8
 } temp_t;
diff --git a/sql/server/rel_schema.c b/sql/server/rel_schema.c
--- a/sql/server/rel_schema.c
+++ b/sql/server/rel_schema.c
@@ -188,7 +188,7 @@ mvc_create_table_as_subquery( mvc *sql, 
        sql_table *t;
        int tt;
 
-       if(temp == SQL_PERSISTED_STREAM || temp == SQL_GLOBAL_TEMP_STREAM) {
+       if(temp == SQL_PERSISTED_STREAM || temp == SQL_LOCAL_TEMP_STREAM) {
                tt =(temp == SQL_PERSISTED_STREAM)?tt_stream_per:tt_stream_temp;
                t = mvc_create_stream_table(sql, s, tname, tt, 0, 
SQL_DECLARED_TABLE, commit_action, -1, window_size, stride);
        } else {
@@ -916,7 +916,7 @@ rel_create_table(mvc *sql, sql_schema *s
        int create = (!instantiate && !deps);
        int tt = (temp == SQL_REMOTE)?tt_remote:
                 (temp == SQL_PERSISTED_STREAM)?tt_stream_per:
-                (temp == SQL_GLOBAL_TEMP_STREAM)?tt_stream_temp:
+                (temp == SQL_LOCAL_TEMP_STREAM)?tt_stream_temp:
                 (temp == SQL_MERGE_TABLE)?tt_merge_table:
                 (temp == SQL_REPLICA_TABLE)?tt_replica_table:tt_table;
        int window_size = DEFAULT_TABLE_WINDOW, stride = DEFAULT_TABLE_STRIDE;
@@ -932,7 +932,7 @@ rel_create_table(mvc *sql, sql_schema *s
        if (temp != SQL_DECLARED_TABLE) {
                if (temp != SQL_PERSIST && temp != SQL_PERSISTED_STREAM && (tt 
== tt_table || tt == tt_stream_temp)) {
                        s = mvc_bind_schema(sql, "tmp");
-                       if (temp == SQL_LOCAL_TEMP && sname && strcmp(sname, 
s->base.name) != 0)
+                       if ((temp == SQL_LOCAL_TEMP || temp == 
SQL_LOCAL_TEMP_STREAM) && sname && strcmp(sname, s->base.name) != 0)
                                return sql_error(sql, 02, SQLSTATE(3F000) 
"CREATE TABLE: local temporary tables should be stored in the '%s' schema", 
s->base.name);
                } else if (s == NULL) {
                        s = ss;
diff --git a/sql/server/sql_parser.y b/sql/server/sql_parser.y
--- a/sql/server/sql_parser.y
+++ b/sql/server/sql_parser.y
@@ -1360,12 +1360,12 @@ table_opt_storage:
 
 opt_temp_stream:
     /* empty */      { $$ = SQL_PERSISTED_STREAM; }
- |  TEMP             { $$ = SQL_GLOBAL_TEMP_STREAM; }
- |  TEMPORARY        { $$ = SQL_GLOBAL_TEMP_STREAM; }
- |  LOCAL TEMPORARY  { $$ = yyerror(m, "LOCAL TEMPORARY STREAM tables not 
supported, only GLOBAL"); $$ = SQL_GLOBAL_TEMP_STREAM; }
- |  LOCAL TEMP       { $$ = yyerror(m, "LOCAL TEMPORARY STREAM tables not 
supported, only GLOBAL"); $$ = SQL_GLOBAL_TEMP_STREAM; }
- |  GLOBAL TEMPORARY { $$ = SQL_GLOBAL_TEMP_STREAM; }
- |  GLOBAL TEMP      { $$ = SQL_GLOBAL_TEMP_STREAM; }
+ |  TEMP             { $$ = SQL_LOCAL_TEMP_STREAM; }
+ |  TEMPORARY        { $$ = SQL_LOCAL_TEMP_STREAM; }
+ |  LOCAL TEMPORARY  { $$ = SQL_LOCAL_TEMP_STREAM; }
+ |  LOCAL TEMP       { $$ = SQL_LOCAL_TEMP_STREAM; }
+ |  GLOBAL TEMPORARY { $$ = yyerror(m, "GLOBAL TEMPORARY STREAM tables not 
supported, only LOCAL"); }
+ |  GLOBAL TEMP      { $$ = yyerror(m, "GLOBAL TEMPORARY STREAM tables not 
supported, only LOCAL"); }
  ;
 
 stream_window_set:
diff --git a/sql/storage/sql_storage.h b/sql/storage/sql_storage.h
--- a/sql/storage/sql_storage.h
+++ b/sql/storage/sql_storage.h
@@ -20,9 +20,9 @@
 #define isNew(x)  ((x)->base.flag == TR_NEW)
 #define isTemp(x) (isNew((x)->t)||((x)->t->persistence!=SQL_PERSIST && 
(x)->t->persistence!=SQL_PERSISTED_STREAM))
 #define isTempTable(x)   ((x)->persistence!=SQL_PERSIST && 
(x)->persistence!=SQL_PERSISTED_STREAM)
-#define isGlobal(x)      ((x)->persistence!=SQL_LOCAL_TEMP && \
+#define isGlobal(x)      ((x)->persistence!=SQL_LOCAL_TEMP && 
(x)->persistence!=SQL_LOCAL_TEMP_STREAM && \
                          (x)->persistence!=SQL_DECLARED_TABLE)
-#define isGlobalTemp(x)  ((x)->persistence==SQL_GLOBAL_TEMP || 
(x)->persistence==SQL_GLOBAL_TEMP_STREAM)
+#define isGlobalTemp(x)  ((x)->persistence==SQL_GLOBAL_TEMP)
 #define isTempSchema(x)  (strcmp((x)->base.name, "tmp") == 0 || \
                          strcmp((x)->base.name, dt_schema) == 0)
 #define isDeclaredTable(x)  ((x)->persistence==SQL_DECLARED_TABLE)
diff --git a/sql/storage/store.c b/sql/storage/store.c
--- a/sql/storage/store.c
+++ b/sql/storage/store.c
@@ -187,7 +187,7 @@ trans_drop_tmp(sql_trans *tr)
                        node *nxt = n->next;
                        sql_table *t = n->data;
 
-                       if (t->persistence == SQL_LOCAL_TEMP)
+                       if (t->persistence == SQL_LOCAL_TEMP || t->persistence 
== SQL_LOCAL_TEMP_STREAM)
                                cs_remove_node(&tmp->tables, n);
                        n = nxt;
                }
@@ -633,7 +633,7 @@ load_table(sql_trans *tr, sql_schema *s,
        if (isPerStream(t))
                t->persistence = SQL_PERSISTED_STREAM;
        if (isTempStream(t))
-               t->persistence = SQL_GLOBAL_TEMP_STREAM;
+               t->persistence = SQL_LOCAL_TEMP_STREAM;
        t->cleared = 0;
        v = table_funcs.column_find_value(tr, find_sql_column(tables, 
"access"),rid);
        t->access = *(sht*)v;   _DELETE(v);
@@ -2589,7 +2589,7 @@ schema_dup(sql_trans *tr, int flag, sql_
                for (n = os->tables.set->h; n; n = n->next) {
                        sql_table *ot = n->data;
 
-                       if (ot->persistence != SQL_LOCAL_TEMP)
+                       if (ot->persistence != SQL_LOCAL_TEMP && 
ot->persistence != SQL_LOCAL_TEMP_STREAM)
                                cs_add(&s->tables, table_dup(tr, flag, ot, s), 
tr_flag(&ot->base, flag));
                }
                if (tr->parent == gtrans)
@@ -4391,7 +4391,7 @@ sql_trans_create_table(sql_trans *tr, sq
        if (isPerStream(t))
                t->persistence = SQL_PERSISTED_STREAM;
        if (isTempStream(t))
-               t->persistence = SQL_GLOBAL_TEMP_STREAM;
+               t->persistence = SQL_LOCAL_TEMP_STREAM;
 
        if (isTable(t)) {
                if (store_funcs.create_del(tr, t) != LOG_OK) {
_______________________________________________
checkin-list mailing list
[email protected]
https://www.monetdb.org/mailman/listinfo/checkin-list

Reply via email to