| Line | Branch | Exec | Source |
|---|---|---|---|
| 1 | /* | ||
| 2 | * Copyright (c) 2026 Tiger Data, Inc. | ||
| 3 | * Licensed under the PostgreSQL License. See LICENSE for details. | ||
| 4 | * | ||
| 5 | * parallel_ctx.c - Parallel context lifecycle over pthreads | ||
| 6 | * | ||
| 7 | * Standalone implementation of the ParallelContext lifecycle the build uses | ||
| 8 | * (see vs_parallel_ctx.h). PG builds use PostgreSQL's ParallelContext, so | ||
| 9 | * this file is compiled only for standalone. | ||
| 10 | */ | ||
| 11 | |||
| 12 | #ifdef VS_STANDALONE | ||
| 13 | |||
| 14 | #include <stdlib.h> | ||
| 15 | #include <string.h> | ||
| 16 | |||
| 17 | #include "core/log.h" | ||
| 18 | #include "core/memory.h" | ||
| 19 | #include "standalone/parallel_ctx.h" | ||
| 20 | |||
| 21 | #define VS_PARALLEL_TOC_MAGIC UINT64_C(0x56535f504152) /* "VS_PAR" */ | ||
| 22 | #define VS_PARALLEL_MAX_WORKERS_REG 8 | ||
| 23 | |||
| 24 | __thread int ParallelWorkerNumber = -1; | ||
| 25 | |||
| 26 | /* Name -> worker-entry registry (populated once at module init). */ | ||
| 27 | static struct | ||
| 28 | { | ||
| 29 | const char *name; | ||
| 30 | VsParallelWorkerFn fn; | ||
| 31 | } vs_worker_registry[VS_PARALLEL_MAX_WORKERS_REG]; | ||
| 32 | |||
| 33 | static int vs_worker_registry_count = 0; | ||
| 34 | |||
| 35 | void | ||
| 36 | 5 | vs_parallel_register_worker(const char *name, VsParallelWorkerFn fn) | |
| 37 | { | ||
| 38 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 5 times.
|
5 | if (vs_worker_registry_count == VS_PARALLEL_MAX_WORKERS_REG) |
| 39 | { | ||
| 40 | ✗ | vs_error("vs_parallel_register_worker: registry full"); | |
| 41 | } | ||
| 42 | 5 | vs_worker_registry[vs_worker_registry_count].name = name; | |
| 43 | 5 | vs_worker_registry[vs_worker_registry_count].fn = fn; | |
| 44 | 5 | vs_worker_registry_count++; | |
| 45 | 5 | } | |
| 46 | |||
| 47 | static VsParallelWorkerFn | ||
| 48 | 84 | vs_worker_lookup(const char *name) | |
| 49 | { | ||
| 50 |
1/2✓ Branch 0 taken 125 times.
✗ Branch 1 not taken.
|
125 | for (int i = 0; i < vs_worker_registry_count; i++) |
| 51 |
2/2✓ Branch 0 taken 84 times.
✓ Branch 1 taken 41 times.
|
125 | if (strcmp(vs_worker_registry[i].name, name) == 0) |
| 52 | 84 | return vs_worker_registry[i].fn; | |
| 53 | |||
| 54 | ✗ | vs_error("vs_worker_lookup: '%s' not registered", name); | |
| 55 | } | ||
| 56 | |||
| 57 | void | ||
| 58 | 84 | EnterParallelMode(void) | |
| 59 | { | ||
| 60 | 84 | } | |
| 61 | |||
| 62 | void | ||
| 63 | 84 | ExitParallelMode(void) | |
| 64 | { | ||
| 65 | 84 | } | |
| 66 | |||
| 67 | ParallelContext * | ||
| 68 | 84 | CreateParallelContext(const char *library, const char *function, int nworkers) | |
| 69 | { | ||
| 70 | 84 | ParallelContext *pcxt = calloc(1, sizeof(ParallelContext)); | |
| 71 | |||
| 72 | (void)library; | ||
| 73 | |||
| 74 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 84 times.
|
84 | if (pcxt == NULL) |
| 75 | { | ||
| 76 | ✗ | vs_error("CreateParallelContext: out of memory"); | |
| 77 | } | ||
| 78 | |||
| 79 | 84 | pcxt->nworkers = nworkers; | |
| 80 | 84 | pcxt->nworkers_launched = 0; | |
| 81 | 84 | pcxt->function_name = function; | |
| 82 | 84 | pcxt->pool = vs_thread_pool_create((uint32_t)nworkers); | |
| 83 | 84 | shm_toc_initialize_estimator(&pcxt->estimator); | |
| 84 | |||
| 85 | /* | ||
| 86 | * Bind the leader's latch now and keep it in the context: workers set it | ||
| 87 | * (via the page queues' receiver) while the leader drains, so it must | ||
| 88 | * outlive them — the context is freed only after they are all joined. | ||
| 89 | */ | ||
| 90 | 84 | InitLatch(&pcxt->leader_latch); | |
| 91 | 84 | vs_latch_attach_self(&pcxt->leader_latch); | |
| 92 | 84 | return pcxt; | |
| 93 | } | ||
| 94 | |||
| 95 | void | ||
| 96 | 84 | InitializeParallelDSM(ParallelContext *pcxt) | |
| 97 | { | ||
| 98 | 84 | pcxt->arena_size = shm_toc_estimate(&pcxt->estimator); | |
| 99 | 84 | pcxt->arena = malloc(pcxt->arena_size); | |
| 100 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 84 times.
|
84 | if (pcxt->arena == NULL) |
| 101 | { | ||
| 102 | ✗ | vs_error("InitializeParallelDSM: out of memory"); | |
| 103 | } | ||
| 104 | 84 | pcxt->toc = shm_toc_create( | |
| 105 | VS_PARALLEL_TOC_MAGIC, pcxt->arena, pcxt->arena_size); | ||
| 106 | |||
| 107 | /* | ||
| 108 | * A non-NULL sentinel: PG sets seg to the DSM segment and the driver | ||
| 109 | * treats NULL as "could not start". Standalone always starts, so seg just | ||
| 110 | * has to be non-NULL; nothing dereferences it. | ||
| 111 | */ | ||
| 112 | 84 | pcxt->seg = (dsm_segment *)pcxt; | |
| 113 | 84 | } | |
| 114 | |||
| 115 | /* | ||
| 116 | * Pool trampoline: runs on a pool worker thread as participant 1..nworkers. | ||
| 117 | * Binds this thread's worker number (ParallelWorkerNumber = participant - 1) | ||
| 118 | * and its stable latch, then runs the registered entry. The pool runs this | ||
| 119 | * once per worker per launch; afterwards the worker waits at the pool barrier | ||
| 120 | * for the leader's join (WaitForParallelWorkersToFinish). | ||
| 121 | */ | ||
| 122 | static void | ||
| 123 | 114 | vs_pool_worker_trampoline(uint32_t participant_id, void *arg) | |
| 124 | { | ||
| 125 | 114 | ParallelContext *pcxt = (ParallelContext *)arg; | |
| 126 | 114 | int worker_number = (int)participant_id - 1; | |
| 127 | |||
| 128 | 114 | ParallelWorkerNumber = worker_number; | |
| 129 | 114 | vs_latch_attach_self(&pcxt->worker_latches[worker_number]); | |
| 130 | |||
| 131 | /* | ||
| 132 | * Each worker thread needs its own current memory context — the analog of | ||
| 133 | * a PG worker process's CurrentMemoryContext — for the entry's allocations | ||
| 134 | * (the worker's per-phase contexts are created under it). Freed when the | ||
| 135 | * worker returns. | ||
| 136 | */ | ||
| 137 | 114 | VsMemCtx wctx = vs_memctx_create(NULL, "vs parallel worker"); | |
| 138 | 114 | VsMemCtx prev = vs_memctx_switch(wctx); | |
| 139 | 114 | pcxt->worker_fn(pcxt->seg, pcxt->toc); | |
| 140 | 114 | vs_memctx_switch(prev); | |
| 141 | 114 | vs_memctx_delete(wctx); | |
| 142 | 114 | } | |
| 143 | |||
| 144 | void | ||
| 145 | 84 | LaunchParallelWorkers(ParallelContext *pcxt) | |
| 146 | { | ||
| 147 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 84 times.
|
84 | if (pcxt->nworkers == 0) |
| 148 | { | ||
| 149 | ✗ | pcxt->nworkers_launched = 0; | |
| 150 | ✗ | return; | |
| 151 | } | ||
| 152 | |||
| 153 | 84 | pcxt->worker_fn = vs_worker_lookup(pcxt->function_name); | |
| 154 | 84 | pcxt->worker_latches = calloc(pcxt->nworkers, sizeof(Latch)); | |
| 155 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 84 times.
|
84 | if (pcxt->worker_latches == NULL) |
| 156 | { | ||
| 157 | ✗ | vs_error("LaunchParallelWorkers: out of memory"); | |
| 158 | } | ||
| 159 |
2/2✓ Branch 0 taken 114 times.
✓ Branch 1 taken 84 times.
|
198 | for (int i = 0; i < pcxt->nworkers; i++) |
| 160 | 114 | InitLatch(&pcxt->worker_latches[i]); | |
| 161 | |||
| 162 | /* | ||
| 163 | * Dispatch the entry onto the pool's worker threads and return: the leader | ||
| 164 | * then participates in k-means and drains the workers' streamed pages | ||
| 165 | * concurrently, joining them in WaitForParallelWorkersToFinish. | ||
| 166 | */ | ||
| 167 | 84 | vs_thread_pool_launch(pcxt->pool, vs_pool_worker_trampoline, pcxt); | |
| 168 | 84 | pcxt->nworkers_launched = pcxt->nworkers; | |
| 169 | } | ||
| 170 | |||
| 171 | void | ||
| 172 | 84 | WaitForParallelWorkersToAttach(ParallelContext *pcxt) | |
| 173 | { | ||
| 174 | /* | ||
| 175 | * No-op: the pool threads are already running the entry and will attach to | ||
| 176 | * the phase barrier. The leader's launch path polls BarrierParticipants | ||
| 177 | * until the party is whole, which is the real attach wait. | ||
| 178 | */ | ||
| 179 | (void)pcxt; | ||
| 180 | 84 | } | |
| 181 | |||
| 182 | void | ||
| 183 | 84 | WaitForParallelWorkersToFinish(ParallelContext *pcxt) | |
| 184 | { | ||
| 185 | 84 | vs_thread_pool_join(pcxt->pool); | |
| 186 | 84 | } | |
| 187 | |||
| 188 | void | ||
| 189 | 84 | DestroyParallelContext(ParallelContext *pcxt) | |
| 190 | { | ||
| 191 | /* Workers have been joined by now, so their latches are safe to free. */ | ||
| 192 | 84 | vs_thread_pool_destroy(pcxt->pool); | |
| 193 | 84 | free(pcxt->worker_latches); | |
| 194 | 84 | free(pcxt->arena); | |
| 195 | 84 | vs_shm_toc_free(pcxt->toc); | |
| 196 | 84 | free(pcxt); | |
| 197 | 84 | } | |
| 198 | |||
| 199 | #endif /* VS_STANDALONE */ | ||
| 200 |