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