Changeset: 76dc5cf78ca2 for MonetDB
URL: http://dev.monetdb.org/hg/MonetDB?cmd=changeset;node=76dc5cf78ca2
Added Files:
        sql/backends/monet5/Tests/cqcreate.sql
Modified Files:
        sql/backends/monet5/sql_cat.c
        sql/backends/monet5/sql_cquery.c
        sql/backends/monet5/sql_cquery.h
        sql/backends/monet5/sql_scenario.c
        sql/include/sql_catalog.h
        sql/scripts/15_querylog.sql
        sql/scripts/25_debug.sql
        sql/scripts/26_sysmon.sql
        sql/server/rel_psm.c
        sql/server/rel_schema.c
        sql/server/sql_parser.y
        sql/server/sql_scan.c
        sql/server/sql_scan.h
Branch: timetrails
Log Message:

First steps into adding continuous queries into the SQL catalog, but I am still 
getting errors while registering :( Can a SQL layer veteran help me?


diffs (truncated from 496 to 300 lines):

diff --git a/sql/backends/monet5/Tests/cqcreate.sql 
b/sql/backends/monet5/Tests/cqcreate.sql
new file mode 100644
--- /dev/null
+++ b/sql/backends/monet5/Tests/cqcreate.sql
@@ -0,0 +1,5 @@
+CREATE stream TABLE testing (a int);
+
+CREATE TABLE results (b int);
+
+CREATE CONTINUOUS QUERY stressing() BEGIN INSERT INTO results SELECT a FROM 
testing; END;
diff --git a/sql/backends/monet5/sql_cat.c b/sql/backends/monet5/sql_cat.c
--- a/sql/backends/monet5/sql_cat.c
+++ b/sql/backends/monet5/sql_cat.c
@@ -18,6 +18,7 @@
 #include "sql_scenario.h"
 #include "sql_mvc.h"
 #include "sql_qc.h"
+#include "sql_cquery.h"
 #include "sql_optimizer.h"
 #include "mal_namespace.h"
 #include "opt_prelude.h"
@@ -446,13 +447,28 @@ static str
 drop_func(mvc *sql, char *sname, char *name, int fid, int type, int action)
 {
        sql_schema *s = NULL;
-       char is_aggr = (type == F_AGGR);
-       char is_func = (type != F_PROC);
-       char *F = is_aggr ? "AGGREGATE" : (is_func ? "FUNCTION" : "PROCEDURE");
-       char *f = is_aggr ? "aggregate" : (is_func ? "function" : "procedure");
+       char *F, *f;
        char *KF = type == F_FILT ? "FILTER " : type == F_UNION ? "UNION " : "";
        char *kf = type == F_FILT ? "filter " : type == F_UNION ? "union " : "";
 
+       switch (type) {
+               case F_AGGR:
+                       F = "AGGREGATE";
+                       f = "aggregate";
+                       break;
+               case F_PROC:
+                       F = "PROCEDURE";
+                       f = "procedure";
+                       break;
+               case F_CONTINUOUS_QUERY:
+                       F = "CONTINUOUS QUERY";
+                       f = "continuous query";
+                       break;
+               default:
+                       F = "FUNCTION";
+                       f = "function";
+       }
+
        if (sname && !(s = mvc_bind_schema(sql, sname)))
                return sql_message("3F000!DROP %s%s: no such schema '%s'", KF, 
F, sname);
        if (!s)
@@ -497,12 +513,24 @@ create_func(mvc *sql, char *sname, char 
 {
        sql_func *nf;
        sql_schema *s = NULL;
-       char is_aggr = (f->type == F_AGGR);
-       char is_func = (f->type != F_PROC);
-       char *F = is_aggr ? "AGGREGATE" : (is_func ? "FUNCTION" : "PROCEDURE");
-       char *KF = f->type == F_FILT ? "FILTER " : f->type == F_UNION ? "UNION 
" : "";
+       char *F, *KF = f->type == F_FILT ? "FILTER " : f->type == F_UNION ? 
"UNION " : "";
 
-       (void)fname;
+       (void) fname;
+
+       switch (f->type) {
+               case F_AGGR:
+                       F = "AGGREGATE";
+                       break;
+               case F_PROC:
+                       F = "PROCEDURE";
+                       break;
+               case F_CONTINUOUS_QUERY:
+                       F = "CONTINUOUS QUERY";
+                       break;
+               default:
+                       F = "FUNCTION";
+       }
+
        if (sname && !(s = mvc_bind_schema(sql, sname)))
                return sql_message("3F000!CREATE %s%s: no such schema '%s'", 
KF, F, sname);
        if (!s)
@@ -546,6 +574,13 @@ create_func(mvc *sql, char *sname, char 
                if (!backend_resolve_function(sql, nf))
                        return sql_message("3F000!CREATE %s%s: external name 
%s.%s not bound", KF, F, nf->mod, nf->base.name);
        }
+       if(f->type == F_CONTINUOUS_QUERY) {
+               Client cntxt = MCgetClient(sql->clientid);
+               char *err = CQregisterInternal(cntxt, (str) sname, 
f->base.name);
+               if (err != NULL) {
+                       return sql_message("3F000!CREATE %s%s: continuous query 
register error: %s", KF, F, err);
+               }
+       }
        return MAL_SUCCEED;
 }
 
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
@@ -416,7 +416,7 @@ IOTprocedureStmt(Client cntxt, MalBlkPtr
     throw(SQL, "cquery.register", "SQL procedure missing");
 }
 
-str
+/*str
 CQregister(Client cntxt, MalBlkPtr mb, MalStkPtr stk, InstrPtr pci)
 {
        int i;
@@ -495,6 +495,94 @@ CQregister(Client cntxt, MalBlkPtr mb, M
        if( pnettop == 0)
                pnstatus = CQSTOP;
        return msg;
+}*/
+
+str
+CQregisterInternal(Client cntxt, str modnme, str fcnnme)
+{
+    int i;
+    InstrPtr sig,q;
+    str msg = MAL_SUCCEED;
+    MalBlkPtr mb, nmb;
+       Module scope;
+       Symbol s = NULL;
+    char buf[IDLENGTH];
+
+       scope = findModule(cntxt->nspace, putName(modnme));
+       if (scope)
+               s = findSymbolInModule(scope, putName(fcnnme));
+
+       if (s == NULL)
+               throw(MAL, "cquery.register", "Could not find SQL procedure");
+
+    if (pnettop == MAXCQ)
+        GDKerror("cquery.register:Too many transitions");
+
+    mb = s->def;
+    sig = getInstrPtr(mb,0);
+    i = CQlocate(getModuleId(sig), getFunctionId(sig));
+    if (i != pnettop)
+        throw(MAL,"cquery.register","Duplicate registration of cquery");
+
+#ifdef DEBUG_CQUERY
+    fprintf(stderr, "#cquery register %s.%s\n", 
getModuleId(sig),getFunctionId(sig));
+       fprintFunction(stderr,mb,0,LIST_MAL_ALL);
+#endif
+    memset((void*) (pnet+pnettop), 0, sizeof(CQnode));
+
+    snprintf(buf,IDLENGTH,"%s_%s",modnme,fcnnme);
+    s = newFunction(userRef, putName(buf), FUNCTIONsymbol);
+    nmb = s->def;
+    setArgType(nmb, nmb->stmt[0],0, TYPE_void);
+    (void) newStmt(nmb, sqlRef, transactionRef);
+    (void) newStmt(nmb, getModuleId(sig),getFunctionId(sig));
+    q = newStmt(nmb, sqlRef, commitRef);
+    setArgType(nmb,q, 0, TYPE_void);
+    pushEndInstruction(nmb);
+    chkProgram(cntxt->fdout, cntxt->nspace, nmb);
+#ifdef DEBUG_CQUERY
+    fprintFunction(stderr, nmb, 0, LIST_MAL_ALL);
+#endif
+
+    MT_lock_set(&ttrLock);
+    if( CQlocate(getModuleId(sig), getFunctionId(sig)) != pnettop){
+        freeSymbol(s);
+        throw(MAL,"cquery.register","Duplicate registration of cquery");
+    }
+    pnet[pnettop].mod = GDKstrdup(modnme);
+    pnet[pnettop].fcn = GDKstrdup(fcnnme);
+    pnet[pnettop].mb = nmb;
+    pnet[pnettop].stk = prepareMALstack(nmb, nmb->vsize);
+
+    pnet[pnettop].cycles = int_nil;
+    pnet[pnettop].beats = lng_nil;
+    pnet[pnettop].run  = lng_nil;
+    pnet[pnettop].seen = *timestamp_nil;
+    pnet[pnettop].status = CQPAUSE;
+    pnettop++;
+
+    msg = CQanalysis(cntxt, mb, pnettop-1);
+    MT_lock_unset(&ttrLock);
+    if( msg != MAL_SUCCEED)
+        // restore the entry
+        CQfree(pnettop);
+    if( pnettop == 0)
+        pnstatus = CQSTOP;
+    return msg;
+}
+
+str
+CQregister(Client cntxt, MalBlkPtr mb, MalStkPtr stk, InstrPtr pci)
+{
+    str msg = MAL_SUCCEED;
+    str modnme = *getArgReference_str(stk, pci, 1);
+    str fcnnme = *getArgReference_str(stk, pci, 2);
+
+    msg = IOTprocedureStmt(cntxt, mb, modnme, fcnnme);
+    if( msg)
+        return msg;
+
+    return CQregisterInternal(cntxt, modnme, fcnnme);
 }
 
 str
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
@@ -64,6 +64,7 @@ sql5_export CQnode pnet[MAXCQ];
 sql5_export int pnettop;
 sql5_export MT_Lock ttrLock;
 
+sql5_export str CQregisterInternal(Client cntxt, str modnme, str fcnnme);
 sql5_export str CQregister(Client cntxt, MalBlkPtr mb, MalStkPtr stk, InstrPtr 
pci);
 sql5_export str CQresume(Client cntxt, MalBlkPtr mb, MalStkPtr stk, InstrPtr 
pci);
 sql5_export str CQpause(Client cntxt, MalBlkPtr mb, MalStkPtr stk, InstrPtr 
pci);
diff --git a/sql/backends/monet5/sql_scenario.c 
b/sql/backends/monet5/sql_scenario.c
--- a/sql/backends/monet5/sql_scenario.c
+++ b/sql/backends/monet5/sql_scenario.c
@@ -1113,7 +1113,7 @@ SQLparser(Client c)
                }
 
                if ((!caching(m) || !cachable(m, r)) && m->emode != m_prepare) {
-                       char *q = query_cleaned(QUERY(m->scanner));
+                       char *q = query_cleaned(MQUERY(m->scanner));
 
                        /* Query template should not be cached */
                        scanner_query_processed(&(m->scanner));
@@ -1128,7 +1128,7 @@ SQLparser(Client c)
                        /* Add the query tree to the SQL query cache
                         * and bake a MAL program for it.
                         */
-                       char *q = query_cleaned(QUERY(m->scanner));
+                       char *q = query_cleaned(MQUERY(m->scanner));
                        char qname[IDLENGTH];
                        (void) snprintf(qname, IDLENGTH, "%c%d_%d", (m->emode 
== m_prepare?'p':'s'), m->qc->id++, m->qc->clientid);
 
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
@@ -282,14 +282,16 @@ typedef struct sql_arg {
 #define F_UNION 5
 #define F_ANALYTIC 6
 #define F_LOADER 7
+#define F_CONTINUOUS_QUERY 8
 
 #define IS_FUNC(f) (f->type == F_FUNC)
-#define IS_PROC(f) (f->type == F_PROC)
+#define IS_PROC(f) (f->type == F_PROC || f->type == F_CONTINUOUS_QUERY) /* To 
be safe */
 #define IS_AGGR(f) (f->type == F_AGGR)
 #define IS_FILT(f) (f->type == F_FILT)
 #define IS_UNION(f) (f->type == F_UNION)
 #define IS_ANALYTIC(f) (f->type == F_ANALYTIC)
 #define IS_LOADER(f) (f->type == F_LOADER)
+#define IS_CONTINUOUS_QUERY(f) (f->type == F_CONTINUOUS_QUERY)
 
 #define FUNC_LANG_INT 0        /* internal */
 #define FUNC_LANG_MAL 1 /* create sql external mod.func */
diff --git a/sql/scripts/15_querylog.sql b/sql/scripts/15_querylog.sql
--- a/sql/scripts/15_querylog.sql
+++ b/sql/scripts/15_querylog.sql
@@ -14,7 +14,7 @@ returns table(
        id oid,
        owner string,
        defined timestamp,
-       query string,
+       "query" string,
        pipe string,
        "plan" string,          -- Name of MAL plan
        mal int,                -- size of MAL plan
diff --git a/sql/scripts/25_debug.sql b/sql/scripts/25_debug.sql
--- a/sql/scripts/25_debug.sql
+++ b/sql/scripts/25_debug.sql
@@ -14,7 +14,7 @@ create function sys.optimizer_stats ()
 -- The SQL query cache returns a table with the query plans kept
 
 create function sys.queryCache()
-       returns table (query string, count int)
+       returns table ("query" string, count int)
        external name sql.dump_cache;
 
 -- Trace the SQL input
diff --git a/sql/scripts/26_sysmon.sql b/sql/scripts/26_sysmon.sql
--- a/sql/scripts/26_sysmon.sql
+++ b/sql/scripts/26_sysmon.sql
@@ -16,7 +16,7 @@ returns table(
        progress int,
        status string,
        tag oid,
-       query string
+       "query" string
 )
 external name sql.sysmon_queue;
 
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
@@ -763,16 +763,35 @@ rel_create_func(mvc *sql, dlist *qname, 
        bit vararg = FALSE;
 
_______________________________________________
checkin-list mailing list
[email protected]
https://www.monetdb.org/mailman/listinfo/checkin-list

Reply via email to