LCOV - code coverage report
Current view: top level - monetdb5/mal - mal_dataflow.c (source / functions) Hit Total Coverage
Test: coverage.info Lines: 430 519 82.9 %
Date: 2024-11-15 19:37:45 Functions: 12 14 85.7 %

          Line data    Source code
       1             : /*
       2             :  * SPDX-License-Identifier: MPL-2.0
       3             :  *
       4             :  * This Source Code Form is subject to the terms of the Mozilla Public
       5             :  * License, v. 2.0.  If a copy of the MPL was not distributed with this
       6             :  * file, You can obtain one at http://mozilla.org/MPL/2.0/.
       7             :  *
       8             :  * Copyright 2024 MonetDB Foundation;
       9             :  * Copyright August 2008 - 2023 MonetDB B.V.;
      10             :  * Copyright 1997 - July 2008 CWI.
      11             :  */
      12             : 
      13             : /*
      14             :  * (author) M Kersten, S Mullender
      15             :  * Dataflow processing only works on a code
      16             :  * sequence that does not include additional (implicit) flow of control
      17             :  * statements and, ideally, consist of expensive BAT operations.
      18             :  * The dataflow portion is identified as a guarded block,
      19             :  * whose entry is controlled by the function language.dataflow();
      20             :  *
      21             :  * The dataflow worker tries to follow the sequence of actions
      22             :  * as laid out in the plan, but abandon this track when it hits
      23             :  * a blocking operator, or an instruction for which not all arguments
      24             :  * are available or resources become scarce.
      25             :  *
      26             :  * The flow graphs is organized such that parallel threads can
      27             :  * access it mostly without expensive locking and dependent
      28             :  * variables are easy to find..
      29             :  */
      30             : #include "monetdb_config.h"
      31             : #include "mal_dataflow.h"
      32             : #include "mal_exception.h"
      33             : #include "mal_private.h"
      34             : #include "mal_internal.h"
      35             : #include "mal_runtime.h"
      36             : #include "mal_resource.h"
      37             : #include "mal_function.h"
      38             : 
      39             : #define DFLOWpending 0                  /* runnable */
      40             : #define DFLOWrunning 1                  /* currently in progress */
      41             : #define DFLOWwrapup  2                  /* done! */
      42             : #define DFLOWretry   3                  /* reschedule */
      43             : #define DFLOWskipped 4                  /* due to errors */
      44             : 
      45             : /* The per instruction status of execution */
      46             : typedef struct FLOWEVENT {
      47             :         struct DATAFLOW *flow;          /* execution context */
      48             :         int pc;                                         /* pc in underlying malblock */
      49             :         int blocks;                                     /* awaiting for variables */
      50             :         sht state;                                      /* of execution */
      51             :         lng clk;
      52             :         sht cost;
      53             :         lng hotclaim;                           /* memory foot print of result variables */
      54             :         lng argclaim;                           /* memory foot print of arguments */
      55             :         lng maxclaim;                           /* memory foot print of largest argument, could be used to indicate result size */
      56             :         struct FLOWEVENT *next;         /* linked list for queues */
      57             : } *FlowEvent, FlowEventRec;
      58             : 
      59             : typedef struct queue {
      60             :         int exitcount;                          /* how many threads should exit */
      61             :         FlowEvent first, last;          /* first and last element of the queue */
      62             :         MT_Lock l;                                      /* it's a shared resource, ie we need locks */
      63             :         MT_Sema s;                                      /* threads wait on empty queues */
      64             : } Queue;
      65             : 
      66             : /*
      67             :  * The dataflow dependency is administered in a graph list structure.
      68             :  * For each instruction we keep the list of instructions that
      69             :  * should be checked for eligibility once we are finished with it.
      70             :  */
      71             : typedef struct DATAFLOW {
      72             :         Client cntxt;                           /* for debugging and client resolution */
      73             :         MalBlkPtr mb;                           /* carry the context */
      74             :         MalStkPtr stk;
      75             :         int start, stop;                        /* guarded block under consideration */
      76             :         FlowEvent status;                       /* status of each instruction */
      77             :         ATOMIC_PTR_TYPE error;          /* error encountered */
      78             :         int *nodes;                                     /* dependency graph nodes */
      79             :         int *edges;                                     /* dependency graph */
      80             :         MT_Lock flowlock;                       /* lock to protect the above */
      81             :         Queue *done;                            /* instructions handled */
      82             :         bool set_qry_ctx;
      83             : } *DataFlow, DataFlowRec;
      84             : 
      85             : struct worker {
      86             :         MT_Id id;
      87             :         enum { WAITING, RUNNING, FREE, EXITED, FINISHING } flag;
      88             :         ATOMIC_PTR_TYPE cntxt;          /* client we do work for (NULL -> any) */
      89             :         MT_Sema s;
      90             :         struct worker *next;
      91             :         char errbuf[GDKMAXERRLEN];      /* GDKerrbuf so that we can allocate before fork */
      92             : };
      93             : /* heads of three mutually exclusive linked lists, all using the .next
      94             :  * field in the worker struct */
      95             : static struct worker *workers;            /* "working" workers */
      96             : static struct worker *exited_workers; /* to be joined threads (.flag==EXITED) */
      97             : static struct worker *free_workers;     /* free workers (.flag==FREE) */
      98             : static int free_count = 0;              /* number of free threads */
      99             : static int free_max = 0;                /* max number of spare free threads */
     100             : 
     101             : static Queue *todo = 0;                 /* pending instructions */
     102             : 
     103             : static ATOMIC_TYPE exiting = ATOMIC_VAR_INIT(0);
     104             : static MT_Lock dataflowLock = MT_LOCK_INITIALIZER(dataflowLock);
     105             : 
     106             : /*
     107             :  * Calculate the size of the dataflow dependency graph.
     108             :  */
     109             : static int
     110             : DFLOWgraphSize(MalBlkPtr mb, int start, int stop)
     111             : {
     112             :         int cnt = 0;
     113             : 
     114    12298729 :         for (int i = start; i < stop; i++)
     115    12191161 :                 cnt += getInstrPtr(mb, i)->argc;
     116      107568 :         return cnt;
     117             : }
     118             : 
     119             : /*
     120             :  * The dataflow execution is confined to a barrier block.
     121             :  * Within the block there are multiple flows, which, in principle,
     122             :  * can be executed in parallel.
     123             :  */
     124             : 
     125             : static Queue *
     126      107882 : q_create(const char *name)
     127             : {
     128      107882 :         Queue *q = GDKzalloc(sizeof(Queue));
     129             : 
     130      107882 :         if (q == NULL)
     131             :                 return NULL;
     132      107882 :         MT_lock_init(&q->l, name);
     133      107882 :         MT_sema_init(&q->s, 0, name);
     134      107882 :         return q;
     135             : }
     136             : 
     137             : static void
     138      107567 : q_destroy(Queue *q)
     139             : {
     140      107567 :         assert(q);
     141      107567 :         MT_lock_destroy(&q->l);
     142      107568 :         MT_sema_destroy(&q->s);
     143      107563 :         GDKfree(q);
     144      107568 : }
     145             : 
     146             : /* keep a simple LIFO queue. It won't be a large one, so shuffles of requeue is possible */
     147             : /* we might actually sort it for better scheduling behavior */
     148             : static void
     149    18864930 : q_enqueue(Queue *q, FlowEvent d)
     150             : {
     151    18864930 :         assert(q);
     152    18864930 :         assert(d);
     153    18864930 :         MT_lock_set(&q->l);
     154    18887430 :         if (q->first == NULL) {
     155     4556820 :                 assert(q->last == NULL);
     156     4556820 :                 q->first = q->last = d;
     157             :         } else {
     158    14330610 :                 assert(q->last != NULL);
     159    14330610 :                 q->last->next = d;
     160    14330610 :                 q->last = d;
     161             :         }
     162    18887430 :         d->next = NULL;
     163    18887430 :         MT_lock_unset(&q->l);
     164    18912466 :         MT_sema_up(&q->s);
     165    18873107 : }
     166             : 
     167             : /*
     168             :  * A priority queue over the hot claims of memory may
     169             :  * be more effective. It prioritizes those instructions
     170             :  * that want to use a big recent result
     171             :  */
     172             : 
     173             : static void
     174           0 : q_requeue(Queue *q, FlowEvent d)
     175             : {
     176           0 :         assert(q);
     177           0 :         assert(d);
     178           0 :         MT_lock_set(&q->l);
     179           0 :         if (q->first == NULL) {
     180           0 :                 assert(q->last == NULL);
     181           0 :                 q->first = q->last = d;
     182           0 :                 d->next = NULL;
     183             :         } else {
     184           0 :                 assert(q->last != NULL);
     185           0 :                 d->next = q->first;
     186           0 :                 q->first = d;
     187             :         }
     188           0 :         MT_lock_unset(&q->l);
     189           0 :         MT_sema_up(&q->s);
     190           0 : }
     191             : 
     192             : static FlowEvent
     193    19095402 : q_dequeue(Queue *q, Client cntxt)
     194             : {
     195    19095402 :         assert(q);
     196    19095402 :         MT_sema_down(&q->s);
     197    19131519 :         if (ATOMIC_GET(&exiting))
     198             :                 return NULL;
     199    19129311 :         MT_lock_set(&q->l);
     200    19114157 :         if (cntxt == NULL && q->exitcount > 0) {
     201      107568 :                 q->exitcount--;
     202      107568 :                 MT_lock_unset(&q->l);
     203      107568 :                 return NULL;
     204             :         }
     205             : 
     206    19006589 :         FlowEvent *dp = &q->first;
     207    19006589 :         FlowEvent pd = NULL;
     208             :         /* if cntxt == NULL, return the first event, if cntxt != NULL, find
     209             :          * the first event in the queue with matching cntxt value and return
     210             :          * that */
     211    19006589 :         if (cntxt != NULL) {
     212    29181454 :                 while (*dp && (*dp)->flow->cntxt != cntxt) {
     213    27728879 :                         pd = *dp;
     214    27728879 :                         dp = &pd->next;
     215             :                 }
     216             :         }
     217    19006589 :         FlowEvent d = *dp;
     218    19006589 :         if (d) {
     219    18844054 :                 *dp = d->next;
     220    18844054 :                 d->next = NULL;
     221    18844054 :                 if (*dp == NULL)
     222     4577082 :                         q->last = pd;
     223             :         }
     224    19006589 :         MT_lock_unset(&q->l);
     225    18983271 :         return d;
     226             : }
     227             : 
     228             : /*
     229             :  * We simply move an instruction into the front of the queue.
     230             :  * Beware, we assume that variables are assigned a value once, otherwise
     231             :  * the order may really create errors.
     232             :  * The order of the instructions should be retained as long as possible.
     233             :  * Delay processing when we run out of memory.  Push the instruction back
     234             :  * on the end of queue, waiting for another attempt. Problem might become
     235             :  * that all threads but one are cycling through the queue, each time
     236             :  * finding an eligible instruction, but without enough space.
     237             :  * Therefore, we wait for a few milliseconds as an initial punishment.
     238             :  *
     239             :  * The process could be refined by checking for cheap operations,
     240             :  * i.e. those that would require no memory at all (aggr.count)
     241             :  * This, however, would lead to a dependency to the upper layers,
     242             :  * because in the kernel we don't know what routines are available
     243             :  * with this property. Nor do we maintain such properties.
     244             :  */
     245             : 
     246             : static void
     247        4708 : DFLOWworker(void *T)
     248             : {
     249        4708 :         struct worker *t = (struct worker *) T;
     250        4708 :         bool locked = false;
     251             : #ifdef _MSC_VER
     252             :         srand((unsigned int) GDKusec());
     253             : #endif
     254        4708 :         GDKsetbuf(t->errbuf);                /* where to leave errors */
     255        4700 :         snprintf(t->s.name, sizeof(t->s.name), "DFLOWsema%04zu", MT_getpid());
     256             : 
     257      109721 :         for (;;) {
     258      109721 :                 DataFlow flow;
     259      109721 :                 FlowEvent fe = 0, fnxt = 0;
     260      109721 :                 str error = 0;
     261      109721 :                 int i;
     262      109721 :                 lng claim;
     263      109721 :                 Client cntxt;
     264      109721 :                 InstrPtr p;
     265             : 
     266      109721 :                 GDKclrerr();
     267             : 
     268      109697 :                 if (t->flag == WAITING) {
     269             :                         /* wait until we are allowed to start working */
     270      107557 :                         MT_sema_down(&t->s);
     271      107564 :                         t->flag = RUNNING;
     272      107564 :                         if (ATOMIC_GET(&exiting)) {
     273             :                                 break;
     274             :                         }
     275             :                 }
     276      109704 :                 assert(t->flag == RUNNING);
     277      109704 :                 cntxt = ATOMIC_PTR_GET(&t->cntxt);
     278    12319641 :                 while (1) {
     279    12319641 :                         MT_thread_set_qry_ctx(NULL);
     280    12291041 :                         if (fnxt == 0) {
     281     7092260 :                                 MT_thread_setworking("waiting for work");
     282     7091462 :                                 cntxt = ATOMIC_PTR_GET(&t->cntxt);
     283     7091462 :                                 fe = q_dequeue(todo, cntxt);
     284     7106177 :                                 if (fe == NULL) {
     285      271994 :                                         if (cntxt) {
     286             :                                                 /* we're not done yet with work for the current
     287             :                                                  * client (as far as we know), so give up the CPU
     288             :                                                  * and let the scheduler enter some more work, but
     289             :                                                  * first compensate for the down we did in
     290             :                                                  * dequeue */
     291      162866 :                                                 MT_sema_up(&todo->s);
     292      162860 :                                                 MT_sleep_ms(1);
     293    12482317 :                                                 continue;
     294             :                                         }
     295             :                                         /* no more work to be done: exit */
     296      109128 :                                         break;
     297             :                                 }
     298     6834183 :                                 if (fe->flow->cntxt && fe->flow->cntxt->mythread)
     299     6833351 :                                         MT_thread_setworking(fe->flow->cntxt->mythread);
     300             :                         } else
     301             :                                 fe = fnxt;
     302    12032428 :                         if (ATOMIC_GET(&exiting)) {
     303             :                                 break;
     304             :                         }
     305    12032428 :                         fnxt = 0;
     306    12032428 :                         assert(fe);
     307    12032428 :                         flow = fe->flow;
     308    12032428 :                         assert(flow);
     309    12032428 :                         MT_thread_set_qry_ctx(flow->set_qry_ctx ? &flow->cntxt->
     310             :                                                                   qryctx : NULL);
     311             : 
     312             :                         /* whenever we have a (concurrent) error, skip it */
     313    12052075 :                         if (ATOMIC_PTR_GET(&flow->error)) {
     314        4179 :                                 q_enqueue(flow->done, fe);
     315        4177 :                                 continue;
     316             :                         }
     317             : 
     318    12047896 :                         p = getInstrPtr(flow->mb, fe->pc);
     319    12047896 :                         claim = fe->argclaim;
     320    24119060 :                         if (p->fcn != (MALfcn) deblockdataflow &&    /* never block on deblockdataflow() */
     321    12047896 :                                 !MALadmission_claim(flow->cntxt, flow->mb, flow->stk, p, claim)) {
     322           0 :                                 fe->hotclaim = 0;    /* don't assume priority anymore */
     323           0 :                                 fe->maxclaim = 0;
     324           0 :                                 MT_lock_set(&todo->l);
     325           0 :                                 FlowEvent last = todo->last;
     326           0 :                                 MT_lock_unset(&todo->l);
     327           0 :                                 if (last == NULL)
     328           0 :                                         MT_sleep_ms(DELAYUNIT);
     329           0 :                                 q_requeue(todo, fe);
     330           0 :                                 continue;
     331             :                         }
     332    12071164 :                         ATOMIC_BASE_TYPE wrks = ATOMIC_INC(&flow->cntxt->workers);
     333    12071164 :                         ATOMIC_BASE_TYPE mwrks = ATOMIC_GET(&flow->mb->workers);
     334    12071491 :                         while (wrks > mwrks) {
     335       31437 :                                 if (ATOMIC_CAS(&flow->mb->workers, &mwrks, wrks))
     336             :                                         break;
     337             :                         }
     338    12071194 :                         error = runMALsequence(flow->cntxt, flow->mb, fe->pc, fe->pc + 1,
     339             :                                                                    flow->stk, 0, 0);
     340    11927223 :                         ATOMIC_DEC(&flow->cntxt->workers);
     341             :                         /* release the memory claim */
     342    11927223 :                         MALadmission_release(flow->cntxt, flow->mb, flow->stk, p, claim);
     343             : 
     344    12015653 :                         MT_lock_set(&flow->flowlock);
     345    12076693 :                         fe->state = DFLOWwrapup;
     346    12076693 :                         MT_lock_unset(&flow->flowlock);
     347    12063188 :                         if (error) {
     348         527 :                                 void *null = NULL;
     349             :                                 /* only collect one error (from one thread, needed for stable testing) */
     350         527 :                                 if (!ATOMIC_PTR_CAS(&flow->error, &null, error))
     351          51 :                                         freeException(error);
     352             :                                 /* after an error we skip the rest of the block */
     353         527 :                                 q_enqueue(flow->done, fe);
     354         527 :                                 continue;
     355             :                         }
     356             : 
     357             :                         /* see if you can find an eligible instruction that uses the
     358             :                          * result just produced. Then we can continue with it right away.
     359             :                          * We are just looking forward for the last block, which means we
     360             :                          * are safe from concurrent actions. No other thread can steal it,
     361             :                          * because we hold the logical lock.
     362             :                          * All eligible instructions are queued
     363             :                          */
     364    12062661 :                         p = getInstrPtr(flow->mb, fe->pc);
     365    12062661 :                         assert(p);
     366    12062661 :                         fe->hotclaim = 0;
     367    12062661 :                         fe->maxclaim = 0;
     368             : 
     369    25082389 :                         for (i = 0; i < p->retc; i++) {
     370    13026937 :                                 lng footprint;
     371    13026937 :                                 footprint = getMemoryClaim(flow->mb, flow->stk, p, i, FALSE);
     372    13019728 :                                 fe->hotclaim += footprint;
     373    13019728 :                                 if (footprint > fe->maxclaim)
     374     4914235 :                                         fe->maxclaim = footprint;
     375             :                         }
     376             : 
     377             : /* Try to get rid of the hot potato or locate an alternative to proceed.
     378             :  */
     379             : #define HOTPOTATO
     380             : #ifdef HOTPOTATO
     381             :                         /* HOT potato choice */
     382    12055452 :                         int last = 0, nxt = -1;
     383    12055452 :                         lng nxtclaim = -1;
     384             : 
     385    12055452 :                         MT_lock_set(&flow->flowlock);
     386    12078143 :                         for (last = fe->pc - flow->start;
     387    49441215 :                                  last >= 0 && (i = flow->nodes[last]) > 0;
     388    37363072 :                                  last = flow->edges[last]) {
     389    37363072 :                                 if (flow->status[i].state == DFLOWpending
     390    37360327 :                                         && flow->status[i].blocks == 1) {
     391             :                                         /* find the one with the largest footprint */
     392    10191306 :                                         if (nxt == -1 || flow->status[i].argclaim > nxtclaim) {
     393     5652580 :                                                 nxt = i;
     394     5652580 :                                                 nxtclaim = flow->status[i].argclaim;
     395             :                                         }
     396             :                                 }
     397             :                         }
     398             :                         /* hot potato can not be removed, use alternative to proceed */
     399    12078143 :                         if (nxt >= 0) {
     400     5248702 :                                 flow->status[nxt].state = DFLOWrunning;
     401     5248702 :                                 flow->status[nxt].blocks = 0;
     402     5248702 :                                 flow->status[nxt].hotclaim = fe->hotclaim;
     403     5248702 :                                 flow->status[nxt].argclaim += fe->hotclaim;
     404     5248702 :                                 if (flow->status[nxt].maxclaim < fe->maxclaim)
     405     2071950 :                                         flow->status[nxt].maxclaim = fe->maxclaim;
     406             :                                 fnxt = flow->status + nxt;
     407             :                         }
     408    12078143 :                         MT_lock_unset(&flow->flowlock);
     409             : #endif
     410             : 
     411    12063364 :                         q_enqueue(flow->done, fe);
     412    12042557 :                         if (fnxt == 0 && profilerStatus) {
     413           0 :                                 profilerHeartbeatEvent("wait");
     414             :                         }
     415             :                 }
     416      109128 :                 MT_lock_set(&dataflowLock);
     417      109641 :                 if (GDKexiting() || ATOMIC_GET(&exiting) || free_count >= free_max) {
     418             :                         locked = true;
     419             :                         break;
     420             :                 }
     421      105612 :                 free_count++;
     422      105612 :                 struct worker **tp = &workers;
     423      695659 :                 while (*tp && *tp != t)
     424      590047 :                         tp = &(*tp)->next;
     425      105612 :                 assert(*tp && *tp == t);
     426      105612 :                 *tp = t->next;
     427      105612 :                 t->flag = FREE;
     428      105612 :                 t->next = free_workers;
     429      105612 :                 free_workers = t;
     430      105612 :                 MT_lock_unset(&dataflowLock);
     431      105612 :                 MT_thread_setworking("idle, waiting for new client");
     432      105612 :                 MT_sema_down(&t->s);
     433      105590 :                 if (GDKexiting() || ATOMIC_GET(&exiting))
     434             :                         break;
     435      105033 :                 assert(t->flag == WAITING);
     436             :         }
     437         551 :         if (!locked)
     438         551 :                 MT_lock_set(&dataflowLock);
     439        4580 :         if (t->flag != FINISHING) {
     440        4020 :                 struct worker **tp = t->flag == FREE ? &free_workers : &workers;
     441       18770 :                 while (*tp && *tp != t)
     442       14750 :                         tp = &(*tp)->next;
     443        4020 :                 assert(*tp && *tp == t);
     444        4020 :                 *tp = t->next;
     445        4020 :                 t->flag = EXITED;
     446        4020 :                 t->next = exited_workers;
     447        4020 :                 exited_workers = t;
     448             :         }
     449        4580 :         MT_lock_unset(&dataflowLock);
     450        4576 :         GDKsetbuf(NULL);
     451        4571 : }
     452             : 
     453             : /*
     454             :  * Create an interpreter pool.
     455             :  * One worker will adaptively be available for each client.
     456             :  * The remainder are taken from the GDKnr_threads argument and
     457             :  * typically is equal to the number of cores
     458             :  * The workers are assembled in a local table to enable debugging.
     459             :  */
     460             : static int
     461         314 : DFLOWinitialize(void)
     462             : {
     463         314 :         int limit;
     464         314 :         int created = 0;
     465             : 
     466         314 :         MT_lock_set(&mal_contextLock);
     467         314 :         MT_lock_set(&dataflowLock);
     468         314 :         if (todo) {
     469             :                 /* somebody else beat us to it */
     470           0 :                 MT_lock_unset(&dataflowLock);
     471           0 :                 MT_lock_unset(&mal_contextLock);
     472           0 :                 return 0;
     473             :         }
     474         314 :         free_max = GDKgetenv_int("dataflow_max_free",
     475             :                                                          GDKnr_threads < 4 ? 4 : GDKnr_threads);
     476         314 :         todo = q_create("todo");
     477         314 :         if (todo == NULL) {
     478           0 :                 MT_lock_unset(&dataflowLock);
     479           0 :                 MT_lock_unset(&mal_contextLock);
     480           0 :                 return -1;
     481             :         }
     482         314 :         limit = GDKnr_threads ? GDKnr_threads - 1 : 0;
     483        2506 :         while (limit > 0) {
     484        2192 :                 limit--;
     485        2192 :                 struct worker *t = GDKmalloc(sizeof(*t));
     486        2192 :                 if (t == NULL) {
     487           0 :                         TRC_CRITICAL(MAL_SERVER, "cannot allocate structure for worker");
     488           0 :                         continue;
     489             :                 }
     490        2192 :                 *t = (struct worker) {
     491             :                         .flag = RUNNING,
     492             :                         .cntxt = ATOMIC_PTR_VAR_INIT(NULL),
     493             :                 };
     494        2192 :                 MT_sema_init(&t->s, 0, "DFLOWsema"); /* placeholder name */
     495        2192 :                 if (MT_create_thread(&t->id, DFLOWworker, t,
     496             :                                                          MT_THR_JOINABLE, "DFLOWworkerXXXX") < 0) {
     497           0 :                         ATOMIC_PTR_DESTROY(&t->cntxt);
     498           0 :                         MT_sema_destroy(&t->s);
     499           0 :                         GDKfree(t);
     500             :                 } else {
     501        2192 :                         t->next = workers;
     502        2192 :                         workers = t;
     503        2192 :                         created++;
     504             :                 }
     505             :         }
     506         314 :         if (created == 0) {
     507             :                 /* no threads created */
     508           0 :                 q_destroy(todo);
     509           0 :                 todo = NULL;
     510           0 :                 MT_lock_unset(&dataflowLock);
     511           0 :                 MT_lock_unset(&mal_contextLock);
     512           0 :                 return -1;
     513             :         }
     514         314 :         MT_lock_unset(&dataflowLock);
     515         314 :         MT_lock_unset(&mal_contextLock);
     516         314 :         return 0;
     517             : }
     518             : 
     519             : /*
     520             :  * The dataflow administration is based on administration of
     521             :  * how many variables are still missing before it can be executed.
     522             :  * For each instruction we keep a list of instructions whose
     523             :  * blocking counter should be decremented upon finishing it.
     524             :  */
     525             : static str
     526      107568 : DFLOWinitBlk(DataFlow flow, MalBlkPtr mb, int size)
     527             : {
     528      107568 :         int pc, i, j, k, l, n, etop = 0;
     529      107568 :         int *assign;
     530      107568 :         InstrPtr p;
     531             : 
     532      107568 :         if (flow == NULL)
     533           0 :                 throw(MAL, "dataflow", "DFLOWinitBlk(): Called with flow == NULL");
     534      107568 :         if (mb == NULL)
     535           0 :                 throw(MAL, "dataflow", "DFLOWinitBlk(): Called with mb == NULL");
     536      107568 :         assign = (int *) GDKzalloc(mb->vtop * sizeof(int));
     537      107568 :         if (assign == NULL)
     538           0 :                 throw(MAL, "dataflow", SQLSTATE(HY013) MAL_MALLOC_FAIL);
     539      107568 :         etop = flow->stop - flow->start;
     540    12191275 :         for (n = 0, pc = flow->start; pc < flow->stop; pc++, n++) {
     541    12083708 :                 p = getInstrPtr(mb, pc);
     542    12083708 :                 if (p == NULL) {
     543           0 :                         GDKfree(assign);
     544           0 :                         throw(MAL, "dataflow",
     545             :                                   "DFLOWinitBlk(): getInstrPtr() returned NULL");
     546             :                 }
     547             : 
     548             :                 /* initial state, ie everything can run */
     549    12083708 :                 flow->status[n].flow = flow;
     550    12083708 :                 flow->status[n].pc = pc;
     551    12083708 :                 flow->status[n].state = DFLOWpending;
     552    12083708 :                 flow->status[n].cost = -1;
     553    12083708 :                 ATOMIC_PTR_SET(&flow->status[n].flow->error, NULL);
     554             : 
     555             :                 /* administer flow dependencies */
     556    56471775 :                 for (j = p->retc; j < p->argc; j++) {
     557             :                         /* list of instructions that wake n-th instruction up */
     558    44388068 :                         if (!isVarConstant(mb, getArg(p, j)) && (k = assign[getArg(p, j)])) {
     559    23459501 :                                 assert(k < pc);      /* only dependencies on earlier instructions */
     560             :                                 /* add edge to the target instruction for wakeup call */
     561    23459501 :                                 k -= flow->start;
     562    23459501 :                                 if (flow->nodes[k]) {
     563             :                                         /* add wakeup to tail of list */
     564   207415736 :                                         for (i = k; flow->edges[i] > 0; i = flow->edges[i])
     565             :                                                 ;
     566    21201973 :                                         flow->nodes[etop] = n;
     567    21201973 :                                         flow->edges[etop] = -1;
     568    21201973 :                                         flow->edges[i] = etop;
     569    21201973 :                                         etop++;
     570    21201973 :                                         (void) size;
     571    21201973 :                                         if (etop == size) {
     572         213 :                                                 int *tmp;
     573             :                                                 /* in case of realloc failure, the original
     574             :                                                  * pointers will be freed by the caller */
     575         213 :                                                 tmp = (int *) GDKrealloc(flow->nodes,
     576             :                                                                                                  sizeof(int) * 2 * size);
     577         212 :                                                 if (tmp == NULL) {
     578           0 :                                                         GDKfree(assign);
     579           0 :                                                         throw(MAL, "dataflow",
     580             :                                                                   SQLSTATE(HY013) MAL_MALLOC_FAIL);
     581             :                                                 }
     582         212 :                                                 flow->nodes = tmp;
     583         212 :                                                 tmp = (int *) GDKrealloc(flow->edges,
     584             :                                                                                                  sizeof(int) * 2 * size);
     585         212 :                                                 if (tmp == NULL) {
     586           0 :                                                         GDKfree(assign);
     587           0 :                                                         throw(MAL, "dataflow",
     588             :                                                                   SQLSTATE(HY013) MAL_MALLOC_FAIL);
     589             :                                                 }
     590         212 :                                                 flow->edges = tmp;
     591         212 :                                                 size *= 2;
     592             :                                         }
     593             :                                 } else {
     594     2257528 :                                         flow->nodes[k] = n;
     595     2257528 :                                         flow->edges[k] = -1;
     596             :                                 }
     597             : 
     598    23459500 :                                 flow->status[n].blocks++;
     599             :                         }
     600             : 
     601             :                         /* list of instructions to be woken up explicitly */
     602    44388067 :                         if (!isVarConstant(mb, getArg(p, j))) {
     603             :                                 /* be careful, watch out for garbage collection interference */
     604             :                                 /* those should be scheduled after all its other uses */
     605    24031751 :                                 l = getEndScope(mb, getArg(p, j));
     606    24031751 :                                 if (l != pc && l < flow->stop && l > flow->start) {
     607             :                                         /* add edge to the target instruction for wakeup call */
     608    13929853 :                                         assert(pc < l);      /* only dependencies on earlier instructions */
     609    13929853 :                                         l -= flow->start;
     610    13929853 :                                         if (flow->nodes[n]) {
     611             :                                                 /* add wakeup to tail of list */
     612    76772856 :                                                 for (i = n; flow->edges[i] > 0; i = flow->edges[i])
     613             :                                                         ;
     614     6910346 :                                                 flow->nodes[etop] = l;
     615     6910346 :                                                 flow->edges[etop] = -1;
     616     6910346 :                                                 flow->edges[i] = etop;
     617     6910346 :                                                 etop++;
     618     6910346 :                                                 if (etop == size) {
     619         209 :                                                         int *tmp;
     620             :                                                         /* in case of realloc failure, the original
     621             :                                                          * pointers will be freed by the caller */
     622         209 :                                                         tmp = (int *) GDKrealloc(flow->nodes,
     623             :                                                                                                          sizeof(int) * 2 * size);
     624         209 :                                                         if (tmp == NULL) {
     625           0 :                                                                 GDKfree(assign);
     626           0 :                                                                 throw(MAL, "dataflow",
     627             :                                                                           SQLSTATE(HY013) MAL_MALLOC_FAIL);
     628             :                                                         }
     629         209 :                                                         flow->nodes = tmp;
     630         209 :                                                         tmp = (int *) GDKrealloc(flow->edges,
     631             :                                                                                                          sizeof(int) * 2 * size);
     632         209 :                                                         if (tmp == NULL) {
     633           0 :                                                                 GDKfree(assign);
     634           0 :                                                                 throw(MAL, "dataflow",
     635             :                                                                           SQLSTATE(HY013) MAL_MALLOC_FAIL);
     636             :                                                         }
     637         209 :                                                         flow->edges = tmp;
     638         209 :                                                         size *= 2;
     639             :                                                 }
     640             :                                         } else {
     641     7019507 :                                                 flow->nodes[n] = l;
     642     7019507 :                                                 flow->edges[n] = -1;
     643             :                                         }
     644    13929853 :                                         flow->status[l].blocks++;
     645             :                                 }
     646             :                         }
     647             :                 }
     648             : 
     649    25137651 :                 for (j = 0; j < p->retc; j++)
     650    13053944 :                         assign[getArg(p, j)] = pc;      /* ensure recognition of dependency on first instruction and constant */
     651             :         }
     652      107567 :         GDKfree(assign);
     653             : 
     654      107567 :         return MAL_SUCCEED;
     655             : }
     656             : 
     657             : /*
     658             :  * Parallel processing is mostly driven by dataflow, but within this context
     659             :  * there may be different schemes to take instructions into execution.
     660             :  * The admission scheme (and wrapup) are the necessary scheduler hooks.
     661             :  * A scheduler registers the functions needed and should release them
     662             :  * at the end of the parallel block.
     663             :  * They take effect after we have ensured that the basic properties for
     664             :  * execution hold.
     665             :  */
     666             : static str
     667      107568 : DFLOWscheduler(DataFlow flow, struct worker *w)
     668             : {
     669      107568 :         int last;
     670      107568 :         int i;
     671      107568 :         int j;
     672      107568 :         InstrPtr p;
     673      107568 :         int tasks = 0, actions = 0;
     674      107568 :         str ret = MAL_SUCCEED;
     675      107568 :         FlowEvent fe, f = 0;
     676             : 
     677      107568 :         if (flow == NULL)
     678           0 :                 throw(MAL, "dataflow", "DFLOWscheduler(): Called with flow == NULL");
     679      107568 :         actions = flow->stop - flow->start;
     680      107568 :         if (actions == 0)
     681           0 :                 throw(MAL, "dataflow", "Empty dataflow block");
     682             :         /* initialize the eligible statements */
     683      107568 :         fe = flow->status;
     684             : 
     685      107568 :         ATOMIC_DEC(&flow->cntxt->workers);
     686      107568 :         MT_lock_set(&flow->flowlock);
     687    12299290 :         for (i = 0; i < actions; i++)
     688    12084156 :                 if (fe[i].blocks == 0) {
     689      572821 :                         p = getInstrPtr(flow->mb, fe[i].pc);
     690      572821 :                         if (p == NULL) {
     691           0 :                                 MT_lock_unset(&flow->flowlock);
     692           0 :                                 ATOMIC_INC(&flow->cntxt->workers);
     693           0 :                                 throw(MAL, "dataflow",
     694             :                                           "DFLOWscheduler(): getInstrPtr(flow->mb,fe[i].pc) returned NULL");
     695             :                         }
     696      572821 :                         fe[i].argclaim = 0;
     697     2799248 :                         for (j = p->retc; j < p->argc; j++)
     698     2226427 :                                 fe[i].argclaim += getMemoryClaim(fe[0].flow->mb,
     699     2226433 :                                                                                                  fe[0].flow->stk, p, j, FALSE);
     700      572815 :                         flow->status[i].state = DFLOWrunning;
     701      572815 :                         q_enqueue(todo, flow->status + i);
     702             :                 }
     703      107567 :         MT_lock_unset(&flow->flowlock);
     704      107567 :         MT_sema_up(&w->s);
     705             : 
     706      107567 :         while (actions != tasks) {
     707    12083123 :                 f = q_dequeue(flow->done, NULL);
     708    12083347 :                 if (ATOMIC_GET(&exiting))
     709             :                         break;
     710    12083347 :                 if (f == NULL) {
     711           0 :                         ATOMIC_INC(&flow->cntxt->workers);
     712           0 :                         throw(MAL, "dataflow",
     713             :                                   "DFLOWscheduler(): q_dequeue(flow->done) returned NULL");
     714             :                 }
     715             : 
     716             :                 /*
     717             :                  * When an instruction is finished we have to reduce the blocked
     718             :                  * counter for all dependent instructions.  for those where it
     719             :                  * drops to zero we can scheduler it we do it here instead of the scheduler
     720             :                  */
     721             : 
     722    12083347 :                 MT_lock_set(&flow->flowlock);
     723    12083019 :                 tasks++;
     724    12083019 :                 for (last = f->pc - flow->start;
     725    49463329 :                          last >= 0 && (i = flow->nodes[last]) > 0; last = flow->edges[last])
     726    37380655 :                         if (flow->status[i].state == DFLOWpending) {
     727    32138068 :                                 flow->status[i].argclaim += f->hotclaim;
     728    32138068 :                                 if (flow->status[i].blocks == 1) {
     729     6262039 :                                         flow->status[i].blocks--;
     730     6262039 :                                         flow->status[i].state = DFLOWrunning;
     731     6262039 :                                         q_enqueue(todo, flow->status + i);
     732             :                                 } else {
     733    25876029 :                                         flow->status[i].blocks--;
     734             :                                 }
     735             :                         }
     736    12190240 :                 MT_lock_unset(&flow->flowlock);
     737             :         }
     738             :         /* release the worker from its specific task (turn it into a
     739             :          * generic worker) */
     740      107567 :         ATOMIC_PTR_SET(&w->cntxt, NULL);
     741      107567 :         ATOMIC_INC(&flow->cntxt->workers);
     742             :         /* wrap up errors */
     743      107567 :         assert(flow->done->last == 0);
     744      107567 :         if ((ret = ATOMIC_PTR_XCG(&flow->error, NULL)) != NULL) {
     745         476 :                 TRC_DEBUG(MAL_SERVER, "Errors encountered: %s\n", ret);
     746             :         }
     747             :         return ret;
     748             : }
     749             : 
     750             : /* called and returns with dataflowLock locked, temporarily unlocks
     751             :  * join the thread associated with the worker and destroy the structure */
     752             : static inline void
     753        4580 : finish_worker(struct worker *t)
     754             : {
     755        4580 :         t->flag = FINISHING;
     756        4580 :         MT_lock_unset(&dataflowLock);
     757        4580 :         MT_join_thread(t->id);
     758        4580 :         MT_sema_destroy(&t->s);
     759        4580 :         ATOMIC_PTR_DESTROY(&t->cntxt);
     760        4580 :         GDKfree(t);
     761        4580 :         MT_lock_set(&dataflowLock);
     762        4580 : }
     763             : 
     764             : /* We create a pool of GDKnr_threads-1 generic workers, that is,
     765             :  * workers that will take on jobs from any clients.  In addition, we
     766             :  * create a single specific worker per client (i.e. each time we enter
     767             :  * here).  This specific worker will only do work for the client for
     768             :  * which it was started.  In this way we can guarantee that there will
     769             :  * always be progress for the client, even if all other workers are
     770             :  * doing something big.
     771             :  *
     772             :  * When all jobs for a client have been done (there are no more
     773             :  * entries for the client in the queue), the specific worker turns
     774             :  * itself into a generic worker.  At the same time, we signal that one
     775             :  * generic worker should exit and this function returns.  In this way
     776             :  * we make sure that there are once again GDKnr_threads-1 generic
     777             :  * workers. */
     778             : str
     779      107561 : runMALdataflow(Client cntxt, MalBlkPtr mb, int startpc, int stoppc,
     780             :                            MalStkPtr stk)
     781             : {
     782      107561 :         DataFlow flow = NULL;
     783      107561 :         str msg = MAL_SUCCEED;
     784      107561 :         int size;
     785      107561 :         bit *ret;
     786      107561 :         struct worker *t;
     787             : 
     788      107561 :         if (stk == NULL)
     789           0 :                 throw(MAL, "dataflow", "runMALdataflow(): Called with stk == NULL");
     790      107561 :         ret = getArgReference_bit(stk, getInstrPtr(mb, startpc), 0);
     791      107561 :         *ret = FALSE;
     792             : 
     793      107561 :         assert(stoppc > startpc);
     794             : 
     795             :         /* check existence of workers */
     796      107561 :         if (todo == NULL) {
     797             :                 /* create thread pool */
     798         314 :                 if (GDKnr_threads <= 1 || DFLOWinitialize() < 0) {
     799             :                         /* no threads created, run serially */
     800           0 :                         *ret = TRUE;
     801           0 :                         return MAL_SUCCEED;
     802             :                 }
     803             :         }
     804      107561 :         assert(todo);
     805             :         /* in addition, create one more worker that will only execute
     806             :          * tasks for the current client to compensate for our waiting
     807             :          * until all work is done */
     808      107561 :         MT_lock_set(&dataflowLock);
     809             :         /* join with already exited threads */
     810      109524 :         while (exited_workers != NULL) {
     811        1956 :                 assert(exited_workers->flag == EXITED);
     812        1956 :                 struct worker *t = exited_workers;
     813        1956 :                 exited_workers = exited_workers->next;
     814        1956 :                 finish_worker(t);
     815             :         }
     816      107568 :         assert(cntxt != NULL);
     817      107568 :         if (free_workers != NULL) {
     818      105044 :                 t = free_workers;
     819      105044 :                 assert(t->flag == FREE);
     820      105044 :                 assert(free_count > 0);
     821      105044 :                 free_count--;
     822      105044 :                 free_workers = t->next;
     823      105044 :                 t->next = workers;
     824      105044 :                 workers = t;
     825      105044 :                 t->flag = WAITING;
     826      105044 :                 ATOMIC_PTR_SET(&t->cntxt, cntxt);
     827      105044 :                 MT_sema_up(&t->s);
     828             :         } else {
     829        2524 :                 t = GDKmalloc(sizeof(*t));
     830        2524 :                 if (t != NULL) {
     831        2524 :                         *t = (struct worker) {
     832             :                                 .flag = WAITING,
     833             :                                 .cntxt = ATOMIC_PTR_VAR_INIT(cntxt),
     834             :                         };
     835        2524 :                         MT_sema_init(&t->s, 0, "DFLOWsema"); /* placeholder name */
     836        2524 :                         if (MT_create_thread(&t->id, DFLOWworker, t,
     837             :                                                                  MT_THR_JOINABLE, "DFLOWworkerXXXX") < 0) {
     838           0 :                                 ATOMIC_PTR_DESTROY(&t->cntxt);
     839           0 :                                 MT_sema_destroy(&t->s);
     840           0 :                                 GDKfree(t);
     841           0 :                                 t = NULL;
     842             :                         } else {
     843        2524 :                                 t->next = workers;
     844        2524 :                                 workers = t;
     845             :                         }
     846             :                 }
     847        2524 :                 if (t == NULL) {
     848             :                         /* cannot start new thread, run serially */
     849           0 :                         *ret = TRUE;
     850           0 :                         MT_lock_unset(&dataflowLock);
     851           0 :                         return MAL_SUCCEED;
     852             :                 }
     853             :         }
     854      107568 :         MT_lock_unset(&dataflowLock);
     855             : 
     856      107568 :         flow = (DataFlow) GDKzalloc(sizeof(DataFlowRec));
     857      107568 :         if (flow == NULL)
     858           0 :                 throw(MAL, "dataflow", SQLSTATE(HY013) MAL_MALLOC_FAIL);
     859             : 
     860      107568 :         size = DFLOWgraphSize(mb, startpc, stoppc);
     861      107568 :         size += stoppc - startpc;
     862             : 
     863      215136 :         *flow = (DataFlowRec) {
     864             :                 .cntxt = cntxt,
     865             :                 .mb = mb,
     866             :                 .stk = stk,
     867      107568 :                 .set_qry_ctx = MT_thread_get_qry_ctx() != NULL,
     868             :                 /* keep real block count, exclude brackets */
     869      107568 :                 .start = startpc + 1,
     870             :                 .stop = stoppc,
     871      107568 :                 .done = q_create("flow->done"),
     872      107568 :                 .status = (FlowEvent) GDKzalloc((stoppc - startpc + 1) *
     873             :                                                                                 sizeof(FlowEventRec)),
     874             :                 .error = ATOMIC_PTR_VAR_INIT(NULL),
     875      107568 :                 .nodes = (int *) GDKzalloc(sizeof(int) * size),
     876      107568 :                 .edges = (int *) GDKzalloc(sizeof(int) * size),
     877             :         };
     878             : 
     879      107568 :         if (flow->done == NULL) {
     880           0 :                 GDKfree(flow->status);
     881           0 :                 GDKfree(flow->nodes);
     882           0 :                 GDKfree(flow->edges);
     883           0 :                 GDKfree(flow);
     884           0 :                 throw(MAL, "dataflow",
     885             :                           "runMALdataflow(): Failed to create flow->done queue");
     886             :         }
     887             : 
     888      107568 :         if (flow->status == NULL || flow->nodes == NULL || flow->edges == NULL) {
     889           0 :                 q_destroy(flow->done);
     890           0 :                 GDKfree(flow->status);
     891           0 :                 GDKfree(flow->nodes);
     892           0 :                 GDKfree(flow->edges);
     893           0 :                 GDKfree(flow);
     894           0 :                 throw(MAL, "dataflow", SQLSTATE(HY013) MAL_MALLOC_FAIL);
     895             :         }
     896             : 
     897      107568 :         MT_lock_init(&flow->flowlock, "flow->flowlock");
     898      107568 :         msg = DFLOWinitBlk(flow, mb, size);
     899             : 
     900      107568 :         if (msg == MAL_SUCCEED)
     901      107568 :                 msg = DFLOWscheduler(flow, t);
     902             : 
     903      107568 :         GDKfree(flow->status);
     904      107568 :         GDKfree(flow->edges);
     905      107568 :         GDKfree(flow->nodes);
     906      107568 :         q_destroy(flow->done);
     907      107568 :         MT_lock_destroy(&flow->flowlock);
     908      107568 :         ATOMIC_PTR_DESTROY(&flow->error);
     909      107568 :         GDKfree(flow);
     910             : 
     911             :         /* we created one worker, now tell one worker to exit again */
     912      107568 :         MT_lock_set(&todo->l);
     913      107568 :         todo->exitcount++;
     914      107568 :         MT_lock_unset(&todo->l);
     915      107568 :         MT_sema_up(&todo->s);
     916             : 
     917      107568 :         return msg;
     918             : }
     919             : 
     920             : str
     921           0 : deblockdataflow(Client cntxt, MalBlkPtr mb, MalStkPtr stk, InstrPtr pci)
     922             : {
     923           0 :         int *ret = getArgReference_int(stk, pci, 0);
     924           0 :         int *val = getArgReference_int(stk, pci, 1);
     925           0 :         (void) cntxt;
     926           0 :         (void) mb;
     927           0 :         *ret = *val;
     928           0 :         return MAL_SUCCEED;
     929             : }
     930             : 
     931             : static void
     932         297 : stopMALdataflow(void)
     933             : {
     934         297 :         ATOMIC_SET(&exiting, 1);
     935         297 :         if (todo) {
     936         297 :                 MT_lock_set(&dataflowLock);
     937             :                 /* first wake up all running threads */
     938         297 :                 int n = 0;
     939         848 :                 for (struct worker *t = free_workers; t; t = t->next)
     940         551 :                         n++;
     941        2370 :                 for (struct worker *t = workers; t; t = t->next)
     942        2073 :                         n++;
     943        2921 :                 while (n-- > 0) {
     944             :                         /* one UP for each thread we know about */
     945        2921 :                         MT_sema_up(&todo->s);
     946             :                 }
     947         848 :                 while (free_workers) {
     948         551 :                         struct worker *t = free_workers;
     949         551 :                         assert(free_count > 0);
     950         551 :                         free_count--;
     951         551 :                         free_workers = free_workers->next;
     952         551 :                         MT_sema_up(&t->s);
     953         551 :                         finish_worker(t);
     954             :                 }
     955         306 :                 while (workers) {
     956           9 :                         struct worker *t = workers;
     957           9 :                         workers = workers->next;
     958           9 :                         finish_worker(t);
     959             :                 }
     960        2361 :                 while (exited_workers) {
     961        2064 :                         struct worker *t = exited_workers;
     962        2064 :                         exited_workers = exited_workers->next;
     963        2064 :                         finish_worker(t);
     964             :                 }
     965         297 :                 MT_lock_unset(&dataflowLock);
     966             :         }
     967         297 : }
     968             : 
     969             : void
     970         297 : mal_dataflow_reset(void)
     971             : {
     972         297 :         stopMALdataflow();
     973         297 :         workers = exited_workers = NULL;
     974         297 :         if (todo) {
     975         297 :                 MT_lock_destroy(&todo->l);
     976         297 :                 MT_sema_destroy(&todo->s);
     977         297 :                 GDKfree(todo);
     978             :         }
     979         297 :         todo = 0;                                       /* pending instructions */
     980         297 :         ATOMIC_SET(&exiting, 0);
     981         297 : }

Generated by: LCOV version 1.14