ProvSQL C/C++ API
Adding support for provenance and uncertainty management to PostgreSQL databases
Loading...
Searching...
No Matches
provsql_mmap.c
Go to the documentation of this file.
1/**
2 * @file provsql_mmap.c
3 * @brief Background worker registration and IPC primitives for mmap-backed storage.
4 *
5 * Implements the PostgreSQL background worker lifecycle functions declared
6 * in @c provsql_mmap.h:
7 * - @c RegisterProvSQLMMapWorker(): registers the worker with the postmaster
8 * during @c _PG_init().
9 * - @c provsql_mmap_worker(): worker entry point; sets up signal handlers
10 * and enters @c provsql_mmap_main_loop().
11 *
12 * The IPC between normal backends and the background worker is handled in
13 * @c MMappedCircuit.cpp. This file provides the PostgreSQL-specific glue
14 * (background worker API, signal handling).
15 *
16 * Also declares the shared write buffer @c buffer[] and position counter
17 * @c bufferpos used by the @c STARTWRITEM / @c ADDWRITEM / @c SENDWRITEM
18 * macros in @c provsql_mmap.h.
19 *
20 * The gate-creation SQL functions (e.g. @c create_gate()) that backends
21 * call are also implemented here; they acquire the IPC lock, write a
22 * message to the background worker, and wait for an acknowledgment.
23 */
24#include "provsql_mmap.h"
25#include "provsql_rmgr.h"
26#include "provsql_shmem.h"
27#include "provsql_utils.h"
28
29#include <errno.h>
30#include <unistd.h>
31#include <poll.h>
32#include <math.h>
33#include <assert.h>
34
35#include "postgres.h"
36#include "access/xact.h"
37#include "postmaster/bgworker.h"
38#include "fmgr.h"
39#include "funcapi.h"
40#include "utils/array.h"
41#include "access/htup_details.h"
42#include "utils/builtins.h"
43
44#include "circuit_cache.h"
45
46#ifdef PROVSQL_INPROCESS_STORE
47
48char *buffer = NULL; // flawfinder: ignore
49unsigned bufferpos = 0;
50size_t buffercap = 0;
51
52void provsql_buffer_ensure(size_t need)
53{
54 if(need > buffercap) {
55 size_t newcap = buffercap ? buffercap * 2 : 4096;
56 while(newcap < need)
57 newcap *= 2;
58 buffer = realloc(buffer, newcap);
59 if(!buffer)
60 provsql_error("ProvSQL: out of memory growing the IPC buffer");
61 buffercap = newcap;
62 }
63}
64
65/* No background worker in the single-process build. */
66
67#else
68
69char buffer[PIPE_BUF]={}; // flawfinder: ignore
70unsigned bufferpos=0;
71
72/* The worker's read buffer. A message is a dozen fields, and a read() per
73 field made the worker slower than the backends that feed it: they then
74 wait on the full pipe. The pipe holds 64 KiB, which is what one read()
75 can bring back. */
76#define WORKER_READ_BUFFER (64 * 1024)
77static char worker_buffer[WORKER_READ_BUFFER]; // flawfinder: ignore
78static size_t worker_buffer_pos = 0, worker_buffer_len = 0;
79
84
85bool provsql_worker_read(void *dst, size_t n)
86{
87 char *p = dst;
88
89 while(n > 0) {
90 size_t avail = worker_buffer_len - worker_buffer_pos;
91 if(avail > 0) {
92 size_t take = avail < n ? avail : n;
93 memcpy(p, worker_buffer + worker_buffer_pos, take);
94 worker_buffer_pos += take;
95 p += take;
96 n -= take;
97 continue;
98 }
99 {
100 ssize_t r = read(provsql_shared_state->pipebmr, worker_buffer, // flawfinder: ignore
102 if(r <= 0)
103 return false;
105 worker_buffer_len = (size_t) r;
106 }
107 }
108 return true;
109}
110
111bool provsql_read_all(int fd, void *dst, size_t n)
112{
113 char *p = dst;
114 size_t remaining = n;
115 while(remaining > 0) {
116 ssize_t r = read(fd, p, remaining); // flawfinder: ignore
117 if(r <= 0)
118 return false;
119 remaining -= r;
120 p += r;
121 }
122 return true;
123}
124
125#if PG_VERSION_NUM >= 190000
126/* PostgreSQL 19 changed the default background-worker SIGTERM handler
127 * from bgworker_die() (immediate FATAL from the signal handler) to the
128 * flag-based die(), which only acts at the next CHECK_FOR_INTERRUPTS().
129 * This worker blocks in read() on the IPC pipe (restarted by
130 * SA_RESTART), so it would never observe the flag and a fast shutdown
131 * would hang on it. Restore the pre-19 semantics: the worker holds no
132 * transaction state, and being interrupted between messages leaves the
133 * store consistent, so exiting mid-read is fine. (Being interrupted
134 * *during* a write is what the write ordering in MMappedCircuit.cpp is
135 * there for; provsql.check_store() reports what an interruption left
136 * behind.) */
137static void provsql_worker_die(SIGNAL_ARGS)
138{
139 ereport(FATAL,
140 (errcode(ERRCODE_ADMIN_SHUTDOWN),
141 errmsg("terminating background worker \"%s\" due to administrator command",
142 MyBgworkerEntry->bgw_type)));
143}
144#endif
145
146PGDLLEXPORT void provsql_mmap_worker(Datum ignored)
147{
148#if PG_VERSION_NUM >= 190000
149 pqsignal(SIGTERM, provsql_worker_die);
150#endif
151 BackgroundWorkerUnblockSignals();
153 close(provsql_shared_state->pipebmw);
154 close(provsql_shared_state->pipembr);
155 provsql_log("%s initialized", MyBgworkerEntry->bgw_name);
156
158
160}
161
163{
164 BackgroundWorker worker;
165
166 snprintf(worker.bgw_name, BGW_MAXLEN, "ProvSQL MMap Worker");
167 snprintf(worker.bgw_type, BGW_MAXLEN, "ProvSQL MMap");
168
169 worker.bgw_flags = BGWORKER_SHMEM_ACCESS;
170 worker.bgw_start_time = BgWorkerStart_PostmasterStart;
171 worker.bgw_restart_time = 1;
172
173 snprintf(worker.bgw_library_name, BGW_MAXLEN, "provsql");
174 snprintf(worker.bgw_function_name, BGW_MAXLEN, "provsql_mmap_worker");
175 worker.bgw_main_arg = (Datum) 0;
176 worker.bgw_notify_pid = 0;
177
178 RegisterBackgroundWorker(&worker);
179}
180
181#endif /* PROVSQL_INPROCESS_STORE */
182
183/* -------------------------------------------------------------------------
184 * Durability of the store across a machine crash
185 *
186 * A backend's writes reach the worker over a FIFO pipe and the worker
187 * applies them to the mmap files, but nothing forces those files to disk
188 * at the moment the writing transaction commits: the heap's WAL record is
189 * fsynced, the circuit's bytes are not. After a crash of the machine
190 * (PostgreSQL crashing is harmless -- the page cache outlives it) a
191 * committed row can reference a gate that never reached the disk, and an
192 * unknown token reads back as an input gate, so the loss is silent.
193 *
194 * Two things narrow that window. The worker forces the files out shortly
195 * after the last write (PROVSQL_STORE_FLUSH_INTERVAL_MS), which bounds the
196 * loss. And, when provsql.synchronous_commit is on, a transaction that
197 * wrote to the store sends a sync request before it commits and waits for
198 * the reply, which closes the window entirely: the reply comes after every
199 * earlier message of this backend has been applied and forced.
200 * ------------------------------------------------------------------------- */
201
203
204/** Whether the current transaction has written anything to the store. */
205static bool store_written = false;
206static bool store_callbacks_registered = false;
207
208/** @brief Send the sync barrier and wait for the worker's acknowledgement. */
210{
211 char ack;
212
213 STARTWRITEM();
214 ADDWRITEM("S", char);
215 ADDWRITEDB();
216
218 if(!SENDWRITEM() || !READB(ack, char)) {
220 provsql_error("Cannot communicate with pipe (message type S)");
221 }
223}
224
225static void provsql_store_xact_callback(XactEvent event, void *arg)
226{
227 (void) arg;
228 switch(event) {
229 case XACT_EVENT_PRE_COMMIT:
230 case XACT_EVENT_PRE_PREPARE:
231 /* Still inside the transaction, so raising here aborts the commit
232 rather than leaving it half-durable. */
235 break;
236 case XACT_EVENT_COMMIT:
237 case XACT_EVENT_ABORT:
238 case XACT_EVENT_PREPARE:
239 case XACT_EVENT_PARALLEL_COMMIT:
240 case XACT_EVENT_PARALLEL_ABORT:
241 store_written = false;
242 break;
243 default:
244 break;
245 }
246}
247
248/** @brief What every store mutation does before it reaches the pipe:
249 * refuse it on a standby, and write it to the WAL. @p data is the
250 * complete message, opcode first. */
251static void provsql_log_store_write(const char *data, size_t len)
252{
254 ereport(ERROR,
255 (errcode(ERRCODE_READ_ONLY_SQL_TRANSACTION),
256 errmsg("cannot write to the ProvSQL circuit store during recovery"),
257 errdetail("The store is maintained on a standby by replaying the "
258 "primary's records; a backend writing to it would make "
259 "the two diverge."),
260 errhint("Provenance queries create gates as they run, reads "
261 "included, so they only work on the primary.")));
263}
264
265/** @brief The same, for a mutation made by a live transaction: it also
266 * arms the at-commit sync barrier. */
267static void provsql_before_store_write(const char *data, size_t len)
268{
269 provsql_log_store_write(data, len);
271}
272
274{
276 RegisterXactCallback(provsql_store_xact_callback, NULL);
278 }
279 store_written = true;
280}
281
283{
284 return store_written;
285}
286
287void provsql_circuit_cleanup_request(bool dry_run, const pg_uuid_t *roots,
288 int64 nb_roots,
290{
291 char flag = dry_run ? 1 : 0;
292 unsigned long n = (unsigned long) nb_roots;
293
294 /* The root set is as large as the number of distinct tokens stored in
295 the database, so it does not fit one atomic pipe write. The lock is
296 held across the whole exchange, which is what makes the sequence of
297 writes one message; nothing else can be talking to the worker anyway,
298 since the caller holds the database exclusively. */
300
301 STARTWRITEM();
302 ADDWRITEM("X", char);
303 ADDWRITEDB();
304 ADDWRITEM(&flag, char);
305 ADDWRITEM(&n, unsigned long);
306 if(!SENDWRITEM()) {
308 provsql_error("Cannot write to pipe (message type X)");
309 }
310
311#ifdef PROVSQL_INPROCESS_STORE
312 /* The in-memory FIFO has no atomicity limit and the dispatch runs
313 inside SENDWRITEM, so the roots must accompany the header. */
315 provsql_error("circuit_cleanup is not available in the single-process build");
316#else
317 {
318 unsigned long per_batch = PIPE_BUF / sizeof(pg_uuid_t);
319 for(unsigned long i = 0; i < n; ) {
320 unsigned long j;
321 STARTWRITEM();
322 for(j = 0; j < per_batch && i < n; ++j, ++i)
323 ADDWRITEM(&roots[i], pg_uuid_t);
324 if(!SENDWRITEM()) {
326 provsql_error("Cannot write to pipe (message type X)");
327 }
328 }
329 }
330
331 if(!READB(out->gates_before, uint64) || !READB(out->gates_after, uint64)
332 || !READB(out->wires_before, uint64) || !READB(out->wires_after, uint64)
333 || !READB(out->extra_before, uint64) || !READB(out->extra_after, uint64)) {
335 provsql_error("Cannot read response from pipe (message type X)");
336 }
338#endif
339}
340
341void provsql_replay_store_message(const char *data, size_t len)
342{
343#ifdef PROVSQL_INPROCESS_STORE
344 (void) data; (void) len;
345#else
346 const char *p = data;
347 size_t left = len;
348
349 if(len == 0)
350 return;
351
352 /* The worker is the single writer, in recovery as in normal running:
353 the message goes back down the same pipe a backend would use. The
354 lock is held across the whole write so the (possibly chunked)
355 message stays one message. */
357
358 while(left > 0) {
359 size_t chunk = left > PIPE_BUF ? PIPE_BUF : left;
360 if(write(provsql_shared_state->pipebmw, p, chunk) == -1) {
362 provsql_error("Cannot replay a store message to the pipe");
363 }
364 p += chunk;
365 left -= chunk;
366 }
367
368 /* Opcodes that answer must be drained, or the reply would be read as
369 the answer to somebody else's later question. */
370 if(data[0] == 'P') {
371 char result;
372 double stored;
373 if(!READB(result, char) || !READB(stored, double)) {
375 provsql_error("Cannot read the reply to a replayed store message");
376 }
377 }
378
380#endif
381}
382
383PG_FUNCTION_INFO_V1(check_store);
384/**
385 * @brief Report what does not add up in this database's circuit store.
386 *
387 * A store nothing has damaged answers zero to every count. A non-zero
388 * one means a write was interrupted at a point the ordering rules do not
389 * cover, or that a set of files was copied at different instants -- a
390 * file-level backup of a running server, or a base backup.
391 * @c provsql.circuit_cleanup() rebuilds the store from what is still
392 * reachable.
393 */
394Datum check_store(PG_FUNCTION_ARGS)
395{
396 char unclean;
397 unsigned long nb_gates, nb_mapping, next_value,
398 dangling, unreferenced, bad_wires, bad_extra;
399 TupleDesc tupdesc;
400 Datum values[8];
401 bool nulls[8] = {false, false, false, false, false, false, false, false};
402
403 STARTWRITEM();
404 ADDWRITEM("k", char);
405 ADDWRITEDB();
406
408 if(!SENDWRITEM()
409 || !READB(unclean, char)
410 || !READB(nb_gates, unsigned long)
411 || !READB(nb_mapping, unsigned long)
412 || !READB(next_value, unsigned long)
413 || !READB(dangling, unsigned long)
414 || !READB(unreferenced, unsigned long)
415 || !READB(bad_wires, unsigned long)
416 || !READB(bad_extra, unsigned long)) {
418 provsql_error("Cannot communicate with pipe (message type k)");
419 }
421
422 if(get_call_result_type(fcinfo, NULL, &tupdesc) != TYPEFUNC_COMPOSITE)
423 provsql_error("check_store: expected composite return type");
424 tupdesc = BlessTupleDesc(tupdesc);
425
426 values[0] = BoolGetDatum(unclean != 0);
427 values[1] = Int64GetDatum((int64) nb_gates);
428 values[2] = Int64GetDatum((int64) nb_mapping);
429 values[3] = Int64GetDatum((int64) next_value);
430 values[4] = Int64GetDatum((int64) dangling);
431 values[5] = Int64GetDatum((int64) unreferenced);
432 values[6] = Int64GetDatum((int64) bad_wires);
433 values[7] = Int64GetDatum((int64) bad_extra);
434
435 PG_RETURN_DATUM(HeapTupleGetDatum(heap_form_tuple(tupdesc, values, nulls)));
436}
437
438PG_FUNCTION_INFO_V1(get_gate_type);
439/** @brief PostgreSQL-callable wrapper for get_gate_type().
440 *
441 * On cache miss this fetches BOTH the gate type and its children from
442 * the worker, in one critical section, then caches them together. If
443 * we cached only the type (with an empty children list), a subsequent
444 * get_children() call for the same token would consult the cache, find
445 * the entry, and return 0 children : never querying the worker for the
446 * real children. provsql.provenance_evaluate hits exactly that pattern
447 * (it calls get_gate_type first, then unnest(get_children(...))) and
448 * silently folds plus/times gates over an empty set.
449 */
450/** @brief Fetch a gate's type and children, cache-first with a worker
451 * round-trip (and cache fill) on a miss. Factored out of the
452 * get_gate_type() wrapper for in-extension callers that walk the circuit
453 * from C (e.g. the annotation-transparent set_prob()). On return
454 * @p *children_out is a @c calloc'd array to be freed by the caller, or
455 * @c NULL when the gate has no children. */
457 unsigned *nb_children_out,
458 pg_uuid_t **children_out)
459{
460 gate_type type;
461 unsigned nb_children = 0;
462 pg_uuid_t *children = NULL;
463
464 type = circuit_cache_get_type(*token);
465 if(type!=gate_invalid) {
466 *nb_children_out = circuit_cache_get_children(*token, children_out);
467 return type;
468 }
469
470 /* Type fetch (message 't'). */
471 STARTWRITEM();
472 ADDWRITEM("t", char);
473 ADDWRITEDB();
474 ADDWRITEM(token, pg_uuid_t);
475
477
478 if(!SENDWRITEM() || !READB(type, gate_type)) {
480 provsql_error("Cannot communicate on pipe (message type t)");
481 }
482
483 /* Children fetch (message 'c'), batched in the same critical
484 * section so the cache entry below is complete. Skipped when the
485 * token is unknown (worker reports gate_invalid). */
486 if(type != gate_invalid) {
487 STARTWRITEM();
488 ADDWRITEM("c", char);
489 ADDWRITEDB();
490 ADDWRITEM(token, pg_uuid_t);
491
492 if(!SENDWRITEM() || !READB(nb_children, unsigned)) {
494 provsql_error("Cannot communicate on pipe (message type c during get_gate_type)");
495 }
496
497 if(nb_children > 0) {
498 children = calloc(nb_children, sizeof(pg_uuid_t));
499 if(!READB_BYTES(children, nb_children * sizeof(pg_uuid_t))) {
501 provsql_error("Cannot read children from pipe (during get_gate_type)");
502 }
503 }
504 }
505
507
508 /* Skip caching the gate_input lazy default: MMappedCircuit::getGateType
509 * returns gate_input both for real input gates and for tokens that are
510 * not yet in the mapping. Caching the latter would poison subsequent
511 * create_gate() calls in this session (the cache hit would short-circuit
512 * the worker IPC, dropping the gate). The cost is one extra IPC per
513 * lookup of a real input gate -- acceptable. */
514 if(!(type == gate_input && nb_children == 0))
515 circuit_cache_create_gate(*token, type, nb_children, children);
516 *nb_children_out = nb_children;
517 *children_out = children;
518 return type;
519}
520
521Datum get_gate_type(PG_FUNCTION_ARGS)
522{
523 pg_uuid_t *token = DatumGetUUIDP(PG_GETARG_DATUM(0));
524 gate_type type;
525 constants_t constants=get_constants(true);
526 unsigned nb_children = 0;
527 pg_uuid_t *children = NULL;
528
529 if(PG_ARGISNULL(0))
530 PG_RETURN_NULL();
531
532 type = provsql_fetch_gate(token, &nb_children, &children);
533 if(children) free(children);
534 PG_RETURN_INT32(constants.GATE_TYPE_TO_OID[type]);
535}
536
537/** @brief Internal entry point behind create_gate(): cache + worker IPC.
538 *
539 * Factored out of the SQL-callable wrapper so in-extension C/C++ code
540 * (e.g. the decomposition-aligned reachability materialiser) can create
541 * gates without Datum marshalling or gate-type-OID lookups. Same
542 * semantics: write-through to the per-session cache, then the C message
543 * to the background worker; MMappedCircuit::createGate is idempotent on
544 * already-mapped tokens. */
546 unsigned nb_children,
547 const pg_uuid_t *children_data)
548{
549 /* Populate the per-session cache, but unconditionally fall through to
550 * the worker IPC: a cache hit only proves "this token has been seen
551 * in this session before" (e.g. by get_gate_type returning the
552 * gate_input lazy default for an unknown token) -- not "the worker
553 * already has a gate for it". Skipping the IPC on a cache hit caused
554 * silently-dropped create_gate calls under concurrent backends.
555 * MMappedCircuit::createGate is idempotent on already-mapped tokens. */
556 circuit_cache_create_gate(*token, type, nb_children, children_data);
557
558 /* The WAL record is the whole logical message, children included, even
559 though the pipe may need several writes for it. */
560 {
561 size_t header = sizeof(char) + 2 * sizeof(Oid) + sizeof(pg_uuid_t)
562 + sizeof(gate_type) + sizeof(unsigned);
563 size_t len = header + nb_children * sizeof(pg_uuid_t);
564 char *msg = palloc(len);
565 char *p = msg;
566 *p++ = 'C';
567 memcpy(p, &MyDatabaseId, sizeof(Oid)); p += sizeof(Oid);
568 memcpy(p, &MyDatabaseTableSpace, sizeof(Oid)); p += sizeof(Oid);
569 memcpy(p, token, sizeof(pg_uuid_t)); p += sizeof(pg_uuid_t);
570 memcpy(p, &type, sizeof(gate_type)); p += sizeof(gate_type);
571 memcpy(p, &nb_children, sizeof(unsigned)); p += sizeof(unsigned);
572 for(unsigned i=0; i<nb_children; ++i) {
573 memcpy(p, &children_data[i], sizeof(pg_uuid_t));
574 p += sizeof(pg_uuid_t);
575 }
577 pfree(msg);
578 }
579
580 STARTWRITEM();
581 ADDWRITEM("C", char);
582 ADDWRITEDB();
583 ADDWRITEM(token, pg_uuid_t);
584 ADDWRITEM(&type, gate_type);
585 ADDWRITEM(&nb_children, unsigned);
586
587#ifdef PROVSQL_INPROCESS_STORE
588 /* The in-memory FIFO has no PIPE_BUF atomicity limit: always send the
589 gate and all its children as a single message. */
590 if(1) {
591#else
592 if(PIPE_BUF-bufferpos>nb_children*sizeof(pg_uuid_t)) {
593#endif
594 // Enough space in the buffer for an atomic write, no need of
595 // exclusive locks
596
597 for(unsigned i=0; i<nb_children; ++i)
598 ADDWRITEM(&children_data[i], pg_uuid_t);
599
601 if(!SENDWRITEM()) {
603 provsql_error("Cannot write to pipe (message type C)");
604 }
606 }
607#ifndef PROVSQL_INPROCESS_STORE
608 else {
609 // Not enough space in buffer, pipe write won't be atomic, we need to
610 // make several writes and use locks
611 unsigned children_per_batch = PIPE_BUF/sizeof(pg_uuid_t);
612
614
615 if(!SENDWRITEM()) {
617 provsql_error("Cannot write to pipe (message type C)");
618 }
619
620 for(unsigned j=0; j<1+(nb_children-1)/children_per_batch; ++j) {
621 STARTWRITEM();
622
623 for(unsigned i=j*children_per_batch; i<(j+1)*children_per_batch && i<nb_children; ++i) {
624 ADDWRITEM(&children_data[i], pg_uuid_t);
625 }
626
627 if(!SENDWRITEM()) {
629 provsql_error("Cannot write to pipe (message type C)");
630 }
631 }
632
634 }
635#endif
636}
637
638/** @brief Send a probability write and read back what the store made of
639 * it. @p tracked arms the at-commit sync barrier; the one caller that
640 * passes false is the rollback path, which runs when the transaction
641 * that would have committed is already gone. */
643 double prob,
644 double *existing,
645 bool tracked)
646{
647 char result;
648 double stored;
649
650 STARTWRITEM();
651 ADDWRITEM("P", char);
652 ADDWRITEDB();
653 ADDWRITEM(token, pg_uuid_t);
654 ADDWRITEM(&prob, double);
655 if(tracked)
657 else
659
661 if(!SENDWRITEM() || !READB(result, char) || !READB(stored, double)) {
663 provsql_error("Cannot communicate with pipe (message type P)");
664 }
666
667 if(existing)
668 *existing = stored;
669 return (provsql_set_prob_result) result;
670}
671
673 double prob,
674 double *existing)
675{
676 return provsql_send_set_prob(token, prob, existing, true);
677}
678
680{
681 provsql_send_set_prob(token, NAN, NULL, false);
682}
683
684bool provsql_internal_get_prob_written(const pg_uuid_t *token, double *prob)
685{
686 char has;
687 double stored;
688
689 STARTWRITEM();
690 ADDWRITEM("q", char);
691 ADDWRITEDB();
692 ADDWRITEM(token, pg_uuid_t);
693
695 if(!SENDWRITEM() || !READB(has, char) || !READB(stored, double)) {
697 provsql_error("Cannot communicate with pipe (message type q)");
698 }
700
701 if(prob)
702 *prob = stored;
703 return has != 0;
704}
705
707 unsigned nb_children,
708 const pg_uuid_t *children,
709 bool has_infos, unsigned info1,
710 unsigned info2, const char *extra)
711{
712 char flag = has_infos ? 1 : 0;
713 unsigned len = extra != NULL ? (unsigned) strlen(extra) : 0;
714 size_t total = sizeof(char) + 2 * sizeof(Oid) + sizeof(pg_uuid_t)
715 + sizeof(gate_type) + sizeof(unsigned)
716 + nb_children * sizeof(pg_uuid_t)
717 + sizeof(char) + 3 * sizeof(unsigned) + len;
718 char *msg = palloc(total);
719 char *p = msg;
720 unsigned i;
721
722 *p++ = 'G';
723 memcpy(p, &MyDatabaseId, sizeof(Oid)); p += sizeof(Oid);
724 memcpy(p, &MyDatabaseTableSpace, sizeof(Oid)); p += sizeof(Oid);
725 memcpy(p, token, sizeof(pg_uuid_t)); p += sizeof(pg_uuid_t);
726 memcpy(p, &type, sizeof(gate_type)); p += sizeof(gate_type);
727 memcpy(p, &nb_children, sizeof(unsigned)); p += sizeof(unsigned);
728 for(i=0; i<nb_children; ++i) {
729 memcpy(p, &children[i], sizeof(pg_uuid_t));
730 p += sizeof(pg_uuid_t);
731 }
732 *p++ = flag;
733 memcpy(p, &info1, sizeof(unsigned)); p += sizeof(unsigned);
734 memcpy(p, &info2, sizeof(unsigned)); p += sizeof(unsigned);
735 memcpy(p, &len, sizeof(unsigned)); p += sizeof(unsigned);
736 if(len > 0)
737 memcpy(p, extra, len);
738
739 circuit_cache_create_gate(*token, type, nb_children, children);
740 provsql_before_store_write(msg, total);
741
742 /* The message goes down the pipe as it is: in one write when it fits, so
743 that it is atomic and the shared lock is enough; in PIPE_BUF pieces
744 under the exclusive lock otherwise. */
745#ifdef PROVSQL_INPROCESS_STORE
747 if(!provsql_inproc_send(msg, total)) {
749 provsql_error("Cannot write to pipe (message type G)");
750 }
752#else
753 if(total <= PIPE_BUF) {
755 if(write(provsql_shared_state->pipebmw, msg, total) != (ssize_t) total) {
757 provsql_error("Cannot write to pipe (message type G)");
758 }
760 } else {
761 size_t left = total;
762 p = msg;
764 while(left > 0) {
765 size_t chunk = left > PIPE_BUF ? PIPE_BUF : left;
766 if(write(provsql_shared_state->pipebmw, p, chunk) != (ssize_t) chunk) {
768 provsql_error("Cannot write to pipe (message type G)");
769 }
770 p += chunk;
771 left -= chunk;
772 }
774 }
775#endif
776 pfree(msg);
777}
778
779
780
781PG_FUNCTION_INFO_V1(create_gate);
782/** @brief PostgreSQL-callable wrapper for create_gate(). */
783Datum create_gate(PG_FUNCTION_ARGS)
784{
785 pg_uuid_t *token = DatumGetUUIDP(PG_GETARG_DATUM(0));
786 Oid oid_type = PG_GETARG_INT32(1);
787 ArrayType *children = PG_ARGISNULL(2)?NULL:PG_GETARG_ARRAYTYPE_P(2);
788 unsigned nb_children = 0;
789 gate_type type = gate_invalid;
790 constants_t constants;
791 pg_uuid_t *children_data;
792
793 if(PG_ARGISNULL(0) || PG_ARGISNULL(1))
794 provsql_error("Invalid NULL value passed to create_gate");
795
796 if(children) {
797 if(ARR_NDIM(children) > 1)
798 provsql_error("Invalid multi-dimensional array passed to create_gate");
799 if(array_contains_nulls(children))
800 provsql_error("create_gate: children array must not contain NULL "
801 "elements (filter them out before calling)");
802 if(ARR_NDIM(children) == 1)
803 nb_children = *ARR_DIMS(children);
804 }
805
806 constants=get_constants(true);
807
808 for(int i=0; i<nb_gate_types; ++i) {
809 if(constants.GATE_TYPE_TO_OID[i]==oid_type) {
810 type = i;
811 break;
812 }
813 }
814 if(type == gate_invalid) {
815 provsql_error("Invalid gate type");
816 }
817
818 if(nb_children>0)
819 children_data = (pg_uuid_t*) ARR_DATA_PTR(children);
820 else
821 children_data = NULL;
822
823 if(PG_NARGS() > 3) {
824 /* create_gate(token, type, children, info1, info2, extra): the gate
825 and what it records in one unanswered message. */
826 bool has_infos = !PG_ARGISNULL(3) || !PG_ARGISNULL(4);
827 unsigned info1 = PG_ARGISNULL(3) ? 0 : (unsigned) PG_GETARG_INT32(3);
828 unsigned info2 = PG_ARGISNULL(4) ? 0 : (unsigned) PG_GETARG_INT32(4);
829 char *extra = PG_ARGISNULL(5) ? NULL : text_to_cstring(PG_GETARG_TEXT_PP(5));
830
831 provsql_internal_create_gate_with(token, type, nb_children, children_data,
832 has_infos, info1, info2, extra);
833 PG_RETURN_VOID();
834 }
835
836 provsql_internal_create_gate(token, type, nb_children, children_data);
837
838 PG_RETURN_VOID();
839}
840
841/* set_infos and set_extra: the SQL functions are gone (a gate is created
842 with what it records), but the install scripts of earlier versions bind
843 these symbols, and an upgrade runs them. What they send is applied
844 write-once and not answered, like everything a backend writes. */
845static void send_unanswered(const char *msg, size_t len)
846{
847#ifdef PROVSQL_INPROCESS_STORE
848 if(!provsql_inproc_send(msg, len))
849 provsql_error("Cannot write to pipe (message type %c)", msg[0]);
850#else
851 if(len > PIPE_BUF)
852 provsql_error("message type %c too long", msg[0]);
854 if(write(provsql_shared_state->pipebmw, msg, len) != (ssize_t) len) {
856 provsql_error("Cannot write to pipe (message type %c)", msg[0]);
857 }
859#endif
860}
861
862PG_FUNCTION_INFO_V1(set_infos);
863/** @brief Entry point of the @c set_infos of earlier versions' scripts. */
864Datum set_infos(PG_FUNCTION_ARGS)
865{
866 pg_uuid_t *token = DatumGetUUIDP(PG_GETARG_DATUM(0));
867 unsigned info1 = PG_ARGISNULL(1) ? 0 : (unsigned) PG_GETARG_INT32(1);
868 unsigned info2 = PG_ARGISNULL(2) ? 0 : (unsigned) PG_GETARG_INT32(2);
869 char msg[sizeof(char) + 2 * sizeof(Oid) + sizeof(pg_uuid_t) + 2 * sizeof(unsigned)];
870 char *p = msg;
871
872 if(PG_ARGISNULL(0))
873 PG_RETURN_VOID();
874 *p++ = 'I';
875 memcpy(p, &MyDatabaseId, sizeof(Oid)); p += sizeof(Oid);
876 memcpy(p, &MyDatabaseTableSpace, sizeof(Oid)); p += sizeof(Oid);
877 memcpy(p, token, sizeof(pg_uuid_t)); p += sizeof(pg_uuid_t);
878 memcpy(p, &info1, sizeof(unsigned)); p += sizeof(unsigned);
879 memcpy(p, &info2, sizeof(unsigned));
880 provsql_before_store_write(msg, sizeof(msg));
881 send_unanswered(msg, sizeof(msg));
882 PG_RETURN_VOID();
883}
884
885PG_FUNCTION_INFO_V1(set_extra);
886/** @brief Entry point of the @c set_extra of earlier versions' scripts. */
887Datum set_extra(PG_FUNCTION_ARGS)
888{
889 pg_uuid_t *token = DatumGetUUIDP(PG_GETARG_DATUM(0));
890 char *str = text_to_cstring(PG_GETARG_TEXT_PP(1));
891 unsigned len = (unsigned) strlen(str);
892 size_t total = sizeof(char) + 2 * sizeof(Oid) + sizeof(pg_uuid_t) + sizeof(unsigned) + len;
893 char *msg = palloc(total);
894 char *p = msg;
895
896 *p++ = 'E';
897 memcpy(p, &MyDatabaseId, sizeof(Oid)); p += sizeof(Oid);
898 memcpy(p, &MyDatabaseTableSpace, sizeof(Oid)); p += sizeof(Oid);
899 memcpy(p, token, sizeof(pg_uuid_t)); p += sizeof(pg_uuid_t);
900 memcpy(p, &len, sizeof(unsigned)); p += sizeof(unsigned);
901 memcpy(p, str, len);
902 provsql_before_store_write(msg, total);
903 send_unanswered(msg, total);
904 pfree(msg);
905 pfree(str);
906 PG_RETURN_VOID();
907}
908
909PG_FUNCTION_INFO_V1(get_extra);
910/** @brief PostgreSQL-callable wrapper for get_extra(). */
911Datum get_extra(PG_FUNCTION_ARGS)
912{
913 pg_uuid_t *token = DatumGetUUIDP(PG_GETARG_DATUM(0));
914 text *result;
915 unsigned len;
916
917 if(PG_ARGISNULL(0))
918 PG_RETURN_NULL();
919
920 STARTWRITEM();
921 ADDWRITEM("e", char);
922 ADDWRITEDB();
923 ADDWRITEM(token, pg_uuid_t);
924
926
927 if(!SENDWRITEM() || !READB(len, unsigned)) {
929 provsql_error("Cannot communicate with pipe (message type e)");
930 }
931
932 result = palloc(len + VARHDRSZ);
933 SET_VARSIZE(result, VARHDRSZ + len);
934
935 if(!READB_BYTES(VARDATA(result), len)) {
937 provsql_error("Cannot communicate with pipe (message type e)");
938 }
939
941
942 PG_RETURN_TEXT_P(result);
943}
944
945PG_FUNCTION_INFO_V1(get_nb_gates);
946/** @brief PostgreSQL-callable wrapper for get_nb_gates(). */
947Datum get_nb_gates(PG_FUNCTION_ARGS)
948{
949 unsigned long nb;
950
951 STARTWRITEM();
952 ADDWRITEM("n", char);
953 ADDWRITEDB();
954
956
957 if(!SENDWRITEM() || !READB(nb, unsigned long)) {
959 provsql_error("Cannot communicate with pipe (message type n)");
960 }
961
963
964 PG_RETURN_INT64((long) nb);
965}
966
967PG_FUNCTION_INFO_V1(get_children);
968/** @brief PostgreSQL-callable wrapper for get_children(). */
969Datum get_children(PG_FUNCTION_ARGS)
970{
971 pg_uuid_t *token = DatumGetUUIDP(PG_GETARG_DATUM(0));
972 ArrayType *result = NULL;
973 unsigned nb_children;
974 pg_uuid_t *children;
975 Datum *children_ptr;
976 constants_t constants;
977
978 if(PG_ARGISNULL(0))
979 PG_RETURN_NULL();
980
981 nb_children = circuit_cache_get_children(*token, &children);
982
983 if(!children) {
984 STARTWRITEM();
985 ADDWRITEM("c", char);
986 ADDWRITEDB();
987 ADDWRITEM(token, pg_uuid_t);
988
990
991 if(!SENDWRITEM()) {
993 provsql_error("Cannot write to pipe (message type c)");
994 }
995
996 if(!READB(nb_children, unsigned)) {
998 provsql_error("Cannot read response from pipe (message type c)");
999 }
1000
1001 children=calloc(nb_children, sizeof(pg_uuid_t));
1002
1003 if(!READB_BYTES(children, nb_children*sizeof(pg_uuid_t))) {
1005 provsql_error("Cannot read from pipe (message type c)");
1006 }
1008
1009 /* Skip caching when the worker reports zero children: we cannot
1010 * distinguish a real zero-child gate (input/zero/one/...) from a
1011 * token unknown to the worker, and caching the latter poisons
1012 * subsequent create_gate() calls in this session. */
1013 if(nb_children > 0)
1014 circuit_cache_create_gate(*token, gate_invalid, nb_children, children);
1015 }
1016
1017 children_ptr = palloc(nb_children * sizeof(Datum));
1018 for(unsigned i=0; i<nb_children; ++i)
1019 children_ptr[i] = UUIDPGetDatum(&children[i]);
1020
1021 constants=get_constants(true);
1022 result = construct_array(
1023 children_ptr,
1024 nb_children,
1025 constants.OID_TYPE_UUID,
1026 16,
1027 false,
1028 'c');
1029 pfree(children_ptr);
1030 free(children);
1031
1032 PG_RETURN_ARRAYTYPE_P(result);
1033}
1034
1035PG_FUNCTION_INFO_V1(get_prob);
1036/** @brief PostgreSQL-callable wrapper for get_prob(). */
1037Datum get_prob(PG_FUNCTION_ARGS)
1038{
1039 pg_uuid_t *token = DatumGetUUIDP(PG_GETARG_DATUM(0));
1040 double result;
1041
1042 if(PG_ARGISNULL(0))
1043 PG_RETURN_NULL();
1044
1045 STARTWRITEM();
1046 ADDWRITEM("p", char);
1047 ADDWRITEDB();
1048 ADDWRITEM(token, pg_uuid_t);
1049
1051
1052 if(!SENDWRITEM() || !READB(result, double)) {
1054 provsql_error("Cannot communicate with pipe (message type p)");
1055 }
1056
1058
1059 if(isnan(result))
1060 PG_RETURN_NULL();
1061 else
1062 PG_RETURN_FLOAT8(result);
1063}
1064
1065PG_FUNCTION_INFO_V1(get_infos);
1066/** @brief PostgreSQL-callable wrapper for get_infos(). */
1067Datum get_infos(PG_FUNCTION_ARGS)
1068{
1069 pg_uuid_t *token = DatumGetUUIDP(PG_GETARG_DATUM(0));
1070 unsigned info1 =0, info2 = 0;
1071
1072 if(PG_ARGISNULL(0))
1073 PG_RETURN_NULL();
1074
1075 STARTWRITEM();
1076 ADDWRITEM("i", char);
1077 ADDWRITEDB();
1078 ADDWRITEM(token, pg_uuid_t);
1079
1081
1082 if(!SENDWRITEM() || !READB(info1, int) || !READB(info2, int)) {
1084 provsql_error("Cannot communicate with pipe (message type i)");
1085 }
1086
1088
1089 {
1090 TupleDesc tupdesc;
1091 Datum values[2];
1092 bool nulls[2] = {false, false};
1093
1094 get_call_result_type(fcinfo,NULL,&tupdesc);
1095 tupdesc = BlessTupleDesc(tupdesc);
1096
1097 values[0] = Int32GetDatum(info1);
1098 values[1] = Int32GetDatum(info2);
1099
1100 PG_RETURN_DATUM(HeapTupleGetDatum(heap_form_tuple(tupdesc, values, nulls)));
1101 }
1102}
void destroy_provsql_mmap()
Unmap and close the mmap files.
void provsql_mmap_main_loop()
Main processing loop of the mmap background worker.
void initialize_provsql_mmap()
Initialise the circuit store.
C-linkage interface to the in-process provenance circuit cache.
gate_type circuit_cache_get_type(pg_uuid_t token)
Retrieve the type of a cached gate.
unsigned circuit_cache_get_children(pg_uuid_t token, pg_uuid_t **children)
Retrieve the children of a cached gate.
bool circuit_cache_create_gate(pg_uuid_t token, gate_type type, unsigned nb_children, const pg_uuid_t *children)
Insert a new gate into the circuit cache.
Datum get_gate_type(PG_FUNCTION_ARGS)
#define provsql_error(fmt,...)
Report a fatal ProvSQL error and abort the current transaction.
#define provsql_log(fmt,...)
Write a ProvSQL message to the server log only.
Datum get_infos(PG_FUNCTION_ARGS)
PostgreSQL-callable wrapper for get_infos().
void provsql_mmap_worker(Datum ignored)
Entry point for the ProvSQL mmap background worker.
bool provsql_store_written(void)
Whether this transaction has written to the circuit store.
void provsql_internal_create_gate(const pg_uuid_t *token, gate_type type, unsigned nb_children, const pg_uuid_t *children_data)
Internal entry point behind create_gate(): cache + worker IPC.
void provsql_replay_store_message(const char *data, size_t len)
Feed a logged store message back to the worker.
#define WORKER_READ_BUFFER
Datum check_store(PG_FUNCTION_ARGS)
Report what does not add up in this database's circuit store.
void provsql_internal_create_gate_with(const pg_uuid_t *token, gate_type type, unsigned nb_children, const pg_uuid_t *children, bool has_infos, unsigned info1, unsigned info2, const char *extra)
Create a gate together with its infos and its text, in one message that is not answered.
void provsql_store_note_write(void)
Note that this transaction has written to the circuit store.
static void provsql_store_xact_callback(XactEvent event, void *arg)
static void provsql_store_sync_barrier(void)
Send the sync barrier and wait for the worker's acknowledgement.
static void provsql_before_store_write(const char *data, size_t len)
The same, for a mutation made by a live transaction: it also arms the at-commit sync barrier.
bool provsql_read_all(int fd, void *dst, size_t n)
Read exactly n bytes from fd into dst; false on EOF/error.
static void send_unanswered(const char *msg, size_t len)
static size_t worker_buffer_pos
static bool store_written
Whether the current transaction has written anything to the store.
Datum set_extra(PG_FUNCTION_ARGS)
Entry point of the set_extra of earlier versions' scripts.
Datum get_nb_gates(PG_FUNCTION_ARGS)
PostgreSQL-callable wrapper for get_nb_gates().
Datum create_gate(PG_FUNCTION_ARGS)
PostgreSQL-callable wrapper for create_gate().
static bool store_callbacks_registered
gate_type provsql_fetch_gate(const pg_uuid_t *token, unsigned *nb_children_out, pg_uuid_t **children_out)
PostgreSQL-callable wrapper for get_gate_type().
Datum set_infos(PG_FUNCTION_ARGS)
Entry point of the set_infos of earlier versions' scripts.
void provsql_circuit_cleanup_request(bool dry_run, const pg_uuid_t *roots, int64 nb_roots, provsql_cleanup_result *out)
Ask the worker to rebuild this database's store, keeping only what roots reach.
bool provsql_worker_buffered(void)
Whether the worker's read buffer holds unread bytes, which poll() cannot see.
bool provsql_worker_read(void *dst, size_t n)
The worker's buffered read of the request pipe: what the pipe holds is read in one call,...
Datum get_gate_type(PG_FUNCTION_ARGS)
static size_t worker_buffer_len
provsql_set_prob_result provsql_internal_set_prob(const pg_uuid_t *token, double prob, double *existing)
Write a gate's probability from in-extension C/C++ code.
char buffer[PIPE_BUF]
Shared write buffer used with STARTWRITEM / ADDWRITEM / SENDWRITEM.
bool provsql_internal_get_prob_written(const pg_uuid_t *token, double *prob)
Report whether a probability has been written on a gate.
void provsql_internal_clear_prob(const pg_uuid_t *token)
Drop a gate's probability, leaving it as it was before anyone wrote one.
Datum get_children(PG_FUNCTION_ARGS)
PostgreSQL-callable wrapper for get_children().
static void provsql_log_store_write(const char *data, size_t len)
What every store mutation does before it reaches the pipe: refuse it on a standby,...
Datum get_extra(PG_FUNCTION_ARGS)
PostgreSQL-callable wrapper for get_extra().
static provsql_set_prob_result provsql_send_set_prob(const pg_uuid_t *token, double prob, double *existing, bool tracked)
Send a probability write and read back what the store made of it.
static char worker_buffer[WORKER_READ_BUFFER]
Datum get_prob(PG_FUNCTION_ARGS)
PostgreSQL-callable wrapper for get_prob().
unsigned bufferpos
Current write position within buffer.
bool provsql_synchronous_commit
Global variable set by the provsql.synchronous_commit run-time configuration parameter: when true,...
void RegisterProvSQLMMapWorker(void)
Register the ProvSQL mmap background worker with PostgreSQL.
Background worker and IPC primitives for mmap-backed circuit storage.
#define ADDWRITEDB()
Append the per-message database header to the write buffer.
#define READB_BYTES(ptr, n)
Read exactly n bytes of a reply from the main-to-background pipe.
provsql_set_prob_result
Outcome of a probability write, mirroring MMappedCircuit::SetProbResult across the IPC boundary.
#define READB(var, type)
Read one value of type from the main-to-background pipe.
#define STARTWRITEM()
Reset the shared write buffer for a new batched write.
#define ADDWRITEM(pvar, type)
Append one value of type to the shared write buffer.
#define SENDWRITEM()
Flush the shared write buffer to the background-to-main pipe atomically.
void provsql_wal_log_store_message(const char *data, size_t len)
Write one store message to the WAL, if WAL logging is on.
bool provsql_store_write_allowed(void)
Whether this process may write to the circuit store.
WAL-logging the circuit store: the entry points other files use.
void provsql_shmem_unlock(void)
Release the ProvSQL LWLock.
void provsql_shmem_lock_exclusive(void)
Acquire the ProvSQL LWLock in exclusive mode.
provsqlSharedState * provsql_shared_state
Pointer to the ProvSQL shared-memory segment (set in provsql_shmem_startup).
void provsql_shmem_lock_shared(void)
Acquire the ProvSQL LWLock in shared mode.
Shared-memory segment and inter-process pipe management.
constants_t get_constants(bool failure_if_not_possible)
Retrieve the cached OID constants for the current database.
Core types, constants, and utilities shared across ProvSQL.
Structure to store the value of various constants.
Oid GATE_TYPE_TO_OID[nb_gate_types]
Array of the OID of each provenance_gate ENUM value.
Oid OID_TYPE_UUID
OID of the uuid TYPE.
What provsql.circuit_cleanup() reports.
uint64 wires_before
Child wires before.
uint64 wires_after
Child wires kept.
uint64 extra_before
Annotation bytes before.
uint64 gates_after
Gate records kept.
uint64 extra_after
Annotation bytes kept.
uint64 gates_before
Gate records before.