36#include "access/xact.h"
37#include "postmaster/bgworker.h"
40#include "utils/array.h"
41#include "access/htup_details.h"
42#include "utils/builtins.h"
46#ifdef PROVSQL_INPROCESS_STORE
52void provsql_buffer_ensure(
size_t need)
54 if(need > buffercap) {
55 size_t newcap = buffercap ? buffercap * 2 : 4096;
60 provsql_error(
"ProvSQL: out of memory growing the IPC buffer");
76#define WORKER_READ_BUFFER (64 * 1024)
92 size_t take = avail < n ? avail : n;
114 size_t remaining = n;
115 while(remaining > 0) {
116 ssize_t r = read(fd, p, remaining);
125#if PG_VERSION_NUM >= 190000
137static void provsql_worker_die(SIGNAL_ARGS)
140 (errcode(ERRCODE_ADMIN_SHUTDOWN),
141 errmsg(
"terminating background worker \"%s\" due to administrator command",
142 MyBgworkerEntry->bgw_type)));
148#if PG_VERSION_NUM >= 190000
149 pqsignal(SIGTERM, provsql_worker_die);
151 BackgroundWorkerUnblockSignals();
155 provsql_log(
"%s initialized", MyBgworkerEntry->bgw_name);
164 BackgroundWorker worker;
166 snprintf(worker.bgw_name, BGW_MAXLEN,
"ProvSQL MMap Worker");
167 snprintf(worker.bgw_type, BGW_MAXLEN,
"ProvSQL MMap");
169 worker.bgw_flags = BGWORKER_SHMEM_ACCESS;
170 worker.bgw_start_time = BgWorkerStart_PostmasterStart;
171 worker.bgw_restart_time = 1;
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;
178 RegisterBackgroundWorker(&worker);
220 provsql_error(
"Cannot communicate with pipe (message type S)");
229 case XACT_EVENT_PRE_COMMIT:
230 case XACT_EVENT_PRE_PREPARE:
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:
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 "
260 errhint(
"Provenance queries create gates as they run, reads "
261 "included, so they only work on the primary.")));
291 char flag = dry_run ? 1 : 0;
292 unsigned long n = (
unsigned long) nb_roots;
311#ifdef PROVSQL_INPROCESS_STORE
315 provsql_error(
"circuit_cleanup is not available in the single-process build");
318 unsigned long per_batch = PIPE_BUF /
sizeof(
pg_uuid_t);
319 for(
unsigned long i = 0; i < n; ) {
322 for(j = 0; j < per_batch && i < n; ++j, ++i)
335 provsql_error(
"Cannot read response from pipe (message type X)");
343#ifdef PROVSQL_INPROCESS_STORE
344 (void) data; (void) len;
346 const char *p = data;
359 size_t chunk = left > PIPE_BUF ? PIPE_BUF : left;
373 if(!
READB(result,
char) || !
READB(stored,
double)) {
375 provsql_error(
"Cannot read the reply to a replayed store message");
397 unsigned long nb_gates, nb_mapping, next_value,
398 dangling, unreferenced, bad_wires, bad_extra;
401 bool nulls[8] = {
false,
false,
false,
false,
false,
false,
false,
false};
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)");
422 if(get_call_result_type(fcinfo, NULL, &tupdesc) != TYPEFUNC_COMPOSITE)
423 provsql_error(
"check_store: expected composite return type");
424 tupdesc = BlessTupleDesc(tupdesc);
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);
435 PG_RETURN_DATUM(HeapTupleGetDatum(heap_form_tuple(tupdesc, values, nulls)));
457 unsigned *nb_children_out,
461 unsigned nb_children = 0;
480 provsql_error(
"Cannot communicate on pipe (message type t)");
494 provsql_error(
"Cannot communicate on pipe (message type c during get_gate_type)");
497 if(nb_children > 0) {
498 children = calloc(nb_children,
sizeof(
pg_uuid_t));
501 provsql_error(
"Cannot read children from pipe (during get_gate_type)");
516 *nb_children_out = nb_children;
517 *children_out = children;
523 pg_uuid_t *token = DatumGetUUIDP(PG_GETARG_DATUM(0));
526 unsigned nb_children = 0;
533 if(children) free(children);
546 unsigned nb_children,
561 size_t header =
sizeof(char) + 2 *
sizeof(Oid) +
sizeof(
pg_uuid_t)
563 size_t len = header + nb_children *
sizeof(
pg_uuid_t);
564 char *msg = palloc(len);
567 memcpy(p, &MyDatabaseId,
sizeof(Oid)); p +=
sizeof(Oid);
568 memcpy(p, &MyDatabaseTableSpace,
sizeof(Oid)); p +=
sizeof(Oid);
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));
587#ifdef PROVSQL_INPROCESS_STORE
597 for(
unsigned i=0; i<nb_children; ++i)
607#ifndef PROVSQL_INPROCESS_STORE
611 unsigned children_per_batch = PIPE_BUF/
sizeof(
pg_uuid_t);
620 for(
unsigned j=0; j<1+(nb_children-1)/children_per_batch; ++j) {
623 for(
unsigned i=j*children_per_batch; i<(j+1)*children_per_batch && i<nb_children; ++i) {
663 provsql_error(
"Cannot communicate with pipe (message type P)");
697 provsql_error(
"Cannot communicate with pipe (message type q)");
707 unsigned nb_children,
709 bool has_infos,
unsigned info1,
710 unsigned info2,
const char *extra)
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)
717 +
sizeof(char) + 3 *
sizeof(
unsigned) + len;
718 char *msg = palloc(total);
723 memcpy(p, &MyDatabaseId,
sizeof(Oid)); p +=
sizeof(Oid);
724 memcpy(p, &MyDatabaseTableSpace,
sizeof(Oid)); p +=
sizeof(Oid);
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));
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);
737 memcpy(p, extra, len);
745#ifdef PROVSQL_INPROCESS_STORE
747 if(!provsql_inproc_send(msg, total)) {
753 if(total <= PIPE_BUF) {
765 size_t chunk = left > PIPE_BUF ? PIPE_BUF : left;
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;
793 if(PG_ARGISNULL(0) || PG_ARGISNULL(1))
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);
819 children_data = (
pg_uuid_t*) ARR_DATA_PTR(children);
821 children_data = NULL;
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));
832 has_infos, info1, info2, extra);
847#ifdef PROVSQL_INPROCESS_STORE
848 if(!provsql_inproc_send(msg, len))
849 provsql_error(
"Cannot write to pipe (message type %c)", msg[0]);
856 provsql_error(
"Cannot write to pipe (message type %c)", msg[0]);
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)];
875 memcpy(p, &MyDatabaseId,
sizeof(Oid)); p +=
sizeof(Oid);
876 memcpy(p, &MyDatabaseTableSpace,
sizeof(Oid)); p +=
sizeof(Oid);
878 memcpy(p, &info1,
sizeof(
unsigned)); p +=
sizeof(unsigned);
879 memcpy(p, &info2,
sizeof(
unsigned));
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);
897 memcpy(p, &MyDatabaseId,
sizeof(Oid)); p +=
sizeof(Oid);
898 memcpy(p, &MyDatabaseTableSpace,
sizeof(Oid)); p +=
sizeof(Oid);
900 memcpy(p, &len,
sizeof(
unsigned)); p +=
sizeof(unsigned);
913 pg_uuid_t *token = DatumGetUUIDP(PG_GETARG_DATUM(0));
929 provsql_error(
"Cannot communicate with pipe (message type e)");
932 result = palloc(len + VARHDRSZ);
933 SET_VARSIZE(result, VARHDRSZ + len);
937 provsql_error(
"Cannot communicate with pipe (message type e)");
942 PG_RETURN_TEXT_P(result);
959 provsql_error(
"Cannot communicate with pipe (message type n)");
964 PG_RETURN_INT64((
long) nb);
971 pg_uuid_t *token = DatumGetUUIDP(PG_GETARG_DATUM(0));
972 ArrayType *result = NULL;
973 unsigned nb_children;
996 if(!
READB(nb_children,
unsigned)) {
998 provsql_error(
"Cannot read response from pipe (message type c)");
1001 children=calloc(nb_children,
sizeof(
pg_uuid_t));
1017 children_ptr = palloc(nb_children *
sizeof(Datum));
1018 for(
unsigned i=0; i<nb_children; ++i)
1019 children_ptr[i] = UUIDPGetDatum(&children[i]);
1022 result = construct_array(
1029 pfree(children_ptr);
1032 PG_RETURN_ARRAYTYPE_P(result);
1039 pg_uuid_t *token = DatumGetUUIDP(PG_GETARG_DATUM(0));
1054 provsql_error(
"Cannot communicate with pipe (message type p)");
1062 PG_RETURN_FLOAT8(result);
1069 pg_uuid_t *token = DatumGetUUIDP(PG_GETARG_DATUM(0));
1070 unsigned info1 =0, info2 = 0;
1084 provsql_error(
"Cannot communicate with pipe (message type i)");
1092 bool nulls[2] = {
false,
false};
1094 get_call_result_type(fcinfo,NULL,&tupdesc);
1095 tupdesc = BlessTupleDesc(tupdesc);
1097 values[0] = Int32GetDatum(info1);
1098 values[1] = Int32GetDatum(info2);
1100 PG_RETURN_DATUM(HeapTupleGetDatum(heap_form_tuple(tupdesc, values, nulls)));
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.