Changeset: 9329e4d43cc9 for MonetDB
URL: https://dev.monetdb.org/hg/MonetDB/rev/9329e4d43cc9
Modified Files:
monetdb5/mal/mal.c
monetdb5/mal/mal_client.c
monetdb5/mal/mal_dataflow.c
monetdb5/mal/mal_embedded.c
monetdb5/mal/mal_import.c
monetdb5/mal/mal_internal.h
monetdb5/mal/mal_interpreter.c
monetdb5/mal/mal_private.h
monetdb5/mal/mal_session.c
monetdb5/modules/mal/tablet.c
Branch: default
Log Message:
Apart from dataflow, threads work for a single context.
Assign the thread-local variable early on and keep it assigned
throughout, except for dataflow which does reassign the context. Also,
size counting only for threads that run for clients, not during
initialization.
diffs (truncated from 502 to 300 lines):
diff --git a/monetdb5/mal/mal.c b/monetdb5/mal/mal.c
--- a/monetdb5/mal/mal.c
+++ b/monetdb5/mal/mal.c
@@ -30,6 +30,7 @@ lng MALdebug;
#include "msabaoth.h"
#include "mal_dataflow.h"
#include "mal_private.h"
+#include "mal_internal.h"
#include "mal_runtime.h"
#include "mal_resource.h"
#include "mal_atom.h"
diff --git a/monetdb5/mal/mal_client.c b/monetdb5/mal/mal_client.c
--- a/monetdb5/mal/mal_client.c
+++ b/monetdb5/mal/mal_client.c
@@ -44,6 +44,7 @@
#include "mal_parser.h"
#include "mal_namespace.h"
#include "mal_private.h"
+#include "mal_internal.h"
#include "mal_interpreter.h"
#include "mal_runtime.h"
#include "mal_authorize.h"
@@ -194,7 +195,7 @@ MCresetProfiler(stream *fdout)
MT_lock_unset(&mal_profileLock);
}
-void
+static void
MCexitClient(Client c)
{
MCresetProfiler(c->fdout);
@@ -221,7 +222,6 @@ MCexitClient(Client c)
&(struct NonMalEvent)
{CLIENT_END, c, Tend, NULL, NULL, 0,
Tend-(c->session)});
}
- setClientContext(NULL);
}
static Client
@@ -311,12 +311,14 @@ MCinitClient(oid user, bstream *fin, str
(void) c_old;
assert(NULL == c_old);
c = MCinitClientRecord(c, user, fin, fout);
+ MT_thread_set_qry_ctx(&c->qryctx);
}
MT_lock_unset(&mal_contextLock);
- profilerEvent(NULL,
- &(struct NonMalEvent)
- {CLIENT_START, c, c->session, NULL, NULL, 0,
0});
+ if (c)
+ profilerEvent(NULL,
+ &(struct NonMalEvent)
+ {CLIENT_START, c, c->session, NULL,
NULL, 0, 0});
return c;
}
@@ -495,6 +497,8 @@ MCcloseClient(Client c)
c->sqlprofiler = 0;
free(c->handshake_options);
c->handshake_options = NULL;
+ setClientContext(NULL);
+ MT_thread_set_qry_ctx(NULL);
MT_sema_destroy(&c->s);
MT_lock_set(&mal_contextLock);
if (shutdowninprogress) {
diff --git a/monetdb5/mal/mal_dataflow.c b/monetdb5/mal/mal_dataflow.c
--- a/monetdb5/mal/mal_dataflow.c
+++ b/monetdb5/mal/mal_dataflow.c
@@ -76,6 +76,7 @@ typedef struct DATAFLOW {
int *edges; /* dependency graph */
MT_Lock flowlock; /* lock to protect the above */
Queue *done; /* instructions handled */
+ bool set_qry_ctx;
} *DataFlow, DataFlowRec;
static struct worker {
@@ -325,6 +326,8 @@ DFLOWworker(void *T)
assert(t->flag == RUNNING);
cntxt = ATOMIC_PTR_GET(&t->cntxt);
while (1) {
+ MT_thread_set_qry_ctx(NULL);
+ setClientContext(NULL);
if (fnxt == 0) {
MT_thread_setworking(NULL);
cntxt = ATOMIC_PTR_GET(&t->cntxt);
@@ -354,6 +357,8 @@ DFLOWworker(void *T)
assert(fe);
flow = fe->flow;
assert(flow);
+ MT_thread_set_qry_ctx(flow->set_qry_ctx ?
&flow->cntxt->qryctx : NULL);
+ setClientContext(flow->cntxt);
/* whenever we have a (concurrent) error, skip it */
if (ATOMIC_PTR_GET(&flow->error)) {
@@ -906,6 +911,7 @@ runMALdataflow(Client cntxt, MalBlkPtr m
flow->cntxt = cntxt;
flow->mb = mb;
flow->stk = stk;
+ flow->set_qry_ctx = MT_thread_get_qry_ctx() != NULL;
/* keep real block count, exclude brackets */
flow->start = startpc + 1;
diff --git a/monetdb5/mal/mal_embedded.c b/monetdb5/mal/mal_embedded.c
--- a/monetdb5/mal/mal_embedded.c
+++ b/monetdb5/mal/mal_embedded.c
@@ -29,6 +29,7 @@
#include "mal_client.h"
#include "mal_dataflow.h"
#include "mal_private.h"
+#include "mal_internal.h"
#include "mal_runtime.h"
#include "mal_atom.h"
#include "mal_resource.h"
@@ -44,6 +45,7 @@ str
malEmbeddedBoot(int workerlimit, int memorylimit, int querytimeout, int
sessiontimeout, bool with_mapi_server)
{
Client c, c_old;
+ QryCtx *qc_old;
str msg = MAL_SUCCEED;
if( embeddedinitialized )
@@ -99,6 +101,7 @@ malEmbeddedBoot(int workerlimit, int mem
initHeartbeat();
// initResource();
c_old = setClientContext(NULL); //save context
+ qc_old = MT_thread_get_qry_ctx();
c = MCinitClient((oid) 0, 0, 0);
if(c == NULL)
throw(MAL, "malEmbeddedBoot", "Failed to initialize client");
@@ -110,22 +113,26 @@ malEmbeddedBoot(int workerlimit, int mem
if(c->usermodule == NULL) {
MCcloseClient(c);
setClientContext(c_old); // restore context
+ MT_thread_set_qry_ctx(qc_old);
throw(MAL, "malEmbeddedBoot", "Failed to initialize client MAL
module");
}
if ( (msg = defaultScenario(c)) ) {
MCcloseClient(c);
setClientContext(c_old); // restore context
+ MT_thread_set_qry_ctx(qc_old);
return msg;
}
if ((msg = MSinitClientPrg(c, "user", "main")) != MAL_SUCCEED) {
MCcloseClient(c);
setClientContext(c_old); // restore context
+ MT_thread_set_qry_ctx(qc_old);
return msg;
}
char *modules[5] = { "embedded", "sql", "generator", "udf" };
if ((msg = malIncludeModules(c, modules, 0, !with_mapi_server, NULL))
!= MAL_SUCCEED) {
MCcloseClient(c);
setClientContext(c_old); // restore context
+ MT_thread_set_qry_ctx(qc_old);
return msg;
}
pushEndInstruction(c->curprg->def);
@@ -133,6 +140,7 @@ malEmbeddedBoot(int workerlimit, int mem
if ( msg != MAL_SUCCEED || (msg= c->curprg->def->errors) != MAL_SUCCEED
) {
MCcloseClient(c);
setClientContext(c_old); // restore context
+ MT_thread_set_qry_ctx(qc_old);
return msg;
}
msg = MALengine(c);
@@ -140,6 +148,7 @@ malEmbeddedBoot(int workerlimit, int mem
embeddedinitialized = true;
MCcloseClient(c);
setClientContext(c_old); // restore context
+ MT_thread_set_qry_ctx(qc_old);
initProfiler();
return msg;
}
diff --git a/monetdb5/mal/mal_import.c b/monetdb5/mal/mal_import.c
--- a/monetdb5/mal/mal_import.c
+++ b/monetdb5/mal/mal_import.c
@@ -32,6 +32,7 @@
#include "mal_parser.h"
#include "mal_authorize.h"
#include "mal_private.h"
+#include "mal_internal.h"
#include "mal_session.h"
#include "mal_utils.h"
@@ -275,6 +276,7 @@ str
compileString(Symbol *fcn, Client cntxt, str s)
{
Client c, c_old;
+ QryCtx *qc_old;
size_t len = strlen(s);
buffer *b;
str msg = MAL_SUCCEED;
@@ -314,11 +316,14 @@ compileString(Symbol *fcn, Client cntxt,
strncpy(fdin->buf, qry, len+1);
c_old = setClientContext(NULL); // save context
+ qc_old = MT_thread_get_qry_ctx();
// compile in context of called for
c = MCinitClient(MAL_ADMIN, fdin, 0);
if( c == NULL){
GDKfree(qry);
GDKfree(b);
+ setClientContext(c_old); // restore context
+ MT_thread_set_qry_ctx(qc_old);
throw(MAL,"mal.eval","Can not create user context");
}
c->curmodule = c->usermodule = cntxt->usermodule;
@@ -331,6 +336,7 @@ compileString(Symbol *fcn, Client cntxt,
c->usermodule= 0;
MCcloseClient(c);
setClientContext(c_old); // restore context
+ MT_thread_set_qry_ctx(qc_old);
return msg;
}
@@ -346,6 +352,7 @@ compileString(Symbol *fcn, Client cntxt,
/* restore IO channel */
MCcloseClient(c);
setClientContext(c_old); // restore context
+ MT_thread_set_qry_ctx(qc_old);
GDKfree(qry);
GDKfree(b);
return msg;
diff --git a/monetdb5/mal/mal_internal.h b/monetdb5/mal/mal_internal.h
--- a/monetdb5/mal/mal_internal.h
+++ b/monetdb5/mal/mal_internal.h
@@ -14,5 +14,7 @@
void setqptimeout(lng usecs)
__attribute__((__visibility__("hidden")));
+Client setClientContext(Client cntxt)
+ __attribute__((__visibility__("hidden")));
extern size_t qsize;
diff --git a/monetdb5/mal/mal_interpreter.c b/monetdb5/mal/mal_interpreter.c
--- a/monetdb5/mal/mal_interpreter.c
+++ b/monetdb5/mal/mal_interpreter.c
@@ -534,15 +534,6 @@ runMALsequence(Client cntxt, MalBlkPtr m
stkpc = startpc;
exceptionVar = -1;
-#ifndef NDEBUG
- /* very short timeout */
- QryCtx qry_ctx_abort = {.querytimeout=100,
.starttime=cntxt->qryctx.starttime};
-#endif
- /* save, in case this function is called recursively */
- QryCtx *qry_ctx_save = MT_thread_get_qry_ctx();
- MT_thread_set_qry_ctx(&cntxt->qryctx);
- Client outer_cntxt = setClientContext(cntxt);
-
while (stkpc < mb->stop && stkpc != stoppc) {
// incomplete block being executed, requires at least signature
and end statement
MT_thread_setalgorithm(NULL);
@@ -570,7 +561,6 @@ runMALsequence(Client cntxt, MalBlkPtr m
if (stk->cmd == 'x' ) {
stk->cmd = 0;
stkpc = mb->stop;
- MT_thread_set_qry_ctx(&qry_ctx_abort);
ret= createException(MAL, "mal.interpreter",
"prematurely stopped client");
break;
}
@@ -1214,7 +1204,6 @@ runMALsequence(Client cntxt, MalBlkPtr m
if (garbage != garbages)
GDKfree(garbage);
ret = yieldFactory(mb, pci, stkpc);
- MT_thread_set_qry_ctx(qry_ctx_save);
return ret;
case RETURNsymbol:
/* Return from factory involves cleanup */
@@ -1255,11 +1244,6 @@ runMALsequence(Client cntxt, MalBlkPtr m
}
}
- /* restore saved values */
- MT_thread_set_qry_ctx(qry_ctx_save);
- setClientContext(outer_cntxt);
-
-
/* if we could not find the exception variable, cascade a new one */
/* don't add 'exception not caught' extra message for MAL sequences
besides main function calls */
if (exceptionVar >= 0 && (ret == MAL_SUCCEED || !pcicaller)) {
diff --git a/monetdb5/mal/mal_private.h b/monetdb5/mal/mal_private.h
--- a/monetdb5/mal/mal_private.h
+++ b/monetdb5/mal/mal_private.h
@@ -14,8 +14,6 @@
#ifdef _MAL_CLIENT_H_
/* _MAL_CLIENT_H_ is defined in the same file as Client */
-void MCexitClient(Client c)
- __attribute__((__visibility__("hidden")));
bool MCinit(void)
__attribute__((__visibility__("hidden")));
int MCinitClientThread(Client c)
@@ -39,9 +37,6 @@ str yieldFactory(MalBlkPtr mb, InstrPtr
__attribute__((__visibility__("hidden")));
str callFactory(Client cntxt, MalBlkPtr mb, ValPtr argv[],char flag)
__attribute__((__visibility__("hidden")));
-
-Client setClientContext(Client cntxt)
- __attribute__((__visibility__("hidden")));
_______________________________________________
checkin-list mailing list -- [email protected]
To unsubscribe send an email to [email protected]